Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 15 additions & 3 deletions api/src/handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -48,10 +49,13 @@ pub trait ValidatorService: Send + Sync {

/// Submit transfer
fn submit_transfer(&self, request: SubmitTransferRequest) -> impl std::future::Future<Output = SubmitTransferResponse> + Send;


/// Submit task for solver execution
fn submit_task(&self, request: SubmitTaskRequest) -> impl std::future::Future<Output = SubmitTaskResponse> + 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<Output = SubmitEventResponse> + Send;

Expand Down Expand Up @@ -132,6 +136,14 @@ pub async fn http_submit_transfer<S: ValidatorService>(
Json(service.submit_transfer(request).await)
}

/// Submit a task for solver execution
pub async fn http_submit_task<S: ValidatorService>(
State(service): State<Arc<S>>,
Json(request): Json<SubmitTaskRequest>,
) -> Json<SubmitTaskResponse> {
Json(service.submit_task(request).await)
}

/// Get transfer status
pub async fn http_get_transfer_status<S: ValidatorService>(
State(service): State<Arc<S>>,
Expand Down
40 changes: 23 additions & 17 deletions consensus/src/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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");

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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");
}
Expand Down Expand Up @@ -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");

Expand Down
36 changes: 23 additions & 13 deletions consensus/src/folder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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
Expand All @@ -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,
Expand Down Expand Up @@ -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);
}
}
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 },
],
};

Expand All @@ -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);
Expand Down
Loading