Skip to main content

elasticctl_api/fleet/
integration_policy_ops.rs

1//! Integration-policy selection, normalization, and local validation.
2
3use 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/// A resolved selector retains the one-object response for later operations.
37#[derive(Debug, Clone, PartialEq)]
38pub(crate) struct ResolvedIntegrationPolicy {
39    pub(crate) summary: IntegrationPolicySummary,
40    pub(crate) item: Map<String, Value>,
41}
42
43/// Strictly reduced package installation state. Fleet's public status decoder
44/// intentionally preserves registry text for agent-policy orchestration; an
45/// integration dependency must not accept an ambiguous state.
46#[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/// Exact package-defined secret variables. The companion `known_*` maps are
59/// deliberately private implementation detail: a configured value is safe
60/// only after Fleet metadata proves it has a definition.
61#[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
87/// A locally decoded import artifact. Its canonical specifications are kept
88/// private so callers can retain a validated source across context setup
89/// without being able to alter what remote planning will use.
90pub 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/// What `plan_import` preflights and `apply_import` rechecks. Only guard
105/// presentation is public: the canonical artifact, effective specifications,
106/// and Fleet snapshots never cross the API boundary.
107#[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    // Exact get-by-id results from before import classification. This covers
122    // every canonical id, so a mutable target or skipped row cannot turn an
123    // absent policy into an existing one (or the reverse).
124    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/// A safe, fully preflighted integration-policy deletion. Only the guard
157/// presentation and count are public: Fleet snapshots and package metadata can
158/// include configuration values, so they remain private to the API layer.
159#[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/// The result of applying an integration-policy deletion plan.
183#[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/// Private execution facts for one stable-id delete. The raw item detects a
208/// Fleet change before deletion; the normalized policy and metadata prove that
209/// no secret or environment-bound configuration has crossed the guard.
210#[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/// One package coordinate is shared by every pending policy that uses it.
227/// `*_snapshot` copies are deliberate: plan validation compares the retained
228/// exact Fleet response with the mutable execution expectation before it ever
229/// starts remote work.
230#[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
246/// Collect all measured pages then sort by stable id locally.
247pub 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
303/// List with local case-insensitive search and post-sort limiting.
304pub 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
331/// Resolve a stable id first, then a unique exact name.
332pub async fn resolve(transport: &Transport, selector: &str) -> Result<IntegrationPolicySummary> {
333    Ok(resolve_item(transport, selector).await?.summary)
334}
335
336/// Resolve a selector and retain the checked one-object response for later
337/// read-only operations.
338pub(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
375/// Bind a one-object response to its route id and, when it came from a list
376/// selection, to that row's safe summary. Do not include raw values in errors:
377/// simplified items can contain user configuration.
378fn 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
397/// Fleet can return parent ids in a different order between its list and
398/// single-item routes. Sort for comparison without removing duplicates, so a
399/// malformed repeated parent remains visible to later validation.
400fn 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
415/// Return a safe integration-policy view. Parent reads are both the attachment
416/// race check and the sole source of affected-agent counts.
417pub 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
458/// Validate every live shape that a safe detail can reason about without
459/// erasing its direct portability blockers. The projection keeps package-owned
460/// configuration intact while replacing only values that a detail reports in
461/// `blocked_by`, so `normalize` remains the single structural validator.
462fn 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
480/// Export selected integrations or every custom integration. A selector is
481/// resolved once and deduplicated by its stable id before any parent, package,
482/// or metadata reads.
483pub 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    // A package policy can only have been compiled from the exact installed
642    // coordinate. Treat a divergent or absent state as a loud conflict.
643    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
1260/// Build a fresh portable policy from a simplified live Fleet response.
1261pub 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
1375/// Keep package-defined input and stream maps open while rebuilding them
1376/// without Fleet-generated ids or compiled content.
1377fn 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
1648/// Read, validate, and retain an integration-policy artifact before context
1649/// or credential construction. The returned value has no public raw fields.
1650pub 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
1665/// Plan a retained integration-policy import without reopening its source or
1666/// sending a write. Apply uses only the returned plan.
1667pub 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    // Read only the requested ids. A conflicting or skipped existing policy is
1682    // deliberately kept raw: normalize can refuse an unsupported object, but
1683    // that object will not be written on those paths.
1684    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    // This phase comes before every parent, package-status, and metadata read.
1717    // A package-coordinate replacement is a different Fleet operation, never
1718    // a package-policy import update.
1719    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    // Every artifact name stays relevant even if its id is skipped. A foreign
1737    // claimant is always a conflict, and this read deliberately happens before
1738    // any skipped object's raw response would be normalized.
1739    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(&current.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
1956/// Plan a canonical integration-policy import without sending a write. This
1957/// convenience wrapper reads the source once, then delegates to the retained
1958/// artifact planner.
1959pub 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
2234/// Decode one JSON or YAML artifact without constructing transport or config.
2235pub 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
2278/// Apply a previously validated import plan. Rows are independent except for
2279/// their exact package group and shared parent-attachment snapshots.
2280pub 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        // A missing package is a shared dependency. Fleet's create path can
2404        // install it even when the policy write fails, so observation is
2405        // mandatory after every create attempt, not only decoded success.
2406        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(&current.item, &plan.space).ok().as_ref() != Some(&current.spec)
2778                {
2779                    return invalid("current integration snapshot does not normalize canonically");
2780                }
2781                let parent_ids = match read_parents(&target.effective.id, &current.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(&current.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(&current.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
2999/// Plan safe, stable-id integration-policy deletion without issuing a
3000/// mutation. Each selected object is retained exactly so `apply_delete` can
3001/// reject a Fleet change before it reaches the single-id delete route.
3002pub 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
3068/// Delete planning must bind an id selector to the id Fleet returned for that
3069/// id route. The general resolver keeps its historical public behavior for
3070/// list, get, and export; a mutation cannot accept a mismatched one-object
3071/// response as a different target.
3072async 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
3369/// Recheck the exact planning snapshots, then delete each independent target.
3370/// An acknowledged wrong-id response is deliberately not treated as a clean
3371/// deletion, and never advances shared parent expectations.
3372pub 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}