package streamgate import ( "errors" "math" "sync" "testing" ) // TestIngressSnapshot_ForeignHandleRejected verifies that a handle issued by // one builder is rejected by another builder, and that same-length payloads // from different builders are also rejected. Zero handles are rejected. func TestIngressSnapshot_ForeignHandleRejected(t *testing.T) { const maxBytes int64 = 1024 b1, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } defer b1.Close() b2, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } defer b2.Close() // Foreign handle from b2 rejected by b1 _, _ = b1.IssueBackingHandle([]byte("hello")) h2, _ := b2.IssueBackingHandle([]byte("hello")) if err := b1.SetCanonicalWithHandle(h2); err != ErrIngressSnapshotInvalidHandle { t.Errorf("foreign SetCanonicalWithHandle: want ErrIngressSnapshotInvalidHandle, got %v", err) } if err := b1.AddTypedView("x", h2); err != ErrIngressSnapshotInvalidHandle { t.Errorf("foreign AddTypedView: want ErrIngressSnapshotInvalidHandle, got %v", err) } // Same-length foreign handle also rejected (not just length check) b1, err = NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } defer b1.Close() _, _ = b1.IssueBackingHandle([]byte("hello")) if err := b1.SetCanonicalWithHandle(h2); err != ErrIngressSnapshotInvalidHandle { t.Errorf("same-length foreign SetCanonicalWithHandle: want ErrIngressSnapshotInvalidHandle, got %v", err) } // Zero handle rejected if err := b1.SetCanonicalWithHandle(BackingHandleZero); err != ErrIngressSnapshotInvalidHandle { t.Errorf("zero SetCanonicalWithHandle: want ErrIngressSnapshotInvalidHandle, got %v", err) } if err := b1.AddTypedView("x", BackingHandleZero); err != ErrIngressSnapshotInvalidHandle { t.Errorf("zero AddTypedView: want ErrIngressSnapshotInvalidHandle, got %v", err) } } // TestIngressSnapshot_SharedHandleUsesSingleOwnedBacking verifies that when the // same handle is used for canonical and typed views, the backing payload is // stored once and retained bytes reflect a single copy. func TestIngressSnapshot_SharedHandleUsesSingleOwnedBacking(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } defer b.Close() payload := []byte("shared-payload-data") h, err := b.IssueBackingHandle(payload) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if b.retainedBytes != int64(len(payload)) { t.Errorf("after canonical: retained=%d, want %d", b.retainedBytes, len(payload)) } if err := b.AddTypedView("view1", h); err != nil { t.Fatalf("AddTypedView view1: %v", err) } if b.retainedBytes != int64(len(payload)) { t.Errorf("after typed view1: retained=%d, want %d", b.retainedBytes, len(payload)) } if err := b.AddTypedView("view2", h); err != nil { t.Fatalf("AddTypedView view2: %v", err) } if b.retainedBytes != int64(len(payload)) { t.Errorf("after typed view2: retained=%d, want %d", b.retainedBytes, len(payload)) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } acc := snap.Accessor() if acc.RetainedCount() != 1 { t.Errorf("retained count: want 1, got %d", acc.RetainedCount()) } if acc.RetainedBytes() != int64(len(payload)) { t.Errorf("retained bytes: want %d, got %d", len(payload), acc.RetainedBytes()) } if acc.HasCanonical() != true { t.Errorf("HasCanonical: want true") } names := acc.TypedViewNames() if len(names) != 2 || names[0] != "view1" || names[1] != "view2" { t.Errorf("typed view names: want [view1 view2], got %v", names) } canon, err := acc.Canonical() if err != nil { t.Fatalf("Canonical: %v", err) } if string(canon) != string(payload) { t.Errorf("canonical payload: want %q, got %q", payload, canon) } // Builder cannot be reused after Build if err := b.SetCanonicalWithHandle(h); err != ErrIngressSnapshotAlreadyBuilt { t.Errorf("post-Build SetCanonical: want ErrIngressSnapshotAlreadyBuilt, got %v", err) } } // TestIngressSnapshot_UniqueBackingCount verifies that unique backing handles // are counted individually while shared handles are counted once. func TestIngressSnapshot_UniqueBackingCount(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } defer b.Close() h1, err := b.IssueBackingHandle([]byte("alpha")) if err != nil { t.Fatalf("IssueBackingHandle h1: %v", err) } if b.retainedBytes != 5 { t.Errorf("after h1: retained=%d, want 5", b.retainedBytes) } h2, err := b.IssueBackingHandle([]byte("beta-gamma")) if err != nil { t.Fatalf("IssueBackingHandle h2: %v", err) } if b.retainedBytes != 15 { t.Errorf("after h2: retained=%d, want 15", b.retainedBytes) } if err := b.SetCanonicalWithHandle(h1); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } // retained unchanged: already accounted in IssueBackingHandle if b.retainedBytes != 15 { t.Errorf("after canonical: retained=%d, want 15", b.retainedBytes) } if err := b.AddTypedView("v1", h2); err != nil { t.Fatalf("AddTypedView v1: %v", err) } // retained unchanged: h2 already accounted if b.retainedBytes != 15 { t.Errorf("after typed v1: retained=%d, want 15", b.retainedBytes) } // Reuse h1 for another typed view (shared backing) if err := b.AddTypedView("v2", h1); err != nil { t.Fatalf("AddTypedView v2: %v", err) } // retained unchanged: h1 already accounted if b.retainedBytes != 15 { t.Errorf("after typed v2 (shared): retained=%d, want 15", b.retainedBytes) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } acc := snap.Accessor() if acc.RetainedCount() != 2 { t.Errorf("retained count: want 2, got %d", acc.RetainedCount()) } if acc.RetainedBytes() != 15 { t.Errorf("retained bytes: want 15, got %d", acc.RetainedBytes()) } } // TestIngressSnapshot_OverflowPoisonsBuilderAndReleases verifies that when a // backing exceeds the remaining budget, the builder transitions to a terminal // failed state: all retained references are released, and subsequent operations // return ErrIngressSnapshotLimitExceeded. func TestIngressSnapshot_OverflowPoisonsBuilderAndReleases(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h1, err := b.IssueBackingHandle([]byte("small")) if err != nil { t.Fatalf("IssueBackingHandle small: %v", err) } if h1 == BackingHandleZero { t.Fatal("expected valid handle for small payload") } if b.retainedBytes != 5 { t.Errorf("after small: retained=%d, want 5", b.retainedBytes) } payloadSize := int64(100) h2, err := b.IssueBackingHandle(make([]byte, payloadSize)) if h2 != BackingHandleZero { t.Fatal("expected zero handle on overflow") } if !errors.Is(err, ErrIngressSnapshotLimitExceeded) { t.Fatalf("expected ErrIngressSnapshotLimitExceeded, got %v", err) } // Builder is poisoned: all subsequent operations return LimitExceeded if err := b.SetCanonicalWithHandle(h1); err != ErrIngressSnapshotLimitExceeded { t.Errorf("post-overflow SetCanonical: want LimitExceeded, got %v", err) } if err := b.AddTypedView("x", h1); err != ErrIngressSnapshotLimitExceeded { t.Errorf("post-overflow AddTypedView: want LimitExceeded, got %v", err) } if _, err := b.Build(); err != ErrIngressSnapshotLimitExceeded { t.Errorf("post-overflow Build: want LimitExceeded, got %v", err) } // Builder zero-state after overflow if b.backings != nil { t.Error("backings map should be nil after overflow") } if b.typedViews != nil { t.Error("typedViews should be nil after overflow") } if b.typedNames != nil { t.Error("typedNames should be nil after overflow") } if b.canonicalHdl.isValid() { t.Error("canonicalHdl should be zero after overflow") } if b.retainedBytes != 0 { t.Errorf("retainedBytes should be 0 after overflow, got %d", b.retainedBytes) } } // TestIngressSnapshot_OverflowExactBoundary verifies exact-limit success and // limit+1 overflow. func TestIngressSnapshot_OverflowExactBoundary(t *testing.T) { const maxBytes int64 = 100 // Exact limit: success b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle(make([]byte, maxBytes)) if h == BackingHandleZero { t.Fatal("expected valid handle at exact limit") } if err != nil { t.Fatalf("IssueBackingHandle at exact limit: %v", err) } if b.retainedBytes != maxBytes { t.Errorf("retained at exact limit: want %d, got %d", maxBytes, b.retainedBytes) } // limit+1: overflow b, err = NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err = b.IssueBackingHandle(make([]byte, maxBytes+1)) if h != BackingHandleZero { t.Fatal("expected zero handle beyond limit") } if !errors.Is(err, ErrIngressSnapshotLimitExceeded) { t.Fatalf("expected ErrIngressSnapshotLimitExceeded on overflow, got %v", err) } if _, err := b.Build(); err != ErrIngressSnapshotLimitExceeded { t.Errorf("post-overflow Build: want LimitExceeded, got %v", err) } } // TestIngressSnapshot_BuildTransfersBackingOwnership verifies that Build moves // the backing store to the snapshot without copying payloads, and that the // builder cannot be reused. func TestIngressSnapshot_BuildTransfersBackingOwnership(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } original := []byte("evidence-payload-data") h, err := b.IssueBackingHandle(original) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("raw", h); err != nil { t.Fatalf("AddTypedView: %v", err) } // Save the backing pointer before Build backingPtr := b.backings[h] snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } // Builder is poisoned after Build if _, err := b.Build(); err != ErrIngressSnapshotAlreadyBuilt { t.Errorf("re-Build: want AlreadyBuilt, got %v", err) } if err := b.SetCanonicalWithHandle(h); err != ErrIngressSnapshotAlreadyBuilt { t.Errorf("post-Build SetCanonical: want AlreadyBuilt, got %v", err) } // Builder backing store is nil (transferred) if b.backings != nil { t.Error("builder backings should be nil after Build") } if b.typedViews != nil { t.Error("builder typedViews should be nil after Build") } if b.typedNames != nil { t.Error("builder typedNames should be nil after Build") } if b.canonicalHdl.isValid() { t.Error("builder canonicalHdl should be zero after Build") } if b.retainedBytes != 0 { t.Errorf("builder retainedBytes should be 0 after Build, got %d", b.retainedBytes) } // Snapshot owns the same backing pointer (single-copy transfer) acc := snap.Accessor() if !acc.HasCanonical() { t.Error("snapshot should have canonical") } snapBacking := snap.backings[h] if snapBacking != backingPtr { t.Error("snapshot backing should be the same pointer as builder backing (single-copy transfer)") } // Snapshot accessor returns defensive copy canon, err := acc.Canonical() if err != nil { t.Fatalf("Canonical: %v", err) } if string(canon) != string(original) { t.Errorf("canonical: want %q, got %q", original, canon) } // Mutating caller's original does not affect snapshot original[0] = 'X' canon2, err := acc.Canonical() if err != nil { t.Fatalf("Canonical after mutation: %v", err) } if string(canon2) != string(canon) { t.Errorf("snapshot should be isolated from caller mutation") } // Snapshot Close releases backing snap.Close() if _, err := acc.Canonical(); err != ErrIngressSnapshotClosed { t.Errorf("post-close Canonical: want Closed, got %v", err) } } // TestIngressSnapshot_ZeroLengthCanonical verifies that a zero-length canonical // payload is preserved when explicitly set (not treated as "no canonical"). func TestIngressSnapshot_ZeroLengthCanonical(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte{}) if h == BackingHandleZero { t.Fatal("expected valid handle for zero-length payload") } if err != nil { t.Fatalf("IssueBackingHandle zero-length: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } acc := snap.Accessor() if !acc.HasCanonical() { t.Error("HasCanonical should be true for explicitly set zero-length canonical") } canon, err := acc.Canonical() if err != nil { t.Fatalf("Canonical: %v", err) } if len(canon) != 0 { t.Errorf("canonical length: want 0, got %d", len(canon)) } // Retained bytes should include the zero-length payload if acc.RetainedBytes() != 0 { t.Errorf("retained bytes: want 0, got %d", acc.RetainedBytes()) } } // TestIngressSnapshot_NoSnapshotOnOverflow verifies that overflow prevents // snapshot creation and isolates the failure from other builders. func TestIngressSnapshot_NoSnapshotOnOverflow(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("abc")) if h == BackingHandleZero { t.Fatal("expected valid handle for small payload") } if err != nil { t.Fatalf("IssueBackingHandle abc: %v", err) } // Overflow via typed view _, _ = b.IssueBackingHandle([]byte("x")) // h=3 bytes, retained=3. typed=1 byte, 3+1=4 <= 100. No overflow yet. // Force overflow: issue a large payload b2, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } _, _ = b2.IssueBackingHandle([]byte("abc")) _, _ = b2.IssueBackingHandle([]byte("x")) // Now overflow bigH, err := b2.IssueBackingHandle(make([]byte, 100)) if bigH != BackingHandleZero { t.Fatal("expected zero handle on overflow") } if !errors.Is(err, ErrIngressSnapshotLimitExceeded) { t.Fatalf("expected ErrIngressSnapshotLimitExceeded, got %v", err) } if _, err := b2.Build(); err != ErrIngressSnapshotLimitExceeded { t.Errorf("post-overflow Build: want LimitExceeded, got %v", err) } // Other builder is unaffected b3, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h3, err := b3.IssueBackingHandle([]byte("abc")) if err != nil { t.Fatalf("IssueBackingHandle abc: %v", err) } h3t, err := b3.IssueBackingHandle([]byte("x")) if err != nil { t.Fatalf("IssueBackingHandle x: %v", err) } if h3 == BackingHandleZero || h3t == BackingHandleZero { t.Fatal("expected valid handles for unaffected builder") } if err := b3.SetCanonicalWithHandle(h3); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if _, err := b3.Build(); err != nil { t.Fatalf("unaffected builder Build: %v", err) } } // TestIngressSnapshot_DuplicateCanonicalRejected verifies that setting canonical // twice returns ErrIngressSnapshotDuplicateCanonical. func TestIngressSnapshot_DuplicateCanonicalRejected(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h1, err := b.IssueBackingHandle([]byte("a")) if err != nil { t.Fatalf("IssueBackingHandle h1: %v", err) } if err := b.SetCanonicalWithHandle(h1); err != nil { t.Fatalf("first SetCanonicalWithHandle: %v", err) } h2, err := b.IssueBackingHandle([]byte("b")) if err != nil { t.Fatalf("IssueBackingHandle h2: %v", err) } if err := b.SetCanonicalWithHandle(h2); err != ErrIngressSnapshotDuplicateCanonical { t.Errorf("second SetCanonicalWithHandle: want DuplicateCanonical, got %v", err) } } // TestIngressSnapshot_DuplicateTypedNameRejected verifies that registering two // typed views with the same name returns ErrIngressSnapshotDuplicateTyped. func TestIngressSnapshot_DuplicateTypedNameRejected(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h1, err := b.IssueBackingHandle([]byte("a")) if err != nil { t.Fatalf("IssueBackingHandle h1: %v", err) } h2, err := b.IssueBackingHandle([]byte("b")) if err != nil { t.Fatalf("IssueBackingHandle h2: %v", err) } if err := b.AddTypedView("raw", h1); err != nil { t.Fatalf("first AddTypedView: %v", err) } if err := b.AddTypedView("raw", h2); err != ErrIngressSnapshotDuplicateTyped { t.Errorf("second AddTypedView same name: want DuplicateTyped, got %v", err) } } // TestIngressSnapshot_BuilderInputMutationIsolation verifies that mutating the // caller's input slice after IssueBackingHandle does not affect the stored payload. func TestIngressSnapshot_BuilderInputMutationIsolation(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } data := []byte("mutation-test-payload") h, err := b.IssueBackingHandle(data) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } // Mutate caller's input after store data[0] = 'Z' snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } acc := snap.Accessor() stored, err := acc.Canonical() if err != nil { t.Fatalf("Canonical: %v", err) } if stored[0] == 'Z' { t.Error("stored payload was mutated by caller input change") } if string(stored) != "mutation-test-payload" { t.Errorf("payload mismatch: want %q, got %q", "mutation-test-payload", stored) } } // TestIngressSnapshot_BuildIdempotent verifies that Build can only succeed once. func TestIngressSnapshot_BuildIdempotent(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("x")) if err != nil { t.Fatalf("IssueBackingHandle x: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("first Build: %v", err) } if _, err := b.Build(); err != ErrIngressSnapshotAlreadyBuilt { t.Errorf("second Build: want AlreadyBuilt, got %v", err) } // Snapshot accessor after close acc := snap.Accessor() if _, err := acc.Canonical(); err != nil { t.Fatalf("Canonical before close: %v", err) } snap.Close() if _, err := acc.Canonical(); err != ErrIngressSnapshotClosed { t.Errorf("post-close Canonical: want Closed, got %v", err) } } // TestIngressSnapshot_ConcurrentCloseAndAccess verifies that concurrent access // and close are safe. func TestIngressSnapshot_ConcurrentCloseAndAccess(t *testing.T) { const maxBytes int64 = 1024 * 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle(make([]byte, 1024)) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } var wg sync.WaitGroup for i := 0; i < 10; i++ { wg.Add(1) go func() { defer wg.Done() acc := snap.Accessor() _, _ = acc.Canonical() _ = acc.RetainedCount() _ = acc.RetainedBytes() _ = acc.HasCanonical() }() } wg.Add(1) go func() { defer wg.Done() snap.Close() }() wg.Wait() } // TestIngressSnapshot_TypedViewsAccessible verifies that typed views are // accessible by name and return defensive copies. func TestIngressSnapshot_TypedViewsAccessible(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h1, err := b.IssueBackingHandle([]byte("payload-alpha")) if err != nil { t.Fatalf("IssueBackingHandle h1: %v", err) } h2, err := b.IssueBackingHandle([]byte("payload-beta-gamma")) if err != nil { t.Fatalf("IssueBackingHandle h2: %v", err) } h3, err := b.IssueBackingHandle([]byte("payload-delta")) if err != nil { t.Fatalf("IssueBackingHandle h3: %v", err) } if err := b.SetCanonicalWithHandle(h1); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("first", h2); err != nil { t.Fatalf("AddTypedView first: %v", err) } if err := b.AddTypedView("second", h3); err != nil { t.Fatalf("AddTypedView second: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } acc := snap.Accessor() names := acc.TypedViewNames() if len(names) != 2 || names[0] != "first" || names[1] != "second" { t.Errorf("typed view names: want [first second], got %v", names) } v1, err := acc.TypedView("first") if err != nil { t.Fatalf("TypedView first: %v", err) } if string(v1) != "payload-beta-gamma" { t.Errorf("TypedView first: want %q, got %q", "payload-beta-gamma", v1) } v2, err := acc.TypedView("second") if err != nil { t.Fatalf("TypedView second: %v", err) } if string(v2) != "payload-delta" { t.Errorf("TypedView second: want %q, got %q", "payload-delta", v2) } // Unknown name returns error if _, err := acc.TypedView("unknown"); err != ErrIngressSnapshotInvalidHandle { t.Errorf("TypedView unknown: want InvalidHandle, got %v", err) } // RetainedBackings backings := acc.RetainedBackings() if len(backings) != 3 { t.Errorf("RetainedBackings: want 3, got %d", len(backings)) } } // TestIngressSnapshot_UnboundedBuilderRejects verifies that a builder with // non-positive maxBytes is rejected. func TestIngressSnapshot_UnboundedBuilderRejects(t *testing.T) { const max int64 = 0 _, err := NewIngressSnapshotBuilder(max) if !errors.Is(err, ErrIngressSnapshotZeroMaxBytes) { t.Errorf("zero maxBytes: want ZeroMaxBytes, got %v", err) } } // TestIngressSnapshot_RaceFixture verifies race safety under the race detector. func TestIngressSnapshot_RaceFixture(t *testing.T) { const maxBytes int64 = 1024 * 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } data := make([]byte, 4096) h, err := b.IssueBackingHandle(data) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("view", h); err != nil { t.Fatalf("AddTypedView: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } var wg sync.WaitGroup for i := 0; i < 50; i++ { wg.Add(1) go func() { defer wg.Done() acc := snap.Accessor() _, _ = acc.Canonical() _, _ = acc.TypedView("view") _ = acc.RetainedCount() _ = acc.RetainedBytes() _ = acc.HasCanonical() _ = acc.TypedViewNames() _ = acc.RetainedBackings() _ = acc.IsClosed() }() } wg.Add(1) go func() { defer wg.Done() snap.Close() }() wg.Wait() } // TestIngressSnapshot_NewBuilderMustBePositive verifies that NewIngressSnapshotBuilder // rejects zero and negative maxBytes. func TestIngressSnapshot_NewBuilderMustBePositive(t *testing.T) { _, err := NewIngressSnapshotBuilder(0) if !errors.Is(err, ErrIngressSnapshotZeroMaxBytes) { t.Errorf("zero maxBytes: want ZeroMaxBytes, got %v", err) } _, err = NewIngressSnapshotBuilder(-1) if !errors.Is(err, ErrIngressSnapshotZeroMaxBytes) { t.Errorf("negative maxBytes: want ZeroMaxBytes, got %v", err) } } // TestIngressSnapshot_ClosedSnapshotAccessor verifies that a closed snapshot // returns ErrIngressSnapshotClosed from all accessor methods. func TestIngressSnapshot_ClosedSnapshotAccessor(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("data")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } acc := snap.Accessor() // Pre-close: all methods work if !acc.HasCanonical() { t.Error("pre-close HasCanonical should be true") } if _, err := acc.Canonical(); err != nil { t.Fatalf("pre-close Canonical: %v", err) } // Close snap.Close() // Post-close: data access methods return Closed if !acc.IsClosed() { t.Error("post-close IsClosed should be true") } // HasCanonical is a metadata query; it reflects whether canonical was set, // not the close state. It remains true because canonical was set before close. if !acc.HasCanonical() { t.Error("post-close HasCanonical should remain true (metadata query)") } if _, err := acc.Canonical(); err != ErrIngressSnapshotClosed { t.Errorf("post-close Canonical: want Closed, got %v", err) } if _, err := acc.TypedView("x"); err != ErrIngressSnapshotClosed { t.Errorf("post-close TypedView: want Closed, got %v", err) } if acc.RetainedCount() != 0 { t.Errorf("post-close RetainedCount: want 0, got %d", acc.RetainedCount()) } if acc.RetainedBytes() != 0 { t.Errorf("post-close RetainedBytes: want 0, got %d", acc.RetainedBytes()) } if len(acc.TypedViewNames()) != 0 { t.Errorf("post-close TypedViewNames: want [], got %v", acc.TypedViewNames()) } if len(acc.RetainedBackings()) != 0 { t.Errorf("post-close RetainedBackings: want [], got %v", acc.RetainedBackings()) } // Idempotent close snap.Close() if !acc.IsClosed() { t.Error("post-idempotent-close IsClosed should be true") } } // TestIngressSnapshot_OverflowAfterBuildPreventsSnapshot verifies that overflow // during IssueBackingHandle prevents Build from producing a snapshot. func TestIngressSnapshot_OverflowAfterBuildPreventsSnapshot(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("abc")) if h == BackingHandleZero { t.Fatal("expected valid handle") } if err != nil { t.Fatalf("IssueBackingHandle abc: %v", err) } // Force overflow big, err := b.IssueBackingHandle(make([]byte, 100)) if big != BackingHandleZero { t.Fatal("expected zero handle on overflow") } if !errors.Is(err, ErrIngressSnapshotLimitExceeded) { t.Fatalf("expected ErrIngressSnapshotLimitExceeded, got %v", err) } // Build should fail if _, err := b.Build(); err != ErrIngressSnapshotLimitExceeded { t.Errorf("post-overflow Build: want LimitExceeded, got %v", err) } // No snapshot was created; other builders are unaffected b2, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h2, _ := b2.IssueBackingHandle([]byte("abc")) _ = b2.SetCanonicalWithHandle(h2) if _, err := b2.Build(); err != nil { t.Errorf("unaffected builder Build: %v", err) } } // TestIngressSnapshot_InvalidHandleAfterBuild verifies that handles issued // before Build are invalid after the builder is built. func TestIngressSnapshot_InvalidHandleAfterBuild(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("data")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } // Snapshot still has the backing acc := snap.Accessor() if !acc.HasCanonical() { t.Error("snapshot should have canonical after Build") } canon, err := acc.Canonical() if err != nil { t.Fatalf("Canonical: %v", err) } if string(canon) != "data" { t.Errorf("canonical: want %q, got %q", "data", canon) } // Builder cannot accept new handles newH, newErr := b.IssueBackingHandle([]byte("new")) if newH != BackingHandleZero { t.Error("builder should return zero handle after Build") } if !errors.Is(newErr, ErrIngressSnapshotAlreadyBuilt) { t.Errorf("expected ErrIngressSnapshotAlreadyBuilt from IssueBackingHandle after Build, got %v", newErr) } } // TestIngressSnapshot_CloseBuilderReleasesMemory verifies that Close on a // non-built, non-failed builder releases backing memory. func TestIngressSnapshot_CloseBuilderReleasesMemory(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("data")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } // Close before Build b.Close() // Builder is closed; subsequent operations return Closed if err := b.SetCanonicalWithHandle(h); err != ErrIngressSnapshotClosed { t.Errorf("post-close SetCanonical: want Closed, got %v", err) } if _, err := b.Build(); err != ErrIngressSnapshotClosed { t.Errorf("post-close Build: want Closed, got %v", err) } // Builder zero-state after Close if b.backings != nil { t.Error("backings should be nil after Close") } if b.typedViews != nil { t.Error("typedViews should be nil after Close") } if b.typedNames != nil { t.Error("typedNames should be nil after Close") } if b.canonicalHdl.isValid() { t.Error("canonicalHdl should be zero after Close") } if b.retainedBytes != 0 { t.Errorf("retainedBytes should be 0 after Close, got %d", b.retainedBytes) } } // TestIngressSnapshot_BuildRequiresCanonical verifies that Build without a // canonical handle returns ErrIngressSnapshotNilCanonical, and that the // builder remains mutable to allow setting canonical before a final Build. func TestIngressSnapshot_BuildRequiresCanonical(t *testing.T) { const maxBytes int64 = 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } // Build without canonical: must return ErrIngressSnapshotNilCanonical _, err = b.Build() if !errors.Is(err, ErrIngressSnapshotNilCanonical) { t.Fatalf("Build without canonical: want ErrIngressSnapshotNilCanonical, got %v", err) } // Builder must still be mutable after failed Build h, err := b.IssueBackingHandle([]byte("canonical-payload")) if err != nil { t.Fatalf("IssueBackingHandle after failed Build: %v", err) } if h == BackingHandleZero { t.Fatal("expected valid handle after failed Build") } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle after failed Build: %v", err) } // Now Build should succeed snap, err := b.Build() if err != nil { t.Fatalf("Build after setting canonical: %v", err) } acc := snap.Accessor() if !acc.HasCanonical() { t.Error("HasCanonical should be true after successful Build") } canon, err := acc.Canonical() if err != nil { t.Fatalf("Canonical: %v", err) } if string(canon) != "canonical-payload" { t.Errorf("canonical payload: want %q, got %q", "canonical-payload", canon) } } // TestIngressSnapshot_IssueBackingStateErrors verifies that IssueBackingHandle // returns explicit errors for overflow, closed, and built states. func TestIngressSnapshot_IssueBackingStateErrors(t *testing.T) { const maxBytes int64 = 100 // Test overflow returns ErrIngressSnapshotLimitExceeded b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } // Issue a small payload first _, err = b.IssueBackingHandle([]byte("small")) if err != nil { t.Fatalf("IssueBackingHandle small: %v", err) } // Overflow: should return (zero, ErrIngressSnapshotLimitExceeded) h, err := b.IssueBackingHandle(make([]byte, 100)) if h != BackingHandleZero { t.Fatal("expected zero handle on overflow") } if !errors.Is(err, ErrIngressSnapshotLimitExceeded) { t.Fatalf("post-overflow IssueBackingHandle: want ErrIngressSnapshotLimitExceeded, got %v", err) } // Re-issue after overflow: should still return ErrIngressSnapshotLimitExceeded h2, err2 := b.IssueBackingHandle([]byte("retry")) if h2 != BackingHandleZero { t.Fatal("expected zero handle on retry after overflow") } if !errors.Is(err2, ErrIngressSnapshotLimitExceeded) { t.Fatalf("retry after overflow: want ErrIngressSnapshotLimitExceeded, got %v", err2) } // Test Close returns ErrIngressSnapshotClosed b2, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } _, _ = b2.IssueBackingHandle([]byte("data")) b2.Close() h3, err3 := b2.IssueBackingHandle([]byte("after-close")) if h3 != BackingHandleZero { t.Fatal("expected zero handle after Close") } if !errors.Is(err3, ErrIngressSnapshotClosed) { t.Fatalf("post-close IssueBackingHandle: want ErrIngressSnapshotClosed, got %v", err3) } // Test Build returns ErrIngressSnapshotAlreadyBuilt b3, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h4, err := b3.IssueBackingHandle([]byte("data")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b3.SetCanonicalWithHandle(h4); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if _, err := b3.Build(); err != nil { t.Fatalf("Build: %v", err) } h5, err5 := b3.IssueBackingHandle([]byte("after-build")) if h5 != BackingHandleZero { t.Fatal("expected zero handle after Build") } if !errors.Is(err5, ErrIngressSnapshotAlreadyBuilt) { t.Fatalf("post-Build IssueBackingHandle: want ErrIngressSnapshotAlreadyBuilt, got %v", err5) } } // TestIngressSnapshot_RebuildPrePeakOverflow verifies that ReserveRebuild fails closed // when retained bytes + expected temporary bytes exceed maxBytes limit. func TestIngressSnapshot_RebuildPrePeakOverflow(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle(make([]byte, 80)) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } defer snap.Close() if snap.MaxBytes() != 100 { t.Errorf("MaxBytes: want 100, got %d", snap.MaxBytes()) } if snap.ReservedTempBytes() != 0 { t.Errorf("ReservedTempBytes before rebuild: want 0, got %d", snap.ReservedTempBytes()) } // 80 retained + 30 expected = 110 > 100 maxBytes => Pre-peak overflow guard, err := snap.ReserveRebuild(30) if guard != nil { t.Fatal("expected nil guard on pre-peak overflow") } if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("ReserveRebuild pre-peak overflow: want ErrIngressSnapshotRebuildOverflow, got %v", err) } if !errors.Is(err, ErrIngressSnapshotLimitExceeded) { t.Fatalf("ReserveRebuild pre-peak overflow should wrap LimitExceeded, got %v", err) } // Reserved temp bytes remains 0 if snap.ReservedTempBytes() != 0 { t.Errorf("ReservedTempBytes after failed reserve: want 0, got %d", snap.ReservedTempBytes()) } } // TestIngressSnapshot_RebuildPostPeakOverflow verifies that when actual rebuild output // exceeds maxBytes during Commit, post-peak overflow triggers, temporary/canonical/typed // references are released in defined order, and no dispatchable object is returned. func TestIngressSnapshot_RebuildPostPeakOverflow(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("initial-retained-data-payload-50-bytes-1234567890")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("view1", h); err != nil { t.Fatalf("AddTypedView: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } defer snap.Close() // Reserve 30 bytes: 50 + 30 = 80 <= 100 => Success guard, err := snap.ReserveRebuild(30) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } if guard == nil { t.Fatal("expected valid guard on successful reserve") } if guard.ReservedBytes() != 30 { t.Errorf("guard ReservedBytes: want 30, got %d", guard.ReservedBytes()) } if snap.ReservedTempBytes() != 30 { t.Errorf("snap ReservedTempBytes: want 30, got %d", snap.ReservedTempBytes()) } // Verify references present before commit failure if !guard.HasCanonicalReference() { t.Error("expected canonical reference before commit") } if !guard.HasTypedReferences() { t.Error("expected typed references before commit") } // Commit actual output of 60 bytes: 50 retained + 60 actual = 110 > 100 => Post-peak overflow rebuiltSnap, err := guard.Commit(make([]byte, 60)) if rebuiltSnap != nil { t.Fatal("expected nil rebuilt snapshot on post-peak overflow") } if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("Commit post-peak overflow: want ErrIngressSnapshotRebuildOverflow, got %v", err) } if !errors.Is(err, ErrIngressSnapshotLimitExceeded) { t.Fatalf("Commit post-peak overflow should wrap LimitExceeded, got %v", err) } // Guard is now in overflow state and released if guard.State() != RebuildStateOverflow { t.Errorf("guard State: want RebuildStateOverflow, got %v", guard.State()) } if !guard.IsReleased() { t.Error("guard IsReleased should be true after overflow") } // Ownership release sequence verified: temporary -> canonical -> typed references dropped if guard.HasTempReference() { t.Error("temporary reference should be dropped after overflow") } if guard.HasCanonicalReference() { t.Error("canonical reference should be dropped after overflow") } if guard.HasTypedReferences() { t.Error("typed references should be dropped after overflow") } // Parent snapshot reserved temp bytes unreserved if snap.ReservedTempBytes() != 0 { t.Errorf("snap ReservedTempBytes after overflow: want 0, got %d", snap.ReservedTempBytes()) } // Dispatchable snapshot returns overflow error if _, err := guard.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Errorf("DispatchableSnapshot after overflow: want ErrIngressSnapshotRebuildOverflow, got %v", err) } } // TestIngressSnapshot_RebuildTwoPhaseCommitSuccess verifies a successful two-phase // reserve and commit, producing a valid dispatchable IngressSnapshot. func TestIngressSnapshot_RebuildTwoPhaseCommitSuccess(t *testing.T) { const maxBytes int64 = 200 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } canonicalData := []byte("canonical-payload") h, err := b.IssueBackingHandle(canonicalData) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("raw", h); err != nil { t.Fatalf("AddTypedView: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } defer snap.Close() // Reserve 50 bytes guard, err := snap.ReserveRebuild(50) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } // Commit actual output with typed view name "rebuilt" rebuiltOutput := []byte("rebuilt-transformed-payload") rebuiltSnap, err := guard.CommitTyped("rebuilt", rebuiltOutput) if err != nil { t.Fatalf("CommitTyped: %v", err) } if rebuiltSnap == nil { t.Fatal("expected non-nil rebuilt snapshot") } defer rebuiltSnap.Close() // Guard state updated to committed if guard.State() != RebuildStateCommitted { t.Errorf("guard State: want RebuildStateCommitted, got %v", guard.State()) } if snap.ReservedTempBytes() != 0 { t.Errorf("snap ReservedTempBytes after commit: want 0, got %d", snap.ReservedTempBytes()) } // Rebuilt snapshot accessor checks acc := rebuiltSnap.Accessor() if acc.MaxBytes() != maxBytes { t.Errorf("rebuilt MaxBytes: want %d, got %d", maxBytes, acc.MaxBytes()) } if !acc.HasCanonical() { t.Error("rebuilt snapshot should have canonical handle") } canon, err := acc.Canonical() if err != nil { t.Fatalf("Canonical: %v", err) } if string(canon) != string(canonicalData) { t.Errorf("Canonical: want %q, got %q", canonicalData, canon) } origView, err := acc.TypedView("raw") if err != nil { t.Fatalf("TypedView raw: %v", err) } if string(origView) != string(canonicalData) { t.Errorf("TypedView raw: want %q, got %q", canonicalData, origView) } rebuiltView, err := acc.TypedView("rebuilt") if err != nil { t.Fatalf("TypedView rebuilt: %v", err) } if string(rebuiltView) != string(rebuiltOutput) { t.Errorf("TypedView rebuilt: want %q, got %q", rebuiltOutput, rebuiltView) } names := acc.TypedViewNames() if len(names) != 2 || names[0] != "raw" || names[1] != "rebuilt" { t.Errorf("TypedViewNames: want [raw rebuilt], got %v", names) } // Dispatchable snapshot accessor matches dispSnap, err := guard.DispatchableSnapshot() if err != nil { t.Fatalf("DispatchableSnapshot: %v", err) } if dispSnap != rebuiltSnap { t.Error("DispatchableSnapshot should return committed rebuilt snapshot") } } // TestIngressSnapshot_RebuildCancelAndReleaseSequence verifies that Cancel releases // temporary, canonical, and typed references in defined order and unreserves bytes. func TestIngressSnapshot_RebuildCancelAndReleaseSequence(t *testing.T) { const maxBytes int64 = 500 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("data")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("view", h); err != nil { t.Fatalf("AddTypedView: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } defer snap.Close() guard, err := snap.ReserveRebuild(100) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } if snap.ReservedTempBytes() != 100 { t.Errorf("ReservedTempBytes: want 100, got %d", snap.ReservedTempBytes()) } if !guard.HasCanonicalReference() || !guard.HasTypedReferences() { t.Fatal("expected active canonical and typed references before cancel") } // Cancel if err := guard.Cancel(); err != nil { t.Fatalf("Cancel: %v", err) } // State and releases if guard.State() != RebuildStateCancelled { t.Errorf("State: want RebuildStateCancelled, got %v", guard.State()) } if !guard.IsReleased() { t.Error("IsReleased should be true after Cancel") } // References dropped in order if guard.HasTempReference() { t.Error("temp reference should be dropped after cancel") } if guard.HasCanonicalReference() { t.Error("canonical reference should be dropped after cancel") } if guard.HasTypedReferences() { t.Error("typed references should be dropped after cancel") } // Reserved temp bytes on parent unreserved if snap.ReservedTempBytes() != 0 { t.Errorf("ReservedTempBytes after cancel: want 0, got %d", snap.ReservedTempBytes()) } // Idempotent cancel & close if err := guard.Cancel(); err != nil { t.Errorf("second Cancel: %v", err) } guard.Close() } // TestIngressSnapshot_RebuildConcurrentCloseAndCancelRace verifies concurrent safety under race detector. func TestIngressSnapshot_RebuildConcurrentCloseAndCancelRace(t *testing.T) { const maxBytes int64 = 1024 * 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle(make([]byte, 1024)) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } guard, err := snap.ReserveRebuild(2048) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } var wg sync.WaitGroup for i := 0; i < 20; i++ { wg.Add(1) go func() { defer wg.Done() _ = guard.State() _ = guard.IsReleased() _ = guard.ReservedBytes() _ = guard.HasTempReference() _ = guard.HasCanonicalReference() _ = guard.HasTypedReferences() _, _ = guard.DispatchableSnapshot() _ = guard.Accessor() }() } wg.Add(1) go func() { defer wg.Done() _, _ = guard.Commit([]byte("concurrent-output-data")) }() wg.Add(1) go func() { defer wg.Done() _ = guard.Cancel() }() wg.Add(1) go func() { defer wg.Done() guard.Close() }() wg.Add(1) go func() { defer wg.Done() snap.Close() }() wg.Wait() } // TestIngressSnapshot_RebuildUseAfterReleaseReturnsStableError verifies that accessing a released // guard returns stable errors. func TestIngressSnapshot_RebuildUseAfterReleaseReturnsStableError(t *testing.T) { const maxBytes int64 = 1000 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("hello")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } defer snap.Close() guard, err := snap.ReserveRebuild(100) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } guard.Close() // Method access after release returns stable errors if _, err := guard.Commit([]byte("data")); !errors.Is(err, ErrIngressSnapshotRebuildInvalidState) && !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("post-release Commit: want InvalidState or Closed, got %v", err) } if _, err := guard.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotClosed) && !errors.Is(err, ErrIngressSnapshotRebuildInvalidState) { t.Errorf("post-release DispatchableSnapshot: want Closed or InvalidState, got %v", err) } acc := guard.Accessor() if !acc.IsClosed() { t.Error("post-release guard Accessor should return closed accessor") } } // TestIngressSnapshot_RebuildPeakArithmeticOverflow verifies that pre-peak and post-peak // byte arithmetic never wraps through int64 MaxInt64, even with math.MaxInt64 inputs. // On overflow the guard returns ErrIngressSnapshotRebuildOverflow, no dispatch object is // produced, and no reservation remains on the parent snapshot. func TestIngressSnapshot_RebuildPeakArithmeticOverflow(t *testing.T) { tests := []struct { name string maxBytes int64 retainedSize int64 reserveExpected int64 commitActualSize int64 wantPreOverflow bool wantPostOverflow bool }{ { name: "math.MaxInt64_reserve_pre_peak", maxBytes: 100, retainedSize: 1, reserveExpected: math.MaxInt64, commitActualSize: 0, wantPreOverflow: true, }, { name: "exact_limit_reserve_pre_peak", maxBytes: 100, retainedSize: 50, reserveExpected: 50, commitActualSize: 50, wantPreOverflow: false, wantPostOverflow: false, }, { name: "exact_limit_plus_one_reserve_pre_peak", maxBytes: 100, retainedSize: 50, reserveExpected: 51, commitActualSize: 0, wantPreOverflow: true, }, { name: "post_peak_overflow", maxBytes: 100, retainedSize: 50, reserveExpected: 30, commitActualSize: 60, wantPreOverflow: false, wantPostOverflow: true, }, } for _, tc := range tests { t.Run(tc.name, func(t *testing.T) { b, err := NewIngressSnapshotBuilder(tc.maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } retainedPayload := make([]byte, tc.retainedSize) h, err := b.IssueBackingHandle(retainedPayload) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if h == BackingHandleZero { t.Fatal("expected valid handle from IssueBackingHandle") } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } defer snap.Close() // Pre-peak: ReserveRebuild guard, err := snap.ReserveRebuild(tc.reserveExpected) if tc.wantPreOverflow { if guard != nil { t.Fatal("expected nil guard on pre-peak overflow") } if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("ReserveRebuild pre-peak: want ErrIngressSnapshotRebuildOverflow, got %v", err) } if snap.ReservedTempBytes() != 0 { t.Errorf("ReservedTempBytes after failed reserve: want 0, got %d", snap.ReservedTempBytes()) } return } if err != nil { t.Fatalf("ReserveRebuild: %v", err) } if guard == nil { t.Fatal("expected non-nil guard") } // Post-peak: Commit rebuilt, err := guard.Commit(make([]byte, tc.commitActualSize)) if tc.wantPostOverflow { if rebuilt != nil { t.Fatal("expected nil rebuilt snapshot on post-peak overflow") } if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("Commit post-peak: want ErrIngressSnapshotRebuildOverflow, got %v", err) } if guard.State() != RebuildStateOverflow { t.Errorf("guard State: want Overflow, got %v", guard.State()) } if !guard.IsReleased() { t.Error("guard should be released after overflow") } if guard.HasTempReference() || guard.HasCanonicalReference() || guard.HasTypedReferences() { t.Error("guard references should be dropped after overflow") } if snap.ReservedTempBytes() != 0 { t.Errorf("snap ReservedTempBytes after overflow: want 0, got %d", snap.ReservedTempBytes()) } if !snap.Accessor().IsClosed() { t.Error("snap should be closed after post-peak overflow") } if snap.Accessor().RetainedCount() != 0 { t.Errorf("snap RetainedCount after overflow: want 0, got %d", snap.Accessor().RetainedCount()) } if snap.Accessor().RetainedBytes() != 0 { t.Errorf("snap RetainedBytes after overflow: want 0, got %d", snap.Accessor().RetainedBytes()) } return } if err != nil { t.Fatalf("Commit: %v", err) } if rebuilt == nil { t.Fatal("expected non-nil rebuilt snapshot on success") } defer rebuilt.Close() if snap.ReservedTempBytes() != 0 { t.Errorf("snap ReservedTempBytes after commit: want 0, got %d", snap.ReservedTempBytes()) } }) } } // TestIngressSnapshot_RebuildOverflowReleasesOwnedSnapshot verifies that post-peak overflow // fully releases the parent snapshot's backing store and marks it closed, so no references // remain on the parent. The overflow guard's state is Overflow and all its references are nil. func TestIngressSnapshot_RebuildOverflowReleasesOwnedSnapshot(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("retained-payload-data-42-bytes-long-here")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("raw", h); err != nil { t.Fatalf("AddTypedView: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } // Verify initial state acc := snap.Accessor() if acc.RetainedCount() != 1 { t.Errorf("initial RetainedCount: want 1, got %d", acc.RetainedCount()) } if acc.RetainedBytes() != 40 { t.Errorf("initial RetainedBytes: want 40, got %d", acc.RetainedBytes()) } if !acc.HasCanonical() { t.Error("initial HasCanonical should be true") } guard, err := snap.ReserveRebuild(30) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } if snap.ReservedTempBytes() != 30 { t.Errorf("ReservedTempBytes after reserve: want 30, got %d", snap.ReservedTempBytes()) } // Commit oversized output to trigger post-peak overflow _, err = guard.Commit(make([]byte, 80)) if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("Commit overflow: want ErrIngressSnapshotRebuildOverflow, got %v", err) } // Guard state and references if guard.State() != RebuildStateOverflow { t.Errorf("guard State: want Overflow, got %v", guard.State()) } if !guard.IsReleased() { t.Error("guard should be released") } if guard.HasTempReference() { t.Error("guard temp reference should be nil") } if guard.HasCanonicalReference() { t.Error("guard canonical reference should be nil") } if guard.HasTypedReferences() { t.Error("guard typed references should be nil") } if guard.RebuiltSnapshot() != nil { t.Error("guard rebuiltSnapshot should be nil after overflow") } // Parent snapshot fully released if snap.ReservedTempBytes() != 0 { t.Errorf("snap ReservedTempBytes: want 0, got %d", snap.ReservedTempBytes()) } if !snap.Accessor().IsClosed() { t.Error("snap should be closed after overflow") } if snap.Accessor().RetainedCount() != 0 { t.Errorf("snap RetainedCount: want 0, got %d", snap.Accessor().RetainedCount()) } if snap.Accessor().RetainedBytes() != 0 { t.Errorf("snap RetainedBytes: want 0, got %d", snap.Accessor().RetainedBytes()) } if _, err := snap.Accessor().Canonical(); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("snap Canonical after overflow: want Closed, got %v", err) } // Idempotent close snap.Close() if !snap.Accessor().IsClosed() { t.Error("snap should remain closed after idempotent close") } } // TestIngressSnapshot_RebuildOverflowInvalidatesSiblingGuard verifies that when one guard's // commit triggers post-peak overflow (closing the parent snapshot), a sibling guard that // has not yet committed returns ErrIngressSnapshotClosed on subsequent operations. func TestIngressSnapshot_RebuildOverflowInvalidatesSiblingGuard(t *testing.T) { const maxBytes int64 = 200 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("retained-payload-data-50-bytes-long-here")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } // Reserve two guards guardA, err := snap.ReserveRebuild(50) if err != nil { t.Fatalf("ReserveRebuild A: %v", err) } guardB, err := snap.ReserveRebuild(50) if err != nil { t.Fatalf("ReserveRebuild B: %v", err) } if snap.ReservedTempBytes() != 100 { t.Errorf("ReservedTempBytes: want 100, got %d", snap.ReservedTempBytes()) } // Guard A commits with oversized output -> triggers overflow -> closes parent _, err = guardA.Commit(make([]byte, 160)) if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("guardA Commit: want overflow, got %v", err) } // Parent is now closed if !snap.Accessor().IsClosed() { t.Error("parent snapshot should be closed after overflow") } // Guard B (sibling, still Reserved) should get closed errors on Commit // because parent snapshot was closed by guardA's overflow. if _, err := guardB.Commit([]byte("data")); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("guardB Commit after sibling overflow: want ErrIngressSnapshotClosed, got %v", err) } if _, err := guardB.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("guardB DispatchableSnapshot after sibling overflow: want Closed, got %v", err) } // guardB.Commit already called releaseLocked(RebuildStateReleased) internally, // so guardB is now in Released state. Cancel is a no-op. if err := guardB.Cancel(); err != nil { t.Fatalf("guardB Cancel after sibling overflow: %v", err) } if guardB.State() != RebuildStateReleased { t.Errorf("guardB State: want Released (set by closed-snap Commit), got %v", guardB.State()) } // Guard A overflow state is stable if guardA.State() != RebuildStateOverflow { t.Errorf("guardA State: want Overflow, got %v", guardA.State()) } if _, err := guardA.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Errorf("guardA DispatchableSnapshot: want Overflow error, got %v", err) } } // TestIngressSnapshot_RebuildCommittedCloseDropsGuardReferences verifies that calling Close // on a committed guard drops all guard-internal references (rebuiltSnapshot, temp, canonical, // typed) without closing the parent snapshot or the rebuilt snapshot. func TestIngressSnapshot_RebuildCommittedCloseDropsGuardReferences(t *testing.T) { const maxBytes int64 = 500 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("canonical-payload-data-20-bytes-here")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("raw", h); err != nil { t.Fatalf("AddTypedView: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } guard, err := snap.ReserveRebuild(100) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } rebuilt, err := guard.CommitTyped("rebuilt", []byte("rebuilt-output-payload-data-25-bytes-long")) if err != nil { t.Fatalf("CommitTyped: %v", err) } if guard.State() != RebuildStateCommitted { t.Errorf("guard State: want Committed, got %v", guard.State()) } // Parent snapshot should still be open and intact if snap.Accessor().IsClosed() { t.Error("parent snapshot should NOT be closed after successful commit") } if snap.Accessor().RetainedCount() != 1 { t.Errorf("parent RetainedCount: want 1, got %d", snap.Accessor().RetainedCount()) } // Rebuilt snapshot should be accessible if rebuilt.Accessor().RetainedCount() != 2 { t.Errorf("rebuilt RetainedCount: want 2, got %d", rebuilt.Accessor().RetainedCount()) } // Close the committed guard guard.Close() // Guard references should be dropped if guard.HasTempReference() { t.Error("guard temp reference should be dropped after Close") } if guard.HasCanonicalReference() { t.Error("guard canonical reference should be dropped after Close") } if guard.HasTypedReferences() { t.Error("guard typed references should be dropped after Close") } if guard.RebuiltSnapshot() != nil { t.Error("guard rebuiltSnapshot should be nil after Close") } // DispatchableSnapshot returns Closed since guard is released if _, err := guard.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("post-close DispatchableSnapshot: want Closed, got %v", err) } // Parent and rebuilt should still be usable if snap.Accessor().IsClosed() { t.Error("parent snapshot should still be open") } if rebuilt.Accessor().IsClosed() { t.Error("rebuilt snapshot should still be open") } // Idempotent Close guard.Close() if guard.RebuiltSnapshot() != nil { t.Error("guard rebuiltSnapshot should still be nil after second Close") } snap.Close() rebuilt.Close() } // TestIngressSnapshot_RebuildOverflowDropsCommittedGuardReferences verifies that when a guard // enters the Overflow state, its rebuiltSnapshot reference is also nil (no leaked committed // snapshot). This covers the case where releaseLocked is called with RebuildStateOverflow. func TestIngressSnapshot_RebuildOverflowDropsCommittedGuardReferences(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("data")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } guard, err := snap.ReserveRebuild(10) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } // Overflow via Commit _, err = guard.Commit(make([]byte, 200)) if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("Commit: want overflow, got %v", err) } // rebuiltSnapshot must be nil even in overflow state if guard.RebuiltSnapshot() != nil { t.Error("rebuiltSnapshot should be nil after overflow release") } // DispatchableSnapshot returns Overflow error if _, err := guard.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Errorf("DispatchableSnapshot: want Overflow, got %v", err) } snap.Close() } // TestIngressSnapshot_RebuildOverflowParentClosedPreventsNewReserve verifies that once the // parent snapshot is closed by an overflow, no new ReserveRebuild calls can succeed. func TestIngressSnapshot_RebuildOverflowParentClosedPreventsNewReserve(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("data")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } guard, err := snap.ReserveRebuild(10) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } // Trigger overflow _, err = guard.Commit(make([]byte, 200)) if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("Commit: want overflow, got %v", err) } // New reserve should fail with Closed _, err = snap.ReserveRebuild(10) if !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("ReserveRebuild after overflow: want Closed, got %v", err) } snap.Close() } // TestIngressSnapshot_RebuildOverflowInvalidatesCommittedAndReservedSiblings verifies that // when one guard's commit triggers post-peak overflow, a previously committed sibling's // rebuilt snapshot is closed and no longer dispatchable, a reserved sibling is released, // and the overflow guard's error remains stable on repeat Commit. func TestIngressSnapshot_RebuildOverflowInvalidatesCommittedAndReservedSiblings(t *testing.T) { const maxBytes int64 = 200 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("retained-payload-for-overflow-test-xx")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } // Reserve three guards guardA, err := snap.ReserveRebuild(50) if err != nil { t.Fatalf("ReserveRebuild A: %v", err) } guardB, err := snap.ReserveRebuild(50) if err != nil { t.Fatalf("ReserveRebuild B: %v", err) } guardC, err := snap.ReserveRebuild(50) if err != nil { t.Fatalf("ReserveRebuild C: %v", err) } if snap.ReservedTempBytes() != 150 { t.Errorf("ReservedTempBytes: want 150, got %d", snap.ReservedTempBytes()) } // Guard A commits successfully rebuiltA, err := guardA.Commit([]byte("rebuilt-output-a-payload-data-xx")) if err != nil { t.Fatalf("guardA Commit: %v", err) } if guardA.State() != RebuildStateCommitted { t.Errorf("guardA State: want Committed, got %v", guardA.State()) } if rebuiltA == nil { t.Fatal("expected non-nil rebuilt snapshot from guardA") } if snap.ReservedTempBytes() != 100 { t.Errorf("ReservedTempBytes after A commit: want 100, got %d", snap.ReservedTempBytes()) } // Guard B commits with oversized output -> triggers overflow _, err = guardB.Commit(make([]byte, 200)) if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("guardB Commit: want overflow, got %v", err) } // Guard A: rebuilt snapshot closed, dispatch returns Closed if !rebuiltA.Accessor().IsClosed() { t.Error("guardA rebuilt snapshot should be closed after sibling overflow") } if _, err := guardA.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("guardA DispatchableSnapshot: want Closed, got %v", err) } if guardA.RebuiltSnapshot() != nil { t.Error("guardA rebuiltSnapshot should be nil after fan-out release") } // Guard B: overflow state, stable error if guardB.State() != RebuildStateOverflow { t.Errorf("guardB State: want Overflow, got %v", guardB.State()) } if !guardB.IsReleased() { t.Error("guardB should be released after overflow") } if guardB.RebuiltSnapshot() != nil { t.Error("guardB rebuiltSnapshot should be nil") } // Guard B repeat Commit returns stable Overflow error _, err = guardB.Commit([]byte("retry")) if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Errorf("guardB repeat Commit: want Overflow, got %v", err) } if _, err := guardB.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Errorf("guardB DispatchableSnapshot: want Overflow, got %v", err) } // Guard C: reserved sibling released if guardC.State() != RebuildStateReleased { t.Errorf("guardC State: want Released, got %v", guardC.State()) } if !guardC.IsReleased() { t.Error("guardC should be released after fan-out") } if _, err := guardC.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("guardC DispatchableSnapshot: want Closed, got %v", err) } // Parent snapshot closed if !snap.Accessor().IsClosed() { t.Error("parent snapshot should be closed after overflow") } if snap.Accessor().RetainedCount() != 0 { t.Errorf("parent RetainedCount: want 0, got %d", snap.Accessor().RetainedCount()) } if snap.Accessor().RetainedBytes() != 0 { t.Errorf("parent RetainedBytes: want 0, got %d", snap.Accessor().RetainedBytes()) } if snap.ReservedTempBytes() != 0 { t.Errorf("ReservedTempBytes: want 0, got %d", snap.ReservedTempBytes()) } // Idempotent close snap.Close() if !snap.Accessor().IsClosed() { t.Error("parent should remain closed after idempotent close") } } // TestIngressSnapshot_RebuildOverflowErrorStableOnRepeat verifies that calling Commit // repeatedly on an overflow guard consistently returns ErrIngressSnapshotRebuildOverflow, // and DispatchableSnapshot also returns the stable overflow error. func TestIngressSnapshot_RebuildOverflowErrorStableOnRepeat(t *testing.T) { const maxBytes int64 = 100 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("data-payload-for-stable-test")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } guard, err := snap.ReserveRebuild(30) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } // First overflow commit _, err = guard.Commit(make([]byte, 200)) if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Fatalf("first Commit: want overflow, got %v", err) } if guard.State() != RebuildStateOverflow { t.Errorf("State after first overflow: want Overflow, got %v", guard.State()) } // Repeat Commit should return stable Overflow error for i := 0; i < 5; i++ { _, err = guard.Commit([]byte("retry-data")) if !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Errorf("repeat Commit #%d: want Overflow, got %v", i+1, err) } if guard.State() != RebuildStateOverflow { t.Errorf("State after repeat #%d: want Overflow, got %v", i+1, guard.State()) } } // DispatchableSnapshot returns stable Overflow error for i := 0; i < 5; i++ { if _, err := guard.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotRebuildOverflow) { t.Errorf("DispatchableSnapshot #%d: want Overflow, got %v", i+1, err) } } // Accessor returns closed accessor acc := guard.Accessor() if !acc.IsClosed() { t.Error("guard Accessor should be closed for overflow guard") } snap.Close() } // TestIngressSnapshot_CloseInvalidatesCommittedGuard verifies that calling Close on the // parent snapshot fans out to all active guards, closes committed siblings' rebuilt // snapshots, and releases reserved siblings. func TestIngressSnapshot_CloseInvalidatesCommittedGuard(t *testing.T) { const maxBytes int64 = 500 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle([]byte("canonical-payload-for-close-inv-test")) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := b.AddTypedView("raw", h); err != nil { t.Fatalf("AddTypedView: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } // Reserve and commit guard A guardA, err := snap.ReserveRebuild(100) if err != nil { t.Fatalf("ReserveRebuild A: %v", err) } rebuiltA, err := guardA.CommitTyped("rebuilt", []byte("rebuilt-output-payload-data-xx")) if err != nil { t.Fatalf("guardA CommitTyped: %v", err) } if guardA.State() != RebuildStateCommitted { t.Errorf("guardA State: want Committed, got %v", guardA.State()) } if rebuiltA == nil { t.Fatal("expected non-nil rebuilt snapshot") } // Reserve guard B (still reserved) guardB, err := snap.ReserveRebuild(100) if err != nil { t.Fatalf("ReserveRebuild B: %v", err) } if guardB.State() != RebuildStateReserved { t.Errorf("guardB State: want Reserved, got %v", guardB.State()) } // Close parent snapshot snap.Close() // Parent is closed if !snap.Accessor().IsClosed() { t.Error("parent snapshot should be closed") } if snap.Accessor().RetainedCount() != 0 { t.Errorf("parent RetainedCount: want 0, got %d", snap.Accessor().RetainedCount()) } if snap.Accessor().RetainedBytes() != 0 { t.Errorf("parent RetainedBytes: want 0, got %d", snap.Accessor().RetainedBytes()) } if snap.ReservedTempBytes() != 0 { t.Errorf("ReservedTempBytes: want 0, got %d", snap.ReservedTempBytes()) } // Guard A: rebuilt snapshot closed, dispatch returns Closed if !rebuiltA.Accessor().IsClosed() { t.Error("guardA rebuilt snapshot should be closed after parent Close") } if _, err := guardA.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("guardA DispatchableSnapshot: want Closed, got %v", err) } if guardA.RebuiltSnapshot() != nil { t.Error("guardA rebuiltSnapshot should be nil after fan-out") } if guardA.State() != RebuildStateReleased { t.Errorf("guardA State: want Released, got %v", guardA.State()) } // Guard B: released if guardB.State() != RebuildStateReleased { t.Errorf("guardB State: want Released, got %v", guardB.State()) } if !guardB.IsReleased() { t.Error("guardB should be released") } if _, err := guardB.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("guardB DispatchableSnapshot: want Closed, got %v", err) } // Idempotent close snap.Close() if !snap.Accessor().IsClosed() { t.Error("parent should remain closed after idempotent close") } } // TestIngressSnapshot_RebuildConcurrentOverflowCloseCancelRace verifies concurrent safety // under the race detector when a committed sibling, overflow commit, reserved sibling // cancel, and parent close all execute concurrently. func TestIngressSnapshot_RebuildConcurrentOverflowCloseCancelRace(t *testing.T) { const maxBytes int64 = 1024 * 1024 b, err := NewIngressSnapshotBuilder(maxBytes) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } h, err := b.IssueBackingHandle(make([]byte, 100)) if err != nil { t.Fatalf("IssueBackingHandle: %v", err) } if err := b.SetCanonicalWithHandle(h); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } snap, err := b.Build() if err != nil { t.Fatalf("Build: %v", err) } // Reserve three guards guardA, err := snap.ReserveRebuild(100) if err != nil { t.Fatalf("ReserveRebuild A: %v", err) } guardB, err := snap.ReserveRebuild(100) if err != nil { t.Fatalf("ReserveRebuild B: %v", err) } guardC, err := snap.ReserveRebuild(100) if err != nil { t.Fatalf("ReserveRebuild C: %v", err) } if snap.ReservedTempBytes() != 300 { t.Errorf("ReservedTempBytes: want 300, got %d", snap.ReservedTempBytes()) } // Pre-commit guard A to obtain a rebuilt snapshot pointer that will be // subject to the concurrent fan-out on terminal transition. rebuiltA, err := guardA.Commit([]byte("rebuilt-output-data-for-guard-a")) if err != nil { t.Fatalf("guardA Commit: %v", err) } if rebuiltA == nil { t.Fatal("guardA rebuilt snapshot is nil") } if rebuiltA.Accessor().IsClosed() { t.Fatal("guardA rebuilt snapshot should not be closed yet") } var wg sync.WaitGroup start := make(chan struct{}) overflowResult := make(chan error, 1) cancelResult := make(chan error, 1) // Goroutine 1: Commit guard B with actual oversized output (guaranteed overflow). wg.Add(1) go func() { defer wg.Done() <-start _, err := guardB.Commit(make([]byte, maxBytes)) overflowResult <- err }() // Goroutine 2: Cancel guard C. wg.Add(1) go func() { defer wg.Done() <-start err := guardC.Cancel() cancelResult <- err }() // Goroutine 3: Close parent snapshot. wg.Add(1) go func() { defer wg.Done() <-start snap.Close() }() close(start) wg.Wait() // Block on result channels to guarantee each goroutine completed and produced a value. // This closes the false-positive path where a nil overflowErr would allow an oversized // commit to appear successful if parent close later cleaned up guard B state. overflowErr := <-overflowResult cancelErr := <-cancelResult // Verify parent is closed. if !snap.Accessor().IsClosed() { t.Error("parent snapshot should be closed") } if snap.Accessor().RetainedCount() != 0 { t.Errorf("parent RetainedCount: want 0, got %d", snap.Accessor().RetainedCount()) } if snap.Accessor().RetainedBytes() != 0 { t.Errorf("parent RetainedBytes: want 0, got %d", snap.Accessor().RetainedBytes()) } if snap.ReservedTempBytes() != 0 { t.Errorf("ReservedTempBytes: want 0, got %d", snap.ReservedTempBytes()) } // Verify rebuiltA is closed (fan-out from overflow or parent close). if !rebuiltA.Accessor().IsClosed() { t.Error("rebuiltA should be closed after concurrent terminal transition") } // Verify all guards are released. if !guardA.IsReleased() { t.Error("guardA should be released") } if !guardB.IsReleased() { t.Error("guardB should be released") } if !guardC.IsReleased() { t.Error("guardC should be released") } // Verify rebuilt pointer is cleared on guardA after fan-out. if guardA.RebuiltSnapshot() != nil { t.Error("guardA rebuiltSnapshot should be nil after release") } if guardB.RebuiltSnapshot() != nil { t.Error("guardB rebuiltSnapshot should be nil") } if guardC.RebuiltSnapshot() != nil { t.Error("guardC rebuiltSnapshot should be nil") } // Verify no data references remain on released guards. if guardA.HasTempReference() || guardA.HasCanonicalReference() || guardA.HasTypedReferences() { t.Error("guardA should have no remaining references") } if guardB.HasTempReference() || guardB.HasCanonicalReference() || guardB.HasTypedReferences() { t.Error("guardB should have no remaining references") } if guardC.HasTempReference() || guardC.HasCanonicalReference() || guardC.HasTypedReferences() { t.Error("guardC should have no remaining references") } // Verify no DispatchableSnapshot from released guards. // guardA must be Released (pre-committed, then fanned-out), so Closed. if _, err := guardA.DispatchableSnapshot(); !errors.Is(err, ErrIngressSnapshotClosed) { t.Errorf("guardA DispatchableSnapshot: want Closed, got %v", err) } // guardB can be Overflow or Released; accept either terminal error. if bSnap, bErr := guardB.DispatchableSnapshot(); bSnap != nil { t.Error("guardB DispatchableSnapshot: expected nil snapshot") } else if !errors.Is(bErr, ErrIngressSnapshotRebuildOverflow) && !errors.Is(bErr, ErrIngressSnapshotClosed) { t.Errorf("guardB DispatchableSnapshot: want RebuildOverflow or Closed, got %v", bErr) } // guardC can be Cancelled or Released; accept either terminal error. if cSnap, cErr := guardC.DispatchableSnapshot(); cSnap != nil { t.Error("guardC DispatchableSnapshot: expected nil snapshot") } else if !errors.Is(cErr, ErrIngressSnapshotRebuildCancelled) && !errors.Is(cErr, ErrIngressSnapshotClosed) { t.Errorf("guardC DispatchableSnapshot: want Cancelled or Closed, got %v", cErr) } // Verify winner-specific terminal states for guard B and C. // Guard B: must be Overflow (actual oversized commit) or Released (parent close won). // The Commit error must be paired with the terminal state: Overflow requires // ErrIngressSnapshotRebuildOverflow; Released requires ErrIngressSnapshotClosed. // Any mismatch or nil error is a false-positive indicator. switch bState := guardB.State(); bState { case RebuildStateOverflow: if !errors.Is(overflowErr, ErrIngressSnapshotRebuildOverflow) { t.Errorf("guardB Commit in Overflow state: want RebuildOverflow, got %v", overflowErr) } case RebuildStateReleased: if !errors.Is(overflowErr, ErrIngressSnapshotClosed) { t.Errorf("guardB Commit in Released state: want Closed, got %v", overflowErr) } default: t.Errorf("guardB State: want Overflow or Released, got %v", bState) } // Guard C: must be Cancelled or Released (parent close won). cState := guardC.State() if cState != RebuildStateCancelled && cState != RebuildStateReleased { t.Errorf("guardC State: want Cancelled or Released, got %v", cState) } // Guard C cancel result must be nil (Cancel is idempotent on already-cancelled). if cancelErr != nil { t.Errorf("guardC Cancel: want nil, got %v", cancelErr) } } func TestIngressSnapshot_OwnedBackingTransferAndCommit(t *testing.T) { canonical := []byte(`{"model":"owned"}`) builder, err := NewIngressSnapshotBuilder(1024) if err != nil { t.Fatalf("NewIngressSnapshotBuilder: %v", err) } defer builder.Close() handle, err := builder.IssueOwnedBackingHandle(canonical) if err != nil { t.Fatalf("IssueOwnedBackingHandle: %v", err) } if err := builder.SetCanonicalWithHandle(handle); err != nil { t.Fatalf("SetCanonicalWithHandle: %v", err) } if err := builder.AddTypedView("semantic", handle); err != nil { t.Fatalf("AddTypedView: %v", err) } snapshot, err := builder.Build() if err != nil { t.Fatalf("Build: %v", err) } defer snapshot.Close() alias, err := snapshot.Accessor().CanonicalAlias() if err != nil { t.Fatalf("CanonicalAlias: %v", err) } if len(alias) == 0 || &alias[0] != &canonical[0] { t.Fatal("owned canonical backing was copied") } copied, err := snapshot.Accessor().Canonical() if err != nil { t.Fatalf("Canonical: %v", err) } if &copied[0] == &canonical[0] { t.Fatal("defensive Canonical API unexpectedly returned an alias") } copied[0] = '!' if alias[0] == '!' { t.Fatal("mutating the defensive copy changed the owned backing") } output := []byte(`{"model":"rebuilt","input":"fixed"}`) guard, err := snapshot.ReserveRebuild(int64(len(output))) if err != nil { t.Fatalf("ReserveRebuild: %v", err) } rebuilt, err := guard.CommitOwnedTyped("rebuilt", output) if err != nil { guard.Close() t.Fatalf("CommitOwnedTyped: %v", err) } rebuiltAlias, err := rebuilt.Accessor().TypedViewAlias("rebuilt") if err != nil { t.Fatalf("TypedViewAlias: %v", err) } if len(rebuiltAlias) == 0 || &rebuiltAlias[0] != &output[0] { t.Fatal("owned rebuild output was copied") } wantRetained := int64(len(canonical) + len(output)) if got := rebuilt.Accessor().RetainedBytes(); got != wantRetained { t.Fatalf("rebuilt retained bytes = %d, want %d", got, wantRetained) } rebuiltAccessor := rebuilt.Accessor() rebuilt.Close() guard.Close() if !guard.IsReleased() { t.Fatal("owned rebuild guard was not released") } if _, err := rebuiltAccessor.TypedViewAlias("rebuilt"); !errors.Is(err, ErrIngressSnapshotClosed) { t.Fatalf("TypedViewAlias after release = %v, want closed", err) } }