Skip to content

refactor(engine): extract IngestOrchestrator -- end-to-end sync-job execution + ObjectIndexer four exits - #175

Merged
zc277584121 merged 2 commits into
zilliztech:mainfrom
code2tan:rebuild_ingest
Jul 16, 2026
Merged

zc277584121 merged 2 commits into
zilliztech:mainfrom
code2tan:rebuild_ingest

Conversation

@code2tan

Copy link
Copy Markdown
Contributor

resole #166

Overview

Consolidates all of Engine's "execute one sync ingest job end-to-end" logic into a new IngestOrchestrator: 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 as ObjectIndexer + 4 Command handlers. Engine shrinks to a thin Facade forwarding add / cancel_job.

Follows the prior InfraStack (#164) and PipelineSupervisor (#171). Zero change to external behavior.

Component boundary adjustment (important)

The engine refactor's original plan split "execute one sync job" across two components: IngestOrchestrator owns the job skeleton (drain_job), and a not-yet-extracted WorkerScheduler owns map-phase execution (run_job / run_job_loop / process_with_retry / classify_error / heartbeat_loop / should_stop). But drain_job calls run_job, and run_job's task processing calls back into _index_object -- an IngestOrchestrator <-> WorkerScheduler cycle the plan bridged with a bind_worker back-reference (Engine acting as its own worker, a circular reference).

This PR instead groups end-to-end: map-phase execution is folded into IngestOrchestrator, so drain_job -> run_job -> process_with_retry -> ObjectIndexer.handle are all in-component internal calls -- the cycle is eliminated at the root and bind_worker is unnecessary. The future WorkerScheduler shrinks to a pure queue/concurrency/reclaim layer that calls back into self.ingest.run_job + finalize_job.

Rationale: run_job / process_with_retry / heartbeat_loop are 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 breaker consec_fail counts object-indexing failures, an Orchestrator concern). The cost is work moving into this PR; the future WorkerScheduler extraction will be smaller.

docs-dev/engine-redesign.md and docs-dev/issue-engine-redesign.md (+.en.md) have been aligned to this boundary adjustment.


Why needed

After #171, Engine still 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:

  1. _index_object's four exits fused in one method: deleted return / renamed return / pipeline return "deferred" / inline tail -- mixed with milvus side effects, stash_finalize, on_object_indexed in a 110-line method; no exit is independently verifiable.
  2. Map-phase execution coupled across "future components": drain_job (Orchestrator) -> run_job (WorkerScheduler) -> process_with_retry -> _index_object (Orchestrator) -- a cyclic call chain the original design bridged with bind_worker, introducing an Engine<->Orchestrator circular reference in the intermediate state.
  3. Tests can only go E2E: _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 + IndexHandler Protocol + IndexContext + 4 handlers.

  • Constructor: IngestOrchestrator(cfg, infra, factory, pipeline, objects, artifacts) -- all forward deps, no Engine back-reference. bind_remover(remove_connector) back-fills the add-failure rollback callable (Engine provides it now; the future ConnectorManager will swap in ConnectorManager.remove).
  • ObjectIndexer.handle(plugin, connector_uri, task): routes _index_object's four exits, control flow line-by-line equivalent (incl. the renamed-empty fallthrough). IndexContext computes stat/okind/indexable once on the fallthrough path; handlers don't re-stat.
  • 4 handlers (bodies taken line-by-line from the original _index_object branches):
    • DeletedHandler -> delete_object_row + milvus.delete_by_object + drop_artifacts + on_object_deleted; returns None
    • RenameHandler -> old_chunks non-empty: chunk_id rewrite + milvus.delete/upsert + rename_artifacts + delete old row + write new row indexed + on_object_indexed, returns None; empty: delete old milvus/row, returns "continue" to fall through
    • PipelineIndexHandler -> 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, returns None
  • Migrated methods (verbatim, only field renames self.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).
  • Module constants: _PER_OBJECT_SKIP_ERRORS / _HEARTBEAT_INTERVAL_S / _normalize_json move with their sole users.

Engine changes (engine/engine.py, 2115 -> 1504)

  • Delete 14 migrated methods + _PER_OBJECT_SKIP_ERRORS / _HEARTBEAT_INTERVAL_S / _normalize_json.
  • __init__: after self.pipeline, add self.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 forwards return await self.ingest.add(...) / self.ingest.cancel_job(...).
  • 5 retained-method call sites rewritten to self.ingest.*: ingest_upload / files_upload (open_sync_job / drain_job), _staging_connector (register_or_get_connector), run_worker_once (run_job / finalize_job).
  • import cleanup: remove get_plugin_cls / chunk_id / time / suppress / TaskStatus (moved with their methods); keep SyncOptions (still used by estimate); add from .ingest import IngestOrchestrator.

Shared thin delegates stay on Engine

_resolve_target / _build_plugin / _resolve_ref / _is_secret_key / _redact_config / _drop_artifacts / _write_object_row are not moved -- they are 2-line delegates shared by probe / estimate / inspect / remove_connector / upload / worker. IngestOrchestrator calls 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

  • New tests/test_ingest_orchestrator.py (9 unit tests): Engine forwarding + bind_remover back-fill; ObjectIndexer's four exits (deleted / renamed-reuse / renamed-empty-fallthrough / pipeline deferred / metadata-only [binary + opted-out + non-pipeline]).
  • Mechanical rewrite of 10 test files, 23 touchpoints: eng._run_job_loop(...) -> eng.ingest.run_job_loop(...) (16), eng.register_or_get_connector(...) -> eng.ingest.register_or_get_connector(...) (7). Public eng.add/eng.cancel_job unchanged; eng._write_object_row/eng._drop_artifacts stay on Engine (no test change).

File manifest

File Change
engine/ingest.py new (+841)
engine/engine.py delete 14 methods + 3 module symbols, rewrite __init__/add/cancel_job, 5 call sites (-611)
tests/test_ingest_orchestrator.py new (+307, 9 tests)
10 tests/test_engine_*.py touchpoint eng._<moved> -> eng.ingest.<moved> (23)
docs-dev/engine-redesign-ingest-orchestrator.md detailed design (existing)
docs-dev/engine-redesign.md / issue-engine-redesign.md(+.en.md) boundary-adjustment alignment

Design decisions

  • End-to-end grouping removes the cycle. The original plan split drain_job (job skeleton) and run_job (map-phase execution) across two components; their mutual callbacks would have formed a cycle, bridged with bind_worker. This PR folds map-phase execution into Orchestrator so drain_job -> run_job is an internal call and the cycle vanishes; the future WorkerScheduler (not yet extracted) becomes a queue layer calling self.ingest.run_job + finalize_job. Cost: work moves into this PR and the future WorkerScheduler extraction shrinks; both steps stay behavior-equivalent.

  • bind_remover, not bind_worker. The only backward touch is add's failure rollback calling remove_connector (owned by the not-yet-extracted ConnectorManager). Injecting a single callable bind_remover(self.remove_connector) is precise, narrow, and temporary -- not an "Engine is its own worker" circular reference. That future extraction swaps in ConnectorManager.remove; the Orchestrator call site is unchanged.

  • _index_object split into 4 handlers, control flow line-by-line equivalent. ObjectIndexer.handle preserves the exact control flow: deleted returns first; renamed-with-old_uri goes to RenameHandler (reuse returns None, empty returns "continue" to fall through to the preamble); the preamble computes stat/okind/indexable once; indexable and routes_to_pipeline -> PipelineIndexHandler returns "deferred", else MetadataOnlyHandler. The original chunk_count==0 / search_status="not_indexed" are constants at the inline tail; MetadataOnlyHandler hardcodes "not_indexed"/0, equivalent.

  • Shared thin delegates stay on Engine; Orchestrator wires direct. _resolve_target / _build_plugin / _drop_artifacts / _write_object_row are shared by probe/estimate/inspect/remove/upload/worker, so they stay on Engine; Orchestrator uses self._factory.* / self._art.* / self._obj.* directly, duplicating no delegate (refactor(engine): extract PipelineSupervisor - process-singleton assembly + GC/recover + finalize hook #171 precedent). _write_object_row no longer has a production caller on Engine (only _index_object used it, now migrated) but is kept for test compatibility (test_engine_orphan_gc calls eng._write_object_row directly).

  • finalize context stays a bare tuple. stash_finalize's (cid, connector_uri, relpath, st, indexable, plugin, task_id) matches PipelineSupervisor._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 (alongside engine/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 at engine/ top level, not components/ (reserved for published repository/factory components).

  • CircuitBreaker / BackoffPolicy / ErrorClassifier value objects stay inline this stage. Breaker state (consec_fail) is in run_job_loop, classification in classify_error, backoff in process_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 original Engine:

  • Verbatim migration: the 14 migrated method bodies change no logic, only field renames (self.objects->self._obj etc.); _index_object -> ObjectIndexer.handle's four-branch control flow is line-by-line equivalent.
  • No circular reference: map-phase execution and the write primitive are in one class; drain_job->run_job->process_with_retry->_indexer.handle are all internal; Engine's run_worker_once calls self.ingest.* one way.
  • _resolve_target unpacking order preserved: add's _, connector_uri, ctype, default_config = r.ctype, r.connector_uri, r.scheme, r.config matches the original _resolve_target return (ctype, connector_uri, scheme, config) position-for-position (including the existing semantics where the ctype variable holds scheme).
  • Single backward touch: add's rollback self._remove_connector(connector_uri) -- bind_remover(self.remove_connector) injects the same Engine.remove_connector bound method; reference-equal.
  • heartbeat_loop ordering: the enumeration-phase task is cancelled in finally before run_job starts its own heartbeat -- no double heartbeat (as before).
  • External signatures: Engine.add / Engine.cancel_job signatures unchanged; api/app.py and __main__.py call sites unchanged.

Test impact

New unit tests (tests/test_ingest_orchestrator.py, 9):

  1. Engine forwarding: eng.ingest non-null; _obj/_pipeline/_factory/_art/_infra reference-equal; bind_remover injected eng.remove_connector.
  2. bind_remover: None before, reference-equal after.
    3-9. ObjectIndexer four exits: deleted (delete row + delete chunk + drop artifacts + on_deleted), renamed-reuse (chunk_id rewrite + upsert + rename + write indexed + on_indexed, returns None), renamed-empty fallthrough (clean old -> metadata-only), pipeline deferred (stash tuple field-for-field assert + pump), metadata-only [binary / opted-out / non-pipeline] (delete chunk + write not_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

  • Risk (low-medium): pure structural extraction, no schema/protocol/SQL change. Splitting _index_object into 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).
  • Deviation from the original plan: end-to-end grouping makes the future WorkerScheduler (not yet extracted) a queue layer; run_job etc. no longer belong to it. Recorded in docs-dev/engine-redesign.md / issue files; this PR states it explicitly.
  • bind_remover temporality: injects Engine.remove_connector as a single callable (not a circular reference); the future ConnectorManager swaps in ConnectorManager.remove; the call site is unchanged.
  • Rollback: git revert this PR's commit; restores the 14 methods + 3 module symbols to Engine, deletes self.ingest / bind_remover / ObjectIndexer / 4 handlers.

code2t 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
@zc277584121
zc277584121 merged commit a9ed790 into zilliztech:main Jul 16, 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.

2 participants