starweaver_model/
adapter.rs1mod 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
21pub 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 #[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 #[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 #[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 #[must_use]
62 pub fn drop_abort_token(&self) -> Option<CancellationToken> {
63 self.drop_abort_token.clone()
64 }
65
66 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}