Skip to content
Open
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
1 change: 1 addition & 0 deletions Cargo.lock

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

32 changes: 30 additions & 2 deletions docs/runbooks/seed-pool-registry.md
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,7 @@ API returned a contract of an unexpected type — investigate before seeding it.

`events-backfill --discover-pools` reads the AMM factory events in a ledger
range from BE's `default.soroban_events` (Aquarius `add_pool`, Phoenix `create`,
Soroswap `new_pair`), runs them through the same `learn_factory` the live
Soroswap `new_pair`, SushiSwap V3 `pool_created` — task 0290), runs them through the same `learn_factory` the live
processor uses, and writes **only the pools `prices.pool_registry` does not
already hold**. No API key, no rewrite of existing rows, no candles. A re-run
writes nothing.
Expand Down Expand Up @@ -183,7 +183,9 @@ scp target/x86_64-unknown-linux-musl/release/events-backfill <prod-host>:~/event
factories (Phoenix and Soroswap leave `signature` NULL), 2-4 s per 320k-ledger
chunk on the shared box. As of 2026-09-17 every missing pool was created after
ledger 63,000,000 (checked per venue over the whole Soroban era), so the
catch-up only needs `63000000` to the tip.
catch-up only needs `63000000` to the tip. **Task 0290 is the exception:**
SushiSwap V3's pools go back to ledger 60,147,305, so its run starts at
`60000000` — the command is [below](#task-0290--sushiswap-v3s-wider-range).

```bash
# On the prod host, under tmux:
Expand All @@ -202,6 +204,32 @@ investigate before the write.
Then drop `--dry-run` to write, and run the dry run once more: it must report
`to_write=0`.

#### Task 0290 — SushiSwap V3's wider range

SushiSwap V3 is the one venue whose pools predate ledger 63,000,000: the first
`pool_created` is at 60,147,305 and the live factory's pools start at ~61.49M,
so the 63M catch-up above misses them. Run it over its own range **once**, then
the 63M catch-up covers it like every other venue:

```bash
# On the prod host, under tmux:
read -rs CH_PW
CLICKHOUSE_PASSWORD="$CH_PW" ~/events-backfill --discover-pools \
--start 60000000 --end <TIP> \
--clickhouse-url http://localhost:8123 --dry-run
```

That is ~4.5M ledgers, so ~14 chunks at the default `--chunk-size 320000` and
2-4 s of `topics_xdr` parsing each — a couple of minutes, well inside a tmux
session. It reads the same factory events as the run above; only `--start`
differs.

Expect `per_venue` to carry a `"sushiswap"` entry. The read has **no emitter
filter**, so it learns every generation's pools, not just the live factory's —
which is what this wider range is for: 99 SushiSwap pools have traded all-time
and three of them come from an earlier factory generation that is still trading.
Confirm with the same `FINAL` count below, then drop `--dry-run` to write.

### Verify

```sql
Expand Down
57 changes: 52 additions & 5 deletions packages/events-backfill/src/discover.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@
//! The reprice reads only the events of contracts already in the registry, so it
//! can never find a pool the registry is missing. This mode reads the factory
//! events themselves — Aquarius `add_pool`, Phoenix `create`, Soroswap
//! `new_pair` — and runs them through the same `learn_factory` the live
//! `new_pair`, SushiSwap V3 `pool_created` (task 0290) — and runs them through the same `learn_factory` the live
//! processor uses, so a row written here is the row live would have learned.
//! Only rows not already in the table are written; a re-run writes nothing.
//!
Expand Down Expand Up @@ -38,8 +38,11 @@ pub struct FactoryEventRow {
/// The factory-event read for `[start, end]`, as text.
///
/// Filters are a superset of what `learn_factory` accepts. BE fills `signature`
/// only when topic[0] is a Symbol, and two of the three factories use Strings:
/// only when topic[0] is a Symbol, and two of the four factories use Strings:
/// - Aquarius `add_pool`: Symbol topics, so `signature = 'add_pool'`;
/// - SushiSwap V3 `pool_created`: Symbol topics, so `signature = 'pool_created'`
/// (checked on production 2026-09-18). Every factory generation emits it, and
/// with no emitter filter one read learns the pools of all four;
/// - Phoenix `create`/`liquidity_pool`: **String** topics, so `signature` is NULL
/// — a `signature`-only filter finds none of the 20 Phoenix pools;
/// - Soroswap `SoroswapFactory`/`new_pair`: String topics, the action in topic[1].
Expand All @@ -59,7 +62,7 @@ pub(crate) fn factory_events_sql(start: u32, end: u32) -> String {
data_xdr \
FROM default.soroban_events \
WHERE ledger_sequence BETWEEN {start} AND {end} \
AND (signature IN ('add_pool', 'create') \
AND (signature IN ('add_pool', 'create', 'pool_created') \
OR (signature IS NULL \
AND (JSONExtractString(topics_xdr, 1, 'value') IN ('add_pool', 'create') \
OR JSONExtractString(topics_xdr, 2, 'value') = 'new_pair'))) \
Expand Down Expand Up @@ -179,12 +182,17 @@ mod tests {
const CREATE_DATA: &str =
r#"{"type":"address","value":"CBHCRSVX3ZZ7EGTSYMKPEFGZNWRVCSESQR3UABET4MIW52N4EVU6BIZX"}"#;

// SushiSwap V3's live factory `CD3KRKGD…GLYF`, ledger 64,116,662.
const POOL_CREATED_TOPICS: &str = r#"[{"type":"sym","value":"pool_created"}]"#;
const POOL_CREATED_DATA: &str = r#"{"type":"map","value":[{"key":{"type":"sym","value":"fee"},"value":{"type":"u32","value":500}},{"key":{"type":"sym","value":"pool_address"},"value":{"type":"address","value":"CBVHBZSZOS6KRDJ4D44FU2YLIENOVSSLM3UGKW6XQMVIFUAMWIWCVH2U"}},{"key":{"type":"sym","value":"sender"},"value":{"type":"address","value":"CD3KRKGDRVWPXVB3VXLUMQKMX6XZ6Q2H334IVZD4XXNAMKSRVQL5GLYF"}},{"key":{"type":"sym","value":"tick_spacing"},"value":{"type":"i32","value":10}},{"key":{"type":"sym","value":"token0"},"value":{"type":"address","value":"CBSJZEIO5C7KC2SF3MKSNXXJSW5G3VTNBX4ATMKUI3B2MR4JKM4R26YF"}},{"key":{"type":"sym","value":"token1"},"value":{"type":"address","value":"CCW67TSZV3SSS2HXMBQ5JFGCKJNXKZM7UQUWUZPUTHXSTZLEO7SJMI75"}}]}"#;

#[test]
fn learns_all_three_factory_shapes_from_real_payloads() {
fn learns_all_four_factory_shapes_from_real_payloads() {
let mut reg = Registries::new();
learn_from_row(&row(NEW_PAIR_TOPICS, NEW_PAIR_DATA), &mut reg);
learn_from_row(&row(ADD_POOL_TOPICS, ADD_POOL_DATA), &mut reg);
learn_from_row(&row(CREATE_TOPICS, CREATE_DATA), &mut reg);
learn_from_row(&row(POOL_CREATED_TOPICS, POOL_CREATED_DATA), &mut reg);

let pair = "CAZ4Z273BBAAFL5NYNQJKEMZDQBRCPKAS4GOXDUFXPSE56M4ONBJUOVD";
assert_eq!(reg.venue.get(pair), Some(&Venue::Soroswap));
Expand All @@ -207,6 +215,45 @@ mod tests {
.get("CBHCRSVX3ZZ7EGTSYMKPEFGZNWRVCSESQR3UABET4MIW52N4EVU6BIZX"),
Some(&Venue::Phoenix)
);
let pool = "CBVHBZSZOS6KRDJ4D44FU2YLIENOVSSLM3UGKW6XQMVIFUAMWIWCVH2U";
assert_eq!(reg.venue.get(pool), Some(&Venue::Sushiswap));
let p = reg.sushiswap.lookup(pool).expect("pool tokens learned");
assert_eq!(
p.token0,
"CBSJZEIO5C7KC2SF3MKSNXXJSW5G3VTNBX4ATMKUI3B2MR4JKM4R26YF"
);
assert_eq!(
p.token1,
"CCW67TSZV3SSS2HXMBQ5JFGCKJNXKZM7UQUWUZPUTHXSTZLEO7SJMI75"
);
}

/// Another protocol also emits `pool_created` (emitter `CBMBKXI7…Q227`,
/// wasm `ED0D122C`, ledger 63,173,214): fixed-denomination pools with no
/// token pair. The read selects it, so the learner must refuse it — a
/// pool registered with empty tokens would price against asset "".
#[test]
fn a_pool_created_without_a_token_pair_learns_nothing() {
let topics = r#"[{"type":"sym","value":"pool_created"},{"type":"address","value":"CAS3J7GYLGXMF6TDJBBYYSE3HQ6BBSMLNUQ34T6TZMYMW2EVH34XOWMA"}]"#;
let data = r#"{"type":"map","value":[{"key":{"type":"sym","value":"denomination"},"value":{"type":"i128","value":"1000000"}},{"key":{"type":"sym","value":"generation"},"value":{"type":"u32","value":0}},{"key":{"type":"sym","value":"pool"},"value":{"type":"address","value":"CC2SYHPYVRQ24IS6BQVA5WUYWDR3GK46W5ADY2UASC2D4IH6J5AXSEQ5"}}]}"#;
let mut reg = Registries::new();
learn_from_row(&row(topics, data), &mut reg);
assert!(reg.venue.is_empty());
}

#[test]
fn a_new_sushiswap_pool_is_written_as_sushiswap() {
let mut reg = Registries::new();
let persisted = snapshot(&reg);
learn_from_row(&row(POOL_CREATED_TOPICS, POOL_CREATED_DATA), &mut reg);

let rows = reg.pool_rows_unpersisted(&persisted);
assert_eq!(rows.len(), 1);
assert_eq!(
rows[0].contract_id,
"CBVHBZSZOS6KRDJ4D44FU2YLIENOVSSLM3UGKW6XQMVIFUAMWIWCVH2U"
);
assert_eq!(rows[0].venue, "sushiswap");
}

#[test]
Expand Down Expand Up @@ -248,7 +295,7 @@ mod tests {
let sql = factory_events_sql(63_000_000, 63_319_999);
assert!(sql.contains("ledger_sequence BETWEEN 63000000 AND 63319999"));
let sql = sql.split_whitespace().collect::<Vec<_>>().join(" ");
assert!(sql.contains("signature IN ('add_pool', 'create')"));
assert!(sql.contains("signature IN ('add_pool', 'create', 'pool_created')"));
// String-topic factories leave `signature` NULL: Phoenix's action is in
// topic[0], Soroswap's in topic[1] (1-based indexes 1 and 2).
assert!(sql.contains(
Expand Down
9 changes: 9 additions & 0 deletions packages/extractors-core/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,13 @@ pub enum Venue {
Soroswap,
Aquarius,
Phoenix,
/// SushiSwap V3 — a Uniswap-v3-style concentrated-liquidity venue (task
/// 0290). Its pool `swap` carries signed `amount0`/`amount1`, the CLMM
/// shape the `soroswap-extractor` crate already decodes, and it resolves
/// its tokens through that crate's pair registry exactly as
/// [`Venue::Soroswap`] does — a separate instance, so the two venues never
/// share a `contract_id`.
Sushiswap,
}

impl Venue {
Expand All @@ -15,6 +22,7 @@ impl Venue {
Venue::Soroswap => "soroswap",
Venue::Aquarius => "aquarius",
Venue::Phoenix => "phoenix",
Venue::Sushiswap => "sushiswap",
}
}

Expand All @@ -25,6 +33,7 @@ impl Venue {
"soroswap" => Some(Venue::Soroswap),
"aquarius" => Some(Venue::Aquarius),
"phoenix" => Some(Venue::Phoenix),
"sushiswap" => Some(Venue::Sushiswap),
_ => None,
}
}
Expand Down
64 changes: 54 additions & 10 deletions packages/ledger-processor/src/dispatch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ use phoenix_extractor::{
PHOENIX_STABLE_EVENT_COUNT, PHOENIX_XYK_MIN_EVENT_COUNT, POOL_TYPE_XYK, PhoenixPoolRegistry,
PhoenixXykExtractor,
};
use soroswap_extractor::{SoroswapPairExtractor, SoroswapPoolRegistry};
use soroswap_extractor::{PairPoolRegistry, PairSwapExtractor, SoroswapPairExtractor};

#[derive(Debug, thiserror::Error)]
pub enum DispatchError {
Expand Down Expand Up @@ -72,16 +72,36 @@ pub fn dispatch_phoenix(
}
}

/// The two pair-backed venues' pool registries, carried as ONE argument with
/// NAMED fields.
///
/// Both are a [`PairPoolRegistry`], so as two adjacent positional parameters
/// they were freely interchangeable: transposing them compiled silently and
/// produced no runtime error — every Soroswap pool would have resolved against
/// SushiSwap's token table and vice versa, emitting candles for the wrong asset
/// pair. Naming the fields makes that mistake something you have to write on
/// purpose, and `Registries::pair_registries` in `prices-ingest-core` is the
/// only place that wires them (task 0290 review).
#[derive(Clone, Copy)]
pub struct PairRegistries<'a> {
pub soroswap: &'a PairPoolRegistry,
pub sushiswap: &'a PairPoolRegistry,
}

/// Top-level dispatcher: routes events by venue, then by pool shape for Phoenix.
///
/// Soroswap requires the pool→tokens registry to resolve token identities; an
/// unresolved pool (created before the indexed window) yields no trades rather
/// than an error. Aquarius and Phoenix carry tokens inline.
/// Soroswap and SushiSwap require the pool→tokens registry to resolve token
/// identities; an unresolved pool (created before the indexed window) yields no
/// trades rather than an error. Aquarius and Phoenix carry tokens inline.
///
/// The two pair-backed venues keep SEPARATE registries so a contract_id can
/// never resolve to the wrong venue's tokens (task 0290); they arrive together
/// in [`PairRegistries`], which names them rather than ordering them.
pub fn dispatch(
rows: &[SorobanEventRow],
venue_registry: &VenueRegistry,
phoenix_registry: &PhoenixPoolRegistry,
soroswap_registry: &SoroswapPoolRegistry,
pairs: PairRegistries<'_>,
) -> Result<Vec<TradeRow>, DispatchError> {
if rows.is_empty() {
return Ok(vec![]);
Expand All @@ -92,10 +112,16 @@ pub fn dispatch(

match venue {
Some(Venue::Phoenix) => dispatch_phoenix(rows, phoenix_registry),
Some(Venue::Soroswap) => match soroswap_registry.lookup(contract_id) {
Some(Venue::Soroswap) => match pairs.soroswap.lookup(contract_id) {
Some(pair) => Ok(SoroswapPairExtractor::new(pair).extract(rows)?.trades),
None => Ok(vec![]),
},
Some(Venue::Sushiswap) => match pairs.sushiswap.lookup(contract_id) {
Some(pair) => Ok(PairSwapExtractor::with_venue(Venue::Sushiswap, pair)
.extract(rows)?
.trades),
None => Ok(vec![]),
},
Some(Venue::Aquarius) => Ok(AquariusPoolExtractor.extract(rows)?.trades),
None => Ok(vec![]),
}
Expand Down Expand Up @@ -127,7 +153,10 @@ mod tests {
&rows,
&venue_reg,
&phoenix_reg,
&SoroswapPoolRegistry::new(),
PairRegistries {
soroswap: &PairPoolRegistry::new(),
sushiswap: &PairPoolRegistry::new(),
},
)
.unwrap();
assert_eq!(trades.len(), 1);
Expand All @@ -145,7 +174,10 @@ mod tests {
&rows,
&venue_reg,
&phoenix_reg,
&SoroswapPoolRegistry::new(),
PairRegistries {
soroswap: &PairPoolRegistry::new(),
sushiswap: &PairPoolRegistry::new(),
},
)
.unwrap();
assert_eq!(trades.len(), 1);
Expand Down Expand Up @@ -196,7 +228,10 @@ mod tests {
&rows,
&venue_reg,
&phoenix_reg,
&SoroswapPoolRegistry::new(),
PairRegistries {
soroswap: &PairPoolRegistry::new(),
sushiswap: &PairPoolRegistry::new(),
},
)
.unwrap();
assert!(trades.is_empty());
Expand All @@ -207,7 +242,16 @@ mod tests {
let phoenix_reg = phoenix_registry_both_wasm_variants();
let venue_reg = venue_registry_phoenix(&[XLM_USDC_POOL]);

let trades = dispatch(&[], &venue_reg, &phoenix_reg, &SoroswapPoolRegistry::new()).unwrap();
let trades = dispatch(
&[],
&venue_reg,
&phoenix_reg,
PairRegistries {
soroswap: &PairPoolRegistry::new(),
sushiswap: &PairPoolRegistry::new(),
},
)
.unwrap();
assert!(trades.is_empty());
}
}
6 changes: 6 additions & 0 deletions packages/prices-clickhouse/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -55,3 +55,9 @@ rustls-pemfile = { version = "2", optional = true }
rustls-pki-types = { version = "1", optional = true }
webpki-roots = { version = "1", optional = true }
reqwest = { version = "0.12", default-features = false, features = ["json"], optional = true }

[dev-dependencies]
# Test-only: the AMM pre-roll's `source IN (...)` filter is pinned against the
# canonical `Venue` list, so a new venue cannot be added without the filter
# moving with it (task 0290 review).
extractors-core = { path = "../extractors-core" }
12 changes: 8 additions & 4 deletions packages/prices-clickhouse/schema/init.sql
Original file line number Diff line number Diff line change
Expand Up @@ -605,10 +605,14 @@ SETTINGS index_granularity = 8192;
-- output of the in-window registry so a partial re-backfill (a mid-history
-- window) or the live processor can LOAD it instead of re-deriving from Soroban
-- activation (this inverts task 0069: registry-as-output, not required-input).
-- venue = 'soroswap' | 'phoenix' | 'aquarius'. token0/token1 are the Soroswap
-- pair tokens (needed because a Soroswap swap event omits them); pool_type /
-- wasm_hash are Phoenix pool details; both default empty for venues that don't
-- use them. ReplacingMergeTree(updated_at) on contract_id collapses re-runs;
-- venue = 'soroswap' | 'phoenix' | 'aquarius' | 'sushiswap' (task 0290).
-- token0/token1 are the pair tokens of the two pair-backed venues — Soroswap
-- (from `new_pair`) and SushiSwap V3 (from `pool_created`) — needed because
-- their swap events omit them; pool_type / wasm_hash are Phoenix pool details;
-- both default empty for venues that don't use them. A pair-backed row with
-- blank tokens does NOT resolve: it loads its venue only, and the pool is
-- reported in `prices.unresolved_pools` rather than priced against an empty
-- asset. ReplacingMergeTree(updated_at) on contract_id collapses re-runs;
-- read with FINAL.
-- ---------------------------------------------------------------------
CREATE TABLE IF NOT EXISTS prices.pool_registry (
Expand Down
Loading
Loading