feat(storage): run family integrity repair

Newly migrated and steady-state identifier families need the same healing path, but remote availability must not gate relay startup. Start one background pass after migration and database initialization, repair through accepted clone URLs, and emit a bounded error for every unresolved family.

Add a durable identifier-scoped request queue and an integrity-check operator command so the live process performs manual checks under the existing family leases. Check-only and repair requests cover every object format and owner or /prs/ view for the identifier, and repeated requests coalesce safely.

Document the family-level integrity boundary and make migration's handoff deliberately small: structural conversion remains fail-closed, then the ordinary new-model pass handles pre-existing damage. This does not add legacy backup archaeology, garbage collection, periodic remote repair, or S3 behavior.

Validated with cargo fmt, strict locked workspace Clippy, the complete locked workspace test suite, focused request-worker tests, and the command help path.
This commit is contained in:
DanConwayDev
2026-08-18 11:40:40 +00:00
parent 4aaf234caa
commit c68ae50e13
6 changed files with 453 additions and 4 deletions
@@ -302,6 +302,42 @@ View lifecycle locks remain responsible for “may this path be removed?” The
family lock is responsible for “is this object inventory and manifest update
atomic?” Neither lock grants authorization.
## Integrity and healing
Integrity is defined for an object-format/identifier family, not for one
owner path. A pass checks the family pack set and object graph, then checks
every owner and `/prs/` view with the same identifier for the correct alternate
and for refs whose targets are available in the family. Multiple independent
histories in one identifier are valid, and unreachable objects are not an
error: retained delete-state and rollback data intentionally remain present.
The relay starts one non-blocking pass after migration, database
initialization, and construction of the hardened outbound Git client. Broken
alternate wiring is repaired locally. For missing OIDs, the pass tries clone
URLs from accepted repository announcements, excluding this service and
applying the same SSRF policy, DNS pinning, credentials, process containment,
and missing-OID behavior used by proactive sync. It then runs the same
integrity check again. An unresolved family produces an `ERROR` log containing
bounded counts and OID/diagnostic samples; it does not make availability depend
on remote servers.
This is also the migration repair path. Migration remains a deterministic,
offline conversion that preserves every Git-readable object and the exact
legacy refs. Once those paths are thin family views, the ordinary family pass
can heal pre-existing missing objects. Unindexed legacy packs remain in the
migration backup; there is no separate legacy repair subsystem.
Operators can queue the same identifier-scoped check in the live process:
```console
ngit-grasp integrity-check --identifier example
ngit-grasp integrity-check --identifier example --repair
```
The command writes a durable request beneath `.grasp/integrity-requests/`.
The server consumes it while holding its normal in-process family locks, so a
manual repair cannot race an object-producing request in another view.
## Security and privacy trade-offs
- Sharing is restricted to one validated identifier, object format, and
+1
View File
@@ -32,6 +32,7 @@ How-to guides are **recipes** that show you how to solve specific problems or ac
- Prepare capacity and a release rollback point
- Run the automatic crash-safe launch migration
- Verify owner and `/prs/` repository views
- Check or repair one identifier family on demand
- Recover safely from an interrupted launch
---
+47
View File
@@ -9,6 +9,12 @@ been validated and before purgatory restoration, background sync, or HTTP
request handling. A failed migration stops startup; restarting resumes from its
fsynced journal.
After runtime database initialization, a non-blocking integrity pass checks
the resulting identifier families and views. It attempts to fetch missing
objects from clone URLs in accepted repository announcements and emits an
`ERROR` log for any family that remains unhealthy. Network repair never holds
up the listening service or changes whether the structural migration commits.
## Before deploying
1. Stop writes to the relay and take a filesystem snapshot of both the Git and
@@ -45,10 +51,51 @@ fetch, and push for both a normal repository and a `/prs/` route. Retain the
external snapshot and `.grasp/migration/backups/` for the rollback window. The
server never deletes those backups automatically.
Watch for the terminal startup-pass summary:
```text
Git identifier-family integrity startup pass completed
```
An `unresolved` or `failed` count above zero is accompanied by an `ERROR` log
for each affected identifier. This reports pre-existing missing data without
putting startup into a network-dependent restart loop.
## Check or repair one identifier on demand
Queue a read-only check for every object format and owner/`/prs/` view sharing
an identifier:
```console
ngit-grasp integrity-check \
--git-data-path /var/lib/ngit-grasp/git \
--identifier example
```
Add `--repair` to repair alternate wiring and try accepted clone servers for
missing OIDs:
```console
ngit-grasp integrity-check \
--git-data-path /var/lib/ngit-grasp/git \
--identifier example \
--repair
```
The command queues a durable request for the running relay rather than opening
or mutating a family from a second process. The worker normally consumes it
within five seconds and writes the result to the service log. A request queued
while the relay is stopped is processed after its next startup integrity pass.
Repeated requests for the same identifier and mode safely coalesce.
## Failure and rollback
- A migration failure is fail-closed. Fix the reported filesystem or Git error
and restart; the journal resumes the safe transition.
- A post-migration integrity repair failure is fail-open because the damage
predates conversion or arose after it. Inspect the identifier's `ERROR` log,
repair or update its listed clone sources, and queue `integrity-check
--repair` again.
- To roll the software release back, stop the service and restore the complete
pre-upgrade Git and relay-data snapshot together. Do not point an older
binary at migrated thin views.
+334 -3
View File
@@ -8,14 +8,75 @@
use std::collections::BTreeSet;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::Arc;
use std::time::Duration;
use anyhow::{anyhow, Context, Result};
use async_trait::async_trait;
use bitcoin_hashes::{sha256, Hash};
use clap::Args;
use nostr_sdk::prelude::{FromBech32, PublicKey};
use serde::{Deserialize, Serialize};
use tokio::task::JoinHandle;
use tracing::{error, info, warn};
use super::storage::{FamilyKey, LocalGitStorage, ObjectFormat};
use super::validate_repository_identifier;
const REQUEST_VERSION: u32 = 1;
const REQUEST_POLL_INTERVAL: Duration = Duration::from_secs(5);
/// Arguments for queueing an identifier-family integrity check in the live
/// relay process.
#[derive(Debug, Args)]
pub struct IntegrityCheckArgs {
/// Repository identifier (`d` tag value). All object formats and views for
/// this identifier are checked together.
#[arg(long)]
pub identifier: String,
/// Attempt repair using accepted repository clone URLs.
#[arg(long, default_value_t = false)]
pub repair: bool,
/// Git data path containing the `.grasp` family storage directory.
#[arg(long, env = "NGIT_GIT_DATA_PATH", default_value = "./data/git")]
pub git_data_path: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct IntegrityRequest {
version: u32,
identifier: String,
repair: bool,
}
/// Queue an integrity request for the running relay. Requests are durable and
/// idempotent; if the service is stopped, the next launch consumes them.
pub fn enqueue_manual_check(args: &IntegrityCheckArgs) -> Result<PathBuf> {
if !validate_repository_identifier(&args.identifier) {
return Err(anyhow!(
"invalid repository identifier {:?}",
args.identifier
));
}
let storage = LocalGitStorage::new(&args.git_data_path);
let directory = request_directory(&storage);
std::fs::create_dir_all(&directory)
.with_context(|| format!("create integrity request directory {}", directory.display()))?;
let request = IntegrityRequest {
version: REQUEST_VERSION,
identifier: args.identifier.clone(),
repair: args.repair,
};
let digest =
sha256::Hash::hash(format!("{}\0{}", request.identifier, request.repair).as_bytes());
let path = directory.join(format!("{digest}.json"));
crate::atomic_file::write(&path, &serde_json::to_vec_pretty(&request)?)
.with_context(|| format!("write integrity request {}", path.display()))?;
Ok(path)
}
/// A view ref whose target is absent from the identifier family.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MissingRefTarget {
@@ -79,6 +140,226 @@ impl FamilyRepairOutcome {
}
}
#[derive(Default)]
struct PassStats {
checked: usize,
healthy: usize,
repaired: usize,
unresolved: usize,
failed: usize,
}
/// Start the non-blocking integrity worker.
///
/// It checks and heals every installed family once after startup, then
/// consumes identifier-scoped requests written by [`enqueue_manual_check`].
pub fn spawn_integrity_worker(
storage: LocalGitStorage,
source: Arc<crate::purgatory::sync::RealSyncContext>,
) -> JoinHandle<()> {
tokio::spawn(async move {
run_startup_pass(&storage, source.as_ref()).await;
let first = tokio::time::Instant::now() + REQUEST_POLL_INTERVAL;
let mut interval = tokio::time::interval_at(first, REQUEST_POLL_INTERVAL);
loop {
interval.tick().await;
process_manual_requests(&storage, source.as_ref()).await;
}
})
}
async fn run_startup_pass<S: FamilyRepairSource + ?Sized>(storage: &LocalGitStorage, source: &S) {
let families = match discover_families(storage) {
Ok(families) => families,
Err(error) => {
error!(%error, "Git identifier-family integrity startup discovery failed");
return;
}
};
info!(
families = families.len(),
"Git identifier-family integrity startup pass started"
);
let mut stats = PassStats::default();
for key in families {
run_one_family(storage, &key, source, true, "startup", &mut stats).await;
tokio::task::yield_now().await;
}
info!(
checked = stats.checked,
healthy = stats.healthy,
repaired = stats.repaired,
unresolved = stats.unresolved,
failed = stats.failed,
"Git identifier-family integrity startup pass completed"
);
}
async fn process_manual_requests<S: FamilyRepairSource + ?Sized>(
storage: &LocalGitStorage,
source: &S,
) {
let requests = match load_requests(storage) {
Ok(requests) => requests,
Err(error) => {
error!(%error, "Read manual Git integrity requests");
return;
}
};
for (path, request) in requests {
let families = match discover_families(storage) {
Ok(families) => families
.into_iter()
.filter(|key| key.identifier == request.identifier)
.collect::<Vec<_>>(),
Err(error) => {
error!(
identifier = %request.identifier,
%error,
"Manual Git integrity family discovery failed"
);
remove_consumed_request(&path);
continue;
}
};
if families.is_empty() {
error!(
identifier = %request.identifier,
"Manual Git integrity request matched no identifier family"
);
remove_consumed_request(&path);
continue;
}
let mut stats = PassStats::default();
for key in families {
run_one_family(storage, &key, source, request.repair, "manual", &mut stats).await;
}
info!(
identifier = %request.identifier,
repair = request.repair,
checked = stats.checked,
healthy = stats.healthy,
repaired = stats.repaired,
unresolved = stats.unresolved,
failed = stats.failed,
"Manual Git identifier-family integrity request completed"
);
remove_consumed_request(&path);
}
}
async fn run_one_family<S: FamilyRepairSource + ?Sized>(
storage: &LocalGitStorage,
key: &FamilyKey,
source: &S,
repair: bool,
trigger: &'static str,
stats: &mut PassStats,
) {
stats.checked += 1;
let outcome = match check_and_repair_family(storage, key, source, repair).await {
Ok(outcome) => outcome,
Err(error) => {
stats.failed += 1;
error!(
identifier = %key.identifier,
object_format = %key.object_format,
trigger,
repair,
%error,
"Git identifier-family integrity check failed"
);
return;
}
};
if outcome.final_report.is_healthy() {
if outcome.repaired() {
stats.repaired += 1;
info!(
identifier = %key.identifier,
object_format = %key.object_format,
trigger,
sources_tried = outcome.sources_tried.len(),
"Git identifier family repaired"
);
} else {
stats.healthy += 1;
}
return;
}
stats.unresolved += 1;
let missing_sample = outcome
.final_report
.missing_oids
.iter()
.take(16)
.cloned()
.collect::<Vec<_>>()
.join(",");
let diagnostic = outcome
.final_report
.fsck_diagnostics
.first()
.map(String::as_str)
.unwrap_or("");
error!(
identifier = %key.identifier,
object_format = %key.object_format,
trigger,
repair,
views = outcome.final_report.views.len(),
missing_oids = outcome.final_report.missing_oids.len(),
missing_sample,
missing_ref_targets = outcome.final_report.missing_ref_targets.len(),
invalid_alternates = outcome.final_report.invalid_alternates.len(),
pack_errors = outcome.final_report.pack_errors.len(),
fsck_diagnostics = outcome.final_report.fsck_diagnostics.len(),
diagnostic,
sources_tried = outcome.sources_tried.len(),
source_failures = outcome.source_failures.len(),
"Git identifier family remains unhealthy after integrity check"
);
}
fn request_directory(storage: &LocalGitStorage) -> PathBuf {
storage.internal_path().join("integrity-requests")
}
fn load_requests(storage: &LocalGitStorage) -> Result<Vec<(PathBuf, IntegrityRequest)>> {
let directory = request_directory(storage);
let entries = match std::fs::read_dir(&directory) {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(error) => return Err(error.into()),
};
let mut paths = entries
.filter_map(|entry| entry.ok())
.filter(|entry| entry.file_type().is_ok_and(|kind| kind.is_file()))
.map(|entry| entry.path())
.filter(|path| path.extension().and_then(|value| value.to_str()) == Some("json"))
.collect::<Vec<_>>();
paths.sort();
let mut requests = Vec::new();
for path in paths {
let request: IntegrityRequest = serde_json::from_slice(&std::fs::read(&path)?)
.with_context(|| format!("parse integrity request {}", path.display()))?;
if request.version != REQUEST_VERSION
|| !validate_repository_identifier(&request.identifier)
{
return Err(anyhow!("invalid integrity request {}", path.display()));
}
requests.push((path, request));
}
Ok(requests)
}
fn remove_consumed_request(path: &Path) {
if let Err(error) = std::fs::remove_file(path) {
warn!(request = %path.display(), %error, "Remove consumed Git integrity request");
}
}
/// Inspect a family, repair local view wiring, fetch missing OIDs from every
/// accepted clone source as needed, then inspect it again.
pub async fn check_and_repair_family<S: FamilyRepairSource + ?Sized>(
@@ -425,12 +706,12 @@ fn inspect_connectivity(
}
fn oid_exists(repo: &Path, oid: &str) -> Result<bool> {
let status = Command::new("git")
let output = Command::new("git")
.args(["cat-file", "-e", oid])
.current_dir(repo)
.status()
.output()
.with_context(|| format!("check object {oid} in {}", repo.display()))?;
Ok(status.success())
Ok(output.status.success())
}
#[cfg(test)]
@@ -606,6 +887,56 @@ mod tests {
assert!(report.pack_errors[0].contains("pack has no index"));
}
#[test]
fn manual_requests_are_durable_and_identifier_scoped() {
let temp = tempfile::tempdir().unwrap();
let args = IntegrityCheckArgs {
identifier: "shared".to_owned(),
repair: true,
git_data_path: temp.path().join("git").to_string_lossy().into_owned(),
};
let path = enqueue_manual_check(&args).unwrap();
let storage = LocalGitStorage::new(&args.git_data_path);
let requests = load_requests(&storage).unwrap();
assert_eq!(requests.len(), 1);
assert_eq!(requests[0].0, path);
assert_eq!(requests[0].1.identifier, "shared");
assert!(requests[0].1.repair);
}
#[test]
fn manual_requests_reject_unsafe_identifiers() {
let temp = tempfile::tempdir().unwrap();
let args = IntegrityCheckArgs {
identifier: "../escape".to_owned(),
repair: false,
git_data_path: temp.path().join("git").to_string_lossy().into_owned(),
};
assert!(enqueue_manual_check(&args).is_err());
}
#[tokio::test]
async fn manual_requests_are_consumed_by_the_family_worker() {
let (temp, storage, _key, _view, _commit) = fixture();
let args = IntegrityCheckArgs {
identifier: "shared".to_owned(),
repair: false,
git_data_path: storage.git_data_path().to_string_lossy().into_owned(),
};
let request = enqueue_manual_check(&args).unwrap();
let source = LocalRepairSource {
url: String::new(),
family_objects: temp.path().join("unused"),
};
process_manual_requests(&storage, &source).await;
assert!(!request.exists());
}
#[tokio::test]
async fn repairs_missing_objects_from_an_identifier_clone_source() {
let (temp, storage, key, view, commit) = fixture();
+24 -1
View File
@@ -36,6 +36,9 @@ enum Cli {
///
/// This is an operator/admin maintenance command and is idempotent.
HoldingEject(nostr::lifecycle::HoldingEjectArgs),
/// Queue an identifier-family integrity check in the running relay.
IntegrityCheck(ngit_grasp::git::integrity::IntegrityCheckArgs),
}
#[tokio::main]
@@ -47,7 +50,13 @@ async fn main() -> Result<()> {
// If not, prepend the implicit "serve" subcommand so that clap routes to Cli::Serve
// and all relay flags are parsed normally (preserving backward compatibility).
let mut args: Vec<String> = std::env::args().collect();
let known_subcommands = ["serve", "cleanup-empty-repos", "holding-eject", "help"];
let known_subcommands = [
"serve",
"cleanup-empty-repos",
"holding-eject",
"integrity-check",
"help",
];
let has_subcommand = args.get(1).is_some_and(|a| {
known_subcommands.contains(&a.as_str())
|| matches!(a.as_str(), "-h" | "--help" | "-V" | "--version")
@@ -59,6 +68,20 @@ async fn main() -> Result<()> {
match Cli::parse_from(args) {
Cli::CleanupEmptyRepos(cleanup_args) => cleanup_empty_repos::run(&cleanup_args).await,
Cli::HoldingEject(eject_args) => nostr::lifecycle::run_holding_eject(eject_args).await,
Cli::IntegrityCheck(integrity_args) => {
let path = ngit_grasp::git::integrity::enqueue_manual_check(&integrity_args)?;
println!(
"Queued {} for identifier '{}' at {}",
if integrity_args.repair {
"integrity check and repair"
} else {
"integrity check"
},
integrity_args.identifier,
path.display()
);
Ok(())
}
Cli::Serve(config) => {
let mut config = *config;
config.relay_owner_nsec = Some(Config::load_relay_owner_key()?);
+11
View File
@@ -439,6 +439,17 @@ impl RelayServer {
outbound_credential_keys,
));
// Check the permanent identifier-family model after migration and
// database initialization. The pass is intentionally non-blocking:
// remote repair must not make availability depend on listed clone
// servers. It also owns filesystem-queued operator requests so all
// repair work remains inside the process holding family write leases.
background_tasks.push(git::integrity::spawn_integrity_worker(
git_storage,
sync_ctx.clone(),
));
info!("Git identifier-family integrity worker started");
// Create throttle manager for rate limiting remote git servers
// Default: 5 concurrent requests per domain, 60 requests per minute per domain
let throttle_manager = Arc::new(ThrottleManager::new(5, 60));