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
6 changes: 4 additions & 2 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 Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ derive_builder = "0.20.0"
thirtyfour = "0.31.0"
test-context = "0.3.0"
portpicker = "0.1.1"
rain-erc = "0.1.1"
rain-erc = "0.1.5"
rain-math-float = "0.1.7"
rain-error-decoding = "0.1.2"
# Direct dep required: the #[wasm_bindgen] macro emits `::wasm_bindgen` paths, and
Expand Down
6 changes: 5 additions & 1 deletion crates/common/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,11 @@ require its `COMMIT_SHA` env var.
prepare batch withdraw calldata; expose WASM‑friendly structs. The `local_db/`
subtree is split into `state.rs` (runtime state, query routing via
`LocalDbState`/`QuerySource`/`SyncReadiness`) and `status.rs` (UI
status‑reporting types).
status‑reporting types). The `markets/` subtree discovers direct quote-token
pairs from active indexed orders and assembles normalized order books, trades,
and rolling statistics without application-specific token semantics. Registry
token data identifies each network's quote token and enriches indexed
metadata; it does not decide which base tokens are listed.
- `dotrain_order` — Parse and validate a DOTRAIN config; compose
scenarios/deployments to Rainlang; fetch authoring metadata and pragma words;
merge additional settings.
Expand Down
1 change: 1 addition & 0 deletions crates/common/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,7 @@ csv = { workspace = true }
chrono = { workspace = true }
futures = { workspace = true }
rain-error-decoding = { workspace = true }
rain-erc = { workspace = true }
rain-interpreter-eval = { workspace = true }
wasm-bindgen = { workspace = true }
wasm-bindgen-utils = { workspace = true }
Expand Down
118 changes: 118 additions & 0 deletions crates/common/src/local_db/query/fetch_latest_trades_per_token/mod.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,118 @@
use crate::local_db::query::fetch_order_trades::LocalDbOrderTrade;
use crate::local_db::query::{SqlBuildError, SqlStatement, SqlValue};
use alloy::primitives::Address;

const QUERY_TEMPLATE: &str = include_str!("query.sql");
const RAINDEXES_CLAUSE: &str = "/*RAINDEXES_CLAUSE*/";
const RAINDEXES_CLAUSE_BODY: &str = "AND tws.raindex_address IN ({list})";
const PAIR_CLAUSE: &str = "/*PAIR_CLAUSE*/";

#[derive(Debug, Clone)]
pub struct FetchLatestTradesPerTokenArgs {
pub chain_id: u32,
pub raindex_addresses: Vec<Address>,
pub quote_token: Address,
pub base_tokens: Vec<Address>,
}

pub fn build_fetch_latest_trades_per_token_stmt(
args: &FetchLatestTradesPerTokenArgs,
) -> Result<SqlStatement, SqlBuildError> {
let mut stmt = SqlStatement::new(QUERY_TEMPLATE);

let quote_partition = push_param(&mut stmt, args.quote_token);
stmt.replace("/*QUOTE_PARTITION*/", &quote_partition)?;

let chain_id = push_param(&mut stmt, args.chain_id);
stmt.replace("/*CHAIN_ID*/", &chain_id)?;

let mut raindexes = args.raindex_addresses.clone();
raindexes.sort_unstable();
raindexes.dedup();
stmt.bind_list_clause(
RAINDEXES_CLAUSE,
RAINDEXES_CLAUSE_BODY,
raindexes.into_iter().map(SqlValue::from),
)?;

let mut base_tokens = args.base_tokens.clone();
base_tokens.sort_unstable();
base_tokens.dedup();
if base_tokens.is_empty() {
stmt.replace(PAIR_CLAUSE, "AND 1 = 0")?;
return Ok(stmt);
}

let input_quote = push_param(&mut stmt, args.quote_token);
let output_bases = push_list(&mut stmt, &base_tokens);
let output_quote = push_param(&mut stmt, args.quote_token);
let input_bases = push_list(&mut stmt, &base_tokens);
stmt.replace(
PAIR_CLAUSE,
&format!(
"AND ((tws.input_token = {input_quote} AND tws.output_token IN ({output_bases})) \
OR (tws.output_token = {output_quote} AND tws.input_token IN ({input_bases})))"
),
)?;
Ok(stmt)
}

fn push_param(stmt: &mut SqlStatement, value: impl Into<SqlValue>) -> String {
let placeholder = format!("?{}", stmt.params().len() + 1);
stmt.push(value.into());
placeholder
}

fn push_list(stmt: &mut SqlStatement, values: &[Address]) -> String {
values
.iter()
.map(|value| push_param(stmt, *value))
.collect::<Vec<_>>()
.join(", ")
}

pub type LatestTradeRow = LocalDbOrderTrade;

#[cfg(test)]
mod tests {
use super::*;
use alloy::primitives::address;

#[test]
fn query_selects_one_latest_trade_per_base_token() {
let quote = address!("2222222222222222222222222222222222222222");
let base_a = address!("1111111111111111111111111111111111111111");
let base_b = address!("3333333333333333333333333333333333333333");
let raindex = address!("4444444444444444444444444444444444444444");

let stmt = build_fetch_latest_trades_per_token_stmt(&FetchLatestTradesPerTokenArgs {
chain_id: 8453,
raindex_addresses: vec![raindex, raindex],
quote_token: quote,
base_tokens: vec![base_b, base_a, base_a],
})
.unwrap();

assert!(stmt.sql.contains("ROW_NUMBER() OVER"));
assert!(stmt.sql.contains("market_rank = 1"));
assert!(stmt.sql.contains("tws.raindex_address IN (?3)"));
assert!(stmt.sql.contains("tws.output_token IN (?5, ?6)"));
assert!(stmt.sql.contains("tws.input_token IN (?8, ?9)"));
assert_eq!(stmt.params().len(), 9);
assert!(!stmt.sql.contains("/*"));
}

#[test]
fn empty_base_tokens_build_a_match_none_query() {
let stmt = build_fetch_latest_trades_per_token_stmt(&FetchLatestTradesPerTokenArgs {
chain_id: 8453,
raindex_addresses: vec![],
quote_token: Address::ZERO,
base_tokens: vec![],
})
.unwrap();

assert!(stmt.sql.contains("AND 1 = 0"));
assert!(!stmt.sql.contains("/*"));
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
WITH ranked_trades AS (
SELECT
tws.*,
ROW_NUMBER() OVER (
PARTITION BY CASE
WHEN tws.input_token = /*QUOTE_PARTITION*/ THEN tws.output_token
ELSE tws.input_token
END
ORDER BY
tws.block_timestamp DESC,
tws.block_number DESC,
tws.log_index DESC,
tws.trade_kind,
tws.trade_side
) AS market_rank
FROM derived_trades tws
WHERE tws.chain_id = /*CHAIN_ID*/
/*RAINDEXES_CLAUSE*/
/*PAIR_CLAUSE*/
)
SELECT
tws.chain_id,
tws.trade_kind,
tws.raindex_address AS raindex,
tws.order_hash,
tws.order_owner,
tws.order_nonce,
tws.transaction_hash,
tws.log_index,
tws.block_number,
tws.block_timestamp,
tws.transaction_sender,
tws.input_vault_id,
tws.input_token,
tok_in.name AS input_token_name,
tok_in.symbol AS input_token_symbol,
tok_in.decimals AS input_token_decimals,
tws.input_delta,
tws.input_running_balance,
tws.output_vault_id,
tws.output_token,
tok_out.name AS output_token_name,
tok_out.symbol AS output_token_symbol,
tok_out.decimals AS output_token_decimals,
tws.output_delta,
tws.output_running_balance,
tws.trade_id
FROM ranked_trades tws
LEFT JOIN erc20_tokens tok_in
ON tok_in.chain_id = tws.chain_id
AND tok_in.raindex_address = tws.raindex_address
AND tok_in.token_address = tws.input_token
LEFT JOIN erc20_tokens tok_out
ON tok_out.chain_id = tws.chain_id
AND tok_out.raindex_address = tws.raindex_address
AND tok_out.token_address = tws.output_token
WHERE tws.market_rank = 1
ORDER BY
tws.block_timestamp DESC,
tws.block_number DESC,
tws.log_index DESC,
tws.trade_kind,
tws.trade_side;
1 change: 1 addition & 0 deletions crates/common/src/local_db/query/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ pub mod fetch_all_tokens;
pub mod fetch_db_metadata;
pub mod fetch_erc20_tokens_by_addresses;
pub mod fetch_last_synced_block;
pub mod fetch_latest_trades_per_token;
pub mod fetch_order_trades;
pub mod fetch_order_trades_count;
pub mod fetch_order_vaults_volume;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
use crate::local_db::query::fetch_latest_trades_per_token::{
build_fetch_latest_trades_per_token_stmt, FetchLatestTradesPerTokenArgs, LatestTradeRow,
};
use crate::local_db::query::{LocalDbQueryError, LocalDbQueryExecutor};
use crate::utils::timing::Timing;

pub async fn fetch_latest_trades_per_token<E: LocalDbQueryExecutor + ?Sized>(
exec: &E,
args: FetchLatestTradesPerTokenArgs,
) -> Result<Vec<LatestTradeRow>, LocalDbQueryError> {
if args.base_tokens.is_empty() {
return Ok(Vec::new());
}
let started = Timing::now();
let base_tokens_count = args.base_tokens.len();
let raindexes_count = args.raindex_addresses.len();
let stmt = build_fetch_latest_trades_per_token_stmt(&args)?;
let trades = exec.query_json::<Vec<LatestTradeRow>>(&stmt).await?;
tracing::info!(
chain_id = args.chain_id,
base_tokens_count,
raindexes_count,
rows = trades.len(),
duration_ms = started.elapsed_ms(),
"local DB latest market trades fetch completed"
);
Ok(trades)
}
1 change: 1 addition & 0 deletions crates/common/src/raindex_client/local_db/query/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ pub mod create_tables;
pub mod fetch_all_tokens;
pub mod fetch_erc20_tokens_by_addresses;
pub mod fetch_last_synced_block;
pub mod fetch_latest_trades_per_token;
pub mod fetch_order_trades;
pub mod fetch_order_trades_count;
pub mod fetch_order_vaults_volume;
Expand Down
Loading
Loading