From 9443284ed78702a5445aac7cbdbdac9c9b4ca011 Mon Sep 17 00:00:00 2001 From: Mathis <154886644+echobt@users.noreply.github.com> Date: Fri, 25 Sep 2026 15:09:54 +0000 Subject: [PATCH 1/8] Require explicit buyback budgets and durable signer ownership --- README.md | 44 +- examples/payment_flow.rs | 16 +- src/bin/buyback.rs | 37 +- src/chain.rs | 40 +- src/config.rs | 3 + src/engine.rs | 987 +++++++++++++++++++++++++++++++++++---- src/state.rs | 40 +- src/store.rs | 284 +++++++++-- tests/localnet.rs | 23 +- 9 files changed, 1312 insertions(+), 162 deletions(-) diff --git a/README.md b/README.md index a8d6045..3bcb1cb 100644 --- a/README.md +++ b/README.md @@ -144,12 +144,16 @@ let master = MasterKey::from_env("BUYBACK_MASTER_KEY")?; let treasury = TreasuryKeySource::EncryptedKeystore("treasury.json".into()).load(Some(&master))?; let mut cfg = Config::new(Network::Finney.url(), 42 /* payment netuid */, keys::parse_ss58("5...treasury hotkey")?); -cfg.auto = AutoBuyback::On { destroy: Destroy::Burn, amount: AutoAmount::PaymentValueBps(10_000) }; +cfg.buyback_budget = Some(state::BuybackBudget { + amount_rao: units::parse_amount("0.5")?, currency: "TAO".into(), + source: "approved-working-capital".into(), netuid: 100, destroy: Destroy::Burn, + hotkey: keys::ss58(&cfg.treasury_hotkey), +}); let engine = Arc::new(Engine::new(chain, store, master, treasury, cfg)); engine.ensure_treasury_hotkey().await?; // try_associate_hotkey if missing -let req: PaymentRequest = engine.create_payment(CreatePayment::default())?; -let status: PaymentStatus = engine.status(&req.id)?; +let req: PaymentRequest = engine.create_payment(CreatePayment::default()).await?; +let status: PaymentStatus = engine.status(&req.id).await?; let mut settled = engine.subscribe(); // in-process callback tokio::spawn({ let e = engine.clone(); async move { e.run().await } }); @@ -227,7 +231,7 @@ receivers must deduplicate on `id`. | `buyback_netuid` | `BUYBACK_NETUID_TARGET` | **100** | | | `slippage_bps` | `BUYBACK_SLIPPAGE_BPS` | 100 (1 %) | `limit_price = spot * (1 + bps/1e4)` | | `allow_partial` | - | false | fill-or-kill by default | -| `auto` | `BUYBACK_AUTO` / `BUYBACK_AUTO_AMOUNT` | off | `keep`/`burn`/`recycle` x `all` / `payment` (spot value of the swept alpha) / fixed amount | +| `auto` | `BUYBACK_AUTO` / `BUYBACK_AUTO_AMOUNT` | off | `keep`/`burn`/`recycle` with fixed TAO amount and `BUYBACK_AUTO_SOURCE` | | `max_attempts` | - | 8 | | Treasury coldkey sources (`TreasuryKeySource`): `EnvUri` (mnemonic or secret URI in an env var), @@ -347,3 +351,35 @@ against the localnet image as a service container. MIT OR Apache-2.0. [`subxt`]: https://github.com/paritytech/subxt + +### Strict automatic jobs and shared signer ownership + +Automatic jobs now require `Config.buyback_budget`: an explicit `BuybackBudget` +with `amount_rao`, `currency="TAO"`, nonempty allocation `source`, target `netuid`, +and `destroy`. This is additional treasury capital, not a sale of deposited alpha. +The budget is frozen at job creation; stores reject later modifications. Missing +budgets, insufficient capital, zero purchases, partial spends, and incomplete burns +block completion. Legacy dynamic `AutoAmount` values do not allocate job capital. +CLI automatic mode requires a numeric `--auto-amount` and `--auto-source`. +Existing jobs without budgets require explicit migration/review, not runtime defaults. + +SQLite reserves each signer exclusively before preparing a transaction. Reservation +ownership is durable and has no lease expiry. The matching pending journal and its +release are committed atomically after verified finality or complete mortality-window +absence. A crash before journaling leaves an orphan reservation: signing remains +blocked pending explicit operator reconciliation; there is no automatic unlock API. +Use one shared SQLite database on a filesystem supporting SQLite locking. Different +database files, unrelated applications using the same key, and separate hosts with +independent stores are **not** protected. Namespace the database per chain. + +`FileStore` and custom `Store` implementations without durable reservations refuse +engine signing. A PostgreSQL adapter must implement equivalent shared ownership and +atomic journal release before enabling signing. Standalone operator methods remain +unjournaled; never use them for automatic payment processing. + +Offline regressions drive the actual `Engine::tick` through funding, sweeping, +purchase and burn, reopening the store and reconstructing the engine after every +lost broadcast response. Separate processes contend for one SQLite signer and +prove orphan ownership survives restart. This is simulated transport, not a live +chain or production readiness claim. `Chain::block_timestamp_ms(hash)` reads exact +block time; historical USD still requires an independently verified provider. diff --git a/examples/payment_flow.rs b/examples/payment_flow.rs index 276f6f7..241bc79 100644 --- a/examples/payment_flow.rs +++ b/examples/payment_flow.rs @@ -18,13 +18,21 @@ async fn main() -> Result<()> { let chain = Chain::connect(&url).await?; chain.verify_metadata().await?; - let store: Arc = Arc::new(FileStore::open("./payments")?); + let store: Arc = Arc::from(open_store("sqlite:payments.sqlite")?); let treasury = TreasuryKeySource::EnvUri("BUYBACK_TREASURY_URI".into()).load(None)?; let mut cfg = Config::new(url, netuid, hotkey); - cfg.auto = AutoBuyback::On { + cfg.buyback_budget = Some(state::BuybackBudget { + amount_rao: units::parse_amount( + &std::env::var("BUYBACK_AUTO_AMOUNT") + .map_err(|_| Error::Config("BUYBACK_AUTO_AMOUNT required".into()))?, + )?, + currency: "TAO".into(), + source: std::env::var("BUYBACK_AUTO_SOURCE") + .map_err(|_| Error::Config("BUYBACK_AUTO_SOURCE required".into()))?, + netuid: 100, destroy: Destroy::Burn, - amount: AutoAmount::PaymentValueBps(10_000), - }; + hotkey: keys::ss58(&hotkey), + }); let engine = Arc::new(Engine::new( chain, store, diff --git a/src/bin/buyback.rs b/src/bin/buyback.rs index 022f51f..fd9eabd 100644 --- a/src/bin/buyback.rs +++ b/src/bin/buyback.rs @@ -76,9 +76,12 @@ struct Common { /// Automatic buyback per settled payment: off | keep | burn | recycle #[arg(long, env = "BUYBACK_AUTO", default_value = "off")] auto: String, - /// Auto amount: all, payment, or a TAO amount + /// Fixed additional TAO allocation per new job; required when auto is enabled #[arg(long, env = "BUYBACK_AUTO_AMOUNT", default_value = "payment")] auto_amount: String, + /// Operator allocation reference for additional treasury capital + #[arg(long, env = "BUYBACK_AUTO_SOURCE")] + auto_source: Option, /// Default settlement webhook #[arg(long, env = "BUYBACK_WEBHOOK_URL")] webhook_url: Option, @@ -198,17 +201,27 @@ async fn engine(c: &Common) -> R { cfg.fee_margin_bps = c.fee_margin_bps; cfg.fee_reserve = units::parse_amount(&c.fee_reserve)?; cfg.webhook_url = c.webhook_url.clone(); - cfg.auto = match c.auto.as_str() { - "off" => AutoBuyback::Off, - d => AutoBuyback::On { - destroy: destroy(d)?, - amount: match c.auto_amount.as_str() { - "all" => AutoAmount::All, - "payment" => AutoAmount::PaymentValueBps(10_000), - a => AutoAmount::Fixed(units::parse_amount(a)?), - }, - }, - }; + if c.auto != "off" { + let amount = units::parse_amount(&c.auto_amount)?; + let destroy = destroy(&c.auto)?; + let budget = bittensor_buyback::state::BuybackBudget { + amount_rao: amount, + currency: "TAO".into(), + hotkey: c.treasury_hotkey.clone(), + source: c + .auto_source + .clone() + .ok_or("--auto-source required for automatic buyback")?, + netuid: c.buyback_netuid, + destroy, + }; + budget.validate()?; + cfg.buyback_budget = Some(budget); + cfg.auto = AutoBuyback::On { + destroy, + amount: AutoAmount::Fixed(amount), + }; + } let chain = Chain::connect(net.url()).await?; let store: Arc = Arc::from(open_store(&c.store)?); #[allow(unused_mut)] diff --git a/src/chain.rs b/src/chain.rs index 5c5ca5c..a272e5a 100644 --- a/src/chain.rs +++ b/src/chain.rs @@ -220,6 +220,9 @@ pub struct FinalizedTx { pub struct PendingTx { /// What the transaction does; decides how its result is applied. pub action: crate::engine::Action, + /// Durable exclusive signer reservation; absent in legacy journals. + #[serde(default)] + pub reservation: Option, pub signer: String, pub nonce: u64, pub tx_hash: String, @@ -442,6 +445,19 @@ impl Chain { } } + /// Timestamp.Now in milliseconds, fetched at an exact block hash. + /// Missing timestamp (including genesis) is an error, never wall-clock time. + pub async fn block_timestamp_ms(&self, hash: subxt::utils::H256) -> Result { + let at = self.api.at_block(hash).await.map_err(chain_err)?; + let value = at + .storage() + .try_fetch(dynamic::storage::<(), u64>("Timestamp", "Now"), ()) + .await + .map_err(chain_err)? + .ok_or_else(|| Error::Chain("block timestamp unavailable".into()))?; + value.decode().map_err(chain_err) + } + pub async fn free_balance(&self, who: &AccountId32) -> Result { Ok(self.account(who).await?.0) } @@ -674,7 +690,7 @@ impl Chain { } /// Build and sign `call` with an explicit nonce, mortal for [`MORTALITY`] blocks. Nothing is - /// sent: persist [`PreparedTx::journal`] first, then [`Chain::broadcast`]. + /// sent: persist a [`PendingTx`] first, then [`Chain::broadcast`]. pub async fn prepare(&self, call: &ChainCall, signer: &Keypair) -> Result { let at = self.api.at_current_block().await.map_err(chain_err)?; let payload = call.payload(); @@ -786,19 +802,15 @@ impl Chain { /// Look for a journaled extrinsic in finalized blocks `birth ..= birth + MORTALITY`. pub async fn find_pending(&self, p: &PendingTx) -> Result { - let signer = crate::keys::parse_ss58(&p.signer)?; - let (_, nonce) = self.account(&signer).await?; let head = self.best_finalized_number().await?; - let horizon = p.birth_block + MORTALITY + 1; - if (nonce as u64) <= p.nonce { - // Nonce unused on finalized state: the tx is not included (yet). - return Ok(if head > horizon { - PendingOutcome::Dead - } else { - PendingOutcome::Wait - }); - } - // Nonce consumed: find which tx used it. + let horizon = p + .birth_block + .checked_add(MORTALITY) + .and_then(|n| n.checked_add(1)) + .ok_or_else(|| Error::Chain("transaction mortality overflow".into()))?; + // Do not infer absence from a separate nonce read: it may be stale or + // from another head. Scan finalized hashes; any RPC error retains intent. + // Search the entire possible inclusion window before declaring death. for n in p.birth_block..=head.min(horizon) { let at = self.api.at_block(n).await.map_err(chain_err)?; let exts = at.extrinsics().fetch().await.map_err(chain_err)?; @@ -813,7 +825,7 @@ impl Chain { } } } - // Nonce consumed by a different transaction (the key is used elsewhere): ours is dead. + // Complete finalized scan proved absence, and mortality has elapsed. Ok(if head > horizon { PendingOutcome::Dead } else { diff --git a/src/config.rs b/src/config.rs index 42bb724..e54b70b 100644 --- a/src/config.rs +++ b/src/config.rs @@ -72,6 +72,8 @@ pub struct Config { /// Let `add_stake_limit` fill partially up to the limit price instead of failing. pub allow_partial: bool, pub auto: AutoBuyback, + /// Required explicit additional capital for automatic jobs; frozen at creation. + pub buyback_budget: Option, /// Attempts per step before a payment goes to `failed`. pub max_attempts: u32, /// Default webhook for settled payments (per-request `callback_url` overrides). @@ -97,6 +99,7 @@ impl Config { slippage_bps: 100, allow_partial: false, auto: AutoBuyback::Off, + buyback_budget: None, max_attempts: 8, webhook_url: None, } diff --git a/src/engine.rs b/src/engine.rs index 36c0640..78ecc93 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -6,17 +6,18 @@ //! compare-and-swap on the record version, and only then submits it. While `pending` is set, the //! record takes no new action. After a crash, a timeout or an RPC error, the pending entry is //! resolved against finalized chain state ([`Chain::find_pending`]): -//! * nonce consumed by *our* tx hash: apply its events (success or dispatch failure); -//! * nonce not consumed and the tx's mortality window (64 blocks) has passed: it can never be -//! included, so it is dropped and the step is rebuilt; +//! * our tx hash found in finalized blocks: apply its events (success or dispatch failure); +//! * complete finalized scan proves absence after the mortality window: drop and rebuild; //! * otherwise: wait. //! //! Treasury funding and buybacks are therefore sent at most once per step, and the sweep moves //! only the stake that is live in the wallet at build time (a replay fails on chain with //! `NotEnoughStakeToWithdraw`, it cannot move funds twice). -use crate::chain::{Chain, ChainCall, EventSummary, PendingOutcome, PendingTx}; -use crate::config::{AutoAmount, AutoBuyback, Config, Destroy}; +use crate::chain::{ + Chain, ChainCall, EventSummary, FinalizedTx, PendingOutcome, PendingTx, PreparedTx, +}; +use crate::config::{AutoBuyback, Config, Destroy}; use crate::keys::{self, DerivationSeed, Keyring, PaymentWallet}; use crate::state::{ BuybackReceipt, PaymentRecord, PaymentRequest, PaymentState, PaymentStatus, TxRef, now, @@ -62,8 +63,8 @@ pub struct CreatePayment { pub callback_url: Option, } -pub struct Engine { - chain: Chain, +pub struct Engine { + chain: C, store: Arc, keyring: Arc, /// Root of deterministic wallets; records with a `derivation_path` are opened from it. @@ -76,9 +77,9 @@ pub struct Engine { webhook: Option, } -impl Engine { +impl Engine { pub fn new( - chain: Chain, + chain: C, store: Arc, keyring: impl Into, treasury: Keypair, @@ -115,7 +116,7 @@ impl Engine { &self.cfg } - pub fn chain(&self) -> &Chain { + pub fn chain(&self) -> &C { &self.chain } @@ -132,6 +133,7 @@ impl Engine { /// Create a payment request backed by a brand-new wallet. pub async fn create_payment(&self, opts: CreatePayment) -> Result { + self.store.require_signing().await?; let wallet = PaymentWallet::generate()?; let id = uuid::Uuid::new_v4().to_string(); let address = keys::ss58(&wallet.account_id()); @@ -157,11 +159,14 @@ impl Engine { swept_alpha: 0, txs: vec![], buyback: None, + buyback_budget: self.cfg.buyback_budget.clone(), + auto_required: Some(self.cfg.buyback_budget.is_some()), attempts: 0, next_attempt_at: 0, last_error: None, failed_from: None, pending: None, + quarantined: None, swept_positions: vec![], consolidated: 0, dust_returned: false, @@ -170,6 +175,12 @@ impl Engine { updated_at: t, derivation_path: None, }; + if matches!(self.cfg.auto, AutoBuyback::On { .. }) && rec.buyback_budget.is_none() { + return Err(Error::Config("explicit job buyback budget required".into())); + } + if let Some(budget) = &rec.buyback_budget { + budget.validate()?; + } self.store.insert(&rec).await?; tracing::info!(id = %id, address = %address, "payment request created"); Ok(PaymentRequest { @@ -197,6 +208,7 @@ impl Engine { expected_address: &str, metadata: Option, ) -> Result { + self.store.require_signing().await?; let seed = self .seed .as_ref() @@ -231,11 +243,14 @@ impl Engine { swept_alpha: 0, txs: vec![], buyback: None, + buyback_budget: self.cfg.buyback_budget.clone(), + auto_required: Some(self.cfg.buyback_budget.is_some()), attempts: 0, next_attempt_at: 0, last_error: None, failed_from: None, pending: None, + quarantined: None, swept_positions: vec![], consolidated: 0, dust_returned: false, @@ -244,6 +259,12 @@ impl Engine { updated_at: t, derivation_path: Some(derivation_path.into()), }; + if matches!(self.cfg.auto, AutoBuyback::On { .. }) && rec.buyback_budget.is_none() { + return Err(Error::Config("explicit job buyback budget required".into())); + } + if let Some(budget) = &rec.buyback_budget { + budget.validate()?; + } self.store.insert(&rec).await?; tracing::info!(id, address = %rec.address, "sweep job created"); Ok(rec.status()) @@ -264,6 +285,11 @@ impl Engine { .get(id) .await? .ok_or_else(|| Error::NotFound(id.into()))?; + if r.quarantined.is_some() { + return Err(Error::Store( + "finalized receipt mismatch requires explicit repair".into(), + )); + } let to = match r.failed_from { Some(PaymentState::Pending) | None => PaymentState::Detected, Some(s) => s, @@ -317,7 +343,7 @@ impl Engine { use PaymentState::*; let recs = self .store - .list(&[Pending, Detected, Funded, Swept, Settled]) + .list(&[Pending, Detected, Funded, Swept, Settled, Failed]) .await?; let t = now(); let mut n = 0; @@ -333,7 +359,9 @@ impl Engine { self.chain.stake_positions_many(&pending).await? }; for r in recs { - if r.state == Settled && r.notified || r.next_attempt_at > t { + if r.pending.is_none() + && (r.state == Failed || r.state == Settled && r.notified || r.next_attempt_at > t) + { continue; } let id = r.id.clone(); @@ -353,6 +381,17 @@ impl Engine { mut r: PaymentRecord, positions: &std::collections::BTreeMap>, ) -> Result { + if r.quarantined.is_some() { + return Ok(false); + } + if r.pending.is_none() + && (r.auto_required.is_none() + || (r.auto_required == Some(true) && r.buyback_budget.is_none())) + { + return Err(Error::Config( + "job policy migration required before processing".into(), + )); + } if let Some(p) = r.pending.clone() { return self.resolve_pending(r, p).await; } @@ -374,6 +413,10 @@ impl Engine { .get(&r.id) .await? .ok_or_else(|| Error::NotFound(r.id.clone()))?; + // An uncertain broadcast is not a failed economic action. Reconcile first. + if cur.pending.is_some() { + return Err(e); + } cur.record_failure(&e.to_string(), self.cfg.max_attempts, now()); let cur = self.store.update(&cur).await?; self.announce_if_final(&cur); @@ -540,8 +583,19 @@ impl Engine { r.dust_returned = true; } // 3. automatic buyback - if let (AutoBuyback::On { destroy, amount }, false) = (self.cfg.auto, r.buyback_done) { - return self.step_auto_buyback(r, destroy, amount).await; + let required = r + .auto_required + .ok_or_else(|| Error::Config("legacy job policy requires explicit migration".into()))?; + if required && r.buyback_budget.is_none() { + return Err(Error::Config("required job budget missing".into())); + } + if required && r.buyback_done { + validate_completed_buyback(r)?; + } + if !r.buyback_done + && (r.buyback_budget.is_some() || matches!(self.cfg.auto, AutoBuyback::On { .. })) + { + return self.step_auto_buyback(r).await; } r.transition(PaymentState::Settled)?; *r = self.store.update(r).await?; @@ -551,21 +605,42 @@ impl Engine { Ok(true) } - async fn step_auto_buyback( - &self, - r: &mut PaymentRecord, - destroy: Destroy, - amount: AutoAmount, - ) -> Result { - let netuid = self.cfg.buyback_netuid; + async fn step_auto_buyback(&self, r: &mut PaymentRecord) -> Result { + let budget = r + .buyback_budget + .clone() + .ok_or_else(|| Error::Config("job has no explicit buyback budget; blocked".into()))?; + budget.validate()?; + let netuid = budget.netuid; + let destroy = budget.destroy; // second half: destroy what was bought if let Some(b) = &r.buyback { - if destroy == Destroy::Keep || b.destroy_tx.is_some() || b.alpha_bought == 0 { + if b.tao_spent != budget.amount_rao || b.netuid != budget.netuid { + return Err(Error::Store( + "purchase receipt disagrees with job budget".into(), + )); + } + if b.alpha_bought == 0 { + return Err(Error::Insufficient( + "buyback produced no alpha; completion blocked".into(), + )); + } + if b.destroy_tx.is_some() && b.alpha_destroyed != b.alpha_bought { + return Err(Error::Store( + "incomplete burn receipt; completion blocked".into(), + )); + } + if destroy == Destroy::Keep || b.destroy_tx.is_some() { r.buyback_done = true; *r = self.store.update(r).await?; return Ok(true); } - let call = destroy_call(destroy, self.cfg.treasury_hotkey, netuid, b.alpha_bought); + let call = destroy_call( + destroy, + keys::parse_ss58(&budget.hotkey)?, + netuid, + b.alpha_bought, + ); let signer = self.treasury.clone(); let _g = self.treasury_lock.lock().await; let action = Action::Destroy { @@ -575,31 +650,15 @@ impl Engine { return self.send(r, action, &call, &signer).await; } let _g = self.treasury_lock.lock().await; - let tao = match amount { - AutoAmount::All => self.spendable().await?, - AutoAmount::Fixed(v) => v, - AutoAmount::PaymentValueBps(bps) => { - let price = self.chain.alpha_price(r.netuid).await?; - let value = - units::alpha_value_in_tao(r.swept_alpha, price).saturating_add(r.detected_tao); - (value as u128 * bps as u128 / units::BPS as u128) as u64 - } - }; - let tao = tao.min(self.spendable().await?); - if tao < MIN_STAKE_RAO { - tracing::warn!(id = %r.id, tao, "auto buyback skipped: amount below minimum stake"); - r.buyback_done = true; - *r = self.store.update(r).await?; - return Ok(true); - } + let tao = strict_buyback_amount(budget.amount_rao, self.spendable().await?)?; let limit_price = units::buy_limit_price(self.chain.alpha_price(netuid).await?, self.cfg.slippage_bps); let call = ChainCall::AddStakeLimit { - hotkey: self.cfg.treasury_hotkey, + hotkey: keys::parse_ss58(&budget.hotkey)?, netuid, tao, limit_price, - allow_partial: self.cfg.allow_partial, + allow_partial: false, }; let signer = self.treasury.clone(); self.send( @@ -631,63 +690,25 @@ impl Engine { call: &ChainCall, signer: &Keypair, ) -> Result { - let prepared = self.chain.prepare(call, signer).await?; - let mut j = r.clone(); - j.pending = Some(PendingTx { - action, - signer: prepared.signer.clone(), - nonce: prepared.nonce, - tx_hash: prepared.tx_hash.clone(), - birth_block: prepared.birth_block, - amount: call.amount(), - }); - // Journal before broadcast: a lost race or a failed write sends nothing. - *r = self.store.update(&j).await?; - let ftx = self.chain.broadcast(&prepared).await?; - self.apply(r, ftx.tx, &ftx.summary)?; - *r = self.store.update(r).await?; - Ok(true) + send_journaled(&self.chain, self.store.as_ref(), r, action, call, signer).await } - /// Resolve a journaled tx against finalized chain state. async fn resolve_pending(&self, mut r: PaymentRecord, p: PendingTx) -> Result { - match self.chain.find_pending(&p).await? { - PendingOutcome::Wait => Ok(false), - PendingOutcome::Dead => { - tracing::warn!(id = %r.id, tx = %p.tx_hash, "journaled tx expired unincluded; rebuilding step"); - r.pending = None; - self.store.update(&r).await?; - Ok(true) - } - PendingOutcome::Included { - block_hash, - summary, - } => { - let tx = TxRef { - action: action_name(&p.action).into(), - tx_hash: p.tx_hash.clone(), - block_hash, - amount: p.amount, - }; - if summary.failed { - r.pending = None; - r.record_failure( - &format!("{} failed on chain in {}", tx.action, tx.block_hash), - self.cfg.max_attempts, - now(), - ); - } else { - self.apply(&mut r, tx, &summary)?; - } - let r = self.store.update(&r).await?; - self.announce_if_final(&r); - Ok(true) - } + let changed = reconcile_journaled( + &self.chain, + self.store.as_ref(), + &mut r, + &p, + self.cfg.max_attempts, + ) + .await?; + if changed { + self.announce_if_final(&r); } + Ok(changed) } - /// Apply the effects of a successful journaled tx to the record. - fn apply(&self, r: &mut PaymentRecord, tx: TxRef, s: &EventSummary) -> Result<()> { + fn apply(r: &mut PaymentRecord, tx: TxRef, s: &EventSummary) -> Result<()> { let action = r .pending .take() @@ -712,6 +733,11 @@ impl Engine { limit_price, } => { let (tao, alpha) = s.stake_added.ok_or(Error::EventMissing("StakeAdded"))?; + if tao != tx.amount || alpha == 0 { + return Err(Error::Store( + "purchase does not match full requested budget".into(), + )); + } r.buyback = Some(BuybackReceipt { netuid: *netuid, tao_spent: tao, @@ -726,14 +752,22 @@ impl Engine { let destroyed = s .alpha_destroyed .ok_or(Error::EventMissing("AlphaBurned/AlphaRecycled"))?; - if let Some(b) = r.buyback.as_mut() { - b.destroy_tx = Some(tx.clone()); - b.alpha_destroyed = destroyed; + let b = r + .buyback + .as_mut() + .ok_or_else(|| Error::Store("burn without purchase receipt".into()))?; + if destroyed != b.alpha_bought || destroyed == 0 { + return Err(Error::Store( + "burn amount does not match purchased alpha".into(), + )); } + b.destroy_tx = Some(tx.clone()); + b.alpha_destroyed = destroyed; r.buyback_done = true; } } r.attempts = 0; + r.next_attempt_at = 0; r.last_error = None; r.txs.push(tx); Ok(()) @@ -917,3 +951,764 @@ fn action_name(a: &Action) -> &'static str { Action::Destroy { .. } => "burn_alpha", } } + +/// Engine transport, injectable for offline state-machine tests. +#[async_trait::async_trait] +pub trait EngineRpc: TransactionRpc { + async fn free_balance(&self, who: &AccountId32) -> Result; + async fn account(&self, who: &AccountId32) -> Result<(u64, u32)>; + async fn alpha_on( + &self, + who: &AccountId32, + netuid: u16, + ) -> Result<(u64, Vec)>; + async fn alpha_of( + &self, + hotkey: &AccountId32, + coldkey: &AccountId32, + netuid: u16, + ) -> Result; + async fn estimate_fee(&self, call: &ChainCall, signer: &Keypair) -> Result; + async fn existential_deposit(&self) -> Result; + async fn alpha_price(&self, netuid: u16) -> Result; + async fn hotkey_owner(&self, hotkey: &AccountId32) -> Result>; + async fn stake_positions_many( + &self, + coldkeys: &[AccountId32], + ) -> Result>>; + async fn submit(&self, call: &ChainCall, signer: &Keypair) -> Result; + async fn finalized_blocks( + &self, + ) -> Result> + Send>>>; +} +#[async_trait::async_trait] +impl EngineRpc for Chain { + async fn free_balance(&self, who: &AccountId32) -> Result { + Chain::free_balance(self, who).await + } + async fn account(&self, who: &AccountId32) -> Result<(u64, u32)> { + Chain::account(self, who).await + } + async fn alpha_on( + &self, + who: &AccountId32, + netuid: u16, + ) -> Result<(u64, Vec)> { + Chain::alpha_on(self, who, netuid).await + } + async fn alpha_of( + &self, + hotkey: &AccountId32, + coldkey: &AccountId32, + netuid: u16, + ) -> Result { + Chain::alpha_of(self, hotkey, coldkey, netuid).await + } + async fn estimate_fee(&self, call: &ChainCall, signer: &Keypair) -> Result { + Chain::estimate_fee(self, call, signer).await + } + async fn existential_deposit(&self) -> Result { + Chain::existential_deposit(self).await + } + async fn alpha_price(&self, netuid: u16) -> Result { + Chain::alpha_price(self, netuid).await + } + async fn hotkey_owner(&self, hotkey: &AccountId32) -> Result> { + Chain::hotkey_owner(self, hotkey).await + } + async fn stake_positions_many( + &self, + coldkeys: &[AccountId32], + ) -> Result>> { + Chain::stake_positions_many(self, coldkeys).await + } + async fn submit(&self, call: &ChainCall, signer: &Keypair) -> Result { + Chain::submit(self, call, signer).await + } + async fn finalized_blocks( + &self, + ) -> Result> + Send>>> { + Ok(Box::pin(Chain::finalized_blocks(self).await?)) + } +} +/// Narrow transaction seam; tests supply no network client or live signer. +#[async_trait::async_trait] +pub trait TransactionRpc: Send + Sync { + async fn prepare(&self, call: &ChainCall, signer: &Keypair) -> Result; + async fn broadcast(&self, tx: &PreparedTx) -> Result; + async fn find_pending(&self, tx: &PendingTx) -> Result; +} +#[async_trait::async_trait] +impl TransactionRpc for Chain { + async fn prepare(&self, call: &ChainCall, signer: &Keypair) -> Result { + Chain::prepare(self, call, signer).await + } + async fn broadcast(&self, tx: &PreparedTx) -> Result { + Chain::broadcast(self, tx).await + } + async fn find_pending(&self, tx: &PendingTx) -> Result { + Chain::find_pending(self, tx).await + } +} +async fn send_journaled( + rpc: &dyn TransactionRpc, + store: &dyn Store, + r: &mut PaymentRecord, + action: Action, + call: &ChainCall, + signer: &Keypair, +) -> Result { + if r.pending.is_some() { + return Err(Error::Store( + "reconcile pending transaction before preparing another".into(), + )); + } + let signer_address = keys::ss58(&signer.public_key().to_account_id()); + let reservation = store.reserve_signer_for_record(&signer_address, r).await?; + // A crash or prepare error leaves an orphan reservation: never expire it automatically. + let prepared = rpc.prepare(call, signer).await?; + if prepared.signer != signer_address { + return Err(Error::Store("prepared signer mismatch".into())); + } + let mut journal = r.clone(); + journal.pending = Some(PendingTx { + action, + reservation: Some(reservation), + signer: prepared.signer.clone(), + nonce: prepared.nonce, + tx_hash: prepared.tx_hash.clone(), + birth_block: prepared.birth_block, + amount: call.amount(), + }); + *r = store.update(&journal).await?; + let finalized = rpc.broadcast(&prepared).await?; + let mut applied = r.clone(); + Engine::::apply(&mut applied, finalized.tx, &finalized.summary)?; + *r = store.update(&applied).await?; + Ok(true) +} +async fn reconcile_journaled( + rpc: &dyn TransactionRpc, + store: &dyn Store, + r: &mut PaymentRecord, + pending: &PendingTx, + max_attempts: u32, +) -> Result { + if r.pending.as_ref() != Some(pending) { + return Err(Error::Store("pending journal mismatch".into())); + } + match rpc.find_pending(pending).await? { + PendingOutcome::Wait => Ok(false), + PendingOutcome::Dead => { + let mut applied = r.clone(); + applied.pending = None; + *r = store.update(&applied).await?; + Ok(true) + } + PendingOutcome::Included { + block_hash, + summary, + } => { + let tx = TxRef { + action: action_name(&pending.action).into(), + tx_hash: pending.tx_hash.clone(), + block_hash, + amount: pending.amount, + }; + let mut applied = r.clone(); + if summary.failed { + applied.pending = None; + applied.record_failure( + "journaled transaction failed on finalized chain", + max_attempts, + now(), + ); + } else { + if let Err(error) = Engine::::apply(&mut applied, tx.clone(), &summary) { + applied = r.clone(); + applied.pending = None; + applied.quarantined = Some(pending.clone()); + applied.txs.push(tx); + applied.record_failure( + &format!("finalized receipt mismatch: {error}"), + 1, + now(), + ); + } + } + *r = store.update(&applied).await?; + Ok(true) + } + } +} + +fn validate_completed_buyback(r: &PaymentRecord) -> Result<()> { + let budget = r + .buyback_budget + .as_ref() + .ok_or_else(|| Error::Store("missing job budget".into()))?; + budget.validate()?; + let receipt = r + .buyback + .as_ref() + .ok_or_else(|| Error::Store("missing purchase receipt".into()))?; + if receipt.tao_spent != budget.amount_rao + || receipt.netuid != budget.netuid + || receipt.alpha_bought == 0 + || receipt.stake_tx.block_hash.is_empty() + || (budget.destroy != Destroy::Keep + && (receipt.alpha_destroyed != receipt.alpha_bought + || receipt.destroy_tx.as_ref().is_none_or(|tx| { + tx.block_hash.is_empty() + || tx.action + != if budget.destroy == Destroy::Burn { + "burn_alpha" + } else { + "recycle_alpha" + } + }))) + { + return Err(Error::Store("incomplete job buyback receipt".into())); + } + Ok(()) +} + +fn strict_buyback_amount(requested: u64, spendable: u64) -> Result { + if requested < MIN_STAKE_RAO { + return Err(Error::Insufficient( + "buyback budget below minimum; blocked, not skipped".into(), + )); + } + if requested > spendable { + return Err(Error::Insufficient( + "full buyback budget unavailable; no partial spend".into(), + )); + } + Ok(requested) +} + +#[cfg(test)] +mod strict_budget_tests { + use super::*; + #[test] + fn no_silent_clamp_or_skipped_success() { + assert!(strict_buyback_amount(0, u64::MAX).is_err()); + assert!(strict_buyback_amount(MIN_STAKE_RAO - 1, u64::MAX).is_err()); + assert!(strict_buyback_amount(MIN_STAKE_RAO, MIN_STAKE_RAO - 1).is_err()); + assert_eq!( + strict_buyback_amount(MIN_STAKE_RAO, MIN_STAKE_RAO).unwrap(), + MIN_STAKE_RAO + ); + assert_eq!( + strict_buyback_amount(MIN_STAKE_RAO, u64::MAX).unwrap(), + MIN_STAKE_RAO + ); + } +} + +#[cfg(all(test, feature = "sqlite"))] +mod journal_rpc_tests { + use super::*; + use crate::{keys::MasterKey, store::SqliteStore}; + use std::sync::atomic::{AtomicUsize, Ordering}; + struct LostReply { + broadcasts: AtomicUsize, + summary: EventSummary, + lookup: AtomicUsize, + dead: bool, + } + #[async_trait::async_trait] + impl TransactionRpc for LostReply { + async fn prepare(&self, call: &ChainCall, signer: &Keypair) -> Result { + Ok(PreparedTx { + call: call.clone(), + signer: keys::ss58(&signer.public_key().to_account_id()), + nonce: 7, + tx_hash: format!("hash-{}", self.broadcasts.load(Ordering::SeqCst)), + birth_block: 100, + bytes: vec![], + }) + } + async fn broadcast(&self, _: &PreparedTx) -> Result { + self.broadcasts.fetch_add(1, Ordering::SeqCst); + Err(Error::Chain("connection lost after inclusion".into())) + } + async fn find_pending(&self, _: &PendingTx) -> Result { + if self.dead { + return Ok(PendingOutcome::Dead); + } + match self.lookup.fetch_add(1, Ordering::SeqCst) { + 0 => return Err(Error::Chain("intermittent RPC".into())), + 1 => return Ok(PendingOutcome::Wait), + _ => {} + } + Ok(PendingOutcome::Included { + block_hash: "finalized-test-block".into(), + summary: self.summary.clone(), + }) + } + } + #[tokio::test] + async fn durable_journal_recovers_lost_replies_without_rebroadcast() { + let dir = std::env::temp_dir().join(format!("buyback-journal-{}", uuid::Uuid::new_v4())); + let store = SqliteStore::open(&dir).unwrap(); + let wallet = PaymentWallet::generate().unwrap(); + let sealed = wallet + .seal_with(&MasterKey::generate().into(), b"test") + .unwrap(); + let mut rec:PaymentRecord=serde_json::from_value(serde_json::json!({ + "id":"simulation","address":keys::ss58(&wallet.account_id()),"netuid":100, + "min_alpha":0,"min_tao":null,"created_at":1,"expires_at":9999999999u64, + "state":"detected","version":0,"sealed_secret":sealed,"metadata":null,"callback_url":null, + "detected_alpha":0,"detected_tao":0,"funded_tao":0,"swept_alpha":0,"txs":[],"buyback":null, + "attempts":0,"next_attempt_at":0,"last_error":null,"failed_from":null,"pending":null, + "swept_positions":[],"consolidated":0,"dust_returned":false,"buyback_done":false,"notified":false,"updated_at":1 + })).unwrap(); + store.insert(&rec).await.unwrap(); + let who = wallet.account_id(); + let steps = [ + ( + Action::Fund, + ChainCall::TransferTao { + dest: who, + amount: 3_000_000, + }, + EventSummary::default(), + ), + ( + Action::Sweep { + hotkey: keys::ss58(&who), + }, + ChainCall::TransferStake { + dest_coldkey: who, + hotkey: who, + netuid: 100, + alpha: 10, + }, + EventSummary { + names: vec![(crate::chain::PALLET.into(), "StakeTransferred".into())], + ..Default::default() + }, + ), + ( + Action::BuyStake { + netuid: 100, + limit_price: 1, + }, + ChainCall::AddStakeLimit { + hotkey: who, + netuid: 100, + tao: MIN_STAKE_RAO, + limit_price: 1, + allow_partial: false, + }, + EventSummary { + stake_added: Some((MIN_STAKE_RAO, 20)), + ..Default::default() + }, + ), + ( + Action::Destroy { + netuid: 100, + recycle: false, + }, + ChainCall::BurnAlpha { + hotkey: who, + netuid: 100, + alpha: 20, + }, + EventSummary { + alpha_destroyed: Some(20), + ..Default::default() + }, + ), + ]; + for (action, call, summary) in steps { + let rpc = LostReply { + broadcasts: AtomicUsize::new(0), + summary, + lookup: AtomicUsize::new(0), + dead: false, + }; + assert!( + send_journaled( + &rpc, + &store, + &mut rec, + action.clone(), + &call, + wallet.keypair() + ) + .await + .is_err() + ); + // A new store instance models process loss; only disk journal survives. + let reopened = SqliteStore::open(&dir).unwrap(); + rec = reopened.get("simulation").await.unwrap().unwrap(); + let pending = rec.pending.clone().unwrap(); + assert_eq!(pending.action, action); + assert_eq!(pending.tx_hash, "hash-0"); + assert!( + send_journaled(&rpc, &reopened, &mut rec, action, &call, wallet.keypair()) + .await + .is_err() + ); + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 1); + let before = rec.version; + assert!( + reconcile_journaled(&rpc, &reopened, &mut rec, &pending, 8) + .await + .is_err() + ); + assert!( + !reconcile_journaled(&rpc, &reopened, &mut rec, &pending, 8) + .await + .unwrap() + ); + assert_eq!(rec.version, before); + assert_eq!( + reopened.get("simulation").await.unwrap().unwrap().pending, + Some(pending.clone()) + ); + assert!( + reconcile_journaled(&rpc, &reopened, &mut rec, &pending, 8) + .await + .unwrap() + ); + assert!(rec.pending.is_none()); + assert!( + reconcile_journaled(&rpc, &reopened, &mut rec, &pending, 8) + .await + .is_err() + ); + if rec.state == PaymentState::Funded && rec.swept_alpha > 0 { + rec.transition(PaymentState::Swept).unwrap(); + rec = reopened.update(&rec).await.unwrap(); + } + } + assert_eq!(rec.txs.len(), 4); + assert_eq!(rec.swept_alpha, 10); + assert_eq!(rec.buyback.as_ref().unwrap().tao_spent, MIN_STAKE_RAO); + assert_eq!(rec.buyback.as_ref().unwrap().alpha_destroyed, 20); + assert!(rec.buyback_done); + rec.transition(PaymentState::Settled).unwrap(); + store.update(&rec).await.unwrap(); + drop(store); + std::fs::remove_file(dir).unwrap(); + } + #[tokio::test] + async fn stale_journal_cannot_broadcast() { + let dir = tempfile::tempdir().unwrap(); + let store = SqliteStore::open(dir.path().join("store.db")).unwrap(); + let mut rec = crate::state::test_record("stale"); + store.insert(&rec).await.unwrap(); + store.update(&rec).await.unwrap(); + let rpc = LostReply { + broadcasts: AtomicUsize::new(0), + summary: EventSummary::default(), + lookup: AtomicUsize::new(2), + dead: false, + }; + let wallet = PaymentWallet::generate().unwrap(); + let call = ChainCall::TransferTao { + dest: wallet.account_id(), + amount: 1, + }; + assert!(matches!( + send_journaled( + &rpc, + &store, + &mut rec, + Action::Fund, + &call, + wallet.keypair() + ) + .await, + Err(Error::Conflict(_)) + )); + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 0); + assert!(rec.pending.is_none()); + assert!( + store + .reserve_signer(&keys::ss58(&wallet.account_id()), "next") + .await + .is_ok() + ); + } + + #[tokio::test] + async fn invalid_finalized_burn_quarantines_job_and_releases_signer() { + let dir = tempfile::tempdir().unwrap(); + let store = SqliteStore::open(dir.path().join("store.db")).unwrap(); + let mut rec = crate::state::test_record("burn"); + rec.state = PaymentState::Swept; + let tx = TxRef { + action: "add_stake_limit".into(), + tx_hash: "buy".into(), + block_hash: "block".into(), + amount: MIN_STAKE_RAO, + }; + rec.buyback = Some(BuybackReceipt { + netuid: 100, + tao_spent: MIN_STAKE_RAO, + alpha_bought: 20, + limit_price: 1, + stake_tx: tx, + destroy_tx: None, + alpha_destroyed: 0, + }); + let pending = PendingTx { + action: Action::Destroy { + netuid: 100, + recycle: false, + }, + reservation: Some(store.reserve_signer("test-only", "burn").await.unwrap()), + signer: "test-only".into(), + nonce: 1, + tx_hash: "burn".into(), + birth_block: 100, + amount: 20, + }; + rec.pending = Some(pending.clone()); + store.insert(&rec).await.unwrap(); + let rpc = LostReply { + broadcasts: AtomicUsize::new(0), + summary: EventSummary { + alpha_destroyed: Some(19), + ..Default::default() + }, + lookup: AtomicUsize::new(2), + dead: false, + }; + assert!( + reconcile_journaled(&rpc, &store, &mut rec, &pending, 8) + .await + .unwrap() + ); + let stored = store.get("burn").await.unwrap().unwrap(); + assert!(stored.pending.is_none()); + assert_eq!(stored.quarantined, Some(pending)); + assert!(!stored.buyback_done); + assert_eq!(stored.state, PaymentState::Failed); + assert_eq!(stored.txs.len(), 1); + assert!( + store + .reserve_signer("test-only", "another-job") + .await + .is_ok() + ); + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 0); + } + + #[tokio::test] + async fn verified_dead_intent_releases_signer_without_completion() { + let store = SqliteStore::in_memory().unwrap(); + let mut rec = crate::state::test_record("dead"); + let pending = PendingTx { + action: Action::Fund, + reservation: Some(store.reserve_signer("signer", "dead").await.unwrap()), + signer: "signer".into(), + nonce: 1, + tx_hash: "dead-hash".into(), + birth_block: 1, + amount: 1, + }; + rec.pending = Some(pending.clone()); + store.insert(&rec).await.unwrap(); + let rpc = LostReply { + broadcasts: AtomicUsize::new(0), + summary: EventSummary::default(), + lookup: AtomicUsize::new(0), + dead: true, + }; + assert!( + reconcile_journaled(&rpc, &store, &mut rec, &pending, 8) + .await + .unwrap() + ); + assert!(rec.pending.is_none()); + assert!(!rec.buyback_done); + assert!(store.reserve_signer("signer", "next").await.is_ok()); + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 0); + } + struct SimulatedChain { + broadcasts: AtomicUsize, + hotkey: AccountId32, + result: std::sync::Mutex, + } + #[async_trait::async_trait] + impl TransactionRpc for SimulatedChain { + async fn prepare(&self, call: &ChainCall, signer: &Keypair) -> Result { + let nonce = self.broadcasts.load(Ordering::SeqCst) as u64; + Ok(PreparedTx { + call: call.clone(), + signer: keys::ss58(&signer.public_key().to_account_id()), + nonce, + tx_hash: format!("simulated-{nonce}"), + birth_block: 100, + bytes: vec![], + }) + } + async fn broadcast(&self, tx: &PreparedTx) -> Result { + self.broadcasts.fetch_add(1, Ordering::SeqCst); + let mut summary = EventSummary::default(); + match &tx.call { + ChainCall::TransferTao { .. } => {} + ChainCall::TransferStake { .. } => summary + .names + .push((crate::chain::PALLET.into(), "StakeTransferred".into())), + ChainCall::AddStakeLimit { + tao, allow_partial, .. + } => { + assert!(!allow_partial); + summary.stake_added = Some((*tao, 20)); + } + ChainCall::BurnAlpha { alpha, .. } => summary.alpha_destroyed = Some(*alpha), + _ => panic!("unexpected simulated action"), + } + *self.result.lock().unwrap() = summary; + Err(Error::Chain("simulated lost broadcast response".into())) + } + async fn find_pending(&self, _: &PendingTx) -> Result { + Ok(PendingOutcome::Included { + block_hash: "simulated-finalized-block".into(), + summary: self.result.lock().unwrap().clone(), + }) + } + } + #[async_trait::async_trait] + impl EngineRpc for Arc { + async fn free_balance(&self, _: &AccountId32) -> Result { + Ok(1_000_000_000) + } + async fn account(&self, _: &AccountId32) -> Result<(u64, u32)> { + Ok((0, 0)) + } + async fn alpha_on( + &self, + _: &AccountId32, + netuid: u16, + ) -> Result<(u64, Vec)> { + Ok(( + 1_000_000_000, + vec![crate::chain::StakePosition { + hotkey: self.hotkey, + netuid, + alpha: 1_000_000_000, + }], + )) + } + async fn alpha_of(&self, _: &AccountId32, _: &AccountId32, _: u16) -> Result { + panic!("consolidation disabled") + } + async fn estimate_fee(&self, _: &ChainCall, _: &Keypair) -> Result { + Ok(100) + } + async fn existential_deposit(&self) -> Result { + Ok(100) + } + async fn alpha_price(&self, _: u16) -> Result { + Ok(1_000_000_000) + } + async fn hotkey_owner(&self, _: &AccountId32) -> Result> { + panic!("no association") + } + async fn stake_positions_many( + &self, + _: &[AccountId32], + ) -> Result>> + { + panic!("starts at detected") + } + async fn submit(&self, _: &ChainCall, _: &Keypair) -> Result { + panic!("unjournaled submit prohibited") + } + async fn finalized_blocks( + &self, + ) -> Result> + Send>>> { + panic!("test drives tick") + } + } + #[async_trait::async_trait] + impl TransactionRpc for Arc { + async fn prepare(&self, c: &ChainCall, s: &Keypair) -> Result { + self.as_ref().prepare(c, s).await + } + async fn broadcast(&self, t: &PreparedTx) -> Result { + self.as_ref().broadcast(t).await + } + async fn find_pending(&self, t: &PendingTx) -> Result { + self.as_ref().find_pending(t).await + } + } + #[tokio::test] + async fn engine_ticks_complete_after_each_lost_response_and_restart() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("engine.db"); + let wallet = PaymentWallet::generate().unwrap(); + let treasury = PaymentWallet::generate().unwrap(); + let keyring: Keyring = MasterKey::from_bytes([42; 32]).into(); + let mut rec = crate::state::test_record("engine"); + rec.netuid = 100; + rec.auto_required = Some(true); + rec.address = keys::ss58(&wallet.account_id()); + rec.sealed_secret = wallet + .seal_with(&keyring, &keys::wallet_aad(&rec.id, &rec.address)) + .unwrap(); + rec.state = PaymentState::Detected; + rec.buyback_budget = Some(crate::state::BuybackBudget { + amount_rao: MIN_STAKE_RAO, + currency: "TAO".into(), + source: "explicit-test-allocation".into(), + hotkey: keys::ss58(&treasury.account_id()), + netuid: 100, + destroy: Destroy::Burn, + }); + SqliteStore::open(&path) + .unwrap() + .insert(&rec) + .await + .unwrap(); + let rpc = Arc::new(SimulatedChain { + broadcasts: AtomicUsize::new(0), + hotkey: treasury.account_id(), + result: std::sync::Mutex::new(EventSummary::default()), + }); + let mut cfg = Config::new("unused", 100, treasury.account_id()); + cfg.consolidate = false; + cfg.return_dust = false; + // Runtime defaults deliberately differ: the persisted job budget wins. + cfg.buyback_budget = None; + for _ in 0..12 { + let store = Arc::new(SqliteStore::open(&path).unwrap()); + let engine = Engine::new( + rpc.clone(), + store.clone(), + MasterKey::from_bytes([42; 32]), + treasury.keypair().clone(), + cfg.clone(), + ); + engine.tick().await.unwrap(); + rec = store.get("engine").await.unwrap().unwrap(); + if rec.state == PaymentState::Settled { + break; + } + } + assert_eq!(rec.state, PaymentState::Settled); + assert!(rec.buyback_done); + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 4); + assert_eq!(rec.txs.len(), 4); + assert_eq!(rec.buyback.as_ref().unwrap().tao_spent, MIN_STAKE_RAO); + assert_eq!(rec.buyback.as_ref().unwrap().alpha_destroyed, 20); + assert!(rec.pending.is_none()); + let store = SqliteStore::open(&path).unwrap(); + assert!( + store + .reserve_signer(&keys::ss58(&treasury.account_id()), "next-job") + .await + .is_ok() + ); + rec.buyback_budget.as_mut().unwrap().amount_rao += 1; + assert!(store.update(&rec).await.is_err()); + } +} diff --git a/src/state.rs b/src/state.rs index 85481dc..a416cba 100644 --- a/src/state.rs +++ b/src/state.rs @@ -88,6 +88,33 @@ pub struct BuybackReceipt { pub alpha_destroyed: u64, } +/// Explicit additional treasury capital, snapshotted when the job is created. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +pub struct BuybackBudget { + pub amount_rao: u64, + pub currency: String, + pub source: String, + pub netuid: u16, + pub destroy: crate::config::Destroy, + pub hotkey: String, +} +impl BuybackBudget { + pub fn validate(&self) -> Result<()> { + crate::keys::parse_ss58(&self.hotkey)?; + if (self.netuid == 0 && self.destroy != crate::config::Destroy::Keep) + || self.amount_rao < crate::engine::MIN_STAKE_RAO + || self.currency != "TAO" + || self.source.trim().is_empty() + || self.source.len() > 128 + { + return Err(Error::Config( + "explicit TAO buyback budget/source required".into(), + )); + } + Ok(()) + } +} + /// Full persisted record. Contains the sealed wallet secret: never return it to clients, use /// [`PaymentStatus`]. #[derive(Clone, Debug, Serialize, Deserialize)] @@ -111,6 +138,11 @@ pub struct PaymentRecord { pub swept_alpha: u64, pub txs: Vec, pub buyback: Option, + #[serde(default)] + pub buyback_budget: Option, + /// None identifies legacy jobs requiring explicit migration. + #[serde(default)] + pub auto_required: Option, pub attempts: u32, pub next_attempt_at: u64, pub last_error: Option, @@ -118,12 +150,15 @@ pub struct PaymentRecord { /// Extrinsic journaled *before* broadcast. While set, no new action is taken for this record /// until the pending one is proven included or dead (see `Engine::resolve_pending`). pub pending: Option, + /// Finalized intent with unexpected events; explicit repair required, never retry blindly. + #[serde(default)] + pub quarantined: Option, /// `(hotkey ss58, alpha)` positions moved to the treasury coldkey by the sweep. pub swept_positions: Vec<(String, u64)>, /// How many of `swept_positions` were moved onto the treasury hotkey. pub consolidated: usize, pub dust_returned: bool, - /// Auto-buyback finished (done, skipped or given up). + /// Requested auto-buyback completed with a validated receipt. pub buyback_done: bool, pub notified: bool, pub updated_at: u64, @@ -261,11 +296,14 @@ pub(crate) fn test_record(id: &str) -> PaymentRecord { swept_alpha: 0, txs: vec![], buyback: None, + buyback_budget: None, + auto_required: Some(false), attempts: 0, next_attempt_at: 0, last_error: None, failed_from: None, pending: None, + quarantined: None, swept_positions: vec![], consolidated: 0, dust_returned: false, diff --git a/src/store.rs b/src/store.rs index 22d87d1..ddcf385 100644 --- a/src/store.rs +++ b/src/store.rs @@ -11,6 +11,26 @@ use std::sync::Mutex; /// Async so a database-backed store (e.g. PostgreSQL) can implement it directly. #[async_trait::async_trait] pub trait Store: Send + Sync + 'static { + async fn require_signing(&self) -> Result<()> { + Err(Error::Store( + "shared signer reservations unsupported; requests disabled".into(), + )) + } + async fn reserve_signer_for_record( + &self, + _signer: &str, + _rec: &PaymentRecord, + ) -> Result { + Err(Error::Store("atomic signer reservation unsupported".into())) + } + + /// Exclusive durable signer ownership, BEFORE nonce selection. No leases or automatic expiry. + /// Implementations must atomically release only when the matching journal is resolved. + async fn reserve_signer(&self, _signer: &str, _job: &str) -> Result { + Err(Error::Store( + "shared durable signer reservations unsupported; signing disabled".into(), + )) + } /// Insert a new record (version 0). Fails if the id exists. async fn insert(&self, rec: &PaymentRecord) -> Result<()>; async fn get(&self, id: &str) -> Result>; @@ -76,7 +96,13 @@ impl FileStore { .map_err(store_err)?; f.sync_all().map_err(store_err)?; } - std::fs::rename(&tmp, path).map_err(store_err) + std::fs::rename(&tmp, path).map_err(store_err)?; + // Persist the directory entry before a journaled transaction can broadcast. + #[cfg(unix)] + std::fs::File::open(&self.dir) + .and_then(|dir| dir.sync_all()) + .map_err(store_err)?; + Ok(()) } } @@ -116,6 +142,9 @@ impl Store for FileStore { let cur = self .read(&p)? .ok_or_else(|| Error::NotFound(rec.id.clone()))?; + if cur.buyback_budget != rec.buyback_budget || cur.auto_required != rec.auto_required { + return Err(Error::Store("immutable job budget".into())); + } if cur.version != rec.version { return Err(Error::Conflict(rec.id.clone())); } @@ -146,7 +175,9 @@ impl SqliteStore { fn init(conn: rusqlite::Connection) -> Result { conn.execute_batch( - "PRAGMA journal_mode=WAL; PRAGMA busy_timeout=5000; + "PRAGMA journal_mode=WAL; PRAGMA synchronous=FULL; PRAGMA busy_timeout=5000; + CREATE TABLE IF NOT EXISTS signer_reservations ( + signer TEXT PRIMARY KEY, job TEXT NOT NULL, token TEXT NOT NULL UNIQUE); CREATE TABLE IF NOT EXISTS payments ( id TEXT PRIMARY KEY, state TEXT NOT NULL, version INTEGER NOT NULL, created_at INTEGER NOT NULL, record TEXT NOT NULL); @@ -157,11 +188,57 @@ impl SqliteStore { conn: Mutex::new(conn), }) } + fn reserve(&self, signer: &str, job: &str, version: Option) -> Result { + let mut c = self.conn.lock().map_err(store_err)?; + let tx = c + .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate) + .map_err(store_err)?; + if let Some(version) = version { + let current: Option = { + use rusqlite::OptionalExtension; + tx.query_row("SELECT version FROM payments WHERE id=?1", [job], |r| { + r.get(0) + }) + .optional() + .map_err(store_err)? + }; + if current != i64::try_from(version).ok() { + return Err(Error::Conflict(job.into())); + } + } + // Legacy journals also fence the signer; migration cannot make uncertainty disappear. + let pending: bool=tx.query_row("SELECT EXISTS(SELECT 1 FROM payments WHERE json_extract(record,'$.pending.signer')=?1)",[signer],|r|r.get(0)).map_err(store_err)?; + if pending { + return Err(Error::Conflict("unresolved signer journal".into())); + } + let token = uuid::Uuid::new_v4().to_string(); + let count = tx + .execute( + "INSERT OR IGNORE INTO signer_reservations(signer,job,token) VALUES (?1,?2,?3)", + rusqlite::params![signer, job, token], + ) + .map_err(store_err)?; + if count != 1 { + return Err(Error::Conflict(format!("signer {signer} reserved"))); + } + tx.commit().map_err(store_err)?; + Ok(token) + } } #[cfg(feature = "sqlite")] #[async_trait::async_trait] impl Store for SqliteStore { + async fn require_signing(&self) -> Result<()> { + Ok(()) + } + async fn reserve_signer_for_record(&self, signer: &str, rec: &PaymentRecord) -> Result { + self.reserve(signer, &rec.id, Some(rec.version)) + } + async fn reserve_signer(&self, signer: &str, job: &str) -> Result { + self.reserve(signer, job, None) + } + async fn insert(&self, rec: &PaymentRecord) -> Result<()> { let c = self.conn.lock().map_err(store_err)?; c.execute( @@ -213,39 +290,79 @@ impl Store for SqliteStore { } async fn update(&self, rec: &PaymentRecord) -> Result { - let c = self.conn.lock().map_err(store_err)?; - let mut next = rec.clone(); - next.version += 1; - next.updated_at = crate::state::now(); - let n = c - .execute( - "UPDATE payments SET state=?1, version=?2, record=?3 WHERE id=?4 AND version=?5", - rusqlite::params![ - next.state.as_str(), - next.version as i64, - serde_json::to_string(&next).map_err(store_err)?, - rec.id, - rec.version as i64 - ], - ) + let mut c = self.conn.lock().map_err(store_err)?; + let tx = c + .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate) .map_err(store_err)?; - match n { - 1 => Ok(next), - _ if self_exists(&c, &rec.id)? => Err(Error::Conflict(rec.id.clone())), - _ => Err(Error::NotFound(rec.id.clone())), + let previous: String = tx + .query_row("SELECT record FROM payments WHERE id=?1", [&rec.id], |r| { + r.get(0) + }) + .map_err(|e| { + if matches!(e, rusqlite::Error::QueryReturnedNoRows) { + Error::NotFound(rec.id.clone()) + } else { + store_err(e) + } + })?; + let previous: PaymentRecord = serde_json::from_str(&previous).map_err(store_err)?; + if previous.version != rec.version { + return Err(Error::Conflict(rec.id.clone())); + } + if previous.buyback_budget != rec.buyback_budget + || previous.auto_required != rec.auto_required + { + return Err(Error::Store("immutable job budget".into())); + } + if let Some(p) = &rec.pending { + let token = p.reservation.as_deref().ok_or_else(|| { + Error::Store("legacy journal requires manual reconciliation".into()) + })?; + let owned: bool = tx.query_row("SELECT EXISTS(SELECT 1 FROM signer_reservations WHERE signer=?1 AND job=?2 AND token=?3)",rusqlite::params![p.signer,rec.id,token],|r|r.get(0)).map_err(store_err)?; + if !owned { + return Err(Error::Store("signer reservation mismatch".into())); + } } + if let Some(p) = &previous.pending { + if rec.pending.is_some() && rec.pending != previous.pending { + return Err(Error::Store("cannot replace unresolved journal".into())); + } + if rec.pending.is_none() { + let token = p.reservation.as_deref().ok_or_else(|| { + Error::Store("legacy journal requires manual reconciliation".into()) + })?; + let removed = tx + .execute( + "DELETE FROM signer_reservations WHERE signer=?1 AND job=?2 AND token=?3", + rusqlite::params![p.signer, rec.id, token], + ) + .map_err(store_err)?; + if removed != 1 { + return Err(Error::Store("signer reservation mismatch".into())); + } + } + } + let mut next = rec.clone(); + next.version = next + .version + .checked_add(1) + .ok_or_else(|| Error::Store("version overflow".into()))?; + next.updated_at = crate::state::now(); + tx.execute( + "UPDATE payments SET state=?1,version=?2,record=?3 WHERE id=?4", + rusqlite::params![ + next.state.as_str(), + next.version as i64, + serde_json::to_string(&next).map_err(store_err)?, + next.id + ], + ) + .map_err(store_err)?; + tx.commit().map_err(store_err)?; + Ok(next) } } -#[cfg(feature = "sqlite")] -fn self_exists(c: &rusqlite::Connection, id: &str) -> Result { - c.query_row("SELECT count(*) FROM payments WHERE id=?1", [id], |r| { - r.get::<_, i64>(0) - }) - .map(|n| n > 0) - .map_err(store_err) -} - /// Open a store from a spec: `sqlite:`, `file:`, or a bare path (sqlite). pub fn open_store(spec: &str) -> Result> { if let Some(dir) = spec.strip_prefix("file:") { @@ -264,6 +381,17 @@ pub fn open_store(spec: &str) -> Result> { #[async_trait::async_trait] impl Store for Box { + async fn require_signing(&self) -> Result<()> { + (**self).require_signing().await + } + async fn reserve_signer_for_record(&self, s: &str, r: &PaymentRecord) -> Result { + (**self).reserve_signer_for_record(s, r).await + } + + async fn reserve_signer(&self, signer: &str, job: &str) -> Result { + (**self).reserve_signer(signer, job).await + } + async fn insert(&self, rec: &PaymentRecord) -> Result<()> { (**self).insert(rec).await } @@ -342,4 +470,100 @@ mod tests { 1 ); } + #[cfg(feature = "sqlite")] + #[test] + fn reservation_child_process() { + let Ok(path) = std::env::var("BUYBACK_RESERVATION_TEST_DB") else { + return; + }; + let job = std::env::var("BUYBACK_RESERVATION_TEST_JOB").unwrap(); + let rt = tokio::runtime::Runtime::new().unwrap(); + let store = SqliteStore::open(&path).unwrap(); + let won = rt + .block_on(store.reserve_signer("shared-signer", &job)) + .is_ok(); + if won { + std::fs::write(format!("{path}.{job}.won"), b"reserved").unwrap(); + } + // Simulates process death immediately after durable reservation, before journal. + } + + #[cfg(feature = "sqlite")] + #[tokio::test] + async fn two_processes_one_signer_orphan_survives_restart() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("shared.db"); + drop(SqliteStore::open(&path).unwrap()); + let exe = std::env::current_exe().unwrap(); + let mut children = vec![]; + for job in ["job-a", "job-b"] { + children.push( + std::process::Command::new(&exe) + .args(["--exact", "store::tests::reservation_child_process"]) + .env("BUYBACK_RESERVATION_TEST_DB", &path) + .env("BUYBACK_RESERVATION_TEST_JOB", job) + .stdout(std::process::Stdio::null()) + .spawn() + .unwrap(), + ); + } + for mut child in children { + assert!(child.wait().unwrap().success()); + } + let winners = ["job-a", "job-b"] + .iter() + .filter(|job| std::path::Path::new(&format!("{}.{}.won", path.display(), job)).exists()) + .count(); + assert_eq!(winners, 1); + let restarted = SqliteStore::open(&path).unwrap(); + for job in ["job-a", "job-b", "job-c"] { + assert!( + restarted + .reserve_signer("shared-signer", job) + .await + .is_err() + ); + } + assert!( + restarted + .reserve_signer("different-signer", "job-c") + .await + .is_ok() + ); + } + #[cfg(feature = "sqlite")] + #[tokio::test] + async fn legacy_pending_signer_is_quarantined_on_upgrade() { + let store = SqliteStore::in_memory().unwrap(); + let mut rec = test_record("legacy"); + rec.pending = Some(crate::chain::PendingTx { + action: crate::engine::Action::Fund, + reservation: None, + signer: "legacy-signer".into(), + nonce: 1, + tx_hash: "hash".into(), + birth_block: 10, + amount: 1, + }); + store.insert(&rec).await.unwrap(); + assert!( + store + .reserve_signer("legacy-signer", "new-job") + .await + .is_err() + ); + assert!( + store + .reserve_signer("other-signer", "new-job") + .await + .is_ok() + ); + assert!( + FileStore::open(tempfile::tempdir().unwrap().path()) + .unwrap() + .require_signing() + .await + .is_err() + ); + } } diff --git a/tests/localnet.rs b/tests/localnet.rs index 2c8f796..708febf 100644 --- a/tests/localnet.rs +++ b/tests/localnet.rs @@ -134,6 +134,14 @@ async fn full_flow() { destroy: Destroy::Burn, amount: AutoAmount::Fixed(RAO_PER_TAO / 2), }; + cfg.buyback_budget = Some(bittensor_buyback::state::BuybackBudget { + amount_rao: RAO_PER_TAO / 2, + currency: "TAO".into(), + source: "localnet-test-allocation".into(), + hotkey: keys::ss58(&id(&treasury_hk)), + netuid: buy_net, + destroy: Destroy::Burn, + }); let master_hex = MasterKey::generate().to_hex(); let engine = Arc::new(Engine::new( chain.clone(), @@ -318,6 +326,10 @@ async fn full_flow() { PaymentState::Detected ); let fee_tao = RAO_PER_TAO / 100; + let reservation = store + .reserve_signer(&keys::ss58(&id(&treasury)), &req2.id) + .await + .unwrap(); let prepared = chain .prepare( &ChainCall::TransferTao { @@ -329,6 +341,7 @@ async fn full_flow() { .await .unwrap(); let journal = Some(bittensor_buyback::chain::PendingTx { + reservation: Some(reservation), action: Action::Fund, signer: keys::ss58(&id(&treasury)), nonce: prepared.nonce, @@ -336,10 +349,10 @@ async fn full_flow() { birth_block: prepared.birth_block, amount: fee_tao, }); - chain.broadcast(&prepared).await.unwrap(); let mut rec = store.get(&req2.id).await.unwrap().unwrap(); rec.pending = journal; store.update(&rec).await.unwrap(); + chain.broadcast(&prepared).await.unwrap(); let engine2 = Arc::new(Engine::new( chain.clone(), store.clone(), @@ -546,6 +559,14 @@ async fn static_address_scan_and_sweep() { destroy: Destroy::Burn, amount: AutoAmount::Fixed(RAO_PER_TAO / 4), }; + cfg.buyback_budget = Some(bittensor_buyback::state::BuybackBudget { + amount_rao: RAO_PER_TAO / 4, + currency: "TAO".into(), + source: "localnet-test-allocation".into(), + hotkey: keys::ss58(&id(&treasury_hk)), + netuid: buy_net, + destroy: Destroy::Burn, + }); let engine = Engine::new( chain.clone(), store.clone(), From 06640dd9ab72cd2555ef617e94fd090840832437 Mon Sep 17 00:00:00 2001 From: Mathis <154886644+echobt@users.noreply.github.com> Date: Fri, 25 Sep 2026 15:28:26 +0000 Subject: [PATCH 2/8] Serialize localnet fixtures sharing Alice and subnet state --- tests/localnet.rs | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/tests/localnet.rs b/tests/localnet.rs index 708febf..79752f9 100644 --- a/tests/localnet.rs +++ b/tests/localnet.rs @@ -18,6 +18,10 @@ use subxt::dynamic::{self, Value}; use subxt::utils::AccountId32; use subxt_signer::sr25519::Keypair; +// ponytail: these fixtures share Alice and global subnet registration state; +// serialize this test binary until each fixture has an isolated local chain. +static LOCALNET: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(()); + fn url() -> String { std::env::var("BUYBACK_LOCALNET_URL").unwrap_or_else(|_| "ws://127.0.0.1:9944".into()) } @@ -88,6 +92,7 @@ async fn new_subnet(chain: &Chain, owner: &Keypair, hotkey: &Keypair, env: &str) #[tokio::test(flavor = "multi_thread")] async fn full_flow() { + let _chain_guard = LOCALNET.lock().await; let _ = tracing_subscriber::fmt() .with_env_filter("info,bittensor_buyback=debug") .try_init(); @@ -440,6 +445,7 @@ async fn full_flow() { /// treasury and burns a buyback. #[tokio::test(flavor = "multi_thread")] async fn static_address_scan_and_sweep() { + let _chain_guard = LOCALNET.lock().await; let _ = tracing_subscriber::fmt() .with_env_filter("info,bittensor_buyback=debug") .try_init(); From c0ddbc83b182df1005fe31aad1b9799622bdf26f Mon Sep 17 00:00:00 2001 From: Mathis <154886644+echobt@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:01:54 +0000 Subject: [PATCH 3/8] Cancel proven prebroadcast preparation failures and migrate pristine legacy policy --- README.md | 12 ++- src/engine.rs | 118 +++++++++++++++++++- src/store.rs | 290 ++++++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 416 insertions(+), 4 deletions(-) diff --git a/README.md b/README.md index 3bcb1cb..88e50e9 100644 --- a/README.md +++ b/README.md @@ -361,13 +361,21 @@ The budget is frozen at job creation; stores reject later modifications. Missing budgets, insufficient capital, zero purchases, partial spends, and incomplete burns block completion. Legacy dynamic `AutoAmount` values do not allocate job capital. CLI automatic mode requires a numeric `--auto-amount` and `--auto-source`. -Existing jobs without budgets require explicit migration/review, not runtime defaults. +Existing jobs without policy require explicit operator review. `Store::migrate_legacy_policy(id, +expected_version, budget)` supports pristine Pending/Detected legacy jobs in SQLite +and the single-process FileStore; `None` explicitly selects no buyback. Existing +policy, reservations, journals, receipts or other execution evidence refuse migration. +Ambiguous/partially executed legacy jobs still require manual chain reconciliation; +never clear their evidence or overwrite a policy. FileStore migration does not enable signing. SQLite reserves each signer exclusively before preparing a transaction. Reservation ownership is durable and has no lease expiry. The matching pending journal and its release are committed atomically after verified finality or complete mortality-window absence. A crash before journaling leaves an orphan reservation: signing remains -blocked pending explicit operator reconciliation; there is no automatic unlock API. +blocked pending explicit operator reconciliation; there is no expiry-based unlock. +A returned preparation error (before broadcast/journal persistence) cancels only +that live caller's exact token, atomically checking no pending journal exists. +Journal-write uncertainty and cancelled futures retain ownership. Use one shared SQLite database on a filesystem supporting SQLite locking. Different database files, unrelated applications using the same key, and separate hosts with independent stores are **not** protected. Namespace the database per chain. diff --git a/src/engine.rs b/src/engine.rs index 78ecc93..df5af43 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -1034,6 +1034,7 @@ impl EngineRpc for Chain { /// Narrow transaction seam; tests supply no network client or live signer. #[async_trait::async_trait] pub trait TransactionRpc: Send + Sync { + /// Reads and local signing only. Never submits or propagates the signed bytes. async fn prepare(&self, call: &ChainCall, signer: &Keypair) -> Result; async fn broadcast(&self, tx: &PreparedTx) -> Result; async fn find_pending(&self, tx: &PendingTx) -> Result; @@ -1065,9 +1066,21 @@ async fn send_journaled( } let signer_address = keys::ss58(&signer.public_key().to_account_id()); let reservation = store.reserve_signer_for_record(&signer_address, r).await?; - // A crash or prepare error leaves an orphan reservation: never expire it automatically. - let prepared = rpc.prepare(call, signer).await?; + // prepare MUST NOT broadcast. A returned preparation error is known pre-send; + // cancellation/crash of this future still leaves its ownership fenced. + let prepared = match rpc.prepare(call, signer).await { + Ok(prepared) => prepared, + Err(error) => { + store + .cancel_preparation(&signer_address, &r.id, &reservation) + .await?; + return Err(error); + } + }; if prepared.signer != signer_address { + store + .cancel_preparation(&signer_address, &r.id, &reservation) + .await?; return Err(Error::Store("prepared signer mismatch".into())); } let mut journal = r.clone(); @@ -1211,6 +1224,107 @@ mod journal_rpc_tests { use super::*; use crate::{keys::MasterKey, store::SqliteStore}; use std::sync::atomic::{AtomicUsize, Ordering}; + struct PrepareFailure; + #[async_trait::async_trait] + impl TransactionRpc for PrepareFailure { + async fn prepare(&self, _: &ChainCall, _: &Keypair) -> Result { + Err(Error::Chain("prepare RPC timeout before submission".into())) + } + async fn broadcast(&self, _: &PreparedTx) -> Result { + panic!("must not broadcast") + } + async fn find_pending(&self, _: &PendingTx) -> Result { + panic!("no journal") + } + } + #[tokio::test] + async fn returned_prepare_failure_releases_only_prebroadcast_ownership() { + let store = SqliteStore::in_memory().unwrap(); + let mut rec = crate::state::test_record("prepare-error"); + store.insert(&rec).await.unwrap(); + let signer = keys::keypair_from_uri("//Alice").unwrap(); + let address = keys::ss58(&signer.public_key().to_account_id()); + let call = ChainCall::TransferTao { + dest: signer.public_key().to_account_id(), + amount: 1, + }; + for _ in 0..2 { + assert!( + send_journaled( + &PrepareFailure, + &store, + &mut rec, + Action::Fund, + &call, + &signer + ) + .await + .is_err() + ); + assert!(store.get(&rec.id).await.unwrap().unwrap().pending.is_none()); + } + store.reserve_signer(&address, "other-job").await.unwrap(); + } + struct StalledPrepare; + #[async_trait::async_trait] + impl TransactionRpc for StalledPrepare { + async fn prepare(&self, _: &ChainCall, _: &Keypair) -> Result { + std::future::pending().await + } + async fn broadcast(&self, _: &PreparedTx) -> Result { + panic!("must not broadcast") + } + async fn find_pending(&self, _: &PendingTx) -> Result { + panic!("no journal") + } + } + #[tokio::test] + async fn cancelled_preparation_and_crash_before_journal_stay_fenced() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("prepare.sqlite"); + let store = SqliteStore::open(&path).unwrap(); + let mut rec = crate::state::test_record("cancelled-prepare"); + store.insert(&rec).await.unwrap(); + let signer = keys::keypair_from_uri("//Alice").unwrap(); + let address = keys::ss58(&signer.public_key().to_account_id()); + let call = ChainCall::TransferTao { + dest: signer.public_key().to_account_id(), + amount: 1, + }; + assert!( + tokio::time::timeout( + std::time::Duration::from_millis(10), + send_journaled( + &StalledPrepare, + &store, + &mut rec, + Action::Fund, + &call, + &signer + ) + ) + .await + .is_err() + ); + drop(store); + let store = SqliteStore::open(&path).unwrap(); + assert!(store.reserve_signer(&address, "other").await.is_err()); + let bob = keys::keypair_from_uri("//Bob").unwrap(); + let bob_address = keys::ss58(&bob.public_key().to_account_id()); + let rpc = LostReply { + broadcasts: AtomicUsize::new(0), + summary: EventSummary::default(), + lookup: AtomicUsize::new(0), + dead: false, + }; + store.reserve_signer(&bob_address, &rec.id).await.unwrap(); + let prepared = rpc.prepare(&call, &bob).await.unwrap(); + assert_eq!(prepared.signer, bob_address); + drop(store); // process disappears after prepare, before journal persistence + let store = SqliteStore::open(&path).unwrap(); + assert!(store.reserve_signer(&bob_address, "other").await.is_err()); + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 0); + } struct LostReply { broadcasts: AtomicUsize, summary: EventSummary, diff --git a/src/store.rs b/src/store.rs index ddcf385..773af95 100644 --- a/src/store.rs +++ b/src/store.rs @@ -31,6 +31,20 @@ pub trait Store: Send + Sync + 'static { "shared durable signer reservations unsupported; signing disabled".into(), )) } + /// Only the live preparing caller may release its token before any broadcast. + /// Never use this for crash recovery, journal-write errors, or pending transactions. + async fn cancel_preparation(&self, _signer: &str, _job: &str, _token: &str) -> Result<()> { + Err(Error::Store("preparation cancellation unsupported".into())) + } + /// Explicit operator migration of a pristine legacy record, never a runtime default. + async fn migrate_legacy_policy( + &self, + _id: &str, + _version: u64, + _budget: Option, + ) -> Result { + Err(Error::Store("legacy policy migration unsupported".into())) + } /// Insert a new record (version 0). Fails if the id exists. async fn insert(&self, rec: &PaymentRecord) -> Result<()>; async fn get(&self, id: &str) -> Result>; @@ -45,6 +59,45 @@ fn store_err(e: impl std::fmt::Display) -> Error { Error::Store(e.to_string()) } +fn migrated_policy( + mut rec: PaymentRecord, + version: u64, + budget: Option, +) -> Result { + if rec.version != version { + return Err(Error::Conflict(rec.id)); + } + if rec.auto_required.is_some() + || rec.buyback_budget.is_some() + || rec.pending.is_some() + || rec.quarantined.is_some() + || !rec.txs.is_empty() + || rec.buyback.is_some() + || rec.buyback_done + || rec.funded_tao != 0 + || rec.swept_alpha != 0 + || !rec.swept_positions.is_empty() + || rec.consolidated != 0 + || rec.dust_returned + || !matches!(rec.state, PaymentState::Pending | PaymentState::Detected) + { + return Err(Error::Store( + "legacy job has policy or execution evidence; manual reconciliation required".into(), + )); + } + if let Some(b) = &budget { + b.validate()?; + } + rec.auto_required = Some(budget.is_some()); + rec.buyback_budget = budget; + rec.version = rec + .version + .checked_add(1) + .ok_or_else(|| Error::Store("version overflow".into()))?; + rec.updated_at = crate::state::now(); + Ok(rec) +} + /// One JSON file per payment in a directory (mode 0700 on unix). Atomic writes via rename. /// Single-process only: CAS is guarded by an in-process mutex. pub struct FileStore { @@ -108,6 +161,21 @@ impl FileStore { #[async_trait::async_trait] impl Store for FileStore { + async fn migrate_legacy_policy( + &self, + id: &str, + version: u64, + budget: Option, + ) -> Result { + let _g = self.lock.lock().map_err(store_err)?; + let path = self.path(id)?; + let previous = self + .read(&path)? + .ok_or_else(|| Error::NotFound(id.into()))?; + let next = migrated_policy(previous, version, budget)?; + self.write(&path, &next)?; + Ok(next) + } async fn insert(&self, rec: &PaymentRecord) -> Result<()> { let _g = self.lock.lock().map_err(store_err)?; let p = self.path(&rec.id)?; @@ -239,6 +307,64 @@ impl Store for SqliteStore { self.reserve(signer, job, None) } + async fn cancel_preparation(&self, signer: &str, job: &str, token: &str) -> Result<()> { + let mut c = self.conn.lock().map_err(store_err)?; + let tx = c + .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate) + .map_err(store_err)?; + let removed = tx.execute( + "DELETE FROM signer_reservations WHERE signer=?1 AND job=?2 AND token=?3 AND NOT EXISTS (SELECT 1 FROM payments WHERE json_extract(record,'$.pending.signer')=?1 OR (id=?2 AND json_type(record,'$.pending')='object'))", + rusqlite::params![signer,job,token], + ).map_err(store_err)?; + if removed != 1 { + return Err(Error::Conflict( + "preparation token or pending journal changed".into(), + )); + } + tx.commit().map_err(store_err) + } + async fn migrate_legacy_policy( + &self, + id: &str, + version: u64, + budget: Option, + ) -> Result { + let mut c = self.conn.lock().map_err(store_err)?; + let tx = c + .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate) + .map_err(store_err)?; + let reserved: bool = tx + .query_row( + "SELECT EXISTS(SELECT 1 FROM signer_reservations WHERE job=?1)", + [id], + |r| r.get(0), + ) + .map_err(store_err)?; + if reserved { + return Err(Error::Conflict("job has a signer reservation".into())); + } + let json: String = tx + .query_row("SELECT record FROM payments WHERE id=?1", [id], |r| { + r.get(0) + }) + .map_err(store_err)?; + let next = migrated_policy( + serde_json::from_str(&json).map_err(store_err)?, + version, + budget, + )?; + tx.execute( + "UPDATE payments SET version=?1,record=?2 WHERE id=?3", + rusqlite::params![ + i64::try_from(next.version).map_err(store_err)?, + serde_json::to_string(&next).map_err(store_err)?, + id + ], + ) + .map_err(store_err)?; + tx.commit().map_err(store_err)?; + Ok(next) + } async fn insert(&self, rec: &PaymentRecord) -> Result<()> { let c = self.conn.lock().map_err(store_err)?; c.execute( @@ -392,6 +518,17 @@ impl Store for Box { (**self).reserve_signer(signer, job).await } + async fn cancel_preparation(&self, s: &str, j: &str, t: &str) -> Result<()> { + (**self).cancel_preparation(s, j, t).await + } + async fn migrate_legacy_policy( + &self, + id: &str, + v: u64, + b: Option, + ) -> Result { + (**self).migrate_legacy_policy(id, v, b).await + } async fn insert(&self, rec: &PaymentRecord) -> Result<()> { (**self).insert(rec).await } @@ -411,6 +548,159 @@ mod tests { use super::*; use crate::state::test_record; + #[tokio::test] + async fn explicit_legacy_policy_migration_preserves_evidence() { + let dir = tempfile::tempdir().unwrap(); + let file = FileStore::open(dir.path()).unwrap(); + #[allow(unused_mut)] + let mut stores: Vec> = vec![Box::new(file)]; + #[cfg(feature = "sqlite")] + stores.push(Box::new(SqliteStore::in_memory().unwrap())); + for store in stores { + let mut legacy = test_record("legacy-policy"); + legacy.auto_required = None; + store.insert(&legacy).await.unwrap(); + assert!( + store + .migrate_legacy_policy(&legacy.id, 1, None) + .await + .is_err() + ); + let migrated = store + .migrate_legacy_policy(&legacy.id, 0, None) + .await + .unwrap(); + assert_eq!(migrated.auto_required, Some(false)); + assert_eq!(migrated.version, 1); + assert!( + store + .migrate_legacy_policy(&legacy.id, 1, None) + .await + .is_err() + ); + let mut evidence = legacy.clone(); + evidence.id = "legacy-evidence".into(); + evidence.funded_tao = 1; + store.insert(&evidence).await.unwrap(); + assert!( + store + .migrate_legacy_policy(&evidence.id, 0, None) + .await + .is_err() + ); + assert_eq!( + store.get(&evidence.id).await.unwrap().unwrap().funded_tao, + 1 + ); + } + } + + #[cfg(feature = "sqlite")] + #[tokio::test] + async fn preparation_cancellation_requires_current_token_and_no_journal() { + let store = SqliteStore::in_memory().unwrap(); + let mut rec = test_record("prepare-cancel"); + rec.auto_required = None; + store.insert(&rec).await.unwrap(); + let token = store + .reserve_signer_for_record("signer", &rec) + .await + .unwrap(); + assert!( + store + .migrate_legacy_policy(&rec.id, rec.version, None) + .await + .is_err() + ); + assert!( + store + .cancel_preparation("signer", &rec.id, "wrong") + .await + .is_err() + ); + store + .cancel_preparation("signer", &rec.id, &token) + .await + .unwrap(); + let replacement = store + .reserve_signer_for_record("signer", &rec) + .await + .unwrap(); + assert!( + store + .cancel_preparation("signer", &rec.id, &token) + .await + .is_err() + ); + rec.pending = Some(crate::chain::PendingTx { + action: crate::engine::Action::Fund, + reservation: Some(replacement.clone()), + signer: "signer".into(), + nonce: 0, + tx_hash: "hash".into(), + birth_block: 1, + amount: 1, + }); + store.update(&rec).await.unwrap(); + assert!( + store + .cancel_preparation("signer", &rec.id, &replacement) + .await + .is_err() + ); + assert!(store.reserve_signer("signer", "other").await.is_err()); + } + + #[cfg(feature = "sqlite")] + #[tokio::test] + async fn legacy_migration_fences_stale_reservers_and_freezes_validated_budget() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("migration.sqlite"); + let first = SqliteStore::open(&path).unwrap(); + let second = SqliteStore::open(&path).unwrap(); + let mut rec = test_record("legacy-budget"); + rec.auto_required = None; + first.insert(&rec).await.unwrap(); + let mut budget = crate::state::BuybackBudget { + amount_rao: 0, + currency: "TAO".into(), + source: "operator-allocation".into(), + netuid: 100, + destroy: crate::config::Destroy::Burn, + hotkey: crate::keys::ss58( + &crate::keys::keypair_from_uri("//Bob") + .unwrap() + .public_key() + .to_account_id(), + ), + }; + assert!( + first + .migrate_legacy_policy(&rec.id, 0, Some(budget.clone())) + .await + .is_err() + ); + budget.amount_rao = crate::engine::MIN_STAKE_RAO; + let migrated = first + .migrate_legacy_policy(&rec.id, 0, Some(budget.clone())) + .await + .unwrap(); + assert_eq!(migrated.auto_required, Some(true)); + assert_eq!( + second.get(&rec.id).await.unwrap().unwrap().buyback_budget, + Some(budget) + ); + assert!( + second + .reserve_signer_for_record("signer", &rec) + .await + .is_err() + ); + let mut changed = migrated; + changed.buyback_budget = None; + assert!(second.update(&changed).await.is_err()); + } + async fn exercise(s: &dyn Store) { let r = test_record("p-1"); s.insert(&r).await.unwrap(); From 9c1d778486d6bf36275da948e252179eb77cc108 Mon Sep 17 00:00:00 2001 From: Mathis <154886644+echobt@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:13:49 +0000 Subject: [PATCH 4/8] Refuse startup association and unjournaled standalone buybacks --- README.md | 31 +++++----- src/engine.rs | 144 ++++++++++++++++++++++++---------------------- tests/localnet.rs | 91 ++++++++++++++--------------- 3 files changed, 134 insertions(+), 132 deletions(-) diff --git a/README.md b/README.md index 88e50e9..b3565d6 100644 --- a/README.md +++ b/README.md @@ -102,8 +102,11 @@ The localnet test covers this directly. It sends the funding transfer, "crashes" it, starts a new engine on the same store, and asserts the payment settles with the fee sent **exactly once**. -Standalone `buyback*()` calls are *not* journaled. They are operator actions. If one errors, -check `buyback balances` before running it again. +Standalone `Engine::buyback*()` calls and the CLI buyback command now fail closed with +`Error::Config` before signing or submission. This is a compatibility change: use a durable +payment/sweep job with an explicit immutable `BuybackBudget` and drive `Engine::tick`. +There is no standalone-job replacement yet. `Chain::prepare`/`broadcast`/`submit` remain +low-level primitives without Store ownership guarantees; never mix them with live workers. ## Static deposit addresses @@ -150,17 +153,14 @@ cfg.buyback_budget = Some(state::BuybackBudget { hotkey: keys::ss58(&cfg.treasury_hotkey), }); let engine = Arc::new(Engine::new(chain, store, master, treasury, cfg)); -engine.ensure_treasury_hotkey().await?; // try_associate_hotkey if missing +engine.ensure_treasury_hotkey().await?; // read-only; fails if hotkey is missing let req: PaymentRequest = engine.create_payment(CreatePayment::default()).await?; let status: PaymentStatus = engine.status(&req.id).await?; let mut settled = engine.subscribe(); // in-process callback tokio::spawn({ let e = engine.clone(); async move { e.run().await } }); -engine.buyback(units::parse_amount("10")?).await?; // buy, keep staked -engine.buyback_and_burn(units::parse_amount("10")?).await?; // buy, burn_alpha -engine.buyback_and_recycle(units::RAO_PER_TAO).await?; // buy, recycle_alpha -engine.buyback_all(Destroy::Burn).await?; // whole balance - fee_reserve +// Standalone buyback wrappers are disabled; the job above carries its budget. # Ok(()) } ``` @@ -171,8 +171,8 @@ engine.buyback_all(Destroy::Burn).await?; // whole balance - f | `Engine::run()` / `Engine::tick()` | follow finalized blocks / one pass | | `Engine::subscribe()` | `broadcast::Receiver` for settled/failed/expired | | `Engine::retry(id)`, `Engine::force_sweep(id)` | operator recovery | -| `Engine::buyback`, `buyback_and_burn`, `buyback_and_recycle`, `buyback_all`, `buyback_with` | treasury buybacks, return `BuybackReceipt` | -| `Engine::ensure_treasury_hotkey()` | create the treasury hotkey account if missing | +| `Engine::buyback`, `buyback_and_burn`, `buyback_and_recycle`, `buyback_all`, `buyback_with` | disabled; return `Error::Config`, migrate to budgeted jobs | +| `Engine::ensure_treasury_hotkey()` | verify a pre-provisioned treasury hotkey; never submit at startup | | `Chain` | `verify_metadata`, `stake_positions`, `alpha_price`, `estimate_fee`, `submit`, `find_pending` | | `Store` trait, `SqliteStore`, `FileStore`, `open_store("sqlite:..."/"file:...")` | persistence with compare-and-swap | | `MasterKey`, `PaymentWallet`, `TreasuryKeySource` | key custody | @@ -192,7 +192,7 @@ export BUYBACK_NETWORK=local BUYBACK_NETUID=2 BUYBACK_TREASURY_HOTKEY=5... buyback create --metadata '{"order":123}' buyback status buyback run --listen 127.0.0.1:8080 # watcher + HTTP (http feature) -buyback buyback 10 --then burn # or `all`, --then keep|burn|recycle +buyback buyback 10 --then burn # disabled: use budgeted durable jobs buyback balances ``` @@ -270,7 +270,7 @@ subnet issuance is tracked | effect | alpha is permanently destroyed and counted as burned supply. Issuance keeps counting it, so the emission schedule (which follows issuance) is unaffected | alpha goes back to the "unissued" pool: issuance and outstanding supply shrink, so it can be emitted again | For a "buy and burn" that permanently removes supply and shows up as burned, use `burn_alpha`, -which is what `buyback_and_burn` does. `buyback_and_recycle` and `Destroy::Recycle` are there if +which is what budgeted jobs with `Destroy::Burn` do. `Destroy::Recycle` is available if you want the alpha returned to future emissions instead. Neither can be undone. Neither works on the root subnet (`CannotBurnOrRecycleOnRootSubnet`), and the hotkey must exist on chain. @@ -325,8 +325,8 @@ alpha to a fresh deposit address. The test then checks: balances before and after; * a simulated restart mid-flow; * a crash-after-broadcast recovery that must not fund twice; -* `buyback`, `buyback_and_burn` and `buyback_and_recycle`, with their balance and event deltas; -* the fee-reserve guard. +* standalone wrappers refused without a treasury balance change; +* durable job budget and fee-reserve guards. It only uses well-known dev accounts (`//Alice`, `//Bob` and their derivations). CI runs it against the localnet image as a service container. @@ -336,8 +336,9 @@ against the localnet image as a service container. 1. Run `buyback verify-metadata --network finney` and confirm every call and event resolves. 2. Generate the master key and store it in your secret manager. Seal the treasury mnemonic into a keystore (`seal-treasury`). -3. Pick the treasury hotkey. Either use a hotkey you own and have registered, or let - `ensure_treasury_hotkey` run `try_associate_hotkey`. It must not be a subnet system account. +3. Provision the treasury hotkey before starting signing workers. + `ensure_treasury_hotkey` is read-only and fails when the hotkey is absent. + It must not be a subnet system account. 4. Fund the treasury coldkey with working TAO for fees, and set `fee_reserve`. 5. Start with `auto = off` and a small `min_alpha`. Do one real payment, check `buyback status`, then enable auto-buyback. diff --git a/src/engine.rs b/src/engine.rs index df5af43..20fcc40 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -815,30 +815,26 @@ impl Engine { // ------------------------------------------------------------------ buybacks - /// Make sure `treasury_hotkey` exists on chain and is owned by the treasury coldkey, creating - /// it with `try_associate_hotkey` if needed. Staking, `move_stake` and `burn_alpha` all fail - /// with `HotKeyAccountNotExists` otherwise. Call once at startup. + /// Verify the treasury hotkey already exists. Never signs or submits at startup. + /// Association must be provisioned separately before starting signing workers. pub async fn ensure_treasury_hotkey(&self) -> Result<()> { let hk = &self.cfg.treasury_hotkey; match self.chain.hotkey_owner(hk).await? { Some(o) if o == self.treasury_account() => Ok(()), Some(o) => { - // Staking to a hotkey you do not own works but delegates to its owner (take - // applies) and `burn_alpha` still works; warn loudly rather than fail. tracing::warn!(hotkey = %keys::ss58(hk), owner = %keys::ss58(&o), "treasury hotkey is owned by another coldkey"); Ok(()) } - None => { - let _g = self.treasury_lock.lock().await; - let call = ChainCall::AssociateHotkey { hotkey: *hk }; - self.chain.submit(&call, &self.treasury).await?; - tracing::info!(hotkey = %keys::ss58(hk), "treasury hotkey associated"); - Ok(()) - } + // ponytail: startup is read-only; provision association separately until + // association has its own durable job and reconciliation protocol. + None => Err(Error::Config( + "treasury hotkey is not associated; provision it before starting signing workers" + .into(), + )), } } - /// Treasury TAO available to `buyback_all`: free minus `fee_reserve` minus existential deposit. + /// Treasury TAO available to budgeted jobs: free minus `fee_reserve` minus existential deposit. pub async fn spendable(&self) -> Result { let free = self.chain.free_balance(&self.treasury_account()).await?; Ok(units::spendable( @@ -848,77 +844,36 @@ impl Engine { )) } - /// Spend `amount_tao` rao of treasury TAO on alpha of `buyback_netuid` (default 100) with - /// `add_stake_limit` at `spot * (1 + slippage)`. The alpha stays staked to the treasury hotkey. + /// Disabled standalone API. Use a durable payment/sweep job with an explicit budget. pub async fn buyback(&self, amount_tao: u64) -> Result { self.buyback_with(amount_tao, Destroy::Keep).await } - /// [`Self::buyback`] then `burn_alpha` of exactly the alpha bought. + /// Disabled standalone API; see [`Self::buyback`]. pub async fn buyback_and_burn(&self, amount_tao: u64) -> Result { self.buyback_with(amount_tao, Destroy::Burn).await } - /// [`Self::buyback`] then `recycle_alpha` of exactly the alpha bought. + /// Disabled standalone API; see [`Self::buyback`]. pub async fn buyback_and_recycle(&self, amount_tao: u64) -> Result { self.buyback_with(amount_tao, Destroy::Recycle).await } - /// Buy back with the whole spendable treasury balance. + /// Disabled: an entire live balance is not an immutable job budget. pub async fn buyback_all(&self, destroy: Destroy) -> Result { - let amount = self.spendable().await?; - self.buyback_with(amount, destroy).await + self.buyback_with(0, destroy).await } - /// Standalone buyback. Not journaled: on an error the caller must check the treasury's - /// stake/balance (e.g. `buyback status`) before retrying. - pub async fn buyback_with(&self, amount_tao: u64, destroy: Destroy) -> Result { - let _g = self.treasury_lock.lock().await; - let netuid = self.cfg.buyback_netuid; - let spendable = self.spendable().await?; - if amount_tao == 0 || amount_tao > spendable { - return Err(Error::Insufficient(format!( - "buyback of {} TAO, spendable {} TAO (fee reserve {})", - units::format_amount(amount_tao), - units::format_amount(spendable), - units::format_amount(self.cfg.fee_reserve) - ))); - } - let price = self.chain.alpha_price(netuid).await?; - let limit_price = units::buy_limit_price(price, self.cfg.slippage_bps); - let hotkey = self.cfg.treasury_hotkey; - let buy = ChainCall::AddStakeLimit { - hotkey, - netuid, - tao: amount_tao, - limit_price, - allow_partial: self.cfg.allow_partial, - }; - let f = self.chain.submit(&buy, &self.treasury).await?; - let (tao_spent, alpha_bought) = f - .summary - .stake_added - .ok_or(Error::EventMissing("StakeAdded"))?; - let mut receipt = BuybackReceipt { - netuid, - tao_spent, - alpha_bought, - limit_price, - stake_tx: f.tx, - destroy_tx: None, - alpha_destroyed: 0, - }; - if destroy != Destroy::Keep && alpha_bought > 0 { - let call = destroy_call(destroy, hotkey, netuid, alpha_bought); - let d = self.chain.submit(&call, &self.treasury).await?; - receipt.alpha_destroyed = d - .summary - .alpha_destroyed - .ok_or(Error::EventMissing("AlphaBurned"))?; - receipt.destroy_tx = Some(d.tx); - } - tracing::info!(?receipt, "buyback done"); - Ok(receipt) + /// Refuses before any RPC/signature. Standalone calls have no durable job journal. + pub async fn buyback_with( + &self, + _amount_tao: u64, + _destroy: Destroy, + ) -> Result { + // ponytail: standalone execution stays disabled until it has durable job semantics. + Err(Error::Config( + "standalone buyback disabled; use a durable payment/sweep job with an explicit buyback budget".into(), + )) } } @@ -1648,6 +1603,7 @@ mod journal_rpc_tests { struct SimulatedChain { broadcasts: AtomicUsize, hotkey: AccountId32, + owner: Option, result: std::sync::Mutex, } #[async_trait::async_trait] @@ -1725,7 +1681,7 @@ mod journal_rpc_tests { Ok(1_000_000_000) } async fn hotkey_owner(&self, _: &AccountId32) -> Result> { - panic!("no association") + Ok(self.owner) } async fn stake_positions_many( &self, @@ -1755,6 +1711,53 @@ mod journal_rpc_tests { self.as_ref().find_pending(t).await } } + #[tokio::test] + async fn startup_missing_hotkey_never_submits_or_releases_ownership() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("startup.db"); + let treasury = PaymentWallet::generate().unwrap(); + let signer = keys::ss58(&treasury.account_id()); + let store = Arc::new(SqliteStore::open(&path).unwrap()); + store.reserve_signer(&signer, "unresolved").await.unwrap(); + let rpc = Arc::new(SimulatedChain { + broadcasts: AtomicUsize::new(0), + hotkey: treasury.account_id(), + owner: None, + result: std::sync::Mutex::new(EventSummary::default()), + }); + let engine = Engine::new( + rpc.clone(), + store, + MasterKey::from_bytes([42; 32]), + treasury.keypair().clone(), + Config::new("unused", 100, treasury.account_id()), + ); + let (a, b) = tokio::join!( + engine.ensure_treasury_hotkey(), + engine.ensure_treasury_hotkey() + ); + assert!(matches!(a, Err(Error::Config(_)))); + assert!(matches!(b, Err(Error::Config(_)))); + for result in [ + engine.buyback(1).await, + engine.buyback_and_burn(1).await, + engine.buyback_and_recycle(1).await, + engine.buyback_all(Destroy::Burn).await, + engine.buyback_with(1, Destroy::Keep).await, + ] { + assert!(matches!(result, Err(Error::Config(_)))); + } + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 0); + drop(engine); + assert!( + SqliteStore::open(&path) + .unwrap() + .reserve_signer(&signer, "other") + .await + .is_err() + ); + } + #[tokio::test] async fn engine_ticks_complete_after_each_lost_response_and_restart() { let dir = tempfile::tempdir().unwrap(); @@ -1786,6 +1789,7 @@ mod journal_rpc_tests { let rpc = Arc::new(SimulatedChain { broadcasts: AtomicUsize::new(0), hotkey: treasury.account_id(), + owner: Some(treasury.account_id()), result: std::sync::Mutex::new(EventSummary::default()), }); let mut cfg = Config::new("unused", 100, treasury.account_id()); diff --git a/tests/localnet.rs b/tests/localnet.rs index 79752f9..a17a2de 100644 --- a/tests/localnet.rs +++ b/tests/localnet.rs @@ -156,6 +156,22 @@ async fn full_flow() { cfg.clone(), )); let mut settled = engine.subscribe(); + // Explicit isolated-localnet provisioning, never engine startup behavior. + if chain + .hotkey_owner(&id(&treasury_hk)) + .await + .unwrap() + .is_none() + { + assert!(engine.ensure_treasury_hotkey().await.is_err()); + raw( + &chain, + &treasury, + "try_associate_hotkey", + vec![Value::from_bytes(id(&treasury_hk).0)], + ) + .await; + } engine.ensure_treasury_hotkey().await.unwrap(); assert_eq!( chain.hotkey_owner(&id(&treasury_hk)).await.unwrap(), @@ -391,53 +407,18 @@ async fn full_flow() { ); assert!(s2.swept_alpha >= RAO_PER_TAO); - // --- explicit buyback (keep) and buyback_and_burn --- - let thk = id(&treasury_hk); - let a0 = chain.alpha_of(&thk, &tre, buy_net).await.unwrap(); - let t0 = chain.free_balance(&tre).await.unwrap(); - let r1 = engine.buyback(RAO_PER_TAO).await.unwrap(); - let a1 = chain.alpha_of(&thk, &tre, buy_net).await.unwrap(); - let t1 = chain.free_balance(&tre).await.unwrap(); - println!( - "--- buyback 1 TAO ---\n{r1:?}\ntreasury TAO {} -> {}, alpha(netuid {buy_net}) {} -> {}", - format_amount(t0), - format_amount(t1), - format_amount(a0), - format_amount(a1) - ); - assert_eq!(r1.tao_spent, RAO_PER_TAO); - assert!(r1.alpha_bought > 0 && r1.destroy_tx.is_none()); - assert!(t1 <= t0 - RAO_PER_TAO); - assert!((a1 - a0).abs_diff(r1.alpha_bought) <= 1, "{a0} -> {a1}"); - - let r2 = engine.buyback_and_burn(RAO_PER_TAO).await.unwrap(); - let a2 = chain.alpha_of(&thk, &tre, buy_net).await.unwrap(); - let t2 = chain.free_balance(&tre).await.unwrap(); - println!( - "--- buyback_and_burn 1 TAO ---\n{r2:?}\ntreasury TAO {} -> {}, alpha(netuid {buy_net}) {} -> {}", - format_amount(t1), - format_amount(t2), - format_amount(a1), - format_amount(a2) - ); - assert!(r2.destroy_tx.is_some()); - assert_eq!(r2.alpha_destroyed, r2.alpha_bought); - assert!(t2 <= t1 - RAO_PER_TAO); - assert!( - a2.abs_diff(a1) <= 1, - "bought alpha was burned: {a1} -> {a2}" - ); - - // recycle variant - let r3 = engine.buyback_and_recycle(RAO_PER_TAO / 2).await.unwrap(); - println!("--- buyback_and_recycle 0.5 TAO ---\n{r3:?}"); - assert_eq!(r3.alpha_destroyed, r3.alpha_bought); - - // Guard: cannot spend beyond the reserve. - assert!(matches!( - engine.buyback(u64::MAX / 2).await, - Err(Error::Insufficient(_)) - )); + // Standalone wrappers cannot bypass the durable job protocol. + let balance = chain.free_balance(&tre).await.unwrap(); + for result in [ + engine.buyback(RAO_PER_TAO).await, + engine.buyback_and_burn(RAO_PER_TAO).await, + engine.buyback_and_recycle(RAO_PER_TAO).await, + engine.buyback_all(Destroy::Burn).await, + engine.buyback_with(RAO_PER_TAO, Destroy::Keep).await, + ] { + assert!(matches!(result, Err(Error::Config(_)))); + } + assert_eq!(chain.free_balance(&tre).await.unwrap(), balance); } /// Static deposit address: a deterministic wallet receives two `transfer_stake`s, the block @@ -581,6 +562,22 @@ async fn static_address_scan_and_sweep() { cfg, ) .with_seed(seed.clone()); + // Explicit isolated-localnet provisioning, never engine startup behavior. + if chain + .hotkey_owner(&id(&treasury_hk)) + .await + .unwrap() + .is_none() + { + assert!(engine.ensure_treasury_hotkey().await.is_err()); + raw( + &chain, + &treasury, + "try_associate_hotkey", + vec![Value::from_bytes(id(&treasury_hk).0)], + ) + .await; + } engine.ensure_treasury_hotkey().await.unwrap(); // wrong expected address is refused before anything is stored assert!( From 07fa61d6e217278cd3b1ed71e86082a57af6acba Mon Sep 17 00:00:00 2001 From: Mathis <154886644+echobt@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:33:37 +0000 Subject: [PATCH 5/8] Reject disabled CLI buyback before initializing network or secrets --- src/bin/buyback.rs | 14 ++++---------- tests/disabled_cli.rs | 26 ++++++++++++++++++++++++++ 2 files changed, 30 insertions(+), 10 deletions(-) create mode 100644 tests/disabled_cli.rs diff --git a/src/bin/buyback.rs b/src/bin/buyback.rs index fd9eabd..56c73e0 100644 --- a/src/bin/buyback.rs +++ b/src/bin/buyback.rs @@ -318,16 +318,10 @@ async fn main() -> R<()> { } e.run().await?; } - Cmd::Buyback { c, amount, then } => { - let e = engine(&c).await?; - e.ensure_treasury_hotkey().await?; - let d = destroy(&then)?; - let r = if amount == "all" { - e.buyback_all(d).await? - } else { - e.buyback_with(units::parse_amount(&amount)?, d).await? - }; - print(&r)?; + Cmd::Buyback { .. } => { + return Err(bittensor_buyback::Error::Config( + "standalone buyback disabled; use a durable payment/sweep job with an explicit buyback budget".into(), + ).into()); } Cmd::Balances { c } => { let e = engine(&c).await?; diff --git a/tests/disabled_cli.rs b/tests/disabled_cli.rs new file mode 100644 index 0000000..d6a6e8e --- /dev/null +++ b/tests/disabled_cli.rs @@ -0,0 +1,26 @@ +//! Disabled CLI commands must not initialize secrets, storage, or network clients. +#[test] +fn standalone_buyback_refuses_without_connecting() { + let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + listener.set_nonblocking(true).unwrap(); + let output = std::process::Command::new(env!("CARGO_BIN_EXE_buyback")) + .env_clear() + .args([ + "buyback", + "1", + "--netuid", + "100", + "--treasury-hotkey", + "unused", + "--network", + &format!("ws://{}", listener.local_addr().unwrap()), + ]) + .output() + .unwrap(); + assert!(!output.status.success()); + assert!(String::from_utf8_lossy(&output.stderr).contains("standalone buyback disabled")); + assert_eq!( + listener.accept().unwrap_err().kind(), + std::io::ErrorKind::WouldBlock + ); +} From b9e5dd954766425f15d259dacbedfee9ce9bb4f3 Mon Sep 17 00:00:00 2001 From: Mathis <154886644+echobt@users.noreply.github.com> Date: Fri, 25 Sep 2026 16:34:28 +0000 Subject: [PATCH 6/8] Document missing evidence blocking partially executed legacy recovery --- README.md | 21 +++++++++++++++++++-- 1 file changed, 19 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index b3565d6..d095fff 100644 --- a/README.md +++ b/README.md @@ -369,6 +369,23 @@ policy, reservations, journals, receipts or other execution evidence refuse migr Ambiguous/partially executed legacy jobs still require manual chain reconciliation; never clear their evidence or overwrite a policy. FileStore migration does not enable signing. +This is a blocked recovery path, not evidence that funds were lost. Records and receipts +are preserved; rejection does not manufacture a persisted quarantine record. +`TxRef` retains action, extrinsic/block hashes and amount, but not signer, nonce or +full call targets/parameters. Those require authenticated finalized archive data. +Neither chain history nor old receipts establish the original auto-buyback intent, +capital source/allocation or destruction policy; mutable runtime defaults cannot recover +that intent. Records also lack a genesis hash, and receipts alone cannot establish that +no additional pre-journal action was emitted before a crash. + +Recovery therefore requires an audited operator manifest identifying network, signers, +intended budget and steps, plus complete finalized extrinsics/events covering execution +and any uncertain interval. An eventual importer must first dry-run hash, signer, nonce, +target, amount and step checks, then apply a version-checked update without emitting any +transaction or deleting evidence. No such importer exists here. Missing evidence means +refusal, not automatic policy assignment, release or retry. + + SQLite reserves each signer exclusively before preparing a transaction. Reservation ownership is durable and has no lease expiry. The matching pending journal and its release are committed atomically after verified finality or complete mortality-window @@ -383,8 +400,8 @@ independent stores are **not** protected. Namespace the database per chain. `FileStore` and custom `Store` implementations without durable reservations refuse engine signing. A PostgreSQL adapter must implement equivalent shared ownership and -atomic journal release before enabling signing. Standalone operator methods remain -unjournaled; never use them for automatic payment processing. +atomic journal release before enabling signing. Standalone Engine buyback methods +are disabled; raw Chain primitives have no durable ownership guarantee. Offline regressions drive the actual `Engine::tick` through funding, sweeping, purchase and burn, reopening the store and reconstructing the engine after every From 3b1b6290ebf9a61165b7a1b3b6cb799b8abaf2cd Mon Sep 17 00:00:00 2001 From: Mathis <154886644+echobt@users.noreply.github.com> Date: Fri, 25 Sep 2026 19:58:54 +0000 Subject: [PATCH 7/8] Recover evidenced legacy jobs through audited archive verification --- .github/workflows/ci.yml | 1 + README.md | 35 +- src/bin/buyback.rs | 48 +++ src/chain.rs | 98 +++++- src/engine.rs | 69 ++++ src/lib.rs | 1 + src/recovery.rs | 715 +++++++++++++++++++++++++++++++++++++++ src/state.rs | 14 + src/store.rs | 74 +++- 9 files changed, 1047 insertions(+), 8 deletions(-) create mode 100644 src/recovery.rs diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9d5adeb..918fb02 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -64,3 +64,4 @@ jobs: echo "localnet did not start"; exit 1 - run: cargo run --example verify_metadata -- local - run: cargo test --all-features --test localnet -- --nocapture + - run: cargo test --all-features --lib archive_dry_run_apply_and_repeated_import -- --nocapture diff --git a/README.md b/README.md index d095fff..9d69f88 100644 --- a/README.md +++ b/README.md @@ -378,12 +378,35 @@ capital source/allocation or destruction policy; mutable runtime defaults cannot that intent. Records also lack a genesis hash, and receipts alone cannot establish that no additional pre-journal action was emitted before a crash. -Recovery therefore requires an audited operator manifest identifying network, signers, -intended budget and steps, plus complete finalized extrinsics/events covering execution -and any uncertain interval. An eventual importer must first dry-run hash, signer, nonce, -target, amount and step checks, then apply a version-checked update without emitting any -transaction or deleting evidence. No such importer exists here. Missing evidence means -refusal, not automatic policy assignment, release or retry. +`buyback recover-legacy --manifest manifest.json --database jobs.sqlite --archive +wss://trusted-archive` performs a read-only dry run. Add `--apply` only during offline +operator maintenance. No signing keys are loaded. The `recovery::LegacyManifest` +requires genesis, treasury, treasury hotkey, explicit consolidation/dust settings, +budget (including explicit null for no buyback), exact extrinsics with block/index, +signer/nonce, receipts and full `ChainCall` parameters. New jobs freeze those identity +and configuration fields; a mismatched or missing identity blocks processing, including +pending reconciliation. The older budget-only migration does not establish identity. + +An operator must attest that ALL old emitters stopped, that `first_block` covers ALL +old emissions, and the historical maximum mortal-signature lifetime. Record the evidence +reference in `operator_attestation`. Unknown mortality or unpublished immortal signatures +preclude recovery. Never infer historical mortality from today's defaults. The archive +is a trusted RPC source, not a light-client finality proof; operator attestations cannot +be proven from on-chain data. This trust assumption is mandatory and residual. + +The importer scans both signers through a finalized anchor beyond that lifetime, +compares every signed call/hash/nonce and its successful events, then reproduces stored +accounting. Any unrelated signer activity, archive gap, missing receipt, unresolved +reservation, journal or quarantine refuses import. Only pristine jobs and completely +explained Funded/Swept prefixes (including Failed at those stages) are supported. +Conservative replay may reject valid histories with skipped consolidation or dust steps. +It never reconstructs missing receipts or clears orphan reservations. + +Apply holds an SQLite write transaction through a fresh archive scan, checks the version, +and commits policy/identity plus the full manifest and finalized anchor in +`legacy_recoveries` atomically. Cancellation/error rolls back; repeated import refuses. +State, receipts and sealed keys remain intact; failed jobs still require explicit retry. +Missing evidence means refusal, not automatic policy assignment, release or retry. SQLite reserves each signer exclusively before preparing a transaction. Reservation diff --git a/src/bin/buyback.rs b/src/bin/buyback.rs index 56c73e0..3e14881 100644 --- a/src/bin/buyback.rs +++ b/src/bin/buyback.rs @@ -115,6 +115,19 @@ enum Cmd { #[arg(long, env = "BUYBACK_NETWORK", default_value = "local")] network: String, }, + /// Verify a legacy manifest against a trusted archive; no signing keys are loaded. + #[cfg(feature = "sqlite")] + RecoverLegacy { + #[arg(long)] + manifest: std::path::PathBuf, + #[arg(long)] + database: std::path::PathBuf, + #[arg(long)] + archive: String, + /// Commit policy and audit evidence after verification. Default is read-only dry run. + #[arg(long)] + apply: bool, + }, /// Create a payment request. Create { #[command(flatten)] @@ -279,6 +292,41 @@ async fn main() -> R<()> { println!("{line}"); } } + #[cfg(feature = "sqlite")] + Cmd::RecoverLegacy { + manifest, + database, + archive, + apply, + } => { + use bittensor_buyback::Store; + use std::io::Read; + let mut bytes = Vec::new(); + std::fs::File::open(manifest)? + .take(4 * 1024 * 1024 + 1) + .read_to_end(&mut bytes)?; + if bytes.len() > 4 * 1024 * 1024 { + return Err("manifest too large".into()); + } + let manifest: bittensor_buyback::recovery::LegacyManifest = + serde_json::from_slice(&bytes)?; + let mut store = if apply { + bittensor_buyback::SqliteStore::open(database)? + } else { + bittensor_buyback::SqliteStore::open_read_only(database)? + }; + let record = store + .get(&manifest.job) + .await? + .ok_or("legacy job not found")?; + let chain = Chain::connect(&archive).await?; + let report = if apply { + store.recover_legacy(&chain, &manifest).await? + } else { + bittensor_buyback::recovery::dry_run(&chain, &record, &manifest).await? + }; + print(&report)?; + } Cmd::Create { c, metadata, diff --git a/src/chain.rs b/src/chain.rs index a272e5a..a2360f7 100644 --- a/src/chain.rs +++ b/src/chain.rs @@ -300,6 +300,9 @@ impl Endpoint { } impl Chain { + pub fn genesis_hash(&self) -> String { + format!("{:?}", self.api.genesis_hash()) + } pub async fn connect(url: &str) -> Result { Self::connect_endpoint(&Endpoint::new(url)).await } @@ -833,6 +836,98 @@ impl Chain { }) } + /// Exhaustive trusted-archive scan for explicit legacy recovery. Never signs. + pub(crate) async fn verify_legacy_archive( + &self, + record: &crate::state::PaymentRecord, + manifest: &crate::recovery::LegacyManifest, + ) -> Result<(Vec, u64, String)> { + let bad = || Error::Chain("incomplete or inconsistent legacy archive evidence".into()); + if format!("{:?}", self.api.genesis_hash()) != manifest.genesis_hash { + return Err(bad()); + } + let (head, hash) = self.finalized_head().await?; + let horizon = manifest + .stopped_at + .checked_add(manifest.maximum_mortality) + .and_then(|n| n.checked_add(1)) + .ok_or_else(bad)?; + // Bound work, never truncate a requested scan or infer absence from a subset. + if head <= horizon || head.checked_sub(manifest.first_block).ok_or_else(bad)? > 100_000 { + return Err(bad()); + } + let treasury = crate::keys::parse_ss58(&manifest.treasury)?; + let deposit = crate::keys::parse_ss58(&record.address)?; + let start = self + .api + .at_block(manifest.first_block - 1) + .await + .map_err(chain_err)?; + let mut next_nonces = std::collections::BTreeMap::new(); + for who in [treasury, deposit] { + let addr = dynamic::storage::<(AccountId32,), AccountInfo>("System", "Account"); + let nonce = match start + .storage() + .try_fetch(addr, (who,)) + .await + .map_err(chain_err)? + { + Some(v) => u64::from(v.decode().map_err(chain_err)?.nonce), + None => 0, + }; + next_nonces.insert(who, nonce); + } + let mut summaries = Vec::new(); + for number in manifest.first_block..=head { + let at = self.api.at_block(number).await.map_err(chain_err)?; + let exts = at.extrinsics().fetch().await.map_err(chain_err)?; + for ext in exts.iter() { + let ext = ext.map_err(chain_err)?; + let Some(address) = ext.address_bytes() else { + continue; + }; + // Unknown address encodings cannot be silently omitted from a completeness claim. + let signer = signer_of(address).ok_or_else(bad)?; + if signer != treasury && signer != deposit { + continue; + } + let expected = manifest.transactions.get(summaries.len()).ok_or_else(|| { + Error::Chain(format!("unexplained signer activity at block {number}, extrinsic {}; multi-job treasury histories are unsupported", ext.index())) + })?; + let nonce = ext + .transaction_extensions() + .and_then(|e| e.nonce()) + .ok_or_else(bad)?; + if expected.block != number + || expected.index != ext.index() + || expected.signer != crate::keys::ss58(&signer) + || expected.nonce != nonce + || next_nonces.get(&signer) != Some(&nonce) + || expected.receipt.tx_hash != format!("{:?}", ext.hash()) + || expected.receipt.block_hash != format!("{:?}", at.block_hash()) + || at + .tx() + .call_data(&expected.call.payload()) + .map_err(chain_err)? + != ext.call_data_bytes() + { + return Err(bad()); + } + next_nonces.insert(signer, nonce.checked_add(1).ok_or_else(bad)?); + let events = ext.events().await.map_err(chain_err)?; + summaries.push(EventSummary::from_events(events.iter())?); + } + } + if summaries.len() != manifest.transactions.len() { + return Err(bad()); + } + // Revalidate the exact finalized anchor after the scan. + if self.block_hash_at(head).await?.as_deref() != Some(hash.as_str()) { + return Err(bad()); + } + Ok((summaries, head, hash)) + } + /// Stream of finalized block numbers. pub async fn finalized_blocks( &self, @@ -953,7 +1048,8 @@ pub fn decode_partial_fee(bytes: &[u8]) -> Result { } /// Every extrinsic this crate sends. -#[derive(Clone, Debug, PartialEq, Eq)] +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] pub enum ChainCall { /// `Balances.transfer_keep_alive(dest, value)` TransferTao { dest: AccountId32, amount: u64 }, diff --git a/src/engine.rs b/src/engine.rs index 20fcc40..93b5d69 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -161,6 +161,7 @@ impl Engine { buyback: None, buyback_budget: self.cfg.buyback_budget.clone(), auto_required: Some(self.cfg.buyback_budget.is_some()), + identity: Some(self.job_identity()), attempts: 0, next_attempt_at: 0, last_error: None, @@ -245,6 +246,7 @@ impl Engine { buyback: None, buyback_budget: self.cfg.buyback_budget.clone(), auto_required: Some(self.cfg.buyback_budget.is_some()), + identity: Some(self.job_identity()), attempts: 0, next_attempt_at: 0, last_error: None, @@ -375,12 +377,27 @@ impl Engine { Ok(n) } + fn job_identity(&self) -> crate::state::JobIdentity { + crate::state::JobIdentity { + genesis_hash: self.chain.genesis_hash(), + treasury: keys::ss58(&self.treasury_account()), + treasury_hotkey: keys::ss58(&self.cfg.treasury_hotkey), + consolidate: self.cfg.consolidate, + return_dust: self.cfg.return_dust, + } + } + /// Advance one payment by at most one chain action. Returns whether anything changed. async fn step( &self, mut r: PaymentRecord, positions: &std::collections::BTreeMap>, ) -> Result { + if r.identity.as_ref() != Some(&self.job_identity()) { + return Err(Error::Config( + "job execution identity missing or mismatched".into(), + )); + } if r.quarantined.is_some() { return Ok(false); } @@ -910,6 +927,7 @@ fn action_name(a: &Action) -> &'static str { /// Engine transport, injectable for offline state-machine tests. #[async_trait::async_trait] pub trait EngineRpc: TransactionRpc { + fn genesis_hash(&self) -> String; async fn free_balance(&self, who: &AccountId32) -> Result; async fn account(&self, who: &AccountId32) -> Result<(u64, u32)>; async fn alpha_on( @@ -938,6 +956,9 @@ pub trait EngineRpc: TransactionRpc { } #[async_trait::async_trait] impl EngineRpc for Chain { + fn genesis_hash(&self) -> String { + Chain::genesis_hash(self) + } async fn free_balance(&self, who: &AccountId32) -> Result { Chain::free_balance(self, who).await } @@ -1648,6 +1669,9 @@ mod journal_rpc_tests { } #[async_trait::async_trait] impl EngineRpc for Arc { + fn genesis_hash(&self) -> String { + "simulated-genesis".into() + } async fn free_balance(&self, _: &AccountId32) -> Result { Ok(1_000_000_000) } @@ -1758,6 +1782,44 @@ mod journal_rpc_tests { ); } + #[tokio::test] + async fn identity_is_persisted_and_mismatch_never_broadcasts() { + let treasury = PaymentWallet::generate().unwrap(); + let rpc = Arc::new(SimulatedChain { + broadcasts: AtomicUsize::new(0), + hotkey: treasury.account_id(), + owner: Some(treasury.account_id()), + result: std::sync::Mutex::new(EventSummary::default()), + }); + let store = Arc::new(SqliteStore::in_memory().unwrap()); + let cfg = Config::new("unused", 100, treasury.account_id()); + let engine = Engine::new( + rpc.clone(), + store.clone(), + MasterKey::from_bytes([42; 32]), + treasury.keypair().clone(), + cfg, + ); + let request = engine + .create_payment(CreatePayment::default()) + .await + .unwrap(); + let record = store.get(&request.id).await.unwrap().unwrap(); + assert_eq!(record.identity.as_ref(), Some(&engine.job_identity())); + for field in 0..4 { + let mut bad = record.clone(); + match field { + 0 => bad.identity = None, + 1 => bad.identity.as_mut().unwrap().genesis_hash = "other".into(), + 2 => bad.identity.as_mut().unwrap().treasury = "other".into(), + _ => bad.identity.as_mut().unwrap().treasury_hotkey = "other".into(), + } + assert!(store.update(&bad).await.is_err()); + assert!(engine.step(bad, &Default::default()).await.is_err()); + } + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 0); + } + #[tokio::test] async fn engine_ticks_complete_after_each_lost_response_and_restart() { let dir = tempfile::tempdir().unwrap(); @@ -1766,6 +1828,13 @@ mod journal_rpc_tests { let treasury = PaymentWallet::generate().unwrap(); let keyring: Keyring = MasterKey::from_bytes([42; 32]).into(); let mut rec = crate::state::test_record("engine"); + rec.identity = Some(crate::state::JobIdentity { + genesis_hash: "simulated-genesis".into(), + treasury: keys::ss58(&treasury.account_id()), + treasury_hotkey: keys::ss58(&treasury.account_id()), + consolidate: false, + return_dust: false, + }); rec.netuid = 100; rec.auto_required = Some(true); rec.address = keys::ss58(&wallet.account_id()); diff --git a/src/lib.rs b/src/lib.rs index 0f389e9..6027179 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -8,6 +8,7 @@ pub mod error; #[cfg(feature = "http")] pub mod http; pub mod keys; +pub mod recovery; pub mod state; pub mod store; pub mod units; diff --git a/src/recovery.rs b/src/recovery.rs new file mode 100644 index 0000000..e03096e --- /dev/null +++ b/src/recovery.rs @@ -0,0 +1,715 @@ +//! Explicit legacy recovery. No signing API is reachable from this module. +//! +//! The archive is a trusted, authenticated RPC source, not a light-client proof. +//! Operators must stop ALL old emitters and attest a historical mortality bound; +//! the chain cannot prove that an unpublished immortal signature does not exist. +use crate::{ + Chain, ChainCall, Error, Result, keys, + state::{BuybackBudget, PaymentRecord, PaymentState, TxRef}, +}; +use serde::{Deserialize, Serialize}; + +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyManifest { + pub job: String, + pub version: u64, + pub genesis_hash: String, + pub treasury: String, + pub treasury_hotkey: String, + pub consolidate: bool, + pub return_dust: bool, + #[serde(deserialize_with = "Deserialize::deserialize")] + pub budget: Option, + /// First block covering every old emission, including those absent from the DB. + pub first_block: u64, + /// Finalized block at which all old signing processes were stopped. + pub stopped_at: u64, + /// Historical maximum lifetime in blocks, NOT the current engine default. + pub maximum_mortality: u64, + /// Audit reference for shutdown, complete scan start, and historical mortality evidence. + pub operator_attestation: String, + /// Explicit acceptance of the off-chain assumption; not an archive verification flag. + pub all_old_emitters_stopped_no_immortal_signatures: bool, + /// Complete chronological emission list for BOTH signers over the scan interval. + pub transactions: Vec, +} + +#[derive(Clone, Debug, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct LegacyTransaction { + pub block: u64, + pub index: usize, + pub signer: String, + pub nonce: u64, + pub receipt: TxRef, + pub call: ChainCall, +} + +#[derive(Clone, Debug, Serialize)] +pub struct RecoveryReport { + pub job: String, + pub previous_version: u64, + pub genesis_hash: String, + pub finalized_block: u64, + pub finalized_hash: String, + pub transactions: usize, + pub applied: bool, +} + +pub(crate) fn validate_manifest(r: &PaymentRecord, m: &LegacyManifest) -> Result<()> { + let bad = || { + Error::Store( + "legacy recovery requires a complete audited manifest and unresolved-free record" + .into(), + ) + }; + if r.id != m.job || r.version != m.version { + return Err(Error::Conflict(m.job.clone())); + } + if r.identity.is_some() + || (r.auto_required.is_some() + && (r.auto_required != Some(m.budget.is_some()) || r.buyback_budget != m.budget)) + || r.pending.is_some() + || r.quarantined.is_some() + || r.notified + || (m.transactions.is_empty() + && !matches!(r.state, PaymentState::Pending | PaymentState::Detected)) + || m.transactions.len() > 4096 + || m.first_block == 0 + || m.first_block > m.stopped_at + || m.maximum_mortality == 0 + || m.maximum_mortality > 65536 + || !m.all_old_emitters_stopped_no_immortal_signatures + || m.operator_attestation.trim().len() < 16 + || m.operator_attestation.len() > 4096 + || !matches!( + r.state, + PaymentState::Pending + | PaymentState::Detected + | PaymentState::Funded + | PaymentState::Swept + | PaymentState::Failed + ) + || (r.state == PaymentState::Failed + && !matches!( + r.failed_from, + Some(PaymentState::Funded | PaymentState::Swept) + )) + { + return Err(bad()); + } + if matches!(r.state, PaymentState::Pending | PaymentState::Detected) + && (!m.transactions.is_empty() + || r.funded_tao != 0 + || r.swept_alpha != 0 + || !r.swept_positions.is_empty() + || r.consolidated != 0 + || r.dust_returned + || r.buyback.is_some() + || r.buyback_done) + { + return Err(bad()); + } + let treasury = keys::parse_ss58(&m.treasury)?; + let deposit = keys::parse_ss58(&r.address)?; + keys::parse_ss58(&m.treasury_hotkey)?; + if treasury == deposit + || keys::ss58(&treasury) != m.treasury + || keys::ss58(&deposit) != r.address + { + return Err(bad()); + } + if let Some(b) = &m.budget { + b.validate()?; + } + let mut last = None; + let mut hashes = std::collections::BTreeSet::new(); + let mut nonces = std::collections::BTreeSet::new(); + for t in &m.transactions { + if t.block < m.first_block + || t.block > m.stopped_at + || last.is_some_and(|p| p >= (t.block, t.index)) + || !hashes.insert(&t.receipt.tx_hash) + || !nonces.insert((&t.signer, t.nonce)) + || (t.signer != m.treasury && t.signer != r.address) + || t.receipt.action != t.call.name() + || t.receipt.amount != t.call.amount() + { + return Err(bad()); + } + last = Some((t.block, t.index)); + } + // ponytail: only complete stored receipts are importable; missing/ambiguous evidence stays blocked. + if m.transactions.iter().map(|t| &t.receipt).ne(r.txs.iter()) { + return Err(bad()); + } + Ok(()) +} + +/// Read-only archive dry run. Applying requires a fresh scan under SQLite signer fences. +pub async fn dry_run( + chain: &Chain, + record: &PaymentRecord, + manifest: &LegacyManifest, +) -> Result { + verify(chain, record, manifest) + .await + .map(|(_, report)| report) +} + +pub(crate) async fn verify( + chain: &Chain, + r: &PaymentRecord, + m: &LegacyManifest, +) -> Result<(PaymentRecord, RecoveryReport)> { + validate_manifest(r, m)?; + let summaries = chain.verify_legacy_archive(r, m).await?; + replay(r, m, summaries) +} + +fn replay( + r: &PaymentRecord, + m: &LegacyManifest, + summaries: (Vec, u64, String), +) -> Result<(PaymentRecord, RecoveryReport)> { + validate_manifest(r, m)?; + if summaries.0.len() != m.transactions.len() { + return Err(Error::Store("missing archived receipts".into())); + } + let treasury = keys::parse_ss58(&m.treasury)?; + let deposit = keys::parse_ss58(&r.address)?; + let hotkey = keys::parse_ss58(&m.treasury_hotkey)?; + let mut funded = 0u64; + let mut positions = Vec::new(); + let mut consolidated = 0usize; + let mut dust = false; + let mut post_sweep = false; + let mut buyback: Option = None; + for (t, s) in m.transactions.iter().zip(&summaries.0) { + let signer = keys::parse_ss58(&t.signer)?; + let invalid = + || Error::Store("legacy call, events or stage disagrees with audited job".into()); + let (p, e) = t.call.expected_event(); + if s.failed || !s.has("System", "ExtrinsicSuccess") || !s.has(p, e) { + return Err(invalid()); + } + match &t.call { + ChainCall::TransferTao { dest, amount } + if signer == treasury + && *dest == deposit + && funded == 0 + && positions.is_empty() + && !dust + && buyback.is_none() + && *amount > 0 => + { + if s.transferred != Some(*amount) { + return Err(invalid()); + } + funded = *amount; + } + ChainCall::TransferStake { + dest_coldkey, + hotkey: source, + netuid, + alpha, + } if signer == deposit + && *dest_coldkey == treasury + && *netuid == r.netuid + && !post_sweep + && consolidated == 0 + && !dust + && buyback.is_none() + && *alpha > 0 => + { + let source = keys::ss58(source); + if positions.iter().any(|(h, _)| h == &source) { + return Err(invalid()); + } + positions.push((source, *alpha)); + } + ChainCall::MoveStake { + from_hotkey, + to_hotkey, + netuid, + alpha, + } if m.consolidate + && signer == treasury + && *to_hotkey == hotkey + && *netuid == r.netuid + && !dust + && buyback.is_none() => + { + if positions.get(consolidated) != Some(&(keys::ss58(from_hotkey), *alpha)) { + return Err(invalid()); + } + post_sweep = true; + consolidated += 1; + } + ChainCall::TransferAll { dest } + if (m.return_dust || r.detected_tao > 0) + && signer == deposit + && *dest == treasury + && !dust + && buyback.is_none() => + { + if m.consolidate && consolidated != positions.len() { + return Err(invalid()); + } + post_sweep = true; + dust = true; + } + ChainCall::AddStakeLimit { + hotkey: target, + netuid, + tao, + limit_price, + allow_partial: false, + } if signer == treasury && buyback.is_none() => { + if (m.consolidate && consolidated != positions.len()) + || ((m.return_dust || r.detected_tao > 0) && !dust) + { + return Err(invalid()); + } + post_sweep = true; + let b = m.budget.as_ref().ok_or_else(invalid)?; + let (spent, bought) = s.stake_added.ok_or_else(invalid)?; + if keys::ss58(target) != b.hotkey + || *netuid != b.netuid + || *tao != b.amount_rao + || spent != *tao + || bought == 0 + { + return Err(invalid()); + } + buyback = Some(crate::state::BuybackReceipt { + netuid: *netuid, + tao_spent: spent, + alpha_bought: bought, + limit_price: *limit_price, + stake_tx: t.receipt.clone(), + destroy_tx: None, + alpha_destroyed: 0, + }); + } + ChainCall::BurnAlpha { + hotkey: target, + netuid, + alpha, + } + | ChainCall::RecycleAlpha { + hotkey: target, + netuid, + alpha, + } if signer == treasury => { + let b = m.budget.as_ref().ok_or_else(invalid)?; + let receipt = buyback.as_mut().ok_or_else(invalid)?; + let destroy = if matches!(t.call, ChainCall::BurnAlpha { .. }) { + crate::Destroy::Burn + } else { + crate::Destroy::Recycle + }; + if b.destroy != destroy + || keys::ss58(target) != b.hotkey + || *netuid != b.netuid + || receipt.destroy_tx.is_some() + || *alpha != receipt.alpha_bought + || s.alpha_destroyed != Some(*alpha) + { + return Err(invalid()); + } + receipt.destroy_tx = Some(t.receipt.clone()); + receipt.alpha_destroyed = *alpha; + } + _ => return Err(invalid()), + } + } + let total = positions + .iter() + .try_fold(0u64, |a, (_, b)| a.checked_add(*b)) + .ok_or_else(|| Error::Store("legacy amount overflow".into()))?; + let done = buyback.as_ref().is_some_and(|b| { + b.destroy_tx.is_some() + || m.budget + .as_ref() + .is_some_and(|p| p.destroy == crate::Destroy::Keep) + }); + let stage = r + .failed_from + .filter(|_| r.state == PaymentState::Failed) + .unwrap_or(r.state); + if (stage == PaymentState::Funded + && (consolidated != 0 || dust || buyback.is_some() || r.buyback_done)) + || funded != r.funded_tao + || positions != r.swept_positions + || total != r.swept_alpha + || consolidated != r.consolidated + || dust != r.dust_returned + || buyback != r.buyback + || (r.buyback_done && !done) + { + return Err(Error::Store( + "legacy receipts do not reproduce stored accounting".into(), + )); + } + let mut next = r.clone(); + next.identity = Some(crate::state::JobIdentity { + genesis_hash: m.genesis_hash.clone(), + treasury: m.treasury.clone(), + treasury_hotkey: m.treasury_hotkey.clone(), + consolidate: m.consolidate, + return_dust: m.return_dust, + }); + next.auto_required = Some(m.budget.is_some()); + next.buyback_budget = m.budget.clone(); + // State, receipts, attempts and sealed wallet are preserved; retry remains an explicit operation. + let report = RecoveryReport { + job: r.id.clone(), + previous_version: r.version, + genesis_hash: m.genesis_hash.clone(), + finalized_block: summaries.1, + finalized_hash: summaries.2, + transactions: m.transactions.len(), + applied: false, + }; + Ok((next, report)) +} + +#[cfg(test)] +mod tests { + use super::*; + fn fixture() -> (PaymentRecord, LegacyManifest, crate::chain::EventSummary) { + let mut r = crate::state::test_record("legacy"); + let treasury = keys::keypair_from_uri("//Bob") + .unwrap() + .public_key() + .to_account_id(); + r.auto_required = None; + r.state = PaymentState::Funded; + r.funded_tao = 123; + let call = ChainCall::TransferTao { + dest: keys::parse_ss58(&r.address).unwrap(), + amount: 123, + }; + let receipt = TxRef { + action: call.name().into(), + tx_hash: "hash".into(), + block_hash: "block".into(), + amount: 123, + }; + r.txs.push(receipt.clone()); + let m = LegacyManifest { + job: r.id.clone(), + version: r.version, + genesis_hash: "test-genesis".into(), + treasury: keys::ss58(&treasury), + treasury_hotkey: keys::ss58(&treasury), + consolidate: false, + return_dust: false, + budget: None, + first_block: 1, + stopped_at: 2, + maximum_mortality: 128, + operator_attestation: "offline shutdown and historical mortality evidence".into(), + all_old_emitters_stopped_no_immortal_signatures: true, + transactions: vec![LegacyTransaction { + block: 2, + index: 0, + signer: keys::ss58(&treasury), + nonce: 0, + receipt, + call, + }], + }; + let s = crate::chain::EventSummary { + names: vec![ + ("System".into(), "ExtrinsicSuccess".into()), + ("Balances".into(), "Transfer".into()), + ], + transferred: Some(123), + ..Default::default() + }; + (r, m, s) + } + #[test] + fn pristine_identity_requires_explicit_history_and_configuration() { + let (mut r, mut m, _) = fixture(); + r.state = PaymentState::Detected; + r.funded_tao = 0; + r.txs.clear(); + m.transactions.clear(); + let (next, _) = replay(&r, &m, (vec![], 200, "anchor".into())).unwrap(); + assert_eq!(next.identity.unwrap().genesis_hash, "test-genesis"); + let mut json = serde_json::to_value(&m).unwrap(); + json.as_object_mut().unwrap().remove("consolidate"); + assert!(serde_json::from_value::(json).is_err()); + m.operator_attestation.clear(); + assert!(validate_manifest(&r, &m).is_err()); + } + + #[test] + fn replay_preserves_evidence_and_rejects_missing_proof() { + let (r, m, s) = fixture(); + let before = serde_json::to_value(&r).unwrap(); + let (next, report) = replay(&r, &m, (vec![s.clone()], 200, "anchor".into())).unwrap(); + assert_eq!(serde_json::to_value(&r).unwrap(), before); + assert_eq!(next.txs, r.txs); + assert_eq!(next.version, r.version); + assert_eq!(next.state, r.state); + assert_eq!(next.identity.unwrap().genesis_hash, m.genesis_hash); + assert!(!report.applied); + assert!(replay(&r, &m, (vec![], 200, "anchor".into())).is_err()); + let mut invalid = s; + invalid.transferred = Some(122); + assert!(replay(&r, &m, (vec![invalid], 200, "anchor".into())).is_err()); + for case in 0..6 { + let mut bad = m.clone(); + match case { + 0 => bad.maximum_mortality = 0, + 1 => bad.all_old_emitters_stopped_no_immortal_signatures = false, + 2 => bad.operator_attestation.clear(), + 3 => bad.transactions.clear(), + 4 => bad.transactions[0].receipt.amount += 1, + _ => bad.version += 1, + } + assert!(validate_manifest(&r, &bad).is_err()); + } + let mut json = serde_json::to_value(&m).unwrap(); + json["verified"] = true.into(); + assert!(serde_json::from_value::(json).is_err()); + } +} + +#[cfg(all(test, feature = "localnet-tests", feature = "sqlite"))] +mod localnet_recovery { + use super::*; + use crate::Store; + #[tokio::test] + async fn archive_dry_run_apply_and_repeated_import() { + // Local endpoint only; never accept an environment override for this signing fixture. + let chain = Chain::connect("ws://127.0.0.1:9944").await.unwrap(); + let alice = keys::keypair_from_uri("//Alice").unwrap(); + let treasury = crate::PaymentWallet::generate().unwrap(); + let deposit = crate::PaymentWallet::generate().unwrap(); + chain + .submit( + &ChainCall::TransferTao { + dest: treasury.account_id(), + amount: 1_000_000_000, + }, + &alice, + ) + .await + .unwrap(); + let mut r = crate::state::test_record("archive-recovery"); + r.address = keys::ss58(&deposit.account_id()); + let keyring: crate::keys::Keyring = crate::MasterKey::from_bytes([42; 32]).into(); + r.sealed_secret = deposit + .seal_with(&keyring, &keys::wallet_aad(&r.id, &r.address)) + .unwrap(); + r.auto_required = None; + r.state = PaymentState::Funded; + r.funded_tao = 10_000_000; + let call = ChainCall::TransferTao { + dest: deposit.account_id(), + amount: r.funded_tao, + }; + let prepared = chain.prepare(&call, treasury.keypair()).await.unwrap(); + let finalized = chain.broadcast(&prepared).await.unwrap(); + r.txs.push(finalized.tx.clone()); + let (stopped, _) = chain.finalized_head().await.unwrap(); + let mut location = None; + for block in prepared.birth_block..=stopped { + let at = chain.api().at_block(block).await.unwrap(); + let exts = at.extrinsics().fetch().await.unwrap(); + for ext in exts.iter() { + let ext = ext.unwrap(); + if format!("{:?}", ext.hash()) == prepared.tx_hash { + location = Some((block, ext.index())); + } + } + } + let (block, index) = location.unwrap(); + let m = LegacyManifest { + job: r.id.clone(), + version: r.version, + genesis_hash: chain.genesis_hash(), + treasury: keys::ss58(&treasury.account_id()), + treasury_hotkey: keys::ss58(&treasury.account_id()), + consolidate: false, + return_dust: false, + budget: None, + first_block: prepared.birth_block, + stopped_at: stopped, + maximum_mortality: 64, + operator_attestation: + "isolated random localnet signer; only prepare uses encoded 64-block mortality" + .into(), + all_old_emitters_stopped_no_immortal_signatures: true, + transactions: vec![LegacyTransaction { + block, + index, + signer: prepared.signer, + nonce: prepared.nonce, + receipt: finalized.tx, + call, + }], + }; + // No further emission from either random signer. Bound waiting rather than weakening evidence. + tokio::time::timeout(std::time::Duration::from_secs(600), async { + while chain.finalized_head().await.unwrap().0 <= stopped + 65 { + tokio::time::sleep(std::time::Duration::from_millis(500)).await; + } + }) + .await + .unwrap(); + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("recovery.db"); + let mut store = crate::SqliteStore::open(&path).unwrap(); + store.insert(&r).await.unwrap(); + let before = serde_json::to_value(store.get(&r.id).await.unwrap()).unwrap(); + let read_only = crate::SqliteStore::open_read_only(&path).unwrap(); + assert!(!dry_run(&chain, &r, &m).await.unwrap().applied); + assert_eq!( + serde_json::to_value(read_only.get(&r.id).await.unwrap()).unwrap(), + before + ); + let mut wrong = m.clone(); + wrong.genesis_hash = "wrong".into(); + assert!(store.recover_legacy(&chain, &wrong).await.is_err()); + assert_eq!( + serde_json::to_value(store.get(&r.id).await.unwrap()).unwrap(), + before + ); + let peer = crate::SqliteStore::open(&path).unwrap(); + let token = peer.reserve_signer(&m.treasury, "other-job").await.unwrap(); + assert!(store.recover_legacy(&chain, &m).await.is_err()); + peer.cancel_preparation(&m.treasury, "other-job", &token) + .await + .unwrap(); + // A canceled import drops the uncommitted SQLite transaction; no policy is installed. + let mut import = Box::pin(store.recover_legacy(&chain, &m)); + assert!(futures::poll!(import.as_mut()).is_pending()); + drop(import); + assert_eq!( + serde_json::to_value(peer.get(&r.id).await.unwrap()).unwrap(), + before + ); + assert!(store.recover_legacy(&chain, &m).await.unwrap().applied); + assert!(store.recover_legacy(&chain, &m).await.is_err()); + assert!(store.update(&r).await.is_err()); + drop(store); + let record = crate::SqliteStore::open(&path) + .unwrap() + .get(&r.id) + .await + .unwrap() + .unwrap(); + assert_eq!(record.version, r.version + 1); + assert_eq!(record.txs, r.txs); + assert_eq!( + record.identity.as_ref().unwrap().genesis_hash, + m.genesis_hash + ); + let nonce_before = chain.account(&treasury.account_id()).await.unwrap().1; + let mut cfg = + crate::Config::new("ws://127.0.0.1:9944", record.netuid, treasury.account_id()); + cfg.consolidate = false; + cfg.return_dust = false; + let store = std::sync::Arc::new(crate::SqliteStore::open(&path).unwrap()); + let engine = crate::Engine::new( + chain, + store.clone(), + crate::MasterKey::from_bytes([42; 32]), + treasury.keypair().clone(), + cfg, + ); + for _ in 0..3 { + engine.tick().await.unwrap(); + } + assert_eq!( + store.get(&r.id).await.unwrap().unwrap().state, + PaymentState::Settled + ); + assert_eq!( + engine + .chain() + .account(&treasury.account_id()) + .await + .unwrap() + .1, + nonce_before + ); + } +} + +#[cfg(all(test, feature = "sqlite", unix))] +mod crash_tests { + use crate::Store; + #[test] + fn recovery_transaction_child() { + let Ok(path) = std::env::var("BUYBACK_RECOVERY_CRASH_DB") else { + return; + }; + let c = rusqlite::Connection::open(&path).unwrap(); + c.execute_batch("BEGIN IMMEDIATE; UPDATE payments SET version=99; CREATE TABLE legacy_recoveries(job TEXT); INSERT INTO legacy_recoveries VALUES ('uncommitted');").unwrap(); + std::fs::write(format!("{path}.ready"), b"ready").unwrap(); + std::thread::sleep(std::time::Duration::from_secs(60)); + panic!("parent must kill the process before commit"); + } + #[tokio::test] + async fn killed_writer_preserves_record_and_signer_ownership() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("crash.db"); + let store = crate::SqliteStore::open(&path).unwrap(); + let r = crate::state::test_record("crash"); + store.insert(&r).await.unwrap(); + let token = store.reserve_signer("signer", &r.id).await.unwrap(); + drop(store); + let mut child = std::process::Command::new(std::env::current_exe().unwrap()) + .args([ + "--exact", + "recovery::crash_tests::recovery_transaction_child", + "--nocapture", + ]) + .env("BUYBACK_RECOVERY_CRASH_DB", &path) + .stdout(std::process::Stdio::null()) + .spawn() + .unwrap(); + let ready = format!("{}.ready", path.display()); + for _ in 0..200 { + if std::path::Path::new(&ready).exists() { + break; + } + std::thread::sleep(std::time::Duration::from_millis(10)); + } + let reached = std::path::Path::new(&ready).exists(); + child.kill().unwrap(); + child.wait().unwrap(); + assert!(reached, "child must reach uncommitted transaction"); + let reopened = crate::SqliteStore::open(&path).unwrap(); + assert_eq!( + serde_json::to_value(reopened.get(&r.id).await.unwrap().unwrap()).unwrap(), + serde_json::to_value(&r).unwrap() + ); + assert!(reopened.reserve_signer("signer", "other").await.is_err()); + let c = rusqlite::Connection::open(&path).unwrap(); + let persisted: String = c + .query_row( + "SELECT token FROM signer_reservations WHERE signer='signer'", + [], + |r| r.get(0), + ) + .unwrap(); + assert_eq!(persisted, token); + let audit: bool = c + .query_row( + "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE name='legacy_recoveries')", + [], + |r| r.get(0), + ) + .unwrap(); + assert!(!audit); + } +} diff --git a/src/state.rs b/src/state.rs index a416cba..69c007d 100644 --- a/src/state.rs +++ b/src/state.rs @@ -115,6 +115,17 @@ impl BuybackBudget { } } +/// Immutable execution identity; absent on legacy records, never inferred during resume. +#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct JobIdentity { + pub genesis_hash: String, + pub treasury: String, + pub treasury_hotkey: String, + pub consolidate: bool, + pub return_dust: bool, +} + /// Full persisted record. Contains the sealed wallet secret: never return it to clients, use /// [`PaymentStatus`]. #[derive(Clone, Debug, Serialize, Deserialize)] @@ -143,6 +154,8 @@ pub struct PaymentRecord { /// None identifies legacy jobs requiring explicit migration. #[serde(default)] pub auto_required: Option, + #[serde(default)] + pub identity: Option, pub attempts: u32, pub next_attempt_at: u64, pub last_error: Option, @@ -298,6 +311,7 @@ pub(crate) fn test_record(id: &str) -> PaymentRecord { buyback: None, buyback_budget: None, auto_required: Some(false), + identity: None, attempts: 0, next_attempt_at: 0, last_error: None, diff --git a/src/store.rs b/src/store.rs index 773af95..afc46fb 100644 --- a/src/store.rs +++ b/src/store.rs @@ -210,7 +210,10 @@ impl Store for FileStore { let cur = self .read(&p)? .ok_or_else(|| Error::NotFound(rec.id.clone()))?; - if cur.buyback_budget != rec.buyback_budget || cur.auto_required != rec.auto_required { + if cur.buyback_budget != rec.buyback_budget + || cur.auto_required != rec.auto_required + || cur.identity != rec.identity + { return Err(Error::Store("immutable job budget".into())); } if cur.version != rec.version { @@ -237,6 +240,16 @@ impl SqliteStore { Self::init(conn) } + /// Open an existing database without schema or pragma writes. + pub fn open_read_only(path: impl AsRef) -> Result { + let conn = + rusqlite::Connection::open_with_flags(path, rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY) + .map_err(store_err)?; + Ok(Self { + conn: Mutex::new(conn), + }) + } + pub fn in_memory() -> Result { Self::init(rusqlite::Connection::open_in_memory().map_err(store_err)?) } @@ -256,6 +269,64 @@ impl SqliteStore { conn: Mutex::new(conn), }) } + /// Apply a freshly verified legacy manifest under an exclusive SQLite transaction. + /// The database write lock fences all new reservations while the archive is read. + /// Existing reservations/journals are never deleted or adopted by this operation. + pub async fn recover_legacy( + &mut self, + chain: &crate::Chain, + manifest: &crate::recovery::LegacyManifest, + ) -> Result { + // ponytail: offline exclusive maintenance; SQLite fences other processes throughout the scan. + let c = self.conn.get_mut().map_err(store_err)?; + let tx = c + .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate) + .map_err(store_err)?; + let json: String = tx + .query_row( + "SELECT record FROM payments WHERE id=?1", + [&manifest.job], + |r| r.get(0), + ) + .map_err(store_err)?; + let record: PaymentRecord = serde_json::from_str(&json).map_err(store_err)?; + crate::recovery::validate_manifest(&record, manifest)?; + let busy: bool = tx.query_row( + "SELECT EXISTS(SELECT 1 FROM signer_reservations WHERE signer IN (?1,?2) OR job=?3) OR EXISTS(SELECT 1 FROM payments WHERE json_extract(record,'$.pending.signer') IN (?1,?2) OR json_extract(record,'$.quarantined.signer') IN (?1,?2))", + rusqlite::params![record.address,manifest.treasury,record.id],|r|r.get(0), + ).map_err(store_err)?; + if busy { + return Err(Error::Conflict( + "legacy recovery signers have unresolved ownership".into(), + )); + } + let (mut next, mut report) = crate::recovery::verify(chain, &record, manifest).await?; + next.version = next + .version + .checked_add(1) + .ok_or_else(|| store_err("version overflow"))?; + next.updated_at = crate::state::now(); + tx.execute_batch("CREATE TABLE IF NOT EXISTS legacy_recoveries (job TEXT NOT NULL, previous_version INTEGER NOT NULL, manifest TEXT NOT NULL, report TEXT NOT NULL, PRIMARY KEY(job,previous_version));").map_err(store_err)?; + report.applied = true; + let changed = tx + .execute( + "UPDATE payments SET version=?1,record=?2 WHERE id=?3 AND version=?4", + rusqlite::params![ + i64::try_from(next.version).map_err(store_err)?, + serde_json::to_string(&next).map_err(store_err)?, + next.id, + i64::try_from(record.version).map_err(store_err)? + ], + ) + .map_err(store_err)?; + if changed != 1 { + return Err(Error::Conflict(record.id)); + } + tx.execute("INSERT INTO legacy_recoveries(job,previous_version,manifest,report) VALUES (?1,?2,?3,?4)",rusqlite::params![record.id,i64::try_from(record.version).map_err(store_err)?,serde_json::to_string(manifest).map_err(store_err)?,serde_json::to_string(&report).map_err(store_err)?]).map_err(store_err)?; + tx.commit().map_err(store_err)?; + Ok(report) + } + fn reserve(&self, signer: &str, job: &str, version: Option) -> Result { let mut c = self.conn.lock().map_err(store_err)?; let tx = c @@ -437,6 +508,7 @@ impl Store for SqliteStore { } if previous.buyback_budget != rec.buyback_budget || previous.auto_required != rec.auto_required + || previous.identity != rec.identity { return Err(Error::Store("immutable job budget".into())); } From c97e40d84018dc31928edfcfe60e7670a41d0206 Mon Sep 17 00:00:00 2001 From: Mathis <154886644+echobt@users.noreply.github.com> Date: Fri, 25 Sep 2026 21:59:04 +0000 Subject: [PATCH 8/8] Use persisted job policy for automatic buyback decisions --- src/engine.rs | 61 ++++++++++++++++++++++++++++++++++++++++++++++++--- 1 file changed, 58 insertions(+), 3 deletions(-) diff --git a/src/engine.rs b/src/engine.rs index 93b5d69..5e514a9 100644 --- a/src/engine.rs +++ b/src/engine.rs @@ -609,9 +609,7 @@ impl Engine { if required && r.buyback_done { validate_completed_buyback(r)?; } - if !r.buyback_done - && (r.buyback_budget.is_some() || matches!(self.cfg.auto, AutoBuyback::On { .. })) - { + if required && !r.buyback_done { return self.step_auto_buyback(r).await; } r.transition(PaymentState::Settled)?; @@ -1820,6 +1818,63 @@ mod journal_rpc_tests { assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 0); } + #[tokio::test] + async fn persisted_auto_off_settles_after_runtime_enables_buyback() { + let treasury = PaymentWallet::generate().unwrap(); + let rpc = Arc::new(SimulatedChain { + broadcasts: AtomicUsize::new(0), + hotkey: treasury.account_id(), + owner: Some(treasury.account_id()), + result: std::sync::Mutex::new(EventSummary::default()), + }); + let store = Arc::new(SqliteStore::in_memory().unwrap()); + let mut cfg = Config::new("unused", 100, treasury.account_id()); + cfg.auto = AutoBuyback::Off; + cfg.consolidate = false; + cfg.return_dust = false; + let engine = Engine::new( + rpc.clone(), + store.clone(), + MasterKey::from_bytes([42; 32]), + treasury.keypair().clone(), + cfg.clone(), + ); + let request = engine + .create_payment(CreatePayment::default()) + .await + .unwrap(); + let mut record = store.get(&request.id).await.unwrap().unwrap(); + assert_eq!(record.auto_required, Some(false)); + assert!(record.buyback_budget.is_none()); + for state in [ + PaymentState::Detected, + PaymentState::Funded, + PaymentState::Swept, + ] { + record.transition(state).unwrap(); + record = store.update(&record).await.unwrap(); + } + drop(engine); + cfg.auto = AutoBuyback::On { + destroy: Destroy::Burn, + amount: crate::AutoAmount::Fixed(MIN_STAKE_RAO), + }; + let engine = Engine::new( + rpc.clone(), + store.clone(), + MasterKey::from_bytes([42; 32]), + treasury.keypair().clone(), + cfg, + ); + assert_eq!(engine.tick().await.unwrap(), 1); + let settled = store.get(&request.id).await.unwrap().unwrap(); + assert_eq!(settled.state, PaymentState::Settled); + assert!(settled.notified); + assert!(settled.buyback.is_none()); + assert_eq!(settled.auto_required, Some(false)); + assert_eq!(rpc.broadcasts.load(Ordering::SeqCst), 0); + } + #[tokio::test] async fn engine_ticks_complete_after_each_lost_response_and_restart() { let dir = tempfile::tempdir().unwrap();