Procházet zdrojové kódy

refactor(gui): session titles ride the generic projection pair; title-snapshot map retired

The manager's titleSnapshots Map and its session/title frame consumption
dissolve into resident per-session ProjectionValueStores (create-on-demand,
outliving instantiation — the same role the snapshot map played): a
session/projection frame lands whether or not the Session exists, list rows
read the store's 'title' key, subscribed baselines truncate phantom rows, and
session-removed drops the store. The fixture converts to the host parallel:
a projections block on the tail page (title + todos units), push frames on
unit-advancing events, and a post-subscribe projection baseline replacing the
bespoke title control frame.
imccyu před 2 měsíci
rodič
revize
f42943a14c

+ 36 - 21
packages/client/connection/src/client/fixture.ts

@@ -252,18 +252,27 @@ function viewFor(event: SessionEvent, log: readonly SessionEvent[]): ToolEventVi
   return undefined
 }
 
-/** Fold the latest fixture title into the host's control-frame projection. */
-function titleFrameOf(id: SessionId, log: readonly SessionEvent[]): Extract<MuxFrame, { type: 'session/title' }> | undefined {
-  const event = log.findLast(item => (item as { type: string }).type === 'session/title')
-  if (event === undefined) return undefined
-  const titleEvent = event as unknown as { seq: number; time: number; data: { title: string } }
-  return {
-    type: 'session/title',
-    sessionId: id,
-    title: titleEvent.data.title,
-    eventSeq: titleEvent.seq,
-    updatedAt: titleEvent.time,
+/** Fixture parallel of the host's projection units: whole current values per key over the full log. */
+function projectionValuesOf(log: readonly SessionEvent[]): Record<string, unknown> {
+  const values: Record<string, unknown> = {}
+  const titleEvent = log.findLast(item => (item as { type: string }).type === 'session/title')
+  if (titleEvent !== undefined) {
+    values['title'] = (titleEvent as unknown as { data: { title: string } }).data.title
   }
+  const todos = backscanTodos(log)
+  if (todos !== undefined) values['todos'] = todos
+  return values
+}
+
+/** Host push-frame parallel: emit one session/projection frame per key the given event advanced. */
+function projectionFramesOf(id: SessionId, log: readonly SessionEvent[], event: SessionEvent): Extract<MuxFrame, { type: 'session/projection' }>[] {
+  const type = (event as { type: string }).type
+  const key = type === 'session/title' ? 'title' : type === 'todo/write' ? 'todos' : undefined
+  if (key === undefined) return []
+  const values = projectionValuesOf(log)
+  /* v8 ignore next -- the advancing event is in the log, so its key always has a value. */
+  if (!Object.hasOwn(values, key)) return []
+  return [{ type: 'session/projection', sessionId: id, key, value: values[key], seq: event.seq }]
 }
 
 /**
@@ -489,10 +498,8 @@ export function createFixtureApi(options: FixtureOptions = {}): ApiProxy {
     emitMux(view === undefined
       ? { type: 'session/event', sessionId: id, event }
       : { type: 'session/event', sessionId: id, event, view })
-    if ((event as { type: string }).type === 'session/title') {
-      // The raw title is already in this log, so the latest-title fold must find it.
-      emitMux(titleFrameOf(id, log) as Extract<MuxFrame, { type: 'session/title' }>)
-    }
+    // Host eager-drive parallel: a unit-advancing event pushes its finished value.
+    for (const frame of projectionFramesOf(id, log, event)) emitMux(frame)
   }
 
   /** At most one in-flight replay per session; cancel clears it. */
@@ -644,14 +651,18 @@ export function createFixtureApi(options: FixtureOptions = {}): ApiProxy {
         const log = logs.get(request.payload.sessionId) ?? []
         // Snapshot at request time, deliver after the transit delay (mirrors a real host under latency).
         const page = pageOf(log, request.payload.beforeSeq, request.payload.maxMessages ?? 50)
-        // Tail page carries the session-level todo projection (host parallel: full-log backscan).
-        const todos = request.payload.beforeSeq === undefined ? backscanTodos(log) : undefined
+        // Tail page carries the projections block (host parallel: one consistent
+        // cut over the registered units, asOfSeq = window tail seq); an empty
+        // log has no cut to stamp, so the block stays absent.
+        const projections = request.payload.beforeSeq === undefined && log.length > 0
+          ? { asOfSeq: log.length - 1, values: projectionValuesOf(log) }
+          : undefined
         const doomed = failNextHistory
         failNextHistory = false
         const delay = historyDelayMs
         if (delay > 0) await new Promise(resolve => setTimeout(resolve, delay))
         if (doomed) throw new Error('fixture: simulated history transport failure')
-        return ok(request, { ...page, ...todos === undefined ? {} : { todos } })
+        return ok(request, { ...page, ...projections === undefined ? {} : { projections } })
       },
       prompt: (request) => {
         const { sessionId: id, mode, content } = request.payload
@@ -853,9 +864,13 @@ export function createFixtureApi(options: FixtureOptions = {}): ApiProxy {
         // Open baseline: subscribed sessions + pending interactions replayed with stable rpcIds.
         for (const s of sessions) {
           if (!s.running) continue
-          conn.push({ rpcId: mint(), payload: { type: 'session/subscribed', sessionId: s.sessionId, lastSeq: (logs.get(s.sessionId)?.length ?? 0) - 1 } })
-          const title = titleFrameOf(s.sessionId, logs.get(s.sessionId) ?? [])
-          if (title !== undefined) conn.push({ rpcId: mint(), payload: title })
+          const log = logs.get(s.sessionId) ?? []
+          conn.push({ rpcId: mint(), payload: { type: 'session/subscribed', sessionId: s.sessionId, lastSeq: log.length - 1 } })
+          // Post-subscribe projection baseline (host parallel: recomputed unit values ride push frames).
+          const values = projectionValuesOf(log)
+          for (const key of Object.keys(values)) {
+            conn.push({ rpcId: mint(), payload: { type: 'session/projection', sessionId: s.sessionId, key, value: values[key], seq: log.length - 1 } })
+          }
         }
         conn.push({
           rpcId: pendingApprovalRpcId,

+ 9 - 7
packages/client/connection/tests/fixture.spec.ts

@@ -179,11 +179,13 @@ describe('createFixtureApi', () => {
     const second = await openOnce()
     expect(first[0]?.payload).toMatchObject({ type: 'session/subscribed', sessionId: 'fx-alpha' })
     expect((first[0]?.payload as { lastSeq: number }).lastSeq).toBeGreaterThan(0)
-    expect(first[1]?.payload).toMatchObject({ type: 'session/title', sessionId: 'fx-alpha', title: 'Fixture 历史会话' })
-    expect(first[2]?.payload).toMatchObject({ type: 'approval/requested', toolName: 'dangerous_tool' })
-    expect(second[2]?.rpcId).toBe(first[2]?.rpcId) // stable rpcId across replays (host replay semantics)
-    expect(first[3]?.payload).toMatchObject({ type: 'question/requested', sessionId: 'fx-alpha' })
-    expect(second[3]?.rpcId).toBe(first[3]?.rpcId)
+    // Projection baseline frames follow the subscribed frame (title + todos units).
+    expect(first[1]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'title', value: 'Fixture 历史会话' })
+    expect(first[2]?.payload).toMatchObject({ type: 'session/projection', sessionId: 'fx-alpha', key: 'todos' })
+    expect(first[3]?.payload).toMatchObject({ type: 'approval/requested', toolName: 'dangerous_tool' })
+    expect(second[3]?.rpcId).toBe(first[3]?.rpcId) // stable rpcId across replays (host replay semantics)
+    expect(first[4]?.payload).toMatchObject({ type: 'question/requested', sessionId: 'fx-alpha' })
+    expect(second[4]?.rpcId).toBe(first[4]?.rpcId)
   })
 
   it('steer with no replay in flight falls through to a fresh queued turn; non-text blocks stringify empty', async () => {
@@ -585,11 +587,11 @@ describe('createFixtureApi', () => {
     hooks.appendTitle('fx-alpha', 'Fixture 修订标题')
     await vi.waitFor(() => {
       expect(seen.some(f => f.type === 'session/event' && JSON.stringify(f.event.data).includes('正常直播'))).toBe(true)
-      expect(seen.some(f => f.type === 'session/title' && f.title === 'Fixture 修订标题')).toBe(true)
+      expect(seen.some(f => f.type === 'session/projection' && f.key === 'title' && f.value === 'Fixture 修订标题')).toBe(true)
     })
     expect(seen.some(f => f.type === 'session/event' && JSON.stringify(f.event.data).includes('静默丢帧'))).toBe(false)
     const rawTitleIndex = seen.findIndex(f => f.type === 'session/event' && (f.event as { type: string }).type === 'session/title')
-    const titleControlIndex = seen.findIndex(f => f.type === 'session/title' && f.title === 'Fixture 修订标题')
+    const titleControlIndex = seen.findIndex(f => f.type === 'session/projection' && f.key === 'title' && f.value === 'Fixture 修订标题')
     expect(titleControlIndex).toBe(rawTitleIndex + 1)
     // But history serves the silent event (the client's repull finds it).
     const repull = await api.sessions.history(req({ sessionId: sid('fx-alpha'), maxMessages: 5 }))

+ 37 - 25
packages/client/runtime/src/client/sessions/manager.ts

@@ -10,6 +10,7 @@ import { mergeOrderedBaseline } from '../ordered-baseline.ts'
 import type { SessionListEntry, TitledSessionSummary } from './lineage.ts'
 import { flattenLineage } from './lineage.ts'
 import { Notifier } from './notifier.ts'
+import { ProjectionValueStore } from './projection-store.ts'
 import { Session } from './session.ts'
 
 /**
@@ -43,12 +44,6 @@ type SessionListMutation =
 /** Per-session cap for pre-instantiation approval/question buffering (low-frequency frames; a few dozen covers any real backlog). */
 const PENDING_BUFFER_CAP = 32
 
-/** Latest title control snapshot retained independently of list/instance arrival. */
-interface SessionTitleSnapshot {
-  title: string
-  eventSeq: number
-  updatedAt: number
-}
 
 /** Instance cluster + frame entry + the session list (see the web client architecture RFC). */
 export class SessionManager {
@@ -58,7 +53,11 @@ export class SessionManager {
    *  drop-and-backfill path; replayed and cleared on instantiation. Bounded per session (these
    *  frames are low-frequency; overflow drops oldest) and dropped on session-removed (audit S7). */
   private readonly pendingBuffers = new Map<SessionId, RpcRequest<MuxFrame>[]>()
-  private readonly titleSnapshots = new Map<SessionId, SessionTitleSnapshot>()
+  /** Per-session projection value stores, retained independently of instance arrival (the
+   *  title-snapshot precedent, generalized): push frames land here whether or not the Session
+   *  is instantiated (list rows read the 'title' key), and an instantiated Session adopts the
+   *  same store so history-baseline seeding and frames converge on one row set. */
+  private readonly projectionStores = new Map<SessionId, ProjectionValueStore>()
   private summaries: SessionSummary[] = []
   private listState: 'idle' | 'loading' | 'error' = 'idle'
   /** Arrival phase; the pending → ready edge fires on the first successful pull (see SessionListPhase). */
@@ -163,9 +162,23 @@ export class SessionManager {
       onEngaged: (engaged) => {
         this.recordMutation({ kind: 'engaged', sessionId: engaged.sessionId })
       },
+      projections: this.projectionStore(sessionId),
     })
   }
 
+  /** Resident per-session projection store (create-on-demand; outlives instantiation). */
+  private projectionStore(sessionId: SessionId): ProjectionValueStore {
+    let store = this.projectionStores.get(sessionId)
+    if (store === undefined) {
+      store = new ProjectionValueStore()
+      // List rows project off store keys (title); any-key changes re-enter
+      // the manager's own batched rebuild channel.
+      store.subscribeAny(() => { this.notifier.markDirty() })
+      this.projectionStores.set(sessionId, store)
+    }
+    return store
+  }
+
   // ---- List surface ----
 
   /** Full refresh via session.list (single-flight: an in-flight call is reused). */
@@ -302,23 +315,20 @@ export class SessionManager {
   handleMuxEnvelope(envelope: RpcRequest<MuxFrame>): void {
     const frame = envelope.payload
     if (frame.type === 'stream/error') return // Controller already treats this as stream failure
-    if (frame.type === 'session/title') {
-      const current = this.titleSnapshots.get(frame.sessionId)
-      if (current !== undefined && current.eventSeq >= frame.eventSeq) return
-      this.titleSnapshots.set(frame.sessionId, {
-        title: frame.title,
-        eventSeq: frame.eventSeq,
-        updatedAt: frame.updatedAt,
-      })
+    if (frame.type === 'session/projection') {
+      // Finished host-computed value: land it in the resident store whether or
+      // not the Session is instantiated (list rows read the 'title' key). The
+      // synchronous markDirty keeps the list snapshot same-tick fresh (the
+      // store's own any-key channel is microtask-batched).
+      this.projectionStore(frame.sessionId).apply(frame.key, frame.value, frame.seq)
       this.notifier.markDirty()
       return
     }
     if (frame.type === 'session/subscribed') {
-      const current = this.titleSnapshots.get(frame.sessionId)
-      if (current !== undefined && current.eventSeq > frame.lastSeq) {
-        this.titleSnapshots.delete(frame.sessionId)
-        this.notifier.markDirty()
-      }
+      // Rows past the host's durable baseline rode state a restart lost; drop
+      // them so last-wins cannot pin a phantom value over recomputed truth.
+      this.projectionStores.get(frame.sessionId)?.truncate(frame.lastSeq)
+      this.notifier.markDirty()
       // New mux-generation baseline: buffered session/queued frames belong to
       // the previous generation and the host is about to resend the live
       // snapshot — drop them, or every reconnect appends a duplicate batch
@@ -377,7 +387,7 @@ export class SessionManager {
         this.recordMutation({ kind: 'remove', sessionId: frame.sessionId })
         this.sessions.get(frame.sessionId)?.handleRemoved() // instance survives (resident-instance rule), only flagged in the snapshot
         this.pendingBuffers.delete(frame.sessionId) // a removed session's buffered frames must not replay on a future instantiation
-        this.titleSnapshots.delete(frame.sessionId)
+        this.projectionStores.delete(frame.sessionId) // removed sessions drop their projection rows with the instance
         return
       }
       case 'host/session-status': {
@@ -402,10 +412,12 @@ export class SessionManager {
 
   private buildListSnapshot(): SessionListSnapshot {
     const merged: TitledSessionSummary[] = this.summaries.map((summary) => {
-      const title = this.titleSnapshots.get(summary.sessionId)
-      return title === undefined
-        ? summary
-        : { ...summary, title: title.title, updatedAt: Math.max(summary.updatedAt, title.updatedAt) }
+      // List rows read the generic 'title' projection key (host-computed unit
+      // value; the bespoke session/title frame is retired).
+      const title = this.projectionStores.get(summary.sessionId)?.get('title')
+      return typeof title === 'string' && title !== ''
+        ? { ...summary, title }
+        : summary
     })
     const fresh = flattenLineage(merged)
     const items = fresh.map((entry) => {

+ 24 - 34
packages/client/runtime/tests/manager.spec.ts

@@ -133,21 +133,18 @@ describe('list lifecycle', () => {
     expect(manager.getListSnapshot().items.map(i => i.sessionId)).toEqual([S2])
   })
 
-  it('retains monotonic title snapshots before list arrival, merges recency, and clears them on removal', async () => {
+  it('retains title projections before list arrival, keeps last-wins by seq, and clears them on removal', async () => {
     const api = new FakeApiClient()
     const manager = new SessionManager(api)
-    manager.handleMuxEnvelope({
-      rpcId: 'title-new' as never,
-      payload: { type: 'session/title', sessionId: S1, title: 'Newest', eventSeq: 4, updatedAt: 300 },
-    })
-    manager.handleMuxEnvelope({
-      rpcId: 'title-stale' as never,
-      payload: { type: 'session/title', sessionId: S1, title: 'Stale', eventSeq: 3, updatedAt: 900 },
-    })
-    manager.handleMuxEnvelope({
-      rpcId: 'title-equal' as never,
-      payload: { type: 'session/title', sessionId: S1, title: 'Equal', eventSeq: 4, updatedAt: 901 },
-    })
+    const titleFrame = (rpcId: string, title: string, seq: number) => {
+      manager.handleMuxEnvelope({
+        rpcId: rpcId as never,
+        payload: { type: 'session/projection', sessionId: S1, key: 'title', value: title, seq } as never,
+      })
+    }
+    titleFrame('title-new', 'Newest', 4)
+    titleFrame('title-stale', 'Stale', 3)
+    titleFrame('title-equal', 'Equal', 4)
     api.onList = () => Promise.resolve(ok({
       items: [summary(S1), summary(S2, { updatedAt: 200 })] as never[],
     }))
@@ -155,7 +152,7 @@ describe('list lifecycle', () => {
 
     const titled = manager.getListSnapshot()
     expect(titled.items.map(item => item.sessionId)).toEqual([S1, S2])
-    expect(titled.items[0]).toMatchObject({ title: 'Newest', updatedAt: 300 })
+    expect(titled.items[0]?.title).toBe('Newest')
     expect(titled.items[1]?.title).toBeUndefined()
 
     manager.handleHostEnvelope({ rpcId: 'removed' as never, payload: { type: 'host/session-removed', sessionId: S1 } })
@@ -163,34 +160,27 @@ describe('list lifecycle', () => {
     expect(manager.getListSnapshot().items.find(item => item.sessionId === S1)?.title).toBeUndefined()
   })
 
-  it('drops a retained title beyond the subscription baseline before accepting its durable replay', async () => {
+  it('drops a projection row beyond the subscription baseline before accepting its durable replay', async () => {
     const api = new FakeApiClient()
     api.onList = () => Promise.resolve(ok({ items: [summary(S1)] as never[] }))
     const manager = new SessionManager(api)
     await manager.refreshList()
-    manager.handleMuxEnvelope({
-      rpcId: 'title-unflushed' as never,
-      payload: { type: 'session/title', sessionId: S1, title: 'Unflushed', eventSeq: 4, updatedAt: 400 },
-    })
+    const frame = (rpcId: string, payload: object) => {
+      manager.handleMuxEnvelope({ rpcId: rpcId as never, payload: payload as never })
+    }
+    frame('title-unflushed', { type: 'session/projection', sessionId: S1, key: 'title', value: 'Unflushed', seq: 4 })
 
-    manager.handleMuxEnvelope({
-      rpcId: 'subscribed-recovered' as never,
-      payload: { type: 'session/subscribed', sessionId: S1, lastSeq: 2 },
-    })
+    // The durable baseline says the host only knows up to seq 2: the phantom
+    // row rode lost state and must drop, or last-wins pins it forever.
+    frame('subscribed-recovered', { type: 'session/subscribed', sessionId: S1, lastSeq: 2 })
     expect(manager.getListSnapshot().items[0]?.title).toBeUndefined()
-    expect(manager.getListSnapshot().items[0]?.updatedAt).toBe(100)
 
-    manager.handleMuxEnvelope({
-      rpcId: 'title-durable' as never,
-      payload: { type: 'session/title', sessionId: S1, title: 'Durable', eventSeq: 2, updatedAt: 200 },
-    })
-    expect(manager.getListSnapshot().items[0]).toMatchObject({ title: 'Durable', updatedAt: 200 })
+    frame('title-durable', { type: 'session/projection', sessionId: S1, key: 'title', value: 'Durable', seq: 2 })
+    expect(manager.getListSnapshot().items[0]?.title).toBe('Durable')
 
-    manager.handleMuxEnvelope({
-      rpcId: 'subscribed-current' as never,
-      payload: { type: 'session/subscribed', sessionId: S1, lastSeq: 2 },
-    })
-    expect(manager.getListSnapshot().items[0]).toMatchObject({ title: 'Durable', updatedAt: 200 })
+    // A baseline at or past the row's seq keeps it (nothing phantom to drop).
+    frame('subscribed-current', { type: 'session/subscribed', sessionId: S1, lastSeq: 2 })
+    expect(manager.getListSnapshot().items[0]?.title).toBe('Durable')
   })
 })
 

+ 1 - 1
packages/client/runtime/tests/sessions-service.spec.ts

@@ -47,7 +47,7 @@ describe('list store projection', () => {
     const b = bench()
     b.svc.handleMuxEnvelope({
       rpcId: 'title' as never,
-      payload: { type: 'session/title', sessionId: sid('s1'), title: 'Durable title', eventSeq: 2, updatedAt: 3 },
+      payload: { type: 'session/projection', sessionId: sid('s1'), key: 'title', value: 'Durable title', seq: 2 } as never,
     })
     await feedList(b, [
       { id: 's1', cwd: '/home/u/proj-a/' },