import { type OwnerChannelConfig } from '../config.js'; import { type AgentSession } from '../session/types.js'; import { type OwnerFleetOps } from './commands.js'; import type { ManagedFleetSpawnResult } from '../fleet-proxy.js'; import { type FleetAuditAttempt, type FleetAuditPresentation, type FleetCommandOutcomeClass } from '../fleet-command-audit.js'; import { type OursOps } from './ours-client.js'; import { type OwnerUpdatePhase } from './notices.js'; import { type OwnerEntry } from './state.js'; import { type OwnerTaskPhase } from './tasks.js'; import { type OwnerBinderDeps, type OwnerBinderLease } from './binder.js'; export interface OwnerChannelOptions { role: string; /** Harness id of the role (e.g. 'claude-code', 'codex'); gates which slash commands may be forwarded. */ harness: string; config: OwnerChannelConfig; session: AgentSession; stateDir: string; env?: Record; log(line: string): void; client?: OursOps; /** Test seam; production uses the detached ours-fleet CLI (`fleetCliOps`). */ fleet?: OwnerFleetOps; /** Read-only restart validation; production uses the shared lifecycle service. */ prepareRestart?: (role: string, mode: 'keep' | 'fresh') => Promise; /** Forwarded to fleet CLI invocations spawned for owner commands. */ configPath?: string; /** Deterministic clock/process seams for binder handoff tests. */ binderDeps?: OwnerBinderDeps; /** Pre-acquired by the runner so the predecessor control socket remains reachable while waiting. */ binderLease?: OwnerBinderLease; recoveryDeps?: { now(): number; setTimer(fn: () => void, ms: number): ReturnType; clearTimer(timer: ReturnType): void; deadlineMs?: number; }; } export interface OwnerChannelHandle { start(): Promise; drain(): Promise; close(): Promise; /** Reattach this same supervisor-owned channel after a daemon generation change. */ recover?(epoch: string): Promise; /** False when shutdown skipped client disposal because quiescence was unproven. */ binderReleaseSafe?(): boolean; manage(request: OwnerChannelManagementRequest): Promise; /** Fleet-owned deterministic lifecycle notice; absent on legacy test doubles. */ notifyFleetSpawn?(event: ManagedFleetSpawnResult): Promise; notifyFleetLifecycle?(presentations: FleetAuditPresentation[]): Promise; beginFleetCommandAudit?(requestId: string, argv: string[]): Promise; finishFleetCommandAudit?(input: { correlationId: string; class: FleetCommandOutcomeClass; exitCode?: number; effect: 'not_started' | 'completed' | 'unknown'; resourceIds?: Record; presentations?: FleetAuditPresentation[]; }): Promise; } export declare const OWNER_RECOVERY_QUEUE_CAPACITY = 32; export declare const OWNER_RECOVERY_DEGRADED = "OWNER_RECOVERY_DEGRADED"; export declare class OwnerRecoveryDegradedError extends Error { readonly code = "OWNER_RECOVERY_DEGRADED"; constructor(); } export declare class OwnerRecoveryTimeoutError extends Error { readonly stage: string; readonly code = "OWNER_RECOVERY_DEGRADED"; constructor(stage: string); } export type OwnerChannelManagementRequest = { action: 'contact_list'; } | { action: 'contact_invite'; name?: string; } | { action: 'contact_add'; invite: string; name?: string; } | { action: 'owner_list'; } | { action: 'owner_authorize'; cid: string; } | { action: 'owner_revoke'; cid: string; } | { action: 'request_update'; requestId: string; phase: OwnerUpdatePhase; message: string; } | { action: 'task_open'; requestId: string; } | { action: 'task_report'; taskId: string; phase: OwnerTaskPhase; message: string; } | { action: 'startup_failure'; }; export type OwnerChannelManagementResult = { action: 'contact_list'; contacts: OwnerContact[]; } | { action: 'contact_invite'; invite: string; } | { action: 'contact_add'; status: 'pending'; contact?: OwnerContact; } | { action: 'owner_list'; integrity: { ok: boolean; error?: string; }; owners: OwnerEntry[]; } | { action: 'owner_authorize' | 'owner_revoke'; owner: OwnerEntry; } | { action: 'request_update'; requestId: string; sequence: number; } | { action: 'task_open'; taskId: string; expiresAt: string; } | { action: 'task_report'; taskId: string; phase: OwnerTaskPhase; sequence: number; state: 'open' | 'closed'; } | { action: 'startup_failure'; status: 'delivered' | 'duplicate'; }; export type { OwnerUpdatePhase } from './notices.js'; export interface OwnerContact { cid: string; name: string; /** Structural, from which daemon collection the row came: established or pending. */ status: string; /** * Retained for the `ours-fleet owner contact list` column. The daemon's typed * contact view has no such field, so it is always absent; it is not inferred * from anything a contact controls. */ kind?: string; human?: { cid?: string; name?: string; }; } /** * Fleet-owned trusted ingress. The agent never binds this identity and never * chooses its reply recipient; both are fixed from authenticated message data. */ export declare class OwnerChannel implements OwnerChannelHandle { private readonly options; private readonly client; private readonly state; private readonly authorizations; private readonly conversations; private readonly tasks; private readonly messageRecovery; private readonly attachmentRecovery; private readonly attachmentConfig; private readonly attachmentRoot; /** * Wire IDs whose turn is still running. They stay OUT of the durable state * (a crash must replay them) but must not be queued twice while live. */ private readonly inFlight; /** SDK handlers entered during the current getMessages call. */ private readonly typedHandlersEntered; /** Wires already NACKed to the managed agent, so a history replay stays quiet. */ private readonly relayNacks; /** * fleet.yaml declares the restart baseline; `/comments on|off` changes only * this process's effective value. The override is deliberately memory-only: a * restart must return to the reviewed, checked-in configuration rather than to * an unreviewable file that could silently keep an owner's channel quiet. */ private readonly commentsBaseline; private commentsEnabled; private stopping; private watchTask?; private watchAbort?; private drainTask?; private drainRequested; private readonly completionTasks; private readonly activeRequests; private managementTail; private recoveryEpoch?; private recoveryTask?; private recoveryToken; private recoveryQueued; private watchGeneration; /** A timed-out client operation that must settle before replacement is safe. */ private recoveryQuiescence?; private releaseSafe; private ready; private startedOnce; private shutdownDegraded; private binder?; private binderOwnedInternally; private readonly fleetOps; private readonly prepareRestart; private readonly commandAudits; private readonly lifecycleOutbox; constructor(options: OwnerChannelOptions); start(): Promise; drain(): Promise; close(): Promise; binderReleaseSafe(): boolean; manage(request: OwnerChannelManagementRequest): Promise; notifyFleetSpawn(event: ManagedFleetSpawnResult): Promise; notifyFleetLifecycle(presentations: FleetAuditPresentation[]): Promise; private flushPendingFleetLifecycle; beginFleetCommandAudit(requestId: string, argv: string[]): Promise; finishFleetCommandAudit(input: { correlationId: string; class: FleetCommandOutcomeClass; exitCode?: number; effect: 'not_started' | 'completed' | 'unknown'; resourceIds?: Record; presentations?: FleetAuditPresentation[]; }): Promise; private flushPendingFleetCommandAudits; private deliverFleetCommandOutcome; recover(epoch: string): Promise; private recoveryStage; private writeShutdownState; private queueManagement; private manageNow; /** * The daemon reports established contacts and pending introductions as two * separate collections, so the status is structural rather than a word parsed * out of a rendered line. Nothing here can be spoofed by a contact's own * display name. */ private contacts; private contact; private assertCid; private assertLabel; private safeMetadata; private sendOwnerUpdate; private sendProactiveMessage; private openOwnerTask; private sendOwnerTaskReport; private safeOwnerUpdate; private safeTaskReport; private safeProactiveMessage; private drainAll; /** * Claim the exact oldest unread SQLite batch before marking it read. * The journal contains only wire IDs and sequence numbers; bodies remain in * the daemon's persistent history and are recovered with getHistoryItem. */ private claimMessages; private messageClaim; private historyMessage; private attachmentMetadata; private attachmentGroups; private handleAttachmentGroup; private handle; /** * Deterministic command path: the message never becomes an agent prompt. * Authorization already happened — the managed-agent relay branch and the * owner-CID check in handle() both run before dispatch, so only an * authenticated owner reaches this: neither ordinary peers nor the managed * agent itself can execute /force-restart, /model, or any other command. */ private handleCommand; /** * Register one typed adapter per primary slash command. The adapter performs * the live CID check before constructing the same slash text and entering the * existing dispatcher; its null SDK result avoids duplicating the ordinary * owner-channel replies that remain the command's result contract. */ private registerTypedCommands; private handleRecoveredTypedCommand; private handleTypedCommand; /** Queue raw slash text to the harness and report the turn's outcome. */ private runHarnessCommand; /** * Confirmation and the durable wire record must both land BEFORE the fleet * CLI is asked to bounce this very process; neither can happen afterwards. */ private restartSelf; /** * The caller may itself be a room member. Persist and acknowledge acceptance * before an external worker starts a saga that can retire this process. */ private closeRoomFromOwner; /** Accept a permanent any-state deletion, acknowledge it, then hand cleanup to a durable worker. */ private deleteTaskFromOwner; private terminalTaskFromOwner; /** Code-point-safe tail of the worklog, or undefined when there is none. */ private readWorklogTail; /** Bound harness-command output to a single outbound message. */ private commandOutput; private acceptedSender; private isAgentSender; private isEffectiveOwner; private relayManagedAgentMessage; /** * Caption and files are admitted as one relay transaction. The authenticated * route and optional source wire are fixed before bytes are retrieved, and no * outbound part is emitted until every file passes admission. A transport * failure after emission starts is durably uncertain and never blind-retried; * the managed agent receives one bounded NACK for the whole transaction. */ private handleManagedAgentAttachmentGroup; private managedAttachmentReplyWire; /** * One bounded NACK per wire: an unroutable or refused relay must be visible * to the authenticated agent, while its history replays stay quiet. NACK * delivery is best-effort — it must never make the failure worse. */ private nackManagedAgent; /** * Daemon contact resolution is case-exact, so a canonical config CID picked * by the sole-owner fallback is translated to the daemon-known contact form * when one exists. Last-inbound routes already carry the daemon form. */ private routableContact; private safeRelayMessage; private warnOwnerOfUnauthorizedSender; private effectiveOwners; private authorizationIntegrity; /** * Report what the session actually did with the prompt, not what the config * asked for. `interrupt: true` used to be reported as "your request * interrupted the previous task" unconditionally; the session now answers * whether anything was cancelled, whether the request is queued behind * earlier prompts, or whether it is held until the current task reaches a * safe stopping point. Backends that report no delivery state keep the old * queuedBehind-based wording. */ private acceptanceNotice; private complete; private commentsState; /** Model-authored commentary only; raw protocol/tool data never reaches here. */ private safeCommentary; private ownerAttachmentPrompt; private ownerPrompt; private outboxDir; private requestId; private send; private sendAttachments; /** Bound message size without splitting Unicode code points. */ private sendFinal; private wireId; private sender; private latestEventSeq; /** Map only event shape and allowlisted status to owner-safe phase text. */ private progressPhase; private watchLoop; /** * `recovered` distinguishes a first-ever start from unreadable persisted * diagnostics. Notification correctness does not depend on this state: every * establishment drains and then replays SDK hints from offset zero. */ private readWatchState; private writeWatchState; private errorText; private logError; }