use super::*;
use crate::{BroadcastResult, EffectObservation, EffectStatus, SignedAttempt};
use alloy::primitives::{keccak256, Bytes, B256, U256};
use newton_core::newton_prover_task_manager::{
INewtonPolicy, INewtonPolicyClient,
INewtonProverTaskManager::{Task, TaskResponse},
NewtonMessage,
};
use newton_submission_protocol::{ExecutionId, ExecutionStatus};
use newton_submission_service::{AdmissionConfig, ChainConfig, ExecutableIntent, RuntimeHealth, SqliteConfig};
use newton_task_submission::{
contract_response_hash, derive_idempotency_key, AdmissionOutcome, EventCursor, SubmissionId, SubmissionState,
TaskChainPolicy, TaskOperation, TaskPlanner, TaskPlannerConfig, TaskPlannerObserver, TaskSubmissionPayload,
TaskSubmissionRequest,
};
use std::{collections::HashMap, sync::Mutex};
#[derive(Debug)]
struct MockBackend {
address: Address,
nonce_error: bool,
permanent_nonce_error: bool,
nonce_counts: Mutex<(u64, u64)>,
balance: Mutex<U256>,
receipt_pending: bool,
receipt_inconclusive: bool,
broadcasts: Mutex<Vec<Bytes>>,
observation_scopes: Mutex<Vec<ObservationScope>>,
}
#[async_trait::async_trait]
impl ChainBackend for MockBackend {
fn address(&self) -> Address {
self.address
}
async fn latest_transaction_count(&self) -> Result<u64, BackendError> {
if self.permanent_nonce_error {
Err(BackendError::Policy("wrong chain configuration".to_string()))
} else if self.nonce_error {
Err(BackendError::RpcTransient("temporary provider outage".to_string()))
} else {
let counts = self.nonce_counts.lock().expect("nonce count lock");
Ok(counts.0)
}
}
async fn pending_transaction_count(&self) -> Result<u64, BackendError> {
if self.permanent_nonce_error {
Err(BackendError::Policy("wrong chain configuration".to_string()))
} else if self.nonce_error {
Err(BackendError::RpcTransient("temporary provider outage".to_string()))
} else {
let counts = self.nonce_counts.lock().expect("nonce count lock");
Ok(counts.1)
}
}
async fn balance(&self) -> Result<U256, BackendError> {
if self.nonce_error {
Err(BackendError::RpcTransient("temporary provider outage".to_string()))
} else {
Ok(*self.balance.lock().expect("balance lock"))
}
}
async fn prepare(
&self,
_intent: &ExecutableIntent,
nonce: u64,
_replacement_fees: Option<(u128, u128)>,
) -> Result<SignedAttempt, BackendError> {
let raw_transaction = Bytes::from_static(b"deterministic-signed-transaction");
Ok(SignedAttempt {
nonce,
gas_limit: 100_000,
max_fee_per_gas: 10,
max_priority_fee_per_gas: 1,
transaction_hash: keccak256(&raw_transaction),
raw_transaction,
})
}
async fn prepare_cancellation(&self, _nonce: u64, _fees: (u128, u128)) -> Result<SignedAttempt, BackendError> {
Err(BackendError::Policy("unexpected cancellation".to_string()))
}
async fn broadcast(&self, raw_transaction: &Bytes) -> Vec<BroadcastResult> {
self.broadcasts
.lock()
.expect("broadcast lock")
.push(raw_transaction.clone());
vec![BroadcastResult {
provider: "mock".to_string(),
result: Ok(keccak256(raw_transaction)),
}]
}
async fn observe(
&self,
transaction_hash: B256,
_receipt_provider: Option<&str>,
scope: ObservationScope,
) -> Result<TransactionObservation, BackendError> {
self.observation_scopes
.lock()
.expect("observation scope lock")
.push(scope);
if self.receipt_pending {
if self.receipt_inconclusive {
return Ok(TransactionObservation::PendingInconclusive {
pending_providers: 1,
failed_providers: 1,
});
}
return Ok(TransactionObservation::Pending);
}
Ok(TransactionObservation::Mined {
provider: "mock".to_string(),
block_number: 100,
block_hash: B256::repeat_byte(8),
confirmations: 2,
receipt: transaction_hash.to_vec(),
})
}
async fn classify_onchain_effects(
&self,
intent: &ExecutableIntent,
) -> Result<Vec<EffectObservation>, BackendError> {
let items = match intent {
ExecutableIntent::BatchCreateAndRespond { items, .. } | ExecutableIntent::BatchRespond { items, .. } => {
items
}
ExecutableIntent::ConfirmGlobalTableRoot { .. } | ExecutableIntent::UpdateOperatorTable { .. } => {
return Err(BackendError::Policy(
"task test backend cannot classify transporter intent".to_string(),
));
}
};
Ok(items
.iter()
.map(|item| EffectObservation {
submission_id: item.submission_id,
status: EffectStatus::Verified,
})
.collect())
}
}
fn chain() -> ChainConfig {
ChainConfig {
chain_id: 31_337,
batch_task_manager: Address::repeat_byte(9),
operator_table_updater: Address::ZERO,
rpc_urls: vec!["http://127.0.0.1:8545".to_string()],
finality_confirmations: 1,
max_create_batch_size: 8,
max_respond_batch_size: 8,
max_task_age_secs: 300,
task_response_window_blocks: 30,
block_time_ms: 1_000,
max_effect_retries: 3,
}
}
fn admission() -> AdmissionConfig {
AdmissionConfig {
max_body_bytes: 2 * 1024 * 1024,
max_live_submissions_per_chain: 10_000,
max_database_bytes: 8 * 1024 * 1024 * 1024,
}
}
fn request() -> TaskSubmissionRequest {
let task_id = B256::repeat_byte(1);
let policy_id = B256::repeat_byte(5);
let intent = NewtonMessage::Intent {
from: Address::repeat_byte(2),
to: Address::repeat_byte(3),
value: U256::ZERO,
data: Bytes::from_static(b"intent"),
chainId: U256::from(31_337),
functionSignature: Bytes::from_static(b"\x01\x02\x03\x04"),
};
let task = Task {
taskId: task_id,
policyClient: Address::repeat_byte(4),
policyId: policy_id,
policyRevision: 1,
taskCreatedBlock: 10,
quorumThresholdPercentage: 67,
intent: intent.clone(),
intentSignature: Bytes::from_static(b"intent-signature"),
policies: vec![INewtonPolicyClient::PolicySpec {
policy: Address::repeat_byte(6),
config: INewtonPolicy::PolicyConfig {
policyParams: Bytes::from_static(b"policy-params"),
expireAfter: 500,
},
}],
wasmArgs: vec![Bytes::from_static(b"component-input-0")],
quorumNumbers: Bytes::from_static(b"\x00"),
initializationTimestamp: U256::from(11),
};
let task_response = TaskResponse {
taskId: task_id,
policyClient: task.policyClient,
policyId: policy_id,
intent,
intentSignature: task.intentSignature.clone(),
policyTaskData: vec![NewtonMessage::PolicyTaskData {
policyId: policy_id,
policyAddress: Address::repeat_byte(6),
policy: Bytes::from_static(b"rego-policy"),
policyData: vec![NewtonMessage::PolicyData {
wasmArgs: Bytes::from_static(b"component-input-0"),
data: Bytes::from_static(b"oracle-output"),
expireBlock: 500,
}],
}],
allowed: true,
initializationTimestamp: U256::from(11),
};
let digest = contract_response_hash(&task_response);
let operation = TaskOperation::CombinedCreateAndRespond;
TaskSubmissionRequest {
producer_id: "gateway".to_string(),
idempotency_key: derive_idempotency_key(b"gateway", 31_337, operation, task_id, digest),
payload: TaskSubmissionPayload {
chain_id: 31_337,
operation,
task_id,
task_response_digest: digest,
task,
task_response,
signature_data: Bytes::from_static(b"bls-signature"),
attestation_data: Bytes::new(),
requested_at: 100,
},
}
}
async fn reserved_job(
backend: Arc<dyn ChainBackend>,
) -> (
tempfile::TempDir,
Arc<Store>,
ManagedSigner,
newton_submission_service::JobRecord,
SubmissionId,
sqlx::SqlitePool,
) {
let directory = tempfile::tempdir().expect("temporary directory");
let path = directory.path().join("submitter.sqlite3");
let store = Arc::new(
Store::open(
&SqliteConfig {
path: path.clone(),
busy_timeout_ms: 5_000,
max_connections: 2,
},
admission(),
)
.await
.expect("store"),
);
let submission_id = match store.admit_task(&request(), &task_policy()).await.expect("admission") {
AdmissionOutcome::Accepted(resource) => resource.submission_id,
other => panic!("expected new admission, got {other:?}"),
};
plan_task(&store).await.expect("job");
let signer = ManagedSigner {
signer_id: SignerId::new("task-test").expect("signer id"),
chain_id: 31_337,
role: SignerRole::Task,
block_time: Duration::from_millis(1),
backend,
};
store
.register_signers(&[signer.record()])
.await
.expect("register signer");
let job = store
.reserve_job(&signer.record())
.await
.expect("reserve")
.expect("job");
let pool = sqlx::sqlite::SqlitePoolOptions::new()
.max_connections(1)
.connect(&format!("sqlite://{}", path.display()))
.await
.expect("open test inspection pool");
(directory, store, signer, job, submission_id, pool)
}
#[derive(Debug, Default)]
struct TestPlannerObserver;
impl TaskPlannerObserver for TestPlannerObserver {
fn heartbeat(&self) {}
fn failed(&self, _error: String) {}
fn planned(&self, _chain_id: u64, _operation: TaskOperation) {}
}
fn task_policy() -> TaskChainPolicy {
let chain = chain();
TaskChainPolicy {
chain_id: chain.chain_id,
max_batch_size: chain.max_respond_batch_size,
max_task_age_secs: chain.max_task_age_secs,
max_effect_retries: chain.max_effect_retries,
}
}
async fn plan_task(store: &Arc<Store>) -> Option<ExecutionId> {
TaskPlanner::new(
store.clone(),
Arc::new(TestPlannerObserver),
TaskPlannerConfig {
batch_interval_ms: 0,
poll_interval_ms: 1,
},
HashMap::from([(chain().chain_id, task_policy())]),
)
.tick_chain(&task_policy())
.await
.expect("task planner tick")
}
fn executor_config() -> ExecutorConfig {
ExecutorConfig {
poll_interval: Duration::from_millis(10),
receipt_poll_interval: Duration::from_millis(1),
receipt_poll_max_interval: Duration::from_millis(8),
watchdog_timeout: Duration::from_secs(30),
cancel_after_bumps: 1,
underfunded_poll_interval: Duration::from_millis(10),
}
}
async fn journal_broadcast_assignment(
backend: &MockBackend,
store: &Store,
signer: &ManagedSigner,
job: &newton_submission_service::JobRecord,
) -> ActiveAssignment {
let signed = backend.prepare(&job.intent, 0, None).await.expect("prepared attempt");
let attempt_number = store
.record_prepared(&PreparedAttempt {
job_id: job.job_id,
signer_id: signer.signer_id.clone(),
signer_address: signer.address(),
chain_id: signer.chain_id,
nonce: signed.nonce,
gas_limit: signed.gas_limit,
max_fee_per_gas: signed.max_fee_per_gas,
max_priority_fee_per_gas: signed.max_priority_fee_per_gas,
raw_transaction: signed.raw_transaction,
transaction_hash: signed.transaction_hash,
kind: AttemptKind::Canonical,
replaces_attempt_number: None,
})
.await
.expect("persist prepared attempt");
let result = vec![BroadcastResult {
provider: "mock".to_string(),
result: Ok(signed.transaction_hash),
}];
store
.record_broadcast(
job.job_id,
attempt_number,
&rmp_serde::to_vec(&result).expect("encode result"),
)
.await
.expect("persist broadcast boundary");
store
.active_assignment(signer.address(), signer.chain_id)
.await
.expect("load assignment")
.expect("active assignment")
}
#[tokio::test]
async fn lifecycle_persists_exact_bytes_before_broadcast_and_classifies_effects() {
let backend = Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: false,
permanent_nonce_error: false,
nonce_counts: Mutex::new((0, 0)),
balance: Mutex::new(U256::MAX),
receipt_pending: false,
receipt_inconclusive: false,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
});
let (_directory, store, signer, job, submission_id, pool) = reserved_job(backend.clone()).await;
assert_eq!(
store
.signer_lifecycle(&signer.signer_id, signer.chain_id)
.await
.expect("reserved signer lifecycle"),
SignerLifecycle::Busy
);
execute_job(&store, &signer, &executor_config(), job, CancellationToken::new())
.await
.expect("execution");
plan_task(&store).await;
let broadcast = {
let broadcasts = backend.broadcasts.lock().expect("broadcast lock");
assert_eq!(broadcasts.len(), 1);
broadcasts[0].clone()
};
let persisted: Vec<u8> = sqlx::query_scalar("SELECT raw_transaction FROM transaction_attempts")
.fetch_one(&pool)
.await
.expect("persisted transaction");
assert_eq!(persisted, broadcast.as_ref());
assert_eq!(
*backend.observation_scopes.lock().expect("observation scope lock"),
vec![ObservationScope::Economical]
);
let resource = store
.submission(submission_id)
.await
.expect("status")
.expect("submission");
assert_eq!(resource.state, SubmissionState::Succeeded);
let events = store
.events_after(EventCursor(0), 100, Duration::ZERO)
.await
.expect("status events");
let terminal = events
.events
.iter()
.find(|event| event.state == SubmissionState::Succeeded)
.expect("terminal success event");
assert_eq!(terminal.task_id, request().payload.task_id);
assert_eq!(terminal.transaction_hash, Some(keccak256(&broadcast)));
store.assert_invariants().await.expect("durable invariants");
}
#[tokio::test]
async fn restart_does_not_rebroadcast_an_attempt_already_journaled_as_broadcast() {
let backend = Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: false,
permanent_nonce_error: false,
nonce_counts: Mutex::new((0, 0)),
balance: Mutex::new(U256::MAX),
receipt_pending: true,
receipt_inconclusive: false,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
});
let (_directory, store, signer, job, _submission_id, _pool) = reserved_job(backend.clone()).await;
let signed = backend.prepare(&job.intent, 0, None).await.expect("prepared attempt");
let attempt_number = store
.record_prepared(&PreparedAttempt {
job_id: job.job_id,
signer_id: signer.signer_id.clone(),
signer_address: signer.address(),
chain_id: signer.chain_id,
nonce: signed.nonce,
gas_limit: signed.gas_limit,
max_fee_per_gas: signed.max_fee_per_gas,
max_priority_fee_per_gas: signed.max_priority_fee_per_gas,
raw_transaction: signed.raw_transaction.clone(),
transaction_hash: signed.transaction_hash,
kind: AttemptKind::Canonical,
replaces_attempt_number: None,
})
.await
.expect("persist prepared attempt");
let initial_broadcast = backend.broadcast(&signed.raw_transaction).await;
store
.record_broadcast(
job.job_id,
attempt_number,
&rmp_serde::to_vec(&initial_broadcast).expect("encode result"),
)
.await
.expect("persist broadcast boundary");
backend.broadcasts.lock().expect("broadcast lock").clear();
let assignment = store
.active_assignment(signer.address(), signer.chain_id)
.await
.expect("load assignment")
.expect("active assignment");
let cancellation = CancellationToken::new();
cancellation.cancel();
recover_prepared_assignment(&store, &signer, &executor_config(), assignment, 0, 0, cancellation)
.await
.expect("restart recovery");
assert!(backend.broadcasts.lock().expect("broadcast lock").is_empty());
}
#[tokio::test]
async fn partial_pending_evidence_allows_same_nonce_replacement_when_nonce_is_unconsumed() {
let backend = Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: false,
permanent_nonce_error: false,
nonce_counts: Mutex::new((0, 0)),
balance: Mutex::new(U256::MAX),
receipt_pending: true,
receipt_inconclusive: true,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
});
let (_directory, store, signer, job, _submission_id, _pool) = reserved_job(backend.clone()).await;
let assignment = journal_broadcast_assignment(&backend, &store, &signer, &job).await;
let mut config = executor_config();
config.watchdog_timeout = Duration::ZERO;
let cancellation = CancellationToken::new();
let (tracking, ()) = tokio::time::timeout(Duration::from_secs(1), async {
tokio::join!(
track_attempt(
&store,
&signer,
&config,
&assignment.job,
assignment.attempts,
cancellation.clone(),
),
async {
loop {
if !backend.broadcasts.lock().expect("broadcast lock").is_empty() {
cancellation.cancel();
break;
}
tokio::task::yield_now().await;
}
}
)
})
.await
.expect("replacement broadcast before timeout");
tracking.expect("tracking cancellation");
assert_eq!(
backend
.observation_scopes
.lock()
.expect("observation scope lock")
.first(),
Some(&ObservationScope::Exhaustive)
);
}
#[tokio::test]
async fn partial_pending_evidence_reconciles_a_consumed_nonce_without_replacement() {
let backend = Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: false,
permanent_nonce_error: false,
nonce_counts: Mutex::new((1, 1)),
balance: Mutex::new(U256::MAX),
receipt_pending: true,
receipt_inconclusive: true,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
});
let (_directory, store, signer, job, _submission_id, _pool) = reserved_job(backend.clone()).await;
let assignment = journal_broadcast_assignment(&backend, &store, &signer, &job).await;
let mut config = executor_config();
config.watchdog_timeout = Duration::ZERO;
let cancellation = CancellationToken::new();
let (tracking, ()) = tokio::join!(
track_attempt(
&store,
&signer,
&config,
&assignment.job,
assignment.attempts,
cancellation.clone(),
),
async {
tokio::time::sleep(Duration::from_millis(25)).await;
cancellation.cancel();
}
);
tracking.expect("tracking cancellation");
assert!(backend.broadcasts.lock().expect("broadcast lock").is_empty());
plan_task(&store).await;
assert_eq!(
store
.submission(_submission_id)
.await
.expect("consumed nonce status")
.expect("consumed nonce submission")
.state,
SubmissionState::Succeeded
);
}
#[tokio::test]
async fn transient_nonce_read_releases_unsigned_reservation_without_quarantine() {
let backend = Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: true,
permanent_nonce_error: false,
nonce_counts: Mutex::new((0, 0)),
balance: Mutex::new(U256::MAX),
receipt_pending: false,
receipt_inconclusive: false,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
});
let (_directory, store, signer, job, submission_id, _pool) = reserved_job(backend.clone()).await;
let error = execute_job(&store, &signer, &executor_config(), job, CancellationToken::new())
.await
.expect_err("provider outage");
assert!(error.is_transient());
assert!(store
.active_assignment(signer.address(), signer.chain_id)
.await
.expect("assignment")
.is_none());
assert!(backend.broadcasts.lock().expect("broadcast lock").is_empty());
let resource = store
.submission(submission_id)
.await
.expect("status")
.expect("submission");
assert_eq!(resource.state, SubmissionState::ReadyForSubmission);
store.assert_invariants().await.expect("durable invariants");
}
#[tokio::test]
async fn underfunded_signer_releases_unsigned_job_and_recovers_after_top_up() {
let backend = Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: false,
permanent_nonce_error: false,
nonce_counts: Mutex::new((0, 0)),
balance: Mutex::new(U256::ZERO),
receipt_pending: false,
receipt_inconclusive: false,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
});
let (_directory, store, signer, job, submission_id, _pool) = reserved_job(backend.clone()).await;
let error = execute_job(&store, &signer, &executor_config(), job, CancellationToken::new())
.await
.expect_err("insufficient signer balance");
assert!(error.is_underfunded());
assert_eq!(
store
.signer_lifecycle(&signer.signer_id, signer.chain_id)
.await
.expect("underfunded lifecycle"),
SignerLifecycle::Underfunded
);
assert!(store
.active_assignment(signer.address(), signer.chain_id)
.await
.expect("released assignment")
.is_none());
assert_eq!(
store
.submission(submission_id)
.await
.expect("status")
.expect("submission")
.state,
SubmissionState::ReadyForSubmission
);
*backend.balance.lock().expect("balance lock") = U256::MAX;
reconcile(&store, &signer, &executor_config(), CancellationToken::new())
.await
.expect("balance recovery");
assert_eq!(
store
.signer_lifecycle(&signer.signer_id, signer.chain_id)
.await
.expect("recovered lifecycle"),
SignerLifecycle::Ready
);
assert!(store
.signer_required_balance(&signer.signer_id, signer.chain_id)
.await
.expect("required balance")
.is_none());
store.assert_invariants().await.expect("durable invariants");
}
#[tokio::test]
async fn consumed_nonce_without_known_receipt_reconciles_contract_effects() {
let backend = Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: false,
permanent_nonce_error: false,
nonce_counts: Mutex::new((0, 0)),
balance: Mutex::new(U256::MAX),
receipt_pending: true,
receipt_inconclusive: false,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
});
let (_directory, store, signer, job, submission_id, _pool) = reserved_job(backend.clone()).await;
let signed = backend.prepare(&job.intent, 0, None).await.expect("prepared attempt");
store
.record_prepared(&PreparedAttempt {
job_id: job.job_id,
signer_id: signer.signer_id.clone(),
signer_address: signer.backend.address(),
chain_id: signer.chain_id,
nonce: signed.nonce,
gas_limit: signed.gas_limit,
max_fee_per_gas: signed.max_fee_per_gas,
max_priority_fee_per_gas: signed.max_priority_fee_per_gas,
raw_transaction: signed.raw_transaction,
transaction_hash: signed.transaction_hash,
kind: AttemptKind::Canonical,
replaces_attempt_number: None,
})
.await
.expect("persist attempt");
*backend.nonce_counts.lock().expect("nonce count lock") = (1, 1);
reconcile(&store, &signer, &executor_config(), CancellationToken::new())
.await
.expect("effect reconciliation");
plan_task(&store).await;
assert_eq!(
store
.submission(submission_id)
.await
.expect("status")
.expect("submission")
.state,
SubmissionState::Succeeded
);
assert_eq!(
store
.signer_lifecycle(&signer.signer_id, signer.chain_id)
.await
.expect("lifecycle"),
SignerLifecycle::Ready
);
store.assert_invariants().await.expect("durable invariants");
}
#[tokio::test]
async fn unowned_pending_nonce_waits_and_recovers_automatically() {
let directory = tempfile::tempdir().expect("temporary directory");
let store = Arc::new(
Store::open(
&SqliteConfig {
path: directory.path().join("submitter.sqlite3"),
busy_timeout_ms: 5_000,
max_connections: 2,
},
admission(),
)
.await
.expect("store"),
);
let submission_id = match store.admit_task(&request(), &task_policy()).await.expect("admission") {
AdmissionOutcome::Accepted(resource) => resource.submission_id,
other => panic!("expected new admission, got {other:?}"),
};
let job_id = plan_task(&store).await.expect("job");
let backend = Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: false,
permanent_nonce_error: false,
nonce_counts: Mutex::new((0, 1)),
balance: Mutex::new(U256::MAX),
receipt_pending: false,
receipt_inconclusive: false,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
});
let signer = ManagedSigner {
signer_id: SignerId::new("task-busy-recovery").expect("signer id"),
chain_id: 31_337,
role: SignerRole::Task,
block_time: Duration::from_millis(1),
backend: backend.clone(),
};
let executor = Executor::new(
store.clone(),
vec![signer.clone()],
executor_config(),
Arc::new(RuntimeHealth::default()),
);
let cancellation = CancellationToken::new();
let runner = tokio::spawn({
let cancellation = cancellation.clone();
async move { executor.run(cancellation).await }
});
tokio::time::sleep(Duration::from_millis(50)).await;
assert_eq!(
store
.signer_lifecycle(&signer.signer_id, signer.chain_id)
.await
.expect("waiting lifecycle"),
SignerLifecycle::Ready
);
assert_eq!(
store
.submission(submission_id)
.await
.expect("waiting status")
.expect("waiting submission")
.state,
SubmissionState::ReadyForSubmission
);
*backend.nonce_counts.lock().expect("nonce count lock") = (1, 1);
tokio::time::timeout(Duration::from_secs(2), async {
loop {
if matches!(
store.execution_status(job_id).await.expect("common submission status"),
ExecutionStatus::Completed(_)
) {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("lane recovers after pending activity drains");
plan_task(&store).await;
assert_eq!(
store
.submission(submission_id)
.await
.expect("recovered status")
.expect("recovered submission")
.state,
SubmissionState::Succeeded
);
assert_eq!(
store
.signer_lifecycle(&signer.signer_id, signer.chain_id)
.await
.expect("recovered lifecycle"),
SignerLifecycle::Ready
);
cancellation.cancel();
runner.await.expect("executor task").expect("executor shutdown");
}
#[tokio::test]
async fn permanent_failure_quarantines_one_lane_without_stopping_siblings() {
let directory = tempfile::tempdir().expect("temporary directory");
let store = Arc::new(
Store::open(
&SqliteConfig {
path: directory.path().join("submitter.sqlite3"),
busy_timeout_ms: 5_000,
max_connections: 2,
},
admission(),
)
.await
.expect("store"),
);
let submission_id = match store.admit_task(&request(), &task_policy()).await.expect("admission") {
AdmissionOutcome::Accepted(resource) => resource.submission_id,
other => panic!("expected new admission, got {other:?}"),
};
let job_id = plan_task(&store).await.expect("job");
let healthy = ManagedSigner {
signer_id: SignerId::new("healthy-lane").expect("signer id"),
chain_id: 31_337,
role: SignerRole::Task,
block_time: Duration::from_millis(1),
backend: Arc::new(MockBackend {
address: Address::repeat_byte(7),
nonce_error: false,
permanent_nonce_error: false,
nonce_counts: Mutex::new((0, 0)),
balance: Mutex::new(U256::MAX),
receipt_pending: false,
receipt_inconclusive: false,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
}),
};
let broken = ManagedSigner {
signer_id: SignerId::new("broken-lane").expect("signer id"),
chain_id: 31_338,
role: SignerRole::Task,
block_time: Duration::from_millis(1),
backend: Arc::new(MockBackend {
address: Address::repeat_byte(8),
nonce_error: false,
permanent_nonce_error: true,
nonce_counts: Mutex::new((0, 0)),
balance: Mutex::new(U256::MAX),
receipt_pending: false,
receipt_inconclusive: false,
broadcasts: Mutex::new(Vec::new()),
observation_scopes: Mutex::new(Vec::new()),
}),
};
let executor = Executor::new(
store.clone(),
vec![broken.clone(), healthy],
executor_config(),
Arc::new(RuntimeHealth::default()),
);
let cancellation = CancellationToken::new();
let runner = tokio::spawn({
let cancellation = cancellation.clone();
async move { executor.run(cancellation).await }
});
tokio::time::timeout(Duration::from_secs(2), async {
loop {
if matches!(
store.execution_status(job_id).await.expect("common submission status"),
ExecutionStatus::Completed(_)
) {
break;
}
tokio::task::yield_now().await;
}
})
.await
.expect("healthy sibling completes its job");
plan_task(&store).await;
assert_eq!(
store
.submission(submission_id)
.await
.expect("submission status")
.expect("submission")
.state,
SubmissionState::Succeeded
);
assert_eq!(
store
.signer_lifecycle(&broken.signer_id, broken.chain_id)
.await
.expect("broken lifecycle"),
SignerLifecycle::Quarantined
);
cancellation.cancel();
runner.await.expect("executor task").expect("executor shutdown");
}
#[test]
fn lane_permanent_errors_quarantine_while_shared_store_errors_remain_fatal() {
let pending = ExecutorError::UnownedPendingTransaction {
signer: Address::ZERO,
latest: 0,
pending: 1,
};
assert!(pending.waits_for_chain());
assert!(!pending.requires_quarantine());
assert!(ExecutorError::NonceJournalAhead { journal: 1, pending: 0 }.waits_for_chain());
assert!(ExecutorError::EffectCoverage.waits_for_chain());
assert!(ExecutorError::AttemptNonceMismatch.requires_quarantine());
assert!(ExecutorError::MissingAttempt.requires_quarantine());
assert!(ExecutorError::Backend(BackendError::Policy("wrong chain".to_string())).requires_quarantine());
assert!(ExecutorError::Store(StoreError::Missing).is_process_fatal());
}