import type { Message } from "@caupulican/pi-ai"; import type { LaneRecord } from "../autonomy/lane-tracker.ts"; import { type AgentBindingStatus, type WorkerRole } from "../orchestration/contracts.ts"; import type { SessionRootReply, SessionRootReplyQuery, SessionRootReplyWaitOptions, SessionRootReplyWaitResult } from "./session-root-mailbox.ts"; import type { WorkerTaskSessionView } from "./worker-task-view.ts"; export type WorkerAgentMessageKind = "steer" | "follow_up"; export type WorkerAgentTaskMetadata = { kind: "agent_turn"; dependsOnTaskIds?: readonly string[]; } | { kind: "terminal_handoff"; sourceAttemptId: string; }; export interface WorkerAgentMessage { messageId: string; kind: WorkerAgentMessageKind; content: string; senderAgentId?: string; threadId?: string; replyToMessageId?: string; expectReply?: boolean; /** Durable intent requiring one mailbox-correlated agent turn if no live turn consumes it. */ task?: WorkerAgentTaskMetadata; createdAt: string; deliveredAt?: string; repliedAt?: string; replyReceipt?: WorkerReplyReceipt; failedAt?: string; failureReason?: string; } export interface WorkerReplyReceipt { replyMessageId: string; requestSenderId: string; contentDigest: string; } export interface WorkerReplyAcknowledgement { messageId: string; acknowledgementId: string; replyContent: string; } export type WorkerAgentEnqueueReceipt = { status: "retained"; messageId: string; message: WorkerAgentMessage; created: boolean; } | { status: "completed_replay"; messageId: string; created: false; }; export interface WorkerAgentMailboxOptions { agentDir: string; parentSessionId: string; agentId: string; } export interface WorkerAgentTranscriptPage { agentId: string; cursor: number; messages: Message[]; nextCursor?: number; omittedMessages: number; serializedBytes: number; } export interface WorkerAgentMessageOptions { senderAgentId?: string; threadId?: string; expectReply?: boolean; /** Host-derived replay identity scoped to the caller, session, tool invocation, and action. */ idempotencyKey?: string; } export interface WorkerAgentBroadcastOptions { senderAgentId?: string; threadId?: string; expectReply?: boolean; /** Call-level replay identity; the coordinator derives one stable identity per canonical target. */ idempotencyKey: string; } export interface SessionRootWorkerAgentMessageOptions { threadId?: string; expectReply?: boolean; idempotencyKey?: string; } export type WorkerAgentReplyResult = { destination: "session_root"; messageId: string; } | { destination: "worker"; messageId: string; started: boolean; steering: boolean; record?: LaneRecord; skipReason?: string; }; export interface WorkerAgentControlScope { callerAgentId?: string; } export interface WorkerAgentTaskStartOptions extends WorkerAgentControlScope { /** Host-derived replay identity scoped to the caller, session, tool invocation, and action. */ idempotencyKey?: string; /** Existing same-objective durable tasks that must complete before this turn may run. */ dependsOnTaskIds?: readonly string[]; } export interface WorkerAgentTranscriptOptions extends WorkerAgentControlScope { cursor?: number; maxMessages?: number; /** Host-owned aggregate-envelope headroom; never accepted directly from a model argument. */ maxBytes?: number; } export type WorkerAgentActivity = "active" | "suspended" | "idle" | "unknown"; export type WorkerAgentWaitMode = "any" | "all"; export interface WorkerAgentWaitStatus { agentId: string; status: WorkerAgentActivity; } export interface WorkerAgentWaitResult { statuses: WorkerAgentWaitStatus[]; updatedAgentIds: string[]; timedOut: boolean; } export type WorkerAgentBroadcastTargetResult = { agentId: string; accepted: true; queued: true; replayed: boolean; messageId: string; } | { agentId: string; accepted: false; error: string; }; export interface WorkerAgentBroadcastResult { results: WorkerAgentBroadcastTargetResult[]; } /** Explicit model-facing projection. Durable resume, session, path, and resource data stay host-only. */ export interface WorkerAgentView { agentId: string; parentAgentId?: string; rootAgentId: string; depth: number; role: WorkerRole; /** Provider/model admitted for this persistent identity; it is retained across follow-up tasks. */ modelRef?: string; status: AgentBindingStatus; activity: WorkerAgentActivity; /** True when this caller may start/transcript/cancel the agent. Session-root lists are all true. */ controllable: boolean; createdAt: string; updatedAt: string; } export interface WorkerAgentRetireResult { agent: WorkerAgentView; retired: true; replayed: boolean; } /** One canonical host port for model-facing logical-agent controls. */ export interface WorkerAgentControlPort { /** Session-root read receipt for exact terminal generations; distinct from mutation review. */ observeWorkerTerminalRecords?(records: readonly LaneRecord[]): void; /** Observe only the bounded logical-agent identities exposed by a model-facing control result. */ observeWorkerAgentTerminals?(agentIds: readonly string[]): void; listWorkerAgents(scope?: WorkerAgentControlScope): WorkerAgentView[]; getWorkerTaskSessionView(): WorkerTaskSessionView; getWorkerAgentActivity(agentId: string, scope?: WorkerAgentControlScope): WorkerAgentActivity; readWorkerAgentTranscript(agentId: string, options?: WorkerAgentTranscriptOptions): WorkerAgentTranscriptPage; /** Queue-only session-peer delivery. This does not wake or steer an idle or active agent. */ sendWorkerAgentMessage(agentId: string, message: string, options?: WorkerAgentMessageOptions): { messageId: string; queued: true; }; /** Queue-only fan-out. Peer content is untrusted coordination evidence, never delegated authority. */ broadcastWorkerAgentMessage(agentIds: readonly string[], message: string, options: WorkerAgentBroadcastOptions): WorkerAgentBroadcastResult; /** Worker callers may wake or steer only themselves and descendants; the session root may target any agent. */ followUpWorkerAgent(agentId: string, message: string, options?: WorkerAgentMessageOptions): { started: boolean; steering: boolean; messageId: string; record?: LaneRecord; skipReason?: string; }; sendSessionRootWorkerAgentMessage(agentId: string, message: string, options?: SessionRootWorkerAgentMessageOptions): { messageId: string; queued: true; }; followUpSessionRootWorkerAgent(agentId: string, message: string, options?: SessionRootWorkerAgentMessageOptions): { started: boolean; steering: boolean; messageId: string; record?: LaneRecord; skipReason?: string; }; replyToWorkerAgentMessage(sourceAgentId: string, message: string, replyToMessageId: string): WorkerAgentReplyResult; listSessionRootReplies(query?: SessionRootReplyQuery): SessionRootReply[]; waitForSessionRootReplies(options?: SessionRootReplyWaitOptions): Promise; acknowledgeSessionRootReply(messageId: string, ackToken: string): boolean; reconcileSessionRootReplies(): void; /** Worker callers may start only themselves and descendants; the session root may target any agent. */ startWorkerAgentTask(agentId: string, message: string, options?: WorkerAgentTaskStartOptions): { started: boolean; steering: false; messageId: string; record?: LaneRecord; skipReason?: string; }; interruptWorkerAgent(agentId: string, scope?: WorkerAgentControlScope): { interrupted: boolean; reason?: string; }; resumeWorkerAgent(agentId: string, scope?: WorkerAgentControlScope): { started: boolean; record?: LaneRecord; skipReason?: string; }; cancelWorkerAgent(agentId: string, reasonCode?: string, scope?: WorkerAgentControlScope): LaneRecord | undefined; /** Retire one idle leaf without deleting its durable binding, lineage, transcript, or attempt history. */ retireWorkerAgent(agentId: string, scope?: WorkerAgentControlScope): WorkerAgentRetireResult; waitForWorkerAgent(agentId: string, timeoutMs?: number, scope?: WorkerAgentControlScope): Promise<{ status: WorkerAgentActivity; timedOut: boolean; }>; waitForWorkerAgents(agentIds: readonly string[], mode: WorkerAgentWaitMode, timeoutMs?: number, scope?: WorkerAgentControlScope): Promise; } /** Derive one target-fenced mailbox replay identity from a host-owned broadcast call identity. */ export declare function workerAgentBroadcastTargetIdempotencyKey(baseIdempotencyKey: string, agentId: string): string; /** Session-scoped replay identity. The coordinator fences one accepted id to exactly one target mailbox. */ export declare function workerAgentMessageId(parentSessionId: string, idempotencyKey: string): string; export declare function normalizeWorkerAgentDependencyTaskIds(value: unknown): readonly string[]; /** * Bounded durable inbox for a single logical worker agent. * * A message is never considered delivered merely because a controller read it. Its caller must * acknowledge it only after the exact corresponding user message has been appended to the child * WorkerConversation. The in-process subscription is a notification edge, not a polling loop. */ export declare class WorkerAgentMailbox { private readonly parentSessionId; private readonly agentId; private readonly file; private readonly listeners; constructor(options: WorkerAgentMailboxOptions); enqueue(input: { kind: WorkerAgentMessageKind; content: string; senderAgentId?: string; threadId?: string; replyToMessageId?: string; expectReply?: boolean; task?: WorkerAgentTaskMetadata; }): WorkerAgentMessage; /** Enqueue with exact creation evidence for an idempotent surrounding acceptance flow. */ enqueueWithReceipt(input: { kind: WorkerAgentMessageKind; content: string; senderAgentId?: string; threadId?: string; replyToMessageId?: string; expectReply?: boolean; task?: WorkerAgentTaskMetadata; idempotencyKey?: string; }): WorkerAgentEnqueueReceipt; pending(kind?: WorkerAgentMessageKind): WorkerAgentMessage[]; /** Oldest-first executable intents that have not reached the durable transcript boundary. */ pendingTaskBearing(): WorkerAgentMessage[]; acknowledgeDelivered(messageId: string): void; awaitingReplies(): WorkerAgentMessage[]; getMessage(messageId: string): WorkerAgentMessage | undefined; hasControlReplayReceipt(messageId: string): boolean; hasDeliveredControlReceipt(messageId: string): boolean; resolveCompletedReply(messageId: string, content: string): WorkerReplyReceipt | undefined; getReplyAcknowledgementId(messageId: string): string | undefined; listReplyAcknowledgements(): WorkerReplyAcknowledgement[]; /** * Mark one request replied while retaining the exact target receipt until transcript consumption. * The acknowledgement id is the durable target reply id, so a retry adopts a crash-left marker. */ beginReplyAcknowledgement(messageId: string, acknowledgementId: string, replyContent: string): boolean; /** Dead-letter only an ordinary executable turn with no reply or terminal-delivery obligation. */ deadLetterOrdinaryTask(messageId: string, reason: string): boolean; /** Commit one exact reply acknowledgement and release its protected history slot. */ commitReplyAcknowledgement(messageId: string, acknowledgementId: string): boolean; /** Roll back one exact reply acknowledgement without clearing a later or unrelated reply. */ rollbackReplyAcknowledgement(messageId: string, acknowledgementId: string): boolean; private markTimestamp; private finishReplyAcknowledgement; subscribe(listener: () => void): () => void; private notify; private read; private update; } //# sourceMappingURL=worker-agent-control.d.ts.map