lean_ctx/core/ocla/
wire_stream.rs1use std::time::{SystemTime, UNIX_EPOCH};
4
5use serde::{Deserialize, Serialize};
6
7use super::types::{CanonicalTokenEnvelopeV1, OclaError, OclaResult};
8use super::wire::MAX_OCLA_WIRE_BYTES;
9
10#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
12#[serde(tag = "type", content = "data")]
13#[serde(rename_all = "snake_case")]
14pub enum StreamFrame {
15 Data(Box<CanonicalTokenEnvelopeV1>),
16 Heartbeat,
17 Cancel,
18 Done,
19}
20
21impl StreamFrame {
22 fn validate(&self) -> OclaResult<()> {
23 if let Self::Data(envelope) = self {
24 envelope.validate()?;
25 }
26 Ok(())
27 }
28}
29
30#[derive(Clone, Copy, Debug, Eq, PartialEq)]
35pub struct StreamConfig {
36 pub deadline_ms: u64,
37 pub max_backpressure_frames: usize,
38}
39
40pub fn encode_frame(frame: &StreamFrame) -> OclaResult<String> {
42 frame.validate()?;
43 let json = serde_json::to_string(frame).map_err(|error| {
44 OclaError::InvalidRequest(format!("cannot encode stream frame: {error}"))
45 })?;
46 let line = format!("{json}\n");
47 if line.len() > MAX_OCLA_WIRE_BYTES {
48 return Err(OclaError::InvalidRequest(format!(
49 "stream frame exceeds {MAX_OCLA_WIRE_BYTES} bytes"
50 )));
51 }
52 Ok(line)
53}
54
55pub fn decode_frame(line: &str) -> OclaResult<StreamFrame> {
57 if line.len() > MAX_OCLA_WIRE_BYTES {
58 return Err(OclaError::InvalidRequest(format!(
59 "stream frame exceeds {MAX_OCLA_WIRE_BYTES} bytes"
60 )));
61 }
62 let line = line.trim_end_matches(['\r', '\n']);
63 let frame: StreamFrame = serde_json::from_str(line).map_err(|error| {
64 OclaError::InvalidRequest(format!("cannot decode stream frame: {error}"))
65 })?;
66 frame.validate()?;
67 Ok(frame)
68}
69
70pub fn validate_deadline(config: &StreamConfig) -> OclaResult<()> {
72 let now_ms = SystemTime::now()
73 .duration_since(UNIX_EPOCH)
74 .map_err(|error| OclaError::InvalidRequest(format!("cannot read system clock: {error}")))?
75 .as_millis();
76 if now_ms >= u128::from(config.deadline_ms) {
77 return Err(OclaError::InvalidRequest("stream deadline expired".into()));
78 }
79 Ok(())
80}
81
82impl StreamConfig {
83 pub fn validate_backpressure(&self, queued_frames: usize) -> OclaResult<()> {
85 if queued_frames > self.max_backpressure_frames {
86 return Err(OclaError::InvalidRequest(format!(
87 "stream backpressure limit exceeded: {queued_frames} queued, limit {}",
88 self.max_backpressure_frames
89 )));
90 }
91 Ok(())
92 }
93}
94
95#[cfg(test)]
96mod tests {
97 use super::*;
98 use crate::core::ocla::{
99 CANONICAL_TOKEN_ENVELOPE_SCHEMA_VERSION, OclaRequestContext, TokenBalanceV1,
100 TokenEnvelopeSurface, TokenFlowDirection,
101 };
102
103 fn envelope() -> CanonicalTokenEnvelopeV1 {
104 CanonicalTokenEnvelopeV1 {
105 schema_version: CANONICAL_TOKEN_ENVELOPE_SCHEMA_VERSION,
106 context: OclaRequestContext {
107 request_id: "request-1".into(),
108 session_id: "session-1".into(),
109 agent_id: "agent-1".into(),
110 content_ref: "blake3:content".into(),
111 tenant_id: None,
112 trace_id: "tr-unit".into(),
113 },
114 surface: TokenEnvelopeSurface::Proxy,
115 direction: TokenFlowDirection::Input,
116 provider: "openai".into(),
117 model: "gpt-5".into(),
118 token_balance: TokenBalanceV1 {
119 original_tokens: 100,
120 materialized_tokens: 80,
121 delivered_tokens: 60,
122 provider_billed_tokens: 60,
123 },
124 route_ref: None,
125 policy_ref: None,
126 idempotency_key: "request-1:input".into(),
127 }
128 }
129
130 #[test]
131 fn data_frame_roundtrips_with_a_trailing_newline() {
132 let original = StreamFrame::Data(Box::new(envelope()));
133 let encoded = encode_frame(&original).expect("encode data frame");
134
135 assert!(encoded.ends_with('\n'));
136 assert_eq!(decode_frame(&encoded).expect("decode data frame"), original);
137 }
138
139 #[test]
140 fn cancel_frame_roundtrips_as_a_payload_free_tag() {
141 let encoded = encode_frame(&StreamFrame::Cancel).expect("encode cancel frame");
142
143 assert_eq!(encoded, "{\"type\":\"cancel\"}\n");
144 assert_eq!(
145 decode_frame(&encoded).expect("decode cancel frame"),
146 StreamFrame::Cancel
147 );
148 }
149
150 #[test]
151 fn expired_deadline_is_rejected_and_future_deadline_is_allowed() {
152 let expired = StreamConfig {
153 deadline_ms: 0,
154 max_backpressure_frames: 1,
155 };
156 assert!(validate_deadline(&expired).is_err());
157
158 let future = StreamConfig {
159 deadline_ms: SystemTime::now()
160 .duration_since(UNIX_EPOCH)
161 .expect("system clock after Unix epoch")
162 .as_millis()
163 .try_into()
164 .expect("test deadline fits u64"),
165 max_backpressure_frames: 1,
166 };
167 assert!(validate_deadline(&future).is_err());
168
169 let future = StreamConfig {
170 deadline_ms: future.deadline_ms.saturating_add(60_000),
171 ..future
172 };
173 assert!(validate_deadline(&future).is_ok());
174 }
175
176 #[test]
177 fn backpressure_limit_allows_boundary_and_rejects_overflow() {
178 let config = StreamConfig {
179 deadline_ms: u64::MAX,
180 max_backpressure_frames: 2,
181 };
182
183 assert!(config.validate_backpressure(2).is_ok());
184 assert!(config.validate_backpressure(3).is_err());
185 }
186
187 #[test]
188 fn invalid_data_frame_is_rejected_before_wire_encoding() {
189 let mut invalid = envelope();
190 invalid.token_balance.delivered_tokens = 81;
191
192 assert!(encode_frame(&StreamFrame::Data(Box::new(invalid))).is_err());
193 }
194}