1use crate::content_codec::{self, ContentFormat};
4use crate::fleet::integration_policies::{
5 self, IntegrationPackageSpec, IntegrationPolicyDetail, IntegrationPolicySpec,
6 IntegrationPolicySummary,
7};
8use crate::fleet::{agent_policies, agent_policy_ops};
9use crate::ops::{ExportOutcome, MutationPlan};
10use elasticctl_core::{Error, ErrorKind, Result, Transport};
11use serde::Serialize;
12use serde_json::{Map, Value, json};
13use std::collections::{BTreeMap, BTreeSet};
14use std::fmt;
15use std::path::{Path, PathBuf};
16
17const PAGE_SIZE: u64 = 1000;
18const IMPORT_RACE_WARNING: &str =
19 "warning Fleet can change after the final recheck and before the write";
20const DELETE_RACE_WARNING: &str =
21 "warning Fleet can change after the final recheck and before the write";
22
23#[derive(Debug, Clone, Default, PartialEq, Eq)]
24pub struct IntegrationPolicyFilter {
25 pub search: Option<String>,
26 pub limit: Option<usize>,
27}
28
29#[derive(Debug, Clone, PartialEq, Serialize)]
30pub struct IntegrationPolicyList {
31 pub total: u64,
32 pub integration_policies: Vec<IntegrationPolicySummary>,
33 pub truncated: bool,
34}
35
36#[derive(Debug, Clone, PartialEq)]
38pub(crate) struct ResolvedIntegrationPolicy {
39 pub(crate) summary: IntegrationPolicySummary,
40 pub(crate) item: Map<String, Value>,
41}
42
43#[derive(Debug, Clone, PartialEq, Eq)]
47enum PackageDependencyState {
48 Installed { version: String },
49 NotInstalled,
50}
51
52#[derive(Debug, Clone, PartialEq, Eq)]
53struct PackageDependencySnapshot {
54 name: String,
55 state: PackageDependencyState,
56}
57
58#[derive(Debug, Clone, Default, PartialEq, Eq)]
62struct SecretSchema {
63 package_vars: BTreeSet<String>,
64 input_vars: BTreeMap<String, BTreeSet<String>>,
65 stream_vars: BTreeMap<(String, String), BTreeSet<String>>,
66}
67
68#[derive(Debug, Clone, Default)]
69struct KnownSchema {
70 package_vars: BTreeSet<String>,
71 input_vars: BTreeMap<String, BTreeSet<String>>,
72 stream_vars: BTreeMap<(String, String), BTreeSet<String>>,
73}
74
75#[derive(Debug, Clone, Default, PartialEq, Eq)]
76struct VariableDefinitions {
77 known: BTreeSet<String>,
78 secrets: BTreeSet<String>,
79}
80
81#[derive(Debug)]
82struct TemplateDefinitions {
83 inputs: BTreeMap<String, String>,
84 datasets: BTreeSet<String>,
85}
86
87pub struct IntegrationPolicyImportArtifact {
91 source: PathBuf,
92 canonical: Vec<IntegrationPolicySpec>,
93}
94
95impl fmt::Debug for IntegrationPolicyImportArtifact {
96 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
97 formatter
98 .debug_struct("IntegrationPolicyImportArtifact")
99 .field("policy_count", &self.canonical.len())
100 .finish()
101 }
102}
103
104#[derive(Clone, PartialEq)]
108pub struct IntegrationPolicyImportPlan {
109 pub preview: MutationPlan,
110 pub skipped: Vec<Value>,
111 pub package_installs: Vec<String>,
112 pub total: usize,
113 source: PathBuf,
114 host: String,
115 space: String,
116 canonical: Vec<IntegrationPolicySpec>,
117 name_owners: BTreeMap<String, BTreeSet<String>>,
118 name_owners_snapshot: BTreeMap<String, BTreeSet<String>>,
119 parent_snapshots: BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
120 skipped_snapshot: Vec<Value>,
121 existing_snapshot: BTreeMap<String, Option<Map<String, Value>>>,
125 targets: Vec<IntegrationPolicyImportTarget>,
126 package_groups: BTreeMap<String, IntegrationPackageGroup>,
127 overwrite: bool,
128 skip_existing: bool,
129}
130
131impl fmt::Debug for IntegrationPolicyImportPlan {
132 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
133 formatter
134 .debug_struct("IntegrationPolicyImportPlan")
135 .field("target_count", &self.preview.targets.len())
136 .field("skipped_count", &self.skipped.len())
137 .field("package_install_count", &self.package_installs.len())
138 .field("existing_snapshot_count", &self.existing_snapshot.len())
139 .field("total", &self.total)
140 .finish()
141 }
142}
143
144#[derive(Debug, Clone, PartialEq, Serialize)]
145pub struct IntegrationPolicyImportReport {
146 pub applied: bool,
147 pub succeeded: Vec<Value>,
148 pub unchanged: Vec<Value>,
149 pub skipped: Vec<Value>,
150 pub failed: Vec<Value>,
151 pub total: usize,
152 pub affected_agents: u64,
153 pub package_installs: Vec<String>,
154}
155
156#[derive(Clone, PartialEq)]
160pub struct IntegrationPolicyDeletePlan {
161 pub preview: MutationPlan,
162 pub total: usize,
163 host: String,
164 host_snapshot: String,
165 space: String,
166 space_snapshot: String,
167 parent_snapshots: BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
168 parent_snapshots_snapshot: BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
169 targets: Vec<IntegrationPolicyDeleteTarget>,
170}
171
172impl fmt::Debug for IntegrationPolicyDeletePlan {
173 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
174 formatter
175 .debug_struct("IntegrationPolicyDeletePlan")
176 .field("target_count", &self.targets.len())
177 .field("total", &self.total)
178 .finish()
179 }
180}
181
182#[derive(Debug, Clone, PartialEq, Serialize)]
184pub struct IntegrationPolicyDeleteReport {
185 pub applied: bool,
186 pub deleted: Vec<Value>,
187 pub failed: Vec<Value>,
188 pub total: usize,
189 pub affected_agents: u64,
190}
191
192#[derive(Debug, Clone, PartialEq)]
193struct IntegrationPolicyImportTarget {
194 effective: IntegrationPolicySpec,
195 current: Option<IntegrationPolicyCurrentSnapshot>,
196 parents: BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
197 replacement_body: Option<Value>,
198}
199
200#[derive(Debug, Clone, PartialEq)]
201struct IntegrationPolicyCurrentSnapshot {
202 item: Map<String, Value>,
203 spec: IntegrationPolicySpec,
204 parent_ids: Vec<String>,
205}
206
207#[derive(Clone, PartialEq)]
211struct IntegrationPolicyDeleteTarget {
212 id: String,
213 name: String,
214 item: Map<String, Value>,
215 item_snapshot: Map<String, Value>,
216 spec: IntegrationPolicySpec,
217 spec_snapshot: IntegrationPolicySpec,
218 parents: BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
219 package: IntegrationPackageSpec,
220 dependency: PackageDependencySnapshot,
221 dependency_snapshot: PackageDependencySnapshot,
222 metadata: Map<String, Value>,
223 metadata_snapshot: Map<String, Value>,
224}
225
226#[derive(Debug, Clone, PartialEq)]
231struct IntegrationPackageGroup {
232 package: IntegrationPackageSpec,
233 state: PackageDependencySnapshot,
234 state_snapshot: PackageDependencySnapshot,
235 metadata: Map<String, Value>,
236 metadata_snapshot: Map<String, Value>,
237}
238
239#[derive(Debug, Clone, Copy, PartialEq, Eq)]
240enum ImportAction {
241 Create,
242 Replace,
243 Unchanged,
244}
245
246pub async fn collect(transport: &Transport) -> Result<Vec<Map<String, Value>>> {
248 let mut page_number = 1;
249 let mut total = None;
250 let mut items = Vec::new();
251 let mut ids = BTreeSet::new();
252 loop {
253 let page = integration_policies::list_page(transport, page_number).await?;
254 if page.page != page_number || page.per_page != PAGE_SIZE {
255 return Err(http(
256 "decoding integration policies list: unexpected page metadata",
257 ));
258 }
259 if page.items.len() as u64 > PAGE_SIZE {
260 return Err(http(
261 "decoding integration policies list: page returned more items than requested",
262 ));
263 }
264 match total {
265 Some(expected) if expected != page.total => {
266 return Err(http(
267 "decoding integration policies list: total changed while paging",
268 ));
269 }
270 Some(_) => {}
271 None => total = Some(page.total),
272 }
273 let page_len = page.items.len() as u64;
274 for item in page.items {
275 let id = required_string(&item, "id", "integration policies list")?;
276 if !ids.insert(id.clone()) {
277 return Err(http(format!(
278 "decoding integration policies list: duplicate integration policy id '{id}'"
279 )));
280 }
281 items.push(item);
282 }
283 let expected = total.expect("first page sets total");
284 if items.len() as u64 >= expected {
285 break;
286 }
287 if page_len != PAGE_SIZE {
288 return Err(http(
289 "decoding integration policies list: page was short before total",
290 ));
291 }
292 page_number += 1;
293 }
294 if items.len() as u64 > total.unwrap_or_default() {
295 return Err(http(
296 "decoding integration policies list: returned more items than total",
297 ));
298 }
299 items.sort_by(|left, right| left["id"].as_str().cmp(&right["id"].as_str()));
300 Ok(items)
301}
302
303pub async fn list_op(
305 transport: &Transport,
306 filter: &IntegrationPolicyFilter,
307) -> Result<IntegrationPolicyList> {
308 let items = collect(transport).await?;
309 let total = items.len() as u64;
310 let needle = filter.search.as_ref().map(|value| value.to_lowercase());
311 let mut integration_policies = Vec::new();
312 for item in &items {
313 let summary = summary_from_item(item)?;
314 if needle.as_ref().is_none_or(|needle| {
315 summary.id.to_lowercase().contains(needle)
316 || summary.name.to_lowercase().contains(needle)
317 }) {
318 integration_policies.push(summary);
319 }
320 }
321 let limit = filter.limit.unwrap_or(usize::MAX);
322 let truncated = integration_policies.len() > limit;
323 integration_policies.truncate(limit);
324 Ok(IntegrationPolicyList {
325 total,
326 integration_policies,
327 truncated,
328 })
329}
330
331pub async fn resolve(transport: &Transport, selector: &str) -> Result<IntegrationPolicySummary> {
333 Ok(resolve_item(transport, selector).await?.summary)
334}
335
336pub(crate) async fn resolve_item(
339 transport: &Transport,
340 selector: &str,
341) -> Result<ResolvedIntegrationPolicy> {
342 match integration_policies::get(transport, selector).await {
343 Ok(policy) => return checked_read(selector, policy.item, None),
344 Err(error) if error.kind == ErrorKind::NotFound => {}
345 Err(error) => return Err(error),
346 }
347 let matches: Vec<IntegrationPolicySummary> = collect(transport)
348 .await?
349 .iter()
350 .filter(|item| item.get("name").and_then(Value::as_str) == Some(selector))
351 .map(summary_from_item)
352 .collect::<Result<_>>()?;
353 match matches.as_slice() {
354 [] => Err(Error::new(
355 ErrorKind::NotFound,
356 format!("no integration policy with id or name '{selector}'"),
357 )),
358 [one] => {
359 let policy = integration_policies::get(transport, &one.id).await?;
360 checked_read(&one.id, policy.item, Some(one))
361 }
362 many => Err(Error::new(
363 ErrorKind::Conflict,
364 format!(
365 "integration policy '{selector}' is ambiguous: {}",
366 many.iter()
367 .map(|policy| policy.id.as_str())
368 .collect::<Vec<_>>()
369 .join(", ")
370 ),
371 )),
372 }
373}
374
375fn checked_read(
379 requested_id: &str,
380 item: Map<String, Value>,
381 selected: Option<&IntegrationPolicySummary>,
382) -> Result<ResolvedIntegrationPolicy> {
383 let summary = summary_from_item(&item)?;
384 if summary.id != requested_id {
385 return Err(http(
386 "decoding integration policy selector read: response id did not match the selector",
387 ));
388 }
389 if selected.is_some_and(|selected| !same_summary(&summary, selected)) {
390 return Err(http(
391 "decoding integration policy selector read: fetched item did not match the selected list summary",
392 ));
393 }
394 Ok(ResolvedIntegrationPolicy { summary, item })
395}
396
397fn same_summary(fetched: &IntegrationPolicySummary, selected: &IntegrationPolicySummary) -> bool {
401 fetched.id == selected.id
402 && fetched.name == selected.name
403 && fetched.namespace == selected.namespace
404 && fetched.description == selected.description
405 && fetched.package == selected.package
406 && sorted_parent_ids(&fetched.policy_ids) == sorted_parent_ids(&selected.policy_ids)
407}
408
409fn sorted_parent_ids(ids: &[String]) -> Vec<&str> {
410 let mut sorted = ids.iter().map(String::as_str).collect::<Vec<_>>();
411 sorted.sort_unstable();
412 sorted
413}
414
415pub async fn get_op(transport: &Transport, selector: &str) -> Result<IntegrationPolicyDetail> {
418 let resolved = resolve_item(transport, selector).await?;
419 let mut blocked_by = live_blocked_by(&resolved.item, &resolved.summary.id, transport.space())?;
420 validate_safe_detail_shape(&resolved.item, transport.space())?;
421 let parents = read_parents(&resolved.summary.id, &resolved.item)?;
422 let parents = read_parent_snapshots(transport, &resolved.summary.id, &parents).await?;
423 for parent in parents.values() {
424 if parent.platform_owned {
425 blocked_by.insert(format!("parent:{}.platform_owned", parent.id));
426 }
427 if parent.protected {
428 blocked_by.insert(format!("parent:{}.is_protected", parent.id));
429 }
430 }
431 if parents
432 .values()
433 .map(|parent| parent.namespace.as_str())
434 .collect::<BTreeSet<_>>()
435 .len()
436 != 1
437 {
438 blocked_by.insert("namespace".into());
439 }
440 if parents
441 .values()
442 .any(|parent| parent.namespace != resolved.summary.namespace)
443 {
444 blocked_by.insert("namespace".into());
445 }
446 Ok(IntegrationPolicyDetail {
447 id: resolved.summary.id,
448 name: resolved.summary.name,
449 namespace: resolved.summary.namespace,
450 description: resolved.summary.description,
451 policy_ids: parents.keys().cloned().collect(),
452 package: resolved.summary.package,
453 affected_agents: parents.values().map(|parent| parent.agents).sum(),
454 blocked_by: blocked_by.into_iter().collect(),
455 })
456}
457
458fn validate_safe_detail_shape(item: &Map<String, Value>, active_space: &str) -> Result<()> {
463 let mut projected = item.clone();
464 projected.insert("enabled".into(), Value::Bool(true));
465 for field in [
466 "is_managed",
467 "supports_agentless",
468 "supports_cloud_connector",
469 ] {
470 projected.insert(field.into(), Value::Bool(false));
471 }
472 for field in ["output_id", "cloud_connector_id", "cloud_connector_name"] {
473 projected.insert(field.into(), Value::Null);
474 }
475 projected.insert("secret_references".into(), Value::Array(Vec::new()));
476 projected.insert("spaceIds".into(), Value::Null);
477 normalize(&projected, active_space).map(|_| ())
478}
479
480pub async fn export(
484 transport: &Transport,
485 selectors: &[String],
486 all_custom: bool,
487 format: ContentFormat,
488) -> Result<ExportOutcome> {
489 if selectors.is_empty() && !all_custom {
490 return Err(Error::new(
491 ErrorKind::Error,
492 "integration-policy export needs selectors or --all-custom",
493 ));
494 }
495 if !selectors.is_empty() && all_custom {
496 return Err(Error::new(
497 ErrorKind::Error,
498 "--all-custom cannot be combined with selectors",
499 ));
500 }
501
502 let mut rows = BTreeMap::new();
503 if all_custom {
504 for item in collect(transport).await? {
505 let summary = summary_from_item(&item)?;
506 if optional_bool(&item, "is_managed", &summary.id)? == Some(true) {
507 continue;
508 }
509 let live = integration_policies::get(transport, &summary.id).await?;
510 let resolved = checked_read(&summary.id, live.item, Some(&summary))?;
511 if optional_bool(&resolved.item, "is_managed", &summary.id)? == Some(true) {
512 continue;
513 }
514 rows.insert(summary.id.clone(), resolved);
515 }
516 } else {
517 for selector in selectors {
518 let resolved = resolve_item(transport, selector).await?;
519 rows.entry(resolved.summary.id.clone()).or_insert(resolved);
520 }
521 }
522
523 let mut specs = Vec::new();
524 for (id, resolved) in rows {
525 let parent_ids = read_parents(&id, &resolved.item)?;
526 let parents = read_parent_snapshots(transport, &id, &parent_ids).await?;
527 if all_custom && parents.values().any(|parent| parent.platform_owned) {
528 continue;
529 }
530 specs.push(effective_spec(transport, &id, &resolved.item, &parents).await?);
531 }
532 specs.sort_by(|left, right| left.id.cmp(&right.id));
533 Ok(ExportOutcome {
534 body: content_codec::encode_sequence(&specs, format)?,
535 exported: specs.len() as u64,
536 missing: Vec::new(),
537 })
538}
539
540fn read_parents(id: &str, item: &Map<String, Value>) -> Result<Vec<String>> {
541 let policy_ids = item
542 .get("policy_ids")
543 .and_then(Value::as_array)
544 .ok_or_else(|| {
545 http(format!(
546 "decoding integration policy '{id}': policy_ids must be an array"
547 ))
548 })?;
549 if policy_ids.is_empty() {
550 return Err(http(format!(
551 "decoding integration policy '{id}': policy_ids must not be empty"
552 )));
553 }
554 let mut ids = Vec::with_capacity(policy_ids.len());
555 for parent in policy_ids {
556 let parent = parent
557 .as_str()
558 .filter(|value| !value.trim().is_empty())
559 .ok_or_else(|| {
560 http(format!(
561 "decoding integration policy '{id}': policy_ids must contain non-empty strings"
562 ))
563 })?;
564 ids.push(parent.to_owned());
565 }
566 ids.sort();
567 if ids.windows(2).any(|pair| pair[0] == pair[1]) {
568 return Err(http(format!(
569 "decoding integration policy '{id}': duplicate policy_ids"
570 )));
571 }
572 Ok(ids)
573}
574
575async fn read_parent_snapshots(
576 transport: &Transport,
577 integration_id: &str,
578 parent_ids: &[String],
579) -> Result<BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>> {
580 let mut parents = BTreeMap::new();
581 for parent_id in parent_ids {
582 let parent = agent_policy_ops::read_parent_snapshot(transport, parent_id).await?;
583 if !parent
584 .attached_integrations
585 .binary_search_by(|attached| attached.as_str().cmp(integration_id))
586 .is_ok()
587 {
588 return Err(http(format!(
589 "decoding integration policy '{integration_id}': parent '{parent_id}' is missing its attachment"
590 )));
591 }
592 parents.insert(parent_id.clone(), parent);
593 }
594 Ok(parents)
595}
596
597async fn effective_spec(
598 transport: &Transport,
599 id: &str,
600 item: &Map<String, Value>,
601 parents: &BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
602) -> Result<IntegrationPolicySpec> {
603 for parent in parents.values() {
604 if parent.platform_owned {
605 return unsupported(format!(
606 "integration policy '{id}' is not portable: parent {} is platform-owned",
607 parent.id
608 ));
609 }
610 if parent.protected {
611 return unsupported(format!(
612 "integration policy '{id}' is not portable: parent {} is_protected",
613 parent.id
614 ));
615 }
616 }
617 let dependency =
618 read_dependencies(transport, &package_coordinate(item, "integration policy")?).await?;
619 let mut spec = normalize(item, transport.space())?;
620 if let Some(namespace) = &spec.namespace {
621 if parents
622 .values()
623 .any(|parent| &parent.namespace != namespace)
624 {
625 return unsupported(format!(
626 "integration policy '{id}' is not portable: namespace does not match every parent"
627 ));
628 }
629 } else {
630 let namespaces: BTreeSet<&str> = parents
631 .values()
632 .map(|parent| parent.namespace.as_str())
633 .collect();
634 if namespaces.len() != 1 {
635 return unsupported(format!(
636 "integration policy '{id}' is not portable: parents have different namespaces"
637 ));
638 }
639 spec.namespace = namespaces.into_iter().next().map(str::to_owned);
640 }
641 match dependency.state {
644 PackageDependencyState::Installed { ref version } if version == &spec.package.version => {}
645 PackageDependencyState::Installed { .. } => {
646 return Err(Error::new(
647 ErrorKind::Conflict,
648 format!(
649 "integration policy '{id}' package {} has a different installed version",
650 dependency.name
651 ),
652 ));
653 }
654 PackageDependencyState::NotInstalled => {
655 return Err(Error::new(
656 ErrorKind::Conflict,
657 format!(
658 "integration policy '{id}' package {} is not installed",
659 dependency.name
660 ),
661 ));
662 }
663 }
664 let metadata = integration_policies::package_metadata(
665 transport,
666 &spec.package.name,
667 &spec.package.version,
668 )
669 .await?;
670 let paths = configured_secret_paths(&spec, &metadata.item)?;
671 if !paths.is_empty() {
672 return unsupported(format!(
673 "integration policy '{id}' is not portable: {}",
674 paths
675 .into_iter()
676 .map(|path| format!("{id}:{path}"))
677 .collect::<Vec<_>>()
678 .join(", ")
679 ));
680 }
681 Ok(spec)
682}
683
684async fn read_dependencies(
685 transport: &Transport,
686 package: &IntegrationPackageSpec,
687) -> Result<PackageDependencySnapshot> {
688 let status = agent_policies::package_status(transport, &package.name).await?;
689 let state = match (status.status.as_str(), status.installed_version) {
690 ("installed", Some(version)) if !version.trim().is_empty() => {
691 PackageDependencyState::Installed { version }
692 }
693 ("not_installed", None) => PackageDependencyState::NotInstalled,
694 _ => {
695 return Err(http(format!(
696 "decoding package dependency '{}': invalid status/version state",
697 package.name
698 )));
699 }
700 };
701 Ok(PackageDependencySnapshot {
702 name: status.name,
703 state,
704 })
705}
706
707fn configured_secret_paths(
708 spec: &IntegrationPolicySpec,
709 metadata: &Map<String, Value>,
710) -> Result<Vec<String>> {
711 let (secrets, known) = secret_schema(metadata)?;
712 let mut paths = BTreeSet::new();
713 configured_vars(
714 spec.vars.as_ref(),
715 &known.package_vars,
716 &secrets.package_vars,
717 &spec.id,
718 "vars",
719 &mut paths,
720 )?;
721 for (input_key, input) in &spec.inputs {
722 let input = input.as_object().ok_or_else(|| {
723 http(format!(
724 "decoding integration policy '{}': inputs.{input_key} must be an object",
725 spec.id
726 ))
727 })?;
728 let known_vars = known.input_vars.get(input_key).ok_or_else(|| {
729 Error::new(
730 ErrorKind::Unsupported,
731 format!(
732 "integration policy '{}': {}:inputs.{input_key} has no matching package definition",
733 spec.id, spec.id
734 ),
735 )
736 })?;
737 let secret_vars = secrets
738 .input_vars
739 .get(input_key)
740 .cloned()
741 .unwrap_or_default();
742 configured_vars(
743 input.get("vars").map(expect_object).transpose()?,
744 known_vars,
745 &secret_vars,
746 &spec.id,
747 &format!("inputs.{input_key}.vars"),
748 &mut paths,
749 )?;
750 if let Some(streams) = input.get("streams") {
751 let streams = streams.as_object().ok_or_else(|| {
752 http(format!(
753 "decoding integration policy '{}': inputs.{input_key}.streams must be an object",
754 spec.id
755 ))
756 })?;
757 for (dataset, stream) in streams {
758 let stream = stream.as_object().ok_or_else(|| {
759 http(format!(
760 "decoding integration policy '{}': inputs.{input_key}.streams.{dataset} must be an object",
761 spec.id
762 ))
763 })?;
764 let key = (input_key.clone(), dataset.clone());
765 let known_vars = known.stream_vars.get(&key).ok_or_else(|| {
766 Error::new(
767 ErrorKind::Unsupported,
768 format!(
769 "integration policy '{}': {}:inputs.{input_key}.streams.{dataset} has no matching package definition",
770 spec.id, spec.id
771 ),
772 )
773 })?;
774 let secret_vars = secrets.stream_vars.get(&key).cloned().unwrap_or_default();
775 configured_vars(
776 stream.get("vars").map(expect_object).transpose()?,
777 known_vars,
778 &secret_vars,
779 &spec.id,
780 &format!("inputs.{input_key}.streams.{dataset}.vars"),
781 &mut paths,
782 )?;
783 }
784 }
785 }
786 Ok(paths.into_iter().collect())
787}
788
789fn expect_object(value: &Value) -> Result<&Map<String, Value>> {
790 value
791 .as_object()
792 .ok_or_else(|| http("decoding integration policy: configured vars must be an object"))
793}
794
795fn configured_vars(
796 configured: Option<&Map<String, Value>>,
797 known: &BTreeSet<String>,
798 secret: &BTreeSet<String>,
799 policy_id: &str,
800 prefix: &str,
801 paths: &mut BTreeSet<String>,
802) -> Result<()> {
803 let Some(configured) = configured else {
804 return Ok(());
805 };
806 for name in configured.keys() {
807 if !known.contains(name) {
808 return Err(Error::new(
809 ErrorKind::Unsupported,
810 format!(
811 "integration policy '{policy_id}' is not portable: {policy_id}:{prefix}.{name} has no matching package definition"
812 ),
813 ));
814 }
815 if secret.contains(name) {
816 paths.insert(format!("{prefix}.{name}"));
817 }
818 }
819 Ok(())
820}
821
822fn secret_schema(metadata: &Map<String, Value>) -> Result<(SecretSchema, KnownSchema)> {
823 let package_name = metadata_name(metadata, "name", "package metadata")?;
824 let package_vars = parse_var_definitions(metadata.get("vars"), "package metadata vars")?;
825 let modern_datasets = parse_modern_data_streams(metadata.get("data_streams"))?;
826
827 let mut secrets = SecretSchema {
828 package_vars: package_vars.secrets.clone(),
829 ..SecretSchema::default()
830 };
831 let mut known = KnownSchema {
832 package_vars: package_vars.known,
833 ..KnownSchema::default()
834 };
835 let mut template_names = BTreeSet::new();
836 let mut templates = Vec::new();
837 let mut legacy_streams = BTreeMap::new();
838 let policy_templates = match metadata.get("policy_templates") {
839 None => &[][..],
840 Some(Value::Array(value)) => value,
841 Some(_) => {
842 return Err(http(
843 "decoding package metadata: policy_templates must be an array",
844 ));
845 }
846 };
847
848 for template in policy_templates {
849 let template = template.as_object().ok_or_else(|| {
850 http("decoding package metadata: policy_templates entry must be an object")
851 })?;
852 let template_name = metadata_name(template, "name", "policy_templates entry")?;
853 if !template_names.insert(template_name.clone()) {
854 return Err(http(format!(
855 "decoding package metadata: duplicate template name '{template_name}'"
856 )));
857 }
858 let datasets = resolve_template_datasets(
859 template.get("data_streams"),
860 &modern_datasets,
861 &package_name,
862 &template_name,
863 )?;
864 let inputs = match template.get("inputs") {
865 None => &[][..],
866 Some(Value::Array(value)) => value,
867 Some(_) => {
868 return Err(http(
869 "decoding package metadata: policy_templates inputs must be an array",
870 ));
871 }
872 };
873 let mut template_inputs = BTreeMap::new();
874 for input in inputs {
875 let input = input.as_object().ok_or_else(|| {
876 http("decoding package metadata: policy_templates inputs entry must be an object")
877 })?;
878 let input_type = metadata_name(input, "type", "policy_templates input")?;
879 let input_key = format!("{template_name}-{input_type}");
880 let input_vars =
881 parse_var_definitions(input.get("vars"), "package metadata input vars")?;
882 if template_inputs
883 .insert(input_type.clone(), input_key.clone())
884 .is_some()
885 || known
886 .input_vars
887 .insert(input_key.clone(), input_vars.known)
888 .is_some()
889 {
890 return Err(http(format!(
891 "decoding package metadata: duplicate input key '{input_key}'"
892 )));
893 }
894 if !input_vars.secrets.is_empty() {
895 secrets
896 .input_vars
897 .insert(input_key.clone(), input_vars.secrets);
898 }
899
900 let streams = match input.get("streams") {
901 None => &[][..],
902 Some(Value::Array(value)) => value,
903 Some(_) => {
904 return Err(http(
905 "decoding package metadata: input streams must be an array",
906 ));
907 }
908 };
909 for stream in streams {
910 let stream = stream.as_object().ok_or_else(|| {
911 http("decoding package metadata: input streams entry must be an object")
912 })?;
913 let data_stream = stream
914 .get("data_stream")
915 .and_then(Value::as_object)
916 .ok_or_else(|| {
917 http("decoding package metadata: stream data_stream must be an object")
918 })?;
919 let dataset = metadata_name(data_stream, "dataset", "stream data_stream")?;
920 let definition =
921 parse_var_definitions(stream.get("vars"), "package metadata stream vars")?;
922 let key = (input_key.clone(), dataset);
923 if legacy_streams.insert(key.clone(), definition).is_some() {
924 return Err(http(format!(
925 "decoding package metadata: duplicate stream key '{}:{}'",
926 key.0, key.1
927 )));
928 }
929 }
930 }
931 templates.push(TemplateDefinitions {
932 inputs: template_inputs,
933 datasets,
934 });
935 }
936
937 let mut modern_streams = BTreeMap::new();
938 for (dataset, streams) in &modern_datasets {
939 for (input_type, definition) in streams {
940 let candidates = templates
941 .iter()
942 .filter(|template| template.datasets.contains(dataset))
943 .filter_map(|template| template.inputs.get(input_type))
944 .cloned()
945 .collect::<Vec<_>>();
946 let input_key = match candidates.as_slice() {
947 [input_key] => input_key.clone(),
948 [] => {
949 return Err(http(format!(
950 "decoding package metadata: stream '{dataset}:{input_type}' has no matching template input"
951 )));
952 }
953 _ => {
954 return Err(http(format!(
955 "decoding package metadata: stream '{dataset}:{input_type}' has multiple matching template inputs"
956 )));
957 }
958 };
959 let key = (input_key, dataset.clone());
960 if modern_streams
961 .insert(key.clone(), definition.clone())
962 .is_some()
963 {
964 return Err(http(format!(
965 "decoding package metadata: duplicate stream key '{}:{}'",
966 key.0, key.1
967 )));
968 }
969 }
970 }
971
972 for (key, definition) in legacy_streams {
973 match modern_streams.get(&key) {
974 Some(modern) if modern != &definition => {
975 return Err(http(format!(
976 "decoding package metadata: conflicting modern and legacy stream definition '{}:{}'",
977 key.0, key.1
978 )));
979 }
980 Some(_) => {}
981 None => {
982 modern_streams.insert(key, definition);
983 }
984 }
985 }
986 for (key, definition) in modern_streams {
987 if !definition.secrets.is_empty() {
988 secrets.stream_vars.insert(key.clone(), definition.secrets);
989 }
990 known.stream_vars.insert(key, definition.known);
991 }
992 Ok((secrets, known))
993}
994
995fn parse_modern_data_streams(
996 value: Option<&Value>,
997) -> Result<BTreeMap<String, BTreeMap<String, VariableDefinitions>>> {
998 let data_streams = match value {
999 None => &[][..],
1000 Some(Value::Array(value)) => value,
1001 Some(_) => {
1002 return Err(http(
1003 "decoding package metadata: data_streams must be an array",
1004 ));
1005 }
1006 };
1007 let mut datasets = BTreeMap::new();
1008 for data_stream in data_streams {
1009 let data_stream = data_stream.as_object().ok_or_else(|| {
1010 http("decoding package metadata: data_streams entry must be an object")
1011 })?;
1012 let dataset = metadata_name(data_stream, "dataset", "data_streams entry")?;
1013 let streams = match data_stream.get("streams") {
1014 None => &[][..],
1015 Some(Value::Array(value)) => value,
1016 Some(_) => {
1017 return Err(http(
1018 "decoding package metadata: data_streams streams must be an array",
1019 ));
1020 }
1021 };
1022 let mut stream_definitions = BTreeMap::new();
1023 for stream in streams {
1024 let stream = stream.as_object().ok_or_else(|| {
1025 http("decoding package metadata: data_streams streams entry must be an object")
1026 })?;
1027 let input = metadata_name(stream, "input", "data_streams stream")?;
1028 let definition =
1029 parse_var_definitions(stream.get("vars"), "package metadata stream vars")?;
1030 if stream_definitions
1031 .insert(input.clone(), definition)
1032 .is_some()
1033 {
1034 return Err(http(format!(
1035 "decoding package metadata: duplicate stream input '{input}' for dataset '{dataset}'"
1036 )));
1037 }
1038 }
1039 if datasets
1040 .insert(dataset.clone(), stream_definitions)
1041 .is_some()
1042 {
1043 return Err(http(format!(
1044 "decoding package metadata: duplicate data stream dataset '{dataset}'"
1045 )));
1046 }
1047 }
1048 Ok(datasets)
1049}
1050
1051fn resolve_template_datasets(
1052 value: Option<&Value>,
1053 datasets: &BTreeMap<String, BTreeMap<String, VariableDefinitions>>,
1054 package_name: &str,
1055 template_name: &str,
1056) -> Result<BTreeSet<String>> {
1057 let Some(value) = value else {
1058 return Ok(datasets.keys().cloned().collect());
1059 };
1060 let selectors = value
1061 .as_array()
1062 .ok_or_else(|| http("decoding package metadata: template data_streams must be an array"))?;
1063 let mut selected = BTreeSet::new();
1064 let mut seen = BTreeSet::new();
1065 for selector in selectors {
1066 let selector = selector
1067 .as_str()
1068 .filter(|selector| !selector.trim().is_empty())
1069 .ok_or_else(|| {
1070 http("decoding package metadata: template data_streams selector must be a non-empty string")
1071 })?;
1072 if !seen.insert(selector) {
1073 return Err(http(format!(
1074 "decoding package metadata: duplicate data_streams selector '{selector}' in template '{template_name}'"
1075 )));
1076 }
1077 let short_name = format!("{package_name}.{selector}");
1078 let candidates = datasets
1079 .keys()
1080 .filter(|dataset| dataset.as_str() == selector || dataset.as_str() == short_name)
1081 .collect::<Vec<_>>();
1082 match candidates.as_slice() {
1083 [dataset] => {
1084 if !selected.insert((**dataset).clone()) {
1085 return Err(http(format!(
1086 "decoding package metadata: duplicate data_streams dataset '{}' in template '{template_name}'",
1087 dataset
1088 )));
1089 }
1090 }
1091 [] => {
1092 return Err(http(format!(
1093 "decoding package metadata: data_streams selector '{selector}' in template '{template_name}' does not match a dataset"
1094 )));
1095 }
1096 _ => {
1097 return Err(http(format!(
1098 "decoding package metadata: data_streams selector '{selector}' in template '{template_name}' matches multiple datasets"
1099 )));
1100 }
1101 }
1102 }
1103 Ok(selected)
1104}
1105
1106fn parse_var_definitions(value: Option<&Value>, context: &str) -> Result<VariableDefinitions> {
1107 let values = match value {
1108 None => return Ok(VariableDefinitions::default()),
1109 Some(Value::Array(values)) => values,
1110 Some(_) => return Err(http(format!("decoding {context}: vars must be an array"))),
1111 };
1112 let mut definitions = VariableDefinitions::default();
1113 for value in values {
1114 let definition = value
1115 .as_object()
1116 .ok_or_else(|| http(format!("decoding {context}: variable must be an object")))?;
1117 let name = metadata_name(definition, "name", context)?;
1118 if !definitions.known.insert(name.clone()) {
1119 return Err(http(format!(
1120 "decoding {context}: duplicate variable name '{name}'"
1121 )));
1122 }
1123 match definition.get("secret") {
1124 None => {}
1125 Some(Value::Bool(true)) => {
1126 definitions.secrets.insert(name);
1127 }
1128 Some(Value::Bool(false)) => {}
1129 Some(_) => {
1130 return Err(http(format!(
1131 "decoding {context}: secret must be a boolean"
1132 )));
1133 }
1134 }
1135 }
1136 Ok(definitions)
1137}
1138
1139fn metadata_name(object: &Map<String, Value>, field: &str, context: &str) -> Result<String> {
1140 object
1141 .get(field)
1142 .and_then(Value::as_str)
1143 .filter(|value| !value.trim().is_empty())
1144 .map(str::to_owned)
1145 .ok_or_else(|| {
1146 http(format!(
1147 "decoding package metadata: {context} {field} must be a non-empty string"
1148 ))
1149 })
1150}
1151
1152fn live_blocked_by(
1153 item: &Map<String, Value>,
1154 id: &str,
1155 active_space: &str,
1156) -> Result<BTreeSet<String>> {
1157 let mut reasons = BTreeSet::new();
1158 match item.get("enabled") {
1159 Some(Value::Bool(true)) => {}
1160 Some(Value::Bool(false)) => {
1161 reasons.insert("enabled".into());
1162 }
1163 _ => {
1164 return Err(http(format!(
1165 "decoding integration policy '{id}': enabled must be true or false"
1166 )));
1167 }
1168 }
1169 for field in [
1170 "is_managed",
1171 "supports_agentless",
1172 "supports_cloud_connector",
1173 ] {
1174 if optional_bool(item, field, id)? == Some(true) {
1175 reasons.insert(field.to_owned());
1176 }
1177 }
1178 for field in ["output_id", "cloud_connector_id", "cloud_connector_name"] {
1179 match item.get(field) {
1180 None | Some(Value::Null) | Some(Value::Bool(false)) => {}
1181 Some(Value::String(_)) => {
1182 reasons.insert(field.to_owned());
1183 }
1184 Some(_) => {
1185 return Err(http(format!(
1186 "decoding integration policy '{id}': {field} must be a string or null"
1187 )));
1188 }
1189 }
1190 }
1191 match item.get("secret_references") {
1192 None | Some(Value::Null) => {}
1193 Some(Value::Array(values)) if values.is_empty() => {}
1194 Some(Value::Array(_)) => {
1195 reasons.insert("secret_references".into());
1196 }
1197 Some(_) => {
1198 return Err(http(format!(
1199 "decoding integration policy '{id}': secret_references must be an array or null"
1200 )));
1201 }
1202 }
1203 let active = if active_space.is_empty() {
1204 "default"
1205 } else {
1206 active_space
1207 };
1208 match item.get("spaceIds") {
1209 None | Some(Value::Null) => {}
1210 Some(Value::Array(spaces)) => {
1211 for space in spaces {
1212 let space = space.as_str().filter(|value| !value.is_empty()).ok_or_else(|| {
1213 http(format!("decoding integration policy '{id}': spaceIds must contain non-empty strings"))
1214 })?;
1215 if space != active {
1216 reasons.insert("spaceIds".into());
1217 }
1218 }
1219 }
1220 Some(_) => {
1221 return Err(http(format!(
1222 "decoding integration policy '{id}': spaceIds must be an array or null"
1223 )));
1224 }
1225 }
1226 Ok(reasons)
1227}
1228
1229const PORTABLE_OPTIONAL: [&str; 6] = [
1230 "description",
1231 "namespace",
1232 "vars",
1233 "var_group_selections",
1234 "condition",
1235 "additional_datastreams_permissions",
1236];
1237
1238const REMOVED_FIELDS: [&str; 19] = [
1239 "agents",
1240 "cloud_connector_id",
1241 "cloud_connector_name",
1242 "created_at",
1243 "created_by",
1244 "elasticsearch",
1245 "enabled",
1246 "is_managed",
1247 "output_id",
1248 "package_agent_version_condition",
1249 "policy_id",
1250 "revision",
1251 "secret_references",
1252 "spaceIds",
1253 "supports_agentless",
1254 "supports_cloud_connector",
1255 "updated_at",
1256 "updated_by",
1257 "version",
1258];
1259
1260pub fn normalize(item: &Map<String, Value>, active_space: &str) -> Result<IntegrationPolicySpec> {
1262 let id = required_string(item, "id", "integration policy")?;
1263 portability_check(item, &id, active_space)?;
1264 reject_unknown_top_level(item, &id)?;
1265
1266 let mut portable = Map::new();
1267 for field in ["id", "name"] {
1268 if let Some(value) = item.get(field) {
1269 portable.insert(field.to_owned(), value.clone());
1270 }
1271 }
1272 portable.insert(
1273 "policy_ids".to_owned(),
1274 Value::Array(
1275 read_parents(&id, item)?
1276 .into_iter()
1277 .map(Value::String)
1278 .collect(),
1279 ),
1280 );
1281 portable.insert("package".to_owned(), normalize_package(item, &id)?);
1282 portable.insert("inputs".to_owned(), normalize_inputs(item, &id)?);
1283 for field in PORTABLE_OPTIONAL {
1284 if let Some(value) = item.get(field)
1285 && !value.is_null()
1286 {
1287 portable.insert(field.to_owned(), value.clone());
1288 }
1289 }
1290 IntegrationPolicySpec::try_from(Value::Object(portable)).map_err(|error| {
1291 http(format!(
1292 "decoding integration policy '{id}': {}",
1293 error.message
1294 ))
1295 })
1296}
1297
1298fn normalize_package(item: &Map<String, Value>, id: &str) -> Result<Value> {
1299 let package = item
1300 .get("package")
1301 .and_then(Value::as_object)
1302 .ok_or_else(|| {
1303 http(format!(
1304 "decoding integration policy '{id}': package must be an object"
1305 ))
1306 })?;
1307 for (field, expected) in [("title", "a string")] {
1308 if let Some(value) = package.get(field)
1309 && !value.is_null()
1310 && !value.is_string()
1311 {
1312 return Err(http(format!(
1313 "decoding integration policy '{id}': package.{field} must be {expected} or null"
1314 )));
1315 }
1316 }
1317 for field in ["requires_root", "fips_compatible"] {
1318 if let Some(value) = package.get(field)
1319 && !value.is_null()
1320 && !value.is_boolean()
1321 {
1322 return Err(http(format!(
1323 "decoding integration policy '{id}': package.{field} must be a boolean or null"
1324 )));
1325 }
1326 }
1327 let mut portable = Map::new();
1328 for field in ["name", "version"] {
1329 if let Some(value) = package.get(field) {
1330 portable.insert(field.to_owned(), value.clone());
1331 }
1332 }
1333 let known: BTreeSet<&str> = [
1334 "name",
1335 "version",
1336 "title",
1337 "requires_root",
1338 "fips_compatible",
1339 ]
1340 .into_iter()
1341 .collect();
1342 if let Some(field) = package
1343 .keys()
1344 .map(String::as_str)
1345 .filter(|field| !known.contains(field))
1346 .min()
1347 {
1348 return Err(Error::new(
1349 ErrorKind::Unsupported,
1350 format!("integration policy '{id}' carries unknown package field '{field}'"),
1351 ));
1352 }
1353 Ok(Value::Object(portable))
1354}
1355
1356fn normalize_inputs(item: &Map<String, Value>, id: &str) -> Result<Value> {
1357 let inputs = item
1358 .get("inputs")
1359 .and_then(Value::as_object)
1360 .ok_or_else(|| {
1361 http(format!(
1362 "decoding integration policy '{id}': inputs must be an object"
1363 ))
1364 })?;
1365 let mut normalized = Map::new();
1366 for (input_id, input) in inputs {
1367 normalized.insert(
1368 input_id.clone(),
1369 normalize_package_map(input, "compiled_input")?,
1370 );
1371 }
1372 Ok(Value::Object(normalized))
1373}
1374
1375fn normalize_package_map(value: &Value, compiled_field: &str) -> Result<Value> {
1378 let object = value
1379 .as_object()
1380 .ok_or_else(|| http("decoding integration policy: input must be an object"))?;
1381 let mut normalized = Map::new();
1382 for (field, value) in object {
1383 if field == "id" {
1384 if !value.is_string() {
1385 return Err(http(
1386 "decoding integration policy: generated id must be a string",
1387 ));
1388 }
1389 continue;
1390 }
1391 if field == compiled_field {
1392 if !value.is_object() {
1393 return Err(http(format!(
1394 "decoding integration policy: {compiled_field} must be an object"
1395 )));
1396 }
1397 continue;
1398 }
1399 if field == "streams" {
1400 let streams = value.as_object().ok_or_else(|| {
1401 http("decoding integration policy: input streams must be an object")
1402 })?;
1403 let mut normalized_streams = Map::new();
1404 for (stream_id, stream) in streams {
1405 normalized_streams.insert(
1406 stream_id.clone(),
1407 normalize_package_map(stream, "compiled_stream")?,
1408 );
1409 }
1410 normalized.insert(field.clone(), Value::Object(normalized_streams));
1411 } else {
1412 normalized.insert(field.clone(), value.clone());
1413 }
1414 }
1415 Ok(Value::Object(normalized))
1416}
1417
1418fn portability_check(item: &Map<String, Value>, id: &str, active_space: &str) -> Result<()> {
1419 let mut reasons = BTreeSet::new();
1420 required_true(item, "enabled", id)?;
1421 if let Some(value) = item.get("elasticsearch")
1422 && !value.is_object()
1423 {
1424 return Err(http(format!(
1425 "decoding integration policy '{id}': elasticsearch must be an object"
1426 )));
1427 }
1428 if let Some(value) = item.get("package_agent_version_condition")
1429 && !value.is_null()
1430 && !value.is_string()
1431 {
1432 return Err(http(format!(
1433 "decoding integration policy '{id}': package_agent_version_condition must be a string or null"
1434 )));
1435 }
1436 for field in [
1437 "is_managed",
1438 "supports_agentless",
1439 "supports_cloud_connector",
1440 ] {
1441 if let Some(true) = optional_bool(item, field, id)? {
1442 reasons.insert(field);
1443 }
1444 }
1445 for field in ["output_id", "cloud_connector_id", "cloud_connector_name"] {
1446 match item.get(field) {
1447 None | Some(Value::Null) | Some(Value::Bool(false)) => {}
1448 Some(Value::String(_)) => {
1449 reasons.insert(field);
1450 }
1451 Some(_) => {
1452 return Err(http(format!(
1453 "decoding integration policy '{id}': {field} must be a string or null"
1454 )));
1455 }
1456 }
1457 }
1458 match item.get("secret_references") {
1459 None | Some(Value::Null) => {}
1460 Some(Value::Array(references)) if references.is_empty() => {}
1461 Some(Value::Array(_)) => {
1462 reasons.insert("secret_references");
1463 }
1464 Some(_) => {
1465 return Err(http(format!(
1466 "decoding integration policy '{id}': secret_references must be an array or null"
1467 )));
1468 }
1469 }
1470 let active = if active_space.is_empty() {
1471 "default"
1472 } else {
1473 active_space
1474 };
1475 match item.get("spaceIds") {
1476 None | Some(Value::Null) => {}
1477 Some(Value::Array(spaces)) => {
1478 for space in spaces {
1479 let space = space.as_str().filter(|space| !space.is_empty()).ok_or_else(|| {
1480 http(format!(
1481 "decoding integration policy '{id}': spaceIds must contain non-empty strings"
1482 ))
1483 })?;
1484 if space != active {
1485 reasons.insert("spaceIds");
1486 }
1487 }
1488 }
1489 Some(_) => {
1490 return Err(http(format!(
1491 "decoding integration policy '{id}': spaceIds must be an array or null"
1492 )));
1493 }
1494 }
1495 if let Some(policy_id) = item.get("policy_id").filter(|value| !value.is_null()) {
1496 let policy_id = policy_id
1497 .as_str()
1498 .filter(|value| !value.is_empty())
1499 .ok_or_else(|| {
1500 http(format!(
1501 "decoding integration policy '{id}': policy_id must be a non-empty string"
1502 ))
1503 })?;
1504 let policy_ids = item
1505 .get("policy_ids")
1506 .and_then(Value::as_array)
1507 .ok_or_else(|| {
1508 http(format!(
1509 "decoding integration policy '{id}': policy_ids must be an array"
1510 ))
1511 })?;
1512 if policy_ids.first().and_then(Value::as_str) != Some(policy_id) {
1513 return Err(http(format!(
1514 "decoding integration policy '{id}': policy_id must equal policy_ids[0]"
1515 )));
1516 }
1517 }
1518 if reasons.is_empty() {
1519 Ok(())
1520 } else {
1521 unsupported(format!(
1522 "integration policy '{id}' is not portable: {}",
1523 reasons.into_iter().collect::<Vec<_>>().join(", ")
1524 ))
1525 }
1526}
1527
1528fn required_true(item: &Map<String, Value>, field: &str, id: &str) -> Result<()> {
1529 match item.get(field) {
1530 Some(Value::Bool(true)) => Ok(()),
1531 Some(Value::Bool(false)) => unsupported(format!(
1532 "integration policy '{id}' is not portable: {field}"
1533 )),
1534 _ => Err(http(format!(
1535 "decoding integration policy '{id}': {field} must be true"
1536 ))),
1537 }
1538}
1539
1540fn optional_bool(item: &Map<String, Value>, field: &str, id: &str) -> Result<Option<bool>> {
1541 match item.get(field) {
1542 None | Some(Value::Null) => Ok(None),
1543 Some(Value::Bool(value)) => Ok(Some(*value)),
1544 Some(_) => Err(http(format!(
1545 "decoding integration policy '{id}': {field} must be a boolean or null"
1546 ))),
1547 }
1548}
1549
1550fn reject_unknown_top_level(item: &Map<String, Value>, id: &str) -> Result<()> {
1551 let known: BTreeSet<&str> = ["id", "name", "policy_ids", "package", "inputs"]
1552 .into_iter()
1553 .chain(PORTABLE_OPTIONAL)
1554 .chain(REMOVED_FIELDS)
1555 .collect();
1556 if let Some(field) = item
1557 .keys()
1558 .map(String::as_str)
1559 .filter(|field| !known.contains(field))
1560 .min()
1561 {
1562 return unsupported(format!(
1563 "integration policy '{id}' carries unknown field '{field}'"
1564 ));
1565 }
1566 Ok(())
1567}
1568
1569fn summary_from_item(item: &Map<String, Value>) -> Result<IntegrationPolicySummary> {
1570 let id = required_string(item, "id", "integration policy")?;
1571 let name = required_string(item, "name", "integration policy")?;
1572 let namespace = required_string(item, "namespace", "integration policy")?;
1573 let description = match item.get("description") {
1574 None | Some(Value::Null) => None,
1575 Some(Value::String(value)) => Some(value.clone()),
1576 Some(_) => {
1577 return Err(http(
1578 "decoding integration policy: description must be a string or null",
1579 ));
1580 }
1581 };
1582 let policy_ids = item
1583 .get("policy_ids")
1584 .and_then(Value::as_array)
1585 .ok_or_else(|| http("decoding integration policy: policy_ids must be an array"))?
1586 .iter()
1587 .map(|value| {
1588 value
1589 .as_str()
1590 .filter(|value| !value.is_empty())
1591 .map(str::to_owned)
1592 .ok_or_else(|| {
1593 http("decoding integration policy: policy_ids must contain non-empty strings")
1594 })
1595 })
1596 .collect::<Result<Vec<_>>>()?;
1597 let package = package_coordinate(item, "integration policy")?;
1598 Ok(IntegrationPolicySummary {
1599 id,
1600 name,
1601 namespace,
1602 description,
1603 policy_ids,
1604 package,
1605 })
1606}
1607
1608fn package_coordinate(item: &Map<String, Value>, context: &str) -> Result<IntegrationPackageSpec> {
1609 let package = item
1610 .get("package")
1611 .and_then(Value::as_object)
1612 .ok_or_else(|| http(format!("decoding {context}: package must be an object")))?;
1613 Ok(IntegrationPackageSpec {
1614 name: package_required_string(package, "name", context)?,
1615 version: package_required_string(package, "version", context)?,
1616 })
1617}
1618
1619fn package_required_string(
1620 package: &Map<String, Value>,
1621 field: &str,
1622 context: &str,
1623) -> Result<String> {
1624 package
1625 .get(field)
1626 .and_then(Value::as_str)
1627 .filter(|value| !value.trim().is_empty())
1628 .map(str::to_owned)
1629 .ok_or_else(|| {
1630 http(format!(
1631 "decoding {context}: package.{field} must be a non-empty string"
1632 ))
1633 })
1634}
1635
1636fn required_string(item: &Map<String, Value>, field: &str, context: &str) -> Result<String> {
1637 item.get(field)
1638 .and_then(Value::as_str)
1639 .filter(|value| !value.trim().is_empty())
1640 .map(str::to_owned)
1641 .ok_or_else(|| {
1642 http(format!(
1643 "decoding {context}: {field} must be a non-empty string"
1644 ))
1645 })
1646}
1647
1648pub fn prepare_import(path: &Path) -> Result<IntegrationPolicyImportArtifact> {
1651 let canonical = validate(path)?;
1652 if canonical.is_empty() {
1653 return Err(Error::new(
1654 ErrorKind::Error,
1655 "integration-policy import needs at least one integration policy",
1656 ));
1657 }
1658 validate_requested_package_versions(&canonical)?;
1659 Ok(IntegrationPolicyImportArtifact {
1660 source: path.to_path_buf(),
1661 canonical,
1662 })
1663}
1664
1665pub async fn plan_prepared_import(
1668 transport: &Transport,
1669 artifact: IntegrationPolicyImportArtifact,
1670 overwrite: bool,
1671 skip_existing: bool,
1672) -> Result<IntegrationPolicyImportPlan> {
1673 let IntegrationPolicyImportArtifact { source, canonical } = artifact;
1674 if overwrite && skip_existing {
1675 return Err(Error::new(
1676 ErrorKind::Error,
1677 "--overwrite and --skip-existing cannot be used together",
1678 ));
1679 }
1680
1681 let mut existing = BTreeMap::new();
1685 let mut conflicts = Vec::new();
1686 for spec in &canonical {
1687 match integration_policies::get(transport, &spec.id).await {
1688 Ok(policy) => {
1689 let returned_id = required_string(&policy.item, "id", "integration policy get")?;
1690 if returned_id != spec.id {
1691 return Err(http(
1692 "decoding integration policy get: response id did not match the request",
1693 ));
1694 }
1695 if !overwrite && !skip_existing {
1696 conflicts.push(spec.id.clone());
1697 }
1698 existing.insert(spec.id.clone(), Some(policy.item));
1699 }
1700 Err(error) if error.kind == ErrorKind::NotFound => {
1701 existing.insert(spec.id.clone(), None);
1702 }
1703 Err(error) => return Err(import_remote_error(error, "planning read")),
1704 }
1705 }
1706 if !conflicts.is_empty() {
1707 return Err(Error::new(
1708 ErrorKind::Conflict,
1709 format!(
1710 "integration policies already exist: {}",
1711 conflicts.join(", ")
1712 ),
1713 ));
1714 }
1715
1716 for spec in &canonical {
1720 if skip_existing && matches!(existing.get(&spec.id), Some(Some(_))) {
1721 continue;
1722 }
1723 let Some(Some(item)) = existing.get(&spec.id) else {
1724 continue;
1725 };
1726 let current = package_coordinate(item, "integration policy")
1727 .map_err(|error| import_remote_error(error, "planning read"))?;
1728 if current != spec.package {
1729 return unsupported(format!(
1730 "integration policy '{}' cannot change package {}@{} to {}@{}",
1731 spec.id, current.name, current.version, spec.package.name, spec.package.version
1732 ));
1733 }
1734 }
1735
1736 let names: BTreeSet<String> = canonical.iter().map(|spec| spec.name.clone()).collect();
1740 let name_owners = relevant_name_owners(transport, &names)
1741 .await
1742 .map_err(|error| import_remote_error(error, "planning names read"))?;
1743 let mut name_conflicts = Vec::new();
1744 for spec in &canonical {
1745 let owners = name_owners
1746 .get(&spec.name)
1747 .expect("requested name has an ownership entry");
1748 for owner in owners.iter().filter(|owner| owner.as_str() != spec.id) {
1749 name_conflicts.push(format!("{} ({owner})", spec.name));
1750 }
1751 }
1752 if !name_conflicts.is_empty() {
1753 return Err(Error::new(
1754 ErrorKind::Conflict,
1755 format!(
1756 "integration policy names already exist: {}",
1757 name_conflicts.join(", ")
1758 ),
1759 ));
1760 }
1761
1762 let mut skipped = Vec::new();
1763 let pending: Vec<IntegrationPolicySpec> = canonical
1764 .iter()
1765 .filter_map(|spec| match existing.get(&spec.id) {
1766 Some(Some(item)) if skip_existing => {
1767 skipped.push(json!({"id": spec.id, "reason": "exists"}));
1768 None
1769 }
1770 _ => Some(spec.clone()),
1771 })
1772 .collect();
1773
1774 let mut targets = Vec::with_capacity(pending.len());
1775 let mut shared_parents = BTreeMap::new();
1776 for spec in pending {
1777 let raw = existing
1778 .get(&spec.id)
1779 .expect("every canonical id was fetched")
1780 .clone();
1781 let current_parent_ids = raw
1782 .as_ref()
1783 .map(|item| {
1784 read_parents(&spec.id, item)
1785 .map_err(|error| import_remote_error(error, "planning read"))
1786 })
1787 .transpose()?;
1788 let parents = read_import_parent_snapshots(
1789 transport,
1790 &spec.id,
1791 current_parent_ids.as_deref().unwrap_or_default(),
1792 &spec.policy_ids,
1793 )
1794 .await
1795 .map_err(|error| import_remote_error(error, "planning parent read"))?;
1796 for (parent_id, parent) in &parents {
1797 match shared_parents.entry(parent_id.clone()) {
1798 std::collections::btree_map::Entry::Vacant(entry) => {
1799 entry.insert(parent.clone());
1800 }
1801 std::collections::btree_map::Entry::Occupied(entry) if entry.get() != parent => {
1802 return Err(Error::new(
1803 ErrorKind::Conflict,
1804 format!(
1805 "agent policy '{parent_id}' changed while planning integration import"
1806 ),
1807 ));
1808 }
1809 std::collections::btree_map::Entry::Occupied(_) => {}
1810 }
1811 }
1812 let effective = effective_import_spec(&spec, &parents)?;
1813 let current = raw
1814 .map(|item| {
1815 let normalized = normalize(&item, transport.space())
1816 .map_err(|error| import_remote_error(error, "planning read"))?;
1817 if normalized.id != spec.id {
1818 return Err(http(
1819 "decoding integration policy: response id did not match the request",
1820 ));
1821 }
1822 Ok(IntegrationPolicyCurrentSnapshot {
1823 item,
1824 spec: normalized,
1825 parent_ids: current_parent_ids.expect("raw policy has parents"),
1826 })
1827 })
1828 .transpose()?;
1829 targets.push(IntegrationPolicyImportTarget {
1830 effective,
1831 current,
1832 parents,
1833 replacement_body: None,
1834 });
1835 }
1836 if targets
1837 .iter()
1838 .any(|target| !target_name_owners_match(target, &name_owners))
1839 {
1840 return Err(Error::new(
1841 ErrorKind::Conflict,
1842 "integration policy name ownership changed while planning",
1843 ));
1844 }
1845
1846 let mut package_coordinates = BTreeMap::new();
1847 for target in &targets {
1848 package_coordinates
1849 .entry(target.effective.package.name.clone())
1850 .or_insert_with(|| target.effective.package.clone());
1851 }
1852 let mut package_groups = BTreeMap::new();
1853 for (name, package) in package_coordinates {
1854 let state = read_dependencies(transport, &package)
1855 .await
1856 .map_err(|error| import_remote_error(error, "planning package read"))?;
1857 match &state.state {
1858 PackageDependencyState::Installed { version } if version == &package.version => {}
1859 PackageDependencyState::Installed { .. } => {
1860 return Err(Error::new(
1861 ErrorKind::Conflict,
1862 format!("integration package {name} has a different installed version"),
1863 ));
1864 }
1865 PackageDependencyState::NotInstalled => {
1866 if targets
1867 .iter()
1868 .any(|target| target.effective.package.name == name && target.current.is_some())
1869 {
1870 return Err(Error::new(
1871 ErrorKind::Conflict,
1872 format!("integration package {name} is not installed"),
1873 ));
1874 }
1875 }
1876 }
1877 let metadata =
1878 integration_policies::package_metadata(transport, &package.name, &package.version)
1879 .await
1880 .map_err(|error| import_remote_error(error, "planning package metadata read"))?
1881 .item;
1882 validate_package_metadata_snapshot(&metadata, &package)
1883 .map_err(|error| import_remote_error(error, "planning package metadata read"))?;
1884 package_groups.insert(
1885 name,
1886 IntegrationPackageGroup {
1887 package,
1888 state: state.clone(),
1889 state_snapshot: state,
1890 metadata_snapshot: metadata.clone(),
1891 metadata,
1892 },
1893 );
1894 }
1895
1896 for target in &targets {
1897 let package = package_groups
1898 .get(&target.effective.package.name)
1899 .expect("every effective package has a group");
1900 validate_effective_input_materialization(&target.effective, &package.metadata)?;
1901 }
1902
1903 let mut secret_paths = BTreeSet::new();
1904 for target in &targets {
1905 let package = package_groups
1906 .get(&target.effective.package.name)
1907 .expect("every effective package has a group");
1908 if let Some(current) = &target.current {
1909 for path in configured_secret_paths(¤t.spec, &package.metadata)? {
1910 secret_paths.insert(format!("{}:{path}", current.spec.id));
1911 }
1912 }
1913 for path in configured_secret_paths(&target.effective, &package.metadata)? {
1914 secret_paths.insert(format!("{}:{path}", target.effective.id));
1915 }
1916 }
1917 if !secret_paths.is_empty() {
1918 return unsupported(format!(
1919 "integration policy import contains configured secrets: {}",
1920 secret_paths.into_iter().collect::<Vec<_>>().join(", ")
1921 ));
1922 }
1923
1924 for target in &mut targets {
1925 if let Some(current) = &target.current
1926 && current.spec != target.effective
1927 {
1928 target.replacement_body = Some(replace_wire_body(&target.effective)?);
1929 }
1930 }
1931 let package_installs = planned_package_installs(&package_groups);
1932 let preview = import_preview(&source, &targets, &package_installs);
1933 let plan = IntegrationPolicyImportPlan {
1934 preview,
1935 skipped_snapshot: skipped.clone(),
1936 skipped,
1937 package_installs,
1938 total: canonical.len(),
1939 source,
1940 host: transport.kibana_url().to_owned(),
1941 space: transport.space().to_owned(),
1942 canonical,
1943 name_owners_snapshot: name_owners.clone(),
1944 name_owners,
1945 parent_snapshots: shared_parents,
1946 existing_snapshot: existing.clone(),
1947 targets,
1948 package_groups,
1949 overwrite,
1950 skip_existing,
1951 };
1952 validate_import_plan(&plan)?;
1953 Ok(plan)
1954}
1955
1956pub async fn plan_import(
1960 transport: &Transport,
1961 path: &Path,
1962 overwrite: bool,
1963 skip_existing: bool,
1964) -> Result<IntegrationPolicyImportPlan> {
1965 let artifact = prepare_import(path)?;
1966 plan_prepared_import(transport, artifact, overwrite, skip_existing).await
1967}
1968
1969fn validate_requested_package_versions(specs: &[IntegrationPolicySpec]) -> Result<()> {
1970 let mut versions = BTreeMap::new();
1971 for spec in specs {
1972 match versions.entry(spec.package.name.as_str()) {
1973 std::collections::btree_map::Entry::Vacant(entry) => {
1974 entry.insert(spec.package.version.as_str());
1975 }
1976 std::collections::btree_map::Entry::Occupied(entry)
1977 if entry.get() != &spec.package.version.as_str() =>
1978 {
1979 return Err(Error::new(
1980 ErrorKind::Conflict,
1981 format!(
1982 "integration package '{}' is requested at more than one version",
1983 spec.package.name
1984 ),
1985 ));
1986 }
1987 std::collections::btree_map::Entry::Occupied(_) => {}
1988 }
1989 }
1990 Ok(())
1991}
1992
1993async fn relevant_name_owners(
1994 transport: &Transport,
1995 names: &BTreeSet<String>,
1996) -> Result<BTreeMap<String, BTreeSet<String>>> {
1997 let mut owners = names
1998 .iter()
1999 .map(|name| (name.clone(), BTreeSet::new()))
2000 .collect::<BTreeMap<_, _>>();
2001 if names.is_empty() {
2002 return Ok(owners);
2003 }
2004 for item in collect(transport).await? {
2005 let Some(name) = item.get("name").and_then(Value::as_str) else {
2006 continue;
2007 };
2008 let Some(owners) = owners.get_mut(name) else {
2009 continue;
2010 };
2011 owners.insert(required_string(&item, "id", "integration policies list")?);
2012 }
2013 Ok(owners)
2014}
2015
2016fn target_name_owners_match(
2017 target: &IntegrationPolicyImportTarget,
2018 owners: &BTreeMap<String, BTreeSet<String>>,
2019) -> bool {
2020 let Some(actual) = owners.get(&target.effective.name) else {
2021 return false;
2022 };
2023 let expected = match &target.current {
2024 None => BTreeSet::new(),
2025 Some(current) if current.spec.name == target.effective.name => {
2026 BTreeSet::from([target.effective.id.clone()])
2027 }
2028 Some(_) => BTreeSet::new(),
2029 };
2030 actual == &expected
2031}
2032
2033async fn read_import_parent_snapshots(
2034 transport: &Transport,
2035 integration_id: &str,
2036 current_parent_ids: &[String],
2037 desired_parent_ids: &[String],
2038) -> Result<BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>> {
2039 let parent_ids = current_parent_ids
2040 .iter()
2041 .chain(desired_parent_ids)
2042 .cloned()
2043 .collect::<BTreeSet<_>>();
2044 let mut parents = BTreeMap::new();
2045 for parent_id in parent_ids {
2046 let parent = agent_policy_ops::read_parent_snapshot(transport, &parent_id).await?;
2047 if current_parent_ids.binary_search(&parent_id).is_ok()
2048 && parent
2049 .attached_integrations
2050 .binary_search_by(|attached| attached.as_str().cmp(integration_id))
2051 .is_err()
2052 {
2053 return Err(http(format!(
2054 "decoding integration policy '{integration_id}': parent '{parent_id}' is missing its attachment"
2055 )));
2056 }
2057 parents.insert(parent_id, parent);
2058 }
2059 Ok(parents)
2060}
2061
2062fn effective_import_spec(
2063 canonical: &IntegrationPolicySpec,
2064 parents: &BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
2065) -> Result<IntegrationPolicySpec> {
2066 canonical.validate()?;
2067 for parent in parents.values() {
2068 if parent.platform_owned {
2069 return unsupported(format!(
2070 "integration policy '{}' is not portable: parent {} is platform-owned",
2071 canonical.id, parent.id
2072 ));
2073 }
2074 if parent.protected {
2075 return unsupported(format!(
2076 "integration policy '{}' is not portable: parent {} is_protected",
2077 canonical.id, parent.id
2078 ));
2079 }
2080 }
2081 let selected = canonical
2082 .policy_ids
2083 .iter()
2084 .map(|id| {
2085 parents.get(id).ok_or_else(|| {
2086 Error::new(
2087 ErrorKind::Error,
2088 format!(
2089 "integration policy '{}' has no parent snapshot for '{id}'",
2090 canonical.id
2091 ),
2092 )
2093 })
2094 })
2095 .collect::<Result<Vec<_>>>()?;
2096 let mut effective = canonical.clone();
2097 if let Some(namespace) = &effective.namespace {
2098 if selected.iter().any(|parent| &parent.namespace != namespace) {
2099 return unsupported(format!(
2100 "integration policy '{}' is not portable: namespace does not match every parent",
2101 canonical.id
2102 ));
2103 }
2104 } else {
2105 let namespaces = selected
2106 .iter()
2107 .map(|parent| parent.namespace.as_str())
2108 .collect::<BTreeSet<_>>();
2109 if namespaces.len() != 1 {
2110 return Err(Error::new(
2111 ErrorKind::Conflict,
2112 format!(
2113 "integration policy '{}' is not portable: parents have different namespaces",
2114 canonical.id
2115 ),
2116 ));
2117 }
2118 effective.namespace = namespaces.into_iter().next().map(str::to_owned);
2119 }
2120 Ok(effective)
2121}
2122
2123fn validate_package_metadata_snapshot(
2124 metadata: &Map<String, Value>,
2125 package: &IntegrationPackageSpec,
2126) -> Result<()> {
2127 let name = metadata_name(metadata, "name", "package metadata")?;
2128 let version = metadata_name(metadata, "version", "package metadata")?;
2129 if name != package.name || version != package.version {
2130 return Err(http(format!(
2131 "decoding package metadata: expected {}@{}, got {name}@{version}",
2132 package.name, package.version
2133 )));
2134 }
2135 secret_schema(metadata).map(|_| ())
2136}
2137
2138fn validate_effective_input_materialization(
2139 spec: &IntegrationPolicySpec,
2140 metadata: &Map<String, Value>,
2141) -> Result<()> {
2142 let (_, known) = secret_schema(metadata)?;
2143 if spec.inputs.is_empty() && !known.input_vars.is_empty() {
2144 return unsupported(format!(
2145 "integration policy '{}' has an empty inputs map but package {}@{} declares inputs",
2146 spec.id, spec.package.name, spec.package.version
2147 ));
2148 }
2149 Ok(())
2150}
2151
2152fn replace_wire_body(spec: &IntegrationPolicySpec) -> Result<Value> {
2153 spec.validate()?;
2154 let mut body = serde_json::to_value(spec)
2155 .map_err(|error| {
2156 Error::new(
2157 ErrorKind::Error,
2158 format!("encoding integration policy: {error}"),
2159 )
2160 })?
2161 .as_object()
2162 .cloned()
2163 .expect("integration policy specs serialize to objects");
2164 body.remove("id");
2165 Ok(Value::Object(body))
2166}
2167
2168fn planned_package_installs(groups: &BTreeMap<String, IntegrationPackageGroup>) -> Vec<String> {
2169 groups
2170 .values()
2171 .filter_map(|group| match group.state.state {
2172 PackageDependencyState::NotInstalled => {
2173 Some(format!("{}@{}", group.package.name, group.package.version))
2174 }
2175 PackageDependencyState::Installed { .. } => None,
2176 })
2177 .collect()
2178}
2179
2180fn import_preview(
2181 path: &Path,
2182 targets: &[IntegrationPolicyImportTarget],
2183 package_installs: &[String],
2184) -> MutationPlan {
2185 let mut details = targets
2186 .iter()
2187 .map(|target| {
2188 let parents = target
2189 .parents
2190 .values()
2191 .map(|parent| format!("{} ({})", parent.id, parent.name))
2192 .collect::<Vec<_>>()
2193 .join(", ");
2194 let agents = target
2195 .parents
2196 .values()
2197 .map(|parent| parent.agents)
2198 .sum::<u64>();
2199 let action = match &target.current {
2200 None => "create".to_owned(),
2201 Some(current) if current.spec == target.effective => "unchanged".to_owned(),
2202 Some(current) if current.spec.name == target.effective.name => "replace".to_owned(),
2203 Some(current) => format!(
2204 "replace {} -> {}",
2205 current.spec.name, target.effective.name
2206 ),
2207 };
2208 format!(
2209 "{} {action} {} parents {parents} agents {agents}",
2210 target.effective.id, target.effective.name
2211 )
2212 })
2213 .collect::<Vec<_>>();
2214 details.extend(
2215 package_installs
2216 .iter()
2217 .map(|package| format!("package install {package}")),
2218 );
2219 details.push(IMPORT_RACE_WARNING.to_owned());
2220 MutationPlan {
2221 preview_action: format!(
2222 "Import {} integration policy(ies) from {}",
2223 targets.len(),
2224 path.display()
2225 ),
2226 preview_details: details,
2227 targets: targets
2228 .iter()
2229 .map(|target| target.effective.id.clone())
2230 .collect(),
2231 }
2232}
2233
2234pub fn validate(path: &Path) -> Result<Vec<IntegrationPolicySpec>> {
2236 let body = std::fs::read_to_string(path).map_err(|error| {
2237 Error::new(
2238 ErrorKind::Error,
2239 format!("reading {}: {error}", path.display()),
2240 )
2241 })?;
2242 let mut specs = content_codec::decode_sequence::<IntegrationPolicySpec>(
2243 &body,
2244 ContentFormat::from_path(path),
2245 "integration policy",
2246 )?;
2247 duplicate_error(&specs, |spec| &spec.id, "ids")?;
2248 duplicate_error(&specs, |spec| &spec.name, "names")?;
2249 specs.sort_by(|left, right| left.id.cmp(&right.id));
2250 Ok(specs)
2251}
2252
2253fn duplicate_error<'a, F>(specs: &'a [IntegrationPolicySpec], key: F, noun: &str) -> Result<()>
2254where
2255 F: Fn(&'a IntegrationPolicySpec) -> &'a String,
2256{
2257 let mut seen = BTreeSet::new();
2258 let mut duplicates = BTreeSet::new();
2259 for spec in specs {
2260 let value = key(spec);
2261 if !seen.insert(value.as_str()) {
2262 duplicates.insert(value.as_str());
2263 }
2264 }
2265 if duplicates.is_empty() {
2266 Ok(())
2267 } else {
2268 Err(Error::new(
2269 ErrorKind::Error,
2270 format!(
2271 "duplicate integration policy {noun}: {}",
2272 duplicates.into_iter().collect::<Vec<_>>().join(", ")
2273 ),
2274 ))
2275 }
2276}
2277
2278pub async fn apply_import(
2281 transport: &Transport,
2282 plan: &IntegrationPolicyImportPlan,
2283) -> Result<IntegrationPolicyImportReport> {
2284 validate_import_plan(plan)?;
2285 if plan.host != transport.kibana_url() || plan.space != transport.space() {
2286 return Err(Error::new(
2287 ErrorKind::Conflict,
2288 "integration import target changed since preview",
2289 ));
2290 }
2291 let mut succeeded = Vec::new();
2292 let mut unchanged = Vec::new();
2293 let mut failed = Vec::new();
2294 let mut expected_groups = plan.package_groups.clone();
2295 let mut expected_parents = plan.parent_snapshots.clone();
2296 let mut blocked_packages = BTreeMap::<String, String>::new();
2297 let mut affected_parents = BTreeMap::<String, u64>::new();
2298 let mut observed_installs = BTreeSet::new();
2299
2300 for target in &plan.targets {
2301 let package_name = &target.effective.package.name;
2302 if let Some(error) = blocked_packages.get(package_name) {
2303 failed.push(import_failed_row(
2304 &target.effective.id,
2305 false,
2306 format!("package dependency is unavailable: {error}"),
2307 ));
2308 continue;
2309 }
2310
2311 let action = match recheck_import_object(transport, target).await {
2312 Ok(action) => action,
2313 Err(error) => {
2314 failed.push(import_failed_row(
2315 &target.effective.id,
2316 false,
2317 error.message,
2318 ));
2319 continue;
2320 }
2321 };
2322 if let Err(error) = recheck_import_name_owner(transport, target, &plan.name_owners).await {
2323 failed.push(import_failed_row(
2324 &target.effective.id,
2325 false,
2326 error.message,
2327 ));
2328 continue;
2329 }
2330 if let Err(error) = recheck_import_parents(transport, target, &expected_parents).await {
2331 failed.push(import_failed_row(
2332 &target.effective.id,
2333 false,
2334 error.message,
2335 ));
2336 continue;
2337 }
2338
2339 let group = expected_groups
2340 .get_mut(package_name)
2341 .expect("validated target package group");
2342 let actual_state = match read_dependencies(transport, &target.effective.package).await {
2343 Ok(state) if state == group.state => state,
2344 Ok(_) => {
2345 let message = "package changed since preview".to_owned();
2346 blocked_packages.insert(package_name.clone(), message.clone());
2347 failed.push(import_failed_row(&target.effective.id, false, message));
2348 continue;
2349 }
2350 Err(error) => {
2351 let message = import_remote_error(error, "apply package read").message;
2352 blocked_packages.insert(package_name.clone(), message.clone());
2353 failed.push(import_failed_row(&target.effective.id, false, message));
2354 continue;
2355 }
2356 };
2357 debug_assert_eq!(actual_state, group.state);
2358
2359 if action == ImportAction::Unchanged {
2360 unchanged.push(json!({"id": target.effective.id}));
2361 continue;
2362 }
2363
2364 let (label, applied, route_error) = match action {
2365 ImportAction::Create => {
2366 match integration_policies::create(transport, &target.effective).await {
2367 Ok(_) => ("created", true, None),
2368 Err(error) => (
2369 "created",
2370 false,
2371 Some(import_remote_error(error, "create request").message),
2372 ),
2373 }
2374 }
2375 ImportAction::Replace => {
2376 let _body = target
2377 .replacement_body
2378 .as_ref()
2379 .expect("validated replacement body");
2380 match integration_policies::update(
2381 transport,
2382 &target.effective.id,
2383 &target.effective,
2384 )
2385 .await
2386 {
2387 Ok(_) => ("replaced", true, None),
2388 Err(error) => (
2389 "replaced",
2390 false,
2391 Some(import_remote_error(error, "update request").message),
2392 ),
2393 }
2394 }
2395 ImportAction::Unchanged => unreachable!("unchanged rows continue above"),
2396 };
2397
2398 if applied {
2399 record_affected_parents(&mut affected_parents, target);
2400 advance_parent_snapshots(&mut expected_parents, target);
2401 }
2402
2403 let mut observed_after_create = None;
2407 let mut package_observation_error = None;
2408 if action == ImportAction::Create
2409 && matches!(group.state.state, PackageDependencyState::NotInstalled)
2410 {
2411 match read_dependencies(transport, &target.effective.package).await {
2412 Ok(after) => {
2413 if is_exact_installed(&after, &target.effective.package) {
2414 group.state = after.clone();
2415 observed_installs.insert(format!(
2416 "{}@{}",
2417 target.effective.package.name, target.effective.package.version
2418 ));
2419 } else if !matches!(after.state, PackageDependencyState::NotInstalled) {
2420 let message =
2421 "package installed a different version after create".to_owned();
2422 blocked_packages.insert(package_name.clone(), message.clone());
2423 package_observation_error = Some(message);
2424 }
2425 observed_after_create = Some(after);
2426 }
2427 Err(error) => {
2428 let message = import_remote_error(error, "post-create package read").message;
2429 blocked_packages.insert(package_name.clone(), message.clone());
2430 package_observation_error = Some(message);
2431 }
2432 }
2433 }
2434
2435 let mut errors = route_error.into_iter().collect::<Vec<_>>();
2436 if applied {
2437 if let Err(error) = verify_import_stored(transport, &target.effective).await {
2438 errors.push(error.message);
2439 }
2440 let package_result = match observed_after_create.as_ref() {
2441 Some(after) => verify_exact_installed(after, &target.effective.package),
2442 None => match read_dependencies(transport, &target.effective.package).await {
2443 Ok(after) => verify_exact_installed(&after, &target.effective.package),
2444 Err(error) => Err(import_remote_error(error, "package verification").message),
2445 },
2446 };
2447 if let Err(message) = package_result {
2448 blocked_packages
2449 .entry(package_name.clone())
2450 .or_insert_with(|| message.clone());
2451 errors.push(message);
2452 }
2453 }
2454 if let Some(error) = package_observation_error
2455 && !errors.contains(&error)
2456 {
2457 errors.push(error);
2458 }
2459
2460 if errors.is_empty() {
2461 succeeded.push(json!({"id": target.effective.id, "action": label}));
2462 } else {
2463 failed.push(import_failed_row(
2464 &target.effective.id,
2465 applied,
2466 errors.join("; "),
2467 ));
2468 }
2469 }
2470
2471 Ok(IntegrationPolicyImportReport {
2472 applied: true,
2473 succeeded,
2474 unchanged,
2475 skipped: plan.skipped.clone(),
2476 failed,
2477 total: plan.total,
2478 affected_agents: affected_parents.values().sum(),
2479 package_installs: observed_installs.into_iter().collect(),
2480 })
2481}
2482
2483async fn recheck_import_object(
2484 transport: &Transport,
2485 target: &IntegrationPolicyImportTarget,
2486) -> Result<ImportAction> {
2487 match &target.current {
2488 None => match integration_policies::get(transport, &target.effective.id).await {
2489 Err(error) if error.kind == ErrorKind::NotFound => Ok(ImportAction::Create),
2490 Ok(_) => Err(Error::new(
2491 ErrorKind::Conflict,
2492 "integration policy appeared since preview",
2493 )),
2494 Err(error) => Err(import_remote_error(error, "apply integration-policy read")),
2495 },
2496 Some(expected) => match integration_policies::get(transport, &target.effective.id).await {
2497 Ok(actual) if actual.item == expected.item => {
2498 if expected.spec == target.effective {
2499 Ok(ImportAction::Unchanged)
2500 } else {
2501 Ok(ImportAction::Replace)
2502 }
2503 }
2504 Ok(_) => Err(Error::new(
2505 ErrorKind::Conflict,
2506 "integration policy changed since preview",
2507 )),
2508 Err(error) if error.kind == ErrorKind::NotFound => Err(Error::new(
2509 ErrorKind::Conflict,
2510 "integration policy disappeared since preview",
2511 )),
2512 Err(error) => Err(import_remote_error(error, "apply integration-policy read")),
2513 },
2514 }
2515}
2516
2517async fn recheck_import_name_owner(
2518 transport: &Transport,
2519 target: &IntegrationPolicyImportTarget,
2520 expected_owners: &BTreeMap<String, BTreeSet<String>>,
2521) -> Result<()> {
2522 let names = BTreeSet::from([target.effective.name.clone()]);
2523 let owners = relevant_name_owners(transport, &names)
2524 .await
2525 .map_err(|error| import_remote_error(error, "apply name read"))?;
2526 if owners.get(&target.effective.name) == expected_owners.get(&target.effective.name) {
2527 Ok(())
2528 } else {
2529 Err(Error::new(
2530 ErrorKind::Conflict,
2531 "integration policy name ownership changed since preview",
2532 ))
2533 }
2534}
2535
2536async fn recheck_import_parents(
2537 transport: &Transport,
2538 target: &IntegrationPolicyImportTarget,
2539 expected_parents: &BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
2540) -> Result<()> {
2541 for parent_id in target.parents.keys() {
2542 let expected = expected_parents.get(parent_id).ok_or_else(|| {
2543 Error::new(
2544 ErrorKind::Error,
2545 "integration import lost a shared parent snapshot",
2546 )
2547 })?;
2548 let actual = agent_policy_ops::read_parent_snapshot(transport, parent_id)
2549 .await
2550 .map_err(|error| import_remote_error(error, "apply parent read"))?;
2551 if actual != *expected {
2552 return Err(Error::new(
2553 ErrorKind::Conflict,
2554 "integration policy parent changed since preview",
2555 ));
2556 }
2557 }
2558 Ok(())
2559}
2560
2561fn record_affected_parents(
2562 affected: &mut BTreeMap<String, u64>,
2563 target: &IntegrationPolicyImportTarget,
2564) {
2565 for parent in target.parents.values() {
2566 affected.entry(parent.id.clone()).or_insert(parent.agents);
2567 }
2568}
2569
2570fn advance_parent_snapshots(
2571 parents: &mut BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
2572 target: &IntegrationPolicyImportTarget,
2573) {
2574 let desired = target
2575 .effective
2576 .policy_ids
2577 .iter()
2578 .map(String::as_str)
2579 .collect::<BTreeSet<_>>();
2580 for parent_id in target.parents.keys() {
2581 let parent = parents
2582 .get_mut(parent_id)
2583 .expect("validated shared parent snapshot");
2584 if desired.contains(parent_id.as_str()) {
2585 if parent
2586 .attached_integrations
2587 .binary_search_by(|attached| attached.as_str().cmp(&target.effective.id))
2588 .is_err()
2589 {
2590 parent
2591 .attached_integrations
2592 .push(target.effective.id.clone());
2593 parent.attached_integrations.sort();
2594 }
2595 } else {
2596 parent
2597 .attached_integrations
2598 .retain(|attached| attached != &target.effective.id);
2599 }
2600 }
2601}
2602
2603async fn verify_import_stored(
2604 transport: &Transport,
2605 desired: &IntegrationPolicySpec,
2606) -> Result<()> {
2607 let stored = integration_policies::get(transport, &desired.id)
2608 .await
2609 .map_err(|error| import_remote_error(error, "stored-policy read"))?;
2610 let stored = normalize(&stored.item, transport.space())
2611 .map_err(|error| import_remote_error(error, "stored-policy read"))?;
2612 if stored == *desired {
2613 Ok(())
2614 } else {
2615 Err(Error::new(
2616 ErrorKind::Http,
2617 "server stored a different integration-policy spec",
2618 ))
2619 }
2620}
2621
2622fn is_exact_installed(
2623 snapshot: &PackageDependencySnapshot,
2624 package: &IntegrationPackageSpec,
2625) -> bool {
2626 snapshot.name == package.name
2627 && matches!(
2628 &snapshot.state,
2629 PackageDependencyState::Installed { version } if version == &package.version
2630 )
2631}
2632
2633fn verify_exact_installed(
2634 snapshot: &PackageDependencySnapshot,
2635 package: &IntegrationPackageSpec,
2636) -> std::result::Result<(), String> {
2637 if is_exact_installed(snapshot, package) {
2638 return Ok(());
2639 }
2640 match &snapshot.state {
2641 PackageDependencyState::Installed { .. } => Err(format!(
2642 "package {} installed a different version",
2643 package.name
2644 )),
2645 PackageDependencyState::NotInstalled => {
2646 Err(format!("package {} is not installed", package.name))
2647 }
2648 }
2649}
2650
2651fn import_remote_error(error: Error, context: &str) -> Error {
2652 let message = format!("integration-policy import {context} failed");
2653 match error.http_status {
2654 Some(status) => Error::with_status(error.kind, status, message),
2655 None => Error::new(error.kind, message),
2656 }
2657}
2658
2659fn import_failed_row(id: &str, applied: bool, error: impl Into<String>) -> Value {
2660 json!({"id": id, "applied": applied, "error": error.into()})
2661}
2662
2663fn validate_import_plan(plan: &IntegrationPolicyImportPlan) -> Result<()> {
2664 let invalid = |message: &str| {
2665 Err(Error::new(
2666 ErrorKind::Error,
2667 format!("invalid integration-policy import plan: {message}"),
2668 ))
2669 };
2670 if plan.overwrite && plan.skip_existing {
2671 return invalid("overwrite and skip-existing cannot both be set");
2672 }
2673 if plan.host.trim().is_empty() {
2674 return invalid("planned Kibana host is empty");
2675 }
2676 if plan.canonical.is_empty() || plan.total != plan.canonical.len() {
2677 return invalid("total does not equal canonical integration policies");
2678 }
2679 let mut canonical_ids = BTreeMap::new();
2680 let mut canonical_names = BTreeSet::new();
2681 let mut previous: Option<&str> = None;
2682 for spec in &plan.canonical {
2683 if spec.validate().is_err() {
2684 return invalid("canonical integration policy is invalid");
2685 }
2686 if previous.is_some_and(|previous| previous >= spec.id.as_str()) {
2687 return invalid("canonical integration policies must be unique and sorted by id");
2688 }
2689 if !canonical_names.insert(spec.name.as_str()) {
2690 return invalid("canonical integration-policy names must be unique");
2691 }
2692 previous = Some(&spec.id);
2693 canonical_ids.insert(spec.id.as_str(), spec);
2694 }
2695 if validate_requested_package_versions(&plan.canonical).is_err() {
2696 return invalid("canonical package requests are inconsistent");
2697 }
2698 let canonical_id_set = canonical_ids.keys().copied().collect::<BTreeSet<_>>();
2699 if plan
2700 .existing_snapshot
2701 .keys()
2702 .map(String::as_str)
2703 .collect::<BTreeSet<_>>()
2704 != canonical_id_set
2705 {
2706 return invalid("existence snapshots do not match canonical integration policies");
2707 }
2708 for (id, existing) in &plan.existing_snapshot {
2709 if let Some(item) = existing
2710 && required_string(item, "id", "existing integration policy")
2711 .ok()
2712 .as_deref()
2713 != Some(id.as_str())
2714 {
2715 return invalid("existence snapshot has an unexpected id");
2716 }
2717 }
2718 if plan.name_owners != plan.name_owners_snapshot {
2719 return invalid("name ownership snapshots do not match");
2720 }
2721 if plan
2722 .name_owners
2723 .keys()
2724 .map(String::as_str)
2725 .collect::<BTreeSet<_>>()
2726 != canonical_names
2727 {
2728 return invalid("name ownership snapshots do not match canonical names");
2729 }
2730 for spec in &plan.canonical {
2731 let Some(owners) = plan.name_owners.get(&spec.name) else {
2732 return invalid("canonical name has no ownership snapshot");
2733 };
2734 if owners
2735 .iter()
2736 .any(|owner| owner.trim().is_empty() || owner != &spec.id)
2737 {
2738 return invalid("name ownership snapshot has a foreign or malformed owner");
2739 }
2740 }
2741
2742 let mut target_ids = BTreeSet::new();
2743 let mut expected_bodies = BTreeMap::new();
2744 let mut expected_group_names = BTreeSet::new();
2745 let mut shared_parents = BTreeMap::new();
2746 let mut previous_target: Option<&str> = None;
2747 for target in &plan.targets {
2748 if target.effective.validate().is_err() {
2749 return invalid("effective integration policy is invalid");
2750 }
2751 if previous_target.is_some_and(|previous| previous >= target.effective.id.as_str()) {
2752 return invalid("pending integration policies must be unique and sorted by id");
2753 }
2754 previous_target = Some(&target.effective.id);
2755 let Some(canonical) = canonical_ids.get(target.effective.id.as_str()) else {
2756 return invalid("pending policy is not in the canonical artifact");
2757 };
2758 if !target_ids.insert(target.effective.id.as_str()) {
2759 return invalid("pending integration policies must be unique and sorted by id");
2760 }
2761 let exact_current = match (
2762 &target.current,
2763 plan.existing_snapshot.get(&target.effective.id),
2764 ) {
2765 (None, Some(None)) => true,
2766 (Some(current), Some(Some(item))) => current.item == *item,
2767 _ => false,
2768 };
2769 if !exact_current {
2770 return invalid("target does not match its plan-time existence snapshot");
2771 }
2772 let current_parent_ids = match &target.current {
2773 None => Vec::new(),
2774 Some(current) => {
2775 if current.spec.validate().is_err()
2776 || current.spec.id != target.effective.id
2777 || normalize(¤t.item, &plan.space).ok().as_ref() != Some(¤t.spec)
2778 {
2779 return invalid("current integration snapshot does not normalize canonically");
2780 }
2781 let parent_ids = match read_parents(&target.effective.id, ¤t.item) {
2782 Ok(parent_ids) if parent_ids == current.parent_ids => parent_ids,
2783 _ => {
2784 return invalid(
2785 "current integration parent snapshot does not match its item",
2786 );
2787 }
2788 };
2789 if package_coordinate(¤t.item, "integration policy").ok()
2790 != Some(target.effective.package.clone())
2791 {
2792 return invalid("current and desired package coordinates differ");
2793 }
2794 parent_ids
2795 }
2796 };
2797 if target.current.is_some() && !plan.overwrite {
2798 return invalid("existing integration target requires overwrite");
2799 }
2800 if !target_name_owners_match(target, &plan.name_owners) {
2801 return invalid("name ownership snapshot does not match target state");
2802 }
2803 let expected_parent_ids = current_parent_ids
2804 .iter()
2805 .chain(&target.effective.policy_ids)
2806 .cloned()
2807 .collect::<BTreeSet<_>>();
2808 if target.parents.keys().cloned().collect::<BTreeSet<_>>() != expected_parent_ids {
2809 return invalid("parent snapshots do not match current and desired parents");
2810 }
2811 for (parent_id, parent) in &target.parents {
2812 if !valid_parent_snapshot(parent_id, parent)
2813 || parent.platform_owned
2814 || parent.protected
2815 {
2816 return invalid("parent snapshot is unsafe or malformed");
2817 }
2818 if current_parent_ids.binary_search(parent_id).is_ok()
2819 && parent
2820 .attached_integrations
2821 .binary_search_by(|attached| attached.as_str().cmp(&target.effective.id))
2822 .is_err()
2823 {
2824 return invalid("current parent snapshot is missing its integration attachment");
2825 }
2826 match shared_parents.entry(parent_id.as_str()) {
2827 std::collections::btree_map::Entry::Vacant(entry) => {
2828 entry.insert(parent);
2829 }
2830 std::collections::btree_map::Entry::Occupied(entry) if *entry.get() != parent => {
2831 return invalid("shared parent snapshots disagree");
2832 }
2833 std::collections::btree_map::Entry::Occupied(_) => {}
2834 }
2835 }
2836 if effective_import_spec(canonical, &target.parents)
2837 .ok()
2838 .as_ref()
2839 != Some(&target.effective)
2840 {
2841 return invalid("effective integration policy does not match canonical parents");
2842 }
2843 if let Some(current) = &target.current {
2844 if current.spec != target.effective {
2845 if !plan.overwrite {
2846 return invalid("replacement plan requires overwrite");
2847 }
2848 let body = match replace_wire_body(&target.effective) {
2849 Ok(body) => body,
2850 Err(_) => return invalid("replacement body cannot be encoded"),
2851 };
2852 if target.replacement_body.as_ref() != Some(&body) {
2853 return invalid("replacement body does not match its effective policy");
2854 }
2855 expected_bodies.insert(target.effective.id.as_str(), body);
2856 } else if target.replacement_body.is_some() {
2857 return invalid("unchanged integration policy carries a replacement body");
2858 }
2859 } else if target.replacement_body.is_some() {
2860 return invalid("planned create carries a replacement body");
2861 }
2862 expected_group_names.insert(target.effective.package.name.as_str());
2863 }
2864
2865 if plan
2866 .package_groups
2867 .keys()
2868 .map(String::as_str)
2869 .collect::<BTreeSet<_>>()
2870 != expected_group_names
2871 {
2872 return invalid("package groups do not match pending integration policies");
2873 }
2874 let expected_parent_snapshots: BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot> =
2875 shared_parents
2876 .iter()
2877 .map(|(id, parent)| ((*id).to_owned(), (*parent).clone()))
2878 .collect();
2879 if plan.parent_snapshots != expected_parent_snapshots {
2880 return invalid("shared parent snapshots do not match pending integrations");
2881 }
2882 for (name, group) in &plan.package_groups {
2883 if group.package.name != *name
2884 || group.package.name.trim().is_empty()
2885 || group.package.version.trim().is_empty()
2886 || group.state != group.state_snapshot
2887 || group.metadata != group.metadata_snapshot
2888 || group.state.name != group.package.name
2889 || !valid_package_state(&group.state)
2890 || validate_package_metadata_snapshot(&group.metadata, &group.package).is_err()
2891 {
2892 return invalid("package group snapshot is malformed or tampered");
2893 }
2894 if !is_exact_installed(&group.state, &group.package)
2895 && !matches!(group.state.state, PackageDependencyState::NotInstalled)
2896 {
2897 return invalid("package group does not hold an exact dependency state");
2898 }
2899 if matches!(group.state.state, PackageDependencyState::NotInstalled)
2900 && plan
2901 .targets
2902 .iter()
2903 .any(|target| target.effective.package.name == *name && target.current.is_some())
2904 {
2905 return invalid("existing integration cannot depend on an absent package");
2906 }
2907 }
2908 for target in &plan.targets {
2909 let Some(group) = plan.package_groups.get(&target.effective.package.name) else {
2910 return invalid("target has no package group");
2911 };
2912 if group.package != target.effective.package {
2913 return invalid("package group coordinate does not match its target");
2914 }
2915 validate_effective_input_materialization(&target.effective, &group.metadata)?;
2916 match configured_secret_paths(&target.effective, &group.metadata) {
2917 Ok(paths) if paths.is_empty() => {}
2918 _ => return invalid("effective integration policy has unsafe configured variables"),
2919 }
2920 if let Some(current) = &target.current {
2921 match configured_secret_paths(¤t.spec, &group.metadata) {
2922 Ok(paths) if paths.is_empty() => {}
2923 _ => {
2924 return invalid("current integration policy has unsafe configured variables");
2925 }
2926 }
2927 }
2928 }
2929 if expected_bodies.len()
2930 != plan
2931 .targets
2932 .iter()
2933 .filter(|target| target.replacement_body.is_some())
2934 .count()
2935 {
2936 return invalid("replacement body set does not match changed integration policies");
2937 }
2938
2939 let expected_skipped_ids = plan
2940 .canonical
2941 .iter()
2942 .filter(|spec| {
2943 plan.skip_existing && matches!(plan.existing_snapshot.get(&spec.id), Some(Some(_)))
2944 })
2945 .map(|spec| spec.id.as_str())
2946 .collect::<BTreeSet<_>>();
2947 let expected_target_ids = canonical_id_set
2948 .iter()
2949 .copied()
2950 .filter(|id| !expected_skipped_ids.contains(id))
2951 .collect::<BTreeSet<_>>();
2952 if target_ids != expected_target_ids {
2953 return invalid("pending integration policies do not match plan-time existence snapshots");
2954 }
2955 let expected_skipped = plan
2956 .canonical
2957 .iter()
2958 .filter(|spec| expected_skipped_ids.contains(spec.id.as_str()))
2959 .map(|spec| json!({"id": spec.id, "reason": "exists"}))
2960 .collect::<Vec<_>>();
2961 if plan.skipped != plan.skipped_snapshot {
2962 return invalid("skipped rows do not match their snapshot");
2963 }
2964 if (!plan.skip_existing && !plan.skipped.is_empty()) || plan.skipped != expected_skipped {
2965 return invalid("skipped rows do not match the canonical artifact");
2966 }
2967 let expected_installs = planned_package_installs(&plan.package_groups);
2968 if plan.package_installs != expected_installs {
2969 return invalid("package install preview does not match package groups");
2970 }
2971 let expected_preview = import_preview(&plan.source, &plan.targets, &plan.package_installs);
2972 if plan.preview != expected_preview {
2973 return invalid("preview does not match the canonical plan");
2974 }
2975 Ok(())
2976}
2977
2978fn valid_parent_snapshot(id: &str, parent: &agent_policy_ops::AgentPolicyParentSnapshot) -> bool {
2979 parent.id == id
2980 && !parent.id.trim().is_empty()
2981 && !parent.name.trim().is_empty()
2982 && !parent.namespace.trim().is_empty()
2983 && parent
2984 .attached_integrations
2985 .windows(2)
2986 .all(|ids| ids[0] < ids[1])
2987}
2988
2989fn valid_package_state(snapshot: &PackageDependencySnapshot) -> bool {
2990 if snapshot.name.trim().is_empty() {
2991 return false;
2992 }
2993 match &snapshot.state {
2994 PackageDependencyState::Installed { version } => !version.trim().is_empty(),
2995 PackageDependencyState::NotInstalled => true,
2996 }
2997}
2998
2999pub async fn plan_delete(
3003 transport: &Transport,
3004 selectors: &[String],
3005) -> Result<IntegrationPolicyDeletePlan> {
3006 if selectors.is_empty() {
3007 return Err(Error::new(
3008 ErrorKind::Error,
3009 "integration-policy delete needs at least one selector",
3010 ));
3011 }
3012 if selectors.iter().any(|selector| selector.trim().is_empty()) {
3013 return Err(Error::new(
3014 ErrorKind::Error,
3015 "integration-policy delete selectors must not be empty",
3016 ));
3017 }
3018
3019 let mut resolved = BTreeMap::new();
3020 for selector in selectors {
3021 let resolved_policy = resolve_delete_item(transport, selector).await?;
3022 let id = required_string(
3023 &resolved_policy.item,
3024 "id",
3025 "integration policy delete planning read",
3026 )?;
3027 if id != resolved_policy.summary.id {
3028 return Err(http(
3029 "decoding integration policy delete planning read: response id did not match its summary",
3030 ));
3031 }
3032 resolved.entry(id).or_insert(resolved_policy);
3033 }
3034
3035 let mut targets = Vec::with_capacity(resolved.len());
3036 let mut issues = Vec::new();
3037 for (id, resolved_policy) in resolved {
3038 match plan_delete_target(transport, &id, resolved_policy.item).await {
3039 Ok(target) => targets.push(target),
3040 Err(error) => issues.push(error),
3041 }
3042 }
3043 if !issues.is_empty() {
3044 return collapse_delete_planning_issues(issues);
3045 }
3046
3047 let parent_snapshots = shared_delete_parents(&targets).map_err(|_| {
3048 Error::new(
3049 ErrorKind::Conflict,
3050 "agent policy changed while planning integration deletion",
3051 )
3052 })?;
3053 let plan = IntegrationPolicyDeletePlan {
3054 preview: delete_preview(&targets),
3055 total: targets.len(),
3056 host_snapshot: transport.kibana_url().to_owned(),
3057 host: transport.kibana_url().to_owned(),
3058 space_snapshot: transport.space().to_owned(),
3059 space: transport.space().to_owned(),
3060 parent_snapshots_snapshot: parent_snapshots.clone(),
3061 parent_snapshots,
3062 targets,
3063 };
3064 validate_delete_plan(&plan)?;
3065 Ok(plan)
3066}
3067
3068async fn resolve_delete_item(
3073 transport: &Transport,
3074 selector: &str,
3075) -> Result<ResolvedIntegrationPolicy> {
3076 match integration_policies::get(transport, selector).await {
3077 Ok(policy) => {
3078 let summary = summary_from_item(&policy.item)?;
3079 if summary.id != selector {
3080 return Err(http(
3081 "decoding integration policy delete planning read: response id did not match the selector",
3082 ));
3083 }
3084 return Ok(ResolvedIntegrationPolicy {
3085 summary,
3086 item: policy.item,
3087 });
3088 }
3089 Err(error) if error.kind == ErrorKind::NotFound => {}
3090 Err(error) => return Err(delete_remote_error(error, "planning integration read")),
3091 }
3092 let matches = collect(transport)
3093 .await
3094 .map_err(|error| delete_remote_error(error, "planning integration list read"))?
3095 .iter()
3096 .filter(|item| item.get("name").and_then(Value::as_str) == Some(selector))
3097 .map(summary_from_item)
3098 .collect::<Result<Vec<_>>>()?;
3099 match matches.as_slice() {
3100 [] => Err(Error::new(
3101 ErrorKind::NotFound,
3102 format!("no integration policy with id or name '{selector}'"),
3103 )),
3104 [summary] => {
3105 let policy = integration_policies::get(transport, &summary.id)
3106 .await
3107 .map_err(|error| delete_remote_error(error, "planning name read"))?;
3108 let returned_id = required_string(
3109 &policy.item,
3110 "id",
3111 "integration policy delete planning read",
3112 )?;
3113 if returned_id != summary.id {
3114 return Err(http(
3115 "decoding integration policy delete planning read: name resolution returned an unexpected id",
3116 ));
3117 }
3118 Ok(ResolvedIntegrationPolicy {
3119 summary: summary.clone(),
3120 item: policy.item,
3121 })
3122 }
3123 many => Err(Error::new(
3124 ErrorKind::Conflict,
3125 format!(
3126 "integration policy '{selector}' is ambiguous: {}",
3127 many.iter()
3128 .map(|policy| policy.id.as_str())
3129 .collect::<Vec<_>>()
3130 .join(", ")
3131 ),
3132 )),
3133 }
3134}
3135
3136async fn plan_delete_target(
3137 transport: &Transport,
3138 id: &str,
3139 item: Map<String, Value>,
3140) -> Result<IntegrationPolicyDeleteTarget> {
3141 if required_string(&item, "id", "integration policy delete planning read")? != id {
3142 return Err(http(
3143 "decoding integration policy delete planning read: response id did not match the request",
3144 ));
3145 }
3146 let spec = normalize(&item, transport.space())?;
3147 if spec.id != id {
3148 return Err(http(
3149 "decoding integration policy delete planning read: normalized id did not match the request",
3150 ));
3151 }
3152 let parent_ids = read_parents(id, &item)?;
3153 let parents = read_parent_snapshots(transport, id, &parent_ids)
3154 .await
3155 .map_err(|error| delete_remote_error(error, "planning parent read"))?;
3156 validate_delete_parent_safety(id, &spec, &parents)?;
3157
3158 let package = package_coordinate(&item, "integration policy delete planning read")?;
3159 if package != spec.package {
3160 return Err(http(
3161 "decoding integration policy delete planning read: package did not normalize canonically",
3162 ));
3163 }
3164 let dependency = read_dependencies(transport, &package)
3165 .await
3166 .map_err(|error| delete_remote_error(error, "planning package read"))?;
3167 ensure_delete_dependency(id, &dependency, &package)?;
3168 let metadata =
3169 integration_policies::package_metadata(transport, &package.name, &package.version)
3170 .await
3171 .map_err(|error| delete_remote_error(error, "planning package metadata read"))?
3172 .item;
3173 validate_package_metadata_snapshot(&metadata, &package)?;
3174 let secret_paths = configured_secret_paths(&spec, &metadata)?;
3175 if !secret_paths.is_empty() {
3176 return unsupported(format!(
3177 "integration policy '{id}' is not portable: {}",
3178 secret_paths
3179 .into_iter()
3180 .map(|path| format!("{id}:{path}"))
3181 .collect::<Vec<_>>()
3182 .join(", ")
3183 ));
3184 }
3185
3186 Ok(IntegrationPolicyDeleteTarget {
3187 id: id.to_owned(),
3188 name: spec.name.clone(),
3189 item_snapshot: item.clone(),
3190 item,
3191 spec_snapshot: spec.clone(),
3192 spec,
3193 parents,
3194 package,
3195 dependency_snapshot: dependency.clone(),
3196 dependency,
3197 metadata_snapshot: metadata.clone(),
3198 metadata,
3199 })
3200}
3201
3202fn collapse_delete_planning_issues(mut issues: Vec<Error>) -> Result<IntegrationPolicyDeletePlan> {
3203 if issues.len() == 1 {
3204 return Err(issues.remove(0));
3205 }
3206 if issues.iter().all(|error| error.kind == ErrorKind::Conflict) {
3207 return Err(Error::new(
3208 ErrorKind::Conflict,
3209 issues
3210 .into_iter()
3211 .map(|error| error.message)
3212 .collect::<Vec<_>>()
3213 .join("; "),
3214 ));
3215 }
3216 Err(issues.remove(0))
3217}
3218
3219fn validate_delete_parent_safety(
3220 id: &str,
3221 spec: &IntegrationPolicySpec,
3222 parents: &BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
3223) -> Result<()> {
3224 if parents.len() != spec.policy_ids.len()
3225 || parents
3226 .keys()
3227 .map(String::as_str)
3228 .ne(spec.policy_ids.iter().map(String::as_str))
3229 {
3230 return Err(http(format!(
3231 "decoding integration policy '{id}': parent snapshots do not match policy_ids"
3232 )));
3233 }
3234 for parent in parents.values() {
3235 if parent.platform_owned {
3236 return unsupported(format!(
3237 "integration policy '{id}' is not portable: parent {} is platform-owned",
3238 parent.id
3239 ));
3240 }
3241 if parent.protected {
3242 return unsupported(format!(
3243 "integration policy '{id}' is not portable: parent {} is_protected",
3244 parent.id
3245 ));
3246 }
3247 if parent
3248 .attached_integrations
3249 .binary_search_by(|attached| attached.as_str().cmp(id))
3250 .is_err()
3251 {
3252 return Err(http(format!(
3253 "decoding integration policy '{id}': parent '{}' is missing its attachment",
3254 parent.id
3255 )));
3256 }
3257 }
3258 let namespaces = parents
3259 .values()
3260 .map(|parent| parent.namespace.as_str())
3261 .collect::<BTreeSet<_>>();
3262 match &spec.namespace {
3263 Some(namespace)
3264 if parents
3265 .values()
3266 .all(|parent| &parent.namespace == namespace) => {}
3267 Some(_) => {
3268 return unsupported(format!(
3269 "integration policy '{id}' is not portable: namespace does not match every parent"
3270 ));
3271 }
3272 None if namespaces.len() == 1 => {}
3273 None => {
3274 return unsupported(format!(
3275 "integration policy '{id}' is not portable: parents have different namespaces"
3276 ));
3277 }
3278 }
3279 Ok(())
3280}
3281
3282fn ensure_delete_dependency(
3283 id: &str,
3284 dependency: &PackageDependencySnapshot,
3285 package: &IntegrationPackageSpec,
3286) -> Result<()> {
3287 match &dependency.state {
3288 PackageDependencyState::Installed { version }
3289 if dependency.name == package.name && version == &package.version =>
3290 {
3291 Ok(())
3292 }
3293 PackageDependencyState::Installed { .. } => Err(Error::new(
3294 ErrorKind::Conflict,
3295 format!(
3296 "integration policy '{id}' package {} has a different installed version",
3297 package.name
3298 ),
3299 )),
3300 PackageDependencyState::NotInstalled => Err(Error::new(
3301 ErrorKind::Conflict,
3302 format!(
3303 "integration policy '{id}' package {} is not installed",
3304 package.name
3305 ),
3306 )),
3307 }
3308}
3309
3310fn shared_delete_parents(
3311 targets: &[IntegrationPolicyDeleteTarget],
3312) -> Result<BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>> {
3313 let mut shared = BTreeMap::new();
3314 for target in targets {
3315 for (id, parent) in &target.parents {
3316 match shared.entry(id.clone()) {
3317 std::collections::btree_map::Entry::Vacant(entry) => {
3318 entry.insert(parent.clone());
3319 }
3320 std::collections::btree_map::Entry::Occupied(entry) if entry.get() != parent => {
3321 return Err(Error::new(
3322 ErrorKind::Conflict,
3323 format!("agent policy '{id}' changed while planning integration deletion"),
3324 ));
3325 }
3326 std::collections::btree_map::Entry::Occupied(_) => {}
3327 }
3328 }
3329 }
3330 Ok(shared)
3331}
3332
3333fn delete_preview(targets: &[IntegrationPolicyDeleteTarget]) -> MutationPlan {
3334 let mut affected = BTreeMap::new();
3335 let mut preview_details = Vec::with_capacity(targets.len() + 2);
3336 for target in targets {
3337 let parents = target
3338 .parents
3339 .values()
3340 .map(|parent| {
3341 affected.entry(parent.id.clone()).or_insert(parent.agents);
3342 format!("{} ({}) agents {}", parent.id, parent.name, parent.agents)
3343 })
3344 .collect::<Vec<_>>();
3345 let agents = target
3346 .parents
3347 .values()
3348 .map(|parent| parent.agents)
3349 .sum::<u64>();
3350 preview_details.push(format!(
3351 "{} {} parents {} agents {agents}",
3352 target.id,
3353 target.name,
3354 parents.join(", ")
3355 ));
3356 }
3357 preview_details.push(format!(
3358 "affected agents {}",
3359 affected.values().sum::<u64>()
3360 ));
3361 preview_details.push(DELETE_RACE_WARNING.to_owned());
3362 MutationPlan {
3363 preview_action: format!("Delete {} integration policy(ies)", targets.len()),
3364 preview_details,
3365 targets: targets.iter().map(|target| target.id.clone()).collect(),
3366 }
3367}
3368
3369pub async fn apply_delete(
3373 transport: &Transport,
3374 plan: &IntegrationPolicyDeletePlan,
3375) -> Result<IntegrationPolicyDeleteReport> {
3376 validate_delete_plan(plan)?;
3377 if plan.host != transport.kibana_url() || plan.space != transport.space() {
3378 return Err(Error::new(
3379 ErrorKind::Conflict,
3380 "integration delete target changed since preview",
3381 ));
3382 }
3383
3384 let mut expected_parents = plan.parent_snapshots.clone();
3385 let mut affected = BTreeMap::new();
3386 let mut deleted = Vec::new();
3387 let mut failed = Vec::new();
3388
3389 for target in &plan.targets {
3390 match integration_policies::get(transport, &target.id).await {
3391 Ok(actual) if actual.item == target.item => {}
3392 Ok(_) => {
3393 failed.push(delete_failed_row(
3394 &target.id,
3395 false,
3396 "integration policy changed since preview",
3397 ));
3398 continue;
3399 }
3400 Err(error) if error.kind == ErrorKind::NotFound => {
3401 failed.push(delete_failed_row(
3402 &target.id,
3403 false,
3404 "integration policy disappeared since preview",
3405 ));
3406 continue;
3407 }
3408 Err(error) => {
3409 failed.push(delete_failed_row(
3410 &target.id,
3411 false,
3412 delete_remote_error(error, "apply integration-policy read").message,
3413 ));
3414 continue;
3415 }
3416 }
3417
3418 if let Err(error) = recheck_delete_parents(transport, target, &expected_parents).await {
3419 failed.push(delete_failed_row(&target.id, false, error.message));
3420 continue;
3421 }
3422 match read_dependencies(transport, &target.package).await {
3423 Ok(actual) if actual == target.dependency => {}
3424 Ok(_) => {
3425 failed.push(delete_failed_row(
3426 &target.id,
3427 false,
3428 "integration policy package changed since preview",
3429 ));
3430 continue;
3431 }
3432 Err(error) => {
3433 failed.push(delete_failed_row(
3434 &target.id,
3435 false,
3436 delete_remote_error(error, "apply package read").message,
3437 ));
3438 continue;
3439 }
3440 }
3441 if let Err(error) = recheck_delete_metadata(transport, target).await {
3442 failed.push(delete_failed_row(&target.id, false, error.message));
3443 continue;
3444 }
3445
3446 match integration_policies::delete(transport, &target.id).await {
3447 Ok(()) => {
3448 record_delete_affected_parents(&mut affected, target, &expected_parents);
3449 advance_delete_parent_snapshots(&mut expected_parents, target);
3450 deleted.push(json!({"id": target.id}));
3451 }
3452 Err(error) => {
3453 let applied = error
3454 .http_status
3455 .is_some_and(|status| (200..300).contains(&status));
3456 let message = if applied {
3457 "integration-policy delete response did not confirm the requested id"
3458 } else {
3459 "integration-policy delete request failed"
3460 };
3461 failed.push(delete_failed_row(&target.id, applied, message));
3462 }
3463 }
3464 }
3465
3466 Ok(IntegrationPolicyDeleteReport {
3467 applied: true,
3468 deleted,
3469 failed,
3470 total: plan.total,
3471 affected_agents: affected.values().sum(),
3472 })
3473}
3474
3475async fn recheck_delete_parents(
3476 transport: &Transport,
3477 target: &IntegrationPolicyDeleteTarget,
3478 expected_parents: &BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
3479) -> Result<()> {
3480 for parent_id in target.parents.keys() {
3481 let expected = expected_parents.get(parent_id).ok_or_else(|| {
3482 Error::new(
3483 ErrorKind::Error,
3484 "integration delete lost a shared parent snapshot",
3485 )
3486 })?;
3487 match agent_policy_ops::read_parent_snapshot(transport, parent_id).await {
3488 Ok(actual) if actual == *expected => {}
3489 Ok(_) => {
3490 return Err(Error::new(
3491 ErrorKind::Conflict,
3492 "integration policy parent changed since preview",
3493 ));
3494 }
3495 Err(error) if error.kind == ErrorKind::NotFound => {
3496 return Err(Error::new(
3497 ErrorKind::NotFound,
3498 "integration policy parent disappeared since preview",
3499 ));
3500 }
3501 Err(error) => return Err(delete_remote_error(error, "apply parent read")),
3502 }
3503 }
3504 Ok(())
3505}
3506
3507async fn recheck_delete_metadata(
3508 transport: &Transport,
3509 target: &IntegrationPolicyDeleteTarget,
3510) -> Result<()> {
3511 let metadata = integration_policies::package_metadata(
3512 transport,
3513 &target.package.name,
3514 &target.package.version,
3515 )
3516 .await
3517 .map_err(|error| delete_remote_error(error, "apply package metadata read"))?
3518 .item;
3519 validate_package_metadata_snapshot(&metadata, &target.package)
3520 .map_err(|error| delete_remote_error(error, "apply package metadata read"))?;
3521 if metadata != target.metadata {
3522 return Err(Error::new(
3523 ErrorKind::Conflict,
3524 "integration policy package metadata changed since preview",
3525 ));
3526 }
3527 Ok(())
3528}
3529
3530fn record_delete_affected_parents(
3531 affected: &mut BTreeMap<String, u64>,
3532 target: &IntegrationPolicyDeleteTarget,
3533 expected_parents: &BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
3534) {
3535 for parent_id in target.parents.keys() {
3536 let parent = expected_parents
3537 .get(parent_id)
3538 .expect("validated delete target parent exists in shared snapshots");
3539 affected.entry(parent.id.clone()).or_insert(parent.agents);
3540 }
3541}
3542
3543fn advance_delete_parent_snapshots(
3544 parents: &mut BTreeMap<String, agent_policy_ops::AgentPolicyParentSnapshot>,
3545 target: &IntegrationPolicyDeleteTarget,
3546) {
3547 for parent_id in target.parents.keys() {
3548 let parent = parents
3549 .get_mut(parent_id)
3550 .expect("validated delete target parent exists in shared snapshots");
3551 parent
3552 .attached_integrations
3553 .retain(|attached| attached != &target.id);
3554 }
3555}
3556
3557fn delete_failed_row(id: &str, applied: bool, error: impl Into<String>) -> Value {
3558 json!({"id": id, "applied": applied, "error": error.into()})
3559}
3560
3561fn delete_remote_error(error: Error, context: &str) -> Error {
3562 let message = format!("integration-policy delete {context} failed");
3563 match error.http_status {
3564 Some(status) => Error::with_status(error.kind, status, message),
3565 None => Error::new(error.kind, message),
3566 }
3567}
3568
3569fn validate_delete_plan(plan: &IntegrationPolicyDeletePlan) -> Result<()> {
3570 let invalid = || Error::new(ErrorKind::Error, "invalid integration-policy delete plan");
3571 if plan.targets.is_empty()
3572 || plan.total != plan.targets.len()
3573 || plan.host.trim().is_empty()
3574 || plan.host != plan.host_snapshot
3575 || plan.space != plan.space_snapshot
3576 {
3577 return Err(invalid());
3578 }
3579
3580 let mut previous: Option<&str> = None;
3581 let mut shared = BTreeMap::new();
3582 for target in &plan.targets {
3583 if target.id.trim().is_empty()
3584 || target.name.trim().is_empty()
3585 || previous.is_some_and(|previous| previous >= target.id.as_str())
3586 || target.item != target.item_snapshot
3587 || target.spec.validate().is_err()
3588 || target.spec != target.spec_snapshot
3589 || target.id != target.spec.id
3590 || target.name != target.spec.name
3591 || required_string(&target.item, "id", "integration policy delete plan")
3592 .ok()
3593 .as_deref()
3594 != Some(target.id.as_str())
3595 || normalize(&target.item, &plan.space).ok().as_ref() != Some(&target.spec)
3596 || package_coordinate(&target.item, "integration policy delete plan")
3597 .ok()
3598 .as_ref()
3599 != Some(&target.package)
3600 || target.package != target.spec.package
3601 || target.dependency != target.dependency_snapshot
3602 || !valid_package_state(&target.dependency)
3603 || target.dependency.name != target.package.name
3604 || !is_exact_installed(&target.dependency, &target.package)
3605 || target.metadata != target.metadata_snapshot
3606 || validate_package_metadata_snapshot(&target.metadata, &target.package).is_err()
3607 || !matches!(configured_secret_paths(&target.spec, &target.metadata), Ok(paths) if paths.is_empty())
3608 {
3609 return Err(invalid());
3610 }
3611
3612 let parent_ids = match read_parents(&target.id, &target.item) {
3613 Ok(ids) => ids,
3614 Err(_) => return Err(invalid()),
3615 };
3616 if parent_ids.iter().collect::<BTreeSet<_>>()
3617 != target.parents.keys().collect::<BTreeSet<_>>()
3618 || validate_delete_parent_safety(&target.id, &target.spec, &target.parents).is_err()
3619 {
3620 return Err(invalid());
3621 }
3622 for (parent_id, parent) in &target.parents {
3623 if !valid_parent_snapshot(parent_id, parent)
3624 || parent.platform_owned
3625 || parent.protected
3626 || parent
3627 .attached_integrations
3628 .binary_search_by(|attached| attached.as_str().cmp(&target.id))
3629 .is_err()
3630 {
3631 return Err(invalid());
3632 }
3633 match shared.entry(parent_id.as_str()) {
3634 std::collections::btree_map::Entry::Vacant(entry) => {
3635 entry.insert(parent);
3636 }
3637 std::collections::btree_map::Entry::Occupied(entry) if *entry.get() != parent => {
3638 return Err(invalid());
3639 }
3640 std::collections::btree_map::Entry::Occupied(_) => {}
3641 }
3642 }
3643 previous = Some(&target.id);
3644 }
3645 if plan.parent_snapshots != plan.parent_snapshots_snapshot
3646 || shared_delete_parents(&plan.targets).ok().as_ref() != Some(&plan.parent_snapshots)
3647 {
3648 return Err(invalid());
3649 }
3650 if plan.preview != delete_preview(&plan.targets) {
3651 return Err(invalid());
3652 }
3653 Ok(())
3654}
3655
3656fn http(message: impl Into<String>) -> Error {
3657 Error::new(ErrorKind::Http, message)
3658}
3659
3660fn unsupported<T>(message: impl Into<String>) -> Result<T> {
3661 Err(Error::new(ErrorKind::Unsupported, message))
3662}
3663
3664#[cfg(test)]
3665mod import_plan_tests {
3666 use super::*;
3667
3668 fn valid_plan() -> IntegrationPolicyImportPlan {
3669 let effective = IntegrationPolicySpec::try_from(json!({
3670 "id": "fresh",
3671 "name": "Fresh integration",
3672 "namespace": "default",
3673 "policy_ids": ["parent-1"],
3674 "package": {"name": "system", "version": "2.0.0"},
3675 "inputs": {}
3676 }))
3677 .expect("valid test policy");
3678 let parent = agent_policy_ops::AgentPolicyParentSnapshot {
3679 id: "parent-1".into(),
3680 name: "Parent 1".into(),
3681 namespace: "default".into(),
3682 agents: 0,
3683 attached_integrations: Vec::new(),
3684 platform_owned: false,
3685 protected: false,
3686 };
3687 let targets = vec![IntegrationPolicyImportTarget {
3688 effective: effective.clone(),
3689 current: None,
3690 parents: BTreeMap::from([(parent.id.clone(), parent)]),
3691 replacement_body: None,
3692 }];
3693 let parent_snapshots = targets[0].parents.clone();
3694 let state = PackageDependencySnapshot {
3695 name: "system".into(),
3696 state: PackageDependencyState::NotInstalled,
3697 };
3698 let metadata = json!({
3699 "name": "system",
3700 "version": "2.0.0",
3701 "vars": [],
3702 "policy_templates": []
3703 })
3704 .as_object()
3705 .expect("metadata object")
3706 .clone();
3707 let package_groups = BTreeMap::from([(
3708 "system".into(),
3709 IntegrationPackageGroup {
3710 package: effective.package.clone(),
3711 state: state.clone(),
3712 state_snapshot: state,
3713 metadata_snapshot: metadata.clone(),
3714 metadata,
3715 },
3716 )]);
3717 let package_installs = vec!["system@2.0.0".into()];
3718 let source = PathBuf::from("fresh.json");
3719 let preview = import_preview(&source, &targets, &package_installs);
3720 IntegrationPolicyImportPlan {
3721 preview,
3722 skipped: Vec::new(),
3723 package_installs,
3724 total: 1,
3725 source,
3726 host: "https://fleet.example.invalid".into(),
3727 space: "default".into(),
3728 canonical: vec![effective.clone()],
3729 name_owners: BTreeMap::from([(effective.name.clone(), BTreeSet::new())]),
3730 name_owners_snapshot: BTreeMap::from([(effective.name.clone(), BTreeSet::new())]),
3731 parent_snapshots,
3732 skipped_snapshot: Vec::new(),
3733 existing_snapshot: BTreeMap::from([("fresh".into(), None)]),
3734 targets,
3735 package_groups,
3736 overwrite: false,
3737 skip_existing: false,
3738 }
3739 }
3740
3741 fn existing_plan_without_overwrite() -> IntegrationPolicyImportPlan {
3742 let mut plan = valid_plan();
3743 let existing = {
3744 let target = plan.targets.first_mut().expect("fresh target");
3745 let mut item = serde_json::to_value(&target.effective)
3746 .expect("serialize current item")
3747 .as_object()
3748 .expect("current item object")
3749 .clone();
3750 item.insert("enabled".into(), Value::Bool(true));
3751 target.current = Some(IntegrationPolicyCurrentSnapshot {
3752 item: item.clone(),
3753 spec: target.effective.clone(),
3754 parent_ids: target.effective.policy_ids.clone(),
3755 });
3756 target
3757 .parents
3758 .get_mut("parent-1")
3759 .expect("parent")
3760 .attached_integrations
3761 .push(target.effective.id.clone());
3762 item
3763 };
3764 plan.existing_snapshot
3765 .insert("fresh".into(), Some(existing));
3766
3767 let state = PackageDependencySnapshot {
3768 name: "system".into(),
3769 state: PackageDependencyState::Installed {
3770 version: "2.0.0".into(),
3771 },
3772 };
3773 let group = plan
3774 .package_groups
3775 .get_mut("system")
3776 .expect("package group");
3777 group.state = state.clone();
3778 group.state_snapshot = state;
3779 plan.package_installs.clear();
3780 plan.name_owners
3781 .get_mut("Fresh integration")
3782 .expect("name owner snapshot")
3783 .insert("fresh".into());
3784 plan.name_owners_snapshot = plan.name_owners.clone();
3785 plan.parent_snapshots = plan.targets[0].parents.clone();
3786 plan.preview = import_preview(&plan.source, &plan.targets, &plan.package_installs);
3787 plan
3788 }
3789
3790 fn valid_skip_existing_plan() -> IntegrationPolicyImportPlan {
3791 let mut plan = existing_plan_without_overwrite();
3792 plan.skip_existing = true;
3793 plan.targets.clear();
3794 plan.parent_snapshots.clear();
3795 plan.package_groups.clear();
3796 plan.skipped = vec![json!({"id": "fresh", "reason": "exists"})];
3797 plan.skipped_snapshot = plan.skipped.clone();
3798 plan.package_installs.clear();
3799 plan.preview = import_preview(&plan.source, &plan.targets, &plan.package_installs);
3800 plan
3801 }
3802
3803 fn valid_replace_plan() -> IntegrationPolicyImportPlan {
3804 let mut plan = existing_plan_without_overwrite();
3805 plan.overwrite = true;
3806 let desired = {
3807 let target = plan.targets.first_mut().expect("existing target");
3808 let mut desired = target.effective.clone();
3809 desired.description = Some("changed".into());
3810 target.effective = desired.clone();
3811 target.replacement_body = Some(replace_wire_body(&desired).expect("replace body"));
3812 desired
3813 };
3814 plan.canonical = vec![desired];
3815 plan.preview = import_preview(&plan.source, &plan.targets, &plan.package_installs);
3816 plan
3817 }
3818
3819 #[test]
3820 fn replacement_body_omits_response_only_enabled_without_changing_input_enabled() {
3821 let spec = IntegrationPolicySpec::try_from(json!({
3822 "id": "replacement",
3823 "name": "Replacement integration",
3824 "namespace": "default",
3825 "policy_ids": ["parent-1"],
3826 "package": {"name": "system", "version": "2.0.0"},
3827 "inputs": {"system-log": {"enabled": true}}
3828 }))
3829 .expect("valid replacement spec");
3830
3831 let body = replace_wire_body(&spec).expect("replacement wire body");
3832 let object = body.as_object().expect("replacement wire object");
3833
3834 assert!(object.get("id").is_none());
3835 assert!(object.get("enabled").is_none());
3836 assert_eq!(object["inputs"]["system-log"]["enabled"], true);
3837 }
3838
3839 fn valid_expanded_inputs_plan() -> IntegrationPolicyImportPlan {
3840 let mut plan = valid_plan();
3841 let inputs = json!({"system-system": {}})
3842 .as_object()
3843 .expect("inputs object")
3844 .clone();
3845 plan.canonical[0].inputs = inputs.clone();
3846 plan.targets[0].effective.inputs = inputs;
3847 let metadata = json!({
3848 "name": "system",
3849 "version": "2.0.0",
3850 "vars": [],
3851 "policy_templates": [{
3852 "name": "system",
3853 "inputs": [{"type": "system"}]
3854 }]
3855 })
3856 .as_object()
3857 .expect("metadata object")
3858 .clone();
3859 let group = plan
3860 .package_groups
3861 .get_mut("system")
3862 .expect("package group");
3863 group.metadata = metadata.clone();
3864 group.metadata_snapshot = metadata;
3865 plan.preview = import_preview(&plan.source, &plan.targets, &plan.package_installs);
3866 plan
3867 }
3868
3869 #[test]
3870 fn import_plan_rejects_a_private_create_name_owner_tamper() {
3871 let mut plan = valid_plan();
3872 plan.name_owners
3873 .get_mut("Fresh integration")
3874 .expect("name owner snapshot")
3875 .insert("fresh".into());
3876
3877 assert!(validate_import_plan(&plan).is_err());
3878 }
3879
3880 #[test]
3881 fn import_plan_rejects_a_coherent_empty_effective_inputs_tamper() {
3882 let mut plan = valid_expanded_inputs_plan();
3883 assert!(validate_import_plan(&plan).is_ok());
3884
3885 plan.canonical[0].inputs.clear();
3886 plan.targets[0].effective.inputs.clear();
3887 plan.preview = import_preview(&plan.source, &plan.targets, &plan.package_installs);
3888
3889 let error = validate_import_plan(&plan)
3890 .expect_err("an empty effective map must not reach import requests");
3891 assert_eq!(error.kind, ErrorKind::Unsupported);
3892 assert_eq!(
3893 error.message,
3894 "integration policy 'fresh' has an empty inputs map but package system@2.0.0 declares inputs"
3895 );
3896 }
3897
3898 #[test]
3899 fn import_plan_rejects_an_existing_target_without_overwrite() {
3900 let plan = existing_plan_without_overwrite();
3901
3902 assert!(validate_import_plan(&plan).is_err());
3903 }
3904
3905 #[test]
3906 fn import_plan_rejects_private_snapshot_body_group_and_order_tampering() {
3907 let replace = valid_replace_plan();
3908 assert!(validate_import_plan(&replace).is_ok());
3909
3910 let mut tampered_body = replace.clone();
3911 tampered_body.targets[0].replacement_body = Some(json!({"tampered": true}));
3912 assert!(validate_import_plan(&tampered_body).is_err());
3913
3914 let mut tampered_current = replace.clone();
3915 tampered_current.targets[0]
3916 .current
3917 .as_mut()
3918 .expect("current snapshot")
3919 .item
3920 .insert("enabled".into(), Value::Bool(false));
3921 assert!(validate_import_plan(&tampered_current).is_err());
3922
3923 let mut tampered_group = valid_plan();
3924 tampered_group
3925 .package_groups
3926 .get_mut("system")
3927 .expect("package group")
3928 .metadata
3929 .insert("version".into(), Value::String("9.9.9".into()));
3930 assert!(validate_import_plan(&tampered_group).is_err());
3931
3932 let mut tampered_order = valid_plan();
3933 tampered_order
3934 .targets
3935 .push(tampered_order.targets[0].clone());
3936 assert!(validate_import_plan(&tampered_order).is_err());
3937
3938 let mut tampered_group_key = valid_plan();
3939 let group = tampered_group_key
3940 .package_groups
3941 .remove("system")
3942 .expect("package group");
3943 tampered_group_key
3944 .package_groups
3945 .insert("other".into(), group);
3946 assert!(validate_import_plan(&tampered_group_key).is_err());
3947
3948 let mut tampered_state = valid_plan();
3949 tampered_state
3950 .package_groups
3951 .get_mut("system")
3952 .expect("package group")
3953 .state = PackageDependencySnapshot {
3954 name: "system".into(),
3955 state: PackageDependencyState::Installed {
3956 version: "2.0.0".into(),
3957 },
3958 };
3959 assert!(validate_import_plan(&tampered_state).is_err());
3960
3961 let mut tampered_coordinate = valid_replace_plan();
3962 tampered_coordinate.canonical[0].package.version = "3.0.0".into();
3963 tampered_coordinate.targets[0].effective.package.version = "3.0.0".into();
3964 tampered_coordinate.targets[0].replacement_body = Some(
3965 replace_wire_body(&tampered_coordinate.targets[0].effective).expect("replace body"),
3966 );
3967 let group = tampered_coordinate
3968 .package_groups
3969 .get_mut("system")
3970 .expect("package group");
3971 group.package.version = "3.0.0".into();
3972 group.state = PackageDependencySnapshot {
3973 name: "system".into(),
3974 state: PackageDependencyState::Installed {
3975 version: "3.0.0".into(),
3976 },
3977 };
3978 group.state_snapshot = group.state.clone();
3979 group.metadata.insert("version".into(), json!("3.0.0"));
3980 group.metadata_snapshot = group.metadata.clone();
3981 tampered_coordinate.preview = import_preview(
3982 &tampered_coordinate.source,
3983 &tampered_coordinate.targets,
3984 &tampered_coordinate.package_installs,
3985 );
3986 assert!(validate_import_plan(&tampered_coordinate).is_err());
3987 }
3988
3989 #[test]
3990 fn import_plan_rejects_a_private_parent_snapshot_tamper_even_with_preview_rebuilt() {
3991 let mut plan = valid_plan();
3992 plan.targets[0]
3993 .parents
3994 .get_mut("parent-1")
3995 .expect("parent snapshot")
3996 .agents = 42;
3997 plan.preview = import_preview(&plan.source, &plan.targets, &plan.package_installs);
3998
3999 assert!(validate_import_plan(&plan).is_err());
4000 }
4001
4002 #[test]
4003 fn import_plan_rejects_target_removal_rebuilt_as_a_skipped_row() {
4004 let mut plan = valid_plan();
4005 plan.skip_existing = true;
4006 plan.targets.clear();
4007 plan.parent_snapshots.clear();
4008 plan.package_groups.clear();
4009 plan.skipped = vec![json!({"id": "fresh", "reason": "exists"})];
4010 plan.skipped_snapshot = plan.skipped.clone();
4011 plan.package_installs.clear();
4012 plan.preview = import_preview(&plan.source, &plan.targets, &plan.package_installs);
4013
4014 assert!(validate_import_plan(&plan).is_err());
4015 }
4016
4017 #[test]
4018 fn import_plan_accepts_a_coherent_skip_existing_snapshot() {
4019 assert!(validate_import_plan(&valid_skip_existing_plan()).is_ok());
4020 }
4021
4022 #[test]
4023 fn import_plan_rejects_existence_snapshot_key_and_target_mismatches() {
4024 let mut missing_current = valid_replace_plan();
4025 missing_current
4026 .existing_snapshot
4027 .insert("fresh".into(), None);
4028 assert!(validate_import_plan(&missing_current).is_err());
4029
4030 let mut changed_current = valid_replace_plan();
4031 changed_current
4032 .existing_snapshot
4033 .get_mut("fresh")
4034 .expect("existing snapshot")
4035 .as_mut()
4036 .expect("existing item")
4037 .insert("description".into(), json!("tampered"));
4038 assert!(validate_import_plan(&changed_current).is_err());
4039
4040 let mut extra_snapshot = valid_plan();
4041 extra_snapshot
4042 .existing_snapshot
4043 .insert("other".into(), None);
4044 assert!(validate_import_plan(&extra_snapshot).is_err());
4045 }
4046
4047 #[test]
4048 fn import_plan_rejects_public_field_tampering_against_private_snapshots() {
4049 let plan = valid_plan();
4050
4051 let mut total = plan.clone();
4052 total.total = 2;
4053 assert!(validate_import_plan(&total).is_err());
4054
4055 let mut preview = plan.clone();
4056 preview.preview.preview_action = "tampered".into();
4057 assert!(validate_import_plan(&preview).is_err());
4058
4059 let mut skipped = plan.clone();
4060 skipped.skipped = vec![json!({"id": "fresh", "reason": "exists"})];
4061 assert!(validate_import_plan(&skipped).is_err());
4062
4063 let mut installs = plan;
4064 installs.package_installs.clear();
4065 assert!(validate_import_plan(&installs).is_err());
4066 }
4067}
4068
4069#[cfg(test)]
4070mod delete_plan_tests {
4071 use super::*;
4072 use elasticctl_core::{Profile, Transport};
4073 use wiremock::matchers::{method, path, query_param};
4074 use wiremock::{Mock, MockServer, ResponseTemplate};
4075
4076 async fn verified_server() -> MockServer {
4077 let server = MockServer::start().await;
4078 Mock::given(method("GET"))
4079 .and(path("/api/status"))
4080 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
4081 "version": {"number": "9.5.1", "build_flavor": "traditional"}
4082 })))
4083 .mount(&server)
4084 .await;
4085 server
4086 }
4087
4088 fn transport_for(server: &MockServer) -> Transport {
4089 Transport::new(&Profile {
4090 kibana_url: server.uri(),
4091 es_url: None,
4092 api_key: Some("essu_test".into()),
4093 username: None,
4094 password: None,
4095 space: "default".into(),
4096 verify: true,
4097 timeout_secs: 5,
4098 })
4099 .expect("transport")
4100 }
4101
4102 fn valid_plan() -> IntegrationPolicyDeletePlan {
4103 let spec = IntegrationPolicySpec::try_from(json!({
4104 "id": "delete-1",
4105 "name": "Delete integration",
4106 "namespace": "default",
4107 "policy_ids": ["parent-1"],
4108 "package": {"name": "system", "version": "2.0.0"},
4109 "inputs": {}
4110 }))
4111 .expect("valid integration policy");
4112 let mut item = serde_json::to_value(&spec)
4113 .expect("serialize integration policy")
4114 .as_object()
4115 .expect("integration policy is an object")
4116 .clone();
4117 item.insert("enabled".into(), Value::Bool(true));
4118 let parent = agent_policy_ops::AgentPolicyParentSnapshot {
4119 id: "parent-1".into(),
4120 name: "Parent 1".into(),
4121 namespace: "default".into(),
4122 agents: 4,
4123 attached_integrations: vec!["delete-1".into()],
4124 platform_owned: false,
4125 protected: false,
4126 };
4127 let parents = BTreeMap::from([(parent.id.clone(), parent)]);
4128 let dependency = PackageDependencySnapshot {
4129 name: "system".into(),
4130 state: PackageDependencyState::Installed {
4131 version: "2.0.0".into(),
4132 },
4133 };
4134 let metadata = json!({
4135 "name": "system",
4136 "version": "2.0.0",
4137 "vars": [],
4138 "policy_templates": []
4139 })
4140 .as_object()
4141 .expect("metadata object")
4142 .clone();
4143 let target = IntegrationPolicyDeleteTarget {
4144 id: spec.id.clone(),
4145 name: spec.name.clone(),
4146 item_snapshot: item.clone(),
4147 item,
4148 spec_snapshot: spec.clone(),
4149 spec,
4150 parents: parents.clone(),
4151 package: IntegrationPackageSpec {
4152 name: "system".into(),
4153 version: "2.0.0".into(),
4154 },
4155 dependency_snapshot: dependency.clone(),
4156 dependency,
4157 metadata_snapshot: metadata.clone(),
4158 metadata,
4159 };
4160 let targets = vec![target];
4161 let parent_snapshots = parents;
4162 IntegrationPolicyDeletePlan {
4163 preview: delete_preview(&targets),
4164 total: targets.len(),
4165 host: "https://fleet.example.invalid".into(),
4166 host_snapshot: "https://fleet.example.invalid".into(),
4167 space: "default".into(),
4168 space_snapshot: "default".into(),
4169 parent_snapshots_snapshot: parent_snapshots.clone(),
4170 parent_snapshots,
4171 targets,
4172 }
4173 }
4174
4175 #[test]
4176 fn delete_plan_accepts_a_coherent_private_snapshot() {
4177 assert!(validate_delete_plan(&valid_plan()).is_ok());
4178 }
4179
4180 #[test]
4181 fn delete_plan_rejects_empty_total_order_and_preview_tampering() {
4182 let plan = valid_plan();
4183
4184 let mut empty = plan.clone();
4185 empty.targets.clear();
4186 empty.total = 0;
4187 empty.parent_snapshots.clear();
4188 empty.parent_snapshots_snapshot.clear();
4189 empty.preview = delete_preview(&empty.targets);
4190 assert!(validate_delete_plan(&empty).is_err());
4191
4192 let mut total = plan.clone();
4193 total.total = 2;
4194 assert!(validate_delete_plan(&total).is_err());
4195
4196 let mut duplicate = plan.clone();
4197 duplicate.targets.push(duplicate.targets[0].clone());
4198 duplicate.total = 2;
4199 duplicate.preview = delete_preview(&duplicate.targets);
4200 assert!(validate_delete_plan(&duplicate).is_err());
4201
4202 let mut preview = plan;
4203 preview.preview.preview_action = "tampered".into();
4204 assert!(validate_delete_plan(&preview).is_err());
4205
4206 let mut host = valid_plan();
4207 host.host = "https://other.example.invalid".into();
4208 assert!(validate_delete_plan(&host).is_err());
4209
4210 let mut space = valid_plan();
4211 space.space = "other".into();
4212 assert!(validate_delete_plan(&space).is_err());
4213 }
4214
4215 #[test]
4216 fn delete_plan_rejects_raw_spec_parent_package_and_metadata_tampering() {
4217 let plan = valid_plan();
4218
4219 let mut raw_and_spec = plan.clone();
4220 raw_and_spec.targets[0]
4221 .item
4222 .insert("description".into(), json!("tampered"));
4223 raw_and_spec.targets[0].spec.description = Some("tampered".into());
4224 raw_and_spec.preview = delete_preview(&raw_and_spec.targets);
4225 assert!(validate_delete_plan(&raw_and_spec).is_err());
4226
4227 let mut parent = plan.clone();
4228 parent.targets[0]
4229 .parents
4230 .get_mut("parent-1")
4231 .expect("parent")
4232 .agents = 99;
4233 parent.preview = delete_preview(&parent.targets);
4234 assert!(validate_delete_plan(&parent).is_err());
4235
4236 let mut parent_snapshot = plan.clone();
4237 parent_snapshot
4238 .parent_snapshots
4239 .get_mut("parent-1")
4240 .expect("parent")
4241 .agents = 99;
4242 assert!(validate_delete_plan(&parent_snapshot).is_err());
4243
4244 let mut dependency = plan.clone();
4245 dependency.targets[0].dependency = PackageDependencySnapshot {
4246 name: "system".into(),
4247 state: PackageDependencyState::Installed {
4248 version: "1.0.0".into(),
4249 },
4250 };
4251 assert!(validate_delete_plan(&dependency).is_err());
4252
4253 let mut metadata = plan;
4254 metadata.targets[0]
4255 .metadata
4256 .insert("version".into(), json!("9.9.9"));
4257 assert!(validate_delete_plan(&metadata).is_err());
4258 }
4259
4260 #[tokio::test]
4261 async fn delete_apply_rereads_metadata_after_coherent_secret_tampering() {
4262 let server = verified_server().await;
4263 let transport = transport_for(&server);
4264 let spec = IntegrationPolicySpec::try_from(json!({
4265 "id": "delete-1",
4266 "name": "Delete integration",
4267 "namespace": "default",
4268 "policy_ids": ["parent-1"],
4269 "package": {"name": "system", "version": "2.0.0"},
4270 "vars": {"package_secret": "live-plaintext-value-must-not-leak"},
4271 "inputs": {}
4272 }))
4273 .expect("valid integration policy");
4274 let mut item = serde_json::to_value(&spec)
4275 .expect("serialize integration policy")
4276 .as_object()
4277 .expect("integration policy is an object")
4278 .clone();
4279 item.insert("enabled".into(), Value::Bool(true));
4280 let parent = agent_policy_ops::AgentPolicyParentSnapshot {
4281 id: "parent-1".into(),
4282 name: "Parent 1".into(),
4283 namespace: "default".into(),
4284 agents: 4,
4285 attached_integrations: vec![spec.id.clone()],
4286 platform_owned: false,
4287 protected: false,
4288 };
4289 let parents = BTreeMap::from([(parent.id.clone(), parent)]);
4290 let dependency = PackageDependencySnapshot {
4291 name: "system".into(),
4292 state: PackageDependencyState::Installed {
4293 version: "2.0.0".into(),
4294 },
4295 };
4296 let original_metadata = json!({
4297 "name": "system",
4298 "version": "2.0.0",
4299 "vars": [{"name": "package_secret", "secret": true}],
4300 "policy_templates": []
4301 })
4302 .as_object()
4303 .expect("metadata object")
4304 .clone();
4305 let mut target = IntegrationPolicyDeleteTarget {
4306 id: spec.id.clone(),
4307 name: spec.name.clone(),
4308 item_snapshot: item.clone(),
4309 item,
4310 spec_snapshot: spec.clone(),
4311 spec,
4312 parents: parents.clone(),
4313 package: IntegrationPackageSpec {
4314 name: "system".into(),
4315 version: "2.0.0".into(),
4316 },
4317 dependency_snapshot: dependency.clone(),
4318 dependency,
4319 metadata_snapshot: original_metadata.clone(),
4320 metadata: original_metadata.clone(),
4321 };
4322 let mut plan = IntegrationPolicyDeletePlan {
4323 preview: delete_preview(std::slice::from_ref(&target)),
4324 total: 1,
4325 host: server.uri(),
4326 host_snapshot: server.uri(),
4327 space: "default".into(),
4328 space_snapshot: "default".into(),
4329 parent_snapshots_snapshot: parents.clone(),
4330 parent_snapshots: parents,
4331 targets: vec![target.clone()],
4332 };
4333 assert!(validate_delete_plan(&plan).is_err());
4334
4335 let forged_metadata = json!({
4336 "name": "system",
4337 "version": "2.0.0",
4338 "vars": [{"name": "package_secret", "secret": false}],
4339 "policy_templates": []
4340 })
4341 .as_object()
4342 .expect("metadata object")
4343 .clone();
4344 target.metadata = forged_metadata.clone();
4345 target.metadata_snapshot = forged_metadata;
4346 target.item_snapshot = target.item.clone();
4347 target.spec_snapshot = target.spec.clone();
4348 plan.targets = vec![target];
4349 plan.parent_snapshots = shared_delete_parents(&plan.targets).expect("shared parents");
4350 plan.parent_snapshots_snapshot = plan.parent_snapshots.clone();
4351 plan.preview = delete_preview(&plan.targets);
4352 assert!(validate_delete_plan(&plan).is_ok());
4353
4354 let item = plan.targets[0].item.clone();
4355 Mock::given(method("GET"))
4356 .and(path("/api/fleet/package_policies/delete-1"))
4357 .and(query_param("format", "simplified"))
4358 .respond_with(ResponseTemplate::new(200).set_body_json(json!({"item": item})))
4359 .expect(1)
4360 .mount(&server)
4361 .await;
4362 Mock::given(method("GET"))
4363 .and(path("/api/fleet/agent_policies/parent-1"))
4364 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
4365 "item": parent_item_for_delete_test("parent-1", "delete-1")
4366 })))
4367 .expect(1)
4368 .mount(&server)
4369 .await;
4370 Mock::given(method("GET"))
4371 .and(path("/api/fleet/epm/packages/system"))
4372 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
4373 "item": {
4374 "name": "system",
4375 "status": "installed",
4376 "installationInfo": {"version": "2.0.0"}
4377 }
4378 })))
4379 .expect(1)
4380 .mount(&server)
4381 .await;
4382 Mock::given(method("GET"))
4383 .and(path("/api/fleet/epm/packages/system/2.0.0"))
4384 .respond_with(ResponseTemplate::new(200).set_body_json(json!({
4385 "item": original_metadata
4386 })))
4387 .expect(1)
4388 .mount(&server)
4389 .await;
4390 Mock::given(method("DELETE"))
4391 .and(path("/api/fleet/package_policies/delete-1"))
4392 .respond_with(ResponseTemplate::new(200).set_body_json(json!({"id": "delete-1"})))
4393 .expect(0)
4394 .mount(&server)
4395 .await;
4396
4397 let report = apply_delete(&transport, &plan)
4398 .await
4399 .expect("metadata race is a row failure");
4400 assert!(report.deleted.is_empty());
4401 assert_eq!(
4402 report.failed,
4403 vec![json!({
4404 "id": "delete-1",
4405 "applied": false,
4406 "error": "integration policy package metadata changed since preview"
4407 })]
4408 );
4409 let requests = server.received_requests().await.expect("recorded requests");
4410 assert_eq!(
4411 requests
4412 .iter()
4413 .filter(|request| request.url.path() == "/api/fleet/epm/packages/system/2.0.0")
4414 .count(),
4415 1
4416 );
4417 assert!(requests.iter().all(|request| request.method != "DELETE"));
4418 assert!(
4419 !report.failed[0]["error"]
4420 .as_str()
4421 .expect("error string")
4422 .contains("live-plaintext-value-must-not-leak")
4423 );
4424 }
4425
4426 fn parent_item_for_delete_test(id: &str, attached: &str) -> Value {
4427 json!({
4428 "id": id,
4429 "name": "Parent 1",
4430 "namespace": "default",
4431 "agents": 4,
4432 "package_policies": [attached],
4433 })
4434 }
4435}