import type { AuthorityMutationV3 } from "../authority-outbox.ts"; import { buildEventV3 } from "../builder.ts"; import { type FingerprintContextV3, fingerprintV3, normalizeNativeIdV3 } from "../canonical.ts"; import type { EventV3, SpanSummaryV3 } from "../contract.ts"; import { eventIdV3 } from "../ids.ts"; import { describePathTargetV3 } from "../targets.ts"; export type CoordinationAuthoritySignalV3 = | "task-changed" | "lifecycle-changed" | "claim-changed" | "identity-attested" | "decision-state-changed" | "wait-started" | "wait-ended"; export interface CoordinationProducerContextV3 { coordRoot: string; root_id: `root_${string}`; instance_id: `inst_${string}`; session_id: `sid_${string}`; generation_id: `gen_${string}`; attestation_id: `att_${string}`; producer_id: `prd_${string}`; boot_id: `boot_${string}`; sequence: number; build_id: `build_${string}`; platform: "linux" | "windows" | "macos" | "unknown"; bridge?: "codex-wsl"; actor_instance_id: `inst_${string}`; subject_instance_id: `inst_${string}`; transaction_id: `txn_${string}`; caused_by?: `evt_${string}`[]; event_id?: `evt_${string}`; observed_at?: string; recorded_at?: string; monotonic_ns?: string; clock_id?: `clk_${string}`; fingerprintContext: FingerprintContextV3; attribution_method: "session_env" | "heartbeat_match" | "explicit_argument"; span_id?: `span_${string}`; parent_span_id?: `span_${string}`; terminal_span?: SpanSummaryV3; } interface CoordinationObservationBaseV3 { native_observation_id: string; } export interface TaskChangedObservationV3 extends CoordinationObservationBaseV3 { state: "set" | "cleared"; prior_state?: string; task?: string; } export interface LifecycleChangedObservationV3 extends CoordinationObservationBaseV3 { state: "active" | "blocked" | "done"; prior_state?: string; reason_code?: string; } export interface ClaimChangedObservationV3 extends CoordinationObservationBaseV3 { operation: "acquired" | "released"; target: unknown; access: "read" | "write"; } export interface IdentityAttestedObservationV3 extends CoordinationObservationBaseV3 { identity_id: string; method: string; } export interface DecisionStateChangedObservationV3 extends CoordinationObservationBaseV3 { decision_id: string; prior_state?: string; outcome: "approved" | "denied" | "deferred"; record_digest: `sha256:${string}`; } export interface WaitStartedObservationV3 extends CoordinationObservationBaseV3 { wait_id: string; kind: "permission" | "approval" | "decision" | "operator_input" | "dependency" | "scheduled"; wake_at?: string; } export interface WaitEndedObservationV3 extends CoordinationObservationBaseV3 { wait_id: string; outcome: | "succeeded" | "failed" | "cancelled" | "timed_out" | "denied" | "interrupted" | "unknown"; } export interface CoordinationObservationBySignalV3 { "task-changed": TaskChangedObservationV3; "lifecycle-changed": LifecycleChangedObservationV3; "claim-changed": ClaimChangedObservationV3; "identity-attested": IdentityAttestedObservationV3; "decision-state-changed": DecisionStateChangedObservationV3; "wait-started": WaitStartedObservationV3; "wait-ended": WaitEndedObservationV3; } export interface NormalizedCoordinationAuthorityV3 { event: EventV3; mutation: AuthorityMutationV3; } /** * Translate one authority-bearing coordination observation without retaining * raw task text or claim targets. The event and mutation are returned as one * unit so callers cannot independently construct contradictory records. */ export function normalizeCoordinationAuthorityV3( signal: S, observation: CoordinationObservationBySignalV3[S], context: CoordinationProducerContextV3, ): NormalizedCoordinationAuthorityV3 { if (context.instance_id !== context.actor_instance_id) { throw new Error("coordination producer instance must be the authority actor"); } const sourceRecordId = normalizeNativeIdV3( context.fingerprintContext, "agent-coord.authority-source", observation.native_observation_id, ); const common = { event_id: context.event_id ?? eventIdV3(), producer: { producer_id: context.producer_id, boot_id: context.boot_id, sequence: context.sequence, component: "agent-coord" as const, build_id: context.build_id, platform: context.platform, ...(context.bridge ? { bridge: context.bridge } : {}), }, scope: { root_id: context.root_id, instance_id: context.instance_id, session_id: context.session_id, generation_id: context.generation_id, }, attestation_id: context.attestation_id, links: { caused_by: context.caused_by ?? [], ...(context.span_id ? { span_id: context.span_id } : {}), ...(context.parent_span_id ? { parent_span_id: context.parent_span_id } : {}), }, provenance: { source_event: `agent-coord.${signal}`, attestation: "derived" as const, confidence: "exact" as const, source_record_id: sourceRecordId, attribution: { method: context.attribution_method, state: "verified" as const, observer_instance_id: context.actor_instance_id, subject_instance_id: context.subject_instance_id, }, }, observed_at: context.observed_at, recorded_at: context.recorded_at, monotonic_ns: context.monotonic_ns, clock_id: context.clock_id, }; const authority = { transaction_id: context.transaction_id }; switch (signal) { case "task-changed": { const input = observation as TaskChangedObservationV3; if (input.state === "set" && !input.task) { throw new Error("task transition to set requires task text for fingerprinting"); } if (input.state === "cleared" && input.task !== undefined) { throw new Error("task transition to cleared must not include task text"); } const taskFingerprint = input.state === "set" ? fingerprintV3(context.fingerprintContext, "coord.task", input.task, "root") : undefined; return { event: buildEventV3("coord.task_changed", { ...common, payload: { actor_instance_id: context.actor_instance_id, subject_instance_id: context.subject_instance_id, ...(input.prior_state ? { prior_state: safeToken(input.prior_state) } : {}), new_state: input.state, reason: "task_transition", ...(taskFingerprint ? { reason_fingerprint: taskFingerprint } : {}), authority, }, }) as EventV3, mutation: { kind: "task.transition", state: input.state, ...(taskFingerprint ? { task_fingerprint: taskFingerprint.digest } : {}), }, }; } case "lifecycle-changed": { const input = observation as LifecycleChangedObservationV3; const reason = input.reason_code ? reasonCode(input.reason_code) : "lifecycle_transition"; return { event: buildEventV3("coord.lifecycle_changed", { ...common, payload: { actor_instance_id: context.actor_instance_id, subject_instance_id: context.subject_instance_id, ...(input.prior_state ? { prior_state: safeToken(input.prior_state) } : {}), new_state: input.state, reason, authority, }, }) as EventV3, mutation: { kind: "lifecycle.transition", state: input.state, reason_code: reason }, }; } case "claim-changed": { const input = observation as ClaimChangedObservationV3; const target = typeof input.target === "string" ? describePathTargetV3({ coordRoot: context.coordRoot, value: input.target, access: input.access, fingerprintContext: context.fingerprintContext, extractorVersion: "agent-coord-claim-v3", }) : { kind: "resource" as const, access: input.access, fingerprint: fingerprintV3( context.fingerprintContext, "semantic-target.resource", input.target, "root", ), extractor_version: "agent-coord-claim-v3", }; return { event: buildEventV3("coord.claim_changed", { ...common, payload: { actor_instance_id: context.actor_instance_id, subject_instance_id: context.subject_instance_id, operation: input.operation, target, access: input.access, authority, }, }) as EventV3, mutation: { kind: "claim.transition", operation: input.operation, target_fingerprint: target.fingerprint.digest as `sha256:${string}`, access: input.access, }, }; } case "identity-attested": { const input = observation as IdentityAttestedObservationV3; const identityId = safeToken(input.identity_id); return { event: buildEventV3("coord.identity_attested", { ...common, payload: { actor_instance_id: context.actor_instance_id, subject_instance_id: context.subject_instance_id, identity_id: identityId, method: safeToken(input.method), authority, }, }) as EventV3, mutation: { kind: "identity.assume", identity_id: identityId }, }; } case "decision-state-changed": { const input = observation as DecisionStateChangedObservationV3; const decisionId = safeToken(input.decision_id); return { event: buildEventV3("decision.state_changed", { ...common, payload: { decision_id: decisionId, ...(input.prior_state ? { prior_state: safeToken(input.prior_state) } : {}), new_state: input.outcome, record_digest: input.record_digest, authority, }, }) as EventV3, mutation: { kind: "decision.resolve", decision_id: decisionId, outcome: input.outcome }, }; } case "wait-started": { const input = observation as WaitStartedObservationV3; const waitId = safeToken(input.wait_id); return { event: buildEventV3("wait.started", { ...common, payload: { wait_id: waitId, kind: input.kind === "operator_input" ? "needs_input" : input.kind === "dependency" ? "unknown" : input.kind, authority_reference: context.transaction_id, ...(input.wake_at ? { wake_at: input.wake_at } : {}), }, }) as EventV3, mutation: { kind: "wait.start", wait_id: waitId, wait_kind: input.kind }, }; } case "wait-ended": { if (!context.terminal_span) throw new Error("wait terminal requires span evidence"); const input = observation as WaitEndedObservationV3; const waitId = safeToken(input.wait_id); return { event: buildEventV3("wait.ended", { ...common, payload: { wait_id: waitId, outcome: input.outcome, resolution_reference: context.transaction_id, span: context.terminal_span, }, }) as EventV3, mutation: { kind: "wait.end", wait_id: waitId, outcome: input.outcome }, }; } } } function safeToken(value: string): string { const normalized = value.normalize("NFC"); if (!/^[a-zA-Z0-9][a-zA-Z0-9._:/+-]{0,127}$/.test(normalized)) { throw new Error("coordination token is invalid"); } return normalized; } function reasonCode(value: string): string { const normalized = value.normalize("NFC"); if (!/^[a-z0-9][a-z0-9._-]{0,79}$/.test(normalized)) { throw new Error("coordination reason code is invalid"); } return normalized; }