diff --git a/pkg/database/proof_artifact_repository.go b/pkg/database/proof_artifact_repository.go index c9e317e6..a8df8f8a 100644 --- a/pkg/database/proof_artifact_repository.go +++ b/pkg/database/proof_artifact_repository.go @@ -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, diff --git a/pkg/database/repository_anchor_batch_record.go b/pkg/database/repository_anchor_batch_record.go new file mode 100644 index 00000000..8150593a --- /dev/null +++ b/pkg/database/repository_anchor_batch_record.go @@ -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) +} diff --git a/pkg/database/repository_anchor_quorum.go b/pkg/database/repository_anchor_quorum.go index 02231ef2..8c241fc9 100644 --- a/pkg/database/repository_anchor_quorum.go +++ b/pkg/database/repository_anchor_quorum.go @@ -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 diff --git a/pkg/database/repository_batch.go b/pkg/database/repository_batch.go index eb3651fd..3337ad71 100644 --- a/pkg/database/repository_batch.go +++ b/pkg/database/repository_batch.go @@ -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 diff --git a/pkg/database/schema_prepare_test.go b/pkg/database/schema_prepare_test.go index 2721f40a..50a831a2 100644 --- a/pkg/database/schema_prepare_test.go +++ b/pkg/database/schema_prepare_test.go @@ -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 } diff --git a/pkg/database/types.go b/pkg/database/types.go index bebd142f..a925a5b2 100644 --- a/pkg/database/types.go +++ b/pkg/database/types.go @@ -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 diff --git a/pkg/server/batch_record_test.go b/pkg/server/batch_record_test.go new file mode 100644 index 00000000..6736ec5e --- /dev/null +++ b/pkg/server/batch_record_test.go @@ -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()) + } +} diff --git a/pkg/server/proof_handlers.go b/pkg/server/proof_handlers.go index 1434c645..64120c63 100644 --- a/pkg/server/proof_handlers.go +++ b/pkg/server/proof_handlers.go @@ -7,6 +7,7 @@ package server import ( "encoding/json" + "errors" "fmt" "log" "net/http" @@ -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)