Skip to main content

systemprompt_sync/jobs/
access_control_sync.rs

1//! Bootstrap job that projects the access-control baseline into the database:
2//! the YAML grants at `access-control/config.yaml` plus the marketplace-access
3//! grants declared in the services config.
4//!
5//! Mirrors [`super::ContentSyncJob`] but with a fixed direction
6//! (config → DB). Disabled by default; operators wire it in via
7//! `scheduler_config.bootstrap_jobs` so it runs once at startup.
8
9use std::path::PathBuf;
10use std::sync::Arc;
11
12use async_trait::async_trait;
13use systemprompt_database::{Database, DbPool};
14use systemprompt_models::AppPaths;
15use systemprompt_security::authz::{AccessControlIngestionService, IngestOptions};
16use systemprompt_traits::{Job, JobContext, JobResult, ProviderError, ProviderResult};
17
18use crate::local::AccessControlLocalSync;
19
20const DEFAULT_YAML_RELATIVE: &str = "access-control/config.yaml";
21
22#[derive(Debug, Clone, Copy)]
23pub struct AccessControlSyncJob;
24
25#[async_trait]
26impl Job for AccessControlSyncJob {
27    fn name(&self) -> &'static str {
28        "access_control_sync"
29    }
30
31    fn description(&self) -> &'static str {
32        "Project services/access-control YAML into access_control_rules"
33    }
34
35    fn schedule(&self) -> &'static str {
36        ""
37    }
38
39    fn tags(&self) -> Vec<&'static str> {
40        vec!["access-control", "sync", "bootstrap"]
41    }
42
43    fn enabled(&self) -> bool {
44        false
45    }
46
47    async fn execute(&self, ctx: &JobContext) -> ProviderResult<JobResult> {
48        let start = std::time::Instant::now();
49
50        let db_pool: &DbPool = ctx.db_pool::<DbPool>().ok_or_else(|| {
51            ProviderError::Configuration("DbPool not available in job context".into())
52        })?;
53
54        let paths = ctx
55            .app_paths::<Arc<AppPaths>>()
56            .ok_or_else(|| {
57                ProviderError::Configuration("AppPaths not available in job context".into())
58            })?
59            .as_ref();
60
61        let yaml_path = resolve_yaml_path(ctx, paths.system().services());
62        let override_existing = bool_param(ctx, "override_existing", true);
63        let delete_orphans = bool_param(ctx, "delete_orphans", true);
64
65        tracing::info!(
66            yaml_path = %yaml_path.display(),
67            override_existing,
68            delete_orphans,
69            "access_control_sync job started",
70        );
71
72        let sync = AccessControlLocalSync::new(Arc::<Database>::clone(db_pool), yaml_path);
73        let acl = sync
74            .sync_to_db(override_existing, delete_orphans)
75            .await
76            .map_err(|e| ProviderError::RenderFailed(e.to_string()))?;
77
78        let services = systemprompt_loader::ConfigLoader::load().map_err(|e| {
79            ProviderError::Configuration(format!("Failed to load services config: {e}"))
80        })?;
81        let svc = AccessControlIngestionService::new(db_pool)
82            .map_err(|e| ProviderError::Configuration(e.to_string()))?;
83        let options = IngestOptions {
84            override_existing,
85            delete_orphans,
86        };
87        let mkt = svc
88            .ingest_marketplace_access(&services.marketplaces, options)
89            .await
90            .map_err(|e| ProviderError::RenderFailed(e.to_string()))?;
91        let slack = svc
92            .ingest_slack_apps(&services.slack_apps, options)
93            .await
94            .map_err(|e| ProviderError::RenderFailed(e.to_string()))?;
95        let teams = svc
96            .ingest_teams_apps(&services.teams_apps, options)
97            .await
98            .map_err(|e| ProviderError::RenderFailed(e.to_string()))?;
99
100        let items_synced = acl.items_synced
101            + mkt.inserted
102            + mkt.updated
103            + slack.inserted
104            + slack.updated
105            + teams.inserted
106            + teams.updated;
107        let items_skipped = acl.items_skipped + mkt.skipped + slack.skipped + teams.skipped;
108        let items_deleted = acl.items_deleted + mkt.deleted + slack.deleted + teams.deleted;
109
110        let duration_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX);
111        tracing::info!(
112            items_synced,
113            items_skipped,
114            items_deleted,
115            duration_ms,
116            "access_control_sync job completed",
117        );
118
119        Ok(JobResult::success()
120            .with_stats(items_synced as u64, acl.errors.len() as u64)
121            .with_duration(duration_ms))
122    }
123}
124
125fn resolve_yaml_path(ctx: &JobContext, services_path: &std::path::Path) -> PathBuf {
126    ctx.parameters().get("yaml_path").map_or_else(
127        || services_path.join(DEFAULT_YAML_RELATIVE),
128        |raw| {
129            let p = std::path::Path::new(raw);
130            if p.is_absolute() {
131                p.to_path_buf()
132            } else {
133                services_path.join(p)
134            }
135        },
136    )
137}
138
139fn bool_param(ctx: &JobContext, key: &str, default: bool) -> bool {
140    ctx.parameters().get(key).map_or(default, |v| {
141        matches!(v.as_str(), "true" | "1" | "yes" | "TRUE" | "True")
142    })
143}
144
145systemprompt_provider_contracts::submit_job!(&AccessControlSyncJob);