Просмотр исходного кода

perf(session-query): materialize live observation events on first read

Dudu-0223 3 недель назад
Родитель
Сommit
5b91dbfba0

+ 2 - 2
.agents/notes/implemented/architecture/2026-08-25-session-observations-and-projection-owned-client-state.i18n.yaml

@@ -2,5 +2,5 @@
 # side as of the last confirmed-consistent state. Both languages carry equal authority;
 # after editing either side, bring the other along and re-record with:
 #   pnpm run verify-translation-pairing --write .agents/notes/implemented/architecture/2026-08-25-session-observations-and-projection-owned-client-state.md
-2026-08-25-session-observations-and-projection-owned-client-state.md: 71c0aa92ddb2ad5bacb238dcc7f6e44e33f84305
-2026-08-25-session-observations-and-projection-owned-client-state.zh.md: f8d4a629501dd354e1210bccab1a3a3862b8c6c3
+2026-08-25-session-observations-and-projection-owned-client-state.md: ff5628c3051086159d532b71ca3e98c1b41a51f8
+2026-08-25-session-observations-and-projection-owned-client-state.zh.md: e433df28fec86564fe1c1e9e88686a5abe6b9549

+ 1 - 1
.agents/notes/implemented/architecture/2026-08-25-session-observations-and-projection-owned-client-state.md

@@ -51,7 +51,7 @@ Every owner disposes its observation. `retain()` creates another lease over the
 
 ### Source resolution and lifetime
 
-An observation binds all returned fields to one lifecycle witness. Callers do not combine a header from corpus listing, events from persistence, and projections from a later live Session. The selected header and event prefix produce the cursor and projection snapshot together.
+An observation binds all returned fields to one lifecycle witness. Callers do not combine a header from corpus listing, events from persistence, and projections from a later live Session. The selected header and event prefix produce the cursor and projection snapshot together. A live observation fixes its cut as the log length at read time and materializes `events` on the first access; the log only appends, so that prefix is identical however late a consumer reads it, and a consumer that needs only the header, cursor, or projections never copies the log.
 
 Live preference is checked both before and after a cold borrow. The second check closes the race in which an Agent attaches while persistence is loading. If persistence itself reports that a live source won but that source has already detached by the time SessionQuery examines it, resolution restarts instead of publishing an unowned reference.
 

+ 1 - 1
.agents/notes/implemented/architecture/2026-08-25-session-observations-and-projection-owned-client-state.zh.md

@@ -51,7 +51,7 @@ flowchart LR
 
 ### 数据源解析与生命周期
 
-一份 observation 把所有返回字段绑定到同一 lifecycle witness。调用方不会把 corpus list 的 header、persistence 的 events 和稍后 live Session 的 projections 拼在一起。选中的 header 与事件前缀共同产生 cursor 和 projection snapshot。
+一份 observation 把所有返回字段绑定到同一 lifecycle witness。调用方不会把 corpus list 的 header、persistence 的 events 和稍后 live Session 的 projections 拼在一起。选中的 header 与事件前缀共同产生 cursor 和 projection snapshot。live observation 在读取时以日志长度固定 cut,并在首次访问时才物化 `events`;日志只会追加,所以无论消费者多晚读取,该前缀都完全相同,而只需要 header、cursor 或 projections 的消费者永远不会复制日志。
 
 系统在 cold borrow 前后都检查 live 优先级。第二次检查封住 persistence 加载期间 Agent 完成 attach 的竞态。如果 persistence 报告由 live source 胜出,但 SessionQuery 检查时该 source 已经 detach,解析会重新开始,而不是发布一份无人持有的引用。
 

+ 2 - 2
packages/session-query/session-query/README.i18n.yaml

@@ -2,5 +2,5 @@
 # side as of the last confirmed-consistent state. Both languages carry equal authority;
 # after editing either side, bring the other along and re-record with:
 #   pnpm run verify-translation-pairing --write packages/session-query/session-query/README.md
-README.md: 1c2af37c26d675da4a4aab3f71557c4072ac9e78
-README.zh.md: 683088c4907346c5ade099081d2811de0b368a44
+README.md: de11c6a21eb9fa394232ce82cf490cff9e8a0757
+README.zh.md: 38a1a8798dd801809e042e98825a1499149bf836

+ 1 - 1
packages/session-query/session-query/README.md

@@ -108,7 +108,7 @@ The decision history lives in the [unified service decision](../../../.agents/no
 
 ### Observation cache
 
-`observeSession` builds point observations without a listing preflight. The cold path stats the stored session first and consults an own bounded cache keyed by the persistence instance and the `stat` revision: an unchanged revision reuses the restored unpublished Session without re-reading the log; a changed revision, or a replaced persistence instance, reloads through the handle seam and replaces the entry. The cache holds `preparedSessionCacheSize` entries with least-recently-used eviction, entries pinned by active observation leases are never evicted, and a session that goes live mid-read retries the live path.
+`observeSession` builds point observations without a listing preflight. A live observation fixes its cut as the current log length and materializes `events` on first read, so header-, cursor-, or projection-only consumers never copy the log; the log only appends, so a late first read still yields exactly that prefix. The cold path stats the stored session first and consults an own bounded cache keyed by the persistence instance and the `stat` revision: an unchanged revision reuses the restored unpublished Session without re-reading the log; a changed revision, or a replaced persistence instance, reloads through the handle seam and replaces the entry. The cache holds `preparedSessionCacheSize` entries with least-recently-used eviction, entries pinned by active observation leases are never evicted, and a session that goes live mid-read retries the live path.
 
 ### Reads and traces
 

+ 1 - 1
packages/session-query/session-query/README.zh.md

@@ -108,7 +108,7 @@ kind: "package-reference"
 
 ### 观察缓存
 
-`observeSession` 不经过列表预检直接构建定点观察。冷路径先对存储会话执行 `stat`,再查询自有的有界缓存,缓存键为持久化实例加 `stat` 修订:修订未变则复用已恢复的未发布 Session,不再重读日志;修订变化或持久化实例被替换则经 handle 缝重新加载并替换条目。缓存保留 `preparedSessionCacheSize` 个条目并按最久未用淘汰,被活跃观察租约钉住的条目从不被淘汰;读取中途转为实时的会话会重试实时路径。
+`observeSession` 不经过列表预检直接构建定点观察。实时观察以当前日志长度固定 cut,并在首次读取时才物化 `events`,因此只需要 header、cursor 或 projection 的消费者永远不会复制日志;日志只会追加,所以延后的首次读取得到的仍然正好是该前缀。冷路径先对存储会话执行 `stat`,再查询自有的有界缓存,缓存键为持久化实例加 `stat` 修订:修订未变则复用已恢复的未发布 Session,不再重读日志;修订变化或持久化实例被替换则经 handle 缝重新加载并替换条目。缓存保留 `preparedSessionCacheSize` 个条目并按最久未用淘汰,被活跃观察租约钉住的条目从不被淘汰;读取中途转为实时的会话会重试实时路径。
 
 ### 读取与追踪
 

+ 15 - 5
packages/session-query/session-query/src/observation.ts

@@ -1,7 +1,7 @@
 /** Shared live/prepared observations for Session page and lifecycle consumers. */
 
 import type { Context } from '@deepseek-ai/cordis'
-import { SessionLogOffset } from '@deepseek-ai/dsh-session'
+import { SessionLogOffset, SessionSeq } from '@deepseek-ai/dsh-session'
 import type { Session, SessionEvent, SessionHeader, SessionId , SessionLogOffset as SessionLogOffsetType , SessionSeqCursor } from '@deepseek-ai/dsh-session'
 import type SessionPersistence from '@deepseek-ai/dsh-session-persistence'
 import type {
@@ -21,7 +21,11 @@ export interface SessionObservation extends Disposable {
   readonly header: SessionHeader
   /** Exact fork-inherited event count paired with {@link header}. */
   readonly inheritedEventCount: SessionLogOffsetType
-  /** Immutable contiguous events at {@link cursor}. */
+  /**
+   * Immutable contiguous events at {@link cursor}. A live observation
+   * materializes this array on first read, so a consumer that reads only the
+   * header, cursor, or projections never copies the log.
+   */
   readonly events: readonly SessionEvent[]
   /** Last observed event seq, or -1 for an empty log. */
   readonly cursor: SessionSeqCursor
@@ -271,7 +275,10 @@ export class SessionObservationReader {
     session: Session,
     projectionMode: NonNullable<SessionObservationOptions['projectionMode']>,
   ): SessionObservation {
-    const events = session.snapshotEvents()
+    // The cut is the log length now. The log only appends, so the prefix
+    // below `seq` is the same array whenever a consumer first reads `events`.
+    const seq = session.seq
+    let materialized: readonly SessionEvent[] | undefined
     const projections = projectionMode === 'none'
       ? undefined
       : this.ctx.get('sessionProjections')?.snapshot(session)
@@ -281,8 +288,11 @@ export class SessionObservationReader {
         source: 'live',
         header: session.header,
         inheritedEventCount: session.inheritedEventCount,
-        events,
-        cursor: events.at(-1)?.seq ?? -1,
+        get events() {
+          materialized ??= session.snapshotEvents(SessionLogOffset(0), seq)
+          return materialized
+        },
+        cursor: seq === 0 ? -1 : SessionSeq(seq - 1),
         ...projections === undefined ? {} : { projections },
         retain: () => {
           if (disposed) throw new Error(`session observation "${session.id}" is disposed`)

+ 35 - 0
packages/session-query/session-query/tests/observation.spec.ts

@@ -166,6 +166,41 @@ describe('SessionObservationReader live path', () => {
     await ctx.fiber.dispose()
   })
 
+  it('materializes live events only on first read and shares them across leases', async () => {
+    const ctx = await readerContext()
+    const session = ctx.sessions.create(SessionId('live-lazy-events'))
+    session.append('turn/start', { turn: 1 })
+    const snapshotEvents = vi.spyOn(session, 'snapshotEvents')
+    const reader = new SessionObservationReader(ctx)
+
+    using observed = await reader.read(session.id, { projectionMode: 'none' })
+    using retained = observed.retain()
+    expect(observed.cursor).toBe(0)
+    expect(snapshotEvents).not.toHaveBeenCalled()
+
+    expect(retained.events).toBe(observed.events)
+    expect(observed.events.map(event => event.type)).toEqual(['turn/start'])
+    expect(snapshotEvents).toHaveBeenCalledOnce()
+    await ctx.fiber.dispose()
+  })
+
+  it('keeps a live cut fixed when the log grows before events are first read', async () => {
+    const ctx = await readerContext()
+    const session = ctx.sessions.create(SessionId('live-fixed-cut'))
+    session.append('turn/start', { turn: 1 })
+    const reader = new SessionObservationReader(ctx)
+
+    using observed = await reader.read(session.id, { projectionMode: 'none' })
+    session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
+    using later = await reader.read(session.id, { projectionMode: 'none' })
+
+    expect(observed.cursor).toBe(0)
+    expect(observed.events.map(event => event.type)).toEqual(['turn/start'])
+    expect(later.cursor).toBe(1)
+    expect(later.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
+    await ctx.fiber.dispose()
+  })
+
   it('reports a missing session when no persistence service is mounted', async () => {
     const ctx = await readerContext()
     await expect(new SessionObservationReader(ctx).read(SessionId('absent'))).rejects.toMatchObject({