ADR-0007: Domain worker and flow structure
Group Engine workers and job contracts by domain, with explicit flow clients and shared connection ownership.
Status: Accepted. Implementation pending.
Date: 2026-09-22
Decision
Group Engine worker apps and shared job definitions by domain. Each worker
remains a separate deployable package. Keep flow definitions and job contracts
in @gigflow/engine-jobs.
Expose the vacancy ingestion workflow through a VacancyClient class in
src/vacancy/client.ts. Its constructor receives an application-owned
JobClient. The flow uses composition and does not extend JobClient.
This ADR records the agreed target structure. It does not move files, change queues, or deploy services.
Context
At this decision, Engine has 25 worker packages directly under
apps/engine/apps/workers/. Package names identify domains, but the directory
tree does not group them. Worker source layouts also differ. Some use a root
processor.ts, some use jobs/, and vacancy embed uses processors/.
The source layout makes related work difficult to find. For example, the
batch-import worker places CSV analysis, dispatch, cancellation, and recovery
beside process setup. Bizzy provisioning uses a general lib/ directory for
both browser login and provisioning lifecycle code.
The original vacancy graph lives in
apps/engine/packages/job-client/src/vacancy-flow.ts. The Engine API,
configuration-harvest worker, and a development script can submit that graph.
Its definition needs a shared owner outside an individual worker app.
apps/engine/packages/jobs already provides a generic JobClient with
enqueue(), enqueueFlow(), and close(). Vacancy embed demonstrates the
class-based processor runtime in @gigflow/worker (then named
@gigflow/worker-new). The older runtime and job contracts remained in use
elsewhere at that time.
Relationship to ADR-0006
This decision refines ADR-0006: Engine job contracts and worker execution. ADR-0006 remains a record of the initial implementation. This ADR takes precedence for the following target design choices:
| Concern | Decision in this ADR |
|---|---|
| Domain layout | Use src/<domain>/ for the flow interface, graph modules, and job schemas. |
| Workflow interface | Use a domain client class, starting with VacancyClient. A shared enqueue.ts does not collect all domain workflows. |
| Public imports | Use explicit subpaths mapped to defining files. Do not extend schema barrels. |
| Queue ownership | One queue per worker is the default. Preserve existing multiple-queue workers when their execution policies require it. |
Retain ADR-0006's generic client, lazy connections, explicit shutdown, processor input/result validation, and separation of domain content from job contracts. The prior schema-barrel exception remains a compatibility detail during migration. It is not the target export structure.
Shared jobs package
Use a domain directory to keep its workflow interface and contracts together. The following tree shows the vacancy target. Add other domains only when their job definitions move into this package.
apps/engine/packages/jobs/
package.json
tsconfig.json
README.md
src/
client.ts
env.ts
types.ts
vacancy/
client.ts
schema.ts
types.ts
flows/
ingestion.ts
schemas/
harvest.ts
preprocess.ts
classify.ts
enrich.ts
geocode.ts
recruiter-match.ts
embed.ts
finalize.tsThe roles are explicit:
| Module | Responsibility |
|---|---|
Root client.ts | Define JobClient. Own Redis connections, queue handles, generic submission, shared queue defaults, and cleanup. |
Root env.ts | Validate producer environment settings. |
Root types.ts | Define generic submission, option, and job-reference types. Re-export domain types. |
vacancy/client.ts | Define VacancyClient. Validate input and submit the ingestion flow or individual jobs. |
vacancy/schema.ts | Define the input contract for starting vacancy ingestion. |
vacancy/types.ts | Define vacancy request, result, and option types. Infer schema-backed types with type-only schema imports. |
vacancy/flows/ingestion.ts | Build the BullMQ graph, including stage payloads, dependency policies, and job identity rules. |
vacancy/schemas/<stage>.ts | Define a stage's job name, queue name, payload schema, and result schema where needed. |
Use singular schema.ts for the workflow entry contract. Use schemas/ for
stage contracts. Do not create a second definition of the same shape. A stage
can reuse the workflow entry schema or another stage's result schema when the
contracts are identical.
Keep generic public types in root src/types.ts. Define vacancy-specific
inferred types in vacancy/types.ts, beside their schemas. This domain-local
type placement is an explicit exception to the root-only public type rule for
the jobs package. Root types.ts uses export * from "./vacancy/types" as an
explicit exception to the barrel rule. Both type import paths remain available
without duplicate definitions. Private graph implementation types can stay in
their owning module. Shared types must not import clients or processors.
The job contracts describe queue messages. Existing domain schemas, such as
vacancy content in @gigflow/engine-schema/vacancy/content, keep their domain
owner. Job schemas can compose them.
Flow interface and connection ownership
The application creates a JobClient and passes it to each domain client.
Each domain client stores that dependency in a private field. Constructors assign
dependencies and do not submit jobs.
The following code defines the target interface. The proposed imports and types become available during implementation.
import type { JobClient } from "../client";
import type { JobReference } from "../types";
import { buildVacancyFlow } from "./flows/ingestion";
import { vacancyFlowInputSchema, vacancyFlowOptionsSchema } from "./schemas";
import type { VacancyFlowInput, VacancyFlowOptions } from "./types";
export class VacancyClient {
readonly #jobs: JobClient;
constructor(jobs: JobClient) {
this.#jobs = jobs;
}
public async enqueueIngestion(
input: VacancyFlowInput,
options?: VacancyFlowOptions,
): Promise<JobReference> {
const payload = vacancyFlowInputSchema.parse(input);
const flowOptions = vacancyFlowOptionsSchema.parse(options ?? {});
const flow = buildVacancyFlow(payload, flowOptions);
return this.#jobs.enqueueFlow({ flow });
}
}VacancyFlowOptions preserves supported workflow options, including the
existing delayMs option. Validate externally supplied options at their entry
point. An ingestion delay applies to the harvest stage, as it does in the
existing graph.
Successful enqueue returns the finalize root's composite job reference. It means that submission succeeded, not that the workflow completed. Callers await submission and use existing status interfaces to observe execution. Validation and submission failures reject the operation.
VacancyClient does not create a hidden JobClient, open its own connections,
or expose close(). The application that creates JobClient closes it.
Several flows in the same process can share that instance.
On shutdown, stop new submissions and drain active work that can still enqueue jobs before closing the producer client. Startup failure cleanup also closes any resources the application created. Do not add a second shutdown framework.
Retain new JobClient() as the generic client's construction contract. This
ADR changes how domain clients receive that client. It does not add dependency
arguments to JobClient itself.
Use VacancyClient for vacancy job submission. enqueueIngestion() submits
the full ingestion flow. Methods such as enqueueHarvest() and
enqueueClassify() submit individual jobs with their domain payload types.
Add named workflow methods only when those workflows exist. Do not add a
forwarding client, a flow registry, or an abstract flow base.
Graph execution and data contracts
The graph builder is a pure function. It receives validated input and options
and returns a BullMQ FlowJob. It does not connect to Redis, query the database,
or call providers. JobClient applies infrastructure defaults and submits the
graph. Domain stage ordering stays in the domain flow module.
Preserve the existing vacancy ingestion sequence:
harvest -> preprocess -> classify -> enrich -> embed -> finalizeBullMQ children run before their parents. The submitted root is finalize, with the earlier stages nested beneath it. The graph must retain that direction. The separate geocode and recruiter-match workers support dynamic children from enrichment. Their existence does not make them additional stages in the static chain. Do not restore the retired claims stage as part of a directory move.
Flow ownership includes the following execution rules:
- Set required dependency failure behavior explicitly. Preserve existing propagation when moving from old submission defaults to the generic client.
- Preserve optional enrichment children and their failure-tolerant behavior.
- Define stable job identity where the workflow requires duplicate prevention. Account for retention and removal. A retained job ID alone does not guarantee that an external effect runs only once.
- Keep retry classification and queue capacity rules in the existing runtime and worker configuration. Do not add a retry loop around flow submission.
- Preserve each processor's input and output contract. Validate processor results before BullMQ records success.
The newer processor runtime merges one object child result with static parent data before input validation. Static parent data takes precedence. Multiple child results require an explicit input mapper. Do not assume BullMQ itself copies child results into parent payloads. Verify dynamic enrichment mappings when migrating that worker to the newer runtime.
Keep graph-building modules internal to the jobs package. Expose one only when a real external composition use requires it.
Public exports
Map each public subpath directly to the file that defines its declarations. The following manifest excerpt shows the pattern. Add an explicit schema entry for every migrated stage that a consumer needs.
{
"exports": {
"./client": "./src/client.ts",
"./types": "./src/types.ts",
"./vacancy/client": "./src/vacancy/client.ts",
"./vacancy/schema": "./src/vacancy/schema.ts",
"./vacancy/types": "./src/vacancy/types.ts",
"./vacancy/schemas/embed": "./src/vacancy/schemas/embed.ts",
"./vacancy/schemas/finalize": "./src/vacancy/schemas/finalize.ts"
}
}Consumers import VacancyClient from
@gigflow/engine-jobs/vacancy/client and vacancy types from
@gigflow/engine-jobs/vacancy/types. Generic types and the domain type re-exports
remain available from @gigflow/engine-jobs/types. Processors import their stage contracts through
the corresponding explicit schema subpath. Within the package, use direct
relative imports.
Do not add a package-root export, wildcard subpath, or another barrel. The root
type re-export above is the explicit exception. Do not remove
existing /enqueue or /schemas exports until their callers migrate. Preserve
compatibility without extending their domain surface or duplicating contracts.
Schema-only and type-only imports must not require environment settings or
create connections.
Worker apps
Use apps/engine/apps/workers/<domain>/<worker-name>/ for deployable workers.
Domain directories are organizational folders. They have no package manifest,
runtime, or shared process of their own.
The existing workers map to the following domain groups:
| Domain folder | Worker folders |
|---|---|
company | import, persist, embed |
legal-entity | import, import-batch, persist |
vacancy | harvest, preprocess, classify, enrich, geocode, recruiter-match, embed, finalize, claims |
configuration | discovery, harvest, noise-filter |
linkedin | profile-search, profile-fetch |
bizzy | account-provisioning, session, health |
storage | logo-mirror, provider-raw-snapshot-archive |
This mapping preserves the existing worker inventory. It does not authorize
removing a worker or adding its queue back to a flow. Worker folder names omit
the domain prefix. Storage worker names remain logo-mirror and
provider-raw-snapshot-archive. Folder renames preserve source contents and
operational identities. Changes to worker internals require separate work.
A simple worker uses the following structure:
apps/engine/apps/workers/vacancy/harvest/
package.json
tsconfig.json
src/
index.ts
env.ts
queue.ts
processors/
harvest-vacancy.tsindex.ts composes dependencies, starts the worker, and installs lifecycle
handlers. env.ts validates environment settings. queue.ts registers
processors and owns concurrency, limits, and queue options.
processors/harvest-vacancy.ts defines HarvestVacancyProcessor.
Use BaseProcessor from the selected shared runtime. Directory changes do not
require all workers to switch runtime in one release. Existing function-based
jobs can retain their execution contract during the structural migration.
For larger workers, group operation-specific code with its processor. For example, the batch-import target can use:
legal-entity/import-batch/src/
index.ts
env.ts
schedulers.ts
queues/
analysis.ts
control.ts
processors/
analyze/
analyze-import.ts
analysis-persistence.ts
csv-analyzer.ts
csv-parser-stream.ts
csv-stream.ts
invalid-report.ts
dispatch-import.ts
cancel-import.ts
sweep/
sweep-imports.ts
job-recovery.tsPreserve its separate analysis and control queues. Their concurrency, locks, retention, and recovery rules serve different operations. A single-queue convention must not remove these rules.
Use specific app folders such as browser/ and provisioning/ for the
corresponding Bizzy provisioning code. Keep existing tests beside the files
they cover. Add nesting only when related files need a common owner. Do not
create an empty folder for every possible role.
Ownership across packages and apps
The two domain trees serve different consumers. Producers need workflow interfaces and contracts. Deployed workers need execution code.
| Location | Owns |
|---|---|
apps/engine/packages/jobs/src/vacancy/ | Submission interface, flow graphs, and shared vacancy job contracts. |
apps/engine/apps/workers/vacancy/ | Process setup, job execution, progress, cancellation, and queue-specific behavior. |
| Existing Engine domain packages | Provider access, domain calculations, and reusable transformations. |
apps/engine/packages/db/ | Database queries, transactions, schema, and persistence rules. |
| Shared worker runtime | Job lifecycle, validation hooks, retry infrastructure, health, and metrics. |
Worker apps can depend on jobs and domain packages. The jobs package must not import worker apps. Keep domain capabilities independent of BullMQ job objects. Do not create a new domain package only to reduce a worker's file count.
The app layout in this ADR does not replace the required package layout for
domain clients, providers, or format processors. Those packages retain root
client.ts, types.ts, and the prescribed role-specific directories.
Deployment and local discovery
Keep each worker's operational identity independent of its source path.
apps/engine/deploy/values.yaml continues to list workers by their existing
names. The Helm template creates a Deployment and KEDA ScaledObject from each
entry. Source grouping does not combine deployments or scaling policies.
For example, preserve the following mapping:
| Item | Value |
|---|---|
| Source directory | apps/engine/apps/workers/company/import |
| Package name | @gigflow/engine-workers-company-import |
| Values entry | company-import |
| Image name | engine-worker-company-import |
| Default queue name | company-import |
| Local service ID | engine.worker.company-import |
Preserve queue names, job names, Redis prefixes, image names, resource names,
health endpoints, environment keys, and per-worker configuration. Existing
queues overrides remain valid, including the two batch-import queues.
The directory move requires coordinated tooling changes:
| Owner | Required change |
|---|---|
Root package.json and Bun lockfile | Discover nested worker packages. Refresh workspace paths through Bun. Retain old patterns while other areas or unmigrated workers need them. |
.github/workflows/images.yml | Discover worker manifests and carry separate package, worker identity, and source-path fields through the build. |
| Per-worker image build | Pass the discovered path as APP_DIR to packages/worker/Dockerfile. Keep existing image names. |
| CI values updates | Match .workers.list[].name with the stable worker identity when updating image tags. |
apps/engine/dev/profile.py and CLI discovery | Discover nested workers while retaining service IDs and environment-consumer identities. |
| CLI workspace graph and environment handling | Include nested manifests and generated environment paths. Preserve manual overrides and worktree isolation. |
| Worker imports and compiler settings | Update relative imports and path aliases affected by the extra directory level. |
Derive the worker identity from the declared package name by removing the
verified @gigflow/engine-workers- prefix. Validate uniqueness. Do not derive
identity from the leaf folder. Later folder-name changes must not change
deployment identities or local port assignments.
The Dockerfile already accepts an APP_DIR argument and runs src/index.ts
from that directory. The shared base build must still include the complete
worker package set through Turbo pruning.
Local discovery also assigns port slots from worker order. Preserve existing
assignments during the move and use stable identities for ordering. Verify
environment generation and stale-path handling through the supported CLI.
Do not hand-edit generated .env files or remove worktree resources.
Migration
Implement the decision in bounded changes. Keep structural changes separate from changes to payloads or execution behavior where possible.
- Record the worker inventory, package names, queue names, callers, and deployment identities. Capture the existing flow behavior and contracts.
- Add the domain layout to
@gigflow/engine-jobs. IntroduceVacancyClient, explicit exports, and schema-inferred public types. Preserve supported delay behavior and the returned root reference. - Reconcile the old and newer stage contracts before switching producers. The original flow and the newer vacancy embed example are not evidence that every stage is already compatible. Migrate schemas, processors, and callers as a coherent execution path.
- Move API, worker, and script callers to the flow class. Create shared producer clients at application composition points and close them through the existing lifecycle. Retain old entry points until their callers migrate.
- Standardize worker internals. Preserve existing tests, queue options, schedules, recovery, and runtime behavior. Migrate processor classes in separately verifiable changes where a runtime change is needed.
- Update package discovery, image builds, local tooling, and path-dependent settings together with worker directory moves. Do not land a move that causes discovery to omit a worker.
- Remove unused legacy definitions and exports only after consumers migrate. Update READMEs, examples, and deployment comments to match the final paths.
No payload compatibility break is implied by this ADR. If implementation needs one, define how existing jobs finish before activating incompatible producers and consumers. Do not purge queues as a directory-migration step.
Validation
The ADR-only change requires documentation checks. Runtime validation below is required when implementing the migration, not evidence of a completed run.
Check changed files with the installed formatter, verify links and paths, and confirm the handbook navigation entry. For implementation, run focused package checks, existing relevant suites, and the required root checks:
bun run --filter @gigflow/engine-jobs lint
bun run --filter @gigflow/engine-jobs typecheck
bun run --filter @gigflow/engine-jobs test
bun run lint
bun run typecheck
dev/cli/dev plan engineAlso check affected worker and producer packages through their actual scripts. Run existing CLI discovery, selection, workspace graph, and environment tests when changing those modules. Build the shared base and representative worker images. Render the Helm chart before and after the directory move and compare identities, queue triggers, resource settings, and environment configuration. Use the repository's established tools and test setup.
Before implementation, record failure cases and expected behavior. Verify the complete workflow through real entry points with the Engine test environment and known fixtures. Cover at least:
- Valid submission, stage order, result handoff, and final persisted outcome.
- Rejected input without enqueuing partial work.
- Required-stage failure and optional enrichment failure.
- Retry and duplicate-submission behavior without repeated domain effects.
- Delayed harvesting and correct root-reference reporting.
- API and worker producers sharing the same workflow contract.
- Shutdown with pending submissions and active workers.
- Discovery of all migrated workers, including repeated leaf names.
- Per-worker image paths, startup, readiness, and unchanged deployment names.
A graph test with mocked processors is not full workflow E2E coverage. Prefer the existing E2E locations and runners. Do not add unit tests after writing implementation code. Preserve and run existing suites.
Every E2E run must produce a repeatable artifact with the tested revision, exact command, setup, fixture or seed, expected and actual results, status, and relevant reports or logs. Exclude secrets. Report blocked or incomplete checks accurately.
Alternatives and consequences
Domain grouping improves navigation while keeping worker deployment and workflow submission separate. It adds one directory level and requires a coordinated update to discovery tools. Keeping both jobs and workers grouped by domain also requires clear ownership rules to prevent duplicated contracts.
The alternatives considered are:
| Alternative | Reason for this decision |
|---|---|
| Keep every worker at one level | Names encode the domain, but related workers remain scattered through one list. |
| Put graphs inside a worker app | API and other worker producers need the same graph without importing app code. |
Add domain methods to generic JobClient | Each workflow would change the infrastructure client's interface. |
Construct JobClient inside each flow | Connection and shutdown ownership becomes distributed across flow instances. |
Make VacancyClient extend JobClient | The flow would expose unrelated queue operations and inherit connection ownership. |
Add a forwarding class around VacancyClient | The domain client already exposes workflow and individual job submission. |
| Put all vacancy contracts in one file | Stage contracts can grow independently and need direct worker imports. |
| Combine domain workers into one process | This would change independent scaling, resource limits, and deployment behavior. |
Midday's BullMQ source separates processors by domain and has domain-specific queue configuration and schema files. This ADR adapts that organization to Engine's independently deployed workers and explicit package exports. It does not require Engine to copy Midday's deployment or legacy barrel structure. See the Midday worker source for the reference inspected on the decision date.