From ea67e1356ff30d206ba853a5f16ca0989bfd9628 Mon Sep 17 00:00:00 2001 From: Brent Rager Date: Wed, 19 Aug 2026 17:16:30 -0400 Subject: [PATCH] clients: assert unknown events stay ignorable, guard action enums, de-time tests MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Test-only follow-up to #505 (shipped in 1.56.0). No changeset: nothing publishable changes, only the tests around it. The drift guard added in #505 has to encode BOTH tiers of the contract, or the obvious way to satisfy it is to make unknown types an error — which would break the other half. Per the stream_reasoning schema, a client that does not recognize an event MUST ignore it, so a frame from a server newer than this build has to be dropped quietly and leave the turn healthy. That behaviour was already correct in all four clients and is untouched here; it was simply untested, so nothing stopped a future "fix" from turning the catch-all into an error path. Each client now asserts an unrecognised type is neither surfaced to consumers nor fatal to the turn. Confirmed by making Go's dispatch loop fail the turn on unknown frames: the new test catches it. Go had no spec/actions guard, unlike the other three — and its ActionType constants had the same drift that hid submit_interaction. Added, derived from spec/actions/*.schema.json the same way. The dispatch assertions were timing-based, which is a flake risk on a box running at load 125. Go now emits every frame plus the terminal, drains the turn's channel to close, and asserts on the collected slice — no timers at all, in either the passing or the failing path, since the terminal always decodes and a dropped frame shows up as an absent entry rather than a hang. Python's and .NET's remaining waits are hang-detectors rather than assertions, so they move from 2s/5s to 30s: observed wall time for the Python file already ranged 0.9s to 3.7s under load. Verified with go test -race -count=10 and 5 consecutive Python runs at load 125. Co-Authored-By: Claude Fable 5 --- dotnet/tests/InteractionTests.cs | 45 ++++- go/protocol/interaction_test.go | 235 ++++++++++++++++++---- python/tests/test_interactions.py | 49 ++++- typescript/test/dispatch-coverage.test.ts | 13 ++ 4 files changed, 296 insertions(+), 46 deletions(-) diff --git a/dotnet/tests/InteractionTests.cs b/dotnet/tests/InteractionTests.cs index c40235de..45e65f22 100644 --- a/dotnet/tests/InteractionTests.cs +++ b/dotnet/tests/InteractionTests.cs @@ -205,7 +205,7 @@ public async Task InteractionRequiredAndInvalidReachTheTurn() transport.Emit(Terminal(reqId)); await turn.Completion; - await iterate.WaitAsync(TimeSpan.FromSeconds(5)); + await iterate.WaitAsync(TimeSpan.FromSeconds(30)); var types = collected.Select(e => e.Type).ToList(); Assert.Contains("interaction_required", types); @@ -259,7 +259,7 @@ public async Task StreamPreambleAndReasoningReachTheTurn() transport.Emit(Terminal(reqId)); await turn.Completion; - await iterate.WaitAsync(TimeSpan.FromSeconds(5)); + await iterate.WaitAsync(TimeSpan.FromSeconds(30)); var types = collected.Select(e => e.Type).ToList(); Assert.Contains("stream_preamble", types); @@ -346,4 +346,45 @@ await client.SubmitInteractionAsync( var result = _validator.ValidateAt(SubmitRef, transport.Sent[^1]); Assert.True(result.IsValid, result.FormatErrors()); } + + /// + /// The OTHER half of the contract, and the reason the drift guard must not be + /// "fixed" by making unknown types an error: the stream_reasoning schema says + /// clients that do not recognize an event MUST ignore it. A frame from a NEWER + /// server that this build predates has to be dropped silently, leaving the turn + /// healthy and still able to complete. + /// + [Fact] + public async Task UnknownEventIsIgnoredNotFatal() + { + var (client, transport) = MakeClient(); + await client.ConnectAsync(); + + var turn = client.SendMessageAsync(new SendMessageAction { SessionId = "sess-1", Message = "hi" }); + var reqId = transport.LastRequestId(); + + var collected = new List(); + var iterate = Task.Run(async () => + { + await foreach (var ev in turn) + collected.Add(ev); + }); + + transport.Emit(Frame( + """{"type":"stream_hologram","requestId":"{rid}","token":"from the future","data":{"requestId":"{rid}","token":"from the future"}}""", + reqId)); + transport.Emit(Frame( + """{"type":"stream_token","requestId":"{rid}","token":"real","data":{"requestId":"{rid}","token":"real"}}""", + reqId)); + transport.Emit(Terminal(reqId)); + + await turn.Completion; + await iterate.WaitAsync(TimeSpan.FromSeconds(30)); + + var types = collected.Select(e => e.Type).ToList(); + Assert.DoesNotContain("stream_hologram", types); + // The unknown frame must not have derailed the turn: the known events still land. + Assert.Contains("stream_token", types); + Assert.Contains("eventual_response", types); + } } diff --git a/go/protocol/interaction_test.go b/go/protocol/interaction_test.go index e68cbda3..4a4d319a 100644 --- a/go/protocol/interaction_test.go +++ b/go/protocol/interaction_test.go @@ -1,12 +1,13 @@ package protocol import ( + "context" "encoding/json" + "errors" "os" "path/filepath" "strings" "testing" - "time" ) // TestEventTypesCoverSpec is the drift guard. It derives the expected discriminator @@ -28,6 +29,73 @@ func TestEventTypesCoverSpec(t *testing.T) { } } +// TestActionTypesCoverSpec is the same guard for the client→server direction: the +// ActionType constants had drifted too (submit_interaction was missing), which is +// why the verb could not be expressed at all. +func TestActionTypesCoverSpec(t *testing.T) { + known := map[string]struct{}{ + string(ActionCreateConversationSession): {}, + string(ActionSendMessage): {}, + string(ActionGetSession): {}, + string(ActionGetConversationMessages): {}, + string(ActionConfirmToolAction): {}, + string(ActionVerifyOTP): {}, + string(ActionSubmitInteraction): {}, + string(ActionCancel): {}, + string(ActionPing): {}, + } + + specActions := specActionDiscriminators(t) + if len(specActions) == 0 { + t.Fatal("no action schemas discovered in spec/actions") + } + for _, disc := range specActions { + if _, ok := known[disc]; !ok { + t.Errorf("spec/actions declares action %q but the ActionType constants omit it", disc) + } + } +} + +// specActionDiscriminators reads every spec/actions/*.schema.json and returns the +// `const` value of the `action` property on its Request definition. +func specActionDiscriminators(t *testing.T) []string { + t.Helper() + dir := filepath.Join(specDir(t), "actions") + entries, err := os.ReadDir(dir) + if err != nil { + t.Fatalf("read spec/actions: %v", err) + } + var out []string + for _, e := range entries { + if e.IsDir() || !strings.HasSuffix(e.Name(), ".schema.json") { + continue + } + raw, err := os.ReadFile(filepath.Join(dir, e.Name())) + if err != nil { + t.Fatalf("read %s: %v", e.Name(), err) + } + // Action schemas nest the frame under $defs/Request. + var schema struct { + Defs struct { + Request struct { + Properties struct { + Action struct { + Const string `json:"const"` + } `json:"action"` + } `json:"properties"` + } `json:"Request"` + } `json:"$defs"` + } + if err := json.Unmarshal(raw, &schema); err != nil { + t.Fatalf("parse %s: %v", e.Name(), err) + } + if c := schema.Defs.Request.Properties.Action.Const; c != "" { + out = append(out, c) + } + } + return out +} + // specEventDiscriminators reads every spec/events/*.schema.json and returns the // `const` value of its `type` property. func specEventDiscriminators(t *testing.T) []string { @@ -88,19 +156,52 @@ func retarget(t *testing.T, raw json.RawMessage, requestID string) map[string]an return m } -// nextEvent pulls one event off the turn, failing if none arrives. A dropped frame -// (the bug this guards) manifests here as a timeout. -func nextEvent(t *testing.T, turn *MessageTurn) ServerEvent { +// drain collects every event the turn surfaces, returning once the terminal event +// closes the channel. No timer: the terminal eventual_response always decodes (it +// was never part of this bug), so the drain always completes and a dropped frame +// shows up as an absent entry in the returned slice rather than as a timeout. That +// keeps the assertion deterministic on a loaded machine. +func drain(turn *MessageTurn) []ServerEvent { + var out []ServerEvent + for ev := range turn.Events() { + out = append(out, ev) + } + return out +} + +// types maps events to their discriminators, for readable assertions. +func types(evs []ServerEvent) []EventType { + out := make([]EventType, len(evs)) + for i, ev := range evs { + out[i] = ev.Type + } + return out +} + +// find returns the first event of the given type, or fails naming what did arrive. +func find(t *testing.T, evs []ServerEvent, want EventType) ServerEvent { t.Helper() - select { - case ev, ok := <-turn.Events(): - if !ok { - t.Fatal("turn events channel closed before the expected event arrived") + for _, ev := range evs { + if ev.Type == want { + return ev } - return ev - case <-time.After(2 * time.Second): - t.Fatal("timed out waiting for event — the frame was dropped by the dispatch loop") - return ServerEvent{} + } + t.Fatalf("%s was dropped by the dispatch loop; got %v", want, types(evs)) + return ServerEvent{} +} + +// terminalFrame is a minimal eventual_response that settles a turn and closes its +// event channel. +func terminalFrame(requestID string) map[string]any { + return map[string]any{ + "type": "eventual_response", "requestId": requestID, "status": 200, + "data": map[string]any{ + "requestId": requestID, "status": 200, + "data": map[string]any{ + "messageId": "66666666-6666-6666-6666-666666666666", + "response": map[string]any{"responseParts": []any{"done"}}, + }, + }, } } @@ -117,13 +218,16 @@ func TestInteractionFixturesReachTheTurn(t *testing.T) { turn := c.SendMessage(SendMessageParams{SessionID: "sess-1", Message: "quote please"}) reqID := turn.RequestID() - // 1. The park. + // The park, the validation rejection (the turn stays parked, so it must arrive + // too), then the terminal that settles the turn and closes the channel. tr.emit(t, retarget(t, fixtures["interaction_required_event"].Instance, reqID)) - ev := nextEvent(t, turn) - if ev.Type != EventInteractionRequired { - t.Fatalf("event type = %q, want interaction_required", ev.Type) - } - req, err := ev.AsInteractionRequired() + tr.emit(t, retarget(t, fixtures["interaction_invalid_event"].Instance, reqID)) + tr.emit(t, terminalFrame(reqID)) + + evs := drain(turn) + + park := find(t, evs, EventInteractionRequired) + req, err := park.AsInteractionRequired() if err != nil { t.Fatalf("AsInteractionRequired: %v", err) } @@ -135,7 +239,19 @@ func TestInteractionFixturesReachTheTurn(t *testing.T) { t.Errorf("interactionId = %q", interactionID) } - // 2. Answer it, and assert the produced frame is spec-valid. + invalidEv := find(t, evs, EventInteractionInvalid) + inv, err := invalidEv.AsInteractionInvalid() + if err != nil { + t.Fatalf("AsInteractionInvalid: %v", err) + } + if len(inv.Data.Data.Errors) != 1 || inv.Data.Data.Errors[0].Field != "email" { + t.Errorf("errors = %+v, want one error on field email", inv.Data.Data.Errors) + } + if invalidEv.IsTerminal() { + t.Error("interaction_invalid must not be terminal — the turn stays parked for a resubmit") + } + + // Answering the park produces a spec-valid frame. if err := c.SubmitInteraction(SubmitInteractionParams{ SessionID: "22222222-2222-2222-2222-222222222222", RequestID: reqID, @@ -156,23 +272,6 @@ func TestInteractionFixturesReachTheTurn(t *testing.T) { t.Error("declined must stay off the wire when not declining") } validateAgainstSpec(t, "actions/submit-interaction.schema.json#/$defs/Request", sent) - - // 3. Server rejects the values — the turn stays parked, so this must arrive too. - tr.emit(t, retarget(t, fixtures["interaction_invalid_event"].Instance, reqID)) - ev = nextEvent(t, turn) - if ev.Type != EventInteractionInvalid { - t.Fatalf("event type = %q, want interaction_invalid", ev.Type) - } - inv, err := ev.AsInteractionInvalid() - if err != nil { - t.Fatalf("AsInteractionInvalid: %v", err) - } - if len(inv.Data.Data.Errors) != 1 || inv.Data.Data.Errors[0].Field != "email" { - t.Errorf("errors = %+v, want one error on field email", inv.Data.Data.Errors) - } - if ev.IsTerminal() { - t.Error("interaction_invalid must not be terminal — the turn stays parked for a resubmit") - } } // TestSubmitInteractionDeclined covers the decline half of the values-or-declined @@ -269,14 +368,17 @@ func TestEphemeralStreamEventsReachTheTurn(t *testing.T) { }, } + for _, tc := range cases { + validateAgainstSpec(t, tc.schema, tc.frame) + tr.emit(t, tc.frame) + } + tr.emit(t, terminalFrame(reqID)) + + evs := drain(turn) + for _, tc := range cases { t.Run(tc.name, func(t *testing.T) { - validateAgainstSpec(t, tc.schema, tc.frame) - tr.emit(t, tc.frame) - ev := nextEvent(t, turn) - if ev.Type != tc.want { - t.Fatalf("event type = %q, want %q", ev.Type, tc.want) - } + ev := find(t, evs, tc.want) if ev.Token == "" { t.Error("envelope Token not populated") } @@ -287,6 +389,55 @@ func TestEphemeralStreamEventsReachTheTurn(t *testing.T) { } } +// TestUnknownEventIsIgnoredNotFatal is the OTHER half of the contract, and the +// reason the drift guard must not be "fixed" by making unknown types an error: +// the stream_reasoning schema says clients that do not recognize an event MUST +// ignore it. A future server sending a frame this client version predates has to +// be dropped silently, leaving the turn healthy and still able to complete. +func TestUnknownEventIsIgnoredNotFatal(t *testing.T) { + c, tr := makeClient(t) + defer c.Close() + + turn := c.SendMessage(SendMessageParams{SessionID: "sess-1", Message: "hi"}) + reqID := turn.RequestID() + + // A plausible frame from a NEWER server, on a type this build has never heard of. + tr.emit(t, map[string]any{ + "type": "stream_hologram", "requestId": reqID, "token": "from the future", + "data": map[string]any{"requestId": reqID, "token": "from the future"}, + }) + tr.emit(t, map[string]any{ + "type": "stream_token", "requestId": reqID, "token": "real", + "data": map[string]any{"requestId": reqID, "token": "real"}, + }) + tr.emit(t, terminalFrame(reqID)) + + evs := drain(turn) + + for _, ev := range evs { + if ev.Type == "stream_hologram" { + t.Error("an unrecognised event type must not be surfaced to consumers") + } + } + // The unknown frame must not have derailed the turn: the known events still land. + find(t, evs, EventStreamToken) + find(t, evs, EventEventualResponse) + + if _, err := turn.Wait(context.Background()); err != nil { + t.Errorf("turn must settle normally despite an unknown frame, got %v", err) + } + + // And the parser reports it as unknown rather than as malformed JSON. + _, err := ParseServerEvent([]byte(`{"type":"stream_hologram","data":{}}`)) + var unknown *UnknownEventError + if !errors.As(err, &unknown) { + t.Fatalf("ParseServerEvent error = %v, want *UnknownEventError", err) + } + if unknown.Type != "stream_hologram" { + t.Errorf("UnknownEventError.Type = %q", unknown.Type) + } +} + // validateAgainstSpec checks an instance against a spec schema ref, so a test that // asserts on a frame also proves the frame is one the protocol actually allows. func validateAgainstSpec(t *testing.T, ref string, instance any) { diff --git a/python/tests/test_interactions.py b/python/tests/test_interactions.py index 25b96f01..47d987d7 100644 --- a/python/tests/test_interactions.py +++ b/python/tests/test_interactions.py @@ -167,7 +167,7 @@ async def iterate() -> None: # Terminate the turn so the iterator completes. transport.emit(_terminal(req_id)) await turn - await asyncio.wait_for(task, timeout=2.0) + await asyncio.wait_for(task, timeout=30.0) types = [e.type for e in collected] assert "interaction_required" in types, f"interaction_required was dropped; got {types}" @@ -221,7 +221,7 @@ async def iterate() -> None: transport.emit(frame) transport.emit(_terminal(req_id)) await turn - await asyncio.wait_for(task, timeout=2.0) + await asyncio.wait_for(task, timeout=30.0) types = [e.type for e in collected] assert "stream_preamble" in types, f"stream_preamble was dropped; got {types}" @@ -286,3 +286,48 @@ async def test_submit_interaction_carries_choices_values(validator: ProtocolVali assert sent["values"] == values, "choices values lost data in transit" result = validator.validate_at(SUBMIT_REF, sent) assert result.valid, format_errors(result.errors) + + +async def test_unknown_event_is_ignored_not_fatal() -> None: + """The OTHER half of the contract, and the reason the drift guard must not be + "fixed" by making unknown types raise: the stream_reasoning schema says clients + that do not recognize an event MUST ignore it. A frame from a NEWER server that + this build predates has to be dropped silently, leaving the turn healthy.""" + client, transport = make_client() + await client.connect() + + turn = client.send_message(session_id="sess-1", message="hi") + req_id = transport.last_sent()["requestId"] + + collected: list = [] + + async def iterate() -> None: + async for ev in turn: + collected.append(ev) + + task = asyncio.create_task(iterate()) + await asyncio.sleep(0) + + transport.emit( + { + "type": "stream_hologram", + "requestId": req_id, + "token": "from the future", + "data": {"requestId": req_id, "token": "from the future"}, + } + ) + transport.emit( + {"type": "stream_token", "requestId": req_id, "token": "real", "data": {"requestId": req_id, "token": "real"}} + ) + transport.emit(_terminal(req_id)) + await turn + await asyncio.wait_for(task, timeout=30.0) + + types = [e.type for e in collected] + assert "stream_hologram" not in types, "an unrecognised event type must not be surfaced to consumers" + # The unknown frame must not have derailed the turn: the known events still land. + assert "stream_token" in types, f"unknown frame derailed the turn; got {types}" + assert "eventual_response" in types, f"turn did not settle normally; got {types}" + + # And the guard rejects it rather than the parser raising on a malformed frame. + assert not is_server_event({"type": "stream_hologram", "data": {}}) diff --git a/typescript/test/dispatch-coverage.test.ts b/typescript/test/dispatch-coverage.test.ts index dbbfb516..df1482af 100644 --- a/typescript/test/dispatch-coverage.test.ts +++ b/typescript/test/dispatch-coverage.test.ts @@ -68,4 +68,17 @@ describe('dispatch unions cover the spec', () => { const rejected = specActions.filter((action) => !isClientAction({ action })); expect(rejected, `isClientAction() would reject these actions: [${rejected.join(', ')}]`).toEqual([]); }); + + /** + * The OTHER half of the contract, and the reason the guard above must not be + * "satisfied" by making unknown types throw: per the stream_reasoning schema, + * clients that do not recognize an event MUST ignore it. A frame from a NEWER + * server that this build predates has to be rejected quietly by the guard, not + * blow up — `handleFrame` returns early on exactly this. + */ + it('ignores an event type this client version predates', () => { + expect(() => isServerEvent({ type: 'stream_hologram', data: {} })).not.toThrow(); + expect(isServerEvent({ type: 'stream_hologram', data: {} })).toBe(false); + expect(isServerEvent({ type: 'stream_token', data: {} })).toBe(true); + }); });