| 1 | //! Safe, ordered Engine owner projection shared by desktop and web hosts. |
| 2 | //! |
| 3 | //! The Engine owns operation and turn facts: it emits |
| 4 | //! `operation_activity_started` / `operation_activity_completed` with an |
| 5 | //! [`OwnerActivityKind`] and [`OwnerOperationOutcome`], never a tool name, |
| 6 | //! argument, command or result. Apps may bind this session-scoped projection |
| 7 | //! to an authenticated account and serve it to another client; the account |
| 8 | //! identity is intentionally not supplied by the Engine. |
| 9 | //! |
| 10 | //! Known limits: |
| 11 | //! - Operation events only carry `Reading`, `Editing`, `Searching`, |
| 12 | //! `Testing`, `Executing`, `Browsing`, `Computer`, `Memory` or `Tool`. |
| 13 | //! `Thinking`, `Responding` and `Delegating` are derived by the reducer |
| 14 | //! from message, reasoning and agent lifecycle events. |
| 15 | //! - [`EngineOwnerProjection`] itself is computed today by the pet reducer |
| 16 | //! (`pet/src/core/pet-engine.ts`, bundled into the TUI pet worker), and |
| 17 | //! validated here on the way back. No Rust producer exists yet; a host |
| 18 | //! that needs it outside the pet must decide whether the runtime produces |
| 19 | //! it in Rust rather than adding a second reducer. |
| 20 | //! - Code-mode (`execute_tools`) nested calls do not report activity yet. |
| 21 | |
| 22 | use serde::{Deserialize, Serialize}; |
| 23 | |
| 24 | use crate::event_msg::TurnOutcomeStatus; |
| 25 | |
| 26 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] |
| 27 | #[serde(rename_all = "snake_case")] |
| 28 | pub enum OwnerActivityKind { |
| 29 | Reading, |
| 30 | Editing, |
| 31 | Searching, |
| 32 | Testing, |
| 33 | Executing, |
| 34 | Browsing, |
| 35 | Computer, |
| 36 | Memory, |
| 37 | Tool, |
| 38 | Thinking, |
| 39 | Responding, |
| 40 | Delegating, |
| 41 | } |
| 42 | |
| 43 | impl OwnerActivityKind { |
| 44 | #[must_use] |
| 45 | pub const fn as_str(self) -> &'static str { |
| 46 | match self { |
| 47 | Self::Reading => "reading", |
| 48 | Self::Editing => "editing", |
| 49 | Self::Searching => "searching", |
| 50 | Self::Testing => "testing", |
| 51 | Self::Executing => "executing", |
| 52 | Self::Browsing => "browsing", |
| 53 | Self::Computer => "computer", |
| 54 | Self::Memory => "memory", |
| 55 | Self::Tool => "tool", |
| 56 | Self::Thinking => "thinking", |
| 57 | Self::Responding => "responding", |
| 58 | Self::Delegating => "delegating", |
| 59 | } |
| 60 | } |
| 61 | } |
| 62 | |
| 63 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] |
| 64 | #[serde(rename_all = "snake_case")] |
| 65 | pub enum OwnerOperationOutcome { |
| 66 | Succeeded, |
| 67 | Failed, |
| 68 | Cancelled, |
| 69 | Denied, |
| 70 | } |
| 71 | |
| 72 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] |
| 73 | #[serde(rename_all = "snake_case")] |
| 74 | pub enum OwnerPresence { |
| 75 | Unknown, |
| 76 | Working, |
| 77 | NeedsYou, |
| 78 | Done, |
| 79 | Idle, |
| 80 | } |
| 81 | |
| 82 | #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] |
| 83 | #[serde(rename_all = "snake_case")] |
| 84 | pub enum OwnerFreshness { |
| 85 | Missing, |
| 86 | Fresh, |
| 87 | Stale, |
| 88 | } |
| 89 | |
| 90 | #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] |
| 91 | #[serde(rename_all = "camelCase", deny_unknown_fields)] |
| 92 | pub struct OwnerActiveSpan { |
| 93 | pub activity_kind: OwnerActivityKind, |
| 94 | pub started_at_ms: f64, |
| 95 | } |
| 96 | |
| 97 | #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] |
| 98 | #[serde(rename_all = "camelCase", deny_unknown_fields)] |
| 99 | pub struct OwnerFailedToolAge { |
| 100 | pub activity_kind: OwnerActivityKind, |
| 101 | pub age_ms: f64, |
| 102 | } |
| 103 | |
| 104 | /// Ordered, account-neutral read model. `cursor` is monotonic within the |
| 105 | /// existing session owner. The Apps control plane supplies account isolation |
| 106 | /// by authenticating the caller and resolving the requested session there. |
| 107 | #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)] |
| 108 | #[serde(rename_all = "camelCase", deny_unknown_fields)] |
| 109 | pub struct EngineOwnerProjection { |
| 110 | pub schema_version: u8, |
| 111 | pub session_id: Option<String>, |
| 112 | pub cursor: u64, |
| 113 | pub observed: bool, |
| 114 | pub freshness: OwnerFreshness, |
| 115 | pub authoritative_presence: OwnerPresence, |
| 116 | pub activity_kind: Option<OwnerActivityKind>, |
| 117 | pub observed_at_ms: Option<f64>, |
| 118 | pub parallel_agent_count: u32, |
| 119 | pub active_spans: Vec<OwnerActiveSpan>, |
| 120 | pub turn_id: Option<String>, |
| 121 | pub turn_outcome: Option<TurnOutcomeStatus>, |
| 122 | pub done_effect_id: Option<String>, |
| 123 | pub failed_tool_age: Option<OwnerFailedToolAge>, |
| 124 | } |
| 125 | |
| 126 | impl EngineOwnerProjection { |
| 127 | #[must_use] |
| 128 | pub fn is_valid(&self) -> bool { |
| 129 | self.schema_version == 1 |
| 130 | && self.session_id.as_ref().is_none_or(|id| { |
| 131 | !id.is_empty() && id.len() <= 256 && id.bytes().all(|byte| !byte.is_ascii_control()) |
| 132 | }) |
| 133 | && self.turn_id.as_ref().is_none_or(|id| { |
| 134 | !id.is_empty() && id.len() <= 256 && id.bytes().all(|byte| !byte.is_ascii_control()) |
| 135 | }) |
| 136 | && self.active_spans.len() <= 4 |
| 137 | && self.parallel_agent_count <= 100_000 |
| 138 | && self |
| 139 | .observed_at_ms |
| 140 | .is_none_or(|time| time.is_finite() && time >= 0.0) |
| 141 | && self.active_spans.iter().all(|span| { |
| 142 | span.started_at_ms.is_finite() |
| 143 | && span.started_at_ms >= 0.0 |
| 144 | && self |
| 145 | .observed_at_ms |
| 146 | .is_none_or(|observed_at| span.started_at_ms <= observed_at) |
| 147 | }) |
| 148 | && self |
| 149 | .failed_tool_age |
| 150 | .as_ref() |
| 151 | .is_none_or(|failure| failure.age_ms.is_finite() && failure.age_ms >= 0.0) |
| 152 | && (self.turn_outcome.is_none() || self.turn_id.is_some()) |
| 153 | && match self.freshness { |
| 154 | OwnerFreshness::Missing => { |
| 155 | !self.observed |
| 156 | && self.observed_at_ms.is_none() |
| 157 | && self.activity_kind.is_none() |
| 158 | && self.active_spans.is_empty() |
| 159 | && self.parallel_agent_count == 0 |
| 160 | && self.authoritative_presence == OwnerPresence::Unknown |
| 161 | && self.turn_id.is_none() |
| 162 | && self.turn_outcome.is_none() |
| 163 | && self.failed_tool_age.is_none() |
| 164 | && self.done_effect_id.is_none() |
| 165 | } |
| 166 | OwnerFreshness::Stale => { |
| 167 | self.observed |
| 168 | && self.session_id.is_some() |
| 169 | && self.observed_at_ms.is_some() |
| 170 | && self.activity_kind.is_none() |
| 171 | && self.active_spans.is_empty() |
| 172 | && self.parallel_agent_count == 0 |
| 173 | && self.authoritative_presence == OwnerPresence::Unknown |
| 174 | && self.failed_tool_age.is_none() |
| 175 | && self.done_effect_id.is_none() |
| 176 | } |
| 177 | OwnerFreshness::Fresh => { |
| 178 | self.observed |
| 179 | && self.session_id.is_some() |
| 180 | && self.observed_at_ms.is_some() |
| 181 | && (self.authoritative_presence != OwnerPresence::Done |
| 182 | || (self.turn_outcome == Some(TurnOutcomeStatus::Completed) |
| 183 | && self.turn_id.is_some() |
| 184 | && self.done_effect_id == self.turn_id)) |
| 185 | && (self.done_effect_id.is_none() |
| 186 | || self.authoritative_presence == OwnerPresence::Done) |
| 187 | } |
| 188 | } |
| 189 | } |
| 190 | } |
| 191 | |
| 192 | #[cfg(test)] |
| 193 | mod tests { |
| 194 | use super::*; |
| 195 | |
| 196 | fn projection() -> EngineOwnerProjection { |
| 197 | EngineOwnerProjection { |
| 198 | schema_version: 1, |
| 199 | session_id: Some("session:opaque".into()), |
| 200 | cursor: 4, |
| 201 | observed: true, |
| 202 | freshness: OwnerFreshness::Fresh, |
| 203 | authoritative_presence: OwnerPresence::Working, |
| 204 | activity_kind: Some(OwnerActivityKind::Reading), |
| 205 | observed_at_ms: Some(100.0), |
| 206 | parallel_agent_count: 0, |
| 207 | active_spans: vec![OwnerActiveSpan { |
| 208 | activity_kind: OwnerActivityKind::Reading, |
| 209 | started_at_ms: 50.0, |
| 210 | }], |
| 211 | turn_id: Some("turn-1".into()), |
| 212 | turn_outcome: None, |
| 213 | done_effect_id: None, |
| 214 | failed_tool_age: None, |
| 215 | } |
| 216 | } |
| 217 | |
| 218 | #[test] |
| 219 | fn projection_accepts_trusted_scoped_activity_and_rejects_bad_clocks() { |
| 220 | assert!(projection().is_valid()); |
| 221 | |
| 222 | let mut future_span = projection(); |
| 223 | future_span.active_spans[0].started_at_ms = 101.0; |
| 224 | assert!(!future_span.is_valid()); |
| 225 | |
| 226 | let mut too_many = projection(); |
| 227 | too_many.active_spans = vec![too_many.active_spans[0].clone(); 5]; |
| 228 | assert!(!too_many.is_valid()); |
| 229 | } |
| 230 | |
| 231 | #[test] |
| 232 | fn projection_needs_authoritative_presence_for_wait_and_completed_turn_for_done() { |
| 233 | let mut waiting = projection(); |
| 234 | waiting.authoritative_presence = OwnerPresence::NeedsYou; |
| 235 | assert!(waiting.is_valid()); |
| 236 | |
| 237 | let mut done = projection(); |
| 238 | done.authoritative_presence = OwnerPresence::Done; |
| 239 | done.activity_kind = None; |
| 240 | done.active_spans.clear(); |
| 241 | done.turn_outcome = Some(TurnOutcomeStatus::Completed); |
| 242 | done.done_effect_id = done.turn_id.clone(); |
| 243 | assert!(done.is_valid()); |
| 244 | |
| 245 | done.turn_outcome = Some(TurnOutcomeStatus::Interrupted); |
| 246 | assert!(!done.is_valid()); |
| 247 | done.turn_outcome = Some(TurnOutcomeStatus::Failed); |
| 248 | assert!(!done.is_valid()); |
| 249 | } |
| 250 | |
| 251 | #[test] |
| 252 | fn missing_and_stale_projection_cannot_claim_specific_activity() { |
| 253 | let missing = EngineOwnerProjection { |
| 254 | schema_version: 1, |
| 255 | session_id: None, |
| 256 | cursor: 0, |
| 257 | observed: false, |
| 258 | freshness: OwnerFreshness::Missing, |
| 259 | authoritative_presence: OwnerPresence::Unknown, |
| 260 | activity_kind: None, |
| 261 | observed_at_ms: None, |
| 262 | parallel_agent_count: 0, |
| 263 | active_spans: Vec::new(), |
| 264 | turn_id: None, |
| 265 | turn_outcome: None, |
| 266 | done_effect_id: None, |
| 267 | failed_tool_age: None, |
| 268 | }; |
| 269 | assert!(missing.is_valid()); |
| 270 | |
| 271 | let mut stale = projection(); |
| 272 | stale.freshness = OwnerFreshness::Stale; |
| 273 | stale.authoritative_presence = OwnerPresence::Unknown; |
| 274 | stale.activity_kind = None; |
| 275 | stale.parallel_agent_count = 0; |
| 276 | stale.active_spans.clear(); |
| 277 | stale.failed_tool_age = None; |
| 278 | stale.done_effect_id = None; |
| 279 | assert!(stale.is_valid()); |
| 280 | |
| 281 | stale.activity_kind = Some(OwnerActivityKind::Reading); |
| 282 | assert!(!stale.is_valid()); |
| 283 | } |
| 284 | } |
| 285 |