返回 CodeWhale
engine_owner.rs
根目录 / crates / protocol / src / engine_owner.rs
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
285 lines RUST