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