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:
committed by
Teknium
parent
81805a97ef
commit
1eb771e2ff
@@ -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()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -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() {
|
||||
|
||||
Reference in New Issue
Block a user