Skip to main content

nautilus_blockchain/cache/
rows.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16use 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/// A data transfer object that maps database rows to token data.
60///
61/// Implements `FromRow` trait to automatically convert PostgreSQL results into `TokenRow`
62/// objects that can be transformed into domain entity `Token` objects.
63#[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/// A data transfer object that maps database rows to block timestamp data.
167#[derive(Debug)]
168pub struct BlockTimestampRow {
169    /// The block number.
170    pub number: u64,
171    /// The block timestamp.
172    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(&timestamp).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
203/// Transforms a database row from the pool events UNION query into a DexPoolData enum variant.
204///
205/// This function directly processes a PostgreSQL row and creates the appropriate DexPoolData
206/// variant based on the event_type discriminator field, using the provided context.
207///
208/// # Errors
209///
210/// Returns an error if row field extraction fails or data validation fails.
211pub 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, // ts_event
294                timestamp, // ts_init (same block timestamp)
295                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            // UNION queries return NUMERIC type, not domain types, so we need to read as strings
337            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, // ts_event
384                timestamp, // ts_init (same block timestamp)
385            );
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            // UNION queries return NUMERIC type, not domain types, so we need to read as strings
398            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, // ts_event
435                timestamp, // ts_init (same block timestamp)
436            );
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, // ts_event
467                timestamp, // ts_init (same block timestamp)
468            );
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            // UNION queries return NUMERIC type, not domain types, so we need to read as strings
491            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, // ts_event
515                timestamp, // ts_init (same block timestamp)
516            );
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            // For flash events, we have flash_amount0, flash_amount1, flash_paid0, flash_paid1
535            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, // ts_event
569                timestamp, // ts_init (same block timestamp)
570                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/// A data transfer object that maps database rows to persisted execution transaction records.
588#[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/// Values persisted when reserving durable ownership of an execution intent.
599#[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/// A durable execution intent which owns the active signer slot and optional client order.
618#[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/// A signed transaction hash associated with an execution intent.
693#[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}