Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
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
369 changes: 286 additions & 83 deletions apps/desktop/src/plugins/hermes-bots/plugin.js

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -193,7 +193,7 @@ function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approva
}
if (method === 'approval.respond') {
approvalResponds.push({ ...params })
return { resolved: true }
return { resolved: 1 }
}
return {}
},
Expand Down

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -191,6 +191,11 @@ for (const expired of [false, true]) {
assert.equal(h.rpc('prompt.submit').length, 5, 'sixth stays queued while four may be running')
assert.equal(h.leases.find(l => l.route.profile === 'bot1').releases, 0)
await h.gc.stopGroupThread('Room', 't1', h.roster); await h.advance(); await pending; await flush()
if (options.unavailable) {
assert.ok(h.leases.some(l => l.releases === 0), 'an interrupt ACK does not retire an unavailable accepted turn')
options.unavailable = false
for (const member of h.roster) await h.gc.harvestStrandedGroupReply('Room', member)
}
assert.ok(h.leases.every(l => l.releases === 1))
})
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -24,11 +24,15 @@ test('capability resume recreation transfers custody to the admitted runtime', a
await flush()
assert.equal(submittedRuntime, `${oldRuntime}-recreated`)
assert.equal(h.rpc('prompt.submit').length, 1)
assert.equal((await h.gc.stopGroupThread('Room', 't1', h.roster)).status, 'stopped')
const stopping = await h.gc.stopGroupThread('Room', 't1', h.roster)
assert.equal(stopping.status, 'stopping', 'interrupt ACK precedes exact terminal retirement')
assert.equal(stopping.pending, 1)
await h.advance()
await pending
assert.equal(h.rpc('session.interrupt').at(-1).params.session_id, submittedRuntime)
assert.equal(h.gc.groupRuntimeSessionOwners.size, 0)
assert.equal(h.activeLeases(), 0)
assert.equal((await h.gc.stopGroupThread('Room', 't1', h.roster)).status, 'stopped')
})

test('Stop during capability resume interrupts the recreated runtime without admission', async () => {
Expand Down Expand Up @@ -126,10 +130,12 @@ test('retry remint collision preserves destination custody and never interrupts
assert.equal(h.gc.groupRuntimeSessionOwners.get(destinationKey), occupant)
assert.equal(h.rpc('session.interrupt').filter(c => c.params.session_id === occupiedRuntime).length, 0)
h.gc.groupRuntimeSessionOwners.delete(destinationKey)
const retained = [...h.gc.groupRuntimeSessionOwners.values()]
assert.equal(retained.length, 1, 'failed original admission retains its unresolved source custody')
assert.equal(retained[0], originalOccurrence)
assert.equal(retained[0].runtime, originalRuntime)
assert.equal(retained[0].sessionLock, originalLock)
assert.equal(h.gc.groupRuntimeSessionOwners.get(originalLock), originalOccurrence)
assert.equal(h.gc.groupRuntimeSessionOwners.size, 0, 'the proven refusal releases only its unused source custody')
assert.equal(h.gc.groupRuntimeSessionOwners.get(originalLock), undefined)
assert.equal(originalOccurrence.released, true)
assert.equal(originalOccurrence.runtime, originalRuntime)
assert.equal(originalOccurrence.sessionLock, originalLock)
assert.equal(h.activeLeases(), 0)
assert.equal(h.leases[0].releases, 1)
assert.equal(h.gc.groupRoomCoordinators.get('Room').active, 0)
})
Original file line number Diff line number Diff line change
Expand Up @@ -224,21 +224,25 @@ test('stopGroupThread with nobody on turn stops the room without any interrupt R
assert.equal(gc.$groupChats.get().Room.epoch, 4)
})

test('stopGroupThread records a stopped activity event visible in the CURRENT run', async () => {
test('stopGroupThread records an honest legacy stop-unconfirmed event in the CURRENT run', async () => {
const gc = load()
seedRoom(gc)

await gc.stopGroupThread('Room', 't1', MEMBERS)

const events = gc.currentGroupActivity('Room')
const stopped = events.find(event => event.kind === 'stopped')
const stopped = events.find(event => event.kind === 'stop-unconfirmed')
assert.ok(stopped, 'stopped event is tagged with the POST-bump epoch, so it survives the epoch filter')
assert.equal(stopped.member, 'You')
assert.equal(stopped.thread, 't1')
// The label comes from the shared GROUP_ACTIVITY_LABELS map (the plugin's
// label pattern) — and stays plain English, no hardcoded localized text.
assert.ok(gc.GROUP_ACTIVITY_LABELS.stopped)
assert.ok(gc.GROUP_ACTIVITY_GLYPHS.stopped)
assert.ok(gc.GROUP_ACTIVITY_LABELS['stop-unconfirmed'])
assert.ok(gc.GROUP_ACTIVITY_GLYPHS['stop-unconfirmed'])
const idle = load()
seedRoom(idle, { turn: null })
assert.equal((await idle.stopGroupThread('Room', 't1', MEMBERS)).status, 'stopped')
assert.ok(idle.currentGroupActivity('Room').some(event => event.kind === 'stopped' && event.epoch === 4))
})

test('stopGroupThread falls back to the durable room roster when called without members', async () => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -267,15 +267,21 @@ test('the per-turn lease is released after the turn — refcount returns to zero
assert.equal(gc.stats().disposals, 1)
})

test('uncertain submit failure retains the per-turn lease until captured Stop', async () => {
test('uncertain submit failure retains the per-turn lease after a captured Stop ACK', async () => {
const fatal = new Error('backend exploded')
const gc = load({ failEverySubmitWith: fatal })

await assert.rejects(() => gc.runGroupChatMemberTurn('Room', ROUTED_MEMBER, 'hi', 't1', []))

assert.equal(gc.stats().refcount, 1, 'unknown acceptance is not vacant capacity')
await gc.stopGroupThread('Room', 't1', [ROUTED_MEMBER])
assert.equal(gc.stats().refcount, 0)
const first = await gc.stopGroupThread('Room', 't1', [ROUTED_MEMBER])
assert.equal(first.status, 'stopping')
assert.equal(first.pending, 1)
assert.equal(gc.stats().refcount, 1, 'ACK cannot prove retirement without admission identity')
const second = await gc.stopGroupThread('Room', 't1', [ROUTED_MEMBER])
assert.equal(second.status, 'stopping')
assert.equal(gc.stats().refcount, 1)
assert.equal(gc.stats().submits, 1, 'unknown acceptance never replays user text')
})

test('hosts without retainProfile still run the turn (feature detection)', async () => {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -436,10 +436,18 @@ test('explicit stop exits the collector and never collects its late complete out
let stopped = false
const h = await harness({ alpha: [async (s, gc) => {
if (!stopped) { stopped = true; await gc.stopGroupThread('Room', 'thread-1', [ALPHA, BETA]) }
return { turn_outcomes: wire(s.ref, 'running', [final('late candidate')]) }
}] })
return { turn_outcomes: wire(s.ref, s.terminalAfterStop ? 'complete' : 'running', [final('late candidate')]) }
}] }, { connectionId: 'pc' })
assert.equal(await run(h), null)
const marker = room(h).stranded.alpha
assert.equal(marker.stop_requested, true)
assert.deepEqual(marker.delivery.accepted_turn, h.sessions.get('alpha').ref)
assert.equal(h.releases(), 0, 'running custody outlives the interrupt ACK')
h.sessions.get('alpha').terminalAfterStop = true
await h.gc.harvestStrandedGroupReply('Room', ALPHA)
await new Promise(resolve => setImmediate(resolve))
assert.equal(room(h).stranded.alpha, undefined)
assert.equal(h.releases(), 1, 'the exact terminal read releases custody after control settles')
assert.equal(posts(h).length, 0)
assert.equal(h.rpc('prompt.submit').length, 1)
})
Expand Down Expand Up @@ -716,12 +724,22 @@ test('rejected read-only poll is visible immediately and next member advances on
})

test('stopped turn discards an in-flight poll rejection without stale failure publication', async () => {
const h = await harness({ alpha: [async (_s, gc) => {
await gc.stopGroupThread('Room', 'thread-1', [ALPHA])
let stopped = false
const h = await harness({ alpha: [async (s, gc) => {
if (!stopped) { stopped = true; await gc.stopGroupThread('Room', 'thread-1', [ALPHA]) }
if (s.terminalAfterStop) return { turn_outcomes: wire(s.ref, 'complete', [final('late rejected candidate')]) }
throw new Error('late lost observer')
}] }, { members: [ALPHA] })
}] }, { members: [ALPHA], connectionId: 'pc' })
assert.equal(await run(h), null)
assert.equal(room(h).stranded.alpha.stop_requested, true)
assert.deepEqual(room(h).stranded.alpha.delivery.accepted_turn, h.sessions.get('alpha').ref)
assert.equal(h.rpc('prompt.submit').length, 1)
assert.equal(h.releases(), 0, 'a rejected observation is not terminal evidence')
h.sessions.get('alpha').terminalAfterStop = true
await h.gc.harvestStrandedGroupReply('Room', ALPHA)
await new Promise(resolve => setImmediate(resolve))
assert.equal(room(h).stranded.alpha, undefined)
assert.equal(h.releases(), 1)
assert.equal(posts(h).length, 0)
assert.equal(h.gc.currentGroupActivity('Room').filter(e => e.kind === 'unavailable').length, 0)
})
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -185,6 +185,7 @@ for (const failsFirst of [false, true]) {
assert.ok(receipts(cold)[0].stop_requested)
options.interruptError = false
await cold.gc.stopGroupThread('Room', 't1', cold.roster)
await flush() // control completion precedes exact terminal observation
assert.deepEqual(cold.rpc('session.interrupt').map(call => call.params.session_id), [accepted.session_id, accepted.session_id])
assert.equal(receipts(cold).length, 0)
}
Expand Down Expand Up @@ -267,6 +268,52 @@ boundedTest('actual UI exposes cold Stop, retains it after failed interrupt and
await hot.gc.stopGroupThread('Room', 't1', hot.roster)
})

for (const cold of [false, true]) {
boundedTest(`actual Stop button projects pending retirement honestly after applied ACK; cold=${cold}`, async () => {
const notices = [], options = { interruptReply: { status: 'interrupted' }, onNotify: value => notices.push(value) }
const hot = await uiHarness(members(1), cold ? {} : options), pending = drive(hot)
await flush()
const h = cold ? await reload(hot, options) : hot
const accepted = clone(receipts(h)[0].delivery.accepted_turn), button = h.stopButton()
assert.ok(button)
button.props.onClick(); await flush()
assert.equal(notices.length, 1)
assert.equal(notices[0].kind, 'info')
assert.match(notices[0].message, /^Stopping Room/)
assert.match(notices[0].message, /waiting for the remaining turns to finish/)
assert.doesNotMatch(notices[0].message, /unconfirmed|Stop can retry/)
assert.deepEqual(receipts(h)[0].delivery.accepted_turn, accepted)
assert.equal(h.rpc('session.interrupt').length, 1)
assert.equal(h.rpc('prompt.submit').length, cold ? 0 : 1)
if (!cold) {
assert.equal(h.activeLeases(), 1)
assert.equal(h.gc.groupRoomCoordinators.get('Room').active, 1)
}
const session = [...h.sessions.values()][0]
session.state = 'interrupted'; session.pending = null
await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await h.advance()
assert.equal(receipts(h).length, 0)
if (cold) await hot.gc.stopGroupThread('Room', 't1', hot.roster)
await hot.advance(); await pending
})
}

boundedTest('actual Stop button keeps failed interruption wording and retry custody', async () => {
const notices = [], h = await uiHarness(members(1), { interruptError: true, onNotify: value => notices.push(value) })
const pending = drive(h); await flush()
const accepted = clone(receipts(h)[0].delivery.accepted_turn)
h.stopButton().props.onClick(); await flush(); await h.advance(); await pending
assert.equal(notices.length, 1)
assert.equal(notices[0].kind, 'info')
assert.match(notices[0].message, /1 interruption\(s\) are unconfirmed.*Stop can retry/)
assert.deepEqual(receipts(h)[0].delivery.accepted_turn, accepted)
assert.equal(h.activeLeases(), 1)
assert.ok(h.stopButton())
h.finish('bot1', '', 'interrupted')
await h.gc.harvestStrandedGroupReply('Room', h.roster[0])
assert.equal(h.activeLeases(), 0)
})

boundedTest('legacy unknown receipt does not manufacture an interrupt target or Stop affordance', async () => {
const h = await uiHarness()
h.gc.updateGroupChat('Room', room => ({ ...room, running: false, turns: [],
Expand Down Expand Up @@ -379,6 +426,7 @@ boundedTest(`unresolved waiting reservations retain capacity; ${invalidation} ca
const secondAdmitted = new Set()
const h = await uiHarness(members(6), { resumeProjection: (session, method, projection) => {
if (method !== 'session.turn.poll' || !session.ref) return projection
if (session.state === 'interrupted') return projection // preserve exact terminal after fake Stop
if (session.submits === 2) secondAdmitted.add(session.profile)
session.state = session.submits === 1 ? 'complete' : 'waiting'
session.text = `FIRST_${session.profile} @all`
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,81 @@
import assert from 'node:assert/strict'
import test from 'node:test'
import { harness, members, deferred, flush, drive } from './stop-custody-harness.mjs'

test('PR127: cancelled proven refusal retains preparation custody, then releases without admission', async () => {
const gate = deferred()
const h = await harness(members(1), {
submitError: s => s.submits === 1 ? Object.assign(new Error('session not found'), { code: 4001 }) : null,
beforeResume(s, sessions) {
if (s.submits !== 1) return
sessions.delete(s.runtime); s.runtime += '-retry'; sessions.set(s.runtime, s)
},
resumeProjection(s, method, projection) {
if (method !== 'session.resume' || s.submits !== 1) return projection
const lazy = { ...projection }; delete lazy.turn_outcomes; return lazy
},
async capabilityPoll() {
await gate.promise
throw Object.assign(new Error('session_id and full accepted_turn identity required'), { code: 4006 })
}
})
const pending = drive(h)
await flush()
const coordinator = h.gc.groupRoomCoordinators.get('Room')
const stopped = await h.gc.stopGroupThread('Room', 't1', h.roster)
assert.equal(stopped.status, 'unconfirmed')
assert.equal(coordinator.active, 1)
assert.equal(h.activeLeases(), 1)
assert.equal(h.gc.groupRuntimeSessionOwners.size, 1)
assert.equal(h.rpc('prompt.submit').length, 1)
gate.resolve(); await flush(); await h.advance(); await pending
assert.equal(h.rpc('prompt.submit').length, 1)
assert.equal(h.gc.groupRuntimeSessionOwners.size, 0)
assert.equal(coordinator.active, 0)
assert.equal(h.activeLeases(), 0)
assert.equal(h.leases[0].releases, 1)
assert.equal(Object.keys(h.room().stranded).length, 0)
assert.equal((await h.gc.stopGroupThread('Room', 't1', h.roster)).status, 'stopped')
})

test('PR127: a refused first submit never excuses unknown acceptance of its retry', async () => {
const h = await harness(members(1), {
submitError: s => s.submits === 1
? Object.assign(new Error('session not found'), { code: 4001 })
: new Error('lost retry admission acknowledgement')
})
await drive(h)
assert.equal(h.rpc('prompt.submit').length, 2)
assert.equal(h.activeLeases(), 1)
assert.equal(h.gc.groupRoomCoordinators.get('Room').active, 1)
const first = await h.gc.stopGroupThread('Room', 't1', h.roster)
assert.equal(first.status, 'stopping')
assert.equal(first.pending, 1)
assert.equal(h.activeLeases(), 1)
assert.equal(h.gc.groupRuntimeSessionOwners.size, 1)
const second = await h.gc.stopGroupThread('Room', 't1', h.roster)
assert.equal(second.status, 'stopping')
assert.equal(h.rpc('prompt.submit').length, 2, 'Stop never retries uncertain text')
assert.equal(h.leases[0].releases, 0)
})

test('PR127: four unknown retry admissions continue to block a six-member room', async () => {
const h = await harness(members(6), {
submitError: s => s.submits === 1
? Object.assign(new Error('session not found'), { code: 4001 })
: new Error('lost retry acknowledgement')
})
const pending = drive(h)
await flush()
assert.equal(h.rpc('prompt.submit').length, 8)
assert.equal(h.gc.groupRoomCoordinators.get('Room').active, 4)
assert.equal(h.gc.groupRoomCoordinators.get('Room').queue.length, 2)
assert.equal(h.activeLeases(), 4)
const stopped = await h.gc.stopGroupThread('Room', 't1', h.roster)
await h.advance(); await pending
assert.equal(stopped.status, 'stopping')
assert.equal(stopped.pending, 4)
assert.equal(h.gc.groupRoomCoordinators.get('Room').active, 4)
assert.equal(h.activeLeases(), 4)
assert.equal(h.rpc('prompt.submit').length, 8)
})
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ async function harness(roster = members(3), options = {}) {
const sourceKey = (route, profile) => `${route?.connectionId || connection}::${route?.targetProfile || profile}`
const handle = async (route, method, params) => {
calls.push({ route: route && { ...route }, method, params: { ...params }, at: now })
const supplied = await options.rpcResponse?.(route, method, params)
if (supplied !== undefined) return supplied
const key = sourceKey(route, params.profile)
if (method === 'session.create') {
if (options.createGate) await options.createGate.promise
Expand Down Expand Up @@ -53,8 +55,10 @@ async function harness(roster = members(3), options = {}) {
if (method === 'clarify.respond' || method === 'approval.respond') {
if (options.answerError?.()) throw new Error('response transport failed')
if (options.answerGate) await options.answerGate.promise
const acknowledgement = method === 'clarify.respond' ? (options.clarifyResult ?? { status: 'ok' }) : {}
const acknowledgement = method === 'clarify.respond'
? (options.clarifyResult ?? { status: 'ok' }) : (options.approvalResult ?? { resolved: 1 })
if (method === 'clarify.respond' && acknowledgement?.status !== 'ok') return acknowledgement
if (method === 'approval.respond' && !(Number.isSafeInteger(acknowledgement?.resolved) && acknowledgement.resolved > 0)) return acknowledgement
if (method === 'clarify.respond' && session.pending?.questions?.length) {
session.questionAnswers ||= new Set()
session.questionAnswers.add(params.question_id)
Expand Down Expand Up @@ -104,7 +108,7 @@ async function harness(roster = members(3), options = {}) {
return () => { lease.releases++; activeLeases-- }
},
state: { profile: atom('default'), gateway: atom(null), connectionId: { get: () => connection, listen: () => () => undefined } },
notify: () => undefined, notifyError: () => undefined }
notify: notice => options.onNotify?.(notice), notifyError: () => undefined }
})
gc.stopGroupChatServerSync()
gc.bindGroupTurnTestStorage({ get: key => clone(storage.get(key) ?? null), set: (key, value) => storage.set(key, clone(value)) })
Expand Down
Loading
Loading