textmachine/platform/internal/pgstore/idempotency_test.go

316 lines
12 KiB
Go

package pgstore
import (
"context"
"crypto/rand"
"encoding/hex"
"errors"
"sync"
"testing"
"time"
)
// idempotency_test.go: the storage half of `Idempotency-Key`, which shipped with no test at all
// while being the thing that decides whether a retry costs a user a second book.
func aKey(user, key string, fingerprint string) IdempotencyKey {
return IdempotencyKey{UserID: user, Method: "POST", Path: "/v0/books", Key: key,
Fingerprint: []byte(fingerprint)}
}
func keyUser(t *testing.T, s *Store, ctx context.Context) string {
t.Helper()
fundedAccount(t, s, ctx, "u1", "10")
return "u1"
}
// The four answers the whole mechanism is: go ahead · replay · in flight · reused.
func TestAClaimAnswersGoAheadThenReplaysWhatTheFirstAttemptSaid(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
now := time.Now().UTC()
k := aKey(user, "k1", "same")
replay, token, err := s.ClaimIdempotency(ctx, k, now)
if err != nil || replay != nil || token == "" {
t.Fatalf("a first claim answered (%v, %q, %v), want the caller to go ahead under a claim of its own", replay, token, err)
}
if _, _, err := s.ClaimIdempotency(ctx, k, now); !errors.Is(err, ErrKeyInFlight) {
t.Errorf("a repeat while the first is running answered %v, want ErrKeyInFlight", err)
}
if err := s.CompleteIdempotency(ctx, k, token, IdempotentResponse{Status: 201, Location: "/v0/books/bk_1",
Body: []byte(`{"id":"bk_1"}`)}, now); err != nil {
t.Fatal(err)
}
replay, _, err = s.ClaimIdempotency(ctx, k, now)
if err != nil {
t.Fatal(err)
}
if replay == nil || replay.Status != 201 || replay.Location != "/v0/books/bk_1" ||
string(replay.Body) != `{"id":"bk_1"}` {
t.Errorf("the replay is %+v, want the first attempt's own answer verbatim", replay)
}
if _, _, err := s.ClaimIdempotency(ctx, aKey(user, "k1", "different"), now); !errors.Is(err, ErrKeyReused) {
t.Errorf("the same key with another request answered %v, want ErrKeyReused", err)
}
}
// ⚠ Two GENUINELY simultaneous first attempts. `for update` on a row that does not exist locks
// nothing, so both reach the insert; without `on conflict do nothing` the loser raised a unique
// violation, which is neither of the contract's conflicts and reached the client as 500.
//
// Mutation caught: removing `on conflict ... do nothing` and the re-read after it.
func TestTwoSimultaneousFirstClaimsAnswerInFlightAndNeverAnUnexpectedError(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
now := time.Now().UTC()
const racers = 8
var wg sync.WaitGroup
start := make(chan struct{})
errs := make([]error, racers)
granted := make([]bool, racers)
for i := range racers {
wg.Add(1)
go func() {
defer wg.Done()
<-start
replay, _, err := s.ClaimIdempotency(ctx, aKey(user, "race", "same"), now)
errs[i], granted[i] = err, err == nil && replay == nil
}()
}
close(start)
wg.Wait()
winners := 0
for i, err := range errs {
switch {
case granted[i]:
winners++
case errors.Is(err, ErrKeyInFlight):
default:
t.Errorf("racer %d answered %v, want the claim or ErrKeyInFlight", i, err)
}
}
if winners != 1 {
t.Errorf("%d racers were granted the claim, want exactly one", winners)
}
}
// ⚠ The scope is (principal, method, path), as the canon ratified: "the same key on another
// operation is ANOTHER key". Not a conflict — a client that mints one key per user action would
// otherwise have two unrelated actions collide.
//
// Mutation caught: narrowing the key to (user, key), which answers `key_reused` there instead.
func TestOneKeyOnAnotherOperationIsAnotherKey(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
fundedAccount(t, s, ctx, "u2", "10")
now := time.Now().UTC()
if _, _, err := s.ClaimIdempotency(ctx, aKey(user, "shared", "books"), now); err != nil {
t.Fatal(err)
}
other := IdempotencyKey{UserID: user, Method: "POST", Path: "/v0/books/bk_1/runs",
Key: "shared", Fingerprint: []byte("runs")}
if replay, _, err := s.ClaimIdempotency(ctx, other, now); err != nil || replay != nil {
t.Errorf("one key on a second operation answered (%v, %v), want a claim of its own", replay, err)
}
// The same key on the same operation with a different request IS the conflict.
if _, _, err := s.ClaimIdempotency(ctx, aKey(user, "shared", "another book"), now); !errors.Is(err, ErrKeyReused) {
t.Errorf("the same key and operation with another request answered %v, want ErrKeyReused", err)
}
if _, _, err := s.ClaimIdempotency(ctx, aKey("u2", "shared", "books"), now); err != nil {
t.Errorf("another account's identical key answered %v, want a claim of its own", err)
}
}
// A path long enough to overflow a btree tuple used to make the claim itself fail, which reached the
// client as 500 on the one case the header exists for.
//
// ⚠ The address is INCOMPRESSIBLE. A repeated byte string of the same length is squeezed inside the
// index tuple and never reaches the limit, so the pin passed with the digest removed — it was
// measuring the compressor, not the bound.
//
// Mutation caught: returning `path` to the key.
func TestAnOverlongPathDoesNotBreakTheClaim(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
k := aKey(user, "long", "same")
noise := make([]byte, 2000)
if _, err := rand.Read(noise); err != nil {
t.Fatal(err)
}
k.Path = "/v0/books/" + hex.EncodeToString(noise) + "/runs"
if _, _, err := s.ClaimIdempotency(ctx, k, time.Now().UTC()); err != nil {
t.Errorf("a claim under a very long path answered %v", err)
}
}
// A claim has an OWNER, and an attempt that lost it may neither answer for it nor take it away.
//
// Both halves were reproduced on a live database. A first attempt outlives the window; a retry takes
// the key over and starts creating the book; then the first one finally returns. Addressed by
// identity alone it would DELETE the successor's row — and a third attempt would then get a fresh
// claim and do the work a second time — or write its OWN receipt over it, so the client replays a
// `Location` for a book that was never created. `finished_at is null` narrows neither: the
// successor's row is legitimately unfinished.
//
// Mutation caught: dropping `and claim_token = $N` from either writer.
func TestAnAttemptThatLostItsClaimCanNeitherAnswerForItNorTakeItAway(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
now := time.Now().UTC()
k := aKey(user, "taken-over", "same")
_, first, err := s.ClaimIdempotency(ctx, k, now)
if err != nil {
t.Fatal(err)
}
replay, second, err := s.ClaimIdempotency(ctx, k, now.Add(ClaimStale+time.Second))
if err != nil || replay != nil {
t.Fatalf("the takeover answered (%v, %v), want a claim of its own", replay, err)
}
if first == second || second == "" {
t.Fatalf("the takeover kept the token %q: this fixture cannot tell the two attempts apart", second)
}
// The attempt that lost the claim finally returns. Neither of its endings may touch the row.
if err := s.CompleteIdempotency(ctx, k, first,
IdempotentResponse{Status: 201, Location: "/v0/books/ghost"}, now); !errors.Is(err, ErrClaimLost) {
t.Errorf("a superseded attempt recorded its own answer: %v, want ErrClaimLost", err)
}
if err := s.ReleaseIdempotency(ctx, k, first); !errors.Is(err, ErrClaimLost) {
t.Errorf("a superseded attempt released a claim it no longer held: %v, want ErrClaimLost", err)
}
// …and the successor is still holding a live claim, which is what the release would have destroyed.
if err := s.CompleteIdempotency(ctx, k, second,
IdempotentResponse{Status: 201, Location: "/v0/books/bk_real"}, now); err != nil {
t.Fatalf("the attempt that HOLDS the claim could not record its answer: %v", err)
}
got, _, err := s.ClaimIdempotency(ctx, k, now)
if err != nil {
t.Fatal(err)
}
if got == nil || got.Location != "/v0/books/bk_real" {
t.Errorf("the replay is %+v, want the answer of the attempt that did the work", got)
}
}
// An attempt that did NOT complete gives the key back, so the retry the contract advises after a
// 408 is a repeat of the first call rather than a second book.
func TestAReleasedKeyMayBePresentedAgain(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
now := time.Now().UTC()
k := aKey(user, "released", "same")
_, token, err := s.ClaimIdempotency(ctx, k, now)
if err != nil {
t.Fatal(err)
}
if err := s.ReleaseIdempotency(ctx, k, token); err != nil {
t.Fatal(err)
}
replay, _, err := s.ClaimIdempotency(ctx, k, now)
if err != nil || replay != nil {
t.Errorf("after a release the key answered (%v, %v), want a fresh claim", replay, err)
}
}
// A claim whose process died is taken over rather than held forever — and only after the window.
func TestAStaleClaimIsTakenOverOnlyAfterItsWindow(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
now := time.Now().UTC()
k := aKey(user, "stale", "same")
if _, _, err := s.ClaimIdempotency(ctx, k, now); err != nil {
t.Fatal(err)
}
if _, _, err := s.ClaimIdempotency(ctx, k, now.Add(ClaimStale-time.Second)); !errors.Is(err, ErrKeyInFlight) {
t.Errorf("inside the window the claim answered %v, want ErrKeyInFlight", err)
}
replay, _, err := s.ClaimIdempotency(ctx, k, now.Add(ClaimStale+time.Second))
if err != nil || replay != nil {
t.Errorf("past the window the claim answered (%v, %v), want a takeover", replay, err)
}
}
// The retention sweep forgets what is past the window and keeps what is not.
func TestTheSweepForgetsOnlyWhatIsPastTheWindow(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
now := time.Now().UTC()
if _, _, err := s.ClaimIdempotency(ctx, aKey(user, "old", "a"), now.Add(-48*time.Hour)); err != nil {
t.Fatal(err)
}
if _, _, err := s.ClaimIdempotency(ctx, aKey(user, "new", "b"), now); err != nil {
t.Fatal(err)
}
n, err := s.SweepIdempotency(ctx, now.Add(-25*time.Hour))
if err != nil {
t.Fatal(err)
}
if n != 1 {
t.Errorf("the sweep forgot %d records, want 1", n)
}
if _, _, err := s.ClaimIdempotency(ctx, aKey(user, "new", "b"), now); !errors.Is(err, ErrKeyInFlight) {
t.Errorf("a record inside the window was forgotten: %v", err)
}
}
// Giving up on a contended key is a CONTRACTUAL answer, and this one line is what the HTTP layer's
// whole mapping rests on: `ErrKeyContended` wraps `ErrKeyInFlight`, so a caller that only knows the
// four documented answers still recognises it and still has a remedy (PD-369).
//
// Cheap, no database, no flake — and it is the property, not a restatement of it: unwrap the sentinel
// and the door answers 500 on a legitimate request again.
//
// Mutation caught: declaring ErrKeyContended with errors.New instead of wrapping.
func TestGivingUpOnAContendedKeyIsOneOfTheContractsOwnAnswers(t *testing.T) {
if !errors.Is(ErrKeyContended, ErrKeyInFlight) {
t.Fatal("ErrKeyContended does not answer as ErrKeyInFlight: a claim that lost the race too" +
" often reaches the wire as an internal error, on a request that was never wrong")
}
}
// A claim can lose a RACE rather than a conflict, and losing a race is not one of the four answers.
// Two attempts arrive together; one wins the insert and then fails fast — a 4xx gives the key back
// immediately — and the loser's re-read, on a fresh READ COMMITTED snapshot, finds nothing at all.
// Answered rather than retried, that is a 500 on precisely the case this header exists for.
//
// Mutation caught: returning the re-read's ErrNoRows instead of starting the claim over.
func TestAClaimThatLostARaceToAReleaseIsRetriedAndNotAnError(t *testing.T) {
s, ctx := testDB(t)
user := keyUser(t, s, ctx)
now := time.Now().UTC()
const racers = 12
var wg sync.WaitGroup
start := make(chan struct{})
errs := make([]error, racers)
for i := range racers {
wg.Add(1)
go func() {
defer wg.Done()
<-start
k := aKey(user, "raced", "same")
replay, token, err := s.ClaimIdempotency(ctx, k, now)
errs[i] = err
if err != nil || replay != nil {
return
}
// Granted, and this attempt fails fast — which is what every 4xx does.
if err := s.ReleaseIdempotency(ctx, k, token); err != nil {
errs[i] = err
}
}()
}
close(start)
wg.Wait()
for i, err := range errs {
switch {
case err == nil, errors.Is(err, ErrKeyInFlight), errors.Is(err, ErrKeyReused):
default:
t.Errorf("racer %d answered %v, want a claim or one of the contract's own conflicts", i, err)
}
}
}