Coverage for src/evutils/io/_dat.py: 80%

185 statements  

« prev     ^ index     » next       coverage.py v7.15.1, created at 2026-07-18 05:24 +0000

1"""Prophesee DAT (.dat) CD-event decoder/encoder. 

2 

3DAT is a fixed-record format: an ASCII ``% ...`` header, two bytes (event type 

4``0x0C`` = CD, event size ``0x08``), then 8-byte little-endian events 

5(``uint32`` timestamp + ``uint32`` data with ``x[0:13]``, ``y[14:27]``, 

6``p[28]``). Decoding is done by the native ``DAT_parse_chunk_soa`` (which also 

7tracks the 32-bit timestamp overflow); encoding is vectorised numpy. 

8""" 

9from __future__ import annotations 

10 

11import io 

12from datetime import datetime 

13 

14import numpy as np 

15 

16from ..types import EventArray, TriggerArray 

17from .common import EventDecoder, EventEncoder 

18from ._native_core import ( 

19 EVUTILS_PARSE_ERROR, 

20 EventSoABuffers, 

21 TriggerSoABuffers, 

22 decode_all_soa, 

23 events_view, 

24 triggers_view, 

25 parse_step, 

26) 

27from ._native_dat import ( 

28 DatInput, 

29 DatParser, 

30) 

31from ._source import ByteSource 

32 

33_EMPTY_EVENTS = EventArray.empty() 

34 

35DAT_EVENT_TYPE_CD = 0x0C 

36DAT_EVENT_SIZE = 0x08 

37 

38class EventDecoder_Dat(EventDecoder): 

39 """Decode Prophesee DAT files into ``EventArray`` chunks. 

40  

41 Parameters 

42 ---------- 

43 source 

44 Byte source to read from. 

45 chunk_size 

46 Maximum number of events produced per :meth:`read_chunk` call (the 

47 native output-buffer capacity). Does not bound the file size. 

48 

49 References 

50 ---------- 

51 [1] Prophesee DAT file format 

52 https://docs.prophesee.ai/stable/data/file_formats/dat.html 

53 

54 """ 

55 

56 #: DAT is a fixed 2-word-per-event format, so the parser fills an output 

57 #: buffer to exactly its capacity -> eligible for EventReader's zero-copy 

58 #: n_events fast path. 

59 _exact_window = True 

60 

61 #: Fixed 8-byte records make event-index <-> byte-offset exact, and the raw 

62 #: 32-bit timestamps are the even-strided words -> seek is direct index math 

63 #: / a searchsorted, no index file needed. Assumes timestamps do not wrap 

64 #: the 32-bit field (~71 min); a wrapped recording would seek approximately. 

65 SUPPORTS_SEEK = True 

66 

67 #: init() slurps the whole payload into memory (or mmaps it), so seek works 

68 #: even over a non-seekable source (e.g. a compressed stream). 

69 _buffers_in_memory = True 

70 

71 def __init__(self, source: ByteSource, chunk_size: int = 1_000_000): 

72 super().__init__(source, chunk_size) 

73 self._header: dict[str, str | int | float] = {} 

74 self._buf: bytes | bytearray | None = None 

75 self._payload_off: int = 0 

76 self._words: "np.ndarray | None" = None # uint32 view of the payload (2 words / event) 

77 self._offset: int = 0 # current uint32-word offset 

78 self._parser: "Callable | None" = None 

79 self._events: "EventArray | None" = None 

80 self._triggers: "TriggerArray | None" = None 

81 

82 # ------------------------------------------------------------------ # 

83 def _parse_header(self, buf: bytes) -> int: 

84 """Scan the leading ``%`` ASCII header; return the byte offset of the 

85 first non-header byte (the event-type byte). 

86 """ 

87 mv = memoryview(buf) 

88 n = len(mv) 

89 off = 0 

90 # Header lines start with "% "; the binary section (event-type byte, 

91 # then records) begins at the first line without that prefix. 

92 while off + 1 < n and mv[off] == 0x25 and mv[off + 1] == 0x20: 

93 window = bytes(mv[off:off + 8192]) 

94 rel = window.find(b"\n") 

95 if rel < 0: 

96 break 

97 self._consume_header_line(window[:rel]) 

98 off += rel + 1 

99 return off 

100 

101 def _consume_header_line(self, line: bytes) -> None: 

102 try: 

103 parts = line.decode("ascii", "ignore").strip().split() 

104 except ValueError: 

105 return 

106 # e.g. ["%", "Width", "1280"] 

107 if len(parts) >= 3: 

108 key = parts[1].lower() 

109 try: 

110 if key == "width": 

111 self._width = int(parts[2]) 

112 elif key == "height": 

113 self._height = int(parts[2]) 

114 except ValueError: 

115 pass 

116 

117 def init(self) -> None: 

118 """Initialize the DAT reader. 

119 

120 Returns 

121 ------- 

122 None 

123 

124 """ 

125 if self._is_initialized: 

126 return 

127 

128 if self._source.mappable(): 

129 self._buf = self._source.buffer() 

130 else: 

131 self._buf = memoryview(self._source.read(-1)) 

132 

133 off = self._parse_header(self._buf) 

134 # Two-byte binary header: event type, event size. 

135 if off + 2 <= len(self._buf): 

136 event_size = self._buf[off + 1] 

137 off += 2 

138 if event_size not in (0, DAT_EVENT_SIZE): 

139 raise NotImplementedError( 

140 f"DAT event size {event_size} not supported (only 8-byte CD)" 

141 ) 

142 self._payload_off = off 

143 

144 n_events = (len(self._buf) - off) // 8 

145 if n_events > 0: 

146 self._words = np.frombuffer( 

147 self._buf, dtype=np.uint32, count=n_events * 2, offset=off 

148 ) 

149 else: 

150 self._words = np.empty(0, dtype=np.uint32) 

151 

152 self._offset = 0 

153 self._parser = DatParser() 

154 self._input_cls = DatInput 

155 self._word_dtype = np.uint32 

156 cap = int(self._chunk_size) 

157 self._events = EventSoABuffers(cap) 

158 self._triggers = TriggerSoABuffers(1) # DAT CD files have no triggers 

159 self._is_initialized = True 

160 

161 def parse_step(self, events: EventSoABuffers, triggers: TriggerSoABuffers) -> int: 

162 """Run the parser once, appending into ``events``; advance the offset. 

163 See :meth:`EventDecoder_EVT.parse_step`. 

164 """ 

165 if not self._is_initialized: 

166 self.init() 

167 if self._words is None or self._offset >= len(self._words): 

168 self._eof = True 

169 return 0 

170 appended, self._offset = parse_step( 

171 self._words, self._offset, DatInput, self._parser, events, triggers, 

172 ) 

173 if self._offset >= len(self._words): 

174 self._eof = True 

175 return appended 

176 

177 def read_chunk(self, delta_t_hint: int | None = None, 

178 n_events_hint: int | None = None) -> 'EventArray | tuple[EventArray, TriggerArray]': 

179 if not self._is_initialized: 

180 self.init() 

181 

182 if self._words is None or self._offset >= len(self._words): 

183 self._eof = True 

184 if self.read_external_triggers: 

185 from ..types import TriggerArray 

186 return _EMPTY_EVENTS, TriggerArray.empty() 

187 return _EMPTY_EVENTS 

188 

189 ev, tr = self._events, self._triggers 

190 ev.reset() 

191 tr.reset() 

192 appended = 0 

193 while appended == 0 and self._offset < len(self._words): 

194 appended = self.parse_step(ev, tr) 

195 

196 n = ev.size 

197 if n == 0: 

198 if self.read_external_triggers: 

199 from ..types import TriggerArray 

200 return _EMPTY_EVENTS, TriggerArray.empty() 

201 return _EMPTY_EVENTS 

202 # Zero-copy view (valid until the next read_chunk); see EVT decoder. 

203 if self.read_external_triggers: 

204 return events_view(ev), triggers_view(tr) 

205 return events_view(ev) 

206 

207 def read_all(self) -> 'EventArray | tuple[EventArray, TriggerArray]': 

208 """Decode the whole remaining payload into one buffer (no per-chunk copy).""" 

209 if not self._is_initialized: 

210 self.init() 

211 if self._words is None or self._offset >= len(self._words): 

212 self._eof = True 

213 return _EMPTY_EVENTS 

214 # Exactly one event per two uint32 words. 

215 out, self._offset = decode_all_soa( 

216 self._words, self._offset, DatInput, self._parser, 

217 est_events_per_word=0.5, 

218 ) 

219 self._eof = True 

220 return out 

221 

222 def seek(self, t: int | None = None, n: int | None = None) -> tuple["SeekResult", "EventArray", "TriggerArray | None"]: 

223 """Seek to an absolute timestamp (µs) or event index. See base class. 

224 

225 Fixed 8-byte records: event ``k`` lives at word offset ``2k`` and its 

226 timestamp is the even-strided word ``_words[2k]``. Time seek is a 

227 ``searchsorted`` over those words; index seek is direct. Assumes the 

228 raw 32-bit timestamps are non-decreasing (no wrap). 

229 """ 

230 from .common import SeekResult 

231 if not self._is_initialized: 

232 self.init() 

233 axis, val = self._seek_axis(t, n) 

234 

235 n_events = 0 if self._words is None else len(self._words) // 2 

236 ts_view = (self._words[0::2] if n_events 

237 else np.empty(0, dtype=np.uint32)) 

238 

239 if not hasattr(self, "_wrap_indices"): 

240 if n_events > 0: 

241 wraps = np.where(ts_view[1:] < ts_view[:-1])[0] + 1 

242 self._wrap_indices = np.concatenate(([0], wraps, [n_events])) 

243 else: 

244 self._wrap_indices = np.array([0, 0]) 

245 

246 if axis == "t": 

247 bucket = int(val >> 32) 

248 val_rem = int(val & 0xFFFFFFFF) 

249 

250 if bucket >= len(self._wrap_indices) - 1: 

251 idx = n_events 

252 else: 

253 start_idx = self._wrap_indices[bucket] 

254 end_idx = self._wrap_indices[bucket + 1] 

255 idx = start_idx + int(np.searchsorted(ts_view[start_idx:end_idx], val_rem, side="left")) 

256 else: 

257 idx = val 

258 idx = max(0, min(idx, n_events)) 

259 

260 self._offset = idx * 2 

261 self._eof = idx >= n_events 

262 if self._parser is not None: 

263 wrap_count = max(0, int(np.searchsorted(self._wrap_indices, idx, side="right")) - 1) 

264 wrap_offset = wrap_count * (1 << 32) 

265 self._parser.reset(wrap_offset) 

266 else: 

267 wrap_offset = max(0, int(np.searchsorted(self._wrap_indices, idx, side="right")) - 1) * (1 << 32) 

268 

269 landed_ts = int(ts_view[idx]) + wrap_offset if idx < n_events else val 

270 return SeekResult(ts=landed_ts, index=idx, eof=self._eof), _EMPTY_EVENTS, None 

271 

272 def reset(self) -> None: 

273 """Reset the DAT reader to the beginning. 

274 

275 Returns 

276 ------- 

277 None 

278 

279 """ 

280 self._offset = 0 

281 self._eof = False 

282 if self._parser is not None: 

283 self._parser.reset() 

284 

285 def tell(self) -> int: 

286 """Get the current byte offset. 

287 

288 Returns 

289 ------- 

290 int 

291 Current byte offset. 

292 

293 """ 

294 return self._payload_off + self._offset * 4 

295 

296 def close(self) -> None: 

297 """Close the DAT reader. 

298 

299 Returns 

300 ------- 

301 None 

302 

303 """ 

304 self._words = None 

305 self._buf = None 

306 

307class EventEncoder_Dat(EventEncoder): 

308 """Encode events into a Prophesee DAT (.dat) CD file. 

309 

310 Parameters 

311 ---------- 

312 writable 

313 Destination stream to write to. 

314 width, height : int 

315 Frame geometry written into the header. 

316 dt : datetime, optional 

317 Recording timestamp (defaults to now). 

318 version : int 

319 DAT version (defaults to 2). 

320 

321 References 

322 ---------- 

323 [1] Prophesee DAT file format 

324 https://docs.prophesee.ai/stable/data/file_formats/dat.html 

325 

326 """ 

327 

328 def __init__(self, writable: io.BufferedWriter, width: int = 1280, height: int = 720, 

329 dt: datetime | None = None, version: int = 2): 

330 super().__init__(writable, width, height, dt) 

331 self._version = version 

332 

333 def init(self) -> None: 

334 """Initialize the DAT writer. 

335 

336 Returns 

337 ------- 

338 None 

339 

340 """ 

341 if self._is_initialized: 

342 return 

343 header = ( 

344 "% Data file containing CD events.\n" 

345 f"% Version {self._version}\n" 

346 f"% Date {self._dt.strftime('%Y-%m-%d %H:%M:%S')}\n" 

347 f"% Height {self._height}\n" 

348 f"% Width {self._width}\n" 

349 ) 

350 self._fd.write(header.encode("ascii")) 

351 self._fd.write(bytes([DAT_EVENT_TYPE_CD, DAT_EVENT_SIZE])) 

352 self._is_initialized = True 

353 

354 def write(self, events: 'np.ndarray | EventArray', triggers: 'np.ndarray | TriggerArray | None' = None) -> int: 

355 """Write events to the DAT file. 

356 

357 Parameters 

358 ---------- 

359 events : np.ndarray or EventArray 

360 Array of events to write. 

361 

362 Returns 

363 ------- 

364 int 

365 Number of written events. 

366 

367 """ 

368 if not self._is_initialized: 

369 self.init() 

370 

371 if isinstance(events, EventArray): 

372 t, x, y, p = events.t, events.x, events.y, events.p 

373 else: 

374 t, x, y, p = events["t"], events["x"], events["y"], events["p"] 

375 

376 n = len(t) 

377 out = np.empty(n * 2, dtype=np.uint32) 

378 out[0::2] = (np.asarray(t).astype(np.uint64) & np.uint64(0xFFFFFFFF)).astype(np.uint32) 

379 out[1::2] = ( 

380 (x.astype(np.uint32) & np.uint32(0x3FFF)) 

381 | ((y.astype(np.uint32) & np.uint32(0x3FFF)) << np.uint32(14)) 

382 | ((p.astype(np.uint32) & np.uint32(0x1)) << np.uint32(28)) 

383 ) 

384 self._fd.write(out.tobytes()) 

385 self._n_written_events += n 

386 return n