How the Agent Graph System Schedules Dependent Work in Apache Maka
The Apache Maka Agent Graph scheduler binds dependent work through the input_ids field and defers execution until all upstream results are durably committed, using a server-side reconciliation loop in stream-graph-schedule-reconcile.ts to validate dependencies before dispatch.
The Agent Graph system in the apache/maka repository provides a durable, server-side scheduler that transforms high-level user requests into concrete operator executions while maintaining strict dependency ordering. When a work item depends on previous graph results, the scheduler ensures those inputs are resolved and committed before provisioning the downstream operator.
Declaring Dependencies via input_ids
The User-Level API Contract
A supervisor agent instructs the graph to add dependent work by submitting an update_agent_graph request with an add_work payload. The critical field for dependency binding is input_ids, an array containing the resultRecordId of every upstream result the new work requires.
{
"operation": "add_work",
"add_work": [
{
"target_kind": "new_agent",
"agent_id": "local-read",
"instruction": "Summarise the article",
"input_ids": ["<upstream-result-record-id>"],
"replacement_mode": "none"
}
]
}
Each string in input_ids represents a durable reference to a specific output record. The runtime-side prompt, defined in packages/runtime/src/graph-mode.ts at line 36, explicitly instructs agents: "When scheduling dependent work, pass each upstream result.resultRecordId in input_ids. The Runtime resolves those references into bounded, source-linked operator handoffs."
Runtime Resolution Guarantees
The scheduler does not execute work immediately upon receipt. Instead, it treats the input_ids as a readiness contract. During the reconciliation phase, the system verifies that every referenced record exists in the committedRecords set before allowing the operator to proceed. This prevents race conditions where a downstream agent might start processing before its inputs are durably persisted.
The Reconciliation Loop
The core scheduling logic resides in the reconcileAgentGraphSchedule function within packages/runtime/src/stream-graph-schedule-reconcile.ts. This durable loop repeatedly evaluates the graph state until all work reaches a terminal condition.
Snapshot and Validation
Each iteration begins by loading a consistent snapshot of the current schedule, topology, and runtime observations. The scheduler first applies any stops—cancelling or terminating work that has been removed from the graph—to ensure the state machine remains consistent with user intent.
Detecting Missing Inputs with missingWorkInputIds
Before provisioning an operator, the scheduler calls missingWorkInputIds to validate that all dependencies are satisfied:
function missingWorkInputIds(
work,
committedRecords,
selectedResultRecords,
): string[] {
return [
...work.inputIds.filter(id => !committedRecords.has(id)),
...(work.selectedResultInputs ?? [])
.filter(sel => !selectedResultRecords.has(selectedResultKey(sel)))
.map(sel => sel.resultId),
];
}
If this function returns any IDs, the work item is deferred with the reason input_not_committed. The scheduler assigns a status of waiting and skips dispatch for this cycle, automatically re-evaluating the deferred items during the next reconciliation run once the upstream work completes and commits its results.
Building Deterministic Intents
For work with satisfied dependencies, the scheduler constructs a scheduledWorkIntent containing a deterministic hash of the dependency context. The policyFingerprint captures the complete input set to ensure reproducibility:
const policyFingerprint = stableHash({
schemaVersion: SCHEDULE_INTENT_SCHEMA_VERSION,
kind: 'supervisor',
graphId,
workId,
target,
inputIds,
...(selectedResultInputs?.length ? { selectedResultInputs } : {}),
});
This fingerprint, generated at lines 64–76 of stream-graph-schedule-reconcile.ts, serves as the identity for the runnable unit, allowing the system to detect changes in dependencies across reconciliation epochs.
Dispatch and Deferral Logic
Once the intent is built, the scheduler calls claimAgentGraphRunnableIntent to atomically claim the work, then dispatches it to the executor via the operator provision system. If deferral occurred due to missing inputs, the scheduler updates the schedule metadata to reflect the waiting state, ensuring no resources are wasted on premature execution.
Resolving Historical Results
When work references historical results from previous graph epochs via selectedResultInputs, the scheduler resolves these through resolveSelectedResultRecords. To prevent redundant resolver calls, the system maintains a cache of both successful and failed resolutions during each reconciliation run:
const pending = group.filter(selected => !cached.has(selected.resultId));
...
records = await input.resolveSelectedResultInputs(pending);
...
resolved.set(selectedResultKey(selected), clonePlain(record));
If a historical input cannot be resolved, the dependent work is deferred with input_not_committed, preventing deadlocks in cyclic or long-running graph configurations.
Complete Execution Flow
The end-to-end pipeline for scheduling dependent work follows this deterministic sequence:
- User Request: The supervisor submits
update_agent_graphwithadd_workcontaininginput_idsreferencing upstream results. - Snapshot Loading:
reconcileAgentGraphScheduleloads the current schedule state and topology bindings. - Dependency Validation: The scheduler invokes
missingWorkInputIds; if inputs are missing, it defers the work immediately. - Operator Provisioning: For ready work, the system provisions dynamic operators and builds execution edges.
- Intent Construction: The
scheduledWorkIntentis created withpolicyFingerprintandreadinessContextFingerprinthashes. - Atomic Claim:
claimAgentGraphRunnableIntentreserves the work to prevent double execution. - Dispatch: The executor receives the runnable intent only after all
input_idsare verified committed. - Result Recording: Completed work commits its
resultRecordId, triggering the next reconciliation cycle to unblock dependent items.
Summary
- Dependency Declaration: Users bind dependent work by listing upstream
resultRecordIdvalues in theinput_idsarray of theadd_workpayload. - Deferred Execution: The
reconcileAgentGraphScheduleloop instream-graph-schedule-reconcile.tsdefers any work lacking committed inputs, assigning awaitingstatus until dependencies resolve. - Deterministic Identity: The
policyFingerprinthash captures the complete dependency context, ensuring reproducible scheduling decisions across graph epochs. - Historical Resolution: The
selectedResultInputsmechanism resolves cross-epoch dependencies with per-reconciliation caching to optimize performance. - Durable Guarantees: Work only dispatches after
missingWorkInputIdsconfirms all references exist incommittedRecords, eliminating race conditions in distributed execution.
Frequently Asked Questions
What happens if an upstream result is not committed when the scheduler runs?
The scheduler defers the dependent work with the reason input_not_committed and marks it with a waiting status. The work remains in the deferred queue until the next reconciliation cycle detects that the upstream record has been added to committedRecords, at which point the scheduler provisions the operator and dispatches the intent.
How does the scheduler handle references to historical results from previous epochs?
Historical results are referenced via selectedResultInputs and resolved once per reconciliation through resolveSelectedResultRecords. The system caches these resolutions to avoid redundant lookups. If a historical reference cannot be resolved, the work is deferred with input_not_committed, ensuring the graph does not proceed with missing data.
What is the purpose of the policyFingerprint in the scheduled work intent?
The policyFingerprint is a deterministic hash generated by stableHash that encodes the inputIds, selectedResultInputs, graphId, and workId. This fingerprint acts as a unique identity for the specific dependency configuration, enabling the scheduler to track whether the readiness context has changed between reconciliation epochs and to ensure idempotent execution.
Can dependent work target different agent types within the same graph?
Yes. The target_kind field in the work description specifies the execution context, such as new_agent for spawning a child operator. The input_ids mechanism is agnostic to the target agent type; the scheduler resolves the dependencies in stream-graph-schedule-reconcile.ts before the operator provision phase determines where to execute the work.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →