From 3b83b5e1fa3ed8a40aa4b95a7789a7743d236084 Mon Sep 17 00:00:00 2001 From: TheTom Date: Tue, 18 Aug 2026 14:10:34 -0500 Subject: [PATCH] fix(gateway): serialize shared bot room updates --- .../desktop/src/plugins/hermes-bots/plugin.js | 304 ++++++++++++++---- .../hermes-bots/tests/group-chat.test.mjs | 240 +++++++++++++- .../tui_gateway/test_profiles_ui_meta_cas.py | 98 ++++++ tui_gateway/methods_profiles.py | 105 ++++-- tui_gateway/server.py | 4 + 5 files changed, 646 insertions(+), 105 deletions(-) create mode 100644 tests/tui_gateway/test_profiles_ui_meta_cas.py diff --git a/apps/desktop/src/plugins/hermes-bots/plugin.js b/apps/desktop/src/plugins/hermes-bots/plugin.js index bfeb792d61..2f614655cb 100644 --- a/apps/desktop/src/plugins/hermes-bots/plugin.js +++ b/apps/desktop/src/plugins/hermes-bots/plugin.js @@ -313,6 +313,7 @@ const GROUP_CHAT_SYNC_META_KEY = 'hermes-bots-groups' const GROUP_CHAT_SYNC_MAX_BYTES = 48000 const GROUP_CHAT_SYNC_MESSAGES = 16 const GROUP_CHAT_SYNC_TEXT_CHARS = 1200 +const GROUP_CHAT_SYNC_IMAGE_CHARS = 24000 let groupChatSyncTimer = null let groupChatSyncRetryTimer = null let groupChatSyncPending = null @@ -364,7 +365,7 @@ function groupChatSyncSnapshot(all = $groupChats.get(), deleted = {}) { .slice(0, 64) ) const envelope = { - version: 1, + version: 2, updatedAt: Date.now(), rooms, ...(Object.keys(boundedDeleted).length ? { deleted: boundedDeleted } : {}) @@ -372,6 +373,7 @@ function groupChatSyncSnapshot(all = $groupChats.get(), deleted = {}) { for (const [name, room] of ranked) { const log = room.log.slice(-GROUP_CHAT_SYNC_MESSAGES).map(entry => ({ + ...(entry?.id ? { id: String(entry.id).slice(0, 160) } : {}), from: { kind: entry?.from?.kind === 'member' ? 'member' : 'user', name: String(entry?.from?.name || (entry?.from?.kind === 'member' ? 'Bot' : 'You')).slice(0, 128), @@ -383,6 +385,7 @@ function groupChatSyncSnapshot(all = $groupChats.get(), deleted = {}) { })) const compact = { log, + revision: Math.max(0, Number(room?.syncRevision ?? room?.revision ?? 0)), members: (Array.isArray(room.members) ? room.members : []).slice(0, GROUP_CHAT_MAX_MEMBERS).map(member => ({ name: String(member?.name || '').slice(0, 128), ...(member?.handle ? { handle: String(member.handle).slice(0, 128) } : {}), @@ -390,13 +393,19 @@ function groupChatSyncSnapshot(all = $groupChats.get(), deleted = {}) { ...(member?.connectionKind ? { connectionKind: String(member.connectionKind).slice(0, 64) } : {}), ...(member?.connectionLabel ? { connectionLabel: String(member.connectionLabel).slice(0, 128) } : {}), ...(member?.sourceScoped ? { sourceScoped: true } : {}) - })) + })), + ...(typeof room?.image === 'string' && room.image.length <= GROUP_CHAT_SYNC_IMAGE_CHARS + ? { image: room.image } + : {}) } rooms[name] = compact while (compact.log.length > 1 && groupChatGatewayJsonSize(envelope) > GROUP_CHAT_SYNC_MAX_BYTES) { compact.log.shift() } + if (compact.image && groupChatGatewayJsonSize(envelope) > GROUP_CHAT_SYNC_MAX_BYTES) { + delete compact.image + } if (groupChatGatewayJsonSize(envelope) > GROUP_CHAT_SYNC_MAX_BYTES) { delete rooms[name] } @@ -406,6 +415,9 @@ function groupChatSyncSnapshot(all = $groupChats.get(), deleted = {}) { } function groupChatSyncEntryKey(entry) { + if (entry?.id) { + return `id:${String(entry.id)}` + } return JSON.stringify([ Number(entry?.at || 0), String(entry?.from?.kind || ''), @@ -426,47 +438,76 @@ function groupChatSyncMemberKey(member) { ]) } +function groupChatSyncDeletedRevision(source, value) { + return Number(source?.version || 0) >= 2 ? Math.max(0, Number(value || 0)) : 0 +} + /** Merge two bounded projections without treating an absent room/message as - * deletion. Explicit tombstones win until a genuinely newer room message - * recreates the same name. */ -function mergeGroupChatSyncSnapshots(remote, local) { + * deletion. Gateway revisions order room identity/membership/picture and + * tombstones; stable message ids make concurrent log union idempotent. */ +function mergeGroupChatSyncSnapshots( + remote, + local, + { changedRooms = [], deletedRooms = [], writeRevision = 0 } = {} +) { const rooms = {} const deleted = {} + const changed = new Set(changedRooms) for (const source of [remote, local]) { for (const [name, at] of Object.entries(source?.deleted || {})) { - deleted[name] = Math.max(Number(deleted[name] || 0), Number(at || 0)) + deleted[name] = Math.max(Number(deleted[name] || 0), groupChatSyncDeletedRevision(source, at)) + } + } + for (const name of deletedRooms) { + deleted[name] = Math.max(Number(deleted[name] || 0), Number(writeRevision || 0)) + } + + const roomNames = new Set([...Object.keys(remote?.rooms || {}), ...Object.keys(local?.rooms || {})]) + for (const name of roomNames) { + const remoteRoom = remote?.rooms?.[name] + const localRoom = local?.rooms?.[name] + if ((!remoteRoom || !Array.isArray(remoteRoom.log)) && (!localRoom || !Array.isArray(localRoom.log))) { + continue + } + const remoteRevision = Math.max(0, Number(remoteRoom?.revision || 0)) + const localRevision = changed.has(name) + ? Math.max(0, Number(writeRevision || 0)) + : Math.max(0, Number(localRoom?.revision || 0)) + const entries = new Map() + for (const entry of [...(remoteRoom?.log || []), ...(localRoom?.log || [])]) { + entries.set(groupChatSyncEntryKey(entry), entry) + } + + let members + let image + if (localRevision > remoteRevision) { + members = [...(localRoom?.members || [])] + image = localRoom?.image + } else if (remoteRevision > localRevision) { + members = [...(remoteRoom?.members || [])] + image = remoteRoom?.image + } else { + const byId = new Map() + for (const member of [...(remoteRoom?.members || []), ...(localRoom?.members || [])]) { + byId.set(groupChatSyncMemberKey(member), member) + } + members = [...byId.values()] + image = Object.prototype.hasOwnProperty.call(localRoom || {}, 'image') ? localRoom.image : remoteRoom?.image + } + rooms[name] = { + log: [...entries.values()].sort((left, right) => { + const byTime = Number(left?.at || 0) - Number(right?.at || 0) + return byTime || groupChatSyncEntryKey(left).localeCompare(groupChatSyncEntryKey(right)) + }), + members, + revision: Math.max(remoteRevision, localRevision), + ...(typeof image === 'string' && image ? { image } : {}) } } - for (const source of [remote, local]) { - for (const [name, room] of Object.entries(source?.rooms || {})) { - if (!room || !Array.isArray(room.log)) { - continue - } - const current = rooms[name] || { log: [], members: [] } - const entries = new Map(current.log.map(entry => [groupChatSyncEntryKey(entry), entry])) - const members = new Map(current.members.map(member => [groupChatSyncMemberKey(member), member])) - - for (const entry of room.log) { - entries.set(groupChatSyncEntryKey(entry), entry) - } - for (const member of Array.isArray(room.members) ? room.members : []) { - members.set(groupChatSyncMemberKey(member), member) - } - rooms[name] = { - log: [...entries.values()].sort((left, right) => { - const byTime = Number(left?.at || 0) - Number(right?.at || 0) - return byTime || groupChatSyncEntryKey(left).localeCompare(groupChatSyncEntryKey(right)) - }), - members: [...members.values()] - } - } - } - - for (const [name, deletedAt] of Object.entries(deleted)) { - const latestMessageAt = Math.max(0, ...(rooms[name]?.log || []).map(entry => Number(entry?.at || 0))) - if (Number(deletedAt || 0) >= latestMessageAt) { + for (const [name, deletedRevision] of Object.entries(deleted)) { + if (Number(deletedRevision || 0) >= Number(rooms[name]?.revision || 0)) { delete rooms[name] } else { delete deleted[name] @@ -480,14 +521,26 @@ function mergeGroupChatSyncSnapshots(remote, local) { * state without discarding local session/watermark/runtime fields. Missing * remote rooms/messages are not deletions; only explicit tombstones remove a * room, and a genuinely newer local message wins over a stale tombstone. */ -function mergeRemoteGroupChatSnapshotIntoRooms(remote, current = $groupChats.get()) { +function mergeRemoteGroupChatSnapshotIntoRooms( + remote, + current = $groupChats.get(), + { preserveRooms = [], deletedRooms = [] } = {} +) { const rooms = { ...(current || {}) } + const preserved = new Set(preserveRooms) + const locallyDeleted = new Set(deletedRooms) for (const [name, projected] of Object.entries(remote?.rooms || {})) { if (!projected || !Array.isArray(projected.log)) { continue } + if (locallyDeleted.has(name)) { + delete rooms[name] + continue + } const existing = rooms[name] || {} + const remoteRevision = Math.max(0, Number(projected.revision || 0)) + const localRevision = Math.max(0, Number(existing.syncRevision || 0)) const entries = new Map( (Array.isArray(existing.log) ? existing.log : []).map(entry => [groupChatSyncEntryKey(entry), entry]) ) @@ -498,8 +551,13 @@ function mergeRemoteGroupChatSnapshotIntoRooms(remote, current = $groupChats.get for (const entry of projected.log) { entries.set(groupChatSyncEntryKey(entry), entry) } - for (const member of Array.isArray(projected.members) ? projected.members : []) { - members.set(groupChatSyncMemberKey(member), { ...member, remoteSource: true }) + if (!preserved.has(name)) { + if (remoteRevision > localRevision) { + members.clear() + } + for (const member of Array.isArray(projected.members) ? projected.members : []) { + members.set(groupChatSyncMemberKey(member), { ...member, remoteSource: true }) + } } const log = assignLegacyThreads( @@ -517,17 +575,29 @@ function mergeRemoteGroupChatSnapshotIntoRooms(remote, current = $groupChats.get sessions: existing.sessions && typeof existing.sessions === 'object' ? existing.sessions : {}, stranded: existing.stranded && typeof existing.stranded === 'object' ? existing.stranded : {}, members: [...members.values()], + image: preserved.has(name) + ? existing.image || null + : remoteRevision >= localRevision && Object.prototype.hasOwnProperty.call(projected, 'image') + ? projected.image || null + : existing.image || null, + syncRevision: preserved.has(name) ? localRevision : Math.max(remoteRevision, localRevision), epoch: Number(existing.epoch || 0), running: Boolean(existing.running) } } for (const [name, deletedAt] of Object.entries(remote?.deleted || {})) { - const latestMessageAt = Math.max(0, ...(rooms[name]?.log || []).map(entry => Number(entry?.at || 0))) - if (Number(deletedAt || 0) >= latestMessageAt) { + if (preserved.has(name)) { + continue + } + const deletedRevision = groupChatSyncDeletedRevision(remote, deletedAt) + if (deletedRevision >= Number(rooms[name]?.syncRevision || 0)) { delete rooms[name] } } + for (const name of locallyDeleted) { + delete rooms[name] + } return rooms } @@ -544,7 +614,9 @@ function durableGroupChatRooms(all = $groupChats.get()) { watermarks: room.watermarks || {}, sessions: room.sessions || {}, stranded: room.stranded || {}, - members: Array.isArray(room.members) ? room.members : [] + members: Array.isArray(room.members) ? room.members : [], + image: room.image || null, + syncRevision: Math.max(0, Number(room.syncRevision || 0)) } } @@ -590,18 +662,26 @@ async function groupChatRemoteSnapshot(job) { const result = await groupChatSyncRequest(job, 'profiles.list', { include_sessions: false }) const profile = (Array.isArray(result?.profiles) ? result.profiles : []).find(row => row?.name === 'default') const snapshot = profile?.ui_meta?.[GROUP_CHAT_SYNC_META_KEY] - return snapshot && typeof snapshot === 'object' && !Array.isArray(snapshot) ? snapshot : null + const supportsCas = Boolean(profile && Object.prototype.hasOwnProperty.call(profile, 'ui_meta_revisions')) + return { + snapshot: snapshot && typeof snapshot === 'object' && !Array.isArray(snapshot) ? snapshot : null, + revision: Math.max(0, Number(profile?.ui_meta_revisions?.[GROUP_CHAT_SYNC_META_KEY] || 0)), + supportsCas + } } /** Pull the shared room projection into this Desktop before it publishes any * local state. This is the receive half of the client-only sync contract. */ async function pullGroupChatServerState(connectionId = groupChatSyncConnectionId()) { - const remote = await groupChatRemoteSnapshot({ connectionId }) + const { snapshot: remote } = await groupChatRemoteSnapshot({ connectionId }) if (!remote) { return false } - const merged = mergeRemoteGroupChatSnapshotIntoRooms(remote, $groupChats.get()) + const merged = mergeRemoteGroupChatSnapshotIntoRooms(remote, $groupChats.get(), { + preserveRooms: groupChatSyncPending?.changedRooms || [], + deletedRooms: groupChatSyncPending?.deletedRooms || [] + }) $groupChats.set(merged) await persistGroupChatRooms(merged) return true @@ -611,6 +691,25 @@ function groupChatSyncBackoff() { return Math.min(30000, 1000 * 2 ** Math.min(groupChatSyncRetryCount, 5)) } +function mergeGroupChatSyncJobs(existing, incoming) { + if (!existing || existing.connectionId !== incoming.connectionId) { + return incoming + } + return { + connectionId: incoming.connectionId, + allowEmpty: Boolean(existing.allowEmpty || incoming.allowEmpty), + changedRooms: [...new Set([...(existing.changedRooms || []), ...(incoming.changedRooms || [])])], + deletedRooms: [...new Set([...(existing.deletedRooms || []), ...(incoming.deletedRooms || [])])] + } +} + +function groupChatSyncPayloadEqual(left, right) { + return ( + JSON.stringify(left?.rooms || {}) === JSON.stringify(right?.rooms || {}) && + JSON.stringify(left?.deleted || {}) === JSON.stringify(right?.deleted || {}) + ) +} + async function flushGroupChatServerSync() { if (groupChatSyncDisposed || groupChatSyncInFlight || !groupChatSyncPending) { return @@ -620,35 +719,66 @@ async function flushGroupChatServerSync() { groupChatSyncInFlight = true try { - const remote = await groupChatRemoteSnapshot(job) - if (remote) { - const mergedRooms = mergeRemoteGroupChatSnapshotIntoRooms(remote, $groupChats.get()) - $groupChats.set(mergedRooms) - await persistGroupChatRooms(mergedRooms) - } - const snapshot = mergeGroupChatSyncSnapshots(remote, job.snapshot) - const result = await groupChatSyncRequest(job, 'profiles.configure', { - name: 'default', - ui_meta: { [GROUP_CHAT_SYNC_META_KEY]: snapshot } + const remoteState = await groupChatRemoteSnapshot(job) + const local = groupChatSyncSnapshot($groupChats.get()) + const writeRevision = remoteState.revision + 1 + const snapshot = mergeGroupChatSyncSnapshots(remoteState.snapshot, local, { + changedRooms: job.changedRooms, + deletedRooms: job.deletedRooms, + writeRevision }) - if (result?.applied && result.applied.ui_meta !== true) { - throw new Error('Gateway rejected group chat ui_meta') + // Reconnect/startup reconciliation often discovers that the gateway + // already holds the exact merged projection. Avoid advancing a revision + // merely because a view reopened. + if (!(job.changedRooms || []).length && !(job.deletedRooms || []).length && groupChatSyncPayloadEqual(snapshot, remoteState.snapshot)) { + if (remoteState.snapshot) { + const mergedRooms = mergeRemoteGroupChatSnapshotIntoRooms(remoteState.snapshot, $groupChats.get(), { + preserveRooms: groupChatSyncPending?.changedRooms || [], + deletedRooms: groupChatSyncPending?.deletedRooms || [] + }) + $groupChats.set(mergedRooms) + await persistGroupChatRooms(mergedRooms) + } + groupChatSyncRetryCount = 0 + return } - const confirmed = await groupChatRemoteSnapshot(job) - if (confirmed && Number(confirmed.updatedAt || 0) !== Number(snapshot.updatedAt || 0)) { - throw new Error('Group chat ui_meta changed before read-back') + const configureParams = { + name: 'default', + ui_meta: { [GROUP_CHAT_SYNC_META_KEY]: snapshot } } - if (confirmed) { - const mergedRooms = mergeRemoteGroupChatSnapshotIntoRooms(confirmed, $groupChats.get()) + if (remoteState.supportsCas) { + configureParams.ui_meta_expected_revisions = { [GROUP_CHAT_SYNC_META_KEY]: remoteState.revision } + } + const result = await groupChatSyncRequest(job, 'profiles.configure', configureParams) + + if (result?.applied?.ui_meta !== true) { + throw new Error('Gateway rejected group chat ui_meta') + } + if ( + remoteState.supportsCas && + Number(result?.applied?.ui_meta_revisions?.[GROUP_CHAT_SYNC_META_KEY] || 0) !== writeRevision + ) { + throw new Error('Gateway did not advance group chat ui_meta revision') + } + + const confirmedState = await groupChatRemoteSnapshot(job) + if (remoteState.supportsCas && confirmedState.revision < writeRevision) { + throw new Error('Group chat ui_meta revision missing after read-back') + } + if (confirmedState.snapshot) { + const mergedRooms = mergeRemoteGroupChatSnapshotIntoRooms(confirmedState.snapshot, $groupChats.get(), { + preserveRooms: groupChatSyncPending?.changedRooms || [], + deletedRooms: groupChatSyncPending?.deletedRooms || [] + }) $groupChats.set(mergedRooms) await persistGroupChatRooms(mergedRooms) } groupChatSyncRetryCount = 0 } catch { if (!groupChatSyncDisposed) { - groupChatSyncPending ||= job + groupChatSyncPending = mergeGroupChatSyncJobs(groupChatSyncPending, job) groupChatSyncRetryCount += 1 if (typeof setTimeout === 'function' && groupChatSyncRetryTimer === null) { groupChatSyncRetryTimer = setTimeout(() => { @@ -680,15 +810,17 @@ function stopGroupChatServerSync() { /** Debounced, pull-merge-write server mirror. Local storage keeps the complete * orchestration log; ui_meta is a bounded cross-client projection. */ -function scheduleGroupChatServerSync(all = $groupChats.get(), { allowEmpty = false, deletedRooms = [] } = {}) { +function scheduleGroupChatServerSync( + all = $groupChats.get(), + { allowEmpty = false, changedRooms = [], deletedRooms = [] } = {} +) { // Browser shells provide timers; source-level VM tests and older embedded // hosts may not. Room persistence must never break the surrounding gateway // lifecycle when the optional mirror cannot be scheduled. if (typeof setTimeout !== 'function') { return } - const deleted = Object.fromEntries(deletedRooms.map(name => [name, Date.now()])) - const snapshot = groupChatSyncSnapshot(all, deleted) + const snapshot = groupChatSyncSnapshot(all) // A newly installed Desktop has no local room cache. Publishing that empty // state on hydrate/reconnect would erase a valid mirror produced elsewhere. // Only an explicit final-room disband is allowed to clear the projection. @@ -702,7 +834,12 @@ function scheduleGroupChatServerSync(all = $groupChats.get(), { allowEmpty = fal clearTimeout(groupChatSyncRetryTimer) groupChatSyncRetryTimer = null } - groupChatSyncPending = { snapshot, connectionId: groupChatSyncConnectionId() } + groupChatSyncPending = mergeGroupChatSyncJobs(groupChatSyncPending, { + connectionId: groupChatSyncConnectionId(), + allowEmpty, + changedRooms, + deletedRooms + }) groupChatSyncTimer = setTimeout(() => { groupChatSyncTimer = null void flushGroupChatServerSync() @@ -4584,7 +4721,7 @@ function trimGroupChatLog(log, watermarks, limit = GROUP_CHAT_HISTORY_LIMIT * 4) } /** Mutate one group's room state through the atom + persist the durable part. */ -function updateGroupChat(group, mutate) { +function updateGroupChat(group, mutate, { sync = true } = {}) { const all = { ...$groupChats.get() } const current = all[group] || { log: [], watermarks: {}, epoch: 0, running: false } const next = mutate({ ...current, log: [...current.log], watermarks: { ...current.watermarks } }) @@ -4621,7 +4758,8 @@ function updateGroupChat(group, mutate) { // Immutable room identity: the member-session title for new rooms. roomId: typeof room.roomId === 'string' && room.roomId ? room.roomId : null, // Room picture (small data URL, same normalization as bot avatars). - image: room.image || null + image: room.image || null, + syncRevision: Math.max(0, Number(room.syncRevision || 0)) } } @@ -4629,7 +4767,9 @@ function updateGroupChat(group, mutate) { } catch { /* storage unavailable — room survives for this window only */ } - scheduleGroupChatServerSync(all) + if (sync) { + scheduleGroupChatServerSync(all, { changedRooms: [group] }) + } return next } @@ -4684,7 +4824,8 @@ async function disbandGroupChat(group, members) { sessions: room.sessions || {}, members: Array.isArray(room.members) ? room.members : [], roomId: typeof room.roomId === 'string' && room.roomId ? room.roomId : null, - image: room.image || null + image: room.image || null, + syncRevision: Math.max(0, Number(room.syncRevision || 0)) } } } @@ -4797,7 +4938,14 @@ async function renameGroupChat(oldName, newName, members) { } // Persist the re-keyed map (updateGroupChat writes the whole durable map). - updateGroupChat(next, r => r) + updateGroupChat(next, r => r, { sync: false }) + // A rename is one revisioned state transition: the new identity is updated + // and the old identity is tombstoned together, so cold hydration cannot + // merge the pre-rename room back into the roster. + scheduleGroupChatServerSync($groupChats.get(), { + changedRooms: [next], + deletedRooms: [oldName] + }) // Follow the open views to the new identity. if ($groupChatWorkspace.get() === oldName) { @@ -4818,8 +4966,21 @@ async function renameGroupChat(oldName, newName, members) { return next } +function groupChatEntryId() { + if (globalThis.crypto && typeof globalThis.crypto.randomUUID === 'function') { + return globalThis.crypto.randomUUID() + } + return `${Date.now().toString(36)}-${Math.random().toString(36).slice(2)}` +} + function appendGroupChatEntry(group, from, text, thread, images) { - const entry = { at: Date.now(), from, text: String(text).trim(), thread: thread || 'legacy' } + const entry = { + id: groupChatEntryId(), + at: Date.now(), + from, + text: String(text).trim(), + thread: thread || 'legacy' + } if (Array.isArray(images) && images.length) { // [{ name, data }] — data URLs. Persisted with the room log so reloads @@ -11318,6 +11479,7 @@ export default { members: Array.isArray(room.members) ? room.members : [], roomId: typeof room.roomId === 'string' && room.roomId ? room.roomId : null, image: typeof room.image === 'string' && room.image ? room.image : null, + syncRevision: Math.max(0, Number(room.syncRevision || 0)), epoch: 0, running: false } diff --git a/apps/desktop/src/plugins/hermes-bots/tests/group-chat.test.mjs b/apps/desktop/src/plugins/hermes-bots/tests/group-chat.test.mjs index b443d8fe7f..04d7ed03aa 100644 --- a/apps/desktop/src/plugins/hermes-bots/tests/group-chat.test.mjs +++ b/apps/desktop/src/plugins/hermes-bots/tests/group-chat.test.mjs @@ -8,7 +8,7 @@ const pluginSource = readFileSync(new URL('../plugin.js', import.meta.url), 'utf /** Load the plugin in a vm with a scripted cli.exec so member turns are * deterministic. `turnScript(profile, prompt)` returns the member's reply * text (or throws to simulate a failed turn). */ -function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approvalUntilResumeCall } = {}) { +function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approvalUntilResumeCall, conflictOnce = false, deferredTimers = false } = {}) { const values = new Map() const atom = initial => { const slot = { get: () => values.get(slot), set: value => values.set(slot, value) } @@ -22,6 +22,8 @@ function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approva const sessions = new Map() const runtimeToStored = new Map() const titleToStored = new Map() + const sharedUiMeta = {} + const uiMetaRevisions = {} let sessionSequence = 0 // busyUntilResumeCall[profile] = N: this profile's session.resume reports // inflight/running for its first N calls, then flips to done — simulating @@ -29,6 +31,7 @@ function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approva // unbounded real-time wait if a caller polls it (used to exercise the // stranded/busy-responder guard). const resumeCallCounts = new Map() + let injectedConflict = false const resolveSession = (profile, target) => { const stored = runtimeToStored.get(target) || (sessions.has(target) ? target : titleToStored.get(`${profile}::${target}`)) @@ -37,6 +40,10 @@ function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approva const context = { atom, setTimeout: fn => { + if (deferredTimers) { + setImmediate(fn) + return 1 + } fn() return 0 }, @@ -47,6 +54,64 @@ function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approva host: { request: async (method, params) => { requests.push({ method, params }) + if (method === 'profiles.list') { + return { + profiles: [ + { + name: 'default', + ui_meta: { ...sharedUiMeta }, + ui_meta_revisions: { ...uiMetaRevisions } + } + ] + } + } + if (method === 'profiles.configure') { + if (conflictOnce && !injectedConflict) { + injectedConflict = true + sharedUiMeta['hermes-bots-groups'] = { + version: 2, + rooms: { + Shared: { + revision: 1, + log: [{ id: 'writer-a:1', from: { kind: 'user', name: 'You' }, text: 'alpha', at: 1 }], + members: [{ name: 'alpha' }] + } + } + } + uiMetaRevisions['hermes-bots-groups'] = 1 + return { + applied: { + ui_meta: false, + ui_meta_conflicts: { 'hermes-bots-groups': { expected: 0, actual: 1 } }, + ui_meta_revisions: { 'hermes-bots-groups': 1 } + } + } + } + const expected = params.ui_meta_expected_revisions || null + if (expected) { + for (const key of Object.keys(params.ui_meta || {})) { + if ((uiMetaRevisions[key] || 0) !== expected[key]) { + return { + applied: { + ui_meta: false, + ui_meta_conflicts: { + [key]: { expected: expected[key], actual: uiMetaRevisions[key] || 0 } + } + } + } + } + } + } + for (const [key, value] of Object.entries(params.ui_meta || {})) { + if (value === null) { + delete sharedUiMeta[key] + } else { + sharedUiMeta[key] = value + } + uiMetaRevisions[key] = (uiMetaRevisions[key] || 0) + 1 + } + return { applied: { ui_meta: true, ui_meta_revisions: { ...uiMetaRevisions } } } + } if (method === 'session.create') { sessionSequence += 1 const stored = `sid-${params.profile}-${sessionSequence}` @@ -126,7 +191,7 @@ function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approva .replace(/^import .* from 'react\/jsx-runtime'\r?\n/m, '') .replace('export default {', 'globalThis.plugin = {') .concat( - '\nglobalThis.__gc = { sendToGroupChat, runGroupChatRounds, harvestStrandedGroupReply, resolveGroupResponders, parseGroupChatMentions, rotateGroupSpeakers, isGroupPassText, formatGroupChatLine, buildGroupChatTurnPrompt, trimGroupChatLog, groupChatSyncSnapshot, groupChatGatewayJsonSize, mergeGroupChatSyncSnapshots, mergeRemoteGroupChatSnapshotIntoRooms, scheduleGroupChatServerSync, disbandGroupChat, updateGroupChat, ensureGroupChatSession, uniqueGroupChatName, liveGroupChatNames, openGroupChat, closeGroupChatMainTab, shouldRenderGroupChatInPane, syncGroupClarify, clearGroupClarify, answerGroupClarify, $groupClarify, $groupChats, $groupNeedsYou, $groupChatWorkspace, $groupMainTabsRev, $botMeta, GROUP_CHAT_MAX_ROUNDS, GROUP_CHAT_MAX_MESSAGES };\n' + '\nglobalThis.__gc = { sendToGroupChat, runGroupChatRounds, harvestStrandedGroupReply, resolveGroupResponders, parseGroupChatMentions, rotateGroupSpeakers, isGroupPassText, formatGroupChatLine, buildGroupChatTurnPrompt, trimGroupChatLog, groupChatSyncSnapshot, groupChatGatewayJsonSize, mergeGroupChatSyncSnapshots, mergeRemoteGroupChatSnapshotIntoRooms, scheduleGroupChatServerSync, disbandGroupChat, renameGroupChat, updateGroupChat, ensureGroupChatSession, uniqueGroupChatName, liveGroupChatNames, openGroupChat, closeGroupChatMainTab, shouldRenderGroupChatInPane, syncGroupClarify, clearGroupClarify, answerGroupClarify, $groupClarify, $groupChats, $groupNeedsYou, $groupChatWorkspace, $groupMainTabsRev, $botMeta, GROUP_CHAT_MAX_ROUNDS, GROUP_CHAT_MAX_MESSAGES };\n' ) vm.runInNewContext(source, context, { filename: 'plugin.js' }) const storageWrites = new Map() @@ -134,7 +199,7 @@ function load(turnScript, { busyUntilResumeCall, clarifyUntilResumeCall, approva storage: { get: () => null, set: (key, value) => storageWrites.set(key, value) }, register: () => undefined }) - return { ...context.__gc, approvalResponds, calls, clarifyResponds, host: context.host, requests, sessions, storageWrites } + return { ...context.__gc, approvalResponds, calls, clarifyResponds, host: context.host, requests, sessions, storageWrites, sharedUiMeta, uiMetaRevisions } } const MEMBERS = [{ name: 'research', title: '' }, { name: 'builder', title: '' }, { name: 'ops', title: 'The Ops' }] @@ -463,7 +528,7 @@ test('group room messages and members mirror through bounded gateway profile met assert.ok(configure, 'room updates are mirrored to the gateway') assert.equal(configure.params.name, 'default') const envelope = configure.params.ui_meta['hermes-bots-groups'] - assert.equal(envelope.version, 1) + assert.equal(envelope.version, 2) assert.equal(envelope.rooms.Research.log[0].text, 'What changed?') assert.equal(JSON.stringify(envelope.rooms.Research.members.map(member => member.name)), JSON.stringify(['research', 'builder'])) assert.ok(gc.groupChatGatewayJsonSize(envelope) <= 48000) @@ -605,19 +670,178 @@ test('gateway projection hydrates a cold Desktop without dropping local runtime test('room deletion tombstone wins over stale history but not a later recreation', () => { const gc = load(() => '(pass)') const stale = gc.mergeGroupChatSyncSnapshots( - { rooms: { Research: { log: [{ from: { kind: 'user', name: 'You' }, text: 'old', at: 10 }] } } }, - { rooms: {}, deleted: { Research: 20 } } + { + version: 2, + rooms: { Research: { revision: 1, log: [{ id: 'old', from: { kind: 'user', name: 'You' }, text: 'old', at: 10 }] } } + }, + { version: 2, rooms: {}, deleted: { Research: 2 } } ) assert.equal(stale.rooms.Research, undefined) - assert.equal(stale.deleted.Research, 20) + assert.equal(stale.deleted.Research, 2) const recreated = gc.mergeGroupChatSyncSnapshots(stale, { - rooms: { Research: { log: [{ from: { kind: 'user', name: 'You' }, text: 'new', at: 30 }] } } + version: 2, + rooms: { Research: { revision: 3, log: [{ id: 'new', from: { kind: 'user', name: 'You' }, text: 'new', at: 1 }] } } }) assert.equal(recreated.rooms.Research.log[0].text, 'new') assert.equal(recreated.deleted, undefined) }) +test('gateway revisions, not device clocks, order room deletion and recreation', () => { + const gc = load(() => '(pass)') + const merged = gc.mergeGroupChatSyncSnapshots( + { + version: 2, + rooms: { + ClockSkewed: { + revision: 8, + log: [{ id: 'future-clock', from: { kind: 'user', name: 'You' }, text: 'old', at: 9999999999999 }] + } + } + }, + { version: 2, rooms: {}, deleted: { ClockSkewed: 9 } } + ) + + assert.equal(merged.rooms.ClockSkewed, undefined) + assert.equal(merged.deleted.ClockSkewed, 9) +}) + +test('a conflicting writer can merge the winner and preserve both stable message ids', () => { + const gc = load(() => '(pass)') + const winner = { + version: 2, + rooms: { + Shared: { + revision: 1, + log: [{ id: 'writer-a:1', from: { kind: 'user', name: 'You' }, text: 'alpha', at: 100 }], + members: [{ name: 'alpha' }] + } + } + } + const loserRetry = gc.mergeGroupChatSyncSnapshots( + winner, + { + version: 2, + rooms: { + Shared: { + revision: 0, + log: [{ id: 'writer-b:1', from: { kind: 'member', name: 'beta' }, text: 'beta', at: 1 }], + members: [{ name: 'beta' }] + } + } + }, + { changedRooms: ['Shared'], writeRevision: 2 } + ) + + assert.equal( + JSON.stringify(loserRetry.rooms.Shared.log.map(entry => entry.id).sort()), + JSON.stringify(['writer-a:1', 'writer-b:1']) + ) + assert.equal(loserRetry.rooms.Shared.revision, 2) +}) + +test('sync worker retries a gateway CAS conflict and publishes the merged room', async () => { + const gc = load(() => '(pass)', { conflictOnce: true, deferredTimers: true }) + gc.$groupChats.set({ + Shared: { + log: [{ id: 'writer-b:1', from: { kind: 'member', name: 'beta' }, text: 'beta', at: 2 }], + members: [{ name: 'beta' }], + watermarks: {}, + sessions: {}, + syncRevision: 0 + } + }) + + gc.scheduleGroupChatServerSync(gc.$groupChats.get(), { changedRooms: ['Shared'] }) + for (let i = 0; i < 20 && (gc.uiMetaRevisions['hermes-bots-groups'] || 0) < 2; i++) { + await new Promise(resolve => setImmediate(resolve)) + } + + const stored = gc.sharedUiMeta['hermes-bots-groups'] + assert.equal(gc.uiMetaRevisions['hermes-bots-groups'], 2) + assert.equal( + JSON.stringify(stored.rooms.Shared.log.map(entry => entry.id).sort()), + JSON.stringify(['writer-a:1', 'writer-b:1']) + ) + assert.equal( + gc.requests.filter(call => call.method === 'profiles.configure').length, + 2 + ) +}) + +test('read-back does not clobber a newer local room edit queued during an in-flight write', () => { + const gc = load(() => '(pass)') + const merged = gc.mergeRemoteGroupChatSnapshotIntoRooms( + { + version: 2, + rooms: { + Shared: { + revision: 5, + log: [{ id: 'remote', from: { kind: 'user', name: 'You' }, text: 'remote message', at: 1 }], + members: [{ name: 'old-member' }], + image: 'data:image/png;base64,old' + } + } + }, + { + Shared: { + syncRevision: 4, + log: [{ id: 'local', from: { kind: 'user', name: 'You' }, text: 'local message', at: 2 }], + members: [{ name: 'new-member' }], + image: 'data:image/png;base64,new' + } + }, + { preserveRooms: ['Shared'] } + ) + + assert.equal(JSON.stringify(merged.Shared.members.map(member => member.name)), JSON.stringify(['new-member'])) + assert.equal(merged.Shared.image, 'data:image/png;base64,new') + assert.equal(merged.Shared.syncRevision, 4) + assert.equal( + JSON.stringify(merged.Shared.log.map(entry => entry.id).sort()), + JSON.stringify(['local', 'remote']) + ) +}) + +test('rename publishes a new room plus old-name tombstone and cold hydrate cannot resurrect it', () => { + const gc = load(() => '(pass)') + const before = { + version: 2, + rooms: { + Old: { + revision: 4, + log: [{ id: 'turn-1', from: { kind: 'user', name: 'You' }, text: 'history', at: 10 }], + members: [{ name: 'research' }], + image: 'data:image/png;base64,room' + } + } + } + const after = gc.mergeGroupChatSyncSnapshots( + before, + { + version: 2, + rooms: { + New: { + revision: 4, + log: before.rooms.Old.log, + members: before.rooms.Old.members, + image: before.rooms.Old.image + } + } + }, + { changedRooms: ['New'], deletedRooms: ['Old'], writeRevision: 5 } + ) + const hydrated = gc.mergeRemoteGroupChatSnapshotIntoRooms(after, {}) + + assert.equal(after.rooms.Old, undefined) + assert.equal(after.deleted.Old, 5) + assert.equal(after.rooms.New.revision, 5) + assert.equal(after.rooms.New.image, 'data:image/png;base64,room') + assert.equal(hydrated.Old, undefined) + assert.equal(hydrated.New.log[0].text, 'history') + assert.equal(hydrated.New.image, 'data:image/png;base64,room') +}) + test('source contract: workspace + main-window door + prompt rules are wired', () => { assert.match(pluginSource, /function GroupChatWorkspace\(/) // Group rows open through the main-window door, feature-detected with the diff --git a/tests/tui_gateway/test_profiles_ui_meta_cas.py b/tests/tui_gateway/test_profiles_ui_meta_cas.py new file mode 100644 index 0000000000..f7de9e1541 --- /dev/null +++ b/tests/tui_gateway/test_profiles_ui_meta_cas.py @@ -0,0 +1,98 @@ +"""Gateway-owned revisions for concurrent profile UI metadata updates. + +Clients use ``profiles.list`` to read a shared key and its revision, merge a +local mutation, then pass that revision back to ``profiles.configure``. The +gateway must let exactly one concurrent writer advance the key and make every +stale writer retry instead of silently replacing somebody else's state. +""" + +from __future__ import annotations + +from concurrent.futures import ThreadPoolExecutor + +import pytest + +import tui_gateway.server as srv + + +@pytest.fixture +def home(tmp_path, monkeypatch): + hermes_home = tmp_path / ".hermes" + hermes_home.mkdir() + monkeypatch.setenv("HERMES_HOME", str(hermes_home)) + return hermes_home + + +def _configure(ui_meta, expected=None): + params = {"name": "default", "ui_meta": ui_meta} + if expected is not None: + params["ui_meta_expected_revisions"] = expected + return srv._methods["profiles.configure"]("configure", params)["result"]["applied"] + + +def _default_profile(): + rows = srv._methods["profiles.list"]( + "list", {"include_sessions": False} + )["result"]["profiles"] + return next(row for row in rows if row["name"] == "default") + + +def test_ui_meta_revision_advances_and_stale_compare_and_swap_is_rejected(home): + first = _configure({"shared-room": {"messages": ["one"]}}) + assert first["ui_meta"] is True + assert first["ui_meta_revisions"] == {"shared-room": 1} + + second = _configure( + {"shared-room": {"messages": ["one", "two"]}}, + {"shared-room": 1}, + ) + assert second["ui_meta"] is True + assert second["ui_meta_revisions"] == {"shared-room": 2} + + stale = _configure( + {"shared-room": {"messages": ["stale replacement"]}}, + {"shared-room": 1}, + ) + assert stale["ui_meta"] is False + assert stale["ui_meta_conflicts"] == { + "shared-room": {"expected": 1, "actual": 2} + } + + row = _default_profile() + assert row["ui_meta"]["shared-room"] == {"messages": ["one", "two"]} + assert row["ui_meta_revisions"]["shared-room"] == 2 + + +def test_ui_meta_revision_survives_key_deletion(home): + _configure({"shared-room": {"messages": ["one"]}}) + deleted = _configure({"shared-room": None}, {"shared-room": 1}) + + assert deleted["ui_meta"] is True + assert deleted["ui_meta_revisions"] == {"shared-room": 2} + row = _default_profile() + assert "shared-room" not in row.get("ui_meta", {}) + assert row["ui_meta_revisions"]["shared-room"] == 2 + + stale_recreate = _configure( + {"shared-room": {"messages": ["resurrected"]}}, + {"shared-room": 1}, + ) + assert stale_recreate["ui_meta"] is False + assert stale_recreate["ui_meta_conflicts"]["shared-room"]["actual"] == 2 + + +def test_two_concurrent_writers_cannot_both_replace_the_same_revision(home): + def write(label): + return _configure( + {"shared-room": {"messages": [label]}}, + {"shared-room": 0}, + ) + + with ThreadPoolExecutor(max_workers=2) as pool: + results = list(pool.map(write, ("alpha", "beta"))) + + assert sum(result["ui_meta"] is True for result in results) == 1 + assert sum(result["ui_meta"] is False for result in results) == 1 + loser = next(result for result in results if result["ui_meta"] is False) + assert loser["ui_meta_conflicts"]["shared-room"]["actual"] == 1 + assert _default_profile()["ui_meta_revisions"]["shared-room"] == 1 diff --git a/tui_gateway/methods_profiles.py b/tui_gateway/methods_profiles.py index def566478f..a7bd780c3c 100644 --- a/tui_gateway/methods_profiles.py +++ b/tui_gateway/methods_profiles.py @@ -243,12 +243,22 @@ def _(rid, params: dict) -> dict: from pathlib import Path as _Path meta_path = _Path(str(p.path)) / "profile.yaml" + # Presence of this field feature-detects gateway-owned CAS, + # including a brand-new profile whose revision map is empty. + row["ui_meta_revisions"] = {} if meta_path.is_file(): with open(meta_path, "r", encoding="utf-8") as f: raw_meta = _yaml.safe_load(f) or {} ui_meta = raw_meta.get("ui_meta") if isinstance(ui_meta, dict) and ui_meta: row["ui_meta"] = ui_meta + revisions = raw_meta.get("_ui_meta_revisions") + if isinstance(revisions, dict) and revisions: + row["ui_meta_revisions"] = { + str(key): max(0, int(value)) + for key, value in revisions.items() + if isinstance(value, int) and not isinstance(value, bool) + } except Exception: pass @@ -693,7 +703,9 @@ def _(rid, params: dict) -> dict: ``model`` + ``provider`` (both required together), ``disabled_skills`` (list[str], replace semantics), ``enabled_toolsets`` (list[str], replace semantics; empty list clears - the pin so every toolset is enabled again). + the pin so every toolset is enabled again), and + ``ui_meta_expected_revisions`` (dict[str, int], optional compare-and-swap + preconditions for keys supplied in ``ui_meta``). Each section is applied independently and best-effort; the result reports per-section success so a UI can surface partial failures. @@ -728,32 +740,73 @@ def _(rid, params: dict) -> dict: else: import yaml as _yaml - meta_path = profile_dir / "profile.yaml" - existing = {} - if meta_path.is_file(): - try: - with open(meta_path, "r", encoding="utf-8") as f: - loaded = _yaml.safe_load(f) or {} - if isinstance(loaded, dict): - existing = loaded - except Exception: - existing = {} - current = existing.get("ui_meta") - if not isinstance(current, dict): - current = {} - for key, value in incoming.items(): - if value is None: - current.pop(key, None) - else: - current[key] = value - if current: - existing["ui_meta"] = current - else: - existing.pop("ui_meta", None) - from utils import atomic_yaml_write + expected = params.get("ui_meta_expected_revisions") + if expected is not None and not isinstance(expected, dict): + raise ValueError("ui_meta_expected_revisions must be an object") - atomic_yaml_write(meta_path, existing, sort_keys=False) - applied["ui_meta"] = True + meta_path = profile_dir / "profile.yaml" + with _profile_ui_meta_lock: + existing = {} + if meta_path.is_file(): + try: + with open(meta_path, "r", encoding="utf-8") as f: + loaded = _yaml.safe_load(f) or {} + if isinstance(loaded, dict): + existing = loaded + except Exception: + existing = {} + + raw_revisions = existing.get("_ui_meta_revisions") + revisions = dict(raw_revisions) if isinstance(raw_revisions, dict) else {} + revisions = { + str(key): max(0, int(value)) + for key, value in revisions.items() + if isinstance(value, int) and not isinstance(value, bool) + } + conflicts = {} + if isinstance(expected, dict): + for key in incoming: + wanted = expected.get(key) + actual = revisions.get(key, 0) + if ( + not isinstance(wanted, int) + or isinstance(wanted, bool) + or wanted < 0 + or wanted != actual + ): + conflicts[key] = {"expected": wanted, "actual": actual} + + if conflicts: + applied["ui_meta"] = False + applied["ui_meta_conflicts"] = conflicts + applied["ui_meta_revisions"] = { + key: revisions.get(key, 0) for key in incoming + } + else: + current = existing.get("ui_meta") + if not isinstance(current, dict): + current = {} + for key, value in incoming.items(): + if value is None: + current.pop(key, None) + else: + current[key] = value + revisions[key] = revisions.get(key, 0) + 1 + if current: + existing["ui_meta"] = current + else: + existing.pop("ui_meta", None) + # Revisions intentionally survive deletion: a + # stale client must not recreate a removed key by + # presenting the initial revision again. + existing["_ui_meta_revisions"] = revisions + from utils import atomic_yaml_write + + atomic_yaml_write(meta_path, existing, sort_keys=False) + applied["ui_meta"] = True + applied["ui_meta_revisions"] = { + key: revisions[key] for key in incoming + } except Exception: applied["ui_meta"] = False diff --git a/tui_gateway/server.py b/tui_gateway/server.py index d9abe90501..7f8ea642b1 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -153,6 +153,10 @@ _db = None _db_error: str | None = None _stdout_lock = threading.Lock() _cfg_lock = threading.Lock() +# Shared profile UI metadata can be updated concurrently by Desktop, mobile, +# and multiple worker-pool RPCs. Its compare/check/write transaction needs a +# dedicated lock rather than the unrelated process-config cache lock. +_profile_ui_meta_lock = threading.Lock() _sessions_lock = threading.RLock() # reentrant: _close_session_by_id may run under callers that already hold it _prompt_lock = threading.Lock() _cfg_cache: dict | None = None