package httpapi import ( "encoding/json" "net/http" "net/http/httptest" "strings" "testing" "textmachine/platform/internal/pgstore" ) // stream_test.go: the book's event stream. // // Everything pinned here is a rule a client's correctness stands on and none of it is observable // from a single frame: that the sequence a resting book produces is FINITE, that a `note` is never // dropped, and that a client which cannot be continued is told so instead of being silently started // from now. func streamOf(t *testing.T, lib *fakeLibrary, lastEventID string) *httptest.ResponseRecorder { t.Helper() h := readingServer(t, lib) r := httptest.NewRequest("GET", "/v0/books/bk_1/events", nil) r.Header.Set("Authorization", "Bearer token") if lastEventID != "" { r.Header.Set("Last-Event-ID", lastEventID) } w := httptest.NewRecorder() h.ServeHTTP(w, r) return w } // frames splits an SSE body into (event, id, data) triples. func frames(t *testing.T, body string) []struct{ Event, ID, Data string } { t.Helper() var out []struct{ Event, ID, Data string } for _, block := range strings.Split(body, "\n\n") { var f struct{ Event, ID, Data string } for _, line := range strings.Split(block, "\n") { switch { case strings.HasPrefix(line, "event: "): f.Event = strings.TrimPrefix(line, "event: ") case strings.HasPrefix(line, "id: "): f.ID = strings.TrimPrefix(line, "id: ") case strings.HasPrefix(line, "data: "): f.Data = strings.TrimPrefix(line, "data: ") } } if f.Event != "" { out = append(out, f) } } return out } // A book at rest must not become an endless open-and-close. The sequence is bounded BY THE ID on // `hello`: `0` on a book with no history is what the client has to present on the reconnect a // browser makes by itself, and that reconnect is the one answered 204. // // Mutation caught: omitting the id from `hello`; sending `end` without one; answering 200 to the // reconnect of a resting book. func TestARestingBookProducesAFiniteSequence(t *testing.T) { lib := &fakeLibrary{stream: pgstore.StreamState{Position: 0, Revision: 5, StructureVersion: 3, AtRest: true}} w := streamOf(t, lib, "") if w.Code != http.StatusOK { t.Fatalf("status %d", w.Code) } got := frames(t, w.Body.String()) if len(got) != 2 || got[0].Event != "hello" || got[1].Event != "end" { t.Fatalf("a resting book produced %+v", got) } for _, f := range got { if f.ID != "0" { t.Errorf("%s carries id %q, want 0 for a book with no history", f.Event, f.ID) } } var hello map[string]any if err := json.Unmarshal([]byte(got[0].Data), &hello); err != nil { t.Fatal(err) } if hello["contract"] != ContractVersion || hello["structure_version"] != float64(3) { t.Errorf("hello: %v", hello) } // And the reconnect the browser makes with that id is what stops it. if again := streamOf(t, lib, "0"); again.Code != http.StatusNoContent { t.Fatalf("the reconnect of a resting book answered %d, want 204", again.Code) } } // A state frame MAY be replaced by a later one of its kind; `note` is an ADDITION and MUST NOT be // coalesced or dropped — a lost one is lost silently and forever. // // Mutation caught: coalescing by event name without exempting `note`; sending every frame. func TestStateFramesCoalesceAndNotesNeverDo(t *testing.T) { lib := &fakeLibrary{ stream: pgstore.StreamState{Position: 5, Revision: 9, StructureVersion: 1, AtRest: true}, frames: []pgstore.Frame{ {Position: 1, Event: pgstore.FrameProgress, Data: json.RawMessage(`{"progress":{"done":1}}`)}, {Position: 2, Event: pgstore.FrameNote, Data: json.RawMessage(`{"note":{"id":"nt_1"}}`)}, {Position: 3, Event: pgstore.FrameProgress, Data: json.RawMessage(`{"progress":{"done":2}}`)}, {Position: 4, Event: pgstore.FrameNote, Data: json.RawMessage(`{"note":{"id":"nt_2"}}`)}, {Position: 5, Event: pgstore.FrameProgress, Data: json.RawMessage(`{"progress":{"done":3}}`)}, }, } got := frames(t, streamOf(t, lib, "1").Body.String()) var notes, progress int for _, f := range got { switch f.Event { case "note": notes++ case "progress": progress++ } } if notes != 2 { t.Errorf("%d note frames survived coalescing, want both", notes) } if progress != 1 { t.Errorf("%d progress frames, want the newest one only", progress) } // The gap that leaves is legal and declared: a client MUST NOT read a skipped number as a lost // frame. What must hold is that the ids are the BOOK's positions, not a per-connection count. for _, f := range got { if f.Event == "progress" && f.ID != "5" { t.Errorf("the surviving progress frame carries id %q, want the position it holds", f.ID) } } } // A client whose id predates the buffer cannot be continued. Replaying history beyond the buffer is // forbidden — a one-off frame like a note would be lost silently — so it is told to re-read. // // Mutation caught: starting such a client from now without a resync frame. func TestAClientTooFarBehindIsToldToResync(t *testing.T) { lib := &fakeLibrary{stream: pgstore.StreamState{Position: 900, Oldest: 800, Revision: 40, AtRest: true}} got := frames(t, streamOf(t, lib, "17").Body.String()) if len(got) != 2 || got[1].Event != "resync_required" { t.Fatalf("frames: %+v", got) } if got[1].ID != "900" { t.Errorf("the connection frame carries id %q, want the last history frame", got[1].ID) } } // The BOUNDARY of that guard, which is where both of its wrong forms live. `last+1 < Oldest` is the // question "is the frame this client needs next still here", and the two forms that read almost the // same are `last < Oldest` — which re-syncs a client that is exactly caught up — and // `last > 0 && last < Oldest`, which excuses id `0`, the id every connection to a fresh book carries. // // Mutation caught: either of those two. The far-behind case above answers all three identically. func TestTheResyncBoundaryIsExact(t *testing.T) { for _, tc := range []struct { name string state pgstore.StreamState frames []pgstore.Frame last string wantResync bool }{ {"one frame short of the buffer", pgstore.StreamState{Position: 801, Oldest: 800, AtRest: true}, nil, "798", true}, {"exactly at the oldest frame", pgstore.StreamState{Position: 801, Oldest: 800, AtRest: true}, []pgstore.Frame{{Position: 800, Event: pgstore.FrameStatus, Data: json.RawMessage(`{}`)}, {Position: 801, Event: pgstore.FrameStatus, Data: json.RawMessage(`{}`)}}, "799", false}, {"inside the buffer", pgstore.StreamState{Position: 801, Oldest: 800, AtRest: true}, []pgstore.Frame{{Position: 801, Event: pgstore.FrameStatus, Data: json.RawMessage(`{}`)}}, "800", false}, // `0` is a legitimate id and not "no id": it is what `hello` carries on a book with no history. {"the legitimate zero", pgstore.StreamState{Position: 2, Oldest: 1, AtRest: true}, []pgstore.Frame{{Position: 1, Event: pgstore.FrameStatus, Data: json.RawMessage(`{}`)}, {Position: 2, Event: pgstore.FrameStatus, Data: json.RawMessage(`{}`)}}, "0", false}, } { t.Run(tc.name, func(t *testing.T) { body := streamOf(t, &fakeLibrary{stream: tc.state, frames: tc.frames}, tc.last).Body.String() if got := strings.Contains(body, "event: resync_required"); got != tc.wantResync { t.Errorf("Last-Event-ID %s against Oldest %d re-synced=%v, want %v: %s", tc.last, tc.state.Oldest, got, tc.wantResync, body) } }) } } // A prune that lands BEFORE the read shows as a hole at the HEAD of the batch rather than as an // Oldest past the watermark — the read returns the frames after the hole and the watermark jumps // over it, so the second test can never fire. // // Mutation caught: comparing only `state.Oldest > from+1` after advancing, which is the form that // shipped: every frame in the hole, `note` included, was skipped in silence. func TestAHoleAtTheHeadOfTheBatchAsksTheClientToResync(t *testing.T) { lib := &fakeLibrary{ stream: pgstore.StreamState{Position: 100, Oldest: 1, Revision: 11}, frames: []pgstore.Frame{{Position: 289, Event: pgstore.FrameNote, Data: json.RawMessage(`{"note":{"id":"nt_1"}}`)}}, } // While the client was not reading, the writers minted 101..289 and the buffer pruned past them. lib.onSecondRead = func(s *pgstore.StreamState) { s.Oldest, s.Position, s.AtRest = 289, 289, true } body := streamOf(t, lib, "100").Body.String() if !strings.Contains(body, "event: resync_required") { t.Errorf("a hole at the head of the batch was skipped silently: %s", body) } } // HEAD reaches this handler too, and a HEAD has no body to stream: it answers the headers a GET // would and ends, rather than running the pump for a request nobody can read (RFC 9110 §9.3.2). // // Mutation caught: dropping the short-circuit — the handler then streams and holds a goroutine. func TestHeadOnTheStreamAnswersTheHeadersAndNothingElse(t *testing.T) { lib := &fakeLibrary{stream: pgstore.StreamState{Position: 4, Revision: 11, AtRest: true}} r := httptest.NewRequest("HEAD", "/v0/books/bk_1/events", nil) r.Header.Set("Authorization", "Bearer token") w := httptest.NewRecorder() readingServer(t, lib).ServeHTTP(w, r) if w.Code != http.StatusOK || w.Header().Get("Content-Type") != "text/event-stream" { t.Fatalf("HEAD answered %d %q", w.Code, w.Header().Get("Content-Type")) } if w.Body.Len() != 0 { t.Errorf("HEAD carried a body of %d bytes", w.Body.Len()) } } // The book was cut again: every pair anchor and every cursor the client holds is invalid, and a // delta cannot express what happened. func TestAReCutEndsTheStreamWithAResync(t *testing.T) { lib := &fakeLibrary{stream: pgstore.StreamState{Position: 3, Revision: 10, StructureVersion: 3}} // The second read of the state — the one the pump makes after sending — sees the new version. lib.onSecondRead = func(s *pgstore.StreamState) { s.StructureVersion = 4 } got := frames(t, streamOf(t, lib, "").Body.String()) if len(got) != 2 || got[1].Event != "resync_required" { t.Fatalf("frames: %+v", got) } } // A request WITHOUT Last-Event-ID always opens a new stream, even on a resting book: otherwise a // client could not begin watching again after a run starts. func TestAFreshRequestAlwaysOpensAStream(t *testing.T) { lib := &fakeLibrary{stream: pgstore.StreamState{Position: 42, Revision: 9, AtRest: true}} w := streamOf(t, lib, "") if w.Code != http.StatusOK { t.Fatalf("status %d, want a stream", w.Code) } if got := frames(t, w.Body.String()); got[0].Event != "hello" || got[0].ID != "42" { t.Fatalf("frames: %+v", got) } } // The pump's two gap tests, at the width a live connection actually produces. `emitFrame` mints one // frame per write, so an ordinary hole is exactly ONE frame wide; the pin for PD-283 runs a hole of // 189 frames, and both boundaries answer a hole that size identically — a shift of either by one is // invisible to it. // // Mutation caught: `>=` for `>` in either check, or `from+2` for `from+1`. func TestTheGapBoundariesAreExactAtOneFrame(t *testing.T) { frame := func(pos int64) pgstore.Frame { return pgstore.Frame{Position: pos, Event: pgstore.FrameNote, Data: json.RawMessage(`{"note":{"id":"nt_1"}}`)} } for _, tc := range []struct { name string // first is where the catch-up read after `Last-Event-ID: 100` begins; oldest is what the // buffer's floor becomes by the state read that follows it. first int64 oldest int64 wantResync bool }{ {"the very next frame is there", 101, 1, false}, {"exactly one frame missing at the head", 102, 1, true}, {"the buffer still holds the next frame", 101, 102, false}, {"the buffer pruned exactly that one frame", 101, 103, true}, } { t.Run(tc.name, func(t *testing.T) { lib := &fakeLibrary{ stream: pgstore.StreamState{Position: 200, Oldest: 1, Revision: 11}, frames: []pgstore.Frame{frame(tc.first)}, } // The state as it stands once the batch has been consumed: at rest, so a stream with no // gap ends instead of polling. lib.onSecondRead = func(s *pgstore.StreamState) { s.Oldest, s.Position, s.AtRest = tc.oldest, tc.first, true } body := streamOf(t, lib, "100").Body.String() if got := strings.Contains(body, "event: resync_required"); got != tc.wantResync { t.Errorf("re-synced=%v, want %v: %s", got, tc.wantResync, body) } }) } } // EVERY frame of a connection carries the book's revision and its structure version — the // connection frames as much as the history ones (canon §EventBase). // // It is asserted separately because the two halves are stamped in two different places and fail // independently: history frames get them from `emitFrame` inside the transaction that caused them // (pinned in pgstore), while the CONNECTION frames — `hello`, `resync_required`, `end` and, since // canon 0.8.0, `session_ended` — get them from `base()` here. The // battery covered neither until this pack, and this side had the thinner cover of the two — the test // above reads `hello`'s contract and structure version and has never looked at its revision, while // nothing at all looked at `end` or `resync_required`. // // What it costs on the wire is not cosmetic: a client applies a frame only when its revision is at // or above the one it holds, so a frame without one is either dropped by every reader or applied out // of order, depending on how forgiving that reader is. // // Mutation caught: dropping either key from `base()`; building `hello` without the two. func TestEveryConnectionFrameCarriesTheBooksRevisionAndStructureVersion(t *testing.T) { for name, tc := range map[string]struct { state pgstore.StreamState last string want []string }{ "a book at rest": { state: pgstore.StreamState{Position: 0, Revision: 5, StructureVersion: 3, AtRest: true}, want: []string{"hello", "end"}, }, "a watermark below the buffer": { // `Oldest` above `last+1` is the pruned case the canon answers with resync_required. state: pgstore.StreamState{Position: 40, Oldest: 20, Revision: 9, StructureVersion: 4, AtRest: true}, last: "1", want: []string{"hello", "resync_required"}, }, } { t.Run(name, func(t *testing.T) { lib := &fakeLibrary{stream: tc.state} got := frames(t, streamOf(t, lib, tc.last).Body.String()) var seen []string for _, f := range got { if f.Event == "" { continue } seen = append(seen, f.Event) var payload map[string]any if err := json.Unmarshal([]byte(f.Data), &payload); err != nil { t.Fatalf("%s does not decode: %v", f.Event, err) } if payload["revision"] != float64(tc.state.Revision) { t.Errorf("%s carries revision %v, want %d: a client cannot order it against what it holds", f.Event, payload["revision"], tc.state.Revision) } if payload["structure_version"] != float64(tc.state.StructureVersion) { t.Errorf("%s carries structure_version %v, want %d: a client cannot tell whether its chapter ids still exist", f.Event, payload["structure_version"], tc.state.StructureVersion) } } if len(seen) != len(tc.want) { t.Fatalf("the connection produced %v, want %v — this fixture is not exercising what it names", seen, tc.want) } for i, want := range tc.want { if seen[i] != want { t.Fatalf("the connection produced %v, want %v", seen, tc.want) } } }) } } // On a RESUME the hello's id is the id the client PRESENTED, never the head of the book's history: // the browser stores hello's id the moment the frame is dispatched, so a head-stamped hello // followed by a drop moved the client's watermark forward across the whole pending replay — notes // lost silently, no resync_required, and the delta remedy disarmed with them (workflow finding, // P9; PD-406). A fresh connection has no replay to protect and presents the head. // // Mutation caught: sending `hello` with state.Position on a resuming connection. func TestAResumingHelloCarriesTheClientsOwnWatermark(t *testing.T) { lib := &fakeLibrary{ stream: pgstore.StreamState{Position: 5, Revision: 9, StructureVersion: 1, AtRest: true}, frames: []pgstore.Frame{ {Position: 2, Event: pgstore.FrameNote, Data: json.RawMessage(`{"note":{"id":"nt_1"}}`)}, {Position: 5, Event: pgstore.FrameProgress, Data: json.RawMessage(`{"progress":{"done":3}}`)}, }, } got := frames(t, streamOf(t, lib, "1").Body.String()) if len(got) == 0 || got[0].Event != "hello" { t.Fatalf("no hello: %+v", got) } if got[0].ID != "1" { t.Fatalf("a resuming hello carries id %q, want the client's own 1", got[0].ID) } fresh := frames(t, streamOf(t, lib, "").Body.String()) if len(fresh) == 0 || fresh[0].ID != "5" { t.Fatalf("a fresh hello must present the head: %+v", fresh) } }