Skip to main content

lenso_service/system_plane/
runtime_observability.rs

1use chrono::{DateTime, Utc};
2use schemars::JsonSchema;
3use serde::{Deserialize, Serialize};
4use serde_json::{Value, json};
5use sha2::{Digest as _, Sha256};
6use utoipa::ToSchema;
7
8pub const RUNTIME_OBSERVABILITY_PROTOCOL: &str = "lenso.system-plane.runtime-observability.v1";
9pub const RUNTIME_OBSERVABILITY_PATH: &str = "/system-plane/v1/runtime-observability";
10pub const RUNTIME_OBSERVABILITY_FEATURE_QUEUE_SUMMARY: &str = "queue-summary";
11pub const RUNTIME_OBSERVABILITY_FEATURE_RECOVERY_FEED: &str = "recovery-feed";
12
13#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
14#[serde(rename_all = "snake_case")]
15pub enum RuntimeObservabilityStatus {
16    Healthy,
17    Degraded,
18    Failing,
19}
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
22#[serde(rename_all = "snake_case")]
23pub enum RuntimeQueueKind {
24    Outbox,
25    Functions,
26}
27
28#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
29#[serde(rename_all = "camelCase", deny_unknown_fields)]
30pub struct RuntimeQueueSummary {
31    pub queue: RuntimeQueueKind,
32    pub pending: u64,
33    pub active: u64,
34    pub completed: u64,
35    pub failed: u64,
36    pub dead: u64,
37    #[serde(skip_serializing_if = "Option::is_none")]
38    pub oldest_pending_age_seconds: Option<u64>,
39    #[serde(skip_serializing_if = "Option::is_none")]
40    pub oldest_failed_age_seconds: Option<u64>,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
44#[serde(rename_all = "camelCase", deny_unknown_fields)]
45pub struct RuntimeObservabilitySnapshot {
46    pub protocol: String,
47    #[schema(min_length = 1)]
48    pub service_id: String,
49    #[schema(min_length = 1)]
50    pub service_revision: String,
51    #[schema(min_length = 1)]
52    pub snapshot_revision: String,
53    #[schema(min_length = 1)]
54    pub schema_digest: String,
55    #[schema(min_length = 1)]
56    pub next_cursor: String,
57    #[schemars(with = "String")]
58    pub observed_at: DateTime<Utc>,
59    pub status: RuntimeObservabilityStatus,
60    pub queues: Vec<RuntimeQueueSummary>,
61}
62
63#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
64#[serde(rename_all = "snake_case")]
65pub enum RuntimeObservationChangeKind {
66    Upserted,
67    Deleted,
68}
69
70#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
71#[serde(rename_all = "camelCase", deny_unknown_fields)]
72pub struct RuntimeObservationChange {
73    pub sequence: u64,
74    pub queue: RuntimeQueueKind,
75    pub resource_id: String,
76    pub change_kind: RuntimeObservationChangeKind,
77    #[schemars(with = "String")]
78    pub recorded_at: DateTime<Utc>,
79}
80
81#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
82#[serde(rename_all = "snake_case")]
83pub enum RuntimeObservationContinuity {
84    Continuous,
85    ResetRequired,
86}
87
88#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
89#[serde(rename_all = "snake_case")]
90pub enum RuntimeObservationGapReason {
91    InvalidCursor,
92    ServiceRevisionChanged,
93    SchemaChanged,
94    RetentionLost,
95}
96
97#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
98#[serde(rename_all = "camelCase", deny_unknown_fields)]
99pub struct RuntimeObservationEvidenceGap {
100    pub reason: RuntimeObservationGapReason,
101    pub message: String,
102    pub required_action: String,
103}
104
105#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, ToSchema)]
106#[serde(rename_all = "camelCase", deny_unknown_fields)]
107pub struct RuntimeObservationFeed {
108    pub protocol: String,
109    pub service_id: String,
110    pub service_revision: String,
111    pub schema_digest: String,
112    #[schemars(with = "String")]
113    pub collected_at: DateTime<Utc>,
114    pub continuity: RuntimeObservationContinuity,
115    #[serde(skip_serializing_if = "Option::is_none")]
116    pub evidence_gap: Option<RuntimeObservationEvidenceGap>,
117    pub changes: Vec<RuntimeObservationChange>,
118    pub next_cursor: String,
119    pub has_more: bool,
120}
121
122#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
123#[serde(tag = "kind", content = "document", rename_all = "snake_case")]
124pub enum RuntimeObservabilityMessage {
125    Snapshot(RuntimeObservabilitySnapshot),
126    Feed(RuntimeObservationFeed),
127}
128
129#[must_use]
130pub fn runtime_observability_schema() -> Value {
131    let mut schema = serde_json::to_value(schemars::schema_for!(RuntimeObservabilityMessage))
132        .expect("Runtime Observability schema serializes");
133    schema["$id"] = Value::String(
134        "https://contracts.lenso.local/system-plane/lenso.system-plane.runtime-observability.v1.schema.json"
135            .to_owned(),
136    );
137    schema["title"] = Value::String("Lenso Runtime Observability Messages".to_owned());
138    for definition in ["RuntimeObservabilitySnapshot", "RuntimeObservationFeed"] {
139        schema["$defs"][definition]["properties"]["protocol"] =
140            json!({ "const": RUNTIME_OBSERVABILITY_PROTOCOL });
141    }
142    for field in [
143        "serviceId",
144        "serviceRevision",
145        "snapshotRevision",
146        "schemaDigest",
147        "nextCursor",
148    ] {
149        schema["$defs"]["RuntimeObservabilitySnapshot"]["properties"][field]["minLength"] =
150            json!(1);
151    }
152    schema
153}
154
155#[must_use]
156pub fn runtime_observability_schema_digest() -> String {
157    let bytes = serde_json::to_vec(&runtime_observability_schema())
158        .expect("Runtime Observability schema serializes to bytes");
159    format!("sha256:{}", hex(&Sha256::digest(bytes)))
160}
161
162fn hex(bytes: &[u8]) -> String {
163    bytes.iter().map(|byte| format!("{byte:02x}")).collect()
164}