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
« prev ^ index » next coverage.py v7.15.1, created at 2026-07-18 05:24 +0000
1"""Prophesee DAT (.dat) CD-event decoder/encoder.
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
11import io
12from datetime import datetime
14import numpy as np
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
33_EMPTY_EVENTS = EventArray.empty()
35DAT_EVENT_TYPE_CD = 0x0C
36DAT_EVENT_SIZE = 0x08
38class EventDecoder_Dat(EventDecoder):
39 """Decode Prophesee DAT files into ``EventArray`` chunks.
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.
49 References
50 ----------
51 [1] Prophesee DAT file format
52 https://docs.prophesee.ai/stable/data/file_formats/dat.html
54 """
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
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
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
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
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
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
117 def init(self) -> None:
118 """Initialize the DAT reader.
120 Returns
121 -------
122 None
124 """
125 if self._is_initialized:
126 return
128 if self._source.mappable():
129 self._buf = self._source.buffer()
130 else:
131 self._buf = memoryview(self._source.read(-1))
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
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)
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
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
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()
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
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)
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)
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
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.
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)
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))
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])
246 if axis == "t":
247 bucket = int(val >> 32)
248 val_rem = int(val & 0xFFFFFFFF)
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))
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)
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
272 def reset(self) -> None:
273 """Reset the DAT reader to the beginning.
275 Returns
276 -------
277 None
279 """
280 self._offset = 0
281 self._eof = False
282 if self._parser is not None:
283 self._parser.reset()
285 def tell(self) -> int:
286 """Get the current byte offset.
288 Returns
289 -------
290 int
291 Current byte offset.
293 """
294 return self._payload_off + self._offset * 4
296 def close(self) -> None:
297 """Close the DAT reader.
299 Returns
300 -------
301 None
303 """
304 self._words = None
305 self._buf = None
307class EventEncoder_Dat(EventEncoder):
308 """Encode events into a Prophesee DAT (.dat) CD file.
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).
321 References
322 ----------
323 [1] Prophesee DAT file format
324 https://docs.prophesee.ai/stable/data/file_formats/dat.html
326 """
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
333 def init(self) -> None:
334 """Initialize the DAT writer.
336 Returns
337 -------
338 None
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
354 def write(self, events: 'np.ndarray | EventArray', triggers: 'np.ndarray | TriggerArray | None' = None) -> int:
355 """Write events to the DAT file.
357 Parameters
358 ----------
359 events : np.ndarray or EventArray
360 Array of events to write.
362 Returns
363 -------
364 int
365 Number of written events.
367 """
368 if not self._is_initialized:
369 self.init()
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"]
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