File size: 6,486 Bytes
7a1ad33
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
/**
 * @license
 * Copyright 2026 Google LLC
 * SPDX-License-Identifier: Apache-2.0
 */

import type {
  AgentProtocol,
  AgentSend,
  AgentEvent,
  Unsubscribe,
} from './types.js';

/**
 * AgentSession is a wrapper around AgentProtocol that provides a more
 * convenient API for consuming agent activity as an AsyncIterable.
 */
export class AgentSession implements AgentProtocol {
  private _protocol: AgentProtocol;

  constructor(protocol: AgentProtocol) {
    this._protocol = protocol;
  }

  async send(payload: AgentSend): Promise<{ streamId: string | null }> {
    return this._protocol.send(payload);
  }

  subscribe(callback: (event: AgentEvent) => void): Unsubscribe {
    return this._protocol.subscribe(callback);
  }

  async abort(): Promise<void> {
    return this._protocol.abort();
  }

  get events(): readonly AgentEvent[] {
    return this._protocol.events;
  }

  /**
   * Sends a payload to the agent and returns an AsyncIterable that yields
   * events for the resulting stream.
   *
   * @param payload The payload to send to the agent.
   */
  async *sendStream(payload: AgentSend): AsyncIterable<AgentEvent> {
    const result = await this._protocol.send(payload);
    const streamId = result.streamId;

    if (streamId === null) {
      return;
    }

    yield* this.stream({ streamId });
  }

  /**
   * Returns an AsyncIterable that yields events from the agent session,
   * optionally replaying events from history or reattaching to an existing stream.
   *
   * @param options Options for replaying or reattaching to the event stream.
   */
  async *stream(
    options: {
      eventId?: string;
      streamId?: string;
    } = {},
  ): AsyncIterable<AgentEvent> {
    let resolve: (() => void) | undefined;
    let next = new Promise<void>((res) => {
      resolve = res;
    });

    let eventQueue: AgentEvent[] = [];
    const earlyEvents: AgentEvent[] = [];
    let done = false;
    let trackedStreamId = options.streamId;
    let started = false;
    let agentActivityStarted = false;

    const queueVisibleEvent = (event: AgentEvent): void => {
      if (trackedStreamId && event.streamId !== trackedStreamId) {
        return;
      }

      if (!agentActivityStarted) {
        if (event.type !== 'agent_start') {
          return;
        }
        trackedStreamId = event.streamId;
        agentActivityStarted = true;
      }

      if (!trackedStreamId) {
        return;
      }

      eventQueue.push(event);
      if (event.type === 'agent_end' && event.streamId === trackedStreamId) {
        done = true;
      }
    };

    // 1. Subscribe early to avoid missing any events that occur during replay setup
    const unsubscribe = this._protocol.subscribe((event) => {
      if (done) return;

      if (!started) {
        earlyEvents.push(event);
        return;
      }

      queueVisibleEvent(event);

      const currentResolve = resolve;
      next = new Promise<void>((r) => {
        resolve = r;
      });
      currentResolve?.();
    });

    try {
      const currentEvents = this._protocol.events;
      let replayStartIndex = -1;

      if (options.eventId) {
        const index = currentEvents.findIndex((e) => e.id === options.eventId);
        if (index === -1) {
          throw new Error(`Unknown eventId: ${options.eventId}`);
        }

        const resumeEvent = currentEvents[index];
        trackedStreamId = resumeEvent.streamId;
        const firstAgentStartIndex = currentEvents.findIndex(
          (event) =>
            event.type === 'agent_start' && event.streamId === trackedStreamId,
        );

        if (resumeEvent.type === 'agent_end') {
          replayStartIndex = index + 1;
          agentActivityStarted = true;
          done = true;
        } else if (
          firstAgentStartIndex !== -1 &&
          firstAgentStartIndex <= index
        ) {
          replayStartIndex = index + 1;
          agentActivityStarted = true;
        } else if (firstAgentStartIndex !== -1) {
          // A pre-agent_start cursor can be resumed once the corresponding
          // agent activity is already present in history. Because stream()
          // yields only agent_start -> agent_end, replay begins at agent_start
          // rather than at the original pre-start event.
          replayStartIndex = firstAgentStartIndex;
        } else {
          // Consumers can only resume by eventId once the corresponding stream
          // has entered the agent_start -> agent_end lifecycle in history.
          // Without a recorded agent_start, this wrapper cannot distinguish
          // "agent activity may start later" from "this send was acknowledged
          // without agent activity" without risking an infinite wait.
          throw new Error(
            `Cannot resume from eventId ${options.eventId} before agent_start for stream ${trackedStreamId}`,
          );
        }
      } else if (options.streamId) {
        const index = currentEvents.findIndex(
          (e) => e.type === 'agent_start' && e.streamId === options.streamId,
        );
        if (index !== -1) {
          replayStartIndex = index;
        }
      } else {
        const activeStarts = currentEvents.filter(
          (e) => e.type === 'agent_start',
        );
        for (let i = activeStarts.length - 1; i >= 0; i--) {
          const start = activeStarts[i];
          if (
            !currentEvents.some(
              (e) => e.type === 'agent_end' && e.streamId === start.streamId,
            )
          ) {
            trackedStreamId = start.streamId;
            replayStartIndex = currentEvents.findIndex(
              (e) => e.id === start.id,
            );
            break;
          }
        }
      }

      if (replayStartIndex !== -1) {
        for (let i = replayStartIndex; i < currentEvents.length; i++) {
          const event = currentEvents[i];
          queueVisibleEvent(event);
          if (done) break;
        }
      }
      started = true;

      // Process events that arrived while we were replaying
      for (const event of earlyEvents) {
        if (done) break;
        queueVisibleEvent(event);
      }

      while (true) {
        if (eventQueue.length > 0) {
          const eventsToYield = eventQueue;
          eventQueue = [];
          for (const event of eventsToYield) {
            yield event;
          }
          continue;
        }

        if (done) break;
        await next;
      }
    } finally {
      unsubscribe();
    }
  }
}