fix(desktop): a push that lands mid-delivery claims its envelope now — lanes outlive the drain
The previous commit claims every outbox before delivering and runs deliveries per target profile, but `drainBusy` still spanned the delivery phase: a `bot_relay.outbox.pending` push that arrived while one lane ran a long turn (up to RELAY_DELIVER_TIMEOUT_MS) only set `drainRerun`, and the new envelope was claimed after that turn — the reporter's step 3 (bot C mails D while A→B runs) still ended in `queued_expired`, because the gateway checks the TTL at the claim. Scope `drainBusy` to the claim phase and make the delivery lanes module state: `relayLanes` maps `target_connection::target_profile` to the tail of that target's in-flight deliveries, so a later drain appends to the running lane (same target stays ordered, one turn at a time) or starts a new one (other targets run now). Lane entries drop once idle; stopBotRelay clears them so a restart begins fresh. Tests: keep the contributor's red-on-base test (claims every outbox first, delivers to different targets concurrently) and replace the ordering-only test — green on base — with one that pins the mid-delivery claim plus the same-target ordering (red on base AND on the previous commit alone). Docs: bot-mode.md states the delivery concurrency contract. Part of #111587 (with the previous commit: Fixes #111587)
This commit is contained in:
@@ -776,11 +776,13 @@ describe('the drain loop does not let one delivery hold every other gateway’s
|
||||
// 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' }
|
||||
type RelayEnvelopeFixture = { id: string; message: string; target_connection: string; target_profile: string }
|
||||
const toB: RelayEnvelopeFixture = { id: 'env-1', message: 'long job', target_connection: 'b', target_profile: 'ops' }
|
||||
const toA: RelayEnvelopeFixture = { 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
|
||||
})
|
||||
@@ -824,22 +826,26 @@ describe('the drain loop does not let one delivery hold every other gateway’s
|
||||
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
|
||||
it('claims an envelope that lands DURING a long turn and delivers it now; the same target still queues behind', async () => {
|
||||
// Issue step 3: the push arrives while b's turn is running. The claim
|
||||
// must not wait for that turn (the TTL clock is running on the gateway),
|
||||
// a different target is delivered immediately, and a second envelope for
|
||||
// the SAME target profile waits for the running turn, in order.
|
||||
let releaseB!: (value: { reply: string }) => void
|
||||
|
||||
const pendingB = new Promise<{ reply: string }>(resolve => {
|
||||
releaseB = resolve
|
||||
})
|
||||
let delivered = 0
|
||||
|
||||
const outbox: Record<string, RelayEnvelopeFixture[]> = { a: [toB], b: [] }
|
||||
|
||||
const calls = respondWith(call => {
|
||||
if (call.method === 'bot_relay.outbox.drain') {
|
||||
return { envelopes: call.connectionId === 'a' ? [toB, { ...toB, id: 'env-3', message: 'second' }] : [] }
|
||||
return { envelopes: outbox[call.connectionId].splice(0) }
|
||||
}
|
||||
|
||||
if (call.method === 'bot_relay.deliver') {
|
||||
delivered += 1
|
||||
|
||||
return delivered === 1 ? first : { reply: 'second done' }
|
||||
return call.params.message === 'long job' ? pendingB : { reply: `${call.params.message} done` }
|
||||
}
|
||||
|
||||
return {}
|
||||
@@ -849,21 +855,32 @@ describe('the drain loop does not let one delivery hold every other gateway’s
|
||||
|
||||
startBotRelay()
|
||||
await pushAndSettle()
|
||||
|
||||
expect(calls.filter(call => call.method === 'bot_relay.deliver').map(call => call.params.message)).toEqual([
|
||||
'long job'
|
||||
])
|
||||
|
||||
releaseFirst({ reply: 'first done' })
|
||||
// b's turn is running; gateway b now queues one envelope for a and one more for b/ops.
|
||||
outbox.b.push(toA, { ...toB, id: 'env-3', message: 'second' })
|
||||
await pushAndSettle()
|
||||
|
||||
expect(calls.filter(call => call.method === 'bot_relay.outbox.drain' && call.connectionId === 'b')).toHaveLength(2)
|
||||
expect(calls.filter(call => call.method === 'bot_relay.deliver').map(call => call.params.message)).toEqual([
|
||||
'long job',
|
||||
'quick one'
|
||||
])
|
||||
|
||||
releaseB({ reply: 'long job done' })
|
||||
await vi.advanceTimersByTimeAsync(10)
|
||||
|
||||
expect(calls.filter(call => call.method === 'bot_relay.deliver').map(call => call.params.message)).toEqual([
|
||||
'long job',
|
||||
'quick one',
|
||||
'second'
|
||||
])
|
||||
expect(calls.filter(call => call.method === 'bot_relay.reply').map(call => call.params.reply)).toEqual([
|
||||
'first done',
|
||||
'second done'
|
||||
expect(calls.filter(call => call.method === 'bot_relay.reply').map(call => call.params.id)).toEqual([
|
||||
'env-2',
|
||||
'env-1',
|
||||
'env-3'
|
||||
])
|
||||
|
||||
stopBotRelay()
|
||||
|
||||
@@ -394,35 +394,18 @@ async function drainRelayOutboxes() {
|
||||
}
|
||||
}
|
||||
|
||||
// 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])
|
||||
}
|
||||
// Phase 2 — deliver, without holding the drain. 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. The lanes outlive this drain: `drainBusy`
|
||||
// covers the claim only, so a push that lands while a turn runs claims
|
||||
// its envelope NOW instead of after that turn — otherwise the TTL clock
|
||||
// kept running against a message the Desktop had not even looked at.
|
||||
for (const { envelope, sender } of queued) {
|
||||
enqueueRelayDelivery(sender, envelope, byId)
|
||||
}
|
||||
|
||||
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
|
||||
|
||||
@@ -440,6 +423,28 @@ interface RelayQueuedEnvelope {
|
||||
sender: RelayConnection
|
||||
}
|
||||
|
||||
// Delivery lanes: `target_connection::target_profile` → the tail of that
|
||||
// target's in-flight deliveries. A lane entry is removed once its tail
|
||||
// settles so an idle target holds no state.
|
||||
const relayLanes = new Map<string, Promise<void>>()
|
||||
|
||||
/** Queue one claimed envelope behind the deliveries already running for the
|
||||
* same target profile; other targets are untouched. */
|
||||
function enqueueRelayDelivery(sender: RelayConnection, envelope: RelayEnvelope, byId: Map<string, RelayConnection>) {
|
||||
const key = `${String(envelope?.target_connection || '')}::${String(envelope?.target_profile || '')}`
|
||||
|
||||
const tail = (relayLanes.get(key) ?? Promise.resolve()).then(() =>
|
||||
relay.disposed ? undefined : deliverRelayEnvelope(sender, envelope, byId)
|
||||
)
|
||||
|
||||
relayLanes.set(key, tail)
|
||||
void tail.finally(() => {
|
||||
if (relayLanes.get(key) === tail) {
|
||||
relayLanes.delete(key)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
/** 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(
|
||||
@@ -564,6 +569,9 @@ export function stopBotRelay() {
|
||||
// A rerun remembered mid-drain must not leak into the next start —
|
||||
// it would fire one stale drain after restart.
|
||||
relay.drainRerun = false
|
||||
// Queued deliveries check `disposed` before they run; forget the lane tails
|
||||
// so a restart starts every target fresh instead of behind stale chains.
|
||||
relayLanes.clear()
|
||||
// Unpin every relay-retained socket (#93594): with the relay stopped the
|
||||
// pooled entries return to dispose-at-refcount-0 semantics.
|
||||
releaseRelayRetention()
|
||||
|
||||
@@ -172,7 +172,7 @@ A failed bot turn or relay delivery carries a machine-readable `reason` code alo
|
||||
Every gateway you register in **Settings → Connections** — local, remote URL, SSH, Hermes Cloud, docker — is a persistent line the Desktop holds open, and Bot Mode uses those lines for messaging automatically. No extra setup:
|
||||
|
||||
- **Rosters propagate on their own.** While the Desktop runs, it periodically tells each connected gateway which agents live on the *other* connections. Every Bot Chat's teammate roster then lists them ("Teammates on OTHER connected machines"), with names, roles, and which machine they're on — and the roster refreshes when agents appear, disappear, or get renamed (capability epoch).
|
||||
- **`message_agent` reaches them directly.** A Bot on your laptop messages the cloud agent with `message_agent(target="moxie", …)` exactly like a local teammate. If the same handle exists on several machines, disambiguate with `target="moxie@<connection>"` (the tool's error tells the Bot the exact forms). Delivery rides the Desktop: the sending gateway queues the message, the Desktop relays it to the target connection's own gateway, the target Bot runs a turn in its canonical Bot Chat, and the reply comes back to the sender as the same background completion notification local DMs use.
|
||||
- **`message_agent` reaches them directly.** A Bot on your laptop messages the cloud agent with `message_agent(target="moxie", …)` exactly like a local teammate. If the same handle exists on several machines, disambiguate with `target="moxie@<connection>"` (the tool's error tells the Bot the exact forms). Delivery rides the Desktop: the sending gateway queues the message, the Desktop relays it to the target connection's own gateway, the target Bot runs a turn in its canonical Bot Chat, and the reply comes back to the sender as the same background completion notification local DMs use. Messages to different Bots are delivered side by side, so one Bot's long turn never delays another Bot's mail (or ages it past `bot_mode.envelope_ttl_seconds`); messages to the *same* Bot are delivered in order, one turn at a time.
|
||||
- **The Desktop is the courier.** Cross-connection delivery works while a Desktop that knows both connections is running (it holds the sockets and the credentials — gateways never see each other's auth). If the Desktop is closed mid-delivery, the sender's Bot is told the reply didn't arrive rather than left hanging. For always-on machine-to-machine messaging with no Desktop in the loop, register a peer (`hermes peer`, below) — the two routes coexist.
|
||||
|
||||
### Bot-initiated DMs across machines (`hermes peer`)
|
||||
|
||||
Reference in New Issue
Block a user