Skip to main content

lean_ctx/core/ocla/
wire_stream.rs

1//! Newline-delimited streaming frames for the public OCLA wire contract.
2
3use 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/// One newline-delimited message in an OCLA stream.
11#[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/// Runtime limits for one OCLA stream.
31///
32/// `deadline_ms` is an absolute Unix timestamp in milliseconds. A stream may
33/// have at most `max_backpressure_frames` frames queued for delivery.
34#[derive(Clone, Copy, Debug, Eq, PartialEq)]
35pub struct StreamConfig {
36    pub deadline_ms: u64,
37    pub max_backpressure_frames: usize,
38}
39
40/// Encodes one stream frame as a newline-delimited JSON record.
41pub 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
55/// Decodes one newline-delimited JSON record into a stream frame.
56pub 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
70/// Rejects a stream whose absolute deadline has passed.
71pub 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    /// Rejects a queue that exceeds the configured backpressure limit.
84    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}