use anyhow::Result;
use pushkin_core::envelope::CheckResult;
use pushkin_core::manifest::Manifest;
use pushkin_core::pipeline::{check_write, WriteRequest};
use pushkin_daemon::protocol::{Request, Response, PROTOCOL_VERSION};
use pushkin_daemon::server::{self, ServerError};
use std::path::{Path, PathBuf};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Mode {
Off,
Probe,
Auto,
}
fn mode() -> Mode {
match std::env::var("PUSHKIN_DAEMON").as_deref() {
Ok("off") => Mode::Off,
Ok("auto") => Mode::Auto,
_ => Mode::Probe,
}
}
enum ServedBy {
Warm,
ColdDisabled,
ColdNoDaemon,
ColdWarmFailed(String),
}
impl ServedBy {
fn notice(&self) -> String {
match self {
Self::Warm => "pushkin: verdict served by the warm daemon.".to_owned(),
Self::ColdDisabled => "pushkin: warm path disabled (PUSHKIN_DAEMON=off); \
verdict served by the cold pipeline."
.to_owned(),
Self::ColdNoDaemon => "pushkin: no warm daemon answering; verdict served by \
the cold pipeline."
.to_owned(),
Self::ColdWarmFailed(reason) => format!(
"pushkin: the warm path was attempted and failed ({reason}); \
verdict served by the cold pipeline."
),
}
}
}
pub fn check_or_cold(manifest: &Manifest, request: &WriteRequest) -> CheckResult {
let (warm, served) = warm_check(request);
eprintln!("{}", served.notice());
let result = warm.unwrap_or_else(|| check_write(manifest, request));
let result = super::gate_read_only(manifest, result, &request.file_path);
super::gate_nested_manifest(result, &request.file_path)
}
fn warm_check(request: &WriteRequest) -> (Option<CheckResult>, ServedBy) {
let root = Path::new(".");
match mode() {
Mode::Off => (None, ServedBy::ColdDisabled),
Mode::Probe => warm_attempt(root, request),
Mode::Auto => auto_attempt(root, request),
}
}
fn warm_attempt(root: &Path, request: &WriteRequest) -> (Option<CheckResult>, ServedBy) {
match warm_request(root, request) {
Ok(result) => (Some(result), ServedBy::Warm),
Err(ServerError::NotRunning) => (None, ServedBy::ColdNoDaemon),
Err(error) => (None, ServedBy::ColdWarmFailed(error.to_string())),
}
}
fn auto_attempt(root: &Path, request: &WriteRequest) -> (Option<CheckResult>, ServedBy) {
match warm_attempt(root, request) {
(None, ServedBy::ColdNoDaemon) => match autostart(root) {
Ok(()) => warm_attempt(root, request),
Err(error) => (
None,
ServedBy::ColdWarmFailed(format!("autostart failed: {error}")),
),
},
other => other,
}
}
fn warm_request(root: &Path, request: &WriteRequest) -> Result<CheckResult, ServerError> {
let response = server::request(
root,
&Request::Check {
v: PROTOCOL_VERSION,
file_path: request.file_path.clone(),
content: request.content.clone(),
},
)?;
match response {
Response::Check { result } => Ok(result),
other => Err(ServerError::Protocol(format!(
"unexpected response to a check: {other:?}"
))),
}
}
fn autostart(root: &Path) -> std::io::Result<()> {
let exe = std::env::current_exe()?;
std::process::Command::new(exe)
.args(["daemon", "serve"])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let socket = server::socket_path(root);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
while std::time::Instant::now() < deadline {
if socket.exists() {
return Ok(());
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
Err(std::io::Error::new(
std::io::ErrorKind::TimedOut,
"daemon socket did not appear within 2s",
))
}
pub fn run_serve() -> Result<i32> {
let manifest = super::load_manifest()?;
startup_sweep();
startup_regen();
let governing = super::manifest_path()?;
server::serve_resolved(
&server::socket_path(Path::new(".")),
manifest,
false,
Some(&governing),
)?;
Ok(0)
}
#[must_use]
pub fn health_line() -> String {
match ping(Path::new(".")) {
Some(info) => format!("daemon: running (pid {})", info.pid),
None => "daemon: not running (warm path off; cold checks remain in force)".to_owned(),
}
}
fn startup_sweep() {
for finding in super::doctor::sweep_findings() {
eprintln!("pushkin daemon: doctor: {finding}");
}
}
fn startup_regen() {
let generated = Path::new("generated");
if !generated.is_dir() {
return;
}
let Ok(owned_manifest) = super::load_manifest() else {
eprintln!("pushkin daemon: epoch probe skipped: manifest failed to load");
return;
};
let stale = match pushkin_daemon::regen::probe_stale(generated, owned_manifest.schema_epoch) {
Ok(stale) => stale,
Err(error) => {
eprintln!("pushkin daemon: epoch probe failed: {error}");
return;
}
};
if stale.is_empty() {
return;
}
for path in &stale {
eprintln!(
"pushkin daemon: stale epoch: generated/{} queued for regeneration",
path.display()
);
}
let outcomes =
pushkin_daemon::regen::run_queue(&stale, std::time::Duration::from_mins(1), move |item| {
let name = item.to_string_lossy();
super::compile::regenerate_one(&owned_manifest, &name)
});
for (path, outcome) in outcomes {
use pushkin_daemon::regen::RegenOutcome;
match outcome {
RegenOutcome::Regenerated => {
eprintln!("pushkin daemon: regenerated generated/{}", path.display());
}
RegenOutcome::TimedOut => {
eprintln!(
"pushkin daemon: regeneration TIMED OUT for generated/{} (queue continued)",
path.display()
);
}
RegenOutcome::Failed(reason) => {
eprintln!(
"pushkin daemon: regeneration FAILED for generated/{}: {reason}",
path.display()
);
}
}
}
}
const CANONICAL_FILE: &str = ".pushkin/daemon.canonical";
fn current_exe_canonical() -> Result<String> {
let exe = std::env::current_exe()?.canonicalize()?;
Ok(exe.to_string_lossy().into_owned())
}
fn guard() -> Result<String> {
let me = current_exe_canonical()?;
let pin_path = Path::new(CANONICAL_FILE);
if !pin_path.exists() {
if let Some(parent) = pin_path.parent() {
std::fs::create_dir_all(parent)?;
}
std::fs::write(pin_path, format!("{me}\n"))?;
return Ok(me);
}
let pinned = std::fs::read_to_string(pin_path)?.trim().to_owned();
if pinned == me {
return Ok(me);
}
anyhow::bail!(
"pushkin daemon: this binary ({me}) is not the canonical one pinned at first start \
({pinned}). Lifecycle verbs are refused for non-canonical copies (spec §8.4). \
Use `pushkin daemon start --read-only` for a read-only daemon on a private socket, \
or have a human update {CANONICAL_FILE}."
)
}
pub fn run_start(read_only: bool) -> Result<i32> {
if read_only {
return start_read_only();
}
if let Err(refusal) = guard() {
eprintln!("{refusal}");
return Ok(1);
}
let root = Path::new(".");
if ping(root).is_some() {
println!("pushkin daemon already running");
return Ok(0);
}
match autostart(root) {
Ok(()) => {
println!("pushkin daemon started");
Ok(0)
}
Err(error) => {
eprintln!("pushkin daemon failed to start: {error}");
Ok(1)
}
}
}
fn start_read_only() -> Result<i32> {
let socket = PathBuf::from(format!(".pushkin/daemon-ro-{}.sock", std::process::id()));
let exe = std::env::current_exe()?;
std::process::Command::new(exe)
.args(["daemon", "serve", "--read-only", "--socket"])
.arg(&socket)
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()?;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
while std::time::Instant::now() < deadline {
if socket.exists() {
let absolute = socket.canonicalize()?;
println!("socket: {}", absolute.display());
println!("read-only daemon started (kill its pid to stop it)");
return Ok(0);
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
eprintln!("read-only daemon socket did not appear within 2s");
Ok(1)
}
#[must_use]
pub fn run_stop() -> i32 {
if let Err(refusal) = guard() {
eprintln!("{refusal}");
return 1;
}
let root = Path::new(".");
match server::request(
root,
&Request::Shutdown {
v: PROTOCOL_VERSION,
},
) {
Ok(Response::ShuttingDown) => {
let socket = server::socket_path(root);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
while std::time::Instant::now() < deadline {
if !socket.exists() {
break;
}
std::thread::sleep(std::time::Duration::from_millis(10));
}
println!("pushkin daemon stopped");
0
}
Err(ServerError::NotRunning) => {
println!("pushkin daemon not running");
0
}
other => {
eprintln!("pushkin daemon stop failed: {other:?}");
1
}
}
}
#[must_use]
pub fn run_restart() -> i32 {
if let Err(refusal) = guard() {
eprintln!("{refusal}");
return 1;
}
let root = Path::new(".");
if ping(root).is_some() {
let code = run_stop();
if code != 0 {
return code;
}
} else {
let _ = std::fs::remove_file(server::socket_path(root));
}
match autostart(root) {
Ok(()) => {
println!("pushkin daemon restarted");
0
}
Err(error) => {
eprintln!("pushkin daemon failed to restart: {error}");
1
}
}
}
#[must_use]
pub fn run_status() -> i32 {
let root = Path::new(".");
if let Some(info) = ping(root) {
println!(
"pushkin daemon running\n pid: {}\n version: {}\n socket: {}\n read-only: {}",
info.pid,
info.version,
server::socket_path(root).display(),
info.read_only,
);
0
} else {
println!("pushkin daemon not running");
1
}
}
fn ping(root: &Path) -> Option<pushkin_daemon::protocol::DaemonInfo> {
match server::request(
root,
&Request::Ping {
v: PROTOCOL_VERSION,
},
) {
Ok(Response::Pong { info }) => Some(info),
_ => None,
}
}
pub fn run_serve_at(socket: &Path, read_only: bool) -> Result<i32> {
let manifest = super::load_manifest()?;
let governing = super::manifest_path()?;
server::serve_resolved(socket, manifest, read_only, Some(&governing))?;
Ok(0)
}