From a5555718c1ade71da3f0482c3cb1f98c5d8d5403 Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 01:32:17 +0200 Subject: [PATCH 01/26] Introduce new logic for syncing raw batch data --- requirements/prod.txt | 1 + src/fetch/orderbook.py | 76 ++++++++++ src/main.py | 15 ++ src/models/tables.py | 1 + src/sql/orderbook/batch_data.sql | 249 +++++++++++++++++++++++++++++++ src/sync/batch_data.py | 43 ++++++ src/sync/common.py | 97 +++++++++++- src/sync/config.py | 10 ++ 8 files changed, 491 insertions(+), 1 deletion(-) create mode 100644 src/sql/orderbook/batch_data.sql create mode 100644 src/sync/batch_data.py diff --git a/requirements/prod.txt b/requirements/prod.txt index 4ab4c0d1..32a1ff67 100644 --- a/requirements/prod.txt +++ b/requirements/prod.txt @@ -7,3 +7,4 @@ ndjson>=0.3.1 py-multiformats-cid>=0.4.4 boto3>=1.26.12 SQLAlchemy>=2.0,<3.0 +web3==6.20.3 diff --git a/src/fetch/orderbook.py b/src/fetch/orderbook.py index 01ce3f93..206542f7 100644 --- a/src/fetch/orderbook.py +++ b/src/fetch/orderbook.py @@ -166,6 +166,82 @@ def get_batch_rewards(cls, block_range: BlockRange) -> DataFrame: return barn.copy() return pd.DataFrame() + @classmethod + def run_batch_data_sql(cls, block_range: BlockRange) -> DataFrame: + """ + Fetches and validates Batch data DataFrame as concatenation from Prod and Staging DB + """ + batch_data_query_prod = ( + open_query("orderbook/batch_data.sql") + .replace("{{start_block}}", str(block_range.block_from)) + .replace("{{end_block}}", str(block_range.block_to)) + .replace( + "{{EPSILON_LOWER}}", "10000000000000000" + ) # lower ETH cap for payment (in WEI) + .replace( + "{{EPSILON_UPPER}}", "12000000000000000" + ) # upper ETH cap for payment (in WEI) + .replace("{{env}}", "prod") + ) + batch_data_query_barn = ( + open_query("orderbook/batch_data.sql") + .replace("{{start_block}}", str(block_range.block_from)) + .replace("{{end_block}}", str(block_range.block_to)) + .replace( + "{{EPSILON_LOWER}}", "10000000000000000" + ) # lower ETH cap for payment (in WEI) + .replace( + "{{EPSILON_UPPER}}", "12000000000000000" + ) # upper ETH cap for payment (in WEI) + .replace("{{env}}", "prod") + ) + data_types = { + # According to this: https://stackoverflow.com/a/11548224 + # capitalized int64 means `Optional` and it appears to work. + "block_number": "Int64", + "block_deadline": "int64", + } + barn, prod = cls._query_both_dbs( + batch_data_query_prod, batch_data_query_barn, data_types + ) + + # Warn if solver appear in both environments. + if not set(prod.solver).isdisjoint(set(barn.solver)): + log.warning( + f"solver overlap in {block_range}: solvers " + f"{set(prod.solver).intersection(set(barn.solver))} part of both prod and barn" + ) + + if not prod.empty and not barn.empty: + return pd.concat([prod, barn]) + if not prod.empty: + return prod.copy() + if not barn.empty: + return barn.copy() + return pd.DataFrame() + + @classmethod + def get_batch_data(cls, block_range: BlockRange) -> DataFrame: + """ + Decomposes the block range into buckets of 10k blocks each, + so as to ensure the batch data query runs fast enough. + At the end, it concatenates everything into one data frame + """ + start = block_range.block_from + end = block_range.block_to + bucket_size = 10000 + res = [] + while start < end: + size = min(end - start, bucket_size) + log.info(f"About to process block range ({start}, {start + size})") + res.append( + cls.run_batch_data_sql( + BlockRange(block_from=start, block_to=start + size) + ) + ) + start = start + size + return pd.concat(res) + @classmethod def get_app_hashes(cls) -> DataFrame: """ diff --git a/src/main.py b/src/main.py index 56b93812..aa2af8c2 100644 --- a/src/main.py +++ b/src/main.py @@ -7,6 +7,7 @@ from dotenv import load_dotenv from dune_client.client import DuneClient +from web3 import Web3 from src.fetch.orderbook import OrderbookFetcher from src.logger import set_log @@ -16,6 +17,7 @@ from src.sync import sync_price_feed from src.sync.config import SyncConfig, AppDataSyncConfig, PriceFeedSyncConfig from src.sync.order_rewards import sync_order_rewards, sync_batch_rewards +from src.sync.batch_data import sync_batch_data log = set_log(__name__) @@ -57,6 +59,7 @@ def main() -> None: request_timeout=float(os.environ.get("DUNE_API_REQUEST_TIMEOUT", 10)), ) orderbook = OrderbookFetcher() + web3 = Web3(Web3.HTTPProvider(os.environ.get("NODE_URL"))) if args.sync_table == SyncTable.APP_DATA: table = os.environ["APP_DATA_TARGET_TABLE"] @@ -98,6 +101,18 @@ def main() -> None: fetcher=orderbook, dry_run=args.dry_run, ) + elif args.sync_table == SyncTable.BATCH_DATA: + table = os.environ["BATCH_DATA_TARGET_TABLE"] + assert table, "BATCH DATA sync needs a BATCH_DATA_TARGET_TABLE env" + asyncio.run( + sync_batch_data( + web3, + orderbook, + dune=dune, + config=PriceFeedSyncConfig(table), + dry_run=args.dry_run, + ) + ) else: log.error(f"unsupported sync_table '{args.sync_table}'") diff --git a/src/models/tables.py b/src/models/tables.py index 94dc69be..51bc4d59 100644 --- a/src/models/tables.py +++ b/src/models/tables.py @@ -8,6 +8,7 @@ class SyncTable(Enum): APP_DATA = "app_data" ORDER_REWARDS = "order_rewards" BATCH_REWARDS = "batch_rewards" + BATCH_DATA = "batch_data" INTERNAL_IMBALANCE = "internal_imbalance" PRICE_FEED = "price_feed" diff --git a/src/sql/orderbook/batch_data.sql b/src/sql/orderbook/batch_data.sql new file mode 100644 index 00000000..355a0107 --- /dev/null +++ b/src/sql/orderbook/batch_data.sql @@ -0,0 +1,249 @@ +WITH observed_settlements AS ( + SELECT + -- settlement + tx_hash, + solver, + s.block_number, + -- settlement_observations + effective_gas_price * gas_used AS execution_cost, + surplus, + s.auction_id + FROM + settlement_observations so + JOIN settlements s ON s.block_number = so.block_number + AND s.log_index = so.log_index + JOIN settlement_scores ss ON s.auction_id = ss.auction_id + WHERE + ss.block_deadline > {{start_block}} + AND ss.block_deadline <= {{end_block}} +), +-- order data +order_data AS ( + SELECT + uid, + sell_token, + buy_token, + sell_amount, + buy_amount, + kind, + app_data + FROM orders + UNION ALL + SELECT + uid, + sell_token, + buy_token, + sell_amount, + buy_amount, + kind, + app_data + FROM jit_orders +), +-- unprocessed trade data +trade_data_unprocessed AS ( + SELECT + ss.winner AS solver, + s.auction_id, + s.tx_hash, + t.order_uid, + od.sell_token, + od.buy_token, + t.sell_amount, -- the total amount the user sends + t.buy_amount, -- the total amount the user receives + oe.surplus_fee AS observed_fee, -- the total discrepancy between what the user sends and what they would have send if they traded at clearing price + od.kind, + CASE + WHEN od.kind = 'sell' THEN od.buy_token + WHEN od.kind = 'buy' THEN od.sell_token + END AS surplus_token, + convert_from(ad.full_app_data, 'UTF8')::JSONB->'metadata'->'partnerFee'->>'recipient' AS partner_fee_recipient, + COALESCE(oe.protocol_fee_amounts[1], 0) AS first_protocol_fee_amount, + COALESCE(oe.protocol_fee_amounts[2], 0) AS second_protocol_fee_amount + FROM + settlements s + JOIN settlement_scores ss -- contains block_deadline + ON s.auction_id = ss.auction_id + JOIN trades t -- contains traded amounts + ON s.block_number = t.block_number -- given the join that follows with the order execution table, this works even when multiple txs appear in the same block + JOIN order_data od -- contains tokens and limit amounts + ON t.order_uid = od.uid + JOIN order_execution oe -- contains surplus fee + ON t.order_uid = oe.order_uid + AND s.auction_id = oe.auction_id + LEFT OUTER JOIN app_data ad -- contains full app data + ON od.app_data = ad.contract_app_data + WHERE + ss.block_deadline > {{start_block}} + AND ss.block_deadline <= {{end_block}} +), +-- processed trade data: +trade_data_processed AS ( + SELECT + auction_id, + solver, + tx_hash, + order_uid, + sell_amount, + buy_amount, + sell_token, + observed_fee, + surplus_token, + second_protocol_fee_amount, + first_protocol_fee_amount + second_protocol_fee_amount AS protocol_fee, + partner_fee_recipient, + CASE + WHEN partner_fee_recipient IS NOT NULL THEN second_protocol_fee_amount + ELSE 0 + END AS partner_fee, + surplus_token AS protocol_fee_token + FROM + trade_data_unprocessed +), +price_data AS ( + SELECT + tdp.auction_id, + tdp.order_uid, + ap_surplus.price / pow(10, 18) AS surplus_token_native_price, + ap_protocol.price / pow(10, 18) AS protocol_fee_token_native_price, + ap_sell.price / pow(10, 18) AS network_fee_token_native_price + FROM + trade_data_processed AS tdp + LEFT OUTER JOIN auction_prices ap_sell -- contains price: sell token + ON tdp.auction_id = ap_sell.auction_id + AND tdp.sell_token = ap_sell.token + LEFT OUTER JOIN auction_prices ap_surplus -- contains price: surplus token + ON tdp.auction_id = ap_surplus.auction_id + AND tdp.surplus_token = ap_surplus.token + LEFT OUTER JOIN auction_prices ap_protocol -- contains price: protocol fee token + ON tdp.auction_id = ap_protocol.auction_id + AND tdp.surplus_token = ap_protocol.token +), +trade_data_processed_with_prices AS ( + SELECT + tdp.auction_id, + tdp.solver, + tdp.tx_hash, + tdp.order_uid, + tdp.surplus_token, + tdp.protocol_fee, + tdp.protocol_fee_token, + tdp.partner_fee, + tdp.partner_fee_recipient, + CASE + WHEN tdp.sell_token != tdp.surplus_token THEN tdp.observed_fee - (tdp.sell_amount - tdp.observed_fee) / tdp.buy_amount * COALESCE(tdp.protocol_fee, 0) + ELSE tdp.observed_fee - COALESCE(tdp.protocol_fee, 0) + END AS network_fee, + tdp.sell_token AS network_fee_token, + surplus_token_native_price, + protocol_fee_token_native_price, + network_fee_token_native_price + FROM + trade_data_processed AS tdp + JOIN price_data pd + ON tdp.auction_id = pd.auction_id + AND tdp.order_uid = pd.order_uid +), +batch_protocol_fees AS ( + SELECT + solver, + tx_hash, + sum(protocol_fee * protocol_fee_token_native_price) AS protocol_fee + FROM + trade_data_processed_with_prices + GROUP BY + solver, + tx_hash +), +batch_network_fees AS ( + SELECT + solver, + tx_hash, + sum(network_fee * network_fee_token_native_price) AS network_fee + FROM + trade_data_processed_with_prices + GROUP BY + solver, + tx_hash +), +reward_data AS ( + SELECT + -- observations + os.tx_hash, + ss.auction_id, + -- TODO - Assuming that `solver == winner` when both not null + -- We will need to monitor that `solver == winner`! + ss.winner as solver, + block_number as settlement_block, + block_deadline, + coalesce(execution_cost, 0) as execution_cost, + coalesce(surplus, 0) as surplus, + -- scores + winning_score, + case + when block_number is not null + and block_number <= block_deadline + 1 then winning_score -- this includes a grace period of one block for settling a batch + else 0 + end as observed_score, + reference_score, + -- protocol_fees + coalesce(cast(protocol_fee as numeric(78, 0)), 0) as protocol_fee, + coalesce( + cast(network_fee as numeric(78, 0)), + 0 + ) as network_fee + FROM + settlement_scores ss + -- outer joins made in order to capture non-existent settlements. + LEFT OUTER JOIN observed_settlements os ON os.auction_id = ss.auction_id + LEFT OUTER JOIN batch_protocol_fees bpf ON bpf.tx_hash = os.tx_hash + LEFT OUTER JOIN batch_network_fees bnf ON bnf.tx_hash = os.tx_hash + WHERE + ss.block_deadline > {{start_block}} + AND ss.block_deadline <= {{end_block}} +), +reward_per_auction as ( + SELECT + tx_hash, + auction_id, + settlement_block, + block_deadline, + solver, + execution_cost, + surplus, + protocol_fee, -- the protocol fee + network_fee, -- the network fee + observed_score - reference_score as uncapped_payment, + -- Capped Reward = CLAMP_[-E, E + exec_cost](uncapped_reward_eth) + LEAST( + GREATEST( + - {{EPSILON_LOWER}}, + observed_score - reference_score + ), + {{EPSILON_UPPER}} + ) as capped_payment, + winning_score, + reference_score + FROM + reward_data +) +SELECT + '{{env}}' as environment, + auction_id, + settlement_block as block_number, + block_deadline, + case + when tx_hash is NULL then NULL + else concat('0x', encode(tx_hash, 'hex')) + end as tx_hash, + concat('0x', encode(solver, 'hex')) as solver, + execution_cost :: text as execution_cost, + surplus :: text as surplus, + protocol_fee :: text as protocol_fee, + network_fee :: text as network_fee, + uncapped_payment :: text as uncapped_payment_eth, + capped_payment :: text as capped_payment, + winning_score :: text as winning_score, + reference_score :: text as reference_score +FROM + reward_per_auction +ORDER BY block_deadline diff --git a/src/sync/batch_data.py b/src/sync/batch_data.py new file mode 100644 index 00000000..4ab31edc --- /dev/null +++ b/src/sync/batch_data.py @@ -0,0 +1,43 @@ +"""Main Entry point for price feed sync""" +from dune_client.client import DuneClient +import web3 + +from src.fetch.orderbook import OrderbookFetcher +from src.logger import set_log +from src.sync.config import BatchDataSyncConfig +from src.sync.common import compute_block_and_month_range +from src.models.block_range import BlockRange + + +log = set_log(__name__) + + +async def sync_batch_data( + node: web3, + orderbook: OrderbookFetcher, + dune: DuneClient, + config: BatchDataSyncConfig, + dry_run: bool, +) -> None: + """Batch data Sync Logic""" + block_range_list, months_list = compute_block_and_month_range(node) + for i in range(len(block_range_list)): + start_block = block_range_list[i][0] + end_block = block_range_list[i][1] + table_name = "test_batch_rewards_ethereum_" + months_list[i] + block_range = BlockRange(block_from=start_block, block_to=end_block) + log.info( + f"About to process block range ({start_block}, {end_block}) for month {months_list[i]}" + ) + batch_data = orderbook.get_batch_data(block_range) + log.info("SQL query successfully executed. About to initiate upload to Dune.") + if not dry_run: + dune.upload_csv( + data=batch_data.to_csv(index=False), + table_name=table_name, # config.table, + description=config.description, + is_private=False, + ) + log.info( + f"batch data sync run completed successfully for month {months_list[i]}" + ) diff --git a/src/sync/common.py b/src/sync/common.py index c3a11195..698b9bcd 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -1,5 +1,8 @@ """Shared methods between both sync scripts.""" - +from datetime import datetime, timezone +import time +from dateutil import tz +from web3 import Web3 from src.logger import set_log from src.models.tables import SyncTable from src.post.aws import AWSClient @@ -18,3 +21,95 @@ def last_sync_block(aws: AWSClient, table: SyncTable, genesis_block: int = 0) -> block_from = genesis_block return block_from + + +def find_block_with_timestamp(node, time_stamp): + """ + This implements binary search and returns the smallest block number + whose timestamp is at least as large as the time_stamp argument passed in the function + """ + block_found = False + end_block_number = node.eth.get_block("finalized").number + start_block_number = 1 + close_in_seconds = 30 + + while not block_found: + mid_block_number = (start_block_number + end_block_number) // 2 + block = node.eth.get_block(mid_block_number) + block_time = block.timestamp + difference_in_seconds = int((time_stamp - block_time)) + + if abs(difference_in_seconds) < close_in_seconds: + block_found = True + continue + + if difference_in_seconds < 0: + end_block_number = mid_block_number - 1 + else: + start_block_number = mid_block_number + 1 + + ## we now brute-force to ensure we have found the right block + for b in range(mid_block_number - 100, mid_block_number + 100): + block = node.eth.get_block(b) + block_time_stamp = block.timestamp + if block_time_stamp >= time_stamp: + return block.number + + +def compute_block_range_from_timestamps(node, start_timestamp, end_timestamp): + start_block = find_block_with_timestamp(node, start_timestamp) + end_block = find_block_with_timestamp(node, end_timestamp) + if node.eth.get_block(end_block).timestamp > end_timestamp: + end_block = end_block - 1 + return start_block, end_block + + +def compute_block_and_month_range(node: Web3): + # We first compute the relevant block range + # Here, we assume that the job runs at least once every 24h + # Because of that, if it is the first day of month, we also + # compute the previous month's table just to be on the safe side + + latest_finalized_block = node.eth.get_block("finalized") + + current_month_end_block = latest_finalized_block.number + current_month_end_timestamp = latest_finalized_block.timestamp + + current_month_end_datetime = datetime.fromtimestamp( + current_month_end_timestamp, tz=timezone.utc + ) + current_month_start_datetime = datetime( + current_month_end_datetime.year, current_month_end_datetime.month, 1, 00, 00 + ) + current_month_start_timestamp = current_month_start_datetime.replace( + tzinfo=timezone.utc + ).timestamp() + + current_month_start_block = find_block_with_timestamp( + node, current_month_start_timestamp + ) + + current_month = ( + f"{current_month_end_datetime.year}_{current_month_end_datetime.month}" + ) + months_list = [current_month] + block_range = [(current_month_start_block, current_month_end_block)] + if current_month_end_datetime.day == 1: + if current_month_end_datetime.month == 1: + previous_month = f"{current_month_end_datetime.year - 1}_12" + previous_month_start_datetime = datetime( + current_month_end_datetime.year - 1, 12, 1, 00, 00 + ) + else: + previous_month = f"{current_month_end_datetime.year}_{current_month_end_datetime.month - 1}" + months_list.append(previous_month) + previous_month_start_timestamp = previous_month_start_datetime.replace( + tzinfo=timezone.utc + ).timestamp() + previous_month_start_block = find_block_with_timestamp( + node, previous_month_start_timestamp + ) + previous_month_end_block = current_month_start_block + block_range.append((previous_month_start_block, previous_month_end_block)) + + return block_range, months_list diff --git a/src/sync/config.py b/src/sync/config.py index 3be39788..6a84d3fb 100644 --- a/src/sync/config.py +++ b/src/sync/config.py @@ -40,3 +40,13 @@ class PriceFeedSyncConfig: description: str = ( "Table containing prices and timestamps from multiple price feeds" ) + + +@dataclass +class BatchDataSyncConfig: + """Configuration for batch data sync.""" + + # The name of the table to upload to + table: str = "batch_data_test" + # Description of the table (for creation) + description: str = "Table containing raw batch data" From fc2c654d55f6cf7dd518f1845ccd08f3a169b467 Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 01:38:02 +0200 Subject: [PATCH 02/26] fix typo --- src/main.py | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/src/main.py b/src/main.py index aa2af8c2..ad9bd8cd 100644 --- a/src/main.py +++ b/src/main.py @@ -15,7 +15,12 @@ from src.post.aws import AWSClient from src.sync import sync_app_data from src.sync import sync_price_feed -from src.sync.config import SyncConfig, AppDataSyncConfig, PriceFeedSyncConfig +from src.sync.config import ( + SyncConfig, + AppDataSyncConfig, + PriceFeedSyncConfig, + BatchDataSyncConfig, +) from src.sync.order_rewards import sync_order_rewards, sync_batch_rewards from src.sync.batch_data import sync_batch_data @@ -109,7 +114,7 @@ def main() -> None: web3, orderbook, dune=dune, - config=PriceFeedSyncConfig(table), + config=BatchDataSyncConfig(table), dry_run=args.dry_run, ) ) From 2b6b5e6b82f95063a656592b7533804ed4402df1 Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 01:39:21 +0200 Subject: [PATCH 03/26] remove redundant code --- src/sync/common.py | 8 -------- 1 file changed, 8 deletions(-) diff --git a/src/sync/common.py b/src/sync/common.py index 698b9bcd..79917f98 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -56,14 +56,6 @@ def find_block_with_timestamp(node, time_stamp): return block.number -def compute_block_range_from_timestamps(node, start_timestamp, end_timestamp): - start_block = find_block_with_timestamp(node, start_timestamp) - end_block = find_block_with_timestamp(node, end_timestamp) - if node.eth.get_block(end_block).timestamp > end_timestamp: - end_block = end_block - 1 - return start_block, end_block - - def compute_block_and_month_range(node: Web3): # We first compute the relevant block range # Here, we assume that the job runs at least once every 24h From 8111a6eef70e1fc4fc95e9176b954b1c3a5e113b Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 01:41:07 +0200 Subject: [PATCH 04/26] fix another typo --- src/fetch/orderbook.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/fetch/orderbook.py b/src/fetch/orderbook.py index 206542f7..345eaa56 100644 --- a/src/fetch/orderbook.py +++ b/src/fetch/orderbook.py @@ -193,7 +193,7 @@ def run_batch_data_sql(cls, block_range: BlockRange) -> DataFrame: .replace( "{{EPSILON_UPPER}}", "12000000000000000" ) # upper ETH cap for payment (in WEI) - .replace("{{env}}", "prod") + .replace("{{env}}", "barn") ) data_types = { # According to this: https://stackoverflow.com/a/11548224 From d7f84dd260749ccdcdf98522b24be185f372294c Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 02:23:02 +0200 Subject: [PATCH 05/26] small fixes --- src/sync/batch_data.py | 9 ++++++--- src/sync/common.py | 8 +++++--- 2 files changed, 11 insertions(+), 6 deletions(-) diff --git a/src/sync/batch_data.py b/src/sync/batch_data.py index 4ab31edc..ada44537 100644 --- a/src/sync/batch_data.py +++ b/src/sync/batch_data.py @@ -1,7 +1,7 @@ """Main Entry point for price feed sync""" from dune_client.client import DuneClient import web3 - +import os from src.fetch.orderbook import OrderbookFetcher from src.logger import set_log from src.sync.config import BatchDataSyncConfig @@ -20,11 +20,14 @@ async def sync_batch_data( dry_run: bool, ) -> None: """Batch data Sync Logic""" + network = os.environ["NETWORK"] + if network == "mainnet": + network = "ethereum" block_range_list, months_list = compute_block_and_month_range(node) for i in range(len(block_range_list)): start_block = block_range_list[i][0] end_block = block_range_list[i][1] - table_name = "test_batch_rewards_ethereum_" + months_list[i] + table_name = config.table + "_" + network + "_" + months_list[i] block_range = BlockRange(block_from=start_block, block_to=end_block) log.info( f"About to process block range ({start_block}, {end_block}) for month {months_list[i]}" @@ -34,7 +37,7 @@ async def sync_batch_data( if not dry_run: dune.upload_csv( data=batch_data.to_csv(index=False), - table_name=table_name, # config.table, + table_name=table_name, description=config.description, is_private=False, ) diff --git a/src/sync/common.py b/src/sync/common.py index 79917f98..aa248b46 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -1,7 +1,5 @@ """Shared methods between both sync scripts.""" from datetime import datetime, timezone -import time -from dateutil import tz from web3 import Web3 from src.logger import set_log from src.models.tables import SyncTable @@ -23,7 +21,7 @@ def last_sync_block(aws: AWSClient, table: SyncTable, genesis_block: int = 0) -> return block_from -def find_block_with_timestamp(node, time_stamp): +def find_block_with_timestamp(node, time_stamp) -> int: """ This implements binary search and returns the smallest block number whose timestamp is at least as large as the time_stamp argument passed in the function @@ -57,6 +55,10 @@ def find_block_with_timestamp(node, time_stamp): def compute_block_and_month_range(node: Web3): + """ + This determines the block range and the relevant months + for which we will compute and upload data on Dune. + """ # We first compute the relevant block range # Here, we assume that the job runs at least once every 24h # Because of that, if it is the first day of month, we also From 7900e3556c6b215db3d780de21f11cf086579793 Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 02:29:20 +0200 Subject: [PATCH 06/26] bug fix --- src/sync/batch_data.py | 4 ++-- src/sync/common.py | 11 ++++++++++- 2 files changed, 12 insertions(+), 3 deletions(-) diff --git a/src/sync/batch_data.py b/src/sync/batch_data.py index ada44537..0dbff89e 100644 --- a/src/sync/batch_data.py +++ b/src/sync/batch_data.py @@ -1,7 +1,7 @@ -"""Main Entry point for price feed sync""" +"""Main Entry point for batch data sync""" +import os from dune_client.client import DuneClient import web3 -import os from src.fetch.orderbook import OrderbookFetcher from src.logger import set_log from src.sync.config import BatchDataSyncConfig diff --git a/src/sync/common.py b/src/sync/common.py index aa248b46..94da1db1 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -95,7 +95,16 @@ def compute_block_and_month_range(node: Web3): current_month_end_datetime.year - 1, 12, 1, 00, 00 ) else: - previous_month = f"{current_month_end_datetime.year}_{current_month_end_datetime.month - 1}" + previous_month = f"""{current_month_end_datetime.year}_ + {current_month_end_datetime.month - 1} + """ + previous_month_start_datetime = datetime( + current_month_end_datetime.year, + current_month_end_datetime.month - 1, + 1, + 00, + 00, + ) months_list.append(previous_month) previous_month_start_timestamp = previous_month_start_datetime.replace( tzinfo=timezone.utc From 02f5089a52b2eca342be111af76863f6830660ae Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 16:40:53 +0200 Subject: [PATCH 07/26] small adaptations and fixes --- src/fetch/orderbook.py | 4 ++-- src/sync/batch_data.py | 17 ++++++++++++----- 2 files changed, 14 insertions(+), 7 deletions(-) diff --git a/src/fetch/orderbook.py b/src/fetch/orderbook.py index 345eaa56..036608ec 100644 --- a/src/fetch/orderbook.py +++ b/src/fetch/orderbook.py @@ -45,7 +45,7 @@ class OrderbookFetcher: def _pg_engine(db_env: OrderbookEnv) -> Engine: """Returns a connection to postgres database""" load_dotenv() - db_url = os.environ[f"{db_env}_DB_URL"] + db_url = os.environ[f"{db_env}_DB_URL"] + "/" + os.environ["NETWORK"] db_string = f"postgresql+psycopg2://{db_url}" return create_engine(db_string) @@ -229,7 +229,7 @@ def get_batch_data(cls, block_range: BlockRange) -> DataFrame: """ start = block_range.block_from end = block_range.block_to - bucket_size = 10000 + bucket_size = 20000 res = [] while start < end: size = min(end - start, bucket_size) diff --git a/src/sync/batch_data.py b/src/sync/batch_data.py index 0dbff89e..c0e69633 100644 --- a/src/sync/batch_data.py +++ b/src/sync/batch_data.py @@ -21,13 +21,20 @@ async def sync_batch_data( ) -> None: """Batch data Sync Logic""" network = os.environ["NETWORK"] + if network == "mainnet": - network = "ethereum" + network_name = "ethereum" + else: + if network == "xdai": + network_name = "gnosis" + else: + network_name = "arbitrum" + block_range_list, months_list = compute_block_and_month_range(node) for i in range(len(block_range_list)): start_block = block_range_list[i][0] end_block = block_range_list[i][1] - table_name = config.table + "_" + network + "_" + months_list[i] + table_name = config.table + "_" + network_name + "_" + months_list[i] block_range = BlockRange(block_from=start_block, block_to=end_block) log.info( f"About to process block range ({start_block}, {end_block}) for month {months_list[i]}" @@ -41,6 +48,6 @@ async def sync_batch_data( description=config.description, is_private=False, ) - log.info( - f"batch data sync run completed successfully for month {months_list[i]}" - ) + log.info( + f"batch data sync run completed successfully for month {months_list[i]}" + ) From 3685c384e0c0de369a168b066175af8bbab765cf Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 17:02:42 +0200 Subject: [PATCH 08/26] small fixes --- src/fetch/orderbook.py | 4 +++- src/main.py | 11 ++++++++++- 2 files changed, 13 insertions(+), 2 deletions(-) diff --git a/src/fetch/orderbook.py b/src/fetch/orderbook.py index 036608ec..13e91625 100644 --- a/src/fetch/orderbook.py +++ b/src/fetch/orderbook.py @@ -45,7 +45,9 @@ class OrderbookFetcher: def _pg_engine(db_env: OrderbookEnv) -> Engine: """Returns a connection to postgres database""" load_dotenv() - db_url = os.environ[f"{db_env}_DB_URL"] + "/" + os.environ["NETWORK"] + db_url = ( + os.environ[f"{db_env}_DB_URL"] + "/" + os.environ.get("NETWORK", "mainnet") + ) db_string = f"postgresql+psycopg2://{db_url}" return create_engine(db_string) diff --git a/src/main.py b/src/main.py index ad9bd8cd..c4d0978d 100644 --- a/src/main.py +++ b/src/main.py @@ -64,7 +64,16 @@ def main() -> None: request_timeout=float(os.environ.get("DUNE_API_REQUEST_TIMEOUT", 10)), ) orderbook = OrderbookFetcher() - web3 = Web3(Web3.HTTPProvider(os.environ.get("NODE_URL"))) + network = os.environ.get("NETWORK", "mainnet") + if network == "mainnet": + node_suffix = "MAINNET" + else: + if network == "xdai": + node_suffix = "GNOSIS" + else: + if network == "arbitrum-one": + node_suffix = "ARBITRUM" + web3 = Web3(Web3.HTTPProvider(os.environ.get("NODE_URL" + "_" + node_suffix))) if args.sync_table == SyncTable.APP_DATA: table = os.environ["APP_DATA_TARGET_TABLE"] From 63e043558c0b89c3333bb31377917cb80413eca3 Mon Sep 17 00:00:00 2001 From: harisang Date: Fri, 22 Nov 2024 17:22:04 +0200 Subject: [PATCH 09/26] hardcode gnosis chain caps for testing purposes --- src/sql/orderbook/batch_data.sql | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/sql/orderbook/batch_data.sql b/src/sql/orderbook/batch_data.sql index 355a0107..a008eda7 100644 --- a/src/sql/orderbook/batch_data.sql +++ b/src/sql/orderbook/batch_data.sql @@ -216,10 +216,10 @@ reward_per_auction as ( -- Capped Reward = CLAMP_[-E, E + exec_cost](uncapped_reward_eth) LEAST( GREATEST( - - {{EPSILON_LOWER}}, + - 10000000000000000000, observed_score - reference_score ), - {{EPSILON_UPPER}} + 20000000000000000000 ) as capped_payment, winning_score, reference_score From 45550a4ccc16a6d8cceb063ea4c847047dd9d5d4 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 01:30:14 +0200 Subject: [PATCH 10/26] add short log --- src/main.py | 1 + 1 file changed, 1 insertion(+) diff --git a/src/main.py b/src/main.py index c4d0978d..b04a92e8 100644 --- a/src/main.py +++ b/src/main.py @@ -65,6 +65,7 @@ def main() -> None: ) orderbook = OrderbookFetcher() network = os.environ.get("NETWORK", "mainnet") + log.info(f"Network is set to: {network}") if network == "mainnet": node_suffix = "MAINNET" else: From 22a8c27d16fa4b11af0f517ee3e12569ce245fb8 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 03:09:07 +0200 Subject: [PATCH 11/26] various small changes --- .env.sample | 18 ++++++++++++++++++ src/fetch/orderbook.py | 4 +++- src/main.py | 13 ++++--------- src/sync/batch_data.py | 13 ++++--------- src/sync/common.py | 12 ++++++++++++ 5 files changed, 41 insertions(+), 19 deletions(-) diff --git a/.env.sample b/.env.sample index 30c8c807..880ea871 100644 --- a/.env.sample +++ b/.env.sample @@ -28,3 +28,21 @@ ETHERSCAN_API_KEY= # Query ID for the aggregate query on Dune AGGREGATE_QUERY_ID= + +# Node for each chain we want to run a sync job on +NODE_URL_MAINNET= +NODE_URL_GNOSIS= +NODE_URL_ARBITRUM= + +# The network which we run a sync job on. +# Current options are: {"mainnet", "xdai", "arbitrum-one"} +NETWORK= + +# The prefix of the dune table where we sync the various mainnet price feeds +PRICE_FEED_TARGET_TABLE + +# The prefix of the dune table where we sync the monthly raw batch data +BATCH_DATA_TARGET_TABLE = "batch_data" + +# Dune api timeout parameter, where we recommend to set it to 600 +DUNE_API_REQUEST_TIMEOUT = 600 diff --git a/src/fetch/orderbook.py b/src/fetch/orderbook.py index 13e91625..5e929cc6 100644 --- a/src/fetch/orderbook.py +++ b/src/fetch/orderbook.py @@ -19,6 +19,7 @@ log = set_log(__name__) MAX_PROCESSING_DELAY = 10 +BUCKET_SIZE = {"mainnet": 10000, "xdai": 30000, "arbitrum-one": 1000000} class OrderbookEnv(Enum): @@ -229,9 +230,10 @@ def get_batch_data(cls, block_range: BlockRange) -> DataFrame: so as to ensure the batch data query runs fast enough. At the end, it concatenates everything into one data frame """ + load_dotenv() start = block_range.block_from end = block_range.block_to - bucket_size = 20000 + bucket_size = BUCKET_SIZE[os.environ.get("NETWORK", "mainnet")] res = [] while start < end: size = min(end - start, bucket_size) diff --git a/src/main.py b/src/main.py index b04a92e8..7276ec1b 100644 --- a/src/main.py +++ b/src/main.py @@ -23,6 +23,7 @@ ) from src.sync.order_rewards import sync_order_rewards, sync_batch_rewards from src.sync.batch_data import sync_batch_data +from src.sync.common import node_suffix log = set_log(__name__) @@ -66,15 +67,9 @@ def main() -> None: orderbook = OrderbookFetcher() network = os.environ.get("NETWORK", "mainnet") log.info(f"Network is set to: {network}") - if network == "mainnet": - node_suffix = "MAINNET" - else: - if network == "xdai": - node_suffix = "GNOSIS" - else: - if network == "arbitrum-one": - node_suffix = "ARBITRUM" - web3 = Web3(Web3.HTTPProvider(os.environ.get("NODE_URL" + "_" + node_suffix))) + web3 = Web3( + Web3.HTTPProvider(os.environ.get("NODE_URL" + "_" + node_suffix(network))) + ) if args.sync_table == SyncTable.APP_DATA: table = os.environ["APP_DATA_TARGET_TABLE"] diff --git a/src/sync/batch_data.py b/src/sync/batch_data.py index c0e69633..385c5658 100644 --- a/src/sync/batch_data.py +++ b/src/sync/batch_data.py @@ -1,11 +1,12 @@ """Main Entry point for batch data sync""" import os +from dotenv import load_dotenv from dune_client.client import DuneClient import web3 from src.fetch.orderbook import OrderbookFetcher from src.logger import set_log from src.sync.config import BatchDataSyncConfig -from src.sync.common import compute_block_and_month_range +from src.sync.common import compute_block_and_month_range, node_suffix from src.models.block_range import BlockRange @@ -20,15 +21,9 @@ async def sync_batch_data( dry_run: bool, ) -> None: """Batch data Sync Logic""" + load_dotenv() network = os.environ["NETWORK"] - - if network == "mainnet": - network_name = "ethereum" - else: - if network == "xdai": - network_name = "gnosis" - else: - network_name = "arbitrum" + network_name = node_suffix(network).lower() block_range_list, months_list = compute_block_and_month_range(node) for i in range(len(block_range_list)): diff --git a/src/sync/common.py b/src/sync/common.py index 94da1db1..1a174855 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -21,6 +21,18 @@ def last_sync_block(aws: AWSClient, table: SyncTable, genesis_block: int = 0) -> return block_from +def node_suffix(network: str) -> str: + if network == "mainnet": + return "MAINNET" + else: + if network == "xdai": + return "GNOSIS" + else: + if network == "arbitrum-one": + return "ARBITRUM" + return "" + + def find_block_with_timestamp(node, time_stamp) -> int: """ This implements binary search and returns the smallest block number From 0b94a49f6bfc77fd2d49401463785cc7022360c6 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 03:19:28 +0200 Subject: [PATCH 12/26] fix some pylint issues --- src/sync/batch_data.py | 2 +- src/sync/common.py | 21 ++++++++++----------- 2 files changed, 11 insertions(+), 12 deletions(-) diff --git a/src/sync/batch_data.py b/src/sync/batch_data.py index 385c5658..9aeeb743 100644 --- a/src/sync/batch_data.py +++ b/src/sync/batch_data.py @@ -26,7 +26,7 @@ async def sync_batch_data( network_name = node_suffix(network).lower() block_range_list, months_list = compute_block_and_month_range(node) - for i in range(len(block_range_list)): + for i, _ in enumerate(block_range_list): start_block = block_range_list[i][0] end_block = block_range_list[i][1] table_name = config.table + "_" + network_name + "_" + months_list[i] diff --git a/src/sync/common.py b/src/sync/common.py index 1a174855..ddb49e17 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -22,14 +22,15 @@ def last_sync_block(aws: AWSClient, table: SyncTable, genesis_block: int = 0) -> def node_suffix(network: str) -> str: + """ + Converts network internal name to name used for nodes and dune tables + """ if network == "mainnet": return "MAINNET" - else: - if network == "xdai": - return "GNOSIS" - else: - if network == "arbitrum-one": - return "ARBITRUM" + if network == "xdai": + return "GNOSIS" + if network == "arbitrum-one": + return "ARBITRUM" return "" @@ -38,20 +39,18 @@ def find_block_with_timestamp(node, time_stamp) -> int: This implements binary search and returns the smallest block number whose timestamp is at least as large as the time_stamp argument passed in the function """ - block_found = False end_block_number = node.eth.get_block("finalized").number start_block_number = 1 close_in_seconds = 30 - while not block_found: + while True: mid_block_number = (start_block_number + end_block_number) // 2 block = node.eth.get_block(mid_block_number) block_time = block.timestamp difference_in_seconds = int((time_stamp - block_time)) if abs(difference_in_seconds) < close_in_seconds: - block_found = True - continue + break if difference_in_seconds < 0: end_block_number = mid_block_number - 1 @@ -59,7 +58,7 @@ def find_block_with_timestamp(node, time_stamp) -> int: start_block_number = mid_block_number + 1 ## we now brute-force to ensure we have found the right block - for b in range(mid_block_number - 100, mid_block_number + 100): + for b in range(mid_block_number - 200, mid_block_number + 200): block = node.eth.get_block(b) block_time_stamp = block.timestamp if block_time_stamp >= time_stamp: From 9b2cea7575f44c0528bccdafd56ece563882e710 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 03:27:09 +0200 Subject: [PATCH 13/26] use old way of setting db url if network variable is missing --- src/fetch/orderbook.py | 11 ++++++++--- 1 file changed, 8 insertions(+), 3 deletions(-) diff --git a/src/fetch/orderbook.py b/src/fetch/orderbook.py index 5e929cc6..90ba1ef0 100644 --- a/src/fetch/orderbook.py +++ b/src/fetch/orderbook.py @@ -46,9 +46,14 @@ class OrderbookFetcher: def _pg_engine(db_env: OrderbookEnv) -> Engine: """Returns a connection to postgres database""" load_dotenv() - db_url = ( - os.environ[f"{db_env}_DB_URL"] + "/" + os.environ.get("NETWORK", "mainnet") - ) + if "NETWORK" in os.environ: + db_url = ( + os.environ[f"{db_env}_DB_URL"] + + "/" + + os.environ.get("NETWORK", "mainnet") + ) + else: + db_url = os.environ[f"{db_env}_DB_URL"] db_string = f"postgresql+psycopg2://{db_url}" return create_engine(db_string) From 986f286c28b9dda452d8c1e2be767df0f9b91ac6 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 03:33:12 +0200 Subject: [PATCH 14/26] fix more pylint issues --- src/sync/common.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/src/sync/common.py b/src/sync/common.py index ddb49e17..10dfbe37 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -63,9 +63,11 @@ def find_block_with_timestamp(node, time_stamp) -> int: block_time_stamp = block.timestamp if block_time_stamp >= time_stamp: return block.number + # fallback if correct block is not found + return mid_block_number + 200 -def compute_block_and_month_range(node: Web3): +def compute_block_and_month_range(node: Web3): # pylint: disable=too-many-locals """ This determines the block range and the relevant months for which we will compute and upload data on Dune. From fe26674421a63da521b5f3757287ba9be067bd00 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 03:39:49 +0200 Subject: [PATCH 15/26] more pylint fixes --- src/sync/batch_data.py | 4 ++-- src/sync/common.py | 6 ++---- 2 files changed, 4 insertions(+), 6 deletions(-) diff --git a/src/sync/batch_data.py b/src/sync/batch_data.py index 9aeeb743..96867d7d 100644 --- a/src/sync/batch_data.py +++ b/src/sync/batch_data.py @@ -2,7 +2,7 @@ import os from dotenv import load_dotenv from dune_client.client import DuneClient -import web3 +from web3 import Web3 from src.fetch.orderbook import OrderbookFetcher from src.logger import set_log from src.sync.config import BatchDataSyncConfig @@ -14,7 +14,7 @@ async def sync_batch_data( - node: web3, + node: Web3, orderbook: OrderbookFetcher, dune: DuneClient, config: BatchDataSyncConfig, diff --git a/src/sync/common.py b/src/sync/common.py index 10dfbe37..de79a561 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -34,7 +34,7 @@ def node_suffix(network: str) -> str: return "" -def find_block_with_timestamp(node, time_stamp) -> int: +def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: """ This implements binary search and returns the smallest block number whose timestamp is at least as large as the time_stamp argument passed in the function @@ -62,9 +62,7 @@ def find_block_with_timestamp(node, time_stamp) -> int: block = node.eth.get_block(b) block_time_stamp = block.timestamp if block_time_stamp >= time_stamp: - return block.number - # fallback if correct block is not found - return mid_block_number + 200 + return int(block.number) def compute_block_and_month_range(node: Web3): # pylint: disable=too-many-locals From 31e1824e833897cfce7aeaece9f2f75ab699397e Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 03:44:54 +0200 Subject: [PATCH 16/26] final pylint fix --- src/sync/common.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/src/sync/common.py b/src/sync/common.py index de79a561..5bd64ffa 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -63,6 +63,9 @@ def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: block_time_stamp = block.timestamp if block_time_stamp >= time_stamp: return int(block.number) + # fallback in case correct block number hasn't been found + # in that case, we will include some more blocks than necessary + return mid_block_number + 200 def compute_block_and_month_range(node: Web3): # pylint: disable=too-many-locals From 14586bb79edf6937562707415225ff1e3442feec Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 03:57:58 +0200 Subject: [PATCH 17/26] mypy fixes --- src/sync/common.py | 15 +++++++++------ 1 file changed, 9 insertions(+), 6 deletions(-) diff --git a/src/sync/common.py b/src/sync/common.py index 5bd64ffa..dd5ed084 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -1,5 +1,6 @@ """Shared methods between both sync scripts.""" from datetime import datetime, timezone +from typing import List, Tuple from web3 import Web3 from src.logger import set_log from src.models.tables import SyncTable @@ -39,14 +40,14 @@ def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: This implements binary search and returns the smallest block number whose timestamp is at least as large as the time_stamp argument passed in the function """ - end_block_number = node.eth.get_block("finalized").number + end_block_number = node.eth.get_block("finalized")["number"] start_block_number = 1 close_in_seconds = 30 while True: mid_block_number = (start_block_number + end_block_number) // 2 block = node.eth.get_block(mid_block_number) - block_time = block.timestamp + block_time = block["timestamp"] difference_in_seconds = int((time_stamp - block_time)) if abs(difference_in_seconds) < close_in_seconds: @@ -60,15 +61,17 @@ def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: ## we now brute-force to ensure we have found the right block for b in range(mid_block_number - 200, mid_block_number + 200): block = node.eth.get_block(b) - block_time_stamp = block.timestamp + block_time_stamp = block["timestamp"] if block_time_stamp >= time_stamp: - return int(block.number) + return int(block["number"]) # fallback in case correct block number hasn't been found # in that case, we will include some more blocks than necessary - return mid_block_number + 200 + return int(mid_block_number + 200) -def compute_block_and_month_range(node: Web3): # pylint: disable=too-many-locals +def compute_block_and_month_range( + node: Web3, +) -> Tuple[List[Tuple[int, int]], List[str]]: # pylint: disable=too-many-locals """ This determines the block range and the relevant months for which we will compute and upload data on Dune. From 0d8ef8600a100e58af51f2a34150103f79cb8cd4 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:00:39 +0200 Subject: [PATCH 18/26] properly disable too many locals pylint error --- src/sync/common.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/sync/common.py b/src/sync/common.py index dd5ed084..d9192211 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -69,9 +69,9 @@ def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: return int(mid_block_number + 200) -def compute_block_and_month_range( +def compute_block_and_month_range( # pylint: disable=too-many-locals node: Web3, -) -> Tuple[List[Tuple[int, int]], List[str]]: # pylint: disable=too-many-locals +) -> Tuple[List[Tuple[int, int]], List[str]]: """ This determines the block range and the relevant months for which we will compute and upload data on Dune. From 17c69f24ae70dca3f8fec77d47d66e675f4e9422 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:03:08 +0200 Subject: [PATCH 19/26] remove redundant whitespaces --- src/sync/common.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/sync/common.py b/src/sync/common.py index d9192211..6811f3f8 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -69,9 +69,9 @@ def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: return int(mid_block_number + 200) -def compute_block_and_month_range( # pylint: disable=too-many-locals +def compute_block_and_month_range( # pylint: disable=too-many-locals node: Web3, -) -> Tuple[List[Tuple[int, int]], List[str]]: +) -> Tuple[List[Tuple[int, int]], List[str]]: """ This determines the block range and the relevant months for which we will compute and upload data on Dune. From 616918b4394302b091e437f84cda6717872852ee Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:06:33 +0200 Subject: [PATCH 20/26] mypy fix --- src/sync/common.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/sync/common.py b/src/sync/common.py index 6811f3f8..020be203 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -83,8 +83,8 @@ def compute_block_and_month_range( # pylint: disable=too-many-locals latest_finalized_block = node.eth.get_block("finalized") - current_month_end_block = latest_finalized_block.number - current_month_end_timestamp = latest_finalized_block.timestamp + current_month_end_block = latest_finalized_block["number"] + current_month_end_timestamp = latest_finalized_block["timestamp"] current_month_end_datetime = datetime.fromtimestamp( current_month_end_timestamp, tz=timezone.utc From 05e9191599205ec1f8f7a776c54b690a197d0ecf Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:10:13 +0200 Subject: [PATCH 21/26] mypy fix --- src/sync/common.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/sync/common.py b/src/sync/common.py index 020be203..5147addf 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -40,7 +40,7 @@ def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: This implements binary search and returns the smallest block number whose timestamp is at least as large as the time_stamp argument passed in the function """ - end_block_number = node.eth.get_block("finalized")["number"] + end_block_number = int(node.eth.get_block("finalized")["number"]) start_block_number = 1 close_in_seconds = 30 @@ -83,7 +83,7 @@ def compute_block_and_month_range( # pylint: disable=too-many-locals latest_finalized_block = node.eth.get_block("finalized") - current_month_end_block = latest_finalized_block["number"] + current_month_end_block = int(latest_finalized_block["number"]) current_month_end_timestamp = latest_finalized_block["timestamp"] current_month_end_datetime = datetime.fromtimestamp( From 51982c41538688a9fa9e57846c6958cf4e262c23 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:14:02 +0200 Subject: [PATCH 22/26] small edits in env.sample --- .env.sample | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/.env.sample b/.env.sample index 880ea871..e5781b36 100644 --- a/.env.sample +++ b/.env.sample @@ -39,10 +39,10 @@ NODE_URL_ARBITRUM= NETWORK= # The prefix of the dune table where we sync the various mainnet price feeds -PRICE_FEED_TARGET_TABLE +PRICE_FEED_TARGET_TABLE= # The prefix of the dune table where we sync the monthly raw batch data -BATCH_DATA_TARGET_TABLE = "batch_data" +BATCH_DATA_TARGET_TABLE="batch_data" # Dune api timeout parameter, where we recommend to set it to 600 -DUNE_API_REQUEST_TIMEOUT = 600 +DUNE_API_REQUEST_TIMEOUT=600 From 65ab782b3020f4646d739a30faf43232f41cc7b4 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:43:45 +0200 Subject: [PATCH 23/26] bug fixing --- .env.sample | 10 +++++++++- src/fetch/orderbook.py | 10 +++++++--- src/main.py | 2 +- src/sync/batch_data.py | 3 ++- src/sync/common.py | 13 ------------- src/utils.py | 13 +++++++++++++ 6 files changed, 32 insertions(+), 19 deletions(-) diff --git a/.env.sample b/.env.sample index e5781b36..d895376d 100644 --- a/.env.sample +++ b/.env.sample @@ -30,7 +30,7 @@ ETHERSCAN_API_KEY= AGGREGATE_QUERY_ID= # Node for each chain we want to run a sync job on -NODE_URL_MAINNET= +NODE_URL_ETHEREUM= NODE_URL_GNOSIS= NODE_URL_ARBITRUM= @@ -46,3 +46,11 @@ BATCH_DATA_TARGET_TABLE="batch_data" # Dune api timeout parameter, where we recommend to set it to 600 DUNE_API_REQUEST_TIMEOUT=600 + +# The following constants define the upper and lower caps per network +EPSILON_LOWER_ETHEREUM= +EPSILON_UPPER_ETHEREUM= +EPSILON_LOWER_GNOSIS= +EPSILON_UPPER_GNOSIS= +EPSILON_LOWER_ARBITRUM= +EPSILON_UPPER_ARBITRUM= \ No newline at end of file diff --git a/src/fetch/orderbook.py b/src/fetch/orderbook.py index 90ba1ef0..3561ce7f 100644 --- a/src/fetch/orderbook.py +++ b/src/fetch/orderbook.py @@ -14,7 +14,7 @@ from src.logger import set_log from src.models.block_range import BlockRange -from src.utils import open_query +from src.utils import open_query, node_suffix log = set_log(__name__) @@ -179,15 +179,19 @@ def run_batch_data_sql(cls, block_range: BlockRange) -> DataFrame: """ Fetches and validates Batch data DataFrame as concatenation from Prod and Staging DB """ + load_dotenv() + network = node_suffix(os.environ["NETWORK"]) + epsilon_upper = str(os.environ[f"EPSILON_UPPER_{network}"]) + epsilon_lower = str(os.environ[f"EPSILON_LOWER_{network}"]) batch_data_query_prod = ( open_query("orderbook/batch_data.sql") .replace("{{start_block}}", str(block_range.block_from)) .replace("{{end_block}}", str(block_range.block_to)) .replace( - "{{EPSILON_LOWER}}", "10000000000000000" + "{{EPSILON_LOWER}}", epsilon_lower ) # lower ETH cap for payment (in WEI) .replace( - "{{EPSILON_UPPER}}", "12000000000000000" + "{{EPSILON_UPPER}}", epsilon_upper ) # upper ETH cap for payment (in WEI) .replace("{{env}}", "prod") ) diff --git a/src/main.py b/src/main.py index 7276ec1b..faed5e6a 100644 --- a/src/main.py +++ b/src/main.py @@ -23,7 +23,7 @@ ) from src.sync.order_rewards import sync_order_rewards, sync_batch_rewards from src.sync.batch_data import sync_batch_data -from src.sync.common import node_suffix +from src.utils import node_suffix log = set_log(__name__) diff --git a/src/sync/batch_data.py b/src/sync/batch_data.py index 96867d7d..6b8bd0e2 100644 --- a/src/sync/batch_data.py +++ b/src/sync/batch_data.py @@ -6,8 +6,9 @@ from src.fetch.orderbook import OrderbookFetcher from src.logger import set_log from src.sync.config import BatchDataSyncConfig -from src.sync.common import compute_block_and_month_range, node_suffix +from src.sync.common import compute_block_and_month_range from src.models.block_range import BlockRange +from src.utils import node_suffix log = set_log(__name__) diff --git a/src/sync/common.py b/src/sync/common.py index 5147addf..eb2f199c 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -22,19 +22,6 @@ def last_sync_block(aws: AWSClient, table: SyncTable, genesis_block: int = 0) -> return block_from -def node_suffix(network: str) -> str: - """ - Converts network internal name to name used for nodes and dune tables - """ - if network == "mainnet": - return "MAINNET" - if network == "xdai": - return "GNOSIS" - if network == "arbitrum-one": - return "ARBITRUM" - return "" - - def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: """ This implements binary search and returns the smallest block number diff --git a/src/utils.py b/src/utils.py index a601e8c9..409452d9 100644 --- a/src/utils.py +++ b/src/utils.py @@ -13,3 +13,16 @@ def open_query(filename: str) -> str: def query_file(filename: str) -> str: """Returns proper path for filename in QUERY_PATH""" return os.path.join(QUERY_PATH, filename) + + +def node_suffix(network: str) -> str: + """ + Converts network internal name to name used for nodes and dune tables + """ + if network == "mainnet": + return "ETHEREUM" + if network == "xdai": + return "GNOSIS" + if network == "arbitrum-one": + return "ARBITRUM" + return "" From 78194c2be5dc031c9d7d55648bff6ddad04ed20e Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:45:09 +0200 Subject: [PATCH 24/26] add new line --- .env.sample | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.env.sample b/.env.sample index d895376d..3adeb977 100644 --- a/.env.sample +++ b/.env.sample @@ -53,4 +53,4 @@ EPSILON_UPPER_ETHEREUM= EPSILON_LOWER_GNOSIS= EPSILON_UPPER_GNOSIS= EPSILON_LOWER_ARBITRUM= -EPSILON_UPPER_ARBITRUM= \ No newline at end of file +EPSILON_UPPER_ARBITRUM= From 00d80ee299074e978605b9a664f76e01f14f042a Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:46:31 +0200 Subject: [PATCH 25/26] small fix --- src/fetch/orderbook.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/src/fetch/orderbook.py b/src/fetch/orderbook.py index 3561ce7f..1de3b621 100644 --- a/src/fetch/orderbook.py +++ b/src/fetch/orderbook.py @@ -200,10 +200,10 @@ def run_batch_data_sql(cls, block_range: BlockRange) -> DataFrame: .replace("{{start_block}}", str(block_range.block_from)) .replace("{{end_block}}", str(block_range.block_to)) .replace( - "{{EPSILON_LOWER}}", "10000000000000000" + "{{EPSILON_LOWER}}", epsilon_lower ) # lower ETH cap for payment (in WEI) .replace( - "{{EPSILON_UPPER}}", "12000000000000000" + "{{EPSILON_UPPER}}", epsilon_upper ) # upper ETH cap for payment (in WEI) .replace("{{env}}", "barn") ) From bced3e599de6b805407489cf64c140a1e1bbf8f2 Mon Sep 17 00:00:00 2001 From: harisang Date: Sat, 23 Nov 2024 04:49:23 +0200 Subject: [PATCH 26/26] remove redundant type cast --- src/sync/common.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/sync/common.py b/src/sync/common.py index eb2f199c..205f0f11 100644 --- a/src/sync/common.py +++ b/src/sync/common.py @@ -53,7 +53,7 @@ def find_block_with_timestamp(node: Web3, time_stamp: float) -> int: return int(block["number"]) # fallback in case correct block number hasn't been found # in that case, we will include some more blocks than necessary - return int(mid_block_number + 200) + return mid_block_number + 200 def compute_block_and_month_range( # pylint: disable=too-many-locals