fix(bot-mode): keep pending group attention bound to its live room
Complete the PR #101452 salvage, preserving sprmn24 authorship. Credit @kokhlo PR #101591 for independently diagnosing the in-flight rename race; use scoped authoritative room bindings rather than persistent aliases. Keep mention attention independent and retire late work after disband.
This commit is contained in:
@@ -19,7 +19,7 @@ import type * as HermesSdk from '@hermes/plugin-sdk'
|
||||
import { fireEvent, render, screen } from '@testing-library/react'
|
||||
import { beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
import { BotRow, GroupRow } from './bot-row'
|
||||
import { BotRow } from './bot-row'
|
||||
import { translateBots } from './i18n-test-helper'
|
||||
import type { RosterRow } from './types'
|
||||
|
||||
@@ -175,23 +175,3 @@ describe('context-menu mutations hydrate the alias first', () => {
|
||||
expect(params).toMatchObject({ name: 'backend-worker', ui_meta: { 'hermes-bots': { pinned: false } } })
|
||||
})
|
||||
})
|
||||
|
||||
describe('GroupRow renders its needs-you badge from the needsYou prop', () => {
|
||||
const noopGroup = () => undefined
|
||||
|
||||
it('shows the question badge when needsYou is true', () => {
|
||||
render(
|
||||
<GroupRow active={false} group="Core" members={[]} needsYou={true} onDisband={noopGroup} onOpen={noopGroup} />
|
||||
)
|
||||
|
||||
expect(screen.getByLabelText('Needs your input')).toBeTruthy()
|
||||
})
|
||||
|
||||
it('hides the question badge when needsYou is false', () => {
|
||||
render(
|
||||
<GroupRow active={false} group="Core" members={[]} needsYou={false} onDisband={noopGroup} onOpen={noopGroup} />
|
||||
)
|
||||
|
||||
expect(screen.queryByLabelText('Needs your input')).toBeNull()
|
||||
})
|
||||
})
|
||||
|
||||
@@ -238,7 +238,7 @@ export async function disbandGroupChat(group: string, members: RosterRow[]) {
|
||||
* rename, so even a member whose sid is later lost falls back to the same
|
||||
* "Group: <roomId>" title lookup instead of a fresh "Group: <new name>".
|
||||
* Returns the new name, or null when the target name is taken. */
|
||||
async function renameGroupChat(oldName: string, newName: string, members: GroupMember[] | null | undefined) {
|
||||
export async function renameGroupChat(oldName: string, newName: string, members: GroupMember[] | null | undefined) {
|
||||
const next = String(newName || '')
|
||||
.trim()
|
||||
.slice(0, 64)
|
||||
|
||||
@@ -9,6 +9,42 @@ import { $groupChats, groupChatRoomKey } from './group-chat'
|
||||
import { botConnectionRoute, botRosterMeta, resolveBotConnectionRoute } from './routing'
|
||||
import type { BotMeta, GroupChat, GroupMember, RosterRow } from './types'
|
||||
|
||||
/** Follow the authoritative room record for one async operation. Rename moves
|
||||
* the record wholesale (including legacy rooms without a roomId); disband
|
||||
* retires this binding permanently, even if the same display name is reused. */
|
||||
export function followGroupChat(group: string, onRename: (name: string) => void) {
|
||||
let live = !$groupChats.get()[group]?.tombstone
|
||||
|
||||
const dispose = $groupChats.listen((rooms, previous) => {
|
||||
const prior = previous?.[group]
|
||||
|
||||
if (!live || !prior) {
|
||||
return
|
||||
}
|
||||
|
||||
const current = rooms[group]
|
||||
|
||||
if (current && !current.tombstone && current.roomId === prior.roomId) {
|
||||
return
|
||||
}
|
||||
|
||||
const moved = Object.entries(rooms).find(([, room]) =>
|
||||
!room.tombstone && (prior.roomId ? room.roomId === prior.roomId : room === prior)
|
||||
)
|
||||
|
||||
if (!moved) {
|
||||
live = false
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
group = moved[0]
|
||||
onRename(group)
|
||||
})
|
||||
|
||||
return { dispose, isLive: () => live }
|
||||
}
|
||||
|
||||
export function groupWorkspaceOwnerKey(group: string) {
|
||||
return `group:${groupChatRoomKey(group, $groupChats.get()[group])}`
|
||||
}
|
||||
|
||||
@@ -21,7 +21,7 @@ import {
|
||||
updateGroupChat
|
||||
} from './group-chat'
|
||||
import type { GroupChatRoom, GroupHoldStamp } from './group-chat'
|
||||
import { durableGroupChatMembers, groupMemberKey } from './group-membership'
|
||||
import { durableGroupChatMembers, followGroupChat, groupMemberKey } from './group-membership'
|
||||
import { harvestStrandedGroupReply, isGroupPassText, runGroupChatMemberTurn } from './group-turns'
|
||||
import { requestForBot } from './routing'
|
||||
import type { Attachment, GroupMember, GroupMessage } from './types'
|
||||
@@ -493,8 +493,9 @@ export async function stopGroupThread(group: string, thread: null | string, memb
|
||||
* Watermarks are per thread+member (`${thread}::${memberKey}`), so parallel
|
||||
* topics never eat each other's deltas. */
|
||||
export async function runGroupChatRounds(group: string, members: GroupMember[], thread: string) {
|
||||
const binding = followGroupChat(group, name => { group = name })
|
||||
const startEpoch = ($groupChats.get()[group] || {}).epoch || 0
|
||||
const isCurrent = () => (($groupChats.get()[group] || {}).epoch || 0) === startEpoch
|
||||
const isCurrent = () => binding.isLive() && (($groupChats.get()[group] || {}).epoch || 0) === startEpoch
|
||||
let posted = 0
|
||||
let continuations = 0
|
||||
// #94478: how this drive ended. 'settled' means quiet consensus (everyone
|
||||
@@ -519,6 +520,10 @@ export async function runGroupChatRounds(group: string, members: GroupMember[],
|
||||
}
|
||||
|
||||
await harvestStrandedGroupReply(group, member)
|
||||
|
||||
if (!binding.isLive()) {
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
const roomLog = (($groupChats.get()[group] || {}).log || []).filter(
|
||||
@@ -676,6 +681,10 @@ export async function runGroupChatRounds(group: string, members: GroupMember[],
|
||||
// during-turn tail is anchored by entry id, not index — the history
|
||||
// trim drops entries from the FRONT, so an index slice could
|
||||
// overshoot after a mid-turn trim and silently commit a stale turn.
|
||||
if (!binding.isLive()) {
|
||||
return
|
||||
}
|
||||
|
||||
const roomNow = $groupChats.get()[group] || {
|
||||
log: []
|
||||
}
|
||||
@@ -911,6 +920,8 @@ export async function runGroupChatRounds(group: string, members: GroupMember[],
|
||||
void harvestStrandedUntilSettled(group, members, thread)
|
||||
}
|
||||
}
|
||||
|
||||
binding.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -919,39 +930,53 @@ export async function runGroupChatRounds(group: string, members: GroupMember[],
|
||||
* stranded, a new loop takes the room over (it harvests on its own), or the
|
||||
* room record disappears (disband). */
|
||||
async function harvestStrandedUntilSettled(group: string, members: GroupMember[], thread: string) {
|
||||
const HARVEST_INTERVAL_MS = 5000
|
||||
const HARVEST_MAX_TRIES = 60
|
||||
const binding = followGroupChat(group, name => { group = name })
|
||||
|
||||
for (let attempt = 0; attempt < HARVEST_MAX_TRIES; attempt++) {
|
||||
await new Promise(resolve => window.setTimeout(resolve, HARVEST_INTERVAL_MS))
|
||||
const room = $groupChats.get()[group]
|
||||
try {
|
||||
const HARVEST_INTERVAL_MS = 5000
|
||||
const HARVEST_MAX_TRIES = 60
|
||||
|
||||
if (!room || room.running) {
|
||||
return
|
||||
}
|
||||
for (let attempt = 0; attempt < HARVEST_MAX_TRIES; attempt++) {
|
||||
await new Promise(resolve => window.setTimeout(resolve, HARVEST_INTERVAL_MS))
|
||||
const room = $groupChats.get()[group]
|
||||
|
||||
const stranded = room.stranded || {}
|
||||
if (!binding.isLive() || !room || room.running) {
|
||||
return
|
||||
}
|
||||
|
||||
if (!Object.keys(stranded).length) {
|
||||
return
|
||||
}
|
||||
const stranded = room.stranded || {}
|
||||
|
||||
for (const member of members) {
|
||||
if (Object.prototype.hasOwnProperty.call(stranded, groupMemberKey(member))) {
|
||||
try {
|
||||
await harvestStrandedGroupReply(group, member)
|
||||
} catch {
|
||||
// Best-effort: the next tick retries; the bound stops runaways.
|
||||
if (!Object.keys(stranded).length) {
|
||||
return
|
||||
}
|
||||
|
||||
for (const member of members) {
|
||||
if (!binding.isLive()) {
|
||||
return
|
||||
}
|
||||
|
||||
if (Object.prototype.hasOwnProperty.call(stranded, groupMemberKey(member))) {
|
||||
try {
|
||||
await harvestStrandedGroupReply(group, member)
|
||||
} catch {
|
||||
// Best-effort: the next tick retries; the bound stops runaways.
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
recordGroupActivity(group, {
|
||||
kind: 'failed',
|
||||
member: null,
|
||||
thread
|
||||
})
|
||||
if (!binding.isLive()) {
|
||||
return
|
||||
}
|
||||
|
||||
recordGroupActivity(group, {
|
||||
kind: 'failed',
|
||||
member: null,
|
||||
thread
|
||||
})
|
||||
} finally {
|
||||
binding.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
/** User send into a group room. `thread` continues that thread (its reply
|
||||
@@ -1026,26 +1051,31 @@ export function sendToGroupChat(
|
||||
thread: target
|
||||
})
|
||||
|
||||
if (!wasRunning) {
|
||||
void runGroupChatRounds(group, members, target).catch(() => {
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.running = false
|
||||
const binding = followGroupChat(group, name => { group = name })
|
||||
|
||||
return r
|
||||
})
|
||||
})
|
||||
} else {
|
||||
// A loop is live; it bails at its next boundary. Chain the fresh loop
|
||||
// after a short settle so exactly one drive owns the room.
|
||||
setTimeout(() => {
|
||||
void runGroupChatRounds(group, members, target).catch(() => {
|
||||
const drive = () => {
|
||||
if (!binding.isLive()) {
|
||||
binding.dispose()
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
void runGroupChatRounds(group, members, target).catch(() => {
|
||||
if (binding.isLive()) {
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.running = false
|
||||
|
||||
return r
|
||||
})
|
||||
})
|
||||
}, 250)
|
||||
}
|
||||
}).finally(binding.dispose)
|
||||
}
|
||||
|
||||
if (!wasRunning) {
|
||||
drive()
|
||||
} else {
|
||||
// Preserve the existing newer-send handoff delay, without pinning its name.
|
||||
setTimeout(drive, 250)
|
||||
}
|
||||
|
||||
return target
|
||||
|
||||
@@ -485,104 +485,194 @@ describe('clarify and approvals (#90694)', () => {
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Other')).toBe(true)
|
||||
})
|
||||
|
||||
it('keeps the derived badge lit until every clarify in the room is answered', async () => {
|
||||
const room = await loadRoom()
|
||||
const research: GroupMember = { name: 'research', title: '' }
|
||||
const ops: GroupMember = { name: 'ops', title: '' }
|
||||
|
||||
room.turns.syncGroupClarify('Core', research, { pending_clarify: CLARIFY })
|
||||
room.turns.syncGroupClarify('Core', ops, { pending_clarify: { ...CLARIFY, request_id: 'req-2' } })
|
||||
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Core')).toBe(true)
|
||||
|
||||
const [first, second] = Object.values(room.chat.$groupClarify.get())
|
||||
await room.turns.answerGroupClarify(first, research, 'staging')
|
||||
|
||||
// One of two questions answered — the room still owes the user a reply.
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Core')).toBe(true)
|
||||
|
||||
await room.turns.answerGroupClarify(second, ops, 'staging')
|
||||
|
||||
// Last pending clarify answered — nothing left to badge.
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Core')).toBe(false)
|
||||
it('keeps pending prompts independent from mention attention through their lifecycle', async () => {
|
||||
const { chat, turns } = await loadRoom()
|
||||
const member = { name: 'research', title: '' }
|
||||
turns.syncGroupClarify('Core', member, { pending_clarify: CLARIFY })
|
||||
expect(chat.$groupNeedsYou.get().Core).toBeFalsy()
|
||||
turns.syncGroupClarify('Core', { name: 'ops' }, { pending_approval: APPROVAL })
|
||||
expect(Object.values(chat.$groupClarify.get())).toHaveLength(2)
|
||||
turns.syncGroupClarify('Core', member, {})
|
||||
expect(turns.groupHasPendingClarify(chat.$groupClarify.get(), 'Core')).toBe(true)
|
||||
turns.syncGroupClarify('Core', { name: 'ops' }, {})
|
||||
expect(turns.groupHasPendingClarify(chat.$groupClarify.get(), 'Core')).toBe(false)
|
||||
chat.appendGroupChatEntry('Core', { kind: 'member', name: 'research' }, '@user please review')
|
||||
turns.syncGroupClarify('Core', member, { pending_approval: APPROVAL })
|
||||
await turns.answerGroupClarify(Object.values(chat.$groupClarify.get())[0], member, 'deny')
|
||||
expect(Object.values(chat.$groupClarify.get())).toHaveLength(0)
|
||||
expect(chat.$groupNeedsYou.get().Core).toBe(true)
|
||||
})
|
||||
|
||||
// $groupNeedsYou (mention attention) and
|
||||
// $groupClarify (pending clarify/approval attention) are independent
|
||||
// sources — resolving one must never clear the other. The roster reads
|
||||
// BOTH: groupNeedsYou[group] || groupHasPendingClarify(clarifies, group).
|
||||
it('resolving the last clarify preserves independent @user mention attention', async () => {
|
||||
const room = await loadRoom()
|
||||
const research: GroupMember = { name: 'research', title: '' }
|
||||
it('keeps late prompt snapshots on the live room and never revives a disbanded room', async ({ onTestFinished }) => {
|
||||
for (const roomId of ['stable-room', undefined]) {
|
||||
for (const disband of [false, true]) {
|
||||
const room = await loadRoom({ turn: ({ n }) => n === 1 ? 'Completed reply' : '(pass)' })
|
||||
const view = await import('./group-chat-view')
|
||||
const member = { name: 'research', title: '' }
|
||||
room.chat.updateGroupChat('Core', current => ({
|
||||
...current, roomId, running: true, epoch: 1,
|
||||
log: [{ id: 'input', at: 1, from: { kind: 'user', name: 'You' }, text: '@research check', thread: 'thread' }]
|
||||
}))
|
||||
let entered!: () => void
|
||||
let release!: () => void
|
||||
const polled = new Promise<void>(resolve => { entered = resolve })
|
||||
const held = new Promise<void>(resolve => { release = resolve })
|
||||
const original = host.request as (method: string, params: Record<string, unknown>) => Promise<any>
|
||||
let submitted = false
|
||||
let answered = false
|
||||
let polls = 0
|
||||
|
||||
// A member addressing @user badges the room — independent of clarify.
|
||||
room.chat.appendGroupChatEntry('Core', { kind: 'member', name: 'research' }, '@user deployment needs a call')
|
||||
expect(room.chat.$groupNeedsYou.get().Core).toBe(true)
|
||||
host.request = async (method: string, params: Record<string, unknown>) => {
|
||||
const result = await original(method, params)
|
||||
|
||||
room.turns.syncGroupClarify('Core', research, { pending_clarify: CLARIFY })
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Core')).toBe(true)
|
||||
if (method === 'prompt.submit') {submitted = true}
|
||||
|
||||
await room.turns.answerGroupClarify(Object.values(room.chat.$groupClarify.get())[0], research, 'staging')
|
||||
if (method === 'clarify.respond') {answered = true}
|
||||
|
||||
// Clarify resolved — but the unrelated mention attention must survive.
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Core')).toBe(false)
|
||||
expect(room.chat.$groupNeedsYou.get().Core).toBe(true)
|
||||
})
|
||||
if (method === 'session.resume' && submitted && !answered) {
|
||||
if (++polls === 1) { entered(); await held }
|
||||
|
||||
it('preserves pending clarify attention across a room rename', async () => {
|
||||
const room = await loadRoom()
|
||||
const research: GroupMember = { name: 'research', title: '' }
|
||||
return { ...result, pending_clarify: CLARIFY }
|
||||
}
|
||||
|
||||
// Core has a pending clarify and no @user mention — attention is
|
||||
// entirely clarify-derived, the case a plain group-key swap loses.
|
||||
room.turns.syncGroupClarify('Core', research, { pending_clarify: CLARIFY })
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Core')).toBe(true)
|
||||
expect(room.chat.$groupNeedsYou.get().Core).toBeFalsy()
|
||||
return result
|
||||
}
|
||||
|
||||
room.turns.renameGroupClarify('Core', 'Renamed')
|
||||
const drive = room.rounds.runGroupChatRounds('Core', [member], 'thread')
|
||||
await polled
|
||||
|
||||
// The prompt now belongs to the new name — not the old one.
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Renamed')).toBe(true)
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Core')).toBe(false)
|
||||
const [prompt] = Object.values(room.chat.$groupClarify.get())
|
||||
expect(prompt.group).toBe('Renamed')
|
||||
if (disband) {
|
||||
await view.disbandGroupChat('Core', [])
|
||||
release()
|
||||
await drive
|
||||
expect(Object.values(room.chat.$groupClarify.get())).toHaveLength(0)
|
||||
expect(room.chat.$groupChats.get().Core === undefined || room.chat.$groupChats.get().Core.tombstone).toBe(true)
|
||||
expect(room.chat.$groupChats.get().Core?.log || []).toHaveLength(0)
|
||||
} else {
|
||||
await view.renameGroupChat('Core', 'Renamed', [])
|
||||
|
||||
// Resolving it under the new name clears attention there.
|
||||
await room.turns.answerGroupClarify(prompt, research, 'staging')
|
||||
expect(room.turns.groupHasPendingClarify(room.chat.$groupClarify.get(), 'Renamed')).toBe(false)
|
||||
})
|
||||
const mirrored = new Promise<void>(resolve => {
|
||||
const stop = room.chat.$groupClarify.listen(entries => {
|
||||
if (Object.keys(entries).length) { stop(); resolve() }
|
||||
})
|
||||
})
|
||||
|
||||
it('lets the renamed room\'s current prompt replace a stale mirror already at the destination', async () => {
|
||||
const room = await loadRoom()
|
||||
|
||||
// Current prompt inserted FIRST, stale destination mirror SECOND — a
|
||||
// single-pass rekey that iterates Object.entries in insertion order
|
||||
// would let the stale entry (visited last) clobber the just-migrated
|
||||
// current one. Two passes must let the live room win regardless of
|
||||
// insertion order.
|
||||
room.turns.syncGroupClarify('Core', { name: 'research', title: '' }, { pending_clarify: CLARIFY })
|
||||
room.chat.$groupClarify.set({
|
||||
...room.chat.$groupClarify.get(),
|
||||
'Renamed::research': {
|
||||
at: Date.now(),
|
||||
choices: [],
|
||||
group: 'Renamed',
|
||||
kind: 'clarify',
|
||||
member: 'research',
|
||||
memberKey: 'research',
|
||||
multiSelect: false,
|
||||
question: 'stale question',
|
||||
questions: null,
|
||||
requestId: 'req-stale',
|
||||
sessionId: null
|
||||
release()
|
||||
await mirrored
|
||||
const [prompt] = Object.values(room.chat.$groupClarify.get())
|
||||
const correctRoom = prompt.group
|
||||
await room.turns.answerGroupClarify(prompt, member, 'staging')
|
||||
await drive
|
||||
expect(correctRoom).toBe('Renamed')
|
||||
expect(Object.keys(room.chat.$groupChats.get())).toEqual(['Renamed'])
|
||||
expect(room.chat.$groupChats.get().Renamed.running).toBe(false)
|
||||
expect(room.chat.$groupChats.get().Renamed.log.filter(entry => entry.from.kind === 'member').map(entry => entry.text)).toEqual(['Completed reply'])
|
||||
expect(Object.values(room.chat.$groupClarify.get())).toHaveLength(0)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
room.turns.renameGroupClarify('Core', 'Renamed')
|
||||
const clock = vi.spyOn(Date, 'now').mockReturnValue(0)
|
||||
onTestFinished(() => clock.mockRestore())
|
||||
|
||||
const migrated = room.chat.$groupClarify.get()['Renamed::research']
|
||||
expect(migrated.requestId).toBe(CLARIFY.request_id)
|
||||
expect(Object.keys(room.chat.$groupClarify.get())).toEqual(['Renamed::research'])
|
||||
for (const recreate of [false, true]) {
|
||||
clock.mockReturnValue(0)
|
||||
const room = await loadRoom()
|
||||
const view = await import('./group-chat-view')
|
||||
const member = { name: 'research', title: '' }
|
||||
room.chat.updateGroupChat('Core', current => ({ ...current, roomId: 'retired-room' }))
|
||||
let entered!: () => void
|
||||
let release!: () => void
|
||||
const polled = new Promise<void>(resolve => { entered = resolve })
|
||||
const held = new Promise<void>(resolve => { release = resolve })
|
||||
const original = host.request as (method: string, params: Record<string, unknown>) => Promise<any>
|
||||
let submitted = false
|
||||
|
||||
host.request = async (method: string, params: Record<string, unknown>) => {
|
||||
if (method === 'session.resume' && submitted) {
|
||||
entered()
|
||||
await held
|
||||
throw new Error('poll rejected after deadline')
|
||||
}
|
||||
|
||||
const result = await original(method, params)
|
||||
|
||||
if (method === 'prompt.submit') { submitted = true }
|
||||
|
||||
return result
|
||||
}
|
||||
|
||||
const turn = room.turns.runGroupChatMemberTurn('Core', member, 'check', 'thread', [])
|
||||
await polled
|
||||
await view.disbandGroupChat('Core', [])
|
||||
|
||||
if (recreate) {
|
||||
room.chat.updateGroupChat('Core', current => ({ ...current, roomId: 'replacement-room' }))
|
||||
room.turns.syncGroupClarify('Core', member, { pending_clarify: CLARIFY })
|
||||
}
|
||||
|
||||
const roomsBefore = structuredClone(room.chat.$groupChats.get())
|
||||
const promptsBefore = structuredClone(room.chat.$groupClarify.get())
|
||||
// Cross even the hard cap while the rejected poll is still in flight.
|
||||
clock.mockReturnValue(24 * 60 * 60 * 1000)
|
||||
release()
|
||||
expect(await turn).toBeNull()
|
||||
expect(room.chat.$groupChats.get()).toEqual(roomsBefore)
|
||||
expect(room.chat.$groupClarify.get()).toEqual(promptsBefore)
|
||||
}
|
||||
|
||||
// Retirement during one background harvest must fence the next member too.
|
||||
const room = await loadRoom()
|
||||
const view = await import('./group-chat-view')
|
||||
const members = [{ name: 'research' }, { name: 'ops' }]
|
||||
room.chat.updateGroupChat('Core', current => ({
|
||||
...current, roomId: 'old-harvest-room', running: true,
|
||||
stranded: { research: 0, ops: 0 }
|
||||
}))
|
||||
let tick!: () => void
|
||||
const previousWindow = globalThis.window
|
||||
vi.stubGlobal('window', { setTimeout: (callback: () => void) => { tick = callback;
|
||||
|
||||
return 0 } })
|
||||
onTestFinished(() => { vi.stubGlobal('window', previousWindow) })
|
||||
let entered!: () => void
|
||||
let release!: () => void
|
||||
const polled = new Promise<void>(resolve => { entered = resolve })
|
||||
const held = new Promise<void>(resolve => { release = resolve })
|
||||
let background = false
|
||||
const backgroundProfiles: unknown[] = []
|
||||
const original = host.request as (method: string, params: Record<string, unknown>) => Promise<any>
|
||||
|
||||
host.request = async (method: string, params: Record<string, unknown>) => {
|
||||
if (method !== 'session.resume') { return original(method, params) }
|
||||
|
||||
if (background) {
|
||||
backgroundProfiles.push(params.profile)
|
||||
|
||||
if (params.profile === 'research') { entered(); await held }
|
||||
}
|
||||
|
||||
return { running: true, pending_clarify: CLARIFY }
|
||||
}
|
||||
|
||||
await room.rounds.runGroupChatRounds('Core', members, 'thread')
|
||||
background = true
|
||||
tick()
|
||||
await polled
|
||||
await view.disbandGroupChat('Core', [])
|
||||
const { setImmediate } = await import('node:timers/promises')
|
||||
await setImmediate()
|
||||
room.chat.updateGroupChat('Core', current => ({
|
||||
...current, roomId: 'new-harvest-room', stranded: { research: 0, ops: 0 }
|
||||
}))
|
||||
const roomsBefore = structuredClone(room.chat.$groupChats.get())
|
||||
const promptsBefore = structuredClone(room.chat.$groupClarify.get())
|
||||
release()
|
||||
// Let the released RPC and its background caller finish their microtasks.
|
||||
await setImmediate()
|
||||
expect(backgroundProfiles).toEqual(['research'])
|
||||
expect(room.chat.$groupChats.get()).toEqual(roomsBefore)
|
||||
expect(room.chat.$groupClarify.get()).toEqual(promptsBefore)
|
||||
})
|
||||
|
||||
it('holds the turn open on a command approval too', async () => {
|
||||
|
||||
@@ -11,7 +11,7 @@ import { host } from '@hermes/plugin-sdk'
|
||||
import { recordGroupActivity } from './group-activity'
|
||||
import { $groupChats, $groupClarify, appendGroupChatEntry, updateGroupChat } from './group-chat'
|
||||
import type { GroupChatRoom } from './group-chat'
|
||||
import { groupMemberKey, groupSessionOwner } from './group-membership'
|
||||
import { followGroupChat, groupMemberKey, groupSessionOwner } from './group-membership'
|
||||
import { botConnectionRoute, requestForBot } from './routing'
|
||||
import type { Attachment, GroupMember, GroupPrompt, GroupPromptQuestion, ProfileRoute } from './types'
|
||||
|
||||
@@ -123,110 +123,128 @@ interface GroupMemberSessionHandle {
|
||||
* it after restarts. Cross-connection members route to their OWN source
|
||||
* via requestForBot; the window's gateway never switches. */
|
||||
export async function ensureGroupChatSession(group: string, member: GroupMember): Promise<GroupMemberSessionHandle> {
|
||||
const room = $groupChats.get()[group] || {}
|
||||
// New rooms title member sessions by their immutable roomId so a
|
||||
// same-name recreate never resumes the old room's sessions by title;
|
||||
// legacy rooms without a roomId fall back to the display name.
|
||||
const title = `Group: ${room.roomId || group}`
|
||||
const key = groupMemberKey(member)
|
||||
const known = room.sessions && room.sessions[key]
|
||||
const binding = followGroupChat(group, name => { group = name })
|
||||
|
||||
// Try resuming what we know (stored sid first, then title lookup).
|
||||
//
|
||||
// FAIL CLOSED on a transient lookup failure — mirrors the sibling fix in
|
||||
// findExistingCanonicalChat (87b645f52c). session.resume signals "this
|
||||
// target genuinely doesn't exist" with JSON-RPC code 4007; every other
|
||||
// failure (network blip, the backend still warming up after a restart,
|
||||
// an oversized-resume refusal) means the real session might still be
|
||||
// there and must not be read as "no session, mint a new one" — that
|
||||
// forks the member's real history, and the fork silently overwrites
|
||||
// room.sessions[key] so the old session becomes unreachable from the
|
||||
// room. Only a genuine 4007 on BOTH targets means there truly is nothing
|
||||
// to resume yet, so the loop falls through to session.create below.
|
||||
for (const target of [known, title]) {
|
||||
if (!target || target === true) {
|
||||
continue
|
||||
}
|
||||
try {
|
||||
const room = $groupChats.get()[group] || {}
|
||||
// New rooms title member sessions by their immutable roomId so a
|
||||
// same-name recreate never resumes the old room's sessions by title;
|
||||
// legacy rooms without a roomId fall back to the display name.
|
||||
const title = `Group: ${room.roomId || group}`
|
||||
const key = groupMemberKey(member)
|
||||
const known = room.sessions && room.sessions[key]
|
||||
|
||||
try {
|
||||
const res = (await requestForBot(member, 'session.resume', {
|
||||
session_id: target,
|
||||
profile: member.name,
|
||||
omit_messages: true
|
||||
})) as GroupSessionSnapshot
|
||||
// Try resuming what we know (stored sid first, then title lookup).
|
||||
//
|
||||
// FAIL CLOSED on a transient lookup failure — mirrors the sibling fix in
|
||||
// findExistingCanonicalChat (87b645f52c). session.resume signals "this
|
||||
// target genuinely doesn't exist" with JSON-RPC code 4007; every other
|
||||
// failure (network blip, the backend still warming up after a restart,
|
||||
// an oversized-resume refusal) means the real session might still be
|
||||
// there and must not be read as "no session, mint a new one" — that
|
||||
// forks the member's real history, and the fork silently overwrites
|
||||
// room.sessions[key] so the old session becomes unreachable from the
|
||||
// room. Only a genuine 4007 on BOTH targets means there truly is nothing
|
||||
// to resume yet, so the loop falls through to session.create below.
|
||||
for (const target of [known, title]) {
|
||||
if (!target || target === true) {
|
||||
continue
|
||||
}
|
||||
|
||||
if (res?.session_id) {
|
||||
// TODO(bot-mode-types): `known` is `room.sessions[key]`, which the
|
||||
// domain model types `string | true` — and the `target === true` skip
|
||||
// above shows the legacy `true` sentinel is expected here. A backend
|
||||
// that answers the title resume without a `session_key` therefore
|
||||
// stores `true` back into room.sessions and hands `true` on as the
|
||||
// durable id, which later rides into `session_id` on the recovery
|
||||
// resume and on session.interrupt. Typed as-written.
|
||||
const stored = res.session_key || known
|
||||
try {
|
||||
const res = (await requestForBot(member, 'session.resume', {
|
||||
session_id: target,
|
||||
profile: member.name,
|
||||
omit_messages: true
|
||||
})) as GroupSessionSnapshot
|
||||
|
||||
if (stored) {
|
||||
updateGroupChat(group, (current: GroupChatRoom) => {
|
||||
current.sessions = {
|
||||
...(current.sessions || {}),
|
||||
[key]: stored
|
||||
}
|
||||
current.sessionOwners = {
|
||||
...(current.sessionOwners || {}),
|
||||
[key]: groupSessionOwner(member)
|
||||
}
|
||||
|
||||
return current
|
||||
})
|
||||
if (!binding.isLive()) {
|
||||
return { runtime: null }
|
||||
}
|
||||
|
||||
return {
|
||||
runtime: res.session_id,
|
||||
stored
|
||||
if (res?.session_id) {
|
||||
// TODO(bot-mode-types): `known` is `room.sessions[key]`, which the
|
||||
// domain model types `string | true` — and the `target === true` skip
|
||||
// above shows the legacy `true` sentinel is expected here. A backend
|
||||
// that answers the title resume without a `session_key` therefore
|
||||
// stores `true` back into room.sessions and hands `true` on as the
|
||||
// durable id, which later rides into `session_id` on the recovery
|
||||
// resume and on session.interrupt. Typed as-written.
|
||||
const stored = res.session_key || known
|
||||
|
||||
if (stored) {
|
||||
updateGroupChat(group, (current: GroupChatRoom) => {
|
||||
current.sessions = {
|
||||
...(current.sessions || {}),
|
||||
[key]: stored
|
||||
}
|
||||
current.sessionOwners = {
|
||||
...(current.sessionOwners || {}),
|
||||
[key]: groupSessionOwner(member)
|
||||
}
|
||||
|
||||
return current
|
||||
})
|
||||
}
|
||||
|
||||
return {
|
||||
runtime: res.session_id,
|
||||
stored
|
||||
}
|
||||
}
|
||||
} catch (error: any) {
|
||||
if (error?.code !== 4007) {
|
||||
const detail = error instanceof Error && error.message ? ` (${error.message})` : ''
|
||||
throw new Error(`Could not check ${member?.name || 'member'}'s group session${detail} — not starting a new one`)
|
||||
}
|
||||
/* genuinely doesn't exist (4007) — try the next target / fall through to create */
|
||||
}
|
||||
} catch (error: any) {
|
||||
if (error?.code !== 4007) {
|
||||
const detail = error instanceof Error && error.message ? ` (${error.message})` : ''
|
||||
throw new Error(`Could not check ${member?.name || 'member'}'s group session${detail} — not starting a new one`)
|
||||
}
|
||||
/* genuinely doesn't exist (4007) — try the next target / fall through to create */
|
||||
}
|
||||
}
|
||||
|
||||
const created = (await requestForBot(member, 'session.create', {
|
||||
profile: member.name,
|
||||
title,
|
||||
// Room member sessions are plumbing — always hidden from the sidebar.
|
||||
hidden: true,
|
||||
// Explicit contracts (PR #97008): room plumbing sessions always rebuild
|
||||
// from the member profile's CURRENT config on resume, never a stale
|
||||
// stored model/provider pin. Older gateways ignore the unknown params;
|
||||
// the server's hidden + "Group: " title fallback then covers legacy.
|
||||
room_plumbing: true,
|
||||
follow_profile_config: true
|
||||
})) as { session_id?: string; stored_session_id?: string }
|
||||
if (!binding.isLive()) {
|
||||
return { runtime: null }
|
||||
}
|
||||
|
||||
const stored = created?.stored_session_id || null
|
||||
const created = (await requestForBot(member, 'session.create', {
|
||||
profile: member.name,
|
||||
title,
|
||||
// Room member sessions are plumbing — always hidden from the sidebar.
|
||||
hidden: true,
|
||||
// Explicit contracts (PR #97008): room plumbing sessions always rebuild
|
||||
// from the member profile's CURRENT config on resume, never a stale
|
||||
// stored model/provider pin. Older gateways ignore the unknown params;
|
||||
// the server's hidden + "Group: " title fallback then covers legacy.
|
||||
room_plumbing: true,
|
||||
follow_profile_config: true
|
||||
})) as { session_id?: string; stored_session_id?: string }
|
||||
|
||||
if (stored) {
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.sessions = {
|
||||
...(r.sessions || {}),
|
||||
[key]: stored
|
||||
}
|
||||
r.sessionOwners = {
|
||||
...(r.sessionOwners || {}),
|
||||
[key]: groupSessionOwner(member)
|
||||
}
|
||||
if (!binding.isLive()) {
|
||||
return { runtime: null }
|
||||
}
|
||||
|
||||
return r
|
||||
})
|
||||
}
|
||||
const stored = created?.stored_session_id || null
|
||||
|
||||
return {
|
||||
runtime: created?.session_id || null,
|
||||
stored
|
||||
if (stored) {
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.sessions = {
|
||||
...(r.sessions || {}),
|
||||
[key]: stored
|
||||
}
|
||||
r.sessionOwners = {
|
||||
...(r.sessionOwners || {}),
|
||||
[key]: groupSessionOwner(member)
|
||||
}
|
||||
|
||||
return r
|
||||
})
|
||||
}
|
||||
|
||||
return {
|
||||
runtime: created?.session_id || null,
|
||||
stored
|
||||
}
|
||||
} finally {
|
||||
binding.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -529,13 +547,8 @@ export function clearGroupClarify(group: string) {
|
||||
}
|
||||
}
|
||||
|
||||
/** Re-key every mirrored clarify belonging to `oldName` onto `newName`
|
||||
* (rename). Unlike disband, the room still exists under its new name — a
|
||||
* pending clarify's attention must follow it, not disappear. Known gap:
|
||||
* an in-flight member turn captured `oldName` as its closure argument
|
||||
* (runGroupChatMemberTurnLeased) and keeps polling under that name, so a
|
||||
* poll that lands after this rename can still re-mirror a NEW entry under
|
||||
* oldName. That race predates this function and isn't fixed here. */
|
||||
/** Move existing prompts with the renamed room; active operations follow
|
||||
* the same room through their scoped followGroupChat binding. */
|
||||
export function renameGroupClarify(oldName: string, newName: string) {
|
||||
const all = $groupClarify.get()
|
||||
const next: Record<string, GroupPrompt> = {}
|
||||
@@ -583,40 +596,51 @@ export async function answerGroupClarify(
|
||||
member: GroupMember,
|
||||
answers: Record<string, string> | string | undefined
|
||||
) {
|
||||
if (entry.kind === 'approval') {
|
||||
await requestForBot(member, 'approval.respond', {
|
||||
session_id: entry.sessionId || undefined,
|
||||
request_id: entry.requestId,
|
||||
choice: typeof answers === 'string' && answers ? answers : 'deny'
|
||||
})
|
||||
} else if (entry.questions && entry.questions.length) {
|
||||
for (const question of entry.questions) {
|
||||
// Question ids are opaque on the wire (`GroupPrompt.questions` types
|
||||
// them `unknown`); the batch card keys its answer bag by exactly them.
|
||||
const qid = (question?.qid ?? question?.id) as string
|
||||
let group = entry.group
|
||||
const binding = followGroupChat(group, name => { group = name })
|
||||
|
||||
try {
|
||||
if (entry.kind === 'approval') {
|
||||
await requestForBot(member, 'approval.respond', {
|
||||
session_id: entry.sessionId || undefined,
|
||||
request_id: entry.requestId,
|
||||
choice: typeof answers === 'string' && answers ? answers : 'deny'
|
||||
})
|
||||
} else if (entry.questions && entry.questions.length) {
|
||||
for (const question of entry.questions) {
|
||||
// Question ids are opaque on the wire (`GroupPrompt.questions` types
|
||||
// them `unknown`); the batch card keys its answer bag by exactly them.
|
||||
const qid = (question?.qid ?? question?.id) as string
|
||||
await requestForBot(member, 'clarify.respond', {
|
||||
request_id: entry.requestId,
|
||||
question_id: qid,
|
||||
answer: (answers as Record<string, string>)?.[qid] ?? ''
|
||||
})
|
||||
}
|
||||
} else {
|
||||
await requestForBot(member, 'clarify.respond', {
|
||||
request_id: entry.requestId,
|
||||
question_id: qid,
|
||||
answer: (answers as Record<string, string>)?.[qid] ?? ''
|
||||
answer: typeof answers === 'string' ? answers : ''
|
||||
})
|
||||
}
|
||||
} else {
|
||||
await requestForBot(member, 'clarify.respond', {
|
||||
request_id: entry.requestId,
|
||||
answer: typeof answers === 'string' ? answers : ''
|
||||
})
|
||||
}
|
||||
|
||||
const all = $groupClarify.get()
|
||||
const key = `${entry.group}::${entry.memberKey}`
|
||||
|
||||
if (all[key]?.requestId === entry.requestId) {
|
||||
const next = {
|
||||
...all
|
||||
if (!binding.isLive()) {
|
||||
return
|
||||
}
|
||||
|
||||
delete next[key]
|
||||
$groupClarify.set(next)
|
||||
const all = $groupClarify.get()
|
||||
const key = `${group}::${entry.memberKey}`
|
||||
|
||||
if (all[key]?.requestId === entry.requestId) {
|
||||
const next = {
|
||||
...all
|
||||
}
|
||||
|
||||
delete next[key]
|
||||
$groupClarify.set(next)
|
||||
}
|
||||
} finally {
|
||||
binding.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -638,12 +662,16 @@ export async function runGroupChatMemberTurn(
|
||||
// lease, every RPC below rides its own request-scoped socket lease; the
|
||||
// socket that minted `runtime` can close between RPCs, the gateway reaps
|
||||
// the runtime session, and prompt.submit dies 4001 — the bot goes silent.
|
||||
const releaseTurnLease = await retainGroupTurnRoute(member)
|
||||
const binding = followGroupChat(group, name => { group = name })
|
||||
let releaseTurnLease: (() => void) | undefined
|
||||
|
||||
try {
|
||||
return await runGroupChatMemberTurnLeased(group, member, prompt, thread, images)
|
||||
releaseTurnLease = await retainGroupTurnRoute(member)
|
||||
|
||||
return binding.isLive() ? await runGroupChatMemberTurnLeased(group, member, prompt, thread, images) : null
|
||||
} finally {
|
||||
releaseTurnLease()
|
||||
releaseTurnLease?.()
|
||||
binding.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -654,281 +682,315 @@ async function runGroupChatMemberTurnLeased(
|
||||
thread: string,
|
||||
images?: Attachment[]
|
||||
): Promise<null | string> {
|
||||
const { runtime, stored } = await ensureGroupChatSession(group, member)
|
||||
|
||||
if (!runtime) {
|
||||
return null
|
||||
}
|
||||
|
||||
// #91868/#94569: remember the epoch this turn was dispatched under so the
|
||||
// poll loop below can tell an explicit stop from ordinary room churn.
|
||||
const dispatchEpoch = ($groupChats.get()[group] || {}).epoch || 0
|
||||
const memberKey = groupMemberKey(member)
|
||||
recordGroupActivity(group, {
|
||||
kind: 'working',
|
||||
member: member.name,
|
||||
thread
|
||||
})
|
||||
|
||||
// Baseline: how many messages exist before our submit.
|
||||
let before = 0
|
||||
// Every runtime id this turn has seen for the member's session. Terminal
|
||||
// frames are keyed by runtime id, and a resume can hand back a fresh one.
|
||||
const runtimeIds = new Set<string>([runtime])
|
||||
const binding = followGroupChat(group, name => { group = name })
|
||||
|
||||
try {
|
||||
const pre = (await requestForBot(member, 'session.resume', {
|
||||
session_id: stored || runtime,
|
||||
profile: member.name
|
||||
})) as GroupSessionSnapshot
|
||||
const { runtime, stored } = await ensureGroupChatSession(group, member)
|
||||
|
||||
before = Array.isArray(pre?.messages) ? pre.messages.length : pre?.message_count || 0
|
||||
|
||||
if (pre?.session_id) {
|
||||
runtimeIds.add(pre.session_id)
|
||||
}
|
||||
} catch {
|
||||
/* lazy session — zero messages */
|
||||
}
|
||||
|
||||
// Stage this delta's attachments into the member's session so the model
|
||||
// receives the actual payload with the prompt — the same attach RPCs the
|
||||
// 1:1 chat uses (they also work cross-connection, where the member's
|
||||
// gateway can't see this machine's files). Images queue as vision tiles,
|
||||
// PDFs render per-page via pdf.attach, and other files materialize in the
|
||||
// session workspace (their @file: refs are appended to the prompt so the
|
||||
// member's file tools can read them). A failed attach degrades that
|
||||
// member to text-only; the transcript line still names the attachment so
|
||||
// the member knows something was shared.
|
||||
const fileRefs: string[] = []
|
||||
|
||||
for (const img of Array.isArray(images) ? images : []) {
|
||||
if (!img || typeof img.data !== 'string' || !img.data) {
|
||||
continue
|
||||
}
|
||||
|
||||
try {
|
||||
if (img.kind === 'pdf') {
|
||||
await requestForBot(member, 'pdf.attach', {
|
||||
session_id: runtime,
|
||||
content_base64: img.data,
|
||||
filename: img.name || 'attachment.pdf'
|
||||
})
|
||||
} else if (img.kind === 'file') {
|
||||
const res = (await requestForBot(member, 'file.attach', {
|
||||
session_id: runtime,
|
||||
data_url: img.data,
|
||||
name: img.name || 'attachment'
|
||||
})) as { ref_text?: string }
|
||||
|
||||
if (res?.ref_text) {
|
||||
fileRefs.push(`${img.name || 'attachment'} → ${res.ref_text}`)
|
||||
}
|
||||
} else {
|
||||
await requestForBot(member, 'image.attach_bytes', {
|
||||
session_id: runtime,
|
||||
content_base64: img.data,
|
||||
filename: img.name || 'attachment.png'
|
||||
})
|
||||
}
|
||||
} catch {
|
||||
/* text-only fallback for this member */
|
||||
}
|
||||
}
|
||||
|
||||
const turnText = fileRefs.length
|
||||
? `${prompt}\n\nAttached files staged in your session workspace:\n${fileRefs.join('\n')}`
|
||||
: prompt
|
||||
|
||||
// #93602: one-shot recovery when the runtime session was reaped between
|
||||
// minting and submitting. Tracks the runtime id the submit landed on so
|
||||
// the poll fallback below targets a live session.
|
||||
const liveRuntime = await submitGroupTurnPrompt(member, runtime, stored, turnText)
|
||||
runtimeIds.add(liveRuntime)
|
||||
const started = Date.now()
|
||||
let deadline = started + GROUP_TURN_TIMEOUT_MS
|
||||
// After the terminal frame fires, the gateway still has to flip
|
||||
// session.running off in its turn `finally` — re-check quickly for a few
|
||||
// beats instead of falling back to the slow backstop cadence.
|
||||
let quickRechecks = 0
|
||||
|
||||
while (Date.now() < deadline) {
|
||||
const signalled = quickRechecks
|
||||
? await waitForTurnSignal([], GROUP_TURN_SETTLE_RECHECK_MS)
|
||||
: await waitForTurnSignal([...runtimeIds], GROUP_TURN_POLL_MS)
|
||||
|
||||
quickRechecks = signalled ? GROUP_TURN_SETTLE_RECHECKS : Math.max(0, quickRechecks - 1)
|
||||
|
||||
// #91868/#94569: an explicit stop (stopGroupThread) bumped the epoch AND
|
||||
// held this member — the member's session was interrupted, so nothing is
|
||||
// coming; abandon the poll instead of grinding until the deadline. Both
|
||||
// conditions on purpose: an ordinary newer send bumps the epoch WITHOUT
|
||||
// a hold, and that turn must keep polling so finished work can still be
|
||||
// delivered (the #93127 commit check decides its fate, not this loop).
|
||||
const roomDuringPoll = $groupChats.get()[group] || {}
|
||||
|
||||
if ((roomDuringPoll.epoch || 0) !== dispatchEpoch && (roomDuringPoll.holds || {})[memberKey]) {
|
||||
if (!runtime || !binding.isLive()) {
|
||||
return null
|
||||
}
|
||||
|
||||
let state: GroupSessionSnapshot | null = null
|
||||
// #91868/#94569: remember the epoch this turn was dispatched under so the
|
||||
// poll loop below can tell an explicit stop from ordinary room churn.
|
||||
const dispatchEpoch = ($groupChats.get()[group] || {}).epoch || 0
|
||||
const memberKey = groupMemberKey(member)
|
||||
recordGroupActivity(group, {
|
||||
kind: 'working',
|
||||
member: member.name,
|
||||
thread
|
||||
})
|
||||
|
||||
// Baseline: how many messages exist before our submit.
|
||||
let before = 0
|
||||
// Every runtime id this turn has seen for the member's session. Terminal
|
||||
// frames are keyed by runtime id, and a resume can hand back a fresh one.
|
||||
const runtimeIds = new Set<string>([runtime])
|
||||
|
||||
try {
|
||||
state = (await requestForBot(member, 'session.resume', {
|
||||
session_id: stored || liveRuntime,
|
||||
const pre = (await requestForBot(member, 'session.resume', {
|
||||
session_id: stored || runtime,
|
||||
profile: member.name
|
||||
})) as GroupSessionSnapshot
|
||||
|
||||
before = Array.isArray(pre?.messages) ? pre.messages.length : pre?.message_count || 0
|
||||
|
||||
if (pre?.session_id) {
|
||||
runtimeIds.add(pre.session_id)
|
||||
}
|
||||
} catch {
|
||||
continue
|
||||
/* lazy session — zero messages */
|
||||
}
|
||||
|
||||
if (state?.session_id) {
|
||||
runtimeIds.add(state.session_id)
|
||||
// Stage this delta's attachments into the member's session so the model
|
||||
// receives the actual payload with the prompt — the same attach RPCs the
|
||||
// 1:1 chat uses (they also work cross-connection, where the member's
|
||||
// gateway can't see this machine's files). Images queue as vision tiles,
|
||||
// PDFs render per-page via pdf.attach, and other files materialize in the
|
||||
// session workspace (their @file: refs are appended to the prompt so the
|
||||
// member's file tools can read them). A failed attach degrades that
|
||||
// member to text-only; the transcript line still names the attachment so
|
||||
// the member knows something was shared.
|
||||
const fileRefs: string[] = []
|
||||
|
||||
for (const img of Array.isArray(images) ? images : []) {
|
||||
if (!img || typeof img.data !== 'string' || !img.data) {
|
||||
continue
|
||||
}
|
||||
|
||||
try {
|
||||
if (img.kind === 'pdf') {
|
||||
await requestForBot(member, 'pdf.attach', {
|
||||
session_id: runtime,
|
||||
content_base64: img.data,
|
||||
filename: img.name || 'attachment.pdf'
|
||||
})
|
||||
} else if (img.kind === 'file') {
|
||||
const res = (await requestForBot(member, 'file.attach', {
|
||||
session_id: runtime,
|
||||
data_url: img.data,
|
||||
name: img.name || 'attachment'
|
||||
})) as { ref_text?: string }
|
||||
|
||||
if (res?.ref_text) {
|
||||
fileRefs.push(`${img.name || 'attachment'} → ${res.ref_text}`)
|
||||
}
|
||||
} else {
|
||||
await requestForBot(member, 'image.attach_bytes', {
|
||||
session_id: runtime,
|
||||
content_base64: img.data,
|
||||
filename: img.name || 'attachment.png'
|
||||
})
|
||||
}
|
||||
} catch {
|
||||
/* text-only fallback for this member */
|
||||
}
|
||||
}
|
||||
|
||||
const messages = Array.isArray(state?.messages) ? state.messages : []
|
||||
const busy = Boolean(state?.inflight || state?.running)
|
||||
// A clarify blocking inside the member's session is a question for the
|
||||
// HUMAN (#90694) — mirror it into the room store so a card renders, and
|
||||
// hold the turn open: the member isn't stalling, it's waiting on us.
|
||||
const awaitingUser = syncGroupClarify(group, member, state)
|
||||
const done = !busy && !awaitingUser
|
||||
if (!binding.isLive()) {
|
||||
return null
|
||||
}
|
||||
|
||||
if (messages.length > before && done) {
|
||||
const replyText = pickGroupTurnReply(messages, before)
|
||||
const turnText = fileRefs.length
|
||||
? `${prompt}\n\nAttached files staged in your session workspace:\n${fileRefs.join('\n')}`
|
||||
: prompt
|
||||
|
||||
// #93602: one-shot recovery when the runtime session was reaped between
|
||||
// minting and submitting. Tracks the runtime id the submit landed on so
|
||||
// the poll fallback below targets a live session.
|
||||
const liveRuntime = await submitGroupTurnPrompt(member, runtime, stored, turnText)
|
||||
|
||||
if (!binding.isLive()) {
|
||||
return null
|
||||
}
|
||||
|
||||
runtimeIds.add(liveRuntime)
|
||||
const started = Date.now()
|
||||
let deadline = started + GROUP_TURN_TIMEOUT_MS
|
||||
// After the terminal frame fires, the gateway still has to flip
|
||||
// session.running off in its turn `finally` — re-check quickly for a few
|
||||
// beats instead of falling back to the slow backstop cadence.
|
||||
let quickRechecks = 0
|
||||
|
||||
while (Date.now() < deadline) {
|
||||
const signalled = quickRechecks
|
||||
? await waitForTurnSignal([], GROUP_TURN_SETTLE_RECHECK_MS)
|
||||
: await waitForTurnSignal([...runtimeIds], GROUP_TURN_POLL_MS)
|
||||
|
||||
quickRechecks = signalled ? GROUP_TURN_SETTLE_RECHECKS : Math.max(0, quickRechecks - 1)
|
||||
|
||||
// #91868/#94569: an explicit stop (stopGroupThread) bumped the epoch AND
|
||||
// held this member — the member's session was interrupted, so nothing is
|
||||
// coming; abandon the poll instead of grinding until the deadline. Both
|
||||
// conditions on purpose: an ordinary newer send bumps the epoch WITHOUT
|
||||
// a hold, and that turn must keep polling so finished work can still be
|
||||
// delivered (the #93127 commit check decides its fate, not this loop).
|
||||
if (!binding.isLive()) {
|
||||
return null
|
||||
}
|
||||
|
||||
const roomDuringPoll = $groupChats.get()[group] || {}
|
||||
|
||||
if ((roomDuringPoll.epoch || 0) !== dispatchEpoch && (roomDuringPoll.holds || {})[memberKey]) {
|
||||
return null
|
||||
}
|
||||
|
||||
let state: GroupSessionSnapshot | null = null
|
||||
|
||||
try {
|
||||
state = (await requestForBot(member, 'session.resume', {
|
||||
session_id: stored || liveRuntime,
|
||||
profile: member.name
|
||||
})) as GroupSessionSnapshot
|
||||
} catch {
|
||||
continue
|
||||
}
|
||||
|
||||
if (!binding.isLive()) {
|
||||
return null
|
||||
}
|
||||
|
||||
if (state?.session_id) {
|
||||
runtimeIds.add(state.session_id)
|
||||
}
|
||||
|
||||
const messages = Array.isArray(state?.messages) ? state.messages : []
|
||||
const busy = Boolean(state?.inflight || state?.running)
|
||||
// A clarify blocking inside the member's session is a question for the
|
||||
// HUMAN (#90694) — mirror it into the room store so a card renders, and
|
||||
// hold the turn open: the member isn't stalling, it's waiting on us.
|
||||
const awaitingUser = syncGroupClarify(group, member, state)
|
||||
const done = !busy && !awaitingUser
|
||||
|
||||
if (messages.length > before && done) {
|
||||
const replyText = pickGroupTurnReply(messages, before)
|
||||
|
||||
if (replyText !== null) {
|
||||
recordGroupActivity(group, {
|
||||
kind: isGroupPassText(replyText) ? 'passed' : 'replied',
|
||||
member: member.name,
|
||||
thread
|
||||
})
|
||||
|
||||
return replyText
|
||||
}
|
||||
|
||||
if (replyText !== null) {
|
||||
recordGroupActivity(group, {
|
||||
kind: isGroupPassText(replyText) ? 'passed' : 'replied',
|
||||
kind: 'passed',
|
||||
member: member.name,
|
||||
thread
|
||||
})
|
||||
|
||||
return replyText
|
||||
return null
|
||||
}
|
||||
|
||||
recordGroupActivity(group, {
|
||||
kind: 'passed',
|
||||
member: member.name,
|
||||
thread
|
||||
})
|
||||
// Still visibly working — or waiting on the user's answer to a clarify:
|
||||
// extend the deadline (never past the hard cap). A pending question must
|
||||
// outlive the base turn timeout or it dies unanswered at 3 minutes.
|
||||
if (busy || awaitingUser) {
|
||||
deadline = Math.min(started + GROUP_TURN_HARD_CAP_MS, Math.max(deadline, Date.now() + GROUP_TURN_TIMEOUT_MS))
|
||||
}
|
||||
}
|
||||
|
||||
if (!binding.isLive()) {
|
||||
return null
|
||||
}
|
||||
|
||||
// Still visibly working — or waiting on the user's answer to a clarify:
|
||||
// extend the deadline (never past the hard cap). A pending question must
|
||||
// outlive the base turn timeout or it dies unanswered at 3 minutes.
|
||||
if (busy || awaitingUser) {
|
||||
deadline = Math.min(started + GROUP_TURN_HARD_CAP_MS, Math.max(deadline, Date.now() + GROUP_TURN_TIMEOUT_MS))
|
||||
}
|
||||
}
|
||||
|
||||
// Timeout — clear any still-mirrored question card (the server-side
|
||||
// clarify timeout runs its own course) and read as a pass, but remember the baseline + thread
|
||||
// (runtime-only) so the finished reply can be posted late into the RIGHT
|
||||
// thread instead of vanishing.
|
||||
recordGroupActivity(group, {
|
||||
kind: 'timed-out',
|
||||
member: member.name,
|
||||
thread
|
||||
})
|
||||
syncGroupClarify(group, member, null)
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.stranded = {
|
||||
...(r.stranded || {}),
|
||||
[groupMemberKey(member)]: {
|
||||
before,
|
||||
thread
|
||||
// Timeout — clear any still-mirrored question card (the server-side
|
||||
// clarify timeout runs its own course) and read as a pass, but remember the baseline + thread
|
||||
// (runtime-only) so the finished reply can be posted late into the RIGHT
|
||||
// thread instead of vanishing.
|
||||
recordGroupActivity(group, {
|
||||
kind: 'timed-out',
|
||||
member: member.name,
|
||||
thread
|
||||
})
|
||||
syncGroupClarify(group, member, null)
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.stranded = {
|
||||
...(r.stranded || {}),
|
||||
[groupMemberKey(member)]: {
|
||||
before,
|
||||
thread
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return r
|
||||
})
|
||||
return r
|
||||
})
|
||||
|
||||
return null
|
||||
return null
|
||||
} finally {
|
||||
binding.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
/** Post a timed-out member's finished reply into the room, if it landed
|
||||
* after we stopped waiting. Called at the member's next turn boundary and
|
||||
* on user sends, so long-running work is delivered late rather than lost. */
|
||||
export async function harvestStrandedGroupReply(group: string, member: GroupMember) {
|
||||
const memberKey = groupMemberKey(member)
|
||||
const room = $groupChats.get()[group] || {}
|
||||
const marker = room.stranded?.[memberKey]
|
||||
// Markers were a bare number before threads; normalize both shapes.
|
||||
const strandedBefore = typeof marker === 'number' ? marker : marker?.before
|
||||
const strandedThread = (typeof marker === 'object' && marker?.thread) || 'legacy'
|
||||
|
||||
if (typeof strandedBefore !== 'number') {
|
||||
return
|
||||
}
|
||||
|
||||
let state: GroupSessionSnapshot | null = null
|
||||
const binding = followGroupChat(group, name => { group = name })
|
||||
|
||||
try {
|
||||
const stored = room.sessions?.[memberKey]
|
||||
state = (await requestForBot(member, 'session.resume', {
|
||||
session_id: stored || `Group: ${room.roomId || group}`,
|
||||
profile: member.name
|
||||
})) as GroupSessionSnapshot
|
||||
} catch {
|
||||
return // source unreachable — leave the marker for the next boundary
|
||||
}
|
||||
const memberKey = groupMemberKey(member)
|
||||
const room = $groupChats.get()[group] || {}
|
||||
const marker = room.stranded?.[memberKey]
|
||||
// Markers were a bare number before threads; normalize both shapes.
|
||||
const strandedBefore = typeof marker === 'number' ? marker : marker?.before
|
||||
const strandedThread = (typeof marker === 'object' && marker?.thread) || 'legacy'
|
||||
|
||||
if (state?.inflight || state?.running) {
|
||||
return // still grinding — keep waiting
|
||||
}
|
||||
|
||||
// A stranded member blocked on a clarify is not "grinding" — surface the
|
||||
// question card (#90694) and keep the marker until it resolves.
|
||||
if (syncGroupClarify(group, member, state)) {
|
||||
return
|
||||
}
|
||||
|
||||
// Done (or dead): the marker is consumed either way.
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
const next = {
|
||||
...(r.stranded || {})
|
||||
if (typeof strandedBefore !== 'number') {
|
||||
return
|
||||
}
|
||||
|
||||
delete next[memberKey]
|
||||
r.stranded = next
|
||||
let state: GroupSessionSnapshot | null = null
|
||||
|
||||
return r
|
||||
})
|
||||
const messages = Array.isArray(state?.messages) ? state.messages : []
|
||||
try {
|
||||
const stored = room.sessions?.[memberKey]
|
||||
state = (await requestForBot(member, 'session.resume', {
|
||||
session_id: stored || `Group: ${room.roomId || group}`,
|
||||
profile: member.name
|
||||
})) as GroupSessionSnapshot
|
||||
} catch {
|
||||
return // source unreachable — leave the marker for the next boundary
|
||||
}
|
||||
|
||||
if (messages.length <= strandedBefore) {
|
||||
return
|
||||
}
|
||||
if (!binding.isLive()) {
|
||||
return
|
||||
}
|
||||
|
||||
const reply = pickGroupTurnReply(messages, strandedBefore)
|
||||
// Pending prompts are authoritative even while the session is running.
|
||||
const awaitingUser = syncGroupClarify(group, member, state)
|
||||
|
||||
if (reply && !isGroupPassText(reply)) {
|
||||
recordGroupActivity(group, {
|
||||
kind: 'delivered',
|
||||
member: member.name,
|
||||
thread: strandedThread
|
||||
})
|
||||
appendGroupChatEntry(
|
||||
group,
|
||||
{
|
||||
kind: 'member',
|
||||
name: member.name,
|
||||
...(member.remoteSource
|
||||
? {
|
||||
source: member.connectionLabel || member.connectionId
|
||||
}
|
||||
: {})
|
||||
},
|
||||
reply,
|
||||
strandedThread
|
||||
)
|
||||
if (state?.inflight || state?.running || awaitingUser) {
|
||||
return
|
||||
}
|
||||
|
||||
// Done (or dead): the marker is consumed either way.
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.watermarks[`${strandedThread}::${memberKey}`] = r.log.length
|
||||
const next = {
|
||||
...(r.stranded || {})
|
||||
}
|
||||
|
||||
delete next[memberKey]
|
||||
r.stranded = next
|
||||
|
||||
return r
|
||||
})
|
||||
const messages = Array.isArray(state?.messages) ? state.messages : []
|
||||
|
||||
if (messages.length <= strandedBefore) {
|
||||
return
|
||||
}
|
||||
|
||||
const reply = pickGroupTurnReply(messages, strandedBefore)
|
||||
|
||||
if (reply && !isGroupPassText(reply)) {
|
||||
recordGroupActivity(group, {
|
||||
kind: 'delivered',
|
||||
member: member.name,
|
||||
thread: strandedThread
|
||||
})
|
||||
appendGroupChatEntry(
|
||||
group,
|
||||
{
|
||||
kind: 'member',
|
||||
name: member.name,
|
||||
...(member.remoteSource
|
||||
? {
|
||||
source: member.connectionLabel || member.connectionId
|
||||
}
|
||||
: {})
|
||||
},
|
||||
reply,
|
||||
strandedThread
|
||||
)
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.watermarks[`${strandedThread}::${memberKey}`] = r.log.length
|
||||
|
||||
return r
|
||||
})
|
||||
}
|
||||
} finally {
|
||||
binding.dispose()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,92 +0,0 @@
|
||||
/**
|
||||
* The needs-you attention badge the roster shows for a Group Chat row is
|
||||
* derived from TWO independent stores: $groupNeedsYou (an @user mention)
|
||||
* and $groupClarify (a pending clarify/approval, via groupHasPendingClarify).
|
||||
* roster-pane.tsx subscribes to both and combines them with ||.
|
||||
*
|
||||
* This harness proves the REAL composition works end-to-end: real stores,
|
||||
* real useValue subscriptions, and the real groupHasPendingClarify helper —
|
||||
* not a re-implementation of the derive expression. Mutating a store here
|
||||
* must be enough to flip what the harness renders, with no manual refresh.
|
||||
*/
|
||||
import { useValue } from '@hermes/plugin-sdk'
|
||||
import { act, render, screen } from '@testing-library/react'
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
|
||||
import { $groupClarify, $groupNeedsYou } from './group-chat'
|
||||
import { groupHasPendingClarify } from './group-turns'
|
||||
import type { GroupPrompt } from './types'
|
||||
|
||||
const GROUP = 'Core'
|
||||
|
||||
function mockPrompt(overrides: Partial<GroupPrompt> = {}): GroupPrompt {
|
||||
return {
|
||||
at: Date.now(),
|
||||
choices: ['staging', 'prod'],
|
||||
group: GROUP,
|
||||
kind: 'clarify',
|
||||
member: 'research',
|
||||
memberKey: 'research',
|
||||
multiSelect: false,
|
||||
question: 'Which env?',
|
||||
questions: null,
|
||||
requestId: 'req-1',
|
||||
sessionId: null,
|
||||
...overrides
|
||||
}
|
||||
}
|
||||
|
||||
function AttentionHarness({ group }: { group: string }) {
|
||||
const mentions = useValue($groupNeedsYou)
|
||||
const clarifies = useValue($groupClarify)
|
||||
const needsYou = Boolean(mentions[group]) || groupHasPendingClarify(clarifies, group)
|
||||
|
||||
return <div data-testid="attention">{String(needsYou)}</div>
|
||||
}
|
||||
|
||||
describe('roster attention badge — real stores, real useValue, real helper', () => {
|
||||
afterEach(() => {
|
||||
$groupNeedsYou.set({})
|
||||
$groupClarify.set({})
|
||||
})
|
||||
|
||||
it('reflects a clarify being added and then removed', () => {
|
||||
act(() => {
|
||||
render(<AttentionHarness group={GROUP} />)
|
||||
})
|
||||
|
||||
expect(screen.getByTestId('attention').textContent).toBe('false')
|
||||
|
||||
act(() => {
|
||||
$groupClarify.set({ 'Core::research': mockPrompt() })
|
||||
})
|
||||
|
||||
expect(screen.getByTestId('attention').textContent).toBe('true')
|
||||
|
||||
act(() => {
|
||||
$groupClarify.set({})
|
||||
})
|
||||
|
||||
expect(screen.getByTestId('attention').textContent).toBe('false')
|
||||
})
|
||||
|
||||
it('keeps attention true when clarify clears but the mention persists', () => {
|
||||
act(() => {
|
||||
render(<AttentionHarness group={GROUP} />)
|
||||
})
|
||||
|
||||
act(() => {
|
||||
$groupNeedsYou.set({ [GROUP]: true })
|
||||
$groupClarify.set({ 'Core::research': mockPrompt() })
|
||||
})
|
||||
|
||||
expect(screen.getByTestId('attention').textContent).toBe('true')
|
||||
|
||||
// Clarify resolves — the unrelated mention must keep the badge lit.
|
||||
act(() => {
|
||||
$groupClarify.set({})
|
||||
})
|
||||
|
||||
expect(screen.getByTestId('attention').textContent).toBe('true')
|
||||
})
|
||||
})
|
||||
@@ -96,7 +96,7 @@ Groups are standalone rows in the same activity-ordered roster as Bot DMs. A Bot
|
||||
|
||||
- **One visible conversation.** Public messages and each member's reply stay readable in arrival order, with the speaker's name and timestamp. Starting another topic does not collapse earlier replies. **Reply in thread** continues that topic without reordering the room; **Activity** is a secondary status view, not a replacement for messages. Private Bot Chats remain separate.
|
||||
- Your message triggers up to **three serial rounds** of member turns. @-mentioned Bots respond (everyone responds when nobody is mentioned); each Bot replies briefly or passes, and the room settles when a full round stays silent.
|
||||
- Bots pull each other in with `@name`, and escalate real judgment calls to you with `@user` — the group row shows a **needs you** badge when that happens.
|
||||
- Bots pull each other in with `@name`, and escalate real judgment calls to you with `@user` — the group row shows a **needs you** badge when that happens. Pending questions and command approvals also light that badge; resolving the last prompt clears only prompt attention, not an independent mention. Prompts follow a renamed room, while disbanding retires them even if a member's in-flight poll arrives later.
|
||||
- Hard caps (10 messages per send, 3 rounds) keep rooms from spinning.
|
||||
- Each member keeps its own persistent `Group: <name>` session, so room context survives like any other conversation.
|
||||
- **Not every Bot replies to every message.** Speaking is each member's own choice — a Bot replies only when it has something new to add and passes otherwise, and @-mentioning specific members scopes the round to them. Expect the members you addressed (or whoever has something to say) to speak, and the rest to stay quiet.
|
||||
|
||||
Reference in New Issue
Block a user