feat: add protocol and reconciliation decisions
This commit is contained in:
@@ -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
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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}
|
||||
}
|
||||
Reference in New Issue
Block a user