textmachine/platform/internal/httpapi/stream_session_test.go

265 lines
11 KiB
Go

package httpapi
import (
"context"
"encoding/json"
"errors"
"log/slog"
"net/http"
"net/http/httptest"
"strconv"
"sync"
"testing"
"time"
"textmachine/platform/internal/auth"
"textmachine/platform/internal/pgstore"
)
// stream_session_test.go: the stream's own authentication, which is NOT the guard's.
//
// The guard admits a caller once, at the door. This route then runs for hours on that one decision,
// and until PD-379 nothing ever asked again — so `logout`, `logout-all`, `tmplatformctl revoke` and
// the session's absolute ceiling all stopped at the door and left the channel behind it open. What
// is pinned here is the thing a unit on "the check was called" would not pin: that a stream ALREADY
// RUNNING and ALREADY DELIVERING frames stops, that it says why, and that it says it at the right
// position.
// revocableSessions is a store whose two questions have two answers, because that is the shape of
// the defect: the row the GUARD read was live, and the row a second later is not.
type revocableSessions struct {
mu sync.Mutex
asked int
// liveFor is how many re-checks answer "live" before the session is gone. 1 means the stream
// runs one full iteration and is cut on the next.
liveFor int
// fail, when set, is answered instead of a verdict: the store could not be ASKED, which is not
// the same as an answer and must not be reported to the client as one.
fail error
}
func (s *revocableSessions) Lookup(context.Context, []byte, time.Time) (auth.Session, error) {
// Always live: the guard is deliberately not the mechanism under test. A stream that only ever
// ended because the DOOR closed would be the defect itself, dressed as a fix.
return auth.Session{
UserID: "u1",
IdleExpiresAt: time.Now().Add(time.Hour),
AbsoluteExpiresAt: time.Now().Add(time.Hour),
}, nil
}
func (s *revocableSessions) Touch(context.Context, []byte, time.Time, time.Duration) error {
return nil
}
func (s *revocableSessions) StillLive(context.Context, []byte, time.Time) error {
s.mu.Lock()
defer s.mu.Unlock()
s.asked++
if s.fail != nil {
return s.fail
}
if s.asked > s.liveFor {
return auth.ErrNoSession
}
return nil
}
func (s *revocableSessions) count() int {
s.mu.Lock()
defer s.mu.Unlock()
return s.asked
}
// streamUnder opens the event stream against a caller whose session behaves as `sessions` says.
// `lastEventID` is the client's watermark, because a FRESH connection starts at the head and would
// be handed no history at all — a fixture that proved nothing about a stream that is running.
func streamUnder(t *testing.T, lib *fakeLibrary, sessions auth.SessionStore, lastEventID string) *httptest.ResponseRecorder {
t.Helper()
h := v0ServerWith(t, Deps{
Library: lib,
Auth: &auth.Authenticator{
Sessions: sessions, IdleTTL: time.Hour,
Deny: ProblemHandler(CodeUnauthenticated),
},
})
// ⚠ A DEADLINE ON THE REQUEST, and it is what makes this a usable pin rather than a trap. Without
// the session check the stream of a book that is not at rest never ends — which is PD-379 exactly
// — so a mutation removing the check does not fail these tests, it HANGS them, and the package
// dies on Go's ten-minute timeout with no test named. Measured: that is what the planted mutation
// actually did. A hang is the weakest possible signal in CI, where it reads as infrastructure
// trouble; with this deadline the same mutation fails in two seconds, on the assertion that says
// what it broke.
ctx, cancel := context.WithTimeout(t.Context(), 2*time.Second)
defer cancel()
r := httptest.NewRequest("GET", "/v0/books/bk_1/events", nil).WithContext(ctx)
r.Header.Set("Authorization", "Bearer token")
if lastEventID != "" {
r.Header.Set("Last-Event-ID", lastEventID)
}
w := httptest.NewRecorder()
h.ServeHTTP(w, r)
return w
}
// A RUNNING stream — one that has already delivered frames — is ended by a revocation, and it is
// ended where the client actually got to.
//
// The last assertion is the one that matters most and the one a plausible fix gets wrong: the
// terminal frame carries the watermark this CONNECTION reached, not the head of the book's history.
// Stamping the head would move a client's `Last-Event-ID` forward across every frame this connection
// never delivered, so signing in again would silently skip them — the same trap `hello` fell into
// (PD-406).
//
// Mutations caught: dropping the re-check from the pump; putting it AFTER the frames; sending `end`
// or `resync_required` instead of a name of its own; stamping the terminal frame with
// `state.Position`.
func TestARevokedSessionEndsAStreamThatIsAlreadyRunning(t *testing.T) {
// A FULL batch, which is what puts the second iteration under test where a test can reach it.
// A short batch leaves the pump in its `select`, waiting out a real one-second tick that this
// package's other stream tests are all written to avoid — and the catch-up path a full batch
// takes is the one that repeats the session check WITHOUT the tick's pacing, so it is the
// unfriendlier of the two to be checking anyway.
//
// It also separates the two numbers the terminal frame has to choose between: the client is
// brought to framesPerRead, while the book's head is beyond it.
var backlog []pgstore.Frame
for i := 1; i <= framesPerRead; i++ {
backlog = append(backlog, pgstore.Frame{
Position: int64(i), Event: pgstore.FrameProgress,
Data: json.RawMessage(`{"progress":{"done":1}}`),
})
}
lib := &fakeLibrary{
stream: pgstore.StreamState{Position: 400, Revision: 9, StructureVersion: 1},
frames: backlog,
}
sessions := &revocableSessions{liveFor: 1}
got := frames(t, streamUnder(t, lib, sessions, "0").Body.String())
var delivered int
for _, f := range got {
if f.Event == "progress" || f.Event == "note" {
delivered++
}
}
if delivered == 0 {
t.Fatalf("the stream delivered no frames at all (%v): this fixture proves nothing about a "+
"RUNNING stream — it would pass against a stream cut at the door", got)
}
// ⚠ THE LITERAL, once, on purpose. Every other assertion here goes through the constant, so
// renaming its VALUE — to `end`, or to `resync_required`, or to anything at all — would leave the
// whole file green while the wire said something else entirely. What a client dispatches on is
// the string, and the string is what the canon is being asked to ratify (0.8.0).
if eventSessionEnded != "session_ended" {
t.Errorf("the terminal frame goes on the wire as %q; the name proposed for the canon is "+
"%q, and a client dispatches on the string and not on our constant", eventSessionEnded, "session_ended")
}
last := got[len(got)-1]
if last.Event != eventSessionEnded {
t.Fatalf("the stream ended with %q, want %q: a revoked caller must be told which of the "+
"three endings this is", last.Event, eventSessionEnded)
}
if want := strconv.Itoa(framesPerRead); last.ID != want {
t.Errorf("%s carries id %q, want %q — the position this connection brought the client to, "+
"not the head of the book's history (400)", eventSessionEnded, last.ID, want)
}
for _, f := range got {
if f.Event == "end" {
t.Error("the stream sent `end`: the BOOK did not finish, the session did, and a client " +
"reads `end` as 'stop reconnecting'")
}
}
if n := sessions.count(); n < 2 {
t.Errorf("the session was re-checked %d time(s): a stream that asks once is the defect", n)
}
}
// The re-check runs BEFORE the frames of its own tick, so a caller whose access was taken away gets
// what was already on the wire and nothing further.
//
// Mutation caught: moving the check below `ReadFrames`, which would hand one more batch of the
// book's contents to a caller who has been signed out.
func TestARevokedCallerIsGivenNoFurtherFrames(t *testing.T) {
lib := &fakeLibrary{
stream: pgstore.StreamState{Position: 2, Revision: 1, StructureVersion: 1},
frames: []pgstore.Frame{
{Position: 1, Event: pgstore.FrameProgress, Data: json.RawMessage(`{"progress":{"done":1}}`)},
{Position: 2, Event: pgstore.FrameProgress, Data: json.RawMessage(`{"progress":{"done":2}}`)},
},
}
// Dead on the very first re-check: the guard let this caller in and the revocation landed
// between the door and the first tick.
got := frames(t, streamUnder(t, lib, &revocableSessions{liveFor: 0}, "0").Body.String())
for _, f := range got {
switch f.Event {
case "hello", eventSessionEnded: // the handshake, and the reason it stopped
default:
t.Errorf("a revoked caller received a %q frame: nothing of the book may follow the "+
"revocation", f.Event)
}
}
}
// A store that cannot be ASKED has not answered. The stream ends — an unaskable session is an
// unrevokable stream — but the client is NOT told it was signed out, because it was not.
//
// Mutation caught: treating a store failure as a revocation (which would sign every watcher out on a
// database blip); treating it as "live" (which restores the defect for as long as the failure lasts).
func TestAnUnaskableSessionEndsTheStreamWithoutClaimingARevocation(t *testing.T) {
lib := &fakeLibrary{
stream: pgstore.StreamState{Position: 1, Revision: 1, StructureVersion: 1},
frames: []pgstore.Frame{
{Position: 1, Event: pgstore.FrameProgress, Data: json.RawMessage(`{"progress":{"done":1}}`)},
},
}
got := frames(t, streamUnder(t, lib, &revocableSessions{fail: context.DeadlineExceeded}, "0").Body.String())
for _, f := range got {
if f.Event == eventSessionEnded {
t.Errorf("the client was told its session ended, but the store only failed to answer: "+
"the two are different facts and only one of them is about the caller (%v)", got)
}
if f.Event == "end" {
t.Error("the client was told the BOOK ended on a store failure")
}
}
}
// A stream mounted without its guard is refused, loudly, before a single byte of the body.
//
// The identity and the means to re-check it are ONE value, so "authenticated but unrevokable" is not
// a state a handler can be handed — and this pins that the missing half is refused rather than
// served. A 200 with an empty body would look like a healthy connection to everything watching, and
// what it would actually be is PD-379 restored.
//
// Mutation caught: reading the caller with anything that tolerates a zero principal.
func TestAStreamWithoutItsGuardIsRefusedRatherThanServed(t *testing.T) {
h := &v0{
lib: &fakeLibrary{stream: pgstore.StreamState{Position: 1, Revision: 1, StructureVersion: 1}},
log: slog.New(slog.NewJSONHandler(&nopWriter{}, nil)),
}
r := httptest.NewRequest("GET", "/v0/books/bk_1/events", nil)
r.SetPathValue("bookId", "bk_1")
w := httptest.NewRecorder()
h.streamEvents(w, r)
if w.Code != http.StatusInternalServerError {
t.Errorf("an unguarded event stream answered %d, want %d: serving it would serve a stream "+
"that revocation cannot reach", w.Code, http.StatusInternalServerError)
}
if w.Header().Get("Content-Type") == "text/event-stream" {
t.Error("an unguarded event stream began streaming before refusing")
}
}
// The zero principal — the one a route outside guard() would produce — answers "not live" rather
// than "live". The direction is the whole point: the only way this can be wrong must be the way
// that refuses.
func TestAPrincipalWithNoSessionBehindItIsNotLive(t *testing.T) {
if err := (auth.Principal{UserID: "u1"}).StillLive(t.Context()); !errors.Is(err, auth.ErrNoSession) {
t.Errorf("a principal with no store behind it answered %v, want %v", err, auth.ErrNoSession)
}
}
type nopWriter struct{}
func (nopWriter) Write(p []byte) (int, error) { return len(p), nil }