Skip to content

Commit 8a8943e

Browse files
fix(forester): restore Clone + compile tests after tracker cleanup
The cleanup commit changed get_ready_to_compress to return Vec<Pubkey> and CTokenCompressor/MintCompressor/PdaCompressor::new to take Arc<Keypair>, but the compressible tests weren't updated. Added get_ready_states / get_ready_states_for_program helpers returning Vec<S> for test ergonomics, wrapped test keypair args in Arc, and re-ran cargo +nightly fmt across the refactored modules. Co-Authored-By: Claude Opus 4.6 (1M context) <[email protected]>
1 parent a9ebc7a commit 8a8943e

23 files changed

Lines changed: 175 additions & 165 deletions

forester/src/compressible/mint/compressor.rs

Lines changed: 18 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -171,20 +171,24 @@ impl<R: Rpc + Indexer> MintCompressor<R> {
171171
};
172172
}
173173

174-
let mint_state =
175-
match compressor.tracker.accounts().get(&pubkey).map(|r| r.clone()) {
176-
Some(state) => state,
177-
None => {
178-
compressor.tracker.unmark_pending(&[pubkey]);
179-
return CompressionOutcome::Failed {
180-
pubkey,
181-
error: CompressionTaskError::Failed(anyhow::anyhow!(
182-
"mint {} removed from tracker before compression",
183-
pubkey
184-
)),
185-
};
186-
}
187-
};
174+
let mint_state = match compressor
175+
.tracker
176+
.accounts()
177+
.get(&pubkey)
178+
.map(|r| r.clone())
179+
{
180+
Some(state) => state,
181+
None => {
182+
compressor.tracker.unmark_pending(&[pubkey]);
183+
return CompressionOutcome::Failed {
184+
pubkey,
185+
error: CompressionTaskError::Failed(anyhow::anyhow!(
186+
"mint {} removed from tracker before compression",
187+
pubkey
188+
)),
189+
};
190+
}
191+
};
188192

189193
match compressor.compress(&mint_state).await {
190194
Ok(sig) => CompressionOutcome::Compressed {

forester/src/compressible/pda/compressor.rs

Lines changed: 18 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -187,20 +187,24 @@ impl<R: Rpc + Indexer> PdaCompressor<R> {
187187
}
188188

189189
// Look up account state from tracker; it may have been removed
190-
let account_state =
191-
match compressor.tracker.accounts().get(&pubkey).map(|r| r.clone()) {
192-
Some(state) => state,
193-
None => {
194-
compressor.tracker.unmark_pending(&[pubkey]);
195-
return CompressionOutcome::Failed {
196-
pubkey,
197-
error: CompressionTaskError::Failed(anyhow::anyhow!(
198-
"account {} removed from tracker before compression",
199-
pubkey
200-
)),
201-
};
202-
}
203-
};
190+
let account_state = match compressor
191+
.tracker
192+
.accounts()
193+
.get(&pubkey)
194+
.map(|r| r.clone())
195+
{
196+
Some(state) => state,
197+
None => {
198+
compressor.tracker.unmark_pending(&[pubkey]);
199+
return CompressionOutcome::Failed {
200+
pubkey,
201+
error: CompressionTaskError::Failed(anyhow::anyhow!(
202+
"account {} removed from tracker before compression",
203+
pubkey
204+
)),
205+
};
206+
}
207+
};
204208

205209
match compressor
206210
.compress(&account_state, &program_config, &cached_config)

forester/src/compressible/pda/state.rs

Lines changed: 17 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,23 @@ impl PdaAccountTracker {
8585
.collect()
8686
}
8787

88+
pub fn get_ready_states_for_program(
89+
&self,
90+
program_id: &Pubkey,
91+
current_slot: u64,
92+
) -> Vec<PdaAccountState> {
93+
let pending = self.pending();
94+
self.accounts()
95+
.iter()
96+
.filter(|entry| {
97+
entry.value().program_id == *program_id
98+
&& entry.value().is_ready_to_compress(current_slot)
99+
&& !pending.contains(entry.key())
100+
})
101+
.map(|entry| entry.value().clone())
102+
.collect()
103+
}
104+
88105
pub fn update_from_account(
89106
&self,
90107
pubkey: Pubkey,

forester/src/compressible/traits.rs

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,7 +60,7 @@ pub enum CompressionOutcome {
6060

6161
pub type CompressionOutcomes = Vec<CompressionOutcome>;
6262

63-
pub trait CompressibleState: Send + Sync {
63+
pub trait CompressibleState: Clone + Send + Sync {
6464
fn pubkey(&self) -> &Pubkey;
6565
fn lamports(&self) -> u64;
6666
fn compressible_slot(&self) -> u64;
@@ -138,6 +138,19 @@ pub trait CompressibleTracker<S: CompressibleState>: Send + Sync {
138138
.map(|entry| *entry.key())
139139
.collect()
140140
}
141+
142+
/// Clone of ready-to-compress states. Prefer `get_ready_to_compress` in
143+
/// production; this helper exists for tests that need full state fields.
144+
fn get_ready_states(&self, current_slot: u64) -> Vec<S> {
145+
let pending = self.pending();
146+
self.accounts()
147+
.iter()
148+
.filter(|entry| {
149+
entry.value().is_ready_to_compress(current_slot) && !pending.contains(entry.key())
150+
})
151+
.map(|entry| entry.value().clone())
152+
.collect()
153+
}
141154
}
142155

143156
/// Allows AccountSubscriber to work with any tracker type.

forester/src/epoch_manager/compression.rs

Lines changed: 5 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -12,18 +12,14 @@ use light_registry::ForesterEpochPda;
1212
use solana_program::pubkey::Pubkey;
1313
use tracing::{debug, error, info, trace, warn};
1414

15-
use crate::{
16-
compressible::{
17-
traits::{
18-
Cancelled, CompressibleState, CompressibleTracker, CompressionOutcome,
19-
CompressionTaskError,
20-
},
21-
CTokenCompressor, CompressibleConfig,
15+
use super::EpochManager;
16+
use crate::compressible::{
17+
traits::{
18+
Cancelled, CompressibleState, CompressibleTracker, CompressionOutcome, CompressionTaskError,
2219
},
20+
CTokenCompressor, CompressibleConfig,
2321
};
2422

25-
use super::EpochManager;
26-
2723
impl<R: Rpc + Indexer> EpochManager<R> {
2824
pub(crate) async fn dispatch_compression(
2925
&self,

forester/src/epoch_manager/context.rs

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -1,24 +1,23 @@
1-
use std::{sync::Arc, time::Duration};
1+
use std::{
2+
sync::{atomic::AtomicU64, Arc},
3+
time::Duration,
4+
};
25

3-
use light_client::{
4-
indexer::Indexer,
5-
rpc::Rpc,
6+
use forester_utils::{
7+
forester_epoch::{Epoch, TreeAccounts},
8+
rpc_pool::SolanaRpcPool,
69
};
10+
use light_client::{indexer::Indexer, rpc::Rpc};
711
use light_compressed_account::TreeType;
8-
use light_registry::utils::get_forester_epoch_pda_from_authority;
12+
use light_registry::{
13+
protocol_config::state::ProtocolConfig, utils::get_forester_epoch_pda_from_authority,
14+
};
915
use solana_sdk::{
1016
address_lookup_table::AddressLookupTableAccount,
1117
signature::{Keypair, Signer},
1218
};
1319
use tokio::sync::Mutex;
1420

15-
use forester_utils::{
16-
forester_epoch::{Epoch, TreeAccounts},
17-
rpc_pool::SolanaRpcPool,
18-
};
19-
use light_registry::protocol_config::state::ProtocolConfig;
20-
use std::sync::atomic::AtomicU64;
21-
2221
use crate::{
2322
logging::ServiceHeartbeat,
2423
priority_fee::PriorityFeeConfig,

forester/src/epoch_manager/mod.rs

Lines changed: 17 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,4 @@
1+
mod compression;
12
pub(crate) mod context;
23
mod monitor;
34
mod pipeline;
@@ -7,53 +8,37 @@ mod reporting;
78
pub(crate) mod tracker;
89
mod v1;
910
mod v2;
10-
mod compression;
1111

1212
use std::{
1313
sync::Arc,
1414
time::{Duration, SystemTime, UNIX_EPOCH},
1515
};
1616

1717
use anyhow::anyhow;
18-
use light_client::{
19-
indexer::Indexer,
20-
rpc::Rpc,
21-
};
18+
use forester_utils::{forester_epoch::TreeAccounts, rpc_pool::SolanaRpcPool};
19+
use light_client::{indexer::Indexer, rpc::Rpc};
2220
use light_compressed_account::TreeType;
23-
use forester_utils::{
24-
forester_epoch::TreeAccounts,
25-
rpc_pool::SolanaRpcPool,
26-
};
2721
use light_registry::protocol_config::state::ProtocolConfig;
28-
use solana_sdk::{
29-
address_lookup_table::AddressLookupTableAccount,
30-
signature::Signer,
31-
};
22+
use solana_sdk::{address_lookup_table::AddressLookupTableAccount, signature::Signer};
3223
use tokio::{
3324
sync::{mpsc, oneshot, watch, Mutex},
3425
task::JoinHandle,
3526
time::{sleep, Instant, MissedTickBehavior},
3627
};
3728
use tracing::{debug, error, info, info_span};
3829

30+
use self::{context::ForesterContext, processor_pool::ProcessorPool, tracker::EpochTracker};
3931
use crate::{
4032
compressible::CTokenAccountTracker,
4133
errors::InitializationError,
4234
logging::ServiceHeartbeat,
4335
processor::tx_cache::ProcessedHashCache,
36+
queue_helpers::QueueItemData,
4437
slot_tracker::SlotTracker,
4538
tree_data_sync::{fetch_protocol_group_authority, fetch_trees},
4639
ForesterConfig, Result,
4740
};
4841

49-
use self::{
50-
context::ForesterContext,
51-
processor_pool::ProcessorPool,
52-
tracker::EpochTracker,
53-
};
54-
55-
use crate::queue_helpers::QueueItemData;
56-
5742
// ── Public re-exports (preserve existing public API) ─────────────────────
5843

5944
/// Timing for a single circuit type (circuit inputs + proof generation)
@@ -438,7 +423,9 @@ fn spawn_heartbeat_task(
438423
let delta_queues_finished = current
439424
.queues_finished
440425
.saturating_sub(previous.queues_finished);
441-
let delta_items = current.items_processed.saturating_sub(previous.items_processed);
426+
let delta_items = current
427+
.items_processed
428+
.saturating_sub(previous.items_processed);
442429
let delta_work_reports = current
443430
.work_reports_sent
444431
.saturating_sub(previous.work_reports_sent);
@@ -1005,7 +992,10 @@ mod tests {
1005992
};
1006993
let work_item = WorkItem {
1007994
tree_account,
1008-
queue_item_data: QueueItemData { hash: [0u8; 32], index: 0 },
995+
queue_item_data: QueueItemData {
996+
hash: [0u8; 32],
997+
index: 0,
998+
},
1009999
};
10101000
assert!(work_item.is_address_tree());
10111001
assert!(!work_item.is_state_tree());
@@ -1022,7 +1012,10 @@ mod tests {
10221012
};
10231013
let work_item = WorkItem {
10241014
tree_account,
1025-
queue_item_data: QueueItemData { hash: [0u8; 32], index: 0 },
1015+
queue_item_data: QueueItemData {
1016+
hash: [0u8; 32],
1017+
index: 0,
1018+
},
10261019
};
10271020
assert!(!work_item.is_address_tree());
10281021
assert!(work_item.is_state_tree());

forester/src/epoch_manager/monitor.rs

Lines changed: 13 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -3,26 +3,20 @@
33
use std::{sync::Arc, time::Duration};
44

55
use anyhow::anyhow;
6-
use forester_utils::forester_epoch::{
7-
get_epoch_phases, TreeAccounts, TreeForesterSchedule,
8-
};
6+
use forester_utils::forester_epoch::{get_epoch_phases, TreeAccounts, TreeForesterSchedule};
97
use light_client::{indexer::Indexer, rpc::Rpc};
108
use solana_program::{native_token::LAMPORTS_PER_SOL, pubkey::Pubkey};
119
use solana_sdk::signature::Signer;
1210
use tokio::sync::mpsc;
1311
use tracing::{debug, error, info, warn};
1412

13+
use super::{context::should_skip_tree, EpochManager};
1514
use crate::{
1615
metrics::update_forester_sol_balance,
1716
slot_tracker::wait_until_slot_reached,
1817
tree_data_sync::{fetch_protocol_group_authority, fetch_trees},
1918
};
2019

21-
use super::{
22-
context::should_skip_tree,
23-
EpochManager,
24-
};
25-
2620
impl<R: Rpc + Indexer> EpochManager<R> {
2721
pub(super) async fn check_sol_balance_periodically(self: Arc<Self>) -> crate::Result<()> {
2822
let interval_duration = Duration::from_secs(300);
@@ -31,7 +25,10 @@ impl<R: Rpc + Indexer> EpochManager<R> {
3125
loop {
3226
interval.tick().await;
3327
match self.ctx.rpc_pool.get_connection().await {
34-
Ok(rpc) => match rpc.get_balance(&self.ctx.config.payer_keypair.pubkey()).await {
28+
Ok(rpc) => match rpc
29+
.get_balance(&self.ctx.config.payer_keypair.pubkey())
30+
.await
31+
{
3532
Ok(balance) => {
3633
let balance_in_sol = balance as f64 / (LAMPORTS_PER_SOL as f64);
3734
update_forester_sol_balance(
@@ -63,7 +60,11 @@ impl<R: Rpc + Indexer> EpochManager<R> {
6360
}
6461

6562
pub(super) async fn discover_trees_periodically(self: Arc<Self>) -> crate::Result<()> {
66-
let interval_secs = self.ctx.config.general_config.tree_discovery_interval_seconds;
63+
let interval_secs = self
64+
.ctx
65+
.config
66+
.general_config
67+
.tree_discovery_interval_seconds;
6768
if interval_secs == 0 {
6869
info!(event = "tree_discovery_disabled", run_id = %self.ctx.run_id, "Tree discovery disabled (interval=0)");
6970
return Ok(());
@@ -429,7 +430,8 @@ impl<R: Rpc + Indexer> EpochManager<R> {
429430
);
430431

431432
if let Err(e) =
432-
wait_until_slot_reached(&mut *rpc, &self.ctx.slot_tracker, wait_target).await
433+
wait_until_slot_reached(&mut *rpc, &self.ctx.slot_tracker, wait_target)
434+
.await
433435
{
434436
error!(
435437
event = "epoch_monitor_wait_for_registration_failed",

forester/src/epoch_manager/pipeline.rs

Lines changed: 2 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -7,14 +7,12 @@ use forester_utils::forester_epoch::{
77
};
88
use light_client::{indexer::Indexer, rpc::Rpc};
99
use light_compressed_account::TreeType;
10-
use light_registry::{
11-
protocol_config::state::EpochState,
12-
ForesterEpochPda,
13-
};
10+
use light_registry::{protocol_config::state::EpochState, ForesterEpochPda};
1411
use solana_sdk::signature::Signer;
1512
use tokio::time::Instant;
1613
use tracing::{debug, error, info, instrument, trace, warn};
1714

15+
use super::{context::should_skip_tree, tracker::RegistrationTracker, EpochManager};
1816
use crate::{
1917
errors::ForesterError,
2018
logging::should_emit_rate_limited_warning,
@@ -23,12 +21,6 @@ use crate::{
2321
ForesterEpochInfo,
2422
};
2523

26-
use super::{
27-
context::should_skip_tree,
28-
tracker::RegistrationTracker,
29-
EpochManager,
30-
};
31-
3224
impl<R: Rpc + Indexer> EpochManager<R> {
3325
#[instrument(level = "debug", skip(self), fields(forester = %self.ctx.config.payer_keypair.pubkey(), epoch = epoch))]
3426
pub(super) async fn process_epoch(self: Arc<Self>, epoch: u64) -> crate::Result<()> {

0 commit comments

Comments
 (0)