use std::sync::Arc;
use std::time::Duration;
use serde_json::Value;
use super::{
AdminApi, ApplyOptions, ApplyOutcome, Error, InProcessAdmin, PackageArtifact, Reporter,
};
use crate::runtime::boot_packages::BootState;
use crate::server::state::AppState;
pub const PRINCIPAL: &str = "system:boot-packages";
const LEASE_TTL_SECS: u64 = 60;
const LEASE_RENEW: Duration = Duration::from_secs(20);
const PEER_POLL: Duration = Duration::from_secs(2);
pub async fn run(state: AppState) {
if state.packages.is_empty() {
return;
}
let files = state.config.packages.apply.clone();
let budget = state.config.packages.apply_timeout_secs;
let work = apply_all(&state, &files);
if budget == 0 {
work.await;
return;
}
if tokio::time::timeout(Duration::from_secs(budget), work)
.await
.is_err()
{
let snapshot = state.packages.snapshot();
let index = snapshot
.iter()
.position(|e| !e.state.is_serving())
.unwrap_or(0);
let error = format!(
"package '{}' failed to apply at startup: not done within \
packages.apply_timeout_secs ({budget}s)",
files[index]
);
tracing::error!("{error}");
state.packages.fail(index, error);
}
}
pub fn check_files(config: &crate::config::PackagesConfig) -> Vec<(usize, String)> {
let mut problems = Vec::new();
for (index, file) in config.apply.iter().enumerate() {
let check = || -> Result<(), Error> {
let artifact = super::read_artifact(file)?;
super::verify_hash(&artifact)?;
let lint = super::lint_artifact(&artifact)?;
if lint.errors.is_empty() {
Ok(())
} else {
Err(format!(
"{} lint error(s): {}",
lint.errors.len(),
lint.errors.join("; ")
)
.into())
}
};
if let Err(e) = check() {
problems.push((index, format!("packages.apply[{index}] '{file}': {e}")));
}
}
if let Some(dir) = &config.signatures_dir
&& !std::path::Path::new(dir).is_dir()
{
problems.push((
0,
format!("packages.signatures_dir '{dir}' is not a directory"),
));
}
problems
}
async fn apply_all(state: &AppState, files: &[String]) {
if let Some((index, problem)) = check_files(&state.config.packages).into_iter().next() {
let error =
format!("[packages] failed to apply at startup: {problem} — nothing was applied");
tracing::error!("{error}");
state.packages.fail(index, error);
return;
}
for (index, file) in files.iter().enumerate() {
state
.packages
.update(index, |entry| entry.state = BootState::Applying);
match apply_one(state, index, file).await {
Ok(outcome) => state.packages.update(index, |entry| entry.state = outcome),
Err(e) => {
let error = format!("package '{file}' failed to apply at startup: {e}");
tracing::error!("{error}");
state.packages.fail(index, error);
return;
}
}
}
tracing::info!(
packages = files.len(),
"every [packages] artifact is applied and serving"
);
}
async fn apply_one(state: &AppState, index: usize, file: &str) -> Result<BootState, Error> {
let mut artifact = super::read_artifact(file)?;
super::verify_hash(&artifact)?;
let package = format!("{}@{}", artifact.package.name, artifact.package.version);
let report = Log {
package: package.clone(),
};
let lint = super::lint_artifact(&artifact)?;
for warning in &lint.warnings {
report.err(&format!("warning: {warning}"));
}
if !lint.errors.is_empty() {
return Err(format!(
"{} lint error(s): {}",
lint.errors.len(),
lint.errors.join("; ")
)
.into());
}
state.packages.update(index, |entry| {
entry.name = artifact.package.name.clone();
entry.version = artifact.package.version.clone();
entry.content_hash = artifact.package.content_hash.clone();
});
let signed = super::sign::attach_from_dir(
&mut artifact,
state.config.packages.signatures_dir.as_deref(),
&report,
)?;
let api = InProcessAdmin::new(state.clone(), PRINCIPAL, format!("package={package} boot"));
let _lease = if state.cluster.enabled {
single_flight(state, &api, &artifact).await?
} else {
None
};
let opts = ApplyOptions {
prune: None,
signed: &signed,
roll_back_superseded: false,
};
let outcome = super::apply(&api, "this node", &artifact, &opts, &report).await?;
Ok(match outcome {
ApplyOutcome::Applied | ApplyOutcome::RolledBack { .. } => BootState::Applied,
ApplyOutcome::AlreadyApplied | ApplyOutcome::Resigned => BootState::AlreadyApplied,
ApplyOutcome::Superseded { .. } => BootState::Superseded,
})
}
async fn single_flight(
state: &AppState,
api: &InProcessAdmin,
artifact: &PackageArtifact,
) -> Result<Option<LeaseRenewal>, Error> {
let gate = Arc::new(crate::cluster::JobLeaseGate::new(
state.cluster.repo.clone(),
state.cluster.instance_id.clone(),
));
let job = format!("package-apply:{}", artifact.package.name);
loop {
if gate.try_acquire(&job, LEASE_TTL_SECS).await {
let (gate, job) = (gate.clone(), job.clone());
return Ok(Some(LeaseRenewal(tokio::spawn(async move {
loop {
tokio::time::sleep(LEASE_RENEW).await;
gate.try_acquire(&job, LEASE_TTL_SECS).await;
}
}))));
}
if applied_on_target(api, artifact).await? {
tracing::info!(
package = %artifact.package.name,
version = %artifact.package.version,
"applied by a peer — reloading this node's generation"
);
crate::runtime::reload_engine(state)
.await
.map_err(|e| format!("reloading after a peer applied the package: {e}"))?;
return Ok(None);
}
tokio::time::sleep(PEER_POLL).await;
}
}
async fn applied_on_target(api: &impl AdminApi, artifact: &PackageArtifact) -> Result<bool, Error> {
let receipts: Option<Value> = api
.get_data_opt(&orion_client::paths::package(&artifact.package.name))
.await?;
Ok(receipts.is_some_and(|receipts| {
receipts["versions"]
.as_array()
.into_iter()
.flatten()
.any(|row| {
row["version"] == artifact.package.version.as_str() && row["state"] == "applied"
})
}))
}
struct LeaseRenewal(tokio::task::JoinHandle<()>);
impl Drop for LeaseRenewal {
fn drop(&mut self) {
self.0.abort();
}
}
struct Log {
package: String,
}
impl Reporter for Log {
fn out(&self, line: &str) {
tracing::info!(package = %self.package, "{line}");
}
fn err(&self, line: &str) {
match line.strip_prefix("error: ") {
Some(line) => tracing::error!(package = %self.package, "{line}"),
None => tracing::warn!(
package = %self.package,
"{}",
line.strip_prefix("warning: ").unwrap_or(line)
),
}
}
}