1use std::collections::HashMap;
2use std::sync::Arc;
3use std::time::Duration;
4
5use futures_util::{SinkExt, StreamExt};
6use tokio::sync::{broadcast, mpsc, oneshot};
7use tokio::task::JoinHandle;
8use tokio_tungstenite::tungstenite::Message;
9
10use super::*;
11
12pub const ORDER_STREAM_AUTH_DOMAIN: &str = "strata:order-command-stream:v2";
13const COMMAND_TIMEOUT: Duration = Duration::from_secs(10);
14const EVENT_BUFFER: usize = 1_024;
15const MAX_SERVER_EVENTS_PER_FRAME: usize = 64;
16
17#[derive(Clone, Debug)]
18pub struct OrderChallengeResult {
19 pub self_trade_prevention: PlatformSelfTradePrevention,
20 pub prevented_order_ids: Vec<String>,
21 pub effective_request: PlatformOrderChallengeRequest,
22 pub response: PlatformOrderChallengeResponse,
23}
24
25struct ActorRequest {
26 command: PlatformOrderCommand,
27 response: oneshot::Sender<Result<PlatformOrderCommandEvent, SdkError>>,
28}
29
30#[derive(Clone)]
34pub struct OrderCommandStream {
35 market_id: String,
36 owner_wallet: String,
37 session_public_key: String,
38 commands: mpsc::Sender<ActorRequest>,
39 events: broadcast::Sender<PlatformOrderCommandEvent>,
40}
41
42impl std::fmt::Debug for OrderCommandStream {
43 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
44 formatter
45 .debug_struct("OrderCommandStream")
46 .field("market_id", &self.market_id)
47 .field("owner_wallet", &self.owner_wallet)
48 .field("session_public_key", &self.session_public_key)
49 .finish_non_exhaustive()
50 }
51}
52
53impl OrderCommandStream {
54 pub(crate) async fn connect<S: SessionSigner + ?Sized>(
55 client: &StrataClient,
56 market_id: &str,
57 owner_wallet: &str,
58 signer: &S,
59 ) -> Result<Self, SdkError> {
60 let market_id = validate_platform_market_id(market_id)?;
61 let owner_wallet = canonical_public_key(owner_wallet, "owner_wallet")?;
62 let session_public_key = canonical_public_key(signer.public_key(), "session_public_key")?;
63 if owner_wallet == session_public_key {
64 return Err(SdkError::InvalidRequest(
65 "session_public_key must be distinct from owner_wallet".to_owned(),
66 ));
67 }
68 let mut url = client.base_url.clone();
69 let scheme = match url.scheme() {
70 "https" => "wss",
71 "http" => "ws",
72 _ => {
73 return Err(SdkError::InvalidBaseUrl(
74 "order command URL must use http or https".to_owned(),
75 ))
76 }
77 };
78 url.set_scheme(scheme).map_err(|_| {
79 SdkError::InvalidBaseUrl("could not select WebSocket scheme".to_owned())
80 })?;
81 let base_path = url.path().trim_end_matches('/');
82 url.set_path(&format!("{base_path}/v2/markets/{market_id}/orders/stream"));
83 url.set_query(None);
84 url.set_fragment(None);
85
86 let (mut socket, _) = tokio_tungstenite::connect_async(url.as_str())
87 .await
88 .map_err(|error| SdkError::Stream(error.to_string()))?;
89 let auth = tokio::time::timeout(COMMAND_TIMEOUT, socket.next())
90 .await
91 .map_err(|_| SdkError::Stream("authentication challenge timed out".to_owned()))?
92 .ok_or_else(|| SdkError::Stream("socket closed before authentication".to_owned()))?
93 .map_err(|error| SdkError::Stream(error.to_string()))?;
94 let Message::Text(auth) = auth else {
95 return Err(SdkError::Stream(
96 "expected a text authentication challenge".to_owned(),
97 ));
98 };
99 let event: PlatformOrderCommandEvent = serde_json::from_str(&auth)
100 .map_err(|error| SdkError::InvalidResponse(error.to_string()))?;
101 let (challenge, server_time_ms, expires_at_ms) = match event {
102 PlatformOrderCommandEvent::AuthChallenge {
103 schema_version,
104 contract_version,
105 market_id: response_market,
106 challenge,
107 server_time_ms,
108 expires_at_ms,
109 } => {
110 validate_platform_version(schema_version, &contract_version)?;
111 if response_market != market_id
112 || challenge.len() != 64
113 || !challenge.bytes().all(|byte| byte.is_ascii_hexdigit())
114 || challenge.bytes().any(|byte| byte.is_ascii_uppercase())
115 {
116 return Err(SdkError::InvalidResponse(
117 "order stream authentication bindings are invalid".to_owned(),
118 ));
119 }
120 (challenge, server_time_ms, expires_at_ms)
121 }
122 _ => {
123 return Err(SdkError::InvalidResponse(
124 "order stream did not begin with authentication".to_owned(),
125 ))
126 }
127 };
128 if expires_at_ms <= server_time_ms {
129 return Err(SdkError::InvalidResponse(
130 "order stream authentication challenge is expired".to_owned(),
131 ));
132 }
133 let auth_message =
134 order_stream_auth_message(&market_id, &owner_wallet, &session_public_key, &challenge);
135 let signature = signer
136 .sign_message(&auth_message)
137 .await
138 .map_err(SdkError::Signer)?;
139 if signature.len() != 64 {
140 return Err(SdkError::Signer(
141 "stream authentication signature must contain 64 bytes".to_owned(),
142 ));
143 }
144 socket
145 .send(Message::Text(
146 serde_json::to_string(&PlatformOrderCommandClientFrame::Authenticate {
147 owner_wallet: owner_wallet.clone(),
148 session_public_key: session_public_key.clone(),
149 signature: bs58::encode(signature).into_string(),
150 })
151 .map_err(|error| SdkError::InvalidRequest(error.to_string()))?
152 .into(),
153 ))
154 .await
155 .map_err(|error| SdkError::Stream(error.to_string()))?;
156
157 let ready = tokio::time::timeout(COMMAND_TIMEOUT, socket.next())
158 .await
159 .map_err(|_| SdkError::Stream("signed authentication timed out".to_owned()))?
160 .ok_or_else(|| SdkError::Stream("socket closed during authentication".to_owned()))?
161 .map_err(|error| SdkError::Stream(error.to_string()))?;
162 let Message::Text(ready) = ready else {
163 return Err(SdkError::Stream(
164 "expected a text authentication result".to_owned(),
165 ));
166 };
167 let mut ready_events = parse_order_command_events(&ready)?;
168 if ready_events.len() != 1 {
169 return Err(SdkError::InvalidResponse(
170 "order stream authentication returned an invalid event batch".to_owned(),
171 ));
172 }
173 let ready = ready_events.pop().expect("one ready event");
174 let (stream_id, sequence) = match &ready {
175 PlatformOrderCommandEvent::Ready {
176 schema_version,
177 contract_version,
178 market_id: response_market,
179 stream_id,
180 sequence,
181 ..
182 } => {
183 validate_platform_version(*schema_version, contract_version)?;
184 if response_market != &market_id
185 || !valid_handle(stream_id, "order_command_stream_")
186 || sequence != "1"
187 {
188 return Err(SdkError::InvalidResponse(
189 "order stream ready bindings are invalid".to_owned(),
190 ));
191 }
192 (stream_id.clone(), 1u64)
193 }
194 _ => {
195 return Err(SdkError::InvalidResponse(
196 "order stream authentication was not accepted".to_owned(),
197 ))
198 }
199 };
200
201 let (commands, receiver) = mpsc::channel(512);
202 let (events, _) = broadcast::channel(EVENT_BUFFER);
203 let _ = events.send(ready);
204 tokio::spawn(run_actor(
205 socket,
206 receiver,
207 events.clone(),
208 market_id.clone(),
209 stream_id,
210 sequence,
211 ));
212 Ok(Self {
213 market_id,
214 owner_wallet,
215 session_public_key,
216 commands,
217 events,
218 })
219 }
220
221 pub fn market_id(&self) -> &str {
222 &self.market_id
223 }
224
225 pub fn owner_wallet(&self) -> &str {
226 &self.owner_wallet
227 }
228
229 pub fn session_public_key(&self) -> &str {
230 &self.session_public_key
231 }
232
233 pub fn subscribe(&self) -> broadcast::Receiver<PlatformOrderCommandEvent> {
236 self.events.subscribe()
237 }
238
239 pub async fn command(
240 &self,
241 command: PlatformOrderCommand,
242 ) -> Result<PlatformOrderCommandEvent, SdkError> {
243 let (response, receiver) = oneshot::channel();
244 self.commands
245 .send(ActorRequest { command, response })
246 .await
247 .map_err(|_| SdkError::Stream("order command socket is closed".to_owned()))?;
248 tokio::time::timeout(COMMAND_TIMEOUT, receiver)
249 .await
250 .map_err(|_| SdkError::Stream("order command timed out".to_owned()))?
251 .map_err(|_| SdkError::Stream("order command actor stopped".to_owned()))?
252 }
253
254 pub async fn challenge(
255 &self,
256 request: PlatformOrderChallengeRequest,
257 self_trade_prevention: PlatformSelfTradePrevention,
258 ) -> Result<OrderChallengeResult, SdkError> {
259 let request = normalize_order_challenge_request(request)?;
260 self.ensure_request_identity(&request)?;
261 match self
262 .command(PlatformOrderCommand::Challenge {
263 request,
264 self_trade_prevention,
265 })
266 .await?
267 {
268 PlatformOrderCommandEvent::ChallengeResult {
269 self_trade_prevention,
270 prevented_order_ids,
271 effective_request,
272 response,
273 ..
274 } => {
275 self.ensure_request_identity(&effective_request)?;
276 validate_challenge_result(&self.market_id, &effective_request, &response)?;
277 Ok(OrderChallengeResult {
278 self_trade_prevention,
279 prevented_order_ids,
280 effective_request,
281 response,
282 })
283 }
284 _ => Err(SdkError::InvalidResponse(
285 "expected an order challenge result".to_owned(),
286 )),
287 }
288 }
289
290 pub async fn prepare(
291 &self,
292 request: PlatformOrderPrepareRequest,
293 ) -> Result<PlatformOrderPrepareResponse, SdkError> {
294 if !valid_handle(&request.challenge_id, "oc_") {
295 return Err(SdkError::InvalidRequest(
296 "order challenge_id is invalid".to_owned(),
297 ));
298 }
299 let request = PlatformOrderPrepareRequest {
300 challenge_id: request.challenge_id,
301 authorization_signature: canonical_signature(
302 &request.authorization_signature,
303 "authorization_signature",
304 )?,
305 };
306 match self
307 .command(PlatformOrderCommand::Prepare { request })
308 .await?
309 {
310 PlatformOrderCommandEvent::PrepareResult { response, .. } => {
311 validate_prepared(&self.market_id, &response)?;
312 Ok(response)
313 }
314 _ => Err(SdkError::InvalidResponse(
315 "expected an order prepare result".to_owned(),
316 )),
317 }
318 }
319
320 pub async fn submit(
321 &self,
322 request: PlatformOrderSubmitRequest,
323 ) -> Result<PlatformOrderSubmitResponse, SdkError> {
324 let request = normalize_submit_request(request)?;
325 let expected_control = request.order_control_id.clone();
326 match self
327 .command(PlatformOrderCommand::Submit { request })
328 .await?
329 {
330 PlatformOrderCommandEvent::SubmitResult { response, .. } => {
331 validate_submit(&self.market_id, &expected_control, &response)?;
332 Ok(response)
333 }
334 _ => Err(SdkError::InvalidResponse(
335 "expected an order submit result".to_owned(),
336 )),
337 }
338 }
339
340 pub async fn status(
341 &self,
342 request: PlatformOrderStatusRequest,
343 ) -> Result<PlatformOrderStatusResponse, SdkError> {
344 let request = normalize_status_request(request)?;
345 let expected_control = request.order_control_id.clone();
346 match self
347 .command(PlatformOrderCommand::Status { request })
348 .await?
349 {
350 PlatformOrderCommandEvent::StatusResult { response, .. } => {
351 validate_status(&self.market_id, &expected_control, &response)?;
352 Ok(response)
353 }
354 _ => Err(SdkError::InvalidResponse(
355 "expected an order status result".to_owned(),
356 )),
357 }
358 }
359
360 pub async fn execute_order<S, V>(
361 &self,
362 operation: &OrderExecuteOperation,
363 signer: &S,
364 verifier: &V,
365 idempotency_key: Option<&str>,
366 self_trade_prevention: PlatformSelfTradePrevention,
367 ) -> Result<PlatformOrderSubmitResponse, SdkError>
368 where
369 S: SessionSigner + ?Sized,
370 V: OrderVerifier + ?Sized,
371 {
372 self.ensure_signer(signer)?;
373 let challenged = self
374 .challenge(
375 operation.challenge_request(self.session_public_key.clone()),
376 self_trade_prevention,
377 )
378 .await?;
379 let (prepared, signed_transaction_base64) = self
380 .authorize_and_prepare(&challenged, signer, verifier)
381 .await?;
382 self.submit(PlatformOrderSubmitRequest {
383 order_control_id: prepared.order_control_id.clone(),
384 signed_transaction_base64,
385 idempotency_key: normalize_idempotency_key(
386 idempotency_key.unwrap_or(&prepared.order_control_id),
387 )?,
388 })
389 .await
390 }
391
392 pub async fn arm_dead_man<S, V>(
396 &self,
397 timeout: Duration,
398 signer: &S,
399 verifier: &V,
400 idempotency_key: Option<&str>,
401 ) -> Result<DeadManGuard, SdkError>
402 where
403 S: SessionSigner + ?Sized,
404 V: OrderVerifier + ?Sized,
405 {
406 let timeout_ms = checked_dead_man_timeout(timeout)?;
407 let (state, _) = self
408 .arm_dead_man_once(timeout_ms, signer, verifier, idempotency_key)
409 .await?;
410 let stream = self.clone();
411 let task = tokio::spawn(async move {
412 let mut interval =
413 tokio::time::interval(Duration::from_millis((timeout_ms / 3).max(250)));
414 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
415 interval.tick().await;
416 loop {
417 interval.tick().await;
418 match stream.heartbeat_dead_man().await {
419 Ok(state) if state.status == PlatformDeadManStatus::Armed => {}
420 _ => break,
421 }
422 }
423 });
424 Ok(DeadManGuard {
425 initial_state: state,
426 stream: self.clone(),
427 task,
428 active: true,
429 })
430 }
431
432 pub async fn maintain_dead_man<S, V>(
437 &self,
438 timeout: Duration,
439 signer: Arc<S>,
440 verifier: Arc<V>,
441 ) -> Result<DeadManGuard, SdkError>
442 where
443 S: SessionSigner + 'static,
444 V: OrderVerifier + 'static,
445 {
446 let timeout_ms = checked_dead_man_timeout(timeout)?;
447 let (state, mut transaction_expires_at_ms) = self
448 .arm_dead_man_once(timeout_ms, signer.as_ref(), verifier.as_ref(), None)
449 .await?;
450 let stream = self.clone();
451 let maintained_stream = stream.clone();
452 let task = tokio::spawn(async move {
453 let mut interval =
454 tokio::time::interval(Duration::from_millis((timeout_ms / 3).max(250)));
455 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
456 interval.tick().await;
457 let refresh_lead_ms = timeout_ms.saturating_mul(2).max(5_000);
458 loop {
459 interval.tick().await;
460 let now = unix_ms().unwrap_or(u64::MAX);
461 if now.saturating_add(refresh_lead_ms) >= transaction_expires_at_ms {
462 match maintained_stream
463 .arm_dead_man_once(timeout_ms, signer.as_ref(), verifier.as_ref(), None)
464 .await
465 {
466 Ok((next, expires)) if next.status == PlatformDeadManStatus::Armed => {
467 transaction_expires_at_ms = expires;
468 }
469 _ => break,
470 }
471 } else {
472 match maintained_stream.heartbeat_dead_man().await {
473 Ok(next) if next.status == PlatformDeadManStatus::Armed => {}
474 _ => break,
475 }
476 }
477 }
478 });
479 Ok(DeadManGuard {
480 initial_state: state,
481 stream,
482 task,
483 active: true,
484 })
485 }
486
487 pub async fn heartbeat_dead_man(&self) -> Result<PlatformDeadManState, SdkError> {
488 self.dead_man_command(PlatformOrderCommand::DeadManHeartbeat)
489 .await
490 }
491
492 pub async fn dead_man_status(&self) -> Result<PlatformDeadManState, SdkError> {
493 self.dead_man_command(PlatformOrderCommand::DeadManStatus)
494 .await
495 }
496
497 pub async fn disarm_dead_man(&self) -> Result<PlatformDeadManState, SdkError> {
498 self.dead_man_command(PlatformOrderCommand::DeadManDisarm)
499 .await
500 }
501
502 async fn dead_man_command(
503 &self,
504 command: PlatformOrderCommand,
505 ) -> Result<PlatformDeadManState, SdkError> {
506 match self.command(command).await? {
507 PlatformOrderCommandEvent::DeadManResult { state, .. } => Ok(state),
508 _ => Err(SdkError::InvalidResponse(
509 "expected a dead-man result".to_owned(),
510 )),
511 }
512 }
513
514 async fn arm_dead_man_once<S, V>(
515 &self,
516 timeout_ms: u64,
517 signer: &S,
518 verifier: &V,
519 idempotency_key: Option<&str>,
520 ) -> Result<(PlatformDeadManState, u64), SdkError>
521 where
522 S: SessionSigner + ?Sized,
523 V: OrderVerifier + ?Sized,
524 {
525 self.ensure_signer(signer)?;
526 let challenged = self
527 .challenge(
528 PlatformOrderChallengeRequest::CancelAll {
529 owner_wallet: self.owner_wallet.clone(),
530 session_public_key: self.session_public_key.clone(),
531 },
532 PlatformSelfTradePrevention::CancelTaker,
533 )
534 .await?;
535 let (prepared, signed_transaction_base64) = self
536 .authorize_and_prepare(&challenged, signer, verifier)
537 .await?;
538 let transaction_expires_at_ms = prepared.expires_at_ms;
539 let request = PlatformOrderSubmitRequest {
540 order_control_id: prepared.order_control_id.clone(),
541 signed_transaction_base64,
542 idempotency_key: normalize_idempotency_key(
543 idempotency_key.unwrap_or(&prepared.order_control_id),
544 )?,
545 };
546 let state = match self
547 .command(PlatformOrderCommand::DeadManArm {
548 timeout_ms,
549 request,
550 })
551 .await?
552 {
553 PlatformOrderCommandEvent::DeadManResult { state, .. } => state,
554 _ => {
555 return Err(SdkError::InvalidResponse(
556 "expected a dead-man result".to_owned(),
557 ))
558 }
559 };
560 if state.status != PlatformDeadManStatus::Armed {
561 return Err(SdkError::InvalidResponse(
562 "dead-man ticket was not armed".to_owned(),
563 ));
564 }
565 Ok((state, transaction_expires_at_ms))
566 }
567
568 async fn authorize_and_prepare<S, V>(
569 &self,
570 challenged: &OrderChallengeResult,
571 signer: &S,
572 verifier: &V,
573 ) -> Result<(PlatformOrderPrepareResponse, String), SdkError>
574 where
575 S: SessionSigner + ?Sized,
576 V: OrderVerifier + ?Sized,
577 {
578 let authorization =
579 validate_order_authorization(&challenged.response, &challenged.effective_request)?;
580 let signature = signer
581 .sign_message(&authorization.bytes)
582 .await
583 .map_err(SdkError::Signer)?;
584 if signature.len() != 64 {
585 return Err(SdkError::Signer(
586 "order authorization signature must contain 64 bytes".to_owned(),
587 ));
588 }
589 let prepared = self
590 .prepare(PlatformOrderPrepareRequest {
591 challenge_id: challenged.response.challenge_id.clone(),
592 authorization_signature: bs58::encode(signature).into_string(),
593 })
594 .await?;
595 validate_order_prepare_binding(&prepared, &challenged.response, &authorization)?;
596 verifier
597 .verify(&OrderVerificationContext {
598 challenge: &challenged.response,
599 prepared: &prepared,
600 owner_wallet: &self.owner_wallet,
601 session_public_key: &self.session_public_key,
602 })
603 .await
604 .map_err(SdkError::Verification)?;
605 let transaction = signer
606 .sign_transaction(&prepared.transaction_base64)
607 .await
608 .map_err(SdkError::Signer)?;
609 Ok((
610 prepared,
611 canonical_base64(&transaction, "signed_transaction_base64")?,
612 ))
613 }
614
615 fn ensure_request_identity(
616 &self,
617 request: &PlatformOrderChallengeRequest,
618 ) -> Result<(), SdkError> {
619 if order_request_owner(request) != self.owner_wallet
620 || order_request_session(request) != self.session_public_key
621 {
622 return Err(SdkError::InvalidRequest(
623 "order command identity does not match the authenticated socket".to_owned(),
624 ));
625 }
626 Ok(())
627 }
628
629 fn ensure_signer<S: SessionSigner + ?Sized>(&self, signer: &S) -> Result<(), SdkError> {
630 if canonical_public_key(signer.public_key(), "session_public_key")?
631 != self.session_public_key
632 {
633 return Err(SdkError::InvalidRequest(
634 "signer does not match the authenticated order command session".to_owned(),
635 ));
636 }
637 Ok(())
638 }
639}
640
641pub struct DeadManGuard {
644 pub initial_state: PlatformDeadManState,
645 stream: OrderCommandStream,
646 task: JoinHandle<()>,
647 active: bool,
648}
649
650impl DeadManGuard {
651 pub async fn disarm(&mut self) -> Result<PlatformDeadManState, SdkError> {
652 self.task.abort();
653 let state = self.stream.disarm_dead_man().await?;
654 self.active = false;
655 Ok(state)
656 }
657}
658
659impl Drop for DeadManGuard {
660 fn drop(&mut self) {
661 self.task.abort();
662 if self.active {
663 }
665 }
666}
667
668pub fn order_stream_auth_message(
669 market_id: &str,
670 owner_wallet: &str,
671 session_public_key: &str,
672 challenge: &str,
673) -> Vec<u8> {
674 format!(
675 "{ORDER_STREAM_AUTH_DOMAIN}\n{market_id}\n{owner_wallet}\n{session_public_key}\n{challenge}"
676 )
677 .into_bytes()
678}
679
680async fn run_actor<S>(
681 mut socket: tokio_tungstenite::WebSocketStream<S>,
682 mut commands: mpsc::Receiver<ActorRequest>,
683 events: broadcast::Sender<PlatformOrderCommandEvent>,
684 market_id: String,
685 stream_id: String,
686 mut server_sequence: u64,
687) where
688 S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static,
689{
690 let mut client_sequence = 0u64;
691 let mut request_counter = 0u64;
692 let mut pending =
693 HashMap::<String, oneshot::Sender<Result<PlatformOrderCommandEvent, SdkError>>>::new();
694 let failure = 'actor: loop {
695 tokio::select! {
696 request = commands.recv() => {
697 let Some(request) = request else {
698 let _ = socket.close(None).await;
699 break "order command handle closed".to_owned();
700 };
701 client_sequence = client_sequence.saturating_add(1);
702 request_counter = request_counter.saturating_add(1);
703 let request_id = format!("rust-{request_counter:x}");
704 let frame = PlatformOrderCommandClientFrame::Command {
705 request_id: request_id.clone(),
706 sequence: client_sequence.to_string(),
707 command: request.command,
708 };
709 let message = match serde_json::to_string(&frame) {
710 Ok(message) => message,
711 Err(error) => {
712 let _ = request.response.send(Err(SdkError::InvalidRequest(error.to_string())));
713 continue;
714 }
715 };
716 if let Err(error) = socket.send(Message::Text(message.into())).await {
717 let text = error.to_string();
718 let _ = request.response.send(Err(SdkError::Stream(text.clone())));
719 break text;
720 }
721 pending.insert(request_id, request.response);
722 }
723 incoming = socket.next() => {
724 let Some(incoming) = incoming else {
725 break "order command socket closed".to_owned();
726 };
727 let message = match incoming {
728 Ok(Message::Text(message)) => message,
729 Ok(Message::Ping(payload)) => {
730 if let Err(error) = socket.send(Message::Pong(payload)).await {
731 break error.to_string();
732 }
733 continue;
734 }
735 Ok(Message::Pong(_)) => continue,
736 Ok(Message::Close(_)) => break "order command socket closed".to_owned(),
737 Ok(_) => continue,
738 Err(error) => break error.to_string(),
739 };
740 let frame_events = match parse_order_command_events(&message) {
741 Ok(events) => events,
742 Err(error) => break error.to_string(),
743 };
744 for event in frame_events {
745 if let Err(error) = validate_event_sequence(
746 &event,
747 &market_id,
748 &stream_id,
749 &mut server_sequence,
750 ) {
751 break 'actor error.to_string();
752 }
753 let request_id = event_request_id(&event).map(str::to_owned);
754 let _ = events.send(event.clone());
755 if let Some(request_id) = request_id {
756 if let Some(response) = pending.remove(&request_id) {
757 let result = match &event {
758 PlatformOrderCommandEvent::CommandError { error, .. } => {
759 Err(SdkError::Command {
760 code: serde_json::to_value(error.code)
761 .ok()
762 .and_then(|value| value.as_str().map(str::to_owned))
763 .unwrap_or_else(|| "command_rejected".to_owned()),
764 message: error.message.clone(),
765 retryable: error.retryable,
766 })
767 }
768 _ => Ok(event),
769 };
770 let _ = response.send(result);
771 }
772 }
773 }
774 }
775 }
776 };
777 for response in pending.into_values() {
778 let _ = response.send(Err(SdkError::Stream(failure.clone())));
779 }
780}
781
782fn parse_order_command_events(message: &str) -> Result<Vec<PlatformOrderCommandEvent>, SdkError> {
783 let value: serde_json::Value = serde_json::from_str(message).map_err(|error| {
784 SdkError::InvalidResponse(format!("invalid order command frame: {error}"))
785 })?;
786 let events = if value.is_array() {
787 serde_json::from_value::<Vec<PlatformOrderCommandEvent>>(value)
788 } else {
789 serde_json::from_value::<PlatformOrderCommandEvent>(value).map(|event| vec![event])
790 }
791 .map_err(|error| SdkError::InvalidResponse(format!("invalid order command event: {error}")))?;
792 if events.is_empty() || events.len() > MAX_SERVER_EVENTS_PER_FRAME {
793 return Err(SdkError::InvalidResponse(
794 "order command event batch is invalid".to_owned(),
795 ));
796 }
797 Ok(events)
798}
799
800fn validate_event_sequence(
801 event: &PlatformOrderCommandEvent,
802 market_id: &str,
803 stream_id: &str,
804 server_sequence: &mut u64,
805) -> Result<(), SdkError> {
806 let Some((schema, contract, market, stream, sequence, previous)) = event_sequence(event) else {
807 return Err(SdkError::InvalidResponse(
808 "order stream restarted without signed authentication".to_owned(),
809 ));
810 };
811 validate_platform_version(schema, contract)?;
812 let sequence = parse_wire_sequence(sequence)?;
813 let previous = parse_wire_sequence(previous)?;
814 if market != market_id
815 || stream != stream_id
816 || previous != *server_sequence
817 || sequence != server_sequence.saturating_add(1)
818 {
819 return Err(SdkError::InvalidResponse(
820 "order command event sequence is not contiguous".to_owned(),
821 ));
822 }
823 *server_sequence = sequence;
824 Ok(())
825}
826
827fn event_sequence(
828 event: &PlatformOrderCommandEvent,
829) -> Option<(u16, &str, &str, &str, &str, &str)> {
830 match event {
831 PlatformOrderCommandEvent::ChallengeResult {
832 schema_version,
833 contract_version,
834 market_id,
835 stream_id,
836 sequence,
837 previous_sequence,
838 ..
839 }
840 | PlatformOrderCommandEvent::PrepareResult {
841 schema_version,
842 contract_version,
843 market_id,
844 stream_id,
845 sequence,
846 previous_sequence,
847 ..
848 }
849 | PlatformOrderCommandEvent::SubmitResult {
850 schema_version,
851 contract_version,
852 market_id,
853 stream_id,
854 sequence,
855 previous_sequence,
856 ..
857 }
858 | PlatformOrderCommandEvent::StatusResult {
859 schema_version,
860 contract_version,
861 market_id,
862 stream_id,
863 sequence,
864 previous_sequence,
865 ..
866 }
867 | PlatformOrderCommandEvent::DeadManResult {
868 schema_version,
869 contract_version,
870 market_id,
871 stream_id,
872 sequence,
873 previous_sequence,
874 ..
875 }
876 | PlatformOrderCommandEvent::CommandError {
877 schema_version,
878 contract_version,
879 market_id,
880 stream_id,
881 sequence,
882 previous_sequence,
883 ..
884 }
885 | PlatformOrderCommandEvent::Heartbeat {
886 schema_version,
887 contract_version,
888 market_id,
889 stream_id,
890 sequence,
891 previous_sequence,
892 ..
893 } => Some((
894 *schema_version,
895 contract_version,
896 market_id,
897 stream_id,
898 sequence,
899 previous_sequence,
900 )),
901 PlatformOrderCommandEvent::AuthChallenge { .. }
902 | PlatformOrderCommandEvent::Ready { .. } => None,
903 }
904}
905
906fn event_request_id(event: &PlatformOrderCommandEvent) -> Option<&str> {
907 match event {
908 PlatformOrderCommandEvent::ChallengeResult { request_id, .. }
909 | PlatformOrderCommandEvent::PrepareResult { request_id, .. }
910 | PlatformOrderCommandEvent::SubmitResult { request_id, .. }
911 | PlatformOrderCommandEvent::StatusResult { request_id, .. }
912 | PlatformOrderCommandEvent::DeadManResult { request_id, .. }
913 | PlatformOrderCommandEvent::CommandError { request_id, .. } => Some(request_id),
914 _ => None,
915 }
916}
917
918fn parse_wire_sequence(value: &str) -> Result<u64, SdkError> {
919 if value.is_empty()
920 || !value.bytes().all(|byte| byte.is_ascii_digit())
921 || (value.len() > 1 && value.starts_with('0'))
922 {
923 return Err(SdkError::InvalidResponse(
924 "order command sequence is not canonical".to_owned(),
925 ));
926 }
927 value
928 .parse()
929 .map_err(|_| SdkError::InvalidResponse("order command sequence exceeds u64".to_owned()))
930}
931
932fn checked_dead_man_timeout(timeout: Duration) -> Result<u64, SdkError> {
933 let timeout_ms = u64::try_from(timeout.as_millis()).unwrap_or(u64::MAX);
934 if !(1_000..=30_000).contains(&timeout_ms) {
935 return Err(SdkError::InvalidRequest(
936 "dead-man timeout must be between one and thirty seconds".to_owned(),
937 ));
938 }
939 Ok(timeout_ms)
940}
941
942fn validate_challenge_result(
943 market_id: &str,
944 request: &PlatformOrderChallengeRequest,
945 response: &PlatformOrderChallengeResponse,
946) -> Result<(), SdkError> {
947 validate_platform_version(response.schema_version, &response.contract_version)?;
948 if response.market_id != market_id
949 || response.action != order_request_action(request)
950 || !valid_handle(&response.challenge_id, "oc_")
951 || response.order_ids.is_empty()
952 || response.order_ids.len() > 12
953 || response.expires_at_ms <= response.server_time_ms
954 || response
955 .order_ids
956 .iter()
957 .any(|id| !valid_handle(id, "order_"))
958 {
959 return Err(SdkError::InvalidResponse(
960 "order challenge bindings are invalid".to_owned(),
961 ));
962 }
963 canonical_base64(
964 &response.authorization_payload_base64,
965 "authorization_payload_base64",
966 )?;
967 Ok(())
968}
969
970fn validate_prepared(
971 market_id: &str,
972 response: &PlatformOrderPrepareResponse,
973) -> Result<(), SdkError> {
974 validate_platform_version(response.schema_version, &response.contract_version)?;
975 if response.market_id != market_id
976 || !valid_handle(&response.order_control_id, "or_")
977 || response.order_ids.is_empty()
978 || response.order_ids.len() > 12
979 || response.expires_at_ms == 0
980 {
981 return Err(SdkError::InvalidResponse(
982 "prepared order control is invalid".to_owned(),
983 ));
984 }
985 canonical_base64(&response.transaction_base64, "transaction_base64")?;
986 canonical_base58_32(&response.recent_blockhash, "recent_blockhash")?;
987 Ok(())
988}
989
990fn normalize_submit_request(
991 request: PlatformOrderSubmitRequest,
992) -> Result<PlatformOrderSubmitRequest, SdkError> {
993 if !valid_handle(&request.order_control_id, "or_") {
994 return Err(SdkError::InvalidRequest(
995 "order_control_id is invalid".to_owned(),
996 ));
997 }
998 Ok(PlatformOrderSubmitRequest {
999 order_control_id: request.order_control_id,
1000 signed_transaction_base64: canonical_base64(
1001 &request.signed_transaction_base64,
1002 "signed_transaction_base64",
1003 )?,
1004 idempotency_key: normalize_idempotency_key(&request.idempotency_key)?,
1005 })
1006}
1007
1008fn normalize_status_request(
1009 request: PlatformOrderStatusRequest,
1010) -> Result<PlatformOrderStatusRequest, SdkError> {
1011 if !valid_handle(&request.order_control_id, "or_") {
1012 return Err(SdkError::InvalidRequest(
1013 "order_control_id is invalid".to_owned(),
1014 ));
1015 }
1016 Ok(PlatformOrderStatusRequest {
1017 order_control_id: request.order_control_id,
1018 idempotency_key: normalize_idempotency_key(&request.idempotency_key)?,
1019 })
1020}
1021
1022fn validate_submit(
1023 market_id: &str,
1024 control_id: &str,
1025 response: &PlatformOrderSubmitResponse,
1026) -> Result<(), SdkError> {
1027 validate_platform_version(response.schema_version, &response.contract_version)?;
1028 if response.market_id != market_id
1029 || response.order_control_id != control_id
1030 || response.status != PlatformOrderSubmissionStatus::Submitted
1031 {
1032 return Err(SdkError::InvalidResponse(
1033 "order submit bindings are invalid".to_owned(),
1034 ));
1035 }
1036 canonical_signature(&response.signature, "signature")?;
1037 Ok(())
1038}
1039
1040fn validate_status(
1041 market_id: &str,
1042 control_id: &str,
1043 response: &PlatformOrderStatusResponse,
1044) -> Result<(), SdkError> {
1045 validate_platform_version(response.schema_version, &response.contract_version)?;
1046 if response.market_id != market_id
1047 || response.order_control_id != control_id
1048 || response.order_ids.is_empty()
1049 || response.order_ids.len() > 12
1050 || response
1051 .order_ids
1052 .iter()
1053 .any(|id| !valid_handle(id, "order_"))
1054 || (response.status == PlatformOrderControlStatus::Failed)
1055 != response
1056 .failure_code
1057 .as_deref()
1058 .is_some_and(|code| !code.is_empty())
1059 {
1060 return Err(SdkError::InvalidResponse(
1061 "order status bindings are invalid".to_owned(),
1062 ));
1063 }
1064 canonical_signature(&response.signature, "signature")?;
1065 Ok(())
1066}
1067
1068#[cfg(test)]
1069mod tests {
1070 use super::*;
1071
1072 #[test]
1073 fn authentication_domain_binds_every_identity() {
1074 let message = String::from_utf8(order_stream_auth_message(
1075 "market_11111111111111111111111111111111",
1076 "owner",
1077 "session",
1078 &"ab".repeat(32),
1079 ))
1080 .unwrap();
1081 assert_eq!(
1082 message,
1083 format!(
1084 "{ORDER_STREAM_AUTH_DOMAIN}\nmarket_11111111111111111111111111111111\nowner\nsession\n{}",
1085 "ab".repeat(32)
1086 )
1087 );
1088 }
1089
1090 #[test]
1091 fn canonical_sequence_rejects_gaps_and_leading_zeroes() {
1092 assert_eq!(parse_wire_sequence("8").unwrap(), 8);
1093 assert!(parse_wire_sequence("08").is_err());
1094 assert!(parse_wire_sequence("-1").is_err());
1095 }
1096
1097 #[test]
1098 fn parses_ordered_event_batches_and_single_frame_fallback() {
1099 let event = serde_json::json!({
1100 "type": "heartbeat",
1101 "schema_version": 2,
1102 "contract_version": "2.0",
1103 "market_id": "market_11111111111111111111111111111111",
1104 "stream_id": "order_command_stream_11111111111111111111111111111111",
1105 "sequence": "2",
1106 "previous_sequence": "1",
1107 "server_time_ms": 1_786_810_000_000u64
1108 });
1109 let single = parse_order_command_events(&event.to_string()).unwrap();
1110 assert_eq!(single.len(), 1);
1111 let batch = parse_order_command_events(
1112 &serde_json::Value::Array(vec![event.clone(), event]).to_string(),
1113 )
1114 .unwrap();
1115 assert_eq!(batch.len(), 2);
1116 assert!(parse_order_command_events("[]").is_err());
1117 }
1118}