Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
112 changes: 106 additions & 6 deletions extensions/gentle-agents.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,12 @@ const NATIVE_PARENT_WAKE_TEXT = "Review the delivered subagent output and contin
// How long a dispatched wake may take to start a parent run before a later
// delivery may send another one.
export const PARENT_WAKE_GRACE_MS = 30_000;
// How often a parent that is busy without a run is re-checked while content is
// held for it. See `armHoldRecheck`.
export const HOLD_RECHECK_MS = 1000;
// How long a user prompt seen at `input` suppresses the idle wake while it is
// still in its pre-run phase. See the `input` handler.
export const INPUT_PROMPT_GRACE_MS = 2000;
// A delivery or wake the host rejected while the parent was idle is retried
// after the grace at most this many times in a row, so a permanently broken
// runtime never retries forever.
Expand Down Expand Up @@ -613,7 +619,7 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
listener = sessionTransport.createListener(registry, sessionId, async (notification) => {
const active = activeSessionTransport;
if (!active || active.generation !== generation || active.sessionManager !== sessionManager || active.sessionId !== sessionId || sessions !== sessionManager || activeSessionId() !== sessionId) throw new Error("stale session transport");
pi.sendMessage({ customType: AGENTS_ORCHESTRATOR_MESSAGE_TYPE, content: `Session message from ${notification.senderSessionId} (correlation ${notification.id}): ${notification.message}`, display: true, details: { gentleAgents: { senderSessionId: notification.senderSessionId, recipientSessionId: sessionId, correlationId: notification.id, direction: "incoming" } } }, { deliverAs: "followUp", triggerTurn: true });
deliverOrchestratorMessage({ customType: AGENTS_ORCHESTRATOR_MESSAGE_TYPE, content: `Session message from ${notification.senderSessionId} (correlation ${notification.id}): ${notification.message}`, display: true, details: { gentleAgents: { senderSessionId: notification.senderSessionId, recipientSessionId: sessionId, correlationId: notification.id, direction: "incoming" } } }, sessionId);
});
client = sessionTransport.createClient(registry, sessionId);
if (sessions !== sessionManager || generation !== transportGeneration || activeSessionId() !== sessionId) {
Expand Down Expand Up @@ -783,6 +789,10 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
let promptStarting = false;
let cancelPromptGrace: (() => void) | undefined;
let cancelBoundaryFlush: (() => void) | undefined;
let cancelHoldRecheck: (() => void) | undefined;
// Completions sit in an opaque queue; this notes that one may be waiting.
// It is cleared by the next flush that reaches a deliverable route.
let completionsMaybePending = false;
// An idle parent has no turn boundary coming, so a failed idle delivery or
// wake is retried by one bounded timer that runs a full flush.
let cancelDeliveryRetry: (() => void) | undefined;
Expand Down Expand Up @@ -822,6 +832,11 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
// A parent that is busy without a run (compaction, or a prompt's pre-run
// compaction) is not streaming, so steer + triggerTurn would also start a
// direct turn. Its content stays queued ("hold") until a later boundary.
//
// Three sources use this one router: child completions, child messages
// (queries and notifications), and incoming orchestrator session messages.
// The route decision is also the single point where a receiver-side
// admission policy for session messages (#1518) would plug in.
type ParentRoute = "idle" | "run" | "hold";
// Throws for a missing or stale parent context, so delivery fails closed.
const parentRoute = (): ParentRoute => {
Expand All @@ -840,14 +855,36 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
cancelPromptGrace?.();
cancelPromptGrace = undefined;
};
const beginPromptStart = () => {
const beginPromptStart = (graceMs = PARENT_WAKE_GRACE_MS) => {
endPromptStart();
promptStarting = true;
cancelPromptGrace = deps.schedule(() => {
cancelPromptGrace = undefined;
promptStarting = false;
requestWake();
}, PARENT_WAKE_GRACE_MS);
}, graceMs);
};

const hasHeldContent = (): boolean =>
heldOrchestratorMessages.length > 0 || wakeOwed || completionsMaybePending || messages.pendingCount() > 0;
const stopHoldRecheck = () => {
cancelHoldRecheck?.();
cancelHoldRecheck = undefined;
};
// Events (`agent_*`, `turn_end`, `session_compact*`, `session_tree`) are the
// fast path that releases held content. They are not a guarantee: a busy
// state without a run can end without any event (a cancelled `/tree`
// summarization emits none) or be added to the host later. So while content
// is held for a parent that is busy without a run, one re-check is armed;
// it re-arms itself only while that state and the held content persist.
// Nothing is armed when nothing is held, and an idle or running parent is
// never polled.
const armHoldRecheck = () => {
if (cancelHoldRecheck || !hasHeldContent()) return;
cancelHoldRecheck = deps.schedule(() => {
cancelHoldRecheck = undefined;
flushAll();
}, HOLD_RECHECK_MS);
};

// Every wake requested during one synchronous delivery pass is coalesced
Expand All @@ -864,6 +901,7 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
try { route = parentRoute(); } catch { return; }
// A run in progress already carries the stored content; a parent busy
// without a run keeps the wake owed until a later boundary flush.
if (route === "hold") armHoldRecheck();
if (route !== "idle") return;
wakeOwed = false;
endDeliveryRetry();
Expand Down Expand Up @@ -894,9 +932,12 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
}
};

const sendToParent = (message: Parameters<ExtensionAPI["sendMessage"]>[0], route: Exclude<ParentRoute, "hold">) => {
// `runMode` is how a message joins a run in progress: child content steers
// (see the #867 rationale above); an orchestrator session message is a
// follow-up, so it never interrupts the turn the parent is working on.
const sendToParent = (message: Parameters<ExtensionAPI["sendMessage"]>[0], route: Exclude<ParentRoute, "hold">, runMode: "steer" | "followUp" = "steer") => {
if (route === "run") {
pi.sendMessage(message, { deliverAs: "steer", triggerTurn: true });
pi.sendMessage(message, { deliverAs: runMode, triggerTurn: true });
return;
}
pi.sendMessage(message, { triggerTurn: false });
Expand All @@ -911,6 +952,35 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
sendToParent({ customType: AGENTS_RESULT_TYPE, content: completionText(task), display: true, details: taskDetails(task) }, route);
};

// Incoming orchestrator session messages. The transport has already
// acknowledged the sender when the listener runs, so a message held for a
// busy parent that is then dropped by a session change is lost, exactly like
// pending child content. A stale or missing parent context throws instead,
// so the sender learns that delivery failed.
interface HeldOrchestratorMessage { message: Parameters<ExtensionAPI["sendMessage"]>[0]; recipientSessionId: string }
let heldOrchestratorMessages: HeldOrchestratorMessage[] = [];
const deliverOrchestratorMessage = (message: Parameters<ExtensionAPI["sendMessage"]>[0], recipientSessionId: string) => {
const route = parentRoute();
if (route === "hold") {
heldOrchestratorMessages.push({ message, recipientSessionId });
armHoldRecheck();
return;
}
sendToParent(message, route, "followUp");
};
const flushOrchestratorMessages = () => {
if (heldOrchestratorMessages.length === 0) return;
const route = deliveryRoute(false);
if (!route) return;
const held = heldOrchestratorMessages;
heldOrchestratorMessages = [];
for (const { message, recipientSessionId } of held) {
if (activeSessionId() !== recipientSessionId) continue;
try { sendToParent(message, route, "followUp"); }
catch { /* Best-effort delivery: at most once, even if forwarding fails. */ }
}
};

// A stale report must not re-enter the LLM conversation, but the parent is
// told never to poll: it gets a compact model-facing notice naming the task,
// and the human gets durable TUI-only content.
Expand Down Expand Up @@ -939,7 +1009,12 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
const deliveryRoute = (rethrow: boolean): Exclude<ParentRoute, "hold"> | undefined => {
try {
const route = parentRoute();
return route === "hold" ? undefined : route;
if (route === "hold") {
armHoldRecheck();
return undefined;
}
stopHoldRecheck();
return route;
} catch (error) {
if (rethrow) throw error;
return undefined;
Expand All @@ -962,6 +1037,7 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
const flushCompletions = () => {
const route = deliveryRoute(false);
if (!route) return;
completionsMaybePending = false;
for (const entry of completions.takeDeliverable(deps.now())) {
try {
if (entry.stale) deliverStale(entry.task, entry.settledAt, route);
Expand All @@ -971,12 +1047,15 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
// lose the completion: it waits for the next boundary instead, and
// an idle parent, which has none coming, gets a bounded retry.
completions.requeue(entry);
// Still pending, so a later held-only re-check must see it.
completionsMaybePending = true;
if (route === "idle") armDeliveryRetry();
}
}
};

const flushAll = () => {
flushOrchestratorMessages();
flushMessages();
flushCompletions();
// A wake left owed by an earlier held or expired attempt is retried here.
Expand All @@ -1003,7 +1082,10 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
wakeVisibilityWarning = false;
parentRunActive = false;
wakeOwed = false;
heldOrchestratorMessages = [];
completionsMaybePending = false;
endPromptStart();
stopHoldRecheck();
cancelBoundaryFlush?.();
cancelBoundaryFlush = undefined;
endDeliveryRetry();
Expand All @@ -1017,6 +1099,7 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
const settleCompletion = (task: TaskRecord) => {
messages.invalidateTask(task.id);
completions.enqueue(task, deps.now());
completionsMaybePending = true;
if (activeAgentRuns === 0) flushAll();
};

Expand Down Expand Up @@ -1056,6 +1139,23 @@ export default function gentleAgents(pi: ExtensionAPI, env: NodeJS.ProcessEnv =
pi.on("turn_end", () => flushAll());
pi.on("session_compact", scheduleBoundaryFlush);
pi.on("session_compact_failed", scheduleBoundaryFlush);
pi.on("session_tree", scheduleBoundaryFlush);
// A user prompt is "starting" from `input`, through input handlers, auth and
// model checks, until `before_agent_start`, while the host still reports
// idle. A wake dispatched in that window races the prompt into
// Agent.prompt(): the user's prompt can be rejected as already processing.
// The starting prompt carries the stored content itself, so no wake is
// needed. An input another extension handles never starts a run and emits
// no event, so this suppression expires quickly (INPUT_PROMPT_GRACE_MS) and
// the owed wake then goes out; a pre-run phase longer than that only
// re-opens the old window. A pre-prompt compaction makes the route "hold",
// and `before_agent_start` re-arms the starting state. A start window that
// is already open is never shortened: a dispatched wake reaches `input`
// too, and its own PARENT_WAKE_GRACE_MS window must keep a second wake out.
pi.on("input", (event) => {
if (parentRunActive || promptStarting || event.streamingBehavior !== undefined) return;
beginPromptStart(INPUT_PROMPT_GRACE_MS);
});

const runner = new AgentRunner(store, loadAgentsConfig({ cwd: process.cwd(), home: deps.home, agentHome }), deps, {
askUser: (_taskId, ask, raw) => answerThroughUi(ui, ask, raw),
Expand Down
Loading