diff --git a/.github/workflows/integration-test.yml b/.github/workflows/integration-test.yml index 5438e1c6..fb0fc996 100644 --- a/.github/workflows/integration-test.yml +++ b/.github/workflows/integration-test.yml @@ -136,5 +136,4 @@ jobs: npx ts-node tests/trade.ts sleep 5 npx ts-node tests/print_orders.ts - npx ts-node tests/transfer.ts npx ts-node tests/put_batch_orders.ts diff --git a/Cargo.lock b/Cargo.lock index 6d8d3332..91ab6fcd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1584,9 +1584,9 @@ dependencies = [ [[package]] name = "num_enum" -version = "0.5.1" +version = "0.5.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "226b45a5c2ac4dd696ed30fa6b94b057ad909c7b7fc2e0d0808192bced894066" +checksum = "3f9bd055fb730c4f8f4f57d45d35cd6b3f0980535b056dc7ff119cee6a66ed6f" dependencies = [ "derivative", "num_enum_derive", @@ -1594,9 +1594,9 @@ dependencies = [ [[package]] name = "num_enum_derive" -version = "0.5.1" +version = "0.5.4" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1c0fd9eba1d5db0994a239e09c1be402d35622277e35468ba891aa5e3188ce7e" +checksum = "486ea01961c4a818096de679a8b740b26d9033146ac5291b1c98557658f8cdd9" dependencies = [ "proc-macro-crate", "proc-macro2", @@ -1625,7 +1625,7 @@ checksum = "624a8340c38c1b80fd549087862da4ba43e08858af025b236e509b6649fc13d5" [[package]] name = "orchestra" version = "0.1.0" -source = "git+https://github.com/gcomte/orchestra.git?branch=master#14745d598fe31ca2c7fd8698f92173f8136335a1" +source = "git+https://github.com/gcomte/orchestra.git?branch=master#c786218142d73acffd24c7d28b60e7bfcef92613" dependencies = [ "prost", "serde 1.0.124", @@ -1880,10 +1880,11 @@ checksum = "ac74c624d6b2d21f425f752262f42188365d7b8ff1aff74c82e45136510a4857" [[package]] name = "proc-macro-crate" -version = "0.1.5" +version = "1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "1d6ea3c4595b96363c13943497db34af4460fb474a95c43f4446ad341b8c9785" +checksum = "1ebace6889caf889b4d3f76becee12e90353f2b8c7d875534a71e5742f8f6f83" dependencies = [ + "thiserror", "toml", ] diff --git a/Cargo.toml b/Cargo.toml index ffdfab85..c8d7356a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -30,7 +30,7 @@ jsonwebtoken = "7.2.0" lazy_static = "1.4.0" log = "0.4.14" nix = "0.20.0" -num_enum = "0.5.1" +num_enum = "0.5.4" orchestra = { git = "https://github.com/gcomte/orchestra.git", branch = "master", features = [ "exchange" ] } paperclip = { git = "https://github.com/fluidex/paperclip.git", features = [ "actix", "chrono", "rust_decimal" ] } qstring = "0.7.2" diff --git a/examples/js/RESTClient.ts b/examples/js/RESTClient.ts index 9658da4a..11960275 100644 --- a/examples/js/RESTClient.ts +++ b/examples/js/RESTClient.ts @@ -14,27 +14,6 @@ class RESTClient { }); } - async internal_txs( - user_id: number | string, - params?: { - limit?: number; - offset?: number; - start_time?: number; - end_time?: number; - order?: "asc" | "desc"; - side?: "from" | "to" | "both"; - } - ) { - let resp = await this.client.get(`/internal_txs/${user_id}`, { - params: _.pickBy(params, _.identity), - }); - if (resp.status === 200) { - return resp.data; - } else { - throw new Error(`request failed with ${resp.status} ${resp.statusText}`); - } - } - async closed_orders(token: string) { if (token !== "") { this.client.defaults.headers.common["Authorization"] = "LoremIpsum"; diff --git a/examples/js/tests/transfer.ts b/examples/js/tests/transfer.ts deleted file mode 100644 index 9c2a1b44..00000000 --- a/examples/js/tests/transfer.ts +++ /dev/null @@ -1,99 +0,0 @@ -import { TestUser } from "../config"; // dotenv -import { defaultClient as client } from "../client"; -import { defaultRESTClient as rest_client } from "../RESTClient"; -import { assertDecimalEqual, sleep } from "../util"; - -import { strict as assert } from "assert"; -import { depositAssets } from "../exchange_helper"; - -async function setupAsset() { - await depositAssets({ ETH: "100.0" }, TestUser.USER1); - - const balance1 = await client.balanceQueryByAsset(TestUser.USER1, "ETH"); - assertDecimalEqual(balance1.available, "100"); - const balance2 = await client.balanceQueryByAsset(TestUser.USER2, "ETH"); - assertDecimalEqual(balance2.available, "0"); -} - -// Test failure with argument delta of value zero -async function failureWithZeroDeltaTest() { - const res = await client.transfer(TestUser.USER1, TestUser.USER2, "ETH", 0); - - assert.equal(res.success, false); - assert.equal(res.asset, "ETH"); - assertDecimalEqual(res.balance_from, "100"); - - const balance1 = await client.balanceQueryByAsset(TestUser.USER1, "ETH"); - assertDecimalEqual(balance1.available, "100"); - const balance2 = await client.balanceQueryByAsset(TestUser.USER2, "ETH"); - assertDecimalEqual(balance2.available, "0"); - - console.log("failureWithZeroDeltaTest passed"); -} - -// Test failure with insufficient balance of from user -async function failureWithInsufficientFromBalanceTest() { - const res = await client.transfer(TestUser.USER1, TestUser.USER2, "ETH", 101); - - assert.equal(res.success, false); - assert.equal(res.asset, "ETH"); - assertDecimalEqual(res.balance_from, "100"); - - const balance1 = await client.balanceQueryByAsset(TestUser.USER1, "ETH"); - assertDecimalEqual(balance1.available, "100"); - const balance2 = await client.balanceQueryByAsset(TestUser.USER2, "ETH"); - assertDecimalEqual(balance2.available, "0"); - - console.log("failureWithInsufficientFromBalanceTest passed"); -} - -// Test success transfer -async function successTransferTest() { - const res = await client.transfer(TestUser.USER1, TestUser.USER2, "ETH", 50); - - assert.equal(res.success, true); - assert.equal(res.asset, "ETH"); - assertDecimalEqual(res.balance_from, "50"); - - const balance1 = await client.balanceQueryByAsset(TestUser.USER1, "ETH"); - assertDecimalEqual(balance1.available, "50"); - const balance2 = await client.balanceQueryByAsset(TestUser.USER2, "ETH"); - assertDecimalEqual(balance2.available, "50"); - - console.log("successTransferTest passed"); -} - -async function listTxs() { - const res1 = (await rest_client.internal_txs(TestUser.USER1))[0]; - const res2 = (await rest_client.internal_txs(TestUser.USER2))[0]; - console.log(res1, res2); - assert.equal(res1.amount, res2.amount); - assert.equal(res1.asset, res2.asset); - assert.equal(res1.time, res2.time); - assert.equal(res1.user_from, res2.user_from); - assert.equal(res1.user_to, res2.user_to); -} - -async function simpleTest() { - await setupAsset(); - await failureWithZeroDeltaTest(); - await failureWithInsufficientFromBalanceTest(); - await successTransferTest(); - await sleep(3 * 1000); - await listTxs(); -} - -async function mainTest() { - await client.debugReset(1); - await simpleTest(); -} - -async function main() { - try { - await mainTest(); - } catch (error) { - console.error("Caught error:", error); - process.exit(1); - } -} -main(); diff --git a/migrations/20210607094808_internal_transfer.sql b/migrations/20210607094808_internal_transfer.sql deleted file mode 100644 index d99d3421..00000000 --- a/migrations/20210607094808_internal_transfer.sql +++ /dev/null @@ -1,13 +0,0 @@ --- Add migration script here -CREATE TABLE internal_tx ( - time TIMESTAMP(0) NOT NULL, - user_from VARCHAR(36) NOT NULL, - user_to VARCHAR(36) NOT NULL, - asset VARCHAR(30) NOT NULL REFERENCES asset(id), - amount DECIMAL(30, 8) CHECK (amount > 0) NOT NULL -); - -CREATE INDEX internal_tx_idx_to_time ON internal_tx (user_to, time DESC); -CREATE INDEX internal_tx_idx_from_time ON internal_tx (user_from, time DESC); - -SELECT create_hypertable('internal_tx', 'time'); diff --git a/orchestra b/orchestra index 14745d59..c7862181 160000 --- a/orchestra +++ b/orchestra @@ -1 +1 @@ -Subproject commit 14745d598fe31ca2c7fd8698f92173f8136335a1 +Subproject commit c786218142d73acffd24c7d28b60e7bfcef92613 diff --git a/src/bin/dump_unify_messages.rs b/src/bin/dump_unify_messages.rs index b1562448..87b1dc8c 100644 --- a/src/bin/dump_unify_messages.rs +++ b/src/bin/dump_unify_messages.rs @@ -18,7 +18,6 @@ use fluidex_common::rdkafka::message::{BorrowedMessage, Message}; fn get_msg_tag_from_topic(t: &str) -> Option<&'static str> { Some(match t { "deposits" => "DepositMessage", - "internaltransfer" => "TransferMessage", "orders" => "OrderMessage", "trades" => "TradeMessage", "withdraws" => "WithdrawMessage", diff --git a/src/bin/persistor.rs b/src/bin/persistor.rs index 4832f6ba..c6717942 100644 --- a/src/bin/persistor.rs +++ b/src/bin/persistor.rs @@ -61,8 +61,6 @@ fn main() { let persistor_balance: DatabaseWriter = DatabaseWriter::new(&write_config).start_schedule(&pool).unwrap(); - let persistor_transfer: DatabaseWriter = DatabaseWriter::new(&write_config).start_schedule(&pool).unwrap(); - let trade_cfg = TopicConfig::::new(message::TRADES_TOPIC) .persist_to(&persistor_kline) .persist_to(&persistor_trade) @@ -76,13 +74,10 @@ fn main() { let balance_cfg = TopicConfig::::new(message::BALANCES_TOPIC).persist_to(&persistor_balance); - let internaltx_cfg = TopicConfig::::new(message::INTERNALTX_TOPIC).persist_to(&persistor_transfer); - let auto_commit = vec![ trade_cfg.auto_commit_start(consumer.clone()), order_cfg.auto_commit_start(consumer.clone()), balance_cfg.auto_commit_start(consumer.clone()), - internaltx_cfg.auto_commit_start(consumer.clone()), ]; let consumer = consumer.as_ref(); @@ -91,7 +86,6 @@ fn main() { .add_topic_config(&trade_cfg).unwrap() .add_topic_config(&order_cfg).unwrap() .add_topic_config(&balance_cfg).unwrap() - .add_topic_config(&internaltx_cfg).unwrap() // .add_topic(message::TRADES_TOPIC, MsgDataPersistor::new(&persistor).handle_message::()) ; @@ -112,7 +106,6 @@ fn main() { persistor_trade.finish(), persistor_order.finish(), persistor_balance.finish(), - persistor_transfer.finish(), ) .expect("all persistor should success finish"); let final_commits: Vec + Send>>> = auto_commit diff --git a/src/bin/restapi.rs b/src/bin/restapi.rs index 5f6c1877..a60967b4 100644 --- a/src/bin/restapi.rs +++ b/src/bin/restapi.rs @@ -4,7 +4,7 @@ use actix_web_httpauth::extractors::AuthenticationError; use actix_web_httpauth::middleware::HttpAuthentication; use dingir_exchange::matchengine::authentication; use dingir_exchange::restapi::manage::market; -use dingir_exchange::restapi::personal_history::{my_internal_txs, my_orders}; +use dingir_exchange::restapi::personal_history::my_orders; use dingir_exchange::restapi::public_history::{order_trades, recent_trades}; use dingir_exchange::restapi::state::{AppCache, AppState}; use dingir_exchange::restapi::tradingview::{chart_config, history, search_symbols, symbols, ticker, unix_timestamp}; @@ -58,7 +58,6 @@ async fn main() -> std::io::Result<()> { .route("/recenttrades/{market}", web::get().to(recent_trades)) .route("/ordertrades/{market}/{order_id}", web::get().to(order_trades)) .route("/closedorders/{market}", web::get().to(my_orders)) - .route("/internal_txs", web::get().to(my_internal_txs)) .route("/ticker_{ticker_inv}/{market}", web::get().to(ticker)) .service( web::scope("/tradingview") diff --git a/src/matchengine/controller.rs b/src/matchengine/controller.rs index fe276f7f..e1ae07d1 100644 --- a/src/matchengine/controller.rs +++ b/src/matchengine/controller.rs @@ -563,106 +563,6 @@ impl Controller { Ok(()) } - pub fn transfer(&mut self, real: bool, req: TransferRequest, user_id: Uuid) -> Result { - if !self.check_service_available() { - return Err(Status::unavailable("")); - } - - let asset = &req.asset; - if !self.balance_manager.asset_manager.asset_exist(asset) { - return Err(Status::invalid_argument("invalid asset")); - } - - let to_user_id = req.to.clone(); - - let balance_manager = &self.balance_manager; - let balance_from = balance_manager.get(user_id, BalanceType::AVAILABLE, asset); - - let zero = Decimal::from(0); - let delta = Decimal::from_str(&req.delta).unwrap_or(zero); - - if delta <= zero || delta > balance_from { - return Ok(TransferResponse { - success: false, - asset: asset.to_owned(), - balance_from: balance_from.to_string(), - }); - } - - let prec = self.balance_manager.asset_manager.asset_prec_show(asset); - let change = delta.round_dp(prec); - - let business = "transfer"; - let timestamp = FTimestamp(current_timestamp()); - let business_id = (timestamp.0 * 1_000_f64) as u64; // milli-seconds - let detail_json: serde_json::Value = if req.memo.is_empty() { - json!({}) - } else { - serde_json::from_str(req.memo.as_str()).map_err(|_| Status::invalid_argument("invalid memo"))? - }; - - // Get market price of requested base asset and quote asset of USDT. - let market_price = self - .asset_market_names - .get(&(asset.to_owned(), "USDT".to_owned())) - .map_or(Decimal::zero(), |market_name| self.markets.get(market_name).unwrap().price); - let persistor = if real { &mut self.persistor } else { &mut self.dummy_persistor }; - self.update_controller - .update_user_balance( - &mut self.balance_manager, - persistor, - BalanceUpdateParams { - balance_type: BalanceType::AVAILABLE, - business_type: BusinessType::Transfer, - user_id, - asset: asset.to_owned(), - business: business.to_owned(), - business_id, - market_price, - change: -change, - detail: detail_json.clone(), - }, - ) - .map_err(|e| Status::invalid_argument(format!("{}", e)))?; - - let persistor = if real { &mut self.persistor } else { &mut self.dummy_persistor }; - self.update_controller - .update_user_balance( - &mut self.balance_manager, - persistor, - BalanceUpdateParams { - balance_type: BalanceType::AVAILABLE, - business_type: BusinessType::Transfer, - user_id: to_user_id.parse().unwrap(), - asset: asset.to_owned(), - business: business.to_owned(), - business_id, - market_price: Decimal::zero(), - change, - detail: detail_json, - }, - ) - .map_err(|e| Status::invalid_argument(format!("{}", e)))?; - - if real { - self.persistor.put_transfer(models::InternalTx { - time: timestamp.into(), - user_from: user_id.to_string(), - user_to: to_user_id, - asset: asset.to_owned(), - amount: change, - }); - - self.append_operation_log(OPERATION_TRANSFER, &req, user_id); - } - - Ok(TransferResponse { - success: true, - asset: asset.to_owned(), - balance_from: (balance_from - change).to_string(), - }) - } - pub async fn debug_reset(&mut self, _req: DebugResetRequest) -> Result { async { log::info!("do full reset: memory and db"); @@ -760,9 +660,6 @@ impl Controller { OPERATION_BATCH_ORDER_PUT => { self.batch_order_put(false, serde_json::from_str(params)?, user_id)?; } - OPERATION_TRANSFER => { - self.transfer(false, serde_json::from_str(params)?, user_id)?; - } _ => bail!("invalid operation {}", method), } Ok(()) diff --git a/src/matchengine/history.rs b/src/matchengine/history.rs index 509e4038..568e7ced 100644 --- a/src/matchengine/history.rs +++ b/src/matchengine/history.rs @@ -7,7 +7,6 @@ use anyhow::Result; use fluidex_common::utils::timeutil::FTimestamp; type BalanceWriter = DatabaseWriter; -type TransferWriter = DatabaseWriter; type OrderWriter = DatabaseWriter; type TradeWriter = DatabaseWriter; @@ -15,7 +14,6 @@ pub trait HistoryWriter: Sync + Send { fn is_block(&self) -> bool; //TODO: don't take the ownership? fn append_balance_history(&mut self, data: models::BalanceHistory); - fn append_internal_transfer(&mut self, data: models::InternalTx); fn append_order_history(&mut self, order: &market::Order); fn append_expired_order_history(&mut self, _order: &market::Order); fn append_pair_user_trade(&mut self, trade: &Trade); @@ -24,7 +22,6 @@ pub trait HistoryWriter: Sync + Send { pub struct DummyHistoryWriter; impl HistoryWriter for DummyHistoryWriter { fn append_balance_history(&mut self, _data: models::BalanceHistory) {} - fn append_internal_transfer(&mut self, _data: models::InternalTx) {} fn append_order_history(&mut self, _order: &market::Order) {} fn append_expired_order_history(&mut self, _order: &market::Order) {} fn append_pair_user_trade(&mut self, _trade: &Trade) {} @@ -35,7 +32,6 @@ impl HistoryWriter for DummyHistoryWriter { pub struct DatabaseHistoryWriter { pub balance_writer: BalanceWriter, - pub transfer_writer: TransferWriter, pub trade_writer: TradeWriter, pub order_writer: OrderWriter, } @@ -44,7 +40,6 @@ impl DatabaseHistoryWriter { pub fn new(config: &DatabaseWriterConfig, pool: &sqlx::Pool) -> Result { Ok(DatabaseHistoryWriter { balance_writer: BalanceWriter::new(config).start_schedule(pool)?, - transfer_writer: TransferWriter::new(config).start_schedule(pool)?, trade_writer: TradeWriter::new(config).start_schedule(pool)?, order_writer: OrderWriter::new(config).start_schedule(pool)?, }) @@ -87,9 +82,6 @@ impl HistoryWriter for DatabaseHistoryWriter { fn append_balance_history(&mut self, data: models::BalanceHistory) { self.balance_writer.append(data).ok(); } - fn append_internal_transfer(&mut self, data: models::InternalTx) { - self.transfer_writer.append(data).ok(); - } fn append_order_history(&mut self, order: &market::Order) { self.order_writer.append(order.into()).ok(); } diff --git a/src/matchengine/persist/persistor.rs b/src/matchengine/persist/persistor.rs index ab8e5765..4f3fd164 100644 --- a/src/matchengine/persist/persistor.rs +++ b/src/matchengine/persist/persistor.rs @@ -20,7 +20,6 @@ pub trait PersistExector: Send + Sync { fn put_balance(&mut self, balance: &BalanceHistory); fn put_deposit(&mut self, balance: &BalanceHistory); fn put_withdraw(&mut self, balance: &BalanceHistory); - fn put_transfer(&mut self, tx: InternalTx); fn put_order(&mut self, order: &Order, at_step: OrderEventType); fn put_trade(&mut self, trade: &Trade); } @@ -41,9 +40,6 @@ impl PersistExector for Box { fn put_withdraw(&mut self, balance: &BalanceHistory) { self.as_mut().put_withdraw(balance) } - fn put_transfer(&mut self, tx: InternalTx) { - self.as_mut().put_transfer(tx) - } fn put_order(&mut self, order: &Order, at_step: OrderEventType) { self.as_mut().put_order(order, at_step) } @@ -68,9 +64,6 @@ impl PersistExector for &mut Box { fn put_withdraw(&mut self, balance: &BalanceHistory) { self.as_mut().put_withdraw(balance) } - fn put_transfer(&mut self, tx: InternalTx) { - self.as_mut().put_transfer(tx) - } fn put_order(&mut self, order: &Order, at_step: OrderEventType) { self.as_mut().put_order(order, at_step) } @@ -100,7 +93,6 @@ impl PersistExector for DummyPersistor { fn put_balance(&mut self, _balance: &BalanceHistory) {} fn put_deposit(&mut self, _balance: &BalanceHistory) {} fn put_withdraw(&mut self, _balance: &BalanceHistory) {} - fn put_transfer(&mut self, _tx: InternalTx) {} fn put_order(&mut self, _order: &Order, _as_step: OrderEventType) {} fn put_trade(&mut self, _trade: &Trade) {} } @@ -112,7 +104,6 @@ impl PersistExector for &mut DummyPersistor { fn put_balance(&mut self, _balance: &BalanceHistory) {} fn put_deposit(&mut self, _balance: &BalanceHistory) {} fn put_withdraw(&mut self, _balance: &BalanceHistory) {} - fn put_transfer(&mut self, _tx: InternalTx) {} fn put_order(&mut self, _order: &Order, _as_step: OrderEventType) {} fn put_trade(&mut self, _trade: &Trade) {} } @@ -146,9 +137,6 @@ impl PersistExector for MemBasedPersistor { fn put_withdraw(&mut self, balance: &BalanceHistory) { self.messages.push(message::Message::WithdrawMessage(Box::new(balance.into()))); } - fn put_transfer(&mut self, tx: InternalTx) { - self.messages.push(message::Message::TransferMessage(Box::new(tx.into()))); - } } ///////////////////////////// FileBasedPersistor //////////////////////////// @@ -189,10 +177,6 @@ impl PersistExector for FileBasedPersistor { let msg = message::Message::WithdrawMessage(Box::new(balance.into())); self.write_msg(msg); } - fn put_transfer(&mut self, tx: InternalTx) { - let msg = message::Message::TransferMessage(Box::new(tx.into())); - self.write_msg(msg); - } } ///////////////////////////// MessengerBasedPersistor //////////////////////////// @@ -224,9 +208,6 @@ impl PersistExector for MessengerBasedPersistor { fn put_withdraw(&mut self, balance: &BalanceHistory) { self.inner.push_withdraw_message(&balance.into()); } - fn put_transfer(&mut self, tx: InternalTx) { - self.inner.push_transfer_message(&tx.into()); - } fn put_order(&mut self, order: &Order, at_step: OrderEventType) { self.inner.push_order_message(&OrderMessage::from_order(order, at_step)); } @@ -264,9 +245,6 @@ impl PersistExector for DBBasedPersistor { fn put_withdraw(&mut self, _balance: &BalanceHistory) { // TODO } - fn put_transfer(&mut self, tx: InternalTx) { - self.inner.append_internal_transfer(tx); - } fn put_order(&mut self, order: &Order, at_step: OrderEventType) { //only persist on finish match at_step { @@ -318,11 +296,6 @@ impl PersistExector for CompositePersistor { p.put_withdraw(balance); } } - fn put_transfer(&mut self, tx: InternalTx) { - for p in &mut self.persistors { - p.put_transfer(tx.clone()); - } - } fn put_order(&mut self, order: &Order, at_step: OrderEventType) { for p in &mut self.persistors { p.put_order(order, at_step); diff --git a/src/matchengine/server.rs b/src/matchengine/server.rs index 74f8c97c..ed720ea7 100644 --- a/src/matchengine/server.rs +++ b/src/matchengine/server.rs @@ -258,16 +258,6 @@ impl matchengine_server::Matchengine for GrpcHandler { Ok(Response::new(SimpleSuccessResponse {})) } - async fn transfer(&self, request: Request) -> Result, Status> { - let user_id = get_user_id_from_request(&request); - let ControllerDispatch(act, rt) = ControllerDispatch::new(move |ctrl: &mut Controller| { - Box::pin(async move { ctrl.transfer(true, request.into_inner(), user_id) }) - }); - - self.task_dispatcher.send(act).await.map_err(map_dispatch_err)?; - map_dispatch_ret(rt.await) - } - // This is the only blocking call of the server #[cfg(debug_assertions)] async fn debug_dump(&self, request: Request) -> Result, Status> { diff --git a/src/message/mod.rs b/src/message/mod.rs index 6b9bf494..7d68485c 100644 --- a/src/message/mod.rs +++ b/src/message/mod.rs @@ -1,10 +1,9 @@ use crate::market::Order; -pub use crate::models::{BalanceHistory, InternalTx}; +pub use crate::models::BalanceHistory; use crate::types::OrderEventType; use uuid::Uuid; use anyhow::Result; -use fluidex_common::utils::timeutil::FTimestamp; use serde::{Deserialize, Serialize}; pub mod consumer; @@ -12,7 +11,7 @@ pub mod persist; pub mod producer; pub use producer::{ - BALANCES_TOPIC, DEPOSITS_TOPIC, INTERNALTX_TOPIC, ORDERS_TOPIC, TRADES_TOPIC, UNIFY_TOPIC, USER_TOPIC, WITHDRAWS_TOPIC, + BALANCES_TOPIC, DEPOSITS_TOPIC, ORDERS_TOPIC, TRADES_TOPIC, UNIFY_TOPIC, USER_TOPIC, WITHDRAWS_TOPIC, }; #[derive(Debug, Serialize, Deserialize, Clone)] @@ -106,27 +105,6 @@ impl From<&BalanceHistory> for WithdrawMessage { } } -#[derive(Debug, Serialize, Deserialize, Clone)] -pub struct TransferMessage { - pub time: f64, - pub user_from: Uuid, - pub user_to: Uuid, - pub asset: String, - pub amount: String, -} - -impl From for TransferMessage { - fn from(tx: InternalTx) -> Self { - Self { - time: FTimestamp::from(&tx.time).into(), - user_from: tx.user_from.parse().unwrap(), - user_to: tx.user_to.parse().unwrap(), - asset: tx.asset.to_string(), - amount: tx.amount.to_string(), - } - } -} - #[derive(Debug, Serialize, Deserialize, Clone)] pub struct OrderMessage { pub event: OrderEventType, @@ -164,7 +142,6 @@ pub trait MessageManager: Sync + Send { fn push_balance_message(&mut self, balance: &BalanceMessage); fn push_deposit_message(&mut self, balance: &DepositMessage); fn push_withdraw_message(&mut self, balance: &WithdrawMessage); - fn push_transfer_message(&mut self, tx: &TransferMessage); } pub struct RdProducerStub { @@ -244,10 +221,6 @@ impl MessageManager for RdProducerStub { let message = serde_json::to_string(&withdraw).unwrap(); self.push_message_and_topic(message, WITHDRAWS_TOPIC) } - fn push_transfer_message(&mut self, tx: &TransferMessage) { - let message = serde_json::to_string(&tx).unwrap(); - self.push_message_and_topic(message, INTERNALTX_TOPIC) - } } pub type SimpleMessageManager = RdProducerStub; @@ -267,7 +240,6 @@ pub enum Message { DepositMessage(Box), OrderMessage(Box), TradeMessage(Box), - TransferMessage(Box), WithdrawMessage(Box), } diff --git a/src/message/persist.rs b/src/message/persist.rs index 54630cd8..69298cb3 100644 --- a/src/message/persist.rs +++ b/src/message/persist.rs @@ -665,15 +665,3 @@ impl MsgDataTransformer for BidTrade { }) } } - -impl<'r> From<&'r super::TransferMessage> for models::InternalTx { - fn from(origin: &'r super::TransferMessage) -> Self { - Self { - time: FTimestamp(origin.time).into(), - user_from: origin.user_from.to_string(), - user_to: origin.user_to.to_string(), - asset: origin.asset.clone(), - amount: DecimalDbType::from_str(&origin.amount).unwrap_or_else(decimal_warning), - } - } -} diff --git a/src/message/producer.rs b/src/message/producer.rs index e4edc018..5476656f 100644 --- a/src/message/producer.rs +++ b/src/message/producer.rs @@ -184,7 +184,6 @@ impl RdProducerContext { pub const BALANCES_TOPIC: &str = "balances"; pub const DEPOSITS_TOPIC: &str = "deposits"; -pub const INTERNALTX_TOPIC: &str = "internaltransfer"; pub const ORDERS_TOPIC: &str = "orders"; pub const TRADES_TOPIC: &str = "trades"; pub const UNIFY_TOPIC: &str = "unifyevents"; @@ -196,7 +195,6 @@ use std::collections::LinkedList; #[derive(Default)] pub struct SimpleMessageScheme { balances_list: LinkedList, - internaltxs_list: LinkedList, orders_list: LinkedList, trades_list: LinkedList, users_list: LinkedList, @@ -213,7 +211,6 @@ impl MessageScheme for SimpleMessageScheme { } fn is_full(&self) -> bool { self.balances_list.len() >= 100 - || self.internaltxs_list.len() >= 100 || self.orders_list.len() >= 100 || self.trades_list.len() >= 100 || self.users_list.len() >= 100 @@ -222,7 +219,6 @@ impl MessageScheme for SimpleMessageScheme { fn on_message(&mut self, title_tip: &'static str, message: String) { let list = match title_tip { BALANCES_TOPIC => &mut self.balances_list, - INTERNALTX_TOPIC => &mut self.internaltxs_list, ORDERS_TOPIC => &mut self.orders_list, TRADES_TOPIC => &mut self.trades_list, USER_TOPIC => &mut self.users_list, @@ -239,12 +235,11 @@ impl MessageScheme for SimpleMessageScheme { let mut topic_name = BALANCES_TOPIC; let mut candi_list = [ - &mut self.internaltxs_list, &mut self.orders_list, &mut self.trades_list, &mut self.users_list, ]; - let iters = [INTERNALTX_TOPIC, ORDERS_TOPIC, TRADES_TOPIC, USER_TOPIC] + let iters = [ORDERS_TOPIC, TRADES_TOPIC, USER_TOPIC] .iter() .zip(&mut candi_list); @@ -309,7 +304,7 @@ impl MessageScheme for FullOrderMessageScheme { fn on_message(&mut self, title_tip: &'static str, message: String) { match title_tip { - DEPOSITS_TOPIC | INTERNALTX_TOPIC | ORDERS_TOPIC | TRADES_TOPIC | USER_TOPIC | WITHDRAWS_TOPIC => { + DEPOSITS_TOPIC | ORDERS_TOPIC | TRADES_TOPIC | USER_TOPIC | WITHDRAWS_TOPIC => { self.ordered_list.push_back((title_tip, message)) } _ => {} diff --git a/src/restapi/personal_history.rs b/src/restapi/personal_history.rs index d4a7f7f6..bc37d81b 100644 --- a/src/restapi/personal_history.rs +++ b/src/restapi/personal_history.rs @@ -1,6 +1,6 @@ use crate::matchengine::authentication::UserExtension; -use crate::models::tablenames::{INTERNALTX, ORDERHISTORY}; -use crate::models::{DateTimeMilliseconds, DecimalDbType, OrderHistory, TimestampDbType}; +use crate::models::tablenames::ORDERHISTORY; +use crate::models::{OrderHistory, TimestampDbType}; use crate::restapi::errors::RpcError; use crate::restapi::state::AppState; use core::cmp::min; @@ -52,16 +52,6 @@ pub async fn my_orders(req: HttpRequest, data: web::Data) -> Result, - #[serde(default, deserialize_with = "u64_timestamp_deserializer")] - end_time: Option, - #[serde(default)] - order: Order, - #[serde(default)] - side: Side, -} - fn u64_timestamp_deserializer<'de, D>(deserializer: D) -> Result, D::Error> where D: Deserializer<'de>, @@ -124,67 +96,3 @@ const fn default_limit() -> usize { const fn default_zero() -> usize { 0 } - -/// `/internal_txs` -#[api_v2_operation] -pub async fn my_internal_txs( - req: HttpRequest, - query: web::Query, - data: web::Data, -) -> Result>, actix_web::Error> { - let user_id = req.extensions().get::().unwrap().user_id; - let limit = min(query.limit, 100); - - let base_query: &'static str = const_format::formatcp!( - r#" -select i.time as time, - af.l2_pubkey as user_from, - at.l2_pubkey as user_to, - i.asset as asset, - i.amount as amount -from {} i -where "#, - INTERNALTX - ); - let (user_condition, args_n) = match query.side { - Side::From => ("i.user_from = $1", 1), - Side::To => ("i.user_to = $1", 1), - Side::Both => ("i.user_from = $1 or i.user_to = $2", 2), - }; - - let time_condition = match (query.start_time, query.end_time) { - (Some(_), Some(_)) => Some(format!("i.time >= ${} and i.time <= ${}", args_n + 1, args_n + 2)), - (Some(_), None) => Some(format!("i.time >= ${}", args_n + 1)), - (None, Some(_)) => Some(format!("i.time <= ${}", args_n + 1)), - (None, None) => None, - }; - - let condition = match time_condition { - Some(time_condition) => format!("({}) and {}", user_condition, time_condition), - None => user_condition.to_string(), - }; - - let constraint = format!("limit {} offset {}", limit, query.offset); - let sql_query = format!("{}{}{}", base_query, condition, constraint); - - let query_as = sqlx::query_as(sql_query.as_str()); - - let query_as = match query.side { - Side::To | Side::From => query_as.bind(user_id.to_string()), - Side::Both => query_as.bind(user_id.to_string()).bind(user_id.to_string()), - }; - - let query_as = match (query.start_time, query.end_time) { - (Some(start_time), Some(end_time)) => query_as.bind(start_time).bind(end_time), - (Some(start_time), None) => query_as.bind(start_time), - (None, Some(end_time)) => query_as.bind(end_time), - (None, None) => query_as, - }; - - let txs: Vec = query_as - .fetch_all(&data.db) - .await - .map_err(|err| actix_web::Error::from(RpcError::from(err)))?; - - Ok(Json(txs)) -} diff --git a/src/storage/models.rs b/src/storage/models.rs index 872f5672..4fb5a628 100644 --- a/src/storage/models.rs +++ b/src/storage/models.rs @@ -225,26 +225,6 @@ use crate::sqlxextend; use crate::types; pub use types::DbType; -/* --------------------- models::InternalTx -----------------------------*/ -impl sqlxextend::TableSchemas for InternalTx { - fn table_name() -> &'static str { - INTERNALTX - } - const ARGN: i32 = 5; -} - -impl sqlxextend::BindQueryArg<'_, DbType> for InternalTx { - fn bind_args<'g, 'q: 'g>(&'q self, arg: &mut impl sqlx::Arguments<'g, Database = DbType>) { - arg.add(self.time); - arg.add(&self.user_from); - arg.add(&self.user_to); - arg.add(&self.asset); - arg.add(self.amount); - } -} - -impl sqlxextend::SqlxAction<'_, sqlxextend::InsertTable, DbType> for InternalTx {} - /* --------------------- models::BalanceHistory -----------------------------*/ impl sqlxextend::TableSchemas for BalanceHistory { fn table_name() -> &'static str {