Files
rvbox/internal/client/spool/output_test.go
T

87 lines
3.5 KiB
Go

package spool
import (
"bytes"
"context"
"errors"
"path/filepath"
"testing"
"time"
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
)
func TestAppendOutputRotatesOnlyUnsequencedRun_HP_OUT_01(t *testing.T) {
t.Parallel()
ctx := context.Background()
limits := QuotaLimits{HardAllocationBytes: 1_000, CommandOutputBytes: 900, CommandTotalBytes: 3_000, ClientTotalBytes: 4_000, CloseoutReserveBytes: 400}
store, err := Open(ctx, Options{DataDir: filepath.Join(t.TempDir(), "spool"), BusyTimeout: time.Second, QuotaLimits: limits})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = store.Close() })
issue := testUUID(t, "019c46f1-1d02-7000-8000-000000000051")
if _, err := store.AcceptCommand(ctx, testCommand(issue, []byte("output")), time.Now().UTC()); err != nil {
t.Fatal(err)
}
first := bytes.Repeat([]byte("a"), 200)
second := bytes.Repeat([]byte("b"), 200)
third := bytes.Repeat([]byte("c"), 200)
for _, raw := range [][]byte{first, second, third} {
if _, err := store.AppendOutput(ctx, issue, OutputInput{Stream: rvboxv1.StreamKind_STREAM_STDOUT, Raw: raw, ObservedAt: time.Now().UTC()}); err != nil {
t.Fatalf("AppendOutput(%q): %v", raw[:1], err)
}
}
assigned, err := store.AssignSendWindow(ctx, issue, 8, 1<<20)
if err != nil {
t.Fatal(err)
}
if len(assigned) != 2 || assigned[0].Kind != EventKindOutputTruncation || assigned[0].LocalOrdinal != 1 || assigned[0].EventSeq != 1 || assigned[1].Kind != EventKindOutput || assigned[1].LocalOrdinal != 3 || assigned[1].EventSeq != 2 {
t.Fatalf("assigned output/loss events = %#v", assigned)
}
marker, err := decodeOutputLossMarker(assigned[0].Payload)
if err != nil {
t.Fatal(err)
}
if marker.RemovedUncompressedBytes != uint64(len(first)+len(second)) || marker.GetRemovedCompressedBytes() == 0 {
t.Fatalf("loss marker = %#v", marker)
}
raw, err := outputChunkRaw(assigned[1].Payload)
if err != nil || !bytes.Equal(raw, third) {
t.Fatalf("stored output = %q, %v", raw, err)
}
report, err := store.Check(ctx)
if err != nil || report.EventsChecked != 2 {
t.Fatalf("output spool Check = %#v, %v", report, err)
}
}
func TestAppendOutputNeverEvictsAssignedEvent_BH_OUTFLOW_02(t *testing.T) {
t.Parallel()
ctx := context.Background()
limits := QuotaLimits{HardAllocationBytes: 1_000, CommandOutputBytes: 600, CommandTotalBytes: 3_000, ClientTotalBytes: 4_000, CloseoutReserveBytes: 400}
store, err := Open(ctx, Options{DataDir: filepath.Join(t.TempDir(), "spool"), BusyTimeout: time.Second, QuotaLimits: limits})
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = store.Close() })
issue := testUUID(t, "019c46f1-1d02-7000-8000-000000000052")
if _, err := store.AcceptCommand(ctx, testCommand(issue, []byte("output")), time.Now().UTC()); err != nil {
t.Fatal(err)
}
raw := bytes.Repeat([]byte("x"), 200)
if _, err := store.AppendOutput(ctx, issue, OutputInput{Stream: rvboxv1.StreamKind_STREAM_STDERR, Raw: raw, ObservedAt: time.Now().UTC()}); err != nil {
t.Fatal(err)
}
if _, err := store.AssignSendWindow(ctx, issue, 1, 1<<20); err != nil {
t.Fatal(err)
}
if _, err := store.AppendOutput(ctx, issue, OutputInput{Stream: rvboxv1.StreamKind_STREAM_STDERR, Raw: raw, ObservedAt: time.Now().UTC()}); !errors.Is(err, ErrPinnedOutput) {
t.Fatalf("output behind assigned event error = %v, want ErrPinnedOutput", err)
}
pending, err := store.PendingEvents(ctx, issue)
if err != nil || len(pending) != 1 || pending[0].EventSeq != 1 || pending[0].Kind != EventKindOutput {
t.Fatalf("assigned event was changed: %#v, %v", pending, err)
}
}