import { randomUUID } from "node:crypto"; import type { Adapter } from "../adapter.ts"; import { EVENT_ADAPTER_IDS_V3 } from "../events/v3/adapter-id.ts"; import type { AuthorityMutationV3 } from "../events/v3/authority-outbox.ts"; import { canonicalJsonV3, normalizeNativeIdV3, sha256V3 } from "../events/v3/canonical.ts"; import { type EventV3WriteMode, readEventV3ControlState } from "../events/v3/control.ts"; import { readCoordinationViewV3, requireAuthoritySafeCoordinationViewV3, } from "../events/v3/coordination-view.ts"; import { fingerprintContextV3 } from "../events/v3/fingerprint-keys.ts"; import { LIVE_HOOK_V3_PRODUCER_ID, liveInstanceIdV3, livePlatformV3, resolveLiveEventLedgerRouteV3, } from "../events/v3/live-routing.ts"; import type { CoordinationAuthoritySignalV3, CoordinationObservationBySignalV3, } from "../events/v3/producers/coordination.ts"; import { type RecordCoordinationAuthorityV3Result, recordCoordinationAuthorityV3, } from "../events/v3/producers/coordination-recorder.ts"; import { type HookProducerStateV3, listHookProducerStateRecordsV3, readHookProducerStateV3, readTerminalHookProducerStateV3, recordHookSignalV3, } from "../events/v3/producers/recorder.ts"; import { canonicalClaimPath } from "./claim-path.ts"; import { acquireClaim, type Heartbeat, readHeartbeat, releaseClaim, setIdentityCache, setLifecycleCache, setTask, } from "./state/heartbeat-writer.ts"; import { liveCoordinationAdapterV3, readLiveCoordinationRow, } from "./state/live-coordination-view.ts"; import { ensureLiveCoordinationHeartbeat } from "./state/live-coordination-writer.ts"; import { recordNameAssumption } from "./state/names.ts"; export const LIVE_COORDINATION_V3_PRODUCER_ID = "prd_agent-coord" as const; export class LiveCoordinationAuthorityV3Error extends Error { constructor(public readonly reason: string) { super(`event_v3_coordination_authority:${reason}`); this.name = "LiveCoordinationAuthorityV3Error"; } } export type LiveCoordinationAuthorityV3Result = | { state: "unchanged" } | { state: "recorded"; result: RecordCoordinationAuthorityV3Result }; export interface ReopenedLiveCoordinationGenerationV3 { state: "reopened"; adapter: Adapter; prior_generation_id: `gen_${string}`; generation_id: `gen_${string}`; } export interface BootstrappedLiveCoordinationAuthorityV3 { state: "created" | "reused"; adapter: Adapter; generation_id: `gen_${string}`; heartbeat: Heartbeat; } export interface RestoredLiveCoordinationStateV3 { state: "unchanged" | "restored"; task: boolean; claims: number; lifecycle: boolean; } interface LiveAuthorityBaseV3 { coordRoot: string; owner: string; subject?: string; nativeSessionId: string; adapter: Adapter; observationId?: string; } /** * Establish the missing generation authority for one adapter-native session, * then reconstruct its disposable heartbeat cache. * * This is intentionally narrower than mid-flight hook onboarding. The caller * must supply an already validated native session, adapter, and current * instance identity. Existing authority is reused only when all three bind to * the same non-terminal producer generation. A terminal generation, an * adapter/session collision, or any disagreement between private producer * control and the public projection fails closed. No turn is synthesized. */ export function bootstrapLiveCoordinationAuthorityV3(input: { coordRoot: string; owner: string; nativeSessionId: string; adapter: Adapter; }): BootstrappedLiveCoordinationAuthorityV3 { if (!/^[a-zA-Z0-9._-]{1,128}$/.test(input.owner)) { throw new LiveCoordinationAuthorityV3Error("invalid_instance_identity"); } if ( input.nativeSessionId.length === 0 || input.nativeSessionId.length > 512 || input.nativeSessionId.includes("\0") ) { throw new LiveCoordinationAuthorityV3Error("invalid_native_session_identity"); } const route = resolveLiveEventLedgerRouteV3(input.coordRoot); if (route.state === "blocked") throw new LiveCoordinationAuthorityV3Error(route.reason); const control = readEventV3ControlState(input.coordRoot); if (control.state !== route.mode) { throw new LiveCoordinationAuthorityV3Error("control_state_mismatch"); } const instanceId = liveInstanceIdV3(input.owner); const rootId = control.genesis.event.scope.root_id as `root_${string}`; const context = fingerprintContextV3( input.coordRoot, rootId, undefined, control.genesis.profile.privacy_key_epoch, ); const sessionId = `sid_${normalizeNativeIdV3( context, `${input.adapter}.session`, input.nativeSessionId, ).slice(4)}`; const view = readAuthoritySafeViewForBootstrap(input.coordRoot); const existing = view.instances[instanceId]; const sessionStates = readSessionProducerStatesForBootstrap( input.coordRoot, input.nativeSessionId, ); const instanceStates = readInstanceProducerStatesForBootstrap(input.coordRoot, instanceId); const crossAdapter = sessionStates.find((state) => state.adapter !== input.adapter); if (crossAdapter) { throw new LiveCoordinationAuthorityV3Error("cross_adapter_identity"); } const crossInstance = sessionStates.find((state) => state.instance_id !== instanceId); if (crossInstance) { throw new LiveCoordinationAuthorityV3Error("cross_instance_identity"); } const direct = sessionStates.find((state) => state.adapter === input.adapter); if ( direct?.terminal || (!existing && terminalGenerationForSession(view, instanceId, sessionId)) ) { throw new LiveCoordinationAuthorityV3Error("terminal_generation_forbidden"); } const liveInstanceStates = instanceStates.filter((state) => !state.terminal); if (liveInstanceStates.length > 1 || sessionStates.length > 1) { throw new LiveCoordinationAuthorityV3Error("generation_ambiguous"); } if (existing) { const observedAdapter = existing.runtime_attestation.adapter; if (observedAdapter.state !== "observed" || observedAdapter.value.id !== input.adapter) { throw new LiveCoordinationAuthorityV3Error("cross_adapter_identity"); } if (existing.session_id !== sessionId) { throw new LiveCoordinationAuthorityV3Error("cross_session_identity"); } if ( !direct || direct.terminal || direct.generation_id !== existing.generation_id || liveInstanceStates.length !== 1 || liveInstanceStates[0]?.generation_id !== existing.generation_id ) { throw new LiveCoordinationAuthorityV3Error("generation_control_mismatch"); } return materializeBootstrappedAuthority(input, existing.generation_id, "reused"); } if (direct || liveInstanceStates.length > 0) { throw new LiveCoordinationAuthorityV3Error("generation_control_mismatch"); } const started = recordHookSignalV3({ coordRoot: input.coordRoot, mode: route.mode, signal: "session-start", payload: { raw: {}, session_id: input.nativeSessionId }, adapter: input.adapter, instance_id: instanceId, producer_id: LIVE_HOOK_V3_PRODUCER_ID, build_id: route.build_id, platform: livePlatformV3(), session_start_derivation: "validated_current_session_heal", }); if (started.state !== "recorded" && started.state !== "already_started") { const reason = "reason" in started ? `:${started.reason}` : ""; throw new LiveCoordinationAuthorityV3Error( `generation_bootstrap_failed:${started.state}${reason}`, ); } const after = readAuthoritySafeViewForBootstrap(input.coordRoot).instances[instanceId]; const producer = readSessionProducerStatesForBootstrap( input.coordRoot, input.nativeSessionId, ).filter((state) => state.adapter === input.adapter && state.instance_id === instanceId); if ( !after || after.session_id !== sessionId || producer.length !== 1 || producer[0]?.terminal || producer[0]?.generation_id !== after.generation_id ) { throw new LiveCoordinationAuthorityV3Error("generation_control_mismatch"); } return materializeBootstrappedAuthority( input, after.generation_id, started.state === "recorded" ? "created" : "reused", ); } function readAuthoritySafeViewForBootstrap(coordRoot: string) { try { return requireAuthoritySafeCoordinationViewV3(readCoordinationViewV3(coordRoot)); } catch (error) { if (error instanceof LiveCoordinationAuthorityV3Error) throw error; throw new LiveCoordinationAuthorityV3Error("coordination_view_unsafe"); } } function readSessionProducerStatesForBootstrap( coordRoot: string, nativeSessionId: string, ): HookProducerStateV3[] { try { return EVENT_ADAPTER_IDS_V3.map((adapter) => readHookProducerStateV3(coordRoot, adapter, nativeSessionId), ).filter((state): state is HookProducerStateV3 => state !== undefined); } catch { throw new LiveCoordinationAuthorityV3Error("producer_state_unsafe"); } } function readInstanceProducerStatesForBootstrap( coordRoot: string, instanceId: `inst_${string}`, ): HookProducerStateV3[] { try { return listHookProducerStateRecordsV3(coordRoot, { includeTerminal: true }) .map((record) => record.state) .filter((state) => state.instance_id === instanceId); } catch { throw new LiveCoordinationAuthorityV3Error("producer_state_unsafe"); } } function terminalGenerationForSession( view: ReturnType, instanceId: `inst_${string}`, sessionId: string, ): boolean { return Object.values(view.terminal_generations).some( (generation) => generation.instance_id === instanceId && generation.session_id === sessionId, ); } function materializeBootstrappedAuthority( input: { coordRoot: string; owner: string; nativeSessionId: string; adapter: Adapter }, generationId: string, state: "created" | "reused", ): BootstrappedLiveCoordinationAuthorityV3 { const heartbeat = ensureLiveCoordinationHeartbeat( input.coordRoot, input.owner, input.nativeSessionId, input.adapter, ); if (!heartbeat || heartbeat.v3_generation_id !== generationId) { throw new LiveCoordinationAuthorityV3Error("generation_materialization_failed"); } return { state, adapter: input.adapter, generation_id: generationId as `gen_${string}`, heartbeat, }; } /** * Reapply private, generation-bound coordination state after a runtime epoch * replaces the canonical ledger. The old heartbeat is only accepted for the * same native session and a different V3 generation. Its task prose and * lifecycle reason remain in the disposable cache; canonical events retain * only the privacy-safe authority signals. */ export function restoreLiveCoordinationStateAfterEpochV3(input: { coordRoot: string; owner: string; nativeSessionId: string; adapter: Adapter; prior: Heartbeat | null; currentGenerationId?: `gen_${string}`; }): RestoredLiveCoordinationStateV3 { const unchanged: RestoredLiveCoordinationStateV3 = { state: "unchanged", task: false, claims: 0, lifecycle: false, }; const prior = input.prior; if (!prior || prior.instance_id !== input.owner) return unchanged; if ( prior.session_id !== input.nativeSessionId && prior.native_session_id !== input.nativeSessionId ) { return unchanged; } // The recorder just committed this hook under the named generation. When // it is the cache's generation, no epoch replacement occurred and the // expensive complete coordination projection cannot change the answer. if (input.currentGenerationId && prior.v3_generation_id === input.currentGenerationId) { return unchanged; } const view = requireAuthoritySafeCoordinationViewV3(readCoordinationViewV3(input.coordRoot)); const priorGenerationStillCurrent = [ ...Object.values(view.instances), ...Object.values(view.terminal_generations), ].some((generation) => generation.generation_id === prior.v3_generation_id); if (priorGenerationStillCurrent) return unchanged; const current = readLiveCoordinationRow(input.coordRoot, input.owner); if (!current || !prior.v3_generation_id || current.v3_generation_id === prior.v3_generation_id) { return unchanged; } let task = false; let claims = 0; let lifecycle = false; if (typeof prior.task === "string" && prior.task.length > 0 && current.task !== prior.task) { recordLiveTaskChangeV3({ coordRoot: input.coordRoot, owner: input.owner, nativeSessionId: input.nativeSessionId, adapter: input.adapter, task: prior.task, }); task = true; } const currentClaims = new Set( (current.files_touched ?? []).map((path) => canonicalClaimPath(input.coordRoot, path)), ); for (const path of [...new Set(prior.files_touched ?? [])].sort()) { const canonical = canonicalClaimPath(input.coordRoot, path); if (currentClaims.has(canonical)) continue; recordLiveClaimChangeV3({ coordRoot: input.coordRoot, owner: input.owner, nativeSessionId: input.nativeSessionId, adapter: input.adapter, operation: "acquired", path: canonical, }); currentClaims.add(canonical); claims += 1; } const priorLifecycle = prior.task_state ?? "active"; if ( priorLifecycle !== "active" && (current.task_state !== priorLifecycle || current.task_state_reason !== prior.task_state_reason) ) { recordLiveLifecycleChangeV3({ coordRoot: input.coordRoot, owner: input.owner, nativeSessionId: input.nativeSessionId, adapter: input.adapter, state: priorLifecycle, ...(prior.task_state_reason ? { reason: prior.task_state_reason } : {}), ...(prior.suggested_session_name ? { suggestedSessionName: prior.suggested_session_name } : {}), }); lifecycle = true; } return task || claims > 0 || lifecycle ? { state: "restored", task, claims, lifecycle } : unchanged; } /** * Open a fresh derived generation for a human-facing session that is executing * again after an authoritative terminal. The terminal generation is never * changed or reused. */ export function reopenLiveCoordinationGenerationV3(input: { coordRoot: string; owner: string; nativeSessionId: string; }): ReopenedLiveCoordinationGenerationV3 { const route = resolveLiveEventLedgerRouteV3(input.coordRoot); if (route.state === "blocked") throw new LiveCoordinationAuthorityV3Error(route.reason); const instanceId = liveInstanceIdV3(input.owner); const terminal = readTerminalHookProducerStateV3( input.coordRoot, input.nativeSessionId, instanceId, ); if (!terminal) { throw new LiveCoordinationAuthorityV3Error("terminal_generation_identity_missing"); } const view = requireAuthoritySafeCoordinationViewV3(readCoordinationViewV3(input.coordRoot)); const terminalView = view.terminal_generations[terminal.generation_id]; if (!terminalView || terminalView.instance_id !== instanceId) { throw new LiveCoordinationAuthorityV3Error("terminal_generation_authority_missing"); } if (terminal.adapter === "openclaw") { throw new LiveCoordinationAuthorityV3Error("lifecycle_not_human_facing"); } if (terminalView.parent_generation_id || terminalView.delegation_id || terminalView.workflow_id) { throw new LiveCoordinationAuthorityV3Error("lifecycle_not_human_facing"); } const reopened = recordHookSignalV3({ coordRoot: input.coordRoot, mode: route.mode, signal: "session-start", payload: { raw: {}, session_id: input.nativeSessionId }, adapter: terminal.adapter, instance_id: instanceId, producer_id: LIVE_HOOK_V3_PRODUCER_ID, build_id: route.build_id, platform: livePlatformV3(), session_start_derivation: "approved_lifecycle_reopen", }); if (reopened.state !== "recorded" && reopened.state !== "already_started") { throw new LiveCoordinationAuthorityV3Error(`generation_reopen_failed:${reopened.state}`); } const heartbeat = ensureLiveCoordinationHeartbeat( input.coordRoot, input.owner, input.nativeSessionId, terminal.adapter, ); if (!heartbeat) { throw new LiveCoordinationAuthorityV3Error("reopened_generation_materialization_failed"); } return { state: "reopened", adapter: terminal.adapter, prior_generation_id: terminal.generation_id, generation_id: heartbeat.v3_generation_id as `gen_${string}`, }; } export function recordLiveTaskChangeV3( input: LiveAuthorityBaseV3 & { task: string }, ): LiveCoordinationAuthorityV3Result { liveCoordinationWriteModeV3(input.coordRoot); const before = requireHeartbeat(input, input.subject ?? input.owner); const cleared = input.task.length === 0; const desired = { ...before, task: undefined, v3_task_state: cleared ? ("cleared" as const) : ("set" as const), }; // A task declaration is also per-turn ritual evidence. Record repeated // declarations, including cleared -> cleared, even when the disposable view // does not change. Otherwise a conversational Cursor remediation turn can // run `set-task ""` exactly as instructed and still loop forever because no // coord.task_changed event reaches the verdict window. return recordLiveAuthority( input, "task-changed", { native_observation_id: input.observationId ?? `task-${randomUUID()}`, state: cleared ? "cleared" : "set", ...(cleared ? {} : { task: input.task }), }, before, desired, () => { // The canonical event records only the privacy-safe set/cleared state. // Task prose and its operator-facing suggested name live exclusively in // this generation-bound disposable cache; they never enter the ledger. if (!setTask(input.coordRoot, input.subject ?? input.owner, input.task)) { throw new LiveCoordinationAuthorityV3Error("task_materialization_failed"); } }, ); } export function recordLiveLifecycleChangeV3( input: LiveAuthorityBaseV3 & { state: "active" | "blocked" | "done"; reason?: string; suggestedSessionName?: string; observedAt?: string; }, ): LiveCoordinationAuthorityV3Result { liveCoordinationWriteModeV3(input.coordRoot); const subject = input.subject ?? input.owner; const before = requireHeartbeat(input, subject); const desired: Heartbeat = { ...before, task_state: input.state, task_state_reason: input.state === "active" ? undefined : input.reason, ...(input.suggestedSessionName ? { suggested_session_name: input.suggestedSessionName } : {}), }; if (coordinationAuthorityStateDigestV3(before) === coordinationAuthorityStateDigestV3(desired)) return { state: "unchanged" }; return recordLiveAuthority( input, "lifecycle-changed", { native_observation_id: input.observationId ?? `lifecycle-${randomUUID()}`, state: input.state, ...(input.reason ? { reason_code: `operator_${input.state}` } : {}), }, before, desired, () => { if ( !setLifecycleCache( input.coordRoot, subject, input.state, input.reason, input.suggestedSessionName, ) ) { throw new LiveCoordinationAuthorityV3Error("lifecycle_materialization_failed"); } }, ); } export function recordLiveClaimChangeV3( input: LiveAuthorityBaseV3 & { operation: "acquired" | "released"; path: string; access?: "read" | "write"; }, ): LiveCoordinationAuthorityV3Result { liveCoordinationWriteModeV3(input.coordRoot); const subject = input.subject ?? input.owner; const before = requireHeartbeat(input, subject); const canonical = canonicalClaimPath(input.coordRoot, input.path); const desiredFiles = input.operation === "acquired" ? [ ...new Set([ ...(before.files_touched ?? []).map((path) => canonicalClaimPath(input.coordRoot, path), ), canonical, ]), ].sort() : (before.files_touched ?? []).filter( (path) => canonicalClaimPath(input.coordRoot, path) !== canonical, ); const desired: Heartbeat = { ...before, files_touched: desiredFiles }; if (coordinationAuthorityStateDigestV3(before) === coordinationAuthorityStateDigestV3(desired)) return { state: "unchanged" }; return recordLiveAuthority( input, "claim-changed", { native_observation_id: input.observationId ?? `claim-${randomUUID()}`, operation: input.operation, target: canonical, access: input.access ?? "write", }, before, desired, () => { const result = input.operation === "acquired" ? acquireClaim(input.coordRoot, subject, canonical) : releaseClaim(input.coordRoot, subject, canonical); if (!result) throw new LiveCoordinationAuthorityV3Error("claim_materialization_failed"); }, ); } export function recordLiveIdentityChangeV3( input: LiveAuthorityBaseV3 & { name: string; identityId: string }, ): LiveCoordinationAuthorityV3Result { liveCoordinationWriteModeV3(input.coordRoot); const subject = input.subject ?? input.owner; const before = requireHeartbeat(input, subject); const desired: Heartbeat = { ...before, name: input.name, agent_id: input.identityId }; if (coordinationAuthorityStateDigestV3(before) === coordinationAuthorityStateDigestV3(desired)) return { state: "unchanged" }; return recordLiveAuthority( input, "identity-attested", { native_observation_id: input.observationId ?? `identity-${randomUUID()}`, identity_id: input.identityId, method: "operator_assumption", }, before, desired, () => { recordNameAssumption(input.coordRoot, subject, input.name, input.identityId, "session"); if (!setIdentityCache(input.coordRoot, subject, input.name, input.identityId)) { throw new LiveCoordinationAuthorityV3Error("identity_materialization_failed"); } }, ); } function recordLiveAuthority( input: LiveAuthorityBaseV3, signal: S, observation: CoordinationObservationBySignalV3[S], expected: Heartbeat, desired: Heartbeat, apply: () => void, ): LiveCoordinationAuthorityV3Result { const route = resolveLiveEventLedgerRouteV3(input.coordRoot); if (route.state === "blocked") throw new LiveCoordinationAuthorityV3Error(route.reason); const adapter = liveCoordinationAdapterV3(input.coordRoot, input.owner); if (!adapter) throw new LiveCoordinationAuthorityV3Error("actor_generation_missing"); const subject = input.subject ?? input.owner; const result = recordAuthority({ coordRoot: input.coordRoot, mode: route.mode, signal, observation, adapter, native_actor_session_id: input.nativeSessionId, actor_instance_id: liveInstanceIdV3(input.owner), subject_instance_id: liveInstanceIdV3(subject), producer_id: LIVE_COORDINATION_V3_PRODUCER_ID, build_id: route.build_id, platform: livePlatformV3(), expected_prior_state_digest: coordinationAuthorityStateDigestV3(expected), desired_state_digest: coordinationAuthorityStateDigestV3(desired), reconciler: { readStateDigest: () => { const heartbeat = readHeartbeat(input.coordRoot, subject); if (!heartbeat || heartbeat.v3_generation_id !== expected.v3_generation_id) { throw new LiveCoordinationAuthorityV3Error(`heartbeat_generation_mismatch:${subject}`); } return coordinationAuthorityStateDigestV3(heartbeat); }, apply: (_mutation: AuthorityMutationV3) => apply(), }, }); if (result.state === "gate_closed" || result.state === "generation_unavailable") { throw new LiveCoordinationAuthorityV3Error(`${result.state}:${result.reason}`); } if (result.state === "pending_transaction") { throw new LiveCoordinationAuthorityV3Error(`pending_transaction:${result.transaction_id}`); } return { state: "recorded", result }; } function recordAuthority( input: Parameters>[0], ): RecordCoordinationAuthorityV3Result { return recordCoordinationAuthorityV3(input); } export function liveCoordinationWriteModeV3(coordRoot: string): EventV3WriteMode { const route = resolveLiveEventLedgerRouteV3(coordRoot); if (route.state === "blocked") throw new LiveCoordinationAuthorityV3Error(route.reason); return route.mode; } function requireHeartbeat(input: LiveAuthorityBaseV3, owner: string): Heartbeat { const heartbeat = ensureLiveCoordinationHeartbeat( input.coordRoot, owner, owner === input.owner ? input.nativeSessionId : owner, input.adapter, ); if (!heartbeat) throw new LiveCoordinationAuthorityV3Error(`heartbeat_missing:${owner}`); return heartbeat; } export function coordinationAuthorityStateDigestV3(heartbeat: Heartbeat): `sha256:${string}` { return sha256V3( canonicalJsonV3({ instance_id: heartbeat.v3_instance_id ?? heartbeat.instance_id, generation_id: heartbeat.v3_generation_id ?? null, task_state: heartbeat.v3_task_state ?? (heartbeat.task ? "set" : "cleared"), lifecycle_state: heartbeat.task_state ?? "active", task_state_reason: heartbeat.task_state_reason ?? null, files_touched: [...new Set(heartbeat.files_touched ?? [])].sort(), identity_id: heartbeat.agent_id ?? null, display_name: heartbeat.name ?? null, }), ); }