From 6b319c85426550d4cbbab38cc1cf45465b6dbab2 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=A2=81=E7=AB=A0=E6=B4=AA?= Date: Sat, 10 Oct 2026 17:15:02 +0800 Subject: [PATCH 1/3] feat(quote): add QuoteTransport to serve quote pull APIs over HTTP Add `Config::quote_transport(QuoteTransport::{WebSocket, Http})` (env `LONGBRIDGE_QUOTE_TRANSPORT=ws|http`, default WebSocket). With `Http`, the 19 `QuoteContext` pull APIs that have a REST equivalent are sent as `POST /quote/*` instead of over the quote WebSocket, so a process that only pulls data never opens a WebSocket connection. Method signatures and return types are unchanged: the request is the same prost message serialised as JSON, and the gateway proto-JSON response is decoded back into the same prost message by a lenient deserializer (int64-as-string, omitted/null fields, camelCase keys), so the existing conversions are reused. Results are aligned with the WebSocket, including US overnight data, which the gateway always returns over REST but the WebSocket only returns when `enable_overnight` is set: on the HTTP path the SDK strips `overnight_quote`, filters overnight intraday lines, drops overnight candlesticks and tops up count/range-capped candlestick windows with offset queries (`quote::overnight`, unit-tested). Expose the option in the C, C++, Java, Node.js and Python layers, mirroring `PushCandlestickMode`, and add `examples/rust/quote_http_compare`, which calls every pull API over both transports and diffs the results. --- .gitignore | 3 + CHANGELOG.md | 7 + c/README.md | 1 + c/cbindgen.toml | 1 + c/csrc/include/longbridge.h | 44 +- c/src/config.rs | 34 +- c/src/types/mod.rs | 2 + c/src/types/quote_transport.rs | 17 + cpp/README.md | 1 + cpp/include/config.hpp | 15 +- cpp/include/types.hpp | 11 + cpp/src/config.cpp | 7 + cpp/src/convert.hpp | 13 + examples/rust/Cargo.toml | 1 + examples/rust/account_asset/src/main.rs | 2 +- examples/rust/grid_trading/src/main.rs | 12 +- examples/rust/quote_http_compare/Cargo.toml | 9 + examples/rust/quote_http_compare/src/main.rs | 354 ++++++++++++++ examples/rust/submit_order/src/main.rs | 3 +- .../rust/subscribe_candlesticks/src/main.rs | 2 +- examples/rust/subscribe_quote/src/main.rs | 2 +- examples/rust/today_orders/src/main.rs | 2 +- java/Makefile.toml | 1 + java/README.md | 1 + .../src/main/java/com/longbridge/Config.java | 19 +- .../java/com/longbridge/QuoteTransport.java | 17 + .../main/java/com/longbridge/SdkNative.java | 2 + java/src/config.rs | 16 +- java/src/init.rs | 1 + java/src/types/enum_types.rs | 6 + nodejs/README.md | 1 + nodejs/index.d.ts | 27 +- nodejs/index.js | 1 + nodejs/src/config.rs | 19 +- nodejs/src/types.rs | 12 + python/README.md | 1 + python/pysrc/longbridge/openapi.pyi | 30 +- python/src/config.rs | 24 +- python/src/lib.rs | 1 + python/src/types.rs | 13 + rust/README.md | 1 + rust/src/config.rs | 70 ++- rust/src/lib.rs | 2 +- rust/src/quote/cmd_code.rs | 27 ++ rust/src/quote/context.rs | 321 ++++++++++--- rust/src/quote/http_json.rs | 409 ++++++++++++++++ rust/src/quote/mod.rs | 2 + rust/src/quote/overnight.rs | 453 ++++++++++++++++++ 48 files changed, 1917 insertions(+), 103 deletions(-) create mode 100644 c/src/types/quote_transport.rs create mode 100644 examples/rust/quote_http_compare/Cargo.toml create mode 100644 examples/rust/quote_http_compare/src/main.rs create mode 100644 java/javasrc/src/main/java/com/longbridge/QuoteTransport.java create mode 100644 rust/src/quote/http_json.rs create mode 100644 rust/src/quote/overnight.rs diff --git a/.gitignore b/.gitignore index ca1b0defd8..8e4b6aa542 100644 --- a/.gitignore +++ b/.gitignore @@ -19,3 +19,6 @@ build *.node .venv/ python/.venv/ +.env +.env.* +__pycache__/ diff --git a/CHANGELOG.md b/CHANGELOG.md index 9231b6f25b..d3048f1d2c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -4,6 +4,13 @@ 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 19 `QuoteContext` pull APIs 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 and `realtime_*` still use the WebSocket. Results are aligned across transports, including US overnight data: the WebSocket only returns it when `Config::enable_overnight` is set while HTTP always does, so the SDK applies the same rule on the HTTP path for US equities (`quote().overnight_quote`, `intraday` lines, candlesticks — count/range-capped candlestick windows are topped up with at most 8 extra offset requests so they cover the same bars; `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), and the REST gateway currently reports business errors (e.g. `301600`, `301607`) as `500`. +- **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 3f0e744716..45b8e1cf37 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 24750a1fd1..b72e6b7491 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 9d281ca68c..6c0e1590a3 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 0b7f20c2cb..8d17f24469 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 dccae30c4d..a87dfb5ff3 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 0000000000..7bf18f1caf --- /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 be58e35aff..f83cc13d03 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 928c5b80ec..320d4e8d99 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 f4114ae69d..8959eeffe2 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 f2c40c9cd8..89033eea80 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 b15920a2a2..1d28123c58 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 11af435b82..d97343ed2c 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/account_asset/src/main.rs b/examples/rust/account_asset/src/main.rs index 7280c1b13d..09f62a29f0 100644 --- a/examples/rust/account_asset/src/main.rs +++ b/examples/rust/account_asset/src/main.rs @@ -1,6 +1,6 @@ use std::sync::Arc; -use longbridge::{Config, oauth::OAuthBuilder, trade::TradeContext}; +use longbridge::{oauth::OAuthBuilder, trade::TradeContext, Config}; use tracing_subscriber::EnvFilter; #[tokio::main] diff --git a/examples/rust/grid_trading/src/main.rs b/examples/rust/grid_trading/src/main.rs index dd10e64c3b..3fa2603a97 100644 --- a/examples/rust/grid_trading/src/main.rs +++ b/examples/rust/grid_trading/src/main.rs @@ -1,13 +1,13 @@ use std::sync::Arc; use longbridge::{ - Config, grid::{ GetGridOrderDetailOptions, GetGridOrdersOptions, GetGridTriggerHistoryOptions, GridContext, GridTradeRule, SubmitGridOrderOptions, }, oauth::OAuthBuilder, trade::TradeContext, + Config, }; use rust_decimal::Decimal; use tracing_subscriber::EnvFilter; @@ -65,7 +65,11 @@ async fn main() -> Result<(), Box> { let list = ctx .list(GetGridOrdersOptions::new().symbol("700.HK").limit(20)) .await?; - println!("grid orders: {} (has_more={})", list.grid_order.len(), list.has_more); + println!( + "grid orders: {} (has_more={})", + list.grid_order.len(), + list.has_more + ); // Detail. let detail = ctx @@ -75,7 +79,9 @@ async fn main() -> Result<(), Box> { // Query by IDs. let by_ids = ctx - .list_by_ids(longbridge::grid::GetGridOrdersByIdsOptions::new([&order_id])) + .list_by_ids(longbridge::grid::GetGridOrdersByIdsOptions::new([ + &order_id, + ])) .await?; println!("grid orders by ids: {}", by_ids.len()); diff --git a/examples/rust/quote_http_compare/Cargo.toml b/examples/rust/quote_http_compare/Cargo.toml new file mode 100644 index 0000000000..fab1ba8977 --- /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 0000000000..7c75b674bb --- /dev/null +++ b/examples/rust/quote_http_compare/src/main.rs @@ -0,0 +1,354 @@ +//! 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). 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/examples/rust/submit_order/src/main.rs b/examples/rust/submit_order/src/main.rs index d35edabb9e..6d7659ebb1 100644 --- a/examples/rust/submit_order/src/main.rs +++ b/examples/rust/submit_order/src/main.rs @@ -1,9 +1,10 @@ use std::sync::Arc; use longbridge::{ - Config, decimal, + decimal, oauth::OAuthBuilder, trade::{OrderSide, OrderType, SubmitOrderOptions, TimeInForceType, TradeContext}, + Config, }; use tracing_subscriber::EnvFilter; diff --git a/examples/rust/subscribe_candlesticks/src/main.rs b/examples/rust/subscribe_candlesticks/src/main.rs index 1034c2fa13..ef268ac7fb 100644 --- a/examples/rust/subscribe_candlesticks/src/main.rs +++ b/examples/rust/subscribe_candlesticks/src/main.rs @@ -1,9 +1,9 @@ use std::sync::Arc; use longbridge::{ - Config, oauth::OAuthBuilder, quote::{Period, QuoteContext, TradeSessions}, + Config, }; use tracing_subscriber::EnvFilter; diff --git a/examples/rust/subscribe_quote/src/main.rs b/examples/rust/subscribe_quote/src/main.rs index 254956931d..432ba19aa9 100644 --- a/examples/rust/subscribe_quote/src/main.rs +++ b/examples/rust/subscribe_quote/src/main.rs @@ -1,9 +1,9 @@ use std::sync::Arc; use longbridge::{ - Config, oauth::OAuthBuilder, quote::{QuoteContext, SubFlags}, + Config, }; use tracing_subscriber::EnvFilter; diff --git a/examples/rust/today_orders/src/main.rs b/examples/rust/today_orders/src/main.rs index a37ba1cfdd..2dad48fec0 100644 --- a/examples/rust/today_orders/src/main.rs +++ b/examples/rust/today_orders/src/main.rs @@ -1,6 +1,6 @@ use std::sync::Arc; -use longbridge::{Config, oauth::OAuthBuilder, trade::TradeContext}; +use longbridge::{oauth::OAuthBuilder, trade::TradeContext, Config}; use tracing_subscriber::EnvFilter; #[tokio::main] diff --git a/java/Makefile.toml b/java/Makefile.toml index fa1a566b54..353711476d 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 9ab211e10d..6321b17b79 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 df9a738f06..12273ba208 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 0000000000..466b6f7a9a --- /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 a3c6b7456a..5e296f3866 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 ca4cbfdf62..fe579f142b 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 855aa87c0c..154b7caa35 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 cdc682da25..7089eccdb1 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 198e59387e..0ae213b5e9 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 5ef828a380..9ea50edaa0 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 efec33ac59..b5d3f4297e 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 ff732cad5c..74ab03f35f 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 7114acea7b..157e60b36b 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 df12f8e229..c773f074da 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 aa4b7708a5..463a16ed58 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 cd5c4b35ea..bd4ef7cf98 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 bab5cefc9e..ccbd516971 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 ecb7f51530..540683ca46 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 81b4aba94d..31473feb3b 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 a26db71d70..e09d9445ec 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" | "websocket" => 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 4f99bc0de6..e8f47ad6cd 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 62563b5ff8..e714d59620 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 d95843378c..7759401b70 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,127 @@ 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`]. Returns the candlesticks and how many were + /// dropped. + fn convert_candlesticks( + &self, + resp: quote::SecurityCandlestickResponse, + ) -> Result<(Vec, usize)> { + if self.aligns_overnight() { + overnight::drop_overnight(resp.candlesticks) + } else { + Ok(( + resp.candlesticks + .into_iter() + .map(TryInto::try_into) + .collect::>()?, + 0, + )) + } + } + + /// Top up a capped US candlestick window after overnight bars were dropped, + /// so the HTTP transport covers the same bars as the WebSocket. `resp` is + /// the raw first page, `target` the number of bars the window should hold + /// and `bound` the New York start date of a date-range query. + #[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, + target: usize, + bound: Option, + ) -> Result> { + let cursor = overnight::raw_edge(&resp.candlesticks, forward); + let (candlesticks, dropped) = self.convert_candlesticks(resp)?; + if dropped == 0 + || trade_sessions != TradeSessions::All + || !overnight::is_us_equity_symbol(symbol) + { + return Ok(candlesticks); + } + 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 + } + + /// 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 +692,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() { + quote.over_night_quote = None; + } + quote.try_into() + }) + .collect() } /// Get quote of option securities @@ -676,7 +857,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 +992,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 +1006,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 +1059,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 +1072,20 @@ impl QuoteContext { }, ) .await?; - let candlesticks = resp - .candlesticks - .into_iter() - .map(TryInto::try_into) - .collect::>>()?; - Ok(candlesticks) + // 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. + let target = count.min(resp.candlesticks.len()); + self.align_candlesticks( + resp, + &symbol, + period, + adjust_type, + trade_sessions, + false, + target, + None, + ) + .await } /// Get security history candlesticks by offset @@ -903,49 +1100,33 @@ 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) + let target = count.min(resp.candlesticks.len()); + self.align_candlesticks( + resp, + &symbol, + period, + adjust_type, + trade_sessions, + forward, + target, + None, + ) + .await } /// Get security history candlesticks by date @@ -958,11 +1139,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 +1175,29 @@ impl QuoteContext { }, ) .await?; - let candlesticks = resp - .candlesticks - .into_iter() - .map(TryInto::try_into) - .collect::>>()?; - Ok(candlesticks) + // A date range is served newest-first and capped at a fixed number of + // bars. With a `start`, the window is capped only if it did not reach + // the first bar of that day (midnight New York, where the overnight + // session begins); without one, a page at least as long as the request + // cap is taken as capped. The response size is then the window the + // WebSocket would have filled. + let capped = match (start, overnight::raw_edge(&resp.candlesticks, false)) { + (Some(start), Some(earliest)) => overnight::new_york_local(earliest) > start.midnight(), + (None, Some(_)) => resp.candlesticks.len() >= overnight::MAX_HISTORY_CANDLESTICKS, + _ => false, + }; + let target = if capped { resp.candlesticks.len() } else { 0 }; + self.align_candlesticks( + resp, + &symbol, + period, + adjust_type, + trade_sessions, + false, + target, + 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 0000000000..f5219d5d7d --- /dev/null +++ b/rust/src/quote/http_json.rs @@ -0,0 +1,409 @@ +//! 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_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 82cb006276..bcdbb1a8a9 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 0000000000..1fb568c376 --- /dev/null +++ b/rust/src/quote/overnight.rs @@ -0,0 +1,453 @@ +//! 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. +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() +} + +/// 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>>, +{ + for _ in 0..MAX_REFILL_ROUNDS { + if candlesticks.len() >= target { + break; + } + let Some(anchor) = cursor else { + break; + }; + let need = target - candlesticks.len(); + // Over-fetch: roughly a third of a US day is overnight, and the anchor + // bar itself comes back and is discarded. + let count = (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 + }); + } + + 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 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")); + } +} From 7a24f57fa8d3ffcc1a5400bb4d040ed97b4c8c82 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=A2=81=E7=AB=A0=E6=B4=AA?= Date: Sat, 10 Oct 2026 18:26:32 +0800 Subject: [PATCH 2/3] fix(quote): address review on HTTP transport overnight alignment - refill: once a page yields no usable bar the cursor is inside the contiguous 20:00-04:00 overnight block; fetch a full page to jump it instead of inching through (regression tests for both directions) - gate overnight dropping on the same US-equity predicate everywhere and compute the window target in one place (CandlestickWindow / window_target, unit-tested) - CHANGELOG: split the entry, list the top-up endpoint, error-variant, 429-retry and overnight-session caveats; 20 methods / 19 commands - drop the undocumented "websocket" env alias; ignore only .env/.env.local; revert unrelated example reformatting --- .gitignore | 2 +- CHANGELOG.md | 4 +- examples/rust/account_asset/src/main.rs | 2 +- examples/rust/grid_trading/src/main.rs | 12 +- examples/rust/quote_http_compare/src/main.rs | 4 +- examples/rust/submit_order/src/main.rs | 3 +- .../rust/subscribe_candlesticks/src/main.rs | 2 +- examples/rust/subscribe_quote/src/main.rs | 2 +- examples/rust/today_orders/src/main.rs | 2 +- rust/src/config.rs | 2 +- rust/src/quote/context.rs | 59 +++---- rust/src/quote/http_json.rs | 21 +++ rust/src/quote/overnight.rs | 164 +++++++++++++++++- 13 files changed, 219 insertions(+), 60 deletions(-) diff --git a/.gitignore b/.gitignore index 8e4b6aa542..16b5a7263a 100644 --- a/.gitignore +++ b/.gitignore @@ -20,5 +20,5 @@ build .venv/ python/.venv/ .env -.env.* +.env.local __pycache__/ diff --git a/CHANGELOG.md b/CHANGELOG.md index d3048f1d2c..4e9420db92 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,7 +8,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### 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 19 `QuoteContext` pull APIs 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 and `realtime_*` still use the WebSocket. Results are aligned across transports, including US overnight data: the WebSocket only returns it when `Config::enable_overnight` is set while HTTP always does, so the SDK applies the same rule on the HTTP path for US equities (`quote().overnight_quote`, `intraday` lines, candlesticks — count/range-capped candlestick windows are topped up with at most 8 extra offset requests so they cover the same bars; `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), and the REST gateway currently reports business errors (e.g. `301600`, `301607`) as `500`. +- **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 diff --git a/examples/rust/account_asset/src/main.rs b/examples/rust/account_asset/src/main.rs index 09f62a29f0..7280c1b13d 100644 --- a/examples/rust/account_asset/src/main.rs +++ b/examples/rust/account_asset/src/main.rs @@ -1,6 +1,6 @@ use std::sync::Arc; -use longbridge::{oauth::OAuthBuilder, trade::TradeContext, Config}; +use longbridge::{Config, oauth::OAuthBuilder, trade::TradeContext}; use tracing_subscriber::EnvFilter; #[tokio::main] diff --git a/examples/rust/grid_trading/src/main.rs b/examples/rust/grid_trading/src/main.rs index 3fa2603a97..dd10e64c3b 100644 --- a/examples/rust/grid_trading/src/main.rs +++ b/examples/rust/grid_trading/src/main.rs @@ -1,13 +1,13 @@ use std::sync::Arc; use longbridge::{ + Config, grid::{ GetGridOrderDetailOptions, GetGridOrdersOptions, GetGridTriggerHistoryOptions, GridContext, GridTradeRule, SubmitGridOrderOptions, }, oauth::OAuthBuilder, trade::TradeContext, - Config, }; use rust_decimal::Decimal; use tracing_subscriber::EnvFilter; @@ -65,11 +65,7 @@ async fn main() -> Result<(), Box> { let list = ctx .list(GetGridOrdersOptions::new().symbol("700.HK").limit(20)) .await?; - println!( - "grid orders: {} (has_more={})", - list.grid_order.len(), - list.has_more - ); + println!("grid orders: {} (has_more={})", list.grid_order.len(), list.has_more); // Detail. let detail = ctx @@ -79,9 +75,7 @@ async fn main() -> Result<(), Box> { // Query by IDs. let by_ids = ctx - .list_by_ids(longbridge::grid::GetGridOrdersByIdsOptions::new([ - &order_id, - ])) + .list_by_ids(longbridge::grid::GetGridOrdersByIdsOptions::new([&order_id])) .await?; println!("grid orders by ids: {}", by_ids.len()); diff --git a/examples/rust/quote_http_compare/src/main.rs b/examples/rust/quote_http_compare/src/main.rs index 7c75b674bb..f3a9b4ca11 100644 --- a/examples/rust/quote_http_compare/src/main.rs +++ b/examples/rust/quote_http_compare/src/main.rs @@ -2,7 +2,9 @@ //! //! Reads credentials from `LONGBRIDGE_*` env (or `.env`). Point //! `LONGBRIDGE_HTTP_URL` / `LONGBRIDGE_QUOTE_WS_URL` at an environment that -//! serves the `/quote/*` routes (canary). For each call prints the arguments, +//! 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}; diff --git a/examples/rust/submit_order/src/main.rs b/examples/rust/submit_order/src/main.rs index 6d7659ebb1..d35edabb9e 100644 --- a/examples/rust/submit_order/src/main.rs +++ b/examples/rust/submit_order/src/main.rs @@ -1,10 +1,9 @@ use std::sync::Arc; use longbridge::{ - decimal, + Config, decimal, oauth::OAuthBuilder, trade::{OrderSide, OrderType, SubmitOrderOptions, TimeInForceType, TradeContext}, - Config, }; use tracing_subscriber::EnvFilter; diff --git a/examples/rust/subscribe_candlesticks/src/main.rs b/examples/rust/subscribe_candlesticks/src/main.rs index ef268ac7fb..1034c2fa13 100644 --- a/examples/rust/subscribe_candlesticks/src/main.rs +++ b/examples/rust/subscribe_candlesticks/src/main.rs @@ -1,9 +1,9 @@ use std::sync::Arc; use longbridge::{ + Config, oauth::OAuthBuilder, quote::{Period, QuoteContext, TradeSessions}, - Config, }; use tracing_subscriber::EnvFilter; diff --git a/examples/rust/subscribe_quote/src/main.rs b/examples/rust/subscribe_quote/src/main.rs index 432ba19aa9..254956931d 100644 --- a/examples/rust/subscribe_quote/src/main.rs +++ b/examples/rust/subscribe_quote/src/main.rs @@ -1,9 +1,9 @@ use std::sync::Arc; use longbridge::{ + Config, oauth::OAuthBuilder, quote::{QuoteContext, SubFlags}, - Config, }; use tracing_subscriber::EnvFilter; diff --git a/examples/rust/today_orders/src/main.rs b/examples/rust/today_orders/src/main.rs index 2dad48fec0..a37ba1cfdd 100644 --- a/examples/rust/today_orders/src/main.rs +++ b/examples/rust/today_orders/src/main.rs @@ -1,6 +1,6 @@ use std::sync::Arc; -use longbridge::{oauth::OAuthBuilder, trade::TradeContext, Config}; +use longbridge::{Config, oauth::OAuthBuilder, trade::TradeContext}; use tracing_subscriber::EnvFilter; #[tokio::main] diff --git a/rust/src/config.rs b/rust/src/config.rs index e09d9445ec..7d361ee121 100644 --- a/rust/src/config.rs +++ b/rust/src/config.rs @@ -144,7 +144,7 @@ impl FromStr for QuoteTransport { fn from_str(s: &str) -> ::std::result::Result { match s.to_ascii_lowercase().as_str() { - "ws" | "websocket" => Ok(QuoteTransport::WebSocket), + "ws" => Ok(QuoteTransport::WebSocket), "http" => Ok(QuoteTransport::Http), _ => Err(()), } diff --git a/rust/src/quote/context.rs b/rust/src/quote/context.rs index 7759401b70..7205bc3505 100644 --- a/rust/src/quote/context.rs +++ b/rust/src/quote/context.rs @@ -307,13 +307,14 @@ impl QuoteContext { } /// Convert a candlestick page, dropping overnight candlesticks when - /// [`Self::aligns_overnight`]. Returns the candlesticks and how many were - /// dropped. + /// [`Self::aligns_overnight`] and `symbol` is a US equity. Returns the + /// candlesticks and how many were dropped. fn convert_candlesticks( &self, + symbol: &str, resp: quote::SecurityCandlestickResponse, ) -> Result<(Vec, usize)> { - if self.aligns_overnight() { + if self.aligns_overnight() && overnight::is_us_equity_symbol(symbol) { overnight::drop_overnight(resp.candlesticks) } else { Ok(( @@ -326,10 +327,11 @@ impl QuoteContext { } } - /// Top up a capped US candlestick window after overnight bars were dropped, - /// so the HTTP transport covers the same bars as the WebSocket. `resp` is - /// the raw first page, `target` the number of bars the window should hold - /// and `bound` the New York start date of a date-range query. + /// 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, @@ -339,17 +341,19 @@ impl QuoteContext { adjust_type: AdjustType, trade_sessions: TradeSessions, forward: bool, - target: usize, - bound: Option, + 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(resp)?; - if dropped == 0 - || trade_sessions != TradeSessions::All - || !overnight::is_us_equity_symbol(symbol) - { + let (candlesticks, dropped) = self.convert_candlesticks(symbol, 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, @@ -1072,9 +1076,6 @@ impl QuoteContext { }, ) .await?; - // 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. - let target = count.min(resp.candlesticks.len()); self.align_candlesticks( resp, &symbol, @@ -1082,8 +1083,7 @@ impl QuoteContext { adjust_type, trade_sessions, false, - target, - None, + overnight::CandlestickWindow::Count(count), ) .await } @@ -1115,7 +1115,6 @@ impl QuoteContext { ), ) .await?; - let target = count.min(resp.candlesticks.len()); self.align_candlesticks( resp, &symbol, @@ -1123,8 +1122,7 @@ impl QuoteContext { adjust_type, trade_sessions, forward, - target, - None, + overnight::CandlestickWindow::Count(count), ) .await } @@ -1175,18 +1173,6 @@ impl QuoteContext { }, ) .await?; - // A date range is served newest-first and capped at a fixed number of - // bars. With a `start`, the window is capped only if it did not reach - // the first bar of that day (midnight New York, where the overnight - // session begins); without one, a page at least as long as the request - // cap is taken as capped. The response size is then the window the - // WebSocket would have filled. - let capped = match (start, overnight::raw_edge(&resp.candlesticks, false)) { - (Some(start), Some(earliest)) => overnight::new_york_local(earliest) > start.midnight(), - (None, Some(_)) => resp.candlesticks.len() >= overnight::MAX_HISTORY_CANDLESTICKS, - _ => false, - }; - let target = if capped { resp.candlesticks.len() } else { 0 }; self.align_candlesticks( resp, &symbol, @@ -1194,8 +1180,7 @@ impl QuoteContext { adjust_type, trade_sessions, false, - target, - start, + overnight::CandlestickWindow::DateRange { start }, ) .await } diff --git a/rust/src/quote/http_json.rs b/rust/src/quote/http_json.rs index f5219d5d7d..aaf12cb4f8 100644 --- a/rust/src/quote/http_json.rs +++ b/rust/src/quote/http_json.rs @@ -394,6 +394,27 @@ mod tests { ); } + #[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 } ] }); diff --git a/rust/src/quote/overnight.rs b/rust/src/quote/overnight.rs index 1fb568c376..eb0d2c6b51 100644 --- a/rust/src/quote/overnight.rs +++ b/rust/src/quote/overnight.rs @@ -15,7 +15,9 @@ use time_tz::{OffsetDateTimeExt, timezones::db::america::NEW_YORK}; use crate::{Result, quote::Candlestick}; -/// Server-side cap on the `count` of a candlestick request. +/// 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. @@ -50,6 +52,41 @@ 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())), + // With a `start`, the window is capped only if it did not reach the + // first bar of that day (midnight New York, where the overnight + // session begins); without one, a page at least as long as the request + // cap is taken as capped. + CandlestickWindow::DateRange { start } => { + let capped = match (start, raw_edge(raw, false)) { + (Some(start), Some(earliest)) => new_york_local(earliest) > start.midnight(), + (None, Some(_)) => raw.len() >= MAX_HISTORY_CANDLESTICKS, + _ => false, + }; + capped.then_some(raw.len()) + } + } +} + /// Earliest (`forward == false`) or latest (`forward == true`) timestamp in a /// raw page. pub(crate) fn raw_edge( @@ -100,6 +137,7 @@ where F: FnMut(PrimitiveDateTime, usize) -> Fut, Fut: Future>>, { + let mut stalled = false; for _ in 0..MAX_REFILL_ROUNDS { if candlesticks.len() >= target { break; @@ -108,9 +146,15 @@ where break; }; let need = target - candlesticks.len(); - // Over-fetch: roughly a third of a US day is overnight, and the anchor - // bar itself comes back and is discarded. - let count = (need + need / 2 + 1).min(MAX_HISTORY_CANDLESTICKS); + // 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); @@ -137,6 +181,7 @@ where }); } + stalled = more.is_empty(); if forward { more.truncate(need); candlesticks.extend(more); @@ -394,6 +439,117 @@ mod tests { 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 + ); + // Earliest bar later than midnight of `start`: capped. + assert_eq!( + window_target(CandlestickWindow::DateRange { start }, &data[20 * 60..]), + Some(4 * 60) + ); + // 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 From c8f0c73fa77ba08571b7388c9b0f48d22c9f5cbd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E8=A2=81=E7=AB=A0=E6=B4=AA?= Date: Sat, 10 Oct 2026 18:35:07 +0800 Subject: [PATCH 3/3] fix(quote): tighten HTTP transport overnight alignment - only drop overnight candlesticks for TradeSessions::All queries, so a page is either left untouched or dropped and topped up (matches Go) - a date-range page shorter than the request cap is never treated as capped; avoids a wasted top-up request for weekend/holiday starts - warn when the top-up returns fewer candlesticks than the WebSocket would - clear overnight_quote only for US equities, like intraday/candlesticks --- rust/src/quote/context.rs | 27 ++++++++++++++++++++++----- rust/src/quote/overnight.rs | 28 +++++++++++++++++----------- 2 files changed, 39 insertions(+), 16 deletions(-) diff --git a/rust/src/quote/context.rs b/rust/src/quote/context.rs index 7205bc3505..685a4cdc9c 100644 --- a/rust/src/quote/context.rs +++ b/rust/src/quote/context.rs @@ -307,14 +307,19 @@ impl QuoteContext { } /// Convert a candlestick page, dropping overnight candlesticks when - /// [`Self::aligns_overnight`] and `symbol` is a US equity. Returns the - /// candlesticks and how many were dropped. + /// [`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() && overnight::is_us_equity_symbol(symbol) { + if self.aligns_overnight() + && trade_sessions == TradeSessions::All + && overnight::is_us_equity_symbol(symbol) + { overnight::drop_overnight(resp.candlesticks) } else { Ok(( @@ -345,7 +350,7 @@ impl QuoteContext { ) -> Result> { let target = overnight::window_target(window, &resp.candlesticks); let cursor = overnight::raw_edge(&resp.candlesticks, forward); - let (candlesticks, dropped) = self.convert_candlesticks(symbol, resp)?; + 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); @@ -379,6 +384,18 @@ impl QuoteContext { }, ) .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 @@ -700,7 +717,7 @@ impl QuoteContext { .into_iter() .map(|mut quote| { // See `overnight`: HTTP always returns overnight quotes. - if self.aligns_overnight() { + if self.aligns_overnight() && overnight::is_us_equity_symbol("e.symbol) { quote.over_night_quote = None; } quote.try_into() diff --git a/rust/src/quote/overnight.rs b/rust/src/quote/overnight.rs index eb0d2c6b51..a25cd5e511 100644 --- a/rust/src/quote/overnight.rs +++ b/rust/src/quote/overnight.rs @@ -72,16 +72,17 @@ pub(crate) fn window_target( // 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())), - // With a `start`, the window is capped only if it did not reach the - // first bar of that day (midnight New York, where the overnight - // session begins); without one, a page at least as long as the request - // cap is taken as capped. + // 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 = match (start, raw_edge(raw, false)) { - (Some(start), Some(earliest)) => new_york_local(earliest) > start.midnight(), - (None, Some(_)) => raw.len() >= MAX_HISTORY_CANDLESTICKS, - _ => false, - }; + 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()) } } @@ -524,10 +525,15 @@ mod tests { window_target(CandlestickWindow::DateRange { start }, &data), None ); - // Earliest bar later than midnight of `start`: capped. + // A short page is never capped, whatever its first bar. assert_eq!( window_target(CandlestickWindow::DateRange { start }, &data[20 * 60..]), - Some(4 * 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());