Skip to main content

starweaver_model/
adapter.rs

1//! Model adapter traits and request context types.
2
3mod context;
4mod error;
5mod guard;
6mod params;
7mod traits;
8
9pub use context::ModelRequestContext;
10pub use error::ModelError;
11pub use guard::{
12    RealModelRequestGuard, allow_real_model_requests, allow_real_model_requests_guard,
13    block_real_model_requests, set_allow_real_model_requests,
14};
15pub use params::{ModelRequestParameters, NativeToolDefinition, ToolDefinition};
16pub use traits::{ModelAdapter, ModelRunSession};
17
18use crate::ModelResponseStreamEvent;
19use starweaver_core::CancellationToken;
20
21/// Receiver for incremental canonical model stream events.
22pub struct ModelResponseEventStream {
23    receiver: tokio::sync::mpsc::Receiver<Result<ModelResponseStreamEvent, ModelError>>,
24    cancellation_token: CancellationToken,
25    drop_abort_token: Option<CancellationToken>,
26}
27
28impl ModelResponseEventStream {
29    /// Build a stream from a channel receiver.
30    #[must_use]
31    pub fn new(
32        receiver: tokio::sync::mpsc::Receiver<Result<ModelResponseStreamEvent, ModelError>>,
33    ) -> Self {
34        Self::new_with_cancellation(receiver, CancellationToken::default())
35    }
36
37    /// Build a stream from a channel receiver and cancellation token.
38    #[must_use]
39    pub const fn new_with_cancellation(
40        receiver: tokio::sync::mpsc::Receiver<Result<ModelResponseStreamEvent, ModelError>>,
41        cancellation_token: CancellationToken,
42    ) -> Self {
43        Self::new_with_cancellation_and_drop_abort(receiver, cancellation_token, None)
44    }
45
46    /// Build a stream from a channel receiver, cancellation token, and drop-abort token.
47    #[must_use]
48    pub const fn new_with_cancellation_and_drop_abort(
49        receiver: tokio::sync::mpsc::Receiver<Result<ModelResponseStreamEvent, ModelError>>,
50        cancellation_token: CancellationToken,
51        drop_abort_token: Option<CancellationToken>,
52    ) -> Self {
53        Self {
54            receiver,
55            cancellation_token,
56            drop_abort_token,
57        }
58    }
59
60    /// Return the transport-local drop-abort token, when the stream owns one.
61    #[must_use]
62    pub fn drop_abort_token(&self) -> Option<CancellationToken> {
63        self.drop_abort_token.clone()
64    }
65
66    /// Receive the next canonical model stream event.
67    pub async fn recv(&mut self) -> Option<Result<ModelResponseStreamEvent, ModelError>> {
68        if self.cancellation_token.is_cancelled() {
69            return Some(Err(ModelError::Cancelled {
70                reason: "model stream cancellation requested".to_string(),
71            }));
72        }
73        tokio::select! {
74            biased;
75            () = self.cancellation_token.cancelled() => Some(Err(ModelError::Cancelled {
76                reason: "model stream cancellation requested".to_string(),
77            })),
78            event = self.receiver.recv() => event,
79        }
80    }
81}
82
83impl Drop for ModelResponseEventStream {
84    fn drop(&mut self) {
85        if let Some(token) = &self.drop_abort_token {
86            token.cancel();
87        }
88    }
89}