commit 5491fe73203eb07ab5b53c3ea37f1f1dfb9bf676
parent 7a3238e4979102c93d2f614849b70491774095da
Author: triesap <tyson@radroots.org>
Date: Tue, 30 Jun 2026 07:02:42 +0000
dvm: harden sdk adoption source guards
- add source guards for SDK-free DVM and processed-job authority
- make subscriber tests use deterministic subscribe hooks under parallel cargo test
- tighten guard checks around lib-owned feedback and transition-result builders
- validated with cargo test --workspace
Diffstat:
2 files changed, 226 insertions(+), 11 deletions(-)
diff --git a/src/features/trade_listing/subscriber.rs b/src/features/trade_listing/subscriber.rs
@@ -11,8 +11,8 @@ use radroots_events::kinds::{
};
use radroots_nostr::prelude::{
RadrootsNostrClient, RadrootsNostrEvent, RadrootsNostrFilter, RadrootsNostrKeys,
- RadrootsNostrKind, RadrootsNostrRelayPoolNotification, RadrootsNostrTag,
- radroots_nostr_tags_resolve,
+ RadrootsNostrKind, RadrootsNostrRelayPoolNotification, RadrootsNostrSubscriptionId,
+ RadrootsNostrTag, radroots_nostr_tags_resolve,
};
use tokio::sync::watch;
use tokio::time::sleep;
@@ -27,6 +27,8 @@ use crate::features::trade_validation_receipt::TradeValidationReceiptProverPolic
#[cfg(test)]
#[derive(Default)]
struct SubscriberTestHooks {
+ subscribe_results: std::collections::VecDeque<Result<RadrootsNostrSubscriptionId, ()>>,
+ unsubscribe_results: std::collections::VecDeque<()>,
notifications: std::collections::VecDeque<Result<RadrootsNostrRelayPoolNotification, ()>>,
delay_before_event_handle: std::collections::VecDeque<bool>,
resolve_tags_results: std::collections::VecDeque<
@@ -46,6 +48,24 @@ fn subscriber_test_hooks() -> &'static std::sync::Mutex<SubscriberTestHooks> {
}
#[cfg(test)]
+fn pop_subscribe_hook() -> Option<Result<RadrootsNostrSubscriptionId, ()>> {
+ subscriber_test_hooks()
+ .lock()
+ .unwrap_or_else(std::sync::PoisonError::into_inner)
+ .subscribe_results
+ .pop_front()
+}
+
+#[cfg(test)]
+fn pop_unsubscribe_hook() -> Option<()> {
+ subscriber_test_hooks()
+ .lock()
+ .unwrap_or_else(std::sync::PoisonError::into_inner)
+ .unsubscribe_results
+ .pop_front()
+}
+
+#[cfg(test)]
fn pop_notification_hook() -> Option<Result<RadrootsNostrRelayPoolNotification, ()>> {
subscriber_test_hooks()
.lock()
@@ -154,6 +174,29 @@ fn map_notification_recv_result(
result.map_err(|_| ())
}
+async fn subscribe_io(
+ client: &RadrootsNostrClient,
+ filter: RadrootsNostrFilter,
+) -> Result<RadrootsNostrSubscriptionId> {
+ #[cfg(test)]
+ if let Some(result) = pop_subscribe_hook() {
+ return result.map_err(|_| anyhow!("trade_listing subscriber subscribe failed"));
+ }
+ let subscription = client.subscribe(filter, None).await?;
+ Ok(subscription.val)
+}
+
+async fn unsubscribe_io(
+ client: &RadrootsNostrClient,
+ subscription_id: &RadrootsNostrSubscriptionId,
+) {
+ #[cfg(test)]
+ if pop_unsubscribe_hook().is_some() {
+ return;
+ }
+ client.unsubscribe(subscription_id).await;
+}
+
async fn handle_event_io(
event: RadrootsNostrEvent,
resolved_tags: Vec<RadrootsNostrTag>,
@@ -286,7 +329,7 @@ pub async fn subscriber(
return Ok(());
}
- let subscription = client.subscribe(filter, None).await?;
+ let subscription_id = subscribe_io(&client, filter).await?;
let mut notifications = client.notifications();
let mut notifications_closed = false;
@@ -323,7 +366,7 @@ pub async fn subscriber(
}
}
- client.unsubscribe(&subscription.val).await;
+ unsubscribe_io(&client, &subscription_id).await;
if notifications_closed {
return Err(anyhow!("trade_listing subscriber notifications closed"));
}
@@ -373,6 +416,16 @@ mod tests {
RadrootsNostrRelayPoolNotification::Shutdown
}
+ fn prime_subscription_hooks() {
+ let mut hooks = subscriber_test_hooks()
+ .lock()
+ .unwrap_or_else(std::sync::PoisonError::into_inner);
+ hooks
+ .subscribe_results
+ .push_back(Ok(RadrootsNostrSubscriptionId::new("sub-hook")));
+ hooks.unsubscribe_results.push_back(());
+ }
+
fn shared_runtime() -> TradeListingRuntime {
TradeListingRuntime::new()
}
@@ -514,11 +567,11 @@ mod tests {
}
#[tokio::test]
- async fn subscriber_can_stop_after_start_when_relay_is_present() {
+ async fn subscriber_can_stop_after_start_when_subscription_is_ready() {
let _guard = test_guard().await;
let keys = RadrootsNostrKeys::generate();
let client = RadrootsNostrClient::new(keys.clone());
- let _ = client.add_relay("wss://relay.example.com").await;
+ prime_subscription_hooks();
let (tx, rx) = watch::channel(false);
let join = tokio::spawn(subscriber(
client,
@@ -537,7 +590,7 @@ mod tests {
let _guard = test_guard().await;
let keys = RadrootsNostrKeys::generate();
let client = RadrootsNostrClient::new(keys.clone());
- let _ = client.add_relay("wss://relay.example.com").await;
+ prime_subscription_hooks();
subscriber_test_hooks()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
@@ -556,7 +609,7 @@ mod tests {
let _guard = test_guard().await;
let keys = RadrootsNostrKeys::generate();
let client = RadrootsNostrClient::new(keys.clone());
- let _ = client.add_relay("wss://relay.example.com").await;
+ prime_subscription_hooks();
subscriber_test_hooks()
.lock()
@@ -583,7 +636,7 @@ mod tests {
let _guard = test_guard().await;
let keys = RadrootsNostrKeys::generate();
let client = RadrootsNostrClient::new(keys.clone());
- let _ = client.add_relay("wss://relay.example.com").await;
+ prime_subscription_hooks();
subscriber_test_hooks()
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
@@ -614,7 +667,7 @@ mod tests {
let _guard = test_guard().await;
let keys = RadrootsNostrKeys::generate();
let client = RadrootsNostrClient::new(keys.clone());
- let _ = client.add_relay("wss://relay.example.com").await;
+ prime_subscription_hooks();
let mut hooks = subscriber_test_hooks()
.lock()
@@ -656,7 +709,7 @@ mod tests {
let _guard = test_guard().await;
let keys = RadrootsNostrKeys::generate();
let client = RadrootsNostrClient::new(keys.clone());
- let _ = client.add_relay("wss://relay.example.com").await;
+ prime_subscription_hooks();
let mut hooks = subscriber_test_hooks()
.lock()
diff --git a/tests/source_guards.rs b/tests/source_guards.rs
@@ -0,0 +1,162 @@
+use std::fs;
+use std::path::Path;
+
+#[test]
+fn rhi_manifest_has_no_sdk_dependency() {
+ let manifest = read_repo_file("Cargo.toml");
+
+ assert!(
+ !manifest.contains("radroots_sdk"),
+ "RHI must not depend on radroots_sdk"
+ );
+}
+
+#[test]
+fn rhi_dvm_transition_and_feedback_paths_use_trade_dvm_contract() {
+ let listing_dvm = read_repo_file("src/features/trade_listing/handlers/dvm.rs");
+ let receipt_worker = read_repo_file("src/features/trade_validation_receipt.rs");
+ let listing_feedback_segment = source_segment(
+ &listing_dvm,
+ "pub async fn handle_error(",
+ "\n#[cfg(test)]\n#[cfg_attr(coverage_nightly, coverage(off))]\nmod tests",
+ );
+
+ assert!(
+ listing_dvm.contains(
+ "use radroots_trade::dvm::{RadrootsTradeDvmFeedbackStatus, build_job_feedback_tags};"
+ ),
+ "trade listing DVM handler must import radroots_trade::dvm feedback contract"
+ );
+ assert!(
+ listing_feedback_segment.contains("build_job_feedback_tags("),
+ "trade listing DVM handler must build feedback tags through radroots_trade::dvm"
+ );
+
+ for required in [
+ "RadrootsTradeTransitionProofRequestEnvelope",
+ "RadrootsTradeTransitionProofResultBinding",
+ ] {
+ assert!(
+ receipt_worker.contains(required),
+ "trade validation receipt worker must use radroots_trade::dvm transition contract `{required}`"
+ );
+ }
+
+ let result_segment = source_segment(
+ &receipt_worker,
+ "fn result_tags_from_dvm(",
+ "fn expected_receipt_binding",
+ );
+ for required in [
+ "parse_transition_proof_request_event(",
+ "build_transition_proof_result_tags(",
+ ] {
+ assert!(
+ result_segment.contains(required),
+ "trade validation receipt worker must delegate result tags through radroots_trade::dvm `{required}`"
+ );
+ }
+ assert!(
+ !result_segment.contains("vec!["),
+ "transition proof result tags must be delegated to radroots_trade::dvm"
+ );
+}
+
+#[test]
+fn rhi_sources_do_not_import_removed_sdk_or_protocol_bypasses() {
+ for (path, source) in rust_sources_under("src") {
+ for forbidden in [
+ "radroots_sdk",
+ "radroots_sdk::protocol::order",
+ "SdkDvmInventoryBinWitness",
+ "TradeProtocolClient",
+ "KIND_TRADE_LISTING_VALIDATE_REQ",
+ "KIND_WORKER_TRADE_TRANSITION_PROOF_REQ",
+ ] {
+ assert!(
+ !source.contains(forbidden),
+ "{path} contains forbidden SDK adoption bypass `{forbidden}`"
+ );
+ }
+ }
+}
+
+#[test]
+fn rhi_processed_job_state_is_durable_workflow_authority() {
+ let state = read_repo_file("src/features/trade_listing/state.rs");
+ let receipt_worker = read_repo_file("src/features/trade_validation_receipt.rs");
+
+ for required in [
+ "rhi_processed_jobs: HashMap<String, RhiProcessedJobState>",
+ "pub fn rhi_processed_job(&self, request_id: &str)",
+ "pub fn upsert_rhi_processed_job(&mut self, job: RhiProcessedJobState)",
+ ] {
+ assert!(
+ state.contains(required),
+ "RHI state must retain processed-job storage contract `{required}`"
+ );
+ }
+
+ for required in [
+ "fn processed_job_for_request(",
+ "async fn processed_job_action(",
+ "mark_job_completed(",
+ "RhiProcessedJobStatus::Completed",
+ ] {
+ assert!(
+ receipt_worker.contains(required),
+ "RHI receipt worker must retain processed-job workflow guard `{required}`"
+ );
+ }
+}
+
+fn read_repo_file(relative_path: &str) -> String {
+ let path = Path::new(env!("CARGO_MANIFEST_DIR")).join(relative_path);
+ fs::read_to_string(path.as_path())
+ .unwrap_or_else(|error| panic!("failed to read {}: {error}", path.display()))
+}
+
+fn source_segment<'a>(source: &'a str, start: &str, end: &str) -> &'a str {
+ let start_index = source.find(start).expect("source segment start");
+ let end_index = source[start_index..]
+ .find(end)
+ .map(|index| start_index + index)
+ .expect("source segment end");
+ &source[start_index..end_index]
+}
+
+fn rust_sources_under(relative_root: &str) -> Vec<(String, String)> {
+ let root = Path::new(env!("CARGO_MANIFEST_DIR"));
+ let mut paths = Vec::new();
+ collect_rust_sources(root.join(relative_root).as_path(), &mut paths);
+ paths.sort();
+ paths
+ .into_iter()
+ .map(|path| {
+ let relative_path = path
+ .strip_prefix(root)
+ .expect("source under manifest root")
+ .to_string_lossy()
+ .replace('\\', "/");
+ let source = fs::read_to_string(path.as_path())
+ .unwrap_or_else(|error| panic!("failed to read {}: {error}", path.display()));
+ (relative_path, source)
+ })
+ .collect()
+}
+
+fn collect_rust_sources(path: &Path, paths: &mut Vec<std::path::PathBuf>) {
+ if path.is_file() {
+ if path.extension().and_then(|extension| extension.to_str()) == Some("rs") {
+ paths.push(path.to_path_buf());
+ }
+ return;
+ }
+
+ for entry in fs::read_dir(path)
+ .unwrap_or_else(|error| panic!("failed to read {}: {error}", path.display()))
+ {
+ let entry = entry.expect("source entry");
+ collect_rust_sources(entry.path().as_path(), paths);
+ }
+}