1use std::net::SocketAddr;
4use std::sync::Arc;
5
6use axum::Extension;
7use axum::extract::{ConnectInfo, Json, Query, State};
8use axum::http::StatusCode;
9use bytes::Bytes;
10use futu_codec::header::ProtoFmtType;
11use futu_server::conn::IncomingRequest;
12use prost::Message;
13use serde::Deserialize;
14use serde_json::Value;
15
16use futu_auth::KeyRecord;
17use futu_core::market::{QotMarketId, qot_market_display_prefix};
18use futu_core::proto_id;
19use futu_core::qot_quote_rights::{
20 SYS_QUERY_GET_QUOTE_CAPABILITY, SYS_QUERY_GET_QUOTE_RIGHTS_PROFILE,
21};
22use futu_core::qot_symbol::{format_qot_symbol, parse_qot_symbol_parts};
23use futu_proto::get_delay_statistics;
24use futu_proto::get_global_state;
25use futu_proto::get_user_info;
26use futu_proto::test_cmd;
27use futu_proto::verification;
28use futu_surface_spec::endpoints::get_delay_statistics::default_request_body_json;
29use futu_backend::proto_internal::futu_token_state;
31
32use crate::adapter::{self, RestState};
33
34type ApiResult = Result<Json<Value>, (StatusCode, Json<Value>)>;
35type RawApiResult = Result<adapter::RawJson, (StatusCode, Json<Value>)>;
36
37pub async fn get_global_state(State(state): State<RestState>) -> RawApiResult {
42 adapter::proto_request_raw::<get_global_state::Request, get_global_state::Response>(
43 &state,
44 proto_id::GET_GLOBAL_STATE,
45 None,
46 )
47 .await
48}
49
50pub async fn get_user_info(State(state): State<RestState>) -> RawApiResult {
52 adapter::proto_request_raw::<get_user_info::Request, get_user_info::Response>(
53 &state,
54 proto_id::GET_USER_INFO,
55 None,
56 )
57 .await
58}
59
60pub async fn verification(
62 State(state): State<RestState>,
63 rec: Option<Extension<Arc<KeyRecord>>>,
64 peer: Option<Extension<ConnectInfo<SocketAddr>>>,
65 Json(body): Json<Value>,
66) -> ApiResult {
67 let legacy_local_mode = rec.is_none();
68 let is_loopback = peer
69 .as_ref()
70 .is_some_and(|Extension(ConnectInfo(addr))| addr.ip().is_loopback());
71 let ctx = crate::caller_context::CallerContext::from_key_record(
72 rec.as_deref().map(|record| record.as_ref()),
73 )
74 .with_transport(is_loopback, legacy_local_mode);
75 adapter::proto_request_with_ctx::<verification::Request, verification::Response>(
76 &state,
77 proto_id::VERIFICATION,
78 Some(body),
79 None,
80 Some(&ctx),
81 )
82 .await
83}
84
85#[derive(Debug, Deserialize)]
86#[serde(deny_unknown_fields)]
87pub struct QuoteRightsQuery {
88 refresh: Option<bool>,
89}
90
91#[derive(Debug, Deserialize)]
92#[serde(deny_unknown_fields)]
93pub struct QuoteCapabilityQuery {
94 symbol: Option<String>,
95 market: Option<String>,
96 code: Option<String>,
97}
98
99pub async fn get_quote_rights(
101 State(state): State<RestState>,
102 Query(query): Query<QuoteRightsQuery>,
103) -> ApiResult {
104 if query.refresh.unwrap_or(false) {
105 let req = test_cmd::Request {
106 c2s: test_cmd::C2s {
107 cmd: "request_highest_quote_right".to_string(),
108 param_str: None,
109 param_bytes: None,
110 },
111 };
112 let resp: test_cmd::Response = dispatch_proto(
113 &state,
114 proto_id::TEST_CMD,
115 req,
116 "request_highest_quote_right",
117 )
118 .await?;
119 if resp.ret_type != 0 {
120 return Err(api_error(
121 StatusCode::BAD_GATEWAY,
122 format_sys_command_error_message(
123 "request_highest_quote_right",
124 resp.ret_type,
125 resp.ret_msg.as_deref(),
126 ),
127 ));
128 }
129 }
130
131 let req = test_cmd::Request {
132 c2s: test_cmd::C2s {
133 cmd: SYS_QUERY_GET_QUOTE_RIGHTS_PROFILE.to_string(),
134 param_str: None,
135 param_bytes: None,
136 },
137 };
138 let resp: test_cmd::Response = dispatch_proto(
139 &state,
140 proto_id::TEST_CMD,
141 req,
142 SYS_QUERY_GET_QUOTE_RIGHTS_PROFILE,
143 )
144 .await?;
145 if resp.ret_type != 0 {
146 return Err(api_error(
147 StatusCode::BAD_GATEWAY,
148 format_sys_command_error_message(
149 SYS_QUERY_GET_QUOTE_RIGHTS_PROFILE,
150 resp.ret_type,
151 resp.ret_msg.as_deref(),
152 ),
153 ));
154 }
155 let json = resp.s2c.and_then(|s| s.result_str).ok_or_else(|| {
156 api_error(
157 StatusCode::BAD_GATEWAY,
158 format!("{SYS_QUERY_GET_QUOTE_RIGHTS_PROFILE}: missing result_str"),
159 )
160 })?;
161 serde_json::from_str::<Value>(&json).map(Json).map_err(|e| {
162 api_error(
163 StatusCode::INTERNAL_SERVER_ERROR,
164 format!("parse {SYS_QUERY_GET_QUOTE_RIGHTS_PROFILE}: {e}"),
165 )
166 })
167}
168
169pub async fn get_quote_capability(
171 State(state): State<RestState>,
172 Query(query): Query<QuoteCapabilityQuery>,
173) -> ApiResult {
174 let symbol = quote_capability_symbol(query)?;
175 let req = test_cmd::Request {
176 c2s: test_cmd::C2s {
177 cmd: SYS_QUERY_GET_QUOTE_CAPABILITY.to_string(),
178 param_str: Some(symbol),
179 param_bytes: None,
180 },
181 };
182 let resp: test_cmd::Response = dispatch_proto(
183 &state,
184 proto_id::TEST_CMD,
185 req,
186 SYS_QUERY_GET_QUOTE_CAPABILITY,
187 )
188 .await?;
189 if resp.ret_type != 0 {
190 return Err(api_error(
191 StatusCode::BAD_GATEWAY,
192 format_sys_command_error_message(
193 SYS_QUERY_GET_QUOTE_CAPABILITY,
194 resp.ret_type,
195 resp.ret_msg.as_deref(),
196 ),
197 ));
198 }
199 let json = resp.s2c.and_then(|s| s.result_str).ok_or_else(|| {
200 api_error(
201 StatusCode::BAD_GATEWAY,
202 format!("{SYS_QUERY_GET_QUOTE_CAPABILITY}: missing result_str"),
203 )
204 })?;
205 serde_json::from_str::<Value>(&json).map(Json).map_err(|e| {
206 api_error(
207 StatusCode::INTERNAL_SERVER_ERROR,
208 format!("parse {SYS_QUERY_GET_QUOTE_CAPABILITY}: {e}"),
209 )
210 })
211}
212
213fn quote_capability_symbol(
214 query: QuoteCapabilityQuery,
215) -> Result<String, (StatusCode, Json<Value>)> {
216 let has_symbol = query
217 .symbol
218 .as_ref()
219 .is_some_and(|symbol| !symbol.trim().is_empty());
220 let has_pair = query.market.as_ref().is_some_and(|m| !m.trim().is_empty())
221 || query.code.as_ref().is_some_and(|c| !c.trim().is_empty());
222 if has_symbol && has_pair {
223 return Err(api_error(
224 StatusCode::BAD_REQUEST,
225 "quote-capability accepts either symbol or market+code, not both".to_string(),
226 ));
227 }
228 if let Some(symbol) = query.symbol {
229 let symbol = symbol.trim().to_string();
230 parse_qot_symbol_parts(&symbol).map_err(|e| {
231 api_error(
232 StatusCode::BAD_REQUEST,
233 format!("invalid quote-capability symbol: {e}"),
234 )
235 })?;
236 return Ok(symbol);
237 }
238
239 let market = query
240 .market
241 .as_deref()
242 .map(str::trim)
243 .filter(|v| !v.is_empty());
244 let code = query
245 .code
246 .as_deref()
247 .map(str::trim)
248 .filter(|v| !v.is_empty());
249 let (Some(market), Some(code)) = (market, code) else {
250 return Err(api_error(
251 StatusCode::BAD_REQUEST,
252 "quote-capability requires symbol=MARKET.CODE or market+code".to_string(),
253 ));
254 };
255
256 let symbol = if let Ok(raw_market) = market.parse::<i32>() {
257 if qot_market_display_prefix(QotMarketId::new(raw_market)).is_none() {
258 return Err(api_error(
259 StatusCode::BAD_REQUEST,
260 format!("invalid quote-capability market {raw_market}"),
261 ));
262 }
263 format_qot_symbol(raw_market, code)
264 } else {
265 format!("{}.{}", market.to_ascii_uppercase(), code)
266 };
267 parse_qot_symbol_parts(&symbol).map_err(|e| {
268 api_error(
269 StatusCode::BAD_REQUEST,
270 format!("invalid quote-capability market/code: {e}"),
271 )
272 })?;
273 Ok(symbol)
274}
275
276fn format_sys_command_error_message(label: &str, ret_type: i32, ret_msg: Option<&str>) -> String {
277 let ret_msg = ret_msg
278 .filter(|msg| !msg.is_empty())
279 .unwrap_or("<missing ret_msg>");
280 format!("{label} ret_type={ret_type} msg={ret_msg}")
281}
282
283async fn dispatch_proto<Req, Rsp>(
284 state: &RestState,
285 proto_id: u32,
286 req: Req,
287 label: &str,
288) -> Result<Rsp, (StatusCode, Json<Value>)>
289where
290 Req: Message,
291 Rsp: Message + Default,
292{
293 let incoming = IncomingRequest::builder(
294 state.next_conn_id(),
295 proto_id,
296 state.next_serial(),
297 ProtoFmtType::Protobuf,
298 Bytes::from(req.encode_to_vec()),
299 )
300 .build();
301 let resp_bytes = state
302 .router
303 .dispatch(incoming.conn_id, &incoming)
304 .await
305 .ok_or_else(|| {
306 api_error(
307 StatusCode::INTERNAL_SERVER_ERROR,
308 format!("{label}: handler returned no response"),
309 )
310 })?;
311 Rsp::decode(Bytes::from(resp_bytes)).map_err(|e| {
312 api_error(
313 StatusCode::INTERNAL_SERVER_ERROR,
314 format!("decode {label}: {e}"),
315 )
316 })
317}
318
319fn api_error(status: StatusCode, message: String) -> (StatusCode, Json<Value>) {
320 (
321 status,
322 Json(serde_json::json!({
323 "ret_type": -1,
324 "ret_msg": message,
325 })),
326 )
327}
328
329pub async fn get_delay_statistics(State(state): State<RestState>) -> RawApiResult {
331 let body = default_request_body_json();
332 adapter::proto_request_raw::<get_delay_statistics::Request, get_delay_statistics::Response>(
333 &state,
334 proto_id::GET_DELAY_STATISTICS,
335 Some(body),
336 )
337 .await
338}
339
340pub async fn get_delay_statistics_post(
346 State(state): State<RestState>,
347 Json(body): Json<Value>,
348) -> RawApiResult {
349 adapter::proto_request_raw::<get_delay_statistics::Request, get_delay_statistics::Response>(
350 &state,
351 proto_id::GET_DELAY_STATISTICS,
352 Some(body),
353 )
354 .await
355}
356
357pub async fn ping(State(state): State<RestState>) -> ApiResult {
365 let ok = adapter::proto_request::<get_global_state::Request, get_global_state::Response>(
368 &state,
369 proto_id::GET_GLOBAL_STATE,
370 None,
371 )
372 .await
373 .is_ok();
374
375 Ok(Json(serde_json::json!({
376 "ok": ok,
377 "version": env!("CARGO_PKG_VERSION"),
378 "gateway": "futu-opend-rs",
379 })))
380}
381
382pub async fn push_subscriber_info(
406 State(state): State<RestState>,
407) -> Result<Json<Value>, (StatusCode, Json<Value>)> {
408 if let Some(ref provider) = state.push_health_snapshot_provider {
410 let health = provider();
411 return Ok(Json(serde_json::json!({
412 "ret_type": 0,
413 "ret_msg": "success",
414 "push_health": health,
415 "recommendations": [
416 {
417 "purpose": "查订阅列表 + 全局 quota",
418 "endpoint": "POST /api/query-subscription -d '{}'",
419 "note": "v1.4.83 起默认 all-conn 视图"
420 },
421 {
422 "purpose": "接收 push 数据(quote / tick / orderbook 等)",
423 "endpoint": "WebSocket /ws (支持 Bearer Token 握手)"
424 }
425 ],
426 })));
427 }
428 Err((
432 StatusCode::SERVICE_UNAVAILABLE,
433 Json(serde_json::json!({
434 "ret_type": -1,
435 "ret_msg": "push health snapshot provider not wired (internal setup bug)",
436 "recommendations": [
437 {
438 "purpose": "查订阅列表 + 全局 quota",
439 "endpoint": "POST /api/query-subscription -d '{}'"
440 },
441 {
442 "purpose": "接收 push 数据",
443 "endpoint": "WebSocket /ws"
444 }
445 ],
446 })),
447 ))
448}
449
450pub async fn unsub_acc_push(
466 State(state): State<RestState>,
467 rec: Option<Extension<Arc<futu_auth::KeyRecord>>>,
468 Json(mut body): Json<Value>,
469) -> ApiResult {
470 if rec.is_none() {
474 futu_auth::audit::reject(
475 "rest",
476 "/api/unsub-acc-push",
477 "<legacy>",
478 "unsub-acc-push not supported in legacy mode (no keys.json)",
479 );
480 return Err((
481 axum::http::StatusCode::FORBIDDEN,
482 Json(serde_json::json!({
483 "error": "/api/unsub-acc-push: legacy mode (no keys.json) does not support per-key sub state. \
484 Configure keys.json and pass Bearer token to enable.",
485 "ret_type": -1,
486 "hint": "v1.4.103 B9: legacy mode previously returned silent success without revoking. Now loud-reject to surface the limitation."
487 })),
488 ));
489 }
490 crate::adapter::normalize_json_keys_snake_case(&mut body);
492
493 let acc_ids = match crate::routes::trd::extract_acc_id_list(&body) {
495 Ok(acc_ids) => acc_ids,
496 Err(reason) => {
497 let key_id = rec
498 .as_deref()
499 .map(|r| r.as_ref().id.clone())
500 .unwrap_or_else(|| "<legacy>".to_string());
501 futu_auth::audit::reject("rest", "/api/unsub-acc-push", &key_id, &reason);
502 return Err((
503 axum::http::StatusCode::BAD_REQUEST,
504 Json(serde_json::json!({
505 "error": format!("/api/unsub-acc-push: {reason}")
506 })),
507 ));
508 }
509 };
510 if let Err(reason) = crate::routes::trd::validate_sub_acc_push_acc_ids(&acc_ids) {
511 let key_id = rec
512 .as_deref()
513 .map(|r| r.as_ref().id.clone())
514 .unwrap_or_else(|| "<legacy>".to_string());
515 futu_auth::audit::reject("rest", "/api/unsub-acc-push", &key_id, reason);
516 return Err((
517 axum::http::StatusCode::BAD_REQUEST,
518 Json(serde_json::json!({
519 "error": format!("/api/unsub-acc-push: {reason}")
520 })),
521 ));
522 }
523
524 crate::routes::trd::check_per_acc_rate_for_caller(
529 &state.counters,
530 rec.as_deref().map(|r| r.as_ref()),
531 &acc_ids,
532 "/api/unsub-acc-push",
533 )?;
534
535 let daemon_resp: Json<serde_json::Value> = Json(serde_json::json!({
540 "ret_type": 0,
541 "ret_msg": serde_json::Value::Null,
542 "err_code": serde_json::Value::Null,
543 "s2c": {}
544 }));
545
546 if let Some(rec_ref) = rec.as_deref() {
552 let key_id = rec_ref.as_ref().id.clone();
553 crate::adapter::with_rest_acc_subscriptions_write(&state.rest_acc_subscriptions, |subs| {
554 let entry = subs.entry(key_id).or_default();
557 for &acc_id in &acc_ids {
558 entry.remove(&acc_id);
559 }
560 });
562 }
563
564 Ok(daemon_resp)
565}
566
567#[derive(Debug, Deserialize, Default)]
591#[serde(deny_unknown_fields)]
592pub struct TokenStateQuery {
593 pub app_id: Option<String>,
595}
596
597pub async fn get_token_state(
598 State(state): State<RestState>,
599 Query(q): Query<TokenStateQuery>,
600 body: Option<Json<Value>>,
601) -> RawApiResult {
602 let body_val = match body {
604 Some(Json(mut v)) => {
605 if let Some(qs_app) = q.app_id.as_ref()
607 && let Some(map) = v.as_object_mut()
608 {
609 let c2s_has = map
610 .get("c2s")
611 .and_then(|c| c.as_object())
612 .is_some_and(|c| c.contains_key("app_id") || c.contains_key("appId"));
613 let top_has = map.contains_key("app_id") || map.contains_key("appId");
614 if !c2s_has && !top_has {
615 map.entry("c2s".to_string())
616 .or_insert_with(|| serde_json::json!({}))
617 .as_object_mut()
618 .map(|c| {
619 c.insert(
620 "app_id".to_string(),
621 serde_json::Value::String(qs_app.clone()),
622 )
623 });
624 }
625 }
626 Some(v)
627 }
628 None => {
629 q.app_id
631 .as_ref()
632 .map(|app_id| serde_json::json!({"c2s": {"app_id": app_id}}))
633 }
634 };
635 adapter::proto_request_raw::<
636 futu_token_state::DaemonGetTokenStateReq,
637 futu_token_state::DaemonGetTokenStateRsp,
638 >(&state, proto_id::GET_TOKEN_STATE, body_val)
639 .await
640}
641
642#[cfg(test)]
643mod tests;