Skip to main content

strata_sdk/
order_stream.rs

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/// Cloneable handle to one authenticated persistent order-command socket.
30/// Commands from clones are sequenced by a single writer task, so concurrent
31/// strategies cannot produce sequence gaps on the wire.
32#[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    /// Subscribe to heartbeats, correlated command results, and pushed chain
228    /// confirmations. Lag is explicit through `broadcast::RecvError::Lagged`.
229    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    /// Prepare and arm a fail-closed cancel-all, then maintain its heartbeat
387    /// for this transaction's lifetime. For indefinite unattended exposure,
388    /// use [`Self::maintain_dead_man`] so the blockhash is refreshed too.
389    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    /// Maintain a dead-man indefinitely by heartbeating the current ticket and
427    /// externally signing a fresh exact cancel-all before its blockhash expires.
428    /// The caller supplies owner-controlled signer/verifier adapters in `Arc`s
429    /// solely because the maintenance task must outlive this method call.
430    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
635/// Keeps the durable dead-man deadline alive. Drop is deliberately fail-closed:
636/// it stops heartbeats and leaves the pre-signed cancel-all armed.
637pub 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            // Intentionally no disarm. A crashed/dropped agent must fail closed.
658        }
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}