fix(gateway): serialize shared bot room updates
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user