diff --git a/docs/dependency-decisions.md b/docs/dependency-decisions.md index 131488c..f73f16b 100644 --- a/docs/dependency-decisions.md +++ b/docs/dependency-decisions.md @@ -9,4 +9,4 @@ All versions are exact in `go.mod`, generated code, or the toolchain image. | protoc 35.0 | Protobuf compiler and well-known includes | Pinned archive with a checked SHA-256; included even though normal generation is driven by Buf. | | protobuf-go 1.36.12 | Go protobuf runtime and generator | Official maintained Go protobuf implementation. | | protoc-gen-go-grpc 1.6.2 | Go gRPC generator | Official maintained gRPC-Go generator. | - +| google/uuid 1.6.0 | Parse canonical UUIDs and verify RFC variant/version bits | Stable maintained package; RVBox owns the monotonic UUIDv7 generator so clock and ordering behavior remain directly testable. | diff --git a/go.mod b/go.mod index ea46290..0e755e2 100644 --- a/go.mod +++ b/go.mod @@ -3,6 +3,7 @@ module github.com/rvbox/rvbox go 1.27.0 require ( + github.com/google/uuid v1.6.0 google.golang.org/grpc v1.83.2 google.golang.org/protobuf v1.36.12 ) diff --git a/internal/domain/errors.go b/internal/domain/errors.go new file mode 100644 index 0000000..73c5e33 --- /dev/null +++ b/internal/domain/errors.go @@ -0,0 +1,83 @@ +package domain + +import ( + "fmt" + "unicode/utf8" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +const MaxProtocolDetailBytes = 4 * 1024 + +// Error is the transport-neutral form of an RVBox structured failure. +type Error struct { + Code rvboxv1.ControlError_Code + Message string + Retryable bool + IssueUUID string + Cause error +} + +func (e *Error) Error() string { + if e.IssueUUID == "" { + return fmt.Sprintf("%s: %s", e.Code, e.Message) + } + return fmt.Sprintf("%s (%s): %s", e.Code, e.IssueUUID, e.Message) +} + +func (e *Error) Unwrap() error { return e.Cause } + +// NewError applies the stable retryability policy and bounds wire-visible text. +func NewError(code rvboxv1.ControlError_Code, message, issueUUID string) *Error { + return &Error{ + Code: code, + Message: boundUTF8(message, MaxProtocolDetailBytes), + Retryable: DefaultRetryable(code), + IssueUUID: issueUUID, + } +} + +func (e *Error) WithCause(cause error) *Error { + e.Cause = cause + return e +} + +func (e *Error) ControlError() *rvboxv1.ControlError { + return &rvboxv1.ControlError{ + Code: e.Code, + Message: boundUTF8(e.Message, MaxProtocolDetailBytes), + Retryable: e.Retryable, + IssueUuid: e.IssueUUID, + } +} + +func DefaultRetryable(code rvboxv1.ControlError_Code) bool { + switch code { + case rvboxv1.ControlError_OFFLINE, + rvboxv1.ControlError_CAPACITY_EXHAUSTED, + rvboxv1.ControlError_TRANSIENT, + rvboxv1.ControlError_INTERNAL: + return true + default: + return false + } +} + +func NewExecutionContextUnavailable(message, issueUUID string) *Error { + return NewError(rvboxv1.ControlError_CODE_EXECUTION_CONTEXT_UNAVAILABLE, message, issueUUID) +} + +func NewElevationUnavailable(message, issueUUID string) *Error { + return NewError(rvboxv1.ControlError_CODE_ELEVATION_UNAVAILABLE, message, issueUUID) +} + +func boundUTF8(value string, limit int) string { + if len(value) <= limit { + return value + } + value = value[:limit] + for !utf8.ValidString(value) { + value = value[:len(value)-1] + } + return value +} diff --git a/internal/domain/errors_test.go b/internal/domain/errors_test.go new file mode 100644 index 0000000..36e54b6 --- /dev/null +++ b/internal/domain/errors_test.go @@ -0,0 +1,52 @@ +package domain + +import ( + "errors" + "strings" + "testing" + "unicode/utf8" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +func TestDomainErrorMapping_HP_PROTO_09(t *testing.T) { + t.Parallel() + + retryable := map[rvboxv1.ControlError_Code]bool{ + rvboxv1.ControlError_OFFLINE: true, + rvboxv1.ControlError_CAPACITY_EXHAUSTED: true, + rvboxv1.ControlError_TRANSIENT: true, + rvboxv1.ControlError_INTERNAL: true, + } + for code := rvboxv1.ControlError_Code(0); code <= rvboxv1.ControlError_CODE_ELEVATION_UNAVAILABLE; code++ { + domainErr := NewError(code, "message", "01890a5d-ac96-7a11-b234-123456789abc") + wire := domainErr.ControlError() + if wire.Code != code || wire.Message != "message" || wire.Retryable != retryable[code] || wire.IssueUuid != domainErr.IssueUUID { + t.Errorf("mapping for %s = %+v", code, wire) + } + } +} + +func TestDomainErrorBoundsAndCause_HP_PROTO_09(t *testing.T) { + t.Parallel() + + cause := errors.New("cause") + domainErr := NewError(rvboxv1.ControlError_INTERNAL, strings.Repeat("界", MaxProtocolDetailBytes), "").WithCause(cause) + if !errors.Is(domainErr, cause) { + t.Fatal("domain error did not unwrap its cause") + } + if len(domainErr.Message) > MaxProtocolDetailBytes || !utf8.ValidString(domainErr.Message) { + t.Fatalf("bounded message has %d bytes or invalid UTF-8", len(domainErr.Message)) + } +} + +func TestWindowsPrelaunchErrorCodes_HP_WINCTX_01(t *testing.T) { + t.Parallel() + + if got := NewExecutionContextUnavailable("normal", "").Code; got != rvboxv1.ControlError_CODE_EXECUTION_CONTEXT_UNAVAILABLE { + t.Fatalf("normal context code = %s", got) + } + if got := NewElevationUnavailable("elevated", "").Code; got != rvboxv1.ControlError_CODE_ELEVATION_UNAVAILABLE { + t.Fatalf("elevation code = %s", got) + } +} diff --git a/internal/domain/id.go b/internal/domain/id.go new file mode 100644 index 0000000..9928fac --- /dev/null +++ b/internal/domain/id.go @@ -0,0 +1,124 @@ +// Package domain contains transport- and storage-independent RVBox rules. +package domain + +import ( + "crypto/rand" + "errors" + "fmt" + "io" + "sync" + "time" + + googleuuid "github.com/google/uuid" +) + +var ( + ErrInvalidUUIDv7 = errors.New("invalid canonical UUIDv7") + ErrUUIDTimeRange = errors.New("UUIDv7 time is outside its 48-bit range") +) + +const maxUUIDv7UnixMilli = uint64(1<<48 - 1) + +// UUID is the canonical binary form used by RVBox domain and storage code. +type UUID [16]byte + +// String returns the lowercase RFC 9562 canonical representation. +func (u UUID) String() string { + return googleuuid.UUID(u).String() +} + +// ParseUUIDv7 accepts only the lowercase, hyphenated RFC 9562 form. +func ParseUUIDv7(value string) (UUID, error) { + if len(value) != 36 { + return UUID{}, fmt.Errorf("%w: expected 36 characters", ErrInvalidUUIDv7) + } + + parsed, err := googleuuid.Parse(value) + if err != nil || parsed.String() != value { + return UUID{}, fmt.Errorf("%w: non-canonical encoding", ErrInvalidUUIDv7) + } + if parsed.Version() != 7 || parsed.Variant() != googleuuid.RFC4122 { + return UUID{}, fmt.Errorf("%w: expected RFC variant and version 7", ErrInvalidUUIDv7) + } + + return UUID(parsed), nil +} + +// UUIDv7Generator emits lexically increasing UUIDv7 values even when the wall +// clock stalls or moves backwards. Calls are safe from multiple goroutines. +type UUIDv7Generator struct { + mu sync.Mutex + now func() time.Time + entropy io.Reader + started bool + lastMS uint64 + sequence uint16 +} + +// NewUUIDv7Generator creates a generator. Nil arguments select the system clock +// and crypto/rand.Reader. Supplying dependencies keeps boundary tests hermetic. +func NewUUIDv7Generator(now func() time.Time, entropy io.Reader) *UUIDv7Generator { + if now == nil { + now = time.Now + } + if entropy == nil { + entropy = rand.Reader + } + return &UUIDv7Generator{now: now, entropy: entropy} +} + +var defaultUUIDv7Generator = NewUUIDv7Generator(nil, nil) + +// NewUUIDv7 uses the process-wide generator. +func NewUUIDv7() (UUID, error) { + return defaultUUIDv7Generator.New() +} + +// New returns a canonical RFC 9562 UUIDv7. +func (g *UUIDv7Generator) New() (UUID, error) { + g.mu.Lock() + defer g.mu.Unlock() + + instant := g.now() + millis := instant.UnixMilli() + if millis < 0 || uint64(millis) > maxUUIDv7UnixMilli { + return UUID{}, ErrUUIDTimeRange + } + + var random [10]byte + if _, err := io.ReadFull(g.entropy, random[:]); err != nil { + return UUID{}, fmt.Errorf("read UUIDv7 entropy: %w", err) + } + + wallMS := uint64(millis) + logicalMS := wallMS + sequence := uint16(random[0]&0x0f)<<8 | uint16(random[1]) + if g.started && wallMS <= g.lastMS { + logicalMS = g.lastMS + sequence = g.sequence + 1 + if sequence > 0x0fff { + if logicalMS == maxUUIDv7UnixMilli { + return UUID{}, ErrUUIDTimeRange + } + logicalMS++ + sequence = 0 + } + } + + var result UUID + result[0] = byte(logicalMS >> 40) + result[1] = byte(logicalMS >> 32) + result[2] = byte(logicalMS >> 24) + result[3] = byte(logicalMS >> 16) + result[4] = byte(logicalMS >> 8) + result[5] = byte(logicalMS) + result[6] = 0x70 | byte(sequence>>8) + result[7] = byte(sequence) + result[8] = 0x80 | random[2]&0x3f + copy(result[9:], random[3:]) + + g.started = true + g.lastMS = logicalMS + g.sequence = sequence + return result, nil +} diff --git a/internal/domain/id_test.go b/internal/domain/id_test.go new file mode 100644 index 0000000..fdf9e32 --- /dev/null +++ b/internal/domain/id_test.go @@ -0,0 +1,135 @@ +package domain + +import ( + "bytes" + "errors" + "sort" + "sync" + "testing" + "time" +) + +func TestUUIDv7CanonicalValidation_HP_IDEM_01(t *testing.T) { + t.Parallel() + + valid := "01890a5d-ac96-7a11-b234-123456789abc" + parsed, err := ParseUUIDv7(valid) + if err != nil { + t.Fatalf("ParseUUIDv7(%q): %v", valid, err) + } + if parsed.String() != valid { + t.Fatalf("round trip = %q, want %q", parsed.String(), valid) + } + + for _, value := range []string{ + "", + "01890A5D-AC96-7A11-B234-123456789ABC", + "01890a5dac967a11b234123456789abc", + "{01890a5d-ac96-7a11-b234-123456789abc}", + "01890a5d-ac96-4a11-b234-123456789abc", + "01890a5d-ac96-7a11-7234-123456789abc", + } { + if _, err := ParseUUIDv7(value); !errors.Is(err, ErrInvalidUUIDv7) { + t.Errorf("ParseUUIDv7(%q) error = %v, want ErrInvalidUUIDv7", value, err) + } + } +} + +func TestUUIDv7GeneratorMonotonicAcrossClockRollback_HP_IDEM_10(t *testing.T) { + t.Parallel() + + times := []time.Time{ + time.UnixMilli(1_700_000_000_000), + time.UnixMilli(1_700_000_000_000), + time.UnixMilli(1_699_999_999_000), + time.UnixMilli(1_700_000_000_001), + } + next := 0 + generator := NewUUIDv7Generator(func() time.Time { + value := times[next] + next++ + return value + }, bytes.NewReader(bytes.Repeat([]byte{0x11}, 40))) + + var previous string + for range times { + id, err := generator.New() + if err != nil { + t.Fatalf("New(): %v", err) + } + if _, err := ParseUUIDv7(id.String()); err != nil { + t.Fatalf("generated invalid UUID %q: %v", id, err) + } + if previous != "" && id.String() <= previous { + t.Fatalf("UUID %q did not sort after %q", id, previous) + } + previous = id.String() + } +} + +func TestUUIDv7GeneratorConcurrentOrderingAndUniqueness_HP_IDEM_10(t *testing.T) { + t.Parallel() + + const count = 512 + generator := NewUUIDv7Generator( + func() time.Time { return time.UnixMilli(1_700_000_000_000) }, + bytes.NewReader(bytes.Repeat([]byte{0x42}, count*10)), + ) + ids := make(chan string, count) + errCh := make(chan error, count) + var group sync.WaitGroup + for range count { + group.Add(1) + go func() { + defer group.Done() + id, err := generator.New() + if err != nil { + errCh <- err + return + } + ids <- id.String() + }() + } + group.Wait() + close(ids) + close(errCh) + for err := range errCh { + t.Errorf("New(): %v", err) + } + + ordered := make([]string, 0, count) + seen := make(map[string]struct{}, count) + for id := range ids { + if _, duplicate := seen[id]; duplicate { + t.Fatalf("duplicate UUID %q", id) + } + seen[id] = struct{}{} + ordered = append(ordered, id) + } + if len(ordered) != count { + t.Fatalf("generated %d UUIDs, want %d", len(ordered), count) + } + sort.Strings(ordered) + for index := 1; index < len(ordered); index++ { + if ordered[index] <= ordered[index-1] { + t.Fatalf("UUID ordering failed at %d", index) + } + } +} + +func TestUUIDv7GeneratorErrors_HP_IDEM_01(t *testing.T) { + t.Parallel() + + t.Run("entropy", func(t *testing.T) { + generator := NewUUIDv7Generator(time.Now, bytes.NewReader(nil)) + if _, err := generator.New(); err == nil { + t.Fatal("New() succeeded with empty entropy") + } + }) + t.Run("before epoch", func(t *testing.T) { + generator := NewUUIDv7Generator(func() time.Time { return time.UnixMilli(-1) }, bytes.NewReader(make([]byte, 10))) + if _, err := generator.New(); !errors.Is(err, ErrUUIDTimeRange) { + t.Fatalf("New() error = %v, want ErrUUIDTimeRange", err) + } + }) +} diff --git a/internal/domain/idempotency.go b/internal/domain/idempotency.go new file mode 100644 index 0000000..bee5e00 --- /dev/null +++ b/internal/domain/idempotency.go @@ -0,0 +1,27 @@ +package domain + +import "crypto/sha256" + +type RequestHash [sha256.Size]byte + +func HashImmutableRequest(canonical []byte) RequestHash { return sha256.Sum256(canonical) } + +type ReplayDecision uint8 + +const ( + ReplayNew ReplayDecision = iota + ReplayExact + ReplayConflict +) + +// ClassifyReplay compares a canonical immutable request against durable state. +// A nil stored hash means no mutation with this request UUID exists. +func ClassifyReplay(stored *RequestHash, incoming RequestHash) ReplayDecision { + if stored == nil { + return ReplayNew + } + if *stored == incoming { + return ReplayExact + } + return ReplayConflict +} diff --git a/internal/domain/idempotency_test.go b/internal/domain/idempotency_test.go new file mode 100644 index 0000000..88013f8 --- /dev/null +++ b/internal/domain/idempotency_test.go @@ -0,0 +1,20 @@ +package domain + +import "testing" + +func TestImmutableRequestReplay_HP_IDEM_02(t *testing.T) { + t.Parallel() + + first := HashImmutableRequest([]byte("method\x00target\x00payload")) + same := HashImmutableRequest([]byte("method\x00target\x00payload")) + other := HashImmutableRequest([]byte("method\x00other\x00payload")) + if got := ClassifyReplay(nil, first); got != ReplayNew { + t.Fatalf("new replay decision = %v", got) + } + if got := ClassifyReplay(&first, same); got != ReplayExact { + t.Fatalf("exact replay decision = %v", got) + } + if got := ClassifyReplay(&first, other); got != ReplayConflict { + t.Fatalf("conflicting replay decision = %v", got) + } +} diff --git a/internal/domain/lifecycle.go b/internal/domain/lifecycle.go new file mode 100644 index 0000000..f49a6ce --- /dev/null +++ b/internal/domain/lifecycle.go @@ -0,0 +1,85 @@ +package domain + +import ( + "fmt" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +// CanTransition reports whether a public lifecycle change is legal. Duplicate +// events are handled by sequence validation and are not state transitions. +func CanTransition(from, to rvboxv1.CommandLifecycle) bool { + switch from { + case rvboxv1.CommandLifecycle_COMMAND_QUEUED: + return to == rvboxv1.CommandLifecycle_COMMAND_DISPATCHED || + to == rvboxv1.CommandLifecycle_COMMAND_CANCELLED || + to == rvboxv1.CommandLifecycle_COMMAND_EXPIRED + case rvboxv1.CommandLifecycle_COMMAND_DISPATCHED: + return to == rvboxv1.CommandLifecycle_COMMAND_QUEUED || + to == rvboxv1.CommandLifecycle_COMMAND_ACCEPTED || + to == rvboxv1.CommandLifecycle_COMMAND_REJECTED || + to == rvboxv1.CommandLifecycle_COMMAND_CANCELLED || + to == rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED + case rvboxv1.CommandLifecycle_COMMAND_ACCEPTED: + return to == rvboxv1.CommandLifecycle_COMMAND_RUNNING || + to == rvboxv1.CommandLifecycle_COMMAND_REJECTED || + to == rvboxv1.CommandLifecycle_COMMAND_CANCELLED || + to == rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED + case rvboxv1.CommandLifecycle_COMMAND_RUNNING: + return to == rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED || + to == rvboxv1.CommandLifecycle_COMMAND_FAILED || + to == rvboxv1.CommandLifecycle_COMMAND_TERMINATED || + to == rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED + default: + return false + } +} + +// ValidateTransition rejects unspecified, repeated, terminal-origin, and +// otherwise impossible public lifecycle changes. +func ValidateTransition(from, to rvboxv1.CommandLifecycle) error { + if CanTransition(from, to) { + return nil + } + return NewError( + rvboxv1.ControlError_PROTOCOL_ERROR, + fmt.Sprintf("illegal command lifecycle transition %s -> %s", from, to), + "", + ) +} + +// IsTerminal reports whether no further public lifecycle transition is legal. +func IsTerminal(state rvboxv1.CommandLifecycle) bool { + switch state { + case rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED, + rvboxv1.CommandLifecycle_COMMAND_FAILED, + rvboxv1.CommandLifecycle_COMMAND_TERMINATED, + rvboxv1.CommandLifecycle_COMMAND_CANCELLED, + rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED, + rvboxv1.CommandLifecycle_COMMAND_EXPIRED, + rvboxv1.CommandLifecycle_COMMAND_REJECTED: + return true + default: + return false + } +} + +// LaunchPhase is durable client-only state between accepted and running. +type LaunchPhase uint8 + +const ( + LaunchPhaseNone LaunchPhase = iota + LaunchPhasePrepared + LaunchPhaseAuthorized +) + +// CanAdvanceLaunchPhase permits only the one-way launch barrier. +func CanAdvanceLaunchPhase(from, to LaunchPhase) bool { + return (from == LaunchPhaseNone && to == LaunchPhasePrepared) || + (from == LaunchPhasePrepared && to == LaunchPhaseAuthorized) +} + +// MayRedispatch is false once command code may have been released. +func MayRedispatch(phase LaunchPhase) bool { + return phase == LaunchPhaseNone || phase == LaunchPhasePrepared +} diff --git a/internal/domain/lifecycle_test.go b/internal/domain/lifecycle_test.go new file mode 100644 index 0000000..df97f6e --- /dev/null +++ b/internal/domain/lifecycle_test.go @@ -0,0 +1,102 @@ +package domain + +import ( + "testing" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +func TestLifecycleTransitionMatrix_HP_CMD_01(t *testing.T) { + t.Parallel() + + states := []rvboxv1.CommandLifecycle{ + rvboxv1.CommandLifecycle_COMMAND_LIFECYCLE_UNSPECIFIED, + rvboxv1.CommandLifecycle_COMMAND_QUEUED, + rvboxv1.CommandLifecycle_COMMAND_DISPATCHED, + rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, + rvboxv1.CommandLifecycle_COMMAND_RUNNING, + rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED, + rvboxv1.CommandLifecycle_COMMAND_FAILED, + rvboxv1.CommandLifecycle_COMMAND_TERMINATED, + rvboxv1.CommandLifecycle_COMMAND_CANCELLED, + rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED, + rvboxv1.CommandLifecycle_COMMAND_EXPIRED, + rvboxv1.CommandLifecycle_COMMAND_REJECTED, + } + legal := map[[2]rvboxv1.CommandLifecycle]bool{ + {rvboxv1.CommandLifecycle_COMMAND_QUEUED, rvboxv1.CommandLifecycle_COMMAND_DISPATCHED}: true, + {rvboxv1.CommandLifecycle_COMMAND_QUEUED, rvboxv1.CommandLifecycle_COMMAND_CANCELLED}: true, + {rvboxv1.CommandLifecycle_COMMAND_QUEUED, rvboxv1.CommandLifecycle_COMMAND_EXPIRED}: true, + {rvboxv1.CommandLifecycle_COMMAND_DISPATCHED, rvboxv1.CommandLifecycle_COMMAND_QUEUED}: true, + {rvboxv1.CommandLifecycle_COMMAND_DISPATCHED, rvboxv1.CommandLifecycle_COMMAND_ACCEPTED}: true, + {rvboxv1.CommandLifecycle_COMMAND_DISPATCHED, rvboxv1.CommandLifecycle_COMMAND_REJECTED}: true, + {rvboxv1.CommandLifecycle_COMMAND_DISPATCHED, rvboxv1.CommandLifecycle_COMMAND_CANCELLED}: true, + {rvboxv1.CommandLifecycle_COMMAND_DISPATCHED, rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED}: true, + {rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_RUNNING}: true, + {rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_REJECTED}: true, + {rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_CANCELLED}: true, + {rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED}: true, + {rvboxv1.CommandLifecycle_COMMAND_RUNNING, rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED}: true, + {rvboxv1.CommandLifecycle_COMMAND_RUNNING, rvboxv1.CommandLifecycle_COMMAND_FAILED}: true, + {rvboxv1.CommandLifecycle_COMMAND_RUNNING, rvboxv1.CommandLifecycle_COMMAND_TERMINATED}: true, + {rvboxv1.CommandLifecycle_COMMAND_RUNNING, rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED}: true, + } + + for _, from := range states { + for _, to := range states { + want := legal[[2]rvboxv1.CommandLifecycle{from, to}] + if got := CanTransition(from, to); got != want { + t.Errorf("CanTransition(%s, %s) = %t, want %t", from, to, got, want) + } + if err := ValidateTransition(from, to); (err == nil) != want { + t.Errorf("ValidateTransition(%s, %s) error = %v, legal = %t", from, to, err, want) + } + } + } +} + +func TestTerminalPredicateExhaustive_HP_CMD_01(t *testing.T) { + t.Parallel() + + terminal := map[rvboxv1.CommandLifecycle]bool{ + rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED: true, + rvboxv1.CommandLifecycle_COMMAND_FAILED: true, + rvboxv1.CommandLifecycle_COMMAND_TERMINATED: true, + rvboxv1.CommandLifecycle_COMMAND_CANCELLED: true, + rvboxv1.CommandLifecycle_COMMAND_INTERRUPTED: true, + rvboxv1.CommandLifecycle_COMMAND_EXPIRED: true, + rvboxv1.CommandLifecycle_COMMAND_REJECTED: true, + } + for value := rvboxv1.CommandLifecycle(0); value <= rvboxv1.CommandLifecycle_COMMAND_REJECTED; value++ { + if got := IsTerminal(value); got != terminal[value] { + t.Errorf("IsTerminal(%s) = %t, want %t", value, got, terminal[value]) + } + } +} + +func TestLaunchBarrier_HP_LAUNCH_01(t *testing.T) { + t.Parallel() + + if !CanAdvanceLaunchPhase(LaunchPhaseNone, LaunchPhasePrepared) || + !CanAdvanceLaunchPhase(LaunchPhasePrepared, LaunchPhaseAuthorized) { + t.Fatal("required launch barrier transitions rejected") + } + for from := LaunchPhaseNone; from <= LaunchPhaseAuthorized; from++ { + for to := LaunchPhaseNone; to <= LaunchPhaseAuthorized; to++ { + want := (from == LaunchPhaseNone && to == LaunchPhasePrepared) || + (from == LaunchPhasePrepared && to == LaunchPhaseAuthorized) + if CanAdvanceLaunchPhase(from, to) != want { + t.Errorf("CanAdvanceLaunchPhase(%d, %d) mismatch", from, to) + } + } + } + if MayRedispatch(LaunchPhaseAuthorized) { + t.Fatal("authorized launch may be redispatched") + } + if !MayRedispatch(LaunchPhaseNone) || !MayRedispatch(LaunchPhasePrepared) { + t.Fatal("pre-authorization phase unexpectedly forbids redispatch") + } + if MayRedispatch(LaunchPhase(255)) { + t.Fatal("unknown launch phase may be redispatched") + } +} diff --git a/internal/domain/revision.go b/internal/domain/revision.go new file mode 100644 index 0000000..12a60e8 --- /dev/null +++ b/internal/domain/revision.go @@ -0,0 +1,63 @@ +package domain + +import ( + "errors" + "math" +) + +var ( + ErrInvalidRevision = errors.New("command revision must be nonzero") + ErrRevisionOverflow = errors.New("command revision overflow") + ErrStaleRevision = errors.New("stale command revision") + ErrFutureRevision = errors.New("unexpected future command revision") +) + +type CommandRevision uint64 + +const InitialCommandRevision CommandRevision = 1 + +type RevisionRelation uint8 + +const ( + RevisionInvalid RevisionRelation = iota + RevisionStale + RevisionCurrent + RevisionFuture +) + +func CompareRevision(current, incoming CommandRevision) RevisionRelation { + if current == 0 || incoming == 0 { + return RevisionInvalid + } + if incoming < current { + return RevisionStale + } + if incoming > current { + return RevisionFuture + } + return RevisionCurrent +} + +func ValidateCurrentRevision(current, incoming CommandRevision) error { + switch CompareRevision(current, incoming) { + case RevisionCurrent: + return nil + case RevisionStale: + return ErrStaleRevision + case RevisionFuture: + return ErrFutureRevision + default: + return ErrInvalidRevision + } +} + +// NextRevision durably orders a new cancel or signal intent. +func NextRevision(current CommandRevision) (CommandRevision, error) { + if current == 0 { + return 0, ErrInvalidRevision + } + if current == math.MaxUint64 { + return 0, ErrRevisionOverflow + } + return current + 1, nil +} diff --git a/internal/domain/revision_test.go b/internal/domain/revision_test.go new file mode 100644 index 0000000..c2608d4 --- /dev/null +++ b/internal/domain/revision_test.go @@ -0,0 +1,47 @@ +package domain + +import ( + "errors" + "math" + "testing" +) + +func TestRevisionComparisonAndAdvance_HP_CMD_01(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + current CommandRevision + incoming CommandRevision + want RevisionRelation + wantErr error + }{ + {"zero current", 0, 1, RevisionInvalid, ErrInvalidRevision}, + {"zero incoming", 1, 0, RevisionInvalid, ErrInvalidRevision}, + {"stale", 4, 3, RevisionStale, ErrStaleRevision}, + {"current", 4, 4, RevisionCurrent, nil}, + {"future", 4, 5, RevisionFuture, ErrFutureRevision}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + t.Parallel() + if got := CompareRevision(test.current, test.incoming); got != test.want { + t.Fatalf("CompareRevision() = %v, want %v", got, test.want) + } + err := ValidateCurrentRevision(test.current, test.incoming) + if !errors.Is(err, test.wantErr) { + t.Fatalf("ValidateCurrentRevision() error = %v, want %v", err, test.wantErr) + } + }) + } + + if next, err := NextRevision(InitialCommandRevision); err != nil || next != 2 { + t.Fatalf("NextRevision(1) = (%d, %v), want (2, nil)", next, err) + } + if _, err := NextRevision(0); !errors.Is(err, ErrInvalidRevision) { + t.Fatalf("NextRevision(0) error = %v", err) + } + if _, err := NextRevision(math.MaxUint64); !errors.Is(err, ErrRevisionOverflow) { + t.Fatalf("NextRevision(max) error = %v", err) + } +} diff --git a/internal/domain/sequence.go b/internal/domain/sequence.go new file mode 100644 index 0000000..9d7e190 --- /dev/null +++ b/internal/domain/sequence.go @@ -0,0 +1,67 @@ +package domain + +import ( + "crypto/sha256" + "errors" + "fmt" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +var ( + ErrInvalidEventSequence = errors.New("event sequence must be nonzero") + ErrEventSequenceGap = errors.New("event sequence gap") + ErrConflictingDuplicate = errors.New("conflicting duplicate event") + ErrMissingDuplicate = errors.New("duplicate event has no durable fingerprint") + ErrInvalidTruncationRange = errors.New("invalid output truncation range") +) + +type EventFingerprint [sha256.Size]byte + +func FingerprintBytes(value []byte) EventFingerprint { return sha256.Sum256(value) } + +type EventSequenceDecision uint8 + +const ( + EventSequenceInvalid EventSequenceDecision = iota + EventSequenceAppend + EventSequenceDuplicate +) + +// ClassifyEventSequence admits exactly the next event, or an immutable replay +// whose stored fingerprint is supplied. Wire event-sequence gaps are never +// repaired by inventing events. +func ClassifyEventSequence(last, incoming uint64, stored *EventFingerprint, candidate EventFingerprint) (EventSequenceDecision, error) { + if incoming == 0 { + return EventSequenceInvalid, ErrInvalidEventSequence + } + if last != ^uint64(0) && incoming == last+1 { + return EventSequenceAppend, nil + } + if incoming > last { + return EventSequenceInvalid, fmt.Errorf("%w: have %d, received %d", ErrEventSequenceGap, last, incoming) + } + if stored == nil { + return EventSequenceInvalid, ErrMissingDuplicate + } + if *stored != candidate { + return EventSequenceInvalid, ErrConflictingDuplicate + } + return EventSequenceDuplicate, nil +} + +// ValidateOutputTruncation checks the two allowed forms: both range endpoints +// absent for pre-sequencing client loss, or a complete nonempty retained range. +func ValidateOutputTruncation(value *rvboxv1.OutputTruncation) error { + if value == nil { + return ErrInvalidTruncationRange + } + first, last := value.FirstRemovedEventSeq, value.LastRemovedEventSeq + if first == nil && last == nil { + return nil + } + if first == nil || last == nil || *first == 0 || *last == 0 || *first > *last { + return ErrInvalidTruncationRange + } + return nil +} diff --git a/internal/domain/sequence_test.go b/internal/domain/sequence_test.go new file mode 100644 index 0000000..5bf114a --- /dev/null +++ b/internal/domain/sequence_test.go @@ -0,0 +1,62 @@ +package domain + +import ( + "errors" + "testing" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +func TestEventSequenceAndDuplicateEquivalence_HP_CLIENT_07(t *testing.T) { + t.Parallel() + + first := FingerprintBytes([]byte("first")) + other := FingerprintBytes([]byte("other")) + if decision, err := ClassifyEventSequence(0, 1, nil, first); err != nil || decision != EventSequenceAppend { + t.Fatalf("first event = (%v, %v)", decision, err) + } + if decision, err := ClassifyEventSequence(1, 1, &first, first); err != nil || decision != EventSequenceDuplicate { + t.Fatalf("exact duplicate = (%v, %v)", decision, err) + } + if _, err := ClassifyEventSequence(1, 1, &first, other); !errors.Is(err, ErrConflictingDuplicate) { + t.Fatalf("conflicting duplicate error = %v", err) + } + if _, err := ClassifyEventSequence(1, 1, nil, first); !errors.Is(err, ErrMissingDuplicate) { + t.Fatalf("missing duplicate error = %v", err) + } + if _, err := ClassifyEventSequence(1, 3, nil, first); !errors.Is(err, ErrEventSequenceGap) { + t.Fatalf("gap error = %v", err) + } + if _, err := ClassifyEventSequence(0, 0, nil, first); !errors.Is(err, ErrInvalidEventSequence) { + t.Fatalf("zero sequence error = %v", err) + } +} + +func TestOutputTruncationRanges_HP_PROTO_01(t *testing.T) { + t.Parallel() + + one, two := uint64(1), uint64(2) + zero := uint64(0) + tests := []struct { + name string + value *rvboxv1.OutputTruncation + valid bool + }{ + {"client pre-sequence loss", &rvboxv1.OutputTruncation{}, true}, + {"server range", &rvboxv1.OutputTruncation{FirstRemovedEventSeq: &one, LastRemovedEventSeq: &two}, true}, + {"nil", nil, false}, + {"first only", &rvboxv1.OutputTruncation{FirstRemovedEventSeq: &one}, false}, + {"last only", &rvboxv1.OutputTruncation{LastRemovedEventSeq: &two}, false}, + {"zero", &rvboxv1.OutputTruncation{FirstRemovedEventSeq: &zero, LastRemovedEventSeq: &two}, false}, + {"reversed", &rvboxv1.OutputTruncation{FirstRemovedEventSeq: &two, LastRemovedEventSeq: &one}, false}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + t.Parallel() + err := ValidateOutputTruncation(test.value) + if (err == nil) != test.valid { + t.Fatalf("ValidateOutputTruncation() error = %v, valid = %t", err, test.valid) + } + }) + } +}