Repository navigation
refactor(engine): extract PipelineSupervisor - process-singleton assembly + GC/recover + finalize hook - #171
Conversation
…er, job_lane, and job_watcher
…tter type checking
|
Running Fix is a mechanical rename in the fakes ( Two smaller notes: (1) this branch is behind |
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
Resolve #166
Overview
Moves three categories of code out of
Engineand into a newPipelineSupervisorcomponent: pipeline process-singleton assembly, thestartup/shutdownpipeline half (orphan GC + Job Lane recovery + watcher start/stop), and the per-object finalize hook.Engineholdsself.pipeline;startup/shutdowncollapse to two lines each: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 ofEngine.startup/shutdownwere purely pipeline semantics; 4b moves them out wholesale, clearing the ground for the laterIngestOrchestrator/WorkerSchedulersplit.Why it's needed
Engineis 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):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 instartup, stopped inshutdown- four sites, with construction-order dependencies (gates before ProducerContext/JobLane;register_on_succeededorder decides finalize priority) expressed only in implicit code ordering.The finalize hook is buried in Engine:
_on_pipeline_object_indexedis a ~50-line atomic method body (claim→won-check →delete-or-write) guarding the core invariant "a chunk exists in Milvus iff a committedobjectsrow 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.Tests can only go E2E: pipeline assembly is coupled to
Engine, soPipelineSupervisorcannot 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 wholeEngineand run the full pipeline.After extraction: assembly + lifecycle + finalize live in 383 independently readable lines; the
startup/shutdownsequences 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
PipelineSupervisorclass +PipelineEmbedConsumer(formerlyengine.py's_PipelineEmbedConsumer, renamed public since it is no longer Engine-internal). Single responsibility: process-singleton assembly + lifecycle + non-blocking pump + atomic finalize hook.PipelineSupervisor(cfg, infra, artifacts, objects, factory)- takesInfraStackplus the three shipped components (ArtifactCacheService/ObjectRepository/ConnectorFactory).self.pipeline.<attr>, no Engine-level forwarding properties):embed_consumer/producer_ctx/job_lane/job_watcher/job_watcher_task.routes_to_pipeline(okind),stash_finalize(full_uri, ctx),pump(...),startup(),shutdown()._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_ctxallNone,_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: buildschunks_q+PipelineEmbedConsumer(4 adapters) + two gates +ProducerContext+JobLane;register_on_succeeded(_on_object_indexed)beforeregister_on_succeeded(job_lane.on_embed_succeeded);consumer.start+job_lane.start.ArtifactStoreAdapteris wired directly toself._art.put_artifactetc., bypassing Engine's thin delegates.pump: non-blocking producer→chunks_qpump; on exceptionpop _pending_finalize[full_uri]then re-raise;finallyGCs themessage_streamraw_recordstemp artifact._on_object_indexed: atomic body unchanged -errorbranch writes failed row +advance_task(FAILED);won==0deletes orphan chunks; elsewrite_object_row+plugin.on_object_indexed.Enginechanges (engine/engine.py, 2460 → 2112)_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_OKINDSconstant / the_PipelineEmbedConsumerclass.__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).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).pipeline/producers/adapters/job_lane/job_watcher/storage.idsimports are dropped;from .pipeline_supervisor import PipelineSupervisoris added.common/embedding.pyTwo private methods of
CachingEmbeddingClientmade public:_key→key,_embed_api→embed_api(including the call sites insidebatch_embed). Lets the supervisor read public API rather than reaching into Engine-private members.engine/pipeline.pyEmbedConsumer.startreturn typeasyncio.Task→Task | None(more accurate; callers no longer assume a Task is always present).Tests
tests/test_pipeline_supervisor.py(429 lines, 11 unit tests): lazy construction, Engine forwarding,routes_to_pipelinetable,ArtifactStoreAdapterwiring,stash_finalize+ success pop, error branch writes failed,won==0purges orphan, pump pops on produce error, recover uses factory, startup order, shutdown order.eng._build_pipeline()→eng.pipeline._build_pipeline()(orawait 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 = ..., readseng._embed_consumer/eng._job_lane→eng.pipeline.embed_consumer/eng.pipeline.job_lane.File list
engine/pipeline_supervisor.pyengine/engine.py__init__/startup/shutdowncollapsed; 13 touch sites rewritten (−348 lines)common/embedding.py_key/_embed_apimade publicengine/pipeline.pystartreturn typeTask | Nonetests/test_pipeline_supervisor.pytests/test_engine_*.pyetc.Design decisions
Public attributes, no Engine-level forwarding properties. The 4b design doc proposed
_embed_consumer/_job_laneprivate + Engine@propertyshims (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 writeself.pipeline.job_lanedirectly. 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_finalizeencapsulated asstash_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 asstash_finalize(full_uri, ctx). The_pending_finalizefield is not exposed via a property; the supervisor fully owns that dict's lifecycle.ArtifactStoreAdapterwired directly toArtifactCacheService. The supervisor builds the adapter withself._art.put_artifact/read_artifact/read_artifact_fresh, bypassing Engine's_put_artifact/_read_artifact/_read_artifact_freshthin delegates (signatures are parameter-compatible). Engine's three thin delegates stay (the_index_objectinline tail still uses_write_object_row; the artifact LRU responsibilities are not migrated in this round).Constructor takes
factory._recover_job_laneneedsconnector_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 beforesupervisor.startup()). It is passed in the constructor, peer toinfra/artifacts/objects. Inside, the supervisor unwrapsBuiltPlugin.plugin._on_object_indexedstays atomic and unsplittable. Theclaim → won-check → delete-or-writesequence stays in one method body, migrated line-for-line, with no split and no intermediateObjectRepositorymethod. Thewon == 0(concurrent cancel) ⇒ delete-orphan branch must sit next toadvance_task(SUCCEEDED, from=RUNNING); the_pending_finalizepop stays at the top as the entry guard (short-circuits onctx is None, skipping the Job Lane dir_summary success which carries no stash)._write_object_rowstays on Engine. A 2-line delegate called from two sites: the supervisor's finalize hook (migrated, callsself._obj.write_object_rowdirectly) and Engine's_index_objectinline tail (stays on Engine, belongs to the future IngestOrchestrator). Migrating it would duplicate a 2-line method in two places.CachingEmbeddingClientprivate 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'sengine/infra.py), notcomponents/-components/currently holds only shipped repo/factory components (ObjectRepository/ArtifactCacheService/ConnectorFactory, each with external importers and unit tests);PipelineSupervisoris an Engine-internal collaborator with no external import.Type-checking & error-handling hardening (commit
38e2a05):_build_pipelineandpumpbuildOptionalfields into locals first and publish last (Pyright does not narrow an instance attribute acrossawait/ calls);producer.produceis declaredasync def -> AsyncIteratorbut the implementations are async generators, so it iscast(AsyncIterator, ...)and iterated directly (noawait);pump'sexceptwidens fromExceptiontoBaseExceptionsoCancelledErroralso runs the cleanup (pop then re-raise); the watcher-task andraw_records-GC barepass/logger.error(e)becomelogger.exception(...), making failures observable.Finalize-path exception protection is out of scope.
docs-dev/fix-pipeline-supervisor-finalize-race.mdproposes wrapping_on_object_indexed's three branches in try/except + compensation logging (_finalize_errors), to cover metadata/Milvus inconsistency whendelete_by_object/write_object_row/advance_task/plugin.on_object_indexedfail. 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 originalEngine:infra.startup→_build_pipeline→_gc_orphan_chunks→_recover_job_lane→ watcher +create_task. After migrationawait self.infra.startup(...)→await self.pipeline.startup(), the latter executing the same steps in the same order.watcher.stop+await task→job_lane.stop→embed_consumer.shutdown→infra.shutdown. After migrationawait self.pipeline.shutdown()(first three) →await self.infra.shutdown(), order preserved. Watcher-task exceptions go from barepasstologger.exception(still does not abort shutdown; observability only)._build_pipelineassembly: of the 4 adapters,EmbedderAdapter/MilvusSinkAdapter/TxCacheAdaptertakeself._infra.*(same instances);ArtifactStoreAdaptertakesself._art.*(signature-compatible, parameter-equivalent to the original Engine thin delegates). The two gates,ProducerContext,build_job_lane, theregister_on_succeededorder (_on_object_indexedbeforejob_lane.on_embed_succeeded), andconsumer.start+job_lane.startare all original order and args. Building into locals and publishing last is a type-checker device, runtime-equivalent._pending_finalizethree 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_indexedwith error → pop; no residue)._on_object_indexedatomic body: migrated line-for-line, logic unchanged.pump: the original_index_via_pipelinebody migrates; themessage_streamraw_recordsGC stays infinally, with barepass→logger.exception(GC failure still does not affect the task; observability only).JobLane/PipelineEmbedConsumerinstance (the ones the supervisor constructs in_build_pipeline); identity preserved.Engine.startup/shutdownexternal signatures and semantics are unchanged; call sites inapi/app.pyand__main__.pyare unchanged.Test impact
New unit tests (
tests/test_pipeline_supervisor.py, 11 tests):_chunks_q/embed_consumer/producer_ctxallNone,_pending_finalize == {}.eng.pipelineis non-None;eng.pipeline._obj is eng.objects,_art is eng.artifacts,_factory is eng.connector_factory,_infra is eng.infra.routes_to_pipelinetable:_PIPELINE_OKINDSvalues return True;imagefollowsdescription.enabled;table_schemafollowssummary.enabled; else False.ArtifactStoreAdapterwiring: asserts the adapter holdsartifacts.put_artifact/read_artifact/read_artifact_fresh(bypassing Engine's thin delegates).stash_finalize+ success pop: afterstash_finalize(uri, ctx)the dict has one entry; after_on_object_indexedfires it is popped._on_object_indexed(error=...)writes the failed row +advance_task(FAILED).won==0purges orphan: whenadvance_taskreturns 0,milvus.delete_by_objectruns and no row is written.producer.produceraises,_pending_finalize.pop(full_uri)runs and the exception re-raises.factory.build_pluginreturns aBuiltPlugin;recover_jobreceivesbuilt.plugin.ConnectorJobWatcher(meta, job_lane)+create_task.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 througheng.pipeline.*). No logic changes;ruff formatpasses 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 raisingValueError) - fails on the baseline.tests/test_feishu_oauth.py/tests/test_enumeration_since.py, 2 collection errors - the optional connector SDKlark_oapiis not installed (documented in CLAUDE.md).Risk & rollback
_on_object_indexedatomic body is preserved line-for-line, thestartup/shutdownsequences step-for-step, and handle identity is preserved. The main risk is a missed site among the 13 mechanical rewrites and the_pending_finalizetouches - covered by the full suite (382 passed)._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 indocs-dev/fix-pipeline-supervisor-finalize-race.md, deferred to a standalone follow-up.git revertthe 7 commits fully reverts the change, with no side effects. Rolling back restores the fields + methods toEngineand removes the public attributes andstash_finalize/pump.