pub mod dexscreener;
pub mod jupiter;
use crate::config;
use anyhow::{anyhow, Result};
use serde::Deserialize;
use serde_json::{json, Value};
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
static RPC_ID: AtomicU64 = AtomicU64::new(0);
pub struct SolClient {
pub http: reqwest::Client,
pub rpc: String,
pub jupiter: String,
limiter: tokio::sync::Semaphore,
pub throttled: AtomicU64,
pub failed: AtomicU64,
started: Instant,
deadline: Duration,
pub past_deadline: AtomicU64,
}
#[derive(Debug, Clone, Deserialize)]
pub struct SigInfo {
pub signature: String,
pub slot: u64,
#[serde(rename = "blockTime")]
pub block_time: Option<i64>,
pub err: Option<Value>,
}
#[derive(Debug, Clone, Default)]
pub struct MintInfo {
pub program: String,
pub decimals: u32,
pub supply: u128,
pub supply_raw: String,
pub mint_authority: Option<String>,
pub freeze_authority: Option<String>,
pub extensions: Vec<Value>,
}
#[derive(Debug, Clone)]
pub struct LargeAccount {
pub token_account: String,
pub amount: u128,
}
#[derive(Debug, Clone, Default)]
pub struct MetaplexMeta {
pub name: String,
pub symbol: String,
pub uri: String,
pub update_authority: String,
pub is_mutable: bool,
}
impl Default for SolClient {
fn default() -> Self {
Self::new()
}
}
impl SolClient {
pub fn new() -> Self {
let http = reqwest::Client::builder()
.timeout(Duration::from_secs(25))
.user_agent("solwatch/0.1")
.build()
.expect("http client");
Self {
http,
rpc: config::rpc_url(),
jupiter: config::jupiter_base(),
limiter: tokio::sync::Semaphore::new(config::rpc_concurrency()),
throttled: AtomicU64::new(0),
failed: AtomicU64::new(0),
started: Instant::now(),
deadline: Duration::from_millis(config::deadline_ms()),
past_deadline: AtomicU64::new(0),
}
}
pub fn past_deadline(&self) -> bool {
!self.deadline.is_zero() && self.started.elapsed() > self.deadline
}
pub fn rpc_pressure(&self) -> Option<String> {
let t = self.throttled.load(Ordering::Relaxed);
let f = self.failed.load(Ordering::Relaxed);
let d = self.past_deadline.load(Ordering::Relaxed);
if d > 0 {
return Some(format!(
"Audit hit its time budget and stopped early — {d} read(s) were skipped, so this scan is partial. Point SOLWATCH_RPC_URL at a dedicated RPC (or set HELIUS_API_KEY) for a faster, complete scan."
));
}
if f == 0 {
return None;
}
Some(format!(
"RPC throttling hit {t} call(s) and {f} failed even after retries — those reads are missing"
))
}
pub async fn rpc(&self, method: &str, params: Value) -> Result<Value> {
if self.past_deadline() {
self.past_deadline.fetch_add(1, Ordering::Relaxed);
return Err(anyhow!("{method}: skipped (audit time budget spent)"));
}
let _permit = self.limiter.acquire().await;
let id = RPC_ID.fetch_add(1, Ordering::Relaxed);
let body = json!({"jsonrpc":"2.0","id":id,"method":method,"params":params});
let mut last = anyhow!("rpc {method} failed");
let mut saw_throttle = false;
for attempt in 0..5u32 {
if attempt > 0 {
if self.past_deadline() {
self.past_deadline.fetch_add(1, Ordering::Relaxed);
return Err(anyhow!(
"{method}: gave up after throttle (time budget spent)"
));
}
let base = 400u64 << (attempt - 1);
let jitter = id.wrapping_mul(2654435761) % 250;
tokio::time::sleep(Duration::from_millis(base + jitter)).await;
}
match self.http.post(&self.rpc).json(&body).send().await {
Ok(resp) => {
let status = resp.status();
if status.as_u16() == 429 || status.is_server_error() {
saw_throttle = true;
last = anyhow!("{method} HTTP {status}");
continue;
}
let v: Value = match resp.json().await {
Ok(v) => v,
Err(e) => {
last = anyhow!("{method} decode: {e}");
continue;
}
};
if let Some(err) = v.get("error") {
let code = err.get("code").and_then(|c| c.as_i64()).unwrap_or(0);
let msg = err
.get("message")
.and_then(|m| m.as_str())
.unwrap_or("error")
.to_string();
if code == 429 || msg.contains("Too many requests") {
saw_throttle = true;
last = anyhow!("{method}: {msg}");
continue;
}
if saw_throttle {
self.throttled.fetch_add(1, Ordering::Relaxed);
}
return Err(anyhow!("{method}: {msg}"));
}
if saw_throttle {
self.throttled.fetch_add(1, Ordering::Relaxed);
}
return Ok(v.get("result").cloned().unwrap_or(Value::Null));
}
Err(e) => last = anyhow!("{method}: {e}"),
}
}
if saw_throttle {
self.throttled.fetch_add(1, Ordering::Relaxed);
}
self.failed.fetch_add(1, Ordering::Relaxed);
Err(last)
}
pub async fn get_slot(&self) -> Result<u64> {
let r = self.rpc("getSlot", json!([])).await?;
r.as_u64().ok_or_else(|| anyhow!("getSlot: bad result"))
}
pub async fn account_info(&self, address: &str) -> Result<Value> {
let r = self
.rpc(
"getAccountInfo",
json!([address, {"encoding": "jsonParsed"}]),
)
.await?;
Ok(r.get("value").cloned().unwrap_or(Value::Null))
}
pub async fn mint_info(&self, mint: &str) -> Result<Option<MintInfo>> {
let v = self.account_info(mint).await?;
if v.is_null() {
return Ok(None);
}
let program = v
.get("owner")
.and_then(|o| o.as_str())
.unwrap_or_default()
.to_string();
let parsed = v.pointer("/data/parsed").cloned().unwrap_or(Value::Null);
if parsed.pointer("/type").and_then(|t| t.as_str()) != Some("mint") {
return Err(anyhow!("{mint} is not a token mint account"));
}
let info = parsed.get("info").cloned().unwrap_or(Value::Null);
let supply_raw = info
.get("supply")
.and_then(|s| s.as_str())
.unwrap_or("0")
.to_string();
Ok(Some(MintInfo {
program,
decimals: info.get("decimals").and_then(|d| d.as_u64()).unwrap_or(0) as u32,
supply: supply_raw.parse().unwrap_or(0),
supply_raw,
mint_authority: info
.get("mintAuthority")
.and_then(|a| a.as_str())
.map(String::from),
freeze_authority: info
.get("freezeAuthority")
.and_then(|a| a.as_str())
.map(String::from),
extensions: info
.get("extensions")
.and_then(|e| e.as_array())
.cloned()
.unwrap_or_default(),
}))
}
pub async fn largest_accounts(&self, mint: &str) -> Result<Vec<LargeAccount>> {
let mut r = self.rpc("getTokenLargestAccounts", json!([mint])).await;
if r.is_err() {
tokio::time::sleep(Duration::from_millis(2500)).await;
r = self.rpc("getTokenLargestAccounts", json!([mint])).await;
}
let r = r?;
let list = r
.get("value")
.and_then(|v| v.as_array())
.cloned()
.unwrap_or_default();
Ok(list
.iter()
.filter_map(|e| {
Some(LargeAccount {
token_account: e.get("address")?.as_str()?.to_string(),
amount: e.get("amount")?.as_str()?.parse().ok()?,
})
})
.collect())
}
pub async fn accounts_raw(&self, addrs: &[String]) -> Result<Vec<Option<Vec<u8>>>> {
if addrs.is_empty() {
return Ok(vec![]);
}
use base64::Engine;
let r = self
.rpc(
"getMultipleAccounts",
json!([addrs, {"encoding": "base64"}]),
)
.await?;
Ok(r.get("value")
.and_then(|v| v.as_array())
.cloned()
.unwrap_or_default()
.iter()
.map(|a| {
a.pointer("/data/0")
.and_then(|d| d.as_str())
.and_then(|s| base64::engine::general_purpose::STANDARD.decode(s).ok())
})
.collect())
}
pub async fn multiple_accounts(&self, addrs: &[String]) -> Result<Vec<Value>> {
if addrs.is_empty() {
return Ok(vec![]);
}
let r = self
.rpc(
"getMultipleAccounts",
json!([addrs, {"encoding": "jsonParsed"}]),
)
.await?;
Ok(r.get("value")
.and_then(|v| v.as_array())
.cloned()
.unwrap_or_default())
}
pub async fn signatures(
&self,
address: &str,
limit: usize,
before: Option<&str>,
) -> Result<Vec<SigInfo>> {
let mut opts = json!({"limit": limit});
if let Some(b) = before {
opts["before"] = json!(b);
}
let r = self
.rpc("getSignaturesForAddress", json!([address, opts]))
.await?;
Ok(serde_json::from_value(r).unwrap_or_default())
}
pub async fn transaction(&self, signature: &str) -> Result<Value> {
self.rpc(
"getTransaction",
json!([signature, {"encoding": "jsonParsed", "maxSupportedTransactionVersion": 0}]),
)
.await
}
pub async fn owner_mint_balance(&self, owner: &str, mint: &str) -> Result<u128> {
let r = self
.rpc(
"getTokenAccountsByOwner",
json!([owner, {"mint": mint}, {"encoding": "jsonParsed"}]),
)
.await?;
let list = r
.get("value")
.and_then(|v| v.as_array())
.cloned()
.unwrap_or_default();
Ok(list
.iter()
.filter_map(|a| {
a.pointer("/account/data/parsed/info/tokenAmount/amount")?
.as_str()?
.parse::<u128>()
.ok()
})
.sum())
}
pub async fn oldest_signatures(
&self,
address: &str,
max_pages: usize,
) -> Result<(Vec<SigInfo>, bool)> {
let mut before: Option<String> = None;
let mut oldest_page: Vec<SigInfo> = vec![];
let mut pages = 0usize;
loop {
let page = self.signatures(address, 1000, before.as_deref()).await?;
pages += 1;
let n = page.len();
if n > 0 {
before = Some(page[n - 1].signature.clone());
oldest_page = page;
}
if n < 1000 {
return Ok((oldest_page, false));
}
if pages >= max_pages {
return Ok((oldest_page, true));
}
}
}
pub async fn first_funding(&self, wallet: &str) -> Result<(Option<String>, usize)> {
let (page, truncated) = self.oldest_signatures(wallet, 3).await?;
let count_floor = if truncated { 3000 } else { page.len() };
if truncated {
return Ok((None, count_floor));
}
let Some(first) = page.last() else {
return Ok((None, 0));
};
let tx = self.transaction(&first.signature).await?;
Ok((first_sol_source(&tx, wallet), count_floor))
}
pub async fn jito_bundle_id(&self, signature: &str) -> Option<String> {
let url = format!("{}/{}", config::JITO_BUNDLES_API, signature);
let resp = self
.http
.get(&url)
.timeout(Duration::from_secs(8))
.send()
.await
.ok()?;
if !resp.status().is_success() {
return None;
}
let v: Value = resp.json().await.ok()?;
v.pointer("/0/bundle_id")
.and_then(|b| b.as_str())
.map(String::from)
}
pub async fn metaplex_metadata(&self, mint: &str) -> Result<Option<MetaplexMeta>> {
let Some(pda) = metadata_pda(mint) else {
return Ok(None);
};
let r = self
.rpc("getAccountInfo", json!([pda, {"encoding": "base64"}]))
.await?;
let Some(data_b64) = r.pointer("/value/data/0").and_then(|d| d.as_str()) else {
return Ok(None);
};
use base64::Engine;
let bytes = base64::engine::general_purpose::STANDARD
.decode(data_b64)
.unwrap_or_default();
Ok(parse_metaplex(&bytes))
}
}
pub fn parse_metaplex(b: &[u8]) -> Option<MetaplexMeta> {
let mut i = 0usize;
let take = |i: &mut usize, n: usize| -> Option<&[u8]> {
let s = b.get(*i..*i + n)?;
*i += n;
Some(s)
};
let _key = take(&mut i, 1)?;
let update_authority = bs58::encode(take(&mut i, 32)?).into_string();
let _mint = take(&mut i, 32)?;
let string = |i: &mut usize| -> Option<String> {
let len = u32::from_le_bytes(take(i, 4)?.try_into().ok()?) as usize;
if len > 4096 {
return None;
}
let raw = take(i, len)?;
Some(
String::from_utf8_lossy(raw)
.trim_end_matches('\0')
.to_string(),
)
};
let name = string(&mut i)?;
let symbol = string(&mut i)?;
let uri = string(&mut i)?;
let _sfbp = take(&mut i, 2)?;
let has_creators = take(&mut i, 1)?[0];
if has_creators == 1 {
let n = u32::from_le_bytes(take(&mut i, 4)?.try_into().ok()?) as usize;
if n > 16 {
return None;
}
take(&mut i, n * 34)?;
}
let _primary_sale = take(&mut i, 1)?;
let is_mutable = take(&mut i, 1)?[0] == 1;
Some(MetaplexMeta {
name,
symbol,
uri,
update_authority,
is_mutable,
})
}
pub fn first_sol_source(tx: &Value, wallet: &str) -> Option<String> {
let check = |ins: &Value| -> Option<String> {
let parsed = ins.get("parsed")?;
let typ = parsed.get("type").and_then(|t| t.as_str())?;
let info = parsed.get("info")?;
let dest_key = match typ {
"transfer" | "transferWithSeed" => "destination",
"createAccount" | "createAccountWithSeed" => "newAccount",
_ => return None,
};
if info.get(dest_key).and_then(|d| d.as_str()) != Some(wallet) {
return None;
}
info.get("source")
.and_then(|s| s.as_str())
.filter(|s| *s != wallet)
.map(String::from)
};
if let Some(list) = tx
.pointer("/transaction/message/instructions")
.and_then(|v| v.as_array())
{
for ins in list {
if let Some(src) = check(ins) {
return Some(src);
}
}
}
if let Some(groups) = tx
.pointer("/meta/innerInstructions")
.and_then(|v| v.as_array())
{
for g in groups {
if let Some(list) = g.get("instructions").and_then(|v| v.as_array()) {
for ins in list {
if let Some(src) = check(ins) {
return Some(src);
}
}
}
}
}
None
}
pub fn find_pda(seeds: &[&[u8]], program: &str) -> Option<String> {
use sha2::{Digest, Sha256};
let program_b = bs58::decode(program).into_vec().ok()?;
if program_b.len() != 32 {
return None;
}
for bump in (0u8..=255).rev() {
let mut h = Sha256::new();
for s in seeds {
h.update(s);
}
h.update([bump]);
h.update(&program_b);
h.update(b"ProgramDerivedAddress");
let out: [u8; 32] = h.finalize().into();
if !is_on_curve(&out) {
return Some(bs58::encode(out).into_string());
}
}
None
}
pub fn metadata_pda(mint: &str) -> Option<String> {
let program = bs58::decode(config::METAPLEX_METADATA).into_vec().ok()?;
let mint_b = bs58::decode(mint).into_vec().ok()?;
if mint_b.len() != 32 {
return None;
}
find_pda(&[b"metadata", &program, &mint_b], config::METAPLEX_METADATA)
}
pub fn pumpfun_bonding_curve_pda(mint: &str) -> Option<String> {
let mint_b = bs58::decode(mint).into_vec().ok()?;
if mint_b.len() != 32 {
return None;
}
find_pda(&[b"bonding-curve", &mint_b], config::PUMPFUN_PROGRAM)
}
pub fn ata_pda(owner: &str, mint: &str, token_program: &str) -> Option<String> {
let owner_b = bs58::decode(owner).into_vec().ok()?;
let mint_b = bs58::decode(mint).into_vec().ok()?;
let prog_b = bs58::decode(token_program).into_vec().ok()?;
if owner_b.len() != 32 || mint_b.len() != 32 || prog_b.len() != 32 {
return None;
}
find_pda(&[&owner_b, &prog_b, &mint_b], config::ASSOCIATED_TOKEN)
}
fn is_on_curve(bytes: &[u8; 32]) -> bool {
curve25519_dalek::edwards::CompressedEdwardsY(*bytes)
.decompress()
.is_some()
}
pub fn is_pubkey(s: &str) -> bool {
(32..=44).contains(&s.len())
&& bs58::decode(s)
.into_vec()
.map(|v| v.len() == 32)
.unwrap_or(false)
}