Skip to main content

systemprompt_api/routes/admin/services/
refresh.rs

1//! `POST /admin/services/refresh`.
2//!
3//! Re-resolves the configured sources, recomposes the tree and — when the
4//! composition changed — projects the new composition into the authz tables
5//! and refreshes the skill inventory, all in-process. The routes that hand
6//! marketplaces, plugins and skills to clients reload the services tree per
7//! request through the `current` link the recompose just swapped, so a
8//! marketplace-only kit is live the moment this returns. `restart=true`
9//! remains an explicit opt-in for the one thing a running process cannot
10//! re-read: the static services config behind governance hooks.
11//!
12//! The pipeline is [`ServicesRefresh`], an in-process handle an extension
13//! router receives as an axum extension (see `extension_mount`), so a console
14//! that has authorised a caller by its own rule — a marketplace participant
15//! syncing their own kit, say — runs the same refresh the admin route runs
16//! without minting an admin token. One process-wide lock guards both. The
17//! handle never restarts the process: an extension router is authorised by
18//! `AuthzPolicy::user()`, so the restart stays on the admin route alone.
19//!
20//! Copyright (c) systemprompt.io — Business Source License 1.1.
21//! See <https://systemprompt.io> for licensing details.
22
23use std::sync::LazyLock;
24use std::time::Duration;
25
26use axum::Json;
27use axum::extract::{Extension, Query, State};
28use serde::Deserialize;
29use systemprompt_config::{ProfileBootstrap, SecretsBootstrap};
30use systemprompt_identifiers::UserId;
31use systemprompt_loader::bundle::bootstrap::baked::BASE_SOURCE_NAME;
32use systemprompt_loader::bundle::{BundleCache, cache_root};
33use systemprompt_loader::services_root::ServicesRootBootstrap;
34use systemprompt_loader::{ConfigLoadError, ConfigLoader, ServicesSourceBootstrap};
35use systemprompt_manifest::services::bundle::ServicesBundleState;
36use systemprompt_models::RequestContext;
37use systemprompt_models::api::ApiError;
38use systemprompt_runtime::managed::inventory::publish_latest;
39use systemprompt_runtime::services_reconcile::{ReconcileOutcome, reconcile_fetched_services};
40use systemprompt_runtime::{AppContext, RuntimeError};
41
42use super::{
43    RefreshLock, ServicesRefreshResponse, composed_hash_of, provenance_view, source_views,
44};
45use crate::error::ApiHttpError;
46
47const RESTART_DELAY: Duration = Duration::from_millis(250);
48const RESTART_REASON: &str = "admin services refresh";
49
50// Why: one lock for the process, not one per router. The admin route and every
51// extension handle share it, so two callers on different routes cannot fetch
52// at once. `RefreshLock` clones share the mutex, so the router's extension is
53// a clone of this same lock.
54static REFRESH_LOCK: LazyLock<RefreshLock> = LazyLock::new(RefreshLock::default);
55
56#[must_use]
57pub fn process_refresh_lock() -> RefreshLock {
58    REFRESH_LOCK.clone()
59}
60
61/// In-process handle to the fetch-verify-compose-reconcile pipeline.
62#[derive(Debug, Clone)]
63pub struct ServicesRefresh {
64    ctx: AppContext,
65    lock: RefreshLock,
66}
67
68impl ServicesRefresh {
69    #[must_use]
70    pub fn new(ctx: &AppContext) -> Self {
71        Self::with_lock(ctx, process_refresh_lock())
72    }
73
74    // Why: the lock is injectable so a test can hold it and prove the second
75    // caller is refused; production always passes the process lock.
76    #[must_use]
77    pub fn with_lock(ctx: &AppContext, lock: RefreshLock) -> Self {
78        Self {
79            ctx: ctx.clone(),
80            lock,
81        }
82    }
83
84    pub async fn run(&self, actor: &UserId) -> Result<ServicesRefreshResponse, ApiHttpError> {
85        let _guard = self.lock.try_acquire().ok_or_else(busy)?;
86        run_refresh(&self.ctx, actor, false).await
87    }
88}
89
90#[derive(Debug, thiserror::Error)]
91enum RefreshError {
92    #[error("recomposed services config is invalid")]
93    RecomposedConfig(#[source] ConfigLoadError),
94    #[error("services reconcile failed")]
95    Reconcile(#[source] RuntimeError),
96}
97
98impl From<RefreshError> for ApiHttpError {
99    fn from(err: RefreshError) -> Self {
100        let context = match &err {
101            RefreshError::RecomposedConfig(_) => "Recomposed services config is invalid",
102            RefreshError::Reconcile(_) => "Services reconcile failed",
103        };
104        ApiError::internal(context, err).into()
105    }
106}
107
108fn busy() -> ApiHttpError {
109    ApiError::conflict("a services refresh is already running").into()
110}
111
112#[derive(Debug, Clone, Copy, Default, Deserialize)]
113pub struct RefreshQuery {
114    #[serde(default)]
115    pub restart: bool,
116}
117
118pub async fn refresh(
119    State(ctx): State<AppContext>,
120    Extension(lock): Extension<RefreshLock>,
121    Extension(req_ctx): Extension<RequestContext>,
122    Query(query): Query<RefreshQuery>,
123) -> Result<Json<ServicesRefreshResponse>, ApiHttpError> {
124    let _guard = lock.try_acquire().ok_or_else(busy)?;
125    run_refresh(&ctx, req_ctx.user_id(), query.restart)
126        .await
127        .map(Json)
128}
129
130async fn run_refresh(
131    ctx: &AppContext,
132    actor: &UserId,
133    restart: bool,
134) -> Result<ServicesRefreshResponse, ApiHttpError> {
135    let profile = ProfileBootstrap::get()?;
136    let secrets = SecretsBootstrap::get()?;
137
138    // Why: the boot-time root is a static; after an in-place import the cache
139    // state names the composition actually being served, so "changed" is
140    // measured against that and a repeat import is a no-op.
141    let cache = BundleCache::new(cache_root(profile));
142    let previous = cache.read_state();
143    let served_hash = (!previous.composed_hash.is_empty())
144        .then_some(previous.composed_hash.clone())
145        .or_else(|| {
146            ServicesRootBootstrap::get()
147                .and_then(composed_hash_of)
148                .map(str::to_owned)
149        });
150
151    let resolved = ServicesSourceBootstrap::resolve(
152        profile,
153        |name| secrets.get(name).cloned(),
154        env!("CARGO_PKG_VERSION"),
155    )
156    .await?;
157
158    let new_hash = composed_hash_of(&resolved).map(str::to_owned);
159    let changed = new_hash != served_hash;
160    // Why: a composition that was swapped in but never projected (a failed
161    // earlier reconcile) is finished by the next import even though nothing
162    // else changed.
163    let unreconciled = new_hash.is_some() && previous.last_reconciled_hash != new_hash;
164    let mut reconciled = false;
165    if changed || unreconciled {
166        // Why: the boot-time root is a static that still names the previous
167        // tree; the recomposed tree is the one whose config is projected.
168        let services =
169            ConfigLoader::reload_from_path(&resolved.path.join("config").join("config.yaml"))
170                .map_err(RefreshError::RecomposedConfig)?;
171        let outcome = reconcile_fetched_services(profile, &resolved, &services, ctx.db_pool())
172            .await
173            .map_err(RefreshError::Reconcile)?;
174        reconciled = outcome == ReconcileOutcome::Projected;
175
176        let system_admin = ctx.system_admin().id().clone();
177        if let Err(error) = publish_latest(ctx, &system_admin, actor).await {
178            tracing::warn!(%error, "Inventory refresh after services import failed; the scheduled pass will retry");
179        }
180    }
181
182    let state = cache.read_state();
183    let restart_recommended = changed && owns_static_config(&cache, &state);
184    let restarting = changed && restart;
185
186    tracing::info!(
187        user_id = %actor,
188        changed,
189        reconciled,
190        restart_recommended,
191        composed_hash = new_hash.as_deref().unwrap_or("none"),
192        provenance = %provenance_view(&resolved.provenance).kind,
193        restarting,
194        "Admin services refresh"
195    );
196
197    if restarting {
198        let restart_ctx = ctx.clone();
199        ctx.background_tasks()
200            .spawn_cancellable("services_refresh_restart", |cancel| async move {
201                tokio::select! {
202                    () = cancel.cancelled() => {},
203                    () = tokio::time::sleep(RESTART_DELAY) => restart_ctx.request_restart(RESTART_REASON),
204                }
205            });
206    }
207
208    Ok(ServicesRefreshResponse {
209        changed,
210        composed_hash: new_hash,
211        sources: source_views(&state),
212        reconciled,
213        restart_recommended,
214        restarting,
215    })
216}
217
218// Why: governance hooks are read once at boot into the static services config;
219// a bundle that ships hooks is the one case an in-process import cannot fully
220// serve, so the caller is told a restart would complete it. A manifest that
221// cannot be read may own hooks, so it recommends the restart too.
222fn owns_static_config(cache: &BundleCache, state: &ServicesBundleState) -> bool {
223    state
224        .sources
225        .iter()
226        .filter(|(name, _)| name.as_str() != BASE_SOURCE_NAME)
227        .any(|(name, fetched)| match cache.read_manifest(name, &fetched.content_hash) {
228            Ok(signed) => !signed.manifest.owns.hooks.is_empty(),
229            Err(error) => {
230                tracing::warn!(source = %name, %error, "Cached bundle manifest unreadable; recommending a restart");
231                true
232            },
233        })
234}