diff --git a/internal/client/spool/events.go b/internal/client/spool/events.go index 3b341ea..e09b038 100644 --- a/internal/client/spool/events.go +++ b/internal/client/spool/events.go @@ -33,6 +33,13 @@ type EventInput struct { CreatedAt time.Time } +const ( + // These match the CommandEvent oneof field numbers, allowing the runtime to + // reconstruct the envelope without a second event-kind translation table. + EventKindOutput uint32 = 5 + EventKindOutputTruncation uint32 = 10 +) + type Event struct { IssueUUID domain.UUID LocalOrdinal uint64 diff --git a/internal/client/spool/output.go b/internal/client/spool/output.go new file mode 100644 index 0000000..a68d5d3 --- /dev/null +++ b/internal/client/spool/output.go @@ -0,0 +1,203 @@ +package spool + +import ( + "bytes" + "context" + "database/sql" + "errors" + "fmt" + "time" + + "github.com/klauspost/compress/zstd" + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" + "github.com/rvbox/rvbox/internal/domain" + "google.golang.org/protobuf/proto" +) + +const MaxOutputChunkBytes = 64 << 10 + +var ErrPinnedOutput = errors.New("client output is pinned awaiting server acknowledgement") + +type OutputInput struct { + Stream rvboxv1.StreamKind + Raw []byte + ObservedAt time.Time +} + +// AppendOutput stores one raw pipe chunk as a Zstandard OutputChunk. Under +// offline pressure it replaces only the oldest unsequenced output with an +// ordered loss marker; assigned wire events are immutable and never evicted. +func (store *Store) AppendOutput(ctx context.Context, issueUUID domain.UUID, input OutputInput) (Event, error) { + if (input.Stream != rvboxv1.StreamKind_STREAM_STDOUT && input.Stream != rvboxv1.StreamKind_STREAM_STDERR) || uint64(len(input.Raw)) > MaxOutputChunkBytes || input.ObservedAt.IsZero() { + return Event{}, errors.New("invalid raw output chunk") + } + payload, err := encodeOutputChunk(input.Stream, input.Raw) + if err != nil { + return Event{}, err + } + eventInput := EventInput{Kind: EventKindOutput, Compression: 2, RawBytes: uint64(len(input.Raw)), Payload: payload, Output: true, CreatedAt: input.ObservedAt} + for { + event, err := store.AppendEvent(ctx, issueUUID, eventInput) + if !errors.Is(err, ErrCapacityExhausted) { + return event, err + } + rotated, rotateErr := store.RotateOldestUnassignedOutput(ctx, issueUUID, input.ObservedAt) + if rotateErr != nil { + return Event{}, rotateErr + } + if !rotated { + return Event{}, err + } + } +} + +// RotateOldestUnassignedOutput replaces exactly one earliest contiguous run. +// Retaining the first removed local ordinal for the marker preserves durable +// command-local order while leaving every assigned event pinned. +func (store *Store) RotateOldestUnassignedOutput(ctx context.Context, issueUUID domain.UUID, observedAt time.Time) (bool, error) { + if !validUUID(issueUUID) || observedAt.IsZero() { + return false, errors.New("invalid output rotation request") + } + tx, err := store.db.BeginTx(ctx, nil) + if err != nil { + return false, err + } + defer tx.Rollback() + var nextOrdinal, outputCharged, totalCharged, closeout uint64 + err = tx.QueryRowContext(ctx, `SELECT next_local_ordinal, output_charged_bytes, total_charged_bytes, closeout_remaining_bytes FROM commands WHERE issue_uuid = ?`, issueUUID[:]).Scan(&nextOrdinal, &outputCharged, &totalCharged, &closeout) + if err == sql.ErrNoRows { + return false, ErrUnknownCommand + } + if err != nil { + return false, err + } + rows, err := tx.QueryContext(ctx, `SELECT local_ordinal, raw_bytes, charged_bytes, payload FROM events WHERE issue_uuid = ? AND event_seq IS NULL AND output = 1 ORDER BY local_ordinal`, issueUUID[:]) + if err != nil { + return false, err + } + var run []outputRow + for rows.Next() { + var current outputRow + if err := rows.Scan(¤t.ordinal, ¤t.rawBytes, ¤t.charge, ¤t.payload); err != nil { + _ = rows.Close() + return false, err + } + if len(run) > 0 && current.ordinal != run[len(run)-1].ordinal+1 { + break + } + run = append(run, current) + } + if err := rows.Close(); err != nil { + return false, err + } + if len(run) == 0 { + return false, ErrPinnedOutput + } + markerPayload, removedRaw, removedCompressed, removedCharge, err := outputLossMarker(run) + if err != nil { + return false, err + } + markerCharge, err := EstimateCharge(ChargeInput{EncodedBytes: uint64(len(markerPayload)), SQLiteRows: 1, IndexEntries: 2}) + if err != nil { + return false, err + } + if removedCharge > totalCharged || removedCharge > outputCharged { + return false, fmt.Errorf("client spool output charge counter mismatch") + } + clientTotal, err := clientTotalCharge(ctx, tx) + if err != nil { + return false, err + } + if removedCharge > clientTotal { + return false, fmt.Errorf("client spool aggregate charge counter mismatch") + } + remainingTotal := totalCharged - removedCharge + remainingOutput := outputCharged - removedCharge + remainingClient := clientTotal - removedCharge + decision, err := CheckReservation(store.quotaLimits, ReservationState{CommandOutputCharged: remainingOutput, CommandTotalCharged: remainingTotal, ClientTotalCharged: remainingClient, CloseoutRemaining: closeout}, ReservationRequest{ChargedBytes: markerCharge}) + if err != nil { + return false, err + } + for _, current := range run { + if _, err := tx.ExecContext(ctx, `DELETE FROM events WHERE issue_uuid = ? AND local_ordinal = ? AND event_seq IS NULL`, issueUUID[:], current.ordinal); err != nil { + return false, err + } + } + digest := immutableDigest(markerPayload) + if _, err := tx.ExecContext(ctx, `INSERT INTO events(issue_uuid, local_ordinal, event_kind, compression, raw_bytes, charged_bytes, output, payload, payload_sha256, created_at) VALUES (?, ?, ?, 2, ?, ?, 0, ?, ?, ?)`, issueUUID[:], run[0].ordinal, EventKindOutputTruncation, removedRaw, markerCharge, markerPayload, digest[:], observedAt.UnixNano()); err != nil { + return false, err + } + if _, err := tx.ExecContext(ctx, `UPDATE commands SET output_charged_bytes = ?, total_charged_bytes = ?, closeout_remaining_bytes = ?, next_local_ordinal = ? WHERE issue_uuid = ?`, decision.CommandOutputCharged, decision.CommandTotalCharged, decision.CloseoutRemaining, nextOrdinal, issueUUID[:]); err != nil { + return false, err + } + if err := updateClientTotalCharge(ctx, tx, decision.ClientTotalCharged); err != nil { + return false, err + } + if err := tx.Commit(); err != nil { + return false, err + } + _ = removedCompressed // retained in the durable marker payload. + return true, nil +} + +type outputRow struct { + ordinal uint64 + rawBytes uint64 + charge uint64 + payload []byte +} + +func encodeOutputChunk(stream rvboxv1.StreamKind, raw []byte) ([]byte, error) { + encoder, err := zstd.NewWriter(nil, zstd.WithEncoderConcurrency(1)) + if err != nil { + return nil, err + } + compressed := encoder.EncodeAll(raw, nil) + encoder.Close() + chunk := &rvboxv1.OutputChunk{Stream: stream, Compression: rvboxv1.Compression_COMPRESSION_ZSTD, Data: compressed, UncompressedSize: uint64(len(raw)), CompressedSize: uint64(len(compressed))} + return proto.Marshal(chunk) +} + +func outputLossMarker(rows []outputRow) ([]byte, uint64, uint64, uint64, error) { + var raw, compressed, charged uint64 + for _, row := range rows { + var chunk rvboxv1.OutputChunk + if err := proto.Unmarshal(row.payload, &chunk); err != nil || chunk.Compression != rvboxv1.Compression_COMPRESSION_ZSTD || chunk.UncompressedSize != row.rawBytes || chunk.CompressedSize != uint64(len(chunk.Data)) { + return nil, 0, 0, 0, ErrStoredPayloadChecksum + } + if raw > ^uint64(0)-row.rawBytes || compressed > ^uint64(0)-chunk.CompressedSize || charged > ^uint64(0)-row.charge { + return nil, 0, 0, 0, ErrCapacityExhausted + } + raw += row.rawBytes + compressed += chunk.CompressedSize + charged += row.charge + } + marker := &rvboxv1.OutputTruncation{RemovedCompressedBytes: &compressed, RemovedUncompressedBytes: raw, Reason: "CLIENT_SPOOL", Source: rvboxv1.OutputTruncationSource_OUTPUT_TRUNCATION_SOURCE_CLIENT_SPOOL} + payload, err := proto.Marshal(marker) + return payload, raw, compressed, charged, err +} + +func decodeOutputLossMarker(payload []byte) (*rvboxv1.OutputTruncation, error) { + var marker rvboxv1.OutputTruncation + if err := proto.Unmarshal(payload, &marker); err != nil || marker.Source != rvboxv1.OutputTruncationSource_OUTPUT_TRUNCATION_SOURCE_CLIENT_SPOOL || marker.FirstRemovedEventSeq != nil || marker.LastRemovedEventSeq != nil || marker.RemovedCompressedBytes == nil { + return nil, ErrStoredPayloadChecksum + } + return &marker, nil +} + +func outputChunkRaw(payload []byte) ([]byte, error) { + var chunk rvboxv1.OutputChunk + if err := proto.Unmarshal(payload, &chunk); err != nil || (chunk.Stream != rvboxv1.StreamKind_STREAM_STDOUT && chunk.Stream != rvboxv1.StreamKind_STREAM_STDERR) || chunk.Compression != rvboxv1.Compression_COMPRESSION_ZSTD || chunk.CompressedSize != uint64(len(chunk.Data)) { + return nil, ErrStoredPayloadChecksum + } + decoder, err := zstd.NewReader(bytes.NewReader(chunk.Data), zstd.WithDecoderConcurrency(1), zstd.WithDecoderMaxMemory(MaxOutputChunkBytes+1)) + if err != nil { + return nil, ErrStoredPayloadChecksum + } + defer decoder.Close() + raw, err := decoder.DecodeAll(chunk.Data, nil) + if err != nil || uint64(len(raw)) != chunk.UncompressedSize || uint64(len(raw)) > MaxOutputChunkBytes { + return nil, ErrStoredPayloadChecksum + } + return raw, nil +} diff --git a/internal/client/spool/output_test.go b/internal/client/spool/output_test.go new file mode 100644 index 0000000..7a8969f --- /dev/null +++ b/internal/client/spool/output_test.go @@ -0,0 +1,86 @@ +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) + } +} diff --git a/internal/client/spool/recovery.go b/internal/client/spool/recovery.go index 6671867..22970a4 100644 --- a/internal/client/spool/recovery.go +++ b/internal/client/spool/recovery.go @@ -83,13 +83,15 @@ OR c.output_charged_bytes != func checkPayloads(ctx context.Context, database *sql.DB, maxScriptBytes uint64) (RecoveryReport, error) { var report RecoveryReport - rows, err := database.QueryContext(ctx, `SELECT issue_uuid, payload, payload_sha256 FROM events ORDER BY issue_uuid, local_ordinal`) + rows, err := database.QueryContext(ctx, `SELECT issue_uuid, event_kind, raw_bytes, payload, payload_sha256 FROM events ORDER BY issue_uuid, local_ordinal`) if err != nil { return report, err } for rows.Next() { var owner, payload, digest []byte - if err := rows.Scan(&owner, &payload, &digest); err != nil { + var kind uint32 + var rawBytes uint64 + if err := rows.Scan(&owner, &kind, &rawBytes, &payload, &digest); err != nil { _ = rows.Close() return report, err } @@ -98,6 +100,20 @@ func checkPayloads(ctx context.Context, database *sql.DB, maxScriptBytes uint64) _ = rows.Close() return report, ErrStoredPayloadChecksum } + switch kind { + case EventKindOutput: + raw, err := outputChunkRaw(payload) + if err != nil || uint64(len(raw)) != rawBytes { + _ = rows.Close() + return report, ErrStoredPayloadChecksum + } + case EventKindOutputTruncation: + marker, err := decodeOutputLossMarker(payload) + if err != nil || marker.RemovedUncompressedBytes != rawBytes { + _ = rows.Close() + return report, ErrStoredPayloadChecksum + } + } report.EventsChecked++ } if err := rows.Close(); err != nil { diff --git a/test/coverage.toml b/test/coverage.toml index b901c10..8ad2f89 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -215,6 +215,18 @@ tests = [ "internal/client/spool/quota_test.go:TestSpoolAcknowledgementReleasesQuota_HP_CLIENT_08", ] +[[requirements]] +id = "HP-OUT-01" +layer = "unit" +status = "implemented" +tests = ["internal/client/spool/output_test.go:TestAppendOutputRotatesOnlyUnsequencedRun_HP_OUT_01"] + +[[requirements]] +id = "BH-OUTFLOW-02" +layer = "unit" +status = "implemented" +tests = ["internal/client/spool/output_test.go:TestAppendOutputNeverEvictsAssignedEvent_BH_OUTFLOW_02"] + [[requirements]] id = "BH-RET-01" layer = "unit"