Skip to main content

platform_system_plane/
module_operations.rs

1//! Managed Service Module inventory, contribution, and descriptor-bound config operations.
2
3use axum::{Extension, Json};
4use lenso_contracts::{
5    AdminSurface, ConsoleActionInputValue, ConsoleContributionAction, ModuleConfigMutability,
6    ModuleDelivery, ModuleHttpMethod, ModuleRelease, digest_json, validate_module_config_value,
7};
8use lenso_service::system_plane::{
9    ActionContributionResolution, ActionContributionResolutionRequest, CapabilityAdvertisement,
10    MODULE_OPERATIONS_FEATURE_CONFIG_READ, MODULE_OPERATIONS_FEATURE_CONFIG_WRITE,
11    MODULE_OPERATIONS_FEATURE_CONTRIBUTIONS_RESOLVE, MODULE_OPERATIONS_FEATURE_INVENTORY_READ,
12    MODULE_OPERATIONS_PATH, MODULE_OPERATIONS_PROTOCOL, ManagedServiceContext,
13    ModuleConfigAuditEvidence, ModuleConfigReadRequest, ModuleConfigReadResponse,
14    ModuleConfigValue, ModuleConfigWriteRequest, ModuleConfigWriteResponse,
15    ModuleInventoryConsoleUi, ModuleInventoryDelivery, ModuleInventoryModule,
16    ModuleInventoryRequest, ModuleInventoryRoute, ModuleInventorySnapshot, ModuleRuntimeStatus,
17    ResolvedActionContribution, module_operations_schema_digest,
18};
19use serde_json::Value;
20use std::{
21    collections::{BTreeMap, BTreeSet},
22    sync::{Arc, RwLock},
23};
24use utoipa_axum::{router::OpenApiRouter, routes};
25
26#[derive(Debug, Clone, PartialEq, Eq)]
27pub struct ModuleOperationsProviderError {
28    pub message: String,
29}
30
31impl std::fmt::Display for ModuleOperationsProviderError {
32    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
33        formatter.write_str(&self.message)
34    }
35}
36
37impl std::error::Error for ModuleOperationsProviderError {}
38
39#[derive(Debug, Clone, Copy, PartialEq, Eq)]
40pub enum ModuleOperationsErrorCode {
41    InvalidRequest,
42    CapabilityDenied,
43    NotFound,
44    Conflict,
45    StoreUnavailable,
46}
47
48#[derive(Debug, Clone, PartialEq, Eq)]
49pub struct ModuleOperationsError {
50    pub code: ModuleOperationsErrorCode,
51    pub message: String,
52}
53
54#[derive(Debug, Default)]
55struct ModuleOperationsState {
56    config: BTreeMap<(String, String), Value>,
57    audit: Vec<ModuleConfigAuditEvidence>,
58    target_revision: String,
59    next_sequence: u64,
60}
61
62#[derive(Debug, Clone)]
63pub struct ModuleOperationsProvider {
64    service_id: String,
65    service_principal: String,
66    service_revision: String,
67    releases: Arc<BTreeMap<String, ModuleRelease>>,
68    statuses: Arc<RwLock<BTreeMap<String, ModuleRuntimeStatus>>>,
69    state: Arc<RwLock<ModuleOperationsState>>,
70}
71
72impl ModuleOperationsProvider {
73    /// Builds the Service-owned Module catalog. Every release is validated
74    /// before it can be advertised or queried through the System Plane.
75    pub fn new(
76        service_id: impl Into<String>,
77        service_principal: impl Into<String>,
78        service_revision: impl Into<String>,
79        releases: impl IntoIterator<Item = ModuleRelease>,
80    ) -> Result<Self, ModuleOperationsProviderError> {
81        let service_id = service_id.into();
82        let service_principal = service_principal.into();
83        let service_revision = service_revision.into();
84        if service_id.trim().is_empty()
85            || service_principal.trim().is_empty()
86            || service_revision.trim().is_empty()
87        {
88            return Err(provider_error(
89                "Module Operations provider requires Service identity, principal, and revision",
90            ));
91        }
92
93        let mut catalog = BTreeMap::new();
94        for release in releases {
95            let issues = release.validate();
96            if !issues.is_empty() {
97                return Err(provider_error(format!(
98                    "Module Release `{}` is invalid: {}",
99                    release.module_id,
100                    issues
101                        .iter()
102                        .map(|issue| format!("{}: {}", issue.path, issue.message))
103                        .collect::<Vec<_>>()
104                        .join("; ")
105                )));
106            }
107            if catalog.insert(release.module_id.clone(), release).is_some() {
108                return Err(provider_error(
109                    "Module Operations catalog contains a duplicate ModuleId",
110                ));
111            }
112        }
113
114        let target_revision = digest_json(&(&service_id, &service_principal, &service_revision))
115            .map_err(|error| {
116                provider_error(format!("could not initialize target revision: {error}"))
117            })?;
118        let mut state = ModuleOperationsState {
119            target_revision,
120            ..ModuleOperationsState::default()
121        };
122        state.next_sequence = 1;
123
124        Ok(Self {
125            service_id,
126            service_principal,
127            service_revision,
128            releases: Arc::new(catalog),
129            statuses: Arc::new(RwLock::new(BTreeMap::new())),
130            state: Arc::new(RwLock::new(state)),
131        })
132    }
133
134    #[must_use]
135    pub fn advertisement() -> CapabilityAdvertisement {
136        CapabilityAdvertisement {
137            contract_id: MODULE_OPERATIONS_PROTOCOL.to_owned(),
138            major_version: 1,
139            feature_ids: BTreeSet::from([
140                MODULE_OPERATIONS_FEATURE_INVENTORY_READ.to_owned(),
141                MODULE_OPERATIONS_FEATURE_CONTRIBUTIONS_RESOLVE.to_owned(),
142                MODULE_OPERATIONS_FEATURE_CONFIG_READ.to_owned(),
143                MODULE_OPERATIONS_FEATURE_CONFIG_WRITE.to_owned(),
144            ]),
145            schema_digest: module_operations_schema_digest(),
146            endpoint: MODULE_OPERATIONS_PATH.to_owned(),
147        }
148    }
149
150    #[must_use]
151    pub fn service_id(&self) -> &str {
152        &self.service_id
153    }
154
155    #[must_use]
156    pub fn service_principal(&self) -> &str {
157        &self.service_principal
158    }
159
160    #[must_use]
161    pub fn service_revision(&self) -> &str {
162        &self.service_revision
163    }
164
165    /// Seeds one descriptor-declared value for integration tests and Service
166    /// bootstrap code. Secret-bearing values are accepted but never returned.
167    pub fn set_config_value(
168        &self,
169        module_id: &str,
170        key: &str,
171        value: Value,
172    ) -> Result<(), ModuleOperationsError> {
173        let field = self.config_field(module_id, key)?;
174        validate_module_config_value(field, &value).map_err(|message| {
175            operation_error(ModuleOperationsErrorCode::InvalidRequest, message)
176        })?;
177        let mut state = self.state.write().map_err(|_| {
178            operation_error(
179                ModuleOperationsErrorCode::StoreUnavailable,
180                "Module configuration store is unavailable",
181            )
182        })?;
183        state
184            .config
185            .insert((module_id.to_owned(), key.to_owned()), value);
186        Ok(())
187    }
188
189    pub fn set_runtime_status(
190        &self,
191        module_id: &str,
192        status: ModuleRuntimeStatus,
193    ) -> Result<(), ModuleOperationsError> {
194        if !self.releases.contains_key(module_id) {
195            return Err(operation_error(
196                ModuleOperationsErrorCode::NotFound,
197                format!("Module `{module_id}` is not installed"),
198            ));
199        }
200        self.statuses
201            .write()
202            .map_err(|_| {
203                operation_error(
204                    ModuleOperationsErrorCode::StoreUnavailable,
205                    "Module status store is unavailable",
206                )
207            })?
208            .insert(module_id.to_owned(), status);
209        Ok(())
210    }
211
212    pub fn inventory(
213        &self,
214        context: &ManagedServiceContext,
215    ) -> Result<ModuleInventorySnapshot, ModuleOperationsError> {
216        self.validate_context(context)?;
217        let statuses = self.statuses.read().map_err(|_| {
218            operation_error(
219                ModuleOperationsErrorCode::StoreUnavailable,
220                "Module status store is unavailable",
221            )
222        })?;
223        let modules = self
224            .releases
225            .values()
226            .map(|release| inventory_module(release, statuses.get(&release.module_id).copied()))
227            .collect::<Result<Vec<_>, _>>()?;
228        let snapshot_revision =
229            digest_json(&(&self.service_revision, &modules)).map_err(|error| {
230                operation_error(
231                    ModuleOperationsErrorCode::StoreUnavailable,
232                    error.to_string(),
233                )
234            })?;
235        Ok(ModuleInventorySnapshot {
236            protocol: MODULE_OPERATIONS_PROTOCOL.to_owned(),
237            context: context.clone(),
238            service_revision: self.service_revision.clone(),
239            snapshot_revision,
240            schema_digest: module_operations_schema_digest(),
241            modules,
242        })
243    }
244
245    pub fn resolve_contributions(
246        &self,
247        request: &ActionContributionResolutionRequest,
248    ) -> Result<ActionContributionResolution, ModuleOperationsError> {
249        self.validate_context(&request.context)?;
250        let slot = self
251            .releases
252            .values()
253            .flat_map(|release| release.manifest.console_slots.iter())
254            .find(|slot| slot.id == request.slot && slot.version == request.slot_version)
255            .ok_or_else(|| {
256                operation_error(
257                    ModuleOperationsErrorCode::NotFound,
258                    format!(
259                        "Console contribution slot `{}` v{} is not installed",
260                        request.slot, request.slot_version
261                    ),
262                )
263            })?;
264        validate_slot_context(slot, &request.slot_context)?;
265
266        let mut contributions = Vec::new();
267        for release in self.releases.values() {
268            for contribution in &release.manifest.console_contributions {
269                if contribution.target != request.slot
270                    || contribution.target_version != request.slot_version
271                {
272                    continue;
273                }
274                let action = validate_action_reference(
275                    &self.releases,
276                    &contribution.action,
277                    &request.slot_context,
278                    &request.context.capabilities,
279                )?;
280                let mut required_capabilities = contribution.required_capabilities.clone();
281                required_capabilities.sort();
282                required_capabilities.dedup();
283                require_context_capabilities(
284                    &request.context.capabilities,
285                    &required_capabilities,
286                )?;
287                contributions.push(ResolvedActionContribution {
288                    contributing_module_id: release.module_id.clone(),
289                    target: contribution.target.clone(),
290                    target_version: contribution.target_version,
291                    label: contribution.label.clone(),
292                    action,
293                    icon: contribution.icon.clone(),
294                    required_capabilities,
295                });
296            }
297        }
298        contributions.sort_by(|left, right| {
299            (&left.contributing_module_id, &left.label, &left.target).cmp(&(
300                &right.contributing_module_id,
301                &right.label,
302                &right.target,
303            ))
304        });
305        Ok(ActionContributionResolution {
306            protocol: MODULE_OPERATIONS_PROTOCOL.to_owned(),
307            context: request.context.clone(),
308            slot: request.slot.clone(),
309            slot_version: request.slot_version,
310            contributions,
311        })
312    }
313
314    pub fn read_config(
315        &self,
316        request: &ModuleConfigReadRequest,
317    ) -> Result<ModuleConfigReadResponse, ModuleOperationsError> {
318        self.validate_context(&request.context)?;
319        self.require_config_namespace(&request.context, &request.module_id)?;
320        let release = self.release(&request.module_id)?;
321        let fields = requested_config_fields(&release.manifest.config.fields, &request.keys)?;
322        let state = self.state.read().map_err(|_| {
323            operation_error(
324                ModuleOperationsErrorCode::StoreUnavailable,
325                "Module configuration store is unavailable",
326            )
327        })?;
328        let values = fields
329            .into_iter()
330            .map(|field| {
331                if let Some(capability) = field.read_capability.as_deref() {
332                    require_context_capabilities(&request.context.capabilities, [capability])?;
333                }
334                let present = state
335                    .config
336                    .contains_key(&(request.module_id.clone(), field.key.clone()));
337                Ok(ModuleConfigValue {
338                    key: field.key.clone(),
339                    field_type: field.field_type,
340                    scope: field.scope,
341                    mutability: field.mutability,
342                    activation: field.activation,
343                    sensitive: field.sensitive,
344                    present,
345                    value: (!field.sensitive)
346                        .then(|| {
347                            state
348                                .config
349                                .get(&(request.module_id.clone(), field.key.clone()))
350                                .cloned()
351                        })
352                        .flatten(),
353                })
354            })
355            .collect::<Result<Vec<_>, ModuleOperationsError>>()?;
356        Ok(ModuleConfigReadResponse {
357            protocol: MODULE_OPERATIONS_PROTOCOL.to_owned(),
358            context: request.context.clone(),
359            module_id: request.module_id.clone(),
360            values,
361        })
362    }
363
364    pub fn write_config(
365        &self,
366        request: &ModuleConfigWriteRequest,
367    ) -> Result<ModuleConfigWriteResponse, ModuleOperationsError> {
368        self.validate_context(&request.context)?;
369        self.require_config_namespace(&request.context, &request.module_id)?;
370        if request.values.is_empty() {
371            return Err(operation_error(
372                ModuleOperationsErrorCode::InvalidRequest,
373                "Module configuration writes require at least one field",
374            ));
375        }
376        let release = self.release(&request.module_id)?;
377        let mut fields = BTreeMap::new();
378        for field in &release.manifest.config.fields {
379            fields.insert(field.key.as_str(), field);
380        }
381        let mut keys = BTreeSet::new();
382        for item in &request.values {
383            if !keys.insert(item.key.as_str()) {
384                return Err(operation_error(
385                    ModuleOperationsErrorCode::InvalidRequest,
386                    format!(
387                        "Configuration key `{}` was written more than once",
388                        item.key
389                    ),
390                ));
391            }
392            let field = fields.get(item.key.as_str()).ok_or_else(|| {
393                operation_error(
394                    ModuleOperationsErrorCode::NotFound,
395                    format!(
396                        "Configuration field `{}` is not declared by Module `{}`",
397                        item.key, request.module_id
398                    ),
399                )
400            })?;
401            if field.mutability == ModuleConfigMutability::Static {
402                return Err(operation_error(
403                    ModuleOperationsErrorCode::Conflict,
404                    format!(
405                        "Configuration field `{}` is static and cannot be changed at runtime",
406                        item.key
407                    ),
408                ));
409            }
410            if let Some(capability) = field.write_capability.as_deref() {
411                require_context_capabilities(&request.context.capabilities, [capability])?;
412            }
413            validate_module_config_value(field, &item.value).map_err(|message| {
414                operation_error(ModuleOperationsErrorCode::InvalidRequest, message)
415            })?;
416        }
417
418        let operation_id = digest_json(request).map_err(|error| {
419            operation_error(
420                ModuleOperationsErrorCode::StoreUnavailable,
421                error.to_string(),
422            )
423        })?;
424        let mut state = self.state.write().map_err(|_| {
425            operation_error(
426                ModuleOperationsErrorCode::StoreUnavailable,
427                "Module configuration store is unavailable",
428            )
429        })?;
430        let target_revision_before = state.target_revision.clone();
431        let mut evidence = Vec::with_capacity(request.values.len());
432        for item in &request.values {
433            let old_value_digest = state
434                .config
435                .get(&(request.module_id.clone(), item.key.clone()))
436                .map(digest_value)
437                .transpose()
438                .map_err(|error| {
439                    operation_error(
440                        ModuleOperationsErrorCode::StoreUnavailable,
441                        error.to_string(),
442                    )
443                })?;
444            state.config.insert(
445                (request.module_id.clone(), item.key.clone()),
446                item.value.clone(),
447            );
448            let sequence = state.next_sequence;
449            state.next_sequence = state.next_sequence.saturating_add(1);
450            let new_value_digest = digest_value(&item.value).map_err(|error| {
451                operation_error(
452                    ModuleOperationsErrorCode::StoreUnavailable,
453                    error.to_string(),
454                )
455            })?;
456            let item_evidence = ModuleConfigAuditEvidence {
457                sequence,
458                operation_id: operation_id.clone(),
459                module_id: request.module_id.clone(),
460                key: item.key.clone(),
461                sensitive: fields[item.key.as_str()].sensitive,
462                old_value_digest,
463                new_value_digest,
464                recorded_at_unix_ms: now_unix_ms(),
465            };
466            state.audit.push(item_evidence.clone());
467            evidence.push(item_evidence);
468        }
469        state.target_revision = digest_json(&(
470            &target_revision_before,
471            &operation_id,
472            &request.module_id,
473            &evidence
474                .iter()
475                .map(|item| item.new_value_digest.clone())
476                .collect::<Vec<_>>(),
477        ))
478        .map_err(|error| {
479            operation_error(
480                ModuleOperationsErrorCode::StoreUnavailable,
481                error.to_string(),
482            )
483        })?;
484        Ok(ModuleConfigWriteResponse {
485            protocol: MODULE_OPERATIONS_PROTOCOL.to_owned(),
486            operation_id,
487            context: request.context.clone(),
488            module_id: request.module_id.clone(),
489            target_revision_before,
490            target_revision_after: state.target_revision.clone(),
491            authorization_digest: request.context.digest().map_err(|error| {
492                operation_error(
493                    ModuleOperationsErrorCode::StoreUnavailable,
494                    error.to_string(),
495                )
496            })?,
497            evidence,
498        })
499    }
500
501    fn release(&self, module_id: &str) -> Result<&ModuleRelease, ModuleOperationsError> {
502        self.releases.get(module_id).ok_or_else(|| {
503            operation_error(
504                ModuleOperationsErrorCode::NotFound,
505                format!("Module `{module_id}` is not installed"),
506            )
507        })
508    }
509
510    fn config_field(
511        &self,
512        module_id: &str,
513        key: &str,
514    ) -> Result<&lenso_contracts::ModuleConfigField, ModuleOperationsError> {
515        self.release(module_id)?
516            .manifest
517            .config
518            .fields
519            .iter()
520            .find(|field| field.key == key)
521            .ok_or_else(|| {
522                operation_error(
523                    ModuleOperationsErrorCode::NotFound,
524                    format!("Configuration field `{key}` is not declared by Module `{module_id}`"),
525                )
526            })
527    }
528
529    fn validate_context(
530        &self,
531        context: &ManagedServiceContext,
532    ) -> Result<(), ModuleOperationsError> {
533        if context.service_id != self.service_id
534            || context.target_service_principal != self.service_principal
535            || context.system_id.trim().is_empty()
536            || context.environment_id.trim().is_empty()
537            || context.caller_module_id.trim().is_empty()
538            || context.delegated_actor_subject.trim().is_empty()
539            || !canonical_digest(&context.delegated_authority_digest)
540        {
541            return Err(operation_error(
542                ModuleOperationsErrorCode::InvalidRequest,
543                "Managed Service context does not identify this Service or its delegated authority",
544            ));
545        }
546        Ok(())
547    }
548
549    fn require_config_namespace(
550        &self,
551        context: &ManagedServiceContext,
552        module_id: &str,
553    ) -> Result<(), ModuleOperationsError> {
554        if context.caller_module_id != module_id {
555            return Err(operation_error(
556                ModuleOperationsErrorCode::CapabilityDenied,
557                format!(
558                    "Module configuration is restricted to the calling Module namespace `{}`",
559                    context.caller_module_id
560                ),
561            ));
562        }
563        Ok(())
564    }
565}
566
567#[must_use]
568pub fn module_operations_router<S>(
569    provider: Option<Arc<ModuleOperationsProvider>>,
570) -> OpenApiRouter<S>
571where
572    S: Clone + Send + Sync + 'static,
573{
574    OpenApiRouter::new()
575        .routes(routes!(module_inventory))
576        .routes(routes!(resolve_action_contributions))
577        .routes(routes!(read_module_config))
578        .routes(routes!(write_module_config))
579        .layer(Extension(provider))
580}
581
582#[utoipa::path(
583    post,
584    path = "/system-plane/v1/modules",
585    request_body = ModuleInventoryRequest,
586    responses(
587        (status = 200, body = ModuleInventorySnapshot),
588        (status = 400, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
589        (status = 401, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
590        (status = 403, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
591        (status = 503, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json")
592    ),
593    security(("bearer_auth" = [])),
594    tag = "system-plane-module-operations"
595)]
596async fn module_inventory(
597    caller: crate::AuthorizedSystemPlaneCaller,
598    Extension(provider): Extension<Option<Arc<ModuleOperationsProvider>>>,
599    Json(request): Json<ModuleInventoryRequest>,
600) -> Result<Json<ModuleInventorySnapshot>, crate::SystemPlaneRejection> {
601    let provider = require_provider(provider)?;
602    authorize_context(&caller, &request.context)?;
603    caller.require_capability(
604        MODULE_OPERATIONS_PROTOCOL,
605        &module_operations_schema_digest(),
606        [MODULE_OPERATIONS_FEATURE_INVENTORY_READ],
607    )?;
608    provider
609        .inventory(&request.context)
610        .map(Json)
611        .map_err(rejection)
612}
613
614#[utoipa::path(
615    post,
616    path = "/system-plane/v1/modules/action-contributions/resolve",
617    request_body = ActionContributionResolutionRequest,
618    responses(
619        (status = 200, body = ActionContributionResolution),
620        (status = 400, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
621        (status = 401, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
622        (status = 403, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
623        (status = 404, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
624        (status = 503, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json")
625    ),
626    security(("bearer_auth" = [])),
627    tag = "system-plane-module-operations"
628)]
629async fn resolve_action_contributions(
630    caller: crate::AuthorizedSystemPlaneCaller,
631    Extension(provider): Extension<Option<Arc<ModuleOperationsProvider>>>,
632    Json(request): Json<ActionContributionResolutionRequest>,
633) -> Result<Json<ActionContributionResolution>, crate::SystemPlaneRejection> {
634    let provider = require_provider(provider)?;
635    authorize_context(&caller, &request.context)?;
636    caller.require_capability(
637        MODULE_OPERATIONS_PROTOCOL,
638        &module_operations_schema_digest(),
639        [MODULE_OPERATIONS_FEATURE_CONTRIBUTIONS_RESOLVE],
640    )?;
641    provider
642        .resolve_contributions(&request)
643        .map(Json)
644        .map_err(rejection)
645}
646
647#[utoipa::path(
648    post,
649    path = "/system-plane/v1/modules/config/read",
650    request_body = ModuleConfigReadRequest,
651    responses(
652        (status = 200, body = ModuleConfigReadResponse),
653        (status = 400, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
654        (status = 401, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
655        (status = 403, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
656        (status = 404, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
657        (status = 503, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json")
658    ),
659    security(("bearer_auth" = [])),
660    tag = "system-plane-module-operations"
661)]
662async fn read_module_config(
663    caller: crate::AuthorizedSystemPlaneCaller,
664    Extension(provider): Extension<Option<Arc<ModuleOperationsProvider>>>,
665    Json(request): Json<ModuleConfigReadRequest>,
666) -> Result<Json<ModuleConfigReadResponse>, crate::SystemPlaneRejection> {
667    let provider = require_provider(provider)?;
668    authorize_context(&caller, &request.context)?;
669    caller.require_capability(
670        MODULE_OPERATIONS_PROTOCOL,
671        &module_operations_schema_digest(),
672        [MODULE_OPERATIONS_FEATURE_CONFIG_READ],
673    )?;
674    provider.read_config(&request).map(Json).map_err(rejection)
675}
676
677#[utoipa::path(
678    post,
679    path = "/system-plane/v1/modules/config/write",
680    request_body = ModuleConfigWriteRequest,
681    responses(
682        (status = 200, body = ModuleConfigWriteResponse),
683        (status = 400, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
684        (status = 401, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
685        (status = 403, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
686        (status = 404, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
687        (status = 409, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json"),
688        (status = 503, body = crate::SystemPlaneErrorBody, content_type = "application/problem+json")
689    ),
690    security(("bearer_auth" = [])),
691    tag = "system-plane-module-operations"
692)]
693async fn write_module_config(
694    caller: crate::AuthorizedSystemPlaneCaller,
695    Extension(provider): Extension<Option<Arc<ModuleOperationsProvider>>>,
696    Json(request): Json<ModuleConfigWriteRequest>,
697) -> Result<Json<ModuleConfigWriteResponse>, crate::SystemPlaneRejection> {
698    let provider = require_provider(provider)?;
699    authorize_context(&caller, &request.context)?;
700    caller.require_capability(
701        MODULE_OPERATIONS_PROTOCOL,
702        &module_operations_schema_digest(),
703        [MODULE_OPERATIONS_FEATURE_CONFIG_WRITE],
704    )?;
705    provider.write_config(&request).map(Json).map_err(rejection)
706}
707
708fn require_provider(
709    provider: Option<Arc<ModuleOperationsProvider>>,
710) -> Result<Arc<ModuleOperationsProvider>, crate::SystemPlaneRejection> {
711    provider.ok_or_else(|| {
712        crate::SystemPlaneRejection::unavailable(
713            "module_operations_unavailable",
714            "Module Operations capability is not configured for this Service",
715            "configure_module_operations",
716        )
717    })
718}
719
720fn authorize_context(
721    caller: &crate::AuthorizedSystemPlaneCaller,
722    context: &ManagedServiceContext,
723) -> Result<(), crate::SystemPlaneRejection> {
724    let core = caller.runtime.registry.document();
725    if context.system_id != caller.enrollment.system_id
726        || context.service_id != core.service_id
727        || context.target_service_principal != core.service_principal
728    {
729        return Err(crate::SystemPlaneRejection::new(
730            axum::http::StatusCode::FORBIDDEN,
731            "system_plane_target_context_mismatch",
732            "Managed Service context does not match the authenticated enrollment target",
733            "use_the_enrolled_service_context",
734        ));
735    }
736    if caller.enrollment.system_id != "system-sandbox" {
737        let granted = caller
738            .enrollment
739            .capabilities
740            .iter()
741            .flat_map(|capability| capability.feature_ids.iter())
742            .collect::<BTreeSet<_>>();
743        if context
744            .capabilities
745            .iter()
746            .any(|capability| !granted.contains(capability))
747        {
748            return Err(crate::SystemPlaneRejection::new(
749                axum::http::StatusCode::FORBIDDEN,
750                "system_plane_context_capability_not_granted",
751                "Managed Service context contains a capability outside the active enrollment grant",
752                "request_the_required_module_capability",
753            ));
754        }
755    }
756    Ok(())
757}
758
759fn rejection(error: ModuleOperationsError) -> crate::SystemPlaneRejection {
760    let (status, code, next_action) = match error.code {
761        ModuleOperationsErrorCode::InvalidRequest => (
762            axum::http::StatusCode::BAD_REQUEST,
763            "module_operations_invalid_request",
764            "send_a_descriptor_bound_module_request",
765        ),
766        ModuleOperationsErrorCode::CapabilityDenied => (
767            axum::http::StatusCode::FORBIDDEN,
768            "module_operations_capability_denied",
769            "request_the_required_module_capability",
770        ),
771        ModuleOperationsErrorCode::NotFound => (
772            axum::http::StatusCode::NOT_FOUND,
773            "module_operations_not_found",
774            "refresh_module_inventory",
775        ),
776        ModuleOperationsErrorCode::Conflict => (
777            axum::http::StatusCode::CONFLICT,
778            "module_operations_conflict",
779            "refresh_module_inventory_and_retry",
780        ),
781        ModuleOperationsErrorCode::StoreUnavailable => (
782            axum::http::StatusCode::SERVICE_UNAVAILABLE,
783            "module_operations_unavailable",
784            "restore_module_operations_store",
785        ),
786    };
787    crate::SystemPlaneRejection::new(status, code, error.message, next_action)
788}
789
790fn operation_error(
791    code: ModuleOperationsErrorCode,
792    message: impl Into<String>,
793) -> ModuleOperationsError {
794    ModuleOperationsError {
795        code,
796        message: message.into(),
797    }
798}
799
800fn provider_error(message: impl Into<String>) -> ModuleOperationsProviderError {
801    ModuleOperationsProviderError {
802        message: message.into(),
803    }
804}
805
806fn inventory_module(
807    release: &ModuleRelease,
808    status: Option<ModuleRuntimeStatus>,
809) -> Result<ModuleInventoryModule, ModuleOperationsError> {
810    let release_digest = digest_json(release).map_err(|error| {
811        operation_error(
812            ModuleOperationsErrorCode::StoreUnavailable,
813            error.to_string(),
814        )
815    })?;
816    let delivery = match release.delivery {
817        ModuleDelivery::Linked(_) => ModuleInventoryDelivery::Linked,
818        ModuleDelivery::Service(_) => ModuleInventoryDelivery::Service,
819    };
820    let mut routes = release
821        .manifest
822        .http_routes
823        .iter()
824        .map(|route| ModuleInventoryRoute {
825            method: match route.method {
826                ModuleHttpMethod::Get => "GET".to_owned(),
827                ModuleHttpMethod::Post => "POST".to_owned(),
828                ModuleHttpMethod::Put => "PUT".to_owned(),
829                ModuleHttpMethod::Patch => "PATCH".to_owned(),
830                ModuleHttpMethod::Delete => "DELETE".to_owned(),
831                _ => "UNKNOWN".to_owned(),
832            },
833            path: route.path.clone(),
834            capability: route.capability.clone(),
835        })
836        .collect::<Vec<_>>();
837    routes.sort_by(|left, right| (&left.method, &left.path).cmp(&(&right.method, &right.path)));
838    let mut dependency_module_ids = release
839        .manifest
840        .requires
841        .iter()
842        .map(|requirement| requirement.module_id.clone())
843        .collect::<Vec<_>>();
844    dependency_module_ids.sort();
845    let runtime_functions = release
846        .manifest
847        .runtime
848        .as_ref()
849        .map(|runtime| {
850            let mut functions = runtime
851                .functions
852                .iter()
853                .map(|function| function.name.clone())
854                .collect::<Vec<_>>();
855            functions.sort();
856            functions
857        })
858        .unwrap_or_default();
859    let console_ui =
860        release
861            .console_ui_artifact
862            .as_ref()
863            .map(|artifact| ModuleInventoryConsoleUi {
864                format: lenso_contracts::CONSOLE_UI_ESM_FORMAT.to_owned(),
865                protocol_major: artifact.protocol_major,
866                artifact_digest: artifact.artifact.digest.clone(),
867                entry: artifact.entry.clone(),
868                style_assets: artifact
869                    .style_assets
870                    .iter()
871                    .map(|asset| asset.path.clone())
872                    .collect(),
873            });
874    Ok(ModuleInventoryModule {
875        module_id: release.module_id.clone(),
876        version: release.version.clone(),
877        release_digest,
878        manifest_digest: release.manifest_digest.clone(),
879        delivery,
880        dependency_module_ids,
881        routes,
882        runtime_functions,
883        runtime_status: status.unwrap_or(ModuleRuntimeStatus::Active),
884        console_ui,
885    })
886}
887
888fn requested_config_fields<'a>(
889    fields: &'a [lenso_contracts::ModuleConfigField],
890    requested: &[String],
891) -> Result<Vec<&'a lenso_contracts::ModuleConfigField>, ModuleOperationsError> {
892    let keys = if requested.is_empty() {
893        fields
894            .iter()
895            .map(|field| field.key.as_str())
896            .collect::<Vec<_>>()
897    } else {
898        let mut keys = requested.iter().map(String::as_str).collect::<Vec<_>>();
899        keys.sort_unstable();
900        if keys.windows(2).any(|pair| pair[0] == pair[1]) {
901            return Err(operation_error(
902                ModuleOperationsErrorCode::InvalidRequest,
903                "Configuration read keys must be unique",
904            ));
905        }
906        keys
907    };
908    keys.into_iter()
909        .map(|key| {
910            fields.iter().find(|field| field.key == key).ok_or_else(|| {
911                operation_error(
912                    ModuleOperationsErrorCode::NotFound,
913                    format!("Configuration field `{key}` is not declared by the Module"),
914                )
915            })
916        })
917        .collect()
918}
919
920fn validate_action_reference(
921    releases: &BTreeMap<String, ModuleRelease>,
922    action: &ConsoleContributionAction,
923    slot_context: &BTreeMap<String, Value>,
924    context_capabilities: &BTreeSet<String>,
925) -> Result<ConsoleContributionAction, ModuleOperationsError> {
926    match action {
927        ConsoleContributionAction::AdminAction {
928            module,
929            name,
930            input_bindings,
931        } => {
932            let release = releases.get(module).ok_or_else(|| {
933                operation_error(
934                    ModuleOperationsErrorCode::NotFound,
935                    format!("Contributed Admin Module `{module}` is not installed"),
936                )
937            })?;
938            let Some(AdminSurface::DeclarativeCustom(surface)) = release.manifest.admin.as_ref()
939            else {
940                return Err(operation_error(
941                    ModuleOperationsErrorCode::InvalidRequest,
942                    format!(
943                        "Contributed Admin Action `{module}:{name}` is not declaratively declared"
944                    ),
945                ));
946            };
947            let admin_action = surface
948                .actions
949                .iter()
950                .find(|candidate| candidate.name == *name)
951                .ok_or_else(|| {
952                    operation_error(
953                        ModuleOperationsErrorCode::NotFound,
954                        format!("Contributed Admin Action `{module}:{name}` is not declared"),
955                    )
956                })?;
957            require_context_capabilities(context_capabilities, [&admin_action.capability])?;
958            for binding in input_bindings {
959                let ConsoleActionInputValue::SlotContext { path } = &binding.value else {
960                    return Err(operation_error(
961                        ModuleOperationsErrorCode::InvalidRequest,
962                        "Only explicit slot context action bindings are supported",
963                    ));
964                };
965                let value = slot_context.get(path).or_else(|| {
966                    path.split('.')
967                        .next()
968                        .and_then(|head| slot_context.get(head))
969                });
970                if value.is_none() {
971                    return Err(operation_error(
972                        ModuleOperationsErrorCode::InvalidRequest,
973                        format!(
974                            "Action input binding `{path}` is not present in the explicit slot context"
975                        ),
976                    ));
977                }
978            }
979            Ok(action.clone())
980        }
981        _ => Err(operation_error(
982            ModuleOperationsErrorCode::InvalidRequest,
983            "Unsupported Console contribution action kind",
984        )),
985    }
986}
987
988fn validate_slot_context(
989    slot: &lenso_contracts::ConsoleSlot,
990    values: &BTreeMap<String, Value>,
991) -> Result<(), ModuleOperationsError> {
992    for context in &slot.context {
993        for field in &context.fields {
994            let key = format!("{}.{}", context.name, field.name);
995            let value = values.get(&key).or_else(|| values.get(&field.name));
996            if field.required && value.is_none() {
997                return Err(operation_error(
998                    ModuleOperationsErrorCode::InvalidRequest,
999                    format!("Required slot context field `{key}` is missing"),
1000                ));
1001            }
1002            if let Some(value) = value {
1003                let valid = match field.field_type {
1004                    lenso_contracts::ConsoleSlotContextFieldType::String => value.is_string(),
1005                    lenso_contracts::ConsoleSlotContextFieldType::Boolean => value.is_boolean(),
1006                    lenso_contracts::ConsoleSlotContextFieldType::Number => value.is_number(),
1007                    lenso_contracts::ConsoleSlotContextFieldType::Timestamp => value.is_string(),
1008                    _ => false,
1009                };
1010                if !valid {
1011                    return Err(operation_error(
1012                        ModuleOperationsErrorCode::InvalidRequest,
1013                        format!("Slot context field `{key}` has the wrong type"),
1014                    ));
1015                }
1016            }
1017        }
1018    }
1019    Ok(())
1020}
1021
1022fn require_context_capabilities(
1023    context: &BTreeSet<String>,
1024    required: impl IntoIterator<Item = impl AsRef<str>>,
1025) -> Result<(), ModuleOperationsError> {
1026    for capability in required {
1027        if !context.contains(capability.as_ref()) {
1028            return Err(operation_error(
1029                ModuleOperationsErrorCode::CapabilityDenied,
1030                format!(
1031                    "Required Module capability `{}` was not delegated",
1032                    capability.as_ref()
1033                ),
1034            ));
1035        }
1036    }
1037    Ok(())
1038}
1039
1040fn digest_value(value: &Value) -> serde_json::Result<String> {
1041    digest_json(value)
1042}
1043
1044fn canonical_digest(value: &str) -> bool {
1045    value.strip_prefix("sha256:").is_some_and(|digest| {
1046        digest.len() == 64
1047            && digest
1048                .bytes()
1049                .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte))
1050    })
1051}
1052
1053fn now_unix_ms() -> u64 {
1054    std::time::SystemTime::now()
1055        .duration_since(std::time::UNIX_EPOCH)
1056        .map_or(0, |duration| {
1057            u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
1058        })
1059}
1060
1061#[cfg(test)]
1062mod tests {
1063    use super::*;
1064    use lenso_contracts::{
1065        AdminAction, AdminActionDangerLevel, AdminDeclarativeSurface, ConsoleActionInputBinding,
1066        ConsoleActionInputValue, ConsoleContribution, ConsoleContributionAction,
1067        ConsoleContributionKind, ConsoleSlot, ConsoleSlotContext, ConsoleSlotContextField,
1068        ConsoleSlotContextFieldType, ModuleConfigActivation, ModuleConfigField,
1069        ModuleConfigFieldType, ModuleConfigMutability, ModuleConfigScope, ModuleRelease,
1070        digest_json,
1071    };
1072
1073    fn release() -> ModuleRelease {
1074        serde_json::from_value(
1075            lenso_contracts::console_contract_vectors()["positive"]["release"].clone(),
1076        )
1077        .expect("positive Console vector should produce a Module Release")
1078    }
1079
1080    fn context(capabilities: impl IntoIterator<Item = &'static str>) -> ManagedServiceContext {
1081        ManagedServiceContext::new(
1082            "system-1",
1083            "service-1",
1084            "production",
1085            "spiffe://lenso/service-1",
1086            "acme/support-console",
1087            "operator-1",
1088            format!("sha256:{}", "f".repeat(64)),
1089            capabilities,
1090        )
1091    }
1092
1093    fn provider() -> ModuleOperationsProvider {
1094        ModuleOperationsProvider::new(
1095            "service-1",
1096            "spiffe://lenso/service-1",
1097            "revision-1",
1098            [release()],
1099        )
1100        .expect("positive Module Release should be accepted")
1101    }
1102
1103    #[test]
1104    fn inventory_is_authoritative_and_reports_esm_delivery() {
1105        let provider = provider();
1106        let snapshot = provider
1107            .inventory(&context([MODULE_OPERATIONS_FEATURE_INVENTORY_READ]))
1108            .unwrap();
1109        assert_eq!(snapshot.protocol, MODULE_OPERATIONS_PROTOCOL);
1110        assert_eq!(snapshot.modules.len(), 1);
1111        assert_eq!(snapshot.modules[0].module_id, "acme/support-console");
1112        assert_eq!(
1113            snapshot.modules[0].console_ui.as_ref().unwrap().format,
1114            "console_ui_esm"
1115        );
1116        assert_eq!(
1117            snapshot.modules[0].runtime_status,
1118            ModuleRuntimeStatus::Active
1119        );
1120    }
1121
1122    #[test]
1123    fn configuration_is_typed_capability_checked_namespace_bound_and_audited() {
1124        let provider = provider();
1125        provider
1126            .set_config_value(
1127                "acme/support-console",
1128                "endpoint",
1129                Value::String("https://old.example".to_owned()),
1130            )
1131            .unwrap();
1132        let request = ModuleConfigReadRequest {
1133            context: context(["support.endpoint.read"]),
1134            module_id: "acme/support-console".to_owned(),
1135            keys: vec!["endpoint".to_owned()],
1136        };
1137        assert_eq!(
1138            provider.read_config(&request).unwrap().values[0].value,
1139            Some(Value::String("https://old.example".to_owned()))
1140        );
1141
1142        let response = provider
1143            .write_config(&ModuleConfigWriteRequest {
1144                context: context(["support.endpoint.write"]),
1145                module_id: "acme/support-console".to_owned(),
1146                values: vec![lenso_service::system_plane::ModuleConfigWriteValue {
1147                    key: "endpoint".to_owned(),
1148                    value: Value::String("https://new.example".to_owned()),
1149                }],
1150            })
1151            .unwrap();
1152        assert_eq!(response.evidence.len(), 1);
1153        assert_eq!(response.evidence[0].key, "endpoint");
1154        assert!(!response.evidence[0].new_value_digest.is_empty());
1155
1156        let wrong_type = provider.write_config(&ModuleConfigWriteRequest {
1157            context: context(["support.endpoint.write"]),
1158            module_id: "acme/support-console".to_owned(),
1159            values: vec![lenso_service::system_plane::ModuleConfigWriteValue {
1160                key: "endpoint".to_owned(),
1161                value: Value::Bool(true),
1162            }],
1163        });
1164        assert_eq!(
1165            wrong_type.unwrap_err().code,
1166            ModuleOperationsErrorCode::InvalidRequest
1167        );
1168
1169        let wrong_namespace = provider.read_config(&ModuleConfigReadRequest {
1170            context: ManagedServiceContext {
1171                caller_module_id: "other/module".to_owned(),
1172                ..context(["support.endpoint.read"])
1173            },
1174            module_id: "acme/support-console".to_owned(),
1175            keys: Vec::new(),
1176        });
1177        assert_eq!(
1178            wrong_namespace.unwrap_err().code,
1179            ModuleOperationsErrorCode::CapabilityDenied
1180        );
1181    }
1182
1183    #[test]
1184    fn contributions_are_data_only_and_require_current_action_capability() {
1185        let mut release = release();
1186        release
1187            .manifest
1188            .capabilities
1189            .push("support.action.execute".to_owned());
1190        release.manifest.capabilities.sort();
1191        release.manifest.console_slots = vec![ConsoleSlot {
1192            id: "support.detail.actions".to_owned(),
1193            version: 1,
1194            label: "Ticket actions".to_owned(),
1195            accepts: vec![ConsoleContributionKind::AdminAction],
1196            context: vec![ConsoleSlotContext {
1197                name: "ticket".to_owned(),
1198                fields: vec![ConsoleSlotContextField {
1199                    name: "id".to_owned(),
1200                    field_type: ConsoleSlotContextFieldType::String,
1201                    required: true,
1202                }],
1203            }],
1204        }];
1205        release.manifest.admin = Some(AdminSurface::DeclarativeCustom(AdminDeclarativeSurface {
1206            pages: Vec::new(),
1207            actions: vec![AdminAction {
1208                name: "reopen".to_owned(),
1209                label: "Reopen".to_owned(),
1210                capability: "support.action.execute".to_owned(),
1211                input_schema: None,
1212                confirmation: None,
1213                danger_level: AdminActionDangerLevel::Low,
1214                operation: None,
1215            }],
1216            fallback_schema: None,
1217        }));
1218        release.manifest.console_contributions = vec![ConsoleContribution {
1219            target: "support.detail.actions".to_owned(),
1220            target_version: 1,
1221            label: "Reopen".to_owned(),
1222            action: ConsoleContributionAction::AdminAction {
1223                module: release.module_id.clone(),
1224                name: "reopen".to_owned(),
1225                input_bindings: vec![ConsoleActionInputBinding {
1226                    input: "ticket_id".to_owned(),
1227                    value: ConsoleActionInputValue::SlotContext {
1228                        path: "ticket.id".to_owned(),
1229                    },
1230                }],
1231            },
1232            icon: None,
1233            required_capabilities: Vec::new(),
1234        }];
1235        release.manifest_digest = digest_json(&release.manifest).unwrap();
1236        assert!(
1237            release.validate().is_empty(),
1238            "fixture release should remain valid"
1239        );
1240        let provider = ModuleOperationsProvider::new(
1241            "service-1",
1242            "spiffe://lenso/service-1",
1243            "revision-1",
1244            [release],
1245        )
1246        .unwrap();
1247
1248        let request = ActionContributionResolutionRequest {
1249            context: context(["support.action.execute"]),
1250            slot: "support.detail.actions".to_owned(),
1251            slot_version: 1,
1252            slot_context: BTreeMap::from([(
1253                "ticket.id".to_owned(),
1254                Value::String("t-1".to_owned()),
1255            )]),
1256        };
1257        let result = provider.resolve_contributions(&request).unwrap();
1258        assert_eq!(result.contributions.len(), 1);
1259        assert_eq!(
1260            result.contributions[0].contributing_module_id,
1261            "acme/support-console"
1262        );
1263
1264        let denied = provider.resolve_contributions(&ActionContributionResolutionRequest {
1265            context: context([]),
1266            ..request
1267        });
1268        assert_eq!(
1269            denied.unwrap_err().code,
1270            ModuleOperationsErrorCode::CapabilityDenied
1271        );
1272    }
1273
1274    #[test]
1275    fn sensitive_configuration_is_write_only_and_audit_evidence_contains_only_digests() {
1276        let mut release = release();
1277        release
1278            .manifest
1279            .capabilities
1280            .push("support.token.write".to_owned());
1281        release.manifest.capabilities.sort();
1282        release.manifest.config.fields.push(ModuleConfigField {
1283            key: "api_token".to_owned(),
1284            field_type: ModuleConfigFieldType::String,
1285            required: false,
1286            scope: ModuleConfigScope::Module,
1287            sensitive: true,
1288            secret_reference: true,
1289            mutability: ModuleConfigMutability::Runtime,
1290            activation: ModuleConfigActivation::None,
1291            read_capability: None,
1292            write_capability: Some("support.token.write".to_owned()),
1293            default: None,
1294            validation: None,
1295        });
1296        release.manifest_digest = digest_json(&release.manifest).unwrap();
1297        assert!(
1298            release.validate().is_empty(),
1299            "fixture release should remain valid"
1300        );
1301        let provider = ModuleOperationsProvider::new(
1302            "service-1",
1303            "spiffe://lenso/service-1",
1304            "revision-1",
1305            [release],
1306        )
1307        .unwrap();
1308        provider
1309            .set_config_value(
1310                "acme/support-console",
1311                "api_token",
1312                Value::String("super-secret".to_owned()),
1313            )
1314            .unwrap();
1315
1316        let read = provider
1317            .read_config(&ModuleConfigReadRequest {
1318                context: context([]),
1319                module_id: "acme/support-console".to_owned(),
1320                keys: vec!["api_token".to_owned()],
1321            })
1322            .unwrap();
1323        assert!(read.values[0].present);
1324        assert!(read.values[0].sensitive);
1325        assert_eq!(read.values[0].value, None);
1326
1327        let write = provider
1328            .write_config(&ModuleConfigWriteRequest {
1329                context: context(["support.token.write"]),
1330                module_id: "acme/support-console".to_owned(),
1331                values: vec![lenso_service::system_plane::ModuleConfigWriteValue {
1332                    key: "api_token".to_owned(),
1333                    value: Value::String("rotated-secret".to_owned()),
1334                }],
1335            })
1336            .unwrap();
1337        assert!(write.evidence[0].sensitive);
1338        assert!(write.evidence[0].old_value_digest.is_some());
1339        assert!(!write.evidence[0].new_value_digest.is_empty());
1340        let evidence = serde_json::to_string(&write.evidence).unwrap();
1341        assert!(!evidence.contains("rotated-secret"));
1342        assert!(!evidence.contains("super-secret"));
1343    }
1344}