use std::{
collections::{HashMap, HashSet},
env, fs,
path::{Path, PathBuf},
time::Duration,
};
use tracing::{error, info, warn};
use tracing_subscriber::EnvFilter;
use tycho_client::{
feed::{component_tracker::ComponentFilter, dto, BlockHeader, FeedMessage},
stream::TychoStreamBuilder,
};
use tycho_common::{
models::{token::Token, Chain},
Bytes,
};
use tycho_simulation::utils::load_all_tokens;
const TYCHO_ENDPOINT: &str = "tycho-beta.propellerheads.xyz";
const TVL_THRESHOLD: f64 = 10.0;
const SNAPSHOT_TIMEOUT_SECS: u64 = 180;
struct FixtureTarget {
label: &'static str,
filename: &'static str,
protocol: &'static str,
min_tokens: usize,
max_tokens: usize,
}
#[tokio::main]
async fn main() {
tracing_subscriber::fmt()
.with_env_filter(EnvFilter::from_default_env().add_directive("info".parse().unwrap()))
.init();
let api_key =
env::var("TYCHO_API_KEY").expect("TYCHO_API_KEY environment variable must be set");
let fixtures_dir = fixtures_path();
fs::create_dir_all(&fixtures_dir).expect("Failed to create fixtures directory");
info!("Connecting to Tycho at {TYCHO_ENDPOINT}...");
let tvl_filter = ComponentFilter::with_tvl_range(TVL_THRESHOLD, TVL_THRESHOLD);
let (_handle, mut rx) = TychoStreamBuilder::new(TYCHO_ENDPOINT, Chain::Ethereum)
.exchange("vm:balancer_v2", tvl_filter.clone())
.exchange("vm:curve", tvl_filter.clone())
.auth_key(Some(api_key.clone()))
.max_messages(1)
.startup_timeout(Duration::from_secs(SNAPSHOT_TIMEOUT_SECS))
.build()
.await
.expect("Failed to build Tycho stream");
info!("Waiting for first snapshot (up to {SNAPSHOT_TIMEOUT_SECS}s)...");
let feed_msg = rx
.recv()
.await
.expect("Stream channel closed before first message — did startup timeout fire?")
.expect("Stream returned an error on first message");
info!("Received snapshot with {} protocol state messages", feed_msg.state_msgs.len());
for (name, msg) in &feed_msg.state_msgs {
info!(
" protocol={name} components={} vm_storage_accounts={}",
msg.snapshots.states.len(),
msg.snapshots.vm_storage.len()
);
}
let dto_feed: dto::FeedMessage<BlockHeader> = feed_msg.into();
let targets: Vec<FixtureTarget> = vec![
FixtureTarget {
label: "Balancer V2 2-token",
filename: "balancer_v2_2token.json",
protocol: "vm:balancer_v2",
min_tokens: 2,
max_tokens: 2,
},
FixtureTarget {
label: "Curve 3-token",
filename: "curve_3token.json",
protocol: "vm:curve",
min_tokens: 3,
max_tokens: 3,
},
FixtureTarget {
label: "Curve 4-token",
filename: "curve_4token.json",
protocol: "vm:curve",
min_tokens: 4,
max_tokens: 4,
},
];
let mut all_referenced_tokens: HashSet<Bytes> = HashSet::new();
let mut captured_count = 0;
for target in &targets {
if capture_target(target, &dto_feed, &fixtures_dir, &mut all_referenced_tokens) {
captured_count += 1;
}
}
if captured_count == 0 {
error!("No fixtures were captured. Check TVL threshold or exchange names.");
std::process::exit(1);
}
write_tokens_fixture(&fixtures_dir, &api_key, &all_referenced_tokens).await;
info!("Done. Captured {}/{} fixtures.", captured_count, targets.len());
}
fn capture_target(
target: &FixtureTarget,
dto_feed: &dto::FeedMessage<BlockHeader>,
fixtures_dir: &Path,
all_referenced_tokens: &mut HashSet<Bytes>,
) -> bool {
let Some(state_msg) = dto_feed.state_msgs.get(target.protocol) else {
warn!("Protocol '{}' not found in snapshot — skipping {}", target.protocol, target.label);
return false;
};
let chosen = choose_component(state_msg, target.min_tokens, target.max_tokens);
let Some((component_id, component_with_state)) = chosen else {
warn!(
"No component with {}-{} tokens found in '{}' — skipping {}",
target.min_tokens, target.max_tokens, target.protocol, target.label
);
return false;
};
info!(
"Selected component '{}' (tokens={}) for {}",
component_id,
component_with_state
.component
.tokens
.len(),
target.label
);
for token in &component_with_state.component.tokens {
all_referenced_tokens.insert(token.clone());
}
let component_contracts: HashSet<Bytes> = component_with_state
.component
.contract_ids
.iter()
.cloned()
.collect();
let filtered_vm_storage: HashMap<Bytes, _> = state_msg
.snapshots
.vm_storage
.iter()
.filter(|(addr, _)| component_contracts.contains(*addr))
.map(|(k, v)| (k.clone(), v.clone()))
.collect();
info!(
" vm_storage entries for component: {} (out of {} total)",
filtered_vm_storage.len(),
state_msg.snapshots.vm_storage.len()
);
let filtered_states = HashMap::from([(component_id, component_with_state.clone())]);
let filtered_state_msg = dto::StateSyncMessage {
header: state_msg.header.clone(),
snapshots: dto::Snapshot { states: filtered_states, vm_storage: filtered_vm_storage },
deltas: None,
removed_components: HashMap::new(),
};
let filtered_feed = dto::FeedMessage {
state_msgs: HashMap::from([(target.protocol.to_string(), filtered_state_msg)]),
sync_states: dto_feed
.sync_states
.iter()
.filter(|(k, _)| k.as_str() == target.protocol)
.map(|(k, v)| (k.clone(), v.clone()))
.collect(),
};
let json =
serde_json::to_string_pretty(&filtered_feed).expect("Failed to serialize fixture to JSON");
let output_path = fixtures_dir.join(target.filename);
fs::write(&output_path, &json).expect("Failed to write fixture file");
info!("Wrote {} ({} bytes) to {:?}", target.filename, json.len(), output_path);
verify_round_trip(&json, target.filename);
true
}
fn choose_component(
state_msg: &dto::StateSyncMessage<BlockHeader>,
min_tokens: usize,
max_tokens: usize,
) -> Option<(String, &dto::ComponentWithState)> {
state_msg
.snapshots
.states
.iter()
.find(|(_, cws)| {
let n = cws.component.tokens.len();
n >= min_tokens && n <= max_tokens
})
.map(|(id, cws)| (id.clone(), cws))
}
fn verify_round_trip(json: &str, filename: &str) {
let dto_result: Result<dto::FeedMessage<BlockHeader>, _> = serde_json::from_str(json);
match dto_result {
Ok(dto_msg) => {
let _runtime_msg: FeedMessage<BlockHeader> = dto_msg.into();
info!("Round-trip OK for {filename}");
}
Err(e) => {
error!("Round-trip FAILED for {filename}: {e}");
std::process::exit(1);
}
}
}
async fn write_tokens_fixture(
fixtures_dir: &Path,
api_key: &str,
referenced_addresses: &HashSet<Bytes>,
) {
info!("Fetching token metadata for {} token addresses...", referenced_addresses.len());
let all_tokens =
load_all_tokens(TYCHO_ENDPOINT, false, Some(api_key), true, Chain::Ethereum, None, None)
.await
.expect("Failed to load token metadata from Tycho");
let relevant: Vec<&Token> = all_tokens
.values()
.filter(|t| referenced_addresses.contains(&t.address))
.collect();
info!("Found {} relevant tokens (out of {} total)", relevant.len(), all_tokens.len());
let json = serde_json::to_string_pretty(&relevant).expect("Failed to serialize tokens to JSON");
let output_path = fixtures_dir.join("tokens.json");
fs::write(&output_path, &json).expect("Failed to write tokens.json");
info!("Wrote tokens.json ({} bytes) to {:?}", json.len(), output_path);
}
fn fixtures_path() -> PathBuf {
PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("benches")
.join("fixtures")
}