1use std::{fmt::Debug, num::ParseIntError, str::FromStr};
17
18use alloy::primitives::{Address, I256, U160, U256};
19use nautilus_core::{
20 UnixNanos,
21 datetime::{NANOSECONDS_IN_MICROSECOND, NANOSECONDS_IN_MILLISECOND, NANOSECONDS_IN_SECOND},
22};
23use nautilus_model::{
24 defi::{
25 PoolLiquidityUpdate, PoolLiquidityUpdateType, PoolSwap, SharedChain, SharedDex,
26 data::{
27 DexPoolData, PoolFeeCollect, PoolFeeProtocolCollect, PoolFeeProtocolUpdate, PoolFlash,
28 },
29 validation::validate_address,
30 },
31 identifiers::InstrumentId,
32};
33use sqlx::{FromRow, Row, postgres::PgRow};
34
35const MAX_UNIX_SECONDS_TIMESTAMP: u64 = 9_999_999_999;
36const MAX_UNIX_MILLISECONDS_TIMESTAMP: u64 = MAX_UNIX_SECONDS_TIMESTAMP * 1_000 + 999;
37const MAX_UNIX_MICROSECONDS_TIMESTAMP: u64 = MAX_UNIX_SECONDS_TIMESTAMP * 1_000_000 + 999_999;
38
39fn decode_address(value: &str, field: &str) -> Result<Address, sqlx::Error> {
40 validate_address(value)
41 .map_err(|e| sqlx::Error::Decode(format!("Invalid {field} address '{value}': {e}").into()))
42}
43
44fn decode_u64(value: i64, field: &str) -> Result<u64, sqlx::Error> {
45 u64::try_from(value)
46 .map_err(|e| sqlx::Error::Decode(format!("Invalid {field} '{value}': {e}").into()))
47}
48
49fn decode_u32(value: i32, field: &str) -> Result<u32, sqlx::Error> {
50 u32::try_from(value)
51 .map_err(|e| sqlx::Error::Decode(format!("Invalid {field} '{value}': {e}").into()))
52}
53
54fn decode_u8(value: i32, field: &str) -> Result<u8, sqlx::Error> {
55 u8::try_from(value)
56 .map_err(|e| sqlx::Error::Decode(format!("Invalid {field} '{value}': {e}").into()))
57}
58
59#[derive(Debug)]
64pub struct TokenRow {
65 pub address: Address,
66 pub name: String,
67 pub symbol: String,
68 pub decimals: u8,
69}
70
71impl<'r> FromRow<'r, PgRow> for TokenRow {
72 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
73 let address_value = row.try_get::<String, _>("address")?;
74 let address = decode_address(&address_value, "token")?;
75 let name = row.try_get::<String, _>("name")?;
76 let symbol = row.try_get::<String, _>("symbol")?;
77 let decimals = decode_u8(row.try_get::<i32, _>("decimals")?, "token decimals")?;
78
79 let token = Self {
80 address,
81 name,
82 symbol,
83 decimals,
84 };
85 Ok(token)
86 }
87}
88
89#[derive(Debug)]
90pub struct PoolRow {
91 pub address: Address,
92 pub pool_identifier: String,
93 pub dex_name: String,
94 pub creation_block: u64,
95 pub creation_block_timestamp: Option<UnixNanos>,
96 pub token0_chain: i32,
97 pub token0_address: Address,
98 pub token1_chain: i32,
99 pub token1_address: Address,
100 pub fee: Option<u32>,
101 pub tick_spacing: Option<u32>,
102 pub initial_tick: Option<i32>,
103 pub initial_sqrt_price_x96: Option<String>,
104 pub hook_address: Option<String>,
105}
106
107impl<'r> FromRow<'r, PgRow> for PoolRow {
108 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
109 let address_value = row.try_get::<String, _>("address")?;
110 let address = decode_address(&address_value, "pool")?;
111 let pool_identifier = row.try_get::<String, _>("pool_identifier")?;
112 let dex_name = row.try_get::<String, _>("dex_name")?;
113 let creation_block = decode_u64(
114 row.try_get::<i64, _>("creation_block")?,
115 "pool creation block",
116 )?;
117 let creation_block_timestamp =
118 row.try_get::<Option<String>, _>("creation_block_timestamp")?;
119 let creation_block_timestamp = creation_block_timestamp
120 .as_deref()
121 .map(parse_cached_block_timestamp)
122 .transpose()
123 .map_err(|e| {
124 sqlx::Error::Decode(
125 format!("Invalid creation block timestamp '{creation_block_timestamp:?}': {e}")
126 .into(),
127 )
128 })?;
129 let token0_chain = row.try_get::<i32, _>("token0_chain")?;
130 let token0_address_value = row.try_get::<String, _>("token0_address")?;
131 let token0_address = decode_address(&token0_address_value, "token0")?;
132 let token1_chain = row.try_get::<i32, _>("token1_chain")?;
133 let token1_address_value = row.try_get::<String, _>("token1_address")?;
134 let token1_address = decode_address(&token1_address_value, "token1")?;
135 let fee = row
136 .try_get::<Option<i32>, _>("fee")?
137 .map(|value| decode_u32(value, "pool fee"))
138 .transpose()?;
139 let tick_spacing = row
140 .try_get::<Option<i32>, _>("tick_spacing")?
141 .map(|value| decode_u32(value, "pool tick spacing"))
142 .transpose()?;
143 let initial_tick = row.try_get::<Option<i32>, _>("initial_tick")?;
144 let initial_sqrt_price_x96 = row.try_get::<Option<String>, _>("initial_sqrt_price_x96")?;
145 let hook_address = row.try_get::<Option<String>, _>("hook_address")?;
146
147 Ok(Self {
148 address,
149 pool_identifier,
150 dex_name,
151 creation_block,
152 creation_block_timestamp,
153 token0_chain,
154 token0_address,
155 token1_chain,
156 token1_address,
157 fee,
158 tick_spacing,
159 initial_tick,
160 initial_sqrt_price_x96,
161 hook_address,
162 })
163 }
164}
165
166#[derive(Debug)]
168pub struct BlockTimestampRow {
169 pub number: u64,
171 pub timestamp: UnixNanos,
173}
174
175impl FromRow<'_, PgRow> for BlockTimestampRow {
176 fn from_row(row: &PgRow) -> Result<Self, sqlx::Error> {
177 let number = decode_u64(row.try_get::<i64, _>("number")?, "block number")?;
178 let timestamp = row.try_get::<String, _>("timestamp")?;
179 let timestamp = parse_cached_block_timestamp(×tamp).map_err(|e| {
180 sqlx::Error::Decode(format!("Invalid block timestamp '{timestamp}': {e}").into())
181 })?;
182 Ok(Self { number, timestamp })
183 }
184}
185
186pub(crate) fn parse_cached_block_timestamp(value: &str) -> Result<UnixNanos, ParseIntError> {
187 let timestamp = value.parse::<u64>()?;
188 if timestamp <= MAX_UNIX_SECONDS_TIMESTAMP {
189 return Ok(UnixNanos::from(timestamp * NANOSECONDS_IN_SECOND));
190 }
191
192 if timestamp <= MAX_UNIX_MILLISECONDS_TIMESTAMP {
193 return Ok(UnixNanos::from(timestamp * NANOSECONDS_IN_MILLISECOND));
194 }
195
196 if timestamp <= MAX_UNIX_MICROSECONDS_TIMESTAMP {
197 return Ok(UnixNanos::from(timestamp * NANOSECONDS_IN_MICROSECOND));
198 }
199
200 Ok(UnixNanos::from(timestamp))
201}
202
203pub fn transform_row_to_dex_pool_data(
212 row: &PgRow,
213 chain: SharedChain,
214 dex: SharedDex,
215 instrument_id: InstrumentId,
216) -> Result<DexPoolData, sqlx::Error> {
217 let event_type = row.try_get::<String, _>("event_type")?;
218 let pool_identifier_str = row.try_get::<String, _>("pool_identifier")?;
219 let pool_identifier = pool_identifier_str
220 .parse()
221 .map_err(|e| sqlx::Error::Decode(format!("Invalid pool identifier: {e}").into()))?;
222 let block = decode_u64(row.try_get::<i64, _>("block")?, "event block")?;
223 let block_hash = row.try_get::<Option<String>, _>("block_hash")?;
224 let transaction_hash = row.try_get::<String, _>("transaction_hash")?;
225 let transaction_index = decode_u32(
226 row.try_get::<i32, _>("transaction_index")?,
227 "transaction index",
228 )?;
229 let log_index = decode_u32(row.try_get::<i32, _>("log_index")?, "log index")?;
230 let block_timestamp = row.try_get::<String, _>("block_timestamp")?;
231 let timestamp = parse_cached_block_timestamp(&block_timestamp).map_err(|e| {
232 sqlx::Error::Decode(format!("Invalid block timestamp '{block_timestamp}': {e}").into())
233 })?;
234
235 match event_type.as_str() {
236 "swap" => {
237 let sender_str = row
238 .try_get::<Option<String>, _>("sender")?
239 .ok_or_else(|| sqlx::Error::Decode("Missing sender for swap event".into()))?;
240 let sender = validate_address(&sender_str)
241 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
242
243 let recipient_str = row
244 .try_get::<Option<String>, _>("recipient")?
245 .ok_or_else(|| sqlx::Error::Decode("Missing recipient for swap event".into()))?;
246 let recipient = validate_address(&recipient_str)
247 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
248
249 let sqrt_price_x96_str = row
250 .try_get::<Option<String>, _>("sqrt_price_x96")?
251 .ok_or_else(|| {
252 sqlx::Error::Decode("Missing sqrt_price_x96 for swap event".into())
253 })?;
254 let sqrt_price_x96 = U160::from_str(&sqrt_price_x96_str).map_err(|e| {
255 sqlx::Error::Decode(
256 format!("Invalid sqrt_price_x96 '{sqrt_price_x96_str}': {e}").into(),
257 )
258 })?;
259
260 let swap_liquidity_str = row.try_get::<String, _>("swap_liquidity")?;
261 let swap_liquidity = u128::from_str(&swap_liquidity_str)
262 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
263
264 let swap_tick = row.try_get::<i32, _>("swap_tick")?;
265
266 let swap_amount0_str = row
267 .try_get::<Option<String>, _>("swap_amount0")?
268 .ok_or_else(|| sqlx::Error::Decode("Missing swap_amount0 for swap event".into()))?;
269 let amount0 = I256::from_str(&swap_amount0_str).map_err(|e| {
270 sqlx::Error::Decode(
271 format!("Invalid swap_amount0 '{swap_amount0_str}': {e}").into(),
272 )
273 })?;
274
275 let swap_amount1_str = row
276 .try_get::<Option<String>, _>("swap_amount1")?
277 .ok_or_else(|| sqlx::Error::Decode("Missing swap_amount1 for swap event".into()))?;
278 let amount1 = I256::from_str(&swap_amount1_str).map_err(|e| {
279 sqlx::Error::Decode(
280 format!("Invalid swap_amount1 '{swap_amount1_str}': {e}").into(),
281 )
282 })?;
283
284 let mut pool_swap = PoolSwap::new(
285 chain,
286 dex,
287 instrument_id,
288 pool_identifier,
289 block,
290 transaction_hash,
291 transaction_index,
292 log_index,
293 timestamp, timestamp, sender,
296 recipient,
297 amount0,
298 amount1,
299 sqrt_price_x96,
300 swap_liquidity,
301 swap_tick,
302 );
303 pool_swap.block_hash = block_hash;
304
305 Ok(DexPoolData::Swap(pool_swap))
306 }
307 "liquidity" => {
308 let kind_str = row
309 .try_get::<Option<String>, _>("liquidity_event_type")?
310 .ok_or_else(|| {
311 sqlx::Error::Decode("Missing liquidity_event_type for liquidity event".into())
312 })?;
313
314 let kind = match kind_str.as_str() {
315 "Mint" => PoolLiquidityUpdateType::Mint,
316 "Burn" => PoolLiquidityUpdateType::Burn,
317 _ => {
318 return Err(sqlx::Error::Decode(
319 format!("Unknown liquidity update type: {kind_str}").into(),
320 ));
321 }
322 };
323
324 let sender = row
325 .try_get::<Option<String>, _>("sender")?
326 .map(|s| validate_address(&s))
327 .transpose()
328 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
329
330 let owner_str = row
331 .try_get::<Option<String>, _>("owner")?
332 .ok_or_else(|| sqlx::Error::Decode("Missing owner for liquidity event".into()))?;
333 let owner = validate_address(&owner_str)
334 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
335
336 let position_liquidity_str = row.try_get::<String, _>("position_liquidity")?;
338 let position_liquidity = position_liquidity_str.parse::<u128>().map_err(|e| {
339 sqlx::Error::Decode(
340 format!("Invalid position_liquidity '{position_liquidity_str}': {e}").into(),
341 )
342 })?;
343
344 let amount0_str = row.try_get::<String, _>("amount0")?;
345 let amount0 = U256::from_str_radix(&amount0_str, 10).map_err(|e| {
346 sqlx::Error::Decode(format!("Invalid amount0 '{amount0_str}': {e}").into())
347 })?;
348
349 let amount1_str = row.try_get::<String, _>("amount1")?;
350 let amount1 = U256::from_str_radix(&amount1_str, 10).map_err(|e| {
351 sqlx::Error::Decode(format!("Invalid amount1 '{amount1_str}': {e}").into())
352 })?;
353
354 let tick_lower = row
355 .try_get::<Option<i32>, _>("tick_lower")?
356 .ok_or_else(|| {
357 sqlx::Error::Decode("Missing tick_lower for liquidity event".into())
358 })?;
359
360 let tick_upper = row
361 .try_get::<Option<i32>, _>("tick_upper")?
362 .ok_or_else(|| {
363 sqlx::Error::Decode("Missing tick_upper for liquidity event".into())
364 })?;
365
366 let mut pool_liquidity_update = PoolLiquidityUpdate::new(
367 chain,
368 dex,
369 instrument_id,
370 pool_identifier,
371 kind,
372 block,
373 transaction_hash,
374 transaction_index,
375 log_index,
376 sender,
377 owner,
378 position_liquidity,
379 amount0,
380 amount1,
381 tick_lower,
382 tick_upper,
383 timestamp, timestamp, );
386 pool_liquidity_update.block_hash = block_hash;
387
388 Ok(DexPoolData::LiquidityUpdate(pool_liquidity_update))
389 }
390 "collect" => {
391 let owner_str = row
392 .try_get::<Option<String>, _>("owner")?
393 .ok_or_else(|| sqlx::Error::Decode("Missing owner for collect event".into()))?;
394 let owner = validate_address(&owner_str)
395 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
396
397 let amount0_str = row.try_get::<String, _>("amount0")?;
399 let amount0 = amount0_str.parse::<u128>().map_err(|e| {
400 sqlx::Error::Decode(format!("Invalid amount0 '{amount0_str}': {e}").into())
401 })?;
402
403 let amount1_str = row.try_get::<String, _>("amount1")?;
404 let amount1 = amount1_str.parse::<u128>().map_err(|e| {
405 sqlx::Error::Decode(format!("Invalid amount1 '{amount1_str}': {e}").into())
406 })?;
407
408 let tick_lower = row
409 .try_get::<Option<i32>, _>("tick_lower")?
410 .ok_or_else(|| {
411 sqlx::Error::Decode("Missing tick_lower for collect event".into())
412 })?;
413
414 let tick_upper = row
415 .try_get::<Option<i32>, _>("tick_upper")?
416 .ok_or_else(|| {
417 sqlx::Error::Decode("Missing tick_upper for collect event".into())
418 })?;
419
420 let mut pool_fee_collect = PoolFeeCollect::new(
421 chain,
422 dex,
423 instrument_id,
424 pool_identifier,
425 block,
426 transaction_hash,
427 transaction_index,
428 log_index,
429 owner,
430 amount0,
431 amount1,
432 tick_lower,
433 tick_upper,
434 timestamp, timestamp, );
437 pool_fee_collect.block_hash = block_hash;
438
439 Ok(DexPoolData::FeeCollect(pool_fee_collect))
440 }
441 "fee_protocol_update" => {
442 let fee_protocol0_new = row.try_get::<i32, _>("fee_protocol0_new")?;
443 let fee_protocol1_new = row.try_get::<i32, _>("fee_protocol1_new")?;
444 let fee_protocol0_new = u32::try_from(fee_protocol0_new).map_err(|e| {
445 sqlx::Error::Decode(
446 format!("Invalid fee_protocol0_new '{fee_protocol0_new}': {e}").into(),
447 )
448 })?;
449 let fee_protocol1_new = u32::try_from(fee_protocol1_new).map_err(|e| {
450 sqlx::Error::Decode(
451 format!("Invalid fee_protocol1_new '{fee_protocol1_new}': {e}").into(),
452 )
453 })?;
454
455 let mut pool_fee_protocol_update = PoolFeeProtocolUpdate::new(
456 chain,
457 dex,
458 instrument_id,
459 pool_identifier,
460 block,
461 transaction_hash,
462 transaction_index,
463 log_index,
464 fee_protocol0_new,
465 fee_protocol1_new,
466 timestamp, timestamp, );
469 pool_fee_protocol_update.block_hash = block_hash;
470
471 Ok(DexPoolData::FeeProtocolUpdate(pool_fee_protocol_update))
472 }
473 "fee_protocol_collect" => {
474 let sender_str = row.try_get::<Option<String>, _>("sender")?.ok_or_else(|| {
475 sqlx::Error::Decode("Missing sender for fee_protocol_collect event".into())
476 })?;
477 let sender = validate_address(&sender_str)
478 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
479
480 let recipient_str =
481 row.try_get::<Option<String>, _>("recipient")?
482 .ok_or_else(|| {
483 sqlx::Error::Decode(
484 "Missing recipient for fee_protocol_collect event".into(),
485 )
486 })?;
487 let recipient = validate_address(&recipient_str)
488 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
489
490 let amount0_str = row.try_get::<String, _>("amount0")?;
492 let amount0 = amount0_str.parse::<u128>().map_err(|e| {
493 sqlx::Error::Decode(format!("Invalid amount0 '{amount0_str}': {e}").into())
494 })?;
495
496 let amount1_str = row.try_get::<String, _>("amount1")?;
497 let amount1 = amount1_str.parse::<u128>().map_err(|e| {
498 sqlx::Error::Decode(format!("Invalid amount1 '{amount1_str}': {e}").into())
499 })?;
500
501 let mut pool_fee_protocol_collect = PoolFeeProtocolCollect::new(
502 chain,
503 dex,
504 instrument_id,
505 pool_identifier,
506 block,
507 transaction_hash,
508 transaction_index,
509 log_index,
510 sender,
511 recipient,
512 amount0,
513 amount1,
514 timestamp, timestamp, );
517 pool_fee_protocol_collect.block_hash = block_hash;
518
519 Ok(DexPoolData::FeeProtocolCollect(pool_fee_protocol_collect))
520 }
521 "flash" => {
522 let sender_str = row
523 .try_get::<Option<String>, _>("sender")?
524 .ok_or_else(|| sqlx::Error::Decode("Missing sender for flash event".into()))?;
525 let sender = validate_address(&sender_str)
526 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
527
528 let recipient_str = row
529 .try_get::<Option<String>, _>("recipient")?
530 .ok_or_else(|| sqlx::Error::Decode("Missing recipient for flash event".into()))?;
531 let recipient = validate_address(&recipient_str)
532 .map_err(|e| sqlx::Error::Decode(e.to_string().into()))?;
533
534 let flash_amount0_str = row.try_get::<String, _>("flash_amount0")?;
536 let amount0 = U256::from_str_radix(&flash_amount0_str, 10).map_err(|e| {
537 sqlx::Error::Decode(
538 format!("Invalid flash_amount0 '{flash_amount0_str}': {e}").into(),
539 )
540 })?;
541
542 let flash_amount1_str = row.try_get::<String, _>("flash_amount1")?;
543 let amount1 = U256::from_str_radix(&flash_amount1_str, 10).map_err(|e| {
544 sqlx::Error::Decode(
545 format!("Invalid flash_amount1 '{flash_amount1_str}': {e}").into(),
546 )
547 })?;
548
549 let flash_paid0_str = row.try_get::<String, _>("flash_paid0")?;
550 let paid0 = U256::from_str_radix(&flash_paid0_str, 10).map_err(|e| {
551 sqlx::Error::Decode(format!("Invalid flash_paid0 '{flash_paid0_str}': {e}").into())
552 })?;
553
554 let flash_paid1_str = row.try_get::<String, _>("flash_paid1")?;
555 let paid1 = U256::from_str_radix(&flash_paid1_str, 10).map_err(|e| {
556 sqlx::Error::Decode(format!("Invalid flash_paid1 '{flash_paid1_str}': {e}").into())
557 })?;
558
559 let mut pool_flash = PoolFlash::new(
560 chain,
561 dex,
562 instrument_id,
563 pool_identifier,
564 block,
565 transaction_hash,
566 transaction_index,
567 log_index,
568 timestamp, timestamp, sender,
571 recipient,
572 amount0,
573 amount1,
574 paid0,
575 paid1,
576 );
577 pool_flash.block_hash = block_hash;
578
579 Ok(DexPoolData::Flash(pool_flash))
580 }
581 _ => Err(sqlx::Error::Decode(
582 format!("Unknown event type: {event_type}").into(),
583 )),
584 }
585}
586
587#[derive(Debug, Clone, PartialEq, Eq)]
589pub struct ExecutionTransactionRow {
590 pub wallet_address: Option<String>,
591 pub nonce: u64,
592 pub transaction_hash: String,
593 pub purpose: String,
594 pub status: String,
595 pub client_order_id: Option<String>,
596}
597
598#[derive(Debug, Clone, PartialEq, Eq)]
600pub struct ExecutionIntentInsert {
601 pub chain_id: u32,
602 pub wallet_address: String,
603 pub purpose: String,
604 pub client_order_id: Option<String>,
605 pub trader_id: Option<String>,
606 pub strategy_id: Option<String>,
607 pub account_id: Option<String>,
608 pub instrument_id: Option<String>,
609 pub pool_address: Option<String>,
610 pub transaction_to: String,
611 pub transaction_input: String,
612 pub transaction_value: String,
613 pub amount_in: Option<String>,
614 pub created_block: u64,
615}
616
617#[derive(Debug, Clone, PartialEq, Eq)]
619pub struct ExecutionIntentRow {
620 pub id: i64,
621 pub schema_version: i16,
622 pub chain_id: u32,
623 pub wallet_address: String,
624 pub nonce: Option<u64>,
625 pub purpose: String,
626 pub status: String,
627 pub client_order_id: Option<String>,
628 pub trader_id: Option<String>,
629 pub strategy_id: Option<String>,
630 pub account_id: Option<String>,
631 pub instrument_id: Option<String>,
632 pub pool_address: Option<String>,
633 pub transaction_to: String,
634 pub transaction_input: String,
635 pub transaction_value: String,
636 pub amount_in: Option<String>,
637 pub created_block: u64,
638 pub acknowledgement_emitted: bool,
639 pub fill_emitted: bool,
640 pub terminal_emitted: bool,
641 pub active: bool,
642}
643
644impl<'r> FromRow<'r, PgRow> for ExecutionIntentRow {
645 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
646 let chain_id_i32 = row.try_get::<i32, _>("chain_id")?;
647 let chain_id = u32::try_from(chain_id_i32).map_err(|_| {
648 sqlx::Error::Decode(format!("Invalid negative chain ID {chain_id_i32}").into())
649 })?;
650 let nonce = row
651 .try_get::<Option<i64>, _>("nonce")?
652 .map(|nonce| {
653 u64::try_from(nonce).map_err(|_| {
654 sqlx::Error::Decode(format!("Invalid negative nonce {nonce}").into())
655 })
656 })
657 .transpose()?;
658 let created_block_i64 = row.try_get::<i64, _>("created_block")?;
659 let created_block = u64::try_from(created_block_i64).map_err(|_| {
660 sqlx::Error::Decode(
661 format!("Invalid negative creation block {created_block_i64}").into(),
662 )
663 })?;
664
665 Ok(Self {
666 id: row.try_get("id")?,
667 schema_version: row.try_get("schema_version")?,
668 chain_id,
669 wallet_address: row.try_get("wallet_address")?,
670 nonce,
671 purpose: row.try_get("purpose")?,
672 status: row.try_get("status")?,
673 client_order_id: row.try_get("client_order_id")?,
674 trader_id: row.try_get("trader_id")?,
675 strategy_id: row.try_get("strategy_id")?,
676 account_id: row.try_get("account_id")?,
677 instrument_id: row.try_get("instrument_id")?,
678 pool_address: row.try_get("pool_address")?,
679 transaction_to: row.try_get("transaction_to")?,
680 transaction_input: row.try_get("transaction_input")?,
681 transaction_value: row.try_get("transaction_value")?,
682 amount_in: row.try_get("amount_in")?,
683 created_block,
684 acknowledgement_emitted: row.try_get("acknowledgement_emitted")?,
685 fill_emitted: row.try_get("fill_emitted")?,
686 terminal_emitted: row.try_get("terminal_emitted")?,
687 active: row.try_get("active")?,
688 })
689 }
690}
691
692#[derive(Clone, PartialEq, Eq)]
694pub struct ExecutionTransactionHashRow {
695 pub id: i64,
696 pub intent_id: i64,
697 pub chain_id: u32,
698 pub transaction_hash: String,
699 pub payload_expected: bool,
700 pub raw_transaction: Option<Vec<u8>>,
701 pub sealed_transaction: Option<Vec<u8>>,
702 pub status: String,
703 pub block_number: Option<u64>,
704 pub block_hash: Option<String>,
705 pub receipt_success: Option<bool>,
706 pub gas_used: Option<u64>,
707 pub effective_gas_price: Option<String>,
708 pub current: bool,
709}
710
711impl Debug for ExecutionTransactionHashRow {
712 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
713 f.debug_struct(stringify!(ExecutionTransactionHashRow))
714 .field("id", &self.id)
715 .field("intent_id", &self.intent_id)
716 .field("chain_id", &self.chain_id)
717 .field("transaction_hash", &self.transaction_hash)
718 .field("payload_expected", &self.payload_expected)
719 .field(
720 "raw_transaction",
721 &self.raw_transaction.as_ref().map(|_| "<redacted>"),
722 )
723 .field(
724 "sealed_transaction",
725 &self.sealed_transaction.as_ref().map(|_| "<redacted>"),
726 )
727 .field("status", &self.status)
728 .field("block_number", &self.block_number)
729 .field("block_hash", &self.block_hash)
730 .field("receipt_success", &self.receipt_success)
731 .field("gas_used", &self.gas_used)
732 .field("effective_gas_price", &self.effective_gas_price)
733 .field("current", &self.current)
734 .finish()
735 }
736}
737
738impl<'r> FromRow<'r, PgRow> for ExecutionTransactionHashRow {
739 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
740 let chain_id_i32 = row.try_get::<i32, _>("chain_id")?;
741 let chain_id = u32::try_from(chain_id_i32).map_err(|_| {
742 sqlx::Error::Decode(format!("Invalid negative chain ID {chain_id_i32}").into())
743 })?;
744 let block_number = row
745 .try_get::<Option<i64>, _>("block_number")?
746 .map(|block| {
747 u64::try_from(block).map_err(|_| {
748 sqlx::Error::Decode(format!("Invalid negative block number {block}").into())
749 })
750 })
751 .transpose()?;
752 let gas_used = row
753 .try_get::<Option<i64>, _>("gas_used")?
754 .map(|gas| {
755 u64::try_from(gas).map_err(|_| {
756 sqlx::Error::Decode(format!("Invalid negative gas used {gas}").into())
757 })
758 })
759 .transpose()?;
760
761 Ok(Self {
762 id: row.try_get("id")?,
763 intent_id: row.try_get("intent_id")?,
764 chain_id,
765 transaction_hash: row.try_get("transaction_hash")?,
766 payload_expected: row.try_get("payload_expected")?,
767 raw_transaction: row.try_get("raw_transaction")?,
768 sealed_transaction: row.try_get("sealed_transaction")?,
769 status: row.try_get("status")?,
770 block_number,
771 block_hash: row.try_get("block_hash")?,
772 receipt_success: row.try_get("receipt_success")?,
773 gas_used,
774 effective_gas_price: row.try_get("effective_gas_price")?,
775 current: row.try_get("current")?,
776 })
777 }
778}
779
780impl<'r> FromRow<'r, PgRow> for ExecutionTransactionRow {
781 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
782 let wallet_address = row.try_get::<Option<String>, _>("wallet_address")?;
783 let nonce_i64 = row.try_get::<i64, _>("nonce")?;
784 let nonce = u64::try_from(nonce_i64).map_err(|_| {
785 sqlx::Error::Decode(format!("Invalid negative nonce {nonce_i64}").into())
786 })?;
787 let transaction_hash = row.try_get::<String, _>("transaction_hash")?;
788 let purpose = row.try_get::<String, _>("purpose")?;
789 let status = row.try_get::<String, _>("status")?;
790 let client_order_id = row.try_get::<Option<String>, _>("client_order_id")?;
791
792 Ok(Self {
793 wallet_address,
794 nonce,
795 transaction_hash,
796 purpose,
797 status,
798 client_order_id,
799 })
800 }
801}
802
803#[cfg(test)]
804mod tests {
805 use nautilus_core::datetime::{
806 NANOSECONDS_IN_MICROSECOND, NANOSECONDS_IN_MILLISECOND, NANOSECONDS_IN_SECOND,
807 };
808 use rstest::rstest;
809
810 use super::*;
811
812 #[rstest]
813 #[case("1700000000", 1_700_000_000 * NANOSECONDS_IN_SECOND)]
814 #[case("9999999999", 9_999_999_999 * NANOSECONDS_IN_SECOND)]
815 #[case("1700000000123", 1_700_000_000_123 * NANOSECONDS_IN_MILLISECOND)]
816 #[case("9999999999999", 9_999_999_999_999 * NANOSECONDS_IN_MILLISECOND)]
817 #[case("1700000000123456", 1_700_000_000_123_456 * NANOSECONDS_IN_MICROSECOND)]
818 #[case("9999999999999999", 9_999_999_999_999_999 * NANOSECONDS_IN_MICROSECOND)]
819 #[case("1700000000123456789", 1_700_000_000_123_456_789)]
820 fn parse_cached_block_timestamp_returns_unix_nanos(#[case] value: &str, #[case] expected: u64) {
821 let timestamp = parse_cached_block_timestamp(value).unwrap();
822
823 assert_eq!(timestamp, UnixNanos::from(expected));
824 }
825
826 #[rstest]
827 fn parse_cached_block_timestamp_rejects_invalid_text() {
828 let result = parse_cached_block_timestamp("not-a-timestamp");
829
830 assert!(result.is_err());
831 }
832
833 #[rstest]
834 fn decode_address_rejects_malformed_database_value() {
835 let error = decode_address("not-an-address", "token").unwrap_err();
836
837 assert!(error.to_string().contains("Invalid token address"));
838 }
839
840 #[rstest]
841 fn decode_unsigned_fields_reject_negative_database_values() {
842 assert!(decode_u64(-1_i64, "block number").is_err());
843 assert!(decode_u32(-1, "log index").is_err());
844 assert!(decode_u8(-1, "token decimals").is_err());
845 assert!(decode_u8(256, "token decimals").is_err());
846 }
847
848 #[rstest]
849 fn execution_transaction_hash_debug_redacts_raw_transaction() {
850 let row = ExecutionTransactionHashRow {
851 id: 1,
852 intent_id: 2,
853 chain_id: 42161,
854 transaction_hash: "0xhash".to_string(),
855 payload_expected: true,
856 raw_transaction: Some(vec![0xde, 0xad, 0xbe, 0xef]),
857 sealed_transaction: Some(vec![0xca, 0xfe]),
858 status: "signed".to_string(),
859 block_number: None,
860 block_hash: None,
861 receipt_success: None,
862 gas_used: None,
863 effective_gas_price: None,
864 current: true,
865 };
866
867 let debug = format!("{row:?}");
868
869 assert!(debug.contains("raw_transaction: Some(\"<redacted>\")"));
870 assert!(debug.contains("sealed_transaction: Some(\"<redacted>\")"));
871 assert!(!debug.contains("[222, 173, 190, 239]"));
872 assert!(!debug.contains("[202, 254]"));
873 }
874}