fix(desktop): the bot relay claims every outbox first and delivers per target, not one envelope at a time

drainRelayOutboxes drained one gateway's outbox and delivered its envelopes
before draining the next gateway, and delivered every envelope one after
another. One long turn — bot_relay.deliver may take up to
RELAY_DELIVER_TIMEOUT_MS, 25 minutes — therefore held every other bot's mail:
a sibling gateway's envelope sat unclaimed in its outbox, and since the
gateway checks the envelope's age against bot_mode.envelope_ttl_seconds
(15 minutes) at the claim, the sender's waiter received queued_expired for a
message nothing was wrong with; envelopes that did get claimed still waited
their turn behind unrelated deliveries, against a finite waiter.

Claim every gateway's outbox first, then deliver in lanes keyed by target
connection and profile: a lane runs its envelopes in order, one turn at a
time (the target gateway serialises that profile's turns behind its turn
lock anyway), and lanes run concurrently.
This commit is contained in:
John Paul Soliva
2026-09-15 15:02:56 +09:00
committed by Teknium
parent 81805a97ef
commit 1eb771e2ff
2 changed files with 221 additions and 76 deletions

View File

@@ -770,3 +770,102 @@ describe('the roster loop forgets a machine that left', () => {
stopBotRelay()
})
})
describe('the drain loop does not let one delivery hold every other gateway’s mail', () => {
// A long turn on b (up to RELAY_DELIVER_TIMEOUT_MS) must not stop b's own
// outbox from being claimed, nor a's delivery from running: the gateway
// checks envelope age against bot_mode.envelope_ttl_seconds at the claim,
// and the sender's waiter is finite.
const toB = { id: 'env-1', message: 'long job', target_connection: 'b', target_profile: 'ops' }
const toA = { id: 'env-2', message: 'quick one', target_connection: 'a', target_profile: 'default' }
it('claims every outbox first and delivers to different targets concurrently', async () => {
let releaseB!: (value: { reply: string }) => void
const pendingB = new Promise<{ reply: string }>(resolve => {
releaseB = resolve
})
const calls = respondWith(call => {
if (call.method === 'bot_relay.outbox.drain') {
return { envelopes: call.connectionId === 'a' ? [toB] : [toA] }
}
if (call.method === 'bot_relay.deliver') {
return call.connectionId === 'b' ? pendingB : { reply: 'done' }
}
return {}
})
const { startBotRelay, stopBotRelay } = await loadRelay()
startBotRelay()
await pushAndSettle()
// b's turn is still running — yet b's outbox was claimed and its envelope delivered on a.
expect(calls.filter(call => call.method === 'bot_relay.outbox.drain').map(call => call.connectionId)).toEqual([
'a',
'b'
])
expect(calls.filter(call => call.method === 'bot_relay.deliver').map(call => call.connectionId)).toEqual(['b', 'a'])
expect(calls.find(call => call.method === 'bot_relay.reply')).toMatchObject({
connectionId: 'b',
params: { id: 'env-2', reply: 'done' }
})
releaseB({ reply: 'finally' })
await vi.advanceTimersByTimeAsync(10)
expect(calls.filter(call => call.method === 'bot_relay.reply').map(call => call.params.id)).toEqual([
'env-2',
'env-1'
])
stopBotRelay()
})
it('keeps envelopes for the same target profile in order, one turn at a time', async () => {
let releaseFirst!: (value: { reply: string }) => void
const first = new Promise<{ reply: string }>(resolve => {
releaseFirst = resolve
})
let delivered = 0
const calls = respondWith(call => {
if (call.method === 'bot_relay.outbox.drain') {
return { envelopes: call.connectionId === 'a' ? [toB, { ...toB, id: 'env-3', message: 'second' }] : [] }
}
if (call.method === 'bot_relay.deliver') {
delivered += 1
return delivered === 1 ? first : { reply: 'second done' }
}
return {}
})
const { startBotRelay, stopBotRelay } = await loadRelay()
startBotRelay()
await pushAndSettle()
expect(calls.filter(call => call.method === 'bot_relay.deliver').map(call => call.params.message)).toEqual([
'long job'
])
releaseFirst({ reply: 'first done' })
await vi.advanceTimersByTimeAsync(10)
expect(calls.filter(call => call.method === 'bot_relay.deliver').map(call => call.params.message)).toEqual([
'long job',
'second'
])
expect(calls.filter(call => call.method === 'bot_relay.reply').map(call => call.params.reply)).toEqual([
'first done',
'second done'
])
stopBotRelay()
})
})

View File

@@ -369,9 +369,16 @@ async function drainRelayOutboxes() {
const byId = new Map(connections.map(connection => [connection.id, connection]))
for (const sender of connections) {
let envelopes: RelayEnvelope[] = []
// Phase 1 — claim every gateway's outbox before delivering anything. The
// gateway checks an envelope's age against bot_mode.envelope_ttl_seconds
// AT the claim, so a claim must never wait behind another gateway's
// delivery turn: draining and delivering sender by sender let one long
// turn (up to RELAY_DELIVER_TIMEOUT_MS) age a sibling gateway's envelope
// past the TTL before it was even claimed, and that sender's waiter got
// `queued_expired` for a message nothing was wrong with.
const queued: RelayQueuedEnvelope[] = []
for (const sender of connections) {
try {
const res = await host.requestProfile<{ envelopes?: RelayEnvelope[] }>(
sender.route,
@@ -379,83 +386,43 @@ async function drainRelayOutboxes() {
{}
)
envelopes = Array.isArray(res?.envelopes) ? res.envelopes : []
for (const envelope of Array.isArray(res?.envelopes) ? res.envelopes : []) {
queued.push({ envelope, sender })
}
} catch {
continue
}
for (const envelope of envelopes) {
if (relay.disposed) {
return
}
const envelopeId = String(envelope?.id || '')
const target = byId.get(String(envelope?.target_connection || ''))
const postReply = async (payload: { error?: string; reason?: string; reply?: string }) => {
try {
await host.requestProfile(sender.route, 'bot_relay.reply', {
id: envelopeId,
...payload
})
} catch {
// Sender gateway unreachable — its waiter times out with guidance.
}
}
if (!envelopeId) {
continue
}
if (!target) {
await postReply({
error: `connection '${envelope?.target_connection}' is not connected to this Desktop right now`
})
continue
}
// Needs-attention hook (#93091 item 3): a delivered background DM is
// this bot's "good turn"; a classified delivery failure badges it.
const attentionKey = `${target.id}::${String(envelope?.target_profile || '')}`
try {
const res = await host.requestProfile<{ reply?: string }>(
target.route,
'bot_relay.deliver',
{
profile: String(envelope?.target_profile || ''),
message: String(envelope?.message || ''),
from_profile: String(envelope?.from_profile || ''),
from_handle: String(envelope?.from_handle || ''),
from_connection: String(sender.id)
},
RELAY_DELIVER_TIMEOUT_MS
)
clearBotAttention(attentionKey)
await postReply({
reply: String(res?.reply || '')
})
} catch (error: any) {
// #93091: bot_relay.deliver classifies the failed turn and ships the
// typed code in the JSON-RPC error's `data.reason`; forward it into
// the sender-side reply file so the waiter (and the sending agent)
// get the machine-readable cause, and prefer it for the badge —
// classified codes beat free-text re-parsing.
const reason = String(error?.data?.reason || '').trim()
noteBotAttention(attentionKey, reason || error?.message || error)
await postReply({
error: String(error?.message || error || 'delivery failed'),
...(reason
? {
reason
}
: {})
})
}
// Older backend without the relay RPCs — skip this connection.
}
}
// Phase 2 — deliver. Envelopes for the same target profile stay in order,
// one turn at a time (the target gateway serialises that profile's turns
// behind its turn lock anyway; a second envelope sent while the first runs
// would only burn its lock-wait budget). Different targets run
// concurrently, so one long turn never holds every other bot's mail.
const lanes = new Map<string, RelayQueuedEnvelope[]>()
for (const item of queued) {
const key = `${String(item.envelope?.target_connection || '')}::${String(item.envelope?.target_profile || '')}`
const lane = lanes.get(key)
if (lane) {
lane.push(item)
} else {
lanes.set(key, [item])
}
}
await Promise.all(
[...lanes.values()].map(async lane => {
for (const { envelope, sender } of lane) {
if (relay.disposed) {
return
}
await deliverRelayEnvelope(sender, envelope, byId)
}
})
)
} finally {
relay.drainBusy = false
@@ -468,6 +435,85 @@ async function drainRelayOutboxes() {
}
}
interface RelayQueuedEnvelope {
envelope: RelayEnvelope
sender: RelayConnection
}
/** Deliver one claimed envelope on the target connection's own socket and post
* the reply (or the error) back to the sender gateway for its waiter. */
async function deliverRelayEnvelope(
sender: RelayConnection,
envelope: RelayEnvelope,
byId: Map<string, RelayConnection>
) {
const envelopeId = String(envelope?.id || '')
const target = byId.get(String(envelope?.target_connection || ''))
const postReply = async (payload: { error?: string; reason?: string; reply?: string }) => {
try {
await host.requestProfile(sender.route, 'bot_relay.reply', {
id: envelopeId,
...payload
})
} catch {
// Sender gateway unreachable — its waiter times out with guidance.
}
}
if (!envelopeId) {
return
}
if (!target) {
await postReply({
error: `connection '${envelope?.target_connection}' is not connected to this Desktop right now`
})
return
}
// Needs-attention hook (#93091 item 3): a delivered background DM is
// this bot's "good turn"; a classified delivery failure badges it.
const attentionKey = `${target.id}::${String(envelope?.target_profile || '')}`
try {
const res = await host.requestProfile<{ reply?: string }>(
target.route,
'bot_relay.deliver',
{
profile: String(envelope?.target_profile || ''),
message: String(envelope?.message || ''),
from_profile: String(envelope?.from_profile || ''),
from_handle: String(envelope?.from_handle || ''),
from_connection: String(sender.id)
},
RELAY_DELIVER_TIMEOUT_MS
)
clearBotAttention(attentionKey)
await postReply({
reply: String(res?.reply || '')
})
} catch (error: any) {
// #93091: bot_relay.deliver classifies the failed turn and ships the
// typed code in the JSON-RPC error's `data.reason`; forward it into
// the sender-side reply file so the waiter (and the sending agent)
// get the machine-readable cause, and prefer it for the badge —
// classified codes beat free-text re-parsing.
const reason = String(error?.data?.reason || '').trim()
noteBotAttention(attentionKey, reason || error?.message || error)
await postReply({
error: String(error?.message || error || 'delivery failed'),
...(reason
? {
reason
}
: {})
})
}
}
/** Push-notified drain (#93091): collapse a burst of pending signals into
* one drain call ~RELAY_PUSH_DEBOUNCE_MS after the first signal. */
function scheduleRelayPushDrain() {