iop/packages/go/streamgate/ingress_snapshot_test.go
toki ae5845cd68 fix: post-commit refinements
- Update runtime loop plan/code-review
- Fix chat_handler, input_estimator_test, openai_request_rebuilder, provider_tunnel, responses_handler
- Fix ingress_snapshot_test
- Add code review and plan logs
2026-07-26 20:53:32 +09:00

2638 lines
79 KiB
Go

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)
}
}