diff --git a/internal/domain/protocol.go b/internal/domain/protocol.go new file mode 100644 index 0000000..376f417 --- /dev/null +++ b/internal/domain/protocol.go @@ -0,0 +1,28 @@ +package domain + +import ( + "errors" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +var ErrNoCompatibleProtocol = errors.New("no compatible protocol version") + +// SelectProtocol chooses the highest mutually supported minor version. +func SelectProtocol(server, client *rvboxv1.ProtocolRange) (*rvboxv1.ProtocolVersion, error) { + if server == nil || client == nil || server.Major == 0 || client.Major == 0 || server.Major != client.Major || server.MinMinor > server.MaxMinor || client.MinMinor > client.MaxMinor { + return nil, ErrNoCompatibleProtocol + } + minimum := server.MinMinor + if client.MinMinor > minimum { + minimum = client.MinMinor + } + maximum := server.MaxMinor + if client.MaxMinor < maximum { + maximum = client.MaxMinor + } + if minimum > maximum { + return nil, ErrNoCompatibleProtocol + } + return &rvboxv1.ProtocolVersion{Major: server.Major, Minor: maximum}, nil +} diff --git a/internal/domain/protocol_test.go b/internal/domain/protocol_test.go new file mode 100644 index 0000000..180847f --- /dev/null +++ b/internal/domain/protocol_test.go @@ -0,0 +1,29 @@ +package domain + +import ( + "errors" + "testing" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +func TestSelectHighestCompatibleProtocol_HP_SES_01(t *testing.T) { + t.Parallel() + + version, err := SelectProtocol(&rvboxv1.ProtocolRange{Major: 1, MinMinor: 1, MaxMinor: 5}, &rvboxv1.ProtocolRange{Major: 1, MinMinor: 3, MaxMinor: 7}) + if err != nil || version.Major != 1 || version.Minor != 5 { + t.Fatalf("SelectProtocol() = (%+v, %v)", version, err) + } + for _, ranges := range [][2]*rvboxv1.ProtocolRange{ + {{Major: 1, MinMinor: 1, MaxMinor: 2}, {Major: 2, MinMinor: 1, MaxMinor: 2}}, + {{Major: 1, MinMinor: 3, MaxMinor: 2}, {Major: 1, MinMinor: 1, MaxMinor: 2}}, + {{Major: 1, MinMinor: 1, MaxMinor: 2}, {Major: 1, MinMinor: 3, MaxMinor: 4}}, + } { + if _, err := SelectProtocol(ranges[0], ranges[1]); !errors.Is(err, ErrNoCompatibleProtocol) { + t.Errorf("SelectProtocol(%+v, %+v) error = %v", ranges[0], ranges[1], err) + } + } + if _, err := SelectProtocol(nil, &rvboxv1.ProtocolRange{Major: 1}); !errors.Is(err, ErrNoCompatibleProtocol) { + t.Fatalf("nil protocol range error = %v", err) + } +} diff --git a/internal/domain/reconcile.go b/internal/domain/reconcile.go new file mode 100644 index 0000000..51e78a3 --- /dev/null +++ b/internal/domain/reconcile.go @@ -0,0 +1,124 @@ +package domain + +import ( + "errors" + "fmt" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +var ErrReconcileContradiction = errors.New("reconciliation evidence contradicts durable state") + +type ClientEvidenceKind uint8 + +const ( + ClientEvidenceAbsent ClientEvidenceKind = iota + ClientEvidenceRetained + ClientEvidenceTombstone +) + +type ReconcileAction uint8 + +const ( + ReconcileNoop ReconcileAction = iota + ReconcileRequeue + ReconcileResumeDelivery + ReconcileInterruptServerSuppressReplay + ReconcileInterruptServerClientStateLoss + ReconcileTerminateLocal + ReconcileDiscardLocalTerminal + ReconcileRecordContradiction +) + +type ReconcileInput struct { + ServerPresent bool + ServerLifecycle rvboxv1.CommandLifecycle + ServerRevision CommandRevision + ServerLastEventSeq uint64 + ServerHasTombstone bool + ClientEvidence ClientEvidenceKind + ClientLifecycle rvboxv1.CommandLifecycle + ClientRevision CommandRevision + ClientLastEventSeq uint64 + ImmutableHashMatches bool +} + +type ReconcileDecision struct { + Action ReconcileAction + EffectiveRevision CommandRevision + RecordIncident bool + SuppressReplay bool +} + +// DecideReconciliation implements the explicit server/client evidence matrix. +// Persistence and message emission remain the caller's atomic responsibility. +func DecideReconciliation(input ReconcileInput) (ReconcileDecision, error) { + if input.ClientEvidence > ClientEvidenceTombstone { + return ReconcileDecision{}, fmt.Errorf("%w: unknown client evidence", ErrReconcileContradiction) + } + if input.ServerPresent { + if input.ServerLifecycle == rvboxv1.CommandLifecycle_COMMAND_LIFECYCLE_UNSPECIFIED || input.ServerRevision == 0 { + return ReconcileDecision{}, fmt.Errorf("%w: invalid server row", ErrReconcileContradiction) + } + } + if input.ClientEvidence != ClientEvidenceAbsent { + if input.ClientRevision == 0 || !input.ImmutableHashMatches { + return ReconcileDecision{}, fmt.Errorf("%w: invalid revision or immutable hash", ErrReconcileContradiction) + } + } + + if !input.ServerPresent { + if input.ClientEvidence == ClientEvidenceAbsent { + return ReconcileDecision{Action: ReconcileNoop}, nil + } + if input.ClientEvidence == ClientEvidenceTombstone || IsTerminal(input.ClientLifecycle) || input.ServerHasTombstone { + return ReconcileDecision{Action: ReconcileDiscardLocalTerminal, EffectiveRevision: input.ClientRevision, SuppressReplay: true}, nil + } + return ReconcileDecision{Action: ReconcileTerminateLocal, EffectiveRevision: input.ClientRevision, RecordIncident: true, SuppressReplay: true}, nil + } + + if IsTerminal(input.ServerLifecycle) { + if input.ClientEvidence == ClientEvidenceAbsent { + return ReconcileDecision{Action: ReconcileNoop, EffectiveRevision: input.ServerRevision}, nil + } + if input.ClientEvidence == ClientEvidenceTombstone || IsTerminal(input.ClientLifecycle) { + if input.ClientEvidence == ClientEvidenceRetained && input.ClientLifecycle != input.ServerLifecycle { + return ReconcileDecision{Action: ReconcileRecordContradiction, EffectiveRevision: maxRevision(input.ServerRevision, input.ClientRevision), RecordIncident: true, SuppressReplay: true}, fmt.Errorf("%w: terminal lifecycle mismatch", ErrReconcileContradiction) + } + return ReconcileDecision{Action: ReconcileDiscardLocalTerminal, EffectiveRevision: maxRevision(input.ServerRevision, input.ClientRevision), SuppressReplay: true}, nil + } + return ReconcileDecision{Action: ReconcileTerminateLocal, EffectiveRevision: maxRevision(input.ServerRevision, input.ClientRevision), RecordIncident: true, SuppressReplay: true}, nil + } + + if input.ClientEvidence == ClientEvidenceTombstone { + return ReconcileDecision{Action: ReconcileInterruptServerSuppressReplay, EffectiveRevision: maxRevision(input.ServerRevision, input.ClientRevision), RecordIncident: true, SuppressReplay: true}, nil + } + if input.ClientEvidence == ClientEvidenceAbsent { + switch input.ServerLifecycle { + case rvboxv1.CommandLifecycle_COMMAND_QUEUED, rvboxv1.CommandLifecycle_COMMAND_DISPATCHED: + return ReconcileDecision{Action: ReconcileRequeue, EffectiveRevision: input.ServerRevision}, nil + case rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_RUNNING: + return ReconcileDecision{Action: ReconcileInterruptServerClientStateLoss, EffectiveRevision: input.ServerRevision, RecordIncident: true, SuppressReplay: true}, nil + default: + return ReconcileDecision{}, fmt.Errorf("%w: unsupported non-terminal server state", ErrReconcileContradiction) + } + } + + if input.ClientLifecycle == rvboxv1.CommandLifecycle_COMMAND_LIFECYCLE_UNSPECIFIED { + return ReconcileDecision{}, fmt.Errorf("%w: invalid client lifecycle", ErrReconcileContradiction) + } + if input.ClientLastEventSeq < input.ServerLastEventSeq { + return ReconcileDecision{Action: ReconcileRecordContradiction, EffectiveRevision: maxRevision(input.ServerRevision, input.ClientRevision), RecordIncident: true, SuppressReplay: true}, fmt.Errorf("%w: client lost server-confirmed events", ErrReconcileContradiction) + } + if input.ServerLifecycle == rvboxv1.CommandLifecycle_COMMAND_QUEUED { + return ReconcileDecision{Action: ReconcileRecordContradiction, EffectiveRevision: maxRevision(input.ServerRevision, input.ClientRevision), RecordIncident: true, SuppressReplay: true}, fmt.Errorf("%w: queued server command exists at client", ErrReconcileContradiction) + } + return ReconcileDecision{Action: ReconcileResumeDelivery, EffectiveRevision: maxRevision(input.ServerRevision, input.ClientRevision)}, nil +} + +func maxRevision(left, right CommandRevision) CommandRevision { + if left > right { + return left + } + return right +} diff --git a/internal/domain/reconcile_test.go b/internal/domain/reconcile_test.go new file mode 100644 index 0000000..af96605 --- /dev/null +++ b/internal/domain/reconcile_test.go @@ -0,0 +1,94 @@ +package domain + +import ( + "errors" + "testing" + + rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1" +) + +func TestReconciliationMatrix_HP_SES_11(t *testing.T) { + t.Parallel() + + tests := []struct { + name string + input ReconcileInput + want ReconcileAction + incident bool + }{ + {"queued absent", serverInput(rvboxv1.CommandLifecycle_COMMAND_QUEUED, ClientEvidenceAbsent), ReconcileRequeue, false}, + {"dispatched absent", serverInput(rvboxv1.CommandLifecycle_COMMAND_DISPATCHED, ClientEvidenceAbsent), ReconcileRequeue, false}, + {"accepted absent", serverInput(rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, ClientEvidenceAbsent), ReconcileInterruptServerClientStateLoss, true}, + {"running absent", serverInput(rvboxv1.CommandLifecycle_COMMAND_RUNNING, ClientEvidenceAbsent), ReconcileInterruptServerClientStateLoss, true}, + {"dispatched retained", retainedInput(rvboxv1.CommandLifecycle_COMMAND_DISPATCHED, rvboxv1.CommandLifecycle_COMMAND_ACCEPTED), ReconcileResumeDelivery, false}, + {"accepted retained", retainedInput(rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_RUNNING), ReconcileResumeDelivery, false}, + {"running retained", retainedInput(rvboxv1.CommandLifecycle_COMMAND_RUNNING, rvboxv1.CommandLifecycle_COMMAND_RUNNING), ReconcileResumeDelivery, false}, + {"nonterminal tombstone", tombstoneInput(rvboxv1.CommandLifecycle_COMMAND_DISPATCHED), ReconcileInterruptServerSuppressReplay, true}, + {"terminal retained", retainedInput(rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED, rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED), ReconcileDiscardLocalTerminal, false}, + {"terminal active", retainedInput(rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED, rvboxv1.CommandLifecycle_COMMAND_RUNNING), ReconcileTerminateLocal, true}, + {"missing active", ReconcileInput{ClientEvidence: ClientEvidenceRetained, ClientLifecycle: rvboxv1.CommandLifecycle_COMMAND_RUNNING, ClientRevision: 1, ImmutableHashMatches: true}, ReconcileTerminateLocal, true}, + {"missing terminal", ReconcileInput{ClientEvidence: ClientEvidenceRetained, ClientLifecycle: rvboxv1.CommandLifecycle_COMMAND_FAILED, ClientRevision: 1, ImmutableHashMatches: true}, ReconcileDiscardLocalTerminal, false}, + } + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + t.Parallel() + decision, err := DecideReconciliation(test.input) + if err != nil { + t.Fatalf("DecideReconciliation(): %v", err) + } + if decision.Action != test.want || decision.RecordIncident != test.incident { + t.Fatalf("decision = %+v, want action=%v incident=%t", decision, test.want, test.incident) + } + }) + } +} + +func TestReconciliationContradictions_BH_SES_10(t *testing.T) { + t.Parallel() + + tests := []ReconcileInput{ + func() ReconcileInput { + value := retainedInput(rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_RUNNING) + value.ImmutableHashMatches = false + return value + }(), + func() ReconcileInput { + value := retainedInput(rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_RUNNING) + value.ClientLastEventSeq = 0 + value.ServerLastEventSeq = 1 + return value + }(), + retainedInput(rvboxv1.CommandLifecycle_COMMAND_QUEUED, rvboxv1.CommandLifecycle_COMMAND_ACCEPTED), + retainedInput(rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED, rvboxv1.CommandLifecycle_COMMAND_FAILED), + } + for _, input := range tests { + decision, err := DecideReconciliation(input) + if !errors.Is(err, ErrReconcileContradiction) { + t.Errorf("error = %v, decision=%+v", err, decision) + } + } +} + +func TestReconciliationKeepsHighestRevision_HP_SES_11(t *testing.T) { + t.Parallel() + + input := retainedInput(rvboxv1.CommandLifecycle_COMMAND_ACCEPTED, rvboxv1.CommandLifecycle_COMMAND_RUNNING) + input.ServerRevision = 4 + input.ClientRevision = 7 + decision, err := DecideReconciliation(input) + if err != nil || decision.EffectiveRevision != 7 { + t.Fatalf("decision = %+v, error = %v", decision, err) + } +} + +func serverInput(state rvboxv1.CommandLifecycle, evidence ClientEvidenceKind) ReconcileInput { + return ReconcileInput{ServerPresent: true, ServerLifecycle: state, ServerRevision: 1, ClientEvidence: evidence} +} + +func retainedInput(server, client rvboxv1.CommandLifecycle) ReconcileInput { + return ReconcileInput{ServerPresent: true, ServerLifecycle: server, ServerRevision: 1, ServerLastEventSeq: 1, ClientEvidence: ClientEvidenceRetained, ClientLifecycle: client, ClientRevision: 1, ClientLastEventSeq: 1, ImmutableHashMatches: true} +} + +func tombstoneInput(server rvboxv1.CommandLifecycle) ReconcileInput { + return ReconcileInput{ServerPresent: true, ServerLifecycle: server, ServerRevision: 1, ClientEvidence: ClientEvidenceTombstone, ClientRevision: 1, ImmutableHashMatches: true} +} diff --git a/test/coverage.toml b/test/coverage.toml index 54b5fe7..1612b8f 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -94,6 +94,27 @@ tests = [ "internal/agentproto/validate_test.go:TestReconciliationRowsAndCapacity_BH_SES_01", ] +[[requirements]] +id = "HP-SES-01" +layer = "unit" +status = "implemented" +tests = ["internal/domain/protocol_test.go:TestSelectHighestCompatibleProtocol_HP_SES_01"] + +[[requirements]] +id = "HP-SES-11" +layer = "unit" +status = "implemented" +tests = [ + "internal/domain/reconcile_test.go:TestReconciliationMatrix_HP_SES_11", + "internal/domain/reconcile_test.go:TestReconciliationKeepsHighestRevision_HP_SES_11", +] + +[[requirements]] +id = "BH-SES-10" +layer = "unit" +status = "implemented" +tests = ["internal/domain/reconcile_test.go:TestReconciliationContradictions_BH_SES_10"] + [[requirements]] id = "BH-IDEM-01" layer = "unit"