Files
rvbox/internal/client/spool/snapshot.go
T

74 lines
2.7 KiB
Go

package spool
import (
"context"
"fmt"
rvboxv1 "github.com/rvbox/rvbox/gen/go/rvbox/v1"
)
// ReconcileSnapshot is a consistent durable view. The caller sends it only
// after welcome and before admitting new dispatches on that session.
func (store *Store) ReconcileSnapshot(ctx context.Context) (*rvboxv1.ReconcileSnapshot, error) {
result := &rvboxv1.ReconcileSnapshot{}
rows, err := store.db.QueryContext(ctx, `SELECT issue_uuid, command_revision, phase, next_event_seq, immutable_sha256 FROM commands ORDER BY issue_uuid`)
if err != nil {
return nil, err
}
for rows.Next() {
command, err := scanReconcileCommand(rows, false)
if err != nil {
_ = rows.Close()
return nil, err
}
result.RetainedCommands = append(result.RetainedCommands, command)
}
if err := rows.Close(); err != nil {
return nil, err
}
tombstones, err := store.db.QueryContext(ctx, `SELECT issue_uuid, command_revision, terminal_lifecycle, immutable_sha256 FROM command_tombstones ORDER BY acknowledged_at, issue_uuid`)
if err != nil {
return nil, err
}
defer tombstones.Close()
for tombstones.Next() {
var owner, digest []byte
var revision uint64
var lifecycle uint32
if err := tombstones.Scan(&owner, &revision, &lifecycle, &digest); err != nil {
return nil, err
}
if _, valid := copiedUUID(owner); !valid || len(digest) != 32 || revision == 0 || !isTerminalPhase(lifecycle) {
return nil, ErrScriptState
}
result.RetainedCommands = append(result.RetainedCommands, &rvboxv1.ReconcileCommandState{IssueUuid: uuidText(owner), CommandRevision: revision, Lifecycle: rvboxv1.CommandLifecycle(lifecycle), Tombstoned: true, ImmutableRequestSha256: append([]byte(nil), digest...)})
}
if err := tombstones.Err(); err != nil {
return nil, err
}
return result, nil
}
type reconcileRow interface {
Scan(...any) error
}
func scanReconcileCommand(row reconcileRow, tombstoned bool) (*rvboxv1.ReconcileCommandState, error) {
var owner, digest []byte
var revision uint64
var lifecycle uint32
var nextSequence uint64
if err := row.Scan(&owner, &revision, &lifecycle, &nextSequence, &digest); err != nil {
return nil, err
}
if _, valid := copiedUUID(owner); !valid || len(digest) != 32 || revision == 0 || lifecycle == 0 || lifecycle > 11 || nextSequence == 0 {
return nil, fmt.Errorf("%w: invalid reconciliation row", ErrScriptState)
}
return &rvboxv1.ReconcileCommandState{IssueUuid: uuidText(owner), CommandRevision: revision, Lifecycle: rvboxv1.CommandLifecycle(lifecycle), LastClientEventSeq: nextSequence - 1, Tombstoned: tombstoned, ImmutableRequestSha256: append([]byte(nil), digest...)}, nil
}
func uuidText(value []byte) string {
parsed, _ := copiedUUID(value)
return parsed.String()
}