Repository navigation
refactor(engine): extract IngestOrchestrator -- end-to-end sync-job execution + ObjectIndexer four exits - #175
Merged
Conversation
added 2 commits
July 15, 2026 11:27
Extract job execution logic from `Engine` into a new `IngestOrchestrator` component. This orchestrator now manages the end-to-end lifecycle of ingestion jobs, including registration, draining, map-phase processing, indexing, and finalization. - Add `IngestOrchestrator` in `server/python/src/mfs_server/engine/ingest.py` - Refactor `Engine` to delegate job loop and connector registration to the orchestrator - Update existing engine tests to use the new orchestrator interface - Add new orchestrator tests in `test_ingest_orchestrator.py` - Update `.gitignore` to include local control-plane state and docs-dev
…te prefixes Renames several internal methods in `IngestOrchestrator` from private to public to allow direct access from the `Engine` and test suites. This aligns with the recent extraction of orchestration logic from the `Engine` class. - Rename `_open_sync_job` to `open_sync_job` - Rename `_drain_job` to `drain_job` - Rename `_run_job` to `run_job` - Rename `_finalize_job` to `finalize_job` - Rename `_run_job_loop` to `run_job_loop` - Update all call sites in `Engine`, `IngestOrchestrator`, and the E2E test suite to use the new public API
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
resole #166
Overview
Consolidates all of
Engine's "execute one sync ingest job end-to-end" logic into a newIngestOrchestrator: connector registration -> open job -> enumerate -> map-phase task processing (claim/retry/circuit-breaker/heartbeat) -> single-object write (_index_object) -> finalize -> cancel._index_object's four exits (deleted / renamed / pipeline 'deferred' / metadata-only) are made explicit asObjectIndexer+ 4 Command handlers.Engineshrinks to a thin Facade forwardingadd/cancel_job.Follows the prior
InfraStack(#164) andPipelineSupervisor(#171). Zero change to external behavior.Component boundary adjustment (important)
The engine refactor's original plan split "execute one sync job" across two components:
IngestOrchestratorowns the job skeleton (drain_job), and a not-yet-extractedWorkerSchedulerowns map-phase execution (run_job/run_job_loop/process_with_retry/classify_error/heartbeat_loop/should_stop). Butdrain_jobcallsrun_job, andrun_job's task processing calls back into_index_object-- anIngestOrchestrator<->WorkerSchedulercycle the plan bridged with abind_workerback-reference (Engine acting as its own worker, a circular reference).This PR instead groups end-to-end: map-phase execution is folded into
IngestOrchestrator, sodrain_job->run_job->process_with_retry->ObjectIndexer.handleare all in-component internal calls -- the cycle is eliminated at the root andbind_workeris unnecessary. The futureWorkerSchedulershrinks to a pure queue/concurrency/reclaim layer that calls back intoself.ingest.run_job+finalize_job.Rationale:
run_job/process_with_retry/heartbeat_loopare tightly bound to_index_object/drain_job(retry wraps indexing, the breaker protects indexing, the heartbeat guards the job) -- not generic scheduling reusable by a different worker; grouping them is both acyclic and semantically more correct (the breakerconsec_failcounts object-indexing failures, an Orchestrator concern). The cost is work moving into this PR; the futureWorkerSchedulerextraction will be smaller.docs-dev/engine-redesign.mdanddocs-dev/issue-engine-redesign.md(+.en.md) have been aligned to this boundary adjustment.Why needed
After #171,
Enginestill carried ~21 ingest write-orchestration methods:add/register_or_get_connector/drain_job/run_job/run_job_loop/process_with_retry/_index_object(the most complex write-path method, 4 branches) /finalize_job/cancel_job/ heartbeat / circuit breaker / error classification. Three concrete problems:_index_object's four exits fused in one method: deleted return / renamed return / pipelinereturn "deferred"/ inline tail -- mixed with milvus side effects,stash_finalize,on_object_indexedin a 110-line method; no exit is independently verifiable.drain_job(Orchestrator) ->run_job(WorkerScheduler) ->process_with_retry->_index_object(Orchestrator) -- a cyclic call chain the original design bridged withbind_worker, introducing an Engine<->Orchestrator circular reference in the intermediate state._index_object's four exits and ObjectIndexer routing can't be unit-tested without running the whole Engine.After extraction: each exit has a dedicated test; Orchestrator has no circular reference (all forward deps);
_index_object's routing (incl. the renamed-empty fallthrough) is explicit and readable.Changes
New
engine/ingest.py(841 lines)IngestOrchestrator+ObjectIndexer+IndexHandlerProtocol +IndexContext+ 4 handlers.IngestOrchestrator(cfg, infra, factory, pipeline, objects, artifacts)-- all forward deps, no Engine back-reference.bind_remover(remove_connector)back-fills theadd-failure rollback callable (Engine provides it now; the futureConnectorManagerwill swap inConnectorManager.remove).ObjectIndexer.handle(plugin, connector_uri, task): routes_index_object's four exits, control flow line-by-line equivalent (incl. the renamed-empty fallthrough).IndexContextcomputesstat/okind/indexableonce on the fallthrough path; handlers don't re-stat._index_objectbranches):DeletedHandler->delete_object_row+milvus.delete_by_object+drop_artifacts+on_object_deleted; returnsNoneRenameHandler-> old_chunks non-empty: chunk_id rewrite +milvus.delete/upsert+rename_artifacts+ delete old row + write new rowindexed+on_object_indexed, returnsNone; empty: delete old milvus/row, returns"continue"to fall throughPipelineIndexHandler->stash_finalize(full_uri, (cid, connector_uri, relpath, st, indexable, plugin, task_id))+pump, returns"deferred"MetadataOnlyHandler->milvus.delete_by_object+write_object_row("not_indexed", 0)+on_object_indexed, returnsNoneself.X->self._X):add/register_or_get_connector/open_sync_job/drain_job/finalize_job/cancel_job/run_job/run_job_loop/process_with_retry/await_map_drained/heartbeat_loop/should_stop/claim_batch/classify_error(staticmethod) /warn_object_failed(staticmethod)._PER_OBJECT_SKIP_ERRORS/_HEARTBEAT_INTERVAL_S/_normalize_jsonmove with their sole users.Enginechanges (engine/engine.py, 2115 -> 1504)_PER_OBJECT_SKIP_ERRORS/_HEARTBEAT_INTERVAL_S/_normalize_json.__init__: afterself.pipeline, addself.ingest = IngestOrchestrator(cfg, self.infra, self.connector_factory, self.pipeline, self.objects, self.artifacts)+self.ingest.bind_remover(self.remove_connector).add/cancel_job: collapsed to thin forwardsreturn await self.ingest.add(...)/self.ingest.cancel_job(...).self.ingest.*:ingest_upload/files_upload(open_sync_job/drain_job),_staging_connector(register_or_get_connector),run_worker_once(run_job/finalize_job).get_plugin_cls/chunk_id/time/suppress/TaskStatus(moved with their methods); keepSyncOptions(still used byestimate); addfrom .ingest import IngestOrchestrator.Shared thin delegates stay on Engine
_resolve_target/_build_plugin/_resolve_ref/_is_secret_key/_redact_config/_drop_artifacts/_write_object_roware not moved -- they are 2-line delegates shared byprobe/estimate/inspect/remove_connector/ upload / worker.IngestOrchestratorcalls the underlying components directly (self._factory.resolve_target/build_plugin,self._art.drop_artifacts,self._obj.write_object_row), duplicating no delegate (per #171 precedent: internal collaborators wire straight to the underlying component; Engine keeps the delegate for its own paths).Tests
tests/test_ingest_orchestrator.py(9 unit tests): Engine forwarding +bind_removerback-fill;ObjectIndexer's four exits (deleted / renamed-reuse / renamed-empty-fallthrough / pipelinedeferred/ metadata-only [binary + opted-out + non-pipeline]).eng._run_job_loop(...)->eng.ingest.run_job_loop(...)(16),eng.register_or_get_connector(...)->eng.ingest.register_or_get_connector(...)(7). Publiceng.add/eng.cancel_jobunchanged;eng._write_object_row/eng._drop_artifactsstay on Engine (no test change).File manifest
engine/ingest.pyengine/engine.py__init__/add/cancel_job, 5 call sites (-611)tests/test_ingest_orchestrator.pytests/test_engine_*.pyeng._<moved>->eng.ingest.<moved>(23)docs-dev/engine-redesign-ingest-orchestrator.mddocs-dev/engine-redesign.md/issue-engine-redesign.md(+.en.md)Design decisions
End-to-end grouping removes the cycle. The original plan split
drain_job(job skeleton) andrun_job(map-phase execution) across two components; their mutual callbacks would have formed a cycle, bridged withbind_worker. This PR folds map-phase execution into Orchestrator sodrain_job->run_jobis an internal call and the cycle vanishes; the futureWorkerScheduler(not yet extracted) becomes a queue layer callingself.ingest.run_job+finalize_job. Cost: work moves into this PR and the future WorkerScheduler extraction shrinks; both steps stay behavior-equivalent.bind_remover, notbind_worker. The only backward touch isadd's failure rollback callingremove_connector(owned by the not-yet-extractedConnectorManager). Injecting a single callablebind_remover(self.remove_connector)is precise, narrow, and temporary -- not an "Engine is its own worker" circular reference. That future extraction swaps inConnectorManager.remove; the Orchestrator call site is unchanged._index_objectsplit into 4 handlers, control flow line-by-line equivalent.ObjectIndexer.handlepreserves the exact control flow: deleted returns first; renamed-with-old_uri goes toRenameHandler(reuse returnsNone, empty returns"continue"to fall through to the preamble); the preamble computesstat/okind/indexableonce;indexable and routes_to_pipeline->PipelineIndexHandlerreturns"deferred", elseMetadataOnlyHandler. The originalchunk_count==0/search_status="not_indexed"are constants at the inline tail;MetadataOnlyHandlerhardcodes"not_indexed"/0, equivalent.Shared thin delegates stay on Engine; Orchestrator wires direct.
_resolve_target/_build_plugin/_drop_artifacts/_write_object_roware shared by probe/estimate/inspect/remove/upload/worker, so they stay on Engine; Orchestrator usesself._factory.*/self._art.*/self._obj.*directly, duplicating no delegate (refactor(engine): extract PipelineSupervisor - process-singleton assembly + GC/recover + finalize hook #171 precedent)._write_object_rowno longer has a production caller on Engine (only_index_objectused it, now migrated) but is kept for test compatibility (test_engine_orphan_gccallseng._write_object_rowdirectly).finalize context stays a bare tuple.
stash_finalize's(cid, connector_uri, relpath, st, indexable, plugin, task_id)matchesPipelineSupervisor._pending_finalize: dict[str, tuple]field-for-field; not upgraded to a dataclass, to avoid mismatching the supervisor's tuple contract (a future incremental improvement).Location
engine/ingest.py(alongsideengine/pipeline_supervisor.py,engine/infra.py), per refactor(engine): extract PipelineSupervisor - process-singleton assembly + GC/recover + finalize hook #171 precedent: Engine-internal collaborators (no external import) live atengine/top level, notcomponents/(reserved for published repository/factory components).CircuitBreaker/BackoffPolicy/ErrorClassifiervalue objects stay inline this stage. Breaker state (consec_fail) is inrun_job_loop, classification inclassify_error, backoff inprocess_with_retry-- all in Orchestrator with map-phase execution. Value-object extraction is a follow-up (behavior equivalence first).Behavior equivalence
IngestOrchestrator's methods are reproduced in the exact same order as the originalEngine:self.objects->self._objetc.);_index_object->ObjectIndexer.handle's four-branch control flow is line-by-line equivalent.drain_job->run_job->process_with_retry->_indexer.handleare all internal; Engine'srun_worker_oncecallsself.ingest.*one way._resolve_targetunpacking order preserved:add's_, connector_uri, ctype, default_config = r.ctype, r.connector_uri, r.scheme, r.configmatches the original_resolve_targetreturn(ctype, connector_uri, scheme, config)position-for-position (including the existing semantics where thectypevariable holdsscheme).add's rollbackself._remove_connector(connector_uri)--bind_remover(self.remove_connector)injects the sameEngine.remove_connectorbound method; reference-equal.heartbeat_loopordering: the enumeration-phase task is cancelled infinallybeforerun_jobstarts its own heartbeat -- no double heartbeat (as before).Engine.add/Engine.cancel_jobsignatures unchanged;api/app.pyand__main__.pycall sites unchanged.Test impact
New unit tests (
tests/test_ingest_orchestrator.py, 9):eng.ingestnon-null;_obj/_pipeline/_factory/_art/_infrareference-equal;bind_removerinjectedeng.remove_connector.bind_remover:Nonebefore, reference-equal after.3-9.
ObjectIndexerfour exits: deleted (delete row + delete chunk + drop artifacts + on_deleted), renamed-reuse (chunk_id rewrite + upsert + rename + writeindexed+ on_indexed, returns None), renamed-empty fallthrough (clean old -> metadata-only), pipelinedeferred(stash tuple field-for-field assert + pump), metadata-only [binary / opted-out / non-pipeline] (delete chunk + writenot_indexed/0 + on_indexed).Existing tests mechanically rewritten: 10 files, 23 touchpoints (
eng.run_job_loop->eng.ingest.run_job_loop,eng.register_or_get_connector->eng.ingest.register_or_get_connector). No logic change; ruff format clean.Regression: full
uv run --extra dev pytest-- 417 passed / 9 skipped (baseline 408/9 + 9 new unit tests, 0 regressions).Risk & rollback
_index_objectinto 4 handlers is "structural" -- the four-exit control flow is verified line-by-line equivalent and covered per-branch by the new unit tests. Main risks are the 23 test touchpoint rewrites and the renamed-empty fallthrough semantics -- covered by the full suite (417 passed).WorkerScheduler(not yet extracted) a queue layer;run_jobetc. no longer belong to it. Recorded indocs-dev/engine-redesign.md/ issue files; this PR states it explicitly.bind_removertemporality: injectsEngine.remove_connectoras a single callable (not a circular reference); the futureConnectorManagerswaps inConnectorManager.remove; the call site is unchanged.git revertthis PR's commit; restores the 14 methods + 3 module symbols toEngine, deletesself.ingest/bind_remover/ObjectIndexer/ 4 handlers.