Эх сурвалжийг харах

Merge pull request #1199 from alvinunreal/perf/stream-event-fast-path

fix(perf): skip stream-delta fan-out and stop consecutive status polls
Alvin 1 өдөр өмнө
parent
commit
f2679131d6

+ 26 - 0
src/hooks/task-session-manager/runtime-status-reconciliation.test.ts

@@ -240,6 +240,32 @@ describe('runtime status reconciliation', () => {
     reconciler.dispose();
   });
 
+  test('routine schedule() during an in-flight lookup does not force an immediate extra pass', async () => {
+    const firstResponse = deferred<unknown>();
+    let lookupCount = 0;
+    const status = mock(() => {
+      lookupCount += 1;
+      if (lookupCount === 1) return firstResponse.promise;
+      return Promise.resolve({ data: { 'child-1': { type: 'busy' } } });
+    });
+    const { board, reconciler } = createReconciler(status);
+
+    const firstReconciliation = reconciler.reconcile();
+    await Promise.resolve();
+    for (let index = 0; index < 20; index += 1) {
+      reconciler.schedule();
+    }
+    firstResponse.resolve({ data: { 'child-1': { type: 'busy' } } });
+    await firstReconciliation;
+
+    expect(status).toHaveBeenCalledTimes(1);
+    expect(board.get('child-1')).toMatchObject({
+      state: 'running',
+      statusUncertain: false,
+    });
+    reconciler.dispose();
+  });
+
   test('does not apply an old status response to a relaunched generation', async () => {
     const response = deferred<unknown>();
     const { board, reconciler } = createReconciler(() => response.promise);

+ 4 - 4
src/hooks/task-session-manager/runtime-status-reconciliation.ts

@@ -58,10 +58,10 @@ export function createRuntimeStatusReconciler(options: {
   function schedule(): void {
     if (disposed) return;
     if (!reconciliationSupported()) return;
-    if (activeReconcile) {
-      rerunRequested = true;
-      return;
-    }
+    // Routine schedule() from the event hook must not force an immediate
+    // extra pass while a lookup is in flight: token-stream deltas would
+    // otherwise collapse the 5s cadence into consecutive host polls.
+    if (activeReconcile) return;
     if (timer) return;
     if (!options.backgroundJobBoard.hasRunningJobs()) {
       return;

+ 30 - 0
src/index.test.ts

@@ -708,6 +708,36 @@ describe('plugin TUI agent activity', () => {
       root: 'fixer',
     });
   });
+
+  test('message.part.delta does not write TUI activity or session model', async () => {
+    await hooks?.['chat.message']?.(
+      {
+        sessionID: 'stream-1',
+        agent: 'orchestrator',
+        model: { providerID: 'openai', modelID: 'gpt-4o' },
+      } as never,
+      {} as never,
+    );
+    const before = readTuiSnapshot(projectDir);
+
+    await hooks?.event?.({
+      event: {
+        type: 'message.part.delta',
+        properties: {
+          sessionID: 'stream-1',
+          messageID: 'msg-1',
+          partID: 'part-1',
+          field: 'text',
+          delta: 'a'.repeat(200),
+        },
+      },
+    } as never);
+
+    const after = readTuiSnapshot(projectDir);
+    expect(after.activeSessions).toEqual(before.activeSessions);
+    expect(after.agentModels).toEqual(before.agentModels);
+    expect(after.updatedAt).toBe(before.updatedAt);
+  });
 });
 
 describe('background task admission model resolution', () => {

+ 15 - 0
src/index.ts

@@ -1239,6 +1239,21 @@ export const OhMyOpenCodeLite: Plugin = async (ctx) => {
     },
 
     event: async (input) => {
+      // Token-stream deltas fire on every reasoning/text chunk. Slim
+      // has no work for them except the multiplexer activity heartbeat
+      // that keeps a child pane from looking idle mid-stream. Skip the
+      // rest of the fan-out. v2 names: session.next.{text,reasoning}.delta.
+      const streamEventType = (input.event as { type?: string } | undefined)
+        ?.type;
+      if (
+        streamEventType === 'message.part.delta' ||
+        streamEventType === 'session.next.text.delta' ||
+        streamEventType === 'session.next.reasoning.delta'
+      ) {
+        await multiplexerSessionManager.onSessionStatus(input.event as never);
+        return;
+      }
+
       await cacheMonitor.event(input);
 
       const event = input.event as {

+ 27 - 17
src/multiplexer/cmux/session-lifecycle.ts

@@ -7,18 +7,22 @@ import { type CmuxSessionRecord, CmuxSessionStore } from './session-state';
 
 export interface CmuxSessionEvent {
   type: string;
-  properties?: {
-    info?: {
-      id?: string;
-      parentID?: string;
-      title?: string;
-      directory?: string;
-      sessionID?: string;
-    };
-    part?: { sessionID?: string };
+  properties?: CmuxSessionEventPayload;
+  /** Live v2 hosts key the payload under `data`. */
+  data?: CmuxSessionEventPayload;
+}
+
+interface CmuxSessionEventPayload {
+  info?: {
+    id?: string;
+    parentID?: string;
+    title?: string;
+    directory?: string;
     sessionID?: string;
-    status?: { type: string };
   };
+  part?: { sessionID?: string };
+  sessionID?: string;
+  status?: { type: string };
 }
 
 interface BackgroundJobs {
@@ -49,6 +53,8 @@ const ACTIVITY_EVENTS = new Set([
   'message.part.updated',
   'message.part.delta',
   'message.part.removed',
+  'session.next.text.delta',
+  'session.next.reasoning.delta',
 ]);
 const MIN_LIFETIME_MS = 10_000;
 const IDLE_CONFIRMATIONS = 3;
@@ -131,8 +137,8 @@ export class CmuxSessionLifecycle {
 
   async onSessionCreated(event: CmuxSessionEvent): Promise<void> {
     if (this.disposed) return;
-    this.claimLatePaneOrphans();
     if (event.type !== 'session.created') return;
+    this.claimLatePaneOrphans();
     const info = event.properties?.info;
     if (!info?.id || !info.parentID) return;
     if (this.permanentlyClosedSessions?.has(info.id)) return;
@@ -171,12 +177,15 @@ export class CmuxSessionLifecycle {
 
   async onSessionStatus(event: CmuxSessionEvent): Promise<void> {
     if (this.disposed) return;
-    this.claimLatePaneOrphans();
+    const isActivity = ACTIVITY_EVENTS.has(event.type);
+    // Token-stream deltas fire per chunk. Claiming orphans is for
+    // lifecycle transitions, not the heartbeat path.
+    if (!isActivity) this.claimLatePaneOrphans();
     const session = this.eventSession(event);
     if (!session) return;
     const owned = this.store.get(session);
     if (!owned || owned.owner !== this.owner) return;
-    if (ACTIVITY_EVENTS.has(event.type)) {
+    if (isActivity) {
       this.activity(session);
       return;
     }
@@ -1017,11 +1026,12 @@ export class CmuxSessionLifecycle {
   }
 
   private eventSession(event: CmuxSessionEvent): string | undefined {
+    const payload = event.data ?? event.properties;
     return (
-      event.properties?.sessionID ??
-      event.properties?.info?.sessionID ??
-      event.properties?.part?.sessionID ??
-      event.properties?.info?.id
+      payload?.sessionID ??
+      payload?.info?.sessionID ??
+      payload?.part?.sessionID ??
+      payload?.info?.id
     );
   }
 

+ 64 - 14
src/v2/interview-bridge.test.ts

@@ -245,7 +245,10 @@ describe('v2 interview bridge', () => {
   });
 
   test('snapshots transcript before downstream part injection', async () => {
-    const bridge = createV2InterviewBridge(createContext());
+    const directory = `.tmp-v2-interview-snap-${Date.now()}`;
+    const bridge = createV2InterviewBridge(createContext(), {
+      outputFolder: directory,
+    } as never);
     const event = {
       sessionID: 'ses_snapshot',
       agent: 'orchestrator',
@@ -256,30 +259,38 @@ describe('v2 interview bridge', () => {
         {
           id: 'answer',
           role: 'user',
-          content: [{ type: 'text', text: 'the answer' }],
+          content: [{ type: 'text', text: markerText('the answer') }],
         },
       ],
     };
 
     await bridge.handleContext(event);
+    const captured = bridge.getTranscript('ses_snapshot');
     event.messages[0].content.push({
       type: 'text',
       text: 'injected by downstream transform',
       synthetic: true,
       metadata: { source: 'bridge-test' },
-    });
-
-    expect(bridge.getTranscript('ses_snapshot')).toEqual([
-      {
-        info: { role: 'user', id: 'answer' },
-        parts: [{ type: 'text', text: 'the answer' }],
-      },
-    ]);
+    } as { type: string; text: string });
+
+    expect(captured[0]?.parts?.[0]?.text).toContain('the answer');
+    expect(
+      captured[0]?.parts?.some(
+        (part) => part.text === 'injected by downstream transform',
+      ),
+    ).toBe(false);
     bridge.dispose();
+    await fs.rm(`${process.cwd()}/${directory}`, {
+      recursive: true,
+      force: true,
+    });
   });
 
   test('projects text events and removes a deleted session', async () => {
-    const bridge = createV2InterviewBridge(createContext());
+    const directory = `.tmp-v2-interview-text-${Date.now()}`;
+    const bridge = createV2InterviewBridge(createContext(), {
+      outputFolder: directory,
+    } as never);
     await bridge.handleContext({
       sessionID: 'ses_text',
       agent: 'orchestrator',
@@ -290,7 +301,7 @@ describe('v2 interview bridge', () => {
         {
           id: 'u',
           role: 'user',
-          content: [{ type: 'text', text: 'hello' }],
+          content: [{ type: 'text', text: markerText('hello') }],
         },
       ],
     });
@@ -316,12 +327,19 @@ describe('v2 interview bridge', () => {
     });
     expect(bridge.getTranscript('ses_text')).toEqual([]);
     bridge.dispose();
+    await fs.rm(`${process.cwd()}/${directory}`, {
+      recursive: true,
+      force: true,
+    });
   });
 
   test('resolves sessionID from live `data`-keyed events (text + deletion)', async () => {
     // Live v2 hosts key the event payload under `data`; reading only
     // `event.properties` left handleEvent dead on live v2 for ALL events.
-    const bridge = createV2InterviewBridge(createContext());
+    const directory = `.tmp-v2-interview-live-${Date.now()}`;
+    const bridge = createV2InterviewBridge(createContext(), {
+      outputFolder: directory,
+    } as never);
     await bridge.handleContext({
       sessionID: 'ses_live',
       agent: 'orchestrator',
@@ -332,7 +350,7 @@ describe('v2 interview bridge', () => {
         {
           id: 'u',
           role: 'user',
-          content: [{ type: 'text', text: 'hello' }],
+          content: [{ type: 'text', text: markerText('hello') }],
         },
       ],
     });
@@ -354,6 +372,38 @@ describe('v2 interview bridge', () => {
     });
     expect(bridge.getTranscript('ses_live')).toEqual([]);
     bridge.dispose();
+    await fs.rm(`${process.cwd()}/${directory}`, {
+      recursive: true,
+      force: true,
+    });
+  });
+
+  test('ignores text streams for sessions that are not interviews', async () => {
+    const bridge = createV2InterviewBridge(createContext());
+    await bridge.handleContext({
+      sessionID: 'ses_plain',
+      agent: 'orchestrator',
+      model: {},
+      system: [],
+      tools: {},
+      messages: [
+        {
+          id: 'u',
+          role: 'user',
+          content: [{ type: 'text', text: 'hello' }],
+        },
+      ],
+    });
+    await bridge.handleEvent({
+      type: 'session.next.text.started',
+      properties: { sessionID: 'ses_plain' },
+    });
+    await bridge.handleEvent({
+      type: 'session.next.text.delta',
+      properties: { sessionID: 'ses_plain', delta: 'ignored' },
+    });
+    expect(bridge.getTranscript('ses_plain')).toEqual([]);
+    bridge.dispose();
   });
 
   test('shares one configured dashboard across multiple v2 sessions', async () => {

+ 22 - 6
src/v2/interview-bridge.ts

@@ -206,15 +206,27 @@ export function createV2InterviewBridge(
     });
   }
 
-  async function handleContext(event: V2SessionContextEvent): Promise<void> {
-    const messages = toInterviewMessages(event);
-    transcripts.set(event.sessionID, messages);
+  function isManagedInterviewSession(sessionID: string): boolean {
+    return (
+      transcripts.has(sessionID) ||
+      Boolean(
+        (dashboardManager?.service ?? service).getActiveInterviewId(sessionID),
+      )
+    );
+  }
 
+  async function handleContext(event: V2SessionContextEvent): Promise<void> {
     const trailing = event.messages.at(-1);
-    if (trailing?.role !== 'user') return;
-    const text = textFromContent(trailing.content);
+    const text =
+      trailing?.role === 'user' ? textFromContent(trailing.content) : '';
     const match = text.match(MARKER_PATTERN);
-    if (!match) return;
+    const managed = isManagedInterviewSession(event.sessionID);
+    if (!match && !managed) return;
+
+    // Capture the current view before executing /interview so resume and
+    // creation can read the history. Ordinary sessions never enter here.
+    transcripts.set(event.sessionID, toInterviewMessages(event));
+    if (!match || trailing?.role !== 'user') return;
 
     const output = {
       parts: [] as Array<{
@@ -281,18 +293,22 @@ export function createV2InterviewBridge(
       ((properties.info as { id?: string } | undefined)?.id ?? '');
     if (!sessionID) return;
 
+    const managed = isManagedInterviewSession(sessionID);
     if (type === 'session.next.text.started') {
+      if (!managed) return;
       activeText.set(sessionID, '');
       beginText(sessionID);
       return;
     }
     if (type === 'session.next.text.delta') {
+      if (!managed) return;
       const text = `${activeText.get(sessionID) ?? ''}${typeof properties.delta === 'string' ? properties.delta : ''}`;
       activeText.set(sessionID, text);
       appendText(sessionID, text);
       return;
     }
     if (type === 'session.next.text.ended') {
+      if (!managed) return;
       const text =
         typeof properties.text === 'string'
           ? properties.text

+ 14 - 3
src/v2/setup.ts

@@ -1771,10 +1771,21 @@ export function createV2Setup(): (ctx: V2Context) => Promise<V2Cleanup> {
               const next = await eventIterator.next();
               if (next.done) break;
               try {
-                // interviewBridge keeps the RAW v2 event; the v1 eventHook
-                // loop iterates raw + synthesized v1 shapes (idle,
-                // early-registration created, message.updated telemetry).
+                // Token-stream deltas: the interview bridge already
+                // gates to managed sessions. Skip permission rules and
+                // v1 synthesis; still deliver the raw event so the
+                // multiplexer heartbeat in the v1 event hook can run.
+                const rawType =
+                  typeof next.value?.type === 'string' ? next.value.type : '';
+                const isStreamDelta =
+                  rawType === 'session.next.text.delta' ||
+                  rawType === 'session.next.reasoning.delta' ||
+                  rawType === 'message.part.delta';
                 await interviewBridge.handleEvent(next.value);
+                if (isStreamDelta) {
+                  if (eventHook) await eventHook({ event: next.value });
+                  continue;
+                }
                 // Child-session permission tightening sees the same RAW
                 // event (before v1-shape synthesis) so it is independent
                 // of v1 event-hook presence.