1use chrono::{DateTime, Utc};
4use serde::{Deserialize, Deserializer, Serialize, Serializer};
5use thiserror::Error;
6use uuid::Uuid;
7
8use crate::{
9 ArtifactId, Capability, CapabilitySet, CommandId, CommandOutcome, ErrorLayer, PrincipalId,
10};
11
12pub const CURRENT_INTERFACE_VERSION: &str = "2026-08-19";
14
15#[derive(Debug, Clone, Copy, Default, Eq, PartialEq, Ord, PartialOrd, Hash)]
17#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
18pub struct InterfaceVersion;
19
20impl InterfaceVersion {
21 pub const CURRENT: Self = Self;
22}
23
24impl TryFrom<&str> for InterfaceVersion {
25 type Error = InterfaceValidationError;
26
27 fn try_from(value: &str) -> Result<Self, Self::Error> {
28 if value == CURRENT_INTERFACE_VERSION {
29 Ok(Self)
30 } else {
31 Err(InterfaceValidationError::UnsupportedInterfaceVersion(
32 value.to_owned(),
33 ))
34 }
35 }
36}
37
38impl TryFrom<String> for InterfaceVersion {
39 type Error = InterfaceValidationError;
40
41 fn try_from(value: String) -> Result<Self, Self::Error> {
42 Self::try_from(value.as_str())
43 }
44}
45
46impl Serialize for InterfaceVersion {
47 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
48 where
49 S: Serializer,
50 {
51 serializer.serialize_str(CURRENT_INTERFACE_VERSION)
52 }
53}
54
55impl<'de> Deserialize<'de> for InterfaceVersion {
56 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
57 where
58 D: Deserializer<'de>,
59 {
60 let value = String::deserialize(deserializer)?;
61 Self::try_from(value).map_err(serde::de::Error::custom)
62 }
63}
64
65#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
67#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
68#[serde(transparent)]
69pub struct CorrelationId(Uuid);
70
71impl CorrelationId {
72 pub fn new() -> Self {
73 Self(Uuid::new_v4())
74 }
75
76 pub fn from_uuid(value: Uuid) -> Self {
77 Self(value)
78 }
79
80 pub fn as_uuid(&self) -> Uuid {
81 self.0
82 }
83}
84
85impl Default for CorrelationId {
86 fn default() -> Self {
87 Self::new()
88 }
89}
90
91#[derive(Debug, Clone, PartialEq, Eq, Hash)]
93#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
94pub struct IdempotencyKey(String);
95
96impl IdempotencyKey {
97 pub fn as_str(&self) -> &str {
98 &self.0
99 }
100}
101
102impl TryFrom<&str> for IdempotencyKey {
103 type Error = InterfaceValidationError;
104
105 fn try_from(value: &str) -> Result<Self, Self::Error> {
106 if !(1..=128).contains(&value.len())
107 || !value
108 .as_bytes()
109 .iter()
110 .all(|byte| (0x20..=0x7e).contains(byte))
111 {
112 return Err(InterfaceValidationError::InvalidIdempotencyKey);
113 }
114 Ok(Self(value.to_owned()))
115 }
116}
117
118impl TryFrom<String> for IdempotencyKey {
119 type Error = InterfaceValidationError;
120
121 fn try_from(value: String) -> Result<Self, Self::Error> {
122 Self::try_from(value.as_str())
123 }
124}
125
126impl Serialize for IdempotencyKey {
127 fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
128 where
129 S: Serializer,
130 {
131 serializer.serialize_str(&self.0)
132 }
133}
134
135impl<'de> Deserialize<'de> for IdempotencyKey {
136 fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
137 where
138 D: Deserializer<'de>,
139 {
140 let value = String::deserialize(deserializer)?;
141 Self::try_from(value).map_err(serde::de::Error::custom)
142 }
143}
144
145#[derive(
146 Debug, Clone, Copy, Default, Eq, PartialEq, Ord, PartialOrd, Hash, Serialize, Deserialize,
147)]
148#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
149#[serde(transparent)]
150pub struct EventCursor(pub u64);
151
152impl EventCursor {
153 pub const ZERO: Self = Self(0);
154}
155
156impl From<u64> for EventCursor {
157 fn from(value: u64) -> Self {
158 Self(value)
159 }
160}
161
162#[derive(Debug, Clone, Serialize, Deserialize)]
163#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
164#[serde(rename_all = "camelCase")]
165pub struct RequestContext {
166 pub interface_version: InterfaceVersion,
167 pub correlation_id: CorrelationId,
168 pub principal_id: PrincipalId,
169 pub capabilities: CapabilitySet,
170 pub deadline: DateTime<Utc>,
171 pub idempotency_key: Option<IdempotencyKey>,
172}
173
174impl RequestContext {
175 pub fn new_for_test(
176 principal_id: PrincipalId,
177 capabilities: impl IntoIterator<Item = Capability>,
178 deadline: DateTime<Utc>,
179 ) -> Self {
180 Self {
181 interface_version: InterfaceVersion::CURRENT,
182 correlation_id: CorrelationId::new(),
183 principal_id,
184 capabilities: CapabilitySet::new(capabilities),
185 deadline,
186 idempotency_key: None,
187 }
188 }
189
190 pub fn validate_at(
192 &self,
193 dispatch_time: DateTime<Utc>,
194 ) -> Result<(), InterfaceValidationError> {
195 if self.deadline <= dispatch_time {
196 return Err(InterfaceValidationError::ExpiredDeadline);
197 }
198 Ok(())
199 }
200}
201
202#[derive(Debug, Clone, Copy, Eq, PartialEq, Hash)]
204pub enum InterfaceOperation {
205 RuntimeInfo,
206 CreateSession,
207 ReadSession,
208 DeleteSession,
209 OpenPage,
210 ReadPage,
211 ClosePage,
212 SubmitCommand,
213 CreateCheckpoint,
214 ReadCheckpoint,
215 RecoverWorkflow,
216 ReadArtifact,
217 ReadContext,
218 CaptureArtifact,
219 SubscribeEvents,
220 SubmitJob,
221 ReadJob,
222 CancelJob,
223 IssuePrincipal,
224 RevokePrincipal,
225}
226
227impl InterfaceOperation {
228 pub const ALL: [Self; 20] = [
229 Self::RuntimeInfo,
230 Self::CreateSession,
231 Self::ReadSession,
232 Self::DeleteSession,
233 Self::OpenPage,
234 Self::ReadPage,
235 Self::ClosePage,
236 Self::SubmitCommand,
237 Self::CreateCheckpoint,
238 Self::ReadCheckpoint,
239 Self::RecoverWorkflow,
240 Self::ReadArtifact,
241 Self::ReadContext,
242 Self::CaptureArtifact,
243 Self::SubscribeEvents,
244 Self::SubmitJob,
245 Self::ReadJob,
246 Self::CancelJob,
247 Self::IssuePrincipal,
248 Self::RevokePrincipal,
249 ];
250
251 pub const fn as_str(self) -> &'static str {
252 match self {
253 Self::RuntimeInfo => "runtimeInfo",
254 Self::CreateSession => "createSession",
255 Self::ReadSession => "readSession",
256 Self::DeleteSession => "deleteSession",
257 Self::OpenPage => "openPage",
258 Self::ReadPage => "readPage",
259 Self::ClosePage => "closePage",
260 Self::SubmitCommand => "submitCommand",
261 Self::CreateCheckpoint => "createCheckpoint",
262 Self::ReadCheckpoint => "readCheckpoint",
263 Self::RecoverWorkflow => "recoverWorkflow",
264 Self::ReadArtifact => "readArtifact",
265 Self::ReadContext => "readContext",
266 Self::CaptureArtifact => "captureArtifact",
267 Self::SubscribeEvents => "subscribeEvents",
268 Self::SubmitJob => "submitJob",
269 Self::ReadJob => "readJob",
270 Self::CancelJob => "cancelJob",
271 Self::IssuePrincipal => "issuePrincipal",
272 Self::RevokePrincipal => "revokePrincipal",
273 }
274 }
275
276 pub const fn required(self) -> &'static [Capability] {
277 match self {
278 Self::RuntimeInfo => &[Capability::SessionRead],
279 Self::CreateSession => &[Capability::SessionWrite],
280 Self::ReadSession => &[Capability::SessionRead],
281 Self::DeleteSession => &[Capability::SessionWrite],
282 Self::OpenPage => &[Capability::PageWrite],
283 Self::ReadPage => &[Capability::PageRead],
284 Self::ClosePage => &[Capability::PageWrite],
285 Self::SubmitCommand => &[Capability::BrowserMutate],
286 Self::CreateCheckpoint => &[Capability::RecoveryWrite],
287 Self::ReadCheckpoint => &[Capability::RecoveryRead],
288 Self::RecoverWorkflow => &[Capability::RecoveryWrite],
289 Self::ReadArtifact => &[Capability::ArtifactRead],
290 Self::ReadContext => &[Capability::ContextRead],
291 Self::CaptureArtifact => &[Capability::ArtifactCapture],
292 Self::SubscribeEvents => &[Capability::SessionRead],
293 Self::SubmitJob => &[Capability::JobSubmit],
294 Self::ReadJob => &[Capability::JobRead],
295 Self::CancelJob => &[Capability::JobCancel],
296 Self::IssuePrincipal => &[Capability::AuthorityAdmin],
297 Self::RevokePrincipal => &[Capability::AuthorityAdmin],
298 }
299 }
300}
301
302#[cfg(test)]
303mod interface_operation_tests {
304 use super::InterfaceOperation;
305 use std::collections::HashSet;
306
307 #[test]
308 fn all_is_exhaustive_and_unique() {
309 let expected_names = [
310 "runtimeInfo",
311 "createSession",
312 "readSession",
313 "deleteSession",
314 "openPage",
315 "readPage",
316 "closePage",
317 "submitCommand",
318 "createCheckpoint",
319 "readCheckpoint",
320 "recoverWorkflow",
321 "readArtifact",
322 "readContext",
323 "captureArtifact",
324 "subscribeEvents",
325 "submitJob",
326 "readJob",
327 "cancelJob",
328 "issuePrincipal",
329 "revokePrincipal",
330 ];
331
332 fn listed(operation: InterfaceOperation) {
333 match operation {
334 InterfaceOperation::RuntimeInfo
335 | InterfaceOperation::CreateSession
336 | InterfaceOperation::ReadSession
337 | InterfaceOperation::DeleteSession
338 | InterfaceOperation::OpenPage
339 | InterfaceOperation::ReadPage
340 | InterfaceOperation::ClosePage
341 | InterfaceOperation::SubmitCommand
342 | InterfaceOperation::CreateCheckpoint
343 | InterfaceOperation::ReadCheckpoint
344 | InterfaceOperation::RecoverWorkflow
345 | InterfaceOperation::ReadArtifact
346 | InterfaceOperation::ReadContext
347 | InterfaceOperation::CaptureArtifact
348 | InterfaceOperation::SubscribeEvents
349 | InterfaceOperation::SubmitJob
350 | InterfaceOperation::ReadJob
351 | InterfaceOperation::CancelJob
352 | InterfaceOperation::IssuePrincipal
353 | InterfaceOperation::RevokePrincipal => {}
354 }
355 assert!(InterfaceOperation::ALL.contains(&operation));
356 }
357
358 for (operation, expected_name) in InterfaceOperation::ALL.into_iter().zip(expected_names) {
359 listed(operation);
360 assert_eq!(operation.as_str(), expected_name);
361 }
362 assert_eq!(
363 InterfaceOperation::ALL
364 .into_iter()
365 .collect::<HashSet<_>>()
366 .len(),
367 InterfaceOperation::ALL.len()
368 );
369 }
370}
371
372#[derive(Debug, Clone, Serialize, Deserialize)]
373#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
374#[serde(tag = "kind", rename_all = "camelCase")]
375pub enum InterfaceEvent {
376 CommandOutcome {
377 cursor: EventCursor,
378 command_id: CommandId,
379 outcome: CommandOutcome,
380 },
381 ArtifactCaptured {
382 cursor: EventCursor,
383 artifact_id: ArtifactId,
384 },
385 EventGap {
386 earliest_available_cursor: EventCursor,
387 },
388}
389
390#[derive(Debug, Clone, Copy, Eq, PartialEq, Serialize, Deserialize)]
391#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
392#[serde(rename_all = "camelCase")]
393pub enum InterfaceErrorCode {
394 InvalidRequest,
395 UnsupportedInterfaceVersion,
396 InvalidIdempotencyKey,
397 IdempotencyConflict,
398 DeadlineExceeded,
399 AuthenticationFailed,
400 TokenExpired,
401 MissingCapability,
402 MalformedScope,
403 ArtifactDenied,
404 UnsupportedOperation,
405 NotFound,
406 ResourceExhausted,
407 EngineUnreachable,
408 Internal,
409}
410
411#[derive(Debug, Clone, Serialize, Deserialize)]
412#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
413#[serde(rename_all = "camelCase")]
414pub struct InterfaceError {
415 pub code: InterfaceErrorCode,
416 pub layer: ErrorLayer,
417 pub message: String,
418 pub correlation_id: CorrelationId,
419 pub command_id: Option<CommandId>,
420 pub retryable: bool,
421 pub retry_after_ms: Option<u64>,
422 pub reconciliation_required: bool,
423 pub required_capability: Option<Capability>,
424}
425
426#[derive(Debug, Clone, Error, Eq, PartialEq)]
427pub enum InterfaceValidationError {
428 #[error("unsupported interface version: {0}")]
429 UnsupportedInterfaceVersion(String),
430 #[error("idempotency key must contain 1-128 printable ASCII characters")]
431 InvalidIdempotencyKey,
432 #[error("request deadline must be after dispatch time")]
433 ExpiredDeadline,
434}