diff --git a/db/migrations/00015_member_outcome_corrections.sql b/db/migrations/00015_member_outcome_corrections.sql new file mode 100644 index 00000000..9744b06f --- /dev/null +++ b/db/migrations/00015_member_outcome_corrections.sql @@ -0,0 +1,20 @@ +-- Corrections to intent outcomes are recorded (RB4-F58). +-- +-- A chain member's recorded outcome is replaced by a later report for the same chain, and the intent's status is +-- derived again from every member (RecordMemberOutcome). The replacement overwrote the member row and, when the +-- intent went from failed to complete, erased failed_at, error_message and failure_class: a published outcome +-- rewritten with nothing left to say it was. On 2026-09-29 intent 000ac79a's base member was recorded failed by a +-- proof cycle a fleet restart broke (RB4-F55) although its settlement landed; repairing it replaces that record. +-- +-- Every change to a recorded member outcome, and every change to an intent's terminal outcome, is now written to +-- evidence_corrections in the same transaction under these two record types: +-- intent_member_outcome record_id /; previous and corrected member rows +-- intent_lifecycle record_id ; previous and corrected terminal status, times, message, class +-- +-- schema: destructive-approved (the evidence_corrections record_type CHECK is replaced by a wider one) + +ALTER TABLE public.evidence_corrections DROP CONSTRAINT evidence_correction_record_type; +ALTER TABLE public.evidence_corrections ADD CONSTRAINT evidence_correction_record_type + CHECK (record_type IN ('anchor_batch', 'layer5', 'certen_anchor_proof', 'proof_artifact', 'anchor_reference', + 'validator_attestation', 'consensus_entry', 'batch_attestation', + 'intent_member_outcome', 'intent_lifecycle')); diff --git a/db/migrations/00016_member_write_backs.sql b/db/migrations/00016_member_write_backs.sql new file mode 100644 index 00000000..26a74092 --- /dev/null +++ b/db/migrations/00016_member_write_backs.sql @@ -0,0 +1,31 @@ +-- A chain member's outcome is written back to Accumulate once (RB4-F59). +-- +-- Phase 9 submitted a write-back for whatever cycle reached it, and the proof bundle was stored before it, so a +-- second proof cycle for a member already written back - a re-driven or re-discovered one - wrote a second +-- Accumulate entry and stored a second proof artifact for the same member. This register is claimed before a +-- write-back is submitted and states its outcome: +-- claimed a cycle is submitting it, or submitted it and its outcome is not known (the submission may have +-- reached Accumulate): no other cycle writes this member back until that is established +-- written written back as write_back_tx: no other cycle writes it back +-- not_sent the claiming cycle failed before anything was sent: another cycle may claim it +-- +-- Expand-only: one table. + +CREATE TABLE public.member_write_backs ( + intent_id character varying(256) NOT NULL, + chain_id bigint NOT NULL, + state character varying(16) NOT NULL, + cycle_id character varying(256) NOT NULL, + validator_id character varying(128) NOT NULL, + write_back_tx character varying(256), + reason text, + claimed_at timestamp with time zone DEFAULT now() NOT NULL, + resolved_at timestamp with time zone, + CONSTRAINT member_write_backs_pkey PRIMARY KEY (intent_id, chain_id), + CONSTRAINT member_write_back_state CHECK (state IN ('claimed', 'written', 'not_sent')), + CONSTRAINT member_write_back_written_has_tx CHECK ((state = 'written') = (write_back_tx IS NOT NULL)), + CONSTRAINT member_write_back_not_sent_says_why CHECK (state <> 'not_sent' OR reason IS NOT NULL) +); + +COMMENT ON TABLE public.member_write_backs IS + 'One row per chain member written back to Accumulate (RB4-F59): claimed before submission, written with its transaction, or not_sent when the claiming cycle failed before sending. A claimed row whose outcome is unknown blocks another write-back of the member.'; diff --git a/db/schema.fingerprint b/db/schema.fingerprint index 7e9cf252..8f37da6d 100644 --- a/db/schema.fingerprint +++ b/db/schema.fingerprint @@ -1 +1 @@ -b9a4c39c2196b0ad4816a0842c22f7c487f932330771820f835d4650ec9c63c2 +0d6659c87b36cf1dd8a8e1b5c1dd705c2add9ebbd6029a3f865e20b38a4829fb diff --git a/main.go b/main.go index 3a3a13ac..f50fcec0 100644 --- a/main.go +++ b/main.go @@ -1837,7 +1837,7 @@ func startValidator( return nil, nil, fmt.Errorf("proof cycle: member outcome outbox: %w", moErr) } (&execution.MemberOutcomeReconciler{ - Outbox: memberOutcomes, Store: batchComponents.Repos.IntentLifecycle, Logf: log.Printf, + Outbox: memberOutcomes, Store: batchComponents.Repos.IntentLifecycle, ValidatorID: cfg.ValidatorID, Logf: log.Printf, }).Start(context.Background()) log.Printf("✅ [Phase 9] Member outcome outbox at %s; reconciler replaying on startup and every minute", memberOutcomes.Dir()) @@ -1962,6 +1962,19 @@ func startValidator( // with properly structured CertenIntent (4-blob canonical) and CertenProof from lite client intentDiscovery.SetBFTConsensus(validator) + // RB4-F55 repair: one decided member's proof cycle is re-driven on request, here, where the orchestrator, its + // keys, its peers and the committed-operation index are (`validator repair member-proof-cycle`). + memberRepairs := &execution.MemberRepairRunner{ + Dir: execution.MemberRepairDir(nsDataDir), ValidatorID: cfg.ValidatorID, DB: dbClient.DB(), + Lifecycle: batchComponents.Repos.IntentLifecycle, Outbox: memberOutcomes, + Observe: unifiedOrchestrator.ObserveSettlement, Arm: validator.ArmMemberRepair, + Reprocess: intentDiscovery.ReprocessIntent, Logf: log.Printf, + } + if err := memberRepairs.Start(context.Background()); err != nil { + return nil, nil, fmt.Errorf("member repair runner: %w", err) + } + log.Printf("✅ [MEMBER-REPAIR] repair requests served from %s", memberRepairs.Dir) + // ENTITLEMENT — wire the epoch snapshot to the two places that consume it. // // The gate inside the ABCI validator only VERIFIES evidence; something has diff --git a/member_repair_command.go b/member_repair_command.go new file mode 100644 index 00000000..af882734 --- /dev/null +++ b/member_repair_command.go @@ -0,0 +1,142 @@ +package main + +import ( + "crypto/rand" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "log" + "os" + "path/filepath" + "strconv" + "strings" + "time" + + "github.com/certen/independant-validator/pkg/config" + "github.com/certen/independant-validator/pkg/execution" +) + +const memberRepairUsage = "usage: certen-validator repair member-proof-cycle --intent ID --chain CHAIN_ID --tx SETTLEMENT_TX [--apply] [--wait DURATION]" + +// runMemberRepairCommand runs `validator repair member-proof-cycle` (RB4-F55 repair; DESIGN_RB4_F55_repair_000ac79a.md): +// it hands one member's re-drive to the RUNNING validator of this container - the process holding the orchestrator, +// its keys, its peers and the committed-operation index - through its data directory, waits for the answer and +// prints it. Without --apply the running validator checks every precondition, re-derives the member's round and +// reports its snapshot, and runs nothing. Run it on the validator that settled the member. +// +// Exit status: 0 when the member was reached (dry run) or repaired, 2 when the request was refused or failed, 1 on +// error (including no answer within --wait). +func runMemberRepairCommand(args []string) int { + req := execution.MemberRepairRequest{} + wait := 20 * time.Minute + for i := 0; i < len(args); i++ { + next := func() (string, bool) { + if i+1 >= len(args) || strings.TrimSpace(args[i+1]) == "" { + return "", false + } + i++ + return strings.TrimSpace(args[i]), true + } + var ok bool + switch args[i] { + case "--apply": + req.Apply, ok = true, true + case "--intent": + req.IntentID, ok = next() + case "--tx": + req.SettlementTx, ok = next() + case "--chain": + var v string + if v, ok = next(); ok { + n, err := strconv.ParseInt(v, 10, 64) + ok = err == nil && n > 0 + req.ChainID = n + } + case "--wait": + var v string + if v, ok = next(); ok { + d, err := time.ParseDuration(v) + ok = err == nil && d > 0 + wait = d + } + } + if !ok { + log.Print(memberRepairUsage) + return 1 + } + } + if req.IntentID == "" || req.ChainID == 0 || req.SettlementTx == "" { + log.Print(memberRepairUsage) + return 1 + } + cfg, err := config.Load() + if err != nil { + log.Printf("load configuration: %v", err) + return 1 + } + dataDir := cfg.DataDir + if dataDir == "" { + dataDir = "data" + } + res, err := requestMemberRepair(execution.MemberRepairDir(dataDir), req, wait, 2*time.Second) + if err != nil { + log.Printf("member repair: %v", err) + return 1 + } + out, _ := json.MarshalIndent(res, "", " ") + fmt.Println(string(out)) + switch res.Outcome { + case execution.MemberRepairReached, execution.MemberRepairRepaired: + return 0 + default: + return 2 + } +} + +// requestMemberRepair writes the request where the running validator serves it and waits for its result. +func requestMemberRepair(dir string, req execution.MemberRepairRequest, wait, poll time.Duration) (*execution.MemberRepairResult, error) { + requests := filepath.Join(dir, "requests") + if st, err := os.Stat(requests); err != nil || !st.IsDir() { + return nil, fmt.Errorf("%s does not exist: the running validator creates it when it serves repairs - is this the container of a validator running this binary?", requests) + } + nonce := make([]byte, 4) + if _, err := rand.Read(nonce); err != nil { + return nil, err + } + req.RequestedAt = time.Now().UTC() + req.ID = fmt.Sprintf("%s-%d-%s", req.RequestedAt.Format("20060102T150405Z"), req.ChainID, hex.EncodeToString(nonce)) + blob, err := json.MarshalIndent(req, "", " ") + if err != nil { + return nil, err + } + tmp := filepath.Join(dir, "."+req.ID+".json.tmp") + if err := os.WriteFile(tmp, blob, 0o600); err != nil { + return nil, err + } + if err := os.Rename(tmp, filepath.Join(requests, req.ID+".json")); err != nil { + return nil, err + } + log.Printf("repair request %s written (intent %s member %d, settlement %s, apply %v); waiting up to %s for the running validator", + req.ID, req.IntentID, req.ChainID, req.SettlementTx, req.Apply, wait) + + result := filepath.Join(dir, "results", req.ID+".json") + deadline := time.Now().Add(wait) + for { + raw, err := os.ReadFile(result) + if err == nil { + var res execution.MemberRepairResult + if err := json.Unmarshal(raw, &res); err != nil { + return nil, fmt.Errorf("result %s cannot be read: %w", result, err) + } + return &res, nil + } + if !errors.Is(err, os.ErrNotExist) { + return nil, err + } + if time.Now().After(deadline) { + return nil, fmt.Errorf("no answer to request %s within %s; it stays in %s and its result will be written to %s", req.ID, wait, requests, result) + } + time.Sleep(poll) + } +} diff --git a/member_repair_command_test.go b/member_repair_command_test.go new file mode 100644 index 00000000..80820e02 --- /dev/null +++ b/member_repair_command_test.go @@ -0,0 +1,54 @@ +package main + +import ( + "encoding/json" + "os" + "path/filepath" + "strings" + "testing" + "time" + + "github.com/certen/independant-validator/pkg/execution" +) + +// RB4-F55 repair: the command hands the request to the running validator through its data directory and prints its +// answer; with no running validator serving repairs it says so rather than waiting on nothing. +func TestTheRepairCommandHandsTheRequestToTheRunningValidator(t *testing.T) { + dir := filepath.Join(t.TempDir(), "member_repairs") + req := execution.MemberRepairRequest{IntentID: "000ac79a", ChainID: 84532, SettlementTx: "0xc409", Apply: true} + + if _, err := requestMemberRepair(dir, req, time.Second, 10*time.Millisecond); err == nil || !strings.Contains(err.Error(), "running validator") { + t.Fatalf("no validator serving repairs: %v", err) + } + + for _, d := range []string{"requests", "results"} { + if err := os.MkdirAll(filepath.Join(dir, d), 0o700); err != nil { + t.Fatal(err) + } + } + // The running validator: answer the request it finds. + go func() { + for i := 0; i < 200; i++ { + entries, _ := os.ReadDir(filepath.Join(dir, "requests")) + for _, e := range entries { + raw, _ := os.ReadFile(filepath.Join(dir, "requests", e.Name())) + var got execution.MemberRepairRequest + if json.Unmarshal(raw, &got) != nil { + continue + } + blob, _ := json.Marshal(execution.MemberRepairResult{Request: got, Outcome: execution.MemberRepairRepaired}) + os.WriteFile(filepath.Join(dir, "results", got.ID+".json"), blob, 0o600) + return + } + time.Sleep(10 * time.Millisecond) + } + }() + res, err := requestMemberRepair(dir, req, 5*time.Second, 10*time.Millisecond) + if err != nil { + t.Fatal(err) + } + if res.Outcome != execution.MemberRepairRepaired || res.Request.IntentID != "000ac79a" || res.Request.ChainID != 84532 || + res.Request.SettlementTx != "0xc409" || !res.Request.Apply || res.Request.ID == "" || res.Request.RequestedAt.IsZero() { + t.Fatalf("the request reached the validator as %+v", res.Request) + } +} diff --git a/pkg/consensus/batch_quorum_prover.go b/pkg/consensus/batch_quorum_prover.go index ad528b36..ab963986 100644 --- a/pkg/consensus/batch_quorum_prover.go +++ b/pkg/consensus/batch_quorum_prover.go @@ -401,6 +401,12 @@ func (bv *BFTValidator) enqueueForBatch( // re-driven intent - discovery rewinding after a restart - arrives here (RB3-F141). bv.logger.Printf("📦 [BATCH-QUEUE] intent %s on chain %d is already decided — not queued again: %v", certenIntent.IntentID, m.chainID, enqErr) + // Unless a repair names it: its proof cycle is then run from this round's snapshot (member_repair.go). + memberLane := "on_cadence" + if plan.onDemand { + memberLane = "on_demand" + } + bv.repairDecidedMember(&memberAtt, m.chainID, memberLane) default: // All-or-nothing: take back what this call queued for the intent's other chains, or the // intent would settle on one chain while being reported refused. diff --git a/pkg/consensus/bft_integration.go b/pkg/consensus/bft_integration.go index ebd175e8..a5e57639 100644 --- a/pkg/consensus/bft_integration.go +++ b/pkg/consensus/bft_integration.go @@ -376,6 +376,9 @@ type BatchEnqueuer interface { // BFTValidator represents a decentralized BFT validator with elected executor consensus // Phase 3: BFTValidator now uses only CometBFT for consensus (no ExecutionConsensus) type BFTValidator struct { + // memberRepairs: members named for a re-driven proof cycle (RB4-F55 repair, member_repair.go). + memberRepairState + engine BFTConsensusEngine anchorManager AnchorManager proofGenerator ProofGenerator diff --git a/pkg/consensus/member_repair.go b/pkg/consensus/member_repair.go new file mode 100644 index 00000000..d5408949 --- /dev/null +++ b/pkg/consensus/member_repair.go @@ -0,0 +1,153 @@ +package consensus + +import ( + "context" + "fmt" + "strings" + "sync" +) + +// Re-driving one decided member's proof cycle (RB4-F55 repair). +// +// A member with a recorded outcome is never queued again (RB3-F141), so nothing re-runs a proof cycle that failed +// for reasons that were not the member's: on 2026-09-29 intent 000ac79a's base member settled on chain and was +// recorded failed by a cycle a fleet restart broke. A repair names the member; its intent is re-processed exactly +// as discovery processes it - its round re-derived and checked against the committed block - and at that member, +// where the round captures its snapshot and finds the member decided, the member's proof cycle is run on its +// settlement from that snapshot. Only a member a repair names, only once, and never from a snapshot that lacks +// what the proof cycle binds. + +// MemberRepair names one member to re-drive. +type MemberRepair struct { + IntentID string + ChainID int64 + SettlementTx string + // Apply runs the proof cycle; without it the member's re-derived snapshot is reported and nothing runs. + Apply bool +} + +// MemberRepairSnapshot is what the re-derived round captured for the member. +type MemberRepairSnapshot struct { + IntentID string + Lane string + BundleID string + GovernanceRoot string + OperationCommitment string + AccumulateHeight uint64 + GovernanceLevel string + GovernanceLevels bool // G0, G1 and G2 are all present + BLSSignature bool + AccountURL string + TransactionHash string +} + +// MemberRepairReach is what happened when the re-derived round reached the named member. +type MemberRepairReach struct { + Snapshot MemberRepairSnapshot + // Started: the member's proof cycle was started (Apply). Its outcome is recorded by the cycle. + Started bool + Err error +} + +type armedRepair struct { + MemberRepair + reach chan MemberRepairReach +} + +func memberRepairKey(intentID string, chainID int64) string { + return fmt.Sprintf("%s|%d", intentID, chainID) +} + +// ArmMemberRepair names a member for re-driving the next time its round reaches it. The channel receives what +// happened then, once; disarm withdraws it if the round never reaches the member. +func (bv *BFTValidator) ArmMemberRepair(r MemberRepair) (<-chan MemberRepairReach, func()) { + a := &armedRepair{MemberRepair: r, reach: make(chan MemberRepairReach, 1)} + key := memberRepairKey(r.IntentID, r.ChainID) + bv.memberRepairsMu.Lock() + if bv.memberRepairs == nil { + bv.memberRepairs = map[string]*armedRepair{} + } + bv.memberRepairs[key] = a + bv.memberRepairsMu.Unlock() + return a.reach, func() { + bv.memberRepairsMu.Lock() + if bv.memberRepairs[key] == a { + delete(bv.memberRepairs, key) + } + bv.memberRepairsMu.Unlock() + } +} + +// takeMemberRepair returns and withdraws the repair naming this member, if any: a repair is used once. +func (bv *BFTValidator) takeMemberRepair(intentID string, chainID int64) *armedRepair { + bv.memberRepairsMu.Lock() + defer bv.memberRepairsMu.Unlock() + key := memberRepairKey(intentID, chainID) + a := bv.memberRepairs[key] + delete(bv.memberRepairs, key) + return a +} + +// memberRepairState is the repair registry and the proof-cycle starter a BFTValidator carries. +type memberRepairState struct { + memberRepairsMu sync.Mutex + memberRepairs map[string]*armedRepair + // memberRepairStart starts a member's proof cycle; nil is RunBatchMemberAttestation (tests replace it). + memberRepairStart func(ctx context.Context, att *PendingAttestation, tx string, chainID int64, lane string) +} + +// repairDecidedMember runs the repair armed for this decided member, if one is: from the round's snapshot, on the +// settlement the repair names. +func (bv *BFTValidator) repairDecidedMember(att *PendingAttestation, chainID int64, lane string) { + a := bv.takeMemberRepair(att.IntentID, chainID) + if a == nil { + return + } + snap := MemberRepairSnapshot{ + IntentID: att.IntentID, Lane: lane, BundleID: att.BundleIDHex, GovernanceRoot: att.GovernanceProofRoot, + OperationCommitment: att.OperationCommitment, AccumulateHeight: att.AccumulateBlockHeight, + GovernanceLevel: att.GovernanceLevel, GovernanceLevels: att.G0Proof != nil && att.G1Proof != nil && att.G2Proof != nil, + BLSSignature: att.BLSSignature != "", AccountURL: att.AccountURL, TransactionHash: att.TransactionHash, + } + var missing []string + for _, m := range []struct { + absent bool + name string + }{ + {att.CertenIntent == nil, "the intent"}, + {att.BundleIDHex == "", "the committed block's bundle id"}, + {att.GovernanceProofRoot == "", "the governance root"}, + {att.OperationCommitment == "", "the operation commitment"}, + {att.G0Proof == nil, "G0"}, {att.G1Proof == nil, "G1"}, {att.G2Proof == nil, "G2"}, + {att.AccountURL == "", "the Accumulate account"}, + {att.TransactionHash == "", "the Accumulate transaction"}, + {a.SettlementTx == "", "the settlement transaction"}, + } { + if m.absent { + missing = append(missing, m.name) + } + } + if len(missing) > 0 { + err := fmt.Errorf("intent %s member %d: the re-derived round lacks %s; its proof cycle is not run", + att.IntentID, chainID, strings.Join(missing, ", ")) + bv.logger.Printf("🛑 [MEMBER-REPAIR] %v", err) + a.reach <- MemberRepairReach{Snapshot: snap, Err: err} + return + } + if !a.Apply { + bv.logger.Printf("🔎 [MEMBER-REPAIR] intent %s member %d reached (dry run): bundle %s, governance root %s, lane %s", + att.IntentID, chainID, att.BundleIDHex, att.GovernanceProofRoot, lane) + a.reach <- MemberRepairReach{Snapshot: snap} + return + } + start := bv.memberRepairStart + if start == nil { + start = func(ctx context.Context, att *PendingAttestation, tx string, chainID int64, lane string) { + bv.RunBatchMemberAttestation(ctx, att, tx, chainID, true, lane) + } + } + bv.logger.Printf("🔧 [MEMBER-REPAIR] intent %s member %d: running its proof cycle on settlement %s (lane %s)", + att.IntentID, chainID, a.SettlementTx, lane) + start(context.Background(), att, a.SettlementTx, chainID, lane) + a.reach <- MemberRepairReach{Snapshot: snap, Started: true} +} diff --git a/pkg/consensus/member_repair_test.go b/pkg/consensus/member_repair_test.go new file mode 100644 index 00000000..1720bc00 --- /dev/null +++ b/pkg/consensus/member_repair_test.go @@ -0,0 +1,159 @@ +package consensus + +import ( + "context" + "fmt" + "strings" + "testing" + + "github.com/certen/independant-validator/pkg/proof" +) + +// RB4-F55 repair. A member whose proof cycle failed has a recorded outcome, so the batch path never queues it +// again (RB3-F141) and nothing re-runs its proof cycle. Intent 000ac79a's base member settled on chain and was +// recorded failed by a cycle a fleet restart broke. The repair re-derives the member's round exactly as discovery +// does (re-processing the intent) and, at that one decided member - only when a repair names it - runs the +// member's proof cycle on its settlement, from the snapshot the round captured. Nothing is invented: a snapshot +// that lacks what the proof cycle binds is refused by name. + +type startedRepair struct { + att *PendingAttestation + tx string + chainID int64 + lane string +} + +func repairValidator(t *testing.T, f *fakeEnqueuer) (*BFTValidator, *[]startedRepair) { + t.Helper() + bv := refusalValidator(f) + var started []startedRepair + bv.memberRepairStart = func(_ context.Context, att *PendingAttestation, tx string, chainID int64, lane string) { + started = append(started, startedRepair{att, tx, chainID, lane}) + } + return bv, &started +} + +// A round's values as consensus would have them: the committed block's bundle, commitment and governance root, +// and G0-G2. +func committedRound(ci *CertenIntent) func(bv *BFTValidator) error { + vb := &ValidatorBlock{BundleID: "0x" + strings.Repeat("0a", 32), OperationCommitment: "0x" + strings.Repeat("0b", 32)} + vb.GovernanceProof.MerkleRoot = "0x" + strings.Repeat("0c", 32) + return func(bv *BFTValidator) error { + return bv.enqueueForBatch(ci, &proof.CertenProof{}, vb, 10007772, &proof.G0Result{}, &proof.G1Result{}, &proof.G2Result{}, + "bls-signature", []string{"validator-signature"}, "G2", 10007772) + } +} + +// repairIntent is an intent as discovery delivers it: with the Accumulate transaction that carries it. +func repairIntent(t *testing.T, id string, chains ...int64) *CertenIntent { + ci := batchableIntent(t, id, chains...) + ci.TransactionHash = "db0236d87fa3c0fd8e6f1c21cba1ce7527bf592c5f727632ad7f329cacbe0e48" + return ci +} + +func decided(f *fakeEnqueuer, chains ...int64) { + for _, c := range chains { + f.addErr[c] = fmt.Errorf("%w: recorded", ErrMemberAlreadyDecided) + } +} + +func TestANamedDecidedMemberIsReDrivenFromItsRederivedRound(t *testing.T) { + f := newFakeEnqueuer() + bv, started := repairValidator(t, f) + ci := repairIntent(t, "000ac79a-repair", 84532, 421614) + decided(f, 84532, 421614) + + reach, disarm := bv.ArmMemberRepair(MemberRepair{IntentID: ci.IntentID, ChainID: 84532, SettlementTx: "0xc409", Apply: true}) + defer disarm() + if err := committedRound(ci)(bv); err != nil { + t.Fatalf("the re-derived round: %v", err) + } + + if len(*started) != 1 { + t.Fatalf("THE regression: the named decided member's proof cycle ran %d times, want 1 (and never the other member)", len(*started)) + } + s := (*started)[0] + if s.chainID != 84532 || s.tx != "0xc409" || s.lane != "on_cadence" || s.att.IntentID != ci.IntentID || + s.att.BundleIDHex == "" || s.att.GovernanceProofRoot == "" || s.att.OperationCommitment == "" || s.att.G2Proof == nil { + t.Fatalf("the proof cycle was not started from the round's snapshot: %+v / %+v", s, s.att) + } + select { + case r := <-reach: + if r.Err != nil || !r.Started || r.Snapshot.BundleID != s.att.BundleIDHex || r.Snapshot.Lane != "on_cadence" { + t.Fatalf("the repair was told %+v", r) + } + default: + t.Fatal("the repair was not told its member was reached") + } +} + +func TestADecidedMemberNoRepairNamesIsNotReDriven(t *testing.T) { + f := newFakeEnqueuer() + bv, started := repairValidator(t, f) + ci := repairIntent(t, "000ac79a-unnamed", 84532) + decided(f, 84532) + if err := committedRound(ci)(bv); err != nil { + t.Fatal(err) + } + if len(*started) != 0 { + t.Fatal("a decided member was re-driven without a repair naming it") + } +} + +func TestARepairIsUsedOnce(t *testing.T) { + f := newFakeEnqueuer() + bv, started := repairValidator(t, f) + ci := repairIntent(t, "000ac79a-once", 84532) + decided(f, 84532) + _, disarm := bv.ArmMemberRepair(MemberRepair{IntentID: ci.IntentID, ChainID: 84532, SettlementTx: "0xc409", Apply: true}) + defer disarm() + round := committedRound(ci) + if err := round(bv); err != nil { + t.Fatal(err) + } + if err := round(bv); err != nil { + t.Fatal(err) + } + if len(*started) != 1 { + t.Fatalf("one repair re-drove the member %d times", len(*started)) + } +} + +func TestADryRunRepairReportsTheSnapshotAndStartsNothing(t *testing.T) { + f := newFakeEnqueuer() + bv, started := repairValidator(t, f) + ci := repairIntent(t, "000ac79a-dry", 84532) + decided(f, 84532) + reach, disarm := bv.ArmMemberRepair(MemberRepair{IntentID: ci.IntentID, ChainID: 84532, SettlementTx: "0xc409"}) + defer disarm() + if err := committedRound(ci)(bv); err != nil { + t.Fatal(err) + } + if len(*started) != 0 { + t.Fatal("a dry run started a proof cycle") + } + r := <-reach + if r.Err != nil || r.Started || r.Snapshot.GovernanceRoot == "" || r.Snapshot.OperationCommitment == "" || !r.Snapshot.GovernanceLevels { + t.Fatalf("a dry run reports %+v", r) + } +} + +func TestARoundWithoutWhatTheProofCycleBindsIsRefusedByName(t *testing.T) { + f := newFakeEnqueuer() + bv, started := repairValidator(t, f) + ci := repairIntent(t, "000ac79a-incomplete", 84532) + decided(f, 84532) + reach, disarm := bv.ArmMemberRepair(MemberRepair{IntentID: ci.IntentID, ChainID: 84532, SettlementTx: "0xc409", Apply: true}) + defer disarm() + // No committed block and no governance levels: nothing to bind the proof cycle to. + if err := bv.enqueueForBatch(ci, nil, nil, 7, nil, nil, nil, "", nil, "", 7); err != nil { + t.Fatal(err) + } + if len(*started) != 0 { + t.Fatal("a proof cycle was started from a round that binds nothing") + } + r := <-reach + if r.Err == nil || r.Started || !strings.Contains(r.Err.Error(), "bundle") || !strings.Contains(r.Err.Error(), "G0") { + t.Fatalf("an incomplete round must be refused naming what it lacks: %+v", r) + } +} diff --git a/pkg/database/intent_failure_class_test.go b/pkg/database/intent_failure_class_test.go index ca08e562..8561f48f 100644 --- a/pkg/database/intent_failure_class_test.go +++ b/pkg/database/intent_failure_class_test.go @@ -63,14 +63,14 @@ func TestAFailedIntentCarriesItsFailureClass(t *testing.T) { // A chain member that did not settle fails its intent as settlement_failed; settling it after clears the class. m := newIntent() - if _, err := repo.RecordMemberOutcome(ctx, MemberOutcome{IntentID: m, ChainID: 84532, MemberChains: []int64{84532}, + if _, err := repo.RecordMemberOutcome(ctx, MemberOutcome{IntentID: m, ReportedBy: "validator-test", ChainID: 84532, MemberChains: []int64{84532}, Settlement: MemberSettlementReverted, ProofCycle: MemberProofCycleWritten, Legs: 1, SettlementTx: "0x" + uuid.NewString()[:8], Reason: "reverted"}); err != nil { t.Fatal(err) } if st, class := read(m); st != "failed" || class == nil || *class != "settlement_failed" { t.Fatalf("a reverted member left the intent %s / %v", st, class) } - if _, err := repo.RecordMemberOutcome(ctx, MemberOutcome{IntentID: m, ChainID: 84532, MemberChains: []int64{84532}, + if _, err := repo.RecordMemberOutcome(ctx, MemberOutcome{IntentID: m, ReportedBy: "validator-test", ChainID: 84532, MemberChains: []int64{84532}, Settlement: MemberSettlementSettled, ProofCycle: MemberProofCycleWritten, Legs: 1, SettlementTx: "0x" + uuid.NewString()[:8]}); err != nil { t.Fatal(err) } diff --git a/pkg/database/intent_member_outcomes.go b/pkg/database/intent_member_outcomes.go index f2d6121b..2ac4f5e1 100644 --- a/pkg/database/intent_member_outcomes.go +++ b/pkg/database/intent_member_outcomes.go @@ -60,6 +60,9 @@ type MemberOutcome struct { // committed none (a native transfer) or they were not assessed, true proven, false provably absent - // the member settled but did not do what the intent committed to, and counts as failed (RB3-F67). EffectsProven *bool + // ReportedBy is the validator reporting the outcome. A report that replaces a recorded outcome is recorded + // as a correction under its name (RB4-F58). + ReportedBy string } // RecordedMemberOutcome is a member's outcome as recorded. @@ -104,6 +107,8 @@ func (o *MemberOutcome) validate() error { return fmt.Errorf("%w: no chain id", ErrMemberOutcomeInvalid) case o.Legs <= 0: return fmt.Errorf("%w: a member carries at least one leg", ErrMemberOutcomeInvalid) + case o.ReportedBy == "": + return fmt.Errorf("%w: no reporting validator", ErrMemberOutcomeInvalid) } inSet := false for _, c := range o.MemberChains { @@ -160,10 +165,15 @@ func (r *IntentLifecycleRepository) RecordMemberOutcome(ctx context.Context, o M } defer tx.Rollback() //nolint:errcheck // a committed transaction ignores it - // Lock the intent row first: every report for this intent derives in turn. + // Lock the intent row first: every report for this intent derives in turn. Its terminal outcome is read + // with it, so a change to it can be recorded (RB4-F58). var stored pq.Int64Array - err = tx.QueryRowContext(ctx, - `SELECT member_chains FROM intent_lifecycle WHERE intent_id = $1 FOR UPDATE`, o.IntentID).Scan(&stored) + var before lifecycleOutcome + err = tx.QueryRowContext(ctx, ` + SELECT member_chains, status, legs_completed, legs_failed, failed_at, completed_at, error_message, failure_class, write_back_tx + FROM intent_lifecycle WHERE intent_id = $1 FOR UPDATE`, o.IntentID). + Scan(&stored, &before.Status, &before.LegsCompleted, &before.LegsFailed, &before.FailedAt, &before.CompletedAt, + &before.ErrorMessage, &before.FailureClass, &before.WriteBackTx) switch { case errors.Is(err, sql.ErrNoRows): derived.Found = false @@ -186,6 +196,27 @@ func (r *IntentLifecycleRepository) RecordMemberOutcome(ctx context.Context, o M } } + // The member's recorded outcome, if any. A member written back under a quorum attestation is never replaced by a + // report that it was not: the write-back is on Accumulate, so that report can only be stale (a broken cycle's + // outbox entry replayed after a repair). Any other change is recorded as a correction. + prior, err := recordedMemberRow(ctx, tx, o.IntentID, o.ChainID) + if err != nil { + return derived, err + } + next := memberRowOf(o) + if prior != nil && prior.ProofCycle == string(MemberProofCycleWritten) && o.ProofCycle != MemberProofCycleWritten { + return derived, fmt.Errorf("%w: intent %s member %d was written back by cycle %s (%s); a later report that it was not (cycle %s) is stale", + ErrMemberOutcomeInvalid, o.IntentID, o.ChainID, prior.CycleID, prior.WriteBackTx, o.CycleID) + } + if prior != nil && !prior.sameAs(next) { + reason := fmt.Sprintf("member outcome replaced by a later report: cycle %s replaces cycle %s", o.CycleID, prior.CycleID) + evidence := map[string]any{"settlement_tx": o.SettlementTx, "write_back_tx": o.WriteBackTx, "cycle_id": o.CycleID} + if _, err := recordCorrection(ctx, tx, "intent_member_outcome", fmt.Sprintf("%s/%d", o.IntentID, o.ChainID), + reason, prior, next, evidence, o.ReportedBy); err != nil { + return derived, fmt.Errorf("record the replaced outcome of %s/%d: %w", o.IntentID, o.ChainID, err) + } + } + if _, err := tx.ExecContext(ctx, ` INSERT INTO intent_member_outcomes (intent_id, chain_id, settlement, proof_cycle, legs, settlement_tx, write_back_tx, cycle_id, reason, effects_proven, recorded_at) @@ -304,9 +335,115 @@ func (r *IntentLifecycleRepository) RecordMemberOutcome(ctx context.Context, o M if err != nil { return derived, fmt.Errorf("derive status of %s: %w", o.IntentID, err) } + + // A terminal outcome that changes is a correction of a published outcome (RB4-F58): keep what it was. + if before.terminal() { + var after lifecycleOutcome + if err := tx.QueryRowContext(ctx, ` + SELECT status, legs_completed, legs_failed, failed_at, completed_at, error_message, failure_class, write_back_tx + FROM intent_lifecycle WHERE intent_id = $1`, o.IntentID). + Scan(&after.Status, &after.LegsCompleted, &after.LegsFailed, &after.FailedAt, &after.CompletedAt, + &after.ErrorMessage, &after.FailureClass, &after.WriteBackTx); err != nil { + return derived, fmt.Errorf("read the derived status of %s: %w", o.IntentID, err) + } + if before.Status != after.Status || before.ErrorMessage != after.ErrorMessage || before.FailureClass != after.FailureClass { + reason := fmt.Sprintf("intent outcome derived again after member %d's outcome was replaced (cycle %s)", o.ChainID, o.CycleID) + evidence := map[string]any{"chain_id": o.ChainID, "settlement_tx": o.SettlementTx, "write_back_tx": o.WriteBackTx, "cycle_id": o.CycleID} + if _, err := recordCorrection(ctx, tx, "intent_lifecycle", o.IntentID, reason, before.view(), after.view(), evidence, o.ReportedBy); err != nil { + return derived, fmt.Errorf("record the replaced outcome of %s: %w", o.IntentID, err) + } + } + } return derived, tx.Commit() } +// memberRow is a member outcome as stored, as a correction records it. +type memberRow struct { + Settlement string `json:"settlement"` + ProofCycle string `json:"proof_cycle"` + Legs int `json:"legs"` + SettlementTx string `json:"settlement_tx"` + WriteBackTx string `json:"write_back_tx"` + CycleID string `json:"cycle_id"` + Reason string `json:"reason"` + EffectsProven *bool `json:"effects_proven"` +} + +func memberRowOf(o MemberOutcome) memberRow { + return memberRow{Settlement: string(o.Settlement), ProofCycle: string(o.ProofCycle), Legs: o.Legs, + SettlementTx: o.SettlementTx, WriteBackTx: o.WriteBackTx, CycleID: o.CycleID, Reason: o.Reason, EffectsProven: o.EffectsProven} +} + +func (m *memberRow) sameAs(n memberRow) bool { + sameEffects := (m.EffectsProven == nil) == (n.EffectsProven == nil) && + (m.EffectsProven == nil || *m.EffectsProven == *n.EffectsProven) + return m.Settlement == n.Settlement && m.ProofCycle == n.ProofCycle && m.Legs == n.Legs && m.SettlementTx == n.SettlementTx && + m.WriteBackTx == n.WriteBackTx && m.CycleID == n.CycleID && m.Reason == n.Reason && sameEffects +} + +// recordedMemberRow reads (and locks) a member's recorded outcome, or nil when none is recorded. +func recordedMemberRow(ctx context.Context, tx *sql.Tx, intentID string, chainID int64) (*memberRow, error) { + var m memberRow + var settlementTx, writeBackTx, cycleID, reason sql.NullString + var effects sql.NullBool + err := tx.QueryRowContext(ctx, ` + SELECT settlement, proof_cycle, legs, settlement_tx, write_back_tx, cycle_id, reason, effects_proven + FROM intent_member_outcomes WHERE intent_id = $1 AND chain_id = $2 FOR UPDATE`, intentID, chainID). + Scan(&m.Settlement, &m.ProofCycle, &m.Legs, &settlementTx, &writeBackTx, &cycleID, &reason, &effects) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("read the recorded outcome of %s/%d: %w", intentID, chainID, err) + } + m.SettlementTx, m.WriteBackTx, m.CycleID, m.Reason = settlementTx.String, writeBackTx.String, cycleID.String, reason.String + if effects.Valid { + v := effects.Bool + m.EffectsProven = &v + } + return &m, nil +} + +// lifecycleOutcome is an intent's recorded outcome, as a correction records it. +type lifecycleOutcome struct { + Status sql.NullString + LegsCompleted, LegsFailed sql.NullInt64 + FailedAt, CompletedAt sql.NullTime + ErrorMessage sql.NullString + FailureClass sql.NullString + WriteBackTx sql.NullString +} + +func (l lifecycleOutcome) terminal() bool { + return l.Status.String == string(IntentLifecycleComplete) || l.Status.String == string(IntentLifecycleFailed) +} + +func (l lifecycleOutcome) view() map[string]any { + opt := func(v sql.NullString) any { + if v.Valid { + return v.String + } + return nil + } + optTime := func(v sql.NullTime) any { + if v.Valid { + return v.Time.UTC().Format(time.RFC3339Nano) + } + return nil + } + optInt := func(v sql.NullInt64) any { + if v.Valid { + return v.Int64 + } + return nil + } + return map[string]any{ + "status": opt(l.Status), "legs_completed": optInt(l.LegsCompleted), "legs_failed": optInt(l.LegsFailed), + "failed_at": optTime(l.FailedAt), "completed_at": optTime(l.CompletedAt), "error_message": opt(l.ErrorMessage), + "failure_class": opt(l.FailureClass), "write_back_tx": opt(l.WriteBackTx), + } +} + func sameChains(a pq.Int64Array, b []int64) bool { x := append([]int64(nil), a...) sort.Slice(x, func(i, j int) bool { return x[i] < x[j] }) diff --git a/pkg/database/member_outcome_corrections_test.go b/pkg/database/member_outcome_corrections_test.go new file mode 100644 index 00000000..d03aaebe --- /dev/null +++ b/pkg/database/member_outcome_corrections_test.go @@ -0,0 +1,181 @@ +package database + +import ( + "context" + "encoding/json" + "errors" + "testing" + + "github.com/google/uuid" +) + +// RB4-F58. A member's recorded outcome is replaced by a later report for the same chain, and the intent's status is +// derived again (RecordMemberOutcome). Replacing it overwrote the member row and, when the intent went from failed +// to complete, erased failed_at, error_message and failure_class - a published failure rewritten with no record. +// Production 2026-09-29: intent 000ac79a's base member was recorded failed by a proof cycle a restart broke (RB4-F55) +// although its settlement landed; its repair must replace that record, and the replacement must be on record. +// +// Every replacement that changes a member row, and every change of an intent's terminal status, now leaves an +// evidence_corrections row in the same transaction: what was recorded, what replaced it, the evidence, and who. +// A member already written back is never replaced by a report that it was not: that report can only be stale. + +type correctionRow struct { + RecordType, RecordID, Reason, CorrectedBy string + Previous, Corrected, Evidence map[string]any +} + +func correctionsFor(t *testing.T, ctx context.Context, recordType, recordID string) []correctionRow { + t.Helper() + rows, err := testDB.QueryContext(ctx, ` + SELECT record_type, record_id, reason, corrected_by, previous, corrected, chain_evidence + FROM evidence_corrections WHERE record_type = $1 AND record_id = $2 ORDER BY corrected_at`, recordType, recordID) + if err != nil { + t.Fatalf("read corrections: %v", err) + } + defer rows.Close() + var out []correctionRow + for rows.Next() { + var c correctionRow + var prev, next, ev []byte + if err := rows.Scan(&c.RecordType, &c.RecordID, &c.Reason, &c.CorrectedBy, &prev, &next, &ev); err != nil { + t.Fatalf("scan correction: %v", err) + } + for _, p := range []struct { + b []byte + m *map[string]any + }{{prev, &c.Previous}, {next, &c.Corrected}, {ev, &c.Evidence}} { + if err := json.Unmarshal(p.b, p.m); err != nil { + t.Fatalf("decode correction json: %v", err) + } + } + out = append(out, c) + } + return out +} + +func newTwoMemberIntent(t *testing.T, ctx context.Context) string { + t.Helper() + id := "f58-" + uuid.NewString() + if _, err := testDB.ExecContext(ctx, `INSERT INTO intent_lifecycle (intent_id, accum_tx_hash, status) VALUES ($1, $2, 'settling')`, + id, uuid.NewString()[:16]); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + bg := context.Background() + _, _ = testDB.ExecContext(bg, `DELETE FROM intent_member_outcomes WHERE intent_id = $1`, id) + _, _ = testDB.ExecContext(bg, `DELETE FROM intent_lifecycle WHERE intent_id = $1`, id) + }) + return id +} + +func TestARepairedMemberLeavesTheFailureOnRecord(t *testing.T) { + requireTestDB(t) + ctx := context.Background() + repo := NewIntentLifecycleRepository(NewClientFromDB(testDB)) + id := newTwoMemberIntent(t, ctx) + members := []int64{84532, 421614} + + // Arbitrum settled and was written back; base was recorded failed by a cycle a restart broke. + if _, err := repo.RecordMemberOutcome(ctx, MemberOutcome{IntentID: id, ChainID: 421614, MemberChains: members, Legs: 1, + Settlement: MemberSettlementSettled, ProofCycle: MemberProofCycleWritten, SettlementTx: "0xarb", WriteBackTx: "acc://wb-arb", + CycleID: "cycle-arb", ReportedBy: "validator-1"}); err != nil { + t.Fatal(err) + } + failed := MemberOutcome{IntentID: id, ChainID: 84532, MemberChains: members, Legs: 1, + Settlement: MemberSettlementUnobserved, ProofCycle: MemberProofCycleFailed, CycleID: "cycle-broken", + Reason: "phase 7 failed: persist chain execution 0xc409: duplicate key", ReportedBy: "validator-6"} + d, err := repo.RecordMemberOutcome(ctx, failed) + if err != nil || d.Status != IntentLifecycleFailed { + t.Fatalf("the broken cycle's report: %v %+v", err, d) + } + if got := correctionsFor(t, ctx, "intent_member_outcome", id+"/84532"); len(got) != 0 { + t.Fatalf("a member's FIRST outcome is not a correction: %d recorded", len(got)) + } + if got := correctionsFor(t, ctx, "intent_lifecycle", id); len(got) != 0 { + t.Fatalf("an intent's first terminal status is not a correction: %d recorded", len(got)) + } + + // The repair: the member's proof cycle re-run, settled and written back. + repaired := MemberOutcome{IntentID: id, ChainID: 84532, MemberChains: members, Legs: 1, + Settlement: MemberSettlementSettled, ProofCycle: MemberProofCycleWritten, SettlementTx: "0xc409", WriteBackTx: "acc://wb-base", + CycleID: "cycle-repair", ReportedBy: "validator-6"} + d, err = repo.RecordMemberOutcome(ctx, repaired) + if err != nil || d.Status != IntentLifecycleComplete { + t.Fatalf("the repair: %v %+v", err, d) + } + + member := correctionsFor(t, ctx, "intent_member_outcome", id+"/84532") + if len(member) != 1 { + t.Fatalf("THE regression: the replaced member outcome left %d corrections, want 1", len(member)) + } + m := member[0] + if m.CorrectedBy != "validator-6" || m.Previous["proof_cycle"] != "failed" || m.Previous["cycle_id"] != "cycle-broken" || + m.Previous["reason"] != failed.Reason || m.Corrected["proof_cycle"] != "written" || m.Corrected["cycle_id"] != "cycle-repair" || + m.Evidence["settlement_tx"] != "0xc409" || m.Evidence["write_back_tx"] != "acc://wb-base" || m.Reason == "" { + t.Fatalf("the member correction does not state what was replaced, by what and on what evidence: %+v", m) + } + + lifecycle := correctionsFor(t, ctx, "intent_lifecycle", id) + if len(lifecycle) != 1 { + t.Fatalf("THE regression: failed -> complete erased the failure with %d corrections, want 1", len(lifecycle)) + } + l := lifecycle[0] + if l.Previous["status"] != "failed" || l.Previous["failure_class"] != "settlement_failed" || l.Previous["error_message"] == nil || + l.Previous["failed_at"] == nil || l.Corrected["status"] != "complete" || l.CorrectedBy != "validator-6" { + t.Fatalf("the lifecycle correction does not keep the failure it replaced: %+v", l) + } + + // The same report again (an outbox replay) changes nothing and records nothing. + if _, err := repo.RecordMemberOutcome(ctx, repaired); err != nil { + t.Fatal(err) + } + if n := len(correctionsFor(t, ctx, "intent_member_outcome", id+"/84532")); n != 1 { + t.Fatalf("an identical report was recorded as a correction: %d", n) + } + if n := len(correctionsFor(t, ctx, "intent_lifecycle", id)); n != 1 { + t.Fatalf("an unchanged status was recorded as a correction: %d", n) + } +} + +func TestAWrittenBackMemberIsNeverReplacedByAStaleFailure(t *testing.T) { + requireTestDB(t) + ctx := context.Background() + repo := NewIntentLifecycleRepository(NewClientFromDB(testDB)) + id := newTwoMemberIntent(t, ctx) + members := []int64{84532} + + written := MemberOutcome{IntentID: id, ChainID: 84532, MemberChains: members, Legs: 1, + Settlement: MemberSettlementSettled, ProofCycle: MemberProofCycleWritten, SettlementTx: "0xc409", WriteBackTx: "acc://wb", + CycleID: "cycle-repair", ReportedBy: "validator-6"} + if _, err := repo.RecordMemberOutcome(ctx, written); err != nil { + t.Fatal(err) + } + // A failure report of the same member left in an outbox by the broken cycle, replayed after the repair. + stale := MemberOutcome{IntentID: id, ChainID: 84532, MemberChains: members, Legs: 1, + Settlement: MemberSettlementUnobserved, ProofCycle: MemberProofCycleFailed, CycleID: "cycle-broken", + Reason: "phase 7 failed", ReportedBy: "validator-6"} + _, err := repo.RecordMemberOutcome(ctx, stale) + if !errors.Is(err, ErrMemberOutcomeInvalid) { + t.Fatalf("a stale failure over a written-back member: want ErrMemberOutcomeInvalid, got %v", err) + } + var cycle, status string + if err := testDB.QueryRowContext(ctx, `SELECT m.proof_cycle, l.status FROM intent_member_outcomes m JOIN intent_lifecycle l USING (intent_id) + WHERE m.intent_id = $1 AND m.chain_id = 84532`, id).Scan(&cycle, &status); err != nil { + t.Fatal(err) + } + if cycle != "written" || status != "complete" { + t.Fatalf("the stale report replaced the written-back member: %s / %s", cycle, status) + } +} + +func TestAMemberOutcomeNamesItsReporter(t *testing.T) { + requireTestDB(t) + ctx := context.Background() + repo := NewIntentLifecycleRepository(NewClientFromDB(testDB)) + id := newTwoMemberIntent(t, ctx) + _, err := repo.RecordMemberOutcome(ctx, MemberOutcome{IntentID: id, ChainID: 84532, MemberChains: []int64{84532}, Legs: 1, + Settlement: MemberSettlementSettled, ProofCycle: MemberProofCycleWritten, CycleID: "c"}) + if !errors.Is(err, ErrMemberOutcomeInvalid) { + t.Fatalf("an outcome without its reporter: want ErrMemberOutcomeInvalid, got %v", err) + } +} diff --git a/pkg/database/member_write_backs.go b/pkg/database/member_write_backs.go new file mode 100644 index 00000000..99e21189 --- /dev/null +++ b/pkg/database/member_write_backs.go @@ -0,0 +1,186 @@ +// Copyright 2026 Certen Protocol +// +// A chain member's outcome is written back to Accumulate once (RB4-F59). + +package database + +import ( + "context" + "database/sql" + "errors" + "fmt" + "time" +) + +// ErrMemberAlreadyWrittenBack: the member's outcome is already on Accumulate; another write-back would be a +// second entry for the same member. +var ErrMemberAlreadyWrittenBack = errors.New("member already written back") + +// ErrMemberWriteBackUnresolved: a cycle claimed the member's write-back and its outcome is not known - the +// submission may have reached Accumulate. No other write-back is submitted until that is established. +var ErrMemberWriteBackUnresolved = errors.New("member write-back outcome unknown") + +// MemberWriteBack is a member's write-back as registered. +type MemberWriteBack struct { + IntentID string + ChainID int64 + State string // claimed | written | not_sent + CycleID string + ValidatorID string + WriteBackTx string + Reason string + ClaimedAt time.Time +} + +// Write-back register states. +const ( + MemberWriteBackClaimed = "claimed" + MemberWriteBackWritten = "written" + MemberWriteBackNotSent = "not_sent" +) + +// MemberWriteBackOf returns a member's registered write-back, or nil when none is registered. +func (r *IntentLifecycleRepository) MemberWriteBackOf(ctx context.Context, intentID string, chainID int64) (*MemberWriteBack, error) { + return memberWriteBackOf(ctx, r.client.db, intentID, chainID, false) +} + +type queryRower interface { + QueryRowContext(ctx context.Context, query string, args ...any) *sql.Row +} + +func memberWriteBackOf(ctx context.Context, q queryRower, intentID string, chainID int64, lock bool) (*MemberWriteBack, error) { + query := `SELECT state, cycle_id, validator_id, write_back_tx, reason, claimed_at + FROM member_write_backs WHERE intent_id = $1 AND chain_id = $2` + if lock { + query += ` FOR UPDATE` + } + w := MemberWriteBack{IntentID: intentID, ChainID: chainID} + var tx, reason sql.NullString + err := q.QueryRowContext(ctx, query, intentID, chainID).Scan(&w.State, &w.CycleID, &w.ValidatorID, &tx, &reason, &w.ClaimedAt) + if errors.Is(err, sql.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, fmt.Errorf("read the write-back of %s/%d: %w", intentID, chainID, err) + } + w.WriteBackTx, w.Reason = tx.String, reason.String + return &w, nil +} + +// refusal names why a registered write-back stops another one, or nil when it does not. +func (w *MemberWriteBack) refusal() error { + switch { + case w == nil || w.State == MemberWriteBackNotSent: + return nil + case w.State == MemberWriteBackWritten: + return fmt.Errorf("%w: intent %s member %d as %s by cycle %s (%s)", + ErrMemberAlreadyWrittenBack, w.IntentID, w.ChainID, w.WriteBackTx, w.CycleID, w.ValidatorID) + default: + return fmt.Errorf("%w: intent %s member %d was claimed by cycle %s (%s) at %s and its outcome was never recorded", + ErrMemberWriteBackUnresolved, w.IntentID, w.ChainID, w.CycleID, w.ValidatorID, w.ClaimedAt.UTC().Format(time.RFC3339)) + } +} + +// writtenBackOutcome refuses a member whose recorded outcome says it was written back. A member written back before +// write-backs were registered has no register row; its recorded outcome states the write-back all the same. +func writtenBackOutcome(ctx context.Context, q queryRower, intentID string, chainID int64) error { + var writeBackTx, cycleID sql.NullString + err := q.QueryRowContext(ctx, ` + SELECT write_back_tx, cycle_id FROM intent_member_outcomes + WHERE intent_id = $1 AND chain_id = $2 AND proof_cycle = 'written'`, intentID, chainID).Scan(&writeBackTx, &cycleID) + if errors.Is(err, sql.ErrNoRows) { + return nil + } + if err != nil { + return fmt.Errorf("read the recorded outcome of %s/%d: %w", intentID, chainID, err) + } + return fmt.Errorf("%w: intent %s member %d as %s by cycle %s (its recorded outcome)", + ErrMemberAlreadyWrittenBack, intentID, chainID, writeBackTx.String, cycleID.String) +} + +// MemberWriteBackAllowed is nil when the member may be written back - nothing registered and no recorded outcome +// saying it was, or a claim that was never sent - and otherwise names why not. +func (r *IntentLifecycleRepository) MemberWriteBackAllowed(ctx context.Context, intentID string, chainID int64) error { + w, err := r.MemberWriteBackOf(ctx, intentID, chainID) + if err != nil { + return err + } + if refusal := w.refusal(); refusal != nil { + return refusal + } + return writtenBackOutcome(ctx, r.client.db, intentID, chainID) +} + +// ClaimMemberWriteBack claims the member's write-back for a cycle, before it is submitted. A member written back, +// or claimed with its outcome unknown, is refused by name; a claim never sent is taken over. +func (r *IntentLifecycleRepository) ClaimMemberWriteBack(ctx context.Context, intentID string, chainID int64, cycleID, validatorID string) error { + if intentID == "" || chainID == 0 || cycleID == "" || validatorID == "" { + return fmt.Errorf("claim a write-back: intent, chain, cycle and validator are required") + } + tx, err := r.client.db.BeginTx(ctx, nil) + if err != nil { + return fmt.Errorf("begin write-back claim: %w", err) + } + defer tx.Rollback() //nolint:errcheck // a committed transaction ignores it + + if err := writtenBackOutcome(ctx, tx, intentID, chainID); err != nil { + return err + } + res, err := tx.ExecContext(ctx, ` + INSERT INTO member_write_backs (intent_id, chain_id, state, cycle_id, validator_id) + VALUES ($1, $2, 'claimed', $3, $4) + ON CONFLICT (intent_id, chain_id) DO NOTHING`, intentID, chainID, cycleID, validatorID) + if err != nil { + return fmt.Errorf("claim the write-back of %s/%d: %w", intentID, chainID, err) + } + if n, _ := res.RowsAffected(); n == 1 { + return tx.Commit() + } + w, err := memberWriteBackOf(ctx, tx, intentID, chainID, true) + if err != nil { + return err + } + if w == nil { + return fmt.Errorf("claim the write-back of %s/%d: the registered claim vanished", intentID, chainID) + } + if refusal := w.refusal(); refusal != nil { + return refusal + } + if _, err := tx.ExecContext(ctx, ` + UPDATE member_write_backs SET state = 'claimed', cycle_id = $3, validator_id = $4, reason = NULL, + claimed_at = now(), resolved_at = NULL + WHERE intent_id = $1 AND chain_id = $2 AND state = 'not_sent'`, intentID, chainID, cycleID, validatorID); err != nil { + return fmt.Errorf("take over the unsent write-back of %s/%d: %w", intentID, chainID, err) + } + return tx.Commit() +} + +// RecordMemberWriteBack records the claimed write-back as written. +func (r *IntentLifecycleRepository) RecordMemberWriteBack(ctx context.Context, intentID string, chainID int64, cycleID, writeBackTx string) error { + if writeBackTx == "" { + return fmt.Errorf("record the write-back of %s/%d: no transaction", intentID, chainID) + } + return r.resolveMemberWriteBack(ctx, intentID, chainID, cycleID, MemberWriteBackWritten, writeBackTx, "") +} + +// ReleaseMemberWriteBack records that the claiming cycle sent nothing, with why, so another cycle may claim it. +func (r *IntentLifecycleRepository) ReleaseMemberWriteBack(ctx context.Context, intentID string, chainID int64, cycleID, reason string) error { + if reason == "" { + return fmt.Errorf("release the write-back of %s/%d: no reason", intentID, chainID) + } + return r.resolveMemberWriteBack(ctx, intentID, chainID, cycleID, MemberWriteBackNotSent, "", reason) +} + +func (r *IntentLifecycleRepository) resolveMemberWriteBack(ctx context.Context, intentID string, chainID int64, cycleID, state, writeBackTx, reason string) error { + res, err := r.client.db.ExecContext(ctx, ` + UPDATE member_write_backs SET state = $4, write_back_tx = NULLIF($5, ''), reason = NULLIF($6, ''), resolved_at = now() + WHERE intent_id = $1 AND chain_id = $2 AND cycle_id = $3 AND state = 'claimed'`, + intentID, chainID, cycleID, state, writeBackTx, reason) + if err != nil { + return fmt.Errorf("record the write-back of %s/%d as %s: %w", intentID, chainID, state, err) + } + if n, _ := res.RowsAffected(); n != 1 { + return fmt.Errorf("record the write-back of %s/%d as %s: cycle %s holds no claim on it", intentID, chainID, state, cycleID) + } + return nil +} diff --git a/pkg/execution/accumulate_submitter.go b/pkg/execution/accumulate_submitter.go index 9f8408d8..fb51f6f1 100644 --- a/pkg/execution/accumulate_submitter.go +++ b/pkg/execution/accumulate_submitter.go @@ -14,6 +14,7 @@ import ( "crypto/ed25519" "encoding/hex" "encoding/json" + "errors" "fmt" "log" "sync" @@ -143,6 +144,10 @@ func NewAccumulateSubmitter(cfg *AccumulateSubmitterConfig) (*AccumulateSubmitte return submitter, nil } +// ErrWriteBackNotSent: a write-back failed before anything was sent to Accumulate - it is certainly not there +// (RB4-F59). Any other submission error may have reached it. +var ErrWriteBackNotSent = errors.New("write-back not sent") + // SubmitTransaction submits a synthetic transaction to Accumulate // Returns the transaction hash on success func (s *AccumulateSubmitterImpl) SubmitTransaction(ctx context.Context, tx *SyntheticTransaction) (string, error) { @@ -154,24 +159,24 @@ func (s *AccumulateSubmitterImpl) SubmitTransaction(ctx context.Context, tx *Syn // Step 1: Check credit balance hasCredits, balance, err := s.creditChecker.HasSufficientCredits(ctx, MinCreditsForWriteData) if err != nil { - return "", fmt.Errorf("failed to check credits: %w", err) + return "", fmt.Errorf("%w: failed to check credits: %w", ErrWriteBackNotSent, err) } if !hasCredits { - return "", fmt.Errorf("insufficient credits: have %d, need %d", balance, MinCreditsForWriteData) + return "", fmt.Errorf("%w: insufficient credits: have %d, need %d", ErrWriteBackNotSent, balance, MinCreditsForWriteData) } s.logger.Printf("✅ Credit check passed: %d credits available", balance) // Step 2: Create the Accumulate Transaction with proper protocol types accTx, err := s.createAccumulateTransaction(tx) if err != nil { - return "", fmt.Errorf("failed to create Accumulate transaction: %w", err) + return "", fmt.Errorf("%w: failed to create Accumulate transaction: %w", ErrWriteBackNotSent, err) } // Step 3: Create and sign the signature using proper Accumulate signing timestamp := uint64(time.Now().UnixMicro()) sig, err := s.createAndSignSignature(ctx, accTx, timestamp) if err != nil { - return "", fmt.Errorf("failed to create signature: %w", err) + return "", fmt.Errorf("%w: failed to create signature: %w", ErrWriteBackNotSent, err) } // Step 4: Create the envelope diff --git a/pkg/execution/bundle_truth_f81_test.go b/pkg/execution/bundle_truth_f81_test.go index 186edaae..a3243198 100644 --- a/pkg/execution/bundle_truth_f81_test.go +++ b/pkg/execution/bundle_truth_f81_test.go @@ -71,6 +71,8 @@ func f81NonSettlementCycle(t *testing.T) *activeCycle { t.Helper() own := nsMember() f, _ := memberFacts(own) + // Each cycle its own member: a member is written back once (RB4-F59). + f.IntentID = fmt.Sprintf("f81-%d", time.Now().UnixNano()) claim, obs, err := observeNonSettlement(context.Background(), nsChainPast(f.Deadline), f, "its batch quorum was never reached") if err != nil { t.Fatal(err) diff --git a/pkg/execution/intent_member_status_test.go b/pkg/execution/intent_member_status_test.go index d1edbefb..77f5ab9c 100644 --- a/pkg/execution/intent_member_status_test.go +++ b/pkg/execution/intent_member_status_test.go @@ -128,16 +128,16 @@ func TestIntentStatus_AMemberSetThatDisagreesIsRefused(t *testing.T) { s1Seed(ctx, t, db, id) t.Cleanup(func() { db.Exec(`DELETE FROM intent_member_outcomes WHERE intent_id=$1`, id) }) - if _, err := repo.RecordMemberOutcome(ctx, database.MemberOutcome{IntentID: id, ChainID: 84532, MemberChains: []int64{84532, 421614}, + if _, err := repo.RecordMemberOutcome(ctx, database.MemberOutcome{IntentID: id, ReportedBy: "validator-test", ChainID: 84532, MemberChains: []int64{84532, 421614}, Settlement: database.MemberSettlementSettled, ProofCycle: database.MemberProofCycleWritten, Legs: 1}); err != nil { t.Fatal(err) } - _, err := repo.RecordMemberOutcome(ctx, database.MemberOutcome{IntentID: id, ChainID: 421614, MemberChains: []int64{421614}, + _, err := repo.RecordMemberOutcome(ctx, database.MemberOutcome{IntentID: id, ReportedBy: "validator-test", ChainID: 421614, MemberChains: []int64{421614}, Settlement: database.MemberSettlementSettled, ProofCycle: database.MemberProofCycleWritten, Legs: 1}) if !errors.Is(err, database.ErrMemberOutcomeInvalid) { t.Fatalf("a report naming another member set was accepted: %v", err) } - _, err = repo.RecordMemberOutcome(ctx, database.MemberOutcome{IntentID: id, ChainID: 11155111, MemberChains: []int64{84532, 421614}, + _, err = repo.RecordMemberOutcome(ctx, database.MemberOutcome{IntentID: id, ReportedBy: "validator-test", ChainID: 11155111, MemberChains: []int64{84532, 421614}, Settlement: database.MemberSettlementSettled, ProofCycle: database.MemberProofCycleWritten, Legs: 1}) if !errors.Is(err, database.ErrMemberOutcomeInvalid) { t.Fatalf("a chain outside the member set was accepted: %v", err) diff --git a/pkg/execution/member_outcome_outbox.go b/pkg/execution/member_outcome_outbox.go index 7c49a233..49052765 100644 --- a/pkg/execution/member_outcome_outbox.go +++ b/pkg/execution/member_outcome_outbox.go @@ -98,10 +98,13 @@ type MemberOutcomeRecorder interface { // MemberOutcomeReconciler replays the outbox into the lifecycle store: at start, then every Interval. type MemberOutcomeReconciler struct { - Outbox MemberOutcomeOutbox - Store MemberOutcomeRecorder - Interval time.Duration // zero: one minute - Logf func(string, ...interface{}) + Outbox MemberOutcomeOutbox + Store MemberOutcomeRecorder + // ValidatorID reports an entry written before outcomes named their reporter (RB4-F58): the outbox is this + // validator's own, so every entry in it is this validator's report. + ValidatorID string + Interval time.Duration // zero: one minute + Logf func(string, ...interface{}) } // MemberOutcomeReconcileReport is what one pass did. @@ -144,8 +147,8 @@ func (r *MemberOutcomeReconciler) passAndLog(ctx context.Context) { // RunOnce replays every pending outcome once. func (r *MemberOutcomeReconciler) RunOnce(ctx context.Context) (*MemberOutcomeReconcileReport, error) { - if r.Outbox == nil || r.Store == nil { - return nil, fmt.Errorf("member outcome reconciler: outbox and store are required") + if r.Outbox == nil || r.Store == nil || r.ValidatorID == "" { + return nil, fmt.Errorf("member outcome reconciler: outbox, store and validator id are required") } entries, err := r.Outbox.List() if err != nil { @@ -160,6 +163,9 @@ func (r *MemberOutcomeReconciler) RunOnce(ctx context.Context) (*MemberOutcomeRe rep.Quarantined++ continue } + if e.Outcome.ReportedBy == "" { + e.Outcome.ReportedBy = r.ValidatorID + } derived, err := r.Store.RecordMemberOutcome(ctx, *e.Outcome) switch { case err == nil: diff --git a/pkg/execution/member_outcome_outbox_test.go b/pkg/execution/member_outcome_outbox_test.go index 8bc1b5be..d4b54e1e 100644 --- a/pkg/execution/member_outcome_outbox_test.go +++ b/pkg/execution/member_outcome_outbox_test.go @@ -27,7 +27,7 @@ func TestARefusedMemberOutcomeIsRecordedOnceTheStoreRecovers(t *testing.T) { t.Fatal(err) } repos := database.NewRepositories(database.NewClientFromDB(db)) - o := &UnifiedOrchestrator{config: &UnifiedOrchestratorConfig{Repos: repos, MemberOutcomes: outbox}} + o := &UnifiedOrchestrator{config: &UnifiedOrchestratorConfig{Repos: repos, MemberOutcomes: outbox, ValidatorID: "validator-test"}} if _, err := db.Exec(`CREATE OR REPLACE FUNCTION f78_refuse() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN RAISE EXCEPTION 'f78: store unavailable'; END $$; CREATE TRIGGER f78_refuse BEFORE INSERT ON intent_member_outcomes FOR EACH ROW EXECUTE FUNCTION f78_refuse();`); err != nil { @@ -52,7 +52,7 @@ func TestARefusedMemberOutcomeIsRecordedOnceTheStoreRecovers(t *testing.T) { } drop() - rep, err := (&MemberOutcomeReconciler{Outbox: outbox, Store: repos.IntentLifecycle, Logf: t.Logf}).RunOnce(ctx) + rep, err := (&MemberOutcomeReconciler{Outbox: outbox, ValidatorID: "validator-test", Store: repos.IntentLifecycle, Logf: t.Logf}).RunOnce(ctx) if err != nil { t.Fatal(err) } @@ -80,11 +80,11 @@ func TestAContradictedMemberOutcomeIsQuarantinedAndATransientOneWaits(t *testing t.Fatal(err) } ctx := context.Background() - rep, err := (&MemberOutcomeReconciler{Outbox: outbox, Store: refusingRecorder{fmt.Errorf("connection refused")}, Logf: t.Logf}).RunOnce(ctx) + rep, err := (&MemberOutcomeReconciler{Outbox: outbox, ValidatorID: "validator-test", Store: refusingRecorder{fmt.Errorf("connection refused")}, Logf: t.Logf}).RunOnce(ctx) if err != nil || rep.Deferred != 1 || rep.Remaining != 1 { t.Fatalf("a transient refusal waits: %+v %v", rep, err) } - rep, err = (&MemberOutcomeReconciler{Outbox: outbox, Store: refusingRecorder{fmt.Errorf("%w: member set differs", database.ErrMemberOutcomeInvalid)}, Logf: t.Logf}).RunOnce(ctx) + rep, err = (&MemberOutcomeReconciler{Outbox: outbox, ValidatorID: "validator-test", Store: refusingRecorder{fmt.Errorf("%w: member set differs", database.ErrMemberOutcomeInvalid)}, Logf: t.Logf}).RunOnce(ctx) if err != nil || rep.Quarantined != 1 || rep.Remaining != 0 { t.Fatalf("a contradiction is quarantined, not retried forever or deleted: %+v %v", rep, err) } @@ -108,3 +108,35 @@ func TestACycleWithoutItsMemberSetIsRefusedBeforeItRuns(t *testing.T) { t.Fatalf("a request with its member set: %v", err) } } + +type capturingRecorder struct{ got []database.MemberOutcome } + +func (r *capturingRecorder) RecordMemberOutcome(_ context.Context, o database.MemberOutcome) (database.DerivedIntentStatus, error) { + r.got = append(r.got, o) + return database.DerivedIntentStatus{}, nil +} + +// RB4-F58: an outcome names the validator that reported it. An outbox entry written before outcomes did is this +// validator's own report (the outbox is local), so the reconciler names it; a reconciler that does not know which +// validator it is refuses to run rather than record reports under no name. +func TestAnOutboxEntryFromBeforeReportersIsRecordedAsThisValidators(t *testing.T) { + outbox, err := NewFileMemberOutcomeOutbox(filepath.Join(t.TempDir(), "outbox")) + if err != nil { + t.Fatal(err) + } + if err := outbox.Put(database.MemberOutcome{IntentID: "i", ChainID: 84532, CycleID: "c", MemberChains: []int64{84532}, Legs: 1}); err != nil { + t.Fatal(err) + } + ctx := context.Background() + if _, err := (&MemberOutcomeReconciler{Outbox: outbox, Store: &capturingRecorder{}, Logf: t.Logf}).RunOnce(ctx); err == nil { + t.Fatal("a reconciler without its validator id ran") + } + rec := &capturingRecorder{} + rep, err := (&MemberOutcomeReconciler{Outbox: outbox, ValidatorID: "validator-6", Store: rec, Logf: t.Logf}).RunOnce(ctx) + if err != nil || rep.Recorded != 1 { + t.Fatalf("replay: %+v %v", rep, err) + } + if len(rec.got) != 1 || rec.got[0].ReportedBy != "validator-6" { + t.Fatalf("the entry was recorded under %+v, want validator-6", rec.got) + } +} diff --git a/pkg/execution/member_repair_runner.go b/pkg/execution/member_repair_runner.go new file mode 100644 index 00000000..3d214f0d --- /dev/null +++ b/pkg/execution/member_repair_runner.go @@ -0,0 +1,478 @@ +// Copyright 2026 Certen Protocol + +package execution + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "sort" + "strconv" + "strings" + "time" + + chain "github.com/certen/independant-validator/pkg/chain/strategy" + "github.com/certen/independant-validator/pkg/consensus" + "github.com/certen/independant-validator/pkg/database" + "github.com/lib/pq" +) + +// Re-driving one chain member's proof cycle, on request (RB4-F55 repair; DESIGN_RB4_F55_repair_000ac79a.md). +// +// `validator repair member-proof-cycle` writes a request into this validator's data directory; the running +// validator - which holds the orchestrator, its keys, its peers and the committed-operation index - picks it up. +// Every precondition is checked first and each failure refused by name. Then the repair is armed +// (consensus.ArmMemberRepair) and the member's intent re-processed from its Directory Network block exactly as +// discovery processes it (IntentDiscovery.ReprocessIntent), so the member's round is re-derived and checked +// against the committed block. Where that round reaches the member, its proof cycle is run on the named settlement +// (Apply) or its snapshot reported (dry run). The outcome is read back from what the proof cycle recorded. Requests +// are moved to done/ and results written to results/; nothing is deleted. + +// MemberRepairRequest is one request, as the repair command writes it. +type MemberRepairRequest struct { + ID string `json:"id"` + IntentID string `json:"intent_id"` + ChainID int64 `json:"chain_id"` + SettlementTx string `json:"settlement_tx"` + Apply bool `json:"apply"` + RequestedAt time.Time `json:"requested_at"` +} + +// MemberRepairCheck is one precondition and whether it held. +type MemberRepairCheck struct { + Name string `json:"name"` + Passed bool `json:"passed"` + Detail string `json:"detail"` +} + +// Outcomes of a repair request. +const ( + MemberRepairRefused = "refused" // a precondition did not hold; nothing ran + MemberRepairReached = "reached" // dry run: the re-derived round reached the member; its snapshot is reported + MemberRepairRepaired = "repaired" // the member's proof cycle ran and recorded it settled and written back + MemberRepairFailed = "failed" // the round did not reach the member, or its proof cycle did not write it back +) + +// MemberRepairResult is what a request came to. +type MemberRepairResult struct { + Request MemberRepairRequest `json:"request"` + Validator string `json:"validator"` + Outcome string `json:"outcome"` + Reason string `json:"reason,omitempty"` + Checks []MemberRepairCheck `json:"checks"` + DNBlock uint64 `json:"dn_block,omitempty"` + Snapshot *consensus.MemberRepairSnapshot `json:"snapshot,omitempty"` + Before map[string]any `json:"member_before,omitempty"` + After map[string]any `json:"member_after,omitempty"` + Intent map[string]any `json:"intent_after,omitempty"` + WriteBack *database.MemberWriteBack `json:"write_back,omitempty"` + ProofIDs []string `json:"proof_ids,omitempty"` + Correction []string `json:"corrections,omitempty"` + StartedAt time.Time `json:"started_at"` + FinishedAt time.Time `json:"finished_at"` +} + +// MemberRepairRunner serves repair requests inside the running validator. +type MemberRepairRunner struct { + Dir string // /member_repairs + ValidatorID string + DB *sql.DB + Lifecycle *database.IntentLifecycleRepository + Outbox MemberOutcomeOutbox + // Observe reads the settlement from its chain, as Phase 7 does. + Observe func(ctx context.Context, chainID int64, tx string) (*chain.ObservationResult, error) + // Arm names the member for its round (consensus.BFTValidator.ArmMemberRepair). + Arm func(consensus.MemberRepair) (<-chan consensus.MemberRepairReach, func()) + // Reprocess re-runs the intent's round (intent.IntentDiscovery.ReprocessIntent). + Reprocess func(ctx context.Context, dnBlock uint64, accumTxHash, intentID string) error + Interval time.Duration // zero: ten seconds + // OutcomeWait bounds the wait for the proof cycle's recorded outcome; zero: twelve minutes (a cycle runs + // under a ten-minute limit). + OutcomeWait time.Duration + Logf func(string, ...interface{}) +} + +// MemberRepairDir is where a validator's repair requests and results live. +func MemberRepairDir(dataDir string) string { return filepath.Join(dataDir, "member_repairs") } + +func (r *MemberRepairRunner) logf(format string, args ...interface{}) { + if r.Logf != nil { + r.Logf(format, args...) + } +} + +func (r *MemberRepairRunner) ready() error { + var missing []string + for _, m := range []struct { + absent bool + name string + }{ + {r.Dir == "", "directory"}, {r.ValidatorID == "", "validator id"}, {r.DB == nil, "database"}, + {r.Lifecycle == nil, "lifecycle repository"}, {r.Outbox == nil, "member outcome outbox"}, {r.Observe == nil, "chain observer"}, + {r.Arm == nil, "consensus"}, {r.Reprocess == nil, "discovery"}, + } { + if m.absent { + missing = append(missing, m.name) + } + } + if len(missing) > 0 { + return fmt.Errorf("member repair runner: missing %s", strings.Join(missing, ", ")) + } + for _, d := range []string{"requests", "done", "results"} { + if err := os.MkdirAll(filepath.Join(r.Dir, d), 0o700); err != nil { + return fmt.Errorf("member repair runner: %w", err) + } + } + return nil +} + +// Start serves requests until ctx ends. A runner that is not fully wired refuses to start, by name. +func (r *MemberRepairRunner) Start(ctx context.Context) error { + if err := r.ready(); err != nil { + return err + } + interval := r.Interval + if interval <= 0 { + interval = 10 * time.Second + } + go func() { + t := time.NewTicker(interval) + defer t.Stop() + for { + if err := r.RunOnce(ctx); err != nil { + r.logf("❌ [MEMBER-REPAIR] %v", err) + } + select { + case <-ctx.Done(): + return + case <-t.C: + } + } + }() + return nil +} + +// RunOnce serves every pending request, oldest first. +func (r *MemberRepairRunner) RunOnce(ctx context.Context) error { + if err := r.ready(); err != nil { + return err + } + entries, err := os.ReadDir(filepath.Join(r.Dir, "requests")) + if err != nil { + return fmt.Errorf("list repair requests: %w", err) + } + var names []string + for _, e := range entries { + if !e.IsDir() && strings.HasSuffix(e.Name(), ".json") { + names = append(names, e.Name()) + } + } + sort.Strings(names) + for _, name := range names { + src := filepath.Join(r.Dir, "requests", name) + raw, err := os.ReadFile(src) + if err != nil { + return fmt.Errorf("read repair request %s: %w", name, err) + } + var req MemberRepairRequest + if err := json.Unmarshal(raw, &req); err != nil || req.ID == "" || req.ID+".json" != name { + res := &MemberRepairResult{Request: req, Validator: r.ValidatorID, Outcome: MemberRepairRefused, + Reason: fmt.Sprintf("the request %s cannot be read as a repair request named by its id (%v)", name, err), + StartedAt: time.Now().UTC()} + req.ID = strings.TrimSuffix(name, ".json") + res.Request.ID = req.ID + if err := r.finish(name, res); err != nil { + return err + } + continue + } + res := r.Serve(ctx, req) + if err := r.finish(name, res); err != nil { + return err + } + } + return nil +} + +// finish writes the result and moves the request to done/. +func (r *MemberRepairRunner) finish(name string, res *MemberRepairResult) error { + res.FinishedAt = time.Now().UTC() + blob, err := json.MarshalIndent(res, "", " ") + if err != nil { + return fmt.Errorf("encode repair result %s: %w", res.Request.ID, err) + } + dst := filepath.Join(r.Dir, "results", res.Request.ID+".json") + tmp := dst + ".tmp" + if err := os.WriteFile(tmp, blob, 0o600); err != nil { + return fmt.Errorf("write repair result %s: %w", res.Request.ID, err) + } + if err := os.Rename(tmp, dst); err != nil { + return fmt.Errorf("write repair result %s: %w", res.Request.ID, err) + } + if err := os.Rename(filepath.Join(r.Dir, "requests", name), filepath.Join(r.Dir, "done", name)); err != nil { + return fmt.Errorf("move repair request %s to done: %w", name, err) + } + r.logf("🔧 [MEMBER-REPAIR] request %s (intent %s member %d, apply %v): %s %s", + res.Request.ID, res.Request.IntentID, res.Request.ChainID, res.Request.Apply, res.Outcome, res.Reason) + return nil +} + +type repairFacts struct { + dnBlock uint64 + accumTxHash string + before map[string]any + cycleBefore string + // dbStart is the database clock when the repair began: what it produced is read back against the clock that + // stamped it, never this host's. + dbStart time.Time +} + +// Serve runs one request. +func (r *MemberRepairRunner) Serve(ctx context.Context, req MemberRepairRequest) *MemberRepairResult { + res := &MemberRepairResult{Request: req, Validator: r.ValidatorID, StartedAt: time.Now().UTC()} + facts, ok := r.preconditions(ctx, req, res) + if !ok { + res.Outcome = MemberRepairRefused + var failed []string + for _, c := range res.Checks { + if !c.Passed { + failed = append(failed, c.Name+": "+c.Detail) + } + } + res.Reason = strings.Join(failed, "; ") + return res + } + res.DNBlock, res.Before = facts.dnBlock, facts.before + if err := r.DB.QueryRowContext(ctx, `SELECT now()`).Scan(&facts.dbStart); err != nil { + res.Outcome, res.Reason = MemberRepairRefused, "the database clock could not be read: "+err.Error() + return res + } + + reach, disarm := r.Arm(consensus.MemberRepair{IntentID: req.IntentID, ChainID: req.ChainID, SettlementTx: req.SettlementTx, Apply: req.Apply}) + defer disarm() + reprocessErr := r.Reprocess(ctx, facts.dnBlock, facts.accumTxHash, req.IntentID) + var reached consensus.MemberRepairReach + select { + case reached = <-reach: + default: + res.Outcome = MemberRepairFailed + res.Reason = fmt.Sprintf("the re-derived round of intent %s did not reach member %d", req.IntentID, req.ChainID) + if reprocessErr != nil { + res.Reason += ": " + reprocessErr.Error() + } + return res + } + res.Snapshot = &reached.Snapshot + if reached.Err != nil { + res.Outcome, res.Reason = MemberRepairFailed, reached.Err.Error() + return res + } + if !req.Apply { + res.Outcome = MemberRepairReached + return res + } + r.readOutcome(ctx, req, facts, res) + return res +} + +func (r *MemberRepairRunner) check(res *MemberRepairResult, name string, passed bool, detail string) bool { + res.Checks = append(res.Checks, MemberRepairCheck{Name: name, Passed: passed, Detail: detail}) + return passed +} + +// preconditions checks everything the repair relies on, in order, and stops at the first that does not hold. +func (r *MemberRepairRunner) preconditions(ctx context.Context, req MemberRepairRequest, res *MemberRepairResult) (*repairFacts, bool) { + if !r.check(res, "request", req.IntentID != "" && req.ChainID != 0 && req.SettlementTx != "", + "intent, chain and settlement transaction are required") { + return nil, false + } + + // 1. The intent exists and this chain is one of its members. + var blockHeight sql.NullInt64 + var accum sql.NullString + var members pq.Int64Array + err := r.DB.QueryRowContext(ctx, `SELECT block_height, accum_tx_hash, member_chains FROM intent_lifecycle WHERE intent_id = $1`, + req.IntentID).Scan(&blockHeight, &accum, &members) + if errors.Is(err, sql.ErrNoRows) { + r.check(res, "intent", false, fmt.Sprintf("intent %s has no lifecycle record", req.IntentID)) + return nil, false + } + if err != nil { + r.check(res, "intent", false, fmt.Sprintf("the lifecycle record of %s could not be read: %v", req.IntentID, err)) + return nil, false + } + inSet := false + for _, c := range members { + inSet = inSet || c == req.ChainID + } + if !r.check(res, "intent", blockHeight.Valid && blockHeight.Int64 > 0 && accum.String != "" && inSet, + fmt.Sprintf("DN block %d, Accumulate transaction %s, members %v", blockHeight.Int64, accum.String, []int64(members))) { + return nil, false + } + facts := &repairFacts{dnBlock: uint64(blockHeight.Int64), accumTxHash: accum.String} + + // 2. The member's proof cycle is recorded failed: a member written back is never re-driven. + var settlement, proofCycle string + var settlementTx, writeBackTx, cycleID, reason sql.NullString + var recordedAt time.Time + err = r.DB.QueryRowContext(ctx, `SELECT settlement, proof_cycle, settlement_tx, write_back_tx, cycle_id, reason, recorded_at + FROM intent_member_outcomes WHERE intent_id = $1 AND chain_id = $2`, req.IntentID, req.ChainID). + Scan(&settlement, &proofCycle, &settlementTx, &writeBackTx, &cycleID, &reason, &recordedAt) + if err != nil { + r.check(res, "member outcome", false, fmt.Sprintf("no recorded outcome for member %d could be read: %v", req.ChainID, err)) + return nil, false + } + facts.before = map[string]any{"settlement": settlement, "proof_cycle": proofCycle, "settlement_tx": settlementTx.String, + "write_back_tx": writeBackTx.String, "cycle_id": cycleID.String, "reason": reason.String, "recorded_at": recordedAt.UTC()} + facts.cycleBefore = cycleID.String + if !r.check(res, "member outcome", proofCycle == string(database.MemberProofCycleFailed), + fmt.Sprintf("recorded %s / %s by cycle %s", settlement, proofCycle, cycleID.String)) { + return nil, false + } + + // 3. No write-back of the member is registered, or may be on Accumulate. + if err := r.Lifecycle.MemberWriteBackAllowed(ctx, req.IntentID, req.ChainID); !r.check(res, "write-back register", err == nil, errText(err, "none registered")) { + return nil, false + } + + // 4. No report of this member waits in the outbox: replayed after the repair it would contradict it. + entries, err := r.Outbox.List() + pending := 0 + for _, e := range entries { + if e.Outcome != nil && e.Outcome.IntentID == req.IntentID && e.Outcome.ChainID == req.ChainID { + pending++ + } + } + if !r.check(res, "member outcome outbox", err == nil && pending == 0, fmt.Sprintf("%d pending report(s) of this member (%v)", pending, err)) { + return nil, false + } + + // 5. This validator observed the settlement: its recorded observation is the one the proof cycle adopts (RB4-F55). + var observer sql.NullString + err = r.DB.QueryRowContext(ctx, `SELECT observer_validator_id FROM chain_execution_results WHERE chain_id = $1 AND lower(tx_hash) = lower($2)`, + strconv.FormatInt(req.ChainID, 10), req.SettlementTx).Scan(&observer) + switch { + case errors.Is(err, sql.ErrNoRows): + r.check(res, "recorded observation", false, fmt.Sprintf("no observation of %s is recorded; run the repair on the validator that settled it", req.SettlementTx)) + return nil, false + case err != nil: + r.check(res, "recorded observation", false, err.Error()) + return nil, false + } + if !r.check(res, "recorded observation", observer.String == r.ValidatorID, + fmt.Sprintf("observed by %s; this is %s", observer.String, r.ValidatorID)) { + return nil, false + } + + // 6. The settlement is final on its chain and executed. + obs, err := r.Observe(ctx, req.ChainID, req.SettlementTx) + if err != nil { + r.check(res, "settlement on chain", false, err.Error()) + return nil, false + } + if !r.check(res, "settlement on chain", obs != nil && obs.IsFinalized && obs.Status == 1, + fmt.Sprintf("finalized %v, status %d, block %d", obs != nil && obs.IsFinalized, statusOf(obs), blockOf(obs))) { + return nil, false + } + return facts, true +} + +func statusOf(o *chain.ObservationResult) uint8 { + if o == nil { + return 0 + } + return o.Status +} + +func blockOf(o *chain.ObservationResult) uint64 { + if o == nil { + return 0 + } + return o.BlockNumber +} + +func errText(err error, ok string) string { + if err != nil { + return err.Error() + } + return ok +} + +// readOutcome waits for the member's proof cycle to record its outcome, and reads back what it recorded. +func (r *MemberRepairRunner) readOutcome(ctx context.Context, req MemberRepairRequest, facts *repairFacts, res *MemberRepairResult) { + wait := r.OutcomeWait + if wait <= 0 { + wait = 12 * time.Minute + } + deadline := time.Now().Add(wait) + for { + var settlement, proofCycle string + var writeBackTx, cycleID, reason sql.NullString + err := r.DB.QueryRowContext(ctx, `SELECT settlement, proof_cycle, write_back_tx, cycle_id, reason + FROM intent_member_outcomes WHERE intent_id = $1 AND chain_id = $2`, req.IntentID, req.ChainID). + Scan(&settlement, &proofCycle, &writeBackTx, &cycleID, &reason) + if err == nil && cycleID.String != facts.cycleBefore { + res.After = map[string]any{"settlement": settlement, "proof_cycle": proofCycle, "write_back_tx": writeBackTx.String, + "cycle_id": cycleID.String, "reason": reason.String} + r.collect(ctx, req, facts, res) + if settlement == string(database.MemberSettlementSettled) && proofCycle == string(database.MemberProofCycleWritten) { + res.Outcome = MemberRepairRepaired + } else { + res.Outcome = MemberRepairFailed + res.Reason = fmt.Sprintf("the proof cycle recorded %s / %s: %s", settlement, proofCycle, reason.String) + } + return + } + if time.Now().After(deadline) { + r.collect(ctx, req, facts, res) + res.Outcome = MemberRepairFailed + res.Reason = fmt.Sprintf("the member's proof cycle recorded no outcome within %s (last read: %v); see the validator log", wait, err) + return + } + select { + case <-ctx.Done(): + res.Outcome, res.Reason = MemberRepairFailed, "stopped before the proof cycle recorded its outcome: "+ctx.Err().Error() + return + case <-time.After(2 * time.Second): + } + } +} + +// collect reads back the intent, the write-back register, the proofs and the corrections the repair produced. +func (r *MemberRepairRunner) collect(ctx context.Context, req MemberRepairRequest, facts *repairFacts, res *MemberRepairResult) { + var status string + var errMsg, class sql.NullString + if err := r.DB.QueryRowContext(ctx, `SELECT status, error_message, failure_class FROM intent_lifecycle WHERE intent_id = $1`, req.IntentID). + Scan(&status, &errMsg, &class); err == nil { + res.Intent = map[string]any{"status": status, "error_message": errMsg.String, "failure_class": class.String} + } + if w, err := r.Lifecycle.MemberWriteBackOf(ctx, req.IntentID, req.ChainID); err == nil { + res.WriteBack = w + } + if rows, err := r.DB.QueryContext(ctx, `SELECT proof_id::text FROM proof_artifacts + WHERE lower(accum_tx_hash) = lower($1) AND created_at >= $2 ORDER BY created_at`, facts.accumTxHash, facts.dbStart); err == nil { + for rows.Next() { + var id string + if rows.Scan(&id) == nil { + res.ProofIDs = append(res.ProofIDs, id) + } + } + rows.Close() + } + if rows, err := r.DB.QueryContext(ctx, `SELECT correction_id::text FROM evidence_corrections + WHERE ((record_type = 'intent_member_outcome' AND record_id = $1) OR (record_type = 'intent_lifecycle' AND record_id = $2)) + AND corrected_at >= $3 ORDER BY corrected_at`, + fmt.Sprintf("%s/%d", req.IntentID, req.ChainID), req.IntentID, facts.dbStart); err == nil { + for rows.Next() { + var id string + if rows.Scan(&id) == nil { + res.Correction = append(res.Correction, id) + } + } + rows.Close() + } +} diff --git a/pkg/execution/member_repair_runner_test.go b/pkg/execution/member_repair_runner_test.go new file mode 100644 index 00000000..2f58b83d --- /dev/null +++ b/pkg/execution/member_repair_runner_test.go @@ -0,0 +1,249 @@ +// Copyright 2026 Certen Protocol + +package execution + +import ( + "context" + "database/sql" + "encoding/json" + "errors" + "fmt" + "os" + "path/filepath" + "strings" + "testing" + "time" + + chain "github.com/certen/independant-validator/pkg/chain/strategy" + "github.com/certen/independant-validator/pkg/consensus" + "github.com/certen/independant-validator/pkg/database" +) + +// RB4-F55 repair runner (DESIGN_RB4_F55_repair_000ac79a.md): a member's proof cycle is re-driven only when every +// precondition holds, each refused by name; a dry run reports the re-derived round's snapshot and runs nothing; +// an applied repair reads back what the proof cycle recorded; requests are kept in done/, results in results/. + +type repairScene struct { + db *sql.DB + repos *database.Repositories + intentID string + accum string + tx string + outbox *FileMemberOutcomeOutbox + armed []consensus.MemberRepair + reproc []string + status uint8 + // onReprocess plays the round (and, when it reaches the member, what its proof cycle records). + onReprocess func(reach chan consensus.MemberRepairReach, apply bool) error + reach chan consensus.MemberRepairReach +} + +func newRepairScene(t *testing.T) *repairScene { + t.Helper() + db := s1OpenDB(t) + ctx := context.Background() + repos := database.NewRepositories(database.NewClientFromDB(db)) + run := fmt.Sprintf("%d", time.Now().UnixNano()) + s := &repairScene{db: db, repos: repos, intentID: "f55-repair-" + run, accum: fmt.Sprintf("%064s", run), + tx: "0x" + fmt.Sprintf("%064s", "c409"+run), status: 1} + if _, err := db.Exec(`INSERT INTO intent_lifecycle (intent_id, accum_tx_hash, status, block_height, created_at, updated_at) + VALUES ($1, $2, 'settling', 10007772, now(), now())`, s.intentID, s.accum); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + db.Exec(`DELETE FROM intent_member_outcomes WHERE intent_id = $1`, s.intentID) + db.Exec(`DELETE FROM member_write_backs WHERE intent_id = $1`, s.intentID) + db.Exec(`DELETE FROM chain_execution_results WHERE tx_hash = $1`, s.tx) + db.Exec(`DELETE FROM intent_lifecycle WHERE intent_id = $1`, s.intentID) + }) + members := []int64{84532, 421614} + for _, o := range []database.MemberOutcome{ + {IntentID: s.intentID, ChainID: 421614, MemberChains: members, Legs: 1, Settlement: database.MemberSettlementSettled, + ProofCycle: database.MemberProofCycleWritten, SettlementTx: "0xarb", WriteBackTx: "acc://wb-arb", CycleID: "cycle-arb", ReportedBy: "validator-1"}, + {IntentID: s.intentID, ChainID: 84532, MemberChains: members, Legs: 1, Settlement: database.MemberSettlementUnobserved, + ProofCycle: database.MemberProofCycleFailed, CycleID: "cycle-broken", Reason: "phase 7 failed: duplicate key", ReportedBy: "validator-6"}, + } { + if _, err := repos.IntentLifecycle.RecordMemberOutcome(ctx, o); err != nil { + t.Fatal(err) + } + } + bn := int64(47451715) + now := time.Now().UTC() + if _, err := repos.Unified.CreateChainExecutionResult(ctx, &database.NewChainExecutionResult{CycleID: "cycle-broken", + ChainPlatform: "evm", ChainID: "84532", NetworkName: "base-sepolia", TxHash: s.tx, BlockNumber: &bn, BlockHash: "0xb33f", + BlockTimestamp: &now, Status: 1, IsFinalized: true, RawReceipt: []byte("null"), Logs: []byte("[]"), PlatformData: []byte("{}"), + ObserverValidatorID: "validator-6", SubmittedAt: &now}); err != nil { + t.Fatal(err) + } + outbox, err := NewFileMemberOutcomeOutbox(filepath.Join(t.TempDir(), "outbox")) + if err != nil { + t.Fatal(err) + } + s.outbox = outbox + return s +} + +func (s *repairScene) runner(t *testing.T, validator string) *MemberRepairRunner { + return &MemberRepairRunner{ + Dir: filepath.Join(t.TempDir(), "member_repairs"), ValidatorID: validator, DB: s.db, Lifecycle: s.repos.IntentLifecycle, Outbox: s.outbox, + Observe: func(_ context.Context, chainID int64, tx string) (*chain.ObservationResult, error) { + return &chain.ObservationResult{TxHash: tx, BlockNumber: 47451715, Status: s.status, IsFinalized: true}, nil + }, + Arm: func(m consensus.MemberRepair) (<-chan consensus.MemberRepairReach, func()) { + s.armed = append(s.armed, m) + s.reach = make(chan consensus.MemberRepairReach, 1) + return s.reach, func() {} + }, + Reprocess: func(_ context.Context, dnBlock uint64, accum, intentID string) error { + s.reproc = append(s.reproc, fmt.Sprintf("%d/%s/%s", dnBlock, accum, intentID)) + if s.onReprocess == nil { + return nil + } + return s.onReprocess(s.reach, s.armed[len(s.armed)-1].Apply) + }, + OutcomeWait: 5 * time.Second, Logf: t.Logf, + } +} + +func (s *repairScene) request(apply bool) MemberRepairRequest { + return MemberRepairRequest{ID: "req-" + s.intentID, IntentID: s.intentID, ChainID: 84532, SettlementTx: s.tx, Apply: apply, RequestedAt: time.Now().UTC()} +} + +var snapshot = consensus.MemberRepairSnapshot{BundleID: "0x0a", GovernanceRoot: "0x0c", OperationCommitment: "0x0b", Lane: "on_demand", GovernanceLevels: true} + +func TestARepairIsRefusedUnlessEveryPreconditionHolds(t *testing.T) { + ctx := context.Background() + for name, tc := range map[string]struct { + mutate func(s *repairScene, req *MemberRepairRequest) + validator string + check string + }{ + "unknown intent": {func(_ *repairScene, r *MemberRepairRequest) { r.IntentID = "no-such-intent" }, "validator-6", "intent"}, + "chain not a member": {func(_ *repairScene, r *MemberRepairRequest) { r.ChainID = 11155111 }, "validator-6", "intent"}, + "member written back": {func(_ *repairScene, r *MemberRepairRequest) { r.ChainID = 421614; r.SettlementTx = "0xarb" }, "validator-6", "member outcome"}, + "observed by another": {func(*repairScene, *MemberRepairRequest) {}, "validator-3", "recorded observation"}, + "no recorded observation": {func(_ *repairScene, r *MemberRepairRequest) { r.SettlementTx = "0x" + strings.Repeat("99", 32) }, "validator-6", "recorded observation"}, + "settlement not executed": {func(s *repairScene, _ *MemberRepairRequest) { s.status = 2 }, "validator-6", "settlement on chain"}, + "a stale report waits": {func(s *repairScene, r *MemberRepairRequest) { + if err := s.outbox.Put(database.MemberOutcome{IntentID: r.IntentID, ChainID: 84532, CycleID: "cycle-broken", MemberChains: []int64{84532, 421614}, Legs: 1}); err != nil { + t.Fatal(err) + } + }, "validator-6", "member outcome outbox"}, + "a write-back registered": {func(s *repairScene, r *MemberRepairRequest) { + if err := s.repos.IntentLifecycle.ClaimMemberWriteBack(ctx, r.IntentID, 84532, "cycle-x", "validator-6"); err != nil { + t.Fatal(err) + } + }, "validator-6", "write-back register"}, + } { + t.Run(name, func(t *testing.T) { + s := newRepairScene(t) + req := s.request(true) + tc.mutate(s, &req) + res := s.runner(t, tc.validator).Serve(ctx, req) + if res.Outcome != MemberRepairRefused || !strings.Contains(res.Reason, tc.check) { + t.Fatalf("want refused on %q, got %s: %s", tc.check, res.Outcome, res.Reason) + } + if len(s.armed) != 0 || len(s.reproc) != 0 { + t.Fatal("a refused repair armed or re-processed something") + } + }) + } +} + +func TestADryRunRepairReachesTheMemberAndRunsNothing(t *testing.T) { + s := newRepairScene(t) + s.onReprocess = func(reach chan consensus.MemberRepairReach, apply bool) error { + reach <- consensus.MemberRepairReach{Snapshot: snapshot} + return nil + } + res := s.runner(t, "validator-6").Serve(context.Background(), s.request(false)) + if res.Outcome != MemberRepairReached || res.Snapshot == nil || res.Snapshot.BundleID != "0x0a" { + t.Fatalf("a dry run: %s %s %+v", res.Outcome, res.Reason, res.Snapshot) + } + if len(s.armed) != 1 || s.armed[0].Apply || s.armed[0].SettlementTx != s.tx { + t.Fatalf("armed %+v", s.armed) + } + if want := fmt.Sprintf("10007772/%s/%s", s.accum, s.intentID); len(s.reproc) != 1 || s.reproc[0] != want { + t.Fatalf("re-processed %v, want %s (the intent's DN block and transaction)", s.reproc, want) + } + var cycle string + s.db.QueryRow(`SELECT cycle_id FROM intent_member_outcomes WHERE intent_id = $1 AND chain_id = 84532`, s.intentID).Scan(&cycle) + if cycle != "cycle-broken" { + t.Fatalf("a dry run changed the member: %s", cycle) + } +} + +func TestAnAppliedRepairReadsBackWhatTheProofCycleRecorded(t *testing.T) { + s := newRepairScene(t) + ctx := context.Background() + s.onReprocess = func(reach chan consensus.MemberRepairReach, apply bool) error { + // What the member's proof cycle records when it runs: its write-back registered, its outcome settled and written. + lc := s.repos.IntentLifecycle + if err := lc.ClaimMemberWriteBack(ctx, s.intentID, 84532, "cycle-repair", "validator-6"); err != nil { + return err + } + if err := lc.RecordMemberWriteBack(ctx, s.intentID, 84532, "cycle-repair", "acc://wb-base"); err != nil { + return err + } + if _, err := lc.RecordMemberOutcome(ctx, database.MemberOutcome{IntentID: s.intentID, ChainID: 84532, MemberChains: []int64{84532, 421614}, + Legs: 1, Settlement: database.MemberSettlementSettled, ProofCycle: database.MemberProofCycleWritten, SettlementTx: s.tx, + WriteBackTx: "acc://wb-base", CycleID: "cycle-repair", ReportedBy: "validator-6"}); err != nil { + return err + } + reach <- consensus.MemberRepairReach{Snapshot: snapshot, Started: true} + return nil + } + res := s.runner(t, "validator-6").Serve(ctx, s.request(true)) + if res.Outcome != MemberRepairRepaired { + t.Fatalf("an applied repair: %s %s", res.Outcome, res.Reason) + } + if res.Before["proof_cycle"] != "failed" || res.After["proof_cycle"] != "written" || res.Intent["status"] != "complete" { + t.Fatalf("before %v after %v intent %v", res.Before, res.After, res.Intent) + } + if res.WriteBack == nil || res.WriteBack.State != database.MemberWriteBackWritten || len(res.Correction) != 2 { + t.Fatalf("write-back %+v, corrections %v; want the write-back registered and two corrections (member and intent)", res.WriteBack, res.Correction) + } +} + +func TestARoundThatNeverReachesTheMemberIsAFailureNamingWhy(t *testing.T) { + s := newRepairScene(t) + s.onReprocess = func(chan consensus.MemberRepairReach, bool) error { return errors.New("committed as another block") } + res := s.runner(t, "validator-6").Serve(context.Background(), s.request(true)) + if res.Outcome != MemberRepairFailed || !strings.Contains(res.Reason, "did not reach") || !strings.Contains(res.Reason, "committed as another block") { + t.Fatalf("got %s: %s", res.Outcome, res.Reason) + } +} + +func TestARepairRequestIsKeptAndAnswered(t *testing.T) { + s := newRepairScene(t) + s.onReprocess = func(reach chan consensus.MemberRepairReach, apply bool) error { + reach <- consensus.MemberRepairReach{Snapshot: snapshot} + return nil + } + r := s.runner(t, "validator-6") + if err := r.ready(); err != nil { + t.Fatal(err) + } + req := s.request(false) + blob, _ := json.Marshal(req) + if err := os.WriteFile(filepath.Join(r.Dir, "requests", req.ID+".json"), blob, 0o600); err != nil { + t.Fatal(err) + } + if err := r.RunOnce(context.Background()); err != nil { + t.Fatal(err) + } + if _, err := os.Stat(filepath.Join(r.Dir, "done", req.ID+".json")); err != nil { + t.Fatalf("the request was not kept in done/: %v", err) + } + raw, err := os.ReadFile(filepath.Join(r.Dir, "results", req.ID+".json")) + if err != nil { + t.Fatal(err) + } + var res MemberRepairResult + if err := json.Unmarshal(raw, &res); err != nil || res.Outcome != MemberRepairReached || res.Request.ID != req.ID { + t.Fatalf("result %+v (%v)", res, err) + } + if err := (&MemberRepairRunner{Dir: r.Dir}).Start(context.Background()); err == nil { + t.Fatal("a runner that is not wired started") + } +} diff --git a/pkg/execution/result_chain_links_test.go b/pkg/execution/result_chain_links_test.go index fc04bd94..83be8048 100644 --- a/pkg/execution/result_chain_links_test.go +++ b/pkg/execution/result_chain_links_test.go @@ -38,6 +38,9 @@ func TestANonSettlementsChainLinkIsPersistedAndContinued(t *testing.T) { t.Fatal(err) } for i := 0; i < 2; i++ { + // Two members on one chain: a member is written back once (RB4-F59). + f := f + f.IntentID = fmt.Sprintf("%s-%s-%d", f.IntentID, validator, i) rec := &NonSettlementRecord{Facts: f, Cause: claim.Cause, MemberChains: []int64{odChain}, MemberLegs: 1} cycle := nonSettlementCycle(rec, claim) cycle.CycleID = fmt.Sprintf("%s-%d", cycle.CycleID, i) diff --git a/pkg/execution/stage1_lifecycle_test.go b/pkg/execution/stage1_lifecycle_test.go index 4504a2a1..4e2e5697 100644 --- a/pkg/execution/stage1_lifecycle_test.go +++ b/pkg/execution/stage1_lifecycle_test.go @@ -40,6 +40,7 @@ func s1Orchestrator(t *testing.T, db *sql.DB) *UnifiedOrchestrator { t.Helper() return &UnifiedOrchestrator{ config: &UnifiedOrchestratorConfig{ + ValidatorID: "validator-test", Repos: &database.Repositories{ IntentLifecycle: database.NewIntentLifecycleRepository(database.NewClientFromDB(db)), }, diff --git a/pkg/execution/unified_orchestrator.go b/pkg/execution/unified_orchestrator.go index ed52b942..176c630d 100644 --- a/pkg/execution/unified_orchestrator.go +++ b/pkg/execution/unified_orchestrator.go @@ -470,6 +470,20 @@ func (o *UnifiedOrchestrator) StartProofCycle(ctx context.Context, req *UnifiedP req.CycleID = uuid.New().String() } + // RB4-F59: a member written back (or whose write-back has an unknown outcome) is not proved again - a second + // cycle would store a second proof bundle and write a second entry. It records nothing: the member's outcome + // is not this cycle's to state. + if chainID, _, _, err := memberSetOf(req); err == nil { + register, rErr := o.memberWriteBackRegister() + if rErr != nil { + return nil, rErr + } + if err := register.MemberWriteBackAllowed(ctx, req.IntentID, chainID); err != nil { + fmt.Printf("🛑 [Phase 9] intent %s member %d: proof cycle %s not started: %v\n", req.IntentID, chainID, req.CycleID, err) + return nil, fmt.Errorf("proof cycle %s not started: %w", req.CycleID, err) + } + } + // Create result result := &UnifiedProofCycleResult{ CycleID: req.CycleID, @@ -740,7 +754,7 @@ func (o *UnifiedOrchestrator) recordMemberOutcome( out := database.MemberOutcome{ IntentID: req.IntentID, ChainID: chainID, MemberChains: chains, Legs: legs, Settlement: settlement, ProofCycle: proofCycle, CycleID: req.CycleID, Reason: reason, - EffectsProven: cycleEffectsProven(cycle), + EffectsProven: cycleEffectsProven(cycle), ReportedBy: o.config.ValidatorID, } if result != nil { out.WriteBackTx = result.WriteBackTxHash @@ -829,6 +843,11 @@ func (o *UnifiedOrchestrator) recordPhaseFailure(ctx context.Context, cycle *act reason := fmt.Sprintf("phase %d failed: %v", phase, err) cycle.Result.Error = reason cycle.Result.FailPhase = phase + if duplicateWriteBack(err) { + // RB4-F59: the member's write-back is on Accumulate, or may be. Its outcome is not this cycle's to state. + fmt.Printf("🛑 [LIFECYCLE] cycle %s: %s - no member outcome recorded for this cycle\n", cycle.CycleID, reason) + return + } if rErr := o.recordMemberOutcome(ctx, cycle, observedSettlement(cycle.Result.ObservationResults), database.MemberProofCycleFailed, reason); rErr != nil { fmt.Printf("❌ [LIFECYCLE] cycle %s failed in phase %d and its failure could not be recorded: %v\n", cycle.CycleID, phase, rErr) } @@ -1932,8 +1951,43 @@ const ( WriteBackWritten = "written" WriteBackRefusedQuorumNotMet = "refused_quorum_not_met" WriteBackFailed = "failed" + // RB4-F59: a member is written back once. + WriteBackRefusedAlreadyWritten = "refused_already_written" // the member's outcome is already on Accumulate + WriteBackRefusedUnresolved = "refused_outcome_unknown" // an earlier write-back of it has an unknown outcome + WriteBackUnresolved = "outcome_unknown" // submitted, and whether it reached Accumulate is unknown ) +// ObserveSettlement reads a settlement from its chain with the strategy Phase 7 observes it with (the RB4-F55 +// repair runner checks a settlement is final and executed before re-driving its member). +func (o *UnifiedOrchestrator) ObserveSettlement(ctx context.Context, chainID int64, tx string) (*chain.ObservationResult, error) { + if o.config.Registry == nil { + return nil, fmt.Errorf("no strategy registry is configured") + } + chainStrategy, _, err := o.config.Registry.GetStrategiesForChain(strconv.FormatInt(chainID, 10)) + if err != nil { + return nil, fmt.Errorf("strategies for chain %d: %w", chainID, err) + } + return chainStrategy.ObserveTransaction(ctx, tx) +} + +// memberWriteBackRegister is where a member's write-back is claimed and recorded (RB4-F59). +func (o *UnifiedOrchestrator) memberWriteBackRegister() (*database.IntentLifecycleRepository, error) { + if o.config.Repos == nil || o.config.Repos.IntentLifecycle == nil { + return nil, fmt.Errorf("no write-back register: a member's write-back cannot be claimed") + } + return o.config.Repos.IntentLifecycle, nil +} + +// duplicateWriteBack is true for a cycle refused because its member's write-back is already on Accumulate or has +// an unknown outcome: such a cycle records no outcome of its own - the member's is not this cycle's to state. +func duplicateWriteBack(err error) bool { + return errors.Is(err, database.ErrMemberAlreadyWrittenBack) || errors.Is(err, database.ErrMemberWriteBackUnresolved) || + errors.Is(err, errWriteBackOutcomeUnknown) +} + +// errWriteBackOutcomeUnknown: this cycle submitted its write-back and does not know whether it reached Accumulate. +var errWriteBackOutcomeUnknown = errors.New("write-back submitted; whether it reached Accumulate is unknown") + func (o *UnifiedOrchestrator) executePhase9(ctx context.Context, cycle *activeCycle) (err error) { cycle.Phase = 9 // Any error below is a write-back that did not happen; say so unless a more specific state @@ -2032,10 +2086,46 @@ func (o *UnifiedOrchestrator) executePhase9(ctx context.Context, cycle *activeCy return fmt.Errorf("add signature: %w", err) } + // RB4-F59: claim the member's write-back before submitting it. A member already written back, or claimed by + // a submission whose outcome is unknown, is refused by name. + memberChain, _, _, err := memberSetOf(cycle.Request) + if err != nil { + return fmt.Errorf("write-back cannot name its member: %w", err) + } + register, err := o.memberWriteBackRegister() + if err != nil { + return err + } + if err := register.ClaimMemberWriteBack(ctx, cycle.Request.IntentID, memberChain, cycle.CycleID, o.config.ValidatorID); err != nil { + switch { + case errors.Is(err, database.ErrMemberAlreadyWrittenBack): + cycle.Result.WriteBackState = WriteBackRefusedAlreadyWritten + case errors.Is(err, database.ErrMemberWriteBackUnresolved): + cycle.Result.WriteBackState = WriteBackRefusedUnresolved + } + return fmt.Errorf("write-back refused: %w", err) + } + // Submit transaction to Accumulate receipt, err := o.config.AccumulateClient.SubmitTransaction(writeBackCtx, tx) if err != nil { - return fmt.Errorf("submit to accumulate: %w", err) + if errors.Is(err, ErrWriteBackNotSent) { + // Nothing reached Accumulate: the member may be written back by another cycle. + if rErr := register.ReleaseMemberWriteBack(ctx, cycle.Request.IntentID, memberChain, cycle.CycleID, err.Error()); rErr != nil { + return fmt.Errorf("submit to accumulate: %w; and releasing its claim failed: %v", err, rErr) + } + return fmt.Errorf("submit to accumulate: %w", err) + } + // The submission may have reached Accumulate. Its claim stays, and blocks another write-back of this member + // until whether it did is established. + cycle.Result.WriteBackState = WriteBackUnresolved + return fmt.Errorf("%w (intent %s member %d, cycle %s; its claim is kept): %v", + errWriteBackOutcomeUnknown, cycle.Request.IntentID, memberChain, cycle.CycleID, err) + } + if rErr := register.RecordMemberWriteBack(ctx, cycle.Request.IntentID, memberChain, cycle.CycleID, receipt); rErr != nil { + // Written, and not registered as written: the claim stays, so no other write-back of the member follows. + fmt.Printf("❌ [Phase 9] intent %s member %d: write-back %s is on Accumulate and could not be registered: %v\n", + cycle.Request.IntentID, memberChain, receipt, rErr) } cycle.Result.WriteBackTxHash = receipt diff --git a/pkg/execution/writeback_once_test.go b/pkg/execution/writeback_once_test.go new file mode 100644 index 00000000..af87dae0 --- /dev/null +++ b/pkg/execution/writeback_once_test.go @@ -0,0 +1,245 @@ +// Copyright 2026 Certen Protocol + +package execution + +import ( + "context" + "crypto/ed25519" + "database/sql" + "errors" + "fmt" + "strings" + "testing" + "time" + + attestation "github.com/certen/independant-validator/pkg/attestation/strategy" + chain "github.com/certen/independant-validator/pkg/chain/strategy" + "github.com/certen/independant-validator/pkg/database" + "github.com/certen/independant-validator/pkg/strategy" +) + +// RB4-F59. Phase 9 submitted a write-back for whatever cycle reached it, and the proof bundle was stored before +// it, so a second proof cycle for a member already written back - re-driven, or re-discovered after a restart - +// wrote a second Accumulate entry and stored a second proof artifact for the same member. A member is written back +// once: a registered write-back stops another; a submission whose outcome is unknown stops another until that is +// established; one that failed before anything was sent does not. + +// scriptedSubmitter answers each submission in turn with the scripted error, or a write-back receipt. +type scriptedSubmitter struct { + errs []error + n int +} + +func (s *scriptedSubmitter) SubmitTransaction(_ context.Context, _ *SyntheticTransaction) (string, error) { + i := s.n + s.n++ + if i < len(s.errs) && s.errs[i] != nil { + return "", s.errs[i] + } + return fmt.Sprintf("acc://%064x@results.acme/data", i+1), nil +} + +func (s *scriptedSubmitter) GetTransactionStatus(context.Context, string) (string, error) { + return "", errors.New("the scripted submitter answers no status queries") +} + +type writeBackFixture struct { + db *sql.DB + repos *database.Repositories + o *UnifiedOrchestrator + sub *scriptedSubmitter + intentID string + cycles int +} + +func newWriteBackFixture(t *testing.T, errs ...error) *writeBackFixture { + t.Helper() + db := s1OpenDB(t) + repos := database.NewRepositories(database.NewClientFromDB(db)) + validator := fmt.Sprintf("f59-validator-%d", time.Now().UnixNano()) + _, key, _ := ed25519.GenerateKey(nil) + sub := &scriptedSubmitter{errs: errs} + f := &writeBackFixture{db: db, repos: repos, sub: sub, intentID: "f59-" + validator} + f.o = &UnifiedOrchestrator{ + config: &UnifiedOrchestratorConfig{ValidatorID: validator, UnifiedRepo: repos.Unified, Repos: repos, + ResultsPrincipal: "acc://results.acme/data", Ed25519Key: key, AccumulateClient: sub}, + resultChains: map[string]*ResultHashChain{}, + txBuilder: NewSyntheticTxBuilder("acc://results.acme/data", validator, key), + } + t.Cleanup(func() { + db.Exec(`DELETE FROM result_hash_chain_links WHERE observer_validator_id=$1`, validator) + db.Exec(`DELETE FROM member_write_backs WHERE intent_id=$1`, f.intentID) + }) + return f +} + +// cycle is a new proof cycle, attested by quorum, for the fixture's one member. +func (f *writeBackFixture) cycle(t *testing.T) *activeCycle { + t.Helper() + facts, _ := memberFacts(nsMember()) + facts.IntentID = f.intentID + claim, obs, err := observeNonSettlement(context.Background(), nsChainPast(facts.Deadline), facts, "its batch quorum was never reached") + if err != nil { + t.Fatal(err) + } + rec := &NonSettlementRecord{Facts: facts, Cause: claim.Cause, MemberChains: []int64{odChain}, MemberLegs: 1} + c := nonSettlementCycle(rec, claim) + f.cycles++ + c.CycleID = fmt.Sprintf("%s-%d", c.CycleID, f.cycles) + c.Request.CycleID = c.CycleID + c.Result.CycleID = c.CycleID + c.Result.ObservationResults = []*chain.ObservationResult{obs} + c.Result.ThresholdMet = true + c.Result.AggregatedAttestation = &attestation.AggregatedAttestation{ + ThresholdMet: true, Verified: true, AchievedWeight: 700, TotalWeight: 700, ParticipantCount: 7} + return c +} + +func (f *writeBackFixture) registered(t *testing.T) *database.MemberWriteBack { + t.Helper() + w, err := f.repos.IntentLifecycle.MemberWriteBackOf(context.Background(), f.intentID, odChain) + if err != nil { + t.Fatal(err) + } + return w +} + +// The defect through Phase 9's own entry point: two cycles for one member submitted two write-backs. +func TestAMemberIsWrittenBackOnce(t *testing.T) { + f := newWriteBackFixture(t) + ctx := context.Background() + if err := f.o.executePhase9(ctx, f.cycle(t)); err != nil { + t.Fatalf("the first write-back: %v", err) + } + second := f.cycle(t) + err := f.o.executePhase9(ctx, second) + if err == nil || f.sub.n != 1 { + t.Fatalf("THE regression: a second cycle for a written-back member submitted a second write-back (submissions %d, err %v)", f.sub.n, err) + } + if second.Result.WriteBackState == WriteBackWritten { + t.Fatalf("the refused cycle reads written") + } +} + +func TestASecondWriteBackIsRefusedByName(t *testing.T) { + f := newWriteBackFixture(t) + ctx := context.Background() + first := f.cycle(t) + if err := f.o.executePhase9(ctx, first); err != nil { + t.Fatal(err) + } + w := f.registered(t) + if w == nil || w.State != database.MemberWriteBackWritten || w.WriteBackTx != first.Result.WriteBackTxHash || w.CycleID != first.CycleID { + t.Fatalf("the write-back is not registered as written by its cycle: %+v", w) + } + second := f.cycle(t) + err := f.o.executePhase9(ctx, second) + if !errors.Is(err, database.ErrMemberAlreadyWrittenBack) || second.Result.WriteBackState != WriteBackRefusedAlreadyWritten { + t.Fatalf("want ErrMemberAlreadyWrittenBack / %s, got %v / %s", WriteBackRefusedAlreadyWritten, err, second.Result.WriteBackState) + } + if !strings.Contains(err.Error(), first.Result.WriteBackTxHash) { + t.Fatalf("the refusal does not name the write-back already on Accumulate: %v", err) + } +} + +func TestAWriteBackNeverSentReleasesTheMember(t *testing.T) { + f := newWriteBackFixture(t, fmt.Errorf("%w: insufficient credits: have 1, need 10", ErrWriteBackNotSent)) + ctx := context.Background() + failed := f.cycle(t) + if err := f.o.executePhase9(ctx, failed); err == nil { + t.Fatal("a write-back that was not sent reads as written") + } + if w := f.registered(t); w == nil || w.State != database.MemberWriteBackNotSent || !strings.Contains(w.Reason, "insufficient credits") { + t.Fatalf("an unsent write-back is registered as %+v; want not_sent, saying why", w) + } + retry := f.cycle(t) + if err := f.o.executePhase9(ctx, retry); err != nil { + t.Fatalf("a retry after a write-back that was never sent: %v", err) + } + if w := f.registered(t); w.State != database.MemberWriteBackWritten || w.CycleID != retry.CycleID { + t.Fatalf("the retry's write-back is registered as %+v", w) + } +} + +func TestAWriteBackWithAnUnknownOutcomeStopsAnother(t *testing.T) { + f := newWriteBackFixture(t, errors.New("failed to submit envelope: context deadline exceeded")) + ctx := context.Background() + unknown := f.cycle(t) + err := f.o.executePhase9(ctx, unknown) + if err == nil || unknown.Result.WriteBackState != WriteBackUnresolved { + t.Fatalf("a submission with no answer: err %v, state %s; want %s", err, unknown.Result.WriteBackState, WriteBackUnresolved) + } + if w := f.registered(t); w == nil || w.State != database.MemberWriteBackClaimed || w.CycleID != unknown.CycleID { + t.Fatalf("the unanswered submission's claim: %+v; want it kept", w) + } + next := f.cycle(t) + err = f.o.executePhase9(ctx, next) + if !errors.Is(err, database.ErrMemberWriteBackUnresolved) || f.sub.n != 1 || next.Result.WriteBackState != WriteBackRefusedUnresolved { + t.Fatalf("a write-back after an unanswered one: err %v, submissions %d, state %s", err, f.sub.n, next.Result.WriteBackState) + } +} + +// A proof cycle for a member already written back does nothing - no observation, attestation, bundle or +// outcome - and says so; it does not record a failure over the member's written-back outcome. +func TestAProofCycleForAWrittenBackMemberDoesNothing(t *testing.T) { + db := s1OpenDB(t) + ctx := context.Background() + o := s1Orchestrator(t, db) + o.config.Registry = strategy.NewRegistry() + id := fmt.Sprintf("f59-start-%d", time.Now().UnixNano()) + s1Seed(ctx, t, db, id) + t.Cleanup(func() { + db.Exec(`DELETE FROM intent_member_outcomes WHERE intent_id = $1`, id) + db.Exec(`DELETE FROM member_write_backs WHERE intent_id = $1`, id) + }) + lifecycle := o.config.Repos.IntentLifecycle + if err := lifecycle.ClaimMemberWriteBack(ctx, id, 84532, "cycle-earlier", "validator-6"); err != nil { + t.Fatal(err) + } + if err := lifecycle.RecordMemberWriteBack(ctx, id, 84532, "cycle-earlier", "acc://"+strings.Repeat("ab", 32)+"@results.acme/data"); err != nil { + t.Fatal(err) + } + req := memberCycle(id, "84532", []int64{84532}, 1, nil).Request + req.ProofClass = string(LaneOnDemand) + req.TxHashes = []string{"0x" + strings.Repeat("cd", 32)} + _, err := o.StartProofCycle(ctx, req) + if !errors.Is(err, database.ErrMemberAlreadyWrittenBack) { + t.Fatalf("a proof cycle for a written-back member: want ErrMemberAlreadyWrittenBack, got %v", err) + } + var n int + if err := db.QueryRow(`SELECT count(*) FROM intent_member_outcomes WHERE intent_id = $1`, id).Scan(&n); err != nil || n != 0 { + t.Fatalf("the refused cycle recorded an outcome (%d, %v)", n, err) + } +} + +// A member written back before write-backs were registered has no register row; its recorded outcome states the +// write-back, and stops another just the same (000ac79a's arbitrum member, written back 2026-09-29, is one). +func TestAMemberWrittenBackBeforeTheRegisterIsNotWrittenBackAgain(t *testing.T) { + f := newWriteBackFixture(t) + ctx := context.Background() + if _, err := f.db.Exec(`INSERT INTO intent_lifecycle (intent_id, accum_tx_hash, status) VALUES ($1, $2, 'settling')`, + f.intentID, "f59-"+f.intentID[len(f.intentID)-12:]); err != nil { + t.Fatal(err) + } + t.Cleanup(func() { + f.db.Exec(`DELETE FROM intent_member_outcomes WHERE intent_id = $1`, f.intentID) + f.db.Exec(`DELETE FROM intent_lifecycle WHERE intent_id = $1`, f.intentID) + }) + earlier := "acc://" + strings.Repeat("ef", 32) + "@results.acme/data" + if _, err := f.repos.IntentLifecycle.RecordMemberOutcome(ctx, database.MemberOutcome{IntentID: f.intentID, ChainID: odChain, + MemberChains: []int64{odChain}, Legs: 1, Settlement: database.MemberSettlementSettled, ProofCycle: database.MemberProofCycleWritten, + SettlementTx: "0x" + strings.Repeat("12", 32), WriteBackTx: earlier, CycleID: "cycle-before-the-register", ReportedBy: "validator-1"}); err != nil { + t.Fatal(err) + } + if w := f.registered(t); w != nil { + t.Fatalf("premise: no register row, got %+v", w) + } + c := f.cycle(t) + err := f.o.executePhase9(ctx, c) + if !errors.Is(err, database.ErrMemberAlreadyWrittenBack) || f.sub.n != 0 || !strings.Contains(err.Error(), earlier) { + t.Fatalf("a member written back before the register: err %v, submissions %d", err, f.sub.n) + } + if err := f.repos.IntentLifecycle.MemberWriteBackAllowed(ctx, f.intentID, odChain); !errors.Is(err, database.ErrMemberAlreadyWrittenBack) { + t.Fatalf("the start check: %v", err) + } +} diff --git a/pkg/intent/discovery.go b/pkg/intent/discovery.go index 7d39aefa..00d5d554 100644 --- a/pkg/intent/discovery.go +++ b/pkg/intent/discovery.go @@ -182,7 +182,9 @@ func (s IntentStatus) String() string { // IntentDiscovery monitors Accumulate blockchain for Certen transaction intents type IntentDiscovery struct { - client accumulate.Client + client accumulate.Client + // reprocess runs one intent's round for ReprocessIntent; nil is processIntent (tests replace it). + reprocess func(intent *CertenIntent, blockHeight uint64) (consensus.TargetChainOutcome, error) accumulateURL string config *IntentDiscoveryConfig ledgerStore LedgerStoreInterface // For persistence @@ -1515,6 +1517,56 @@ func (id *IntentDiscovery) handleRetryJob(job *intentRetryJob) { // STAGE 1: returns the target-chain outcome alongside the error. A nil error means // CONSENSUS committed; it says nothing about whether the chain write landed, and // the caller used to read it as though it did. +// ReprocessIntent processes one intent again exactly as discovery processes it: found in its Directory Network +// block by the Accumulate transaction that carries it, converted as discovery converts it, and run through +// processIntent - so its round is re-derived and checked against the committed block. It exists for a repair that +// must reach one of the intent's members (RB4-F55; consensus.ArmMemberRepair). Nothing else in the block is +// touched; an intent being processed now, or permanently invalid, is refused; a transaction carrying another intent +// than the one named is refused. +func (id *IntentDiscovery) ReprocessIntent(ctx context.Context, dnBlock uint64, accumTxHash, intentID string) error { + txs, err := id.client.SearchCertenTransactions(ctx, int64(dnBlock)) + if err != nil { + return fmt.Errorf("search DN block %d: %w", dnBlock, err) + } + want := strings.ToLower(strings.TrimPrefix(accumTxHash, "0x")) + for _, tx := range txs { + if tx.BlockHeight != int64(dnBlock) || strings.ToLower(strings.TrimPrefix(tx.Hash, "0x")) != want { + continue + } + ci, err := id.convertCertenTransactionToIntent(tx) + if err != nil { + return fmt.Errorf("transaction %s in DN block %d: %w", tx.Hash, dnBlock, err) + } + if ci.IntentID != intentID { + return fmt.Errorf("transaction %s in DN block %d carries intent %s, not %s", tx.Hash, dnBlock, ci.IntentID, intentID) + } + id.mu.Lock() + switch id.intentStatus[intentID] { + case IntentStatusInProgress: + id.mu.Unlock() + return fmt.Errorf("intent %s is being processed now; not processed twice", intentID) + case IntentStatusFailedPermanent: + id.mu.Unlock() + return fmt.Errorf("intent %s is permanently invalid; it is not processed again", intentID) + } + id.intentStatus[intentID] = IntentStatusInProgress + id.mu.Unlock() + + id.logger.Printf("🔧 [REPROCESS] intent %s from DN block %d (transaction %s)", intentID, dnBlock, tx.Hash) + run := id.reprocess + if run == nil { + run = id.processIntent + } + if _, err := run(ci, dnBlock); err != nil { + id.markFailedClassified(intentID, err) + return fmt.Errorf("reprocess intent %s: %w", intentID, err) + } + id.markCompleted(intentID) + return nil + } + return fmt.Errorf("transaction %s is not in DN block %d", accumTxHash, dnBlock) +} + func (id *IntentDiscovery) processIntent(intent *CertenIntent, blockHeight uint64) (consensus.TargetChainOutcome, error) { id.logger.Printf("🚀 Processing Certen intent: %s", intent.IntentID) diff --git a/pkg/intent/reprocess_test.go b/pkg/intent/reprocess_test.go new file mode 100644 index 00000000..f7751be9 --- /dev/null +++ b/pkg/intent/reprocess_test.go @@ -0,0 +1,128 @@ +// Copyright 2026 Certen Protocol + +package intent + +import ( + "context" + "errors" + "io" + "log" + "strings" + "testing" + + "github.com/certen/independant-validator/pkg/accumulate" + "github.com/certen/independant-validator/pkg/consensus" +) + +// RB4-F55 repair: one intent is processed again exactly as discovery processes it - from its Directory Network +// block, converted from the same transaction - so the round a repair needs is re-derived and checked against the +// committed block (consensus.ArmMemberRepair). Nothing else in the block is touched, an intent being processed now +// is not processed twice, and an intent named by the wrong transaction is refused. + +type blockClient struct { + accumulate.Client + txs []*accumulate.CertenTransaction +} + +func (c *blockClient) SearchCertenTransactions(_ context.Context, h int64) ([]*accumulate.CertenTransaction, error) { + var out []*accumulate.CertenTransaction + for _, tx := range c.txs { + if tx.BlockHeight == h { + out = append(out, tx) + } + } + return out, nil +} + +func blockTx(hash, intentID string, height int64) *accumulate.CertenTransaction { + return &accumulate.CertenTransaction{ + Hash: hash, AccountURL: "acc://harbor.acme/data", BlockHeight: height, + Partition: "acc://dn.acme", ProofPartition: "bvn1", ProofBlockIndex: 8270001, + IntentData: map[string]interface{}{ + "intentData": map[string]interface{}{"intent_id": intentID, "proof_class": "on_demand"}, + "crossChainData": map[string]interface{}{}, + "governanceData": map[string]interface{}{"organizationAdi": "acc://harbor.acme"}, + "replayData": map[string]interface{}{"expires_at": 1790600000}, + }, + } +} + +func reprocessDiscovery(t *testing.T, txs ...*accumulate.CertenTransaction) (*IntentDiscovery, *[]*CertenIntent) { + t.Helper() + id := &IntentDiscovery{client: &blockClient{txs: txs}, logger: log.New(io.Discard, "", 0), intentStatus: map[string]IntentStatus{}} + var processed []*CertenIntent + id.reprocess = func(ci *CertenIntent, height uint64) (consensus.TargetChainOutcome, error) { + processed = append(processed, ci) + return consensus.TargetChainPending, nil + } + return id, &processed +} + +func intentIDOf(t *testing.T, tx *accumulate.CertenTransaction) string { + t.Helper() + ci, err := (&IntentDiscovery{logger: log.New(io.Discard, "", 0)}).convertCertenTransactionToIntent(tx) + if err != nil { + t.Fatal(err) + } + return ci.IntentID +} + +func TestReprocessingAnIntentTouchesOnlyThatIntent(t *testing.T) { + target := blockTx(strings.Repeat("db", 32), "repair-target", 10007772) + other := blockTx(strings.Repeat("0e", 32), "repair-other", 10007772) + id, processed := reprocessDiscovery(t, other, target) + want := intentIDOf(t, target) + if err := id.ReprocessIntent(context.Background(), 10007772, target.Hash, want); err != nil { + t.Fatal(err) + } + if len(*processed) != 1 || (*processed)[0].IntentID != want || (*processed)[0].TransactionHash != target.Hash { + t.Fatalf("processed %d intents; want exactly the named one", len(*processed)) + } + if (*processed)[0].ProofPartition != "bvn1" { + t.Fatalf("the intent was not converted as discovery converts it: partition %q", (*processed)[0].ProofPartition) + } + if id.getIntentStatus(want) != IntentStatusCompleted { + t.Fatalf("status after: %v", id.getIntentStatus(want)) + } +} + +func TestReprocessingRefusesWhatItCannotName(t *testing.T) { + target := blockTx(strings.Repeat("db", 32), "repair-target", 10007772) + want := intentIDOf(t, target) + ctx := context.Background() + + id, processed := reprocessDiscovery(t, target) + if err := id.ReprocessIntent(ctx, 10007773, target.Hash, want); err == nil || !strings.Contains(err.Error(), "not in") { + t.Fatalf("an intent looked for in the wrong block: %v", err) + } + if err := id.ReprocessIntent(ctx, 10007772, target.Hash, "another-intent"); err == nil || !strings.Contains(err.Error(), want) { + t.Fatalf("a transaction carrying another intent than the one named: %v", err) + } + id.intentStatus[want] = IntentStatusInProgress + if err := id.ReprocessIntent(ctx, 10007772, target.Hash, want); err == nil || !strings.Contains(err.Error(), "being processed") { + t.Fatalf("an intent being processed now: %v", err) + } + id.intentStatus[want] = IntentStatusFailedPermanent + if err := id.ReprocessIntent(ctx, 10007772, target.Hash, want); err == nil || !strings.Contains(err.Error(), "permanently invalid") { + t.Fatalf("a permanently invalid intent: %v", err) + } + if len(*processed) != 0 { + t.Fatalf("a refused reprocess processed %d intents", len(*processed)) + } +} + +func TestAReprocessingFailureIsReturned(t *testing.T) { + target := blockTx(strings.Repeat("db", 32), "repair-target", 10007772) + id, _ := reprocessDiscovery(t, target) + want := intentIDOf(t, target) + boom := errors.New("committed as another block") + id.reprocess = func(*CertenIntent, uint64) (consensus.TargetChainOutcome, error) { + return consensus.TargetChainPending, boom + } + if err := id.ReprocessIntent(context.Background(), 10007772, target.Hash, want); !errors.Is(err, boom) { + t.Fatalf("the round's failure was not returned: %v", err) + } + if id.getIntentStatus(want) != IntentStatusFailed { + t.Fatalf("status after a failure: %v", id.getIntentStatus(want)) + } +} diff --git a/repair_command.go b/repair_command.go index daf3c7ae..f982bcd1 100644 --- a/repair_command.go +++ b/repair_command.go @@ -20,7 +20,7 @@ import ( "github.com/certen/independant-validator/pkg/execution" ) -const repairUsage = "usage: certen-validator repair anchor-blocks [--apply] [--min-depth N] | repair projections [--apply] | repair consensus-records --rpc ADDR [--apply]" +const repairUsage = "usage: certen-validator repair anchor-blocks [--apply] [--min-depth N] | repair projections [--apply] | repair consensus-records --rpc ADDR [--apply] | repair member-proof-cycle --intent ID --chain CHAIN_ID --tx SETTLEMENT_TX [--apply] [--wait DURATION]" // runRepairCommand runs `validator repair anchor-blocks`: it reads every canonical anchor's verify and // create transactions back from their chain - locating the create transaction where the row does not name @@ -39,6 +39,10 @@ const repairUsage = "usage: certen-validator repair anchor-blocks [--apply] [--m // // Exit status: 0 when nothing was refused, 2 when something was refused or could not be read, 1 on error. func runRepairCommand(args []string) int { + // Served by the running validator, not by this process (member_repair_command.go). + if len(args) > 0 && args[0] == "member-proof-cycle" { + return runMemberRepairCommand(args[1:]) + } if len(args) == 0 || (args[0] != "anchor-blocks" && args[0] != "projections" && args[0] != "consensus-records") { log.Print(repairUsage) return 1