diff --git a/.gitignore b/.gitignore index ca1b0defd..16b5a7263 100644 --- a/.gitignore +++ b/.gitignore @@ -19,3 +19,6 @@ build *.node .venv/ python/.venv/ +.env +.env.local +__pycache__/ diff --git a/CHANGELOG.md b/CHANGELOG.md index 9231b6f25..4e9420db9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,15 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +### Added + +- **Rust core:** `QuoteTransport` config option (`Config::quote_transport` / `Config::set_quote_transport`, env `LONGBRIDGE_QUOTE_TRANSPORT=ws|http`, default `WebSocket`). With `QuoteTransport::Http`, the 20 `QuoteContext` pull methods backed by the 19 WebSocket commands that have a REST equivalent (`static_info`, `quote`, `option_quote`, `warrant_quote`, `depth`, `brokers`, `participants`, `trades`, `intraday`, `candlesticks`, `history_candlesticks_by_offset`, `history_candlesticks_by_date`, `option_chain_expiry_date_list`, `warrant_issuers`, `warrant_list`, `trading_session`, `trading_days`, `capital_flow`, `capital_distribution`, `calc_indexes`) are served over `POST /quote/*` instead of the quote WebSocket, so a process that only pulls data never opens a WebSocket connection. Public method signatures and return types are unchanged: the gateway proto-JSON is decoded into the same protobuf messages as the WebSocket path (tolerating int64-as-string, omitted/`null` fields and camelCase keys) and the existing conversions are reused. Subscriptions, push events, `realtime_*`, `member_id` / `quote_level` / `quote_package_details` still use the WebSocket. + - **US overnight alignment:** the WebSocket only returns US overnight data when `Config::enable_overnight` is set, while REST always does. On the HTTP path (and only there) the SDK applies the WebSocket's rule for US equities: `quote().overnight_quote` is cleared, overnight `intraday` lines are dropped, and for `TradeSessions::All` candlestick queries overnight bars are dropped and count/range-capped windows are topped up with `POST /quote/history-candlesticks` offset queries (≤1000 bars per request, ≤8 extra requests; a failure in a top-up request fails the call). `trades` and `calc_indexes` are passed through unchanged. + - **Known differences:** HTTP errors are `Error::HttpClient(OpenApi { .. })` rather than `Error::WsClient(ResponseError { .. })` (`Error::openapi_error_code()` / `into_simple_error()` cover both, but a non-JSON gateway error surfaces as `HttpClientError::UnexpectedHttpResponse` with no code); the HTTP client retries `429` with back-off while the WebSocket does not; the REST gateway currently reports business errors (e.g. `301600`, `301607`) as `500`; during the US overnight session itself (20:00–04:00 New York) REST `quote` / `trades` values may reflect overnight trading that the WebSocket (without `enable_overnight`) does not report — not verified in this release. +- **C / C++ / Java / Node.js / Python:** the `QuoteTransport` enum and the corresponding config option are exposed in every binding layer. + ## [5.2.0] - 2026-09-30 ### Added diff --git a/c/README.md b/c/README.md index 3f0e74471..45b8e1cf3 100644 --- a/c/README.md +++ b/c/README.md @@ -160,6 +160,7 @@ setx LONGBRIDGE_ACCESS_TOKEN "Access Token get from user center" | LONGBRIDGE_TRADE_WS_URL | Trade websocket endpoint url (Default: `wss://openapi-trade.longbridge.com/v2`) | | LONGBRIDGE_ENABLE_OVERNIGHT | Enable overnight quote, `true` or `false` (Default: `false`) | | LONGBRIDGE_PUSH_CANDLESTICK_MODE | `realtime` or `confirmed` (Default: `realtime`) | +| LONGBRIDGE_QUOTE_TRANSPORT | Transport for quote pull APIs, `ws` or `http` (Default: `ws`) | | LONGBRIDGE_PRINT_QUOTE_PACKAGES | Print quote packages when connected, `true` or `false` (Default: `true`) | | LONGBRIDGE_LOG_PATH | Set the path of the log files (Default: `no logs`) | | LONGBRIDGE_PAPERTRADING | Enable paper trading mode, `true` or `false` (Default: `false`). See [Paper Trading](#paper-trading). | diff --git a/c/cbindgen.toml b/c/cbindgen.toml index 24750a1fd..b72e6b749 100644 --- a/c/cbindgen.toml +++ b/c/cbindgen.toml @@ -9,6 +9,7 @@ cpp_compat = true "CHttpResult" = "lb_http_result_t" "CLanguage" = "lb_language_t" "CPushCandlestickMode" = "lb_push_candlestick_mode_t" +"CQuoteTransport" = "lb_quote_transport_t" "CMarket" = "lb_market_t" "CDecimal" = "lb_decimal_t" "CAsyncCallback" = "lb_async_callback_t" diff --git a/c/csrc/include/longbridge.h b/c/csrc/include/longbridge.h index 9d281ca68..6c0e1590a 100644 --- a/c/csrc/include/longbridge.h +++ b/c/csrc/include/longbridge.h @@ -317,6 +317,22 @@ typedef enum lb_push_candlestick_mode_t { PushCandlestickMode_Confirmed, } lb_push_candlestick_mode_t; +/** + * Transport used by the quote pull APIs + */ +typedef enum lb_quote_transport_t { + /** + * Send pull requests over the quote WebSocket connection (default) + */ + QuoteTransport_WebSocket, + /** + * Send pull requests over HTTP (`POST /quote/...`) where the API has a + * REST equivalent, falling back to the WebSocket for the rest. Using + * only HTTP-backed APIs never opens a WebSocket connection. + */ + QuoteTransport_Http, +} lb_quote_transport_t; + /** * DCA investment frequency */ @@ -14318,9 +14334,10 @@ void lb_calendar_context_finance_calendar(const struct lb_calendar_context_t *ct * Optional environment variables are read automatically: * `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_QUOTE_WS_URL`, * `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_ENABLE_OVERNIGHT`, - * `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, - * `LONGBRIDGE_LOG_PATH`. Use the corresponding `lb_config_set_*` functions - * to override any of these values after construction. + * `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, + * `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`. Use the + * corresponding `lb_config_set_*` functions to override any of these values + * after construction. * * @param app_key App key * @param app_secret App secret @@ -14339,7 +14356,8 @@ struct lb_config_t *lb_config_from_apikey(const char *app_key, * `LONGBRIDGE_ACCESS_TOKEN`, `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_QUOTE_WS_URL`, * `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_LANGUAGE`, * `LONGBRIDGE_ENABLE_OVERNIGHT`, `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, - * `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH` + * `LONGBRIDGE_QUOTE_TRANSPORT`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, + * `LONGBRIDGE_LOG_PATH` */ struct lb_config_t *lb_config_from_apikey_env(struct lb_error_t **error); @@ -14349,9 +14367,10 @@ struct lb_config_t *lb_config_from_apikey_env(struct lb_error_t **error); * Optional environment variables are read automatically: * `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_QUOTE_WS_URL`, * `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_ENABLE_OVERNIGHT`, - * `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, - * `LONGBRIDGE_LOG_PATH`. Use the corresponding `lb_config_set_*` functions - * to override any of these values after construction. + * `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, + * `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`. Use the + * corresponding `lb_config_set_*` functions to override any of these values + * after construction. * * Does **not** take ownership of `oauth`. The caller must free `oauth` with * `lb_oauth_free` after this call returns. @@ -14408,6 +14427,17 @@ void lb_config_enable_overnight(struct lb_config_t *config); void lb_config_set_push_candlestick_mode(struct lb_config_t *config, enum lb_push_candlestick_mode_t mode); +/** + * Set the transport used by the quote pull APIs + * + * Default: `QuoteTransport_WebSocket` (or `LONGBRIDGE_QUOTE_TRANSPORT` = + * `ws` / `http`) + * + * @param config Config object + * @param transport Quote transport + */ +void lb_config_set_quote_transport(struct lb_config_t *config, enum lb_quote_transport_t transport); + /** * Disable printing of quote packages on connection * diff --git a/c/src/config.rs b/c/src/config.rs index 0b7f20c2c..8d17f2446 100644 --- a/c/src/config.rs +++ b/c/src/config.rs @@ -11,7 +11,7 @@ use crate::{ async_call::{CAsyncCallback, execute_async}, error::{CError, set_error}, oauth::COAuth, - types::{CLanguage, CPushCandlestickMode, CString}, + types::{CLanguage, CPushCandlestickMode, CQuoteTransport, CString}, }; /// Configuration options for Longbridge SDK @@ -22,9 +22,10 @@ pub struct CConfig(pub(crate) Config); /// Optional environment variables are read automatically: /// `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_QUOTE_WS_URL`, /// `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_ENABLE_OVERNIGHT`, -/// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, -/// `LONGBRIDGE_LOG_PATH`. Use the corresponding `lb_config_set_*` functions -/// to override any of these values after construction. +/// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, +/// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`. Use the +/// corresponding `lb_config_set_*` functions to override any of these values +/// after construction. /// /// @param app_key App key /// @param app_secret App secret @@ -54,7 +55,8 @@ pub unsafe extern "C" fn lb_config_from_apikey( /// `LONGBRIDGE_ACCESS_TOKEN`, `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_QUOTE_WS_URL`, /// `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_LANGUAGE`, /// `LONGBRIDGE_ENABLE_OVERNIGHT`, `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, -/// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH` +/// `LONGBRIDGE_QUOTE_TRANSPORT`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, +/// `LONGBRIDGE_LOG_PATH` #[unsafe(no_mangle)] pub unsafe extern "C" fn lb_config_from_apikey_env(error: *mut *mut CError) -> *mut CConfig { match Config::from_apikey_env() { @@ -74,9 +76,10 @@ pub unsafe extern "C" fn lb_config_from_apikey_env(error: *mut *mut CError) -> * /// Optional environment variables are read automatically: /// `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_QUOTE_WS_URL`, /// `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_ENABLE_OVERNIGHT`, -/// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, -/// `LONGBRIDGE_LOG_PATH`. Use the corresponding `lb_config_set_*` functions -/// to override any of these values after construction. +/// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, +/// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`. Use the +/// corresponding `lb_config_set_*` functions to override any of these values +/// after construction. /// /// Does **not** take ownership of `oauth`. The caller must free `oauth` with /// `lb_oauth_free` after this call returns. @@ -163,6 +166,21 @@ pub unsafe extern "C" fn lb_config_set_push_candlestick_mode( (*config).0.set_push_candlestick_mode(mode.into()); } +/// Set the transport used by the quote pull APIs +/// +/// Default: `QuoteTransport_WebSocket` (or `LONGBRIDGE_QUOTE_TRANSPORT` = +/// `ws` / `http`) +/// +/// @param config Config object +/// @param transport Quote transport +#[unsafe(no_mangle)] +pub unsafe extern "C" fn lb_config_set_quote_transport( + config: *mut CConfig, + transport: CQuoteTransport, +) { + (*config).0.set_quote_transport(transport.into()); +} + /// Disable printing of quote packages on connection /// /// @param config Config object diff --git a/c/src/types/mod.rs b/c/src/types/mod.rs index dccae30c4..a87dfb5ff 100644 --- a/c/src/types/mod.rs +++ b/c/src/types/mod.rs @@ -6,6 +6,7 @@ mod language; mod market; mod option; mod push_candlestick_mode; +mod quote_transport; mod string; use std::{ffi::CStr, os::raw::c_char}; @@ -18,6 +19,7 @@ pub(crate) use language::CLanguage; pub(crate) use market::CMarket; pub(crate) use option::COption; pub(crate) use push_candlestick_mode::CPushCandlestickMode; +pub(crate) use quote_transport::CQuoteTransport; pub(crate) use string::CString; pub(crate) trait ToFFI { diff --git a/c/src/types/quote_transport.rs b/c/src/types/quote_transport.rs new file mode 100644 index 000000000..7bf18f1ca --- /dev/null +++ b/c/src/types/quote_transport.rs @@ -0,0 +1,17 @@ +use longbridge_c_macros::CEnum; + +/// Transport used by the quote pull APIs +#[derive(Debug, Copy, Clone, Eq, PartialEq, CEnum)] +#[c(remote = "longbridge::QuoteTransport")] +#[allow(clippy::enum_variant_names, non_camel_case_types)] +#[repr(C)] +pub enum CQuoteTransport { + /// Send pull requests over the quote WebSocket connection (default) + #[c(remote = "WebSocket")] + QuoteTransport_WebSocket, + /// Send pull requests over HTTP (`POST /quote/...`) where the API has a + /// REST equivalent, falling back to the WebSocket for the rest. Using + /// only HTTP-backed APIs never opens a WebSocket connection. + #[c(remote = "Http")] + QuoteTransport_Http, +} diff --git a/cpp/README.md b/cpp/README.md index be58e35af..f83cc13d0 100644 --- a/cpp/README.md +++ b/cpp/README.md @@ -155,6 +155,7 @@ setx LONGBRIDGE_ACCESS_TOKEN "Access Token get from user center" | LONGBRIDGE_TRADE_WS_URL | Trade websocket endpoint url (Default: `wss://openapi-trade.longbridge.com/v2`) | | LONGBRIDGE_ENABLE_OVERNIGHT | Enable overnight quote, `true` or `false` (Default: `false`) | | LONGBRIDGE_PUSH_CANDLESTICK_MODE | `realtime` or `confirmed` (Default: `realtime`) | +| LONGBRIDGE_QUOTE_TRANSPORT | Transport for quote pull APIs, `ws` or `http` (Default: `ws`) | | LONGBRIDGE_PRINT_QUOTE_PACKAGES | Print quote packages when connected, `true` or `false` (Default: `true`) | | LONGBRIDGE_LOG_PATH | Set the path of the log files (Default: `no logs`) | | LONGBRIDGE_PAPERTRADING | Enable paper trading mode, `true` or `false` (Default: `false`). See [Paper Trading](#paper-trading). | diff --git a/cpp/include/config.hpp b/cpp/include/config.hpp index 928c5b80e..320d4e8d9 100644 --- a/cpp/include/config.hpp +++ b/cpp/include/config.hpp @@ -32,7 +32,8 @@ class Config /// Optional environment variables are read automatically: /// `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_QUOTE_WS_URL`, /// `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_ENABLE_OVERNIGHT`, - /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, + /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, + /// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, /// `LONGBRIDGE_LOG_PATH`. Use the chainable `set_*` methods to override any /// of these values. /// @@ -48,7 +49,8 @@ class Config /// Variables: `LONGBRIDGE_APP_KEY`, `LONGBRIDGE_APP_SECRET`, /// `LONGBRIDGE_ACCESS_TOKEN`, `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_QUOTE_WS_URL`, /// `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_ENABLE_OVERNIGHT`, - /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, + /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, + /// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, /// `LONGBRIDGE_LOG_PATH` static Config from_apikey_env(Status& status); @@ -57,7 +59,8 @@ class Config /// Optional environment variables are read automatically: /// `LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_QUOTE_WS_URL`, /// `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_ENABLE_OVERNIGHT`, - /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, + /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, + /// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, /// `LONGBRIDGE_LOG_PATH`. Use the chainable `set_*` methods to override any /// of these values. /// @@ -96,6 +99,12 @@ class Config /// Set the push candlestick mode Config& set_push_candlestick_mode(PushCandlestickMode mode); + /// Set the transport used by the quote pull APIs + /// + /// Default: `QuoteTransport::WebSocket` (or `LONGBRIDGE_QUOTE_TRANSPORT` = + /// `ws` / `http`) + Config& set_quote_transport(QuoteTransport transport); + /// Disable printing of quote packages on connection Config& disable_print_quote_packages(); diff --git a/cpp/include/types.hpp b/cpp/include/types.hpp index f4114ae69..8959eeffe 100644 --- a/cpp/include/types.hpp +++ b/cpp/include/types.hpp @@ -47,6 +47,17 @@ enum class PushCandlestickMode Confirmed, }; +/// Transport used by the quote pull APIs +enum class QuoteTransport +{ + /// Send pull requests over the quote WebSocket connection (default) + WebSocket, + /// Send pull requests over HTTP (`POST /quote/*`) where the API has a REST + /// equivalent, falling back to the WebSocket for the rest. Using only + /// HTTP-backed APIs never opens a WebSocket connection. + Http, +}; + /// Market enum class Market { diff --git a/cpp/src/config.cpp b/cpp/src/config.cpp index f2c40c9cd..89033eea8 100644 --- a/cpp/src/config.cpp +++ b/cpp/src/config.cpp @@ -112,6 +112,13 @@ Config::set_push_candlestick_mode(PushCandlestickMode mode) return *this; } +Config& +Config::set_quote_transport(QuoteTransport transport) +{ + lb_config_set_quote_transport(config_, convert::convert(transport)); + return *this; +} + Config& Config::disable_print_quote_packages() { diff --git a/cpp/src/convert.hpp b/cpp/src/convert.hpp index b15920a2a..1d28123c5 100644 --- a/cpp/src/convert.hpp +++ b/cpp/src/convert.hpp @@ -229,6 +229,19 @@ convert(PushCandlestickMode mode) } } +inline lb_quote_transport_t +convert(QuoteTransport transport) +{ + switch (transport) { + case QuoteTransport::WebSocket: + return QuoteTransport_WebSocket; + case QuoteTransport::Http: + return QuoteTransport_Http; + default: + throw std::invalid_argument("unreachable"); + } +} + inline Market convert(lb_market_t market) { diff --git a/examples/rust/Cargo.toml b/examples/rust/Cargo.toml index 11af435b8..d97343ed2 100644 --- a/examples/rust/Cargo.toml +++ b/examples/rust/Cargo.toml @@ -9,4 +9,5 @@ members = [ "subscribe_candlesticks", "http_client", "grid_trading", + "quote_http_compare", ] diff --git a/examples/rust/quote_http_compare/Cargo.toml b/examples/rust/quote_http_compare/Cargo.toml new file mode 100644 index 000000000..fab1ba897 --- /dev/null +++ b/examples/rust/quote_http_compare/Cargo.toml @@ -0,0 +1,9 @@ +[package] +name = "quote_http_compare" +version = "0.1.0" +edition.workspace = true + +[dependencies] +longbridge = { path = "../../../rust" } +tokio = { version = "1.19", features = ["rt-multi-thread", "macros"] } +time = { version = "0.3", features = ["macros"] } diff --git a/examples/rust/quote_http_compare/src/main.rs b/examples/rust/quote_http_compare/src/main.rs new file mode 100644 index 000000000..f3a9b4ca1 --- /dev/null +++ b/examples/rust/quote_http_compare/src/main.rs @@ -0,0 +1,356 @@ +//! Compare every quote pull API over the WebSocket vs HTTP transports. +//! +//! Reads credentials from `LONGBRIDGE_*` env (or `.env`). Point +//! `LONGBRIDGE_HTTP_URL` / `LONGBRIDGE_QUOTE_WS_URL` at an environment that +//! serves the `/quote/*` routes (canary). Dates below are hard-coded for the +//! October 2026 verification runs and must be moved along before reuse. For +//! each call prints the arguments, +//! both results (`Debug`), latency and whether WS == HTTP. +use std::{sync::Arc, time::Instant}; + +use longbridge::{ + quote::{ + AdjustType, CalcIndex, Period, QuoteContext, SortOrderType, TradeSessions, WarrantSortBy, + }, + Config, Market, QuoteTransport, +}; +use time::macros::{date, datetime}; + +macro_rules! compare { + ($ws:ident, $http:ident, $label:expr, $args:expr, |$ctx:ident| $call:expr) => {{ + let t = Instant::now(); + let a = { + let $ctx = &$ws; + $call.await.map_err(|e| e.into_simple_error()) + }; + let ws_ms = t.elapsed().as_millis(); + let t = Instant::now(); + let b = { + let $ctx = &$http; + $call.await.map_err(|e| e.into_simple_error()) + }; + let http_ms = t.elapsed().as_millis(); + // Errors are compared by business code only: WS and HTTP differ in + // the error variant / trace_id by construction. + let (a, b) = match (&a, &b) { + (Err(x), Err(y)) => ( + format!("Err(code={:?}, message={:?})", x.code(), x.message()), + format!("Err(code={:?}, message={:?})", y.code(), y.message()), + ), + _ => (format!("{a:#?}"), format!("{b:#?}")), + }; + println!("=== {}", $label); + println!("REQ {}", $args); + println!("WS_MS {ws_ms}\nHTTP_MS {http_ms}"); + println!("EQUAL {}", a == b); + println!("WS<<\n{a}\n>>WS\nHTTP<<\n{b}\n>>HTTP"); + }}; +} + +#[tokio::main] +async fn main() -> Result<(), Box> { + let base = Config::from_apikey_env()?; + let (ws, _) = QuoteContext::new(Arc::new( + base.clone().quote_transport(QuoteTransport::WebSocket), + )); + let (http, _) = QuoteContext::new(Arc::new(base.quote_transport(QuoteTransport::Http))); + + // Resolve an option and a warrant symbol to use below. + let expiries = ws.option_chain_expiry_date_list("AAPL.US").await?; + let Some(&first_expiry) = expiries.first() else { + return Err("no option expiry dates for AAPL.US".into()); + }; + // `option_chain_info_by_date` is not deployed on every environment; fall + // back to building an at-the-money-ish contract symbol from the expiry. + let option_symbol = match ws + .option_chain_info_by_date("AAPL.US", first_expiry, false) + .await + { + Ok(chain) if !chain.is_empty() => chain[chain.len() / 2].symbol.clone(), + _ => { + let d = first_expiry; + format!( + "AAPL{:02}{:02}{:02}C340000.US", + d.year() % 100, + d.month() as u8, + d.day() + ) + } + }; + let warrants = ws + .warrant_list( + "700.HK", + WarrantSortBy::LastDone, + SortOrderType::Descending, + None, + None, + None, + None, + None, + ) + .await?; + let Some(warrant_symbol) = warrants.first().map(|w| w.symbol.clone()) else { + return Err("no warrants for 700.HK".into()); + }; + + compare!(ws, http, "static_info", r#"["700.HK","AAPL.US"]"#, |c| c + .static_info(["700.HK", "AAPL.US"])); + compare!(ws, http, "quote", r#"["700.HK","AAPL.US"]"#, |c| c + .quote(["700.HK", "AAPL.US"])); + compare!( + ws, + http, + "option_quote", + format!("[{option_symbol:?}]"), + |c| c.option_quote([option_symbol.clone()]) + ); + compare!( + ws, + http, + "warrant_quote", + format!("[{warrant_symbol:?}]"), + |c| c.warrant_quote([warrant_symbol.clone()]) + ); + for symbol in ["700.HK", "AAPL.US"] { + compare!(ws, http, format!("depth {symbol}"), symbol, |c| c + .depth(symbol)); + } + compare!(ws, http, "brokers 700.HK", "700.HK", |c| c + .brokers("700.HK")); + compare!(ws, http, "participants", "-", |c| c.participants()); + compare!(ws, http, "trades 700.HK", "700.HK, count=20", |c| c + .trades("700.HK", 20)); + compare!(ws, http, "intraday AAPL.US", "AAPL.US, All", |c| c + .intraday("AAPL.US", TradeSessions::All)); + compare!( + ws, + http, + "candlesticks AAPL.US", + "AAPL.US, Day, 20, NoAdjust, All", + |c| c.candlesticks( + "AAPL.US", + Period::Day, + 20, + AdjustType::NoAdjust, + TradeSessions::All + ) + ); + compare!( + ws, + http, + "candlesticks 700.HK 1m", + "700.HK, OneMinute, 30, ForwardAdjust, Intraday", + |c| c.candlesticks( + "700.HK", + Period::OneMinute, + 30, + AdjustType::ForwardAdjust, + TradeSessions::Intraday + ) + ); + compare!( + ws, + http, + "candlesticks AAPL.US 1m All", + "AAPL.US, OneMinute, 1000, NoAdjust, All", + |c| c.candlesticks( + "AAPL.US", + Period::OneMinute, + 1000, + AdjustType::NoAdjust, + TradeSessions::All + ) + ); + compare!( + ws, + http, + "history_candlesticks_by_date AAPL.US 1m All", + "AAPL.US, OneMinute, NoAdjust, 2026-10-08..2026-10-09, All", + |c| c.history_candlesticks_by_date( + "AAPL.US", + Period::OneMinute, + AdjustType::NoAdjust, + Some(date!(2026 - 10 - 08)), + Some(date!(2026 - 10 - 09)), + TradeSessions::All + ) + ); + for forward in [false, true] { + compare!( + ws, + http, + format!("history_candlesticks_by_offset AAPL.US 1m All forward={forward}"), + format!("AAPL.US, OneMinute, NoAdjust, forward={forward}, 2026-10-08 12:00, 600, All"), + |c| c.history_candlesticks_by_offset( + "AAPL.US", + Period::OneMinute, + AdjustType::NoAdjust, + forward, + Some(datetime!(2026-10-08 12:00)), + 600, + TradeSessions::All + ) + ); + } + compare!( + ws, + http, + "history_candlesticks_by_date AAPL.US 1m All (uncapped)", + "AAPL.US, OneMinute, NoAdjust, 2026-10-09..2026-10-09, All", + |c| c.history_candlesticks_by_date( + "AAPL.US", + Period::OneMinute, + AdjustType::NoAdjust, + Some(date!(2026 - 10 - 09)), + Some(date!(2026 - 10 - 09)), + TradeSessions::All + ) + ); + compare!(ws, http, "intraday .VIX.US (index)", ".VIX.US, All", |c| c + .intraday(".VIX.US", TradeSessions::All)); + compare!(ws, http, "intraday 700.HK", "700.HK, Intraday", |c| c + .intraday("700.HK", TradeSessions::Intraday)); + compare!(ws, http, "trades AAPL.US", "AAPL.US, count=50", |c| c + .trades("AAPL.US", 50)); + compare!( + ws, + http, + "history_candlesticks_by_offset 700.HK", + "700.HK, Day, NoAdjust, forward=false, 2026-06-01 00:00, 10, Intraday", + |c| c.history_candlesticks_by_offset( + "700.HK", + Period::Day, + AdjustType::NoAdjust, + false, + Some(datetime!(2026-06-01 00:00)), + 10, + TradeSessions::Intraday + ) + ); + compare!( + ws, + http, + "history_candlesticks_by_date AAPL.US", + "AAPL.US, Week, ForwardAdjust, 2026-01-01..2026-03-01, All", + |c| c.history_candlesticks_by_date( + "AAPL.US", + Period::Week, + AdjustType::ForwardAdjust, + Some(date!(2026 - 01 - 01)), + Some(date!(2026 - 03 - 01)), + TradeSessions::All + ) + ); + compare!( + ws, + http, + "option_chain_expiry_date_list AAPL.US", + "AAPL.US", + |c| c.option_chain_expiry_date_list("AAPL.US") + ); + compare!(ws, http, "warrant_issuers", "-", |c| c.warrant_issuers()); + compare!( + ws, + http, + "warrant_list 700.HK", + "700.HK, ChangeRate, Ascending, no filters", + |c| c.warrant_list( + "700.HK", + WarrantSortBy::ChangeRate, + SortOrderType::Ascending, + None, + None, + None, + None, + None + ) + ); + compare!(ws, http, "trading_session", "-", |c| c.trading_session()); + compare!( + ws, + http, + "trading_days HK", + "HK, 2026-10-01..2026-10-31", + |c| c.trading_days(Market::HK, date!(2026 - 10 - 01), date!(2026 - 10 - 31)) + ); + compare!(ws, http, "capital_flow 700.HK", "700.HK", |c| c + .capital_flow("700.HK")); + compare!(ws, http, "capital_distribution 700.HK", "700.HK", |c| c + .capital_distribution("700.HK")); + compare!( + ws, + http, + "calc_indexes", + "[700.HK, AAPL.US, option, warrant], all indexes", + |c| c.calc_indexes( + [ + "700.HK".to_string(), + "AAPL.US".to_string(), + option_symbol.clone(), + warrant_symbol.clone() + ], + [ + CalcIndex::LastDone, + CalcIndex::ChangeValue, + CalcIndex::ChangeRate, + CalcIndex::Volume, + CalcIndex::Turnover, + CalcIndex::YtdChangeRate, + CalcIndex::TurnoverRate, + CalcIndex::TotalMarketValue, + CalcIndex::CapitalFlow, + CalcIndex::Amplitude, + CalcIndex::VolumeRatio, + CalcIndex::PeTtmRatio, + CalcIndex::PbRatio, + CalcIndex::DividendRatioTtm, + CalcIndex::FiveDayChangeRate, + CalcIndex::TenDayChangeRate, + CalcIndex::HalfYearChangeRate, + CalcIndex::FiveMinutesChangeRate, + CalcIndex::ExpiryDate, + CalcIndex::StrikePrice, + CalcIndex::UpperStrikePrice, + CalcIndex::LowerStrikePrice, + CalcIndex::OutstandingQty, + CalcIndex::OutstandingRatio, + CalcIndex::Premium, + CalcIndex::ItmOtm, + CalcIndex::ImpliedVolatility, + CalcIndex::WarrantDelta, + CalcIndex::CallPrice, + CalcIndex::ToCallPrice, + CalcIndex::EffectiveLeverage, + CalcIndex::LeverageRatio, + CalcIndex::ConversionRatio, + CalcIndex::BalancePoint, + CalcIndex::OpenInterest, + CalcIndex::Delta, + CalcIndex::Gamma, + CalcIndex::Theta, + CalcIndex::Vega, + CalcIndex::Rho, + ] + ) + ); + + // Error paths: the business error code must match across transports. + compare!(ws, http, "ERR depth NOTEXIST.XX", "NOTEXIST.XX", |c| c + .depth("NOTEXIST.XX")); + compare!(ws, http, "ERR quote []", "[]", |c| c + .quote(Vec::::new())); + compare!( + ws, + http, + "ERR candlesticks count=5000", + "700.HK, Day, 5000", + |c| c.candlesticks( + "700.HK", + Period::Day, + 5000, + AdjustType::NoAdjust, + TradeSessions::Intraday + ) + ); + Ok(()) +} diff --git a/java/Makefile.toml b/java/Makefile.toml index fa1a566b5..353711476 100644 --- a/java/Makefile.toml +++ b/java/Makefile.toml @@ -47,6 +47,7 @@ args = [ "javasrc/src/main/java/com/longbridge/OAuthBuilder.java", "javasrc/src/main/java/com/longbridge/OpenApiException.java", "javasrc/src/main/java/com/longbridge/PushCandlestickMode.java", + "javasrc/src/main/java/com/longbridge/QuoteTransport.java", "javasrc/src/main/java/com/longbridge/SdkNative.java", "javasrc/src/main/java/com/longbridge/HttpClient.java", diff --git a/java/README.md b/java/README.md index 9ab211e10..6321b17b7 100644 --- a/java/README.md +++ b/java/README.md @@ -150,6 +150,7 @@ setx LONGBRIDGE_ACCESS_TOKEN "Access Token get from user center" | LONGBRIDGE_TRADE_WS_URL | Trade websocket endpoint url (Default: `wss://openapi-trade.longbridge.com/v2`) | | LONGBRIDGE_ENABLE_OVERNIGHT | Enable overnight quote, `true` or `false` (Default: `false`) | | LONGBRIDGE_PUSH_CANDLESTICK_MODE | `realtime` or `confirmed` (Default: `realtime`) | +| LONGBRIDGE_QUOTE_TRANSPORT | Transport for quote pull APIs, `ws` or `http` (Default: `ws`) | | LONGBRIDGE_PRINT_QUOTE_PACKAGES | Print quote packages when connected, `true` or `false` (Default: `true`) | | LONGBRIDGE_LOG_PATH | Set the path of the log files (Default: `no logs`) | | LONGBRIDGE_PAPERTRADING | Enable paper trading mode, `true` or `false` (Default: `false`). See [Paper Trading](#paper-trading). | diff --git a/java/javasrc/src/main/java/com/longbridge/Config.java b/java/javasrc/src/main/java/com/longbridge/Config.java index df9a738f0..12273ba20 100644 --- a/java/javasrc/src/main/java/com/longbridge/Config.java +++ b/java/javasrc/src/main/java/com/longbridge/Config.java @@ -32,7 +32,8 @@ private long raw() { * {@code LONGBRIDGE_HTTP_URL}, {@code LONGBRIDGE_LANGUAGE}, * {@code LONGBRIDGE_QUOTE_WS_URL}, {@code LONGBRIDGE_TRADE_WS_URL}, * {@code LONGBRIDGE_ENABLE_OVERNIGHT}, {@code LONGBRIDGE_PUSH_CANDLESTICK_MODE}, - * {@code LONGBRIDGE_PRINT_QUOTE_PACKAGES}, {@code LONGBRIDGE_LOG_PATH}. + * {@code LONGBRIDGE_QUOTE_TRANSPORT}, {@code LONGBRIDGE_PRINT_QUOTE_PACKAGES}, + * {@code LONGBRIDGE_LOG_PATH}. * Use the chainable setter methods (e.g. {@link #httpUrl}) to override any of * these values. * @@ -69,6 +70,8 @@ public static Config fromApikey(String appKey, String appSecret, String accessTo * or {@code false} (Default: {@code false}) *
  • {@code LONGBRIDGE_PUSH_CANDLESTICK_MODE} - {@code realtime} or * {@code confirmed} (Default: {@code realtime})
  • + *
  • {@code LONGBRIDGE_QUOTE_TRANSPORT} - Transport for the quote pull APIs, + * {@code ws} or {@code http} (Default: {@code ws})
  • *
  • {@code LONGBRIDGE_PRINT_QUOTE_PACKAGES} - Print quote packages when * connected, {@code true} or {@code false} (Default: {@code true})
  • *
  • {@code LONGBRIDGE_LOG_PATH} - Set the path of the log files (Default: no @@ -92,7 +95,8 @@ public static Config fromApikeyEnv() throws OpenApiException { * {@code LONGBRIDGE_HTTP_URL}, {@code LONGBRIDGE_LANGUAGE}, * {@code LONGBRIDGE_QUOTE_WS_URL}, {@code LONGBRIDGE_TRADE_WS_URL}, * {@code LONGBRIDGE_ENABLE_OVERNIGHT}, {@code LONGBRIDGE_PUSH_CANDLESTICK_MODE}, - * {@code LONGBRIDGE_PRINT_QUOTE_PACKAGES}, {@code LONGBRIDGE_LOG_PATH}. + * {@code LONGBRIDGE_QUOTE_TRANSPORT}, {@code LONGBRIDGE_PRINT_QUOTE_PACKAGES}, + * {@code LONGBRIDGE_LOG_PATH}. * Use the chainable setter methods (e.g. {@link #httpUrl}) to override any of * these values. * @@ -175,6 +179,17 @@ public synchronized Config pushCandlestickMode(PushCandlestickMode mode) { return this; } + /** + * Set the transport used by the quote pull APIs. + * + * @param transport Transport (Default: {@link QuoteTransport#WebSocket}) + * @return this object + */ + public synchronized Config quoteTransport(QuoteTransport transport) { + this.raw = SdkNative.configSetQuoteTransport(raw(), transport); + return this; + } + /** * Disable printing quote packages when connected to the server. * diff --git a/java/javasrc/src/main/java/com/longbridge/QuoteTransport.java b/java/javasrc/src/main/java/com/longbridge/QuoteTransport.java new file mode 100644 index 000000000..466b6f7a9 --- /dev/null +++ b/java/javasrc/src/main/java/com/longbridge/QuoteTransport.java @@ -0,0 +1,17 @@ +package com.longbridge; + +/** + * Transport used by the quote pull APIs + */ +public enum QuoteTransport { + /** + * Send pull requests over the quote WebSocket connection (default) + */ + WebSocket, + /** + * Send pull requests over HTTP ({@code POST /quote/*}) where the API has a + * REST equivalent, falling back to the WebSocket for the rest. Using only + * HTTP-backed APIs never opens a WebSocket connection. + */ + Http, +} diff --git a/java/javasrc/src/main/java/com/longbridge/SdkNative.java b/java/javasrc/src/main/java/com/longbridge/SdkNative.java index a3c6b7456..5e296f386 100644 --- a/java/javasrc/src/main/java/com/longbridge/SdkNative.java +++ b/java/javasrc/src/main/java/com/longbridge/SdkNative.java @@ -50,6 +50,8 @@ public static native long newHttpClientFromApikey(String appKey, String appSecre public static native long configSetPushCandlestickMode(long config, PushCandlestickMode mode); + public static native long configSetQuoteTransport(long config, QuoteTransport transport); + public static native long configSetEnablePrintQuotePackages(long config, boolean enable); public static native long configSetLogPath(long config, String logPath); diff --git a/java/src/config.rs b/java/src/config.rs index ca4cbfdf6..fe579f142 100644 --- a/java/src/config.rs +++ b/java/src/config.rs @@ -3,7 +3,7 @@ use jni::{ objects::{JClass, JObject, JString}, sys::{jboolean, jlong}, }; -use longbridge::{Config, Language, PushCandlestickMode}; +use longbridge::{Config, Language, PushCandlestickMode, QuoteTransport}; use time::OffsetDateTime; use crate::{async_util, error::jni_result, types::FromJValue}; @@ -137,6 +137,20 @@ pub unsafe extern "system" fn Java_com_longbridge_SdkNative_configSetPushCandles }) } +#[unsafe(no_mangle)] +pub unsafe extern "system" fn Java_com_longbridge_SdkNative_configSetQuoteTransport( + mut env: JNIEnv, + _class: JClass, + config: jlong, + transport: JObject, +) -> jlong { + jni_result(&mut env, config, |env| { + let transport = QuoteTransport::from_jvalue(env, transport.into())?; + (*(config as *mut Config)).set_quote_transport(transport); + Ok(config) + }) +} + #[unsafe(no_mangle)] pub unsafe extern "system" fn Java_com_longbridge_SdkNative_configSetEnablePrintQuotePackages( mut env: JNIEnv, diff --git a/java/src/init.rs b/java/src/init.rs index 855aa87c0..154b7caa3 100644 --- a/java/src/init.rs +++ b/java/src/init.rs @@ -86,6 +86,7 @@ pub extern "system" fn Java_com_longbridge_SdkNative_init<'a>( longbridge::SimpleErrorKind, longbridge::Language, longbridge::PushCandlestickMode, + longbridge::QuoteTransport, longbridge::Market, longbridge::quote::TradeStatus, longbridge::quote::TradeSession, diff --git a/java/src/types/enum_types.rs b/java/src/types/enum_types.rs index cdc682da2..7089eccdb 100644 --- a/java/src/types/enum_types.rs +++ b/java/src/types/enum_types.rs @@ -30,6 +30,12 @@ impl_java_enum!( [Realtime, Confirmed] ); +impl_java_enum!( + "com/longbridge/QuoteTransport", + longbridge::QuoteTransport, + [WebSocket, Http] +); + impl_java_enum!( "com/longbridge/Market", longbridge::Market, diff --git a/nodejs/README.md b/nodejs/README.md index 198e59387..0ae213b5e 100644 --- a/nodejs/README.md +++ b/nodejs/README.md @@ -143,6 +143,7 @@ setx LONGBRIDGE_ACCESS_TOKEN "Access Token get from user center" | LONGBRIDGE_TRADE_WS_URL | Trade websocket endpoint url (Default: `wss://openapi-trade.longbridge.com/v2`) | | LONGBRIDGE_ENABLE_OVERNIGHT | Enable overnight quote, `true` or `false` (Default: `false`) | | LONGBRIDGE_PUSH_CANDLESTICK_MODE | `realtime` or `confirmed` (Default: `realtime`) | +| LONGBRIDGE_QUOTE_TRANSPORT | Transport for quote pull APIs, `ws` or `http` (Default: `ws`) | | LONGBRIDGE_PRINT_QUOTE_PACKAGES | Print quote packages when connected, `true` or `false` (Default: `true`) | | LONGBRIDGE_LOG_PATH | Set the path of the log files (Default: `no logs`) | | LONGBRIDGE_PAPERTRADING | Enable paper trading mode, `true` or `false` (Default: `false`). See [Paper Trading](#paper-trading). | diff --git a/nodejs/index.d.ts b/nodejs/index.d.ts index 5ef828a38..9ea50edaa 100644 --- a/nodejs/index.d.ts +++ b/nodejs/index.d.ts @@ -367,7 +367,8 @@ export declare class Config { * (`LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, * `LONGBRIDGE_QUOTE_WS_URL`, `LONGBRIDGE_TRADE_WS_URL`, * `LONGBRIDGE_ENABLE_OVERNIGHT`, `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, - * `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`). Fields + * `LONGBRIDGE_QUOTE_TRANSPORT`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, + * `LONGBRIDGE_LOG_PATH`). Fields * set in `extra` override the corresponding environment variables. * * @param appKey Application key @@ -407,6 +408,8 @@ export declare class Config { * `false` (Default: `false`) * - `LONGBRIDGE_PUSH_CANDLESTICK_MODE` - `realtime` or `confirmed` * (Default: `realtime`) + * - `LONGBRIDGE_QUOTE_TRANSPORT` - Transport for quote pull APIs, `ws` or + * `http` (Default: `ws`) * - `LONGBRIDGE_PRINT_QUOTE_PACKAGES` - Print quote packages when * connected, `true` or `false` (Default: `true`) * - `LONGBRIDGE_LOG_PATH` - Log file directory (Default: no logs) @@ -422,7 +425,8 @@ export declare class Config { * (`LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, * `LONGBRIDGE_QUOTE_WS_URL`, `LONGBRIDGE_TRADE_WS_URL`, * `LONGBRIDGE_ENABLE_OVERNIGHT`, `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, - * `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`). Fields + * `LONGBRIDGE_QUOTE_TRANSPORT`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, + * `LONGBRIDGE_LOG_PATH`). Fields * set in `extra` override the corresponding environment variables. * * @param oauth OAuth handle obtained from `OAuth.build(...)` @@ -6215,6 +6219,14 @@ export interface ExtraConfigParams { enableOvernight?: boolean /** Push candlesticks mode (default: PushCandlestickMode.Realtime) */ pushCandlestickMode?: PushCandlestickMode + /** + * Transport for quote pull (request/response) APIs (default: + * QuoteTransport.WebSocket) + * + * Subscriptions and push events always use the WebSocket connection; this + * only selects how request/response quote APIs reach the server. + */ + quoteTransport?: QuoteTransport /** * Enable printing the opened quote packages when connected to the server * (default: true). Set to `false` to suppress the output. @@ -7912,6 +7924,17 @@ export interface QuestionOption { description: string } +export declare const enum QuoteTransport { + /** Send pull requests over the quote WebSocket connection (default) */ + WebSocket = 0, + /** + * Send pull requests over HTTP (`POST /quote/*`) where the API has a REST + * equivalent, falling back to the WebSocket for the rest. Using only + * HTTP-backed APIs never opens a WebSocket connection. + */ + Http = 1 +} + /** Rank categories response. */ export interface RankCategoriesResponse { /** All top-level rank categories */ diff --git a/nodejs/index.js b/nodejs/index.js index efec33ac5..b5d3f4297 100644 --- a/nodejs/index.js +++ b/nodejs/index.js @@ -772,6 +772,7 @@ module.exports.OutsideRTH = nativeBinding.OutsideRTH module.exports.Period = nativeBinding.Period module.exports.PinnedMode = nativeBinding.PinnedMode module.exports.PushCandlestickMode = nativeBinding.PushCandlestickMode +module.exports.QuoteTransport = nativeBinding.QuoteTransport module.exports.SecuritiesUpdateMode = nativeBinding.SecuritiesUpdateMode module.exports.SecurityBoard = nativeBinding.SecurityBoard module.exports.SecurityListCategory = nativeBinding.SecurityListCategory diff --git a/nodejs/src/config.rs b/nodejs/src/config.rs index ff732cad5..74ab03f35 100644 --- a/nodejs/src/config.rs +++ b/nodejs/src/config.rs @@ -4,7 +4,7 @@ use napi::Result; use crate::{ error::ErrorNewType, oauth::OAuth, - types::{Language, PushCandlestickMode}, + types::{Language, PushCandlestickMode, QuoteTransport}, utils::from_datetime, }; @@ -26,6 +26,12 @@ pub struct ExtraConfigParams { pub enable_overnight: Option, /// Push candlesticks mode (default: PushCandlestickMode.Realtime) pub push_candlestick_mode: Option, + /// Transport for quote pull (request/response) APIs (default: + /// QuoteTransport.WebSocket) + /// + /// Subscriptions and push events always use the WebSocket connection; this + /// only selects how request/response quote APIs reach the server. + pub quote_transport: Option, /// Enable printing the opened quote packages when connected to the server /// (default: true). Set to `false` to suppress the output. pub enable_print_quote_packages: Option, @@ -67,6 +73,9 @@ fn apply_extra( if let Some(mode) = extra.push_candlestick_mode { config.set_push_candlestick_mode(mode.into()); } + if let Some(quote_transport) = extra.quote_transport { + config.set_quote_transport(quote_transport.into()); + } if let Some(false) = extra.enable_print_quote_packages { config.set_dont_print_quote_packages(); } @@ -92,7 +101,8 @@ impl Config { /// (`LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, /// `LONGBRIDGE_QUOTE_WS_URL`, `LONGBRIDGE_TRADE_WS_URL`, /// `LONGBRIDGE_ENABLE_OVERNIGHT`, `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, - /// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`). Fields + /// `LONGBRIDGE_QUOTE_TRANSPORT`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, + /// `LONGBRIDGE_LOG_PATH`). Fields /// set in `extra` override the corresponding environment variables. /// /// @param appKey Application key @@ -140,6 +150,8 @@ impl Config { /// `false` (Default: `false`) /// - `LONGBRIDGE_PUSH_CANDLESTICK_MODE` - `realtime` or `confirmed` /// (Default: `realtime`) + /// - `LONGBRIDGE_QUOTE_TRANSPORT` - Transport for quote pull APIs, `ws` or + /// `http` (Default: `ws`) /// - `LONGBRIDGE_PRINT_QUOTE_PACKAGES` - Print quote packages when /// connected, `true` or `false` (Default: `true`) /// - `LONGBRIDGE_LOG_PATH` - Log file directory (Default: no logs) @@ -159,7 +171,8 @@ impl Config { /// (`LONGBRIDGE_HTTP_URL`, `LONGBRIDGE_LANGUAGE`, /// `LONGBRIDGE_QUOTE_WS_URL`, `LONGBRIDGE_TRADE_WS_URL`, /// `LONGBRIDGE_ENABLE_OVERNIGHT`, `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, - /// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`). Fields + /// `LONGBRIDGE_QUOTE_TRANSPORT`, `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, + /// `LONGBRIDGE_LOG_PATH`). Fields /// set in `extra` override the corresponding environment variables. /// /// @param oauth OAuth handle obtained from `OAuth.build(...)` diff --git a/nodejs/src/types.rs b/nodejs/src/types.rs index 7114acea7..157e60b36 100644 --- a/nodejs/src/types.rs +++ b/nodejs/src/types.rs @@ -41,6 +41,18 @@ pub enum PushCandlestickMode { Confirmed, } +#[napi_derive::napi] +#[derive(Debug, JsEnum, Hash, Eq, PartialEq)] +#[js(remote = "longbridge::QuoteTransport")] +pub enum QuoteTransport { + /// Send pull requests over the quote WebSocket connection (default) + WebSocket, + /// Send pull requests over HTTP (`POST /quote/*`) where the API has a REST + /// equivalent, falling back to the WebSocket for the rest. Using only + /// HTTP-backed APIs never opens a WebSocket connection. + Http, +} + #[napi_derive::napi] #[derive(Debug, JsEnum, Hash, Eq, PartialEq, Copy, Clone)] #[js(remote = "longbridge::portfolio::types::FlowDirection")] diff --git a/python/README.md b/python/README.md index df12f8e22..c773f074d 100644 --- a/python/README.md +++ b/python/README.md @@ -179,6 +179,7 @@ setx LONGBRIDGE_ACCESS_TOKEN "Access Token get from user center" | LONGBRIDGE_TRADE_WS_URL | Trade websocket endpoint url (Default: `wss://openapi-trade.longbridge.com/v2`) | | LONGBRIDGE_ENABLE_OVERNIGHT | Enable overnight quote, `true` or `false` (Default: `false`) | | LONGBRIDGE_PUSH_CANDLESTICK_MODE | `realtime` or `confirmed` (Default: `realtime`) | +| LONGBRIDGE_QUOTE_TRANSPORT | Transport for quote pull APIs, `ws` or `http` (Default: `ws`) | | LONGBRIDGE_PRINT_QUOTE_PACKAGES | Print quote packages when connected, `true` or `false` (Default: `true`) | | LONGBRIDGE_LOG_PATH | Set the path of the log files (Default: `no logs`) | | LONGBRIDGE_PAPERTRADING | Enable paper trading mode, `true` or `false` (Default: `false`). See [Paper Trading](#paper-trading). | diff --git a/python/pysrc/longbridge/openapi.pyi b/python/pysrc/longbridge/openapi.pyi index aa4b7708a..463a16ed5 100644 --- a/python/pysrc/longbridge/openapi.pyi +++ b/python/pysrc/longbridge/openapi.pyi @@ -198,6 +198,23 @@ class PushCandlestickMode: Confirmed """ +class QuoteTransport: + """ + Transport used by the quote pull APIs + """ + + class WebSocket(QuoteTransport): + """ + Send pull requests over the quote WebSocket connection (default) + """ + + class Http(QuoteTransport): + """ + Send pull requests over HTTP (``POST /quote/*``) where the API has a + REST equivalent, falling back to the WebSocket for the rest. Using + only HTTP-backed APIs never opens a WebSocket connection. + """ + class OAuth: """ OAuth 2.0 client handle for Longbridge OpenAPI. @@ -291,6 +308,8 @@ class Config: enable_print_quote_packages: Print opened quote packages on connect (default: ``True``) log_path: Path for log files (default: no logs) + quote_transport: Transport used by the quote pull APIs + (default: ``QuoteTransport.WebSocket``) """ @staticmethod @@ -307,6 +326,7 @@ class Config: enable_print_quote_packages: bool = True, log_path: Optional[str] = None, enable_papertrading: bool = False, + quote_transport: Optional[Type[QuoteTransport]] = None, ) -> Config: """ Create a new ``Config`` using API Key authentication. @@ -316,7 +336,7 @@ class Config: ``LONGBRIDGE_QUOTE_WS_URL``, ``LONGBRIDGE_TRADE_WS_URL``, ``LONGBRIDGE_ENABLE_OVERNIGHT``, ``LONGBRIDGE_PUSH_CANDLESTICK_MODE``, ``LONGBRIDGE_PRINT_QUOTE_PACKAGES``, ``LONGBRIDGE_LOG_PATH``, - ``LONGBRIDGE_PAPERTRADING``). + ``LONGBRIDGE_PAPERTRADING``, ``LONGBRIDGE_QUOTE_TRANSPORT``). Any explicit parameter overrides the corresponding env variable. Args: @@ -341,6 +361,9 @@ class Config: the token: if it belongs to a real-money account the server returns an error. Default: ``False`` (no restriction imposed by server). + quote_transport: Transport used by the quote pull APIs (reads + ``LONGBRIDGE_QUOTE_TRANSPORT`` from env if omitted; default: + ``QuoteTransport.WebSocket``) """ @classmethod @@ -369,6 +392,7 @@ class Config: - ``LONGBRIDGE_PRINT_QUOTE_PACKAGES`` - ``true`` or ``false`` (Default: ``true``) - ``LONGBRIDGE_LOG_PATH`` - Log file directory (Default: no logs) + - ``LONGBRIDGE_QUOTE_TRANSPORT`` - ``ws`` or ``http`` (Default: ``ws``) """ @classmethod @@ -384,6 +408,7 @@ class Config: enable_print_quote_packages: Optional[bool] = None, log_path: Optional[str] = None, enable_papertrading: Optional[bool] = None, + quote_transport: Optional[Type[QuoteTransport]] = None, ) -> Config: """ Create a new ``Config`` for OAuth 2.0 authentication. @@ -396,7 +421,7 @@ class Config: ``LONGBRIDGE_QUOTE_WS_URL``, ``LONGBRIDGE_TRADE_WS_URL``, ``LONGBRIDGE_ENABLE_OVERNIGHT``, ``LONGBRIDGE_PUSH_CANDLESTICK_MODE``, ``LONGBRIDGE_PRINT_QUOTE_PACKAGES``, ``LONGBRIDGE_LOG_PATH``, - ``LONGBRIDGE_PAPERTRADING``). + ``LONGBRIDGE_PAPERTRADING``, ``LONGBRIDGE_QUOTE_TRANSPORT``). Any explicit parameter overrides the corresponding env variable. Args: @@ -420,6 +445,7 @@ class Config: the token: if it belongs to a real-money account the server returns an error. Default: ``None`` (no restriction imposed by server). + quote_transport: Transport used by the quote pull APIs (optional) Returns: Config object diff --git a/python/src/config.rs b/python/src/config.rs index cd5c4b35e..bd4ef7cf9 100644 --- a/python/src/config.rs +++ b/python/src/config.rs @@ -4,7 +4,7 @@ use crate::{ error::ErrorNewType, oauth::OAuth, time::PyOffsetDateTimeWrapper, - types::{Language, PushCandlestickMode}, + types::{Language, PushCandlestickMode, QuoteTransport}, }; #[pyclass(name = "Config")] @@ -18,7 +18,8 @@ impl Config { /// (``LONGBRIDGE_HTTP_URL``, ``LONGBRIDGE_LANGUAGE``, /// ``LONGBRIDGE_QUOTE_WS_URL``, ``LONGBRIDGE_TRADE_WS_URL``, /// ``LONGBRIDGE_ENABLE_OVERNIGHT``, ``LONGBRIDGE_PUSH_CANDLESTICK_MODE``, - /// ``LONGBRIDGE_PRINT_QUOTE_PACKAGES``, ``LONGBRIDGE_LOG_PATH``). + /// ``LONGBRIDGE_PRINT_QUOTE_PACKAGES``, ``LONGBRIDGE_LOG_PATH``, + /// ``LONGBRIDGE_QUOTE_TRANSPORT``). /// Any explicit parameter passed to this method overrides the /// corresponding environment variable. /// @@ -39,6 +40,9 @@ impl Config { /// mode enable_print_quote_packages: Print opened quote packages on /// connect (default: ``True``) /// log_path: Path for log files (default: no logs) + /// quote_transport: Transport used by the quote pull APIs (reads + /// ``LONGBRIDGE_QUOTE_TRANSPORT`` from env if omitted; default: + /// ``QuoteTransport.WebSocket``) #[staticmethod] #[pyo3(signature = ( app_key, @@ -53,6 +57,7 @@ impl Config { enable_print_quote_packages = true, log_path = None, enable_papertrading = false, + quote_transport = None, ))] #[allow(clippy::too_many_arguments)] fn from_apikey( @@ -68,6 +73,7 @@ impl Config { enable_print_quote_packages: bool, log_path: Option, enable_papertrading: bool, + quote_transport: Option, ) -> Self { let mut config = longbridge::Config::from_apikey(app_key, app_secret, access_token); @@ -96,6 +102,9 @@ impl Config { if enable_papertrading { config.set_enable_papertrading(); } + if let Some(quote_transport) = quote_transport { + config.set_quote_transport(quote_transport.into()); + } Self(config) } @@ -122,6 +131,8 @@ impl Config { /// - ``LONGBRIDGE_PRINT_QUOTE_PACKAGES`` - ``true`` or ``false`` /// (Default: ``true``) /// - ``LONGBRIDGE_LOG_PATH`` - Log file directory (Default: no logs) + /// - ``LONGBRIDGE_QUOTE_TRANSPORT`` - ``ws`` or ``http`` (Default: + /// ``ws``) #[classmethod] fn from_apikey_env(_cls: Bound) -> PyResult { Ok(Self( @@ -138,7 +149,8 @@ impl Config { /// (``LONGBRIDGE_HTTP_URL``, ``LONGBRIDGE_LANGUAGE``, /// ``LONGBRIDGE_QUOTE_WS_URL``, ``LONGBRIDGE_TRADE_WS_URL``, /// ``LONGBRIDGE_ENABLE_OVERNIGHT``, ``LONGBRIDGE_PUSH_CANDLESTICK_MODE``, - /// ``LONGBRIDGE_PRINT_QUOTE_PACKAGES``, ``LONGBRIDGE_LOG_PATH``). + /// ``LONGBRIDGE_PRINT_QUOTE_PACKAGES``, ``LONGBRIDGE_LOG_PATH``, + /// ``LONGBRIDGE_QUOTE_TRANSPORT``). /// Any explicit parameter passed to this method overrides the /// corresponding environment variable. /// @@ -158,6 +170,7 @@ impl Config { /// enable_print_quote_packages: Print opened quote packages on connect /// (optional) /// log_path: Path for log files (optional) + /// quote_transport: Transport used by the quote pull APIs (optional) /// /// Returns: /// Config object @@ -173,6 +186,7 @@ impl Config { enable_print_quote_packages = None, log_path = None, enable_papertrading = None, + quote_transport = None, ))] #[allow(clippy::too_many_arguments)] fn from_oauth( @@ -187,6 +201,7 @@ impl Config { enable_print_quote_packages: Option, log_path: Option, enable_papertrading: Option, + quote_transport: Option, ) -> Self { let mut config = longbridge::Config::from_oauth(oauth.0.clone()); @@ -217,6 +232,9 @@ impl Config { if let Some(true) = enable_papertrading { config.set_enable_papertrading(); } + if let Some(quote_transport) = quote_transport { + config.set_quote_transport(quote_transport.into()); + } Self(config) } diff --git a/python/src/lib.rs b/python/src/lib.rs index bab5cefc9..ccbd51697 100644 --- a/python/src/lib.rs +++ b/python/src/lib.rs @@ -41,6 +41,7 @@ fn longbridge(py: Python<'_>, m: Bound) -> PyResult<()> { openapi.add_class::()?; openapi.add_class::()?; openapi.add_class::()?; + openapi.add_class::()?; openapi.add_class::()?; openapi.add_class::()?; agent::register_types(&openapi)?; diff --git a/python/src/types.rs b/python/src/types.rs index ecb7f5153..540683ca4 100644 --- a/python/src/types.rs +++ b/python/src/types.rs @@ -42,3 +42,16 @@ pub(crate) enum PushCandlestickMode { /// Confirmed mode Confirmed, } + +#[pyclass(eq, eq_int, from_py_object)] +#[derive(Debug, PyEnum, Copy, Clone, Hash, Eq, PartialEq)] +#[allow(non_camel_case_types)] +#[py(remote = "longbridge::QuoteTransport")] +pub(crate) enum QuoteTransport { + /// Send pull requests over the quote WebSocket connection (default) + WebSocket, + /// Send pull requests over HTTP (``POST /quote/*``) where the API has a + /// REST equivalent, falling back to the WebSocket for the rest. Using only + /// HTTP-backed APIs never opens a WebSocket connection. + Http, +} diff --git a/rust/README.md b/rust/README.md index 81b4aba94..31473feb3 100644 --- a/rust/README.md +++ b/rust/README.md @@ -173,6 +173,7 @@ setx LONGBRIDGE_ACCESS_TOKEN "Access Token get from user center" | LONGBRIDGE_TRADE_WS_URL | Trade websocket endpoint url (Default: `wss://openapi-trade.longbridge.com/v2`) | | LONGBRIDGE_ENABLE_OVERNIGHT | Enable overnight quote, `true` or `false` (Default: `false`) | | LONGBRIDGE_PUSH_CANDLESTICK_MODE | `realtime` or `confirmed` (Default: `realtime`) | +| LONGBRIDGE_QUOTE_TRANSPORT | Transport for quote pull APIs, `ws` or `http` (Default: `ws`) | | LONGBRIDGE_PRINT_QUOTE_PACKAGES | Print quote packages when connected, `true` or `false` (Default: `true`) | | LONGBRIDGE_LOG_PATH | Set the path of the log files (Default: `no logs`) | | LONGBRIDGE_PAPERTRADING | Enable paper trading mode, `true` or `false` (Default: `false`). See [Paper Trading](#paper-trading). | diff --git a/rust/src/config.rs b/rust/src/config.rs index a26db71d7..7d361ee12 100644 --- a/rust/src/config.rs +++ b/rust/src/config.rs @@ -119,6 +119,38 @@ impl fmt::Debug for AuthMode { } } +/// Transport used by the [`QuoteContext`](crate::quote::QuoteContext) pull +/// (request/response) APIs. +/// +/// Subscriptions, push events, `realtime_*`, `member_id` / `quote_level` / +/// `quote_package_details` and the few APIs without a REST equivalent always +/// use the WebSocket connection; this only selects how the other pull APIs +/// reach the server. Results are the same either way: with `Http` the SDK +/// applies the WebSocket's US overnight rule (overnight data only when +/// [`Config::enable_overnight`] is set) client-side. +#[derive(Debug, Default, Copy, Clone, PartialEq, Eq)] +pub enum QuoteTransport { + /// Send pull requests over the quote WebSocket connection (default). + #[default] + WebSocket, + /// Send pull requests over HTTP (`POST /quote/*`) where the API has a REST + /// equivalent, falling back to the WebSocket for the rest. Using only + /// HTTP-backed APIs never opens a WebSocket connection. + Http, +} + +impl FromStr for QuoteTransport { + type Err = (); + + fn from_str(s: &str) -> ::std::result::Result { + match s.to_ascii_lowercase().as_str() { + "ws" => Ok(QuoteTransport::WebSocket), + "http" => Ok(QuoteTransport::Http), + _ => Err(()), + } + } +} + /// Configuration options for Longbridge SDK #[derive(Debug, Clone)] pub struct Config { @@ -134,6 +166,7 @@ pub struct Config { /// Extra headers injected into every HTTP and WebSocket upgrade request. pub(crate) custom_headers: HashMap, pub(crate) enable_papertrading: bool, + pub(crate) quote_transport: QuoteTransport, } /// Reads an env var by trying `LONGBRIDGE_` first, then falling back @@ -166,6 +199,7 @@ struct ConfigExtras { enable_print_quote_packages: bool, log_path: Option, enable_papertrading: bool, + quote_transport: QuoteTransport, } impl ConfigExtras { @@ -191,6 +225,15 @@ impl ConfigExtras { enable_print_quote_packages, log_path: env_var("LOG_PATH").map(PathBuf::from), enable_papertrading, + quote_transport: env_var("QUOTE_TRANSPORT") + .and_then(|v| { + let parsed = v.parse().ok(); + if parsed.is_none() { + tracing::warn!(value = %v, "invalid LONGBRIDGE_QUOTE_TRANSPORT, using ws"); + } + parsed + }) + .unwrap_or_default(), } } } @@ -201,7 +244,7 @@ impl Config { /// All optional environment variables (`LONGBRIDGE_HTTP_URL`, /// `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_QUOTE_WS_URL`, /// `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_ENABLE_OVERNIGHT`, - /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, + /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, /// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`) are read from /// the environment (or `.env` file) and applied automatically if set. /// @@ -230,6 +273,7 @@ impl Config { log_path: extras.log_path, custom_headers: Default::default(), enable_papertrading: extras.enable_papertrading, + quote_transport: extras.quote_transport, } } @@ -238,7 +282,7 @@ impl Config { /// All optional environment variables (`LONGBRIDGE_HTTP_URL`, /// `LONGBRIDGE_LANGUAGE`, `LONGBRIDGE_QUOTE_WS_URL`, /// `LONGBRIDGE_TRADE_WS_URL`, `LONGBRIDGE_ENABLE_OVERNIGHT`, - /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, + /// `LONGBRIDGE_PUSH_CANDLESTICK_MODE`, `LONGBRIDGE_QUOTE_TRANSPORT`, /// `LONGBRIDGE_PRINT_QUOTE_PACKAGES`, `LONGBRIDGE_LOG_PATH`) are read from /// the environment (or `.env` file) and applied automatically if set. /// @@ -280,6 +324,7 @@ impl Config { log_path: extras.log_path, custom_headers: Default::default(), enable_papertrading: extras.enable_papertrading, + quote_transport: extras.quote_transport, } } @@ -305,6 +350,8 @@ impl Config { /// `false` (Default: `false`) /// - `LONGBRIDGE_PUSH_CANDLESTICK_MODE` - `realtime` or `confirmed` /// (Default: `realtime`) + /// - `LONGBRIDGE_QUOTE_TRANSPORT` - Transport for the quote pull APIs, `ws` + /// or `http` (Default: `ws`) /// - `LONGBRIDGE_PRINT_QUOTE_PACKAGES` - Print quote packages when /// connected, `true` or `false` (Default: `true`) /// - `LONGBRIDGE_LOG_PATH` - Set the path of the log files (Default: `no @@ -336,6 +383,7 @@ impl Config { log_path: extras.log_path, custom_headers: Default::default(), enable_papertrading: extras.enable_papertrading, + quote_transport: extras.quote_transport, }) } @@ -428,6 +476,17 @@ impl Config { self } + /// Specifies the transport used by the quote pull APIs. + /// + /// Default: `QuoteTransport::WebSocket` (or `LONGBRIDGE_QUOTE_TRANSPORT` + /// = `ws` / `http`) + pub fn quote_transport(self, quote_transport: QuoteTransport) -> Self { + Self { + quote_transport, + ..self + } + } + /// Create metadata for auth/reconnect request pub fn create_metadata(&self) -> HashMap { let mut metadata = HashMap::new(); @@ -681,6 +740,13 @@ impl Config { self.enable_papertrading = true; } + /// Set the quote transport in place. + /// + /// See [`Config::quote_transport`] for full semantics. + pub fn set_quote_transport(&mut self, quote_transport: QuoteTransport) { + self.quote_transport = quote_transport; + } + /// Set the log path in place. pub fn set_log_path(&mut self, path: impl Into) { self.log_path = Some(path.into()); diff --git a/rust/src/lib.rs b/rust/src/lib.rs index 4f99bc0de..e8f47ad6c 100644 --- a/rust/src/lib.rs +++ b/rust/src/lib.rs @@ -41,7 +41,7 @@ pub use agent::AgentContext; pub use alert::AlertContext; pub use asset::AssetContext; pub use calendar::CalendarContext; -pub use config::{Config, Language, PushCandlestickMode}; +pub use config::{Config, Language, PushCandlestickMode, QuoteTransport}; pub use content::ContentContext; pub use dca::DCAContext; pub use error::{Error, Result, SimpleError, SimpleErrorKind}; diff --git a/rust/src/quote/cmd_code.rs b/rust/src/quote/cmd_code.rs index 62563b5ff..e714d5962 100644 --- a/rust/src/quote/cmd_code.rs +++ b/rust/src/quote/cmd_code.rs @@ -75,3 +75,30 @@ pub(crate) const PUSH_REALTIME_BROKERS: u8 = 103; /// Push Real-time Trades pub(crate) const PUSH_REALTIME_TRADES: u8 = 104; + +/// Returns the REST path (`POST`) that serves the same request/response +/// messages as the WebSocket command `command_code`, if there is one. +pub(crate) fn http_path(command_code: u8) -> Option<&'static str> { + Some(match command_code { + GET_TRADING_SESSION => "/quote/markets/trading-sessions", + GET_TRADING_DAYS => "/quote/markets/trading-days", + GET_BASIC_INFO => "/quote/static-info", + GET_REALTIME_QUOTE => "/quote/quotes", + GET_REALTIME_OPTION_QUOTE => "/quote/options/quotes", + GET_REALTIME_WARRANT_QUOTE => "/quote/warrants/quotes", + GET_SECURITY_DEPTH => "/quote/depth", + GET_SECURITY_BROKERS => "/quote/brokers", + GET_BROKER_IDS => "/quote/participants", + GET_SECURITY_TRADES => "/quote/trades", + GET_SECURITY_INTRADAY => "/quote/intraday", + GET_SECURITY_CANDLESTICKS => "/quote/candlesticks", + GET_OPTION_CHAIN_EXPIRY_DATE_LIST => "/quote/options/expiry-dates", + GET_WARRANT_ISSUER_IDS => "/quote/warrants/issuers", + GET_FILTERED_WARRANT => "/quote/warrants", + GET_CAPITAL_FLOW_INTRADAY => "/quote/capital-flow", + GET_SECURITY_CAPITAL_DISTRIBUTION => "/quote/capital-distribution", + GET_CALC_INDEXES => "/quote/calc-indexes", + GET_SECURITY_HISTORY_CANDLESTICKS => "/quote/history-candlesticks", + _ => return None, + }) +} diff --git a/rust/src/quote/context.rs b/rust/src/quote/context.rs index d95843378..685a4cdc9 100644 --- a/rust/src/quote/context.rs +++ b/rust/src/quote/context.rs @@ -8,13 +8,13 @@ use longbridge_httpcli::{DcRegion, HttpClient, Json, Method}; use longbridge_proto::quote; use longbridge_wscli::WsClientError; use rust_decimal::Decimal; -use serde::{Deserialize, Serialize}; +use serde::{Deserialize, Serialize, de::DeserializeOwned}; use time::{Date, PrimitiveDateTime}; use tokio::sync::{mpsc, oneshot}; use tracing::{Subscriber, dispatcher, instrument::WithSubscriber}; use crate::{ - Config, Error, Language, Market, Result, + Config, Error, Language, Market, QuoteTransport, Result, quote::{ AdjustType, CalcIndex, Candlestick, CapitalDistributionResponse, CapitalFlowLine, FilingItem, HistoryMarketTemperatureResponse, IntradayLine, IssuerInfo, MarketTemperature, @@ -29,6 +29,7 @@ use crate::{ cache::{Cache, CacheWithKey}, cmd_code, core::{Command, Core, UserProfile}, + http_json, overnight, sub_flags::SubFlags, types::{ FilterWarrantExpiryDate, FilterWarrantInOutBoundsType, PinnedMode, @@ -42,6 +43,47 @@ use crate::{ const RETRY_COUNT: usize = 3; const PARTICIPANT_INFO_CACHE_TIMEOUT: Duration = Duration::from_secs(30 * 60); +fn history_offset_request( + symbol: String, + period: Period, + adjust_type: AdjustType, + forward: bool, + time: Option, + count: usize, + trade_sessions: TradeSessions, +) -> quote::SecurityHistoryCandlestickRequest { + quote::SecurityHistoryCandlestickRequest { + symbol, + period: period.into(), + adjust_type: adjust_type.into(), + query_type: quote::HistoryCandlestickQueryType::QueryByOffset.into(), + offset_request: Some(quote::security_history_candlestick_request::OffsetQuery { + direction: if forward { + quote::Direction::Forward + } else { + quote::Direction::Backward + } + .into(), + date: time + .map(|time| { + format!( + "{:04}{:02}{:02}", + time.year(), + time.month() as u8, + time.day() + ) + }) + .unwrap_or_default(), + minute: time + .map(|time| format!("{:02}{:02}", time.hour(), time.minute())) + .unwrap_or_default(), + count: count as i32, + }), + date_request: None, + trade_session: trade_sessions as i32, + } +} + /// Convert a Unix-seconds string (or integer string) to an RFC 3339 timestamp. /// If parsing fails, the original string is returned unchanged. fn unix_secs_to_rfc3339(s: &str) -> String { @@ -60,6 +102,8 @@ const TRADING_SESSION_CACHE_TIMEOUT: Duration = Duration::from_secs(60 * 60 * 2) struct InnerQuoteContext { language: Language, + quote_transport: QuoteTransport, + enable_overnight: bool, http_cli: HttpClient, command_tx: mpsc::UnboundedSender, /// Kept alive only so the background `Core::run` task can observe the @@ -101,6 +145,8 @@ impl QuoteContext { }); let language = config.language; + let quote_transport = config.quote_transport; + let enable_overnight = config.enable_overnight.unwrap_or_default(); let http_cli = config.create_http_client(); let (command_tx, command_rx) = mpsc::unbounded_channel(); let (push_tx, push_rx) = mpsc::unbounded_channel(); @@ -119,6 +165,8 @@ impl QuoteContext { ( QuoteContext(Arc::new(InnerQuoteContext { language, + quote_transport, + enable_overnight, http_cli, command_tx, _shutdown_tx: shutdown_tx, @@ -219,11 +267,20 @@ impl QuoteContext { } /// Send a request `T` to get a response `R` + /// + /// With [`QuoteTransport::Http`], a command that has a REST equivalent is + /// sent as `POST /quote/*` with the same message as a JSON body, and the + /// JSON response is decoded into the same message type, so callers are + /// transport-agnostic. async fn request(&self, command_code: u8, req: T) -> Result where - T: prost::Message, - R: prost::Message + Default, + T: prost::Message + Serialize, + R: prost::Message + Default + DeserializeOwned, { + if let Some(path) = self.http_path(command_code) { + return self.http_request(path, req).await; + } + let resp = self.request_raw(command_code, req.encode_to_vec()).await?; Ok(R::decode(&*resp)?) } @@ -231,12 +288,148 @@ impl QuoteContext { /// Send a request to get a response `R` async fn request_without_body(&self, command_code: u8) -> Result where - R: prost::Message + Default, + R: prost::Message + Default + DeserializeOwned, { + if let Some(path) = self.http_path(command_code) { + return self + .http_request(path, serde_json::Value::Object(Default::default())) + .await; + } + let resp = self.request_raw(command_code, vec![]).await?; Ok(R::decode(&*resp)?) } + /// Whether this context must drop US overnight data to match what the + /// WebSocket returns: see [`overnight`]. + fn aligns_overnight(&self) -> bool { + self.0.quote_transport == QuoteTransport::Http && !self.0.enable_overnight + } + + /// Convert a candlestick page, dropping overnight candlesticks when + /// [`Self::aligns_overnight`], `symbol` is a US equity and the query asked + /// for all sessions (the only case that both drops and tops up). Returns + /// the candlesticks and how many were dropped. + fn convert_candlesticks( + &self, + symbol: &str, + trade_sessions: TradeSessions, + resp: quote::SecurityCandlestickResponse, + ) -> Result<(Vec, usize)> { + if self.aligns_overnight() + && trade_sessions == TradeSessions::All + && overnight::is_us_equity_symbol(symbol) + { + overnight::drop_overnight(resp.candlesticks) + } else { + Ok(( + resp.candlesticks + .into_iter() + .map(TryInto::try_into) + .collect::>()?, + 0, + )) + } + } + + /// Convert a raw candlestick page and, on the HTTP transport, top up a + /// capped US window after overnight bars were dropped so it covers the + /// same bars as the WebSocket (see [`overnight`]). The top-up pages with + /// `POST /quote/history-candlesticks` offset queries; a failure there + /// fails the call. + #[allow(clippy::too_many_arguments)] + async fn align_candlesticks( + &self, + resp: quote::SecurityCandlestickResponse, + symbol: &str, + period: Period, + adjust_type: AdjustType, + trade_sessions: TradeSessions, + forward: bool, + window: overnight::CandlestickWindow, + ) -> Result> { + let target = overnight::window_target(window, &resp.candlesticks); + let cursor = overnight::raw_edge(&resp.candlesticks, forward); + let (candlesticks, dropped) = self.convert_candlesticks(symbol, trade_sessions, resp)?; + let (Some(target), true) = (target, dropped > 0 && trade_sessions == TradeSessions::All) + else { + return Ok(candlesticks); + }; + let bound = match window { + overnight::CandlestickWindow::DateRange { start } => start, + overnight::CandlestickWindow::Count(_) => None, + }; + overnight::refill( + forward, + candlesticks, + cursor, + target, + bound, + |anchor, count| { + let request = history_offset_request( + symbol.to_string(), + period, + adjust_type, + forward, + Some(anchor), + count, + trade_sessions, + ); + async move { + let resp: quote::SecurityCandlestickResponse = self + .request(cmd_code::GET_SECURITY_HISTORY_CANDLESTICKS, request) + .await?; + Ok(resp.candlesticks) + } + }, + ) + .await + .inspect(|candlesticks| { + if candlesticks.len() < target { + dispatcher::with_default(&self.0.log_subscriber.clone().into(), || { + tracing::warn!( + symbol, + target, + got = candlesticks.len(), + "overnight top-up returned fewer candlesticks than the WebSocket would" + ); + }); + } + }) + } + + /// Returns the REST path to use for `command_code`, if the HTTP transport + /// is selected and the command has a REST equivalent. + fn http_path(&self, command_code: u8) -> Option<&'static str> { + match self.0.quote_transport { + QuoteTransport::Http => cmd_code::http_path(command_code), + QuoteTransport::WebSocket => None, + } + } + + /// `POST` `body` to `path` and decode the proto-JSON response into `R`. + /// + /// Absent sub-messages (`Option::None`) are omitted from the body instead + /// of being sent as `null`, matching how a proto-JSON client would encode + /// them. + async fn http_request(&self, path: &'static str, body: B) -> Result + where + B: Serialize, + R: DeserializeOwned, + { + let mut body = serde_json::to_value(body)?; + http_json::strip_nulls(&mut body); + let Json(value) = self + .0 + .http_cli + .request(Method::POST, path) + .body(Json(body)) + .response::>() + .send() + .await?; + Ok(http_json::from_value(value)?) + } + /// Subscribe /// /// Reference: @@ -520,7 +713,16 @@ impl QuoteContext { }, ) .await?; - resp.secu_quote.into_iter().map(TryInto::try_into).collect() + resp.secu_quote + .into_iter() + .map(|mut quote| { + // See `overnight`: HTTP always returns overnight quotes. + if self.aligns_overnight() && overnight::is_us_equity_symbol("e.symbol) { + quote.over_night_quote = None; + } + quote.try_into() + }) + .collect() } /// Get quote of option securities @@ -676,7 +878,7 @@ impl QuoteContext { let current = self.0.http_cli.dc_region().await; if !current.allows(longbridge_httpcli::DcRegion::Ap) { return Err(longbridge_httpcli::HttpClientError::DcRegionRestricted { - path: "quote/brokers (WebSocket)".to_string(), + path: "quote/brokers".to_string(), required: longbridge_httpcli::DcRegion::Ap, current, } @@ -811,11 +1013,13 @@ impl QuoteContext { symbol: impl Into, trade_sessions: TradeSessions, ) -> Result> { + let symbol = symbol.into(); + let drop_overnight = self.aligns_overnight() && overnight::is_us_equity_symbol(&symbol); let resp: quote::SecurityIntradayResponse = self .request( cmd_code::GET_SECURITY_INTRADAY, quote::SecurityIntradayRequest { - symbol: symbol.into(), + symbol, trade_session: trade_sessions as i32, }, ) @@ -823,7 +1027,12 @@ impl QuoteContext { let lines = resp .lines .into_iter() - .map(TryInto::try_into) + .map(TryInto::::try_into) + .filter(|line| { + // See `overnight`: HTTP always returns the overnight session. + !drop_overnight + || !matches!(line, Ok(line) if overnight::is_us_overnight(line.timestamp)) + }) .collect::>>()?; Ok(lines) } @@ -871,11 +1080,12 @@ impl QuoteContext { adjust_type: AdjustType, trade_sessions: TradeSessions, ) -> Result> { + let symbol = symbol.into(); let resp: quote::SecurityCandlestickResponse = self .request( cmd_code::GET_SECURITY_CANDLESTICKS, quote::SecurityCandlestickRequest { - symbol: symbol.into(), + symbol: symbol.clone(), period: period.into(), count: count as i32, adjust_type: adjust_type.into(), @@ -883,12 +1093,16 @@ impl QuoteContext { }, ) .await?; - let candlesticks = resp - .candlesticks - .into_iter() - .map(TryInto::try_into) - .collect::>>()?; - Ok(candlesticks) + self.align_candlesticks( + resp, + &symbol, + period, + adjust_type, + trade_sessions, + false, + overnight::CandlestickWindow::Count(count), + ) + .await } /// Get security history candlesticks by offset @@ -903,49 +1117,31 @@ impl QuoteContext { count: usize, trade_sessions: TradeSessions, ) -> Result> { + let symbol = symbol.into(); let resp: quote::SecurityCandlestickResponse = self .request( cmd_code::GET_SECURITY_HISTORY_CANDLESTICKS, - quote::SecurityHistoryCandlestickRequest { - symbol: symbol.into(), - period: period.into(), - adjust_type: adjust_type.into(), - query_type: quote::HistoryCandlestickQueryType::QueryByOffset.into(), - offset_request: Some( - quote::security_history_candlestick_request::OffsetQuery { - direction: if forward { - quote::Direction::Forward - } else { - quote::Direction::Backward - } - .into(), - date: time - .map(|time| { - format!( - "{:04}{:02}{:02}", - time.year(), - time.month() as u8, - time.day() - ) - }) - .unwrap_or_default(), - minute: time - .map(|time| format!("{:02}{:02}", time.hour(), time.minute())) - .unwrap_or_default(), - count: count as i32, - }, - ), - date_request: None, - trade_session: trade_sessions as i32, - }, + history_offset_request( + symbol.clone(), + period, + adjust_type, + forward, + time, + count, + trade_sessions, + ), ) .await?; - let candlesticks = resp - .candlesticks - .into_iter() - .map(TryInto::try_into) - .collect::>>()?; - Ok(candlesticks) + self.align_candlesticks( + resp, + &symbol, + period, + adjust_type, + trade_sessions, + forward, + overnight::CandlestickWindow::Count(count), + ) + .await } /// Get security history candlesticks by date @@ -958,11 +1154,12 @@ impl QuoteContext { end: Option, trade_sessions: TradeSessions, ) -> Result> { + let symbol = symbol.into(); let resp: quote::SecurityCandlestickResponse = self .request( cmd_code::GET_SECURITY_HISTORY_CANDLESTICKS, quote::SecurityHistoryCandlestickRequest { - symbol: symbol.into(), + symbol: symbol.clone(), period: period.into(), adjust_type: adjust_type.into(), query_type: quote::HistoryCandlestickQueryType::QueryByDate.into(), @@ -993,12 +1190,16 @@ impl QuoteContext { }, ) .await?; - let candlesticks = resp - .candlesticks - .into_iter() - .map(TryInto::try_into) - .collect::>>()?; - Ok(candlesticks) + self.align_candlesticks( + resp, + &symbol, + period, + adjust_type, + trade_sessions, + false, + overnight::CandlestickWindow::DateRange { start }, + ) + .await } /// Get option chain expiry date list diff --git a/rust/src/quote/http_json.rs b/rust/src/quote/http_json.rs new file mode 100644 index 000000000..aaf12cb4f --- /dev/null +++ b/rust/src/quote/http_json.rs @@ -0,0 +1,430 @@ +//! Lenient decoding of the gateway's proto-JSON responses into the `prost` +//! message types shared with the WebSocket transport. +//! +//! The protobuf types derive plain `serde` impls, which expect every field to +//! be present with its native JSON type. The gateway's proto-JSON differs in a +//! few ways that this deserializer absorbs, so the existing +//! `TryFrom` conversions can be reused unchanged: +//! +//! - 64-bit integers (and occasionally other scalars) are encoded as strings +//! - fields holding their default value may be omitted or `null` +//! - keys may use the lowerCamelCase JSON name instead of the proto name + +use serde::{ + Deserializer, + de::{ + self, DeserializeOwned, DeserializeSeed, IntoDeserializer, MapAccess, SeqAccess, Visitor, + }, + forward_to_deserialize_any, +}; +use serde_json::{Map, Value}; + +type Error = serde_json::Error; + +/// Recursively remove `null` object members, so an absent `Option` +/// is omitted from a request body rather than sent as `null`. +pub(crate) fn strip_nulls(value: &mut Value) { + match value { + Value::Object(map) => { + map.retain(|_, v| !v.is_null()); + map.values_mut().for_each(strip_nulls); + } + Value::Array(arr) => arr.iter_mut().for_each(strip_nulls), + _ => {} + } +} + +/// Deserialize `T` from a gateway proto-JSON value. +pub(crate) fn from_value(value: Value) -> Result { + T::deserialize(Lenient(value)) +} + +/// A `serde_json::Value` deserializer that treats `null` as the target type's +/// default value and accepts numbers / booleans encoded as strings. +struct Lenient(Value); + +fn invalid(value: &Value, exp: &dyn de::Expected) -> Error { + de::Error::invalid_type( + match value { + Value::Null => de::Unexpected::Unit, + Value::Bool(b) => de::Unexpected::Bool(*b), + Value::Number(_) => de::Unexpected::Other("number"), + Value::String(s) => de::Unexpected::Str(s), + Value::Array(_) => de::Unexpected::Seq, + Value::Object(_) => de::Unexpected::Map, + }, + exp, + ) +} + +macro_rules! deserialize_int { + ($($method:ident => $visit:ident: $ty:ty),* $(,)?) => { + $( + fn $method>(self, visitor: V) -> Result { + let n: $ty = match &self.0 { + Value::Null => 0 as $ty, + Value::Number(n) => n + .to_string() + .parse() + .map_err(|_| invalid(&self.0, &visitor))?, + Value::String(s) if s.is_empty() => 0 as $ty, + Value::String(s) => s.parse().map_err(|_| invalid(&self.0, &visitor))?, + Value::Bool(b) => *b as u8 as $ty, + _ => return Err(invalid(&self.0, &visitor)), + }; + visitor.$visit(n) + } + )* + }; +} + +impl<'de> Deserializer<'de> for Lenient { + type Error = Error; + + deserialize_int! { + deserialize_i8 => visit_i8: i8, + deserialize_i16 => visit_i16: i16, + deserialize_i32 => visit_i32: i32, + deserialize_i64 => visit_i64: i64, + deserialize_u8 => visit_u8: u8, + deserialize_u16 => visit_u16: u16, + deserialize_u32 => visit_u32: u32, + deserialize_u64 => visit_u64: u64, + deserialize_f32 => visit_f32: f32, + deserialize_f64 => visit_f64: f64, + } + + fn deserialize_any>(self, visitor: V) -> Result { + match self.0 { + Value::Array(arr) => visitor.visit_seq(Seq(arr.into_iter())), + Value::Object(map) => visitor.visit_map(Fields::new(map, &[])), + other => other.deserialize_any(visitor), + } + } + + fn deserialize_bool>(self, visitor: V) -> Result { + match &self.0 { + Value::Null => visitor.visit_bool(false), + Value::Bool(b) => visitor.visit_bool(*b), + Value::String(s) if s == "true" || s == "false" => visitor.visit_bool(s == "true"), + Value::Number(n) => visitor.visit_bool(n.as_f64() != Some(0.0)), + _ => Err(invalid(&self.0, &visitor)), + } + } + + fn deserialize_str>(self, visitor: V) -> Result { + self.deserialize_string(visitor) + } + + fn deserialize_string>(self, visitor: V) -> Result { + match self.0 { + Value::Null => visitor.visit_string(String::new()), + Value::String(s) => visitor.visit_string(s), + Value::Number(n) => visitor.visit_string(n.to_string()), + Value::Bool(b) => visitor.visit_string(b.to_string()), + other => Err(invalid(&other, &visitor)), + } + } + + fn deserialize_option>(self, visitor: V) -> Result { + match self.0 { + Value::Null => visitor.visit_none(), + other => visitor.visit_some(Lenient(other)), + } + } + + fn deserialize_newtype_struct>( + self, + _name: &'static str, + visitor: V, + ) -> Result { + visitor.visit_newtype_struct(self) + } + + fn deserialize_seq>(self, visitor: V) -> Result { + match self.0 { + Value::Null => visitor.visit_seq(Seq(Vec::new().into_iter())), + Value::Array(arr) => visitor.visit_seq(Seq(arr.into_iter())), + other => Err(invalid(&other, &visitor)), + } + } + + fn deserialize_tuple>( + self, + _len: usize, + visitor: V, + ) -> Result { + self.deserialize_seq(visitor) + } + + fn deserialize_tuple_struct>( + self, + _name: &'static str, + _len: usize, + visitor: V, + ) -> Result { + self.deserialize_seq(visitor) + } + + fn deserialize_map>(self, visitor: V) -> Result { + match self.0 { + Value::Null => visitor.visit_map(Fields::new(Map::new(), &[])), + Value::Object(map) => visitor.visit_map(Fields::new(map, &[])), + other => Err(invalid(&other, &visitor)), + } + } + + fn deserialize_struct>( + self, + _name: &'static str, + fields: &'static [&'static str], + visitor: V, + ) -> Result { + match self.0 { + Value::Null => visitor.visit_map(Fields::new(Map::new(), fields)), + Value::Object(map) => visitor.visit_map(Fields::new(map, fields)), + other => Err(invalid(&other, &visitor)), + } + } + + fn deserialize_enum>( + self, + name: &'static str, + variants: &'static [&'static str], + visitor: V, + ) -> Result { + self.0.deserialize_enum(name, variants, visitor) + } + + fn deserialize_ignored_any>(self, visitor: V) -> Result { + visitor.visit_unit() + } + + forward_to_deserialize_any! { + char bytes byte_buf unit unit_struct identifier + } +} + +struct Seq(std::vec::IntoIter); + +impl<'de> SeqAccess<'de> for Seq { + type Error = Error; + + fn next_element_seed>( + &mut self, + seed: T, + ) -> Result, Error> { + self.0 + .next() + .map(|v| seed.deserialize(Lenient(v))) + .transpose() + } +} + +/// Map access over an object's entries. For a struct, keys are normalized to +/// the expected field names and every expected field that is missing is +/// yielded as `null`, so it decodes to its default value. +struct Fields { + entries: std::vec::IntoIter<(String, Value)>, + value: Option, +} + +impl Fields { + fn new(map: Map, fields: &'static [&'static str]) -> Self { + let mut entries: Vec<(String, Value)> = Vec::with_capacity(map.len() + fields.len()); + // camelCase aliases are applied after all exact names, so a payload + // carrying both spellings never yields a duplicate: the proto name + // wins. + let mut aliases: Vec<(String, Value)> = Vec::new(); + for (key, value) in map { + if fields.is_empty() || fields.contains(&key.as_str()) { + entries.push((key, value)); + continue; + } + let snake = camel_to_snake(&key); + if fields.contains(&snake.as_str()) { + aliases.push((snake, value)); + } else { + entries.push((key, value)); + } + } + for (key, value) in aliases { + if !entries.iter().any(|(k, _)| *k == key) { + entries.push((key, value)); + } + } + for field in fields { + if !entries.iter().any(|(key, _)| key == field) { + entries.push(((*field).to_string(), Value::Null)); + } + } + Self { + entries: entries.into_iter(), + value: None, + } + } +} + +impl<'de> MapAccess<'de> for Fields { + type Error = Error; + + fn next_key_seed>( + &mut self, + seed: K, + ) -> Result, Error> { + match self.entries.next() { + Some((key, value)) => { + self.value = Some(value); + seed.deserialize(key.into_deserializer()).map(Some) + } + None => Ok(None), + } + } + + fn next_value_seed>(&mut self, seed: V) -> Result { + seed.deserialize(Lenient(self.value.take().unwrap_or(Value::Null))) + } +} + +fn camel_to_snake(s: &str) -> String { + let mut out = String::with_capacity(s.len() + 4); + for ch in s.chars() { + if ch.is_ascii_uppercase() { + out.push('_'); + out.push(ch.to_ascii_lowercase()); + } else { + out.push(ch); + } + } + out +} + +#[cfg(test)] +mod tests { + use longbridge_proto::quote::{SecurityDepthResponse, SecurityQuoteResponse}; + use serde_json::json; + + use super::*; + + #[test] + fn decode_quote_with_string_int64_and_missing_fields() { + // int64 as string, enum as int, camelCase keys, omitted fields, null + // sub-message, and an extra field unknown to the proto. + let resp: SecurityQuoteResponse = from_value(json!({ + "secuQuote": [{ + "symbol": "700.HK", + "lastDone": "432.200", + "prev_close": "430.000", + "volume": "12345678", + "turnover": "5300000000.000", + "timestamp": "1700000000", + "trade_status": 0, + "volume_str": "12345678", + "pre_market_quote": null, + }] + })) + .unwrap(); + let q = &resp.secu_quote[0]; + assert_eq!(q.symbol, "700.HK"); + assert_eq!(q.last_done, "432.200"); + assert_eq!(q.prev_close, "430.000"); + assert_eq!(q.volume, 12345678); + assert_eq!(q.timestamp, 1700000000); + assert_eq!(q.trade_status, 0); + assert_eq!(q.open, ""); + assert!(q.pre_market_quote.is_none()); + } + + #[test] + fn decode_depth() { + let resp: SecurityDepthResponse = from_value(json!({ + "symbol": "700.HK", + "ask": [{ "position": 1, "price": "432.400", "volume": "1000", "order_num": "3", "volume_str": "1000" }], + "bid": [{ "position": "1", "price": "432.200", "volume": 2000, "orderNum": 5 }], + })) + .unwrap(); + assert_eq!(resp.ask[0].volume, 1000); + assert_eq!(resp.ask[0].order_num, 3); + assert_eq!(resp.bid[0].position, 1); + assert_eq!(resp.bid[0].order_num, 5); + } + + #[test] + fn decode_edge_cases() { + use longbridge_proto::quote::{Brokers, StrikePriceInfo}; + // null inside arrays, mixed string/number ints, bool spellings. + let b: Brokers = from_value(json!({ + "position": "1", "broker_ids": ["1", 2, null, "-3"] + })) + .unwrap(); + assert_eq!(b.broker_ids, vec![1, 2, 0, -3]); + for (v, want) in [ + (json!("true"), true), + (json!(true), true), + (json!(1), true), + (json!(null), false), + ] { + let s: StrikePriceInfo = from_value(json!({ "price": "1", "standard": v })).unwrap(); + assert_eq!(s.standard, want); + } + // camelCase and proto name both present: proto name wins, no duplicate. + for order in [ + json!({ "orderNum": 9, "order_num": "3", "position": 1 }), + json!({ "order_num": "3", "orderNum": 9, "position": 1 }), + ] { + let d: SecurityDepthResponse = from_value(json!({ "ask": [order] })).unwrap(); + assert_eq!(d.ask[0].order_num, 3); + } + // Values that must not be silently accepted. + assert!( + from_value::(json!({ "ask": [{ "volume": "1.5" }] })).is_err() + ); + assert!( + from_value::( + json!({ "ask": [{ "volume": "99999999999999999999" }] }) + ) + .is_err() + ); + // Enum names are not accepted (the gateway sends numbers). + assert!( + from_value::(json!({ + "secu_quote": [{ "trade_status": "Normal" }] + })) + .is_err() + ); + } + + #[test] + fn strip_nulls_on_history_request() { + use longbridge_proto::quote::{ + SecurityHistoryCandlestickRequest, security_history_candlestick_request::OffsetQuery, + }; + let req = SecurityHistoryCandlestickRequest { + symbol: "AAPL.US".into(), + offset_request: Some(OffsetQuery { + count: 10, + ..Default::default() + }), + date_request: None, + ..Default::default() + }; + let mut v = serde_json::to_value(&req).unwrap(); + strip_nulls(&mut v); + assert!(v.get("date_request").is_none(), "{v}"); + assert_eq!(v["offset_request"]["count"], 10); + assert_eq!(v["offset_request"]["date"], ""); + } + + #[test] + fn strip_nulls_omits_absent_messages() { + let mut v = json!({ "a": null, "b": { "c": null, "d": 1 }, "e": [ { "f": null } ] }); + strip_nulls(&mut v); + assert_eq!(v, json!({ "b": { "d": 1 }, "e": [ {} ] })); + } + + #[test] + fn decode_empty_object() { + let resp: SecurityDepthResponse = from_value(json!({})).unwrap(); + assert!(resp.ask.is_empty() && resp.bid.is_empty()); + } +} diff --git a/rust/src/quote/mod.rs b/rust/src/quote/mod.rs index 82cb00627..bcdbb1a8a 100644 --- a/rust/src/quote/mod.rs +++ b/rust/src/quote/mod.rs @@ -4,6 +4,8 @@ mod cache; mod cmd_code; mod context; mod core; +mod http_json; +mod overnight; mod push_types; mod store; mod sub_flags; diff --git a/rust/src/quote/overnight.rs b/rust/src/quote/overnight.rs new file mode 100644 index 000000000..a25cd5e51 --- /dev/null +++ b/rust/src/quote/overnight.rs @@ -0,0 +1,615 @@ +//! Keeps HTTP-transport quote data aligned with the WebSocket's US overnight +//! rule. +//! +//! The quote WebSocket only returns the US overnight session when the +//! connection opted in (`need_over_night_quote`, i.e. +//! `Config::enable_overnight`), and it applies that *before* count / range +//! caps. The REST endpoints always return overnight data and offer no opt-out, +//! so with [`QuoteTransport::Http`](crate::QuoteTransport::Http) and overnight +//! disabled the SDK drops overnight bars client-side and, for capped +//! candlestick windows, pages further so both transports yield the same bars. + +use longbridge_proto::quote; +use time::{Date, OffsetDateTime, PrimitiveDateTime}; +use time_tz::{OffsetDateTimeExt, timezones::db::america::NEW_YORK}; + +use crate::{Result, quote::Candlestick}; + +/// Server-side cap on the `count` of a candlestick request (`count = 5000` is +/// rejected with `301607`; date-range queries observed returning up to 1440 +/// bars, so a page of at least this size is treated as capped). +pub(crate) const MAX_HISTORY_CANDLESTICKS: usize = 1000; + +/// Hard limit on extra requests one call may issue while topping up a window. +pub(crate) const MAX_REFILL_ROUNDS: usize = 8; + +/// `true` for US equities / ETFs (`AAPL.US`), whose overnight session is what +/// the WebSocket gates. Indices (`.VIX.US`) trade on their own schedule and +/// are left alone. +pub(crate) fn is_us_equity_symbol(symbol: &str) -> bool { + symbol.ends_with(".US") && !symbol.starts_with('.') +} + +/// Whether `timestamp` falls in the US overnight session (20:00–04:00 New +/// York). +pub(crate) fn is_us_overnight(timestamp: OffsetDateTime) -> bool { + longbridge_candlesticks::markets::US.trade_session(timestamp) + == Some(longbridge_candlesticks::TRADE_SESSION_OVERNIGHT) +} + +fn is_overnight_candlestick(candlestick: "e::Candlestick) -> bool { + candlestick.trade_session == quote::TradeSession::OvernightTrade as i32 +} + +/// New York wall-clock time of `timestamp`, as used by the offset query. +pub(crate) fn new_york_local(timestamp: OffsetDateTime) -> PrimitiveDateTime { + let local = timestamp.to_timezone(NEW_YORK); + PrimitiveDateTime::new(local.date(), local.time()) +} + +/// New York calendar date of `timestamp`. +pub(crate) fn new_york_date(timestamp: OffsetDateTime) -> Date { + timestamp.to_timezone(NEW_YORK).date() +} + +/// How many bars a candlestick window should hold once overnight bars are +/// dropped, i.e. what the WebSocket (which drops them before capping) returns. +#[derive(Debug, Clone, Copy)] +pub(crate) enum CandlestickWindow { + /// A count-limited query (`count` requested). + Count(usize), + /// A date-range query, served newest-first and capped at a fixed size. + DateRange { start: Option }, +} + +/// The number of bars the WebSocket would return for `window` given the raw +/// page, or `None` if the page is not capped (nothing to top up). +pub(crate) fn window_target( + window: CandlestickWindow, + raw: &[quote::Candlestick], +) -> Option { + match window { + // A page shorter than `count` is either all the data there is or the + // server's cap; either way it is the window the WebSocket would return. + CandlestickWindow::Count(count) => Some(count.min(raw.len())), + // A page shorter than the request cap cannot have been cut off. A + // full-size page is capped unless it already reaches the first bar of + // `start` (midnight New York, where the overnight session begins); + // without a `start` it is taken as capped. + CandlestickWindow::DateRange { start } => { + let capped = raw.len() >= MAX_HISTORY_CANDLESTICKS + && match (start, raw_edge(raw, false)) { + (Some(start), Some(earliest)) => new_york_local(earliest) > start.midnight(), + (None, Some(_)) => true, + _ => false, + }; + capped.then_some(raw.len()) + } + } +} + +/// Earliest (`forward == false`) or latest (`forward == true`) timestamp in a +/// raw page. +pub(crate) fn raw_edge( + candlesticks: &[quote::Candlestick], + forward: bool, +) -> Option { + let timestamps = candlesticks.iter().map(|candlestick| candlestick.timestamp); + if forward { + timestamps.max() + } else { + timestamps.min() + } + .and_then(|timestamp| OffsetDateTime::from_unix_timestamp(timestamp).ok()) +} + +/// Convert a raw page, dropping overnight candlesticks. Returns the kept +/// candlesticks and how many were dropped. +pub(crate) fn drop_overnight( + candlesticks: Vec, +) -> Result<(Vec, usize)> { + let total = candlesticks.len(); + let kept = candlesticks + .into_iter() + .filter(|candlestick| !is_overnight_candlestick(candlestick)) + .map(TryInto::try_into) + .collect::>>()?; + let dropped = total - kept.len(); + Ok((kept, dropped)) +} + +/// Top up a capped candlestick window after overnight bars were dropped. +/// +/// `candlesticks` is the already-filtered first page and `cursor` the edge of +/// its *raw* window (overnight bars included). Pages are fetched beyond the +/// cursor with `fetch(anchor_in_new_york, count)`, where the server treats the +/// anchor as inclusive, until `target` bars are collected, the data runs out, +/// `bound` (a New York date, for date-range queries) is crossed, or +/// [`MAX_REFILL_ROUNDS`] is reached. +pub(crate) async fn refill( + forward: bool, + mut candlesticks: Vec, + mut cursor: Option, + target: usize, + bound: Option, + mut fetch: F, +) -> Result> +where + F: FnMut(PrimitiveDateTime, usize) -> Fut, + Fut: Future>>, +{ + let mut stalled = false; + for _ in 0..MAX_REFILL_ROUNDS { + if candlesticks.len() >= target { + break; + } + let Some(anchor) = cursor else { + break; + }; + let need = target - candlesticks.len(); + // Over-fetch a little for the anchor bar (returned and discarded) and + // scattered overnight bars. The overnight session is one contiguous + // 20:00–04:00 block, so once a page yields nothing the cursor is + // inside it: jump with a full page rather than inching through it. + let count = if stalled { + MAX_HISTORY_CANDLESTICKS + } else { + (need + need / 2 + 1).min(MAX_HISTORY_CANDLESTICKS) + }; + let raw = fetch(new_york_local(anchor), count).await?; + let exhausted = raw.len() < count; + let next_cursor = raw_edge(&raw, forward); + + let (mut more, _) = drop_overnight(raw)?; + more.retain(|candlestick| { + if forward { + candlestick.timestamp > anchor + } else { + candlestick.timestamp < anchor + } + }); + let mut crossed_bound = false; + if let Some(bound) = bound { + more.retain(|candlestick| { + let date = new_york_date(candlestick.timestamp); + let inside = if forward { + date <= bound + } else { + date >= bound + }; + crossed_bound |= !inside; + inside + }); + } + + stalled = more.is_empty(); + if forward { + more.truncate(need); + candlesticks.extend(more); + } else { + let skip = more.len().saturating_sub(need); + candlesticks.splice(0..0, more.into_iter().skip(skip)); + } + + if exhausted || crossed_bound || next_cursor.is_none() || next_cursor == cursor { + break; + } + cursor = next_cursor; + } + + Ok(candlesticks) +} + +#[cfg(test)] +mod tests { + use std::{cell::Cell, rc::Rc}; + + use time::macros::datetime; + use time_tz::PrimitiveDateTimeExt; + + use super::*; + + /// One-minute bars for New York 2026-10-09 (EDT) from `start` for `n` + /// minutes, with the trade session derived from the clock. + fn timeline(start: OffsetDateTime, n: usize) -> Vec { + (0..n as i64) + .map(|i| { + let timestamp = start + time::Duration::minutes(i); + let local = timestamp.to_timezone(NEW_YORK).time(); + let session = match local.hour() { + 0..=3 | 20..=23 => quote::TradeSession::OvernightTrade, + 4..=8 => quote::TradeSession::PreTrade, + 9 if local.minute() < 30 => quote::TradeSession::PreTrade, + 9..=15 => quote::TradeSession::NormalTrade, + _ => quote::TradeSession::PostTrade, + }; + quote::Candlestick { + close: "1".into(), + open: "1".into(), + low: "1".into(), + high: "1".into(), + volume: 1, + turnover: "1".into(), + timestamp: timestamp.unix_timestamp(), + trade_session: session as i32, + } + }) + .collect() + } + + /// A fake server: `count` bars from the anchor (inclusive) in `forward` + /// direction, ascending. + fn server( + data: Vec, + forward: bool, + calls: Rc>, + ) -> impl FnMut(PrimitiveDateTime, usize) -> std::future::Ready>> + { + move |anchor, count| { + calls.set(calls.get() + 1); + let anchor = anchor + .assume_timezone(NEW_YORK) + .unwrap_first() + .unix_timestamp(); + let page: Vec<_> = if forward { + data.iter() + .filter(|c| c.timestamp >= anchor) + .take(count) + .cloned() + .collect() + } else { + let mut page: Vec<_> = data + .iter() + .rev() + .filter(|c| c.timestamp <= anchor) + .take(count) + .cloned() + .collect(); + page.reverse(); + page + }; + std::future::ready(Ok(page)) + } + } + + fn minutes(candlesticks: &[Candlestick]) -> Vec { + candlesticks + .iter() + .map(|c| { + let local = c.timestamp.to_timezone(NEW_YORK); + format!("{:02}:{:02}", local.hour(), local.minute()) + }) + .collect() + } + + fn assert_sorted_unique(candlesticks: &[Candlestick]) { + for pair in candlesticks.windows(2) { + assert!( + pair[0].timestamp < pair[1].timestamp, + "not ascending/unique" + ); + } + } + + const DAY: OffsetDateTime = datetime!(2026-10-09 00:00 -04:00); + + #[test] + fn backward_fills_after_dropping_overnight() { + let data = timeline(DAY, 24 * 60); + // Raw page 19:00–20:39 NY: 60 post + 40 overnight. + let end = 20 * 60 + 40; + let page = data[end - 100..end].to_vec(); + let cursor = raw_edge(&page, false); + let (kept, dropped) = drop_overnight(page).unwrap(); + assert_eq!((kept.len(), dropped), (60, 40)); + + let calls = Rc::new(Cell::new(0)); + let rt = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let out = rt + .block_on(refill( + false, + kept, + cursor, + 100, + None, + server(data, false, calls.clone()), + )) + .unwrap(); + assert_eq!(out.len(), 100); + assert_sorted_unique(&out); + let m = minutes(&out); + assert_eq!( + (m.first().unwrap().as_str(), m.last().unwrap().as_str()), + ("18:20", "19:59") + ); + assert_eq!(calls.get(), 1); + } + + #[test] + fn forward_fills_after_dropping_overnight() { + let data = timeline(DAY, 24 * 60); + // First 100 raw bars from 03:00 NY: 60 overnight + 40 pre. + let start = 3 * 60; + let page = data[start..start + 100].to_vec(); + let cursor = raw_edge(&page, true); + let (kept, dropped) = drop_overnight(page).unwrap(); + assert_eq!((kept.len(), dropped), (40, 60)); + + let calls = Rc::new(Cell::new(0)); + let rt = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let out = rt + .block_on(refill( + true, + kept, + cursor, + 100, + None, + server(data, true, calls.clone()), + )) + .unwrap(); + assert_eq!(out.len(), 100); + assert_sorted_unique(&out); + let m = minutes(&out); + assert_eq!( + (m.first().unwrap().as_str(), m.last().unwrap().as_str()), + ("04:00", "05:39") + ); + assert_eq!(calls.get(), 1); + } + + #[test] + fn page_of_only_overnight_bars_still_makes_progress() { + let data = timeline(DAY, 24 * 60); + // Raw page 21:00–22:39: every bar is overnight. + let start = 21 * 60; + let page = data[start..start + 100].to_vec(); + let cursor = raw_edge(&page, false); + let (kept, dropped) = drop_overnight(page).unwrap(); + assert_eq!((kept.len(), dropped), (0, 100)); + + let calls = Rc::new(Cell::new(0)); + let rt = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let out = rt + .block_on(refill( + false, + kept, + cursor, + 100, + None, + server(data, false, calls.clone()), + )) + .unwrap(); + assert_eq!(out.len(), 100); + assert_sorted_unique(&out); + let m = minutes(&out); + assert_eq!( + (m.first().unwrap().as_str(), m.last().unwrap().as_str()), + ("18:20", "19:59") + ); + assert!(calls.get() <= 3, "took {} rounds", calls.get()); + } + + #[test] + fn stops_at_date_bound_and_when_data_runs_out() { + // Only 2026-10-09 exists; a range starting that day cannot reach back. + let data = timeline(DAY, 24 * 60); + let page = data[..100].to_vec(); // 00:00–01:39, all overnight + let cursor = raw_edge(&page, false); + let (kept, _) = drop_overnight(page).unwrap(); + + let calls = Rc::new(Cell::new(0)); + let rt = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let out = rt + .block_on(refill( + false, + kept, + cursor, + 100, + Some(DAY.date()), + server(data.clone(), false, calls.clone()), + )) + .unwrap(); + assert!(out.is_empty()); + assert_eq!(calls.get(), 1); + + // Without a bound the data simply runs out: terminates, returns what + // exists. + let page = data[..100].to_vec(); + let cursor = raw_edge(&page, false); + let (kept, _) = drop_overnight(page).unwrap(); + let calls = Rc::new(Cell::new(0)); + let out = rt + .block_on(refill( + false, + kept, + cursor, + 100, + None, + server(data, false, calls.clone()), + )) + .unwrap(); + assert!(out.is_empty()); + assert_eq!(calls.get(), 1); + } + + #[test] + fn small_need_inside_overnight_block_backward() { + // Raw page 03:50–05:29: 10 overnight + 90 pre. Only 10 bars are + // needed, but they lie on the far side of the 8-hour overnight block. + let data = timeline(DAY - time::Duration::days(1), 2 * 24 * 60); + let start = 24 * 60 + 3 * 60 + 50; + let page = data[start..start + 100].to_vec(); + let cursor = raw_edge(&page, false); + let (kept, dropped) = drop_overnight(page).unwrap(); + assert_eq!((kept.len(), dropped), (90, 10)); + let calls = Rc::new(Cell::new(0)); + let rt = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let out = rt + .block_on(refill( + false, + kept, + cursor, + 100, + None, + server(data, false, calls.clone()), + )) + .unwrap(); + assert_eq!(out.len(), 100); + assert_sorted_unique(&out); + let m = minutes(&out); + assert_eq!( + (m[0].as_str(), m[9].as_str(), m[10].as_str()), + ("19:50", "19:59", "04:00") + ); + assert!(calls.get() <= 3, "took {} rounds", calls.get()); + } + + #[test] + fn small_need_inside_overnight_block_forward() { + // Raw page 19:40–21:19: 20 post + 80 overnight; the next bars are at + // 04:00 the following day. + let data = timeline(DAY, 2 * 24 * 60); + let start = 19 * 60 + 40; + let page = data[start..start + 100].to_vec(); + let cursor = raw_edge(&page, true); + let (kept, dropped) = drop_overnight(page).unwrap(); + assert_eq!((kept.len(), dropped), (20, 80)); + let calls = Rc::new(Cell::new(0)); + let rt = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let out = rt + .block_on(refill( + true, + kept, + cursor, + 100, + None, + server(data, true, calls.clone()), + )) + .unwrap(); + assert_eq!(out.len(), 100); + assert_sorted_unique(&out); + let m = minutes(&out); + assert_eq!( + (m[19].as_str(), m[20].as_str(), m[99].as_str()), + ("19:59", "04:00", "05:19") + ); + assert!(calls.get() <= 3, "took {} rounds", calls.get()); + } + + #[test] + fn window_targets() { + let data = timeline(DAY, 24 * 60); + assert_eq!( + window_target(CandlestickWindow::Count(1000), &data[..300]), + Some(300) + ); + assert_eq!( + window_target(CandlestickWindow::Count(100), &data[..300]), + Some(100) + ); + let start = Some(DAY.date()); + // Reaches midnight of `start`: not capped. + assert_eq!( + window_target(CandlestickWindow::DateRange { start }, &data), + None + ); + // A short page is never capped, whatever its first bar. + assert_eq!( + window_target(CandlestickWindow::DateRange { start }, &data[20 * 60..]), + None + ); + // A full-size page starting after midnight of `start`: capped. + assert_eq!( + window_target(CandlestickWindow::DateRange { start }, &data[4 * 60..]), + Some(1200) + ); + // Earlier `start` the page did not reach: capped. + let earlier = Some(DAY.date().previous_day().unwrap()); + assert_eq!( + window_target(CandlestickWindow::DateRange { start: earlier }, &data), + Some(1440) + ); + // No start: only a full-size page counts as capped. + assert_eq!( + window_target(CandlestickWindow::DateRange { start: None }, &data[..999]), + None + ); + assert_eq!( + window_target(CandlestickWindow::DateRange { start: None }, &data), + Some(1440) + ); + assert_eq!( + window_target(CandlestickWindow::DateRange { start }, &[]), + None + ); + } + + #[test] + fn round_limit_is_enforced() { + // A pathological server that always returns a full page of overnight + // bars. + let calls = Rc::new(Cell::new(0)); + let c2 = calls.clone(); + let mut t = DAY.unix_timestamp(); + let fetch = move |_anchor: PrimitiveDateTime, count: usize| { + c2.set(c2.get() + 1); + t -= 60 * count as i64; + let page = (0..count as i64) + .map(|i| quote::Candlestick { + close: "1".into(), + open: "1".into(), + low: "1".into(), + high: "1".into(), + volume: 1, + turnover: "1".into(), + timestamp: t + i * 60, + trade_session: quote::TradeSession::OvernightTrade as i32, + }) + .collect(); + std::future::ready(Ok(page)) + }; + let rt = tokio::runtime::Builder::new_current_thread() + .build() + .unwrap(); + let out = rt + .block_on(refill(false, vec![], Some(DAY), 50, None, fetch)) + .unwrap(); + assert!(out.is_empty()); + assert_eq!(calls.get(), MAX_REFILL_ROUNDS); + } + + #[test] + fn us_overnight_window() { + // EDT (2026-10-09) and EST (2026-01-09) boundaries. + for day in [ + datetime!(2026-10-09 00:00 -04:00), + datetime!(2026-01-09 00:00 -05:00), + ] { + let at = |h: i64, m: i64| day + time::Duration::minutes(h * 60 + m); + assert!(is_us_overnight(at(0, 0))); + assert!(is_us_overnight(at(3, 59))); + assert!(!is_us_overnight(at(4, 0))); + assert!(!is_us_overnight(at(9, 30))); + assert!(!is_us_overnight(at(16, 0))); + assert!(!is_us_overnight(at(19, 59))); + assert!(is_us_overnight(at(20, 0))); + assert!(is_us_overnight(at(23, 59))); + } + assert!(is_us_equity_symbol("AAPL.US")); + assert!(is_us_equity_symbol("AAPL251017C340000.US")); + assert!(!is_us_equity_symbol(".VIX.US")); + assert!(!is_us_equity_symbol("700.HK")); + } +}