Skip to content

Commit 775afad

Browse files
authored
fix(server): cast finished listener call, remove stale TODO, document execute contract (#553)
# Description Three small SDK quality fixes surfaced while reviewing the codebase: ## SDK changes - **`execution_event_bus`**: Cast `listener.call(this)` for `'finished'` listeners (which are invoked with no args despite the typed `Listener` signature). Clears 2 TS strict-mode issues; `.betterer.results` regression baseline shrinks accordingly. - **`jsonrpc_transport_handler`**: Removed stale TODO about mid-stream SSE errors; the Express layer already catches and writes `formatSSEErrorEvent` (covered by `express_app.spec.ts`). - **`AgentExecutor.execute` JSDoc**: Documented the contract that every streaming turn must begin with a `task` or `message` event (including follow-ups). The server already enforces this; the JSDoc just makes it explicit to implementors.
1 parent 8fa6237 commit 775afad

7 files changed

Lines changed: 96 additions & 76 deletions

File tree

.betterer.results

Lines changed: 1 addition & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -4,10 +4,5 @@
44
// https://phenomnomnominal.github.io/betterer/docs/results-file/#merge
55
//
66
exports[`TypeScript Strict Mode`] = {
7-
value: `{
8-
"src/server/events/execution_event_bus.ts:1102247656": [
9-
[233, 15, 4, "tsc: Expected 2 arguments, but got 1.", "2087764327"],
10-
[254, 15, 4, "tsc: Expected 2 arguments, but got 1.", "2087764327"]
11-
]
12-
}`
7+
value: `{}`
138
};

src/compat/v0_3/server/transports/jsonrpc/jsonrpc_transport_handler.ts

Lines changed: 6 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -111,21 +111,17 @@ export class LegacyJsonRpcTransportHandler {
111111
context
112112
);
113113

114+
// Errors thrown by `agentEventStream` propagate out of the
115+
// generator; the Express layer catches them, logs the failure,
116+
// and writes a final SSE `event: error` frame carrying the
117+
// JSON-RPC error envelope before closing the stream.
114118
return (async function* legacyJsonRpcEventStream(): AsyncGenerator<
115119
LegacyJSONRPCResponse,
116120
void,
117121
undefined
118122
> {
119-
try {
120-
for await (const event of agentEventStream) {
121-
yield toCompatStreamResponse(event, requestId);
122-
}
123-
} catch (streamError) {
124-
console.error(
125-
`Error in agent event stream for ${method} (request ${requestId}):`,
126-
streamError
127-
);
128-
throw streamError;
123+
for await (const event of agentEventStream) {
124+
yield toCompatStreamResponse(event, requestId);
129125
}
130126
})();
131127
}

src/samples/extensions/extensions.ts

Lines changed: 26 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,8 @@ import {
33
RequestContext,
44
ExecutionEventBus,
55
AgentExecutionEvent,
6+
EventListener,
7+
FinishedListener,
68
} from '../../server/index.js';
79

810
export const EXTENSION_URI = 'https://github.com/a2aproject/a2a-js/src/samples/extensions/v1';
@@ -72,18 +74,36 @@ class TimestampingEventQueue implements ExecutionEventBus {
7274
this._delegate.finished();
7375
}
7476

75-
on(eventName: 'event' | 'finished', listener: (event: AgentExecutionEvent) => void): this {
76-
this._delegate.on(eventName, listener);
77+
on(eventName: 'event', listener: EventListener): this;
78+
on(eventName: 'finished', listener: FinishedListener): this;
79+
on(eventName: 'event' | 'finished', listener: EventListener | FinishedListener): this {
80+
if (eventName === 'event') {
81+
this._delegate.on('event', listener as EventListener);
82+
} else {
83+
this._delegate.on('finished', listener as FinishedListener);
84+
}
7785
return this;
7886
}
7987

80-
off(eventName: 'event' | 'finished', listener: (event: AgentExecutionEvent) => void): this {
81-
this._delegate.off(eventName, listener);
88+
off(eventName: 'event', listener: EventListener): this;
89+
off(eventName: 'finished', listener: FinishedListener): this;
90+
off(eventName: 'event' | 'finished', listener: EventListener | FinishedListener): this {
91+
if (eventName === 'event') {
92+
this._delegate.off('event', listener as EventListener);
93+
} else {
94+
this._delegate.off('finished', listener as FinishedListener);
95+
}
8296
return this;
8397
}
8498

85-
once(eventName: 'event' | 'finished', listener: (event: AgentExecutionEvent) => void): this {
86-
this._delegate.once(eventName, listener);
99+
once(eventName: 'event', listener: EventListener): this;
100+
once(eventName: 'finished', listener: FinishedListener): this;
101+
once(eventName: 'event' | 'finished', listener: EventListener | FinishedListener): this {
102+
if (eventName === 'event') {
103+
this._delegate.once('event', listener as EventListener);
104+
} else {
105+
this._delegate.once('finished', listener as FinishedListener);
106+
}
87107
return this;
88108
}
89109

src/server/agent_execution/agent_executor.ts

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,11 @@ export interface AgentExecutor {
55
/**
66
* Executes the agent logic and publishes events to the bus.
77
*
8+
* Every call MUST publish either a `task` or a `message` event as its
9+
* first event — including follow-up turns where `requestContext.task`
10+
* is already set. The server enforces this ordering and rejects
11+
* streams that begin with a `statusUpdate` or `artifactUpdate`.
12+
*
813
* Multi-tenant implementations can read the tenant identifier from
914
* `requestContext.context.tenant`.
1015
*/

src/server/events/execution_event_bus.ts

Lines changed: 46 additions & 33 deletions
Original file line numberDiff line numberDiff line change
@@ -52,11 +52,20 @@ export function assertUnreachableEvent(event: never): never {
5252
/** Event names supported by {@link ExecutionEventBus}. */
5353
export type ExecutionEventName = 'event' | 'finished';
5454

55+
/** Listener for `'event'` notifications, invoked with the published event. */
56+
export type EventListener = (event: AgentExecutionEvent) => void;
57+
58+
/** Listener for `'finished'` notifications, invoked with no arguments. */
59+
export type FinishedListener = () => void;
60+
5561
export interface ExecutionEventBus {
5662
publish(event: AgentExecutionEvent): void;
57-
on(eventName: ExecutionEventName, listener: (event: AgentExecutionEvent) => void): this;
58-
off(eventName: ExecutionEventName, listener: (event: AgentExecutionEvent) => void): this;
59-
once(eventName: ExecutionEventName, listener: (event: AgentExecutionEvent) => void): this;
63+
on(eventName: 'event', listener: EventListener): this;
64+
on(eventName: 'finished', listener: FinishedListener): this;
65+
off(eventName: 'event', listener: EventListener): this;
66+
off(eventName: 'finished', listener: FinishedListener): this;
67+
once(eventName: 'event', listener: EventListener): this;
68+
once(eventName: 'finished', listener: FinishedListener): this;
6069
removeAllListeners(eventName?: ExecutionEventName): this;
6170
finished(): void;
6271
}
@@ -74,7 +83,6 @@ const CustomEventImpl: typeof CustomEvent =
7483
}
7584
} as typeof CustomEvent);
7685

77-
type Listener = (event: AgentExecutionEvent) => void;
7886
type WrappedListener = (e: Event) => void;
7987

8088
// Should always pass for 'event' type events since we control the dispatch
@@ -92,11 +100,11 @@ function isAgentExecutionCustomEvent(e: Event): e is CustomEvent<AgentExecutionE
92100
* `listenerCount`, `rawListeners`, etc. are not available.
93101
*/
94102
export class DefaultExecutionEventBus extends EventTarget implements ExecutionEventBus {
95-
// Separate storage for each event type — both use the interface's
96-
// Listener type but are invoked differently (with event payload vs. no
97-
// arguments).
98-
private readonly eventListeners: Map<Listener, WrappedListener[]> = new Map();
99-
private readonly finishedListeners: Map<Listener, WrappedListener[]> = new Map();
103+
// Separate storage so each event type can hold listeners of its own
104+
// signature: 'event' listeners receive a payload, 'finished' listeners
105+
// are invoked with no arguments.
106+
private readonly eventListeners: Map<EventListener, WrappedListener[]> = new Map();
107+
private readonly finishedListeners: Map<FinishedListener, WrappedListener[]> = new Map();
100108

101109
publish(event: AgentExecutionEvent): void {
102110
this.dispatchEvent(new CustomEventImpl('event', { detail: event }));
@@ -106,29 +114,35 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
106114
this.dispatchEvent(new Event('finished'));
107115
}
108116

109-
on(eventName: ExecutionEventName, listener: (event: AgentExecutionEvent) => void): this {
117+
on(eventName: 'event', listener: EventListener): this;
118+
on(eventName: 'finished', listener: FinishedListener): this;
119+
on(eventName: ExecutionEventName, listener: EventListener | FinishedListener): this {
110120
if (eventName === 'event') {
111-
this.addEventListenerInternal(listener);
121+
this.addEventListenerInternal(listener as EventListener);
112122
} else {
113-
this.addFinishedListenerInternal(listener);
123+
this.addFinishedListenerInternal(listener as FinishedListener);
114124
}
115125
return this;
116126
}
117127

118-
off(eventName: ExecutionEventName, listener: (event: AgentExecutionEvent) => void): this {
128+
off(eventName: 'event', listener: EventListener): this;
129+
off(eventName: 'finished', listener: FinishedListener): this;
130+
off(eventName: ExecutionEventName, listener: EventListener | FinishedListener): this {
119131
if (eventName === 'event') {
120-
this.removeEventListenerInternal(listener);
132+
this.removeEventListenerInternal(listener as EventListener);
121133
} else {
122-
this.removeFinishedListenerInternal(listener);
134+
this.removeFinishedListenerInternal(listener as FinishedListener);
123135
}
124136
return this;
125137
}
126138

127-
once(eventName: ExecutionEventName, listener: (event: AgentExecutionEvent) => void): this {
139+
once(eventName: 'event', listener: EventListener): this;
140+
once(eventName: 'finished', listener: FinishedListener): this;
141+
once(eventName: ExecutionEventName, listener: EventListener | FinishedListener): this {
128142
if (eventName === 'event') {
129-
this.addEventListenerOnceInternal(listener);
143+
this.addEventListenerOnceInternal(listener as EventListener);
130144
} else {
131-
this.addFinishedListenerOnceInternal(listener);
145+
this.addFinishedListenerOnceInternal(listener as FinishedListener);
132146
}
133147
return this;
134148
}
@@ -157,9 +171,9 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
157171

158172
// Listener tracking helpers.
159173

160-
private trackListener(
161-
listenerMap: Map<Listener, WrappedListener[]>,
162-
listener: Listener,
174+
private trackListener<L>(
175+
listenerMap: Map<L, WrappedListener[]>,
176+
listener: L,
163177
wrapped: WrappedListener
164178
): void {
165179
const existing = listenerMap.get(listener);
@@ -170,9 +184,9 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
170184
}
171185
}
172186

173-
private untrackWrappedListener(
174-
listenerMap: Map<Listener, WrappedListener[]>,
175-
listener: Listener,
187+
private untrackWrappedListener<L>(
188+
listenerMap: Map<L, WrappedListener[]>,
189+
listener: L,
176190
wrapped: WrappedListener
177191
): void {
178192
const wrappedList = listenerMap.get(listener);
@@ -189,7 +203,7 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
189203

190204
// 'event' listeners.
191205

192-
private addEventListenerInternal(listener: Listener): void {
206+
private addEventListenerInternal(listener: EventListener): void {
193207
const wrapped: WrappedListener = (e: Event) => {
194208
if (!isAgentExecutionCustomEvent(e)) {
195209
throw new Error('Internal error: expected CustomEvent for "event" type');
@@ -201,7 +215,7 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
201215
this.addEventListener('event', wrapped);
202216
}
203217

204-
private removeEventListenerInternal(listener: Listener): void {
218+
private removeEventListenerInternal(listener: EventListener): void {
205219
const wrappedList = this.eventListeners.get(listener);
206220
if (wrappedList && wrappedList.length > 0) {
207221
const wrapped = wrappedList.pop()!;
@@ -212,7 +226,7 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
212226
}
213227
}
214228

215-
private addEventListenerOnceInternal(listener: Listener): void {
229+
private addEventListenerOnceInternal(listener: EventListener): void {
216230
const wrapped: WrappedListener = (e: Event) => {
217231
if (!isAgentExecutionCustomEvent(e)) {
218232
throw new Error('Internal error: expected CustomEvent for "event" type');
@@ -225,11 +239,10 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
225239
this.addEventListener('event', wrapped, { once: true });
226240
}
227241

228-
// 'finished' listeners. The interface declares listeners as taking an
229-
// `AgentExecutionEvent`, but for 'finished' they're invoked with no
230-
// arguments (matching EventEmitter behaviour).
242+
// 'finished' listeners. Invoked with no arguments; the interface
243+
// declares them as `FinishedListener` so callers get a precise type.
231244

232-
private addFinishedListenerInternal(listener: Listener): void {
245+
private addFinishedListenerInternal(listener: FinishedListener): void {
233246
const wrapped: WrappedListener = () => {
234247
listener.call(this);
235248
};
@@ -238,7 +251,7 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
238251
this.addEventListener('finished', wrapped);
239252
}
240253

241-
private removeFinishedListenerInternal(listener: Listener): void {
254+
private removeFinishedListenerInternal(listener: FinishedListener): void {
242255
const wrappedList = this.finishedListeners.get(listener);
243256
if (wrappedList && wrappedList.length > 0) {
244257
const wrapped = wrappedList.pop()!;
@@ -249,7 +262,7 @@ export class DefaultExecutionEventBus extends EventTarget implements ExecutionEv
249262
}
250263
}
251264

252-
private addFinishedListenerOnceInternal(listener: Listener): void {
265+
private addFinishedListenerOnceInternal(listener: FinishedListener): void {
253266
const wrapped: WrappedListener = () => {
254267
this.untrackWrappedListener(this.finishedListeners, listener, wrapped);
255268
listener.call(this);

src/server/index.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,10 @@ export { RequestContext } from './agent_execution/request_context.js';
88
export type {
99
AgentExecutionEvent,
1010
AgentExecutionEventKind,
11+
EventListener,
1112
ExecutionEventBus,
1213
ExecutionEventName,
14+
FinishedListener,
1315
} from './events/execution_event_bus.js';
1416
export {
1517
AgentEvent,

src/server/transports/jsonrpc/jsonrpc_transport_handler.ts

Lines changed: 10 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -117,32 +117,21 @@ export class JsonRpcTransportHandler {
117117
: this.requestHandler.resubscribe(SubscribeToTaskRequest.fromJSON(params), context);
118118

119119
// Wrap the agent event stream into a JSON-RPC result stream.
120+
// Errors thrown by `agentEventStream` propagate out of the
121+
// generator; the Express layer catches them, logs the failure,
122+
// and writes a final SSE `event: error` frame carrying the
123+
// JSON-RPC error envelope before closing the stream.
120124
return (async function* jsonRpcEventStream(): AsyncGenerator<
121125
JSONRPCResponse,
122126
void,
123127
undefined
124128
> {
125-
try {
126-
for await (const event of agentEventStream) {
127-
yield {
128-
jsonrpc: '2.0',
129-
id: requestId,
130-
result: StreamResponse.toJSON(event),
131-
};
132-
}
133-
} catch (streamError) {
134-
// If the underlying agent stream throws an error, we need to yield a JSONRPCErrorResponse.
135-
// However, an AsyncGenerator is expected to yield JSONRPCResult.
136-
// This indicates an issue with how errors from the agent's stream are propagated.
137-
// For now, log it. The Express layer will handle the generator ending.
138-
console.error(
139-
`Error in agent event stream for ${method} (request ${requestId}):`,
140-
streamError
141-
);
142-
// Ideally, the Express layer should catch this and send a final error to the client if the stream breaks.
143-
// Or, the agentEventStream itself should yield a final error event that gets wrapped.
144-
// For now, we re-throw so it can be caught by the Express layer streaming support.
145-
throw streamError;
129+
for await (const event of agentEventStream) {
130+
yield {
131+
jsonrpc: '2.0',
132+
id: requestId,
133+
result: StreamResponse.toJSON(event),
134+
};
146135
}
147136
})();
148137
} else {

0 commit comments

Comments
 (0)