Merge remote-tracking branch 'origin/main' into codex/81234-live-main-final
This commit is contained in:
@@ -15,6 +15,7 @@ import {
|
||||
backendScopeKey,
|
||||
backendScopePrefix,
|
||||
buildAgentRoster,
|
||||
connectionDialFieldsChanged,
|
||||
connectionIdForLabel,
|
||||
labelKey,
|
||||
labelSlug,
|
||||
@@ -523,3 +524,31 @@ test('upsertConnection replaces by id and appends new ids', () => {
|
||||
assert.equal(registry.connections.filter(c => c.id === a.id).length, 1)
|
||||
assert.equal(registry.connections.find(c => c.id === a.id)?.url, 'http://a:2')
|
||||
})
|
||||
|
||||
// --- connectionDialFieldsChanged (edit → recycle decision) ---
|
||||
|
||||
test('connectionDialFieldsChanged: label-only edits do not recycle', () => {
|
||||
const before = { id: 'homelab', kind: 'remote', label: 'Homelab', url: 'http://10.0.0.5:9119', authMode: 'token', token: { encoding: 'safeStorage', value: 'abc' } } as const
|
||||
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before, label: 'Home lab (renamed)' }), false)
|
||||
// Identity edit is also a no-op.
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before }), false)
|
||||
})
|
||||
|
||||
test('connectionDialFieldsChanged: url / auth / token changes recycle', () => {
|
||||
const before = { id: 'homelab', kind: 'remote', label: 'Homelab', url: 'http://10.0.0.5:9119', authMode: 'token', token: { encoding: 'safeStorage', value: 'abc' } } as const
|
||||
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before, url: 'http://10.0.0.9:9119' }), true)
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before, authMode: 'oauth', token: undefined }), true)
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before, token: { encoding: 'safeStorage', value: 'NEW' } }), true)
|
||||
})
|
||||
|
||||
test('connectionDialFieldsChanged: ssh routing fields recycle, kind change recycles', () => {
|
||||
const before = { id: 'box', kind: 'ssh', label: 'Box', host: 'box.lan', user: 'me', port: 22 } as const
|
||||
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before, label: 'Box 2' }), false)
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before, host: 'other.lan' }), true)
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before, port: 2222 }), true)
|
||||
assert.equal(connectionDialFieldsChanged(before, { ...before, remoteProfile: 'work' }), true)
|
||||
assert.equal(connectionDialFieldsChanged(before, { id: 'box', kind: 'remote', label: 'Box', url: 'http://x:1' }), true)
|
||||
})
|
||||
|
||||
@@ -414,6 +414,43 @@ export function mergeConnectionInput(input: ConnectionInput, existing?: null | R
|
||||
return merged
|
||||
}
|
||||
|
||||
/**
|
||||
* True when an edit changes how a connection is DIALED — endpoint, auth, or
|
||||
* ssh routing fields — as opposed to a cosmetic label rename. Callers use
|
||||
* this to decide whether live pooled backends / renderer sockets for the
|
||||
* connection must be recycled after a save: a label-only edit keeps traffic
|
||||
* flowing, while a url/token/host change means everything currently open
|
||||
* points at the OLD target and must be torn down and re-dialed.
|
||||
*/
|
||||
export function connectionDialFieldsChanged(before: RegistryConnection, after: RegistryConnection): boolean {
|
||||
if (before.kind !== after.kind) {
|
||||
return true
|
||||
}
|
||||
|
||||
const fields: (keyof RegistryConnection)[] = [
|
||||
'url',
|
||||
'authMode',
|
||||
'org',
|
||||
'host',
|
||||
'user',
|
||||
'port',
|
||||
'keyPath',
|
||||
'remoteHermesPath',
|
||||
'remoteProfile'
|
||||
]
|
||||
|
||||
for (const field of fields) {
|
||||
if ((before[field] ?? null) !== (after[field] ?? null)) {
|
||||
return true
|
||||
}
|
||||
}
|
||||
|
||||
// Token envelopes are opaque here (main.ts encrypts). An edit that carries
|
||||
// no new token inherits the stored envelope verbatim, so structural
|
||||
// equality is exact for the label-only case.
|
||||
return JSON.stringify(before.token ?? null) !== JSON.stringify(after.token ?? null)
|
||||
}
|
||||
|
||||
// ── Registry-level operations (all pure: return a new registry) ────────────
|
||||
|
||||
function localEntry(label = 'This device'): RegistryConnection {
|
||||
|
||||
@@ -86,6 +86,7 @@ import {
|
||||
backendScopeKey,
|
||||
backendScopePrefix,
|
||||
buildAgentRoster,
|
||||
connectionDialFieldsChanged,
|
||||
mergeConnectionInput,
|
||||
migrateV1ToRegistry,
|
||||
normalizeConnectionInput,
|
||||
@@ -212,6 +213,7 @@ import {
|
||||
electronProcessStartMarker,
|
||||
parentWatchdogEnv
|
||||
} from './parent-process-identity'
|
||||
import { selectPoolEvictions } from './pool-eviction'
|
||||
import { createKeepAwake } from './power-save'
|
||||
import { FirstRunSetupResetError, runPrimaryBackendStartup } from './primary-backend-startup'
|
||||
import { rehomePrimaryConnection } from './primary-connection-rehome'
|
||||
@@ -8061,6 +8063,16 @@ function saveRegistryConnection(input: any = {}) {
|
||||
|
||||
writeDesktopConnectionsRegistry(upsertConnection(registry, entry))
|
||||
|
||||
// A dial-material edit (endpoint/auth/ssh routing — NOT a label rename)
|
||||
// leaves pooled backends under `conn:<id>::*` and renderer sockets pointing
|
||||
// at the OLD target while the UI shows the new one. Recycle them: stop this
|
||||
// connection's pooled backends/tunnels and tell renderers to dispose+redial
|
||||
// their secondaries for this connection id.
|
||||
if (existing && connectionDialFieldsChanged(existing, entry)) {
|
||||
stopRegistryConnectionBackends(entry.id)
|
||||
broadcastConnectionsChanged({ connectionId: entry.id, reason: 'updated' })
|
||||
}
|
||||
|
||||
return sanitizeRegistryConnection(entry)
|
||||
}
|
||||
|
||||
@@ -9143,6 +9155,21 @@ function sendConnectionApplied() {
|
||||
webContents.send('hermes:connection:applied')
|
||||
}
|
||||
|
||||
// Registry lifecycle push: a connection was removed or materially edited, so
|
||||
// every window must tear down (and, for edits, re-dial) its secondary sockets
|
||||
// scoped to that connection. Without this, a removed remote/cloud source keeps
|
||||
// its renderer WebSocket open and streaming as a ghost, and an edited one
|
||||
// keeps talking to the OLD endpoint until idle-reap.
|
||||
function broadcastConnectionsChanged(payload: { connectionId: string; reason: 'removed' | 'updated' }) {
|
||||
for (const win of BrowserWindow.getAllWindows()) {
|
||||
const { webContents } = win
|
||||
|
||||
if (webContents && !webContents.isDestroyed()) {
|
||||
webContents.send('hermes:connections:changed', payload)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async function waitForBackendExit(child, timeoutMs = 5000) {
|
||||
if (!child || child.exitCode !== null || child.signalCode !== null) {
|
||||
return
|
||||
@@ -9420,31 +9447,22 @@ function touchPoolBackend(profile) {
|
||||
}
|
||||
}
|
||||
|
||||
// Evict least-recently-used pool backends until at most `keep` remain — but only
|
||||
// ever evict backends without a live renderer socket (stale beyond the keepalive
|
||||
// window). When every backend is actively kept alive we let the pool exceed the
|
||||
// soft cap rather than kill a running session.
|
||||
// Evict least-recently-used SPAWNED pool backends until at most `keep` remain —
|
||||
// but only ever evict backends without a live renderer socket (stale beyond the
|
||||
// keepalive window). When every backend is actively kept alive we let the pool
|
||||
// exceed the soft cap rather than kill a running session. Process-less
|
||||
// descriptor entries (remote/cloud registry sources, per-profile remote
|
||||
// overrides — `entry.process === null`) are excluded from the cap entirely:
|
||||
// they hold no local process, so counting them used to let a roster refresh
|
||||
// across N registered remote connections LRU-evict a REAL local backend that
|
||||
// was merely idle past the keepalive window. Descriptors are still reclaimed
|
||||
// by the idle reaper.
|
||||
function evictLruPoolBackends(keep) {
|
||||
if (backendPool.size <= keep) {
|
||||
return
|
||||
}
|
||||
|
||||
const now = Date.now()
|
||||
|
||||
const evictable = [...backendPool.entries()]
|
||||
.filter(([, entry]) => now - (entry.lastActiveAt || 0) > POOL_KEEPALIVE_FRESH_MS)
|
||||
.sort((a, b) => (a[1].lastActiveAt || 0) - (b[1].lastActiveAt || 0))
|
||||
|
||||
let removable = backendPool.size - Math.max(0, keep)
|
||||
|
||||
for (const [profile] of evictable) {
|
||||
if (removable <= 0) {
|
||||
break
|
||||
}
|
||||
const evictions = selectPoolEvictions(backendPool.entries(), Math.max(0, keep), Date.now(), POOL_KEEPALIVE_FRESH_MS)
|
||||
|
||||
for (const profile of evictions) {
|
||||
rememberLog(`Evicting idle profile backend "${profile}" (LRU cap ${POOL_MAX_BACKENDS})`)
|
||||
stopPoolBackend(profile)
|
||||
removable -= 1
|
||||
}
|
||||
}
|
||||
|
||||
@@ -11915,6 +11933,10 @@ ipcMain.handle('hermes:connections:remove', async (_event, id) => {
|
||||
// Tear down anything the removed connection still had running: pooled
|
||||
// backends under its composite keys and any ssh tunnel scopes it owned.
|
||||
stopRegistryConnectionBackends(key)
|
||||
// And the renderer side: without this push, secondaries scoped to the
|
||||
// removed connection keep their WebSocket open (remote/cloud have no local
|
||||
// process to kill) and stream ghost events until page reload.
|
||||
broadcastConnectionsChanged({ connectionId: key, reason: 'removed' })
|
||||
|
||||
return { ok: true, registry: sanitizeConnectionsRegistry(registry) }
|
||||
})
|
||||
|
||||
@@ -0,0 +1,85 @@
|
||||
/**
|
||||
* Tests for electron/pool-eviction.ts — LRU cap accounting for the desktop
|
||||
* backend pool. The cap exists to bound SPAWNED local backends (real child
|
||||
* processes); process-less descriptor entries (remote/cloud registry sources,
|
||||
* per-profile remote overrides) must not count against it, or a roster
|
||||
* refresh across N registered remote connections evicts a real local backend
|
||||
* that was merely idle past the keepalive window.
|
||||
*/
|
||||
|
||||
import assert from 'node:assert/strict'
|
||||
|
||||
import { test } from 'vitest'
|
||||
|
||||
import { selectPoolEvictions } from './pool-eviction'
|
||||
|
||||
const NOW = 1_000_000
|
||||
const FRESH_MS = 90_000
|
||||
|
||||
/** A spawned local backend entry (has a child process). */
|
||||
const spawned = (idleMs: number) => ({ process: { pid: 123 }, lastActiveAt: NOW - idleMs })
|
||||
|
||||
/** A process-less remote/cloud descriptor entry. */
|
||||
const descriptor = (idleMs: number) => ({ process: null, lastActiveAt: NOW - idleMs })
|
||||
|
||||
test('process-less descriptors do not count toward the cap', () => {
|
||||
// 1 real spawned backend idle beyond the keepalive window + 3 remote
|
||||
// descriptors: total size (4) exceeds keep (2), but only ONE entry holds a
|
||||
// process, so nothing may be evicted. This is the roster-refresh regression:
|
||||
// the old size-based accounting evicted the real local backend here.
|
||||
const entries: [string, ReturnType<typeof spawned>][] = [
|
||||
['default', spawned(120_000)],
|
||||
['conn:homelab::default', descriptor(0)],
|
||||
['conn:office::default', descriptor(0)],
|
||||
['conn:cloud-a::default', descriptor(0)]
|
||||
]
|
||||
|
||||
assert.deepEqual(selectPoolEvictions(entries, 2, NOW, FRESH_MS), [])
|
||||
})
|
||||
|
||||
test('spawned backends over the cap are still LRU-evicted', () => {
|
||||
const entries: [string, ReturnType<typeof spawned>][] = [
|
||||
['a', spawned(500_000)],
|
||||
['b', spawned(300_000)],
|
||||
['c', spawned(100_000)],
|
||||
// Descriptors interleaved: must neither inflate the count nor be evicted.
|
||||
['conn:x::a', descriptor(999_000)]
|
||||
]
|
||||
|
||||
// keep=2 → one spawned backend over; evict the least-recently-used ('a').
|
||||
assert.deepEqual(selectPoolEvictions(entries, 2, NOW, FRESH_MS), ['a'])
|
||||
})
|
||||
|
||||
test('fresh spawned backends are spared even over the cap', () => {
|
||||
const entries: [string, ReturnType<typeof spawned>][] = [
|
||||
['a', spawned(1_000)],
|
||||
['b', spawned(2_000)],
|
||||
['c', spawned(3_000)]
|
||||
]
|
||||
|
||||
// All within the keepalive window → the pool may exceed the soft cap.
|
||||
assert.deepEqual(selectPoolEvictions(entries, 1, NOW, FRESH_MS), [])
|
||||
})
|
||||
|
||||
test('evicts only enough stale spawned backends to reach the cap', () => {
|
||||
const entries: [string, ReturnType<typeof spawned>][] = [
|
||||
['a', spawned(500_000)],
|
||||
['b', spawned(400_000)],
|
||||
['c', spawned(300_000)],
|
||||
['d', spawned(1_000)]
|
||||
]
|
||||
|
||||
// 4 spawned, keep 2 → remove 2, oldest first.
|
||||
assert.deepEqual(selectPoolEvictions(entries, 2, NOW, FRESH_MS), ['a', 'b'])
|
||||
})
|
||||
|
||||
test('descriptor-only pools never evict', () => {
|
||||
const entries: [string, ReturnType<typeof descriptor>][] = [
|
||||
['conn:a::p', descriptor(999_000)],
|
||||
['conn:b::p', descriptor(999_000)],
|
||||
['conn:c::p', descriptor(999_000)],
|
||||
['conn:d::p', descriptor(999_000)]
|
||||
]
|
||||
|
||||
assert.deepEqual(selectPoolEvictions(entries, 2, NOW, FRESH_MS), [])
|
||||
})
|
||||
@@ -0,0 +1,58 @@
|
||||
// LRU cap accounting for the desktop backend pool.
|
||||
//
|
||||
// The pool holds two very different kinds of entries under one Map:
|
||||
// 1. SPAWNED local profile backends — a real child process each (the thing
|
||||
// the POOL_MAX_BACKENDS cap exists to bound).
|
||||
// 2. Process-less connection DESCRIPTORS — remote/cloud registry sources and
|
||||
// per-profile remote overrides (`entry.process === null`). These hold no
|
||||
// local process; their only cost is a cached descriptor.
|
||||
//
|
||||
// Counting both kinds against the cap meant a roster refresh across N
|
||||
// registered remote connections could push the Map size over the cap and
|
||||
// LRU-evict a REAL spawned backend that had merely been idle past the
|
||||
// keepalive window. Cap accounting (and cap-driven eviction) therefore only
|
||||
// considers entries with a live child process; descriptor entries remain
|
||||
// subject to the idle reaper, just not to the process cap.
|
||||
|
||||
export interface PoolEvictionEntry {
|
||||
lastActiveAt?: null | number
|
||||
process?: unknown
|
||||
}
|
||||
|
||||
/**
|
||||
* Pick which pool keys the LRU cap should evict so that at most `keep`
|
||||
* SPAWNED backends remain. Only entries with a live child process count
|
||||
* toward the cap or are eligible for cap eviction, and — as before — only
|
||||
* entries idle beyond `freshMs` may be evicted (an actively kept-alive pool
|
||||
* may exceed the soft cap rather than kill a running session).
|
||||
*/
|
||||
export function selectPoolEvictions<K>(
|
||||
entries: Iterable<[K, PoolEvictionEntry]>,
|
||||
keep: number,
|
||||
now: number,
|
||||
freshMs: number
|
||||
): K[] {
|
||||
const spawned = [...entries].filter(([, entry]) => Boolean(entry.process))
|
||||
|
||||
if (spawned.length <= keep) {
|
||||
return []
|
||||
}
|
||||
|
||||
const evictable = spawned
|
||||
.filter(([, entry]) => now - (entry.lastActiveAt || 0) > freshMs)
|
||||
.sort((a, b) => (a[1].lastActiveAt || 0) - (b[1].lastActiveAt || 0))
|
||||
|
||||
let removable = spawned.length - Math.max(0, keep)
|
||||
const evictions: K[] = []
|
||||
|
||||
for (const [key] of evictable) {
|
||||
if (removable <= 0) {
|
||||
break
|
||||
}
|
||||
|
||||
evictions.push(key)
|
||||
removable -= 1
|
||||
}
|
||||
|
||||
return evictions
|
||||
}
|
||||
@@ -138,7 +138,16 @@ contextBridge.exposeInMainWorld('hermesDesktop', {
|
||||
setPrimary: id => ipcRenderer.invoke('hermes:connections:set-primary', id),
|
||||
test: id => ipcRenderer.invoke('hermes:connections:test', id),
|
||||
// Fan out `hermes update` to every eligible registered connection.
|
||||
updateAll: () => ipcRenderer.invoke('hermes:connections:update-all')
|
||||
updateAll: () => ipcRenderer.invoke('hermes:connections:update-all'),
|
||||
// Registry lifecycle push (main → renderer): a connection was removed or
|
||||
// materially edited, so secondaries scoped to it must be disposed (and,
|
||||
// for edits, re-dialed at the new target).
|
||||
onChanged: callback => {
|
||||
const listener = (_event, payload) => callback(payload)
|
||||
ipcRenderer.on('hermes:connections:changed', listener)
|
||||
|
||||
return () => ipcRenderer.removeListener('hermes:connections:changed', listener)
|
||||
}
|
||||
},
|
||||
sshConfigHosts: () => ipcRenderer.invoke('hermes:ssh-config:hosts'),
|
||||
sshResolveHost: host => ipcRenderer.invoke('hermes:ssh-config:resolve', host),
|
||||
|
||||
@@ -17,6 +17,7 @@ import {
|
||||
$gateway,
|
||||
closeSecondaryGateways,
|
||||
configureGatewayRegistry,
|
||||
disposeSecondariesForConnection,
|
||||
ensureGatewayForProfile,
|
||||
pruneSecondaryGateways,
|
||||
reconnectSecondaryGateways,
|
||||
@@ -434,6 +435,18 @@ export function useGatewayBoot({
|
||||
const offPowerResume = desktop.onPowerResume?.(() => reconnectNow())
|
||||
const offConnectionApplied = desktop.onConnectionApplied?.(() => void softSwitch())
|
||||
|
||||
// Registry lifecycle: a removed connection's secondaries must close NOW
|
||||
// (remote/cloud have no local process whose death would drop the socket —
|
||||
// they'd keep streaming ghost events); a materially edited one is
|
||||
// disposed AND re-dialed so its sockets target the new endpoint.
|
||||
const offConnectionsChanged = desktop.connections?.onChanged?.(payload => {
|
||||
if (!payload || typeof payload.connectionId !== 'string') {
|
||||
return
|
||||
}
|
||||
|
||||
disposeSecondariesForConnection(payload.connectionId, { redial: payload.reason === 'updated' })
|
||||
})
|
||||
|
||||
const onOnline = () => reconnectNow()
|
||||
|
||||
const onVisible = () => {
|
||||
@@ -636,6 +649,7 @@ export function useGatewayBoot({
|
||||
document.removeEventListener('visibilitychange', onVisible)
|
||||
offPowerResume?.()
|
||||
offConnectionApplied?.()
|
||||
offConnectionsChanged?.()
|
||||
offState()
|
||||
offEvent()
|
||||
offExit()
|
||||
|
||||
Vendored
+7
@@ -143,6 +143,13 @@ declare global {
|
||||
// Fan out `hermes update` to every eligible registered connection;
|
||||
// cloud entries are skipped (platform-managed), each row independent.
|
||||
updateAll?: () => Promise<{ ok: boolean; results: DesktopConnectionUpdateResult[] }>
|
||||
// Registry lifecycle push: fired when a connection is removed or
|
||||
// materially edited so the renderer can dispose (and re-dial) the
|
||||
// secondary gateways scoped to it. Optional: older Electron mains
|
||||
// don't emit it.
|
||||
onChanged?: (
|
||||
callback: (payload: { connectionId: string; reason: 'removed' | 'updated' }) => void
|
||||
) => () => void
|
||||
}
|
||||
sshConfigHosts: () => Promise<DesktopSshHostsResult>
|
||||
sshResolveHost: (host: string) => Promise<DesktopSshResolveResult>
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { afterEach, describe, expect, it } from 'vitest'
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
import { createClientSessionState } from '@/lib/chat-runtime'
|
||||
import { $gatewayState } from '@/store/session'
|
||||
@@ -6,6 +6,118 @@ import { $sessionStates, dropSessionState, publishSessionState } from '@/store/s
|
||||
|
||||
import { host } from './index'
|
||||
|
||||
// Plugins read app state exclusively through host.state — and before this
|
||||
// contract existed, only the PRIMARY workspace tab was reachable
|
||||
// ($activeSessionId). Clicking a tile never moved any plugin-visible atom,
|
||||
// and tile focus is pure renderer state that gateway RPC can never see, so
|
||||
// no plugin-side workaround was possible. These atoms are the plugin door to
|
||||
// the same focused-session signals the core statusbar reads
|
||||
// (use-statusbar-items.tsx).
|
||||
|
||||
describe('host.state focused-session atoms', () => {
|
||||
beforeEach(() => {
|
||||
window.localStorage.clear()
|
||||
vi.resetModules()
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
vi.resetModules()
|
||||
})
|
||||
|
||||
async function setup() {
|
||||
const { host } = await import('@/sdk/index')
|
||||
const states = await import('@/store/session-states')
|
||||
const session = await import('@/store/session')
|
||||
|
||||
return { host, states, session }
|
||||
}
|
||||
|
||||
it('exposes readonly atoms for the focused session (runtime id, stored id, usage)', async () => {
|
||||
const { host } = await setup()
|
||||
|
||||
for (const key of ['focusedSessionId', 'focusedStoredSessionId', 'focusedUsage'] as const) {
|
||||
const store = host.state[key]
|
||||
expect(store, key).toBeDefined()
|
||||
expect(typeof store.get, key).toBe('function')
|
||||
expect(typeof store.listen, key).toBe('function')
|
||||
expect(typeof store.subscribe, key).toBe('function')
|
||||
}
|
||||
})
|
||||
|
||||
it('mirrors the primary session while no tile is focused', async () => {
|
||||
const { host, states } = await setup()
|
||||
|
||||
expect(host.state.focusedSessionId.get()).toBe(states.$focusedRuntimeId.get())
|
||||
expect(host.state.focusedStoredSessionId.get()).toBe(states.$focusedStoredSessionId.get())
|
||||
})
|
||||
|
||||
it('focusedUsage projects the focused session usage, null while unresolved', async () => {
|
||||
const { host, states } = await setup()
|
||||
|
||||
const focused = states.$focusedSessionState.get()
|
||||
expect(host.state.focusedUsage.get()).toBe(focused?.usage ?? null)
|
||||
})
|
||||
|
||||
it('follows the interacted tile while the primary-only atom stays put', async () => {
|
||||
const { host, session, states } = await setup()
|
||||
const tree = await import('@/components/pane-shell/tree/store')
|
||||
const model = await import('@/components/pane-shell/tree/model')
|
||||
const { registry } = await import('@/contrib/registry')
|
||||
|
||||
// A second chat zone holding a session tile, next to the main workspace.
|
||||
for (const id of ['workspace', 'session-tile:tile-a']) {
|
||||
registry.register({
|
||||
area: 'panes',
|
||||
data: id === 'workspace' ? { placement: 'main', uncloseable: true } : { placement: 'main' },
|
||||
id,
|
||||
render: () => null,
|
||||
title: id
|
||||
})
|
||||
}
|
||||
|
||||
tree.declareDefaultTree(
|
||||
model.split('row', [
|
||||
model.group(['workspace'], { active: 'workspace', id: 'grp-main' }),
|
||||
model.group(['session-tile:tile-a'], { active: 'session-tile:tile-a', id: 'grp-side' })
|
||||
])
|
||||
)
|
||||
|
||||
const primaryBefore = session.$activeSessionId.get()
|
||||
const primarySelection = session.$selectedStoredSessionId.get()
|
||||
|
||||
// Bind the tile to a live runtime with its own usage, so the readout
|
||||
// atoms (focusedSessionId / focusedUsage) — not just the navigation id —
|
||||
// are proven to follow the tile.
|
||||
const tileUsage = {
|
||||
calls: 3,
|
||||
input: 1200,
|
||||
output: 300,
|
||||
total: 1500,
|
||||
context_used: 42000,
|
||||
context_max: 200000,
|
||||
context_percent: 21,
|
||||
cost_usd: 0.0123
|
||||
}
|
||||
|
||||
states.$sessionTiles.set([{ storedSessionId: 'tile-a', runtimeId: 'runtime-tile-a' }])
|
||||
states.$sessionStates.set({
|
||||
'runtime-tile-a': { storedSessionId: 'tile-a', usage: tileUsage } as never
|
||||
})
|
||||
|
||||
// Focusing the tile zone moves the focused atoms onto the tile's session…
|
||||
tree.noteActiveTreeGroup('grp-side')
|
||||
expect(host.state.focusedStoredSessionId.get()).toBe('tile-a')
|
||||
expect(host.state.focusedSessionId.get()).toBe('runtime-tile-a')
|
||||
expect(host.state.focusedUsage.get()).toBe(tileUsage)
|
||||
// …while the primary-only atom a plugin used to rely on does not move.
|
||||
expect(host.state.activeSessionId.get()).toBe(primaryBefore)
|
||||
|
||||
// Focusing back homes to the primary's selection, whatever it is.
|
||||
tree.noteActiveTreeGroup('grp-main')
|
||||
expect(host.state.focusedStoredSessionId.get()).toBe(primarySelection)
|
||||
})
|
||||
})
|
||||
|
||||
describe('host.state busy vs gateway', () => {
|
||||
afterEach(() => {
|
||||
$sessionStates.set({})
|
||||
|
||||
@@ -45,8 +45,14 @@ import {
|
||||
$gatewayState,
|
||||
$selectedStoredSessionId
|
||||
} from '@/store/session'
|
||||
import { $focusedSessionState, $focusedStoredSessionId, $sessionStates } from '@/store/session-states'
|
||||
import {
|
||||
$focusedRuntimeId,
|
||||
$focusedSessionState,
|
||||
$focusedStoredSessionId,
|
||||
$sessionStates
|
||||
} from '@/store/session-states'
|
||||
import { runGatewayRestart } from '@/store/system-actions'
|
||||
import type { UsageStats } from '@/types/hermes'
|
||||
|
||||
// -- state: readonly views over the app's live atoms -------------------------
|
||||
|
||||
@@ -109,6 +115,10 @@ if (typeof window !== 'undefined') {
|
||||
$narrowViewport.listen(refresh)
|
||||
}
|
||||
|
||||
/** Live usage of the FOCUSED session, projected out of the streamed session
|
||||
* state — the same readout the core statusbar's context chip paints. */
|
||||
const $focusedUsage = computed($focusedSessionState, state => state?.usage ?? null)
|
||||
|
||||
export const host = {
|
||||
state: {
|
||||
/** Runtime id of the active chat session (null on a fresh draft). */
|
||||
@@ -126,6 +136,19 @@ export const host = {
|
||||
busyBySession: readonlyAtom<Record<string, boolean>>($busyBySession),
|
||||
/** Active workspace cwd ('' when detached). */
|
||||
cwd: readonlyAtom<string>($currentCwd),
|
||||
/** Runtime id of the FOCUSED chat session — the interacted tile, else the
|
||||
* primary. Prefer this over `activeSessionId` for any readout that
|
||||
* should follow the user between tiles (context, tokens, cost). */
|
||||
focusedSessionId: readonlyAtom<null | string>($focusedRuntimeId),
|
||||
/** Stored (durable) id of the focused session — for navigation and
|
||||
* session-list matching, where runtime ids don't survive reloads. */
|
||||
focusedStoredSessionId: readonlyAtom<null | string>($focusedStoredSessionId),
|
||||
/** Live usage snapshot of the focused session (`context_used` /
|
||||
* `context_max` / `context_percent`, token counts, `cost_usd`) —
|
||||
* streamed by the backend, no RPC needed. Null while unresolved.
|
||||
* The UsageStats-optional fields (context_*, cost_usd) arrive as the
|
||||
* backend reports them, so read them with a fallback. */
|
||||
focusedUsage: readonlyAtom<null | UsageStats>($focusedUsage),
|
||||
/** Gateway socket state: 'idle' | 'connecting' | 'open' | …. Not turn-busy. */
|
||||
gateway: readonlyAtom<string>($gatewayState),
|
||||
/** Current main model slug. */
|
||||
|
||||
@@ -0,0 +1,195 @@
|
||||
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'
|
||||
|
||||
// Connection lifecycle for registry-scoped secondary gateways:
|
||||
//
|
||||
// 1. Removing a connection must dispose its secondaries — remote/cloud
|
||||
// sources have no local process whose death would drop the socket, so
|
||||
// without an explicit dispose the WebSocket stays open and streams ghost
|
||||
// events until page reload.
|
||||
// 2. A materially edited connection re-dials so fresh sockets target the
|
||||
// NEW endpoint.
|
||||
// 3. When the Electron main reports the connection no longer exists
|
||||
// (`No connection with id`), the reconnect loop fail-stops and evicts
|
||||
// the entry instead of retrying forever.
|
||||
|
||||
const gatewayMocks = vi.hoisted(() => {
|
||||
const instances: { close: ReturnType<typeof vi.fn>; connectionState: string }[] = []
|
||||
|
||||
return {
|
||||
connect: vi.fn(async (_wsUrl: string): Promise<void> => undefined),
|
||||
instances
|
||||
}
|
||||
})
|
||||
|
||||
vi.mock('@/hermes', () => ({
|
||||
HermesGateway: class {
|
||||
connectionState = 'closed'
|
||||
close = vi.fn(() => {
|
||||
this.connectionState = 'closed'
|
||||
})
|
||||
connect = async (wsUrl: string): Promise<void> => {
|
||||
await gatewayMocks.connect(wsUrl)
|
||||
this.connectionState = 'open'
|
||||
}
|
||||
onEvent = vi.fn(() => () => {})
|
||||
onState = vi.fn(() => () => {})
|
||||
constructor() {
|
||||
gatewayMocks.instances.push(this as never)
|
||||
}
|
||||
}
|
||||
}))
|
||||
vi.mock('@/store/session', () => ({
|
||||
setConnection: vi.fn(),
|
||||
setGatewayState: vi.fn()
|
||||
}))
|
||||
vi.mock('@/store/notify-baseline', () => ({ markNativeNotifyBaseline: vi.fn() }))
|
||||
|
||||
const {
|
||||
closeSecondaryGateways,
|
||||
configureGatewayRegistry,
|
||||
disposeSecondariesForConnection,
|
||||
ensureActiveGatewayOpen,
|
||||
ensureGatewayForAgent,
|
||||
setPrimaryGateway
|
||||
} = await import('./gateway')
|
||||
|
||||
function installDesktop(stub: Record<string, unknown>): void {
|
||||
;(window as unknown as { hermesDesktop: unknown }).hermesDesktop = stub
|
||||
}
|
||||
|
||||
function descriptorFor(connectionId: string, profile: string) {
|
||||
return {
|
||||
authMode: 'token',
|
||||
baseUrl: `https://${connectionId}.invalid`,
|
||||
mode: 'remote',
|
||||
profile,
|
||||
token: 'fake-test-token',
|
||||
wsUrl: `wss://${connectionId}.invalid/api/ws?profile=${profile}`
|
||||
}
|
||||
}
|
||||
|
||||
beforeEach(() => {
|
||||
configureGatewayRegistry({ onEvent: vi.fn() } as never)
|
||||
setPrimaryGateway({ connectionState: 'open' } as never, 'default')
|
||||
})
|
||||
|
||||
afterEach(() => {
|
||||
closeSecondaryGateways()
|
||||
gatewayMocks.instances.length = 0
|
||||
vi.clearAllMocks()
|
||||
vi.useRealTimers()
|
||||
delete (window as unknown as { hermesDesktop?: unknown }).hermesDesktop
|
||||
})
|
||||
|
||||
describe('disposeSecondariesForConnection', () => {
|
||||
it('closes and evicts every secondary scoped to the removed connection', async () => {
|
||||
const getConnectionFor = vi.fn(async ({ connectionId, profile }: { connectionId: string; profile: string }) =>
|
||||
descriptorFor(connectionId, profile)
|
||||
)
|
||||
|
||||
installDesktop({ getConnectionFor })
|
||||
|
||||
await ensureGatewayForAgent('homelab', 'default')
|
||||
await ensureGatewayForAgent('homelab', 'work')
|
||||
await ensureGatewayForAgent('office', 'default')
|
||||
|
||||
expect(gatewayMocks.instances).toHaveLength(3)
|
||||
|
||||
disposeSecondariesForConnection('homelab')
|
||||
|
||||
// Both homelab sockets closed; the office socket untouched.
|
||||
expect(gatewayMocks.instances[0].close).toHaveBeenCalledOnce()
|
||||
expect(gatewayMocks.instances[1].close).toHaveBeenCalledOnce()
|
||||
expect(gatewayMocks.instances[2].close).not.toHaveBeenCalled()
|
||||
|
||||
// No redial for a removal.
|
||||
expect(getConnectionFor).toHaveBeenCalledTimes(3)
|
||||
})
|
||||
|
||||
it('re-dials disposed secondaries when redial is requested (material edit)', async () => {
|
||||
const getConnectionFor = vi.fn(async ({ connectionId, profile }: { connectionId: string; profile: string }) =>
|
||||
descriptorFor(connectionId, profile)
|
||||
)
|
||||
|
||||
installDesktop({ getConnectionFor })
|
||||
|
||||
await ensureGatewayForAgent('homelab', 'default')
|
||||
expect(gatewayMocks.connect).toHaveBeenCalledTimes(1)
|
||||
|
||||
disposeSecondariesForConnection('homelab', { redial: true })
|
||||
|
||||
// The redial runs async through the normal open path — flush it.
|
||||
await vi.waitFor(() => {
|
||||
expect(gatewayMocks.connect).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
// Old socket closed, fresh descriptor fetched (would carry the new URL).
|
||||
expect(gatewayMocks.instances[0].close).toHaveBeenCalledOnce()
|
||||
expect(getConnectionFor).toHaveBeenCalledTimes(2)
|
||||
})
|
||||
|
||||
it('is a no-op for blank or unknown connection ids', async () => {
|
||||
installDesktop({
|
||||
getConnectionFor: vi.fn(async ({ connectionId, profile }: { connectionId: string; profile: string }) =>
|
||||
descriptorFor(connectionId, profile)
|
||||
)
|
||||
})
|
||||
|
||||
await ensureGatewayForAgent('homelab', 'default')
|
||||
|
||||
disposeSecondariesForConnection('')
|
||||
disposeSecondariesForConnection('ghost')
|
||||
|
||||
expect(gatewayMocks.instances[0].close).not.toHaveBeenCalled()
|
||||
})
|
||||
})
|
||||
|
||||
describe('reconnect fail-stop on a removed connection', () => {
|
||||
it('evicts the entry instead of retrying when the registry no longer knows the id', async () => {
|
||||
const getConnectionFor = vi
|
||||
.fn()
|
||||
.mockResolvedValueOnce(descriptorFor('homelab', 'default'))
|
||||
.mockRejectedValue(new Error('No connection with id "homelab".'))
|
||||
|
||||
installDesktop({ getConnectionFor })
|
||||
|
||||
await ensureGatewayForAgent('homelab', 'default')
|
||||
expect(gatewayMocks.instances).toHaveLength(1)
|
||||
|
||||
// Simulate the socket dropping after the connection was removed.
|
||||
const socket = gatewayMocks.instances[0] as unknown as { connectionState: string }
|
||||
socket.connectionState = 'closed'
|
||||
|
||||
// ensureActiveGatewayOpen drives reconnectSecondary for the active scope.
|
||||
const result = await ensureActiveGatewayOpen()
|
||||
|
||||
expect(result).toBeNull()
|
||||
// Fail-stop: the entry was disposed + evicted, so a second drive finds
|
||||
// nothing to retry (no further getConnectionFor calls).
|
||||
const callsAfterFailStop = getConnectionFor.mock.calls.length
|
||||
await ensureActiveGatewayOpen()
|
||||
expect(getConnectionFor.mock.calls.length).toBe(callsAfterFailStop)
|
||||
})
|
||||
|
||||
it('keeps retrying on ordinary transport failures', async () => {
|
||||
const getConnectionFor = vi
|
||||
.fn()
|
||||
.mockResolvedValueOnce(descriptorFor('homelab', 'default'))
|
||||
.mockRejectedValueOnce(new Error('ECONNREFUSED'))
|
||||
.mockResolvedValue(descriptorFor('homelab', 'default'))
|
||||
|
||||
installDesktop({ getConnectionFor })
|
||||
|
||||
await ensureGatewayForAgent('homelab', 'default')
|
||||
|
||||
const socket = gatewayMocks.instances[0] as unknown as { connectionState: string }
|
||||
socket.connectionState = 'closed'
|
||||
|
||||
// First drive fails with a transport error → entry survives.
|
||||
await ensureActiveGatewayOpen()
|
||||
// Second drive succeeds against the surviving entry.
|
||||
const reopened = await ensureActiveGatewayOpen()
|
||||
|
||||
expect(reopened).not.toBeNull()
|
||||
})
|
||||
})
|
||||
@@ -235,8 +235,18 @@ async function reconnectSecondary(entry: Secondary): Promise<void> {
|
||||
try {
|
||||
await openSecondary(entry)
|
||||
entry.reconnectAttempt = 0
|
||||
} catch {
|
||||
// Transport failure → fall through to the backoff below.
|
||||
} catch (error) {
|
||||
// The registry no longer knows this connection (removed while we were
|
||||
// backing off). Retrying forever can never succeed — fail-stop: dispose
|
||||
// the entry and evict it instead of an infinite 15s-cap retry loop.
|
||||
if (entry.connectionId && isMissingConnectionError(error)) {
|
||||
entry.reconnecting = false
|
||||
disposeSecondary(entry)
|
||||
g.secondaries.delete(entry.scope)
|
||||
|
||||
return
|
||||
}
|
||||
// Other transport failure → fall through to the backoff below.
|
||||
} finally {
|
||||
entry.reconnecting = false
|
||||
|
||||
@@ -246,6 +256,15 @@ async function reconnectSecondary(entry: Secondary): Promise<void> {
|
||||
}
|
||||
}
|
||||
|
||||
// Electron's getConnectionFor rejects with `No connection with id "…"` when
|
||||
// the registry entry is gone. That is a permanent condition for the scoped
|
||||
// socket, unlike transient transport errors.
|
||||
function isMissingConnectionError(error: unknown): boolean {
|
||||
const message = error instanceof Error ? error.message : String(error ?? '')
|
||||
|
||||
return message.includes('No connection with id')
|
||||
}
|
||||
|
||||
function createSecondary(profile: string, connectionId: null | string = null): Secondary {
|
||||
const gateway = new HermesGateway()
|
||||
const scope = backendScopeKey(connectionId, profile)
|
||||
@@ -533,6 +552,39 @@ export function closeSecondaryGateways(): void {
|
||||
restoreActiveToPrimaryIfEvicted()
|
||||
}
|
||||
|
||||
// Registry lifecycle: a connection was removed or materially edited. Dispose
|
||||
// every secondary scoped to it (a removed remote/cloud source has no local
|
||||
// process to die, so without this its WebSocket stays open streaming ghost
|
||||
// events). With `redial` (the edit case) each disposed profile is re-dialed
|
||||
// through the normal open path so the fresh socket targets the NEW endpoint;
|
||||
// the active scope re-activates so the foreground keeps painting.
|
||||
export function disposeSecondariesForConnection(connectionId: string, opts: { redial?: boolean } = {}): void {
|
||||
const id = String(connectionId || '').trim()
|
||||
|
||||
if (!id) {
|
||||
return
|
||||
}
|
||||
|
||||
for (const [key, entry] of [...g.secondaries]) {
|
||||
if (entry.connectionId !== id) {
|
||||
continue
|
||||
}
|
||||
|
||||
const wasActive = key === g.activeKey
|
||||
|
||||
disposeSecondary(entry)
|
||||
g.secondaries.delete(key)
|
||||
|
||||
if (opts.redial) {
|
||||
const reopen = wasActive
|
||||
? ensureGatewayForAgent(entry.connectionId, entry.profile)
|
||||
: openGatewayForAgent(entry.connectionId, entry.profile)
|
||||
|
||||
void reopen.catch(() => undefined)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Self-accept so editing this module (or a fan-out that lands here) is an
|
||||
// in-place hot update instead of a full page reload — the live sockets in `g`
|
||||
// survive the swap. Dev-only: production strips import.meta.hot.
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
lepetitprince716-prog
|
||||
@@ -56,11 +56,16 @@ The ONLY import surface is `@hermes/plugin-sdk` (plus `react` /
|
||||
|
||||
- `host.state.*` — readonly reactive atoms: `activeSessionId`, `busy`,
|
||||
`awaitingResponse`, `busyBySession`, `cwd`, `gateway` (socket state, not
|
||||
turn-busy), `model`, `profile`, `viewport`.
|
||||
`busy` is true while the focused chat is working after a send (thinking
|
||||
and streaming). `awaitingResponse` is true until the first assistant
|
||||
payload. `busyBySession` maps runtime session id → mid-turn, for rosters
|
||||
that watch every session. Read with `.get()` in handlers,
|
||||
turn-busy), `model`, `profile`, `viewport`, plus the tile-aware focused
|
||||
session atoms: `focusedSessionId` (runtime id — key for `session.*` RPC),
|
||||
`focusedStoredSessionId` (durable id — navigation / list matching), and
|
||||
`focusedUsage` (live streamed `UsageStats` of the focused session, no RPC
|
||||
needed). `busy` is true while the focused chat is working after a send
|
||||
(thinking and streaming). `awaitingResponse` is true until the first
|
||||
assistant payload. `busyBySession` maps runtime session id → mid-turn,
|
||||
for rosters that watch every session. Prefer the focused atoms for any
|
||||
readout that should follow the user between tiles. Read with `.get()`
|
||||
in handlers,
|
||||
`useValue(atom)` in components.
|
||||
- `host.request(method, params)` — gateway JSON-RPC (sessions, config,
|
||||
skills, cron — everything the app uses).
|
||||
|
||||
@@ -378,6 +378,9 @@ host.state.activeSessionId // ReadableAtom<string | null>
|
||||
host.state.awaitingResponse // ReadableAtom<boolean> true until the first assistant payload
|
||||
host.state.busy // ReadableAtom<boolean> focused chat is working after a send
|
||||
host.state.busyBySession // ReadableAtom<Record<string, boolean>> runtime id → mid-turn
|
||||
host.state.focusedSessionId // ReadableAtom<string | null> (runtime id of the FOCUSED session — tile-aware; prefer for session.* RPC)
|
||||
host.state.focusedStoredSessionId // ReadableAtom<string | null> (durable id — navigation / session-list matching)
|
||||
host.state.focusedUsage // ReadableAtom<UsageStats | null> (live streamed usage of the focused session, no RPC needed)
|
||||
host.state.cwd // ReadableAtom<string>
|
||||
host.state.gateway // ReadableAtom<string> socket state ('idle' | 'connecting' | 'open' | …)
|
||||
host.state.model // ReadableAtom<string>
|
||||
|
||||
Reference in New Issue
Block a user