fix(bot-mode): keep group follow-ups ordered and late answers visible
Serialize room drives through their actual member completion, freeze input watermarks by retained entry identity, and share the same completion path with handoff continuations. Stop discards queued work without releasing an active owner early; rename follows the existing room binding. Observe stranded replies for the hard-cap duration plus grace after the foreground wait, and retain unresolved failures in collapsed Activity. Never automatically retry an ambiguous failed submit within the same drive. Adapted from the queue and boundary approach in #92041 by @enwaiax and harvest-budget approach in #107193 by @Finn763; #106502 by @wadib identified failed-submit watermark consumption. The implementation retains current numeric watermark storage, room lifecycle bindings and serial round limits. Related: #92003, #105247, #100026
This commit is contained in:
176
apps/desktop/e2e/group-turn-integrity.spec.ts
Normal file
176
apps/desktop/e2e/group-turn-integrity.spec.ts
Normal file
@@ -0,0 +1,176 @@
|
||||
import { MOCK_REPLY } from '../../../tests-js/scripts/mock-server'
|
||||
|
||||
import { type MockBackendFixture, setupMockBackend, waitForAppReady } from './fixtures'
|
||||
import { expect, test } from './test'
|
||||
|
||||
let fixture: MockBackendFixture | null = null
|
||||
|
||||
async function openBots(page: MockBackendFixture['page']): Promise<void> {
|
||||
const tab = page.getByRole('button', { name: 'Bots', exact: true }).or(page.getByRole('tab', { name: 'Bots', exact: true })).first()
|
||||
await tab.click()
|
||||
await expect(page.getByRole('button', { name: 'New bot or group chat' })).toBeVisible()
|
||||
}
|
||||
|
||||
async function createAgent(page: MockBackendFixture['page'], name: string, title: string): Promise<void> {
|
||||
await page.getByRole('button', { name: 'New bot or group chat' }).click()
|
||||
await page.getByRole('menuitem', { name: 'New Bot' }).click()
|
||||
|
||||
const dialog = page.getByRole('dialog', { name: 'New Bot' })
|
||||
await dialog.getByPlaceholder('inbox-triage').fill(name)
|
||||
await dialog.getByPlaceholder('Inbox Triage').fill(title)
|
||||
await dialog.getByRole('button', { name: 'Create Bot' }).click()
|
||||
await expect(dialog).toBeHidden({ timeout: 30_000 })
|
||||
await expect(page.getByRole('button', { name: new RegExp(`^${title}\\b`) }).first()).toBeVisible({ timeout: 30_000 })
|
||||
}
|
||||
|
||||
async function createRoom(page: MockBackendFixture['page']) {
|
||||
await openBots(page)
|
||||
await createAgent(page, 'programmer', 'Programmer')
|
||||
await createAgent(page, 'reviewer', 'Reviewer')
|
||||
|
||||
await page.getByRole('button', { name: 'New bot or group chat' }).click()
|
||||
await page.getByRole('menuitem', { name: 'New Group Chat' }).click()
|
||||
|
||||
const dialog = page.getByRole('dialog', { name: 'New Group Chat' })
|
||||
|
||||
for (const title of ['Programmer', 'Reviewer']) {
|
||||
await dialog.getByText(title, { exact: true }).locator('xpath=ancestor::label').getByRole('checkbox').click()
|
||||
}
|
||||
|
||||
await dialog.getByRole('textbox', { name: 'Group name' }).fill('Programmer, Reviewer')
|
||||
await dialog.getByRole('button', { name: 'Create Group (2)' }).click()
|
||||
|
||||
const groupTab = page.getByRole('tab', { name: /Programmer, Reviewer Close/ })
|
||||
const groupComposer = page.getByRole('textbox', { name: 'Message Programmer, Reviewer' }).filter({ visible: true })
|
||||
await expect(groupTab).toBeVisible({ timeout: 20_000 })
|
||||
await expect(groupTab).toHaveAttribute('aria-selected', 'true')
|
||||
await expect(groupComposer).toBeVisible()
|
||||
|
||||
|
||||
return groupComposer
|
||||
}
|
||||
|
||||
test.beforeEach(async () => {
|
||||
fixture = await setupMockBackend({ mockServer: { holdFirstCompletionContaining: 'LANE_A_FIRST' } })
|
||||
await waitForAppReady(fixture, 120_000)
|
||||
})
|
||||
|
||||
test.afterEach(async () => {
|
||||
await fixture?.cleanup()
|
||||
fixture = null
|
||||
})
|
||||
|
||||
test('group follow-up waits for its active member and retains the reply', async () => {
|
||||
test.setTimeout(240_000)
|
||||
const page = fixture!.page
|
||||
const groupComposer = await createRoom(page)
|
||||
await groupComposer.fill('@programmer LANE_A_FIRST')
|
||||
await groupComposer.press('Enter')
|
||||
await fixture!.mock.waitForHeldCompletion()
|
||||
console.log('Held first inference; no second submit is allowed until release.')
|
||||
await page.screenshot({ path: '/tmp/botmode-campaign/lane-a-held.png' })
|
||||
await page.getByRole('button', { name: 'Reply in thread', exact: true }).click()
|
||||
const replyComposer = page.getByRole('textbox', { name: 'Reply in thread', exact: true })
|
||||
await replyComposer.fill('@programmer LANE_A_FOLLOWUP')
|
||||
await replyComposer.press('Enter')
|
||||
await page.waitForTimeout(3000)
|
||||
const overlapping = fixture!.mock.receivedPrompts.filter(p => p.includes('LANE_A_FOLLOWUP'))
|
||||
console.log('OVERLAPPING', overlapping)
|
||||
fixture!.mock.releaseHeldStream()
|
||||
expect(overlapping).toHaveLength(0)
|
||||
await expect.poll(() => fixture!.mock.receivedPrompts.some(p => p.includes('LANE_A_FOLLOWUP')), { timeout: 60000 }).toBe(true)
|
||||
const delivered = fixture!.mock.receivedPrompts.find(p => p.includes('LANE_A_FOLLOWUP'))!
|
||||
expect(delivered).toContain(MOCK_REPLY)
|
||||
expect(delivered).not.toContain('LANE_A_FIRST')
|
||||
expect(delivered.indexOf('LANE_A_FOLLOWUP')).toBeLessThan(delivered.indexOf(MOCK_REPLY))
|
||||
await expect(page.getByText(MOCK_REPLY, { exact: true }).first()).toBeVisible()
|
||||
await page.screenshot({ path: '/tmp/botmode-campaign/lane-a-followup-after.png' })
|
||||
})
|
||||
|
||||
test('quiet group still harvests a late answer after sixty observation ticks', async () => {
|
||||
test.setTimeout(240_000)
|
||||
const page = fixture!.page
|
||||
const groupComposer = await createRoom(page)
|
||||
|
||||
await groupComposer.fill('@programmer LANE_A_FIRST')
|
||||
await groupComposer.press('Enter')
|
||||
await fixture!.mock.waitForHeldCompletion()
|
||||
// Jump only the deadline clock, not WebSocket heartbeats or reconnect timers.
|
||||
await page.evaluate(() => {
|
||||
const now = Date.now
|
||||
Date.now = () => now() + 21 * 60_000
|
||||
|
||||
const timeout = window.setTimeout.bind(window)
|
||||
|
||||
;(window as any).__harvestTicks = 0
|
||||
window.setTimeout = ((handler: TimerHandler, delay?: number, ...args: any[]) => {
|
||||
if (delay === 5000) {
|
||||
return timeout(() => {
|
||||
(window as any).__harvestTicks++
|
||||
|
||||
if (typeof handler === 'function') { handler(...args) }
|
||||
}, 100)
|
||||
}
|
||||
|
||||
return timeout(handler, delay, ...args)
|
||||
}) as typeof window.setTimeout
|
||||
})
|
||||
await expect.poll(() => page.evaluate(() => (window as any).__harvestTicks), { timeout: 60000 }).toBeGreaterThanOrEqual(60)
|
||||
await new Promise(resolve => setTimeout(resolve, 1500))
|
||||
console.log('Harvest ticks before release:', await page.evaluate(() => (window as any).__harvestTicks))
|
||||
console.log('Activity before release:', await page.getByRole('button', { name: /^Activity/ }).textContent())
|
||||
fixture!.mock.releaseHeldStream()
|
||||
await expect(page.getByText(MOCK_REPLY, { exact: true }).first()).toBeVisible({ timeout: 5000 })
|
||||
await page.screenshot({ path: '/tmp/botmode-campaign/lane-a-late-after.png' })
|
||||
})
|
||||
|
||||
test('a rejected member turn stays visible when the room settles', async () => {
|
||||
test.setTimeout(240_000)
|
||||
const page = fixture!.page
|
||||
await page.evaluate(() => {
|
||||
const send = WebSocket.prototype.send
|
||||
|
||||
WebSocket.prototype.send = function(data) {
|
||||
const frame = JSON.parse(String(data))
|
||||
|
||||
if (frame.method === 'prompt.submit' && JSON.stringify(frame.params).includes('LANE_A_FAILURE')) {
|
||||
(window as any).__rejected = ((window as any).__rejected || 0) + 1
|
||||
queueMicrotask(() => this.dispatchEvent(new MessageEvent('message', { data: JSON.stringify({ jsonrpc: '2.0', id: frame.id, error: { code: 4003, message: 'Controlled member admission refusal' } }) })))
|
||||
} else {
|
||||
send.call(this, data)
|
||||
}
|
||||
}
|
||||
})
|
||||
const groupComposer = await createRoom(page)
|
||||
|
||||
await groupComposer.fill('@programmer LANE_A_FAILURE')
|
||||
await groupComposer.press('Enter')
|
||||
await expect.poll(() => page.evaluate(() => (window as any).__rejected), { timeout: 30000 }).toBe(1)
|
||||
await expect(page.getByRole('button', { name: 'Stop', exact: true })).toHaveCount(0)
|
||||
await expect(page.getByRole('button', { name: /^Activity/ })).toContainText('Programmer hit an error')
|
||||
await page.screenshot({ path: '/tmp/botmode-campaign/lane-a-error-after.png' })
|
||||
})
|
||||
|
||||
test('Stop clears a queued follow-up and a direct mention resumes the held member', async () => {
|
||||
test.setTimeout(240_000)
|
||||
const page = fixture!.page
|
||||
const groupComposer = await createRoom(page)
|
||||
await groupComposer.fill('@programmer LANE_A_FIRST')
|
||||
await groupComposer.press('Enter')
|
||||
await fixture!.mock.waitForHeldCompletion()
|
||||
await page.getByRole('button', { name: 'Reply in thread', exact: true }).click()
|
||||
const replyComposer = page.getByRole('textbox', { name: 'Reply in thread', exact: true })
|
||||
await replyComposer.fill('@programmer LANE_A_CANCELLED')
|
||||
await replyComposer.press('Enter')
|
||||
await page.getByRole('button', { name: 'Stop', exact: true }).click()
|
||||
fixture!.mock.releaseHeldStream()
|
||||
await expect(page.getByRole('button', { name: 'Stop', exact: true })).toHaveCount(0)
|
||||
await page.waitForTimeout(2000)
|
||||
expect(fixture!.mock.receivedPrompts.some(p => p.includes('LANE_A_CANCELLED'))).toBe(false)
|
||||
console.log('Stop: queued inference count = 0; paused status:', await page.getByText(/Paused:/).allTextContents())
|
||||
await groupComposer.fill('@programmer LANE_A_RESUME')
|
||||
await groupComposer.press('Enter')
|
||||
await expect.poll(() => fixture!.mock.receivedPrompts.some(p => p.includes('LANE_A_RESUME')), { timeout: 60000 }).toBe(true)
|
||||
await expect(page.getByText(MOCK_REPLY, { exact: true }).first()).toBeVisible()
|
||||
await page.screenshot({ path: '/tmp/botmode-campaign/lane-a-stop-resume.png' })
|
||||
})
|
||||
@@ -153,46 +153,25 @@ describe('turn arc', () => {
|
||||
})
|
||||
|
||||
describe('epoch scoping', () => {
|
||||
it('a newer send interrupts the previous run and records cancelled in the CURRENT epoch', async () => {
|
||||
const gates = new Map<number, { promise: Promise<string>; resolve: (value: string) => void }>()
|
||||
|
||||
const gate = (n: number) => {
|
||||
const existing = gates.get(n)
|
||||
|
||||
if (existing) {
|
||||
return existing.promise
|
||||
}
|
||||
|
||||
let resolve!: (value: string) => void
|
||||
|
||||
const promise = new Promise<string>(settle => {
|
||||
resolve = settle
|
||||
})
|
||||
|
||||
gates.set(n, { promise, resolve })
|
||||
|
||||
return promise
|
||||
}
|
||||
|
||||
const room = await loadRoom({ turn: ({ n }) => gate(n) })
|
||||
it('queues follow-ups without cancelling the active turn or losing its reply delta', async () => {
|
||||
let release!: (reply: string) => void
|
||||
const first = new Promise<string>(resolve => { release = resolve })
|
||||
const room = await loadRoom({ turn: ({ n }) => n === 1 ? first : '(pass)' })
|
||||
const member: GroupMember[] = [{ name: 'research', title: '' }]
|
||||
|
||||
room.rounds.sendToGroupChat('Busy', member, 'first ask')
|
||||
const thread = room.rounds.sendToGroupChat('Busy', member, 'first ask')!
|
||||
await drain(() => room.gateway.calls.length < 1, 50)
|
||||
room.rounds.sendToGroupChat('Busy', member, 'second ask, supersede')
|
||||
await drain(() => room.gateway.calls.length < 2, 50)
|
||||
|
||||
gates.get(2)?.resolve('from the new run')
|
||||
await drain(() => Boolean(room.chat.$groupChats.get().Busy?.running))
|
||||
gates.get(1)?.resolve('late from the old run')
|
||||
const epoch = room.chat.$groupChats.get().Busy.epoch
|
||||
room.rounds.sendToGroupChat('Busy', member, 'follow-up', thread)
|
||||
await drain(() => false)
|
||||
|
||||
const epoch = room.chat.$groupChats.get().Busy?.epoch || 0
|
||||
|
||||
expect(feed(room, 'Busy').some(event => event.kind === 'cancelled')).toBe(true)
|
||||
// The view shows only the current run: every visible event is this epoch.
|
||||
expect(room.activity.currentGroupActivity('Busy').every(event => (event.epoch || 0) === epoch)).toBe(true)
|
||||
expect(room.activity.currentGroupActivity('Busy').some(event => event.kind === 'cancelled')).toBe(true)
|
||||
expect(room.gateway.calls).toHaveLength(1)
|
||||
expect(room.chat.$groupChats.get().Busy.epoch).toBe(epoch)
|
||||
release('first reply')
|
||||
await drain(() => room.gateway.calls.length < 2)
|
||||
await drain(() => Boolean(room.chat.$groupChats.get().Busy?.running))
|
||||
expect(room.gateway.calls).toHaveLength(2)
|
||||
expect(room.gateway.calls[1].prompt).toMatch(/follow-up[\s\S]*first reply/)
|
||||
expect(room.gateway.calls[1].prompt).not.toContain('first ask')
|
||||
expect(feed(room, 'Busy').some(event => event.kind === 'cancelled')).toBe(false)
|
||||
})
|
||||
|
||||
it('epoch filtering drops events from a superseded run', async () => {
|
||||
|
||||
@@ -710,6 +710,23 @@ export function GroupChatWorkspace({ group, members, onBack, visible = true }: G
|
||||
// Events are epoch-tagged, so a superseded run's history drops out of view.
|
||||
const activityEvents: GroupActivityEntry[] = currentGroupActivity(group)
|
||||
const latestActivity = activityEvents.length ? activityEvents[activityEvents.length - 1] : null
|
||||
// A later "settled" event must not hide a member that failed to answer.
|
||||
// Successful completion for that member clears its unresolved warning.
|
||||
const unresolvedFailures = new Map<string, GroupActivityEntry>()
|
||||
|
||||
for (const event of activityEvents) {
|
||||
const key = event.member || ''
|
||||
|
||||
if (event.kind === 'failed' || event.kind === 'timed-out') {
|
||||
unresolvedFailures.set(key, event)
|
||||
} else if (event.kind === 'replied' || event.kind === 'passed' || event.kind === 'delivered') {
|
||||
unresolvedFailures.delete(key)
|
||||
}
|
||||
}
|
||||
|
||||
const summaryActivity = !room.running && unresolvedFailures.size
|
||||
? [...unresolvedFailures.values()].at(-1)!
|
||||
: latestActivity
|
||||
|
||||
// #94570 shell rewired onto the real primitive (#91868/#94569): the button
|
||||
// must stop the ROUND, not just spray per-member interrupts — without the
|
||||
@@ -735,8 +752,8 @@ export function GroupChatWorkspace({ group, members, onBack, visible = true }: G
|
||||
>
|
||||
<Codicon className="shrink-0 text-[0.65rem]" name={activityOpen ? 'chevron-down' : 'chevron-right'} />
|
||||
<span className="shrink-0 font-medium">{b.group.activity}</span>
|
||||
{latestActivity ? (
|
||||
<span className="min-w-0 flex-1 truncate">{`${groupActivityLabel(latestActivity)} · ${relativeTime(latestActivity.at)}`}</span>
|
||||
{summaryActivity ? (
|
||||
<span className={cn('min-w-0 flex-1 truncate', groupActivityTone(summaryActivity.kind))}>{`${groupActivityLabel(summaryActivity)} · ${relativeTime(summaryActivity.at)}`}</span>
|
||||
) : null}
|
||||
</RowButton>
|
||||
{room.running ? (
|
||||
|
||||
@@ -23,6 +23,7 @@ export interface GroupRoundMemberContext {
|
||||
startEpoch: number
|
||||
binding: { isLive(): boolean }
|
||||
isCurrent(): boolean
|
||||
failedMembers?: Set<string>
|
||||
}
|
||||
|
||||
/** #93129: a held member's skip must consume its delta exactly once —
|
||||
@@ -129,6 +130,11 @@ export async function runGroupRoundMember(
|
||||
member: GroupMember
|
||||
): Promise<boolean | null> {
|
||||
const { thread, startEpoch, binding } = context
|
||||
|
||||
if (context.failedMembers?.has(groupMemberKey(member))) {
|
||||
return false
|
||||
}
|
||||
|
||||
const prepared = prepareGroupRoundMember(context, member)
|
||||
|
||||
if (!prepared) {
|
||||
@@ -136,10 +142,13 @@ export async function runGroupRoundMember(
|
||||
}
|
||||
|
||||
const { room, markKey, prompt, deltaImages } = prepared
|
||||
const anchorId = room.log.at(-1)?.id ?? null
|
||||
let reply: null | string = null
|
||||
let accepted = false
|
||||
|
||||
try {
|
||||
reply = await runVisibleMemberTurn(context, member, prompt, deltaImages)
|
||||
accepted = true
|
||||
|
||||
// Needs-attention hook (#93091 item 3): a turn that produced a real
|
||||
// reply (or an explicit pass) is a good turn — clear the badge.
|
||||
@@ -165,6 +174,7 @@ export async function runGroupRoundMember(
|
||||
: {})
|
||||
})
|
||||
noteBotAttention(groupMemberKey(member), reason || error?.message || error)
|
||||
context.failedMembers?.add(groupMemberKey(member))
|
||||
reply = null // a failed turn is a pass, never a room error
|
||||
}
|
||||
|
||||
@@ -188,7 +198,6 @@ export async function runGroupRoundMember(
|
||||
}
|
||||
|
||||
const epochNow = roomNow.epoch || 0
|
||||
const anchorId = room.log.length ? room.log[room.log.length - 1].id : null
|
||||
const anchorIdx = anchorId === null ? -1 : roomNow.log.findIndex((e: GroupMessage) => e.id === anchorId)
|
||||
// Anchor trimmed away ⇒ every pre-turn entry was dropped, so every
|
||||
// surviving entry is newer — scanning the whole log stays exact.
|
||||
@@ -198,7 +207,10 @@ export async function runGroupRoundMember(
|
||||
(e: GroupMessage) => e.from?.kind === 'user' && groupThreadOf(e) === thread
|
||||
)
|
||||
|
||||
if (!shouldCommitMemberTurn(startEpoch, epochNow, newerUserEntryInThread)) {
|
||||
if (
|
||||
(epochNow !== startEpoch && roomNow.holds?.[groupMemberKey(member)]) ||
|
||||
!shouldCommitMemberTurn(startEpoch, epochNow, newerUserEntryInThread)
|
||||
) {
|
||||
recordGroupActivity(context.group, {
|
||||
kind: 'cancelled',
|
||||
member: member.name,
|
||||
@@ -208,12 +220,16 @@ export async function runGroupRoundMember(
|
||||
return null
|
||||
}
|
||||
|
||||
// The member has now seen everything up to the pre-reply log length.
|
||||
updateGroupChat(context.group, (r: GroupChatRoom) => {
|
||||
r.watermarks[markKey] = r.log.length
|
||||
// Resolve the frozen submit boundary against the retained log. If it was
|
||||
// trimmed away, every surviving entry is still unseen. Throws do not
|
||||
// acknowledge input, and a timed-out turn keeps its submitted boundary.
|
||||
if (accepted) {
|
||||
updateGroupChat(context.group, (r: GroupChatRoom) => {
|
||||
r.watermarks[markKey] = anchorIdx + 1
|
||||
|
||||
return r
|
||||
})
|
||||
return r
|
||||
})
|
||||
}
|
||||
|
||||
if (reply !== null && !isGroupPassText(reply)) {
|
||||
appendGroupChatEntry(
|
||||
@@ -230,107 +246,11 @@ export async function runGroupRoundMember(
|
||||
reply,
|
||||
thread
|
||||
)
|
||||
// Its own message counts as seen too.
|
||||
// A reply cannot acknowledge user entries that arrived during inference.
|
||||
updateGroupChat(context.group, (r: GroupChatRoom) => {
|
||||
r.watermarks[markKey] = r.log.length
|
||||
|
||||
return r
|
||||
})
|
||||
|
||||
return true
|
||||
}
|
||||
|
||||
return false
|
||||
}
|
||||
|
||||
async function runGroupContinuationMember(
|
||||
context: GroupRoundMemberContext,
|
||||
member: GroupMember
|
||||
): Promise<boolean | null> {
|
||||
const { members, thread, binding, isCurrent } = context
|
||||
|
||||
const room = $groupChats.get()[context.group] || {
|
||||
log: [],
|
||||
watermarks: {}
|
||||
}
|
||||
|
||||
const memberKey = groupMemberKey(member)
|
||||
const markKey = `${thread}::${memberKey}`
|
||||
const seen = room.watermarks[markKey] || 0
|
||||
const delta = room.log.slice(seen).filter((e: GroupMessage) => groupThreadOf(e) === thread)
|
||||
|
||||
// A cited member always has delta here (the citing reply IS in
|
||||
// its tail); skip defensively anyway so an empty prompt never
|
||||
// fires.
|
||||
if (!delta.length) {
|
||||
return false
|
||||
}
|
||||
|
||||
const heldEntry = (room.holds || {})[memberKey]
|
||||
|
||||
if (heldEntry) {
|
||||
return false // holds still apply to continuation turns (#93129)
|
||||
}
|
||||
|
||||
const prompt = buildGroupChatTurnPrompt({
|
||||
groupName: context.group,
|
||||
members,
|
||||
viewer: member,
|
||||
// The continuation prompt centers on what the member missed:
|
||||
// everything since its watermark, which includes the reply
|
||||
// that cites it.
|
||||
deltaLines: delta.slice(-GROUP_CHAT_HISTORY_LIMIT).map((e: GroupMessage) => formatGroupChatLine(e, member))
|
||||
})
|
||||
|
||||
let continuationReply: null | string = null
|
||||
|
||||
try {
|
||||
continuationReply = await runVisibleMemberTurn(context, member, prompt)
|
||||
|
||||
if (continuationReply !== null) {
|
||||
clearBotAttention(memberKey)
|
||||
}
|
||||
} catch (error: any) {
|
||||
if (!binding.isLive()) {
|
||||
return null
|
||||
}
|
||||
|
||||
recordGroupActivity(context.group, {
|
||||
kind: 'failed',
|
||||
member: member.name,
|
||||
thread
|
||||
})
|
||||
noteBotAttention(memberKey, error?.message || error)
|
||||
continuationReply = null
|
||||
}
|
||||
|
||||
if (!isCurrent()) {
|
||||
return null
|
||||
}
|
||||
|
||||
updateGroupChat(context.group, (r: GroupChatRoom) => {
|
||||
r.watermarks[markKey] = r.log.length
|
||||
|
||||
return r
|
||||
})
|
||||
|
||||
if (continuationReply !== null && !isGroupPassText(continuationReply)) {
|
||||
appendGroupChatEntry(
|
||||
context.group,
|
||||
{
|
||||
kind: 'member',
|
||||
name: member.name,
|
||||
...(member.remoteSource
|
||||
? {
|
||||
source: member.connectionLabel || member.connectionId
|
||||
}
|
||||
: {})
|
||||
},
|
||||
continuationReply,
|
||||
thread
|
||||
)
|
||||
updateGroupChat(context.group, (r: GroupChatRoom) => {
|
||||
r.watermarks[markKey] = r.log.length
|
||||
if (r.watermarks[markKey] === r.log.length - 1) {
|
||||
r.watermarks[markKey] = r.log.length
|
||||
}
|
||||
|
||||
return r
|
||||
})
|
||||
@@ -365,7 +285,7 @@ export async function runGroupContinuationMembers(
|
||||
break
|
||||
}
|
||||
|
||||
const result = await runGroupContinuationMember(context, member)
|
||||
const result = await runGroupRoundMember(context, member)
|
||||
|
||||
if (result === null) {
|
||||
return null
|
||||
|
||||
@@ -230,6 +230,19 @@ describe('round lifecycle', () => {
|
||||
expect(posted.length).toBeLessThanOrEqual(room.chat.GROUP_CHAT_MAX_MESSAGES)
|
||||
})
|
||||
|
||||
it('does not retry ambiguous member admission in later rounds or continuations', async () => {
|
||||
const room = await loadRoom({ turn: ({ profile }) => {
|
||||
if (profile === 'builder') { throw new Error('Ambiguous admission failure') }
|
||||
|
||||
return '@builder please investigate'
|
||||
} })
|
||||
|
||||
room.rounds.sendToGroupChat('Failure', MEMBERS.slice(0, 2), '@research start')
|
||||
await settle(room, 'Failure')
|
||||
expect(room.gateway.calls.filter(call => call.profile === 'builder')).toHaveLength(1)
|
||||
expect(Object.keys(room.chat.$groupChats.get().Failure.watermarks).some(key => key.endsWith('::builder'))).toBe(false)
|
||||
})
|
||||
|
||||
it('treats a failed member turn as a pass, not a room error', async () => {
|
||||
const room = await loadRoom({
|
||||
turn: ({ profile }) => {
|
||||
@@ -292,6 +305,25 @@ describe('round lifecycle', () => {
|
||||
})
|
||||
|
||||
describe('per-member delta', () => {
|
||||
it('retained-log trimming cannot acknowledge messages appended during inference', async () => {
|
||||
let release!: (reply: string) => void
|
||||
const held = new Promise<string>(resolve => { release = resolve })
|
||||
const room = await loadRoom({ turn: ({ n }) => n === 1 ? held : '(pass)' })
|
||||
const members = [MEMBERS[0]]
|
||||
const thread = room.rounds.sendToGroupChat('Trim', members, 'delivered')!
|
||||
await drain(() => room.gateway.calls.length < 1)
|
||||
|
||||
for (let i = 0; i < 100; i++) {
|
||||
room.chat.appendGroupChatEntry('Trim', { kind: 'user', name: 'You' }, `unseen-${i}`, thread)
|
||||
}
|
||||
|
||||
release('(pass)')
|
||||
await settle(room, 'Trim')
|
||||
expect(room.chat.$groupChats.get().Trim.watermarks[`${thread}::research`]).toBe(0)
|
||||
await room.rounds.runGroupChatRounds('Trim', members, thread)
|
||||
expect(room.gateway.calls.at(-1)?.prompt).toContain('unseen-99')
|
||||
})
|
||||
|
||||
it('feeds a second send only the NEW messages', async () => {
|
||||
const room = await loadRoom()
|
||||
const member: GroupMember[] = [{ name: 'research', title: '' }]
|
||||
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
GROUP_CHAT_MAX_CONTINUATIONS,
|
||||
GROUP_CHAT_MAX_MESSAGES,
|
||||
GROUP_CHAT_MAX_ROUNDS,
|
||||
groupChatRoomKey,
|
||||
groupThreadOf,
|
||||
mintGroupThreadId,
|
||||
updateGroupChat
|
||||
@@ -20,7 +21,7 @@ import type { GroupChatRoom, GroupHoldStamp } from './group-chat'
|
||||
import { durableGroupChatMembers, followGroupChat, groupMemberKey } from './group-membership'
|
||||
import { runGroupContinuationMembers, runGroupRoundMember } from './group-round-members'
|
||||
import { rejectGroupSlashCommand } from './group-slash'
|
||||
import { harvestStrandedGroupReply } from './group-turns'
|
||||
import { GROUP_TURN_HARD_CAP_MS, harvestStrandedGroupReply } from './group-turns'
|
||||
import { requestForBot } from './routing'
|
||||
import type { Attachment, GroupMember, GroupMessage } from './types'
|
||||
|
||||
@@ -353,6 +354,7 @@ export async function stopGroupThread(group: string, thread: null | string, memb
|
||||
const room = $groupChats.get()[group] || {}
|
||||
const roster = Array.isArray(members) && members.length ? members : room.members || []
|
||||
const onTurn = room.turn || null
|
||||
groupChatDrives.get(groupChatRoomKey(group, room))?.pending.clear()
|
||||
|
||||
const stamp: GroupHoldStamp = {
|
||||
at: Date.now(),
|
||||
@@ -412,8 +414,8 @@ export async function stopGroupThread(group: string, thread: null | string, memb
|
||||
}
|
||||
|
||||
/** Drive one bounded round-robin turn for ONE THREAD. Serial — one member at
|
||||
* a time. A newer user send bumps the room epoch; this loop notices at the
|
||||
* next member boundary, bails, and the newest send's own loop takes over.
|
||||
* a time. User follow-ups queue behind this drive; Stop invalidates its
|
||||
* epoch and discards queued continuations.
|
||||
* 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) {
|
||||
@@ -431,6 +433,7 @@ export async function runGroupChatRounds(group: string, members: GroupMember[],
|
||||
members,
|
||||
thread,
|
||||
startEpoch,
|
||||
failedMembers: new Set<string>(),
|
||||
binding,
|
||||
isCurrent
|
||||
}
|
||||
@@ -592,7 +595,8 @@ export async function runGroupChatRounds(group: string, members: GroupMember[],
|
||||
}
|
||||
|
||||
/** Bounded background harvest for members whose replies outlived the turn
|
||||
* loop. Polls every 5s for up to 5 minutes; stops early when nothing is
|
||||
* loop. Watches for a further hard-cap duration plus a minute of grace
|
||||
* after foreground polling ends; stops early when nothing is
|
||||
* 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) {
|
||||
@@ -602,7 +606,7 @@ async function harvestStrandedUntilSettled(group: string, members: GroupMember[]
|
||||
|
||||
try {
|
||||
const HARVEST_INTERVAL_MS = 5000
|
||||
const HARVEST_MAX_TRIES = 60
|
||||
const HARVEST_MAX_TRIES = Math.ceil((GROUP_TURN_HARD_CAP_MS + 60000) / HARVEST_INTERVAL_MS)
|
||||
|
||||
for (let attempt = 0; attempt < HARVEST_MAX_TRIES; attempt++) {
|
||||
await new Promise(resolve => window.setTimeout(resolve, HARVEST_INTERVAL_MS))
|
||||
@@ -649,8 +653,8 @@ async function harvestStrandedUntilSettled(group: string, members: GroupMember[]
|
||||
|
||||
/** User send into a group room. `thread` continues that thread (its reply
|
||||
* box); omitted/null mints a NEW thread — the main composer's Slack shape.
|
||||
* Appends, bumps the room epoch (supersedes any running loop at its next
|
||||
* member boundary), and starts the turn drive for the target thread.
|
||||
* Appends and queues the target thread behind the active room drive,
|
||||
* without redirecting a member whose inference is still in flight.
|
||||
* Returns the thread id the message landed in. */
|
||||
export function sendToGroupChat(
|
||||
group: string,
|
||||
@@ -696,10 +700,7 @@ export function sendToGroupChat(
|
||||
attached
|
||||
)
|
||||
|
||||
const wasRunning = ($groupChats.get()[group] || {}).running === true
|
||||
updateGroupChat(group, (room: GroupChatRoom) => {
|
||||
room.epoch = (room.epoch || 0) + 1
|
||||
room.running = true
|
||||
// #93129: user text is the ONLY input that changes member holds. An
|
||||
// explicit "stop @member" sets a sticky hold; "@member resume" (or
|
||||
// @all resume, or any direct non-stop mention of the held member)
|
||||
@@ -724,36 +725,59 @@ export function sendToGroupChat(
|
||||
thread: target
|
||||
})
|
||||
|
||||
const binding = followGroupChat(group, name => {
|
||||
group = name
|
||||
})
|
||||
|
||||
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
|
||||
})
|
||||
}
|
||||
})
|
||||
.finally(binding.dispose)
|
||||
}
|
||||
|
||||
if (!wasRunning) {
|
||||
drive()
|
||||
} else {
|
||||
// Preserve the existing newer-send handoff delay, without pinning its name.
|
||||
setTimeout(drive, 250)
|
||||
}
|
||||
queueGroupChatDrive(group, members, target)
|
||||
|
||||
return target
|
||||
}
|
||||
|
||||
interface GroupChatDrive {
|
||||
pending: Map<string, GroupMember[]>
|
||||
binding: ReturnType<typeof followGroupChat>
|
||||
}
|
||||
|
||||
// Keep the owner until its awaited member releases, even after Stop. A
|
||||
// rename follows the room identity; disband retires the binding permanently.
|
||||
const groupChatDrives = new Map<string, GroupChatDrive>()
|
||||
|
||||
function queueGroupChatDrive(group: string, members: GroupMember[], thread: string) {
|
||||
let key = groupChatRoomKey(group, $groupChats.get()[group])
|
||||
const active = groupChatDrives.get(key)
|
||||
|
||||
if (active?.binding.isLive()) {
|
||||
active.pending.set(thread, members)
|
||||
|
||||
return
|
||||
}
|
||||
|
||||
const binding = followGroupChat(group, name => {
|
||||
groupChatDrives.delete(key)
|
||||
group = name
|
||||
key = groupChatRoomKey(group, $groupChats.get()[group])
|
||||
groupChatDrives.set(key, drive)
|
||||
})
|
||||
|
||||
const drive: GroupChatDrive = { pending: new Map([[thread, members]]), binding }
|
||||
groupChatDrives.set(key, drive)
|
||||
|
||||
void (async () => {
|
||||
try {
|
||||
while (binding.isLive() && drive.pending.size) {
|
||||
const [nextThread, nextMembers] = drive.pending.entries().next().value!
|
||||
drive.pending.delete(nextThread)
|
||||
updateGroupChat(group, room => ({ ...room, epoch: (room.epoch || 0) + 1, running: true }))
|
||||
await runGroupChatRounds(group, nextMembers, nextThread)
|
||||
}
|
||||
} catch {
|
||||
if (binding.isLive()) {
|
||||
recordGroupActivity(group, { kind: 'failed', member: null, thread })
|
||||
updateGroupChat(group, room => ({ ...room, running: false, turn: null }))
|
||||
}
|
||||
} finally {
|
||||
binding.dispose()
|
||||
|
||||
if (groupChatDrives.get(key) === drive) {
|
||||
groupChatDrives.delete(key)
|
||||
}
|
||||
}
|
||||
})()
|
||||
}
|
||||
|
||||
@@ -432,7 +432,7 @@ async function submitGroupTurnPrompt(
|
||||
// timeout alone silently dropped long real turns: a 7-minute research run
|
||||
// timed out at 3 minutes, read as a pass, and its finished result never
|
||||
// reached the room (db's Aug 2026 report).
|
||||
const GROUP_TURN_HARD_CAP_MS = 20 * 60000
|
||||
export const GROUP_TURN_HARD_CAP_MS = 20 * 60000
|
||||
|
||||
/** Mirror a member's pending prompt — clarify question OR command approval —
|
||||
* from its resume snapshot into the room store, keyed
|
||||
@@ -1069,7 +1069,11 @@ export async function harvestStrandedGroupReply(group: string, member: GroupMemb
|
||||
strandedThread
|
||||
)
|
||||
updateGroupChat(group, (r: GroupChatRoom) => {
|
||||
r.watermarks[`${strandedThread}::${memberKey}`] = r.log.length
|
||||
const markKey = `${strandedThread}::${memberKey}`
|
||||
|
||||
if (r.watermarks[markKey] === r.log.length - 1) {
|
||||
r.watermarks[markKey] = r.log.length
|
||||
}
|
||||
|
||||
return r
|
||||
})
|
||||
|
||||
2
package-lock.json
generated
2
package-lock.json
generated
@@ -31,7 +31,7 @@
|
||||
},
|
||||
"apps/bootstrap-installer": {
|
||||
"name": "@hermes/bootstrap-installer",
|
||||
"version": "0.0.1",
|
||||
"version": "0.21.1",
|
||||
"dependencies": {
|
||||
"@nous-research/ui": "0.18.2",
|
||||
"@tailwindcss/typography": "0.5.20",
|
||||
|
||||
@@ -86,6 +86,16 @@ Routines are plain [Hermes cron jobs](./features/cron.md) namespaced `[bot:<name
|
||||
|
||||
## Groups and group chats
|
||||
|
||||
Messages sent during a member turn queue behind the active room drive, including
|
||||
replies in the same thread. They do not interrupt that member or mark unseen
|
||||
messages as read. Stop cancels queued continuations and holds the members until
|
||||
resumed. A quiet room watches timed-out members for another 21 minutes after the
|
||||
foreground wait ends; this observation window does not extend the turn itself.
|
||||
Unresolved member failures remain visible in the collapsed Activity summary after
|
||||
the room settles. Expand Activity for the turn sequence; re-address the member to
|
||||
try again. An ambiguous submit failure is not automatically resubmitted.
|
||||
|
||||
|
||||
Right-click a local Bot → **Manage groups** to add or remove it from any number of group chats. Pick existing groups independently or create one inline. Local membership is stored in the Bot's backend-synced profile metadata, so it follows that profile across desktops; older profiles with one legacy group continue to work. Connections Bots join through the New Group Chat picker and remain source-qualified in the room's shared state.
|
||||
|
||||
**Rooms follow your gateways, not one Desktop.** Each room's recent transcript, members, picture, and name are mirrored into the shared profile metadata of **every** gateway your Desktop is connected to, with per-gateway versioning so two Desktops writing at once merge instead of overwriting each other. Open Hermes Desktop on another machine against the same gateway (local network, Tailscale, anywhere) and the room appears with its history; gateway-only clients see it too. Rooms carry a durable internal identity, so renaming one changes just its display name everywhere, disbanding one removes it permanently on every client — even ones that were offline at the time — and recreating a same-name group starts a genuinely fresh room. If a gateway dies or is removed, nothing is lost: every connected Desktop keeps the full room locally and re-seeds any gateway it reconnects to. (The full orchestration log stays in each Desktop's local storage; the shared mirror is a bounded recent-history projection.)
|
||||
|
||||
Reference in New Issue
Block a user