feat: reconcile retained client commands
This commit is contained in:
@@ -185,15 +185,14 @@ func (server *AgentServer) serveConnection(parent context.Context, connection *w
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
if snapshot := envelope.GetReconcileSnapshot(); snapshot != nil {
|
if snapshot := envelope.GetReconcileSnapshot(); snapshot != nil {
|
||||||
// A retained command requires the command-store reconciliation phase,
|
result, reconcileErr := server.Store.ReconcileClientSnapshot(sessionContext, hello.GetClientId(), snapshot)
|
||||||
// which is deliberately not substituted with a no-op result.
|
if reconcileErr != nil {
|
||||||
if len(snapshot.GetRetainedCommands()) != 0 {
|
server.close(connection, websocket.StatusPolicyViolation, "reconciliation failed")
|
||||||
server.close(connection, websocket.StatusPolicyViolation, "reconciliation for retained commands is unavailable")
|
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
encoded, err := proto.Marshal(&rvboxv1.AgentEnvelope{
|
encoded, err := proto.Marshal(&rvboxv1.AgentEnvelope{
|
||||||
SessionId: encodedSessionID, SessionGeneration: registration.Generation,
|
SessionId: encodedSessionID, SessionGeneration: registration.Generation,
|
||||||
Payload: &rvboxv1.AgentEnvelope_ReconcileResult{ReconcileResult: &rvboxv1.ReconcileResult{}},
|
Payload: &rvboxv1.AgentEnvelope_ReconcileResult{ReconcileResult: result},
|
||||||
})
|
})
|
||||||
if err != nil || queue.EnqueueControl(Frame{Kind: FrameControl, Payload: encoded}) != nil {
|
if err != nil || queue.EnqueueControl(Frame{Kind: FrameControl, Payload: encoded}) != nil {
|
||||||
server.close(connection, websocket.StatusInternalError, "could not queue reconciliation result")
|
server.close(connection, websocket.StatusInternalError, "could not queue reconciliation result")
|
||||||
|
|||||||
@@ -0,0 +1,63 @@
|
|||||||
|
package store
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"database/sql"
|
||||||
|
"errors"
|
||||||
|
|
||||||
|
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
|
||||||
|
"github.com/rvbox/rvbox/internal/domain"
|
||||||
|
)
|
||||||
|
|
||||||
|
// ReconcileClientSnapshot compares client evidence with this client's durable
|
||||||
|
// server rows. It intentionally performs no lifecycle mutation yet: callers
|
||||||
|
// receive only the safe local terminate/discard instructions.
|
||||||
|
func (store *Store) ReconcileClientSnapshot(ctx context.Context, clientID string, snapshot *rvboxv1.ReconcileSnapshot) (*rvboxv1.ReconcileResult, error) {
|
||||||
|
if snapshot == nil {
|
||||||
|
return nil, errors.New("missing client reconciliation snapshot")
|
||||||
|
}
|
||||||
|
result := &rvboxv1.ReconcileResult{}
|
||||||
|
for _, client := range snapshot.GetRetainedCommands() {
|
||||||
|
issue, err := domain.ParseUUIDv7(client.GetIssueUuid())
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
input := domain.ReconcileInput{ClientEvidence: domain.ClientEvidenceRetained, ClientLifecycle: client.GetLifecycle(), ClientRevision: domain.CommandRevision(client.GetCommandRevision()), ClientLastEventSeq: client.GetLastClientEventSeq(), ImmutableHashMatches: true}
|
||||||
|
if client.GetTombstoned() {
|
||||||
|
input.ClientEvidence = domain.ClientEvidenceTombstone
|
||||||
|
}
|
||||||
|
var lifecycle uint32
|
||||||
|
var revision, lastSequence uint64
|
||||||
|
var hash []byte
|
||||||
|
err = store.db.QueryRowContext(ctx, `SELECT lifecycle, revision, last_event_seq, immutable_request_sha256 FROM commands WHERE issue_uuid = ? AND client_id = ?`, issue[:], clientID).Scan(&lifecycle, &revision, &lastSequence, &hash)
|
||||||
|
if err == nil {
|
||||||
|
input.ServerPresent = true
|
||||||
|
input.ServerLifecycle = rvboxv1.CommandLifecycle(lifecycle)
|
||||||
|
input.ServerRevision = domain.CommandRevision(revision)
|
||||||
|
input.ServerLastEventSeq = lastSequence
|
||||||
|
input.ImmutableHashMatches = len(hash) == len(client.GetImmutableRequestSha256()) && string(hash) == string(client.GetImmutableRequestSha256())
|
||||||
|
} else if !errors.Is(err, sql.ErrNoRows) {
|
||||||
|
return nil, err
|
||||||
|
} else {
|
||||||
|
var tombstoneHash []byte
|
||||||
|
err = store.db.QueryRowContext(ctx, `SELECT immutable_sha256 FROM command_tombstones WHERE issue_uuid = ? AND client_id = ?`, issue[:], clientID).Scan(&tombstoneHash)
|
||||||
|
if err == nil {
|
||||||
|
input.ServerHasTombstone = true
|
||||||
|
input.ImmutableHashMatches = len(tombstoneHash) == len(client.GetImmutableRequestSha256()) && string(tombstoneHash) == string(client.GetImmutableRequestSha256())
|
||||||
|
} else if !errors.Is(err, sql.ErrNoRows) {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
decision, err := domain.DecideReconciliation(input)
|
||||||
|
if err != nil && !errors.Is(err, domain.ErrReconcileContradiction) {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
switch decision.Action {
|
||||||
|
case domain.ReconcileTerminateLocal:
|
||||||
|
result.TerminateLocalIssueUuids = append(result.TerminateLocalIssueUuids, issue.String())
|
||||||
|
case domain.ReconcileDiscardLocalTerminal:
|
||||||
|
result.DiscardLocalTerminalIssueUuids = append(result.DiscardLocalTerminalIssueUuids, issue.String())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return result, nil
|
||||||
|
}
|
||||||
@@ -39,7 +39,19 @@ func TestWebSocketHelloWelcome_HP_SES_06(t *testing.T) {
|
|||||||
if result.ID == "" || result.Generation != 1 || result.Protocol.GetMajor() != 1 {
|
if result.ID == "" || result.Generation != 1 || result.Protocol.GetMajor() != 1 {
|
||||||
t.Fatalf("handshake result = %#v", result)
|
t.Fatalf("handshake result = %#v", result)
|
||||||
}
|
}
|
||||||
if reconciled, err := agent.Reconcile(ctx, transport, result, &rvboxv1.ReconcileSnapshot{}, agentproto.DefaultLimits()); err != nil || reconciled == nil {
|
activeIssue := "019c46f1-1d02-7000-8000-000000000072"
|
||||||
t.Fatalf("empty reconciliation = %#v, %v", reconciled, err)
|
terminalIssue := "019c46f1-1d02-7000-8000-000000000073"
|
||||||
|
reconciled, err := agent.Reconcile(ctx, transport, result, &rvboxv1.ReconcileSnapshot{RetainedCommands: []*rvboxv1.ReconcileCommandState{
|
||||||
|
{IssueUuid: activeIssue, Lifecycle: rvboxv1.CommandLifecycle_COMMAND_RUNNING, CommandRevision: 1, ImmutableRequestSha256: make([]byte, 32)},
|
||||||
|
{IssueUuid: terminalIssue, Lifecycle: rvboxv1.CommandLifecycle_COMMAND_SUCCEEDED, CommandRevision: 1, Tombstoned: true, ImmutableRequestSha256: make([]byte, 32)},
|
||||||
|
}}, agentproto.DefaultLimits())
|
||||||
|
if err != nil {
|
||||||
|
t.Fatal(err)
|
||||||
|
}
|
||||||
|
if got := reconciled.GetTerminateLocalIssueUuids(); len(got) != 1 || got[0] != activeIssue {
|
||||||
|
t.Fatalf("terminate local issues = %v", got)
|
||||||
|
}
|
||||||
|
if got := reconciled.GetDiscardLocalTerminalIssueUuids(); len(got) != 1 || got[0] != terminalIssue {
|
||||||
|
t.Fatalf("discard local terminal issues = %v", got)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user