Skip to content
Closed
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
11 changes: 6 additions & 5 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

68 changes: 55 additions & 13 deletions crates/gitlawb-node/src/db/mod.rs
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
use std::time::Duration;

use anyhow::{Context, Result};
use chrono::{DateTime, Utc};
use serde::{Deserialize, Serialize};
use sqlx::{postgres::PgPoolOptions, PgPool, Row};
use std::time::Duration;
use tracing::info;
use uuid::Uuid;

Expand Down Expand Up @@ -2158,19 +2159,16 @@ impl Db {
// ── Pinned CIDs ───────────────────────────────────────────────────────────────

impl Db {
pub async fn is_pinned(&self, sha256_hex: &str) -> Result<bool> {
let row = sqlx::query("SELECT COUNT(*) as cnt FROM pinned_cids WHERE sha256_hex = $1")
.bind(sha256_hex)
.fetch_one(&self.pool)
.await?;
Ok(row.get::<i64, _>("cnt") > 0)
}

/// Record the local IPFS CID for a git object.
/// If a Pinata-only row already exists (cid IS NULL), this updates it with
/// the real IPFS CID so that `has_ipfs_cid` correctly reflects IPFS state.
pub async fn record_pinned_cid(&self, sha256_hex: &str, cid: &str) -> Result<()> {
sqlx::query(
"INSERT INTO pinned_cids (sha256_hex, cid, pinned_at)
VALUES ($1, $2, $3)
ON CONFLICT(sha256_hex) DO NOTHING",
ON CONFLICT(sha256_hex) DO UPDATE SET
cid = COALESCE(pinned_cids.cid, EXCLUDED.cid),
pinned_at = EXCLUDED.pinned_at",
Comment on lines +2162 to +2171

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🗄️ Data Integrity & Integration | 🟠 Major | 🏗️ Heavy lift

Backfill legacy Pinata rows before using cid as local-IPFS state.

This only fixes new rows: record_pinata_cid binds cid = NULL on insert, but existing rows from the old Pinata path can still have Pinata’s CID in cid. Those rows are treated as locally pinned by has_ipfs_cid and filter_ipfs_pinned_oids, while record_pinned_cid’s COALESCE refuses to replace the stale value. Add a migration/backfill or explicit local-IPFS state before relying on these predicates. This is the same unresolved issue noted in the previous review.

#!/usr/bin/env bash
rg -n 'ALTER TABLE.*pinned_cids|UPDATE pinned_cids|pinata_cid|record_pinata_cid' \
  crates --glob '*.sql' --glob '*.rs' || true

Also applies to: 2268-2278, 2306-2319, 2322-2333

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@crates/gitlawb-node/src/db/mod.rs` around lines 2162 - 2171, Separate legacy
Pinata CIDs from local-IPFS state before the predicates and upsert logic use
pinned_cids.cid: add a migration/backfill or explicit local-IPFS indicator that
identifies existing rows created by record_pinata_cid and clears or excludes
their stale cid values. Update record_pinned_cid, has_ipfs_cid, and
filter_ipfs_pinned_oids to use that corrected state, ensuring legacy Pinata-only
rows can be replaced by a real local CID.

)
.bind(sha256_hex)
.bind(cid)
Expand Down Expand Up @@ -2267,6 +2265,17 @@ impl Db {
.collect())
}

/// Returns true if this object already has a local IPFS CID recorded.
pub async fn has_ipfs_cid(&self, sha256_hex: &str) -> Result<bool> {
let row = sqlx::query(
"SELECT COUNT(*) as cnt FROM pinned_cids WHERE sha256_hex = $1 AND cid IS NOT NULL",
)
.bind(sha256_hex)
.fetch_one(&self.pool)
.await?;
Ok(row.get::<i64, _>("cnt") > 0)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

/// Returns true if this object already has a Pinata CID recorded.
pub async fn has_pinata_cid(&self, sha256_hex: &str) -> Result<bool> {
let row = sqlx::query(
Expand All @@ -2278,17 +2287,50 @@ impl Db {
Ok(row.get::<i64, _>("cnt") > 0)
}

/// Given a list of sha256_hex values, returns the subset that already have
/// a Pinata CID recorded. Used by the reconciliation sweep to skip objects
/// that Pinata has already handled.
pub async fn filter_pinata_pinned_oids(&self, oids: &[String]) -> Result<Vec<String>> {
if oids.is_empty() {
return Ok(Vec::new());
}
let rows = sqlx::query(
"SELECT sha256_hex FROM pinned_cids WHERE sha256_hex = ANY($1) AND pinata_cid IS NOT NULL",
)
.bind(oids)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|r| r.get("sha256_hex")).collect())
}

/// Given a list of sha256_hex values, returns the subset that already have
/// a local IPFS CID recorded. Used by the reconciliation sweep to skip objects
/// that local IPFS has already handled.
pub async fn filter_ipfs_pinned_oids(&self, oids: &[String]) -> Result<Vec<String>> {
if oids.is_empty() {
return Ok(Vec::new());
}
let rows = sqlx::query(
"SELECT sha256_hex FROM pinned_cids WHERE sha256_hex = ANY($1) AND cid IS NOT NULL",
)
.bind(oids)
.fetch_all(&self.pool)
.await?;
Ok(rows.into_iter().map(|r| r.get("sha256_hex")).collect())
}

/// Record the Pinata CID for a git object.
/// Inserts the row if it doesn't exist (objects pinned directly to Pinata
/// without a prior local IPFS pin get cid = pinata_cid).
/// `cid` is left NULL for new rows so that `has_ipfs_cid` (which checks
/// `cid IS NOT NULL`) correctly distinguishes local IPFS state from
/// Pinata-only state.
pub async fn record_pinata_cid(&self, sha256_hex: &str, pinata_cid: &str) -> Result<()> {
sqlx::query(
"INSERT INTO pinned_cids (sha256_hex, cid, pinned_at, pinata_cid)
VALUES ($1, $2, $3, $4)
ON CONFLICT(sha256_hex) DO UPDATE SET pinata_cid = EXCLUDED.pinata_cid",
)
.bind(sha256_hex)
.bind(pinata_cid) // fallback local cid if row is new
.bind(Option::<&str>::None) // cid is NULL for Pinata-only new rows
.bind(Utc::now().to_rfc3339())
.bind(pinata_cid)
.execute(&self.pool)
Expand Down
7 changes: 4 additions & 3 deletions crates/gitlawb-node/src/ipfs_pin.rs
Original file line number Diff line number Diff line change
Expand Up @@ -107,12 +107,13 @@ pub async fn pin_new_objects(
let mut pinned = Vec::new();

for sha in object_list {
// Skip if already pinned
match db.is_pinned(&sha).await {
// Skip if already pinned to local IPFS (checks cid column,
// which is NULL for Pinata-only rows so those will be retried).
match db.has_ipfs_cid(&sha).await {
Ok(true) => continue,
Ok(false) => {}
Err(e) => {
tracing::warn!(sha = %sha, err = %e, "DB error checking pinned status");
tracing::warn!(sha = %sha, err = %e, "DB error checking IPFS pinned status");
continue;
}
}
Expand Down
14 changes: 14 additions & 0 deletions crates/gitlawb-node/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ mod operator;
mod p2p;
mod pinata;
mod rate_limit;
mod reconciliation;
mod server;
mod state;
mod sync;
Expand Down Expand Up @@ -498,6 +499,19 @@ async fn main() -> Result<()> {
info!("auto-sync worker started");
}

// Periodic reconciliation sweep: re-derives pin/seal sets and fills gaps
// so a dropped replication job never means data loss.
{
let db = state.db.clone();
let config = Arc::clone(&state.config);
let http_client = Arc::clone(&state.http_client);
let node_keypair = Arc::clone(&state.node_keypair);
let node_did = state.node_did.clone();
let shutdown_rx = state.subscribe_shutdown();
reconciliation::spawn(db, config, http_client, node_keypair, node_did, shutdown_rx);
info!("reconciliation sweep worker started");
}

// On-chain operator setup: verify stake + spawn heartbeat loop
if !state.config.contract_node_staking.is_empty()
&& !state.config.operator_private_key.is_empty()
Expand Down
47 changes: 45 additions & 2 deletions crates/gitlawb-node/src/metrics.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,9 @@
//! `gitlawb_pack_size_bytes`
//! * a single `gitlawb_info{version, did}` gauge = 1, for joins/dashboards
//! * currently-connected peer count — `gitlawb_peers_connected`
//! * reconciliation sweep gaps found and filled —
//! `gitlawb_reconciliation_gaps_found_total` /
//! `gitlawb_reconciliation_gaps_filled_total`
//!
//! All metrics live in a single process-wide registry initialized by
//! [`init`]. Increment helpers (`record_push`, `record_auth_failure`, ...)
Expand All @@ -33,8 +36,8 @@
use std::sync::OnceLock;

use prometheus::{
Encoder, Histogram, HistogramOpts, IntCounterVec, IntGauge, IntGaugeVec, Opts, Registry,
TextEncoder,
Encoder, Histogram, HistogramOpts, IntCounter, IntCounterVec, IntGauge, IntGaugeVec, Opts,
Registry, TextEncoder,
};

/// The single, process-wide metrics registry. Initialized by [`init`].
Expand All @@ -51,6 +54,8 @@ static SYNC_PROCESSED: OnceLock<IntCounterVec> = OnceLock::new();
static WEBHOOK_DELIVERIES: OnceLock<IntCounterVec> = OnceLock::new();
static PACK_SIZE: OnceLock<Histogram> = OnceLock::new();
static PEERS_CONNECTED: OnceLock<IntGauge> = OnceLock::new();
static RECONCILIATION_GAPS_FOUND: OnceLock<IntCounter> = OnceLock::new();
static RECONCILIATION_GAPS_FILLED: OnceLock<IntCounter> = OnceLock::new();

/// One-time initializer. Builds the registry, registers every metric,
/// and sets the constant `gitlawb_info` gauge. Idempotent — calling
Expand Down Expand Up @@ -197,6 +202,30 @@ pub fn init(version: &str, node_did: &str) {
.set(peers_connected)
.expect("set PEERS_CONNECTED once");

let gaps_found = IntCounter::with_opts(Opts::new(
"gitlawb_reconciliation_gaps_found_total",
"Total reconciliation sweep gaps detected (objects that should be pinned but are not)",
))
.expect("gitlawb_reconciliation_gaps_found_total definition");
registry
.register(Box::new(gaps_found.clone()))
.expect("register gitlawb_reconciliation_gaps_found_total");
RECONCILIATION_GAPS_FOUND
.set(gaps_found)
.expect("set RECONCILIATION_GAPS_FOUND once");

let gaps_filled = IntCounter::with_opts(Opts::new(
"gitlawb_reconciliation_gaps_filled_total",
"Total reconciliation sweep gaps successfully filled (objects pinned by the sweep)",
))
.expect("gitlawb_reconciliation_gaps_filled_total definition");
registry
.register(Box::new(gaps_filled.clone()))
.expect("register gitlawb_reconciliation_gaps_filled_total");
RECONCILIATION_GAPS_FILLED
.set(gaps_filled)
.expect("set RECONCILIATION_GAPS_FILLED once");

REGISTRY
.set(registry)
.expect("set REGISTRY once (init must be called exactly once)");
Expand Down Expand Up @@ -264,6 +293,20 @@ pub fn set_peers_connected(count: i64) {
}
}

/// Record reconciliation sweep gaps found (objects that should be pinned but are not).
pub fn record_reconciliation_gaps_found(count: u64) {
if let Some(c) = RECONCILIATION_GAPS_FOUND.get() {
c.inc_by(count);
}
}

/// Record reconciliation sweep gaps filled (objects successfully pinned by the sweep).
pub fn record_reconciliation_gaps_filled(count: u64) {
if let Some(c) = RECONCILIATION_GAPS_FILLED.get() {
c.inc_by(count);
}
}

/// Encode the registry as the standard Prometheus text exposition format.
/// Returns an error if `init` was never called.
pub fn encode() -> Result<String, prometheus::Error> {
Expand Down
Loading
Loading