Engine¶
frames2py.Engine(sensor_size, kernel="event_count", *, snapshot_interval_ms=16.0) is the
live runtime: it accumulates events from one producer thread and publishes snapshots for any
number of consumers.
import numpy as np
import frames2py
# 10,000 synthetic events on a 640x480 sensor, one every microsecond.
rng = np.random.default_rng(seed=0)
events = np.zeros(10_000, dtype=frames2py.EVENT_DTYPE)
events["t"] = np.arange(10_000) # timestamps, µs
events["x"] = rng.integers(0, 640, 10_000) # column
events["y"] = rng.integers(0, 480, 10_000) # row
events["p"] = rng.integers(0, 2, 10_000) # polarity: 0 is OFF, anything else ON
engine = frames2py.Engine((640, 480), "event_count") # sensor_size is (width, height)
engine.ingest(events) # the first ingest() always publishes
snapshot = engine.snapshot() # the latest publication: shared, read-only
print(snapshot.frame.shape, snapshot.frame.dtype)
print("events counted:", int(snapshot.frame.sum()))
print("watermark:", snapshot.meta.watermark, "sequence:", snapshot.meta.sequence)
print("ingested:", engine.stats.events_ingested, "out of bounds:", engine.stats.events_out_of_bounds)
Constructor¶
sensor_size:(width, height).kernel: a configured kernel instance, or"event_count"(the default),"polarity"or"time_surface".ExpDecayandTimestampDecaytake a parameter, so they are always passed as instances;"exp_decay"as a name raisesValueError.snapshot_interval_ms(keyword-only): the publication cadence, below.
On a free-threaded CPython build running with the GIL disabled, construction raises
RuntimeError unless that minor version has been verified; today that is 3.14 only (see
Lifecycle and threads).
ingest(events)¶
Accumulates one call's events and publishes if the cadence allows, all on the caller's thread, before returning. The steps are on the Architecture page and the validation rules in the event contract.
- Raises
RuntimeErrorif called from a thread other than the producer's, changing nothing. The firstingest()that doesn't raise makes its thread the producer, for the Engine's lifetime. - Raises
TypeErrorfor a malformed array andValueErrorif any event hast >= 2**63, in both cases before any state or statistic changes. - A no-op while the Engine is stopped.
Publication cadence¶
snapshot_interval_ms=0publishes on everyingest().- A positive interval publishes at most once per interval, on the first
ingest()after it has elapsed since the last publication (measured withtime.monotonic_ns()). - The first
ingest()always publishes, and so does the first afterreset(). That includes an empty call, or one whose events are all out of bounds: the snapshot then carrieswatermark=Noneuntil an in-bounds event arrives. stop()publishes if in-bounds events were accumulated since the last publication.- Nothing else publishes. There is no timer thread, so if the producer stops calling
ingest()mid-window, that window's events stay unpublished until the nextingest()orstop().
For windowed kernels (event_count, polarity) each publication closes the window: the
next snapshot counts only the events after it. Running kernels are not changed by
publication.
The 16 ms default is about one publication per frame of a 60 Hz display. A consumer that wants every publication should read at about the same cadence; one that reads less often simply sees a later snapshot.
snapshot()¶
Returns the latest published Snapshot, or None before the first
publication and after reset() until the next. It returns the published object itself, not
a copy: every consumer that reads the same publication gets the same object, with a
read-only frame. Takes no lock. Safe from any thread.
The Engine has no watermark attribute of its own: the watermark at each publication is
snapshot.meta.watermark.
Seeing every publication¶
Publication happens only inside ingest() and stop(),
which run on the threads that call them. So code that reads snapshot() right after each
of its own ingest() calls, and once after stop(), sees every publication, provided no
other thread calls reset() meanwhile. A consumer on another thread that polls at its own
pace may skip publications; that is by design. Summing the windows of a windowed kernel
this way gives totals over the whole stream (an Accumulator gives them
without publication):
import numpy as np
import frames2py
from frames2py.adapters import evt
# Whole-file totals through an Engine with a windowed kernel: read the snapshot after every
# ingest() and after stop(), and add each new publication's window once (by its sequence).
# Any snapshot_interval_ms works; 0 here only makes every call publish, so the output below
# doesn't depend on timing. The totals are uint64, so they can't wrap as uint32 counts can.
totals, last = None, None
def collect(engine):
global totals, last
snapshot = engine.snapshot()
if snapshot is not None and snapshot.meta.sequence != last:
totals = snapshot.frame.astype(np.uint64) if totals is None else totals + snapshot.frame
last = snapshot.meta.sequence
with evt.open("sparklers_100k.evt2.raw", sensor_size=(640, 480), batch_size=10_000) as reader:
engine = frames2py.Engine(reader.sensor_size, "polarity", snapshot_interval_ms=0)
for events in reader:
engine.ingest(events)
collect(engine)
engine.stop()
collect(engine)
print("publications:", last, "events:", int(totals.sum()))
print("OFF:", int(totals[..., 0].sum()), "ON:", int(totals[..., 1].sum()))
print("last window only:", int(engine.snapshot().frame.sum()), "events")
stats¶
An EngineStats, a frozen dataclass, built when you read the property. Takes no lock.
| field | meaning |
|---|---|
events_ingested |
every event of every accepted ingest() call, out-of-bounds ones included |
events_out_of_bounds |
the subset the bounds check rejected |
snapshots_published |
publications since construction or the last reset() |
uptime_ns |
nanoseconds since the Engine was constructed; reset() doesn't change it |
The accumulated count is events_ingested - events_out_of_bounds; there is no separate
field. Every event of an accepted call to a running Engine is either accumulated or counted
as out of bounds. A call rejected by validation or by the timestamp-range check changes no
statistic, and neither does ingest() while stopped.
stats is a diagnostic view: its fields are read one after another while the producer may
be running, so on a live Engine they are not one instant-consistent snapshot of all
counters.
start(), stop(), reset()¶
In short: stop() publishes any pending window and makes ingest() a no-op, leaving the
last snapshot readable; start() resumes; reset() clears the kernel state, counters,
watermark and published snapshot. They may be called from any thread. The details are in
Lifecycle and threads.