From 5d4aa4fcb23ff2cd35f2658322a0d459e4f0f0ae Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 02:42:33 -0700 Subject: [PATCH] fix(desktop): group chat rooms are serial again; keep only the push-woken turn poll MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit #101112 made round members take their turns concurrently. That changed what a group chat IS: later speakers in a round no longer saw earlier speakers' replies, so bots answered the user independently instead of building on each other. Group rooms are serial round-robin by design — this restores the pre-#101112 round engine (group-rounds.ts, group-chat.ts, group-chat-view.tsx, their tests, and the docs) byte-for-byte. What stays from #101112: the per-turn poll wakes on the member session's terminal frame (message.complete / error via host.onEvent) instead of sleeping a fixed 2s between session.resume reads; 5s timer kept as backstop. That is a pure latency fix with no change to room semantics. Live A/B (real tui_gateway over WS, 4 members, one serial round): 2s poll 32.5s -> push-woken 22.5s. The remaining time is model latency. Refs #92760 --- .../plugins/hermes-bots/group-chat-view.tsx | 8 +- .../src/plugins/hermes-bots/group-chat.ts | 5 +- .../plugins/hermes-bots/group-rounds.test.ts | 49 +- .../src/plugins/hermes-bots/group-rounds.ts | 656 +++++++++--------- website/docs/user-guide/bot-mode.md | 3 +- 5 files changed, 339 insertions(+), 382 deletions(-) diff --git a/apps/desktop/src/plugins/hermes-bots/group-chat-view.tsx b/apps/desktop/src/plugins/hermes-bots/group-chat-view.tsx index 65b0858573..e9749dc3c3 100644 --- a/apps/desktop/src/plugins/hermes-bots/group-chat-view.tsx +++ b/apps/desktop/src/plugins/hermes-bots/group-chat-view.tsx @@ -1256,12 +1256,8 @@ export function GroupChatWorkspace({ group, members, onBack, visible = true }: G
{roomClarifies.length ? b.group.waitingForAnswer - : Object.keys(room.turns || {}).length - ? b.group.memberThinking( - Object.values(room.turns || {}) - .map(name => groupSpeakerLabel(name)) - .join(', ') - ) + : room.turn + ? b.group.memberThinking(groupSpeakerLabel(room.turn)) : b.group.roomWorking}
) : null} diff --git a/apps/desktop/src/plugins/hermes-bots/group-chat.ts b/apps/desktop/src/plugins/hermes-bots/group-chat.ts index 84015a5c50..9949c57824 100644 --- a/apps/desktop/src/plugins/hermes-bots/group-chat.ts +++ b/apps/desktop/src/plugins/hermes-bots/group-chat.ts @@ -1375,13 +1375,12 @@ export interface GroupHoldStamp extends GroupHold { } /** The room record as the coordination engine handles it: `GroupChat` plus - * `turns`, the runtime-only memberKey → name map of members currently - * mid-turn (several at once — a round's turns run concurrently). Like + * `turn`, the runtime-only name of the member currently mid-turn. Like * `running`/`epoch` it never persists, so it has no place in the durable * shape. Holds carry the fuller live stamp. */ export interface GroupChatRoom extends GroupChat { holds?: Record - turns?: Record + turn?: null | string } /** Set or clear a group chat's room picture (small data URL, normalized by diff --git a/apps/desktop/src/plugins/hermes-bots/group-rounds.test.ts b/apps/desktop/src/plugins/hermes-bots/group-rounds.test.ts index fbdc740b2e..fd9ec778bf 100644 --- a/apps/desktop/src/plugins/hermes-bots/group-rounds.test.ts +++ b/apps/desktop/src/plugins/hermes-bots/group-rounds.test.ts @@ -236,49 +236,6 @@ describe('round lifecycle', () => { }) }) -describe('concurrent rounds', () => { - it('runs every responder of a round at once, so a round takes as long as its slowest member', async () => { - // Each turn holds until ALL three members have submitted: under the old - // serial loop the first member's turn could never finish (nobody else - // submits until it does) and this would deadlock at the drain bound. - let submitted = 0 - const release: Array<() => void> = [] - - const room = await loadRoom({ - turn: ({ profile, prompt }) => - new Promise(resolve => { - submitted += 1 - // Round 1 (the user's question is in the delta): speak. Later - // rounds (only sibling replies in the delta): pass. - release.push(() => resolve(prompt.includes('status?') ? `${profile} here` : '(pass)')) - - if (submitted % 3 === 0) { - for (const fn of release.splice(0)) { - fn() - } - } - }) - }) - - room.rounds.sendToGroupChat('Fast', MEMBERS, 'everyone, status?') - await settle(room, 'Fast') - - const replies = log(room, 'Fast').filter(entry => entry.from.kind === 'member') - - expect(replies.map(entry => entry.from.name).sort()).toEqual(['builder', 'ops', 'research']) - // Round 2 delivers every sibling's round-1 reply to each member exactly - // once — concurrent commits never eat or duplicate each other's deltas — - // and never echoes a member's own reply back to it. - expect(room.gateway.calls).toHaveLength(6) - - for (const call of room.gateway.calls.slice(3)) { - for (const name of MEMBERS.map(member => member.name)) { - expect(call.prompt.split(`${name} here`)).toHaveLength(name === call.profile ? 1 : 2) - } - } - }) -}) - describe('per-member delta', () => { it('feeds a second send only the NEW messages', async () => { const room = await loadRoom() @@ -769,7 +726,7 @@ describe('stopGroupThread (#91868/#94569)', () => { members: STOP_MEMBERS, running: true, sessions: { alpha: 'live-alpha-sid' }, - turns: turn ? { [turn]: turn } : {}, + turn, watermarks: {} } } as unknown as Record) @@ -785,7 +742,7 @@ describe('stopGroupThread (#91868/#94569)', () => { expect(state.epoch).toBe(4) expect(state.running).toBe(false) - expect(state.turns).toEqual({}) + expect(state.turn).toBeNull() for (const member of STOP_MEMBERS) { expect(state.holds?.[member.name]).toBeTruthy() @@ -801,7 +758,7 @@ describe('stopGroupThread (#91868/#94569)', () => { const interrupts = room.gateway.rpcFor('session.interrupt') - // Exactly one — only alpha is mid-turn in the seeded room. + // Exactly one — the serial loop has one member in flight. expect(interrupts).toHaveLength(1) expect(interrupts[0].params.session_id).toBe('live-alpha-sid') }) diff --git a/apps/desktop/src/plugins/hermes-bots/group-rounds.ts b/apps/desktop/src/plugins/hermes-bots/group-rounds.ts index 3758c1625a..ff65ff2745 100644 --- a/apps/desktop/src/plugins/hermes-bots/group-rounds.ts +++ b/apps/desktop/src/plugins/hermes-bots/group-rounds.ts @@ -415,19 +415,18 @@ export function unaddressedGroupMentions(group: string, members: GroupMember[], * 2. Sets a #93129 hold for EVERY member — future turns stay skipped until * the user explicitly releases (resume / @all resume / direct mention), * the exact contract user-typed "@all stop" already has. - * 3. Sends session.interrupt to every member currently ON TURN - * (room.turns, runtime-only — a round's turns run concurrently) via its - * own route, so the in-flight model calls actually die instead of - * grinding to completion in the background. Best-effort: an unreachable - * member still leaves the room stopped — the poll loop's staleness check - * (epoch moved AND member held) abandons the turn. + * 3. Sends session.interrupt to the member currently ON TURN (room.turn, + * runtime-only) via its own route, so the in-flight model call actually + * dies instead of grinding to completion in the background. Best-effort: + * an unreachable member still leaves the room stopped — the poll loop's + * staleness check (epoch moved AND member held) abandons the turn. * * `members` is the live roster when the caller has one (the workspace); * falls back to the room's durable roster so a two-arg call still works. */ export async function stopGroupThread(group: string, thread: null | string, members: GroupMember[] | null = null) { const room = $groupChats.get()[group] || {} const roster = Array.isArray(members) && members.length ? members : room.members || [] - const onTurnKeys = Object.keys(room.turns || {}) + const turnName = room.turn || null const stamp: GroupHoldStamp = { at: Date.now(), @@ -438,7 +437,7 @@ export async function stopGroupThread(group: string, thread: null | string, memb updateGroupChat(group, (r: GroupChatRoom) => { r.epoch = (r.epoch || 0) + 1 r.running = false - r.turns = {} + r.turn = null // Same hold shape applyGroupHoldDirective mints for "@all stop" — the // held-skip path (watermark consume + 'held' activity note) and every @@ -471,262 +470,28 @@ export async function stopGroupThread(group: string, thread: null | string, memb thread: thread || null }) - // Interrupt every member actually mid-turn. room.turns is runtime-only; - // a settled room has none. - await Promise.all( - roster - .filter((member: GroupMember) => onTurnKeys.includes(groupMemberKey(member))) - .map(async (onTurn: GroupMember) => { - const sessionId = (room.sessions || {})[groupMemberKey(onTurn)] + // Interrupt the member actually mid-turn. room.turn is runtime-only and + // names exactly one member (the loop is serial); a settled room has none. + const onTurn = turnName ? roster.find((member: GroupMember) => member?.name === turnName) : null + const sessionId = onTurn ? (room.sessions || {})[groupMemberKey(onTurn)] : null - if (!sessionId) { - return - } - - try { - await requestForBot(onTurn, 'session.interrupt', { - session_id: sessionId - }) - } catch { - /* best-effort — the epoch/hold legs above already stopped the room; - the abandoned poll loop exits on its staleness check */ - } + if (onTurn && sessionId) { + try { + await requestForBot(onTurn, 'session.interrupt', { + session_id: sessionId }) - ) + } catch { + /* best-effort — the epoch/hold legs above already stopped the room; + the abandoned poll loop exits on its staleness check */ + } + } } -/** Why one member's turn ended without a committed reply. */ -type GroupTurnOutcome = 'cancelled' | 'passed' | 'skipped' | 'spoke' - -/** One member's full turn against the room: compute its unseen delta, skip - * (held / nothing new), run the turn, then commit the result under the - * #93127 staleness check. Pure with respect to its siblings — several - * members' turns run CONCURRENTLY within a round (they share nothing but - * the room log, which appends are serialized through the atom), and each - * member still sees only what was in the room when its turn started. */ -async function takeMemberTurn( - group: string, - members: GroupMember[], - member: GroupMember, - thread: string, - startEpoch: number, - attachImages: boolean -): Promise { - const room = $groupChats.get()[group] || { - log: [], - watermarks: {} - } - - const memberKey = groupMemberKey(member) - const markKey = `${thread}::${memberKey}` - const seen = room.watermarks[markKey] || 0 - - // Delta: NEW room entries, narrowed to this thread — the member's turn sees - // only the conversation it's part of — minus its OWN replies. Those already - // live in its session as assistant messages; echoing them back costs a - // turn that can only pass. (Concurrent rounds make this matter: a member's - // reply lands beside its siblings', so an index watermark can't cleanly - // step over "just mine".) - const delta = room.log - .slice(seen) - .filter((e: GroupMessage) => groupThreadOf(e) === thread && !isOwnGroupEntry(e, member)) - - if (!delta.length) { - return 'skipped' - } - - // #93129: a member the user told to stop is HELD — no turn until an - // explicit release (resume / @all resume / a direct non-stop mention). - // Consume the delta exactly once (watermark past the current log) so the - // same entries never re-trigger this skip, and surface WHY the bot is - // silent in the activity feed the first time. - const heldEntry = (room.holds || {})[memberKey] - - if (heldEntry) { - const advance = heldMemberWatermarkAdvance(seen, room.log.length) - updateGroupChat(group, (r: GroupChatRoom) => { - if (advance !== null) { - r.watermarks[markKey] = advance - } - - if (r.holds?.[memberKey] && !r.holds[memberKey].noted) { - r.holds = { - ...r.holds, - [memberKey]: { - ...r.holds[memberKey], - noted: true - } - } - } - - return r - }) - - if (!heldEntry.noted) { - recordGroupActivity(group, { - kind: 'held', - member: member.name, - thread - }) - } - - return 'skipped' - } - - const prompt = buildGroupChatTurnPrompt({ - groupName: group, - members, - viewer: member, - deltaLines: delta.slice(-GROUP_CHAT_HISTORY_LIMIT).map((e: GroupMessage) => formatGroupChatLine(e, member.name)) - }) - - // Images riding this delta (user attachments — member entries don't carry - // images today, but flatMap keeps this future-proof) get staged into the - // member's session so the model sees the pixels, not just the transcript's - // [attached image: …] marker. Continuation turns are text-only. - const deltaImages = attachImages - ? delta.flatMap((e: GroupMessage) => (Array.isArray(e.images) ? e.images : [])) - : undefined - - // Surface WHO is on turn (runtime-only, like running/epoch) so the room - // shows "Radar is thinking…" — several members can be mid-turn at once. - updateGroupChat(group, (r: GroupChatRoom) => { - r.turns = { - ...(r.turns || {}), - [memberKey]: member.name - } - - return r - }) - let reply: null | string = null - - try { - reply = await runGroupChatMemberTurn(group, member, prompt, thread, deltaImages) - - // Needs-attention hook (#93091 item 3): a turn that produced a real - // reply (or an explicit pass) is a good turn — clear the badge. A - // timed-out turn also returns null but never threw; leaving any prior - // badge in place there is the conservative choice. - if (reply !== null) { - clearBotAttention(memberKey) - } - } catch (error: any) { - const reason = String(error?.data?.reason || '').trim() - recordGroupActivity(group, { - kind: 'failed', - member: member.name, - thread, - ...(reason - ? { - reason - } - : {}) - }) - noteBotAttention(memberKey, reason || error?.message || error) - reply = null // a failed turn is a pass, never a room error - } finally { - updateGroupChat(group, (r: GroupChatRoom) => { - const next = { - ...(r.turns || {}) - } - - delete next[memberKey] - r.turns = next - - return r - }) - } - - // #93127: the turn may have finished AFTER a newer user send bumped the - // room epoch. That newer send's loop re-drives this member with the full - // delta, so committing this stale result (watermark advance + append) - // would double-deliver the same reply. Drop it here — BEFORE the watermark - // advance and BEFORE the append. Only a newer USER entry in THIS thread - // makes the re-drive premise true: a cross-thread send bumps the epoch - // too, but its loop filters this thread out and would never regenerate - // the finished reply. The 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. - const roomNow = $groupChats.get()[group] || { - log: [] - } - - 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. - const turnTail = anchorIdx >= 0 ? roomNow.log.slice(anchorIdx + 1) : roomNow.log - - const newerUserEntryInThread = turnTail.some( - (e: GroupMessage) => e.from?.kind === 'user' && groupThreadOf(e) === thread - ) - - if (!shouldCommitMemberTurn(startEpoch, epochNow, newerUserEntryInThread)) { - recordGroupActivity(group, { - kind: 'cancelled', - member: member.name, - thread - }) - - return 'cancelled' - } - - // The member has now seen everything up to the PRE-TURN log length — not - // the current one: a sibling's concurrent reply that landed while this - // member was thinking is genuinely unseen and must reach it next round. - // (Its own reply, appended below, is excluded from deltas by author.) - const seenThrough = room.log.length - updateGroupChat(group, (r: GroupChatRoom) => { - r.watermarks[markKey] = Math.min(seenThrough, r.log.length) - - return r - }) - - if (reply === null || isGroupPassText(reply)) { - return 'passed' - } - - appendGroupChatEntry( - group, - { - kind: 'member', - name: member.name, - ...(member.remoteSource - ? { - source: member.connectionLabel || member.connectionId - } - : {}) - }, - reply, - thread - ) - - return 'spoke' -} - -/** A room entry this member authored (same kind, name and — for - * cross-connection members — same source device). */ -function isOwnGroupEntry(entry: GroupMessage, member: GroupMember): boolean { - if (entry.from?.kind !== 'member' || entry.from.name !== member.name) { - return false - } - - const source = member.remoteSource ? member.connectionLabel || member.connectionId : undefined - - return (entry.from.source || undefined) === (source || undefined) -} - -/** Drive one bounded set of rounds for ONE THREAD. Within a round, every - * responder takes its turn CONCURRENTLY — a room of N bots answers in the - * time of the slowest one, not the sum of all N. Rounds stay serial: each - * round's prompts include the previous round's replies, so bots build on - * each other. A newer user send bumps the room epoch; this loop notices at - * the next round boundary (and every turn's commit check), bails, and the - * newest send's own loop takes over. Watermarks are per thread+member - * (`${thread}::${memberKey}`), so parallel topics never eat each other's - * deltas. */ +/** 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. + * 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 startEpoch = ($groupChats.get()[group] || {}).epoch || 0 const isCurrent = () => (($groupChats.get()[group] || {}).epoch || 0) === startEpoch @@ -737,53 +502,23 @@ export async function runGroupChatRounds(group: string, members: GroupMember[], // cap forced the exit — the activity feed must tell those apart. let exitKind: 'capped' | 'settled' = 'settled' - /** Run one set of members concurrently. Returns how many spoke, or null - * when a newer send superseded this drive mid-round. */ - const runConcurrentTurns = async (responders: GroupMember[], attachImages: boolean): Promise => { - // The message cap is enforced per ROUND: a round admits at most the - // remaining budget worth of speakers, and its concurrent turns can't - // overshoot it by more than the round size. - const budget = GROUP_CHAT_MAX_MESSAGES - posted - - if (budget <= 0) { - return 0 - } - - const outcomes = await Promise.all( - responders.slice(0, budget).map(member => takeMemberTurn(group, members, member, thread, startEpoch, attachImages)) - ) - - if (!isCurrent() || outcomes.includes('cancelled')) { - return null - } - - const spoke = outcomes.filter(outcome => outcome === 'spoke').length - posted += spoke - - return spoke - } - try { for (let round = 0; round < GROUP_CHAT_MAX_ROUNDS; round++) { // Deliver any replies that finished after their turn timed out — // every member, not just this round's responders, so long work is // late, never lost. - if (!isCurrent()) { - recordGroupActivity(group, { - kind: 'cancelled', - member: null, - thread - }) + for (const member of members) { + if (!isCurrent()) { + recordGroupActivity(group, { + kind: 'cancelled', + member: null, + thread + }) - return - } + return + } - await Promise.all(members.map(member => harvestStrandedGroupReply(group, member))) - - if (posted >= GROUP_CHAT_MAX_MESSAGES) { - exitKind = 'capped' // message cap, not consensus (#94478) - - return + await harvestStrandedGroupReply(group, member) } const roomLog = (($groupChats.get()[group] || {}).log || []).filter( @@ -808,16 +543,195 @@ export async function runGroupChatRounds(group: string, members: GroupMember[], (member: GroupMember) => !Object.prototype.hasOwnProperty.call(strandedNow, groupMemberKey(member)) ) - let spokeThisRound = await runConcurrentTurns(responders, true) + let spokeThisRound = 0 - if (spokeThisRound === null) { - recordGroupActivity(group, { - kind: 'cancelled', - member: null, - thread + for (const member of responders) { + if (!isCurrent() || posted >= GROUP_CHAT_MAX_MESSAGES) { + if (!isCurrent()) { + recordGroupActivity(group, { + kind: 'cancelled', + member: null, + thread + }) + } else { + exitKind = 'capped' // message cap, not consensus (#94478) + } + + return + } + + const room = $groupChats.get()[group] || { + log: [], + watermarks: {} + } + + const memberKey = groupMemberKey(member) + const markKey = `${thread}::${memberKey}` + const seen = room.watermarks[markKey] || 0 + // Delta: NEW room entries, narrowed to this thread — the member's + // turn sees only the conversation it's part of. + const delta = room.log.slice(seen).filter((e: GroupMessage) => groupThreadOf(e) === thread) + + if (!delta.length) { + continue + } + + // #93129: a member the user told to stop is HELD — no turn until an + // explicit release (resume / @all resume / a direct non-stop + // mention). Consume the delta exactly once (watermark past the + // current log) so the same entries never re-trigger this skip, and + // surface WHY the bot is silent in the activity feed the first time. + const heldEntry = (room.holds || {})[memberKey] + + if (heldEntry) { + const advance = heldMemberWatermarkAdvance(seen, room.log.length) + updateGroupChat(group, (r: GroupChatRoom) => { + if (advance !== null) { + r.watermarks[markKey] = advance + } + + if (r.holds?.[memberKey] && !r.holds[memberKey].noted) { + r.holds = { + ...r.holds, + [memberKey]: { + ...r.holds[memberKey], + noted: true + } + } + } + + return r + }) + + if (!heldEntry.noted) { + recordGroupActivity(group, { + kind: 'held', + member: member.name, + thread + }) + } + + continue + } + + const prompt = buildGroupChatTurnPrompt({ + groupName: group, + members, + viewer: member, + deltaLines: delta + .slice(-GROUP_CHAT_HISTORY_LIMIT) + .map((e: GroupMessage) => formatGroupChatLine(e, member.name)) }) - return + // Images riding this delta (user attachments — member entries don't + // carry images today, but flatMap keeps this future-proof) get staged + // into the member's session so the model sees the pixels, not just + // the transcript's [attached image: …] marker. + const deltaImages = delta.flatMap((e: GroupMessage) => (Array.isArray(e.images) ? e.images : [])) + + // Surface WHO is on turn (runtime-only, like running/epoch) so the + // room shows "Radar is thinking…" instead of a generic working line — + // long model turns otherwise read as the room being stuck. + updateGroupChat(group, (r: GroupChatRoom) => { + r.turn = member.name + + return r + }) + let reply: null | string = null + + try { + reply = await runGroupChatMemberTurn(group, member, prompt, thread, deltaImages) + + // Needs-attention hook (#93091 item 3): a turn that produced a real + // reply (or an explicit pass) is a good turn — clear the badge. + // A timed-out turn also returns null but never threw; leaving any + // prior badge in place there is the conservative choice. + if (reply !== null) { + clearBotAttention(groupMemberKey(member)) + } + } catch (error: any) { + const reason = String(error?.data?.reason || '').trim() + recordGroupActivity(group, { + kind: 'failed', + member: member.name, + thread, + ...(reason + ? { + reason + } + : {}) + }) + noteBotAttention(groupMemberKey(member), reason || error?.message || error) + reply = null // a failed turn is a pass, never a room error + } + + // #93127: the turn may have finished AFTER a newer user send bumped + // the room epoch. That newer send's loop re-drives this member with + // the full delta, so committing this stale result (watermark advance + // + append) would double-deliver the same reply. Drop it here — + // BEFORE the watermark advance and BEFORE the append. Only a newer + // USER entry in THIS thread makes the re-drive premise true: a + // cross-thread send bumps the epoch too, but its loop filters this + // thread out and would never regenerate the finished reply. The + // 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. + const roomNow = $groupChats.get()[group] || { + log: [] + } + + 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. + const turnTail = anchorIdx >= 0 ? roomNow.log.slice(anchorIdx + 1) : roomNow.log + + const newerUserEntryInThread = turnTail.some( + (e: GroupMessage) => e.from?.kind === 'user' && groupThreadOf(e) === thread + ) + + if (!shouldCommitMemberTurn(startEpoch, epochNow, newerUserEntryInThread)) { + recordGroupActivity(group, { + kind: 'cancelled', + member: member.name, + thread + }) + + return + } + + // The member has now seen everything up to the pre-reply log length. + updateGroupChat(group, (r: GroupChatRoom) => { + r.watermarks[markKey] = r.log.length + + return r + }) + + if (reply !== null && !isGroupPassText(reply)) { + appendGroupChatEntry( + group, + { + kind: 'member', + name: member.name, + ...(member.remoteSource + ? { + source: member.connectionLabel || member.connectionId + } + : {}) + }, + reply, + thread + ) + // Its own message counts as seen too. + updateGroupChat(group, (r: GroupChatRoom) => { + r.watermarks[markKey] = r.log.length + + return r + }) + posted += 1 + spokeThisRound += 1 + } } if (spokeThisRound === 0) { @@ -836,27 +750,119 @@ export async function runGroupChatRounds(group: string, members: GroupMember[], continuations += 1 if (pendingKeys.length && continuations <= GROUP_CHAT_MAX_CONTINUATIONS) { - const strandedNow = ($groupChats.get()[group] || {}).stranded || {} + const citedMembers = members.filter((member: GroupMember) => pendingKeys.includes(groupMemberKey(member))) - const citedMembers = members.filter( - (member: GroupMember) => - pendingKeys.includes(groupMemberKey(member)) && - !Object.prototype.hasOwnProperty.call(strandedNow, groupMemberKey(member)) - ) + if (citedMembers.length && posted < GROUP_CHAT_MAX_MESSAGES) { + const strandedNow = ($groupChats.get()[group] || {}).stranded || {} - // The continuation prompt centers on what each cited member - // missed: everything since its watermark, which includes the - // reply that cites it. Holds still apply (#93129). - const continued = citedMembers.length ? await runConcurrentTurns(citedMembers, false) : 0 + const continuationResponders = citedMembers.filter( + (member: GroupMember) => !Object.prototype.hasOwnProperty.call(strandedNow, groupMemberKey(member)) + ) - if (continued === null) { - return + for (const member of continuationResponders) { + if (!isCurrent() || posted >= GROUP_CHAT_MAX_MESSAGES || continuations > GROUP_CHAT_MAX_CONTINUATIONS) { + break + } + + const room = $groupChats.get()[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) { + continue + } + + const heldEntry = (room.holds || {})[memberKey] + + if (heldEntry) { + continue // holds still apply to continuation turns (#93129) + } + + const prompt = buildGroupChatTurnPrompt({ + groupName: 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.name)) + }) + + updateGroupChat(group, (r: GroupChatRoom) => { + r.turn = member.name + + return r + }) + let continuationReply: null | string = null + + try { + continuationReply = await runGroupChatMemberTurn(group, member, prompt, thread) + + if (continuationReply !== null) { + clearBotAttention(memberKey) + } + } catch (error: any) { + recordGroupActivity(group, { + kind: 'failed', + member: member.name, + thread + }) + noteBotAttention(memberKey, error?.message || error) + continuationReply = null + } + + if (!isCurrent()) { + return + } + + updateGroupChat(group, (r: GroupChatRoom) => { + r.watermarks[markKey] = r.log.length + + return r + }) + + if (continuationReply !== null && !isGroupPassText(continuationReply)) { + appendGroupChatEntry( + group, + { + kind: 'member', + name: member.name, + ...(member.remoteSource + ? { + source: member.connectionLabel || member.connectionId + } + : {}) + }, + continuationReply, + thread + ) + updateGroupChat(group, (r: GroupChatRoom) => { + r.watermarks[markKey] = r.log.length + + return r + }) + posted += 1 + + // The continuation's own reply may cite someone else — fall + // through to the normal loop so the next round handles it via + // the same responder machinery. Reaching here means the loop + // continues rather than settling; the outer for-loop's next + // iteration re-evaluates everything. + spokeThisRound += 1 + } + } } - - // The continuation's own replies may cite someone else — fall - // through to the normal loop so the next round handles it via the - // same responder machinery. - spokeThisRound = continued } if (spokeThisRound === 0) { @@ -889,7 +895,7 @@ export async function runGroupChatRounds(group: string, members: GroupMember[], }) updateGroupChat(group, (r: GroupChatRoom) => { r.running = false - r.turns = {} + r.turn = null return r }) diff --git a/website/docs/user-guide/bot-mode.md b/website/docs/user-guide/bot-mode.md index 42ff758372..0e50502a3f 100644 --- a/website/docs/user-guide/bot-mode.md +++ b/website/docs/user-guide/bot-mode.md @@ -83,8 +83,7 @@ Groups are standalone rows in the same activity-ordered roster as Bot DMs. A Bot **Open chat** on any group row (2–6 Bots) opens a shared room where the whole group coordinates: -- Your message triggers up to **three rounds** of member turns. Within a round every responding Bot thinks **at the same time**, so a room of five answers about as fast as one; rounds run in sequence so each Bot sees what the others just said before it speaks again. @-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. -- Replies land the moment a Bot finishes — the room listens for each member session's completion event rather than polling on a timer (a slow 5s poll remains as a backstop for older gateways). +- 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. - Hard caps (10 messages per send, 3 rounds) keep rooms from spinning. - Each member keeps its own persistent `Group: ` session, so room context survives like any other conversation.