diff --git a/apps/desktop/src/plugins/hermes-bots/plugin.js b/apps/desktop/src/plugins/hermes-bots/plugin.js index 53d7edb03c942..519e205382020 100644 --- a/apps/desktop/src/plugins/hermes-bots/plugin.js +++ b/apps/desktop/src/plugins/hermes-bots/plugin.js @@ -67,6 +67,15 @@ import { import { useEffect, useMemo, useRef, useState } from 'react' import { jsx, jsxs } from 'react/jsx-runtime' +// Explicit renderer ports keep the shipped module importable with offline dependencies. +function createGroupTurnPorts(ports = globalThis) { + return { Date: ports.Date ?? globalThis.Date, setTimeout: ports.setTimeout ?? globalThis.setTimeout, + clearTimeout: ports.clearTimeout ?? globalThis.clearTimeout, + setInterval: ports.setInterval ?? globalThis.setInterval, clearInterval: ports.clearInterval ?? globalThis.clearInterval, + document: ports.document ?? globalThis.document } +} +let { Date, setTimeout, clearTimeout, setInterval, clearInterval, document } = createGroupTurnPorts() + const { McpTab, ToolsetConfigPanel } = sdk // Keep optional exports feature-detected; test harnesses may strip the SDK namespace. const SkillsView = typeof sdk === 'undefined' ? undefined : sdk.SkillsView @@ -679,7 +688,8 @@ function groupChatSyncSnapshot(all = $groupChats.get(), deleted = {}) { text: String(entry?.text || '').slice(0, GROUP_CHAT_SYNC_TEXT_CHARS), at: Number(entry?.at || 0), ...(entry?.thread ? { thread: String(entry.thread).slice(0, 128) } : {}), - ...(groupTurnDeliveryKey(entry?.delivery) ? { delivery: entry.delivery } : {}) + ...(groupTurnDeliveryKey(entry?.delivery) ? { delivery: entry.delivery } : {}), + ...(groupTurnDeliveryKey(entry?.answer_to) ? { answer_to: entry.answer_to } : {}) })) const compact = { name: String(name).slice(0, 64), @@ -975,7 +985,7 @@ function mergeRemoteGroupChatSnapshotIntoRooms( // watermark deltas for members that already saw it (phantom rounds). if (!entries.has(entryKey)) { entries.set(entryKey, entry) - if (entry.from?.kind === 'user') { + if (groupIsUserInstruction(entry)) { const thread = groupThreadOf(entry) threadInputVersions[thread] = (threadInputVersions[thread] || 0) + 1 } @@ -1079,6 +1089,8 @@ function durableGroupChatRooms(all = $groupChats.get()) { recoveryOrigin: room.recoveryOrigin || null, sessions: room.sessions || {}, stranded: room.stranded || {}, + driveContexts: durableGroupDriveContexts(room), + coordinationId: room.coordinationId || null, members: Array.isArray(room.members) ? room.members : [], // Immutable room identity: without this, a room merged in via the // remote-sync path (the only caller of this function) loses its @@ -7128,6 +7140,8 @@ function updateGroupChat(group, mutate, { sync = true } = {}) { // with the pre-turn message baseline. Survives reloads so finished // work is still harvested after a window restart. stranded: room.stranded || {}, + driveContexts: durableGroupDriveContexts(room), + coordinationId: room.coordinationId || null, // #93129: sticky per-member stop holds. Watermarks persist, so holds // must too — otherwise a window restart silently releases a bot the // user explicitly stopped. @@ -7224,6 +7238,8 @@ async function disbandGroupChat(group, members) { consumedInputs: groupConsumedInputs(room), threadInputVersions: room.threadInputVersions || {}, stranded: room.stranded || {}, + driveContexts: durableGroupDriveContexts(room), + coordinationId: room.coordinationId || null, holds: room.holds || {}, sessions: room.sessions || {}, sessionOwners: room.sessionOwners || {}, @@ -7393,14 +7409,15 @@ function normalizeGroupChatText(text) { return trimmed === GROUP_EMPTY_SENTINEL ? GROUP_EMPTY_FRIENDLY : trimmed } -function appendGroupChatEntry(group, from, text, thread, images, delivery) { +function appendGroupChatEntry(group, from, text, thread, images, delivery, answerTo) { const entry = { id: groupChatEntryId(), at: Date.now(), from, text: normalizeGroupChatText(text), thread: thread || 'legacy', - ...(groupTurnDeliveryKey(delivery) ? { delivery } : {}) + ...(groupTurnDeliveryKey(delivery) ? { delivery } : {}), + ...(groupTurnDeliveryKey(answerTo) ? { answer_to: answerTo } : {}) } if (Array.isArray(images) && images.length) { @@ -7426,13 +7443,17 @@ function appendGroupChatEntry(group, from, text, thread, images, delivery) { updateGroupChat(group, room => { room.log.push(entry) - if (from.kind === 'user') { + if (groupIsUserInstruction(entry)) { room.threadInputVersions = { ...(room.threadInputVersions || {}), [entry.thread]: (room.threadInputVersions?.[entry.thread] || 0) + 1 } } return room }) + if (groupIsUserInstruction(entry)) { + for (const pump of [...(groupRoomCoordinators.get(group)?.budgetPumps || [])]) pump() + } + // Needs-you: a member addressing @user badges the group header. if (from.kind === 'member' && /@user\b/i.test(entry.text)) { $groupNeedsYou.set({ ...$groupNeedsYou.get(), [group]: true }) @@ -7850,8 +7871,10 @@ function groupTurnMarkerBlocksDispatch(room, memberKey, coordinator) { !collectingGroupTurnMarkers.has(marker) && !stoppingOwner } -function consumeGroupTurnMarker(group, memberKey, marker) { +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) + const coldDrive = owned ? null : restoreGroupDrive(group, marker) let consumed = false updateGroupChat(group, room => { if (room.stranded?.[memberKey] === marker) { @@ -7863,9 +7886,16 @@ function consumeGroupTurnMarker(group, memberKey, marker) { return room }, { sync: false }) if (consumed) { + if (coldDrive) { + coldDrive.budget.reserved = Math.max(0, coldDrive.budget.reserved - 1) + if (published) coldDrive.budget.posted++ + persistGroupDrive(group, coldDrive) + } const parked = [...(groupRoomCoordinators.get(group)?.occurrences || [])].find(o => o.id === marker?.occurrence_id && o.collectorDone) if (parked && !parked.answerPromise) finishGroupOccurrence(parked) + const coordinator = groupRoomCoordinators.get(group) + if (coordinator) pumpGroupOccurrences(coordinator) } return consumed } @@ -7919,7 +7949,9 @@ function syncGroupClarify(group, member, state, requestMember = member) { // Same request already mirrored — keep the object identity so the card // doesn't lose its draft to a re-render. - if (current?.requestId === requestId) { + if (current?.requestId === requestId && current.sessionId === (state?.session_id || null) && + groupTurnDeliveryKey(current.receipt?.delivery) === + groupTurnDeliveryKey($groupChats.get()[group]?.stranded?.[groupMemberKey(member)]?.delivery)) { return true } @@ -7929,6 +7961,9 @@ function syncGroupClarify(group, member, state, requestMember = member) { member: member.name, memberKey: groupMemberKey(member), requestMember, + receipt: $groupChats.get()[group]?.stranded?.[groupMemberKey(member)], + roomId: $groupChats.get()[group]?.roomId || null, + roomToken: $groupChats.get()[group]?.coordinationId || null, thread: $groupChats.get()[group]?.stranded?.[groupMemberKey(member)]?.thread || 'legacy', // approval.respond keys on the session, not just the request — carry the // runtime id the snapshot came from. @@ -7997,7 +8032,23 @@ function clearGroupClarify(group) { * - approval: `approval.respond` with the choice (once/session/always/deny), * keyed by session + request_id — the same wire the 1:1 approval card * and native notifications use. */ +function requireCurrentGroupQuestion(entry) { + const room = $groupChats.get()[entry.group] + const card = $groupClarify.get()[`${entry.group}::${entry.memberKey}`] + if (!room || room.tombstone || room.holds?.[entry.memberKey] || + (entry.roomId !== undefined && entry.roomId !== (room.roomId || null)) || + (entry.roomToken !== undefined && entry.roomToken !== (room.coordinationId || null)) || + (entry.receipt && ( + groupTurnDeliveryKey(room.stranded?.[entry.memberKey]?.delivery) !== groupTurnDeliveryKey(entry.receipt.delivery) || + room.stranded?.[entry.memberKey]?.stop_requested || + !groupTurnMarkerIntentIsCurrent(room, room.stranded?.[entry.memberKey]))) || + card?.requestId !== entry.requestId || card?.sessionId !== entry.sessionId) { + throw new Error('Member question was stopped or replaced') + } +} + function answerGroupClarify(entry, member, answers) { + try { requireCurrentGroupQuestion(entry) } catch (error) { return Promise.reject(error) } member = entry.requestMember || member const occurrence = [...(groupRoomCoordinators.get(entry.group)?.occurrences || [])].find(o => o.memberKey === entry.memberKey && o.runtime === entry.sessionId) @@ -8044,10 +8095,7 @@ async function respondGroupClarify(entry, member, answers, occurrence, fence) { try { if (occurrence) await reserveGroupResumeWorker(occurrence) if (occurrence && !groupOccurrenceCurrent(occurrence)) throw new Error('Member occurrence was stopped or replaced') - const currentCard = $groupClarify.get()[`${entry.group}::${entry.memberKey}`] - if (occurrence && (currentCard?.requestId !== entry.requestId || currentCard?.sessionId !== entry.sessionId)) { - throw new Error('Member question was replaced before its response could start') - } + requireCurrentGroupQuestion(entry) attempted = true if (entry.kind === 'approval') { await requestForBot(member, 'approval.respond', { @@ -8124,7 +8172,7 @@ function groupRoomCoordinator(group) { let coordinator = groupRoomCoordinators.get(group) if (!coordinator || coordinator.roomId !== (room.roomId || null) || coordinator.roomToken !== (room.coordinationId || null)) { coordinator = { group, roomId: room.roomId || null, roomToken: room.coordinationId || null, resumeQueue: [], occurrences: new Set(), - queue: [], active: 0, members: new Map(), drives: new Map() } + queue: [], budgetPumps: new Set(), active: 0, members: new Map(), drives: new Map(), capturedDrives: new Map() } groupRoomCoordinators.set(group, coordinator) } return coordinator @@ -8152,7 +8200,7 @@ function registerGroupOccurrence(group, member, thread, deliveryResult = {}, pre occurrence.memberLock = groupSourceSessionKey(captured, `room:${room.roomId || 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 => e.from?.kind === 'user' && groupThreadOf(e) === occurrence.thread).map(groupChatSyncEntryKey) + occurrence.userIds = (room.log || []).filter(e => groupIsUserInstruction(e) && groupThreadOf(e) === occurrence.thread).map(groupChatSyncEntryKey) const seen = new Set(groupConsumedInputs(room)[`${occurrence.thread}::${occurrence.memberKey}`] || []) occurrence.consumedIds = (room.log || []).filter(e => groupThreadOf(e) === occurrence.thread && !seen.has(groupChatSyncEntryKey(e))).map(groupChatSyncEntryKey) coordinator.occurrences.add(occurrence) // before retain/session preparation @@ -8169,6 +8217,7 @@ function groupOccurrenceCurrent(occurrence) { function groupOccurrenceCanSubmit(occurrence) { if (!groupOccurrenceCurrent(occurrence)) return false const room = $groupChats.get()[occurrence.group] + if (occurrence.drive?.onlyHandoffs && !groupDriveMemberIsCurrent(room, occurrence.captured)) return false return !Array.isArray(occurrence.userIds) || !groupTurnHasNewerUser(room, occurrence.thread, occurrence.userIds, occurrence.inputEndId, occurrence.inputVersion) } @@ -8284,7 +8333,19 @@ function reserveGroupResumeWorker(occurrence) { return pending } +// Reload does not prove that an accepted execution released its worker. A +// cold marker without a live occurrence conservatively occupies one until its +// exact terminal is consumed; no runtime job or session lock is replayed. +function groupOccupiedWorkers(coordinator) { + const room = $groupChats.get()[coordinator.group] + if (!room || (room.roomId || null) !== coordinator.roomId || + (room.coordinationId || null) !== coordinator.roomToken) return coordinator.active + const owned = new Set([...coordinator.occurrences].filter(o => !o.released).map(o => o.id)) + return coordinator.active + Object.values(room.stranded || {}).filter(m => !owned.has(m?.occurrence_id)).length +} + function pumpGroupOccurrences(coordinator) { + for (const pump of [...coordinator.budgetPumps]) pump() for (const waiter of [...coordinator.resumeQueue]) { if (!groupOccurrenceCurrent(waiter.occurrence)) { coordinator.resumeQueue.splice(coordinator.resumeQueue.indexOf(waiter), 1) @@ -8292,7 +8353,7 @@ function pumpGroupOccurrences(coordinator) { waiter.reject(new Error('Member occurrence was stopped or replaced')) continue } - if (coordinator.active >= GROUP_CHAT_PARALLEL_CEILING) break + if (groupOccupiedWorkers(coordinator) >= GROUP_CHAT_PARALLEL_CEILING) break coordinator.resumeQueue.splice(coordinator.resumeQueue.indexOf(waiter), 1) coordinator.active++ waiter.occurrence.ownsWorker = true @@ -8304,7 +8365,7 @@ function pumpGroupOccurrences(coordinator) { // still perform authoritative admission; four is never a slots-free probe. // Waiting collectors keep sockets/acceptance custody, not admission workers. for (const occurrence of [...coordinator.queue]) { - if (coordinator.active >= GROUP_CHAT_PARALLEL_CEILING) break + if (groupOccupiedWorkers(coordinator) >= GROUP_CHAT_PARALLEL_CEILING) break if (!coordinator.queue.includes(occurrence)) continue if (coordinator.members.has(occurrence.memberLock)) continue coordinator.queue.splice(coordinator.queue.indexOf(occurrence), 1) @@ -8424,8 +8485,10 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima let marker = { delivery: { accepted_turn: null, member_key: memberKey, owner }, thread: thread || 'legacy', epoch: occurrence?.epoch ?? dispatchEpoch, room_id: occurrence?.roomId ?? roomAtDispatch.roomId ?? null, room_token: occurrence?.roomToken, occurrence_id: occurrence?.id, + drive_key: occurrence?.drive?.key, consumed_ids: occurrence?.consumedIds, user_ids: occurrence?.userIds, input_version: occurrence?.inputVersion, anchor_id: occurrence?.inputEndId || roomAtDispatch.log?.at(-1)?.id || null } + if (occurrence?.drive) persistGroupDrive(group, occurrence.drive) let settleCollector const collector = { settled: new Promise(resolve => { settleCollector = resolve }) } collectingGroupTurnMarkers.set(marker, collector) @@ -8619,6 +8682,7 @@ async function runGroupChatMemberTurnLeased(group, captured, prompt, thread, ima if (!fence || (fence === resumeFenceAtPoll && responseAcknowledgedAtPoll)) { occurrence.resumeFence = null occurrence.phase = 'waiting' + occurrence.wasWaiting = true releaseGroupOccurrenceWorker(occurrence) pumpGroupOccurrences(occurrence.coordinator) occurrence.markReady() // parked, but still owns its lease, marker and reservation @@ -8689,7 +8753,8 @@ async function harvestStrandedGroupReply(group, member) { if (ownedOccurrence?.answerPromise) return const requested = member === ownedOccurrence?.captured.member ? ownedOccurrence.captured : captureGroupTurnMember(member) - if (groupTurnOwnerKey(marker?.delivery?.owner) !== groupTurnOwnerKey(requested.owner)) { + if (groupTurnOwnerKey(marker?.delivery?.owner) !== groupTurnOwnerKey(requested.owner) || + (room.sessionOwners?.[memberKey] && groupTurnOwnerKey(room.sessionOwners[memberKey]) !== groupTurnOwnerKey(requested.owner))) { reportUnavailableGroupTurn(group, member, marker, 'Member source or target profile changed') return } @@ -8766,7 +8831,10 @@ async function harvestStrandedGroupReply(group, member) { consumeGroupTurnMarker(group, memberKey, marker) return } - if (!consumeGroupTurnMarker(group, memberKey, marker)) return + const continuationDrive = ownedOccurrence?.drive || restoreGroupDrive(group, marker) + const publishes = outcome.state === 'complete' && !isGroupPassText(outcome.reply) + if (!consumeGroupTurnMarker(group, memberKey, marker, publishes)) return + if (ownedOccurrence) settleGroupPublication(ownedOccurrence, publishes) if (Array.isArray(marker.consumed_ids)) consumeGroupInputs(group, `${marker.thread}::${memberKey}`, marker.consumed_ids, marker.anchor_id, current.roomId || null) if (outcome.state === 'error' || outcome.state === 'interrupted') { @@ -8780,6 +8848,7 @@ async function harvestStrandedGroupReply(group, member) { appendGroupChatEntry(group, { kind: 'member', name: captured.member.name, ...(captured.member.remoteSource ? { source: captured.member.connectionLabel || captured.member.connectionId } : {}) }, outcome.reply, marker.thread, undefined, marker.delivery) + scheduleGroupDriveContinuation(group, continuationDrive) } } @@ -9018,15 +9087,19 @@ function consumeGroupInputs(group, markKey, ids, inputEndId, roomId) { }) } +function groupIsUserInstruction(entry) { + return entry?.from?.kind === 'user' && !groupTurnDeliveryKey(entry.answer_to) +} + function groupTurnHasNewerUser(room, thread, userIds, anchorId, inputVersion) { if (inputVersion !== undefined && (room.threadInputVersions?.[thread] || 0) !== inputVersion) return true if (Array.isArray(userIds)) { - return room.log.some(e => e.from?.kind === 'user' && groupThreadOf(e) === thread && + return room.log.some(e => groupIsUserInstruction(e) && groupThreadOf(e) === thread && !userIds.includes(groupChatSyncEntryKey(e))) } const anchor = room.log.findIndex(e => e.id === anchorId) return (anchor >= 0 ? room.log.slice(anchor + 1) : room.log).some( - e => e.from?.kind === 'user' && groupThreadOf(e) === thread) + e => groupIsUserInstruction(e) && groupThreadOf(e) === thread) } /** Stop owns ALL captured occurrences, including queue/lease/session races. @@ -9039,8 +9112,9 @@ async function stopGroupThread(group, thread, members = null) { ? [...coordinator.occurrences] : [] const roster = Array.isArray(members) && members.length ? members : room.members || [] const stamp = { at: Date.now(), byMessageId: null, thread: thread || null } + for (const occurrence of occurrences) occurrence.cancelled = true for (const occurrence of occurrences) { - occurrence.cancelled = true + if (occurrence.phase === 'budget-queued') finishGroupOccurrence(occurrence) markGroupOccurrenceStop(occurrence) syncGroupClarify(group, occurrence.captured.member, null) } @@ -9072,6 +9146,11 @@ async function stopGroupThread(group, thread, members = null) { // identify exact runtime and source, never a display-name roster guess. for (const [memberKey, marker] of Object.entries(room.stranded || {})) { if (occurrences.some(o => o.memberKey === memberKey)) continue + const cardKey = `${group}::${memberKey}` + const card = $groupClarify.get()[cardKey] + if (card?.receipt === marker || card?.sessionId === marker?.delivery?.accepted_turn?.session_id) { + const cards = { ...$groupClarify.get() }; delete cards[cardKey]; $groupClarify.set(cards) + } updateGroupChat(group, r => { if (r.stranded?.[memberKey] === marker) { marker.stop_requested = true; r.stranded = { ...r.stranded } } return r @@ -9117,15 +9196,85 @@ async function stopGroupThread(group, thread, members = null) { return result } +function groupRoomCanStop(room, clarifies = []) { + return Boolean(room?.running || clarifies.length || room?.turns?.length || + Object.values(room?.stranded || {}).some(marker => groupTurnDeliveryKey(marker?.delivery))) +} + +// Only accepted/pending receipts retain continuation context on disk. Original +// inputs and source pins survive reload; neither prompts nor runtime jobs replay. +function durableGroupDriveContexts(room) { + const keys = new Set(Object.values(room.stranded || {}).map(marker => marker?.drive_key).filter(Boolean)) + return Object.fromEntries(Object.entries(room.driveContexts || {}).filter(([key]) => keys.has(key))) +} + +function groupDriveKey(drive) { + return JSON.stringify([drive.roomId, drive.roomToken, drive.thread, drive.epoch, drive.userEntryId]) +} + +function persistGroupDrive(group, drive) { + const room = $groupChats.get()[group] + if (!room || (room.roomId || null) !== drive.roomId || (room.coordinationId || null) !== drive.roomToken) return + updateGroupChat(group, current => { + current.driveContexts = { ...durableGroupDriveContexts(current), [drive.key]: { + version: 1, roomId: drive.roomId, roomToken: drive.roomToken, epoch: drive.epoch, + thread: drive.thread, userEntryId: drive.userEntryId, userIds: drive.userIds, + inputIds: drive.inputIds, inputVersion: drive.inputVersion, members: drive.members, budget: drive.budget + } } + return current + }, { sync: false }) +} + +function restoreGroupDrive(group, marker) { + const room = $groupChats.get()[group], saved = room?.driveContexts?.[marker.drive_key] + if (!saved || saved.version !== 1 || groupDriveKey(saved) !== marker.drive_key || + saved.roomId !== (room.roomId || null) || saved.roomToken !== (room.coordinationId || null) || + saved.thread !== marker.thread || saved.epoch !== marker.epoch || + !Array.isArray(saved.members) || !saved.members.length || saved.members.length > GROUP_CHAT_MAX_MEMBERS || + !saved.members.every(c => typeof c?.member?.name === 'string' && c.member.name && + typeof groupMemberKey(c.member) === 'string' && groupTurnOwnerKey(c.owner) && + c.requestMember?.name === c.member.name && + groupTurnOwnerKey(groupSessionOwner(c.requestMember)) === groupTurnOwnerKey(c.owner)) || + new Set(saved.members.map(c => groupMemberKey(c.member))).size !== saved.members.length || + !['inputIds', 'userIds'].every(key => Array.isArray(saved[key]) && saved[key].length <= GROUP_CHAT_HISTORY_LIMIT * 16 && + saved[key].every(id => typeof id === 'string' && id)) || + typeof saved.thread !== 'string' || !saved.thread || !Number.isSafeInteger(saved.epoch) || saved.epoch < 0 || + !Number.isSafeInteger(saved.inputVersion) || saved.inputVersion < 0 || + !(saved.userEntryId === null || typeof saved.userEntryId === 'string') || !saved.budget || + !['posted', 'reserved', 'rounds', 'continuations'].every(key => Number.isSafeInteger(saved.budget[key]) && saved.budget[key] >= 0) || + saved.budget.posted + saved.budget.reserved > GROUP_CHAT_MAX_MESSAGES || + saved.budget.rounds > GROUP_CHAT_MAX_ROUNDS || saved.budget.continuations > GROUP_CHAT_MAX_CONTINUATIONS) return null + const source = saved.members.find(c => groupMemberKey(c.member) === marker.delivery.member_key) + if (!source || groupTurnOwnerKey(source.owner) !== groupTurnOwnerKey(marker.delivery.owner)) return null + const coordinator = groupRoomCoordinator(group) + if (coordinator.capturedDrives.has(marker.drive_key)) return coordinator.capturedDrives.get(marker.drive_key) + // Unstarted in-memory queue jobs vanish on reload. Every producer writes its + // marker before session preparation; unknown markers keep their reservation. + const reserved = Object.values(room.stranded || {}).filter(m => m?.drive_key === marker.drive_key).length + if (saved.budget.posted + reserved > GROUP_CHAT_MAX_MESSAGES) return null + const drive = { ...saved, key: marker.drive_key, budget: { ...saved.budget, reserved }, onlyHandoffs: true } + coordinator.capturedDrives.set(drive.key, drive) + return drive +} + +function groupDriveMemberIsCurrent(room, captured) { + const key = groupMemberKey(captured.member) + const current = (room.members || []).find(member => groupMemberKey(member) === key) + if (!current || groupTurnOwnerKey(captureGroupTurnMember(current).owner) !== groupTurnOwnerKey(captured.owner)) return false + const sessionOwner = room.sessionOwners?.[key] + return !sessionOwner || groupTurnOwnerKey(sessionOwner) === groupTurnOwnerKey(captured.owner) +} + function captureGroupDrive(group, members, thread) { groupRoomCoordinator(group) const room = $groupChats.get()[group] || {} - const user = (room.log || []).findLast(e => e.from?.kind === 'user' && groupThreadOf(e) === thread) + const user = (room.log || []).findLast(e => groupIsUserInstruction(e) && groupThreadOf(e) === thread) return { roomId: room.roomId || null, roomToken: room.coordinationId || null, epoch: room.epoch || 0, thread, userEntryId: user ? groupChatSyncEntryKey(user) : null, + inputIds: (room.log || []).map(groupChatSyncEntryKey), inputVersion: room.threadInputVersions?.[thread] || 0, members: members.map(member => captureGroupTurnMember(member)), - userIds: (room.log || []).filter(e => e.from?.kind === 'user' && groupThreadOf(e) === thread).map(groupChatSyncEntryKey) } + userIds: (room.log || []).filter(e => groupIsUserInstruction(e) && groupThreadOf(e) === thread).map(groupChatSyncEntryKey) } } /** A frozen round, never a raw Promise.all over roster members. The room @@ -9133,18 +9282,60 @@ function captureGroupDrive(group, members, thread) { function runGroupChatRounds(group, members, thread, capturedDrive) { const drive = capturedDrive || captureGroupDrive(group, members, thread) const coordinator = groupRoomCoordinator(group) - const key = JSON.stringify([drive.roomId, drive.roomToken, drive.thread, drive.epoch, drive.userEntryId]) + const key = groupDriveKey(drive) + drive.key = key + drive.group = group const prior = coordinator.drives.get(key) if (prior) return prior // Defer the body one microtask so drive ownership exists before first await. + drive.budget ||= { posted: 0, reserved: 0, continuations: 0, rounds: 0 } + coordinator.capturedDrives.set(key, drive) + persistGroupDrive(group, drive) + drive.running = true const running = Promise.resolve().then(() => driveFrozenGroupRounds(group, members, drive, coordinator)) coordinator.drives.set(key, running) void running.finally(() => { 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)) { + coordinator.capturedDrives.delete(key) + const room = $groupChats.get()[group] + if (room && !room.tombstone && (room.roomId || null) === drive.roomId && + (room.coordinationId || null) === drive.roomToken) { + updateGroupChat(group, current => { current.driveContexts = durableGroupDriveContexts(current); return current }, { sync: false }) + } + } + if (drive.continuationRequested) { + drive.continuationRequested = false + scheduleGroupDriveContinuation(group, drive) + } }).catch(() => undefined) return running } +function settleGroupPublication(job, published = false) { + if (!job?.reservation) return + job.reservation = false + job.drive.budget.reserved-- + if (published) job.drive.budget.posted++ + persistGroupDrive(job.group, job.drive) + job.onBudgetRelease?.() +} + +function scheduleGroupDriveContinuation(group, drive) { + const room = $groupChats.get()[group] + if (!drive || !room || room.tombstone || (room.roomId || null) !== drive.roomId || + (room.coordinationId || null) !== drive.roomToken || + groupTurnHasNewerUser(room, drive.thread, drive.userIds, drive.userEntryId, drive.inputVersion) || + drive.budget.continuations >= GROUP_CHAT_MAX_CONTINUATIONS || + drive.budget.posted + drive.budget.reserved >= GROUP_CHAT_MAX_MESSAGES || + !unaddressedGroupMentions(group, drive.members.map(c => c.member), drive.thread).length) return + if (drive.running) { drive.continuationRequested = true; return } + drive.onlyHandoffs = true + // The original captured drive owns identity, consumed inputs and all budgets. + void runGroupChatRounds(group, drive.members.map(c => c.member), drive.thread, drive).catch(() => undefined) +} + async function driveFrozenGroupRounds(group, members, drive, coordinator) { const { thread } = drive members = drive.members.map(captured => captured.member) @@ -9153,22 +9344,26 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { return room && !room.tombstone && (room.roomId || null) === drive.roomId && (room.coordinationId || null) === drive.roomToken } - const isCurrent = () => sameRoom() && ($groupChats.get()[group].epoch || 0) === drive.epoch - let posted = 0, reserved = 0, continuations = 0 + const isCurrent = () => sameRoom() && (drive.onlyHandoffs + ? !groupTurnHasNewerUser($groupChats.get()[group], thread, drive.userIds, drive.userEntryId, drive.inputVersion) + : ($groupChats.get()[group].epoch || 0) === drive.epoch) + const budget = drive.budget let exitKind = 'settled' const freeze = responders => { const room = $groupChats.get()[group] const consumed = groupConsumedInputs(room) - const userIds = room.log.filter(e => e.from?.kind === 'user' && groupThreadOf(e) === thread).map(groupChatSyncEntryKey) + const userIds = room.log.filter(e => groupIsUserInstruction(e) && groupThreadOf(e) === thread).map(groupChatSyncEntryKey) const inputEndId = room.log.length ? groupChatSyncEntryKey(room.log.at(-1)) : null const jobs = [] for (const member of responders) { const memberKey = groupMemberKey(member) + if (drive.onlyHandoffs && !groupDriveMemberIsCurrent(room, drive.members.find(c => groupMemberKey(c.member) === memberKey))) continue const markKey = `${thread}::${memberKey}` const seen = new Set(consumed[markKey] || []) const delta = room.log.filter(e => groupThreadOf(e) === thread && !seen.has(groupChatSyncEntryKey(e)) && // A member's own final is never fresh peer input; its siblings' are. - !(e.delivery?.member_key === memberKey)) + !(e.delivery?.member_key === memberKey) && + !(drive.onlyHandoffs && drive.inputIds?.includes(groupChatSyncEntryKey(e)))) if (!delta.length || groupTurnMarkerBlocksDispatch(room, memberKey, coordinator) || [...coordinator.occurrences].some(o => o.memberKey === memberKey && o.phase === 'waiting' && !o.cancelled)) continue if (room.holds?.[memberKey]) { @@ -9188,6 +9383,8 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { occurrence.inputVersion = room.threadInputVersions?.[thread] || 0 occurrence.userIds = userIds occurrence.markKey = markKey + occurrence.drive = drive + occurrence.phase = 'budget-queued' occurrence.prompt = buildGroupChatTurnPrompt({ groupName: group, members, viewer: member, deltaLines: delta.slice(-GROUP_CHAT_HISTORY_LIMIT).map(e => formatGroupChatLine(e, member.name)) }) occurrence.images = delta.flatMap(e => Array.isArray(e.images) ? e.images : []).map(image => Object.freeze({ ...image })) @@ -9198,14 +9395,30 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { const dispatch = async responders => { const jobs = freeze(responders) let spoke = 0 - for (const job of jobs) { - if (posted + reserved >= GROUP_CHAT_MAX_MESSAGES) { - exitKind = 'capped' - finishGroupOccurrence(job) - continue + const pending = [...jobs] + const admitted = [] + const pumpPublication = () => { + for (const job of [...pending]) { + if (!pending.includes(job)) continue + if (job.released || budget.posted >= GROUP_CHAT_MAX_MESSAGES || !groupOccurrenceCanSubmit(job) || !sameRoom() || + groupTurnHasNewerUser($groupChats.get()[group], thread, drive.userIds, drive.userEntryId, drive.inputVersion)) { + pending.splice(pending.indexOf(job), 1) + finishGroupOccurrence(job) + continue + } + if (budget.posted + budget.reserved >= GROUP_CHAT_MAX_MESSAGES) break + pending.splice(pending.indexOf(job), 1) + admitted.push(job) + start(job) } - reserved++ // reserve potential publication BEFORE async admission + if (!pending.length) coordinator.budgetPumps.delete(pumpPublication) + } + const start = job => { + budget.reserved++ // reserve potential publication BEFORE async admission job.reservation = true + persistGroupDrive(group, drive) + job.onBudgetRelease = pumpPublication + job.phase = 'queued' const complete = runGroupChatMemberTurn(group, job.captured.member, job.prompt, thread, job.images, job.deliveryResult, job).then(reply => { const room = $groupChats.get()[group] @@ -9224,7 +9437,8 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { appendGroupChatEntry(group, { kind: 'member', name: job.captured.member.name, ...(job.captured.member.remoteSource ? { source: job.captured.member.connectionLabel || job.captured.member.connectionId } : {}) }, reply, thread, undefined, job.deliveryResult.value) - posted++; spoke++ + settleGroupPublication(job, true); spoke++ + if (job.wasWaiting) scheduleGroupDriveContinuation(group, drive) } }, error => { const room = $groupChats.get()[group] @@ -9239,19 +9453,31 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { noteBotAttention(job.memberKey, reason || error?.message || error) }).finally(() => { if (job.reservation && !Object.prototype.hasOwnProperty.call($groupChats.get()[group]?.stranded || {}, job.memberKey)) { - reserved--; job.reservation = false + settleGroupPublication(job) } job.markReady() }) // Prevent parked collector failures from becoming unhandled rejections. void complete.catch(() => undefined) } - // Only a round barrier; execution admission belongs to the coordinator. - for (const job of jobs) await job.ready + coordinator.budgetPumps.add(pumpPublication) + pumpPublication() + // Admit frozen excess jobs whenever terminal non-publications release budget. + // Waiting/unknown keeps its reservation; the round may quiesce with excess queued. + let observed = 0 + while (observed < admitted.length) { + const batch = admitted.slice(observed); observed = admitted.length + for (const job of batch) await job.ready + pumpPublication() + } + if (pending.length) exitKind = 'capped' return spoke } try { - for (let round = 0; round < GROUP_CHAT_MAX_ROUNDS; round++) { + while (drive.onlyHandoffs ? budget.continuations < GROUP_CHAT_MAX_CONTINUATIONS : budget.rounds < GROUP_CHAT_MAX_ROUNDS) { + const round = budget.rounds + if (!drive.onlyHandoffs) budget.rounds++ + persistGroupDrive(group, drive) if (!sameRoom()) return const room = $groupChats.get()[group] // A callback captured for an older thread may still run its FIRST round. @@ -9261,17 +9487,25 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { for (const member of members) await harvestStrandedGroupReply(group, member) if (!sameRoom() || groupTurnHasNewerUser($groupChats.get()[group], thread, drive.userIds, drive.userEntryId, drive.inputVersion)) return const log = $groupChats.get()[group].log.filter(e => groupThreadOf(e) === thread) - let spoke = await dispatch(rotateGroupSpeakers(resolveGroupResponders(log, members), round)) + let spoke + if (drive.onlyHandoffs) { + const pendingKeys = unaddressedGroupMentions(group, members, thread) + if (!pendingKeys.length || budget.continuations >= GROUP_CHAT_MAX_CONTINUATIONS) return + budget.continuations++ + persistGroupDrive(group, drive) + spoke = await dispatch(members.filter(m => pendingKeys.includes(groupMemberKey(m)))) + } else spoke = await dispatch(rotateGroupSpeakers(resolveGroupResponders(log, members), round)) if (!isCurrent()) return - if (posted + reserved >= GROUP_CHAT_MAX_MESSAGES) { exitKind = 'capped'; return } + if (budget.posted + budget.reserved >= GROUP_CHAT_MAX_MESSAGES) { exitKind = 'capped'; return } if (!spoke) { + if (drive.onlyHandoffs) return const pendingKeys = unaddressedGroupMentions(group, members, thread) - continuations++ - if (pendingKeys.length && continuations <= GROUP_CHAT_MAX_CONTINUATIONS) { + if (pendingKeys.length && budget.continuations >= GROUP_CHAT_MAX_CONTINUATIONS) { exitKind = 'capped'; return } + if (pendingKeys.length) { budget.continuations++; persistGroupDrive(group, drive) } + if (pendingKeys.length) { spoke = await dispatch(members.filter(m => pendingKeys.includes(groupMemberKey(m)))) } if (!spoke) { - if (pendingKeys.length && continuations > GROUP_CHAT_MAX_CONTINUATIONS) exitKind = 'capped' return } } @@ -9288,7 +9522,9 @@ async function driveFrozenGroupRounds(group, members, drive, coordinator) { } else { recordGroupActivity(group, { kind: exitKind, member: null, thread }) } - updateGroupChat(group, r => { r.running = false; r.turn = null; return r }) + if ((current.epoch || 0) === drive.epoch) { + updateGroupChat(group, r => { r.running = false; r.turn = null; return r }) + } if (Object.keys($groupChats.get()[group].stranded || {}).length && typeof window !== 'undefined') { void harvestStrandedUntilSettled(group, members, thread) } @@ -13826,7 +14062,8 @@ 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') + appendGroupChatEntry(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}` }) } finally { @@ -14381,7 +14618,7 @@ function GroupChatWorkspace({ group, members, onBack, visible = true }) { : null ] }), - (room.running || roomClarifies.length || room.turns?.length) + groupRoomCanStop(room, roomClarifies) ? jsx('button', { type: 'button', title: 'Stop this run — interrupts all active members and holds queued turns', @@ -16486,7 +16723,26 @@ function BotsPane() { // ── plugin ─────────────────────────────────────────────────────────────────── +/** Group room host boundary shared by production callers and offline behavior tests. */ +const groupTurnRuntime = { + runGroupChatMemberTurn, runGroupChatRounds, harvestStrandedGroupReply, stopGroupChatServerSync, + currentGroupActivity, groupActivityLabel, $groupChats, $groupClarify, $botAttention, + appendGroupChatEntry, syncGroupClarify, answerGroupClarify, groupChatSyncSnapshot, + groupChatSyncEntryKey, stopGroupThread, mergeRemoteGroupChatSnapshotIntoRooms, + durableGroupChatRooms, sendToGroupChat, groupRoomCoordinators, groupRuntimeSessionOwners, + groupMemberKey, updateGroupChat, groupBlockedMembers, GroupBlockedNotice, + CreateGroupChatDialog, createFreshGroupChat, groupComposerDraftKey, + groupComposerDraftSnapshot, updateGroupComposerDraft, GroupChatWorkspace, GroupClarifyCard, + groupRoomCanStop, + bindGroupTurnPorts(ports) { + ({ Date, setTimeout, clearTimeout, setInterval, clearInterval, document } = createGroupTurnPorts(ports)) + }, + bindGroupTurnStorage(storage) { pluginCtx = { ...pluginCtx, storage } } +} + + export default { + groupTurnRuntime, id: ID, name: 'Bots', description: 'Bot Mode — a one-chat-per-agent roster with avatars, routines, group chats, and bot-to-bot messaging. Ships with the app; disable here if unwanted.', @@ -16637,6 +16893,8 @@ export default { sessionOwners: room.sessionOwners && typeof room.sessionOwners === 'object' ? room.sessionOwners : {}, recoveryOrigin: room.recoveryOrigin && typeof room.recoveryOrigin === 'object' ? room.recoveryOrigin : null, stranded: room.stranded && typeof room.stranded === 'object' ? room.stranded : {}, + driveContexts: room.driveContexts && typeof room.driveContexts === 'object' ? room.driveContexts : {}, + coordinationId: typeof room.coordinationId === 'string' ? room.coordinationId : null, // #93129: rehydrate sticky stop holds with the same shape // guard as the other maps — a held bot stays held across // window restarts until explicitly released. @@ -17081,3 +17339,5 @@ export default { }) } } + + 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 71b86cfa0ae51..d646bc6b7a18a 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 @@ -1901,6 +1901,7 @@ test('older backends without pending_clarify never mirror a question', () => { test('answerGroupClarify routes clarify.respond and clears the mirror', async () => { const gc = load(() => '(pass)') + gc.updateGroupChat('Core', room => room) const member = { name: 'research', title: '' } gc.syncGroupClarify('Core', member, { session_id: 'runtime-research', pending_clarify: CLARIFY_PAYLOAD }) @@ -1914,6 +1915,7 @@ test('answerGroupClarify routes clarify.respond and clears the mirror', async () test('answerGroupClarify sends one respond per batch question, in order', async () => { const gc = load(() => '(pass)') + gc.updateGroupChat('Core', room => room) const member = { name: 'research', title: '' } const batch = { request_id: 'req-batch-1', @@ -2023,6 +2025,7 @@ test('an approval without a server choice set falls back to once/deny', () => { test('answerGroupClarify routes approvals through approval.respond with session + choice', async () => { const gc = load(() => '(pass)') + gc.updateGroupChat('Core', room => room) const member = { name: 'research', title: '' } gc.syncGroupClarify('Core', member, { session_id: 'rt-research-1', pending_approval: APPROVAL_PAYLOAD }) 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 1927848fa85b2..cfc16a188af51 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 @@ -300,8 +300,10 @@ test('an ordinary newer-send epoch bump WITHOUT a hold does not abandon the poll test('Stop button: workspace renders it while the room is running and wires it to stopGroupThread', () => { const workspace = pluginSource.slice(pluginSource.indexOf('function GroupChatWorkspace')) - // Visible while a round is running… - assert.match(workspace, /room\.running \|\| roomClarifies\.length \|\| room\.turns\?\.length\)\s*\?\s*jsx\('button'/, 'Stop button is gated on room.running') + assert.match(workspace, /groupRoomCanStop\(room, roomClarifies\)\s*\?\s*jsx\('button'/, 'Stop button uses the shared room custody predicate') + const predicate = pluginSource.slice(pluginSource.indexOf('function groupRoomCanStop'), pluginSource.indexOf('function captureGroupDrive')) + assert.match(predicate, /room\?\.running \|\| clarifies\.length \|\| room\?\.turns\?\.length/, 'running work, question cards and pending turns expose Stop') + assert.match(predicate, /Object\.values\(room\?\.stranded \|\| \{\}\)\.some\(marker => groupTurnDeliveryKey\(marker\?\.delivery\)\)/, 'durable accepted custody also exposes Stop') // …and wired to the real primitive, not a per-member interrupt spray. assert.match(workspace, /stopGroupThread\(/, 'button calls the stopGroupThread primitive') assert.doesNotMatch(workspace, /stopAllBots/, 'the #94570 interrupt-only shell was rewired, not kept') diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-turn-test-loader.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-turn-test-loader.mjs index fbed552ff359f..e87ca9a85b811 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/group-turn-test-loader.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-turn-test-loader.mjs @@ -1,7 +1,6 @@ import { registerHooks } from 'node:module' -// Import the whole real plugin with inert renderer dependencies. The loader -// exposes its existing private engine seam; tests make no source-text assertions. +// Import the unmodified shipped plugin; mock only its explicit dependencies. const fixtures = new Map() globalThis.__groupTurnTestFixtures = fixtures let sequence = 0 @@ -37,25 +36,6 @@ registerHooks({ } } return nextResolve(specifier, context) - }, - load(url, context, nextLoad) { - const result = nextLoad(url, context) - const target = new URL(url) - if (target.pathname === pluginURL.pathname && target.searchParams.has('groupTurnFixture')) { - return { - ...result, - source: `const { Date, setTimeout, clearTimeout } = globalThis.__groupTurnTestFixtures.get(${JSON.stringify(target.searchParams.get('groupTurnFixture'))});\nconst fixturePorts = globalThis.__groupTurnTestFixtures.get(${JSON.stringify(target.searchParams.get('groupTurnFixture'))});\n const document = fixturePorts.document ?? globalThis.document;\n const setInterval = fixturePorts.setInterval ?? globalThis.setInterval;\n const clearInterval = fixturePorts.clearInterval ?? globalThis.clearInterval;\n${result.source}\nexport { runGroupChatMemberTurn, runGroupChatRounds, - harvestStrandedGroupReply, stopGroupChatServerSync, currentGroupActivity, groupActivityLabel, - $groupChats, $groupClarify, $botAttention, appendGroupChatEntry, syncGroupClarify, answerGroupClarify, - groupChatSyncSnapshot, groupChatSyncEntryKey, stopGroupThread, - mergeRemoteGroupChatSnapshotIntoRooms, durableGroupChatRooms, sendToGroupChat, - groupRoomCoordinators, groupRuntimeSessionOwners, groupMemberKey, updateGroupChat, - groupBlockedMembers, GroupBlockedNotice, CreateGroupChatDialog, createFreshGroupChat, - groupComposerDraftKey, groupComposerDraftSnapshot, updateGroupComposerDraft };\n - export const groupRecoveryTestAPI = {\n createFreshGroupChat: typeof createFreshGroupChat === 'function' ? createFreshGroupChat : undefined,\n GroupBlockedNotice: typeof GroupBlockedNotice === 'function' ? GroupBlockedNotice : undefined,\n groupBlockedMembers: typeof groupBlockedMembers === 'function' ? groupBlockedMembers : undefined,\n updateGroupComposerDraft, GroupChatWorkspace, CreateGroupChatDialog, groupComposerDraftSnapshot, groupComposerDraftKey };\n export function bindGroupTurnTestStorage(storage) { pluginCtx = { storage }; }\n` - } - } - return result } }) @@ -65,7 +45,12 @@ export async function importGroupTurnPlugin(fixture) { const url = new URL(pluginURL) url.searchParams.set('groupTurnFixture', id) try { - return await import(url.href) + const module = await import(url.href) + const runtime = module.default.groupTurnRuntime + runtime.bindGroupTurnPorts(fixture) + return { ...module, ...runtime, + bindGroupTurnTestStorage: runtime.bindGroupTurnStorage, + groupRecoveryTestAPI: runtime } } finally { fixtures.delete(id) } diff --git a/apps/desktop/src/plugins/hermes-bots/tests/hide-bot-chats.runtime.test.ts b/apps/desktop/src/plugins/hermes-bots/tests/hide-bot-chats.runtime.test.ts index da3f06f2bc8c9..f966a08130e02 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/hide-bot-chats.runtime.test.ts +++ b/apps/desktop/src/plugins/hermes-bots/tests/hide-bot-chats.runtime.test.ts @@ -53,11 +53,13 @@ afterEach(() => { gatewayState.set('closed') vi.clearAllMocks() vi.useRealTimers() + plugin.groupTurnRuntime.bindGroupTurnPorts(globalThis) }) describe('Bot Mode hidden-session reconciliation lifecycle', () => { it('uses persisted REST on load/reconnect and stops with plugin disposal', async () => { vi.useFakeTimers() + plugin.groupTurnRuntime.bindGroupTurnPorts(globalThis) const disposers: Array<() => void> = [] plugin.register(createPluginContext(plugin.id, dispose => disposers.push(dispose))) 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 new file mode 100644 index 0000000000000..f7ec4c18992d0 --- /dev/null +++ b/apps/desktop/src/plugins/hermes-bots/tests/pr118-postmerge.test.mjs @@ -0,0 +1,439 @@ +import assert from 'node:assert/strict' +import test from 'node:test' +import { harness, members, deferred, flush, drive } from './stop-custody-harness.mjs' + +const boundedTest = (name, body) => test(name, { timeout: 10000 }, body) +const clone = value => JSON.parse(JSON.stringify(value)) +const receipts = h => Object.values(h.room().stranded || {}) +const cards = h => Object.values(h.gc.$groupClarify.get()) +const nodes = tree => { + const result = [] + const walk = value => { + if (Array.isArray(value)) value.forEach(walk) + else if (value && typeof value === 'object') { result.push(value); walk(value.props?.children) } + } + walk(tree) + return result +} + +// Render the shipped components with stateful React ports. Event handlers, +// answer ownership and the Stop visibility predicate remain production code. +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: () => '' }, + document: { getElementById: () => true, addEventListener: () => {}, removeEventListener: () => {} }, + setInterval: () => 1, clearInterval: () => {}, + render: (component, props, fresh = false) => { + cursor = 0 + if (fresh) states.length = 0 + return component(props) + } + } +} + +async function uiHarness(roster = members(1), options = {}) { + const ui = renderer(), resumeProjection = options.resumeProjection + Object.assign(options, ui, { + resumeProjection: (session, method, projection) => { + if (session.approval && session.state === 'waiting') projection.pending_approval = session.approval + return resumeProjection?.(session, method, projection) || projection + } }) + const h = await harness(roster, options) + h.ui = ui + h.workspace = () => ui.render(h.gc.GroupChatWorkspace, + { group: 'Room', members: h.roster, onBack: () => {} }, true) + h.stopButton = () => nodes(h.workspace()).find(node => node.type === 'button' && + node.props?.title?.startsWith('Stop this run')) + return h +} + +function preparedAnswer(h, entry = cards(h)[0]) { + assert.ok(entry, 'the actual pending projection exposes a card') + const props = { entry, members: h.roster } + let tree = h.ui.render(h.gc.GroupClarifyCard, props, true) + if (entry.kind === 'approval') { + const choice = nodes(tree).find(node => node.props?.children === 'once' && node.props?.onClick) + assert.ok(choice); choice.props.onClick() + } else { + const inputs = nodes(tree).filter(node => node.props?.['aria-label'] === `Answer @${entry.member}`) + assert.equal(inputs.length, entry.questions?.length || 1) + inputs.forEach((node, index) => node.props.onChange({ target: { value: `ANSWER_${index + 1}` } })) + } + tree = h.ui.render(h.gc.GroupClarifyCard, props) + const submit = nodes(tree).find(node => ['Answer', 'Respond'].includes(node.props?.children) && node.props?.onClick) + assert.ok(submit); assert.equal(submit.props.disabled, false) + return () => submit.props.onClick() +} + +async function park(h, kind = 'single') { + const pending = drive(h) + await flush() + const session = [...h.sessions.values()][0] + assert.ok(session) + session.state = 'waiting' + session.pending = kind === 'batch' + ? { request_id: 'question-1', questions: [{ qid: 'first', question: 'First?' }, { qid: 'second', question: 'Second?' }] } + : { request_id: 'question-1', question: 'Choose?', choices: [] } + if (kind === 'approval') { + session.pending = null + session.approval = { request_id: 'question-1', command: 'echo offline', choices: ['once', 'deny'] } + } + await h.advance(); await pending + if (kind === 'approval' && !cards(h).length) { + h.gc.syncGroupClarify('Room', h.roster[0], { + session_id: session.runtime, pending_approval: { request_id: 'question-1', command: 'echo offline', choices: ['once', 'deny'] } + }, h.roster[0]) + } + assert.equal(receipts(h).length, 1) + return session +} + +for (const kind of ['single', 'batch', 'approval']) { + boundedTest(`actual ${kind} card answer preserves accepted ownership and delivers one final`, async () => { + const h = await uiHarness(), session = await park(h, kind) + const accepted = clone(session.ref), beforeVersion = h.room().threadInputVersions?.t1 || 0 + session.text = `FINAL_${kind}` + preparedAnswer(h)(); await flush() + await h.until(() => h.posts().some(entry => entry.text === `FINAL_${kind}`)) + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await flush() + assert.equal(h.rpc(kind === 'approval' ? 'approval.respond' : 'clarify.respond').length, kind === 'batch' ? 2 : 1) + assert.deepEqual(session.ref, accepted) + assert.equal(h.room().threadInputVersions?.t1 || 0, beforeVersion) + assert.equal(h.posts().filter(entry => entry.text === `FINAL_${kind}`).length, 1) + assert.equal(h.rpc('prompt.submit').length, 1, 'accepted prompt is never replayed') + assert.ok(h.room().log.some(entry => entry.from.kind === 'user' && entry.text.includes(kind === 'approval' ? 'once' : 'ANSWER_1'))) + }) +} + +for (const kind of ['single', 'batch', 'approval']) { + boundedTest(`failed ${kind} card answer produces no successful echo`, async () => { + const h = await uiHarness(members(1), { answerError: () => true }) + await park(h, kind) + const before = clone(h.room().log) + preparedAnswer(h)(); await flush(); await h.advance() + assert.deepEqual(h.room().log, before) + assert.equal(h.posts().length, 0) + assert.equal(receipts(h).length, 1) + await h.gc.stopGroupThread('Room', 't1', h.roster) + }) +} + +boundedTest('remote projected answer echo preserves ownership; genuine remote input still supersedes', async () => { + const h = await uiHarness(members(1), { resumeRunning: true }), session = await park(h) + const preAnswer = { ...h.room(), log: [...h.room().log] }, accepted = clone(session.ref) + preparedAnswer(h)(); await flush() + const snapshot = h.gc.groupChatSyncSnapshot() + const echo = Object.values(snapshot.rooms)[0].log.find(entry => entry.id !== 'user-1' && entry.from.kind === 'user') + assert.ok(echo, 'sync projection retains the visible exchange') + const merged = h.gc.mergeRemoteGroupChatSnapshotIntoRooms(snapshot, { Room: preAnswer }) + assert.equal(merged.Room.threadInputVersions?.t1 || 0, preAnswer.threadInputVersions?.t1 || 0) + h.gc.$groupChats.set(merged) + session.state = 'complete'; session.text = 'AFTER_REMOTE_ECHO' + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await h.advance() + assert.equal(h.posts().filter(entry => entry.text === 'AFTER_REMOTE_ECHO').length, 1) + assert.deepEqual(session.ref, accepted) + + const control = await uiHarness(), old = await park(control) + const remote = control.gc.groupChatSyncSnapshot(), epoch = control.room().epoch + Object.values(remote.rooms)[0].log.push({ id: 'remote-genuine', at: 100050, + from: { kind: 'user', name: 'You' }, text: 'GENUINE_NEW_INSTRUCTION', thread: 't1' }) + control.gc.$groupChats.set(control.gc.mergeRemoteGroupChatSnapshotIntoRooms(remote, control.gc.$groupChats.get())) + assert.equal(control.room().epoch, epoch, 'remote supersession is independent of epoch') + old.state = 'complete'; old.pending = null; old.text = 'STALE_RESULT' + await control.gc.harvestStrandedGroupReply('Room', control.roster[0]); await control.advance() + assert.equal(control.posts().filter(entry => entry.text === 'STALE_RESULT').length, 0) +}) + +async function reload(hot, options = {}) { + const h = await uiHarness(hot.roster, options), durable = clone(hot.gc.durableGroupChatRooms()) + for (const [id, session] of hot.sessions) h.sessions.set(id, clone(session)) + h.gc.$groupChats.set({}) + h.gc.default.register({ storage: { get: key => key === 'group-chats' ? durable : null, + set: (key, value) => h.storage.set(key, clone(value)) }, register: () => {}, onDispose: () => {} }) + await flush(); h.gc.stopGroupChatServerSync() + assert.equal(h.room().running, false) + assert.equal(h.room().turns?.length || 0, 0) + assert.equal(cards(h).length, 0) + assert.equal(h.gc.currentGroupActivity('Room').length, 0) + return h +} + +for (const failsFirst of [false, true]) { + boundedTest(`cold Stop clears card and fences captured Answer; retry=${failsFirst}`, async () => { + const hot = await uiHarness(); await park(hot) + const options = { interruptError: failsFirst }, cold = await reload(hot, options) + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]) + assert.equal(cards(cold).length, 1) + const staleClick = preparedAnswer(cold), accepted = clone(receipts(cold)[0].delivery.accepted_turn) + await cold.gc.stopGroupThread('Room', 't1', cold.roster) + assert.equal(cards(cold).length, 0) + staleClick(); await flush() + assert.equal(cold.rpc('clarify.respond').length, 0, 'a captured stale handler cannot authorize an RPC') + assert.equal(receipts(cold).length, failsFirst ? 1 : 0) + if (failsFirst) { + assert.ok(receipts(cold)[0].stop_requested) + options.interruptError = false + await cold.gc.stopGroupThread('Room', 't1', cold.roster) + assert.deepEqual(cold.rpc('session.interrupt').map(call => call.params.session_id), [accepted.session_id, accepted.session_id]) + assert.equal(receipts(cold).length, 0) + } + await hot.gc.stopGroupThread('Room', 't1', hot.roster) + }) +} + +boundedTest('cold current card accepts equivalent cloned receipt metadata after fresh waiting sync', async () => { + const hot = await uiHarness(); await park(hot) + const cold = await reload(hot) + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]) + const original = receipts(cold)[0], accepted = clone(original.delivery.accepted_turn) + cold.gc.updateGroupChat('Room', room => ({ ...room, stranded: { + ...room.stranded, 'local::bot1': { ...clone(original), reported: true, reason: 'temporarily unavailable' } + } })) + assert.notEqual(receipts(cold)[0], original, 'metadata reconciliation clones the receipt object') + assert.deepEqual(receipts(cold)[0].delivery.accepted_turn, accepted) + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]) + const session = [...cold.sessions.values()][0] + session.text = 'CLONED_RECEIPT_RESULT' + preparedAnswer(cold)(); await flush() + assert.equal(cold.rpc('clarify.respond').length, 1, 'same accepted identity keeps its current card actionable') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]) + assert.equal(cold.posts().filter(entry => entry.text === 'CLONED_RECEIPT_RESULT').length, 1) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) +}) + +for (const changed of ['accepted tuple', 'source']) { + boundedTest(`replacement cold receipt with changed ${changed} fences captured Answer`, async () => { + const hot = await uiHarness(); await park(hot) + const cold = await reload(hot) + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]) + const staleClick = preparedAnswer(cold), replacement = clone(receipts(cold)[0]) + if (changed === 'accepted tuple') replacement.delivery.accepted_turn.request_id = 'different-accepted-request' + else replacement.delivery.owner.route.connectionId = 'different-source' + cold.gc.updateGroupChat('Room', room => ({ ...room, + stranded: { ...room.stranded, 'local::bot1': replacement } })) + staleClick(); await flush() + assert.equal(cold.rpc('clarify.respond').length, 0) + assert.equal(cold.room().log.filter(entry => entry.from.kind === 'user').length, 1) + assert.equal(receipts(cold).length, 1, 'replacement acceptance custody survives rejection') + await hot.gc.stopGroupThread('Room', 't1', hot.roster) + }) +} + +boundedTest('old answer completion preserves an unrelated replacement question card', async () => { + const answerAckGate = deferred(), h = await uiHarness(members(1), { answerAckGate }) + const session = await park(h) + preparedAnswer(h)(); await flush() + assert.equal(h.rpc('clarify.respond').length, 1) + session.state = 'waiting' + session.pending = { request_id: 'replacement-question', question: 'A distinct question?' } + h.gc.syncGroupClarify('Room', h.roster[0], { + session_id: session.runtime, pending_clarify: session.pending + }, h.roster[0]) + const replacement = cards(h)[0] + assert.equal(replacement.requestId, 'replacement-question') + answerAckGate.resolve(); await flush(); await h.advance() + assert.equal(cards(h)[0], replacement) + assert.equal(cards(h)[0].requestId, 'replacement-question') + await h.gc.stopGroupThread('Room', 't1', h.roster) +}) + +boundedTest('actual UI exposes cold Stop, retains it after failed interrupt and reload, then removes retry', async () => { + const hot = await uiHarness(); await park(hot) + const cold = await reload(hot, { interruptError: true }), first = cold.stopButton() + assert.ok(first, 'durable accepted custody exposes the actual Stop button without transient state') + const accepted = clone(receipts(cold)[0].delivery.accepted_turn) + first.props.onClick(); await flush() + assert.equal(receipts(cold).length, 1) + assert.ok(cold.room().holds['local::bot1']) + const retried = await reload(cold), retry = retried.stopButton() + assert.ok(retry, 'unconfirmed held custody remains actionable after reload') + retry.props.onClick(); await flush() + assert.deepEqual(retried.rpc('session.interrupt').map(call => [call.route.connectionId, call.params.session_id]), + [['local', accepted.session_id]]) + assert.equal(receipts(retried).length, 0) + assert.ok(retried.room().holds['local::bot1']) + assert.equal(retried.stopButton(), undefined) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) +}) + +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: [], + stranded: { 'local::bot1': { before: 1, thread: 't1', reported: true } } })) + assert.equal(h.stopButton(), undefined) + await h.gc.stopGroupThread('Room', 't1', h.roster) + assert.equal(h.rpc('session.interrupt').length, 0) + assert.equal(receipts(h).length, 1) +}) + +boundedTest('late accepted answer drives its mentioned teammate once with peer delta', async () => { + const h = await uiHarness(members(2)) + h.gc.updateGroupChat('Room', room => { + room.log[0].text = '@bot1 FIRST_INPUT' + room.log[0].images = [{ kind: 'file', data: 'ORIGINAL_ATTACHMENT', name: 'original.txt' }] + return room + }) + const session = await park(h) + assert.equal(h.rpc('prompt.submit').length, 1) + assert.equal(h.room().running, false, 'first drive is quiescent while the accepted turn waits') + session.text = 'LATE_PEER_REPLY @bot2' + preparedAnswer(h)(); await flush() + await h.until(() => h.rpc('prompt.submit').some(call => call.route.profile === 'bot2')) + h.finish('bot2') + await h.advance(); await h.gc.harvestStrandedGroupReply('Room', h.roster[0]) + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await h.advance() + const teammate = h.rpc('prompt.submit').filter(call => call.route.profile === 'bot2') + assert.equal(teammate.length, 1) + assert.match(teammate[0].params.text, /LATE_PEER_REPLY/) + assert.doesNotMatch(teammate[0].params.text, /FIRST_INPUT/) + assert.equal(h.calls.filter(call => call.method.includes('attach') && call.route.profile === 'bot2').length, 0) + assert.equal(h.rpc('prompt.submit').filter(call => call.route.profile === 'bot1').length, 1) + assert.ok(h.maxLeases() <= 4) + assert.ok(h.posts().length <= 10) +}) + +boundedTest('new user send in another thread preserves the parked thread late-answer handoff', async () => { + const h = await uiHarness(members(2)) + h.gc.updateGroupChat('Room', room => { room.log[0].text = '@bot1 FIRST_INPUT'; return room }) + const session = await park(h), originalEpoch = h.room().epoch + h.gc.sendToGroupChat('Room', h.roster, '@bot1 OTHER_THREAD_INSTRUCTION', 't2') + await h.advance(250) + assert.ok(h.room().epoch > originalEpoch, 'ordinary Send advances the room drive generation') + assert.equal(h.rpc('prompt.submit').length, 1, 'parked A keeps its accepted custody while t2 is blocked') + session.text = 'OLDER_THREAD_PEER_REPLY @bot2' + preparedAnswer(h)(); await flush() + await h.until(() => h.rpc('prompt.submit').some(call => call.route.profile === 'bot2')) + h.finish('bot2'); await h.advance() + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await h.advance() + const teammate = h.rpc('prompt.submit').filter(call => call.route.profile === 'bot2') + assert.equal(teammate.length, 1) + assert.match(teammate[0].params.text, /OLDER_THREAD_PEER_REPLY/) + assert.doesNotMatch(teammate[0].params.text, /FIRST_INPUT|OTHER_THREAD_INSTRUCTION/) + assert.equal(h.posts().filter(entry => entry.text === 'OLDER_THREAD_PEER_REPLY @bot2' && entry.thread === 't1').length, 1) +}) + +for (const invalidation of ['Stop', 'newer input', 'replacement']) { + boundedTest(`${invalidation} prevents stale late-answer teammate continuation`, async () => { + const h = await uiHarness(members(2), { resumeRunning: true }) + h.gc.updateGroupChat('Room', room => { room.log[0].text = '@bot1 FIRST_INPUT'; return room }) + const session = await park(h) + preparedAnswer(h)(); await flush() + if (invalidation === 'Stop') await h.gc.stopGroupThread('Room', 't1', h.roster) + else if (invalidation === 'newer input') h.gc.appendGroupChatEntry('Room', { kind: 'user', name: 'You' }, + 'GENUINE_NEW_INSTRUCTION', 't1') + else h.gc.updateGroupChat('Room', room => ({ ...room, roomId: 'replacement-room', epoch: room.epoch + 1 })) + session.state = 'complete'; session.pending = null; session.text = 'STALE_HANDOFF @bot2' + await h.gc.harvestStrandedGroupReply('Room', h.roster[0]); await h.advance(5000) + assert.equal(h.posts().filter(entry => entry.text.includes('STALE_HANDOFF')).length, 0) + assert.equal(h.rpc('prompt.submit').filter(call => call.route.profile === 'bot2').length, 0) + assert.equal(h.rpc('prompt.submit').filter(call => call.route.profile === 'bot1').length, 1) + }) +} + +for (const terminal of ['complete', 'error']) { + boundedTest(`publication capacity backfills after four ${terminal === 'complete' ? 'passes' : 'terminal failures'}`, async () => { + const secondAdmitted = new Set(), maxWorkers = { value: 0 } + let h + const options = { resumeProjection: (session, method, projection) => { + if (method !== 'session.turn.poll' || !session.ref) return projection + maxWorkers.value = Math.max(maxWorkers.value, h.gc.groupRoomCoordinators.get('Room').active) + if (session.submits === 2) secondAdmitted.add(session.profile) + const omitted = session.submits === 2 && [...secondAdmitted].indexOf(session.profile) < 4 + session.state = omitted ? terminal : 'complete' + session.text = session.submits === 1 ? `ROUND_ONE_${session.profile} @all` + : omitted || session.submits > 2 ? '(pass)' : `BACKFILL_${session.profile}` + projection.running = false + projection.turn_outcomes.turns[0] = { accepted_turn: { ...session.ref }, state: session.state, + finalized: [{ status: session.state, text: session.text, ...(session.state === 'error' ? { error: 'offline terminal failure' } : {}) }] } + return projection + } } + h = await uiHarness(members(6), options) + const pending = drive(h) + await h.until(() => !h.room().running); await pending + assert.equal(secondAdmitted.size, 6, 'all frozen second-round jobs get their freed-capacity admission') + assert.equal(h.posts().filter(entry => entry.text.startsWith('ROUND_ONE_')).length, 6) + assert.equal(h.posts().filter(entry => entry.text.startsWith('BACKFILL_')).length, 2) + assert.ok(maxWorkers.value <= 4) + assert.ok(h.posts().length <= 10) + for (const session of h.sessions.values()) { + const prompts = h.rpc('prompt.submit').filter(call => call.params.session_id === session.runtime) + assert.ok(prompts.length >= 2) + assert.doesNotMatch(prompts[1].params.text, /FIRST_INPUT/, 'backfill receives its frozen peer delta without replay') + } + }) +} + +for (const invalidation of ['Stop', 'newer input']) { +boundedTest(`unresolved waiting reservations retain capacity; ${invalidation} cancels frozen excess jobs`, async () => { + 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.submits === 2) secondAdmitted.add(session.profile) + session.state = session.submits === 1 ? 'complete' : 'waiting' + session.text = `FIRST_${session.profile} @all` + if (session.state === 'waiting') { + session.pending = { request_id: `waiting-${session.profile}`, question: 'Continue?' } + projection.pending_clarify = session.pending + } + projection.running = false + projection.turn_outcomes.turns[0] = { accepted_turn: { ...session.ref }, state: session.state, + finalized: session.state === 'complete' ? [{ status: 'complete', text: session.text }] : [] } + return projection + } }) + const pending = drive(h) + await h.until(() => secondAdmitted.size === 4) + await h.advance(); await pending + assert.equal(secondAdmitted.size, 4, 'waiting may release workers, but its publication custody stays reserved') + assert.equal(h.posts().length, 6) + assert.equal(receipts(h).length, 4) + const accepted = receipts(h).map(marker => clone(marker.delivery.accepted_turn)) + if (invalidation === 'Stop') await h.gc.stopGroupThread('Room', 't1', h.roster) + else h.gc.appendGroupChatEntry('Room', { kind: 'user', name: 'You' }, 'NEWER_GENUINE_INPUT', 't1') + await h.advance() + assert.equal(secondAdmitted.size, 4, 'superseded budget-pending frozen jobs never start') + assert.equal(h.rpc('prompt.submit').length, 10) + assert.equal(h.posts().length, 6) + const coordinator = h.gc.groupRoomCoordinators.get('Room') + assert.equal(coordinator.occurrences.size, invalidation === 'Stop' ? 0 : 4, + 'only unresolved accepted occurrences retain custody') + if (invalidation === 'newer input') { + assert.deepEqual(receipts(h).map(marker => marker.delivery.accepted_turn), accepted) + assert.equal(h.rpc('session.interrupt').length, 0, 'newer intent does not prove old accepted work ended') + await h.gc.stopGroupThread('Room', 't1', h.roster) + } +}) +} + +boundedTest('chatty continuations share the cumulative ten-message cap and four-worker ceiling', async () => { + let h, maxWorkers = 0 + h = await uiHarness(members(6), { resumeProjection: (session, method, projection) => { + if (method !== 'session.turn.poll' || !session.ref) return projection + maxWorkers = Math.max(maxWorkers, h.gc.groupRoomCoordinators.get('Room').active) + session.state = 'complete'; session.text = `CHATTER_${session.profile}_${session.submits} @all` + projection.running = false + projection.turn_outcomes.turns[0] = { accepted_turn: { ...session.ref }, state: 'complete', + finalized: [{ status: 'complete', text: session.text }] } + return projection + } }) + const pending = drive(h) + await h.until(() => !h.room().running); await pending + assert.equal(h.posts().length, 10) + assert.equal(h.rpc('prompt.submit').length, 10) + assert.ok(maxWorkers <= 4) + assert.equal(h.gc.groupRoomCoordinators.get('Room').occurrences.size, 0, + 'the message cap releases never-admitted frozen occurrences') + await h.advance(5000) + assert.equal(h.posts().length, 10) + assert.equal(h.rpc('prompt.submit').length, 10) +}) diff --git a/apps/desktop/src/plugins/hermes-bots/tests/pr123-continuation.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/pr123-continuation.test.mjs new file mode 100644 index 0000000000000..a73ccc6187192 --- /dev/null +++ b/apps/desktop/src/plugins/hermes-bots/tests/pr123-continuation.test.mjs @@ -0,0 +1,411 @@ +import assert from 'node:assert/strict' +import test from 'node:test' +import { harness, members, deferred, flush, drive } from './stop-custody-harness.mjs' + +const boundedTest = (name, body) => test(name, { timeout: 12000 }, body) +const clone = value => JSON.parse(JSON.stringify(value)) +const receipts = h => Object.values(h.room().stranded || {}) +const submits = (h, name) => h.rpc('prompt.submit').filter(call => call.route?.profile === name) +const nodes = tree => { + const out = [] + const walk = value => { + if (Array.isArray(value)) value.forEach(walk) + else if (value && typeof value === 'object') { out.push(value); walk(value.props?.children) } + } + walk(tree); return out +} + +async function uiHarness(roster = members(2), options = {}) { + const states = []; let cursor = 0 + Object.assign(options, { + 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: () => '' }, + document: { getElementById: () => true, addEventListener: () => {}, removeEventListener: () => {} } + }) + const h = await harness(roster, options) + h.answer = () => { + const entry = Object.values(h.gc.$groupClarify.get()).find(card => card.member === 'bot1') + assert.ok(entry, 'production waiting projection exposes the actual card') + states.length = 0; cursor = 0 + const props = { entry, members: h.roster } + let tree = h.gc.GroupClarifyCard(props) + nodes(tree).find(node => node.props?.['aria-label'] === 'Answer @bot1').props.onChange({ target: { value: 'ANSWER' } }) + cursor = 0; tree = h.gc.GroupClarifyCard(props) + const button = nodes(tree).find(node => node.props?.children === 'Answer' && node.props?.onClick) + assert.ok(button); assert.equal(button.props.disabled, false) + button.props.onClick() + } + return h +} + +function terminalProjection(session, projection, text) { + session.state = 'complete'; session.text = text + projection.running = false + delete projection.pending_clarify + projection.turn_outcomes.turns[0] = { accepted_turn: { ...session.ref }, state: 'complete', finalized: [{ status: 'complete', text }] } + return projection +} + +async function thirdRoundWaiting({ chain = false } = {}) { + let h + h = await uiHarness(members(2), { resumeProjection: (session, method, projection) => { + if (method !== 'session.turn.poll' || !session.ref) return projection + if (session.profile === 'bot1' && session.submits === 3 && !session.answeredRequest) { + session.state = 'waiting'; session.pending = { request_id: 'third-round-question', question: 'Continue?' } + projection.pending_clarify = session.pending + projection.running = false + projection.turn_outcomes.turns[0] = { accepted_turn: { ...session.ref }, state: 'waiting', finalized: [] } + return projection + } + if (session.profile === 'bot1' && session.submits === 3 && session.answeredRequest) return projection + if (chain && session.submits === 4) return terminalProjection(session, projection, + session.profile === 'bot2' ? 'CONTINUATION_ONE @bot1' : 'CONTINUATION_TWO @bot2') + return terminalProjection(session, projection, session.submits < 3 ? `NORMAL_${session.profile}_${session.submits} @all` : '(pass)') + } }) + const pending = drive(h) + await h.until(() => submits(h, 'bot1').length === 3 && receipts(h).length === 1) + await h.advance(); await pending + assert.equal(submits(h, 'bot2').length, 3) + assert.equal(h.room().running, false) + return h +} + +async function firstRoundWaiting(roster = members(2)) { + const h = await uiHarness(roster) + h.gc.updateGroupChat('Room', room => { + room.log[0].text = '@bot1 ORIGINAL_INSTRUCTION' + room.log[0].images = [{ kind: 'file', data: 'ORIGINAL_ATTACHMENT', name: 'original.txt' }] + return room + }) + const pending = drive(h); await flush() + const session = [...h.sessions.values()][0] + session.state = 'waiting'; session.pending = { request_id: 'cold-question', question: 'Continue?' } + await h.advance(); await pending + assert.equal(receipts(h).length, 1) + return h +} + +async function reload(hot, options = {}) { + const h = await uiHarness(hot.roster, options), durable = clone(hot.gc.durableGroupChatRooms()) + // A fresh fake wire starts its sequence at zero. Namespace copied runtime + // handles so newly-created continuation sessions cannot collide with them. + const ids = new Map() + for (const session of hot.sessions.values()) { + const copied = clone(session), runtime = `cold:${session.runtime}`, stored = `cold:${session.stored}` + ids.set(session.runtime, runtime); ids.set(session.stored, stored) + copied.runtime = runtime; copied.stored = stored + if (copied.ref) copied.ref.session_id = runtime + h.sessions.set(runtime, copied) + } + for (const room of Object.values(durable)) { + for (const key of Object.keys(room.sessions || {})) room.sessions[key] = ids.get(room.sessions[key]) || room.sessions[key] + for (const marker of Object.values(room.stranded || {})) { + if (marker.delivery?.accepted_turn) marker.delivery.accepted_turn.session_id = ids.get(marker.delivery.accepted_turn.session_id) || marker.delivery.accepted_turn.session_id + } + } + h.gc.$groupChats.set({}) + h.gc.default.register({ storage: { get: key => key === 'group-chats' ? durable : null, + set: (key, value) => h.storage.set(key, clone(value)) }, register: () => {}, onDispose: () => {} }) + await flush(); h.gc.stopGroupChatServerSync() + assert.equal(h.room().running, false) + assert.equal(h.room().turns?.length || 0, 0) + return h +} + +boundedTest('third normal round waiting answer schedules a bounded fourth teammate turn', async () => { + const h = await thirdRoundWaiting(), session = [...h.sessions.values()].find(s => s.profile === 'bot1') + session.text = 'HOT_LATE @bot2' + h.answer(); await flush() + await h.until(() => submits(h, 'bot2').length === 4) + await h.advance() + assert.equal(submits(h, 'bot1').length, 3) + assert.equal(submits(h, 'bot2').length, 4) + assert.match(submits(h, 'bot2')[3].params.text, /HOT_LATE/) + assert.doesNotMatch(submits(h, 'bot2')[3].params.text, /FIRST_INPUT/) + assert.equal(h.posts().filter(entry => entry.text === 'HOT_LATE @bot2').length, 1) +}) + +boundedTest('cold hydrated accepted terminal harvest continues its teammate once using peer delta', async () => { + const hot = await firstRoundWaiting(), cold = await reload(hot) + cold.finish('bot1', 'COLD_LATE @bot2') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]); await flush() + await cold.until(() => submits(cold, 'bot2').length === 1) + const teammate = submits(cold, 'bot2')[0] + assert.match(teammate.params.text, /COLD_LATE/) + assert.doesNotMatch(teammate.params.text, /ORIGINAL_INSTRUCTION|ORIGINAL_ATTACHMENT/) + assert.equal(cold.calls.filter(call => call.method.includes('attach') && call.route?.profile === 'bot2').length, 0) + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]) + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]); await cold.advance() + assert.equal(submits(cold, 'bot2').length, 1, 'duplicate harvest cannot schedule a second teammate') + cold.finish('bot2'); await cold.advance() + assert.equal(cold.posts().filter(entry => entry.text === 'COLD_LATE @bot2').length, 1) + assert.equal(submits(cold, 'bot1').length, 0) + const reloaded = await reload(cold) + await reloaded.gc.harvestStrandedGroupReply('Room', reloaded.roster[0]); await reloaded.advance() + assert.equal(reloaded.rpc('prompt.submit').length, 0, 'consumed cold publication cannot replay after another reload') + assert.equal(reloaded.posts().filter(entry => entry.text === 'COLD_LATE @bot2').length, 1) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) +}) + +for (const change of ['other thread', 'newer same thread', 'Stop', 'replaced room', 'changed owner', 'changed recipient']) { + boundedTest(`cold continuation control: ${change}`, async () => { + const hot = await firstRoundWaiting(), cold = await reload(hot) + if (change === 'other thread') { + cold.gc.sendToGroupChat('Room', cold.roster, '@bot1 OTHER_THREAD_INPUT', 't2') + await cold.advance(250) + } else if (change === 'newer same thread') cold.gc.appendGroupChatEntry('Room', { kind: 'user', name: 'You' }, 'NEWER_INPUT', 't1') + else if (change === 'Stop') await cold.gc.stopGroupThread('Room', 't1', cold.roster) + else if (change === 'replaced room') cold.gc.updateGroupChat('Room', room => ({ ...room, roomId: 'replacement-room' })) + else if (change === 'changed owner') cold.gc.updateGroupChat('Room', room => ({ ...room, sessionOwners: { + ...room.sessionOwners, 'local::bot1': { ...room.sessionOwners['local::bot1'], + route: { ...room.sessionOwners['local::bot1'].route, connectionId: 'replacement-source' } } + } })) + else cold.gc.updateGroupChat('Room', room => ({ ...room, members: room.members.map(member => member.name !== 'bot2' + ? member : { ...member, connectionId: 'replacement-peer', route: { ...member.route, connectionId: 'replacement-peer' } }) })) + cold.finish('bot1', 'CONTROL_LATE @bot2') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]); await cold.advance() + if (change === 'other thread') { + await cold.until(() => submits(cold, 'bot2').length === 1) + assert.match(submits(cold, 'bot2')[0].params.text, /CONTROL_LATE/) + assert.doesNotMatch(submits(cold, 'bot2')[0].params.text, /OTHER_THREAD_INPUT|ORIGINAL_INSTRUCTION/) + cold.finish('bot2'); await cold.advance() + } else { + assert.equal(submits(cold, 'bot2').length, 0) + assert.equal(cold.posts().filter(entry => entry.text === 'CONTROL_LATE @bot2').length, + change === 'changed recipient' ? 1 : 0, + 'a valid accepted source reply may publish while the changed recipient loses admission authority') + } + assert.equal(cold.rpc('prompt.submit').filter(call => call.route?.connectionId === 'replacement-source').length, 0) + assert.equal(cold.rpc('prompt.submit').filter(call => call.route?.connectionId === 'replacement-peer').length, 0) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) + }) +} + +boundedTest('hot resumed handoffs retain the two-continuation bound after normal rounds are exhausted', async () => { + const h = await thirdRoundWaiting({ chain: true }), session = [...h.sessions.values()].find(s => s.profile === 'bot1') + session.text = 'HOT_CHAIN_START @bot2' + h.answer(); await flush() + await h.until(() => submits(h, 'bot1').length === 4) + await h.advance(5000); await h.advance(5000) + assert.equal(submits(h, 'bot1').length, 4) + assert.equal(submits(h, 'bot2').length, 4, 'third late handoff is beyond the continuation budget') + assert.equal(h.posts().length, 7) + assert.ok(h.posts().length <= 10) + assert.ok(h.maxLeases() <= 4) +}) + +boundedTest('cold handoff chain preserves its two-continuation budget across repeated reloads', async () => { + const hot = await firstRoundWaiting(members(4)) + let cold = await reload(hot) + for (const [from, to] of [['bot1', 'bot2'], ['bot2', 'bot3']]) { + cold.finish(from, `COLD_CHAIN_${from} @${to}`) + await cold.gc.harvestStrandedGroupReply('Room', cold.roster.find(member => member.name === from)) + await cold.until(() => submits(cold, to).length === 1) + const session = [...cold.sessions.values()].find(item => item.profile === to) + session.state = 'waiting'; session.pending = { request_id: `chain-${to}`, question: 'Continue?' } + await cold.advance(); await cold.until(() => !cold.room().running) + assert.equal(receipts(cold).length, 1) + cold = await reload(cold) + } + cold.finish('bot3', 'COLD_CHAIN_THIRD @bot4') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[2]); await cold.advance(5000) + assert.equal(cold.posts().length, 3) + assert.equal(submits(cold, 'bot4').length, 0, 'a third handoff cannot reset the durable continuation budget') + await hot.gc.stopGroupThread('Room', 't1', hot.roster) +}) + +boundedTest('cold resumed targeted work uses at most four execution workers across five teammates', async () => { + const hot = await firstRoundWaiting(members(6)), cold = await reload(hot) + cold.finish('bot1', 'COLD_BROADCAST @bot2 @bot3 @bot4 @bot5 @bot6') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]); await flush() + await cold.until(() => cold.rpc('prompt.submit').length >= 1) + assert.ok(cold.gc.groupRoomCoordinators.get('Room').active <= 4) + for (let attempt = 0; attempt < 20 && (cold.rpc('prompt.submit').length < 5 || + cold.gc.groupRoomCoordinators.get('Room').active); attempt++) { + assert.ok(cold.gc.groupRoomCoordinators.get('Room').active <= 4) + for (const name of ['bot2', 'bot3', 'bot4', 'bot5', 'bot6']) cold.finish(name) + await cold.advance() + } + await cold.until(() => cold.rpc('prompt.submit').length === 5 && + cold.gc.groupRoomCoordinators.get('Room').active === 0) + assert.equal(cold.rpc('prompt.submit').length, 5) + assert.ok(cold.maxLeases() <= 4) + assert.equal(cold.posts().length, 1) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) +}) + +async function capacityWaiting(unknown = false) { + const h = await uiHarness(members(6), { resumeProjection: (session, method, projection) => { + if (method !== 'session.turn.poll' || !session.ref) return projection + if (session.submits === 1) return terminalProjection(session, projection, `INITIAL_${session.profile} @all`) + if (session.profile === 'bot2') { + session.state = 'waiting'; session.pending = { request_id: 'capacity-question', question: 'Continue?' } + projection.pending_clarify = session.pending; projection.running = false + projection.turn_outcomes.turns[0] = { accepted_turn: { ...session.ref }, state: 'waiting', finalized: [] } + return projection + } + if (unknown) { + session.state = 'running'; projection.running = true + projection.turn_outcomes.availability = 'unavailable'; projection.turn_outcomes.turns = [] + return projection + } + return terminalProjection(session, projection, `SECOND_${session.profile} @all`) + } }) + const pending = drive(h) + await h.until(() => h.rpc('prompt.submit').length === 10) + await h.until(() => !h.room().running); await pending + assert.equal(receipts(h).length, unknown ? 4 : 1) + assert.equal(h.posts().length, unknown ? 6 : 9) + return h +} + +boundedTest('cold terminal publication respects the cumulative ten-message cap', async () => { + const hot = await capacityWaiting(), cold = await reload(hot) + cold.finish('bot2', 'TENTH_PUBLICATION @bot6') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[1]); await cold.advance(5000) + assert.equal(cold.posts().length, 10) + assert.equal(cold.posts().filter(entry => entry.text === 'TENTH_PUBLICATION @bot6').length, 1) + assert.equal(cold.rpc('prompt.submit').length, 0, 'no eleventh publication is reserved through a cold handoff') + const reloaded = await reload(cold) + await reloaded.gc.harvestStrandedGroupReply('Room', reloaded.roster[1]); await reloaded.advance() + assert.equal(reloaded.posts().length, 10) + assert.equal(reloaded.rpc('prompt.submit').length, 0) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) +}) + +boundedTest('cold continuation retains publication capacity for three unknown accepted receipts', async () => { + const hot = await capacityWaiting(true) + const cold = await reload(hot, { resumeProjection: (session, method, projection) => { + if (method === 'session.turn.poll' && session.profile !== 'bot2' && session.submits === 2) { + projection.turn_outcomes.availability = 'unavailable'; projection.turn_outcomes.turns = [] + } + return projection + } }) + const unknown = receipts(cold).filter(marker => marker.delivery.member_key !== 'local::bot2') + .map(marker => clone(marker.delivery.accepted_turn)) + cold.finish('bot2', 'UNKNOWN_CAPACITY_LATE @bot6') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[1]); await cold.advance(5000) + assert.equal(cold.posts().length, 7) + assert.equal(cold.rpc('prompt.submit').length, 0, 'unknown accepted work reserves the remaining three slots') + assert.deepEqual(receipts(cold).map(marker => marker.delivery.accepted_turn), unknown) + assert.equal(cold.rpc('session.interrupt').length, 0) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) +}) + +boundedTest('three cold unknown accepted workers permit one fresh teammate, then terminal releases backfill the second', async () => { + const hot = await uiHarness(members(6)) + hot.gc.updateGroupChat('Room', room => { + room.log[0].text = '@bot1 @bot2 @bot3 @bot4 ORIGINAL_WORKER_INPUT' + room.log[0].images = [{ kind: 'file', data: 'ORIGINAL_WORKER_ATTACHMENT', name: 'original.txt' }] + return room + }) + const hotDrive = drive(hot); await flush() + assert.equal(hot.rpc('prompt.submit').length, 4) + assert.equal(receipts(hot).length, 4) + const unresolved = new Set(['bot2', 'bot3', 'bot4']) + const cold = await reload(hot, { resumeProjection: (session, method, projection) => { + if (method === 'session.turn.poll' && unresolved.has(session.profile)) { + projection.turn_outcomes.availability = 'unavailable'; projection.turn_outcomes.turns = [] + } + return projection + } }) + cold.finish('bot1', 'FRESH_HANDOFF @bot5 @bot6') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]); await flush() + await cold.until(() => cold.rpc('prompt.submit').length >= 1) + const coordinator = cold.gc.groupRoomCoordinators.get('Room') + assert.equal(cold.posts().length, 1, 'publication capacity is far below ten') + assert.equal(cold.rpc('prompt.submit').length, 1, 'three uncertain cold workers leave one fresh worker slot') + assert.equal(coordinator.active, 1) + assert.equal(receipts(cold).filter(marker => unresolved.has(marker.delivery.owner.name)).length, 3) + await cold.advance(5000) + assert.equal(cold.rpc('prompt.submit').length, 1, 'uncertain occupancy does not expire into free capacity') + for (const name of ['bot2', 'bot3', 'bot4']) { + unresolved.delete(name); cold.finish(name) + await cold.gc.harvestStrandedGroupReply('Room', cold.roster.find(member => member.name === name)); await flush() + assert.ok(coordinator.active + unresolved.size <= 4, + 'cold accepted custody and fresh active workers share the same four-worker ceiling') + } + await cold.until(() => cold.rpc('prompt.submit').length === 2) + assert.equal(coordinator.active, 2, 'freed cold capacity backfills the pending fresh teammate') + assert.deepEqual(cold.rpc('prompt.submit').map(call => call.route.profile).sort(), ['bot5', 'bot6']) + for (const call of cold.rpc('prompt.submit')) { + assert.match(call.params.text, /FRESH_HANDOFF/) + assert.doesNotMatch(call.params.text, /ORIGINAL_WORKER_INPUT|ORIGINAL_WORKER_ATTACHMENT/) + } + assert.equal(cold.calls.filter(call => call.method.includes('attach')).length, 0) + for (const name of ['bot5', 'bot6']) cold.finish(name) + await cold.advance(); await cold.until(() => coordinator.active === 0) + assert.equal(cold.rpc('prompt.submit').length, 2) + assert.ok(cold.maxLeases() <= 4) + await hot.gc.stopGroupThread('Room', 't1', hot.roster); await hot.advance(); await hotDrive +}) + +for (const change of ['delete', 'same-name replacement']) { + boundedTest(`cold active continuation finalizer respects room ${change}`, async () => { + const hot = await firstRoundWaiting(), cold = await reload(hot) + cold.finish('bot1', 'ACTIVE_BEFORE_REMOVAL @bot2') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]); await flush() + await cold.until(() => submits(cold, 'bot2').length === 1) + const coordinator = cold.gc.groupRoomCoordinators.get('Room'), running = [...coordinator.drives.values()] + assert.equal(coordinator.active, 1) + // Direct atom deletion/replacement is the room-removal state transition; + // no runtime lifecycle or profile mutation is exercised by this fixture. + const replacement = { roomId: 'same-name-new-room', coordinationId: 'same-name-new-token', epoch: 17, + running: false, log: [{ id: 'new-room-user', at: 200000, thread: 'new-thread', + from: { kind: 'user', name: 'You' }, text: 'REPLACEMENT_ROOM_INPUT' }], + members: clone(cold.roster), sessions: {}, sessionOwners: {}, stranded: {}, holds: {}, watermarks: {} } + const expected = change === 'delete' ? {} : { Room: clone(replacement) } + cold.gc.$groupChats.set(clone(expected)) + cold.finish('bot2', 'OLD_CONTINUATION_FINAL @bot1') + await cold.advance(5000); await Promise.all(running); await flush() + assert.deepEqual(cold.gc.$groupChats.get(), expected, 'old finalization cannot recreate or mutate another room lifetime') + await cold.advance(5000) + assert.deepEqual(cold.gc.$groupChats.get(), expected) + assert.equal(cold.rpc('prompt.submit').length, 1) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) + }) +} + +for (const change of ['owner changed', 'recipient removed']) { + boundedTest(`cold teammate retain race: ${change} prevents old owner submission`, async () => { + const hot = await firstRoundWaiting(), retainGate = deferred(), cold = await reload(hot, { retainGate }) + cold.finish('bot1', 'RETAIN_RACE_HANDOFF @bot2') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]); await flush() + await cold.until(() => cold.gc.groupRoomCoordinators.get('Room')?.active === 1) + const coordinator = cold.gc.groupRoomCoordinators.get('Room') + assert.equal(submits(cold, 'bot2').length, 0, 'old owner has not passed its pending retain boundary') + cold.gc.updateGroupChat('Room', room => ({ ...room, members: change === 'recipient removed' + ? room.members.filter(member => member.name !== 'bot2') + : room.members.map(member => member.name !== 'bot2' ? member : { ...member, + connectionId: 'changed-during-retain', route: { ...member.route, connectionId: 'changed-during-retain' } }) })) + retainGate.resolve(); await flush(); await cold.advance() + await cold.until(() => coordinator.active === 0) + assert.equal(submits(cold, 'bot2').length, 0) + assert.equal(cold.rpc('prompt.submit').length, 0) + assert.equal(cold.posts().filter(entry => entry.text === 'RETAIN_RACE_HANDOFF @bot2').length, 1, + 'the valid accepted source publication remains delivered') + assert.ok(cold.leases.every(lease => lease.releases === 1)) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) + }) +} + +boundedTest('malformed durable captured request owner cannot route a handoff but preserves accepted source publication', async () => { + const hot = await firstRoundWaiting(), cold = await reload(hot) + const marker = receipts(cold)[0] + cold.gc.updateGroupChat('Room', room => { + const contexts = clone(room.driveContexts), captured = contexts[marker.drive_key].members.find(item => item.member.name === 'bot2') + captured.requestMember.route.connectionId = 'malformed-request-owner' + return { ...room, driveContexts: contexts } + }) + cold.finish('bot1', 'MALFORMED_CONTEXT_FINAL @bot2') + await cold.gc.harvestStrandedGroupReply('Room', cold.roster[0]); await cold.advance(5000) + assert.equal(cold.posts().filter(entry => entry.text === 'MALFORMED_CONTEXT_FINAL @bot2').length, 1) + assert.equal(cold.rpc('prompt.submit').length, 0) + assert.equal(cold.calls.filter(call => call.route?.connectionId === 'malformed-request-owner').length, 0) + assert.equal(receipts(cold).length, 0) + await hot.gc.stopGroupThread('Room', 't1', hot.roster) +}) 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 c8978d8ffa6e1..caf570c4a3c1e 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 @@ -89,6 +89,8 @@ async function harness(roster = members(3), options = {}) { return projection } const gc = await importGroupTurnPlugin({ atom, + react: options.react, sdk: options.sdk, document: options.document, + setInterval: () => 1, clearInterval: () => {}, Date: class extends Date { static now() { return now } }, setTimeout: (fn, delay = 0) => { const id = ++sequence; timers.set(id, { fn, at: now + delay }); return id }, clearTimeout: id => timers.delete(id), @@ -105,7 +107,7 @@ async function harness(roster = members(3), options = {}) { notify: () => undefined, notifyError: () => undefined } }) gc.stopGroupChatServerSync() - gc.bindGroupTurnTestStorage({ set: (key, value) => storage.set(key, clone(value)) }) + gc.bindGroupTurnTestStorage({ get: key => clone(storage.get(key) ?? null), set: (key, value) => storage.set(key, clone(value)) }) const input = { id: 'user-1', at: now, from: { kind: 'user', name: 'You' }, text: '@all evaluate FIRST_INPUT', thread: 't1' } gc.$groupChats.set({ Room: { roomId: 'room1', epoch: 1, running: true, log: [input], watermarks: {}, sessions: {}, stranded: {}, holds: {}, members: roster } })