Skip to content

Commit 37b3cd5

Browse files
committed
fix(protocol): add question and approval lifecycle events to the schema
The question/approval services published seven event types through unchecked casts that were never registered in agentEventSchema, so journals containing them failed validation on read (error 50001) and bricked session loading. Register the schemas, drop the casts, and cover the real journal shapes with regression tests.
1 parent 58d17ac commit 37b3cd5

5 files changed

Lines changed: 240 additions & 7 deletions

File tree

‎packages/protocol/src/__tests__/ws-control.test.ts‎

Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -750,4 +750,82 @@ describe('ws-control — operation registry', () => {
750750
}).success,
751751
).toBe(true);
752752
});
753+
754+
it('accepts journaled question and approval events', () => {
755+
const sessionId = 'session_test';
756+
const epoch = 'epoch_test';
757+
const message = (
758+
type: string,
759+
seq: number,
760+
fields: Record<string, unknown>,
761+
) => ({
762+
type,
763+
seq,
764+
epoch,
765+
session_id: sessionId,
766+
timestamp: TS,
767+
payload: {
768+
type,
769+
agentId: 'main',
770+
sessionId,
771+
...fields,
772+
},
773+
});
774+
const approvalRequest = {
775+
approval_id: 'approval_test',
776+
session_id: sessionId,
777+
turn_id: 1,
778+
tool_call_id: 'tool_call_test',
779+
tool_name: 'shell.run',
780+
action: 'Run a command',
781+
tool_input_display: { kind: 'generic', summary: 'test' },
782+
created_at: TS,
783+
expires_at: '2026-06-04T10:31:00.000Z',
784+
};
785+
786+
const events = [
787+
message('event.question.requested', 1, {
788+
question_id: 'question_test',
789+
session_id: sessionId,
790+
questions: [{
791+
id: 'q_0',
792+
question: 'Which option?',
793+
options: [{ id: 'opt_0_0', label: 'A', description: 'First option' }],
794+
header: 'Choice',
795+
allow_other: true,
796+
other_label: 'Other',
797+
}],
798+
created_at: TS,
799+
expires_at: '2026-06-04T10:31:00.000Z',
800+
}),
801+
message('event.question.answered', 2, {
802+
question_id: 'question_test',
803+
answers: { q_0: 'opt_0_0' },
804+
resolved_at: TS,
805+
}),
806+
message('event.question.dismissed', 3, {
807+
question_id: 'question_test',
808+
dismissed_at: TS,
809+
}),
810+
message('event.question.expired', 4, {
811+
question_id: 'question_test',
812+
}),
813+
message('event.approval.requested', 5, approvalRequest),
814+
message('event.approval.resolved', 6, {
815+
approval_id: 'approval_test',
816+
decision: 'approved',
817+
scope: 'session',
818+
feedback: 'Proceed',
819+
selected_label: 'Approve',
820+
resolved_at: TS,
821+
}),
822+
message('event.approval.expired', 7, {
823+
approval_id: 'approval_test',
824+
}),
825+
];
826+
827+
for (const event of events) {
828+
expect(sessionEventMessageSchema.safeParse(event).success).toBe(true);
829+
}
830+
});
753831
});

‎packages/protocol/src/events.ts‎

Lines changed: 95 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,13 @@ import { sessionSchema, sessionStatusSchema, type Session, type SessionStatus }
66
import { isoDateTimeSchema } from './time';
77
import { configResponseSchema, type ConfigResponse } from './rest/config';
88
import { workspaceSchema, type Workspace } from './workspace';
9+
import {
10+
approvalRequestSchema,
11+
approvalResponseSchema,
12+
type ApprovalRequest,
13+
type ApprovalResponse,
14+
} from './approval';
15+
import { questionRequestSchema, type QuestionRequest } from './question';
916

1017
export interface TokenUsage {
1118
readonly inputOther: number;
@@ -355,6 +362,43 @@ export interface SessionStatusChangedEvent {
355362
readonly current_prompt_id?: string;
356363
}
357364

365+
export interface QuestionRequestedEvent extends QuestionRequest {
366+
readonly type: 'event.question.requested';
367+
}
368+
369+
export interface QuestionAnsweredEvent {
370+
readonly type: 'event.question.answered';
371+
readonly question_id: string;
372+
readonly answers: unknown;
373+
readonly resolved_at: string;
374+
}
375+
376+
export interface QuestionDismissedEvent {
377+
readonly type: 'event.question.dismissed';
378+
readonly question_id: string;
379+
readonly dismissed_at: string;
380+
}
381+
382+
export interface QuestionExpiredEvent {
383+
readonly type: 'event.question.expired';
384+
readonly question_id: string;
385+
}
386+
387+
export interface ApprovalRequestedEvent extends ApprovalRequest {
388+
readonly type: 'event.approval.requested';
389+
}
390+
391+
export interface ApprovalResolvedEvent extends ApprovalResponse {
392+
readonly type: 'event.approval.resolved';
393+
readonly approval_id: string;
394+
readonly resolved_at: string;
395+
}
396+
397+
export interface ApprovalExpiredEvent {
398+
readonly type: 'event.approval.expired';
399+
readonly approval_id: string;
400+
}
401+
358402
export interface ConfigChangedEvent {
359403
readonly type: 'event.config.changed';
360404
readonly changedFields: string[];
@@ -693,6 +737,13 @@ export type AgentEvent =
693737
| WorkspaceUpdatedEvent
694738
| WorkspaceDeletedEvent
695739
| SessionStatusChangedEvent
740+
| QuestionRequestedEvent
741+
| QuestionAnsweredEvent
742+
| QuestionDismissedEvent
743+
| QuestionExpiredEvent
744+
| ApprovalRequestedEvent
745+
| ApprovalResolvedEvent
746+
| ApprovalExpiredEvent
696747
| ConfigChangedEvent
697748
| GoalUpdatedEvent
698749
| SkillActivatedEvent
@@ -1082,6 +1133,43 @@ export const sessionStatusChangedEventSchema = z.object({
10821133
current_prompt_id: z.string().min(1).optional(),
10831134
}) satisfies z.ZodType<SessionStatusChangedEvent>;
10841135

1136+
export const questionRequestedEventSchema = questionRequestSchema.extend({
1137+
type: z.literal('event.question.requested'),
1138+
}) satisfies z.ZodType<QuestionRequestedEvent>;
1139+
1140+
export const questionAnsweredEventSchema = z.object({
1141+
type: z.literal('event.question.answered'),
1142+
question_id: z.string().min(1),
1143+
answers: z.unknown().nullable(),
1144+
resolved_at: isoDateTimeSchema,
1145+
}) satisfies z.ZodType<QuestionAnsweredEvent>;
1146+
1147+
export const questionDismissedEventSchema = z.object({
1148+
type: z.literal('event.question.dismissed'),
1149+
question_id: z.string().min(1),
1150+
dismissed_at: isoDateTimeSchema,
1151+
}) satisfies z.ZodType<QuestionDismissedEvent>;
1152+
1153+
export const questionExpiredEventSchema = z.object({
1154+
type: z.literal('event.question.expired'),
1155+
question_id: z.string().min(1),
1156+
}) satisfies z.ZodType<QuestionExpiredEvent>;
1157+
1158+
export const approvalRequestedEventSchema = approvalRequestSchema.extend({
1159+
type: z.literal('event.approval.requested'),
1160+
}) satisfies z.ZodType<ApprovalRequestedEvent>;
1161+
1162+
export const approvalResolvedEventSchema = approvalResponseSchema.extend({
1163+
type: z.literal('event.approval.resolved'),
1164+
approval_id: z.string().min(1),
1165+
resolved_at: isoDateTimeSchema,
1166+
}) satisfies z.ZodType<ApprovalResolvedEvent>;
1167+
1168+
export const approvalExpiredEventSchema = z.object({
1169+
type: z.literal('event.approval.expired'),
1170+
approval_id: z.string().min(1),
1171+
}) satisfies z.ZodType<ApprovalExpiredEvent>;
1172+
10851173
export const configChangedEventSchema = z.object({
10861174
type: z.literal('event.config.changed'),
10871175
changedFields: z.array(z.string()),
@@ -1406,6 +1494,13 @@ export const agentEventSchema = z.discriminatedUnion('type', [
14061494
workspaceUpdatedEventSchema,
14071495
workspaceDeletedEventSchema,
14081496
sessionStatusChangedEventSchema,
1497+
questionRequestedEventSchema,
1498+
questionAnsweredEventSchema,
1499+
questionDismissedEventSchema,
1500+
questionExpiredEventSchema,
1501+
approvalRequestedEventSchema,
1502+
approvalResolvedEventSchema,
1503+
approvalExpiredEventSchema,
14091504
goalUpdatedEventSchema,
14101505
skillActivatedEventSchema,
14111506
advisorStatusEventSchema,

‎packages/server/src/services/approval/approvalService.ts‎

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -106,7 +106,7 @@ export class ApprovalService extends Disposable implements IApprovalService {
106106
sessionId: req.sessionId,
107107
agentId: req.agentId,
108108
...protocolRequest,
109-
} as unknown as Event;
109+
};
110110

111111
this.eventService.publish(event);
112112

@@ -160,7 +160,7 @@ export class ApprovalService extends Disposable implements IApprovalService {
160160
feedback: response.feedback,
161161
selected_label: response.selectedLabel,
162162
resolved_at: resolvedAt,
163-
} as unknown as Event;
163+
};
164164
this.eventService.publish(resolvedEvent);
165165

166166
p.resolve(response);
@@ -217,7 +217,7 @@ export class ApprovalService extends Disposable implements IApprovalService {
217217
sessionId: p.sessionId,
218218
agentId: 'main',
219219
approval_id: p.approvalId,
220-
} as unknown as Event;
220+
};
221221
this.eventService.publish(expiredEvent);
222222

223223
p.reject(new ApprovalExpiredError(p.approvalId, this._timeoutMs));

‎packages/server/src/services/question/questionService.ts‎

Lines changed: 4 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -104,7 +104,7 @@ export class QuestionService extends Disposable implements IQuestionService {
104104
sessionId: req.sessionId,
105105
agentId: req.agentId,
106106
...protocolRequest,
107-
} as unknown as Event;
107+
};
108108
this.eventService.publish(event);
109109

110110
this.logger.info(
@@ -154,7 +154,7 @@ export class QuestionService extends Disposable implements IQuestionService {
154154
question_id: p.questionId,
155155
answers: response === null ? null : (response as { answers?: unknown }).answers ?? response,
156156
resolved_at: resolvedAt,
157-
} as unknown as Event;
157+
};
158158
this.eventService.publish(answeredEvent);
159159

160160
p.resolve(response);
@@ -174,7 +174,7 @@ export class QuestionService extends Disposable implements IQuestionService {
174174
agentId: 'main',
175175
question_id: p.questionId,
176176
dismissed_at: dismissedAt,
177-
} as unknown as Event;
177+
};
178178
this.eventService.publish(dismissedEvent);
179179

180180
p.resolve(questionDismissedResult());
@@ -230,7 +230,7 @@ export class QuestionService extends Disposable implements IQuestionService {
230230
sessionId: p.sessionId,
231231
agentId: 'main',
232232
question_id: p.questionId,
233-
} as unknown as Event;
233+
};
234234
this.eventService.publish(expiredEvent);
235235

236236
p.reject(new QuestionExpiredError(p.questionId, this._timeoutMs));

‎packages/server/test/services.test.ts‎

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -747,6 +747,66 @@ describe('WSBroadcastService (WS transport pump)', () => {
747747
bus2.dispose();
748748
});
749749

750+
it('reopens and replays a journal containing a question request event', async () => {
751+
const sessionId = 'sid_question_journal';
752+
const requestedEnvelope = {
753+
type: 'event.question.requested',
754+
seq: 2,
755+
epoch: journalEpoch,
756+
session_id: sessionId,
757+
timestamp: '2026-08-02T00:00:00.000Z',
758+
payload: {
759+
type: 'event.question.requested',
760+
agentId: 'main',
761+
sessionId,
762+
question_id: 'question_test',
763+
session_id: sessionId,
764+
questions: [{
765+
id: 'q_0',
766+
question: 'Which option?',
767+
options: [{ id: 'opt_0_0', label: 'A', description: 'First option' }],
768+
header: 'Choice',
769+
allow_other: true,
770+
other_label: 'Other',
771+
}],
772+
created_at: '2026-08-02T00:00:00.000Z',
773+
expires_at: '2026-08-02T00:15:00.000Z',
774+
},
775+
};
776+
writeJournal(sessionId, [
777+
JSON.stringify(header()),
778+
JSON.stringify(eventLine(sessionId, 1)),
779+
JSON.stringify({ kind: 'event', seq: 2, envelope: requestedEnvelope }),
780+
]);
781+
782+
const bus = new EventService();
783+
const broadcast = new WSBroadcastService(
784+
bus,
785+
testLogger,
786+
new FakeSessionClients(),
787+
new FakeConnectionRegistry(),
788+
makeEnv(),
789+
);
790+
791+
await expect(broadcast.getCursor(sessionId)).resolves.toMatchObject({
792+
seq: 2,
793+
epoch: journalEpoch,
794+
});
795+
await expect(broadcast.getBufferedSince(sessionId, { seq: 0, epoch: journalEpoch })).resolves.toMatchObject({
796+
events: [expect.objectContaining({ seq: 1 }), expect.objectContaining({ seq: 2 })],
797+
currentSeq: 2,
798+
epoch: journalEpoch,
799+
});
800+
await expect(broadcast.getSnapshotState(sessionId)).resolves.toMatchObject({
801+
seq: 2,
802+
epoch: journalEpoch,
803+
});
804+
805+
await broadcast.closeJournals();
806+
broadcast.dispose();
807+
bus.dispose();
808+
});
809+
750810
it('volatile events ride the watermark, are flagged, and are not journaled or replayed', async () => {
751811
const clients = new FakeSessionClients();
752812
const c = fakeConn();

0 commit comments

Comments
 (0)