From b4e6c5859e27b6495d7387279fdefd2012bc2ff1 Mon Sep 17 00:00:00 2001 From: Josh Stevenson Date: Sat, 3 Oct 2026 23:38:22 -0700 Subject: [PATCH 1/2] fix(desktop): preserve group turn custody through Stop and rename races --- .../desktop/src/plugins/hermes-bots/plugin.js | 341 +++++++--- .../tests/group-lifecycle-boundaries.test.mjs | 580 ++++++++++++++++++ .../hermes-bots/tests/group-parallel.test.mjs | 5 + .../tests/pr118-postmerge.test.mjs | 48 ++ .../tests/stop-custody-harness.mjs | 8 +- .../hermes-bots/tests/stop-custody.test.mjs | 27 +- 6 files changed, 922 insertions(+), 87 deletions(-) create mode 100644 apps/desktop/src/plugins/hermes-bots/tests/group-lifecycle-boundaries.test.mjs diff --git a/apps/desktop/src/plugins/hermes-bots/plugin.js b/apps/desktop/src/plugins/hermes-bots/plugin.js index 519e20538202..cf6fa1635fed 100644 --- a/apps/desktop/src/plugins/hermes-bots/plugin.js +++ b/apps/desktop/src/plugins/hermes-bots/plugin.js @@ -529,6 +529,7 @@ const GROUP_ACTIVITY_LABELS = { capped: 'turn stopped at the round/message cap', delivered: 'delivered a late reply', held: 'is held (stopped by you) — @mention it or say resume to release', + stopping: 'held the room — interruption was applied; waiting for turn retirement', 'stop-unconfirmed': 'held the room — interruption is unconfirmed; Stop can retry', stopped: 'stopped the room — remaining turns are held until resumed' } @@ -549,6 +550,7 @@ const GROUP_ACTIVITY_GLYPHS = { capped: 'debug-step-over', delivered: 'mail-read', held: 'debug-pause', + stopping: 'debug-stop', 'stop-unconfirmed': 'error', stopped: 'debug-stop' } @@ -1285,6 +1287,7 @@ async function pullGroupChatServerState(connectionId = groupChatSyncConnectionId preserveRooms: pending?.changedRooms || [], deletedRooms: pending?.deletedRooms || [] }) + rehomeGroupRuntimeOwners(merged) $groupChats.set(merged) await persistGroupChatRooms(merged) return true @@ -1376,6 +1379,7 @@ async function flushGroupChatServerSync(connectionId) { preserveRooms: pending?.changedRooms || [], deletedRooms: pending?.deletedRooms || [] }) + rehomeGroupRuntimeOwners(mergedRooms) $groupChats.set(mergedRooms) await persistGroupChatRooms(mergedRooms) } @@ -1412,6 +1416,7 @@ async function flushGroupChatServerSync(connectionId) { preserveRooms: pending?.changedRooms || [], deletedRooms: pending?.deletedRooms || [] }) + rehomeGroupRuntimeOwners(mergedRooms) $groupChats.set(mergedRooms) await persistGroupChatRooms(mergedRooms) } @@ -7285,9 +7290,9 @@ function setGroupChatImage(group, image) { }) } -/** Rename a group chat. The group's NAME is its identity everywhere — the - * room-map key, each local member's ui_meta membership list, and derived - * state — so a rename re-keys all of them. Member gateway sessions are kept +/** Rename a group chat. The group's name keys its presentation — the + * room map, each local member's ui_meta membership list, and derived + * state — so a rename re-keys them while retaining its room lifetime. Member gateway sessions are kept * as-is: stored sids keep resuming, so no history is lost. The room's * immutable roomId (the member-session title) is preserved across the * rename, so even a member whose sid is later lost falls back to the same @@ -7333,6 +7338,7 @@ async function renameGroupChat(oldName, newName, members) { all[next] = room } + rehomeGroupRuntimeOwners(all) $groupChats.set(all) const needs = { ...$groupNeedsYou.get() } @@ -7343,13 +7349,13 @@ async function renameGroupChat(oldName, newName, members) { $groupNeedsYou.set(needs) } - // Mirrored clarify cards key by group name; drop the old room's — the - // next poll re-mirrors any still-blocking question under the new name. - clearGroupClarify(oldName) + // Runtime owners/cards follow the same immutable lifetime above. A display + // rename neither replaces the accepted occurrence nor erases its question. // Local memberships: swap the name inside each member's canonical groups // list (syncs cross-machine via ui_meta). Remote members' seating lives in // the room record we just moved. + const metadataPersistence = [] for (const member of members || []) { if (!member?.name) { continue @@ -7358,11 +7364,20 @@ async function renameGroupChat(oldName, newName, members) { const meta = botRosterMeta(member, $botMeta.get()) || {} const groups = [...new Set(botGroups(meta).map(g => (g === oldName ? next : g)))] - await saveBotMeta(member, { groups, group: groups[0] || null }) + // saveBotMeta applies its local atom before yielding. Start every member's + // transition now so a later rename cannot leave old-name writes in this loop. + metadataPersistence.push(saveBotMeta(member, { groups, group: groups[0] || null })) } // Persist the re-keyed map (updateGroupChat writes the whole durable map). - updateGroupChat(next, r => r, { sync: false }) + const renamedRoom = updateGroupChat(next, r => { + // Idle legacy/metadata-only rooms need the same existing lifetime token + // as a runtime coordinator before this operation first yields. + if (!r.roomId && !r.coordinationId) r.coordinationId = groupChatEntryId() + return r + }, { sync: false }) + const lifetime = { group: next, roomId: renamedRoom.roomId || null, + roomToken: renamedRoom.coordinationId || null } // A rename is one revisioned state transition: the new identity is updated // and the old identity is tombstoned together, so cold hydration cannot // merge the pre-rename room back into the roster. @@ -7381,13 +7396,18 @@ async function renameGroupChat(oldName, newName, members) { openGroupChat(next) } + // Room, membership, persistence and navigation intent are all committed + // before waiting for metadata acknowledgement. No captured name is written + // by this continuation after an overlapping local/remote rename or disband. + await Promise.all(metadataPersistence) + // Same convergence as disband: drop the pre-rename roster snapshot so the // old name can't linger anywhere the fence doesn't cover. if (typeof queryClient !== 'undefined' && queryClient?.invalidateQueries) { queryClient.invalidateQueries({ queryKey: ROSTER_KEY }) } - return next + return groupNameForLifetime(lifetime) } function groupChatEntryId() { @@ -7532,7 +7552,7 @@ async function ensureGroupChatSession(group, member, requestMember = member, occ const stored = res.session_key || known if (stored) { - updateGroupChat(group, current => { + updateGroupChat(occurrence?.group || group, current => { if (occurrence && !groupOccurrenceCurrent(occurrence)) return current current.sessions = { ...(current.sessions || {}), [key]: stored } current.sessionOwners = { ...(current.sessionOwners || {}), [key]: groupSessionOwner(requestMember) } @@ -7560,7 +7580,7 @@ async function ensureGroupChatSession(group, member, requestMember = member, occ const stored = created?.stored_session_id || null if (stored) { - updateGroupChat(group, r => { + updateGroupChat(occurrence?.group || group, r => { if (occurrence && !groupOccurrenceCurrent(occurrence)) return r r.sessions = { ...(r.sessions || {}), [key]: stored } r.sessionOwners = { ...(r.sessionOwners || {}), [key]: groupSessionOwner(requestMember) } @@ -7845,7 +7865,7 @@ async function waitForGroupTurnCollector(group, memberKey, collector, occurrence const epoch = $groupChats.get()[group]?.epoch || 0 let timer, unbind const cancelled = () => { - const room = $groupChats.get()[group] + const room = $groupChats.get()[occurrence?.group || group] return !room || room.tombstone || Boolean(room.holds?.[memberKey]) || (occurrence?.userIds ? !groupOccurrenceCanSubmit(occurrence) : (room.epoch || 0) !== epoch) } @@ -7874,6 +7894,10 @@ function groupTurnMarkerBlocksDispatch(room, memberKey, coordinator) { function consumeGroupTurnMarker(group, memberKey, marker, published = false) { if ($groupChats.get()[group]?.stranded?.[memberKey] !== marker) return false const owned = [...(groupRoomCoordinators.get(group)?.occurrences || [])].find(o => o.id === marker?.occurrence_id) + // session.interrupt targets the runtime's current turn when applied. Even + // exact predecessor terminal proof cannot let a successor reuse that runtime + // while an older hot generation or a cold marker's control is still pending. + if (groupOccurrenceHasPendingInterrupt(owned) || interruptingGroupTurnMarkers.has(marker)) return false const coldDrive = owned ? null : restoreGroupDrive(group, marker) let consumed = false updateGroupChat(group, room => { @@ -7901,7 +7925,8 @@ function consumeGroupTurnMarker(group, memberKey, marker, published = false) { } function groupTurnMarkerIntentIsCurrent(room, marker) { - if (!room || room.tombstone || room.holds?.[marker.delivery.member_key]) return false + if (!room || room.tombstone || marker.stop_requested || marker.hold_requested || + room.holds?.[marker.delivery.member_key]) return false const newerUser = groupTurnHasNewerUser(room, marker.thread, marker.user_ids, marker.anchor_id, marker.input_version) return shouldCommitMemberTurn(marker.epoch ?? 0, room.epoch || 0, newerUser, marker.input_version !== undefined || Array.isArray(marker.user_ids)) @@ -8004,7 +8029,7 @@ function syncGroupClarify(group, member, state, requestMember = member) { return true } -/** Drop every mirrored clarify belonging to `group` (disband/rename). */ +/** Drop every mirrored clarify belonging to `group` (disband/Stop). */ function clearGroupClarify(group) { const all = $groupClarify.get() const next = {} @@ -8098,11 +8123,15 @@ async function respondGroupClarify(entry, member, answers, occurrence, fence) { requireCurrentGroupQuestion(entry) attempted = true if (entry.kind === 'approval') { - await requestForBot(member, 'approval.respond', { + const result = await requestForBot(member, 'approval.respond', { session_id: entry.sessionId || undefined, request_id: entry.requestId, choice: typeof answers === 'string' && answers ? answers : 'deny' }) + if (!Number.isSafeInteger(result?.resolved) || result.resolved <= 0) { + throw new Error(result?.resolved === 0 ? 'The approval is no longer pending.' + : 'The approval acknowledgement is unconfirmed.') + } } else if (entry.questions && entry.questions.length) { for (const question of entry.questions) { const qid = question?.qid ?? question?.id @@ -8164,6 +8193,49 @@ const groupRoomCoordinators = new Map() const groupRuntimeSessionOwners = new Map() const GROUP_CHAT_PARALLEL_CEILING = 4 +function groupNameForLifetime(owner, rooms = $groupChats.get()) { + const matches = room => room && !room.tombstone && + (room.roomId || null) === owner.roomId && (room.coordinationId || null) === owner.roomToken + if (matches(rooms[owner.group])) return owner.group + // Legacy owners also need their minted coordination token. Never follow a + // display-name reuse or an unidentifiable room into a replacement lifetime. + if (!owner.roomId && !owner.roomToken) return null + const names = Object.keys(rooms).filter(name => matches(rooms[name])) + return names.length === 1 ? names[0] : null +} + +function rehomeGroupRuntimeOwners(rooms) { + for (const [oldName, coordinator] of [...groupRoomCoordinators]) { + const next = groupNameForLifetime(coordinator, rooms) + if (!next || next === oldName) continue + if (groupRoomCoordinators.get(oldName) === coordinator) groupRoomCoordinators.delete(oldName) + coordinator.group = next + groupRoomCoordinators.set(next, coordinator) + for (const occurrence of coordinator.occurrences) occurrence.group = next + for (const drive of coordinator.capturedDrives.values()) drive.group = next + } + const cards = {}, needs = { ...$groupNeedsYou.get() }, activity = { ...$groupActivity.get() } + let cardsChanged = false, needsChanged = false, activityChanged = false + for (const entry of Object.values($groupClarify.get())) { + const next = groupNameForLifetime({ group: entry.group, roomId: entry.roomId, roomToken: entry.roomToken }, rooms) + if (next && next !== entry.group) { + entry.group = next // preserve pending response/card identity + cardsChanged = true + } + cards[`${entry.group}::${entry.memberKey}`] = entry + } + for (const [oldName, room] of Object.entries($groupChats.get())) { + const next = groupNameForLifetime({ group: oldName, roomId: room.roomId || null, + roomToken: room.coordinationId || null }, rooms) + if (!next || next === oldName) continue + if (oldName in needs) { needs[next] = needs[oldName]; delete needs[oldName]; needsChanged = true } + if (oldName in activity) { activity[next] = activity[oldName]; delete activity[oldName]; activityChanged = true } + } + if (cardsChanged) $groupClarify.set(cards) + if (needsChanged) $groupNeedsYou.set(needs) + if (activityChanged) $groupActivity.set(activity) +} + function groupRoomCoordinator(group) { let room = $groupChats.get()[group] || {} if (!room.roomId && !room.coordinationId) { @@ -8197,7 +8269,7 @@ function registerGroupOccurrence(group, member, thread, deliveryResult = {}, pre ready: new Promise(resolve => { ready = resolve }) } occurrence.settle = settle occurrence.markReady = ready - occurrence.memberLock = groupSourceSessionKey(captured, `room:${room.roomId || group}`) + occurrence.memberLock = groupSourceSessionKey(captured, `room:${room.roomId || room.coordinationId || group}`) occurrence.inputEndId = room.log?.length ? groupChatSyncEntryKey(room.log.at(-1)) : null occurrence.inputVersion = room.threadInputVersions?.[occurrence.thread] || 0 occurrence.userIds = (room.log || []).filter(e => groupIsUserInstruction(e) && groupThreadOf(e) === occurrence.thread).map(groupChatSyncEntryKey) @@ -8207,11 +8279,16 @@ function registerGroupOccurrence(group, member, thread, deliveryResult = {}, pre return occurrence } -function groupOccurrenceCurrent(occurrence) { +function groupOccurrenceRoomIsCurrent(occurrence) { const room = $groupChats.get()[occurrence.group] - return !occurrence.cancelled && room && !room.tombstone && + return room && !room.tombstone && (room.roomId || null) === occurrence.roomId && - (room.coordinationId || null) === occurrence.roomToken && !room.holds?.[occurrence.memberKey] + (room.coordinationId || null) === occurrence.roomToken +} + +function groupOccurrenceCurrent(occurrence) { + return !occurrence.cancelled && groupOccurrenceRoomIsCurrent(occurrence) && + !$groupChats.get()[occurrence.group].holds?.[occurrence.memberKey] } function groupOccurrenceCanSubmit(occurrence) { @@ -8228,7 +8305,7 @@ function paintGroupOccurrences(coordinator) { (current.coordinationId || null) !== coordinator.roomToken) return const turns = [...coordinator.occurrences].filter(o => !o.released).map(o => ({ id: o.id, memberKey: o.memberKey, member: o.captured.member.name, - phase: o.cancelled ? (groupOccurrenceStopConfirmed(o) ? 'stopping' : 'stop-unconfirmed') : o.phase, thread: o.thread })) + phase: o.cancelled ? (groupOccurrenceInterruptApplied(o) ? 'stopping' : 'stop-unconfirmed') : o.phase, thread: o.thread })) updateGroupChat(coordinator.group, room => { room.turns = turns room.turn = turns.find(o => o.phase === 'running' || o.phase === 'starting')?.member || null @@ -8237,16 +8314,26 @@ function paintGroupOccurrences(coordinator) { } function groupInterruptConfirmed(reply) { + // This confirms application of cancellation, not execution retirement. return reply?.status === 'interrupted' } function interruptStoppedGroupMarker(group, memberKey, marker, target) { const prior = interruptingGroupTurnMarkers.get(marker) if (prior) return prior + const room = $groupChats.get()[group] + const lifetime = { group, roomId: room?.roomId || null, roomToken: room?.coordinationId || null } const pending = Promise.resolve().then(() => requestForBot(target, 'session.interrupt', { session_id: marker.delivery.accepted_turn.session_id })).then(reply => { const confirmed = groupInterruptConfirmed(reply) - if (confirmed) consumeGroupTurnMarker(group, memberKey, marker) + const currentGroup = groupNameForLifetime(lifetime) + if (confirmed && currentGroup) updateGroupChat(currentGroup, room => { + if (room.stranded?.[memberKey] === marker) { + marker.interrupt_applied = true + room.stranded = { ...room.stranded } + } + return room + }, { sync: false }) return confirmed }, () => false).finally(() => { if (interruptingGroupTurnMarkers.get(marker) === pending) interruptingGroupTurnMarkers.delete(marker) @@ -8256,10 +8343,19 @@ function interruptStoppedGroupMarker(group, memberKey, marker, target) { } function groupOccurrenceStopConfirmed(occurrence) { - if (occurrence.submissionPending || occurrence.answerPromise) return false + if (occurrence.submissionPending || occurrence.answerPromise || groupOccurrenceHasPendingInterrupt(occurrence)) return false if (occurrence.terminalObserved) return true if (!occurrence.submitAttempted) return Boolean(occurrence.collectorDone) - return occurrence.interrupts.get(`${occurrence.runtime}::${occurrence.admissionVersion || 0}`)?.confirmed === true + return false // A matching interrupt ACK still needs exact terminal evidence. +} + +function groupOccurrenceHasPendingInterrupt(occurrence) { + return [...(occurrence?.interrupts?.values() || [])].some(attempt => attempt.pending) +} + +function groupOccurrenceInterruptApplied(occurrence) { + return !occurrence.submissionPending && !occurrence.answerPromise && + occurrence.interrupts.get(`${occurrence.runtime}::${occurrence.admissionVersion || 0}`)?.confirmed === true } function markGroupOccurrenceStop(occurrence) { @@ -8300,7 +8396,7 @@ async function interruptGroupOccurrence(occurrence, runtime = occurrence.runtime function finishGroupOccurrence(occurrence) { const c = occurrence.coordinator - if (occurrence.released) return + if (occurrence.released || groupOccurrenceHasPendingInterrupt(occurrence)) return occurrence.released = true occurrence.releaseLease?.() if (c.members.get(occurrence.memberLock) === occurrence) c.members.delete(occurrence.memberLock) @@ -8399,7 +8495,12 @@ function retainUnresolvedGroupOccurrence(occurrence) { return false } markGroupOccurrenceStop(occurrence) - } else if (!groupOccurrenceCurrent(occurrence) || marker?.occurrence_id !== occurrence.id) return false + if (groupOccurrenceRoomIsCurrent(occurrence) && groupTurnDeliveryKey(marker?.delivery)) { + // Late admission/collector settlement may supply the first exact receipt + // after Stop. Reconcile it without making control completion wait on a read. + void harvestStrandedGroupReply(occurrence.group, occurrence.captured.member).catch(() => undefined) + } + } else if (!groupOccurrenceRoomIsCurrent(occurrence) || marker?.occurrence_id !== occurrence.id) return false // A timeout/unavailable projection is not capacity evidence. Waiting has // already surrendered only its worker; running/unknown keeps that worker. occurrence.markReady() @@ -8462,6 +8563,7 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima return null } while (true) { + group = occurrence?.group || group const current = $groupChats.get()[group] || {} const prior = current.stranded?.[memberKey] const priorCollector = prior && collectingGroupTurnMarkers.get(prior) @@ -8480,6 +8582,7 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima // Another waiter may have reserved the member. Recheck synchronously // before any session preparation, attachments, or prompt admission. } + group = occurrence?.group || group const roomAtDispatch = $groupChats.get()[group] || {} const dispatchEpoch = roomAtDispatch.epoch || 0 let marker = { delivery: { accepted_turn: null, member_key: memberKey, owner }, @@ -8499,6 +8602,7 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima let submitAttempted = false try { const prepared = await ensureGroupChatSession(group, member, requestMember, occurrence) + group = occurrence?.group || group let { runtime } = prepared const { stored } = prepared if (occurrence) { @@ -8513,6 +8617,7 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima if (occurrence.cancelled) await interruptGroupOccurrence(occurrence) } const beforeSubmit = () => { + group = occurrence?.group || group const room = $groupChats.get()[group] || {} return (!occurrence || groupOccurrenceCanSubmit(occurrence)) && room.stranded?.[memberKey] === marker && !room.tombstone && !((room.epoch || 0) !== dispatchEpoch && room.holds?.[memberKey]) @@ -8522,6 +8627,7 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima return discarded() } runtime = await requireGroupTurnProtocol(requestMember, prepared) + group = occurrence?.group || group if (occurrence && runtime !== occurrence.runtime) { const sessionLock = groupSourceSessionKey(captured, runtime) const prior = groupRuntimeSessionOwners.get(sessionLock) @@ -8583,6 +8689,7 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima if (occurrence.cancelled) { markGroupOccurrenceStop(occurrence); await interruptGroupOccurrence(occurrence) } } } + group = occurrence?.group || group const previous = marker marker = { ...marker, runtime: submitted.runtime, delivery: { ...marker.delivery, accepted_turn: submitted.acceptedTurn } } @@ -8605,15 +8712,11 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima let progress = 'working' while (Date.now() < deadline) { await new Promise(resolve => setTimeout(resolve, GROUP_TURN_POLL_MS)) + group = occurrence?.group || group const roomDuringPoll = $groupChats.get()[group] || {} if (roomDuringPoll.stranded?.[memberKey] !== marker) return discarded() - if (occurrence && !groupOccurrenceCurrent(occurrence)) { + if (occurrence && (occurrence.cancelled || !groupOccurrenceRoomIsCurrent(occurrence))) { if (occurrence.cancelled) markGroupOccurrenceStop(occurrence) - else consumeGroupTurnMarker(group, memberKey, marker) - return discarded() - } - if ((roomDuringPoll.epoch || 0) !== dispatchEpoch && roomDuringPoll.holds?.[memberKey]) { - consumeGroupTurnMarker(group, memberKey, marker) return discarded() } let state @@ -8624,28 +8727,21 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima session_id: submitted.acceptedTurn.session_id, profile: member.name, accepted_turn: submitted.acceptedTurn }) } catch (error) { + group = occurrence?.group || group const roomAfterError = $groupChats.get()[group] || {} if (roomAfterError.stranded?.[memberKey] !== marker) return discarded() if (occurrence?.cancelled) { markGroupOccurrenceStop(occurrence); return discarded() } - if ((roomAfterError.epoch || 0) !== dispatchEpoch && roomAfterError.holds?.[memberKey]) { - consumeGroupTurnMarker(group, memberKey, marker) - return discarded() - } if (!groupTurnMarkerIntentIsCurrent(roomAfterError, marker)) return discarded() // Observation failure grants no replay or resume authority. Surface it // now and retain the accepted receipt for an explicit later harvest. throw groupTurnOutcomeError({ state: 'unavailable', reason: `Could not observe member turn: ${error?.message || 'gateway poll failed'}` }) } + group = occurrence?.group || group const roomAfterResume = $groupChats.get()[group] || {} if (roomAfterResume.stranded?.[memberKey] !== marker) return discarded() - if (occurrence && !groupOccurrenceCurrent(occurrence)) { + if (occurrence && (occurrence.cancelled || !groupOccurrenceRoomIsCurrent(occurrence))) { if (occurrence.cancelled) markGroupOccurrenceStop(occurrence) - else consumeGroupTurnMarker(group, memberKey, marker) - return discarded() - } - if ((roomAfterResume.epoch || 0) !== dispatchEpoch && roomAfterResume.holds?.[memberKey]) { - consumeGroupTurnMarker(group, memberKey, marker) return discarded() } const outcome = readGroupTurnOutcome(state, marker.delivery) @@ -8689,7 +8785,9 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima } paintGroupOccurrences(occurrence.coordinator) } else if (occurrence) { - if (occurrence.phase === 'waiting') await reserveGroupResumeWorker(occurrence) + if (occurrence.phase === 'waiting' && !['complete', 'error', 'interrupted'].includes(outcome.state)) { + await reserveGroupResumeWorker(occurrence) + } if (occurrence.resumeFence === resumeFenceAtPoll && responseAcknowledgedAtPoll) occurrence.resumeFence = null paintGroupOccurrences(occurrence.coordinator) } @@ -8713,6 +8811,7 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima // Observation expiry cannot invalidate a retained unavailable child card. return null } catch (error) { + group = occurrence?.group || group if (!submitAttempted || error?.code === 4090 || error?.code === 'POOL_CAPACITY_EXCEEDED') { if (error?.code === 4090 || error?.code === 'POOL_CAPACITY_EXCEEDED') { if (occurrence) occurrence.submitAttempted = false @@ -8742,9 +8841,11 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima async function harvestStrandedGroupReply(group, member) { const memberKey = groupMemberKey(member) const room = $groupChats.get()[group] || {} + const lifetime = { group, roomId: room.roomId || null, roomToken: room.coordinationId || null } const marker = room.stranded?.[memberKey] - if (marker === undefined || collectingGroupTurnMarkers.has(marker) || - (room.holds?.[memberKey] && !marker?.stop_requested)) return + const reconciliationOnly = marker?.stop_requested || marker?.hold_requested + if (marker === undefined || (collectingGroupTurnMarkers.has(marker) && !reconciliationOnly) || + (room.holds?.[memberKey] && !reconciliationOnly)) return if (marker?.room_token && marker.room_token !== room.coordinationId) { reportUnavailableGroupTurn(group, member, marker, 'Legacy room lifetime is unavailable after reload or replacement') return @@ -8772,23 +8873,25 @@ async function harvestStrandedGroupReply(group, member) { session_id: marker.delivery.accepted_turn.session_id, profile: captured.member.name, accepted_turn: marker.delivery.accepted_turn }) } catch (error) { + group = groupNameForLifetime(lifetime) || group const current = $groupChats.get()[group] if (current?.stranded?.[memberKey] !== marker || !groupTurnMarkerIntentIsCurrent(current, marker)) return reportUnavailableGroupTurn(group, member, marker, `Could not observe member turn: ${error?.message || 'gateway poll failed'}`) return } + group = groupNameForLifetime(lifetime) || group if ($groupChats.get()[group]?.stranded?.[memberKey] !== marker || ownedOccurrence?.answerPromise) return const outcome = readGroupTurnOutcome(state, marker.delivery) const freshResumeRead = !ownedOccurrence?.resumeFence || (ownedOccurrence.resumeFence === resumeFenceAtRead && responseAcknowledgedAtRead) - if (marker.stop_requested || ownedOccurrence?.cancelled) { - // A stopped receipt is reconciliation only, including after reload. No + if (reconciliationOnly || ownedOccurrence?.cancelled) { + // A stopped/held receipt is reconciliation only, including after reload. No // reply, attention, input consumption, clarify card or prompt replay. if (['complete', 'error', 'interrupted'].includes(outcome.state) && !ownedOccurrence?.submissionPending && freshResumeRead) { if (ownedOccurrence) ownedOccurrence.terminalObserved = true - consumeGroupTurnMarker(group, memberKey, marker) + if (consumeGroupTurnMarker(group, memberKey, marker) && ownedOccurrence) settleGroupPublication(ownedOccurrence) } else if (ownedOccurrence && outcome.state === 'waiting' && !ownedOccurrence.submissionPending && freshResumeRead) { ownedOccurrence.resumeFence = null @@ -9107,6 +9210,7 @@ function groupTurnHasNewerUser(room, thread, userIds, anchorId, inputVersion) { * settle, so an old cleanup can never interrupt a replacement occurrence. */ async function stopGroupThread(group, thread, members = null) { const room = $groupChats.get()[group] || {} + const lifetime = { group, roomId: room.roomId || null, roomToken: room.coordinationId || null } const coordinator = groupRoomCoordinators.get(group) const occurrences = coordinator && coordinator.roomId === (room.roomId || null) && coordinator.roomToken === (room.coordinationId || null) ? [...coordinator.occurrences] : [] @@ -9141,7 +9245,7 @@ async function stopGroupThread(group, thread, members = null) { } if (coordinator) pumpGroupOccurrences(coordinator) const interrupts = occurrences.map(occurrence => interruptGroupOccurrence(occurrence, occurrence.runtime, true)) - let coldUnconfirmed = 0 + let legacyUnconfirmed = 0 // A cold reload has no runtime registry: accepted durable receipts still // identify exact runtime and source, never a display-name roster guess. for (const [memberKey, marker] of Object.entries(room.stranded || {})) { @@ -9155,11 +9259,10 @@ async function stopGroupThread(group, thread, members = null) { if (r.stranded?.[memberKey] === marker) { marker.stop_requested = true; r.stranded = { ...r.stranded } } return r }, { sync: false }) - if (!groupTurnDeliveryKey(marker?.delivery)) { coldUnconfirmed++; continue } + if (!groupTurnDeliveryKey(marker?.delivery)) continue const owner = marker.delivery.owner const target = Object.freeze({ ...owner, ...(owner.route ? { route: Object.freeze({ ...owner.route }) } : {}) }) - interrupts.push(interruptStoppedGroupMarker(group, memberKey, marker, target) - .then(confirmed => { if (!confirmed) coldUnconfirmed++ })) + interrupts.push(interruptStoppedGroupMarker(group, memberKey, marker, target)) } // Legacy transient room.turn has no acceptance proof. Retain its old Stop // behavior only when the name is unambiguous and no occurrences exist. @@ -9169,30 +9272,67 @@ async function stopGroupThread(group, thread, members = null) { const member = candidates[0] const sid = room.sessions?.[groupMemberKey(member)] if (sid) interrupts.push(Promise.resolve().then(() => requestForBot(member, 'session.interrupt', { session_id: sid })) - .then(reply => { if (!groupInterruptConfirmed(reply)) coldUnconfirmed++ }, () => { coldUnconfirmed++ })) + .then(() => { legacyUnconfirmed++ }, () => { legacyUnconfirmed++ })) } } for (const interrupt of interrupts) await interrupt + const missingLifetime = () => { + const pending = Math.max(1, occurrences.filter(o => !o.released && !groupOccurrenceStopConfirmed(o)).length) + return { status: 'unconfirmed', unconfirmed: pending, pending } + } + group = groupNameForLifetime(lifetime) + if (!group) return missingLifetime() + // The ACK only applied cancellation. The existing exact-turn collector owns + // retirement; do not infer a vacant worker from a successful control write. + for (const [memberKey, marker] of Object.entries($groupChats.get()[group]?.stranded || {})) { + if (!marker?.stop_requested || !groupTurnDeliveryKey(marker.delivery)) continue + if ($groupChats.get()[group]?.stranded?.[memberKey] !== marker) continue + const occurrence = occurrences.find(o => o.id === marker.occurrence_id) + const owner = marker.delivery.owner + const target = occurrence?.captured.member || Object.freeze({ ...owner, + ...(owner.route ? { route: Object.freeze({ ...owner.route }) } : {}) }) + // A slow/unavailable poll is observation, not part of the interrupt ACK. + // Its exact terminal may retire custody later through the same collector. + void harvestStrandedGroupReply(group, target).catch(() => undefined) + } for (const occurrence of occurrences) { - const marker = $groupChats.get()[group]?.stranded?.[occurrence.memberKey] if (occurrence.collectorDone) { if (occurrence.answerPromise) await occurrence.answerPromise.catch(() => undefined) await interruptGroupOccurrence(occurrence) + group = groupNameForLifetime(lifetime) + if (!group) return missingLifetime() + const settledMarker = $groupChats.get()[group]?.stranded?.[occurrence.memberKey] + if (settledMarker?.occurrence_id === occurrence.id && groupTurnDeliveryKey(settledMarker.delivery)) { + // A late human response advanced the generation after the first read + // was fenced. Observe only after its producer and final interrupt settle. + void harvestStrandedGroupReply(group, occurrence.captured.member).catch(() => undefined) + } if (groupOccurrenceStopConfirmed(occurrence)) { + const marker = $groupChats.get()[group]?.stranded?.[occurrence.memberKey] if (marker?.occurrence_id === occurrence.id) consumeGroupTurnMarker(group, occurrence.memberKey, marker) finishGroupOccurrence(occurrence) } } } - const unconfirmed = coldUnconfirmed + occurrences.filter(o => !o.released && !groupOccurrenceStopConfirmed(o)).length + const pendingMarkers = Object.values($groupChats.get()[group]?.stranded || {}).filter(marker => marker?.stop_requested) + const unconfirmed = legacyUnconfirmed + pendingMarkers.filter(marker => { + const occurrence = occurrences.find(o => o.id === marker.occurrence_id) + return occurrence ? !groupOccurrenceInterruptApplied(occurrence) : !marker.interrupt_applied + }).length + occurrences.filter(o => !o.released && !groupOccurrenceStopConfirmed(o) && + !pendingMarkers.some(marker => marker.occurrence_id === o.id) && !groupOccurrenceInterruptApplied(o)).length + const pending = pendingMarkers.length + occurrences.filter(o => !o.released && !groupOccurrenceStopConfirmed(o) && + !pendingMarkers.some(marker => marker.occurrence_id === o.id)).length if (coordinator) paintGroupOccurrences(coordinator) - const result = { status: unconfirmed ? 'unconfirmed' : 'stopped', unconfirmed } + const result = { status: unconfirmed ? 'unconfirmed' : pending ? 'stopping' : 'stopped', unconfirmed, pending } const current = $groupChats.get()[group] if (current && (current.roomId || null) === (room.roomId || null) && (current.coordinationId || null) === (room.coordinationId || null)) { - recordGroupActivity(group, { kind: unconfirmed ? 'stop-unconfirmed' : 'stopped', + recordGroupActivity(group, { kind: unconfirmed ? 'stop-unconfirmed' : pending ? 'stopping' : 'stopped', member: 'You', thread: thread || null, epoch: stoppedRoom.epoch }) } + if (pending && typeof window !== 'undefined') { + void harvestStrandedUntilSettled(group, occurrences.map(o => o.captured.member).concat(roster), thread) + } return result } @@ -9280,6 +9420,10 @@ function captureGroupDrive(group, members, thread) { /** A frozen round, never a raw Promise.all over roster members. The room * coordinator is the sole start owner; each completion publishes immediately. */ function runGroupChatRounds(group, members, thread, capturedDrive) { + if (capturedDrive) { + group = groupNameForLifetime({ ...capturedDrive, group }) + if (!group) return Promise.resolve() // deletion/replacement cannot revive a captured drive + } const drive = capturedDrive || captureGroupDrive(group, members, thread) const coordinator = groupRoomCoordinator(group) const key = groupDriveKey(drive) @@ -9295,6 +9439,7 @@ function runGroupChatRounds(group, members, thread, capturedDrive) { const running = Promise.resolve().then(() => driveFrozenGroupRounds(group, members, drive, coordinator)) coordinator.drives.set(key, running) void running.finally(() => { + group = coordinator.group if (coordinator.drives.get(key) === running) coordinator.drives.delete(key) drive.running = false if (!Object.values($groupChats.get()[group]?.stranded || {}).some(m => m?.drive_key === key)) { @@ -9340,6 +9485,7 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { const { thread } = drive members = drive.members.map(captured => captured.member) const sameRoom = () => { + group = coordinator.group const room = $groupChats.get()[group] return room && !room.tombstone && (room.roomId || null) === drive.roomId && (room.coordinationId || null) === drive.roomToken @@ -9421,6 +9567,7 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { job.phase = 'queued' const complete = runGroupChatMemberTurn(group, job.captured.member, job.prompt, thread, job.images, job.deliveryResult, job).then(reply => { + if (!sameRoom()) return const room = $groupChats.get()[group] if (job.cancelled) return // Stop owns the confirmation/unconfirmed activity. if (!sameRoom() || job.deliveryResult.discarded || @@ -9441,6 +9588,7 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { if (job.wasWaiting) scheduleGroupDriveContinuation(group, drive) } }, error => { + if (!sameRoom()) return const room = $groupChats.get()[group] if (!sameRoom() || job.cancelled || !shouldCommitMemberTurn(job.epoch, room.epoch || 0, groupTurnHasNewerUser(room, thread, job.userIds, job.inputEndId, job.inputVersion), true)) return @@ -9539,10 +9687,15 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { async function harvestStrandedUntilSettled(group, members, thread) { const HARVEST_INTERVAL_MS = 5000 const HARVEST_MAX_TRIES = 60 + const initial = $groupChats.get()[group] + if (!initial) return + const lifetime = { group, roomId: initial.roomId || null, roomToken: initial.coordinationId || null } for (let attempt = 0; attempt < HARVEST_MAX_TRIES; attempt++) { await new Promise(resolve => window.setTimeout(resolve, HARVEST_INTERVAL_MS)) + group = groupNameForLifetime(lifetime) + if (!group) return const room = $groupChats.get()[group] if (!room || room.running) { @@ -9610,29 +9763,35 @@ function sendToGroupChat(group, members, text, thread, images) { { at: sent?.at, byMessageId: sent?.id, thread: target }, members.map(member => groupMemberKey(member)) ) + for (const [memberKey, marker] of Object.entries(room.stranded || {})) { + if (room.holds?.[memberKey] && marker && typeof marker === 'object') { + // A future-turn hold suppresses this result; it is not death evidence. + // Keep the live collector/CAS identity and persist reconciliation intent. + marker.hold_requested = true + room.stranded = { ...room.stranded } + } + } return room }) recordGroupActivity(group, { kind: 'queued', member: 'You', thread: target }) const capturedDrive = captureGroupDrive(group, members, target) + const clearFailedDrive = () => { + const currentGroup = groupNameForLifetime({ ...capturedDrive, group }) + if (!currentGroup) return + updateGroupChat(currentGroup, room => { + if ((room.epoch || 0) === capturedDrive.epoch) room.running = false + return room + }) + } if (!wasRunning) { - void runGroupChatRounds(group, members, target, capturedDrive).catch(() => { - updateGroupChat(group, r => { - if ((r.roomId || null) === capturedDrive.roomId && (r.epoch || 0) === capturedDrive.epoch) r.running = false - return r - }) - }) + void runGroupChatRounds(group, members, target, capturedDrive).catch(clearFailedDrive) } else { // A loop is live; it bails at its next boundary. Chain the fresh loop // after a short settle so exactly one drive owns the room. setTimeout(() => { - void runGroupChatRounds(group, members, target, capturedDrive).catch(() => { - updateGroupChat(group, r => { - if ((r.roomId || null) === capturedDrive.roomId && (r.epoch || 0) === capturedDrive.epoch) r.running = false - return r - }) - }) + void runGroupChatRounds(group, members, target, capturedDrive).catch(clearFailedDrive) }, 250) } @@ -13498,7 +13657,11 @@ function GroupImageControls({ image, onImage, seedName, seedMembers }) { * the room record. Both apply on Save so a cancelled dialog changes nothing. */ function GroupChatSettingsDialog({ group, members, open, onClose, onRenamed }) { const rooms = useValue($groupChats) - const current = (rooms[group] || {}).image || null + const displayedRoom = rooms[group] + const displayedLifetime = { group, roomId: displayedRoom?.roomId || null, + roomToken: displayedRoom?.coordinationId || null } + const metadataOnlyGroup = !displayedRoom && knownGroups($botMeta.get()).includes(group) + const current = (displayedRoom || {}).image || null const [name, setName] = useState(group) const [image, setImage] = useState(current) @@ -13511,11 +13674,31 @@ function GroupChatSettingsDialog({ group, members, open, onClose, onRenamed }) { }, [open, group]) const save = async () => { - const finalName = await renameGroupChat(group, name, members) + let savingGroup = group + if (displayedRoom) { + savingGroup = groupNameForLifetime(displayedLifetime) + // A button rendered for a deleted/replaced room cannot adopt today's + // same-name record. Unidentified legacy rooms require the same object + // until their existing canonical lifetime token is minted below. + if (!savingGroup || (!displayedLifetime.roomId && !displayedLifetime.roomToken && + $groupChats.get()[group] !== displayedRoom)) return + } else if (!metadataOnlyGroup || $groupChats.get()[group] || !knownGroups($botMeta.get()).includes(group)) { + return + } + const original = $groupChats.get()[savingGroup] + if (original?.tombstone) return + const savingRoom = original?.roomId || original?.coordinationId ? original : updateGroupChat(savingGroup, room => { + room.coordinationId = groupChatEntryId() + return room + }, { sync: false }) + const lifetime = { group: savingGroup, roomId: savingRoom.roomId || null, roomToken: savingRoom.coordinationId || null } + const renamed = await renameGroupChat(savingGroup, name, members) - if (finalName === null) { - return // collision — dialog stays open for a different name + if (renamed === null) { + return // collision or removed lifetime — the dialog cannot update another room } + const finalName = groupNameForLifetime(lifetime) + if (finalName === null) return if (image !== current) { setGroupChatImage(finalName, image) @@ -14009,7 +14192,6 @@ function GroupMentionInput({ members, onChange, onSubmitDraft, value, ...inputPr * (once/session/always/deny) as buttons — no free text; approvals are a * closed choice. Answer sends via the member's own source. */ function GroupClarifyCard({ entry, members }) { - const { group } = entry const isApproval = entry.kind === 'approval' const member = members.find(m => groupMemberKey(m) === entry.memberKey) || members.find(m => m.name === entry.member) const [drafts, setDrafts] = useState({}) @@ -14062,7 +14244,10 @@ function GroupClarifyCard({ entry, members }) { : questions .map(q => (questions.length > 1 ? `${q.question}: ${answerFor(q)}` : answerFor(q))) .join('\n') - appendGroupChatEntry(group, { kind: 'user', name: 'You' }, summary, entry.thread || 'legacy', + const currentRoom = $groupChats.get()[entry.group] + if (!currentRoom || currentRoom.tombstone || (currentRoom.roomId || null) !== entry.roomId || + (currentRoom.coordinationId || null) !== entry.roomToken) return + appendGroupChatEntry(entry.group, { kind: 'user', name: 'You' }, summary, entry.thread || 'legacy', undefined, undefined, entry.receipt?.delivery) } catch (err) { host.notify({ kind: 'error', message: `Could not send the answer to @${botHandle(entry.member, member)}: ${err?.message || err}` }) @@ -14587,7 +14772,9 @@ function GroupChatWorkspace({ group, members, onBack, visible = true }) { const result = await stopGroupThread(group, latestActivity?.thread || null, memberDescriptors()) host.notify(result.status === 'stopped' ? { kind: 'success', message: `Stopped ${group} — remaining turns are held until you resume` } - : { kind: 'info', message: `Held ${group} — ${result.unconfirmed} interruption(s) are unconfirmed. Stop can retry.` }) + : result.status === 'stopping' + ? { kind: 'info', message: `Stopping ${group} — waiting for the remaining turns to finish.` } + : { kind: 'info', message: `Held ${group} — ${result.unconfirmed} interruption(s) are unconfirmed. Stop can retry.` }) } const activityPanel = jsxs('div', { @@ -16733,7 +16920,7 @@ const groupTurnRuntime = { groupMemberKey, updateGroupChat, groupBlockedMembers, GroupBlockedNotice, CreateGroupChatDialog, createFreshGroupChat, groupComposerDraftKey, groupComposerDraftSnapshot, updateGroupComposerDraft, GroupChatWorkspace, GroupClarifyCard, - groupRoomCanStop, + groupRoomCanStop, renameGroupChat, pullGroupChatServerState, GroupChatSettingsDialog, $groupChatWorkspace, bindGroupTurnPorts(ports) { ({ Date, setTimeout, clearTimeout, setInterval, clearInterval, document } = createGroupTurnPorts(ports)) }, @@ -17339,5 +17526,3 @@ export default { }) } } - - diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-lifecycle-boundaries.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-lifecycle-boundaries.test.mjs new file mode 100644 index 000000000000..f8e78e010676 --- /dev/null +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-lifecycle-boundaries.test.mjs @@ -0,0 +1,580 @@ +import assert from 'node:assert/strict' +import test from 'node:test' +import { harness, members, deferred, flush, drive } from './stop-custody-harness.mjs' + +// Offline only: the shipped plugin runs against fake gateway/profile ports and +// a bounded fake clock. No provider, backend process or installed state is used. +const bounded = (name, body) => test(name, { timeout: 10000 }, body) +const clone = value => JSON.parse(JSON.stringify(value)) +const receipt = (h, member = h.roster[0], name = 'Room') => + h.gc.$groupChats.get()[name]?.stranded?.[h.gc.groupMemberKey(member)] +const coordinator = (h, name = 'Room') => h.gc.groupRoomCoordinators.get(name) +const posts = (h, name = 'Room') => h.gc.$groupChats.get()[name]?.log.filter(e => e.from.kind === 'member') || [] +const sessionFor = (h, profile = 'bot1') => [...h.sessions.values()].find(s => s.profile === profile) +const cards = h => Object.values(h.gc.$groupClarify.get()) + +async function settle(h, pending, name = 'Room') { + for (let i = 0; i < 30 && coordinator(h, name)?.occurrences.size; i++) { + for (const session of h.sessions.values()) { session.state = 'complete'; session.pending = null } + await h.advance() + for (const member of h.roster) await h.gc.harvestStrandedGroupReply(name, member) + } + if (pending) await pending +} + +for (const terminal of ['complete', 'error', 'interrupted']) { + bounded(`GC-TERM: saturated waiting collector retires ${terminal} without a worker reservation`, async () => { + const h = await harness(members(6)), pending = drive(h) + await flush() + assert.equal(h.rpc('prompt.submit').length, 4) + const parked = sessionFor(h), accepted = clone(parked.ref) + parked.state = 'waiting'; parked.pending = { request_id: 'question-1', question: 'Choose?' } + await h.advance() + assert.equal(h.rpc('prompt.submit').length, 5, 'the fifth member uses the positively waiting slot') + assert.equal(coordinator(h).active, 4) + assert.equal(coordinator(h).queue.length, 1) + parked.state = terminal; parked.pending = null; parked.text = 'TERMINAL_WITH_ALL_SLOTS_BUSY' + await h.advance(); await flush() + assert.equal(receipt(h), undefined, 'exact terminal does not wait behind the four busy workers') + assert.equal(h.leases.find(l => l.route.targetProfile === 'bot1').releases, 1) + assert.equal(coordinator(h).active, 4) + assert.equal(h.rpc('prompt.submit').length, 5, 'terminal collection does not manufacture a sixth worker') + assert.deepEqual(parked.ref, accepted) + assert.equal(posts(h).filter(e => e.text === 'TERMINAL_WITH_ALL_SLOTS_BUSY').length, terminal === 'complete' ? 1 : 0) + await settle(h, pending) + assert.equal(coordinator(h).active, 0) + assert.ok(h.leases.every(l => l.releases === 1)) + }) +} + +bounded('GC-HOLD: typed hold retains running custody and release never replays or publishes the held turn', async () => { + const h = await harness(members(1)), pending = drive(h) + await flush() + const first = sessionFor(h), accepted = clone(first.ref), lease = h.leases[0] + h.gc.sendToGroupChat('Room', h.roster, '@bot1 stop', 'hold-thread') + await h.advance(); await pending + assert.deepEqual(receipt(h)?.delivery.accepted_turn, accepted) + assert.equal(lease.releases, 0, 'a future-turn hold is not terminal evidence') + assert.equal(coordinator(h).active, 1) + assert.equal(h.rpc('session.interrupt').length, 0) + assert.equal(posts(h).length, 0) + h.gc.sendToGroupChat('Room', h.roster, '@bot1 continue with NEW_INPUT', 'new-thread') + await h.advance() + assert.equal(h.rpc('prompt.submit').length, 1, 'accepted work owns the member until exact retirement') + assert.equal(lease.releases, 0) + first.state = 'complete'; first.pending = null; first.text = 'HELD_OLD_RESULT' + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await flush(); await h.advance() + assert.equal(posts(h).filter(e => e.text === 'HELD_OLD_RESULT').length, 0) + assert.equal(lease.releases, 1) + // An instruction shown as blocked requires a fresh explicit drive. It is + // not silently replayed when an unknown old turn later becomes terminal. + void h.gc.runGroupChatRounds('Room', h.roster, 'new-thread') + await flush(); await h.advance() + assert.equal(h.rpc('prompt.submit').length, 2, 'the new instruction admits once after retirement') + assert.ok(h.rpc('prompt.submit')[1].params.text.includes('NEW_INPUT')) + await settle(h) +}) + +for (const unavailable of [false, true]) { + bounded(`GC-ACK: application ACK retains hot capacity while outcome is ${unavailable ? 'unavailable' : 'running'}`, async () => { + const options = { interruptReply: { status: 'interrupted' }, unavailable } + const h = await harness(members(6), options), pending = drive(h) + await flush() + const accepted = [...h.sessions.values()].map(s => clone(s.ref)) + const result = await h.gc.stopGroupThread('Room', 't1', h.roster) + await h.advance(); await pending + assert.equal(result.status, 'stopping') + assert.equal(result.unconfirmed, 0) + assert.equal(Object.keys(h.room().stranded).length, 4) + assert.equal(coordinator(h).active, 4) + assert.equal(h.activeLeases(), 4) + assert.ok(h.leases.every(l => l.releases === 0)) + const repeats = await h.gc.stopGroupThread('Room', 't1', h.roster) + assert.equal(repeats.status, 'stopping') + assert.equal(h.rpc('session.interrupt').length, 4, 'confirmed application need not repeat the control write') + h.gc.sendToGroupChat('Room', h.roster, '@all resume NEW_INPUT', 't-new') + await h.advance() + assert.equal(h.rpc('prompt.submit').length, 4, 'all four unknown workers still occupy capacity') + options.unavailable = false + for (const session of h.sessions.values()) { session.state = 'interrupted'; session.pending = null } + for (const member of h.roster) await h.gc.harvestStrandedGroupReply('Room', member) + await h.advance() + assert.ok(h.rpc('prompt.submit').length > 4, 'exact terminal releases capacity for the new instruction') + assert.equal(posts(h).length, 0, 'Stop reconciliation never publishes the old result') + for (const ref of accepted) assert.ok(h.rpc('session.turn.poll').some(c => + c.params.accepted_turn?.request_id === ref.request_id && c.params.accepted_turn?.host_boot_id === ref.host_boot_id)) + await settle(h) + assert.equal(coordinator(h).active, 0) + assert.ok(h.leases.every(l => l.releases === 1)) + }) +} + +for (const cold of [false, true]) { + bounded(`GC-ACK: ${cold ? 'cold' : 'hot'} lost ACK retries the captured generation and still awaits terminal`, async () => { + const options = { interruptError: true, interruptReply: { status: 'interrupted' } } + const hot = await harness(members(1), cold ? {} : options), pending = drive(hot) + await flush() + let h = hot + if (cold) { + h = await harness(hot.roster, options) + for (const [id, session] of hot.sessions) h.sessions.set(id, clone(session)) + const rooms = clone(hot.gc.durableGroupChatRooms()) + rooms.Room.running = false; rooms.Room.turns = []; rooms.Room.turn = null + h.gc.$groupChats.set(rooms) + } + const ref = clone(receipt(h).delivery.accepted_turn) + const first = await h.gc.stopGroupThread('Room', 't1', h.roster) + await h.advance() + assert.equal(first.status, 'unconfirmed') + assert.deepEqual(receipt(h)?.delivery.accepted_turn, ref) + options.interruptError = false + const second = await h.gc.stopGroupThread('Room', 't1', h.roster) + assert.equal(second.status, 'stopping') + assert.deepEqual(h.rpc('session.interrupt').map(c => c.params.session_id), [ref.session_id, ref.session_id]) + assert.deepEqual(receipt(h)?.delivery.accepted_turn, ref) + sessionFor(h).state = 'interrupted'; sessionFor(h).pending = null + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await h.advance() + assert.equal(receipt(h), undefined) + assert.equal(posts(h).length, 0) + if (!cold) assert.equal(h.activeLeases(), 0) + await settle(hot, pending) + }) +} + +function renderer() { + const states = []; let cursor = 0 + return { react: { + useState: initial => { + const index = cursor++ + if (!(index in states)) states[index] = typeof initial === 'function' ? initial() : initial + return [states[index], value => { states[index] = typeof value === 'function' ? value(states[index]) : value }] + }, useRef: current => ({ current }), useEffect: () => {}, useMemo: factory => factory() + }, sdk: { useValue: atom => atom.get(), cn: (...values) => values.filter(Boolean).join(' '), + profileColor: () => '#000000', relativeTime: () => '' }, + render: (component, props, fresh = false) => { + cursor = 0; if (fresh) states.length = 0 + return component(props) + } } +} +function nodes(tree) { + const result = [] + const walk = node => { + if (Array.isArray(node)) node.forEach(walk) + else if (node && typeof node === 'object') { result.push(node); walk(node.props?.children) } + } + walk(tree); return result +} +function preparedAnswer(h, entry, approvalChoice = 'once') { + const props = { entry, members: h.roster } + let tree = h.ui.render(h.gc.GroupClarifyCard, props, true) + if (entry.kind === 'approval') nodes(tree).find(n => n.props?.children === approvalChoice && n.props?.onClick).props.onClick() + else nodes(tree).find(n => n.props?.['aria-label'] === `Answer @${entry.member}`).props.onChange({ target: { value: 'ANSWER' } }) + tree = h.ui.render(h.gc.GroupClarifyCard, props) + const button = nodes(tree).find(n => ['Answer', 'Respond'].includes(n.props?.children) && n.props?.onClick) + assert.equal(button.props.disabled, false) + return () => button.props.onClick() +} +async function park(kind = 'approval', options = {}) { + const ui = renderer(), notices = [], originalProjection = options.resumeProjection + Object.assign(options, ui, { onNotify: notice => notices.push(notice), + resumeProjection: (session, method, projection) => { + if (session.approval && session.state === 'waiting') projection.pending_approval = session.approval + return originalProjection?.(session, method, projection) || projection + } }) + const h = await harness(members(1), options), pending = drive(h) + h.ui = ui; h.notices = notices; await flush() + const session = sessionFor(h) + session.state = 'waiting' + if (kind === 'approval') session.approval = { request_id: 'approval-1', command: 'echo offline', choices: ['once', 'deny'] } + else session.pending = { request_id: 'clarify-1', question: 'Choose?' } + await h.advance(); await pending + assert.equal(cards(h).length, 1) + return { h, session, options } +} + +for (const approvalResult of [{ resolved: 0 }, {}, { resolved: '1' }, { resolved: -1 }]) { + bounded(`GC-APPROVAL: ${JSON.stringify(approvalResult)} produces no successful echo`, async () => { + const { h, session } = await park('approval', { approvalResult }) + const accepted = clone(session.ref), log = clone(h.room().log), entry = cards(h)[0] + await preparedAnswer(h, entry)(); await flush() + assert.deepEqual(h.room().log, log) + assert.equal(cards(h)[0], entry, 'a negative or unknown acknowledgement keeps the card') + assert.equal(h.notices.filter(n => n.kind === 'error').length, 1) + assert.deepEqual(receipt(h)?.delivery.accepted_turn, accepted) + assert.equal(h.rpc('prompt.submit').length, 1) + await settle(h) + }) +} +for (const choice of ['once', 'deny']) { + bounded(`GC-APPROVAL: resolved=1 confirms ${choice} and delivers only the exact accepted turn`, async () => { + const { h, session } = await park('approval', { approvalResult: { resolved: 1 } }) + const accepted = clone(session.ref) + session.text = `AFTER_${choice}` + await preparedAnswer(h, cards(h)[0], choice)(); await h.advance() + assert.equal(cards(h).length, 0) + assert.equal(h.room().log.filter(e => e.from.kind === 'user' && e.answer_to).length, 1) + assert.equal(posts(h).filter(e => e.text === `AFTER_${choice}`).length, 1) + assert.deepEqual(session.ref, accepted) + assert.equal(h.rpc('prompt.submit').length, 1) + await settle(h) + }) +} +bounded('GC-APPROVAL: lost response ACK retains card, receipt and honest error', async () => { + const { h } = await park('approval', { answerError: () => true }) + const entry = cards(h)[0], log = clone(h.room().log) + await preparedAnswer(h, entry)(); await flush() + assert.deepEqual(h.room().log, log) + assert.equal(cards(h)[0], entry) + assert.equal(h.notices.filter(n => n.kind === 'error').length, 1) + assert.ok(receipt(h)) + await settle(h) +}) + +for (const legacy of [false, true]) { + bounded(`GC-RENAME: ${legacy ? 'legacy token' : 'room ID'} retains running owner through repeated local renames`, async t => { + const h = await harness(members(1)) + if (!h.gc.renameGroupChat) return t.skip('Exact baseline lacks the rename test port; final source must run this gate') + if (legacy) delete h.room().roomId + const pending = drive(h); await flush() + const owner = coordinator(h), session = sessionFor(h), accepted = clone(session.ref), title = session.title + await h.gc.renameGroupChat('Room', 'Renamed', h.roster) + await h.gc.renameGroupChat('Renamed', 'RenamedAgain', h.roster) + assert.equal(coordinator(h, 'RenamedAgain'), owner) + assert.equal(coordinator(h), undefined) + assert.equal([...owner.occurrences][0].group, 'RenamedAgain') + assert.deepEqual(receipt(h, h.roster[0], 'RenamedAgain')?.delivery.accepted_turn, accepted) + session.state = 'complete'; session.text = 'FINAL_AFTER_RENAME' + await h.advance(); await pending + assert.equal(posts(h, 'RenamedAgain').filter(e => e.text === 'FINAL_AFTER_RENAME').length, 1) + assert.equal(h.gc.$groupChats.get().Room, undefined) + assert.equal(h.gc.$groupChats.get().Renamed, undefined) + assert.equal(h.gc.$groupChats.get().RenamedAgain.running, false) + assert.equal(session.title, title) + assert.equal(h.rpc('session.create').length, 1) + assert.equal(h.rpc('prompt.submit').length, 1) + await settle(h, null, 'RenamedAgain') + }) +} +bounded('GC-RENAME: pending card response and echo follow the same lifetime across rename', async t => { + const answerAckGate = deferred(), { h, session } = await park('clarify', { answerAckGate }) + if (!h.gc.renameGroupChat) return t.skip('Exact baseline lacks the rename test port; final source must run this gate') + const entry = cards(h)[0], owner = coordinator(h), accepted = clone(session.ref) + session.text = 'FINAL_AFTER_RENAMED_ANSWER' + const answering = preparedAnswer(h, entry)(); await flush() + await h.gc.renameGroupChat('Room', 'AnsweredRoom', h.roster) + assert.equal(cards(h)[0], entry) + assert.equal(entry.group, 'AnsweredRoom') + answerAckGate.resolve(); await answering; await h.advance() + const room = h.gc.$groupChats.get().AnsweredRoom + assert.equal(h.gc.$groupChats.get().Room, undefined) + assert.equal(room.log.filter(e => e.answer_to).length, 1) + assert.equal(posts(h, 'AnsweredRoom').filter(e => e.text === 'FINAL_AFTER_RENAMED_ANSWER').length, 1) + assert.equal(coordinator(h, 'AnsweredRoom'), owner) + assert.deepEqual(session.ref, accepted) + assert.equal(h.rpc('prompt.submit').length, 1) + await settle(h, null, 'AnsweredRoom') +}) +bounded('GC-RENAME: in-flight Stop follows rename and retains custody until exact terminal', async t => { + const interruptGate = deferred(), h = await harness(members(1), { interruptGate, interruptReply: { status: 'interrupted' } }) + if (!h.gc.renameGroupChat) return t.skip('Exact baseline lacks the rename test port; final source must run this gate') + const pending = drive(h); await flush() + const accepted = clone(sessionFor(h).ref) + const stopping = h.gc.stopGroupThread('Room', 't1', h.roster); await flush() + await h.gc.renameGroupChat('Room', 'StoppedRoom', h.roster) + interruptGate.resolve(); const result = await stopping; await h.advance(); await pending + assert.equal(result.status, 'stopping') + assert.deepEqual(receipt(h, h.roster[0], 'StoppedRoom')?.delivery.accepted_turn, accepted) + assert.equal(h.activeLeases(), 1) + sessionFor(h).state = 'interrupted' + await h.gc.harvestStrandedGroupReply('StoppedRoom', h.roster[0]); await flush() + assert.equal(receipt(h, h.roster[0], 'StoppedRoom'), undefined) + assert.equal(h.activeLeases(), 0) + assert.equal(h.gc.$groupChats.get().Room, undefined) + assert.equal(posts(h, 'StoppedRoom').length, 0) +}) +bounded('GC-RENAME: remote receive rehomes running owner before room listeners see the new name', async t => { + const options = { rpcResponse: (_route, method) => method === 'profiles.list' + ? { profiles: [{ name: 'default', ui_meta: { 'hermes-bots-groups': options.remote } }] } + : method === 'profiles.configure' ? { applied: { ui_meta: true } } : undefined } + const h = await harness(members(1), options) + if (!h.gc.pullGroupChatServerState || !h.gc.renameGroupChat) return t.skip('Exact baseline lacks the rename test port; final source must run this gate') + const pending = drive(h); await flush() + const owner = coordinator(h), session = sessionFor(h), accepted = clone(session.ref) + const snapshot = h.gc.groupChatSyncSnapshot(), [roomKey, sourceRoom] = Object.entries(snapshot.rooms)[0] + const projected = clone(sourceRoom) + assert.ok(projected) + options.remote = { ...snapshot, rooms: { [roomKey]: { ...projected, name: 'RemoteRoom', revision: projected.revision + 10 } } } + // This receive starts after local writes have settled. An intentionally + // pending local edit has separate preserveRooms authority over the name. + h.gc.stopGroupChatServerSync() + await h.gc.pullGroupChatServerState('local') + assert.equal(coordinator(h, 'RemoteRoom'), owner) + assert.equal(h.gc.$groupChats.get().Room, undefined) + assert.deepEqual(receipt(h, h.roster[0], 'RemoteRoom')?.delivery.accepted_turn, accepted) + session.state = 'complete'; session.text = 'REMOTE_RENAMED_FINAL' + await h.advance(); await pending + assert.equal(posts(h, 'RemoteRoom').filter(e => e.text === 'REMOTE_RENAMED_FINAL').length, 1) + assert.equal(h.gc.$groupChats.get().RemoteRoom.running, false) + assert.equal(h.rpc('prompt.submit').length, 1) + await settle(h, null, 'RemoteRoom') +}) +bounded('GC-RENAME: same display name with a new room lifetime cannot adopt the old final', async () => { + const h = await harness(members(1)), pending = drive(h) + await flush() + const session = sessionFor(h) + h.gc.$groupChats.set({ Room: { ...clone(h.room()), roomId: 'replacement-room', epoch: 1, running: false, + sessions: {}, sessionOwners: {}, stranded: {}, log: [], turns: [], turn: null } }) + session.state = 'complete'; session.text = 'OLD_LIFETIME_FINAL' + await h.advance(); await pending + assert.equal(posts(h).length, 0) + assert.equal(h.room().roomId, 'replacement-room') + assert.equal(h.rpc('prompt.submit').length, 1) +}) +bounded('GC-RENAME: late Stop ACK cannot poll or update a replacement room at the same name', async () => { + const interruptGate = deferred(), h = await harness(members(1), { interruptGate, interruptReply: { status: 'interrupted' } }) + const pending = drive(h); await flush() + const stopping = h.gc.stopGroupThread('Room', 't1', h.roster); await flush() + const before = h.rpc('session.turn.poll').length + const replacement = { ...clone(h.room()), roomId: 'new-lifetime', running: true, epoch: 20, + stranded: {}, sessions: {}, sessionOwners: {}, holds: {}, log: [], turns: [], turn: null } + h.gc.$groupChats.set({ Room: replacement }) + interruptGate.resolve(); const result = await stopping; await h.advance(); await pending + assert.equal(result.status, 'unconfirmed') + assert.equal(h.rpc('session.turn.poll').length, before) + assert.equal(h.room().roomId, 'new-lifetime') + assert.equal(h.room().epoch, 20) + assert.equal(h.room().running, true) + assert.equal(h.room().log.length, 0) +}) +bounded('GC-HOLD: waiting hold retains the exact question custody while suppressing stale Answer', async () => { + const { h, session } = await park('clarify') + const accepted = clone(session.ref), staleAnswer = preparedAnswer(h, cards(h)[0]) + assert.equal(coordinator(h).active, 0) + h.gc.sendToGroupChat('Room', h.roster, '@bot1 pause', 'hold-thread') + await h.advance(); await staleAnswer(); await flush() + assert.equal(h.rpc('clarify.respond').length, 0) + assert.deepEqual(receipt(h)?.delivery.accepted_turn, accepted) + assert.equal(h.activeLeases(), 1) + assert.equal(coordinator(h).active, 0, 'verified waiting keeps no execution worker') + assert.equal(h.rpc('session.interrupt').length, 0) + session.state = 'complete'; session.pending = null; session.text = 'HELD_QUESTION_FINAL' + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await h.advance() + assert.equal(receipt(h), undefined) + assert.equal(h.activeLeases(), 0) + assert.equal(posts(h).length, 0) +}) +bounded('GC-ACK: unavailable timeout and repeated Stop retain four occupied workers', async () => { + const options = { unavailable: true, interruptReply: { status: 'interrupted' } } + const h = await harness(members(6), options), pending = drive(h) + await flush(); await h.advance(20 * 60000) + assert.equal(h.rpc('prompt.submit').length, 4) + assert.equal(coordinator(h).active, 4) + assert.equal(h.activeLeases(), 4) + const first = await h.gc.stopGroupThread('Room', 't1', h.roster); await h.advance(); await pending + assert.equal(first.status, 'stopping') + const second = await h.gc.stopGroupThread('Room', 't1', h.roster) + assert.equal(second.status, 'stopping') + assert.equal(Object.keys(h.room().stranded).length, 4) + assert.equal(coordinator(h).active, 4) + assert.equal(h.rpc('prompt.submit').length, 4) + options.unavailable = false + for (const session of h.sessions.values()) session.state = 'interrupted' + for (const member of h.roster) await h.gc.harvestStrandedGroupReply('Room', member) + assert.equal(h.activeLeases(), 0) + assert.equal(coordinator(h).active, 0) + assert.equal(posts(h).length, 0) +}) + +for (const cold of [false, true]) { + for (const failed of [false, true]) { + bounded(`GC-ACK-02: ${cold ? 'cold' : 'parked hot'} terminal retains custody through ${failed ? 'rejected' : 'applied'} pending interrupt`, async () => { + const interruptGate = deferred(), options = { unavailable: true, interruptGate, interruptError: failed } + const hot = await harness(members(1), cold ? { unavailable: true } : options), initial = drive(hot) + await flush(); await hot.advance(); await initial + const oldOwner = [...coordinator(hot).occurrences][0] + assert.equal(oldOwner.collectorDone, true, 'the collector has parked on unavailable outcome') + let h = hot + if (cold) { + h = await harness(hot.roster, options) + h.gc.$groupChats.set(clone(hot.gc.durableGroupChatRooms())) + for (const [id, session] of hot.sessions) h.sessions.set(id, clone(session)) + } + const marker = receipt(h), accepted = clone(marker.delivery.accepted_turn) + const stopping = h.gc.stopGroupThread('Room', 't1', h.roster) + const repeated = h.gc.stopGroupThread('Room', 't1', h.roster) + await flush() + options.unavailable = false + const session = sessionFor(h) + session.state = 'complete'; session.pending = null; session.text = 'OLD_STOPPED_FINAL' + h.gc.sendToGroupChat('Room', h.roster, '@all resume REPLACEMENT_INPUT', 'replacement-thread') + await h.advance(250) + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await flush() + assert.equal(receipt(h), marker, 'exact terminal cannot erase a pending session control target') + assert.equal(h.rpc('prompt.submit').length, cold ? 0 : 1, 'no successor may reuse the runtime before control settles') + assert.equal(h.rpc('session.interrupt').length, 1, 'concurrent Stop shares the exact pending request') + assert.deepEqual(session.ref, accepted) + if (!cold) { + assert.equal(oldOwner.released, undefined) + assert.equal(coordinator(h).members.get(oldOwner.memberLock), oldOwner) + assert.equal(h.gc.groupRuntimeSessionOwners.get(oldOwner.sessionLock), oldOwner) + assert.equal(h.leases[0].releases, 0) + assert.equal(coordinator(h).active, 1) + assert.equal(oldOwner.reservation, true, 'pending retirement keeps its publication reservation') + } + interruptGate.resolve() + await stopping; await repeated; await flush() + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await h.advance() + // A cold blocked instruction can require a fresh explicit drive. That is + // intentional; it cannot replay the uncertain predecessor automatically. + if (h.rpc('prompt.submit').length === (cold ? 0 : 1)) { + void h.gc.runGroupChatRounds('Room', h.roster, 'replacement-thread') + await flush(); await h.advance() + } + assert.equal(h.rpc('prompt.submit').length, cold ? 1 : 2) + assert.ok(h.rpc('prompt.submit').at(-1).params.text.includes('REPLACEMENT_INPUT')) + assert.equal(h.rpc('session.interrupt').length, 1, 'no delayed cleanup targets the admitted replacement') + assert.equal(posts(h).filter(e => e.text === 'OLD_STOPPED_FINAL').length, 0) + if (!cold) assert.equal(h.leases[0].releases, 1) + await settle(h) + await settle(hot) + }) + } +} + +async function delayedRenameHarness(extra = {}) { + const metadataGate = deferred(), ui = renderer(), options = { ...ui, ...extra, + rpcResponse: async (route, method, params) => { + const supplied = await extra.rpcResponse?.(route, method, params) + if (supplied !== undefined) return supplied + if (method === 'profiles.configure') { + if (params.name === 'bot1' && params.ui_meta?.['hermes-bots']?.groups?.includes('Renamed')) { + await metadataGate.promise + } + return { applied: { ui_meta: true } } + } + return undefined + } } + const h = await harness(members(2), options) + h.ui = ui + // Seed memberships through the actual creation handler, then use the same + // room/metadata owners as the real rename and settings actions. + h.gc.$groupChats.set({}) + assert.equal(h.gc.createFreshGroupChat('Room', h.roster), 'Room') + await flush() + h.gc.updateGroupChat('Room', room => { room.log = [clone(h.input)]; return room }, { sync: false }) + h.gc.$groupChatWorkspace.set('Room') + return { h, metadataGate, options } +} + +bounded('GC-RENAME-02: overlapping local renames cannot resurrect a name or stale later-member membership', async () => { + const { h, metadataGate } = await delayedRenameHarness() + const roomId = h.room().roomId + const first = h.gc.renameGroupChat('Room', 'Renamed', h.roster); await flush() + const second = await h.gc.renameGroupChat('Renamed', 'RenamedAgain', h.roster) + metadataGate.resolve() + assert.equal(await first, 'RenamedAgain', 'the late return names the same current lifetime') + assert.equal(second, 'RenamedAgain') + assert.equal(h.gc.$groupChats.get().Room, undefined) + assert.equal(h.gc.$groupChats.get().Renamed, undefined) + assert.equal(h.gc.$groupChats.get().RenamedAgain.roomId, roomId) + assert.equal(h.gc.$groupChatWorkspace.get(), 'RenamedAgain') + for (const member of h.roster) { + const last = h.rpc('profiles.configure').filter(c => c.params.name === member.name).at(-1) + assert.deepEqual(last.params.ui_meta['hermes-bots'].groups, ['RenamedAgain']) + } + assert.equal(h.rpc('session.create').length, 0) + assert.equal(h.rpc('prompt.submit').length, 0) +}) + +function prepareSettings(h, renamed = []) { + let closed = 0 + const props = { group: 'Room', members: h.roster, open: true, onClose: () => { closed++ }, onRenamed: name => renamed.push(name) } + let tree = h.ui.render(h.gc.GroupChatSettingsDialog, props, true) + nodes(tree).find(n => n.props?.['aria-label'] === 'Group name').props.onChange({ target: { value: 'Renamed' } }) + nodes(tree).find(n => n.props?.onImage).props.onImage('OFFLINE_ROOM_IMAGE') + tree = h.ui.render(h.gc.GroupChatSettingsDialog, props) + const click = nodes(tree).find(n => n.props?.children === 'Save' && n.props?.onClick).props.onClick + return { click, renamed, closed: () => closed } +} +function submitSettings(h, renamed = []) { + const settings = prepareSettings(h, renamed) + settings.click() + return settings +} + +bounded('GC-RENAME-02: the actual settings image and callback follow an overlapping rename', async () => { + const { h, metadataGate } = await delayedRenameHarness(), settings = submitSettings(h) + await flush() + await h.gc.renameGroupChat('Renamed', 'RenamedAgain', h.roster) + metadataGate.resolve(); await flush() + assert.equal(h.gc.$groupChats.get().Renamed, undefined) + assert.equal(h.gc.$groupChats.get().RenamedAgain.image, 'OFFLINE_ROOM_IMAGE') + assert.deepEqual(settings.renamed, ['RenamedAgain']) + assert.equal(settings.closed(), 1) +}) + +bounded('GC-RENAME-02: remote rename during delayed member persistence keeps the current lifetime only', async () => { + let remote + const { h, metadataGate } = await delayedRenameHarness({ rpcResponse: (_route, method) => method === 'profiles.list' + ? { profiles: [{ name: 'default', ui_meta: { 'hermes-bots-groups': remote } }] } : undefined }) + const first = h.gc.renameGroupChat('Room', 'Renamed', h.roster); await flush() + const snapshot = h.gc.groupChatSyncSnapshot(), [key, room] = Object.entries(snapshot.rooms)[0] + remote = { ...snapshot, rooms: { [key]: { ...clone(room), name: 'RemoteRenamed', revision: room.revision + 10 } } } + // Clear only the fake mirror's pending edit, matching a settled/received + // remote revision. The real receive handler retains its existing CAS rule. + h.gc.stopGroupChatServerSync() + await h.gc.pullGroupChatServerState('local') + metadataGate.resolve() + assert.equal(await first, 'RemoteRenamed') + assert.equal(h.gc.$groupChats.get().Room, undefined) + assert.equal(h.gc.$groupChats.get().Renamed, undefined) + assert.equal(h.gc.$groupChats.get().RemoteRenamed.roomId, room.roomId) + assert.equal(h.rpc('prompt.submit').length, 0) +}) + +for (const replaced of [false, true]) { + bounded(`GC-RENAME-02: delayed actual settings cannot update a ${replaced ? 'replacement' : 'deleted'} room`, async () => { + const { h, metadataGate } = await delayedRenameHarness(), settings = submitSettings(h) + await flush() + const replacement = { ...clone(h.gc.$groupChats.get().Renamed), roomId: 'replacement-lifetime', + log: [], sessions: {}, sessionOwners: {}, stranded: {}, image: 'REPLACEMENT_IMAGE' } + h.gc.$groupChats.set(replaced ? { Renamed: replacement } : {}) + metadataGate.resolve(); await flush() + assert.equal(h.gc.$groupChats.get().Room, undefined) + assert.deepEqual(settings.renamed, []) + assert.equal(settings.closed(), 0) + if (replaced) { + assert.equal(h.gc.$groupChats.get().Renamed, replacement) + assert.equal(replacement.image, 'REPLACEMENT_IMAGE') + } else assert.equal(h.gc.$groupChats.get().Renamed, undefined) + }) +} + +for (const replaced of [false, true]) { + bounded(`GC-RENAME-02: an already rendered Save cannot adopt a ${replaced ? 'replacement' : 'deleted'} lifetime`, async () => { + const { h, metadataGate } = await delayedRenameHarness(), settings = prepareSettings(h) + const replacement = { ...clone(h.room()), roomId: 'replacement-before-click', image: 'REPLACEMENT_IMAGE', log: [] } + const before = h.rpc('profiles.configure').length + h.gc.$groupChats.set(replaced ? { Room: replacement } : {}) + settings.click(); metadataGate.resolve(); await flush() + assert.deepEqual(h.gc.$groupChats.get(), replaced ? { Room: replacement } : {}) + assert.equal(h.rpc('profiles.configure').length, before, 'a stale button cannot rename members of a vanished room') + assert.deepEqual(settings.renamed, []) + assert.equal(settings.closed(), 0) + }) +} + +bounded('GC-RENAME-02: actual settings still create and rename a legitimate metadata-only group', async () => { + const { h, metadataGate } = await delayedRenameHarness() + h.gc.$groupChats.set({}) + const settings = submitSettings(h); await flush() + metadataGate.resolve(); await flush() + const room = h.gc.$groupChats.get().Renamed + assert.ok(room) + assert.ok(room.coordinationId, 'the existing legacy lifetime token fences the save') + assert.equal(room.image, 'OFFLINE_ROOM_IMAGE') + assert.equal(h.gc.$groupChats.get().Room, undefined) + assert.deepEqual(settings.renamed, ['Renamed']) + assert.equal(settings.closed(), 1) + assert.equal(h.rpc('prompt.submit').length, 0) +}) diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-parallel.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-parallel.test.mjs index 35bfe3a108c2..e2871f592f5d 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/group-parallel.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-parallel.test.mjs @@ -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)) }) } diff --git a/apps/desktop/src/plugins/hermes-bots/tests/pr118-postmerge.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/pr118-postmerge.test.mjs index f7ec4c18992d..4950349ab049 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/pr118-postmerge.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/pr118-postmerge.test.mjs @@ -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) } @@ -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: [], @@ -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` diff --git a/apps/desktop/src/plugins/hermes-bots/tests/stop-custody-harness.mjs b/apps/desktop/src/plugins/hermes-bots/tests/stop-custody-harness.mjs index caf570c4a3c1..ada7cbac2a58 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/stop-custody-harness.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/stop-custody-harness.mjs @@ -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 @@ -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) @@ -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)) }) diff --git a/apps/desktop/src/plugins/hermes-bots/tests/stop-custody.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/stop-custody.test.mjs index 309f7634c050..1d5f98373443 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/stop-custody.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/stop-custody.test.mjs @@ -42,7 +42,7 @@ for (const mode of ['running', 'waiting', 'expired']) { if (mode === 'expired') await h.advance(21 * 60 * 1000) const count = mode === 'waiting' ? 6 : 4, workers = mode === 'waiting' ? 2 : 4 const result = await settleStop(h, pending) - assert.deepEqual(result, { status: 'unconfirmed', unconfirmed: count }) + assert.deepEqual(result, { status: 'unconfirmed', unconfirmed: count, pending: count }) assertCustody(h, count, workers) assert.equal(h.rpc('session.interrupt').length, count) assert.equal(Object.keys(h.gc.$groupClarify.get()).length, 0) @@ -56,7 +56,10 @@ for (const mode of ['running', 'waiting', 'expired']) { assert.deepEqual(h.rpc('session.interrupt').slice(count).map(c => [c.route, c.params.session_id]), firstTargets) assertCustody(h, count, workers) options.interruptError = false - assert.deepEqual(await h.gc.stopGroupThread('Room', 't1', h.roster), { status: 'stopped', unconfirmed: 0 }) + const stopped = await h.gc.stopGroupThread('Room', 't1', h.roster) + await flush() + assert.equal(stopped.unconfirmed, 0) + assert.ok(['stopping', 'stopped'].includes(stopped.status)) assertReleased(h) await h.gc.stopGroupThread('Room', 't1', h.roster) assert.equal(h.rpc('session.interrupt').length, count * 3) @@ -73,16 +76,18 @@ for (const reply of [{}, { interrupted: false }, { interrupted: true }, { status assert.equal([...h.sessions.values()][0].state, 'running') options.interruptReply = undefined await h.gc.stopGroupThread('Room', 't1', h.roster) + await flush() assertReleased(h) }) } -test('lost ACK after backend applied interrupt retains custody until exact interrupted outcome', async () => { +test('lost ACK after backend applied interrupt can retire only from exact interrupted outcome', async () => { const h = await harness(members(1), { interruptBehavior: session => { session.state = 'interrupted'; throw new Error('ACK was lost after apply') } }), pending = drive(h); await flush() - assert.equal((await settleStop(h, pending)).status, 'unconfirmed') - assertCustody(h, 1, 1) + await settleStop(h, pending) + assert.equal(h.rpc('session.turn.poll').at(-1).params.accepted_turn.request_id, + [...h.sessions.values()][0].ref.request_id, 'retirement used the accepted request despite a lost ACK') await h.gc.harvestStrandedGroupReply('Room', h.roster[0]) assertReleased(h) assert.equal(h.posts().length, 0) @@ -129,6 +134,7 @@ test('exact waiting reconciliation releases only worker, preserves stopped recei assert.equal(Object.keys(h.gc.$groupClarify.get()).length, 0) options.interruptError = false await h.gc.stopGroupThread('Room', 't1', h.roster) + await flush() assertReleased(h) }) @@ -202,7 +208,9 @@ test('late submit rejection after earlier interrupt ACK requires new-generation assert.equal(h.leases[0].releases, 0) options.interruptBehavior = undefined await h.gc.stopGroupThread('Room', 't1', h.roster) - assertReleased(h) + assert.equal(receipts(h).length, 1, 'a later ACK cannot reconstruct the lost accepted identity') + assert.equal(coordinator(h).active, 1) + assert.equal(h.leases[0].releases, 0) }) test('Stop during late session acquisition cannot submit; lease closes once after producer settles', async () => { @@ -267,7 +275,12 @@ test('reload Stop failures retain exact receipt, coalesce concurrent retry and l const a = cold.gc.stopGroupThread('Room', 't1', cold.roster), b = cold.gc.stopGroupThread('Room', 't1', cold.roster) await flush() assert.equal(cold.rpc('session.interrupt').length, 2, 'two concurrent retries issue one exact RPC') - gate.resolve(); assert.equal((await a).status, 'stopped'); assert.equal((await b).status, 'stopped') + gate.resolve() + for (const result of [await a, await b]) { + assert.equal(result.unconfirmed, 0) + assert.ok(['stopping', 'stopped'].includes(result.status)) + } + await flush() assert.equal(receipts(cold).length, 0) assert.equal(cold.rpc('prompt.submit').length, 0) }) From b4b752d2e62222163a232595f021e2a5f1e3e255 Mon Sep 17 00:00:00 2001 From: Josh Stevenson Date: Sun, 4 Oct 2026 02:33:20 -0700 Subject: [PATCH 2/2] fix(desktop): reconcile group cancellation preparation and admission refusal --- .../desktop/src/plugins/hermes-bots/plugin.js | 32 ++++++-- .../hermes-bots/tests/group-chat.test.mjs | 2 +- .../tests/group-protocol-custody.test.mjs | 20 +++-- .../tests/group-stop-thread.test.mjs | 12 ++- .../tests/group-turn-lease.test.mjs | 12 ++- .../tests/group-turn-outcomes.test.mjs | 28 +++++-- .../tests/pr127-ci-boundaries.test.mjs | 81 +++++++++++++++++++ 7 files changed, 160 insertions(+), 27 deletions(-) create mode 100644 apps/desktop/src/plugins/hermes-bots/tests/pr127-ci-boundaries.test.mjs diff --git a/apps/desktop/src/plugins/hermes-bots/plugin.js b/apps/desktop/src/plugins/hermes-bots/plugin.js index cf6fa1635fed..a217c6cb4ebb 100644 --- a/apps/desktop/src/plugins/hermes-bots/plugin.js +++ b/apps/desktop/src/plugins/hermes-bots/plugin.js @@ -7693,8 +7693,21 @@ async function retainGroupTurnRoute(member) { * exactly once more. Returns the runtime id the submit actually landed on so * the poll loop keeps a live fallback target. */ async function submitGroupTurnPrompt(member, runtime, stored, text, occurrence, canSubmit) { + const submit = async target => { + if (occurrence) occurrence.admissionRefused = false + try { + return await requestForBot(member, 'prompt.submit', { session_id: target, text }) + } catch (error) { + // The gateway rejects a missing runtime before prompt admission. A + // cancelled remint/probe after that refusal owns no accepted turn. + if (occurrence && (error?.code === 4001 || error?.code === 4090 || error?.code === 'POOL_CAPACITY_EXCEEDED')) { + occurrence.admissionRefused = true + } + throw error + } + } try { - const ack = await requestForBot(member, 'prompt.submit', { session_id: runtime, text }) + const ack = await submit(runtime) return { runtime, acceptedTurn: groupAcceptedTurn(ack?.accepted_turn, runtime) } } catch (error) { @@ -7730,7 +7743,7 @@ async function submitGroupTurnPrompt(member, runtime, stored, text, occurrence, if (occurrence?.cancelled) await interruptGroupOccurrence(occurrence) throw error } - const ack = await requestForBot(member, 'prompt.submit', { session_id: fresh, text }) + const ack = await submit(fresh) return { runtime: fresh, acceptedTurn: groupAcceptedTurn(ack?.accepted_turn, fresh) } } @@ -8343,7 +8356,7 @@ function interruptStoppedGroupMarker(group, memberKey, marker, target) { } function groupOccurrenceStopConfirmed(occurrence) { - if (occurrence.submissionPending || occurrence.answerPromise || groupOccurrenceHasPendingInterrupt(occurrence)) return false + if (occurrence.preparationPending || occurrence.submissionPending || occurrence.answerPromise || groupOccurrenceHasPendingInterrupt(occurrence)) return false if (occurrence.terminalObserved) return true if (!occurrence.submitAttempted) return Boolean(occurrence.collectorDone) return false // A matching interrupt ACK still needs exact terminal evidence. @@ -8354,7 +8367,7 @@ function groupOccurrenceHasPendingInterrupt(occurrence) { } function groupOccurrenceInterruptApplied(occurrence) { - return !occurrence.submissionPending && !occurrence.answerPromise && + return !occurrence.preparationPending && !occurrence.submissionPending && !occurrence.answerPromise && occurrence.interrupts.get(`${occurrence.runtime}::${occurrence.admissionVersion || 0}`)?.confirmed === true } @@ -8509,6 +8522,7 @@ function retainUnresolvedGroupOccurrence(occurrence) { } async function executeGroupOccurrence(occurrence) { + occurrence.preparationPending = true try { const release = await retainGroupTurnRoute(occurrence.captured.requestMember) let released = false @@ -8520,6 +8534,7 @@ async function executeGroupOccurrence(occurrence) { return await runGroupChatMemberTurnLeased(occurrence.group, occurrence.captured, occurrence.prompt, occurrence.thread, occurrence.images, occurrence.deliveryResult, occurrence) } finally { + occurrence.preparationPending = false // Stop owns any interrupt through its acknowledgement, including a runtime // that became available during acquisition. Successors cannot start yet. if (occurrence.answerPromise) await occurrence.answerPromise.catch(() => undefined) @@ -8647,7 +8662,10 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima consumeGroupTurnMarker(group, memberKey, marker) return discarded() } - if (occurrence) { occurrence.phase = 'running'; paintGroupOccurrences(occurrence.coordinator) } + if (occurrence) { + occurrence.preparationPending = false + occurrence.phase = 'running'; paintGroupOccurrences(occurrence.coordinator) + } recordGroupActivity(group, { kind: 'working', member: member.name, thread }) const fileRefs = [] for (const img of Array.isArray(images) ? images : []) { @@ -8812,8 +8830,8 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima return null } catch (error) { group = occurrence?.group || group - if (!submitAttempted || error?.code === 4090 || error?.code === 'POOL_CAPACITY_EXCEEDED') { - if (error?.code === 4090 || error?.code === 'POOL_CAPACITY_EXCEEDED') { + if (!submitAttempted || occurrence?.admissionRefused || error?.code === 4090 || error?.code === 'POOL_CAPACITY_EXCEEDED') { + if (occurrence?.admissionRefused || error?.code === 4090 || error?.code === 'POOL_CAPACITY_EXCEEDED') { if (occurrence) occurrence.submitAttempted = false error.data = { ...error.data, outcomeState: 'admission-refused', reason: error.message } } diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-chat.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-chat.test.mjs index d646bc6b7a18..50b2c7388712 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/group-chat.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-chat.test.mjs @@ -193,7 +193,7 @@ function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approva } if (method === 'approval.respond') { approvalResponds.push({ ...params }) - return { resolved: true } + return { resolved: 1 } } return {} }, diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-protocol-custody.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-protocol-custody.test.mjs index 9c27b8f5cc49..c90b0cd72271 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/group-protocol-custody.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-protocol-custody.test.mjs @@ -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 () => { @@ -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) }) diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-stop-thread.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-stop-thread.test.mjs index cfc16a188af5..8aade666d9b6 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/group-stop-thread.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-stop-thread.test.mjs @@ -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 () => { diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-turn-lease.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-turn-lease.test.mjs index 6dccbe76a205..60e754d2a261 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/group-turn-lease.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-turn-lease.test.mjs @@ -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 () => { diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-turn-outcomes.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-turn-outcomes.test.mjs index 66691229b61f..05cde7594a7b 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/group-turn-outcomes.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-turn-outcomes.test.mjs @@ -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) }) @@ -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) }) diff --git a/apps/desktop/src/plugins/hermes-bots/tests/pr127-ci-boundaries.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/pr127-ci-boundaries.test.mjs new file mode 100644 index 000000000000..9457018451ad --- /dev/null +++ b/apps/desktop/src/plugins/hermes-bots/tests/pr127-ci-boundaries.test.mjs @@ -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) +})