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
Next Next commit
fix(opencode): bound s2s delivery failures
  • Loading branch information
iceteaSA committed Aug 22, 2026
commit 336eb3c18da0ac0b03ca8a9facbb4e5391a0709f
5 changes: 4 additions & 1 deletion packages/opencode/src/messaging/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -300,7 +300,10 @@ export const layer = Layer.effect(
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)
// S2S has its own hourly sender cap and durable recipient inbox cap;
// this budget is specifically for spawned-subagent abuse and must not
// add a second, process-local limit to consented peer traffic.
if (input.source !== "sibling-session" && used >= INBOX_OUTBOUND_BUDGET)
return yield* new AbuseError({ detail: `per-agent outbound budget (${INBOX_OUTBOUND_BUDGET}) reached` })
// Lazily init the recipient queue. The inbox no longer depends on a prior
// registerSlug/registerLocal having pre-created it — s2s addresses peers
Expand Down
18 changes: 17 additions & 1 deletion packages/opencode/src/s2s/poller.ts
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,9 @@ import { registerWakeBody } from "@/s2s/wake-registry"

const REAPER_WINDOW_MS_DEFAULT = 60_000
const MIN_REAPER_MS = 1
const MAX_ROW_FAILURES = 3
const rowFailures = new Map<string, number>()
const abandonedRows = new Set<string>()

export interface Interface {
readonly pollOnce: () => Effect.Effect<void>
Expand All @@ -65,6 +68,7 @@ export class Service extends Context.Service<Service, Interface>()("@opencode/S2
// `claimForServices` returns rows in the order the SQL engine picks; we
// don't care about that here, only that each row is processed once.
const processRow = Effect.fn("S2SPoller.processRow")(function* (row: S2SStore.InboxRow) {
if (abandonedRows.has(row.id)) return
const messaging = yield* Messaging.Service
const store = yield* S2SStore.Service

Expand Down Expand Up @@ -158,7 +162,19 @@ export const pollOnceImpl = Effect.fn("S2SPoller.pollOnce")(function* () {
// prevent subsequent rows from being processed. Failures are caught
// and logged so the loop always continues to the next row.
yield* processRow(row).pipe(
Effect.catch((e) => Effect.logWarning("S2SPoller: row processing failed", e)),
Effect.catch((e) =>
Effect.gen(function* () {
const failures = (rowFailures.get(row.id) ?? 0) + 1
rowFailures.set(row.id, failures)
if (failures < MAX_ROW_FAILURES) return
abandonedRows.add(row.id)
yield* Effect.logWarning("S2SPoller: giving up on row for this process lifetime", {
rowID: row.id,
failures,
error: e,
})
}),
),
)
}
})
Expand Down
21 changes: 21 additions & 0 deletions packages/opencode/test/messaging/inbox.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,27 @@ it.instance("over-budget send (M) fails with AbuseError", () =>
}),
)

it.instance("sibling-session enqueue is not limited by the subagent outbound budget", () =>
Effect.gen(function* () {
const m = yield* Messaging.Service
const to = SessionID.make("ses_sibling_target_xxxxxxxxxx")
const from = SessionID.make("ses_sibling_sender_xxxxxxxxxx")
yield* m.registerSlug("sibling-target", to)

for (let i = 0; i < 21; i++) {
yield* m.enqueue({
target: to,
from,
fromSlug: "sibling-peer",
body: `sibling-${i}`,
source: "sibling-session",
})
}

expect((yield* m.drain(to)).length).toBe(21)
}),
)

it.instance("awaitInbox returns true when inbox already has items", () =>
Effect.gen(function* () {
const m = yield* Messaging.Service
Expand Down
70 changes: 70 additions & 0 deletions packages/opencode/test/s2s/poller-retry.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
import { expect } from "bun:test"
import { Effect, Layer } from "effect"
import { Database } from "@opencode-ai/core/database/database"
import { Messaging } from "../../src/messaging"
import { pollOnceImpl } from "../../src/s2s/poller"
import { S2SStore } from "../../src/s2s/store"
import { SessionID } from "../../src/session/schema"
import { encodeCapsule } from "../../src/s2s/capsule"
import { testEffectIsolatedShared } from "../lib/effect"

const target = SessionID.make("ses_retry_target_xxxxxxxxxxxx")
const sender = SessionID.make("ses_retry_sender_xxxxxxxxxxxx")
let attempts = 0
const capsule = encodeCapsule({
version: 1,
id: "0190abcd-7abc-7abc-8abc-0190abcdef01",
sender_slug: "retry-peer",
sender_session_id: String(sender),
timestamp: 1_700_000_000_000,
body: "BOUNDED-RETRY-PAYLOAD",
})

const it = testEffectIsolatedShared(
Layer.mergeAll(
S2SStore.defaultLayer,
Layer.succeed(Messaging.Service, {
send: () => Effect.die("unexpected Messaging.send"),
reply: () => Effect.die("unexpected Messaging.reply"),
reject: () => Effect.die("unexpected Messaging.reject"),
list: () => Effect.die("unexpected Messaging.list"),
registerSlug: () => Effect.die("unexpected Messaging.registerSlug"),
resolveSlug: () => Effect.die("unexpected Messaging.resolveSlug"),
setAllow: () => Effect.die("unexpected Messaging.setAllow"),
getAllow: () => Effect.die("unexpected Messaging.getAllow"),
slugFor: () => Effect.die("unexpected Messaging.slugFor"),
enqueue: () => {
attempts++
return Effect.fail(new Messaging.AbuseError({ detail: "deterministic test failure" }))
},
drain: () => Effect.die("unexpected Messaging.drain"),
awaitInbox: () => Effect.die("unexpected Messaging.awaitInbox"),
registerLocal: () => Effect.die("unexpected Messaging.registerLocal"),
isLocal: () => Effect.die("unexpected Messaging.isLocal"),
localSet: () => Effect.succeed([target]),
}),
).pipe(Layer.provide(Database.layerFromPath(":memory:"))),
)

it.instance("stops retrying a deterministically failing row within the process lifetime", () =>
Effect.gen(function* () {
const store = yield* S2SStore.Service
attempts = 0
yield* store.insertInbox({
id: "inb_bounded_retry",
targetSessionID: target,
fromSessionID: sender,
fromSlug: "retry-peer",
capsule,
timeCreated: 1,
})

for (let i = 0; i < 6; i++) {
yield* (pollOnceImpl() as unknown as Effect.Effect<void, S2SStore.S2SStoreError>)
yield* store.reapStale(Date.now() + 10 ** 9)
}

expect(attempts).toBeLessThan(6)
expect(attempts).toBeGreaterThan(0)
}),
)