Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions db/migrations/00015_member_outcome_corrections.sql
Original file line number Diff line number Diff line change
@@ -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 <intent_id>/<chain_id>; previous and corrected member rows
-- intent_lifecycle record_id <intent_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'));
31 changes: 31 additions & 0 deletions db/migrations/00016_member_write_backs.sql
Original file line number Diff line number Diff line change
@@ -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.';
2 changes: 1 addition & 1 deletion db/schema.fingerprint
Original file line number Diff line number Diff line change
@@ -1 +1 @@
b9a4c39c2196b0ad4816a0842c22f7c487f932330771820f835d4650ec9c63c2
0d6659c87b36cf1dd8a8e1b5c1dd705c2add9ebbd6029a3f865e20b38a4829fb
15 changes: 14 additions & 1 deletion main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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())

Expand Down Expand Up @@ -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
Expand Down
142 changes: 142 additions & 0 deletions member_repair_command.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
54 changes: 54 additions & 0 deletions member_repair_command_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
6 changes: 6 additions & 0 deletions pkg/consensus/batch_quorum_prover.go
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
3 changes: 3 additions & 0 deletions pkg/consensus/bft_integration.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading
Loading