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
2 changes: 1 addition & 1 deletion pkg/database/proof_artifact_repository.go
Original file line number Diff line number Diff line change
Expand Up @@ -1314,7 +1314,7 @@ func (r *ProofArtifactRepository) GetProofsModifiedSince(ctx context.Context, si
func (r *ProofArtifactRepository) GetBatchProofStats(ctx context.Context, batchID uuid.UUID) (*BatchProofStats, error) {
query := `
SELECT
$1 as batch_id,
$1::uuid as batch_id,
COUNT(*) as proof_count,
(SELECT COUNT(*) FROM validator_attestations WHERE batch_id = $1) as attestation_count,
COUNT(*) FILTER (WHERE verification_status = 'verified') as verified_count,
Expand Down
173 changes: 173 additions & 0 deletions pkg/database/repository_anchor_batch_record.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,173 @@
// Copyright 2025 Certen Protocol
//
// The anchor batch record as the validators write it (RecordAnchorQuorum, repository_anchor_quorum.go): the anchored root, the chain and its create and verify transactions, and the
// quorum evidence over it.

package database

import (
"context"
"database/sql"
"encoding/hex"
"encoding/json"
"fmt"
"time"

"github.com/google/uuid"
)

// AnchorBatchRecord is one anchor_batches row. A column the validators have not written is null, never a
// default: validator_id, for one, is NULL on every row the quorum path writes (a batch is the quorum's, not one
// validator's), and the chain's transactions are unknown until the validator that sent them records them.
type AnchorBatchRecord struct {
BatchID uuid.UUID `json:"batch_id"`
BatchType string `json:"batch_type"`
Status string `json:"status"`
Lane *string `json:"lane"`
MerkleRoot *string `json:"merkle_root"`
TransactionCount int `json:"transaction_count"`
TargetChain string `json:"target_chain"`
ChainID *int64 `json:"chain_id"`
BundleID *string `json:"bundle_id"`
BatchOperationID *string `json:"batch_operation_id"`
ValidatorID *string `json:"validator_id"`

AnchorCreateTx *string `json:"anchor_create_tx"`
AnchorTxHash *string `json:"anchor_tx_hash"`
AnchorBlockNumber *int64 `json:"anchor_block_number"`
AnchorCreateSender *string `json:"anchor_create_sender"`
VerifyTx *string `json:"verify_tx"`
VerifyBlock *int64 `json:"verify_block"`
VerifySender *string `json:"verify_sender"`
GasUsed *int64 `json:"gas_used"`

MessageHash *string `json:"message_hash"`
QuorumReached bool `json:"quorum_reached"`
AttestationCount int `json:"attestation_count"`
SignedVotingPower *string `json:"signed_voting_power"`
TotalVotingPower *string `json:"total_voting_power"`
Signers json.RawMessage `json:"signers"`
AggregatedSignature *string `json:"aggregated_signature"`
AggregatedPublicKey *string `json:"aggregated_public_key"`
EvidenceSource *string `json:"evidence_source"`
ConsensusCompletedAt *time.Time `json:"consensus_completed_at"`

AccumulateBlockHeight *int64 `json:"accumulate_block_height"`
AccumulateBlockHash *string `json:"accumulate_block_hash"`
BPTRoot *string `json:"bpt_root"`
GovernanceRoot *string `json:"governance_root"`
ErrorMessage *string `json:"error_message"`

CreatedAt time.Time `json:"created_at"`
UpdatedAt time.Time `json:"updated_at"`
StartedAt *time.Time `json:"batch_start_time"`
EndedAt *time.Time `json:"batch_end_time"`
ClosedAt *time.Time `json:"closed_at"`
AnchoredAt *time.Time `json:"anchored_at"`
ConfirmedAt *time.Time `json:"confirmed_at"`
}

// anchorBatchRecordColumns is the column list every batch reader selects, in scanAnchorBatchRecord's order.
const anchorBatchRecordColumns = `id, batch_type, status, lane, merkle_root, transaction_count, target_chain, chain_id,
bundle_id, batch_operation_id, validator_id,
anchor_create_tx, anchor_tx_hash, anchor_block_num, anchor_create_sender, verify_tx, verify_block,
verify_sender, gas_used,
message_hash, COALESCE(quorum_reached, FALSE), COALESCE(attestation_count, 0),
signed_voting_power::text, total_voting_power::text, signers, aggregated_signature,
aggregated_public_key, evidence_source, consensus_completed_at,
accumulate_block_height, accumulate_block_hash, bpt_root, governance_root, error_message,
created_at, updated_at, batch_start_time, batch_end_time, closed_at, anchored_at, confirmed_at`

func nullStringPtr(v sql.NullString) *string {
if !v.Valid {
return nil
}
return &v.String
}

func nullInt64Ptr(v sql.NullInt64) *int64 {
if !v.Valid {
return nil
}
return &v.Int64
}

func nullTimePtr(v sql.NullTime) *time.Time {
if !v.Valid {
return nil
}
return &v.Time
}

func hexPtr(b []byte) *string {
if b == nil {
return nil
}
s := hex.EncodeToString(b)
return &s
}

// scanAnchorBatchRecord scans one row selected with anchorBatchRecordColumns. A column the validators have not
// written stays null. The scan's own error (sql.ErrNoRows included) is returned as is.
func scanAnchorBatchRecord(scan func(...any) error) (*AnchorBatchRecord, error) {
var (
rec AnchorBatchRecord
lane, bundleID, opID, validatorID, createTx, anchorTx, createSender sql.NullString
verifyTx, verifySender, messageHash, signedPower, totalPower, evidence sql.NullString
accumHash, errorMessage sql.NullString
chainID, anchorBlock, verifyBlock, gasUsed, accumHeight sql.NullInt64
merkleRoot, aggSig, aggPub, bptRoot, govRoot, signers []byte
consensusAt, startedAt, endedAt, closedAt, anchoredAt, confirmedAt sql.NullTime
)
if err := scan(
&rec.BatchID, &rec.BatchType, &rec.Status, &lane, &merkleRoot, &rec.TransactionCount, &rec.TargetChain,
&chainID, &bundleID, &opID, &validatorID,
&createTx, &anchorTx, &anchorBlock, &createSender, &verifyTx, &verifyBlock, &verifySender, &gasUsed,
&messageHash, &rec.QuorumReached, &rec.AttestationCount,
&signedPower, &totalPower, &signers, &aggSig, &aggPub, &evidence, &consensusAt,
&accumHeight, &accumHash, &bptRoot, &govRoot, &errorMessage,
&rec.CreatedAt, &rec.UpdatedAt, &startedAt, &endedAt, &closedAt, &anchoredAt, &confirmedAt,
); err != nil {
return nil, err
}
rec.Lane, rec.BundleID, rec.BatchOperationID, rec.ValidatorID = nullStringPtr(lane), nullStringPtr(bundleID), nullStringPtr(opID), nullStringPtr(validatorID)
rec.AnchorCreateTx, rec.AnchorTxHash, rec.AnchorCreateSender = nullStringPtr(createTx), nullStringPtr(anchorTx), nullStringPtr(createSender)
rec.VerifyTx, rec.VerifySender, rec.MessageHash = nullStringPtr(verifyTx), nullStringPtr(verifySender), nullStringPtr(messageHash)
rec.SignedVotingPower, rec.TotalVotingPower, rec.EvidenceSource = nullStringPtr(signedPower), nullStringPtr(totalPower), nullStringPtr(evidence)
rec.AccumulateBlockHash, rec.ErrorMessage = nullStringPtr(accumHash), nullStringPtr(errorMessage)
rec.ChainID, rec.AnchorBlockNumber, rec.VerifyBlock = nullInt64Ptr(chainID), nullInt64Ptr(anchorBlock), nullInt64Ptr(verifyBlock)
rec.GasUsed, rec.AccumulateBlockHeight = nullInt64Ptr(gasUsed), nullInt64Ptr(accumHeight)
rec.MerkleRoot, rec.AggregatedSignature, rec.AggregatedPublicKey = hexPtr(merkleRoot), hexPtr(aggSig), hexPtr(aggPub)
rec.BPTRoot, rec.GovernanceRoot = hexPtr(bptRoot), hexPtr(govRoot)
if signers != nil {
rec.Signers = json.RawMessage(signers)
} else {
rec.Signers = json.RawMessage("null")
}
rec.ConsensusCompletedAt = nullTimePtr(consensusAt)
rec.StartedAt, rec.EndedAt, rec.ClosedAt = nullTimePtr(startedAt), nullTimePtr(endedAt), nullTimePtr(closedAt)
rec.AnchoredAt, rec.ConfirmedAt = nullTimePtr(anchoredAt), nullTimePtr(confirmedAt)
return &rec, nil
}

// oneAnchorBatch reads the single row a query selects; ErrBatchNotFound if it selects none.
func (r *BatchRepository) oneAnchorBatch(ctx context.Context, what, query string, args ...any) (*AnchorBatchRecord, error) {
rec, err := scanAnchorBatchRecord(r.client.QueryRowContext(ctx, query, args...).Scan)
if err == sql.ErrNoRows {
return nil, ErrBatchNotFound
}
if err != nil {
return nil, fmt.Errorf("failed to read %s: %w", what, err)
}
return rec, nil
}

// GetBatch reads one anchor batch as the validators wrote it; ErrBatchNotFound if there is none. It read 13 of
// the row's columns into a type whose validator_id could not hold the NULL the quorum path writes, so it failed on
// every batch the validators record.
func (r *BatchRepository) GetBatch(ctx context.Context, batchID uuid.UUID) (*AnchorBatchRecord, error) {
return r.oneAnchorBatch(ctx, "anchor batch "+batchID.String(), `
SELECT `+anchorBatchRecordColumns+`
FROM anchor_batches
WHERE id = $1`, batchID)
}
2 changes: 1 addition & 1 deletion pkg/database/repository_anchor_quorum.go
Original file line number Diff line number Diff line change
Expand Up @@ -290,7 +290,7 @@ func (r *BatchRepository) RecordAnchorQuorum(
}

// AnchorQuorumRow is a canonical anchor as stored: the identity the contract uses, and the quorum proven
// over it. Distinct from AnchorBatch, which models the legacy collector's row.
// over it: the fields the quorum identity needs. The whole row is AnchorBatchRecord (GetBatch).
type AnchorQuorumRow struct {
BatchID uuid.UUID
ChainID int64
Expand Down
29 changes: 1 addition & 28 deletions pkg/database/repository_batch.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,34 +27,7 @@ func NewBatchRepository(client *Client) *BatchRepository {
// ANCHOR BATCH OPERATIONS
// ============================================================================

// GetBatch retrieves a batch by ID
func (r *BatchRepository) GetBatch(ctx context.Context, batchID uuid.UUID) (*AnchorBatch, error) {
query := `
SELECT id, batch_type, merkle_root, transaction_count,
batch_start_time, batch_end_time, accumulate_block_height,
accumulate_block_hash, validator_id, status, error_message,
created_at, updated_at
FROM anchor_batches
WHERE id = $1`

batch := &AnchorBatch{}
err := r.client.QueryRowContext(ctx, query, batchID).Scan(
&batch.BatchID, &batch.BatchType, &batch.MerkleRoot, &batch.TxCount,
&batch.StartTime, &batch.EndTime, &batch.AccumHeight,
&batch.AccumHash, &batch.ValidatorID, &batch.Status, &batch.ErrorMessage,
&batch.CreatedAt, &batch.UpdatedAt,
)

if err == sql.ErrNoRows {
// F.4 remediation: Return explicit error instead of nil, nil
return nil, ErrBatchNotFound
}
if err != nil {
return nil, fmt.Errorf("failed to get batch: %w", err)
}

return batch, nil
}
// GetBatch, the complete record of one batch, is in repository_anchor_batch_record.go.

// ============================================================================
// BATCH TRANSACTION OPERATIONS
Expand Down
12 changes: 8 additions & 4 deletions pkg/database/schema_prepare_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -330,10 +330,14 @@ func looksLikeRepositorySQL(value string) bool {
if strings.Contains(value, "%") {
return false
}
for _, prefix := range []string{"SELECT ", "INSERT ", "UPDATE ", "DELETE ", "WITH "} {
if strings.HasPrefix(value, prefix) {
return true
}
// The keyword may be followed by any whitespace: a statement written "SELECT" then a newline was skipped.
words := strings.Fields(value)
if len(words) < 2 {
return false
}
switch words[0] {
case "SELECT", "INSERT", "UPDATE", "DELETE", "WITH":
return true
}
return false
}
Expand Down
21 changes: 4 additions & 17 deletions pkg/database/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -46,23 +46,10 @@ const (
BatchStatusFailed BatchStatus = "failed" // Anchoring failed
)

// AnchorBatch represents a batch of transactions anchored together
// Maps to: anchor_batches table
type AnchorBatch struct {
BatchID uuid.UUID `db:"batch_id" json:"batch_id"`
BatchType BatchType `db:"batch_type" json:"batch_type"`
MerkleRoot []byte `db:"merkle_root" json:"merkle_root"` // 32 bytes SHA256
TxCount int `db:"transaction_count" json:"transaction_count"`
StartTime time.Time `db:"batch_start_time" json:"batch_start_time"`
EndTime sql.NullTime `db:"batch_end_time" json:"batch_end_time,omitempty"`
AccumHeight sql.NullInt64 `db:"accumulate_block_height" json:"accumulate_block_height,omitempty"`
AccumHash sql.NullString `db:"accumulate_block_hash" json:"accumulate_block_hash,omitempty"`
ValidatorID string `db:"validator_id" json:"validator_id"`
Status BatchStatus `db:"status" json:"status"`
ErrorMessage sql.NullString `db:"error_message" json:"error_message,omitempty"`
CreatedAt time.Time `db:"created_at" json:"created_at"`
UpdatedAt time.Time `db:"updated_at" json:"updated_at"`
}
// AnchorBatch is an anchor_batches row: the complete record (repository_anchor_batch_record.go). It was a 13-column
// struct whose validator_id and batch_start_time could not hold the NULLs the quorum path writes, so GetBatch failed
// on every batch the validators record.
type AnchorBatch = AnchorBatchRecord

// ============================================================================
// BATCH TRANSACTION TYPES
Expand Down
95 changes: 95 additions & 0 deletions pkg/server/batch_record_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
// Copyright 2026 Certen Protocol

package server

import (
"context"
"crypto/sha256"
"database/sql"
"encoding/hex"
"encoding/json"
"net/http"
"net/http/httptest"
"os"
"testing"

"github.com/google/uuid"
_ "github.com/lib/pq"

schema "github.com/certen/independant-validator/db"
"github.com/certen/independant-validator/pkg/database"
)

// RB4-F29/F31: the batch endpoints read the row the quorum path writes. /api/batches/{id} scanned validator_id
// (NULL on every quorum row) into a string and failed on every batch the validators record; the stats query at
// /api/v1/batches/{id}/stats could never run (an untyped $1), and would have answered zero counts for a batch
// that does not exist.
func TestBatchEndpointsReadTheQuorumRow(t *testing.T) {
conn := os.Getenv("CERTEN_TEST_DB")
if conn == "" {
t.Fatal("CERTEN_TEST_DB is required (a skipped gate is not a green gate)")
}
db, err := sql.Open("postgres", conn)
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { db.Close() })
ctx := context.Background()
if err := (schema.Runner{DB: db}).Up(ctx, "server-test"); err != nil {
t.Fatalf("migrate: %v", err)
}
id := uuid.New()
sum := sha256.Sum256([]byte("quorum-row-" + id.String()))
createTx := "0x" + hex.EncodeToString(sum[:])
// The columns RecordAnchorQuorum writes (repository_anchor_quorum.go), validator_id NULL as it writes it.
if _, err := db.ExecContext(ctx, `
INSERT INTO anchor_batches (id, batch_type, status, merkle_root, target_chain, validator_id, transaction_count, tx_count,
chain_id, bundle_id, anchor_create_tx, anchor_tx_hash, anchor_block_num, quorum_reached, attestation_count,
signed_voting_power, total_voting_power, evidence_source, lane, proof_data_included, anchored_at, confirmed_at, closed_at)
VALUES ($1, 'on_demand', 'confirmed', $2, 'ethereum', NULL, 1, 1, 11155111, $3, $4, $4, 900, TRUE, 5, 5, 7,
'chain_event', 'on_demand', TRUE, now(), now(), now())`, id, sum[:], "bundle-"+id.String()[:8], createTx); err != nil {
t.Fatal(err)
}
repos := database.NewRepositories(database.NewClientFromDB(db))

rr := httptest.NewRecorder()
NewBatchHandlers(repos, "t", nil).HandleBatchStatus(rr, httptest.NewRequest(http.MethodGet, "/api/batches/"+id.String(), nil))
if rr.Code != http.StatusOK {
t.Fatalf("/api/batches/{id}: status %d: %s", rr.Code, rr.Body.String())
}
var got struct {
BatchID string `json:"batch_id"`
ValidatorID *string `json:"validator_id"`
ChainID int64 `json:"chain_id"`
MerkleRoot string `json:"merkle_root"`
AnchorCreateTx string `json:"anchor_create_tx"`
QuorumReached bool `json:"quorum_reached"`
Attestations int `json:"attestation_count"`
SignedPower string `json:"signed_voting_power"`
}
if err := json.Unmarshal(rr.Body.Bytes(), &got); err != nil {
t.Fatal(err)
}
if got.BatchID != id.String() || got.ValidatorID != nil || got.ChainID != 11155111 || got.MerkleRoot != hex.EncodeToString(sum[:]) ||
got.AnchorCreateTx != createTx || !got.QuorumReached || got.Attestations != 5 || got.SignedPower != "5" {
t.Fatalf("/api/batches/{id}: not the record as written: %s", rr.Body.String())
}

proofs := NewProofHandlers(repos, "t", nil)
rr = httptest.NewRecorder()
proofs.HandleGetBatchStats(rr, httptest.NewRequest(http.MethodGet, "/api/v1/batches/"+id.String()+"/stats", nil))
if rr.Code != http.StatusOK {
t.Fatalf("/api/v1/batches/{id}/stats: status %d: %s", rr.Code, rr.Body.String())
}
rr = httptest.NewRecorder()
proofs.HandleGetBatchStats(rr, httptest.NewRequest(http.MethodGet, "/api/v1/batches/"+uuid.NewString()+"/stats", nil))
var e struct {
Error struct {
Code string `json:"code"`
} `json:"error"`
}
_ = json.Unmarshal(rr.Body.Bytes(), &e)
if rr.Code != http.StatusNotFound || e.Error.Code != "BATCH_NOT_FOUND" {
t.Fatalf("stats of a batch that does not exist: want 404 BATCH_NOT_FOUND, got %d %s", rr.Code, rr.Body.String())
}
}
14 changes: 14 additions & 0 deletions pkg/server/proof_handlers.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ package server

import (
"encoding/json"
"errors"
"fmt"
"log"
"net/http"
Expand Down Expand Up @@ -482,7 +483,20 @@ func (h *ProofHandlers) HandleGetBatchStats(w http.ResponseWriter, r *http.Reque
return
}

if h.repos == nil {
h.writeError(w, http.StatusServiceUnavailable, "DATABASE_UNAVAILABLE", "Database not available")
return
}
ctx := r.Context()
// The counts of a batch that does not exist are not zero; the batch is not found (they were zero for any ID).
if _, err := h.repos.Batches.GetBatch(ctx, batchID); errors.Is(err, database.ErrBatchNotFound) {
h.writeError(w, http.StatusNotFound, "BATCH_NOT_FOUND", "No anchor batch with ID "+batchID.String())
return
} else if err != nil {
h.logger.Printf("Error getting batch %s: %v", batchID, err)
h.writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR", "Failed to retrieve batch")
return
}
stats, err := h.repos.ProofArtifacts.GetBatchProofStats(ctx, batchID)
if err != nil {
h.logger.Printf("Error getting batch stats: %v", err)
Expand Down
Loading