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
59 changes: 59 additions & 0 deletions crates/dig-node-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10227,6 +10227,65 @@ mod tests {
);
}

/// **Proves:** `funded.observed_at` dates the CONSULTATION (the oldest per-item chain read),
/// never the assembly (dig_ecosystem#3323). The handler's pre-loop clock predates every read
/// `port.distributor_report` performs, so stamping it there understates staleness; the fix
/// folds the reports' own `observed_at` and keeps the oldest (a collection is only as fresh
/// as its stalest member). `claimable.observed_at` is untouched — no read happens for it, so
/// the handler's own clock remains the honest answer for that `NotConsulted` arm.
#[test]
fn list_reward_distributors_funded_observed_at_is_the_oldest_report_stamp() {
let state_dir = tempfile::tempdir().unwrap();
let registry =
crate::rewards::funded::FundedDistributorRegistry::with_state_dir(state_dir.path());
for launcher_id in [[0x55u8; 32], [0x66u8; 32]] {
assert_eq!(
registry.record(&crate::rewards::funded::FundedDistributor {
launcher_id,
store_id: None,
}),
crate::rewards::funded::RecordOutcome::Recorded
);
}
let (node, _td) = test_node(None);
assert!(node.install_funded_distributor_registry(registry));

let report_a = sample_distributor_report(0x55, vec![]);
let report_b = sample_distributor_report(0x66, vec![]);
assert!(
node.install_reward_chain_port(Arc::new(FakeRewardsChainPort {
reports: std::collections::HashMap::from([
([0x55u8; 32], Ok(report_a.clone())),
([0x66u8; 32], Ok(report_b.clone())),
]),
}))
);

let resp = rt().block_on(handle_rpc(
&node,
json!({"jsonrpc":"2.0","id":1,"method":"dig.listRewardDistributors"}),
crate::download::ReadOrigin::Local,
crate::download::RequestProvenance::FirstParty,
));
assert_eq!(resp["result"]["funded"]["outcome"], json!("consulted"));
assert_eq!(
resp["result"]["funded"]["observed_at"],
json!(1_700_200_000u64 + 0x55)
);
assert_eq!(
resp["result"]["funded"]["observed_at"],
json!(report_a.observed_at)
);
assert_ne!(
resp["result"]["funded"]["observed_at"],
json!(report_b.observed_at)
);
assert_eq!(
resp["result"]["claimable"]["outcome"],
json!("not_consulted")
);
}

/// **Proves:** a funded identity whose per-item chain report fails refuses the WHOLE call
/// (the same `ChainPortError` response the sibling reward-distributor handlers use), rather
/// than emitting a partial list or a fabricated ref (dig_ecosystem#3308/#3309).
Expand Down
94 changes: 91 additions & 3 deletions crates/dig-node-core/src/rewards/funded.rs
Original file line number Diff line number Diff line change
Expand Up @@ -213,7 +213,7 @@ impl FundedDistributorRegistry {
},
Ok(None) => FundedDistributorsRead::NotConfigured(self.absent_record_reason()),
Err(LoadFailure::Corrupt(reason)) => self.report_corrupt(&path, &reason),
Err(LoadFailure::Io(error)) => FundedDistributorsRead::IoFailed { path, error },
Err(LoadFailure::Io(error)) => self.report_io_failed(path, error),
}
}

Expand All @@ -231,7 +231,7 @@ impl FundedDistributorRegistry {
},
Ok(None) => Vec::new(),
Err(LoadFailure::Corrupt(reason)) => return self.refuse_corrupt(&path, &reason),
Err(LoadFailure::Io(error)) => return RecordOutcome::IoFailed { path, error },
Err(LoadFailure::Io(error)) => return self.refuse_io_failed(path, error),
};

let outcome = match merge(&mut set, distributor) {
Expand All @@ -243,7 +243,7 @@ impl FundedDistributorRegistry {
}
match self.save(&path, &set) {
Ok(()) => outcome,
Err(error) => RecordOutcome::IoFailed { path, error },
Err(error) => self.refuse_io_failed(path, error),
}
}

Expand Down Expand Up @@ -327,6 +327,31 @@ impl FundedDistributorRegistry {
}
}

/// Log and report an unreadable/unwritable record to a READER. Mirrors [`Self::report_corrupt`]
/// so an operator sees the same signal for either fault: `PersistedStateCorrupt` has logged
/// since v0.257.0, but `IoFailed` was constructed bare beside it, and `dispatch.rs` collapses
/// both into one wire `NotConsulted` with no log of its own (dig_ecosystem#3324) — so the
/// operator saw nothing for a permissions error, a missing mount, or any other I/O fault.
fn report_io_failed(&self, path: PathBuf, error: String) -> FundedDistributorsRead {
tracing::error!(
path = %path.display(),
error,
"the funded-distributor record could not be read"
);
FundedDistributorsRead::IoFailed { path, error }
}

/// Log and report an unreadable/unwritable record to a WRITER. Mirrors [`Self::refuse_corrupt`];
/// see [`Self::report_io_failed`] for why this arm was silent (dig_ecosystem#3324).
fn refuse_io_failed(&self, path: PathBuf, error: String) -> RecordOutcome {
tracing::error!(
path = %path.display(),
error,
"the funded-distributor record could not be read or written"
);
RecordOutcome::IoFailed { path, error }
}

/// Copy the corrupt record beside itself for the operator, leaving the original in place.
/// `None` when the copy failed — the corrupt verdict does not depend on it.
fn quarantine(&self, path: &Path) -> Option<PathBuf> {
Expand Down Expand Up @@ -817,6 +842,69 @@ mod tests {
);
}

/// **Proves:** an `IoFailed` read logs the path and the error, mirroring
/// [`FundedDistributorRegistry::report_corrupt`]'s `tracing::error!` (dig_ecosystem#3324 — the
/// corrupt arm has logged since v0.257.0; the I/O arm was constructed bare, so an operator saw
/// nothing for a permissions error or a missing mount). Forces `IoFailed` by making the record
/// PATH a directory, so `read_to_string` fails with a non-`NotFound` error — the wildcard
/// `PersistedStateCorrupt`/`IoFailed` collapse this fixture must not trip is
/// `dispatch.rs`'s wire mapping, untouched here; this test drives the registry directly (not
/// `handle_rpc`), because `rt()` may run the handler on a worker thread a thread-scoped
/// subscriber never sees.
#[test]
fn an_io_failed_read_is_logged_with_its_path_and_error() {
/// An in-memory sink a `tracing_subscriber::fmt` layer writes formatted records into.
#[derive(Clone, Default)]
struct LogCapture(std::sync::Arc<std::sync::Mutex<Vec<u8>>>);

impl std::io::Write for LogCapture {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}

impl<'a> tracing_subscriber::fmt::MakeWriter<'a> for LogCapture {
type Writer = LogCapture;
fn make_writer(&'a self) -> Self::Writer {
self.clone()
}
}

let dir = TempDir::new().expect("temp dir");
std::fs::create_dir(record_path(&dir)).expect("make the record path a directory");
let registry = FundedDistributorRegistry::with_state_dir(dir.path());

let buffer = LogCapture::default();
let subscriber = tracing_subscriber::fmt()
.with_ansi(false) // plain text: the assertions read the fields as an operator would
.with_writer(buffer.clone())
.finish();

// Scoped to this thread only, deliberately: `read()` runs synchronously here, never on a
// worker thread a thread-scoped subscriber would miss (dig_ecosystem#3324's note on why
// this test does not go through `handle_rpc`).
let read = tracing::subscriber::with_default(subscriber, || registry.read());

assert!(
matches!(read, FundedDistributorsRead::IoFailed { .. }),
"a directory at the record path must fail as IoFailed, got {read:?}"
);
let logged = String::from_utf8(buffer.0.lock().unwrap().clone()).unwrap();
let expected_path = record_path(&dir);
assert!(
logged.contains(&expected_path.display().to_string()),
"log did not name the record path: {logged}"
);
assert!(
logged.contains("error"),
"log did not name the error field: {logged}"
);
}

/// **Catches:** a regression on the not-an-answer half of [`FundedDistributorsRead`] listed
/// below reporting a renderable set. It does NOT by itself catch a future variant added to the
/// enum — that would need adding here too. The guard that actually forces the issue is
Expand Down
15 changes: 14 additions & 1 deletion crates/dig-node-core/src/seams/dig_rpc/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1167,6 +1167,12 @@ impl RpcDispatch for Node {
};

let mut funded_refs = Vec::with_capacity(identities.len());
// `observed_at` must date the CONSULTATION, never the assembly (dig_ecosystem#3323,
// dig-rpc-protocol SPEC §4.4.2's first bullet): each `report.observed_at` postdates
// the pre-loop `now` above (the chain port stamps it after its own read, uncached),
// so stamping the handler's clock here understates staleness. Fold the reports' own
// stamps and keep the OLDEST — a collection is only as fresh as its stalest member.
let mut oldest_observed_at: Option<u64> = None;
for identity in identities {
let Some(port) = node.reward_chain_port() else {
return reward_chain_port_absent_response(&id);
Expand All @@ -1179,6 +1185,10 @@ impl RpcDispatch for Node {
Ok(report) => report,
Err(e) => return reward_chain_port_error_response(&id, &e),
};
oldest_observed_at = Some(match oldest_observed_at {
Some(oldest) => oldest.min(report.observed_at),
None => report.observed_at,
});
funded_refs.push(dig_rpc_protocol::types::RewardDistributorRef {
launcher_id: hex::encode(report.launcher_id),
store_id: hex::encode(report.store_id),
Expand All @@ -1188,7 +1198,10 @@ impl RpcDispatch for Node {

let result = dig_rpc_protocol::types::ListRewardDistributorsResult {
funded: dig_rpc_protocol::types::Half::Consulted {
observed_at: now,
// No reads happened for `FundsNothing` (empty `identities`), so the
// pre-loop handler clock is the honest stamp for that case — it IS the
// consultation. Both `NotConsulted` arms above keep `now` unchanged.
observed_at: oldest_observed_at.unwrap_or(now),
items: funded_refs,
},
claimable: dig_rpc_protocol::types::Half::NotConsulted { observed_at: now },
Expand Down
Loading