From 53694d90034a61f84d5d20ffa70c03c89455b8de Mon Sep 17 00:00:00 2001 From: luoye520ww <100058663+luoye520ww@users.noreply.github.com> Date: Tue, 14 Jul 2026 05:53:50 +0800 Subject: [PATCH 1/2] feat(sse): define sequenced event contract --- src/shared/sse-sequence.test.ts | 46 +++++++++++++++++++++++++++++++++ src/shared/sse-sequence.ts | 11 ++++++++ 2 files changed, 57 insertions(+) create mode 100644 src/shared/sse-sequence.test.ts create mode 100644 src/shared/sse-sequence.ts diff --git a/src/shared/sse-sequence.test.ts b/src/shared/sse-sequence.test.ts new file mode 100644 index 000000000..d3b3cb061 --- /dev/null +++ b/src/shared/sse-sequence.test.ts @@ -0,0 +1,46 @@ +import { describe, expect, it } from 'vitest' +import { SequencedRuntimeEventSchema } from './sse-sequence' + +describe('SequencedRuntimeEventSchema', () => { + it('accepts a stable stream sequence envelope', () => { + const event = { + streamId: 'stream_123', + sequence: 42, + eventId: 'evt_42', + occurredAt: '2026-07-14T00:00:00.000Z' + } + + expect(SequencedRuntimeEventSchema.parse(event)).toEqual(event) + }) + + it('accepts sequence zero for a newly-created stream', () => { + expect(SequencedRuntimeEventSchema.parse({ + streamId: 'stream_123', + sequence: 0, + eventId: 'evt_0', + occurredAt: '2026-07-14T00:00:00+00:00' + }).sequence).toBe(0) + }) + + it('rejects unsafe sequences, invalid timestamps, and unknown fields', () => { + expect(() => SequencedRuntimeEventSchema.parse({ + streamId: 'stream_123', + sequence: Number.MAX_SAFE_INTEGER + 1, + eventId: 'evt', + occurredAt: '2026-07-14T00:00:00.000Z' + })).toThrow() + expect(() => SequencedRuntimeEventSchema.parse({ + streamId: 'stream_123', + sequence: 1, + eventId: 'evt', + occurredAt: 'yesterday' + })).toThrow() + expect(() => SequencedRuntimeEventSchema.parse({ + streamId: 'stream_123', + sequence: 1, + eventId: 'evt', + occurredAt: '2026-07-14T00:00:00.000Z', + payload: {} + })).toThrow() + }) +}) diff --git a/src/shared/sse-sequence.ts b/src/shared/sse-sequence.ts new file mode 100644 index 000000000..55aab68f5 --- /dev/null +++ b/src/shared/sse-sequence.ts @@ -0,0 +1,11 @@ +import { z } from 'zod' + +/** Metadata required to detect duplicate, stale, and missing SSE events. */ +export const SequencedRuntimeEventSchema = z.object({ + streamId: z.string().min(1).max(128), + sequence: z.number().int().nonnegative().safe(), + eventId: z.string().min(1).max(128), + occurredAt: z.string().datetime({ offset: true }) +}).strict() + +export type SequencedRuntimeEvent = z.infer From da71ecc09177bd535d3f36dbc9806a57a1008e45 Mon Sep 17 00:00:00 2001 From: luoye520ww <100058663+luoye520ww@users.noreply.github.com> Date: Tue, 14 Jul 2026 07:23:49 +0800 Subject: [PATCH 2/2] feat(sse): detect sequence gaps and stale events --- src/shared/sse-gap.test.ts | 51 +++++++++++++++++++++ src/shared/sse-gap.ts | 90 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 141 insertions(+) create mode 100644 src/shared/sse-gap.test.ts create mode 100644 src/shared/sse-gap.ts diff --git a/src/shared/sse-gap.test.ts b/src/shared/sse-gap.test.ts new file mode 100644 index 000000000..8bf472728 --- /dev/null +++ b/src/shared/sse-gap.test.ts @@ -0,0 +1,51 @@ +import { describe, expect, it } from 'vitest' +import { inspectSseSequence, SseSequenceDecision, type SseSequenceCursor } from './sse-gap' + +const event = (sequence: number, streamId = 'stream_1') => ({ + streamId, + sequence, + eventId: `event_${sequence}`, + occurredAt: '2026-07-14T00:00:00.000Z' +}) + +describe('inspectSseSequence', () => { + it('accepts the first event and creates a cursor', () => { + const result = inspectSseSequence(null, event(0)) + expect(result.decision).toBe(SseSequenceDecision.Accepted) + expect(result.nextCursor).toEqual({ streamId: 'stream_1', lastSequence: 0 }) + }) + + it('accepts the next contiguous event and advances only the returned cursor', () => { + const cursor: SseSequenceCursor = { streamId: 'stream_1', lastSequence: 0 } + const result = inspectSseSequence(cursor, event(1)) + expect(result.decision).toBe(SseSequenceDecision.Accepted) + expect(result.nextCursor?.lastSequence).toBe(1) + expect(cursor.lastSequence).toBe(0) + }) + + it('classifies duplicate and out-of-order events without moving the cursor', () => { + const cursor = { streamId: 'stream_1', lastSequence: 4 } + expect(inspectSseSequence(cursor, event(4)).decision).toBe(SseSequenceDecision.Duplicate) + expect(inspectSseSequence(cursor, event(2)).decision).toBe(SseSequenceDecision.OutOfOrder) + expect(inspectSseSequence(cursor, event(4)).nextCursor).toEqual(cursor) + }) + + it('reports a gap and leaves projection at the last confirmed sequence', () => { + const cursor = { streamId: 'stream_1', lastSequence: 4 } + const result = inspectSseSequence(cursor, event(7)) + expect(result.decision).toBe(SseSequenceDecision.Gap) + expect(result.expectedSequence).toBe(5) + expect(result.nextCursor).toEqual(cursor) + }) + + it('rejects events from an old or replaced stream', () => { + const cursor = { streamId: 'stream_2', lastSequence: 4 } + const result = inspectSseSequence(cursor, event(5, 'stream_1')) + expect(result.decision).toBe(SseSequenceDecision.OldStream) + expect(result.nextCursor).toEqual(cursor) + }) + + it('validates event metadata before making a projection decision', () => { + expect(() => inspectSseSequence(null, { ...event(0), eventId: '' })).toThrow() + }) +}) diff --git a/src/shared/sse-gap.ts b/src/shared/sse-gap.ts new file mode 100644 index 000000000..78a31f3a8 --- /dev/null +++ b/src/shared/sse-gap.ts @@ -0,0 +1,90 @@ +import { + SequencedRuntimeEventSchema, + type SequencedRuntimeEvent +} from './sse-sequence' + +export type SseSequenceCursor = { + streamId: string + lastSequence: number +} + +export const SseSequenceDecision = { + Accepted: 'accepted', + Duplicate: 'duplicate', + Gap: 'gap', + OutOfOrder: 'out-of-order', + OldStream: 'old-stream' +} as const +export type SseSequenceDecision = typeof SseSequenceDecision[keyof typeof SseSequenceDecision] + +export type SseSequenceInspection = { + decision: SseSequenceDecision + streamId: string + receivedSequence: number + expectedSequence: number | null + nextCursor: SseSequenceCursor | null +} + +/** + * Compares one validated event with a renderer cursor without mutating it. + * A gap must stop projection until a bounded replay or authoritative reload runs. + */ +export function inspectSseSequence( + cursor: SseSequenceCursor | null, + rawEvent: SequencedRuntimeEvent +): SseSequenceInspection { + const event = SequencedRuntimeEventSchema.parse(rawEvent) + if (!cursor) { + return { + decision: SseSequenceDecision.Accepted, + streamId: event.streamId, + receivedSequence: event.sequence, + expectedSequence: event.sequence, + nextCursor: { streamId: event.streamId, lastSequence: event.sequence } + } + } + if (event.streamId !== cursor.streamId) { + return { + decision: SseSequenceDecision.OldStream, + streamId: event.streamId, + receivedSequence: event.sequence, + expectedSequence: null, + nextCursor: cursor + } + } + const expectedSequence = cursor.lastSequence + 1 + if (event.sequence === expectedSequence) { + return { + decision: SseSequenceDecision.Accepted, + streamId: event.streamId, + receivedSequence: event.sequence, + expectedSequence, + nextCursor: { streamId: cursor.streamId, lastSequence: event.sequence } + } + } + if (event.sequence === cursor.lastSequence) { + return { + decision: SseSequenceDecision.Duplicate, + streamId: event.streamId, + receivedSequence: event.sequence, + expectedSequence, + nextCursor: cursor + } + } + if (event.sequence > expectedSequence) { + return { + decision: SseSequenceDecision.Gap, + streamId: event.streamId, + receivedSequence: event.sequence, + expectedSequence, + nextCursor: cursor + } + } + return { + decision: SseSequenceDecision.OutOfOrder, + streamId: event.streamId, + receivedSequence: event.sequence, + expectedSequence, + nextCursor: cursor + } +}