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