Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
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
Prev Previous commit
feat(opencode): coordinator-messaging — sibling/coordinator subagent …
…communication

Lets the sibling subagents one parent spawns message each other directly
instead of routing everything through the parent. The surface is one
composable primitive — a per-child allow-list (`message_allow`) — so the
parent builds whatever graph it wants: hub, mesh, or chain.

Peer sends are fire-and-forget; the synchronous round-trip stays parent-only,
which is what keeps siblings deadlock-free. Slugs are parent-owned handles
resolved through a registry, authorized at send time against both the
allow-list and true sibling-hood. Recipients drain a FIFO inbox at their own
turn boundary, batched so an M-member graph costs O(1) turns per drain.

The interrupt tools also accept a slug task_id, resolved through the same
registry and still subject to the full ancestry check.

Off by default behind OPENCODE_EXPERIMENTAL_AGENT_MESSAGING.
  • Loading branch information
iceteaSA committed Aug 22, 2026
commit 13af75e32260f9057798e57d4fa0f43bf05f4c30
4 changes: 3 additions & 1 deletion packages/core/src/cross-spawn-spawner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import * as Effect from "effect/Effect"
import * as Exit from "effect/Exit"
import * as FileSystem from "effect/FileSystem"
import * as Layer from "effect/Layer"
import { LayerNode } from "./effect/layer-node"
import * as Path from "effect/Path"
import * as PlatformError from "effect/PlatformError"
import * as Predicate from "effect/Predicate"
Expand Down Expand Up @@ -497,11 +498,12 @@ export const make = Effect.gen(function* () {
return makeSpawner(spawnCommand)
})

const layer: Layer.Layer<ChildProcessSpawner, never, FileSystem.FileSystem | Path.Path> = Layer.effect(
export const layer: Layer.Layer<ChildProcessSpawner, never, FileSystem.FileSystem | Path.Path> = Layer.effect(
ChildProcessSpawner,
make,
)

export const node = makeGlobalNode({ service: ChildProcessSpawner, layer, deps: [filesystem, path] })
export const defaultLayer = layer

export * as CrossSpawnSpawner from "./cross-spawn-spawner"
3 changes: 2 additions & 1 deletion packages/opencode/src/agent/agent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ export class Service extends Context.Service<Service, Interface>()("@opencode/Ag

export const use = serviceUse(Service)

const layer = Layer.effect(
export const layer = Layer.effect(
Service,
Effect.gen(function* () {
const config = yield* Config.Service
Expand Down Expand Up @@ -449,5 +449,6 @@ export const node = LayerNode.make({
layer: layer,
deps: [Config.node, Auth.node, Plugin.node, Skill.node, Provider.node, locationServiceMapNode],
})
export const defaultLayer = layer

export * as Agent from "./agent"
3 changes: 2 additions & 1 deletion packages/opencode/src/config/config.ts
Original file line number Diff line number Diff line change
Expand Up @@ -172,7 +172,7 @@ function writableGlobal(info: Info) {
return next
}

const layer = Layer.effect(
export const layer = Layer.effect(
Service,
Effect.gen(function* () {
const fs = yield* FSUtil.Service
Expand Down Expand Up @@ -677,5 +677,6 @@ export const node = LayerNode.make({
layer: layer,
deps: [FSUtil.node, Auth.node, Account.node, Env.node, Npm.node, httpClient],
})
export const defaultLayer = layer

export * as Config from "./config"
1 change: 1 addition & 0 deletions packages/opencode/src/event-v2-bridge.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,5 +67,6 @@ const layer = Layer.effect(
)

export const node = LayerNode.make({ service: Service, layer: layer, deps: [EventV2.node] })
export const defaultLayer = layer

export * as EventV2Bridge from "./event-v2-bridge"
156 changes: 155 additions & 1 deletion packages/opencode/src/messaging/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,10 @@ export const Rejected = MessagingEvent.Rejected

const ROUND_TRIP_CAP = 8
const DEFAULT_TIMEOUT = Duration.seconds(300)
export const INBOX_OUTBOUND_BUDGET = 20
export const INBOX_CAP = 50
export const DEDUP_WINDOW = 100
export const TREE_MESSAGE_CAP = 2000

export class RejectedError extends Schema.TaggedErrorClass<RejectedError>()("Messaging.RejectedError", {}) {
override get message() {
Expand Down Expand Up @@ -46,9 +50,23 @@ interface ChildCounters {
roundTrips: number
}

export interface InboxItem {
from: SessionID
fromSlug: string
body: string
time: number
}

interface State {
pending: Map<SessionID, PendingReply>
counters: Map<SessionID, ChildCounters>
registry: Map<string, SessionID>
allow: Map<SessionID, string[]>
inbox: Map<SessionID, InboxItem[]>
dedup: Map<SessionID, string[]>
outbound: Map<SessionID, number>
waiters: Map<SessionID, Deferred.Deferred<void>>
treeTotal: { count: number }
}

export interface Interface {
Expand All @@ -67,6 +85,22 @@ export interface Interface {
}) => Effect.Effect<void, NotFoundError>
readonly reject: (childSessionID: SessionID) => Effect.Effect<void, NotFoundError>
readonly list: () => Effect.Effect<ReadonlyArray<PendingReply>>
readonly registerSlug: (slug: string, sessionID: SessionID) => Effect.Effect<void>
readonly resolveSlug: (slug: string) => Effect.Effect<Option.Option<SessionID>>
readonly setAllow: (sessionID: SessionID, slugs: string[]) => Effect.Effect<void>
readonly getAllow: (sessionID: SessionID) => Effect.Effect<string[]>
readonly slugFor: (sessionID: SessionID) => Effect.Effect<string>
readonly enqueue: (input: {
target: SessionID
from: SessionID
fromSlug: string
body: string
}) => Effect.Effect<void, AbuseError | NotFoundError>
readonly drain: (sessionID: SessionID) => Effect.Effect<ReadonlyArray<InboxItem>>
readonly awaitInbox: (
sessionID: SessionID,
opts: { timeoutMs: number },
) => Effect.Effect<boolean>
}

export class Service extends Context.Service<Service, Interface>()("@opencode/Messaging") {}
Expand All @@ -80,6 +114,13 @@ export const layer = Layer.effect(
const state: State = {
pending: new Map<SessionID, PendingReply>(),
counters: new Map<SessionID, ChildCounters>(),
registry: new Map<string, SessionID>(),
allow: new Map<SessionID, string[]>(),
inbox: new Map<SessionID, InboxItem[]>(),
dedup: new Map<SessionID, string[]>(),
outbound: new Map<SessionID, number>(),
waiters: new Map<SessionID, Deferred.Deferred<void>>(),
treeTotal: { count: 0 },
}
yield* Effect.addFinalizer(() =>
Effect.gen(function* () {
Expand All @@ -88,6 +129,13 @@ export const layer = Layer.effect(
}
state.pending.clear()
state.counters.clear()
state.registry.clear()
state.allow.clear()
state.inbox.clear()
state.dedup.clear()
state.outbound.clear()
state.waiters.clear()
state.treeTotal.count = 0
}),
)
return state
Expand Down Expand Up @@ -193,10 +241,116 @@ export const layer = Layer.effect(
return Array.from(value.pending.values())
})

return Service.of({ send, reply, reject, list })
const registerSlug: Interface["registerSlug"] = Effect.fn("Messaging.registerSlug")(function* (slug, sessionID) {
const value = yield* InstanceState.get(state)
value.registry.set(slug, sessionID)
if (!value.inbox.has(sessionID)) value.inbox.set(sessionID, [])
})

const resolveSlug: Interface["resolveSlug"] = Effect.fn("Messaging.resolveSlug")(function* (slug) {
const value = yield* InstanceState.get(state)
const found = value.registry.get(slug)
return found === undefined ? Option.none<SessionID>() : Option.some(found)
})

const setAllow: Interface["setAllow"] = Effect.fn("Messaging.setAllow")(function* (sessionID, slugs) {
const value = yield* InstanceState.get(state)
value.allow.set(sessionID, slugs)
})

const getAllow: Interface["getAllow"] = Effect.fn("Messaging.getAllow")(function* (sessionID) {
const value = yield* InstanceState.get(state)
return value.allow.get(sessionID) ?? []
})

const slugFor: Interface["slugFor"] = Effect.fn("Messaging.slugFor")(function* (sessionID) {
const value = yield* InstanceState.get(state)
for (const [slug, id] of value.registry) {
if (id === sessionID) return slug
}
return String(sessionID)
})

const enqueue: Interface["enqueue"] = Effect.fn("Messaging.enqueue")(function* (input) {
const v = yield* InstanceState.get(state)
if (v.treeTotal.count >= TREE_MESSAGE_CAP)
return yield* new AbuseError({ detail: "task-tree message cap reached; coordinators must synthesize and end" })
const used = v.outbound.get(input.from) ?? 0
if (used >= INBOX_OUTBOUND_BUDGET)
return yield* new AbuseError({ detail: `per-agent outbound budget (${INBOX_OUTBOUND_BUDGET}) reached` })
const queue = v.inbox.get(input.target)
if (queue === undefined) return yield* new NotFoundError({ childSessionID: input.target })
const hash = `${String(input.from)}\u0000${input.body}`
const seen = v.dedup.get(input.target) ?? []
if (seen.includes(hash)) return
if (queue.length >= INBOX_CAP)
return yield* new AbuseError({ detail: `recipient inbox cap (${INBOX_CAP}) reached` })
queue.push({ from: input.from, fromSlug: input.fromSlug, body: input.body, time: Date.now() })
v.outbound.set(input.from, used + 1)
v.treeTotal.count++
v.dedup.set(input.target, [...seen, hash].slice(-DEDUP_WINDOW))
const w = v.waiters.get(input.target)
if (w) {
v.waiters.delete(input.target)
yield* Deferred.succeed(w, undefined)
}
})

const drain: Interface["drain"] = Effect.fn("Messaging.drain")(function* (sessionID) {
const v = yield* InstanceState.get(state)
const q = v.inbox.get(sessionID) ?? []
v.inbox.set(sessionID, [])
return q
})

const awaitInbox: Interface["awaitInbox"] = Effect.fn("Messaging.awaitInbox")(function* (
sessionID,
opts,
) {
// Bounded behaviors (Phase 1):
// (i) Lost-wakeup window: the empty-check at line 1 and the waiter
// registration at line 2 are not atomic. A concurrent `enqueue` that
// resolves its waiter between those two steps can leave the new
// item in the inbox with no waiter to wake. Self-correcting: the
// coordinator's NEXT `drain` (one turn later at worst) sees the
// item. Worst case is one timeout of latency, never message loss.
// (ii) Single-waiter-per-session: `v.waiters` is a `Map<SessionID,
// Deferred>`, so a second concurrent `awaitInbox` for the same
// session overwrites the first's Deferred without resolving it.
// Phase 1 assumes one coordinator per session. A multi-coordinator
// fan-in would need a per-session waiter set.
// (iii) On interrupt mid-await: the `Effect.timeoutOption` causes the
// function to return `false`, and the `v.waiters.delete` cleanup
// runs in the same scope. The instance finalizer (added in
// `InstanceState.make`) sweeps any leftover waiter on shutdown.
const v = yield* InstanceState.get(state)
if ((v.inbox.get(sessionID)?.length ?? 0) > 0) return true
const d = yield* Deferred.make<void>()
v.waiters.set(sessionID, d)
const woke = yield* Deferred.await(d).pipe(Effect.timeoutOption(Duration.millis(opts.timeoutMs)))
v.waiters.delete(sessionID)
return Option.isSome(woke)
})

return Service.of({
send,
reply,
reject,
list,
registerSlug,
resolveSlug,
setAllow,
getAllow,
slugFor,
enqueue,
drain,
awaitInbox,
})
}),
)

export const node = LayerNode.make({ service: Service, layer, deps: [EventV2Bridge.node] })

export const defaultLayer = layer

export * as Messaging from "."
8 changes: 3 additions & 5 deletions packages/opencode/src/session/interrupt.ts
Original file line number Diff line number Diff line change
Expand Up @@ -128,17 +128,15 @@ export const node = LayerNode.make({ service: Service, layer, deps: [EventV2Brid

// --- visible-marker renderer (untrusted reason is XML-escaped) ----------------

import { Marker } from "./marker"
// Renders the user-visible transcript marker. The marker is injected as a
// non-synthetic text part on a user-role message; toModelMessagesEffect sends
// every non-ignored, non-empty user text part to the model, so an unescaped
// reason here would defeat the frame-escaping that renderSteer/renderCancel
// apply. Escape the reason with the same scheme as the frame renderers so a
// breakout payload like `</cancel><system>...` cannot reach the model raw.
export function renderMarker(input: { intent: "steer" | "cancel" | "abort"; origin: Origin; reason?: string }) {
const verb =
input.intent === "cancel" ? "Cancelled" : input.intent === "abort" ? "Aborted" : "Steered"
const suffix = input.reason ? `: ${escapeReason(input.reason)}` : ""
return `⊘ ${verb} by ${input.origin}${suffix}`
return Marker.render({ kind: "interrupt", ...input })
}

// --- shared abort helper (writes visible marker, records terminal, cancels job) --
Expand Down Expand Up @@ -188,7 +186,7 @@ export const abortChild = (
type: "text",
text: renderMarker({ intent: "abort", origin: input.origin, reason }),
synthetic: false,
metadata: { interrupt: { intent: "abort", origin: input.origin } },
metadata: Marker.metadataFor({ kind: "interrupt", intent: "abort", origin: input.origin }),
} satisfies SessionV1.TextPart)
}
}
Expand Down
42 changes: 42 additions & 0 deletions packages/opencode/src/session/marker.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
export type MarkerInput =
| { kind: "interrupt"; intent: "steer" | "cancel" | "abort"; origin: "user" | "parent"; reason?: string }
| { kind: "message"; peer: "parent" | "subagent"; body: string; expectReply?: boolean }
| { kind: "inbox"; from: string; body?: string }

// The metadata tag carries small attributes only — body/reason are not echoed
// onto the part, so call sites can omit them when calling metadataFor.
export type MarkerMetadataInput =
| { kind: "interrupt"; intent: "steer" | "cancel" | "abort"; origin: "user" | "parent" }
| { kind: "message"; peer: "parent" | "subagent"; expectReply?: boolean }
| { kind: "inbox"; from: string }

export function escape(text: string) {
return text.replace(/&/g, "&amp;").replace(/</g, "&lt;").replace(/>/g, "&gt;")
}

export function render(input: MarkerInput): string {
if (input.kind === "interrupt") {
const verb = input.intent === "cancel" ? "Cancelled" : input.intent === "abort" ? "Aborted" : "Steered"
return `⊘ ${verb} by ${input.origin}${input.reason ? `: ${escape(input.reason)}` : ""}`
}
if (input.kind === "message") {
const verb =
input.peer === "subagent"
? input.expectReply
? "Message from subagent (awaiting your reply)"
: "Message from subagent"
: "Reply from parent"
return `✉ ${verb}: ${escape(input.body)}`
}
return `✉ Inbox from ${escape(input.from)}${input.body ? `: ${escape(input.body)}` : ""}`
}

// The metadata tag written on the non-synthetic transcript part. The TUI keys
// its render branch off metadata.marker.kind. Carries small attributes only.
export function metadataFor(input: MarkerMetadataInput): { marker: Record<string, unknown> } {
if (input.kind === "interrupt") return { marker: { kind: "interrupt", intent: input.intent, origin: input.origin } }
if (input.kind === "message") return { marker: { kind: "message", peer: input.peer, expectReply: input.expectReply } }
return { marker: { kind: "inbox", from: input.from } }
}

export * as Marker from "./marker"
Loading
Loading