feat: add command domain invariants
This commit is contained in:
@@ -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. |
|
||||
|
||||
@@ -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
|
||||
)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user