Skip to content

refactor(engine): extract PipelineSupervisor - process-singleton assembly + GC/recover + finalize hook - #171

Merged
zc277584121 merged 10 commits into
zilliztech:mainfrom
code2tan:rebuild_pipeline_supervisor
Jul 14, 2026
Merged

zc277584121 merged 10 commits into
zilliztech:mainfrom
code2tan:rebuild_pipeline_supervisor

Conversation

@code2tan

Copy link
Copy Markdown
Contributor

Resolve #166

Overview

Moves three categories of code out of Engine and into a new PipelineSupervisor component: pipeline process-singleton assembly, the startup/shutdown pipeline half (orphan GC + Job Lane recovery + watcher start/stop), and the per-object finalize hook. Engine holds self.pipeline; startup/shutdown collapse to two lines each:

async def startup(self, *, preload_local_models=False):
    await self.infra.startup(preload_local_models=preload_local_models)
    await self.pipeline.startup()

async def shutdown(self):
    await self.pipeline.shutdown()
    await self.infra.shutdown()

This is the second half of Engine redesign phase 4 (4b), continuing from 4a's InfraStack (the infra layer). Once 4a consolidated the 8 infra clients, the remaining lines of Engine.startup/shutdown were purely pipeline semantics; 4b moves them out wholesale, clearing the ground for the later IngestOrchestrator / WorkerScheduler split.


Why it's needed

Engine is a god class of ~60 methods with 9 coupled responsibility areas. On the pipeline side it has three specific problems (mirror of 4a's infra layer):

  1. Assembly & lifecycle scattered across four sites: process-singleton fields (_chunks_q / _embed_consumer / _producer_ctx / _job_lane / two concurrency gates / watcher) are declared in __init__, constructed in _build_pipeline, started in startup, stopped in shutdown - four sites, with construction-order dependencies (gates before ProducerContext/JobLane; register_on_succeeded order decides finalize priority) expressed only in implicit code ordering.

  2. The finalize hook is buried in Engine: _on_pipeline_object_indexed is a ~50-line atomic method body (claim → won-check → delete-or-write) guarding the core invariant "a chunk exists in Milvus iff a committed objects row points at it" - yet it sits in the same class as drain/retry/run_job orchestration, hard to verify its three exits (success / error / won==0) in isolation.

  3. Tests can only go E2E: pipeline assembly is coupled to Engine, so PipelineSupervisor cannot be built in isolation to verify the startup call sequence, the pump error-pop, or the finalize hook's three branches - the only option is to spin up the whole Engine and run the full pipeline.

After extraction: assembly + lifecycle + finalize live in 383 independently readable lines; the startup/shutdown sequences are independently unit-testable (recording mocks, asserting build → gc → recover → watcher order); _on_object_indexed's three exits each get a dedicated test.


Changes

New engine/pipeline_supervisor.py (383 lines)

The PipelineSupervisor class + PipelineEmbedConsumer (formerly engine.py's _PipelineEmbedConsumer, renamed public since it is no longer Engine-internal). Single responsibility: process-singleton assembly + lifecycle + non-blocking pump + atomic finalize hook.

  • Constructor: PipelineSupervisor(cfg, infra, artifacts, objects, factory) - takes InfraStack plus the three shipped components (ArtifactCacheService / ObjectRepository / ConnectorFactory).
  • Public attributes (no underscore; Engine reads via self.pipeline.<attr>, no Engine-level forwarding properties): embed_consumer / producer_ctx / job_lane / job_watcher / job_watcher_task.
  • Public methods: routes_to_pipeline(okind), stash_finalize(full_uri, ctx), pump(...), startup(), shutdown().
  • Internal methods (migrated line-for-line): _build_pipeline, _gc_orphan_chunks, _recover_job_lane, _on_object_indexed (formerly _on_pipeline_object_indexed).
  • __init__: lazy singleton fields (_chunks_q/embed_consumer/producer_ctx all None, _pending_finalize = {}).
  • startup(): _build_pipeline → _gc_orphan_chunks → _recover_job_lane → ConnectorJobWatcher(meta, job_lane) + create_task(watcher.run).
  • shutdown(): watcher.stop + await watcher_task → job_lane.stop → embed_consumer.shutdown (set to None).
  • _build_pipeline: builds chunks_q + PipelineEmbedConsumer (4 adapters) + two gates + ProducerContext + JobLane; register_on_succeeded(_on_object_indexed) before register_on_succeeded(job_lane.on_embed_succeeded); consumer.start + job_lane.start. ArtifactStoreAdapter is wired directly to self._art.put_artifact etc., bypassing Engine's thin delegates.
  • pump: non-blocking producer→chunks_q pump; on exception pop _pending_finalize[full_uri] then re-raise; finally GCs the message_stream raw_records temp artifact.
  • _on_object_indexed: atomic body unchanged - error branch writes failed row + advance_task(FAILED); won==0 deletes orphan chunks; else write_object_row + plugin.on_object_indexed.

Engine changes (engine/engine.py, 2460 → 2112)

  • Removed: _build_pipeline / _routes_to_pipeline / _gc_orphan_chunks / _recover_job_lane / _on_pipeline_object_indexed / _index_via_pipeline / the finalize call site in _write_object_row / the _PIPELINE_OKINDS constant / the _PipelineEmbedConsumer class.
  • __init__: 10 pipeline singleton fields → self.pipeline = PipelineSupervisor(cfg, self.infra, self.artifacts, self.objects, self.connector_factory).
  • startup / shutdown: collapsed to two lines each (infra + pipeline).
  • 13 external touch sites rewritten to self.pipeline.<...>: _drain_job (register_job / routes_to_pipeline / on_yield_object_change / on_sync_done), _finalize_job (evict_job), cancel_job (mark_job_cancelled), _process_with_retry (on_task_retry), _run_job (await_done), _index_object (routes_to_pipeline / stash_finalize / pump).
  • Import cleanup: the migrated pipeline / producers / adapters / job_lane / job_watcher / storage.ids imports are dropped; from .pipeline_supervisor import PipelineSupervisor is added.

common/embedding.py

Two private methods of CachingEmbeddingClient made public: _key → key, _embed_api → embed_api (including the call sites inside batch_embed). Lets the supervisor read public API rather than reaching into Engine-private members.

engine/pipeline.py

EmbedConsumer.start return type asyncio.Task → Task | None (more accurate; callers no longer assume a Task is always present).

Tests

  • New tests/test_pipeline_supervisor.py (429 lines, 11 unit tests): lazy construction, Engine forwarding, routes_to_pipeline table, ArtifactStoreAdapter wiring, stash_finalize + success pop, error branch writes failed, won==0 purges orphan, pump pops on produce error, recover uses factory, startup order, shutdown order.
  • Mechanical rewrite across 12 test files: eng._build_pipeline() → eng.pipeline._build_pipeline() (or await eng.pipeline.startup()), eng._pending_finalize[...] = (...) → eng.pipeline.stash_finalize(...), eng._on_pipeline_object_indexed(...) → eng.pipeline._on_object_indexed(...), eng._job_lane = ... → eng.pipeline.job_lane = ..., reads eng._embed_consumer / eng._job_lane → eng.pipeline.embed_consumer / eng.pipeline.job_lane.

File list

File Change
engine/pipeline_supervisor.py new (+383)
engine/engine.py removed 6 methods + class + constant; __init__/startup/shutdown collapsed; 13 touch sites rewritten (−348 lines)
common/embedding.py _key/_embed_api made public
engine/pipeline.py start return type Task | None
tests/test_pipeline_supervisor.py new (+429, 11 tests)
12 tests/test_engine_*.py etc. touch-site contract rewrite

Design decisions

  • Public attributes, no Engine-level forwarding properties. The 4b design doc proposed _embed_consumer/_job_lane private + Engine @property shims (so the 9 field accesses wouldn't change). The actual implementation goes further: supervisor fields are public (embed_consumer/job_lane/job_watcher/producer_ctx), and Engine keeps no forwarding properties - call sites write self.pipeline.job_lane directly. This matches 4a's "direct access, no shim" decision: zero boilerplate, a single access path, ownership visible at every call site. The cost is 13 mechanical replacements, covered by tests.

  • _pending_finalize encapsulated as stash_finalize. The current touches are "2 pops (internal to the supervisor) + 1 write site (_index_object)". After migration both pops become internal; the only external touch is the write site, wrapped as stash_finalize(full_uri, ctx). The _pending_finalize field is not exposed via a property; the supervisor fully owns that dict's lifecycle.

  • ArtifactStoreAdapter wired directly to ArtifactCacheService. The supervisor builds the adapter with self._art.put_artifact / read_artifact / read_artifact_fresh, bypassing Engine's _put_artifact / _read_artifact / _read_artifact_fresh thin delegates (signatures are parameter-compatible). Engine's three thin delegates stay (the _index_object inline tail still uses _write_object_row; the artifact LRU responsibilities are not migrated in this round).

  • Constructor takes factory. _recover_job_lane needs connector_factory.build_plugin(...) to rebuild a crashed job's plugin; that is a startup responsibility of the supervisor and cannot be back-filled by Engine (that would require a half-initialized state before supervisor.startup()). It is passed in the constructor, peer to infra/artifacts/objects. Inside, the supervisor unwraps BuiltPlugin.plugin.

  • _on_object_indexed stays atomic and unsplittable. The claim → won-check → delete-or-write sequence stays in one method body, migrated line-for-line, with no split and no intermediate ObjectRepository method. The won == 0 (concurrent cancel) ⇒ delete-orphan branch must sit next to advance_task(SUCCEEDED, from=RUNNING); the _pending_finalize pop stays at the top as the entry guard (short-circuits on ctx is None, skipping the Job Lane dir_summary success which carries no stash).

  • _write_object_row stays on Engine. A 2-line delegate called from two sites: the supervisor's finalize hook (migrated, calls self._obj.write_object_row directly) and Engine's _index_object inline tail (stays on Engine, belongs to the future IngestOrchestrator). Migrating it would duplicate a 2-line method in two places.

  • CachingEmbeddingClient private methods made public. _key / _embed_api → key / embed_api, so the supervisor reads public API rather than Engine-private members - a cleaner boundary.

  • Location engine/pipeline_supervisor.py (sibling to 4a's engine/infra.py), not components/ - components/ currently holds only shipped repo/factory components (ObjectRepository / ArtifactCacheService / ConnectorFactory, each with external importers and unit tests); PipelineSupervisor is an Engine-internal collaborator with no external import.

  • Type-checking & error-handling hardening (commit 38e2a05): _build_pipeline and pump build Optional fields into locals first and publish last (Pyright does not narrow an instance attribute across await / calls); producer.produce is declared async def -> AsyncIterator but the implementations are async generators, so it is cast(AsyncIterator, ...) and iterated directly (no await); pump's except widens from Exception to BaseException so CancelledError also runs the cleanup (pop then re-raise); the watcher-task and raw_records-GC bare pass / logger.error(e) become logger.exception(...), making failures observable.

  • Finalize-path exception protection is out of scope. docs-dev/fix-pipeline-supervisor-finalize-race.md proposes wrapping _on_object_indexed's three branches in try/except + compensation logging (_finalize_errors), to cover metadata/Milvus inconsistency when delete_by_object / write_object_row / advance_task / plugin.on_object_indexed fail. That is a behavior change (exception paths go from silently swallowed to explicitly recorded + compensated), not a structural extraction, so it is deliberately deferred to a standalone follow-up PR after the migration lands.


Behavioral equivalence

PipelineSupervisor's methods replicate the lifecycle and logic in exactly the same order as the original Engine:

  • startup sequence: original infra.startup → _build_pipeline → _gc_orphan_chunks → _recover_job_lane → watcher + create_task. After migration await self.infra.startup(...) → await self.pipeline.startup(), the latter executing the same steps in the same order.
  • shutdown sequence: original watcher.stop + await task → job_lane.stop → embed_consumer.shutdown → infra.shutdown. After migration await self.pipeline.shutdown() (first three) → await self.infra.shutdown(), order preserved. Watcher-task exceptions go from bare pass to logger.exception (still does not abort shutdown; observability only).
  • _build_pipeline assembly: of the 4 adapters, EmbedderAdapter / MilvusSinkAdapter / TxCacheAdapter take self._infra.* (same instances); ArtifactStoreAdapter takes self._art.* (signature-compatible, parameter-equivalent to the original Engine thin delegates). The two gates, ProducerContext, build_job_lane, the register_on_succeeded order (_on_object_indexed before job_lane.on_embed_succeeded), and consumer.start + job_lane.start are all original order and args. Building into locals and publishing last is a type-checker device, runtime-equivalent.
  • _pending_finalize three touches: write site → stash_finalize (same tuple value); the two pops migrate inside the supervisor, accessing the same dict. The cancel path is unchanged (mark_job_cancelled → consumer _fail_batch → fires _on_object_indexed with error → pop; no residue).
  • _on_object_indexed atomic body: migrated line-for-line, logic unchanged.
  • pump: the original _index_via_pipeline body migrates; the message_stream raw_records GC stays in finally, with bare pass → logger.exception (GC failure still does not affect the task; observability only).
  • 13 touch sites: what forwards through public attributes is the same JobLane / PipelineEmbedConsumer instance (the ones the supervisor constructs in _build_pipeline); identity preserved.
  • Engine.startup / shutdown external signatures and semantics are unchanged; call sites in api/app.py and __main__.py are unchanged.

Test impact

New unit tests (tests/test_pipeline_supervisor.py, 11 tests):

  1. Lazy construction: _chunks_q / embed_consumer / producer_ctx all None, _pending_finalize == {}.
  2. Engine forwarding: eng.pipeline is non-None; eng.pipeline._obj is eng.objects, _art is eng.artifacts, _factory is eng.connector_factory, _infra is eng.infra.
  3. routes_to_pipeline table: _PIPELINE_OKINDS values return True; image follows description.enabled; table_schema follows summary.enabled; else False.
  4. ArtifactStoreAdapter wiring: asserts the adapter holds artifacts.put_artifact / read_artifact / read_artifact_fresh (bypassing Engine's thin delegates).
  5. stash_finalize + success pop: after stash_finalize(uri, ctx) the dict has one entry; after _on_object_indexed fires it is popped.
  6. error branch writes failed: _on_object_indexed(error=...) writes the failed row + advance_task(FAILED).
  7. won==0 purges orphan: when advance_task returns 0, milvus.delete_by_object runs and no row is written.
  8. pump pops on produce error: when producer.produce raises, _pending_finalize.pop(full_uri) runs and the exception re-raises.
  9. recover uses factory: factory.build_plugin returns a BuiltPlugin; recover_job receives built.plugin.
  10. startup order: build → gc → recover → ConnectorJobWatcher(meta, job_lane) + create_task.
  11. shutdown order: watcher.stop → await task → job_lane.stop → embed_consumer.shutdown, set to None.

Existing tests mechanically rewritten: 12 files of touch-site contract change (eng._build_pipeline() / eng._pending_finalize[...] / eng._on_pipeline_object_indexed(...) / eng._job_lane = ... go through eng.pipeline.*). No logic changes; ruff format passes throughout.

Regression result: full uv run --extra dev pytest - 382 passed / 11 skipped, behavior-equivalent to before the refactor.

Pre-existing failures (unrelated to this PR):

  • tests/test_connector_factory.py::TestValidateConfig (config validation not raising ValueError) - fails on the baseline.
  • tests/test_feishu_oauth.py / tests/test_enumeration_since.py, 2 collection errors - the optional connector SDK lark_oapi is not installed (documented in CLAUDE.md).

Risk & rollback

  • Risk: pure structural extraction, with no schema / protocol / data changes. The _on_object_indexed atomic body is preserved line-for-line, the startup/shutdown sequences step-for-step, and handle identity is preserved. The main risk is a missed site among the 13 mechanical rewrites and the _pending_finalize touches - covered by the full suite (382 passed).
  • Finalize-path exception protection (not done): failures inside _on_object_indexed (write_object_row / delete_by_object / advance_task / plugin.on_object_indexed) still fail silently inconsistent (same behavior as before the refactor, not introduced here). The fix is in docs-dev/fix-pipeline-supervisor-finalize-race.md, deferred to a standalone follow-up.
  • Rollback: git revert the 7 commits fully reverts the change, with no side effects. Rolling back restores the fields + methods to Engine and removes the public attributes and stash_finalize / pump.

@zc277584121

Copy link
Copy Markdown
Collaborator

⚠️ CI-equivalent test run is red on the current head (efe9ce2)

Running uv run --extra dev pytest on this branch gives 25 failed / 357 passed, not the 382 passed the description claims. Root cause: this PR promotes CachingEmbeddingClient._embed_api → embed_api and _key → key (public) and updates the call sites in pipeline_supervisor.py, but the 10 _FakeEmbed test doubles across the test suite still define the old private _embed_api / _key, so every test that swaps in _FakeEmbed and builds the pipeline hits AttributeError: '_FakeEmbed' object has no attribute 'embed_api'.

Fix is a mechanical rename in the fakes (_embed_api→embed_api, _key→key) — with it applied the suite is back to 382 passed, 11 skipped. Files: test_pipeline_supervisor.py, test_artifact_adapter.py, and the 7 test_engine_*_e2e.py / test_engine_claim_global.py.

Two smaller notes: (1) this branch is behind main (missing #172/#173) and will need a rebase; (2) it doesn't touch the eng().meta → eng().infra.meta breakage in api/app.py (that's #170) — worth landing #170 first.

code2t added 2 commits July 13, 2026 14:20
  Rename _FakeEmbed._key -> key and _FakeEmbed._embed_api ->
  across the e2e test fakes so they mirror the public surface of the
  CachingEmbeddingClient. Update the EmbedderAdapter docstring to
  embed_api accordingly. This finishes the earlier "use public
  and methods" refactor by removing the last private-method references
  test doubles.
# Conflicts:
#	server/python/src/mfs_server/engine/engine.py
@zc277584121
zc277584121 merged commit 2deee33 into zilliztech:main Jul 14, 2026
8 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

refactor(engine): split the Engine monolithic class into a Facade + single-responsibility collaborating components

2 participants