|
|
|
@@ -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<GroupTurnOutcome> {
|
|
|
|
|
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<null | number> => {
|
|
|
|
|
// 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
|
|
|
|
|
})
|
|
|
|
|