1use 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 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 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}