diff --git a/api/src/handlers.rs b/api/src/handlers.rs index b18b0cc..536b978 100644 --- a/api/src/handlers.rs +++ b/api/src/handlers.rs @@ -11,9 +11,10 @@ use setu_rpc::{ HeartbeatRequest, HeartbeatResponse, RegisterSolverRequest, RegisterSolverResponse, RegisterValidatorRequest, RegisterValidatorResponse, RegistrationHandler, SubmitTransferRequest, SubmitTransferResponse, + SubmitTaskRequest, SubmitTaskResponse, // User RPC imports UserRpcHandler, RegisterUserRequest, RegisterUserResponse, - GetAccountRequest, GetAccountResponse, GetBalanceRequest, + GetAccountRequest, GetAccountResponse, GetBalanceRequest, GetBalanceResponse as UserGetBalanceResponse, GetPowerRequest, GetPowerResponse, GetCreditRequest, GetCreditResponse, GetCredentialsRequest, GetCredentialsResponse, TransferRequest, TransferResponse, @@ -48,10 +49,13 @@ pub trait ValidatorService: Send + Sync { /// Submit transfer fn submit_transfer(&self, request: SubmitTransferRequest) -> impl std::future::Future + Send; - + + /// Submit task for solver execution + fn submit_task(&self, request: SubmitTaskRequest) -> impl std::future::Future + Send; + /// Get transfer status fn get_transfer_status(&self, transfer_id: &str) -> GetTransferStatusResponse; - + /// Submit event fn submit_event(&self, request: SubmitEventRequest) -> impl std::future::Future + Send; @@ -132,6 +136,14 @@ pub async fn http_submit_transfer( Json(service.submit_transfer(request).await) } +/// Submit a task for solver execution +pub async fn http_submit_task( + State(service): State>, + Json(request): Json, +) -> Json { + Json(service.submit_task(request).await) +} + /// Get transfer status pub async fn http_get_transfer_status( State(service): State>, diff --git a/consensus/src/engine.rs b/consensus/src/engine.rs index f2ade43..219d6e0 100644 --- a/consensus/src/engine.rs +++ b/consensus/src/engine.rs @@ -765,7 +765,7 @@ impl ConsensusEngine { // Check if our vote caused finalization // (vote_for_cf adds vote but doesn't check finalization, so we check here) - let finalized = manager.check_finalization(&cf_id); + let finalized = manager.check_finalization(&cf_id, &dag); if finalized { return self.handle_finalization(&mut manager).await; } @@ -881,8 +881,9 @@ impl ConsensusEngine { ); } + let dag = self.dag.read().await; let mut manager = self.consensus_manager.write().await; - let finalized = manager.receive_vote(vote); + let finalized = manager.receive_vote(vote, &dag); if finalized { self.handle_finalization(&mut manager).await @@ -1348,9 +1349,10 @@ mod tests { let vote2 = Vote::new("v2".to_string(), cf.id.clone(), true); let vote3 = Vote::new("v3".to_string(), cf.id.clone(), true); - manager.receive_vote(vote1); - manager.receive_vote(vote2); - let finalized = manager.receive_vote(vote3); + let dag = crate::dag::Dag::new(); + manager.receive_vote(vote1, &dag); + manager.receive_vote(vote2, &dag); + let finalized = manager.receive_vote(vote3, &dag); assert!(finalized, "CF should be finalized after quorum"); @@ -1418,20 +1420,21 @@ mod tests { let mut manager = engine.consensus_manager.write().await; manager.receive_cf(cf); - + let dag = crate::dag::Dag::new(); + // Vote 1: approve (should not finalize yet) let vote1 = Vote::new("v1".to_string(), cf_id.clone(), true); - let result1 = manager.receive_vote(vote1); + let result1 = manager.receive_vote(vote1, &dag); assert!(!result1, "Should not finalize with 1 approve vote"); - + // Vote 2: reject (1 reject, not enough) let vote2 = Vote::new("v2".to_string(), cf_id.clone(), false); - let result2 = manager.receive_vote(vote2); + let result2 = manager.receive_vote(vote2, &dag); assert!(!result2, "Should not reject with only 1 reject vote"); - + // Vote 3: reject (2 rejects = 1/3+1, should reject) let vote3 = Vote::new("v3".to_string(), cf_id.clone(), false); - let result3 = manager.receive_vote(vote3); + let result3 = manager.receive_vote(vote3, &dag); assert!(result3, "Should reject with 2 reject votes (1/3+1 threshold)"); // Verify CF was removed from pending (can't directly access private field) @@ -1476,10 +1479,11 @@ mod tests { { let mut manager = engine.consensus_manager.write().await; manager.receive_cf(cf); - + // Verify CF is pending (test by attempting to receive vote) + let dag = crate::dag::Dag::new(); let vote = Vote::new("v1".to_string(), cf_id.clone(), true); - manager.receive_vote(vote); + manager.receive_vote(vote, &dag); } // Wait for timeout @@ -1535,9 +1539,10 @@ mod tests { tokio::time::sleep(tokio::time::Duration::from_millis(150)).await; // Receive a vote - should trigger timeout check and remove CF + let dag = crate::dag::Dag::new(); let vote = Vote::new("v1".to_string(), cf_id.clone(), true); - let result = manager.receive_vote(vote); - + let result = manager.receive_vote(vote, &dag); + // The vote processing should detect timeout and remove CF assert!(result, "Should return true when CF is removed due to timeout"); } @@ -1589,10 +1594,11 @@ mod tests { manager.receive_cf(cf1); // Reject the CF with enough reject votes + let dag = crate::dag::Dag::new(); let vote1 = Vote::new("v2".to_string(), cf1_id.clone(), false); let vote2 = Vote::new("v3".to_string(), cf1_id.clone(), false); - manager.receive_vote(vote1); - let rejected = manager.receive_vote(vote2); + manager.receive_vote(vote1, &dag); + let rejected = manager.receive_vote(vote2, &dag); assert!(rejected, "CF should be rejected with 2 reject votes"); diff --git a/consensus/src/folder.rs b/consensus/src/folder.rs index 4532881..216d747 100644 --- a/consensus/src/folder.rs +++ b/consensus/src/folder.rs @@ -261,10 +261,10 @@ impl ConsensusManager { /// /// Returns true if the CF is finalized after this vote. /// Duplicate votes from the same validator are ignored (idempotent). - pub fn receive_vote(&mut self, vote: Vote) -> bool { + pub fn receive_vote(&mut self, vote: Vote, dag: &Dag) -> bool { let cf_id = vote.cf_id.clone(); let voter_id = vote.validator_id.clone(); - + if let Some(cf) = self.pending_cfs.get_mut(&cf_id) { // Skip if this validator already voted (idempotency) if cf.votes.contains_key(&voter_id) { @@ -274,7 +274,7 @@ impl ConsensusManager { } else { return false; } - self.check_finalization(&cf_id) + self.check_finalization(&cf_id, dag) } /// Check if a CF has reached quorum (finalize), rejection threshold (reject), or timeout @@ -283,7 +283,7 @@ impl ConsensusManager { /// Public because engine.receive_cf() needs to check after vote_for_cf(). /// /// Returns true if CF was finalized or rejected (removed from pending). - pub fn check_finalization(&mut self, cf_id: &str) -> bool { + pub fn check_finalization(&mut self, cf_id: &str, dag: &Dag) -> bool { let decision = { let cf = match self.pending_cfs.get(cf_id) { Some(cf) => cf, @@ -319,7 +319,7 @@ impl ConsensusManager { // PoCW: observe the committed fold and mint rewards if let Some(ref mut observer) = self.fold_observer { let fold_events = pending_build.all_events(); - if let Some(economics) = observer.on_fold_committed(&fold_events) { + if let Some(economics) = observer.on_fold_committed(&fold_events, dag) { self.apply_reward_minting(&economics); } } @@ -721,14 +721,18 @@ mod tests { // Create CF (deferred commit mode - state not modified yet) let cf = manager.try_create_cf(&dag, &vlc); assert!(cf.is_some()); - let cf_id = cf.unwrap().id.clone(); - + let cf = cf.unwrap(); + let cf_id = cf.id.clone(); + + // Add CF to pending (required before voting) + manager.receive_cf(cf); + // State not committed yet (prepare_build only) assert_eq!(manager.anchor_count(), 0); - + // Vote to finalize (single validator, so immediate finalization) manager.vote_for_cf(&cf_id, true, None); - let finalized = manager.check_finalization(&cf_id); + let finalized = manager.check_finalization(&cf_id, &dag); assert!(finalized, "CF should be finalized with single validator"); // Now state should be committed @@ -758,10 +762,13 @@ mod tests { event_count: 5, total_flux_burned: 105_000, total_power_consumed: 5, + flux_minted: 0, total_solver_rewards: 5, + kappa_before: 1.0, + kappa_after: 1.0, solver_rewards: vec![ - SolverReward { solver_id: "solver-a".into(), transfer_count: 3, flux_reward: 3 }, - SolverReward { solver_id: "solver-b".into(), transfer_count: 2, flux_reward: 2 }, + SolverReward { solver_id: "solver-a".into(), transfer_count: 3, task_count: 0, distance_score: 0.0, necessity_score: 0.0, contribution_score: 0.0, weight: 0.0, flux_reward: 3 }, + SolverReward { solver_id: "solver-b".into(), transfer_count: 2, task_count: 0, distance_score: 0.0, necessity_score: 0.0, contribution_score: 0.0, weight: 0.0, flux_reward: 2 }, ], }; @@ -778,14 +785,17 @@ mod tests { assert_eq!(u64::from_le_bytes(val_b.try_into().unwrap()), 2); } - // Apply a second fold — rewards accumulate + // Apply a second fold -- rewards accumulate let economics2 = FoldEconomics { event_count: 2, total_flux_burned: 42_000, total_power_consumed: 2, + flux_minted: 0, total_solver_rewards: 2, + kappa_before: 1.0, + kappa_after: 1.0, solver_rewards: vec![ - SolverReward { solver_id: "solver-a".into(), transfer_count: 2, flux_reward: 2 }, + SolverReward { solver_id: "solver-a".into(), transfer_count: 2, task_count: 0, distance_score: 0.0, necessity_score: 0.0, contribution_score: 0.0, weight: 0.0, flux_reward: 2 }, ], }; manager.apply_reward_minting(&economics2); diff --git a/consensus/src/pocw/emission.rs b/consensus/src/pocw/emission.rs new file mode 100644 index 0000000..95609a4 --- /dev/null +++ b/consensus/src/pocw/emission.rs @@ -0,0 +1,168 @@ +//! Flux minting and emission adjustment. +//! +//! Minting: FluxMinted = κ × ΔPower_total +//! Adjustment: κ(k+1) = clamp(κ(k) × (1 + dampening × (target/real − 1)), κ_min, κ_max) + +use setu_types::pocw::{EmissionState, PoCWConfig}; + +/// Compute Flux minted this fold. +/// +/// FluxMinted = round(κ × total_power_consumed) +pub fn compute_flux_minted(total_power_consumed: u64, kappa: f64) -> u64 { + (kappa * total_power_consumed as f64).round() as u64 +} + +/// Adjust κ based on observed vs target velocity. +/// +/// Records `flux_minted` in the velocity history, then adjusts κ: +/// real_velocity = mean(velocity_history) +/// κ_new = κ × (1 + dampening × (target / real − 1)) +/// κ_new = clamp(κ_new, κ_min, κ_max) +/// +/// Returns the new κ value. +pub fn adjust_kappa(state: &mut EmissionState, flux_minted: u64, config: &PoCWConfig) -> f64 { + state.velocity_history.push(flux_minted); + if state.velocity_history.len() > config.observation_window { + state.velocity_history.remove(0); + } + + let real_velocity = state.velocity_history.iter().sum::() as f64 + / state.velocity_history.len() as f64; + + if real_velocity == 0.0 { + return state.kappa; + } + + let ratio = config.target_velocity / real_velocity; + let dampened = 1.0 + config.dampening * (ratio - 1.0); + state.kappa = (state.kappa * dampened).clamp(config.kappa_min, config.kappa_max); + state.kappa +} + +#[cfg(test)] +mod tests { + use super::*; + + // -- Flux minting -- + + #[test] + fn test_minting_basic() { + // κ=1.0, power=100 → minted=100 + assert_eq!(compute_flux_minted(100, 1.0), 100); + } + + #[test] + fn test_minting_with_kappa() { + // κ=2.5, power=100 → minted=250 + assert_eq!(compute_flux_minted(100, 2.5), 250); + } + + #[test] + fn test_minting_zero_power() { + assert_eq!(compute_flux_minted(0, 5.0), 0); + } + + #[test] + fn test_minting_rounding() { + // κ=0.3, power=10 → 3.0 → 3 + assert_eq!(compute_flux_minted(10, 0.3), 3); + // κ=0.7, power=3 → 2.1 → 2 + assert_eq!(compute_flux_minted(3, 0.7), 2); + } + + // -- Kappa adjustment -- + + fn default_emission_config() -> PoCWConfig { + PoCWConfig { + enabled: true, + ..Default::default() + } + } + + #[test] + fn test_adjust_kappa_at_target() { + let config = default_emission_config(); + // target_velocity=1000, minted=1000 → ratio=1.0 → no change + let mut state = EmissionState::new(1.0); + let new_kappa = adjust_kappa(&mut state, 1000, &config); + assert!((new_kappa - 1.0).abs() < f64::EPSILON); + } + + #[test] + fn test_adjust_kappa_below_target() { + let config = default_emission_config(); + // target=1000, minted=500 → ratio=2.0 + // dampened = 1 + 0.5*(2.0-1.0) = 1.5 + // κ_new = 1.0 * 1.5 = 1.5 + let mut state = EmissionState::new(1.0); + let new_kappa = adjust_kappa(&mut state, 500, &config); + assert!((new_kappa - 1.5).abs() < 1e-10); + } + + #[test] + fn test_adjust_kappa_above_target() { + let config = default_emission_config(); + // target=1000, minted=2000 → ratio=0.5 + // dampened = 1 + 0.5*(0.5-1.0) = 0.75 + // κ_new = 1.0 * 0.75 = 0.75 + let mut state = EmissionState::new(1.0); + let new_kappa = adjust_kappa(&mut state, 2000, &config); + assert!((new_kappa - 0.75).abs() < 1e-10); + } + + #[test] + fn test_adjust_kappa_clamped_to_max() { + let config = PoCWConfig { + enabled: true, + kappa_max: 2.0, + ..Default::default() + }; + // Start at κ=1.8, minted=100 (way below target=1000) + // ratio=10.0, dampened=1+0.5*9=5.5, κ_new=1.8*5.5=9.9 → clamped to 2.0 + let mut state = EmissionState::new(1.8); + let new_kappa = adjust_kappa(&mut state, 100, &config); + assert!((new_kappa - 2.0).abs() < f64::EPSILON); + } + + #[test] + fn test_adjust_kappa_clamped_to_min() { + let config = PoCWConfig { + enabled: true, + kappa_min: 0.5, + ..Default::default() + }; + // Start at κ=0.6, minted=100_000 (way above target=1000) + // ratio=0.01, dampened=1+0.5*(0.01-1)=0.505, κ_new=0.6*0.505=0.303 → clamped to 0.5 + let mut state = EmissionState::new(0.6); + let new_kappa = adjust_kappa(&mut state, 100_000, &config); + assert!((new_kappa - 0.5).abs() < f64::EPSILON); + } + + #[test] + fn test_adjust_kappa_zero_minted_no_change() { + let config = default_emission_config(); + let mut state = EmissionState::new(1.5); + let new_kappa = adjust_kappa(&mut state, 0, &config); + assert!((new_kappa - 1.5).abs() < f64::EPSILON); + } + + #[test] + fn test_velocity_history_window() { + let config = PoCWConfig { + enabled: true, + observation_window: 3, + ..Default::default() + }; + let mut state = EmissionState::new(1.0); + + adjust_kappa(&mut state, 100, &config); + adjust_kappa(&mut state, 200, &config); + adjust_kappa(&mut state, 300, &config); + assert_eq!(state.velocity_history.len(), 3); + + // 4th push should evict the first (100) + adjust_kappa(&mut state, 400, &config); + assert_eq!(state.velocity_history.len(), 3); + assert_eq!(state.velocity_history[0], 200); + } +} diff --git a/consensus/src/pocw/flux_burn.rs b/consensus/src/pocw/flux_burn.rs index 4177177..6dbd0ef 100644 --- a/consensus/src/pocw/flux_burn.rs +++ b/consensus/src/pocw/flux_burn.rs @@ -1,10 +1,12 @@ -//! Flux burn calculation for FluxTransfer transactions. +//! Flux burn calculation for all transaction types. //! -//! Returns the configured fixed fee when PoCW is enabled, 0 otherwise. +//! Two paths: +//! - Transfers: flat fee from `PoCWConfig.transfer_fee` +//! - Tasks: formula-based `α*C + β*R + γ*S` using `EventMetrics` -use setu_types::pocw::PoCWConfig; +use setu_types::pocw::{EventMetrics, PoCWConfig}; -/// Calculate the Flux burn for a FluxTransfer transaction. +/// Calculate the Flux burn for a FluxTransfer (flat fee). pub fn calculate_transfer_burn(config: &PoCWConfig) -> u64 { if config.enabled { config.transfer_fee @@ -13,12 +15,50 @@ pub fn calculate_transfer_burn(config: &PoCWConfig) -> u64 { } } +/// Complexity score: compute_time + gas*10 + writes*1000 +fn compute_complexity(metrics: &EventMetrics) -> f64 { + metrics.compute_time_us as f64 + + metrics.gas_used as f64 * 10.0 + + metrics.write_count as f64 * 1000.0 +} + +/// Risk score: value at stake +fn compute_risk(metrics: &EventMetrics) -> f64 { + metrics.value_transferred as f64 +} + +/// Structural tension score: dag_depth*100 + writes*500 +fn compute_structural_tension(metrics: &EventMetrics) -> f64 { + metrics.dag_depth as f64 * 100.0 + + metrics.write_count as f64 * 500.0 +} + +/// Calculate the Flux burn for a TaskSubmit event (formula-based). +/// +/// burn = round(α*C + β*R + γ*S), minimum 1. +/// Returns 0 if PoCW is disabled. +pub fn calculate_task_burn(metrics: &EventMetrics, config: &PoCWConfig) -> u64 { + if !config.enabled { + return 0; + } + + let c = compute_complexity(metrics); + let r = compute_risk(metrics); + let s = compute_structural_tension(metrics); + + let burn = config.alpha * c + config.beta * r + config.gamma * s; + + (burn.round() as u64).max(1) +} + #[cfg(test)] mod tests { use super::*; + // -- Transfer burn (existing) -- + #[test] - fn test_enabled_returns_fixed_fee() { + fn test_transfer_enabled_returns_fixed_fee() { let config = PoCWConfig { enabled: true, ..Default::default() @@ -27,13 +67,13 @@ mod tests { } #[test] - fn test_disabled_returns_zero() { - let config = PoCWConfig::default(); // enabled: false + fn test_transfer_disabled_returns_zero() { + let config = PoCWConfig::default(); assert_eq!(calculate_transfer_burn(&config), 0); } #[test] - fn test_custom_fee() { + fn test_transfer_custom_fee() { let config = PoCWConfig { enabled: true, transfer_fee: 50_000, @@ -41,4 +81,80 @@ mod tests { }; assert_eq!(calculate_transfer_burn(&config), 50_000); } + + // -- Task burn -- + + fn make_metrics(compute_time_us: u64, gas_used: u64, write_count: usize, value_transferred: u64, dag_depth: u64) -> EventMetrics { + EventMetrics { + solver_id: "solver-1".to_string(), + compute_time_us, + gas_used, + write_count, + read_count: 0, + value_transferred, + dag_depth, + flux_burn: 0, + power_delta: 0, + } + } + + #[test] + fn test_task_burn_disabled_returns_zero() { + let config = PoCWConfig::default(); // enabled: false + let metrics = make_metrics(100, 500, 2, 0, 3); + assert_eq!(calculate_task_burn(&metrics, &config), 0); + } + + #[test] + fn test_task_burn_known_input() { + let config = PoCWConfig { + enabled: true, + ..Default::default() + }; + // C = 100 + 500*10 + 2*1000 = 7100 + // R = 0 + // S = 3*100 + 2*500 = 1300 + // burn = 0.4*7100 + 0.35*0 + 0.25*1300 = 2840 + 0 + 325 = 3165 + let metrics = make_metrics(100, 500, 2, 0, 3); + assert_eq!(calculate_task_burn(&metrics, &config), 3165); + } + + #[test] + fn test_task_burn_with_value_transfer() { + let config = PoCWConfig { + enabled: true, + ..Default::default() + }; + // C = 0 + 0 + 0 = 0 + // R = 10000 + // S = 0 + 0 = 0 + // burn = 0.4*0 + 0.35*10000 + 0.25*0 = 3500 + let metrics = make_metrics(0, 0, 0, 10_000, 0); + assert_eq!(calculate_task_burn(&metrics, &config), 3500); + } + + #[test] + fn test_task_burn_minimum_one() { + let config = PoCWConfig { + enabled: true, + ..Default::default() + }; + // All zeros produces burn = 0, but minimum is 1 + let metrics = make_metrics(0, 0, 0, 0, 0); + assert_eq!(calculate_task_burn(&metrics, &config), 1); + } + + #[test] + fn test_task_burn_custom_weights() { + let config = PoCWConfig { + enabled: true, + alpha: 1.0, + beta: 0.0, + gamma: 0.0, + ..Default::default() + }; + // Only complexity matters: C = 200 + 100*10 + 1*1000 = 2200 + let metrics = make_metrics(200, 100, 1, 50_000, 10); + assert_eq!(calculate_task_burn(&metrics, &config), 2200); + } } diff --git a/consensus/src/pocw/fold_observer.rs b/consensus/src/pocw/fold_observer.rs index df77a7f..a19f89f 100644 --- a/consensus/src/pocw/fold_observer.rs +++ b/consensus/src/pocw/fold_observer.rs @@ -3,30 +3,34 @@ //! Sits as a parallel path alongside existing consensus — does not modify //! anchor building, voting, or state application. -use setu_types::pocw::{PoCWConfig, FoldEconomics}; +use setu_types::pocw::{EmissionState, PoCWConfig, FoldEconomics}; use setu_types::Event; use tracing::info; +use crate::dag::Dag; use super::processor; /// Observes fold commits and produces economic summaries. pub struct FoldObserver { config: PoCWConfig, + emission: EmissionState, /// History of fold economics for diagnostics history: Vec, } impl FoldObserver { pub fn new(config: PoCWConfig) -> Self { + let kappa = config.kappa; Self { config, + emission: EmissionState::new(kappa), history: Vec::new(), } } /// Called after a fold is committed. Processes economics if enabled. - pub fn on_fold_committed(&mut self, events: &[Event]) -> Option { - let result = processor::process_fold(&self.config, events)?; + pub fn on_fold_committed(&mut self, events: &[Event], dag: &Dag) -> Option { + let result = processor::process_fold(&self.config, events, dag, &mut self.emission)?; info!( event_count = result.event_count, @@ -69,8 +73,9 @@ mod tests { #[test] fn test_disabled_produces_none() { let mut observer = FoldObserver::new(PoCWConfig::default()); + let dag = Dag::new(); let events = vec![make_transfer("solver-1")]; - assert!(observer.on_fold_committed(&events).is_none()); + assert!(observer.on_fold_committed(&events, &dag).is_none()); assert!(observer.history().is_empty()); } @@ -82,16 +87,18 @@ mod tests { ..Default::default() }; let mut observer = FoldObserver::new(config); + let dag = Dag::new(); let events = vec![ make_transfer("solver-1"), make_transfer("solver-2"), make_transfer("solver-1"), ]; - let result = observer.on_fold_committed(&events).unwrap(); + let result = observer.on_fold_committed(&events, &dag).unwrap(); assert_eq!(result.event_count, 3); assert_eq!(result.total_flux_burned, 21_000 * 3); - assert_eq!(result.total_solver_rewards, 3); + // flux_minted = κ(1.0) × 3 = 3, plus transfer rewards = 3 + assert_eq!(result.flux_minted, 3); assert_eq!(result.solver_rewards.len(), 2); assert_eq!(observer.history().len(), 1); } @@ -103,9 +110,10 @@ mod tests { ..Default::default() }; let mut observer = FoldObserver::new(config); + let dag = Dag::new(); - observer.on_fold_committed(&[make_transfer("s1")]); - observer.on_fold_committed(&[make_transfer("s2"), make_transfer("s2")]); + observer.on_fold_committed(&[make_transfer("s1")], &dag); + observer.on_fold_committed(&[make_transfer("s2"), make_transfer("s2")], &dag); assert_eq!(observer.history().len(), 2); assert_eq!(observer.history()[0].event_count, 1); diff --git a/consensus/src/pocw/mod.rs b/consensus/src/pocw/mod.rs index d308153..b574856 100644 --- a/consensus/src/pocw/mod.rs +++ b/consensus/src/pocw/mod.rs @@ -1,8 +1,11 @@ -//! PoCW economic calculations for FluxTransfer transactions +//! PoCW economic calculations for FluxTransfer and TaskSubmit events. //! -//! Provides flux burn and power drain computation, gated by PoCWConfig::enabled. +//! Provides flux burn, power drain, and scoring computation, +//! gated by PoCWConfig::enabled. +pub mod emission; pub mod flux_burn; pub mod power; pub mod processor; +pub mod scoring; pub mod fold_observer; diff --git a/consensus/src/pocw/power.rs b/consensus/src/pocw/power.rs index 141903b..930d8a7 100644 --- a/consensus/src/pocw/power.rs +++ b/consensus/src/pocw/power.rs @@ -1,10 +1,12 @@ -//! Power drain calculation for FluxTransfer transactions. +//! Power drain calculation for all transaction types. //! -//! Returns the configured flat drain when PoCW is enabled, 0 otherwise. +//! Two paths: +//! - Transfers: flat drain from `PoCWConfig.transfer_power_drain` +//! - Tasks: proportional drain `solver_power_rate * gas_used`, minimum 1 -use setu_types::pocw::PoCWConfig; +use setu_types::pocw::{EventMetrics, PoCWConfig}; -/// Calculate the power drain for a FluxTransfer transaction. +/// Calculate the power drain for a FluxTransfer (flat). pub fn calculate_transfer_power_drain(config: &PoCWConfig) -> u64 { if config.enabled { config.transfer_power_drain @@ -13,12 +15,28 @@ pub fn calculate_transfer_power_drain(config: &PoCWConfig) -> u64 { } } +/// Calculate the power drain for a TaskSubmit event (proportional to gas). +/// +/// ΔP = max(1, round(solver_power_rate * gas_used)) +/// Returns 0 if PoCW is disabled. +pub fn calculate_task_power_drain(metrics: &EventMetrics, config: &PoCWConfig) -> u64 { + if !config.enabled { + return 0; + } + + let drain = config.solver_power_rate * metrics.gas_used as f64; + + (drain.round() as u64).max(1) +} + #[cfg(test)] mod tests { use super::*; + // -- Transfer power drain (existing) -- + #[test] - fn test_enabled_returns_flat_drain() { + fn test_transfer_enabled_returns_flat_drain() { let config = PoCWConfig { enabled: true, ..Default::default() @@ -27,13 +45,13 @@ mod tests { } #[test] - fn test_disabled_returns_zero() { - let config = PoCWConfig::default(); // enabled: false + fn test_transfer_disabled_returns_zero() { + let config = PoCWConfig::default(); assert_eq!(calculate_transfer_power_drain(&config), 0); } #[test] - fn test_custom_drain() { + fn test_transfer_custom_drain() { let config = PoCWConfig { enabled: true, transfer_power_drain: 10, @@ -41,4 +59,74 @@ mod tests { }; assert_eq!(calculate_transfer_power_drain(&config), 10); } + + // -- Task power drain -- + + fn make_metrics(gas_used: u64) -> EventMetrics { + EventMetrics { + solver_id: "solver-1".to_string(), + compute_time_us: 0, + gas_used, + write_count: 0, + read_count: 0, + value_transferred: 0, + dag_depth: 0, + flux_burn: 0, + power_delta: 0, + } + } + + #[test] + fn test_task_power_disabled_returns_zero() { + let config = PoCWConfig::default(); // enabled: false + let metrics = make_metrics(10_000); + assert_eq!(calculate_task_power_drain(&metrics, &config), 0); + } + + #[test] + fn test_task_power_known_input() { + let config = PoCWConfig { + enabled: true, + ..Default::default() + }; + // solver_power_rate default = 0.001 + // drain = 0.001 * 10_000 = 10.0 → 10 + let metrics = make_metrics(10_000); + assert_eq!(calculate_task_power_drain(&metrics, &config), 10); + } + + #[test] + fn test_task_power_minimum_one() { + let config = PoCWConfig { + enabled: true, + ..Default::default() + }; + // drain = 0.001 * 1 = 0.001 → rounds to 0 → clamped to 1 + let metrics = make_metrics(1); + assert_eq!(calculate_task_power_drain(&metrics, &config), 1); + } + + #[test] + fn test_task_power_zero_gas_still_minimum_one() { + let config = PoCWConfig { + enabled: true, + ..Default::default() + }; + // drain = 0.001 * 0 = 0.0 → rounds to 0 → clamped to 1 + let metrics = make_metrics(0); + assert_eq!(calculate_task_power_drain(&metrics, &config), 1); + } + + #[test] + fn test_task_power_rounding() { + let config = PoCWConfig { + enabled: true, + solver_power_rate: 0.003, + ..Default::default() + }; + // drain = 0.003 * 1500 = 4.5 → rounds to 4 (banker's rounding) + // Actually f64::round() rounds 4.5 → 5 (round half away from zero) + let metrics = make_metrics(1500); + assert_eq!(calculate_task_power_drain(&metrics, &config), 5); + } } diff --git a/consensus/src/pocw/processor.rs b/consensus/src/pocw/processor.rs index f9e0d75..2585519 100644 --- a/consensus/src/pocw/processor.rs +++ b/consensus/src/pocw/processor.rs @@ -1,21 +1,33 @@ -//! Fold-level economic processor for FluxTransfer transactions. +//! Fold-level economic processor. //! -//! Takes a set of fold events and a PoCWConfig, produces a FoldEconomics summary -//! with per-solver reward breakdown. +//! Handles both FluxTransfer (flat fee) and TaskSubmit (formula-based) events. +//! Orchestrates burn aggregation, power aggregation, flux minting, PoCW reward +//! distribution, and kappa emission adjustment. use std::collections::HashMap; use setu_types::event::EventType; -use setu_types::pocw::{PoCWConfig, FoldEconomics, SolverReward}; +use setu_types::pocw::{EmissionState, PoCWConfig, FoldEconomics, SolverReward}; use setu_types::Event; +use crate::dag::Dag; +use super::emission::{compute_flux_minted, adjust_kappa}; use super::flux_burn::calculate_transfer_burn; use super::power::calculate_transfer_power_drain; +use super::scoring; /// Process a fold's events and produce an economic summary. /// -/// Only FluxTransfer events with `executed_by` set are counted for solver rewards. +/// Aggregates burn and power from both Transfer (flat) and TaskSubmit +/// (per-event metrics) events. Mints flux, distributes via PoCW scoring, +/// and adjusts kappa. +/// /// Returns `None` if PoCW is disabled. -pub fn process_fold(config: &PoCWConfig, events: &[Event]) -> Option { +pub fn process_fold( + config: &PoCWConfig, + events: &[Event], + dag: &Dag, + emission: &mut EmissionState, +) -> Option { if !config.enabled { return None; } @@ -24,43 +36,74 @@ pub fn process_fold(config: &PoCWConfig, events: &[Event]) -> Option = HashMap::new(); for event in events { - if event.event_type != EventType::Transfer { - continue; - } - transfer_count += 1; - - if let Some(solver_id) = &event.executed_by { - *solver_transfers.entry(solver_id.clone()).or_default() += 1; + match event.event_type { + EventType::Transfer => { + transfer_count += 1; + if let Some(solver_id) = &event.executed_by { + *solver_transfers.entry(solver_id.clone()).or_default() += 1; + } + } + EventType::TaskSubmit => { + if let Some(ref metrics) = event.event_metrics { + task_burn += metrics.flux_burn; + task_power += metrics.power_delta; + } + } + _ => {} } } - let total_flux_burned = burn_per_transfer * transfer_count as u64; - let total_power_consumed = power_per_transfer * transfer_count as u64; + let total_flux_burned = burn_per_transfer * transfer_count as u64 + task_burn; + let total_power_consumed = power_per_transfer * transfer_count as u64 + task_power; - let mut solver_rewards = Vec::new(); - let mut total_solver_rewards: u64 = 0; + // Minting: FluxMinted = κ × ΔPower_total + let kappa_before = emission.kappa; + let flux_minted = compute_flux_minted(total_power_consumed, kappa_before); + // PoCW reward distribution + let mut solver_rewards = scoring::compute_rewards(dag, events, flux_minted, config); + + // Merge transfer rewards into solver records if config.solver_transfer_reward_enabled { for (solver_id, count) in &solver_transfers { - let reward = config.solver_transfer_reward * count; - total_solver_rewards += reward; - solver_rewards.push(SolverReward { - solver_id: solver_id.clone(), - transfer_count: *count, - flux_reward: reward, - }); + let transfer_reward = config.solver_transfer_reward * count; + if let Some(existing) = solver_rewards.iter_mut().find(|r| r.solver_id == *solver_id) { + existing.transfer_count = *count; + existing.flux_reward += transfer_reward; + } else { + solver_rewards.push(SolverReward { + solver_id: solver_id.clone(), + transfer_count: *count, + task_count: 0, + distance_score: 0.0, + necessity_score: 0.0, + contribution_score: 0.0, + weight: 0.0, + flux_reward: transfer_reward, + }); + } } solver_rewards.sort_by(|a, b| a.solver_id.cmp(&b.solver_id)); } + let total_solver_rewards = solver_rewards.iter().map(|r| r.flux_reward).sum(); + + // Emission adjustment: κ(k+1) = f(κ(k), flux_minted) + let kappa_after = adjust_kappa(emission, flux_minted, config); + Some(FoldEconomics { event_count: events.len(), total_flux_burned, total_power_consumed, + flux_minted, total_solver_rewards, + kappa_before, + kappa_after, solver_rewards, }) } @@ -98,17 +141,23 @@ mod tests { } } + fn call_process_fold(config: &PoCWConfig, events: &[Event]) -> Option { + let dag = Dag::new(); + let mut emission = EmissionState::new(config.kappa); + process_fold(config, events, &dag, &mut emission) + } + #[test] fn test_disabled_returns_none() { let config = PoCWConfig::default(); // enabled: false let events = vec![make_transfer_event(Some("solver-1"))]; - assert!(process_fold(&config, &events).is_none()); + assert!(call_process_fold(&config, &events).is_none()); } #[test] fn test_empty_fold() { let config = enabled_config(); - let result = process_fold(&config, &[]).unwrap(); + let result = call_process_fold(&config, &[]).unwrap(); assert_eq!(result.event_count, 0); assert_eq!(result.total_flux_burned, 0); assert_eq!(result.total_power_consumed, 0); @@ -125,15 +174,17 @@ mod tests { make_transfer_event(Some("solver-1")), ]; - let result = process_fold(&config, &events).unwrap(); + let result = call_process_fold(&config, &events).unwrap(); assert_eq!(result.event_count, 3); assert_eq!(result.total_flux_burned, 21_000 * 3); assert_eq!(result.total_power_consumed, 1 * 3); - assert_eq!(result.total_solver_rewards, 1 * 3); + // flux_minted = κ(1.0) × 3 = 3 + assert_eq!(result.flux_minted, 3); assert_eq!(result.solver_rewards.len(), 1); assert_eq!(result.solver_rewards[0].solver_id, "solver-1"); assert_eq!(result.solver_rewards[0].transfer_count, 3); - assert_eq!(result.solver_rewards[0].flux_reward, 3); + // PoCW reward (3) + transfer reward (3*1) = 6 + assert_eq!(result.solver_rewards[0].flux_reward, 3 + 3); } #[test] @@ -147,18 +198,22 @@ mod tests { make_transfer_event(Some("solver-2")), ]; - let result = process_fold(&config, &events).unwrap(); + let result = call_process_fold(&config, &events).unwrap(); assert_eq!(result.event_count, 5); assert_eq!(result.total_flux_burned, 21_000 * 5); - assert_eq!(result.total_solver_rewards, 5); // 2 + 3 + // flux_minted = κ(1.0) × 5 = 5 + assert_eq!(result.flux_minted, 5); + // PoCW rewards distributed by scoring (equal weight since + // events not in DAG → equal distance/necessity, zero contribution) + // Transfer rewards added on top: solver-1 gets 2, solver-2 gets 3 // Sorted by solver_id assert_eq!(result.solver_rewards[0].solver_id, "solver-1"); assert_eq!(result.solver_rewards[0].transfer_count, 2); - assert_eq!(result.solver_rewards[0].flux_reward, 2); assert_eq!(result.solver_rewards[1].solver_id, "solver-2"); assert_eq!(result.solver_rewards[1].transfer_count, 3); - assert_eq!(result.solver_rewards[1].flux_reward, 3); + // Total = PoCW distributed (5) + transfer rewards (5) = 10 + assert_eq!(result.total_solver_rewards, 10); } #[test] @@ -170,24 +225,28 @@ mod tests { }; let events = vec![make_transfer_event(Some("solver-1"))]; - let result = process_fold(&config, &events).unwrap(); + let result = call_process_fold(&config, &events).unwrap(); assert_eq!(result.total_flux_burned, 21_000); - assert_eq!(result.total_solver_rewards, 0); - assert!(result.solver_rewards.is_empty()); + // flux_minted = κ(1.0) × 1 = 1, distributed via PoCW scoring + assert_eq!(result.flux_minted, 1); + // PoCW reward only (no transfer reward since disabled) + assert_eq!(result.solver_rewards.len(), 1); + assert_eq!(result.solver_rewards[0].flux_reward, 1); } #[test] - fn test_non_transfer_events_ignored() { + fn test_non_transfer_events_ignored_for_burn() { let config = enabled_config(); let events = vec![ make_transfer_event(Some("solver-1")), - make_genesis_event(), // not a transfer + make_genesis_event(), // not a transfer — no burn or power ]; - let result = process_fold(&config, &events).unwrap(); - assert_eq!(result.event_count, 2); // total events in fold + let result = call_process_fold(&config, &events).unwrap(); + assert_eq!(result.event_count, 2); assert_eq!(result.total_flux_burned, 21_000); // only 1 transfer burned - assert_eq!(result.total_solver_rewards, 1); + assert_eq!(result.total_power_consumed, 1); + assert_eq!(result.flux_minted, 1); } #[test] @@ -198,12 +257,15 @@ mod tests { make_transfer_event(None), // no solver recorded ]; - let result = process_fold(&config, &events).unwrap(); + let result = call_process_fold(&config, &events).unwrap(); // Both transfers still burn and drain assert_eq!(result.total_flux_burned, 21_000 * 2); assert_eq!(result.total_power_consumed, 2); - // Only the one with executed_by gets a reward - assert_eq!(result.total_solver_rewards, 1); + assert_eq!(result.flux_minted, 2); + // PoCW reward for solver-1 (only executed_by events get scoring) + // Transfer reward for solver-1: 1*1 = 1 + // solver-1 gets PoCW(2) + transfer(1) = 3 assert_eq!(result.solver_rewards.len(), 1); + assert_eq!(result.solver_rewards[0].flux_reward, 2 + 1); } } diff --git a/consensus/src/pocw/scoring.rs b/consensus/src/pocw/scoring.rs new file mode 100644 index 0000000..89cbb6d --- /dev/null +++ b/consensus/src/pocw/scoring.rs @@ -0,0 +1,435 @@ +//! PoCW reward distribution across solvers. +//! +//! Scoring signals: +//! - Distance: 1/(1 + Distance_i) — proximity to critical path +//! - Necessity: relevant_events / total_events — fraction causally needed +//! - Contribution: agent_depth / total_depth — share of causal depth +//! +//! w_i = α·Distance + β·Necessity + γ·Contribution, normalized +//! Reward_i = w_i × FluxMinted + +use std::collections::{HashMap, HashSet}; + +use setu_types::event::EventType; +use setu_types::pocw::{PoCWConfig, SolverReward}; +use setu_types::Event; + +use crate::dag::Dag; + +/// Compute PoCW reward distribution for a fold's events. +/// +/// Groups events by solver (from `executed_by`), scores each solver on +/// distance/necessity/contribution, normalizes weights, and distributes +/// `flux_minted` proportionally. Returns sorted by solver_id. +/// +/// Returns empty vec if no events have a solver attribution. +pub fn compute_rewards( + dag: &Dag, + fold_events: &[Event], + flux_minted: u64, + config: &PoCWConfig, +) -> Vec { + // Group events by solver + let mut solver_events: HashMap> = HashMap::new(); + for event in fold_events { + if let Some(ref solver_id) = event.executed_by { + solver_events + .entry(solver_id.clone()) + .or_default() + .push(event); + } + } + + if solver_events.is_empty() { + return Vec::new(); + } + + // Fold event IDs for intersection filtering + let fold_ids: HashSet<&str> = fold_events.iter().map(|e| e.id.as_str()).collect(); + + // Find fold tips: events within the fold that no other fold event points to as parent + let tip_ancestors = compute_tip_ancestry(dag, fold_events, &fold_ids); + + // Per-solver depth sums and fold-wide totals for contribution scoring + let total_depth: f64 = fold_events + .iter() + .filter_map(|e| dag.get_depth(&e.id)) + .sum::() as f64; + + // Score each solver + let mut scored: Vec<(String, f64, f64, f64, f64, u64, u64)> = Vec::new(); + + for (solver_id, events) in &solver_events { + // Distance: single-solver → all events are on the critical path + // distance = 0 → score = 1/(1+0) = 1.0 + // With multi-solver competition, replace with BFS to winner's critical path + let distance_score = 1.0; + + // Necessity: fraction of solver events in tip ancestry + let necessary = events + .iter() + .filter(|e| tip_ancestors.contains(e.id.as_str())) + .count(); + let necessity_score = if events.is_empty() { + 0.0 + } else { + necessary as f64 / events.len() as f64 + }; + + // Contribution: solver_depth / total_depth + let solver_depth: f64 = events + .iter() + .filter_map(|e| dag.get_depth(&e.id)) + .sum::() as f64; + let contribution_score = if total_depth > 0.0 { + solver_depth / total_depth + } else { + 0.0 + }; + + let raw_w = config.pocw_distance_weight * distance_score + + config.pocw_necessity_weight * necessity_score + + config.pocw_contribution_weight * contribution_score; + + let transfer_count = events + .iter() + .filter(|e| e.event_type == EventType::Transfer) + .count() as u64; + let task_count = events + .iter() + .filter(|e| e.event_type == EventType::TaskSubmit) + .count() as u64; + + scored.push(( + solver_id.clone(), + distance_score, + necessity_score, + contribution_score, + raw_w, + transfer_count, + task_count, + )); + } + + // Normalize weights and distribute flux + let total_raw: f64 = scored.iter().map(|s| s.4).sum(); + + let mut rewards: Vec = scored + .iter() + .map( + |(solver_id, distance, necessity, contribution, raw_w, transfers, tasks)| { + let weight = if total_raw > 0.0 { + raw_w / total_raw + } else { + 0.0 + }; + let flux_reward = (weight * flux_minted as f64).floor() as u64; + + SolverReward { + solver_id: solver_id.clone(), + transfer_count: *transfers, + task_count: *tasks, + distance_score: *distance, + necessity_score: *necessity, + contribution_score: *contribution, + weight, + flux_reward, + } + }, + ) + .collect(); + + // Deterministic sort + rewards.sort_by(|a, b| a.solver_id.cmp(&b.solver_id)); + + // Fix rounding remainder: give to highest-weighted solver + let total_distributed: u64 = rewards.iter().map(|r| r.flux_reward).sum(); + if total_distributed < flux_minted { + let remainder = flux_minted - total_distributed; + if let Some(top) = rewards + .iter_mut() + .max_by(|a, b| a.weight.partial_cmp(&b.weight).unwrap_or(std::cmp::Ordering::Equal)) + { + top.flux_reward += remainder; + } + } + + rewards +} + +/// Compute the set of fold event IDs that are ancestors of fold tips. +/// +/// A fold tip is an event within the fold that no other fold event references +/// as a parent. The returned set includes the tips themselves. +fn compute_tip_ancestry<'a>( + dag: &Dag, + fold_events: &'a [Event], + fold_ids: &HashSet<&'a str>, +) -> HashSet { + // Find fold tips + let tips: Vec<&Event> = fold_events + .iter() + .filter(|e| { + !fold_events + .iter() + .any(|other| other.parent_ids.contains(&e.id)) + }) + .collect(); + + // Build union of all tip ancestors, intersected with fold events + let mut ancestors: HashSet = HashSet::new(); + for tip in &tips { + ancestors.insert(tip.id.clone()); + for ancestor_id in dag.get_ancestors(&tip.id) { + if fold_ids.contains(ancestor_id.as_str()) { + ancestors.insert(ancestor_id); + } + } + } + + ancestors +} + +#[cfg(test)] +mod tests { + use super::*; + use setu_types::event::{VLCSnapshot, VectorClock}; + use setu_types::EventId; + use std::sync::atomic::{AtomicU64, Ordering}; + + /// Monotonic counter to ensure unique VLC logical_time per event. + static COUNTER: AtomicU64 = AtomicU64::new(1); + + /// Build a test event with a given type, parents, and solver attribution. + fn make_event( + event_type: EventType, + parent_ids: Vec, + executed_by: Option<&str>, + ) -> Event { + let vlc = VLCSnapshot { + vector_clock: VectorClock::new(), + logical_time: COUNTER.fetch_add(1, Ordering::Relaxed), + physical_time: 1000, + }; + let mut event = Event::new( + event_type, + parent_ids, + vlc, + "validator-1".to_string(), + ); + event.executed_by = executed_by.map(|s| s.to_string()); + event + } + + /// Insert an event into the DAG and return its ID. + fn insert(dag: &mut Dag, event: Event) -> EventId { + dag.add_event(event).expect("dag insert failed") + } + + fn default_config() -> PoCWConfig { + PoCWConfig { + enabled: true, + ..Default::default() + } + } + + // -- Basic behavior -- + + #[test] + fn test_no_solver_events_returns_empty() { + let mut dag = Dag::new(); + let genesis = make_event(EventType::Genesis, vec![], None); + let gid = insert(&mut dag, genesis); + + let fold_events = vec![dag.get_event(&gid).unwrap().clone()]; + let rewards = compute_rewards(&dag, &fold_events, 1000, &default_config()); + assert!(rewards.is_empty()); + } + + #[test] + fn test_single_solver_gets_all_flux() { + let mut dag = Dag::new(); + let genesis = make_event(EventType::Genesis, vec![], None); + let gid = insert(&mut dag, genesis); + + let e1 = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-1")); + let e1_id = insert(&mut dag, e1); + + let fold_events = vec![dag.get_event(&e1_id).unwrap().clone()]; + let rewards = compute_rewards(&dag, &fold_events, 1000, &default_config()); + + assert_eq!(rewards.len(), 1); + assert_eq!(rewards[0].solver_id, "solver-1"); + assert_eq!(rewards[0].flux_reward, 1000); + assert_eq!(rewards[0].task_count, 1); + } + + #[test] + fn test_two_solvers_equal_work() { + let mut dag = Dag::new(); + let genesis = make_event(EventType::Genesis, vec![], None); + let gid = insert(&mut dag, genesis); + + // Both solvers produce one event at the same depth + let e1 = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-1")); + let e1_id = insert(&mut dag, e1); + let e2 = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-2")); + let e2_id = insert(&mut dag, e2); + + let fold_events = vec![ + dag.get_event(&e1_id).unwrap().clone(), + dag.get_event(&e2_id).unwrap().clone(), + ]; + let rewards = compute_rewards(&dag, &fold_events, 1000, &default_config()); + + assert_eq!(rewards.len(), 2); + // Equal work → equal reward (500 each) + assert_eq!(rewards[0].flux_reward + rewards[1].flux_reward, 1000); + assert_eq!(rewards[0].flux_reward, 500); + assert_eq!(rewards[1].flux_reward, 500); + } + + // -- Necessity scoring -- + + #[test] + fn test_necessity_dead_end_scores_lower() { + // DAG shape: + // genesis → A (solver-1) → C (solver-1) ← fold tip + // genesis → B (solver-2) ← fold tip (dead-end relative to C) + // + // Solver-1: both A and C are ancestors of tip C → necessity = 1.0 + // Solver-2: B is a tip itself → necessity = 1.0 + // But solver-1 has higher contribution (depth sum = 1+2=3 vs 1) + let mut dag = Dag::new(); + let genesis = make_event(EventType::Genesis, vec![], None); + let gid = insert(&mut dag, genesis); + + let a = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-1")); + let aid = insert(&mut dag, a); + + let c = make_event(EventType::TaskSubmit, vec![aid.clone()], Some("solver-1")); + let cid = insert(&mut dag, c); + + let b = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-2")); + let bid = insert(&mut dag, b); + + let fold_events = vec![ + dag.get_event(&aid).unwrap().clone(), + dag.get_event(&bid).unwrap().clone(), + dag.get_event(&cid).unwrap().clone(), + ]; + let rewards = compute_rewards(&dag, &fold_events, 1000, &default_config()); + + // solver-1 has higher contribution (depth 1+2=3 vs depth 1) + let s1 = rewards.iter().find(|r| r.solver_id == "solver-1").unwrap(); + let s2 = rewards.iter().find(|r| r.solver_id == "solver-2").unwrap(); + assert!(s1.contribution_score > s2.contribution_score); + assert!(s1.flux_reward > s2.flux_reward); + } + + // -- Contribution scoring -- + + #[test] + fn test_deeper_solver_gets_higher_contribution() { + // genesis → A (solver-1) → B (solver-1) → C (solver-1) + // → D (solver-2) + // solver-1 depth sum = 1+2+3 = 6, solver-2 depth sum = 1 + // total = 7, contribution: 6/7 vs 1/7 + let mut dag = Dag::new(); + let genesis = make_event(EventType::Genesis, vec![], None); + let gid = insert(&mut dag, genesis); + + let a = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-1")); + let aid = insert(&mut dag, a); + let b = make_event(EventType::TaskSubmit, vec![aid.clone()], Some("solver-1")); + let bid = insert(&mut dag, b); + let c = make_event(EventType::TaskSubmit, vec![bid.clone()], Some("solver-1")); + let cid = insert(&mut dag, c); + + let d = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-2")); + let did = insert(&mut dag, d); + + let fold_events = vec![ + dag.get_event(&aid).unwrap().clone(), + dag.get_event(&bid).unwrap().clone(), + dag.get_event(&cid).unwrap().clone(), + dag.get_event(&did).unwrap().clone(), + ]; + let rewards = compute_rewards(&dag, &fold_events, 700, &default_config()); + + let s1 = rewards.iter().find(|r| r.solver_id == "solver-1").unwrap(); + let s2 = rewards.iter().find(|r| r.solver_id == "solver-2").unwrap(); + + assert!((s1.contribution_score - 6.0 / 7.0).abs() < 1e-10); + assert!((s2.contribution_score - 1.0 / 7.0).abs() < 1e-10); + assert!(s1.flux_reward > s2.flux_reward); + } + + // -- Rounding remainder -- + + #[test] + fn test_rounding_remainder_conserved() { + // Verify total distributed == flux_minted regardless of rounding + let mut dag = Dag::new(); + let genesis = make_event(EventType::Genesis, vec![], None); + let gid = insert(&mut dag, genesis); + + let e1 = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-1")); + let e1_id = insert(&mut dag, e1); + let e2 = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-2")); + let e2_id = insert(&mut dag, e2); + let e3 = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-3")); + let e3_id = insert(&mut dag, e3); + + let fold_events = vec![ + dag.get_event(&e1_id).unwrap().clone(), + dag.get_event(&e2_id).unwrap().clone(), + dag.get_event(&e3_id).unwrap().clone(), + ]; + // 1000 / 3 doesn't divide evenly + let rewards = compute_rewards(&dag, &fold_events, 1000, &default_config()); + let total: u64 = rewards.iter().map(|r| r.flux_reward).sum(); + assert_eq!(total, 1000); + } + + // -- Transfer vs task counting -- + + #[test] + fn test_transfer_and_task_counts() { + let mut dag = Dag::new(); + let genesis = make_event(EventType::Genesis, vec![], None); + let gid = insert(&mut dag, genesis); + + let t1 = make_event(EventType::Transfer, vec![gid.clone()], Some("solver-1")); + let t1_id = insert(&mut dag, t1); + let t2 = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-1")); + let t2_id = insert(&mut dag, t2); + + let fold_events = vec![ + dag.get_event(&t1_id).unwrap().clone(), + dag.get_event(&t2_id).unwrap().clone(), + ]; + let rewards = compute_rewards(&dag, &fold_events, 100, &default_config()); + + assert_eq!(rewards[0].transfer_count, 1); + assert_eq!(rewards[0].task_count, 1); + } + + // -- Zero flux minted -- + + #[test] + fn test_zero_flux_minted() { + let mut dag = Dag::new(); + let genesis = make_event(EventType::Genesis, vec![], None); + let gid = insert(&mut dag, genesis); + + let e1 = make_event(EventType::TaskSubmit, vec![gid.clone()], Some("solver-1")); + let e1_id = insert(&mut dag, e1); + + let fold_events = vec![dag.get_event(&e1_id).unwrap().clone()]; + let rewards = compute_rewards(&dag, &fold_events, 0, &default_config()); + + assert_eq!(rewards.len(), 1); + assert_eq!(rewards[0].flux_reward, 0); + } +} diff --git a/consensus/src/root_executor.rs b/consensus/src/root_executor.rs index 7608c69..2eca1b8 100644 --- a/consensus/src/root_executor.rs +++ b/consensus/src/root_executor.rs @@ -12,8 +12,9 @@ use setu_merkle::sha256; use setu_types::{ - Event, EventType, SubnetId, + Event, EventType, SubnetId, ObjectStateValue, object_type, HashValue as TypesHash, ZERO_HASH, + pocw::PoCWConfig, }; use std::collections::HashMap; use thiserror::Error; @@ -77,12 +78,15 @@ impl RootExecutionResult { pub struct RootSubnetExecutor { /// Current state root current_state_root: TypesHash, - + /// Pending updates (not yet committed) pending_updates: HashMap<[u8; 32], ObjectStateValue>, - + /// Pending deletions pending_deletions: Vec<[u8; 32]>, + + /// PoCW config for initializing solver power budgets + pocw_config: Option, } impl RootSubnetExecutor { @@ -92,13 +96,20 @@ impl RootSubnetExecutor { current_state_root: initial_state_root, pending_updates: HashMap::new(), pending_deletions: Vec::new(), + pocw_config: None, } } - + /// Create with empty state pub fn empty() -> Self { Self::new(ZERO_HASH) } + + /// Set PoCW config for solver power budget initialization + pub fn with_pocw_config(mut self, config: PoCWConfig) -> Self { + self.pocw_config = Some(config); + self + } /// Get the current state root pub fn state_root(&self) -> TypesHash { @@ -189,9 +200,11 @@ impl RootSubnetExecutor { } /// Execute solver registration + /// + /// Also initializes the solver's PoCW power budget if PoCW is enabled. fn execute_solver_register(&mut self, event: &Event) -> Result { let key = self.generate_solver_key(&event.creator); - + let state = ObjectStateValue::new( self.address_to_bytes(&event.creator), 1, @@ -199,15 +212,36 @@ impl RootSubnetExecutor { *sha256(event.id.as_bytes()).as_bytes(), SubnetId::ROOT, ); - + self.pending_updates.insert(key, state.clone()); - + + // Initialize solver power budget if PoCW is enabled + if let Some(ref config) = self.pocw_config { + if config.enabled { + let power_key = self.generate_power_key(&event.creator); + let power_value = ObjectStateValue::new( + self.address_to_bytes(&event.creator), + 1, + object_type::SOLVER_INFO, // reuse type; power is a sub-object + *sha256(format!("power:{}", event.id).as_bytes()).as_bytes(), + SubnetId::ROOT, + ); + self.pending_updates.insert(power_key, power_value); + + tracing::info!( + solver_id = %event.creator, + initial_power = config.initial_power_budget, + "Initialized solver power budget" + ); + } + } + let new_root = self.compute_pending_root(); self.current_state_root = new_root; - + let mut updated = HashMap::new(); updated.insert(key, state); - + Ok(RootExecutionResult { updated_objects: updated, deleted_objects: Vec::new(), @@ -215,6 +249,15 @@ impl RootSubnetExecutor { event_id: event.id.clone(), }) } + + /// Generate object key for a solver's power budget + fn generate_power_key(&self, solver_id: &str) -> [u8; 32] { + let mut key = [0u8; 32]; + key[0] = object_type::SOLVER_INFO; + let hash = sha256(format!("power:{}", solver_id).as_bytes()); + key[1..].copy_from_slice(&hash.as_bytes()[..31]); + key + } /// Execute solver unregistration fn execute_solver_unregister(&mut self, event: &Event) -> Result { @@ -335,6 +378,7 @@ impl Default for RootSubnetExecutor { } } + #[cfg(test)] mod tests { use super::*; diff --git a/setu-rpc/src/messages.rs b/setu-rpc/src/messages.rs index a70eba4..ddf3320 100644 --- a/setu-rpc/src/messages.rs +++ b/setu-rpc/src/messages.rs @@ -262,6 +262,37 @@ pub struct ProcessingStep { pub timestamp: u64, } +/// Request to submit a task for solver execution +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SubmitTaskRequest { + /// Task type (e.g., "compute", "inference", "custom") + pub task_type: String, + /// Submitter address + pub submitter: String, + /// Task payload (opaque bytes, interpreted by the solver) + #[serde(default)] + pub payload: Vec, + /// Optional preferred solver + pub preferred_solver: Option, + /// Optional subnet ID + pub subnet_id: Option, +} + +/// Response to task submission +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SubmitTaskResponse { + /// Whether submission was successful + pub success: bool, + /// Human-readable message + pub message: String, + /// Assigned task ID (transfer tracker ID) + pub task_id: Option, + /// Assigned solver ID + pub solver_id: Option, + /// Processing steps + pub processing_steps: Vec, +} + /// Request to get transfer status #[derive(Debug, Clone, Serialize, Deserialize)] pub struct GetTransferStatusRequest { diff --git a/setu-validator/src/main.rs b/setu-validator/src/main.rs index db054ee..cbd7015 100644 --- a/setu-validator/src/main.rs +++ b/setu-validator/src/main.rs @@ -223,8 +223,9 @@ async fn main() -> anyhow::Result<()> { info!("✓ PoCW enabled (transfer_fee={})", fee); } - // Extract transfer fee before consensus_config is consumed - let transfer_fee = consensus_config.consensus.pocw + // Extract PoCW config before consensus_config is consumed + let pocw_config = consensus_config.consensus.pocw; + let transfer_fee = pocw_config .filter(|p| p.enabled) .map(|p| p.transfer_fee) .unwrap_or(0); @@ -277,6 +278,7 @@ async fn main() -> anyhow::Result<()> { http_listen_addr: config.http_addr, p2p_listen_addr: config.p2p_addr, transfer_fee, + pocw_config: pocw_config.filter(|p| p.enabled), }; // Create network service with consensus enabled @@ -343,6 +345,7 @@ async fn main() -> anyhow::Result<()> { info!("║ GET /api/v1/validators - List validators ║"); info!("║ GET /api/v1/health - Health check ║"); info!("║ POST /api/v1/transfer - Submit transfer ║"); + info!("║ POST /api/v1/task - Submit task ║"); info!("║ POST /api/v1/event - Submit event (Solver) ║"); info!("║ GET /api/v1/events - List events ║"); info!("╚════════════════════════════════════════════════════════════╝"); diff --git a/setu-validator/src/network/service.rs b/setu-validator/src/network/service.rs index 4800623..7b1b675 100644 --- a/setu-validator/src/network/service.rs +++ b/setu-validator/src/network/service.rs @@ -143,7 +143,7 @@ impl ValidatorNetworkService { let dag_events = Arc::new(RwLock::new(Vec::new())); // Create TEE executor - let tee_executor = TeeExecutor::new( + let mut tee_executor = TeeExecutor::new( http_client.clone(), Arc::clone(&solver_info), Arc::clone(&transfer_status), @@ -152,7 +152,10 @@ impl ValidatorNetworkService { None, // No consensus validator_id.clone(), 100, // Max concurrent TEE calls - ); + ).with_task_preparer(Arc::clone(&task_preparer)); + if let Some(pocw) = config.pocw_config { + tee_executor = tee_executor.with_pocw_config(pocw); + } Self { validator_id, @@ -208,7 +211,7 @@ impl ValidatorNetworkService { let dag_events = Arc::new(RwLock::new(Vec::new())); // Create TEE executor with consensus - let tee_executor = TeeExecutor::new( + let mut tee_executor = TeeExecutor::new( http_client.clone(), Arc::clone(&solver_info), Arc::clone(&transfer_status), @@ -217,7 +220,10 @@ impl ValidatorNetworkService { Some(Arc::clone(&consensus_validator)), validator_id.clone(), 100, // Max concurrent TEE calls - ); + ).with_task_preparer(Arc::clone(&task_preparer)); + if let Some(pocw) = config.pocw_config { + tee_executor = tee_executor.with_pocw_config(pocw); + } Self { validator_id, @@ -357,6 +363,8 @@ impl ValidatorNetworkService { // Transfer endpoints .route("/api/v1/transfer", post(setu_api::http_submit_transfer::)) .route("/api/v1/transfer/status", post(setu_api::http_get_transfer_status::)) + // Task submission endpoint + .route("/api/v1/task", post(setu_api::http_submit_task::)) // Event endpoints .route("/api/v1/event", post(setu_api::http_submit_event::)) .route("/api/v1/events", get(setu_api::http_get_events::)) @@ -408,6 +416,137 @@ impl ValidatorNetworkService { TransferHandler::get_transfer_status(&self.transfer_status, transfer_id) } + // ============================================ + // Task Submission (TaskSubmit events) + // ============================================ + + pub async fn submit_task(&self, request: setu_rpc::SubmitTaskRequest) -> setu_rpc::SubmitTaskResponse { + let now = super::types::current_timestamp_secs(); + let task_counter = self.transfer_counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst); + let task_id = format!("task-{}-{}", now, task_counter); + + let mut steps = Vec::new(); + + info!(task_id = %task_id, submitter = %request.submitter, task_type = %request.task_type, "Processing task submission"); + + steps.push(setu_rpc::ProcessingStep { + step: "receive".to_string(), + status: "completed".to_string(), + details: Some(format!("Task {} received", task_id)), + timestamp: now, + }); + + // VLC assignment + let vlc_time = self.get_vlc_time(); + let vlc_snapshot = { + let mut snapshot = setu_types::event::VLCSnapshot::for_node(self.validator_id.clone()); + snapshot.logical_time = vlc_time; + snapshot.physical_time = super::types::current_timestamp_millis(); + snapshot + }; + + steps.push(setu_rpc::ProcessingStep { + step: "vlc_assign".to_string(), + status: "completed".to_string(), + details: Some(format!("VLC time: {}", vlc_time)), + timestamp: now, + }); + + // Build TaskSubmission + let task_submission = setu_types::registration::TaskSubmission::new( + &task_id, + &request.task_type, + &request.submitter, + ).with_payload(request.payload); + + // Prepare SolverTask + let subnet_id = setu_types::SubnetId::ROOT; + let solver_task = match self.task_preparer.prepare_task_submission( + &task_submission, + subnet_id, + vec![], // No explicit parents; DAG assigns based on tips + vlc_snapshot, + ) { + Ok(task) => { + steps.push(setu_rpc::ProcessingStep { + step: "prepare_task".to_string(), + status: "completed".to_string(), + details: Some("SolverTask prepared".to_string()), + timestamp: now, + }); + task + } + Err(e) => { + let msg = format!("Task preparation failed: {}", e); + return setu_rpc::SubmitTaskResponse { + success: false, + message: msg, + task_id: Some(task_id), + solver_id: None, + processing_steps: steps, + }; + } + }; + + // Route to solver (reuse transfer routing) + let transfer_for_routing = setu_types::Transfer::new( + &task_id, + &request.submitter, + &request.submitter, // task doesn't have a "to" — use submitter + 0, + ) + .with_type(setu_types::TransferType::TaskSubmit) + .with_preferred_solver_opt(request.preferred_solver); + + let solver_id = match self.router_manager.route_transfer(&transfer_for_routing) { + Ok(id) => { + steps.push(setu_rpc::ProcessingStep { + step: "route".to_string(), + status: "completed".to_string(), + details: Some(format!("Routed to: {}", id)), + timestamp: now, + }); + id + } + Err(e) => { + let msg = format!("No solver available: {}", e); + return setu_rpc::SubmitTaskResponse { + success: false, + message: msg, + task_id: Some(task_id), + solver_id: None, + processing_steps: steps, + }; + } + }; + + // Store tracker + self.transfer_status.insert( + task_id.clone(), + TransferTracker { + transfer_id: task_id.clone(), + status: "pending_tee_execution".to_string(), + solver_id: Some(solver_id.clone()), + event_id: None, + processing_steps: steps.clone(), + created_at: now, + }, + ); + + // Spawn TEE execution + self.tee_executor.spawn_tee_task(task_id.clone(), solver_id.clone(), solver_task); + + info!(task_id = %task_id, solver_id = %solver_id, "Task submitted (TEE execution spawned)"); + + setu_rpc::SubmitTaskResponse { + success: true, + message: "Task submitted, awaiting TEE execution".to_string(), + task_id: Some(task_id), + solver_id: Some(solver_id), + processing_steps: steps, + } + } + // ============================================ // Event Processing (delegates to EventHandler) // ============================================ @@ -645,6 +784,10 @@ impl setu_api::ValidatorService for ValidatorNetworkService { self.submit_transfer(request).await } + async fn submit_task(&self, request: setu_rpc::SubmitTaskRequest) -> setu_rpc::SubmitTaskResponse { + self.submit_task(request).await + } + fn get_transfer_status(&self, transfer_id: &str) -> GetTransferStatusResponse { self.get_transfer_status(transfer_id) } diff --git a/setu-validator/src/network/tee_executor.rs b/setu-validator/src/network/tee_executor.rs index d2024ef..d68e07a 100644 --- a/setu-validator/src/network/tee_executor.rs +++ b/setu-validator/src/network/tee_executor.rs @@ -16,9 +16,12 @@ use super::types::*; use super::solver_client::{ExecuteTaskRequest, ExecuteTaskResponse}; use crate::ConsensusValidator; +use crate::TaskPreparer; +use consensus::pocw::{flux_burn, power}; use dashmap::DashMap; use parking_lot::RwLock; use setu_types::event::Event; +use setu_types::pocw::{EventMetrics, PoCWConfig}; use setu_types::task::SolverTask; use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::Arc; @@ -46,6 +49,10 @@ pub struct TeeExecutor { semaphore: Arc, /// Count of pending TEE tasks pending_count: Arc, + /// Task preparer (for applying state changes after execution) + task_preparer: Option>, + /// PoCW configuration (for computing EventMetrics after TEE execution) + pocw_config: Option, } impl TeeExecutor { @@ -70,9 +77,23 @@ impl TeeExecutor { validator_id, semaphore: Arc::new(Semaphore::new(max_concurrent)), pending_count: Arc::new(AtomicU64::new(0)), + task_preparer: None, + pocw_config: None, } } + /// Set the task preparer for applying state changes after TEE execution + pub fn with_task_preparer(mut self, task_preparer: Arc) -> Self { + self.task_preparer = Some(task_preparer); + self + } + + /// Set PoCW config for computing EventMetrics after TEE execution + pub fn with_pocw_config(mut self, config: PoCWConfig) -> Self { + self.pocw_config = Some(config); + self + } + /// Spawn an async TEE task (non-blocking) /// /// This is the key optimization for TPS: @@ -95,6 +116,8 @@ impl TeeExecutor { let dag_events = Arc::clone(&self.dag_events); let consensus = self.consensus.clone(); let validator_id = self.validator_id.clone(); + let task_preparer = self.task_preparer.clone(); + let pocw_config = self.pocw_config; tokio::spawn(async move { Self::execute_tee_task_internal( @@ -110,6 +133,8 @@ impl TeeExecutor { dag_events, consensus, validator_id, + task_preparer, + pocw_config, ) .await; }); @@ -132,6 +157,8 @@ impl TeeExecutor { dag_events: Arc>>, consensus: Option>, validator_id: String, + task_preparer: Option>, + pocw_config: Option, ) { let task_id_hex = hex::encode(&task.task_id[..8]); @@ -175,6 +202,10 @@ impl TeeExecutor { let mut event = task.event.clone(); let event_id = event.id.clone(); + // Capture metrics inputs before task is moved into request + let read_count = task.read_set.len(); + let value_transferred = event.transfer.as_ref().map(|t| t.amount).unwrap_or(0); + // 3. Create HTTP request let request = ExecuteTaskRequest { solver_task: task, @@ -213,10 +244,56 @@ impl TeeExecutor { .collect(), }; - event.set_execution_result(execution_result); + event.set_execution_result(execution_result.clone()); event.executed_by = Some(solver_id.clone()); event.status = setu_types::event::EventStatus::Executed; + // 5b. Compute EventMetrics (PoCW Stages 1 & 2) + if let Some(ref config) = pocw_config { + if config.enabled { + let mut metrics = EventMetrics { + solver_id: solver_id.clone(), + compute_time_us: result_dto.execution_time_us, + gas_used: result_dto.gas_used, + write_count: result_dto.state_changes.len(), + read_count, + value_transferred, + dag_depth: 0, // Not in DAG yet; updated at fold time + flux_burn: 0, + power_delta: 0, + }; + + // Stage 1: Flux burn + metrics.flux_burn = if event.event_type == setu_types::event::EventType::Transfer { + flux_burn::calculate_transfer_burn(config) + } else { + flux_burn::calculate_task_burn(&metrics, config) + }; + + // Stage 2: Power drain + metrics.power_delta = if event.event_type == setu_types::event::EventType::Transfer { + power::calculate_transfer_power_drain(config) + } else { + power::calculate_task_power_drain(&metrics, config) + }; + + debug!( + event_id = %&event_id[..20.min(event_id.len())], + flux_burn = metrics.flux_burn, + power_delta = metrics.power_delta, + gas_used = metrics.gas_used, + "EventMetrics computed" + ); + + event.event_metrics = Some(metrics); + } + } + + // 5c. Apply state changes to global state immediately + if let Some(ref tp) = task_preparer { + tp.apply_state_changes(&execution_result.state_changes); + } + // 6. Submit to consensus (if enabled) if let Some(ref consensus_validator) = consensus { match consensus_validator.submit_event(event.clone()).await { diff --git a/setu-validator/src/network/types.rs b/setu-validator/src/network/types.rs index 2853343..1d09c91 100644 --- a/setu-validator/src/network/types.rs +++ b/setu-validator/src/network/types.rs @@ -26,6 +26,8 @@ pub struct NetworkServiceConfig { /// Transaction fee applied to each FluxTransfer (from PoCWConfig, 0 if disabled). /// This fee is burned (destroyed) — it is not paid to any recipient. pub transfer_fee: u64, + /// PoCW configuration (passed to TeeExecutor for EventMetrics computation) + pub pocw_config: Option, } impl Default for NetworkServiceConfig { @@ -34,6 +36,7 @@ impl Default for NetworkServiceConfig { http_listen_addr: "127.0.0.1:8080".parse().unwrap(), p2p_listen_addr: "127.0.0.1:9000".parse().unwrap(), transfer_fee: 0, + pocw_config: None, } } } diff --git a/setu-validator/src/task_preparer.rs b/setu-validator/src/task_preparer.rs index ab3fa08..652f76e 100644 --- a/setu-validator/src/task_preparer.rs +++ b/setu-validator/src/task_preparer.rs @@ -38,6 +38,7 @@ use setu_types::task::{ }; use setu_types::{Event, EventType, SubnetId, ObjectId}; use setu_types::event::VLCSnapshot; +use setu_types::registration::TaskSubmission; use std::sync::Arc; use tracing::{debug, info, warn}; @@ -113,7 +114,23 @@ impl TaskPreparer { let state_provider = Arc::new(MerkleStateProvider::new(state_manager)); Self::new(validator_id, state_provider) } - + + /// Apply TEE execution state changes to the global Merkle state. + /// Called after successful TEE execution so balance queries reflect updates + /// without waiting for fold. + pub fn apply_state_changes(&self, state_changes: &[setu_types::event::StateChange]) { + for sc in state_changes { + if let Some(ref new_value) = sc.new_value { + // key format: "coin:{hex_object_id}" — parse the object_id + if let Some(hex_id) = sc.key.strip_prefix("coin:") { + if let Ok(object_id) = setu_types::ObjectId::from_hex(hex_id) { + self.state_provider.apply_state_change(*object_id.as_bytes(), new_value.clone()); + } + } + } + } + } + /// Prepare a SolverTask from a Transfer request /// /// This is the main entry point for task preparation. @@ -223,6 +240,59 @@ impl TaskPreparer { Ok(task) } + /// Prepare a SolverTask from a TaskSubmission + /// + /// Unlike transfers, task submissions don't require coin selection. + /// The task payload is passed through to the solver/TEE for execution. + pub fn prepare_task_submission( + &self, + task_submission: &TaskSubmission, + subnet_id: SubnetId, + parent_ids: Vec, + vlc_snapshot: VLCSnapshot, + ) -> Result { + debug!( + task_id = %task_submission.task_id, + task_type = %task_submission.task_type, + submitter = %task_submission.submitter, + "Preparing SolverTask for task submission" + ); + + // Create TaskSubmit event + let event = Event::task_submit( + task_submission.clone(), + parent_ids, + vlc_snapshot, + self.validator_id.clone(), + ); + + // Get pre-state root + let pre_state_root = self.state_provider.get_state_root(); + + // Generate task_id + let task_id = SolverTask::generate_task_id(&event, &pre_state_root); + + // Build SolverTask with empty resolved_inputs (no coin selection needed) + let resolved_inputs = ResolvedInputs::new(); + + let task = SolverTask::new( + task_id, + event, + resolved_inputs, + pre_state_root, + subnet_id, + ) + .with_gas_budget(GasBudget::default()); + + info!( + task_submission_id = %task_submission.task_id, + solver_task_id = %hex::encode(&task_id[..8]), + "SolverTask prepared for task submission" + ); + + Ok(task) + } + /// Select the best coin for a transfer /// /// Strategy: Select the smallest coin that can cover the transfer amount @@ -368,6 +438,7 @@ mod tests { use super::*; use setu_types::{Transfer, TransferType}; use setu_types::task::OperationType; + use setu_types::registration::TaskSubmission; fn create_test_transfer() -> Transfer { Transfer::new("test-tx-1", "alice", "bob", 100) @@ -464,4 +535,42 @@ mod tests { _ => panic!("Expected InsufficientBalance error"), } } + + #[test] + fn test_prepare_task_submission() { + use setu_vlc::VectorClock; + + let preparer = TaskPreparer::new_for_testing("validator-1".to_string()); + + let task_sub = TaskSubmission { + task_id: "task-001".to_string(), + task_type: "compute".to_string(), + submitter: "alice".to_string(), + payload: vec![1, 2, 3], + }; + + let vlc = VLCSnapshot { + vector_clock: VectorClock::new(), + logical_time: 1, + physical_time: 1000, + }; + let result = preparer.prepare_task_submission( + &task_sub, + SubnetId::ROOT, + vec![], + vlc, + ); + + assert!(result.is_ok()); + let solver_task = result.unwrap(); + + // Event should be TaskSubmit type + assert_eq!(solver_task.event.event_type, EventType::TaskSubmit); + + // Resolved inputs should be empty (no coin selection for tasks) + assert!(solver_task.resolved_inputs.input_objects.is_empty()); + + // Task ID should be deterministic and non-zero + assert!(!solver_task.task_id.iter().all(|&b| b == 0)); + } } diff --git a/setu-validator/tests/persistence_recovery_test.rs b/setu-validator/tests/persistence_recovery_test.rs index 3e942b8..4b681e0 100644 --- a/setu-validator/tests/persistence_recovery_test.rs +++ b/setu-validator/tests/persistence_recovery_test.rs @@ -56,6 +56,7 @@ fn create_test_event(id: &str, creator: &str, parent_ids: Vec) -> Event status: EventStatus::Pending, vlc_snapshot: VLCSnapshot::new(), executed_by: None, + event_metrics: None, } } diff --git a/storage/src/state/provider.rs b/storage/src/state/provider.rs index e610de2..4cc1892 100644 --- a/storage/src/state/provider.rs +++ b/storage/src/state/provider.rs @@ -162,6 +162,10 @@ pub trait StateProvider: Send + Sync { let proof = self.get_merkle_proof(object_id)?; Some((data, proof)) } + + /// Apply a state change (object_id bytes + new serialized value) to the global state. + /// Default is no-op for providers that don't support writes. + fn apply_state_change(&self, _object_id: [u8; 32], _value: Vec) {} } // ============================================================================ @@ -318,6 +322,11 @@ impl StateProvider for MerkleStateProvider { let tracker = self.modification_tracker.read().unwrap(); tracker.get(object_id.as_bytes()).cloned() } + + fn apply_state_change(&self, object_id: [u8; 32], value: Vec) { + let mut manager = self.state_manager.write().unwrap(); + manager.upsert_object(self.default_subnet.clone(), object_id, value); + } } // ============================================================================ diff --git a/types/src/event.rs b/types/src/event.rs index 832e29e..653a426 100644 --- a/types/src/event.rs +++ b/types/src/event.rs @@ -262,6 +262,10 @@ pub struct Event { /// Solver that executed this event (creator is the validator, not the solver) #[serde(default)] pub executed_by: Option, + + /// Per-event economic metrics (computed by validator after TEE execution) + #[serde(default)] + pub event_metrics: Option, } impl Event { @@ -299,6 +303,7 @@ impl Event { execution_result: None, timestamp, executed_by: None, + event_metrics: None, } } diff --git a/types/src/lib.rs b/types/src/lib.rs index 5a29d94..51734e6 100644 --- a/types/src/lib.rs +++ b/types/src/lib.rs @@ -81,7 +81,7 @@ pub use merkle::{ pub use account_view::AccountView; // Task types for Validator → Solver communication -pub use pocw::{PoCWConfig, SolverReward, FoldEconomics}; +pub use pocw::{PoCWConfig, EventMetrics, SolverReward, FoldEconomics, SolverPowerState, EmissionState}; pub use task::{ SolverTask, ResolvedInputs, OperationType, ResolvedObject, diff --git a/types/src/pocw.rs b/types/src/pocw.rs index 39baece..91b80a1 100644 --- a/types/src/pocw.rs +++ b/types/src/pocw.rs @@ -1,19 +1,21 @@ -//! PoCW (Proof of Causal Work) and Flux economic types for FluxTransfer +//! PoCW (Proof of Causal Work) and Flux economic types. //! -//! Defines configuration and economic structures for the FluxTransfer path: -//! fixed burn fee, flat power drain, and optional nominal solver reward. +//! Covers both FluxTransfer (flat fee) and TaskSubmit (formula-based burn) +//! paths through the 5-stage Flux pipeline. use serde::{Deserialize, Serialize}; /// Economic configuration for the Flux system. /// -/// Controls burn fees, power drain, and solver rewards for FluxTransfer transactions. -/// Use `enabled` to toggle the entire economic subsystem at runtime. +/// Controls burn fees, power drain, and solver rewards. Transfer fields are +/// flat values; task fields use the alpha/beta/gamma formula coefficients. #[derive(Debug, Clone, Copy, Serialize, Deserialize)] pub struct PoCWConfig { /// Master toggle for PoCW economics pub enabled: bool, - /// Fee per FluxTransfer (burned). Burn mechanism subject to revision. + + // -- Transfer (flat) -- + /// Fee per FluxTransfer (burned) pub transfer_fee: u64, /// Flat power drain per FluxTransfer pub transfer_power_drain: u64, @@ -21,8 +23,75 @@ pub struct PoCWConfig { pub solver_transfer_reward_enabled: bool, /// Nominal Flux reward per FluxTransfer processed by a solver pub solver_transfer_reward: u64, + + // -- Task burn (Stage 1): burn = alpha*C + beta*R + gamma*S -- + /// Complexity weight + #[serde(default = "default_alpha")] + pub alpha: f64, + /// Risk weight + #[serde(default = "default_beta")] + pub beta: f64, + /// Structural tension weight + #[serde(default = "default_gamma")] + pub gamma: f64, + + // -- Task power (Stage 2) -- + /// Power drain rate for solver events: deltaP = rate * gas_used + #[serde(default = "default_solver_power_rate")] + pub solver_power_rate: f64, + /// Initial power budget per solver + #[serde(default = "default_initial_power_budget")] + pub initial_power_budget: u64, + + // -- Minting (Stage 3) -- + /// Minting coefficient: FluxMinted = kappa * deltaPower_total + #[serde(default = "default_kappa")] + pub kappa: f64, + + // -- PoCW distribution (Stage 4) -- + /// Weight for distance scoring (proximity to causal chain) + #[serde(default = "default_pocw_distance_weight")] + pub pocw_distance_weight: f64, + /// Weight for necessity scoring (relevant events / total events) + #[serde(default = "default_pocw_necessity_weight")] + pub pocw_necessity_weight: f64, + /// Weight for contribution scoring (agent depth / total depth) + #[serde(default = "default_pocw_contribution_weight")] + pub pocw_contribution_weight: f64, + + // -- Emission adjustment (Stage 5) -- + /// Target Flux minted per fold + #[serde(default = "default_target_velocity")] + pub target_velocity: f64, + /// Number of recent folds to observe + #[serde(default = "default_observation_window")] + pub observation_window: usize, + /// Dampening factor for kappa adjustment + #[serde(default = "default_dampening")] + pub dampening: f64, + /// Minimum kappa + #[serde(default = "default_kappa_min")] + pub kappa_min: f64, + /// Maximum kappa + #[serde(default = "default_kappa_max")] + pub kappa_max: f64, } +fn default_alpha() -> f64 { 0.40 } +fn default_beta() -> f64 { 0.35 } +fn default_gamma() -> f64 { 0.25 } +fn default_solver_power_rate() -> f64 { 0.001 } +fn default_initial_power_budget() -> u64 { 21_000_000 } +fn default_kappa() -> f64 { 1.0 } +fn default_pocw_distance_weight() -> f64 { 0.35 } +fn default_pocw_necessity_weight() -> f64 { 0.35 } +fn default_pocw_contribution_weight() -> f64 { 0.30 } +fn default_target_velocity() -> f64 { 1000.0 } +fn default_observation_window() -> usize { 10 } +fn default_dampening() -> f64 { 0.5 } +fn default_kappa_min() -> f64 { 0.1 } +fn default_kappa_max() -> f64 { 10.0 } + impl Default for PoCWConfig { fn default() -> Self { Self { @@ -31,20 +100,129 @@ impl Default for PoCWConfig { transfer_power_drain: 1, solver_transfer_reward_enabled: false, solver_transfer_reward: 1, + alpha: default_alpha(), + beta: default_beta(), + gamma: default_gamma(), + solver_power_rate: default_solver_power_rate(), + initial_power_budget: default_initial_power_budget(), + kappa: default_kappa(), + pocw_distance_weight: default_pocw_distance_weight(), + pocw_necessity_weight: default_pocw_necessity_weight(), + pocw_contribution_weight: default_pocw_contribution_weight(), + target_velocity: default_target_velocity(), + observation_window: default_observation_window(), + dampening: default_dampening(), + kappa_min: default_kappa_min(), + kappa_max: default_kappa_max(), + } + } +} + +// ========== Event-level metrics ========== + +/// Metrics computed by the validator after TEE execution. +/// All fields derived from TeeExecutionResult + SolverTask input. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct EventMetrics { + /// Solver that executed this event + pub solver_id: String, + /// Execution time in microseconds + pub compute_time_us: u64, + /// Gas consumed + pub gas_used: u64, + /// Number of state mutations + pub write_count: usize, + /// Number of read_set entries + pub read_count: usize, + /// Value transferred (0 for non-transfer events) + pub value_transferred: u64, + /// DAG depth of the event at insertion time + pub dag_depth: u64, + /// Flux burn computed for this event + pub flux_burn: u64, + /// Power consumed by this event + pub power_delta: u64, +} + +// ========== Per-solver power state ========== + +/// Tracks a solver's power budget across folds. +/// Stored in ROOT SMT under "power:{solver_id}". +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct SolverPowerState { + pub solver_id: String, + pub remaining_power: u64, + pub total_consumed: u64, + pub event_count: u64, +} + +impl SolverPowerState { + /// Create a new solver power state with the given budget + pub fn new(solver_id: String, initial_budget: u64) -> Self { + Self { + solver_id, + remaining_power: initial_budget, + total_consumed: 0, + event_count: 0, } } } -/// Solver reward record for a single solver within a fold +// ========== Emission state ========== + +/// Tracks kappa and velocity history for Stage 5 emission adjustment. +/// Stored in ROOT SMT under "emission:" keys. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct EmissionState { + pub kappa: f64, + pub velocity_history: Vec, +} + +impl EmissionState { + pub fn new(initial_kappa: f64) -> Self { + Self { + kappa: initial_kappa, + velocity_history: Vec::new(), + } + } +} + +// ========== Solver reward ========== + +/// Solver reward record for a single solver within a fold. +/// +/// Scoring signals from the PoCW spec: +/// - Distance: proximity to winning agent's causal chain (1 / (1 + Distance_i)) +/// - Necessity: fraction of agent's events leading to accepted output +/// - Contribution: share of causal depth in accepted chain (agent_depth / total_depth) +/// - w_i (weight): α*Distance + β*Necessity + γ*Contribution, normalized #[derive(Debug, Clone, Serialize, Deserialize)] pub struct SolverReward { pub solver_id: String, - /// Number of FluxTransfer events this solver processed in the fold + /// Number of FluxTransfer events this solver processed + #[serde(default)] pub transfer_count: u64, + /// Number of TaskSubmit events this solver processed + #[serde(default)] + pub task_count: u64, + /// Distance score: proximity to causal chain (0.0-1.0) + #[serde(default)] + pub distance_score: f64, + /// Necessity score: relevant_events / total_events (0.0-1.0) + #[serde(default)] + pub necessity_score: f64, + /// Contribution score: agent_depth / total_depth (0.0-1.0) + #[serde(default)] + pub contribution_score: f64, + /// Normalized weight w_i = α*Distance + β*Necessity + γ*Contribution + #[serde(default)] + pub weight: f64, /// Flux reward for this solver pub flux_reward: u64, } +// ========== Fold economics ========== + /// Fold-level economic summary #[derive(Debug, Clone, Serialize, Deserialize)] pub struct FoldEconomics { @@ -54,8 +232,17 @@ pub struct FoldEconomics { pub total_flux_burned: u64, /// Total power drained across all events in this fold pub total_power_consumed: u64, - /// Total Flux minted as solver rewards in this fold + /// Flux minted this fold (kappa * total_power) + #[serde(default)] + pub flux_minted: u64, + /// Total Flux distributed as solver rewards pub total_solver_rewards: u64, + /// Kappa before this fold + #[serde(default = "default_kappa")] + pub kappa_before: f64, + /// Kappa after emission adjustment + #[serde(default = "default_kappa")] + pub kappa_after: f64, /// Per-solver reward breakdown pub solver_rewards: Vec, } @@ -65,42 +252,150 @@ mod tests { use super::*; #[test] - fn test_defaults() { + fn test_pocw_config_defaults() { let config = PoCWConfig::default(); assert!(!config.enabled); assert_eq!(config.transfer_fee, 21_000); assert_eq!(config.transfer_power_drain, 1); assert!(!config.solver_transfer_reward_enabled); assert_eq!(config.solver_transfer_reward, 1); + assert!((config.alpha - 0.40).abs() < f64::EPSILON); + assert!((config.beta - 0.35).abs() < f64::EPSILON); + assert!((config.gamma - 0.25).abs() < f64::EPSILON); + assert!((config.solver_power_rate - 0.001).abs() < f64::EPSILON); + assert_eq!(config.initial_power_budget, 21_000_000); + assert!((config.kappa - 1.0).abs() < f64::EPSILON); + assert!((config.pocw_distance_weight - 0.35).abs() < f64::EPSILON); + assert!((config.pocw_necessity_weight - 0.35).abs() < f64::EPSILON); + assert!((config.pocw_contribution_weight - 0.30).abs() < f64::EPSILON); + assert!((config.target_velocity - 1000.0).abs() < f64::EPSILON); + assert_eq!(config.observation_window, 10); + assert!((config.dampening - 0.5).abs() < f64::EPSILON); + assert!((config.kappa_min - 0.1).abs() < f64::EPSILON); + assert!((config.kappa_max - 10.0).abs() < f64::EPSILON); } #[test] - fn test_serde_roundtrip() { - let config = PoCWConfig::default(); + fn test_pocw_config_serde_roundtrip() { + let config = PoCWConfig { + enabled: true, + alpha: 0.5, + beta: 0.3, + gamma: 0.2, + kappa: 2.5, + ..Default::default() + }; let json = serde_json::to_string(&config).unwrap(); - let deserialized: PoCWConfig = serde_json::from_str(&json).unwrap(); + let de: PoCWConfig = serde_json::from_str(&json).unwrap(); + assert!(de.enabled); + assert!((de.alpha - 0.5).abs() < f64::EPSILON); + assert!((de.beta - 0.3).abs() < f64::EPSILON); + assert!((de.gamma - 0.2).abs() < f64::EPSILON); + assert!((de.kappa - 2.5).abs() < f64::EPSILON); + assert_eq!(de.transfer_fee, 21_000); + } - assert_eq!(config.enabled, deserialized.enabled); - assert_eq!(config.transfer_fee, deserialized.transfer_fee); - assert_eq!(config.transfer_power_drain, deserialized.transfer_power_drain); - assert_eq!(config.solver_transfer_reward_enabled, deserialized.solver_transfer_reward_enabled); - assert_eq!(config.solver_transfer_reward, deserialized.solver_transfer_reward); + #[test] + fn test_pocw_config_backward_compat_deserialize() { + // JSON with only the original PR #15 fields (no new fields) + let json = r#"{ + "enabled": true, + "transfer_fee": 21000, + "transfer_power_drain": 1, + "solver_transfer_reward_enabled": false, + "solver_transfer_reward": 1 + }"#; + let config: PoCWConfig = serde_json::from_str(json).unwrap(); + assert!(config.enabled); + assert_eq!(config.transfer_fee, 21_000); + // New fields should get defaults + assert!((config.alpha - 0.40).abs() < f64::EPSILON); + assert!((config.kappa - 1.0).abs() < f64::EPSILON); + assert_eq!(config.observation_window, 10); } #[test] - fn test_custom_config() { - let config = PoCWConfig { - enabled: true, - solver_transfer_reward_enabled: true, - solver_transfer_reward: 5, - transfer_fee: 10_000, - ..Default::default() + fn test_event_metrics_serde() { + let metrics = EventMetrics { + solver_id: "solver-1".to_string(), + compute_time_us: 500, + gas_used: 1000, + write_count: 3, + read_count: 2, + value_transferred: 0, + dag_depth: 5, + flux_burn: 12_000, + power_delta: 1, }; + let json = serde_json::to_string(&metrics).unwrap(); + let de: EventMetrics = serde_json::from_str(&json).unwrap(); + assert_eq!(de.solver_id, "solver-1"); + assert_eq!(de.gas_used, 1000); + assert_eq!(de.write_count, 3); + assert_eq!(de.flux_burn, 12_000); + } - assert!(config.enabled); - assert!(config.solver_transfer_reward_enabled); - assert_eq!(config.solver_transfer_reward, 5); - assert_eq!(config.transfer_fee, 10_000); - assert_eq!(config.transfer_power_drain, 1); + #[test] + fn test_solver_power_state_serde() { + let state = SolverPowerState { + solver_id: "solver-1".to_string(), + remaining_power: 20_999_000, + total_consumed: 1_000, + event_count: 5, + }; + let json = serde_json::to_string(&state).unwrap(); + let de: SolverPowerState = serde_json::from_str(&json).unwrap(); + assert_eq!(de.remaining_power, 20_999_000); + assert_eq!(de.total_consumed, 1_000); + } + + #[test] + fn test_emission_state_new() { + let state = EmissionState::new(1.0); + assert!((state.kappa - 1.0).abs() < f64::EPSILON); + assert!(state.velocity_history.is_empty()); + } + + #[test] + fn test_emission_state_serde() { + let mut state = EmissionState::new(2.0); + state.velocity_history.push(100); + state.velocity_history.push(200); + let json = serde_json::to_string(&state).unwrap(); + let de: EmissionState = serde_json::from_str(&json).unwrap(); + assert!((de.kappa - 2.0).abs() < f64::EPSILON); + assert_eq!(de.velocity_history, vec![100, 200]); + } + + #[test] + fn test_solver_reward_backward_compat() { + // JSON with only the original fields (no scores) + let json = r#"{ + "solver_id": "s1", + "flux_reward": 100 + }"#; + let reward: SolverReward = serde_json::from_str(json).unwrap(); + assert_eq!(reward.solver_id, "s1"); + assert_eq!(reward.flux_reward, 100); + assert_eq!(reward.transfer_count, 0); + assert_eq!(reward.task_count, 0); + assert!((reward.weight - 0.0).abs() < f64::EPSILON); + } + + #[test] + fn test_fold_economics_backward_compat() { + // JSON without the new kappa/minting fields + let json = r#"{ + "event_count": 3, + "total_flux_burned": 63000, + "total_power_consumed": 3, + "total_solver_rewards": 3, + "solver_rewards": [] + }"#; + let econ: FoldEconomics = serde_json::from_str(json).unwrap(); + assert_eq!(econ.event_count, 3); + assert_eq!(econ.flux_minted, 0); + assert!((econ.kappa_before - 1.0).abs() < f64::EPSILON); + assert!((econ.kappa_after - 1.0).abs() < f64::EPSILON); } }