- 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
2638 lines
79 KiB
Go
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)
|
|
}
|
|
}
|