Orchestrate operations
The operations library exposes one orchestrator per op family. The persistence-bearing baseline orchestrators take your adapter (the port implementation from Implement an adapter) plus the wire request, run the protocol gates, load and persist data through your ports, and return a SpokeResult — never a throw for expected rejects. Every orchestrator is an async entrypoint: call it with await (TypeScript) or .await inside an async fn (Rust). The baseline calls run unchanged against a RemoteAdapter over a consumer Transport or a multi-peer router — either drops in as the BaselinePorts implementation. The optional extraction path takes its own standalone ExtractionPort plus the async extractor callback, loads the referenced input through that port, runs the extractor once, and returns provisional candidates for later admission through promote; extraction over connect is delegated as a whole-operation remote op.
The orchestrators
| Orchestrator | Request | Response | What it runs |
|---|---|---|---|
orchestrateUpsert(ports, request) | UpsertRequest | UpsertResponse | validate → status gate → batch uniqueness → OCC put |
orchestratePromote(ports, request) | PromoteRequest | PromoteResponse | promote acceptance gates → merge target → OCC put |
orchestrateRelate(ports, request) | RelateRequest | RelateResponse | validate → create/update OCC put |
orchestrateCheck(ports, request, runChecker) | CheckRequest | CheckResponse | resolve rules → load scope → run checker → persist findings |
orchestrateAssemble(ports, request) | AssembleRequest | AssembleResponse | load scope → filter → build AssemblePacket |
orchestrateProject(ports, request) — l2-computable | ProjectRequest | ProjectResponse | validate → ComputablePort.project |
orchestrateCompute(ports, request) — l2-computable | ComputeRequest | ComputeResponse | validate → ComputablePort.compute |
orchestrateForkCheck / orchestrateForkAssemble — l5-fork | fork-scoped requests | same response shapes | require scope.fork_id → fork timeline reads |
orchestrateExtract(ports, request, runExtractor) — ke-extraction | ExtractRequest | ExtractResponse | validate → ExtractionPort.loadExtractionInput → runExtractor → provisional gate |
Upsert — create or update entries
import { orchestrateUpsert } from "@42ch/spoke-operations";
import type { UpsertRequest } from "@42ch/spoke-schemas";
async function runUpsert() {
const result = await orchestrateUpsert(adapter, {
knowledge_entries: [mira, harbor],
});
if (result.ok) {
console.log(result.value.knowledge_entries.map((e) => e.entry_id));
}
}UpsertRequest carries 1..n entries plus an optional idempotency_key (an opaque hint — wire semantics are product-side). The orchestrator validates each entry (MISSING_REQUIRED_FIELD, EMPTY_CANONICAL_NAME, …), gates status transitions when the entry already exists, checks active-uniqueness against the batch, and persists with the correct expected base revision.
Promote — admit a candidate to durable storage
import { orchestratePromote } from "@42ch/spoke-operations";
import type { PromoteRequest } from "@42ch/spoke-schemas";
async function runPromote() {
const result = await orchestratePromote(adapter, {
candidate: provisionalEntry, // typically status "provisional"
target_entry_id: "kb_existing", // optional merge target
});
}Promote covers admission: it admits one candidate to durable storage, and extraction output reaches durability through this step. Promote runs the acceptance gates (CANDIDATE_NOT_PROVISIONAL, CANDIDATE_TERMINAL_STATUS, …) and the revision gate, applies the acceptance transition, and persists through putKnowledgeEntry. With a target_entry_id, the response carries superseded_id for the merged-away entry.
Relate — typed directed edges
import { orchestrateRelate } from "@42ch/spoke-operations";
import type { RelateRequest } from "@42ch/spoke-schemas";
async function runRelate() {
const result = await orchestrateRelate(adapter, {
relation: {
schema_version: 1,
relation_id: "rel_mira_harbor",
relation_type: "located_in",
from_id: "kb_mira",
to_id: "kb_harbor",
extensions: {},
},
});
}Relation validation distinguishes create vs update (RELATION_SELF_EDGE, RELATION_MISSING_ENDPOINT, …), and the OCC-aware put handles revision assignment in your adapter.
Check — run a checker over a scope
orchestrateCheck loads the scoped rules and data first, then hands you a CheckRunInput — your checker callback returns Finding[], and the orchestrator persists them:
import { orchestrateCheck, spokeOk, type CheckRunInput } from "@42ch/spoke-operations";
import type { CheckRequest } from "@42ch/spoke-schemas";
async function runCheck() {
const result = await orchestrateCheck(adapter, checkRequest, (input: CheckRunInput) => {
// input: { request, entries, events, rules }
const findings = myChecker(input.entries, input.rules);
return spokeOk(findings); // or spokeReject(SpokeRejectCode.INVALID_INPUT, "...")
});
}Rules resolve from rule_refs via RuleQueryPort, with embedded rules[] overriding by rule_id. check returns findings only — use assemble for context packets.
Assemble — build a context packet
import { orchestrateAssemble } from "@42ch/spoke-operations";
import type { AssembleRequest } from "@42ch/spoke-schemas";
async function runAssemble() {
const result = await orchestrateAssemble(adapter, {
scope: { scope_id: "book-harbor", entry_types: ["character"] },
max_entries: 20, // optional entry limit hint
});
}The orchestrator loads the scope, applies scope filters, and builds a wire-only AssemblePacket with order-preserving truncation. Assembly itself — ranking, retrieval, token budgets — is product-side.
Extract — propose provisional candidates
orchestrateExtract(ports, request, runExtractor) is the optional ke-extraction path. ExtractionPort is a standalone optional family passed directly to orchestrateExtract, together with your own async extractor callback:
import { orchestrateExtract, spokeOk, type ExtractionPort, type RunExtractor } from "@42ch/spoke-operations";
import type { ExtractRequest } from "@42ch/spoke-schemas";
const extractionPort: ExtractionPort = {
async loadExtractionInput(request: ExtractRequest) {
// host-local source loading: the returned value stays in-process
return spokeOk(await readSources(request.sources));
},
};
const runExtractor: RunExtractor = async ({ request, input }) => {
const candidates = await myExtractionService.propose(input);
return spokeOk({ candidates, method: "llm-v1" });
};
async function runExtract() {
const result = await orchestrateExtract(extractionPort, extractRequest, runExtractor);
if (result.ok) {
// result.value.candidates — each entry carries status "provisional"
}
}The orchestrator validates run_id and sources, loads the referenced material through the port, then invokes your extractor exactly once with { request, input }. The loaded value stays in-process and opaque; the wire carries the request's references and the response's candidates. The run returns candidates plus run metadata — provisional proposals whose durable admission happens later through promote. A missing port at a dynamic boundary rejects as CAPABILITY_PORT_MISSING with details.capability = "ke-extraction".
Rust spells the same entrypoint orchestrate_extract, with ExtractRunInput / ExtractionResult structs and the callback as a generic F: FnOnce(ExtractRunInput) -> Fut.
Handle rejects
Every orchestrator returns SpokeResult:
import { SpokeRejectCode } from "@42ch/spoke-operations";
if (!result.ok) {
switch (result.code) {
case SpokeRejectCode.REVISION_CONFLICT:
case SpokeRejectCode.STORED_REVISION_STALE:
// reload and retry with the fresh revision
break;
case SpokeRejectCode.CAPABILITY_PORT_MISSING:
// the adapter does not implement the optional port this op needs
break;
default:
// validation / state rejects — surface to the caller
}
}The wire responses follow the same one-failure dialect as the request/response envelopes: a response is either the success payload or { "error": ErrorEnvelope } — never both. The library's rejects map to ErrorEnvelope shapes your transport carries.
The purity boundary
The orchestrators run the protocol gates and drive your ports — they do not touch storage, LLM calls, ranking, retrieval, or transport directly. All of that is supplied by your product through the injected adapter. Finding and promote lifecycles (status transitions, acceptance gates) are pure, pre-persist rules in the library; persistence happens through your ports.
Next steps
- Ops wire reference — request/response envelope field tables and
Scope. - Implement an adapter — the port contract behind every orchestrator.
- Walk the ToyWorld reference adapter — orchestrator usage in the committed fixture graph.