systemprompt_api/routes/admin/services/
refresh.rs1use 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
50static 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#[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 #[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 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 let unreconciled = new_hash.is_some() && previous.last_reconciled_hash != new_hash;
164 let mut reconciled = false;
165 if changed || unreconciled {
166 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
218fn 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}