feat: rotate unsequenced client output safely
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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 {
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user