fix(coding-agent): restored side-channel IRC auto-reply for busy recipients when async execution is disabled

A subagent's send await:true to Main during a blocking task spawn was a structural deadlock: deliverIrcMessage queues mid-turn messages as step-boundary asides, but Main's next boundary requires the sender's own batch to finish, so the sender always burned the full irc.timeoutMs. Awaited sends now pass expectsReply through IrcBus.send; a mid-turn recipient with async.enabled off generates an ephemeral no-tools reply via runEphemeralTurn, records an irc:autoreply aside in its own history, and delivers it back over the bus with replyTo threading so the sender's waiter resolves.
This commit is contained in:
can1357
2026-06-11 18:12:22 +02:00
parent abc152e5ca
commit 926dfdea39
10 changed files with 185 additions and 21 deletions
+5 -4
View File
@@ -11,6 +11,7 @@
- `packages/coding-agent/src/registry/agent-lifecycle.ts` — revival of parked recipients on direct send.
- `packages/coding-agent/src/session/agent-session.ts` — `deliverIrcMessage(...)`: recipient-side injection and wake turns.
- `packages/coding-agent/src/prompts/system/irc-incoming.md` — incoming-message rendering for the recipient.
- `packages/coding-agent/src/prompts/system/irc-autoreply.md` — prompt for the ephemeral auto-reply side turn (busy recipient, async disabled).
- `packages/coding-agent/src/config/settings-schema.ts` — `irc.timeoutMs`.
- `packages/coding-agent/src/modes/controllers/event-controller.ts` — renders IRC events into chat UI.
@@ -43,11 +44,11 @@
4. `op: "send"` validates `to`/`message`, rejects self-sends, and rejects `await` with `to: "all"`.
5. Target resolution: broadcasts fan out to `registry.listVisibleTo(senderId)` (live peers only — `running`/`idle`; reviving every parked agent on a broadcast would be a stampede). Direct sends go through the bus unfiltered, so a parked recipient is revived.
6. `IrcBus.send(...)` is fire-and-forget — it never blocks on the recipient generating anything. Delivery by recipient status:
- `running` → message enqueued and injected as a non-interrupting aside at the recipient's next step boundary (`AgentSession.deliverIrcMessage`, rendered from `irc-incoming.md`, persisted as an `irc:incoming` custom message) — receipt `injected`;
- `running` → message enqueued and injected as a non-interrupting aside at the recipient's next step boundary (`AgentSession.deliverIrcMessage`, rendered from `irc-incoming.md`, persisted as an `irc:incoming` custom message) — receipt `injected`. If the sender awaits a reply (`expectsReply` from `await: true`) and the recipient has `async.enabled` off, the recipient also generates an ephemeral no-tools auto-reply (`runEphemeralTurn`, the `/btw` pipeline) and sends it back over the bus with `replyTo` set, recording an `irc:autoreply` aside in its own history — a recipient blocked in a synchronous task spawn can never reach a step boundary before the sender's timeout otherwise;
- `idle` (live session) → enqueued and a real turn is started — the message wakes the agent — receipt `woken`;
- `parked` → `AgentLifecycleManager.global().ensureLive(to)` revives the session first, then the wake path — receipt `revived`;
- resolution/revival failure → receipt `failed` with the error; other recipients still complete.
7. `send` with `await: true` then calls `IrcBus.wait(senderId, { from: to }, timeoutMs, signal)` and appends the reply (or a no-reply note suggesting `inbox`/`wait`) to the result.
7. `send` with `await: true` then calls `IrcBus.wait(senderId, { from: to }, timeoutMs, signal)` and appends the reply (or a no-reply note suggesting `inbox`/`wait`) to the result. Awaited sends pass `{ expectsReply: true }` to `IrcBus.send` so a busy recipient can auto-reply (see step 6).
8. `op: "wait"` blocks until a message for the caller (optionally filtered by `from`) arrives, consumes it, and returns it. Timeout returns a clean "no message" result, not an error.
9. `op: "inbox"` drains pending messages (or peeks with `peek: true`) without blocking.
10. Timeouts resolve as `params.timeoutMs ?? irc.timeoutMs`, normalized: `0` disables the timeout, negative/non-finite values fall back to the default `120_000`, positive values are truncated and clamped to ≥ 1 ms.
@@ -56,7 +57,7 @@
- `list`: enumerate peers with status (`running`/`idle`/`parked`), unread counts, and last activity.
- `send` direct: one exact peer id; wakes idle peers, revives parked ones.
- `send` broadcast: `to: "all"` to every live peer; parked peers are skipped.
- `send` + `await: true`: round-trip convenience — send, then wait for the next message from that peer. Replaces the old `awaitReply` auto-reply semantics without a fake reply.
- `send` + `await: true`: round-trip convenience — send, then wait for the next message from that peer. Marks the send `expectsReply`, enabling the busy-recipient auto-reply path when async execution is disabled.
- `wait`: block for an incoming message, optionally filtered by sender.
- `inbox`: non-blocking drain or peek.
@@ -93,7 +94,7 @@
## Notes
- This is IRC-like naming only: no servers, sockets, channels, or join/part state. Addressing is by exact registry agent id.
- Replies are real turns by the recipient — the old ephemeral no-tools auto-reply (`awaitReply` / `respondAsBackground`) no longer exists. A recipient may keep working before answering; check `inbox` or `wait` again rather than re-sending.
- Replies are real turns by the recipient, with one exception: an awaited send to a mid-turn recipient with `async.enabled` off triggers an ephemeral no-tools auto-reply (the old `respondAsBackground` path), because a recipient blocked in a synchronous task spawn whose batch includes the sender can never run a real turn before the sender's timeout. A recipient may otherwise keep working before answering; check `inbox` or `wait` again rather than re-sending.
- Wake-on-message is the only resume primitive: messaging a parked agent revives it (same `ensureLive` path as the Agent Hub). The task tool has no `resume` parameter.
- Message ids are Snowflakes; pass them as `replyTo` to thread an answer to a specific message.
- Persistence is per recipient history: the sender gets receipts in the tool result; the recipient sees the injected `irc:incoming` message in its own transcript (visible via `history://<id>`).
+1
View File
@@ -17,6 +17,7 @@
- Routed LSP hover code rendering through the shared cached highlighter instead of calling the native highlighter directly on each render.
- Kept smooth assistant streaming renders transient until `message_end` so catch-up frames do not synchronously re-highlight still-growing code blocks.
- Fixed the `job` tool's poll header repeating the running count ("waiting on 19 of 19 19 running"): the title now reads "waiting on N jobs" (or "waiting on N of M jobs" when some settled) and the dim meta lists only settled categories (done/failed/cancelled)
- Fixed `irc` `send await:true` to a busy peer stranding the sender until timeout when async execution is disabled: the recipient (e.g. Main blocked in a synchronous task spawn awaiting the sender's own batch) can never reach a step boundary to reply, so it now generates an ephemeral side-channel auto-reply from its context (the pre-mailbox `respondAsBackground` behavior), records it as an `irc:autoreply` aside in its own history, and delivers it back over the bus with `replyTo` threading
## [15.11.1] - 2026-06-11
### Added
+14 -3
View File
@@ -7,7 +7,11 @@
* AgentLifecycleManager, idle agents are woken with a real turn, and busy
* agents receive the message as a non-interrupting aside at the next step
* boundary (see AgentSession.deliverIrcMessage). Replies are real turns by
* the recipient, observed via `wait`.
* the recipient, observed via `wait` — with one exception: when the sender
* awaits a reply and the recipient is mid-turn with async execution
* disabled, the recipient session generates an ephemeral side-channel
* auto-reply (it may be blocked in a synchronous task spawn whose batch
* includes the sender, so a real turn could never happen in time).
*/
import { logger, Snowflake } from "@oh-my-pi/pi-utils";
@@ -80,8 +84,15 @@ export class IrcBus {
* context, so buffering it too would double-deliver via a later
* `wait`/`inbox` and inflate unread counts. Only a failed live hand-off
* is buffered for the recipient to drain later.
*
* `opts.expectsReply` marks sends whose caller is blocked on an answer
* (`send await:true`). It is forwarded to the recipient session so a
* mid-turn recipient that cannot reach a step boundary (async execution
* disabled — e.g. blocked in a synchronous task spawn awaiting the
* sender's own batch) can generate an ephemeral side-channel auto-reply
* instead of stranding the sender until timeout.
*/
async send(msg: Omit<IrcMessage, "id" | "ts">): Promise<IrcDeliveryReceipt> {
async send(msg: Omit<IrcMessage, "id" | "ts">, opts?: { expectsReply?: boolean }): Promise<IrcDeliveryReceipt> {
const message: IrcMessage = { ...msg, id: Snowflake.next(), ts: Date.now() };
const ref = this.#registry.get(message.to);
if (!ref || ref.status === "aborted") {
@@ -118,7 +129,7 @@ export class IrcBus {
}
try {
const delivery = await session.deliverIrcMessage(message);
const delivery = await session.deliverIrcMessage(message, opts);
this.#relayToMainUi(message);
return { to: message.to, outcome: revived ? "revived" : delivery };
} catch (error) {
@@ -191,7 +191,11 @@ export class UiHelpers {
this.ctx.chatContainer.addChild(component);
break;
}
if (message.customType === "irc:incoming" || message.customType === "irc:relay") {
if (
message.customType === "irc:incoming" ||
message.customType === "irc:autoreply" ||
message.customType === "irc:relay"
) {
const details = (
message as CustomMessage<{
from?: string;
@@ -201,13 +205,18 @@ export class UiHelpers {
replyTo?: string;
}>
).details;
const incoming = message.customType === "irc:incoming";
const kind =
message.customType === "irc:incoming"
? ("incoming" as const)
: message.customType === "irc:autoreply"
? ("autoreply" as const)
: ("relay" as const);
const card = createIrcMessageCard(
{
kind: incoming ? "incoming" : "relay",
kind,
from: details?.from,
to: details?.to,
body: incoming ? details?.message : details?.body,
body: kind === "incoming" ? details?.message : details?.body,
replyTo: details?.replyTo,
timestamp: message.timestamp,
},
@@ -0,0 +1,6 @@
<irc>
You received an IRC message from agent `{{from}}`{{#if replyTo}} (replying to {{replyTo}}){{/if}} while you are busy mid-task. This is a side-channel turn: reply briefly and directly using the conversation context already available to you. NEVER call tools. The text you write is delivered back to `{{from}}` as your answer.
Message:
{{message}}
</irc>
@@ -3,5 +3,5 @@ Incoming IRC message from agent `{{from}}`{{#if replyTo}} (replying to {{replyTo
{{message}}
If a response is expected, reply with the `irc` tool (`op: "send"`, `to: "{{from}}"`) — you may finish your current step first. Nobody replies on your behalf.
{{#if autoReplied}}You are mid-task, so a side-channel auto-reply was generated from your context and delivered to `{{from}}` on your behalf (recorded after this message). Follow up with the `irc` tool (`op: "send"`, `to: "{{from}}"`) only if that auto-reply needs correcting.{{else}}If a response is expected, reply with the `irc` tool (`op: "send"`, `to: "{{from}}"`) — you may finish your current step first. Nobody replies on your behalf.{{/if}}
</irc>
@@ -9,7 +9,7 @@ Sends short text messages to other agents in this process and receives theirs.
- `op: "wait"` — block until a message arrives (optionally only `from` a specific peer); consumes and returns it. A timeout is a clean "no message" result, not an error.
- `op: "inbox"` — drain pending messages without blocking (`peek: true` to leave them unread).
- `replyTo` — set it to the id of the message you are answering so the sender can correlate.
- Nobody answers on a peer's behalf anymore: a reply only arrives when the recipient actually sends one. For background on what a peer has been doing, `read` `history://<id>` instead of interrogating them.
- Nobody answers on a peer's behalf — a reply normally arrives only when the recipient sends one — with one exception: `send` with `await: true` to a peer that is mid-turn and cannot reach a step boundary (async execution disabled, e.g. blocked in a synchronous task spawn) gets a side-channel auto-reply generated from that peer's context. For background on what a peer has been doing, `read` `history://<id>` instead of interrogating them.
</instruction>
<when_to_use>
@@ -171,7 +171,7 @@ import { GoalRuntime } from "../goals/runtime";
import type { Goal, GoalModeState } from "../goals/state";
import type { HindsightSessionState } from "../hindsight/state";
import { type LocalProtocolOptions, resolveLocalUrlToPath } from "../internal-urls";
import type { IrcMessage } from "../irc/bus";
import { IrcBus, type IrcMessage } from "../irc/bus";
import { resolveMemoryBackend } from "../memory-backend";
import { getMnemopiSessionState, type MnemopiSessionState, setMnemopiSessionState } from "../mnemopi/state";
import { containsOrchestrate, ORCHESTRATE_NOTICE } from "../modes/orchestrate";
@@ -185,6 +185,7 @@ import type { PlanModeState } from "../plan-mode/state";
import autoContinuePrompt from "../prompts/system/auto-continue.md" with { type: "text" };
import eagerTodoPrompt from "../prompts/system/eager-todo.md" with { type: "text" };
import emptyStopRetryTemplate from "../prompts/system/empty-stop-retry.md" with { type: "text" };
import ircAutoReplyTemplate from "../prompts/system/irc-autoreply.md" with { type: "text" };
import ircIncomingTemplate from "../prompts/system/irc-incoming.md" with { type: "text" };
import planModeActivePrompt from "../prompts/system/plan-mode-active.md" with { type: "text" };
import planModeReferencePrompt from "../prompts/system/plan-mode-reference.md" with { type: "text" };
@@ -9151,11 +9152,20 @@ export class AgentSession {
* → "woken".
*
* Never blocks on the recipient's turn: the wake turn is fire-and-forget.
*
* When the sender expects a reply (`send await:true`) and this session is
* mid-turn with async execution disabled, the next step boundary may be
* gated on the sender's own batch finishing (blocking task spawns), so a
* real reply turn can never happen in time. In that case an ephemeral
* side-channel auto-reply is generated from the current context (the old
* `respondAsBackground` path) and sent back over the bus on this agent's
* behalf.
*/
async deliverIrcMessage(msg: IrcMessage): Promise<"injected" | "woken"> {
async deliverIrcMessage(msg: IrcMessage, opts?: { expectsReply?: boolean }): Promise<"injected" | "woken"> {
if (this.#isDisposed) {
throw new Error("Recipient session is disposed.");
}
const autoReply = (opts?.expectsReply ?? false) && this.isStreaming && !this.settings.get("async.enabled");
const record: CustomMessage = {
role: "custom",
customType: "irc:incoming",
@@ -9163,6 +9173,7 @@ export class AgentSession {
from: msg.from,
message: msg.body,
replyTo: msg.replyTo ?? "",
autoReplied: autoReply,
}),
display: true,
details: { id: msg.id, from: msg.from, message: msg.body, ...(msg.replyTo ? { replyTo: msg.replyTo } : {}) },
@@ -9172,6 +9183,7 @@ export class AgentSession {
void this.#emitSessionEvent({ type: "irc_message", message: record });
if (this.isStreaming) {
this.#pendingIrcAsides.push(record);
if (autoReply) void this.#runIrcAutoReply(msg);
return "injected";
}
// Idle: same wake primitive the yield queue uses for async-result
@@ -9182,6 +9194,50 @@ export class AgentSession {
return "woken";
}
/**
* Generate and deliver an ephemeral auto-reply to `msg` on this agent's
* behalf: a no-tools side-channel turn over the current history (same
* pipeline as `/btw`), recorded into this session as an `irc:autoreply`
* aside so the model knows what was said for it, and sent back to the
* sender as a regular bus message (`replyTo: msg.id`) so their parked
* `wait`/`await:true` resolves. Failures only log — the sender then hits
* its normal wait timeout.
*/
async #runIrcAutoReply(msg: IrcMessage): Promise<void> {
try {
const { replyText } = await this.runEphemeralTurn({
promptText: prompt.render(ircAutoReplyTemplate, {
from: msg.from,
message: msg.body,
replyTo: msg.replyTo ?? "",
}),
});
const body = replyText.trim();
if (!body || this.#isDisposed) return;
const record: CustomMessage = {
role: "custom",
customType: "irc:autoreply",
content: `[IRC you → \`${msg.from}\` (auto)]\n\n${body}`,
display: true,
details: { to: msg.from, body, replyTo: msg.id },
attribution: "agent",
timestamp: Date.now(),
};
void this.#emitSessionEvent({ type: "irc_message", message: record });
// Asides drain at the next step boundary; anything left over is
// flushed at the start of the next prompt (#flushPendingIrcAsides).
this.#pendingIrcAsides.push(record);
// `from` must be the id the sender addressed (msg.to) so their
// from-filtered waiter matches.
const receipt = await IrcBus.global().send({ from: msg.to, to: msg.from, body, replyTo: msg.id });
if (receipt.outcome === "failed") {
logger.warn("IRC auto-reply delivery failed", { to: msg.from, error: receipt.error });
}
} catch (error) {
logger.warn("IRC auto-reply turn failed", { from: msg.from, error: String(error) });
}
}
/**
* Emit an IRC relay observation event on this session for UI rendering only.
* Does not persist the record to history. Called by the IrcBus to surface
+16 -4
View File
@@ -234,7 +234,15 @@ export class IrcTool implements AgentTool<typeof ircSchema, IrcDetails> {
// through the bus unfiltered so parked recipients are revived.
const targets = isBroadcast ? registry.listVisibleTo(senderId).map(ref => ref.id) : [to];
const receipts = await Promise.all(
targets.map(target => bus.send({ from: senderId, to: target, body: message, replyTo: params.replyTo })),
targets.map(target =>
bus.send(
{ from: senderId, to: target, body: message, replyTo: params.replyTo },
// Awaited sends mark the sender as blocked on an answer so a
// busy recipient that cannot reach a step boundary (async
// disabled) auto-replies instead of stranding the sender.
params.await ? { expectsReply: true } : undefined,
),
),
);
const lines: string[] = [];
@@ -457,13 +465,14 @@ function callMeta(args: IrcRenderArgs | undefined): string[] {
/**
* Display-only transcript card for live IRC traffic: `irc:incoming` DMs
* delivered to this session and `irc:relay` observations of agent↔agent
* delivered to this session, `irc:autoreply` side-channel replies sent on
* this session's behalf, and `irc:relay` observations of agent↔agent
* traffic. Shares the tool renderer's glyph + quote-border conventions so
* cards and `irc` tool output look identical in the transcript.
*/
export function createIrcMessageCard(
card: {
kind: "incoming" | "relay";
kind: "incoming" | "autoreply" | "relay";
from?: string;
to?: string;
body?: string;
@@ -477,9 +486,12 @@ export function createIrcMessageCard(
const title =
card.kind === "incoming"
? `IRC ${uiTheme.nav.back} ${from}`
: `IRC ${from} ${uiTheme.nav.selected} ${card.to?.trim() || "?"}`;
: card.kind === "autoreply"
? `IRC ${uiTheme.nav.selected} ${card.to?.trim() || "?"}`
: `IRC ${from} ${uiTheme.nav.selected} ${card.to?.trim() || "?"}`;
const body = card.body ?? "";
const meta: string[] = [];
if (card.kind === "autoreply") meta.push("auto");
if (card.replyTo) meta.push("reply");
const age = messageAge(card.timestamp);
if (age) meta.push(age);
+70 -2
View File
@@ -1,6 +1,7 @@
import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
import { Agent } from "@oh-my-pi/pi-agent-core";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import type { SettingPath } from "@oh-my-pi/pi-coding-agent/config/settings-schema";
import { IrcBus, type IrcMessage } from "@oh-my-pi/pi-coding-agent/irc/bus";
import { AgentLifecycleManager } from "@oh-my-pi/pi-coding-agent/registry/agent-lifecycle";
import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry";
@@ -73,7 +74,10 @@ function makeToolSession(registry: AgentRegistry, agentId: string): ToolSession
};
}
function createRealSession(): { session: AgentSession; sessionManager: SessionManager } {
function createRealSession(overrides: Partial<Record<SettingPath, unknown>> = {}): {
session: AgentSession;
sessionManager: SessionManager;
} {
const sessionManager = SessionManager.inMemory("/tmp");
const session = new AgentSession({
agent: new Agent({
@@ -84,7 +88,7 @@ function createRealSession(): { session: AgentSession; sessionManager: SessionMa
},
}),
sessionManager,
settings: Settings.isolated({ "compaction.enabled": false }),
settings: Settings.isolated({ "compaction.enabled": false, ...overrides }),
modelRegistry: {} as never,
});
return { session, sessionManager };
@@ -602,5 +606,69 @@ describe("IRC", () => {
expect(outcome).toBe("injected");
expect(promptSpy).not.toHaveBeenCalled();
});
it("auto-replies via an ephemeral side turn when the sender awaits and async execution is disabled", async () => {
const { session } = createRealSession({ "async.enabled": false });
sessions.push(session);
registry.register({ id: "Main", displayName: "main", kind: "main", session });
const sub = makeFakeSession();
registry.register({ id: "0-Sub", displayName: "task", kind: "sub", session: sub.session });
Object.defineProperty(session, "isStreaming", { value: true, configurable: true });
const ephemeralSpy = vi
.spyOn(session, "runEphemeralTurn")
.mockResolvedValue({ replyText: "auto answer", assistantMessage: {} as never });
const autoReplyEvent = new Promise<CustomMessage>(resolve => {
session.subscribe(event => {
if (event.type === "irc_message" && event.message.customType === "irc:autoreply") {
resolve(event.message);
}
});
});
// The sender parks a waiter (the `await: true` path), then sends with
// the expectsReply hint — exactly what the irc tool does.
const waiting = bus.wait("0-Sub", { from: "Main" }, 1000);
const receipt = await bus.send(
{ from: "0-Sub", to: "Main", body: "which PR did you mean?" },
{ expectsReply: true },
);
expect(receipt).toEqual({ to: "Main", outcome: "injected" });
// The side-channel reply resolves the sender's waiter as a real bus
// message threaded to the original send.
const reply = await waiting;
expect(reply?.from).toBe("Main");
expect(reply?.body).toBe("auto answer");
expect(reply?.replyTo).toBeTruthy();
expect(ephemeralSpy.mock.calls[0]?.[0]?.promptText).toContain("which PR did you mean?");
// The recipient records what was said on its behalf.
const record = await autoReplyEvent;
expect(record.details).toMatchObject({ to: "0-Sub", body: "auto answer" });
});
it("does not auto-reply when async execution is enabled or the sender does not await", async () => {
const enabled = createRealSession({ "async.enabled": true });
sessions.push(enabled.session);
registry.register({ id: "Main", displayName: "main", kind: "main", session: enabled.session });
Object.defineProperty(enabled.session, "isStreaming", { value: true, configurable: true });
const enabledSpy = vi
.spyOn(enabled.session, "runEphemeralTurn")
.mockResolvedValue({ replyText: "nope", assistantMessage: {} as never });
const awaited = await bus.send({ from: "0-Sub", to: "Main", body: "q?" }, { expectsReply: true });
expect(awaited.outcome).toBe("injected");
expect(enabledSpy).not.toHaveBeenCalled();
const disabled = createRealSession({ "async.enabled": false });
sessions.push(disabled.session);
registry.register({ id: "Main2", displayName: "main", kind: "main", session: disabled.session });
Object.defineProperty(disabled.session, "isStreaming", { value: true, configurable: true });
const disabledSpy = vi
.spyOn(disabled.session, "runEphemeralTurn")
.mockResolvedValue({ replyText: "nope", assistantMessage: {} as never });
const fireAndForget = await bus.send({ from: "0-Sub", to: "Main2", body: "fyi" });
expect(fireAndForget.outcome).toBe("injected");
expect(disabledSpy).not.toHaveBeenCalled();
});
});
});