diff --git a/Cargo.lock b/Cargo.lock index 3686cab4..8a4ad48e 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -9792,7 +9792,6 @@ dependencies = [ "serde_json", "sp-core", "sp-runtime", - "storage-client", "storage-primitives", "subxt", "subxt-signer", diff --git a/provider-node/Cargo.toml b/provider-node/Cargo.toml index d9eacdf3..57e1212c 100644 --- a/provider-node/Cargo.toml +++ b/provider-node/Cargo.toml @@ -9,7 +9,6 @@ description = "Off-chain provider node for scalable Web3 storage" [dependencies] provider-negotiation = { workspace = true } -storage-client = { workspace = true } storage-primitives = { workspace = true, features = ["serde", "std"] } tokio = { workspace = true } axum = { workspace = true } diff --git a/provider-node/src/chain_state_coordinator.rs b/provider-node/src/chain_state_coordinator.rs index b162d0c8..3fea4875 100644 --- a/provider-node/src/chain_state_coordinator.rs +++ b/provider-node/src/chain_state_coordinator.rs @@ -12,28 +12,28 @@ //! chain's replay window. `None` until the provider is registered. //! //! [`ChainStateCoordinator`] is the **only writer** for all four fields. It -//! drives a [`BlockSubscriberStream`] in a reconnect loop; on every relevant -//! provider event it re-fetches the full `ProviderInfo` so `committed_bytes`, -//! `stake`, and all settings stay current — no field-patching, no partial -//! updates, no second writer. +//! drives a finalized-block subscription on its own subxt connection in a +//! reconnect loop; on every relevant provider event it re-fetches the full +//! `ProviderInfo` so `committed_bytes`, `stake`, and all settings stay +//! current — no field-patching, no partial updates, no second writer. use crate::negotiate::NonceCounter; use crate::storage::{NonceStore, NullNonceStore}; +use crate::types::ProviderInfo; +use crate::Error; use async_trait::async_trait; use parking_lot::RwLock; -use sp_core::H256; use sp_runtime::AccountId32; use std::sync::atomic::AtomicU32; use std::sync::Arc; use std::time::Duration; -use storage_client::discovery::ProviderInfo; -use storage_client::{ - BlockSubscriberStream, ClientConfig, ClientError, EventParser, ProviderClient, StorageEvent, - StorageProviderEventParser, -}; -use subxt::ext::futures::StreamExt; +use subxt::ext::scale_value::{At, Composite, Primitive, Value, ValueDef, Variant}; +use subxt::{OnlineClient, PolkadotConfig}; use tokio::task::JoinHandle; +/// Pallet whose storage, constants, and events the coordinator follows. +const PALLET_NAME: &str = "StorageProvider"; + // ── ChainState ──────────────────────────────────────────────────────────────── /// Live chain state kept in sync with the runtime by the chain-state coordinator. @@ -89,48 +89,140 @@ pub struct PalletConstants { #[async_trait] pub trait ChainStateChainClient: Send + Sync { /// Full on-chain `ProviderInfo`, or `None` if the provider is not registered. - async fn get_provider_info( - &self, - who: &AccountId32, - ) -> Result, ClientError>; + async fn get_provider_info(&self, who: &AccountId32) -> Result, Error>; /// Provider's replay-window head sequence (`hsn`), or `None` if no replay /// state exists yet (the provider has never signed any terms). - async fn fetch_replay_hsn(&self, who: &AccountId32) -> Result, ClientError>; + async fn fetch_replay_hsn(&self, who: &AccountId32) -> Result, Error>; /// `StorageProvider::RequestTimeout` runtime constant, or `None` if absent /// from the node's metadata. - async fn fetch_request_timeout(&self) -> Result, ClientError>; + async fn fetch_request_timeout(&self) -> Result, Error>; } -/// Production [`ChainStateChainClient`] backed by a connected [`ProviderClient`]. -/// -/// `get_provider_info` reuses the already-connected client; the two read-only -/// queries are associated functions that open their own short-lived connection, -/// so they only need the WS URL. +/// Production [`ChainStateChainClient`] running dynamic storage queries on the +/// coordinator's own subxt connection (shared with the block subscription). struct RealChainStateClient { - client: ProviderClient, - chain_ws_url: String, + api: OnlineClient, } -#[async_trait] -impl ChainStateChainClient for RealChainStateClient { - async fn get_provider_info( +impl RealChainStateClient { + async fn fetch_value( &self, + entry: &str, who: &AccountId32, - ) -> Result, ClientError> { - self.client.get_provider_info(who).await + ) -> Result>, Error> { + let addr = subxt::dynamic::storage( + PALLET_NAME, + entry, + vec![Value::from_bytes(who.as_ref() as &[u8])], + ); + let Some(thunk) = self + .api + .storage() + .at_latest() + .await + .map_err(|e| Error::Internal(format!("Failed to get storage: {e}")))? + .fetch(&addr) + .await + .map_err(|e| Error::Internal(format!("Failed to fetch {entry}: {e}")))? + else { + return Ok(None); + }; + thunk + .to_value() + .map(Some) + .map_err(|e| Error::Internal(format!("Failed to decode {entry}: {e}"))) + } +} + +#[async_trait] +impl ChainStateChainClient for RealChainStateClient { + async fn get_provider_info(&self, who: &AccountId32) -> Result, Error> { + match self.fetch_value("Providers", who).await? { + Some(value) => decode_provider_info(&value).map(Some), + None => Ok(None), + } } - async fn fetch_replay_hsn(&self, who: &AccountId32) -> Result, ClientError> { - ProviderClient::fetch_replay_hsn(&self.chain_ws_url, who).await + async fn fetch_replay_hsn(&self, who: &AccountId32) -> Result, Error> { + Ok(self + .fetch_value("ProviderReplayStates", who) + .await? + .as_ref() + .and_then(|value| named_field(value, "hsn")) + .and_then(|v| v.as_u128()) + .map(|h| h as u64)) } - async fn fetch_request_timeout(&self) -> Result, ClientError> { - ProviderClient::fetch_request_timeout(&self.chain_ws_url).await + async fn fetch_request_timeout(&self) -> Result, Error> { + let value = self + .api + .constants() + .at(&subxt::dynamic::constant(PALLET_NAME, "RequestTimeout")) + .map_err(|e| Error::Internal(format!("Failed to read RequestTimeout: {e}")))? + .to_value() + .map_err(|e| Error::Internal(format!("Failed to decode RequestTimeout: {e}")))?; + + Ok(value.as_u128().map(|v| v as u32)) } } +// ── provider lifecycle events ───────────────────────────────────────────────── + +/// Minimal decoded view of a `StorageProvider` provider-lifecycle event. +/// +/// The coordinator re-fetches the full provider state on any relevant event, +/// so only the affected provider account — and whether the event is a +/// confirmed deregistration — needs decoding. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum ProviderLifecycleEvent { + /// `ProviderRegistered`, `ProviderSettingsUpdated`, + /// `ProviderMultiaddrUpdated`, `DeregisterAnnounced`, or + /// `DeregisterCancelled`. + Updated { provider: AccountId32 }, + /// Confirmed `ProviderDeregistered`. + Deregistered { provider: AccountId32 }, +} + +impl ProviderLifecycleEvent { + /// The provider account the event concerns. + pub fn provider(&self) -> &AccountId32 { + match self { + Self::Updated { provider } | Self::Deregistered { provider } => provider, + } + } +} + +/// Decode a finalized block's events down to the provider-lifecycle events. +fn parse_provider_lifecycle_events( + events: &subxt::events::Events, +) -> Vec { + events + .iter() + .filter_map(|event| event.ok()) + .filter(|event| event.pallet_name() == PALLET_NAME) + .filter_map(|event| { + let deregistered = match event.variant_name() { + "ProviderDeregistered" => true, + "ProviderRegistered" + | "ProviderSettingsUpdated" + | "ProviderMultiaddrUpdated" + | "DeregisterAnnounced" + | "DeregisterCancelled" => false, + _ => return None, + }; + let fields = event.field_values().ok()?; + let provider = decode_account(fields.at("provider")?)?; + Some(if deregistered { + ProviderLifecycleEvent::Deregistered { provider } + } else { + ProviderLifecycleEvent::Updated { provider } + }) + }) + .collect() +} + // ── ChainStateCoordinator ───────────────────────────────────────────────────── /// Builds and starts the live chain-state synchronisation for a single provider. @@ -191,21 +283,16 @@ impl ChainStateCoordinator { /// Connect to the chain, bootstrap initial state, then drive the finalized-block /// stream until it ends. Returns `Err` if connecting fails; `Ok(())` if the /// stream terminates cleanly — either way the caller reconnects. - async fn connect_and_follow(&self) -> Result<(), ClientError> { - let mut stream = BlockSubscriberStream::connect(&self.chain_ws_url).await?; - - let mut client = ProviderClient::new( - ClientConfig { - chain_ws_url: self.chain_ws_url.clone(), - ..Default::default() - }, - self.provider_account.to_string(), - )?; - client.connect().await?; - let chain = RealChainStateClient { - client, - chain_ws_url: self.chain_ws_url.clone(), - }; + async fn connect_and_follow(&self) -> Result<(), Error> { + let api = OnlineClient::::from_url(&self.chain_ws_url) + .await + .map_err(|e| Error::Internal(format!("Failed to connect to chain: {e}")))?; + let mut blocks = api + .blocks() + .subscribe_finalized() + .await + .map_err(|e| Error::Internal(format!("Failed to subscribe to blocks: {e}")))?; + let chain = RealChainStateClient { api }; tracing::info!("chain-state coordinator: connected; following finalized blocks"); @@ -217,8 +304,14 @@ impl ChainStateCoordinator { // rather than waiting for the next relevant event. refresh_provider_state(&chain, &self.chain_state, &self.provider_account).await; - while let Some(block) = stream.next().await { - let block_hash = H256::from_slice(block.hash().as_ref()); + while let Some(next) = blocks.next().await { + let block = match next { + Ok(block) => block, + Err(e) => { + tracing::warn!("chain-state coordinator: block subscription error: {e}"); + break; + } + }; let block_number = block.number(); tracing::debug!("Finalized block: {}", block_number); @@ -227,12 +320,7 @@ impl ChainStateCoordinator { .store(block_number, std::sync::atomic::Ordering::Relaxed); let parsed = match block.events().await { - Ok(events) => parse_pallet_events::( - &events, - storage_client::substrate::PALLET_NAME, - block_hash, - block_number, - ), + Ok(events) => parse_provider_lifecycle_events(&events), Err(e) => { tracing::warn!( "chain-state coordinator: failed to fetch events for block {block_number}: {e}" @@ -252,7 +340,7 @@ impl ChainStateCoordinator { async fn process_provider_events( &self, chain: &dyn ChainStateChainClient, - parsed: &[StorageEvent], + parsed: &[ProviderLifecycleEvent], block_number: u32, ) { refresh_if_relevant_event( @@ -365,7 +453,7 @@ pub async fn refresh_if_relevant_event( chain: &dyn ChainStateChainClient, chain_state: &ChainState, provider_account: &AccountId32, - events: &[StorageEvent], + events: &[ProviderLifecycleEvent], block_number: u32, ) { let relevant = events @@ -396,7 +484,7 @@ pub async fn refresh_if_relevant_event( // reconnect/bootstrap and non-finalized reads). This preserves the watermark // as a backstop on every path that is not a real deregistration. let deregistered = events.iter().any(|e| { - matches!(e, StorageEvent::ProviderDeregistered { provider, .. } if provider == provider_account) + matches!(e, ProviderLifecycleEvent::Deregistered { provider } if provider == provider_account) }); if deregistered { chain_state.nonce_store.reset(); @@ -406,17 +494,13 @@ pub async fn refresh_if_relevant_event( /// Whether `event` is a provider lifecycle event for `provider_account` — i.e. one /// that should trigger a [`refresh_provider_state`]. Settings, multiaddr, and the /// (de)registration events all change state `/negotiate` depends on; everything -/// else (checkpoints, challenges, agreements, other providers) is ignored. -pub fn is_relevant_provider_event(event: &StorageEvent, provider_account: &AccountId32) -> bool { - match event { - StorageEvent::ProviderRegistered { provider, .. } - | StorageEvent::ProviderSettingsUpdated { provider, .. } - | StorageEvent::ProviderMultiaddrUpdated { provider, .. } - | StorageEvent::DeregisterAnnounced { provider, .. } - | StorageEvent::ProviderDeregistered { provider, .. } - | StorageEvent::DeregisterCancelled { provider, .. } => provider == provider_account, - _ => false, - } +/// else (checkpoints, challenges, agreements, other providers) is filtered out +/// at parse time already. +pub fn is_relevant_provider_event( + event: &ProviderLifecycleEvent, + provider_account: &AccountId32, +) -> bool { + event.provider() == provider_account } // ── ChainStateCoordinatorHandle ─────────────────────────────────────────────── @@ -436,19 +520,147 @@ impl ChainStateCoordinatorHandle { } } -/// Filter a block's events down to a single pallet and parse them with `P`. -fn parse_pallet_events>( - events: &subxt::events::Events, - pallet_name: &str, - block_hash: H256, - block_number: u32, -) -> Vec { - events +// ── dynamic-value decoding ──────────────────────────────────────────────────── + +/// Decode a `StorageProvider::Providers` storage value into [`ProviderInfo`]. +fn decode_provider_info(value: &Value) -> Result { + let missing = |field: &str| Error::Internal(format!("Missing '{field}' in ProviderInfo")); + + let multiaddr = named_field(value, "multiaddr") + .map(|v| String::from_utf8_lossy(&decode_byte_vec(v)).into_owned()) + .unwrap_or_default(); + + let stake = named_field(value, "stake") + .and_then(|v| v.as_u128()) + .ok_or_else(|| missing("stake"))?; + + let committed_bytes = named_field(value, "committed_bytes") + .and_then(|v| v.as_u128()) + .ok_or_else(|| missing("committed_bytes"))? as u64; + + let settings = named_field(value, "settings").ok_or_else(|| missing("settings"))?; + + let replica_sync_price = + named_field(settings, "replica_sync_price").and_then(|v| match &v.value { + ValueDef::Variant(Variant { name, values }) if name == "Some" => { + values.values().next().and_then(|v| v.as_u128()) + } + _ => None, + }); + + let stats = named_field(value, "stats"); + let agreements_total = stats + .and_then(|s| named_field(s, "agreements_total")) + .and_then(|v| v.as_u128()) + .unwrap_or(0) as u32; + let challenges_failed = stats + .and_then(|s| named_field(s, "challenges_failed")) + .and_then(|v| v.as_u128()) + .unwrap_or(0) as u32; + + Ok(ProviderInfo { + multiaddr, + stake, + committed_bytes, + max_capacity: named_field(settings, "max_capacity") + .and_then(|v| v.as_u128()) + .ok_or_else(|| missing("max_capacity"))? as u64, + min_duration: named_field(settings, "min_duration") + .and_then(|v| v.as_u128()) + .ok_or_else(|| missing("min_duration"))? as u32, + max_duration: named_field(settings, "max_duration") + .and_then(|v| v.as_u128()) + .ok_or_else(|| missing("max_duration"))? as u32, + price_per_byte: named_field(settings, "price_per_byte") + .and_then(|v| v.as_u128()) + .ok_or_else(|| missing("price_per_byte"))?, + accepting_primary: named_field(settings, "accepting_primary") + .and_then(|v| v.as_bool()) + .ok_or_else(|| missing("accepting_primary"))?, + replica_sync_price, + accepting_extensions: named_field(settings, "accepting_extensions") + .and_then(|v| v.as_bool()) + .ok_or_else(|| missing("accepting_extensions"))?, + agreements_total, + challenges_failed, + deregister_at: named_field(value, "deregister_at").and_then(|v| match &v.value { + ValueDef::Variant(Variant { name, values }) if name == "Some" => values + .values() + .next() + .and_then(|v| v.as_u128()) + .map(|n| n as u32), + _ => None, + }), + }) +} + +/// Look up a named field in a scale_value composite. +fn named_field<'a>(value: &'a Value, field: &str) -> Option<&'a Value> { + match &value.value { + ValueDef::Composite(Composite::Named(fields)) => { + fields.iter().find(|(n, _)| n == field).map(|(_, v)| v) + } + _ => None, + } +} + +/// Decode a `Vec` / `BoundedVec` from a scale_value composite. +/// +/// `BoundedVec` serializes its `TypeInfo` as a 1-field unnamed composite +/// wrapping the inner `Vec`, so scale_value surfaces it as +/// `Composite::Unnamed([inner_vec])`. This helper drills through that wrapper +/// if present, then collects the bytes. +fn decode_byte_vec(value: &Value) -> Vec { + let ValueDef::Composite(Composite::Unnamed(items)) = &value.value else { + return Vec::new(); + }; + // Direct sequence of byte primitives. + let bytes: Vec = items .iter() - .filter_map(|event| event.ok()) - .filter(|event| event.pallet_name() == pallet_name) - .filter_map(|event| P::parse_event_detail(&event, block_hash, block_number)) - .collect() + .filter_map(|b| b.as_u128().map(|n| n as u8)) + .collect(); + if !items.is_empty() && bytes.len() == items.len() { + return bytes; + } + // BoundedVec wrapper: single inner field holds the actual sequence. + if items.len() == 1 { + return decode_byte_vec(&items[0]); + } + Vec::new() +} + +/// Decode an [`AccountId32`] from a SCALE value (a possibly-nested composite of +/// 32 byte primitives). +fn decode_account(v: &Value) -> Option { + let mut bytes = [0u8; 32]; + if collect_bytes(v, &mut bytes, 0) == 32 { + Some(AccountId32::new(bytes)) + } else { + None + } +} + +/// Recursively collect raw bytes from a SCALE value into `buf` starting at +/// `offset`, returning the new offset. +fn collect_bytes(v: &Value, buf: &mut [u8; 32], offset: usize) -> usize { + match &v.value { + ValueDef::Primitive(Primitive::U128(n)) => { + if offset < 32 { + buf[offset] = *n as u8; + offset + 1 + } else { + offset + } + } + ValueDef::Composite(Composite::Unnamed(items)) => { + let mut pos = offset; + for item in items { + pos = collect_bytes(item, buf, pos); + } + pos + } + _ => offset, + } } // ── tests ───────────────────────────────────────────────────────────────────── @@ -522,4 +734,135 @@ mod tests { }); assert_eq!(cs.constants.read().as_ref().unwrap().request_timeout, 100); } + + #[test] + fn lifecycle_event_relevance_matches_on_provider() { + let me = AccountId32::new([1u8; 32]); + let other = AccountId32::new([2u8; 32]); + let mine = ProviderLifecycleEvent::Updated { + provider: me.clone(), + }; + let theirs = ProviderLifecycleEvent::Deregistered { provider: other }; + assert!(is_relevant_provider_event(&mine, &me)); + assert!(!is_relevant_provider_event(&theirs, &me)); + } + + // ── dynamic-value decoders ──────────────────────────────────────────── + + /// Build a `Providers`-storage-shaped value the way subxt surfaces it: + /// a named composite with nested `settings`/`stats` composites, `Option` + /// fields as `Some`/`None` variants, and the `multiaddr` `BoundedVec` + /// wrapped in the single-field unnamed composite scale_value produces. + fn provider_info_value( + replica_sync_price: Option, + deregister_at: Option, + ) -> Value { + let opt = |val: Option| match val { + Some(v) => Value::unnamed_variant("Some", vec![Value::u128(v)]), + None => Value::unnamed_variant("None", Vec::>::new()), + }; + let settings = Value::named_composite([ + ("max_capacity", Value::u128(10_000)), + ("min_duration", Value::u128(10)), + ("max_duration", Value::u128(100)), + ("price_per_byte", Value::u128(5)), + ("accepting_primary", Value::bool(true)), + ("accepting_extensions", Value::bool(true)), + ("replica_sync_price", opt(replica_sync_price)), + ]); + let stats = Value::named_composite([ + ("agreements_total", Value::u128(3)), + ("challenges_failed", Value::u128(1)), + ]); + // `BoundedVec` surfaces as a 1-field unnamed composite wrapping + // the byte sequence. + let multiaddr = Value::unnamed_composite([Value::from_bytes("/ip4/1.2.3.4/tcp/3333")]); + Value::named_composite([ + ("multiaddr", multiaddr), + ("stake", Value::u128(1_000)), + ("committed_bytes", Value::u128(500)), + ("settings", settings), + ("stats", stats), + ("deregister_at", opt(deregister_at.map(u128::from))), + ]) + .map_context(|_| 0u32) + } + + #[test] + fn decode_provider_info_full() { + let info = decode_provider_info(&provider_info_value(Some(7), Some(42))).unwrap(); + assert_eq!(info.multiaddr, "/ip4/1.2.3.4/tcp/3333"); + assert_eq!(info.stake, 1_000); + assert_eq!(info.committed_bytes, 500); + assert_eq!(info.max_capacity, 10_000); + assert_eq!(info.min_duration, 10); + assert_eq!(info.max_duration, 100); + assert_eq!(info.price_per_byte, 5); + assert!(info.accepting_primary); + assert!(info.accepting_extensions); + assert_eq!(info.replica_sync_price, Some(7)); + assert_eq!(info.agreements_total, 3); + assert_eq!(info.challenges_failed, 1); + assert_eq!(info.deregister_at, Some(42)); + } + + #[test] + fn decode_provider_info_none_options() { + let info = decode_provider_info(&provider_info_value(None, None)).unwrap(); + assert_eq!(info.replica_sync_price, None); + assert_eq!(info.deregister_at, None); + } + + #[test] + fn decode_provider_info_missing_required_field_errors() { + // Everything present except the required `stake` field. + let value = Value::named_composite([("multiaddr", Value::from_bytes("/ip4/1.2.3.4"))]) + .map_context(|_| 0u32); + let err = decode_provider_info(&value).unwrap_err(); + assert!( + matches!(&err, Error::Internal(msg) if msg.contains("stake")), + "unexpected error: {err:?}" + ); + } + + #[test] + fn named_field_finds_and_misses() { + let value = Value::named_composite([("present", Value::u128(1))]).map_context(|_| 0u32); + assert!(named_field(&value, "present").is_some()); + assert!(named_field(&value, "absent").is_none()); + // Not a named composite → always None. + let prim = Value::u128(9).map_context(|_| 0u32); + assert!(named_field(&prim, "present").is_none()); + } + + #[test] + fn decode_byte_vec_handles_direct_and_wrapped_and_other() { + // Direct byte sequence (e.g. `Vec`). + let direct = Value::from_bytes(b"hello").map_context(|_| 0u32); + assert_eq!(decode_byte_vec(&direct), b"hello"); + // Single-field unnamed wrapper (e.g. `BoundedVec`). + let wrapped = Value::unnamed_composite([Value::from_bytes(b"hi")]).map_context(|_| 0u32); + assert_eq!(decode_byte_vec(&wrapped), b"hi"); + // Non-composite → empty. + let prim = Value::u128(5).map_context(|_| 0u32); + assert!(decode_byte_vec(&prim).is_empty()); + } + + #[test] + fn decode_account_from_flat_and_nested_bytes() { + // Flat 32-byte sequence. + let flat = Value::from_bytes([7u8; 32]).map_context(|_| 0u32); + assert_eq!(decode_account(&flat), Some(AccountId32::new([7u8; 32]))); + // `[u8; 32]` newtype nests the sequence one level deeper. + let nested = Value::unnamed_composite([Value::from_bytes([9u8; 32])]).map_context(|_| 0u32); + assert_eq!(decode_account(&nested), Some(AccountId32::new([9u8; 32]))); + } + + #[test] + fn decode_account_rejects_wrong_length() { + let short = Value::from_bytes([0u8; 31]).map_context(|_| 0u32); + assert_eq!(decode_account(&short), None); + let long = Value::from_bytes([0u8; 33]).map_context(|_| 0u32); + assert_eq!(decode_account(&long), None); + } } diff --git a/provider-node/src/lib.rs b/provider-node/src/lib.rs index b7585ed5..3aef0ab2 100644 --- a/provider-node/src/lib.rs +++ b/provider-node/src/lib.rs @@ -34,7 +34,7 @@ pub use api::create_router; pub use chain_state_coordinator::{ is_relevant_provider_event, refresh_if_relevant_event, refresh_provider_state, sync_constants, ChainState, ChainStateChainClient, ChainStateCoordinator, ChainStateCoordinatorHandle, - PalletConstants, + PalletConstants, ProviderLifecycleEvent, }; pub use challenge_responder::{ ChallengeChainClient, ChallengeResponder, ChallengeResponderConfig, ChallengeResponderHandle, diff --git a/provider-node/src/negotiate.rs b/provider-node/src/negotiate.rs index bf35a63b..88da8973 100644 --- a/provider-node/src/negotiate.rs +++ b/provider-node/src/negotiate.rs @@ -21,9 +21,9 @@ use crate::error::Error; use crate::storage::{NonceStore, NullNonceStore}; +use crate::types::ProviderInfo; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::Arc; -use storage_client::discovery::ProviderInfo; // Wire types are shared with the SDK so client + server agree on serde shape. pub use provider_negotiation::{AgreementTermsOf, NegotiateRequest, SignedTerms}; diff --git a/provider-node/src/types.rs b/provider-node/src/types.rs index 8a307158..fbed0388 100644 --- a/provider-node/src/types.rs +++ b/provider-node/src/types.rs @@ -3,9 +3,46 @@ //! API types for the provider node. use serde::{Deserialize, Serialize}; -use storage_client::discovery::ProviderInfo; use storage_primitives::BucketId; +// ───────────────────────────────────────────────────────────────────────────── +// On-chain Provider Info +// ───────────────────────────────────────────────────────────────────────────── + +/// The node's view of its on-chain provider registration. +/// +/// Decoded from the `StorageProvider::Providers` storage entry by the +/// chain-state coordinator; consumed by `/negotiate` validation and `/info`. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ProviderInfo { + /// Network address for connecting. + pub multiaddr: String, + /// Total stake locked. + pub stake: u128, + /// Currently committed bytes. + pub committed_bytes: u64, + /// Maximum capacity (0 = unlimited). + pub max_capacity: u64, + /// Minimum agreement duration. + pub min_duration: u32, + /// Maximum agreement duration. + pub max_duration: u32, + /// Price per byte per block. + pub price_per_byte: u128, + /// Whether accepting primary agreements. + pub accepting_primary: bool, + /// Replica sync price (None if not accepting replicas). + pub replica_sync_price: Option, + /// Whether accepting extensions. + pub accepting_extensions: bool, + /// Total agreements ever. + pub agreements_total: u32, + /// Failed challenges count. + pub challenges_failed: u32, + /// Block at which deregistration becomes finalisable (`None` = not deregistering). + pub deregister_at: Option, +} + // ───────────────────────────────────────────────────────────────────────────── // Node Upload/Download Types // ───────────────────────────────────────────────────────────────────────────── diff --git a/provider-node/tests/chain_state_integration.rs b/provider-node/tests/chain_state_integration.rs index 7478e3b4..9e3b553f 100644 --- a/provider-node/tests/chain_state_integration.rs +++ b/provider-node/tests/chain_state_integration.rs @@ -13,8 +13,8 @@ //! we assert the resulting [`ChainState`]. //! //! 2. **Event relevance.** [`is_relevant_provider_event`] decides which block -//! events trigger a refresh; tested across the provider lifecycle variants, -//! wrong-account events, and unrelated events. +//! events trigger a refresh; tested across the provider lifecycle variants +//! and wrong-account events. //! //! 3. **Resilience.** [`ChainStateCoordinator::start`] drives a reconnect loop. //! Pointed at an unreachable chain it must stay up, never panic, leave @@ -22,18 +22,14 @@ //! shut down cleanly when stopped. use async_trait::async_trait; -use sp_core::H256; use sp_runtime::AccountId32; use std::sync::atomic::Ordering; use std::sync::Arc; use std::time::Duration; -use storage_client::discovery::ProviderInfo; -use storage_client::{ClientError, ProviderSettings, StorageEvent}; -use storage_primitives::Commitment; use storage_provider_node::{ is_relevant_provider_event, refresh_if_relevant_event, refresh_provider_state, sync_constants, - ChainState, ChainStateChainClient, ChainStateCoordinator, NonceCounter, NonceStore, - PalletConstants, + ChainState, ChainStateChainClient, ChainStateCoordinator, Error, NonceCounter, NonceStore, + PalletConstants, ProviderInfo, ProviderLifecycleEvent, }; /// A WS URL that refuses immediately: port 1 on loopback is never listening, so @@ -177,7 +173,7 @@ async fn coordinator_releases_shared_state_after_stop() { /// Canned [`ChainStateChainClient`] for driving the synchronisation logic /// without a chain. Each read is either `Ok(value)` or, when its `*_err` flag is -/// set, a `ClientError` — so every branch of `sync_constants` / +/// set, an `Error` — so every branch of `sync_constants` / /// `refresh_provider_state` is reachable. #[derive(Default)] struct MockChainClient { @@ -191,28 +187,23 @@ struct MockChainClient { #[async_trait] impl ChainStateChainClient for MockChainClient { - async fn get_provider_info( - &self, - _who: &AccountId32, - ) -> Result, ClientError> { + async fn get_provider_info(&self, _who: &AccountId32) -> Result, Error> { if self.info_err { - return Err(ClientError::Chain("mock get_provider_info failure".into())); + return Err(Error::Internal("mock get_provider_info failure".into())); } Ok(self.info.clone()) } - async fn fetch_replay_hsn(&self, _who: &AccountId32) -> Result, ClientError> { + async fn fetch_replay_hsn(&self, _who: &AccountId32) -> Result, Error> { if self.hsn_err { - return Err(ClientError::Chain("mock fetch_replay_hsn failure".into())); + return Err(Error::Internal("mock fetch_replay_hsn failure".into())); } Ok(self.hsn) } - async fn fetch_request_timeout(&self) -> Result, ClientError> { + async fn fetch_request_timeout(&self) -> Result, Error> { if self.request_timeout_err { - return Err(ClientError::Chain( - "mock fetch_request_timeout failure".into(), - )); + return Err(Error::Internal("mock fetch_request_timeout failure".into())); } Ok(self.request_timeout) } @@ -423,50 +414,12 @@ async fn refresh_completes_pending_bootstrap_when_replay_state_appears() { #[test] fn lifecycle_events_for_self_are_relevant() { let me = provider_account(); - let bh = H256::zero(); let events = [ - StorageEvent::ProviderRegistered { + ProviderLifecycleEvent::Updated { provider: me.clone(), - stake: 0, - block_hash: bh, - block_number: 1, }, - StorageEvent::ProviderSettingsUpdated { + ProviderLifecycleEvent::Deregistered { provider: me.clone(), - block_hash: bh, - block_number: 1, - provider_settings: ProviderSettings { - price_per_byte: 5, - min_duration: 10, - max_duration: 100, - accepting_primary: true, - replica_sync_price: None, - accepting_extensions: true, - max_capacity: 0, - }, - }, - StorageEvent::ProviderMultiaddrUpdated { - provider: me.clone(), - multiaddr: "/ip4/1.2.3.4/tcp/3333".to_string(), - block_hash: bh, - block_number: 1, - }, - StorageEvent::DeregisterAnnounced { - provider: me.clone(), - complete_after: 10, - block_hash: bh, - block_number: 1, - }, - StorageEvent::ProviderDeregistered { - provider: me.clone(), - stake_returned: 0, - block_hash: bh, - block_number: 1, - }, - StorageEvent::DeregisterCancelled { - provider: me.clone(), - block_hash: bh, - block_number: 1, }, ]; @@ -480,39 +433,17 @@ fn lifecycle_events_for_self_are_relevant() { #[test] fn lifecycle_event_for_other_provider_is_irrelevant() { - let event = StorageEvent::ProviderRegistered { + let event = ProviderLifecycleEvent::Updated { provider: provider_account_2(), - stake: 0, - block_hash: H256::zero(), - block_number: 1, }; // Same event shape, different account → not ours, ignore it. assert!(!is_relevant_provider_event(&event, &provider_account())); } -#[test] -fn non_lifecycle_event_is_irrelevant() { - // A checkpoint event names no `provider` we filter on; it must never trigger - // a provider-state refresh. - let event = StorageEvent::BucketCheckpointed { - bucket_id: 1, - commitment: Commitment::default(), - providers: vec![provider_account()], - block_hash: H256::zero(), - block_number: 1, - }; - assert!(!is_relevant_provider_event(&event, &provider_account())); -} - // ── refresh_if_relevant_event (block-event dispatch) ────────────────────────── -fn registered_event(provider: AccountId32) -> StorageEvent { - StorageEvent::ProviderRegistered { - provider, - stake: 0, - block_hash: H256::zero(), - block_number: 1, - } +fn registered_event(provider: AccountId32) -> ProviderLifecycleEvent { + ProviderLifecycleEvent::Updated { provider } } #[tokio::test] @@ -547,12 +478,8 @@ async fn irrelevant_block_events_do_not_refresh() { }; let events = [ registered_event(provider_account_2()), - StorageEvent::BucketCheckpointed { - bucket_id: 1, - commitment: Commitment::default(), - providers: vec![provider_account()], - block_hash: H256::zero(), - block_number: 1, + ProviderLifecycleEvent::Deregistered { + provider: provider_account_2(), }, ]; @@ -651,11 +578,8 @@ async fn deregister_event_resets_persisted_nonce_store() { ..Default::default() }; - let deregister_event = StorageEvent::ProviderDeregistered { + let deregister_event = ProviderLifecycleEvent::Deregistered { provider: provider_account(), - stake_returned: 0, - block_hash: H256::zero(), - block_number: 1, }; // Chain reports provider not registered (after the event). let chain = MockChainClient::default(); // info=None diff --git a/provider-node/tests/negotiate_integration.rs b/provider-node/tests/negotiate_integration.rs index 8ff0fbfc..92e34544 100644 --- a/provider-node/tests/negotiate_integration.rs +++ b/provider-node/tests/negotiate_integration.rs @@ -15,8 +15,8 @@ use sp_core::{sr25519, Pair}; use sp_runtime::{AccountId32, MultiSignature}; use std::net::SocketAddr; use std::sync::Arc; -use storage_client::discovery::ProviderInfo; use storage_primitives::ReplicaTerms; +use storage_provider_node::ProviderInfo; use storage_provider_node::{ create_router, DiskStorage, NegotiateRequest, NonceCounter, NonceStore, NullNonceStore, PalletConstants, ProviderState, SignedTerms, Storage,