diff --git a/v2ecoli/composites/ecoli_baseline.py b/v2ecoli/composites/ecoli_baseline.py index 382377cd8..7b3075896 100644 --- a/v2ecoli/composites/ecoli_baseline.py +++ b/v2ecoli/composites/ecoli_baseline.py @@ -129,7 +129,7 @@ def _listener_leaf_paths(listeners: dict, *, prefix: str = "listeners"): def _single_cell_xarray_config(*, out_uri: str, metadata: dict | None = None, - buffer_size: int = 3) -> dict: + buffer_size: int = 600) -> dict: """Build the STATIC XArrayEmitter ``config`` skeleton for a single-cell, agent-relative, in-document capture. @@ -162,7 +162,9 @@ def _single_cell_xarray_config(*, out_uri: str, metadata: dict | None = None, out_uri: zarr store path/URI. metadata: non-empty run-identity metadata (experiment_id / variant / lineage_seed). Falls back to a non-empty placeholder if omitted. - buffer_size: transducer buffer size (streaming, bounded). Default 3. + buffer_size: transducer buffer size (streaming, bounded), in emit steps. + Default 600 — matches the viva-emitters library default; flushes a + handful of times per generation rather than every few steps. Returns: The static XArrayEmitter config skeleton (no ``view`` / diff --git a/v2ecoli/library/vivarium_ecoli_engine.py b/v2ecoli/library/vivarium_ecoli_engine.py index b9a3fc5f8..d881eef69 100644 --- a/v2ecoli/library/vivarium_ecoli_engine.py +++ b/v2ecoli/library/vivarium_ecoli_engine.py @@ -802,7 +802,9 @@ def run_vivarium_ecoli_pbg_multigen( em = _build_emitter( core=core, store_path=store_path, view=view, metadata_base=metadata_base, generation=gen + 1, # 1-indexed to match run_multigen_xarray (v2ecoli side) - agent_id=partition_agent_id, output_metadata={}, buffer_size=3) + # Inherit build_emitter_config's buffer_size default (600): flush a + # handful of times per generation, not every few steps. + agent_id=partition_agent_id, output_metadata={}) steps = 1 divided = False diff --git a/v2ecoli/library/xarray_run.py b/v2ecoli/library/xarray_run.py index a9d58fb9a..69672336c 100644 --- a/v2ecoli/library/xarray_run.py +++ b/v2ecoli/library/xarray_run.py @@ -386,7 +386,7 @@ def build_emitter_config( metadata_base: dict, generation: int, agent_id: str, - buffer_size: int = 4, + buffer_size: int = 600, output_metadata: dict | None = None, writer: dict | None = None, predicate: list | None = None, @@ -466,7 +466,7 @@ def run_multigen_xarray( chunk: int = 60, initial_agent_id: str = "0", overwrite: bool = True, - buffer_size: int = 3, + buffer_size: int = 600, single_daughters: bool = False, division_detector: Callable[[set[str], set[str]], tuple[bool, str | None]] | None = None, provenance: dict | None = None, @@ -485,12 +485,17 @@ def run_multigen_xarray( chunk: how many ticks between emitter updates. initial_agent_id: agent_id to start following. overwrite: if True, delete ``store_path`` before starting. - buffer_size: XArrayEmitter transducer buffer size. Default 3 (NOT the - config builder's 4): the installed pbg-emitters trips an - ``assert not include_static`` in its ``flush(final=True)`` path when the - buffer is exactly full at close (i.e. ``n_updates % buffer_size == 0``); - 3 is the value the single-generation runner uses to dodge it. Override - only if you know the emit count won't land on a multiple of it. + buffer_size: XArrayEmitter transducer buffer size, in *emit steps*, held + in memory before each flush to the zarr store. Default 600, matching the + viva-emitters library default — sized to flush only a handful of times + per generation (see :py:class:`viva_emitters.xarray_emitter.transducer. + XarrayTransducer`). A tiny buffer (the old default of 3) forces a flush + every few emit steps, degrading latency and compression by ~2 orders of + magnitude; it was originally chosen to dodge a ``flush(final=True)`` + assertion in an early vendored emitter, which current viva-emitters has + since fixed (the buffer-full/partial/aligned close paths are all + regression-tested). The ``except AssertionError`` guards around the + closes below are now belt-and-suspenders against that fixed quirk. division_detector: optional ``(prev_ids, curr_ids) -> (divided?, daughter_id|None)``. Default: detect division when ``len(curr) > len(prev)`` and pick the first new agent_id sorted. diff --git a/v2ecoli/workflow/lineage.py b/v2ecoli/workflow/lineage.py index 56caf5ab8..184a08bd1 100644 --- a/v2ecoli/workflow/lineage.py +++ b/v2ecoli/workflow/lineage.py @@ -270,7 +270,10 @@ def _open_xarray_emitter(self, emit_cell): raw_view = [dict(e, root=tuple(e["root"])) for e in raw_view] transducer = arg.get("transducer") or {} buf = ((transducer.get("buffer") or {}).get("size")) - buf = max(3, int(buf or 4)) # transducer requires buffer.size > 2 + # Default 600 (viva-emitters library default: a handful of flushes per + # generation, not one every few steps); floor 3 since the transducer + # requires buffer.size > 2. + buf = max(3, int(buf or 600)) predicate = transducer.get("predicate") writer = arg.get("writer") out_dir = arg.get("out_dir") or self.config["out_dir"]