Skip to content
Merged
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
13 changes: 13 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion cranker/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -25,4 +25,4 @@ stake-deposit-interceptor = { path = "../stake_deposit_interceptor" }
thiserror = "1.0.65"
tokio = { version = "1.41.0", features = ["full"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] } # Add env-filter feature
tracing-subscriber = { version = "0.3", features = ["env-filter", "json"] }
93 changes: 33 additions & 60 deletions cranker/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -79,14 +79,12 @@ impl InterceptorCranker {
}

pub async fn start(&self) {
info!("Starting InterceptorCranker service");
let mut interval_timer = time::interval(self.interval);
let mut tick: u64 = 0;
info!("Set interval timer to {} seconds", self.interval.as_secs());

loop {
interval_timer.tick().await;
info!("Tick: Starting new processing cycle");
info!("Cranker tick tick={tick}");
match emit_heartbeat(self.rpc_client.clone(), tick, &self.cluster_name).await {
Ok(_) => tick += 1,
Comment thread
aoikurokawa marked this conversation as resolved.
Err(e) => emit_error(format!("Failed to emit heartbeat: {e}"), &self.cluster_name),
Expand All @@ -108,19 +106,18 @@ impl InterceptorCranker {
}

match self.process_expired_receipts().await {
Ok(_) => info!("Successfully processed expired receipts"),
Ok(_) => info!("Processed expired receipts tick={tick}"),
Err(e) => emit_error(
format!("Error processing receipts: {e}"),
format!("Failed to process expired receipts: {e}"),
&self.cluster_name,
),
}
}
}

async fn process_expired_receipts(&self) -> Result<(), CrankerError> {
info!("Starting to process expired receipts");
let receipts = self.get_deposit_receipts().await?;
info!("Found {} deposit receipts", receipts.len());
info!("Fetched deposit receipts count={}", receipts.len());

let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
Expand All @@ -133,50 +130,42 @@ impl InterceptorCranker {
let mut claimed_receipts: u64 = 0;

for receipt in receipts {
// Get raw bytes using bytemuck and interpret as little-endian
let deposit_time = u64::from(receipt.deposit_time);
let cool_down = u64::from(receipt.cool_down_seconds);

info!(
"Receipt {} raw bytes:\n\
Interpreted values:\n\
deposit_time: {}\n\
cool_down: {}\n\
current_time: {}",
"Processing deposit receipt base={} deposit_time={} cool_down={} now={}",
receipt.base, deposit_time, cool_down, now
);
emit_deposit_receipt(&receipt, &self.cluster_name);

if deposit_time > now {
info!(
"Receipt {} not yet expired (future deposit time). Current time: {}, Deposit time: {}",
receipt.base,
now,
deposit_time
);
"Skipping receipt status=future_deposit base={} deposit_time={} now={}",
receipt.base, deposit_time, now
);
future_deposits += 1;
continue;
}

// Safe addition check
match deposit_time.checked_add(cool_down) {
Some(expiry_time) => {
if now > expiry_time {
info!(
"Receipt {} is expired. Current time: {}, Expiry time: {}",
receipt.base, now, expiry_time
"Receipt expired base={} expiry_time={} now={}",
receipt.base, expiry_time, now
);
match self.claim_pool_tokens(&receipt).await {
Ok(_) => {
info!("Successfully claimed tokens for receipt {}", receipt.base);
info!("Claimed pool tokens base={}", receipt.base);
let mut metrics = self.metrics.lock().unwrap();
metrics.successful_claims += 1;
claimed_receipts += 1;
}
Err(e) => {
emit_error(
format!(
"Failed to claim tokens for receipt {}: {}",
"Failed to claim pool tokens base={}: {}",
receipt.base, e
),
&self.cluster_name,
Expand All @@ -187,19 +176,20 @@ impl InterceptorCranker {
}
} else {
info!(
"Receipt {} not yet expired. Current time: {}, Expiry time: {}",
receipt.base, now, expiry_time
"Skipping receipt status=not_yet_expired base={} expiry_time={} now={}",
receipt.base, expiry_time, now
);
not_yet_expired_receipts += 1;
}
}
None => {
emit_error(format!(
"Receipt {} has invalid timing values - would overflow. Deposit time: {}, Cool down: {}",
receipt.base,
deposit_time,
cool_down
), &self.cluster_name);
emit_error(
format!(
"Skipping receipt status=overflow base={} deposit_time={} cool_down={}",
receipt.base, deposit_time, cool_down
),
&self.cluster_name,
);
}
}
}
Expand All @@ -217,7 +207,6 @@ impl InterceptorCranker {

async fn get_deposit_receipts(&self) -> Result<Vec<DepositReceipt>, CrankerError> {
let discriminator = StakeDepositInterceptorDiscriminators::DepositReceipt as u8;
info!("Searching for deposit receipts");

let accounts = self
.rpc_client
Expand All @@ -239,7 +228,7 @@ impl InterceptorCranker {
.await
.map_err(CrankerError::RpcError)?;

info!("Found {} raw accounts", accounts.len());
info!("Fetched program accounts count={}", accounts.len());

Ok(accounts
.into_iter()
Expand All @@ -248,11 +237,7 @@ impl InterceptorCranker {
match DepositReceipt::try_from_slice_unchecked(account_data.as_slice()) {
Ok(receipt) => {
info!(
"Found receipt:\n\
Account pubkey: {}\n\
Receipt base: {}\n\
Receipt stake pool: {}\n\
Derived PDA: {}",
"Decoded deposit receipt account={} base={} stake_pool={} derived_pda={}",
pubkey,
receipt.base,
receipt.stake_pool,
Expand All @@ -263,12 +248,11 @@ impl InterceptorCranker {
)
.0
);

Some(*receipt)
}
Err(e) => {
emit_error(
format!("Failed to deserialize receipt for {pubkey}: {e}"),
format!("Failed to deserialize receipt account={pubkey}: {e}"),
&self.cluster_name,
);
None
Expand All @@ -279,7 +263,7 @@ impl InterceptorCranker {
}

async fn claim_pool_tokens(&self, receipt: &DepositReceipt) -> Result<(), CrankerError> {
info!("Starting detailed claim debug for receipt {}", receipt.base);
info!("Claiming pool tokens base={}", receipt.base);

let stake_pool_deposit_authority = self
.get_stake_pool_deposit_authority(&receipt.stake_pool_deposit_stake_authority)
Expand All @@ -288,13 +272,12 @@ impl InterceptorCranker {
let owner_ata =
get_associated_token_address(&receipt.owner, &stake_pool_deposit_authority.pool_mint);

// Check if account exists
match self.rpc_client.get_account(&owner_ata).await {
Ok(_) => {
info!("Owner token account exists: {owner_ata}");
info!("Owner token account exists ata={owner_ata}");
}
Err(_) => {
info!("Creating owner token account: {owner_ata}");
info!("Creating owner token account ata={owner_ata}");
let create_ata_ix = create_associated_token_account(
&self.payer.pubkey(),
&receipt.owner,
Expand All @@ -313,7 +296,7 @@ impl InterceptorCranker {
self.rpc_client
.send_and_confirm_transaction(&create_ata_tx)
.await?;
info!("Created owner ata token account");
info!("Created owner token account ata={owner_ata}");
}
}

Expand All @@ -322,19 +305,12 @@ impl InterceptorCranker {
&stake_pool_deposit_authority.pool_mint,
);

// Check if account exists
match self.rpc_client.get_account(&fee_wallet_token_account).await {
Ok(_) => {
info!(
"Fee wallet token account exists: {}",
fee_wallet_token_account
);
info!("Fee wallet token account exists ata={fee_wallet_token_account}");
}
Err(_) => {
info!(
"Creating fee wallet token account: {}",
fee_wallet_token_account
);
info!("Creating fee wallet token account ata={fee_wallet_token_account}");
let create_ata_ix = create_associated_token_account(
&self.payer.pubkey(),
&stake_pool_deposit_authority.fee_wallet,
Expand All @@ -353,7 +329,7 @@ impl InterceptorCranker {
self.rpc_client
.send_and_confirm_transaction(&create_ata_tx)
.await?;
info!("Created fee wallet token account");
info!("Created fee wallet token account ata={fee_wallet_token_account}");
}
}

Expand Down Expand Up @@ -389,17 +365,14 @@ impl InterceptorCranker {
{
Ok(sig) => {
info!(
"Successfully claimed pool tokens for receipt {}. Transaction signature: {}",
"Claimed pool tokens base={} signature={}",
receipt.base, sig
);
Ok(())
}
Err(e) => {
emit_error(
format!(
"Failed to claim pool tokens for receipt {}. Error: {}",
receipt.base, e
),
format!("Failed to claim pool tokens base={}: {}", receipt.base, e),
&self.cluster_name,
);
Err(CrankerError::RpcError(e))
Expand Down
Loading
Loading