use std::cell::Cell;
use std::collections::{HashMap, HashSet};
use std::fs::{self, File, OpenOptions};
use std::io::{self, BufRead, BufReader, Read, Write};
use std::os::unix::fs::PermissionsExt;
use std::os::unix::net::UnixStream as StdUnixStream;
use std::os::unix::process::CommandExt;
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::time::{Duration, Instant};
use fs2::FileExt;
use serde_json::json;
use tokio::io::{AsyncBufReadExt, AsyncReadExt, AsyncWriteExt, BufReader as AsyncBufReader};
use tokio::net::{UnixListener, UnixStream};
use tokio::sync::{Notify, mpsc, oneshot};
use crate::clock::{Clock, FileClock, SystemClock};
use crate::config::{Config, ConfigError, Repository, TicketSourceConfig};
use crate::db::{Db, StoreError};
use crate::frontmatter::FrontmatterError;
use crate::ids::IdError;
use crate::protocol::{
Capability, ErrorBody, ErrorCode, Request, RequestEnvelope, RequestId, ResponseEnvelope,
};
use crate::run_ref::{RandomRunIds, RunIdSource};
use crate::run_store::RunStore;
use crate::runner::local::{process_identity_matches, process_start_time};
use crate::vendor::{CatalogError, VendorErrorClassifier};
use crate::work_state::TicketFeeder;
use crate::work_state::local::LocalSqlite;
use super::dispatcher::{
DaemonControl, DispatcherMessage, DispatcherState, RequestOrigin, internal, protocol_error,
run_dispatcher, unauthorized,
};
use super::logging::{LogLevel, OperationalLog};
use super::recovery::recover_inflight_runs;
use super::scheduler::{
index_projects, reconcile_merged_ticket_triggers, reconcile_tickets,
restore_reported_output_stalls,
};
const MAX_ENVELOPE_BYTES: u64 = 1024 * 1024;
const STARTUP_TIMEOUT: Duration = Duration::from_secs(5);
const LOCK_GRACE: Duration = Duration::from_secs(2);
const LOCK_POLL: Duration = Duration::from_millis(20);
const CLIENT_TIMEOUT: Duration = Duration::from_secs(5);
const DISPATCH_CHANNEL_CAPACITY: usize = 64;
const EVENT_RETENTION: i64 = 10_000;
static NEXT_REQUEST_ID: AtomicU64 = AtomicU64::new(1);
pub struct ClientResponse {
pub response: ResponseEnvelope,
pub started: bool,
}
pub fn request(request: Request) -> Result<ClientResponse, DaemonError> {
let cwd = std::env::current_dir().map_err(DaemonError::CurrentDirectory)?;
let repository = Repository::discover(&cwd)?;
Config::validate_client_essentials(&repository)?;
if matches!(&request, Request::Post(_)) {
Config::load(&repository)?;
}
if let Ok(response) = send_existing(&repository, request.clone()) {
return Ok(ClientResponse {
response,
started: false,
});
}
Config::load(&repository)?;
spawn_daemon(&repository)?;
let deadline = Instant::now() + STARTUP_TIMEOUT;
loop {
match send_existing(&repository, request.clone()) {
Ok(response) => {
return Ok(ClientResponse {
response,
started: true,
});
}
Err(error) if Instant::now() >= deadline => return Err(error),
Err(_) => std::thread::sleep(Duration::from_millis(20)),
}
}
}
pub fn request_running(request: Request) -> Result<Option<ResponseEnvelope>, DaemonError> {
let cwd = std::env::current_dir().map_err(DaemonError::CurrentDirectory)?;
let repository = Repository::discover(&cwd)?;
Config::validate_client_essentials(&repository)?;
match send_existing(&repository, request) {
Ok(response) => Ok(Some(response)),
Err(DaemonError::Connect(_)) => Ok(None),
Err(error) => Err(error),
}
}
pub fn serve_current_repository() -> Result<(), DaemonError> {
let executable = std::env::current_exe().map_err(DaemonError::CurrentExecutable)?;
loop {
match serve_current_repository_once()? {
ServeExit::Stopped => return Ok(()),
ServeExit::Restart { root, daemon_log } => {
if let Ok(log) = OperationalLog::open(&daemon_log) {
log.emit(LogLevel::Info, "sloop::daemon", "restart_exec");
}
let error = Command::new(&executable)
.args(["daemon", "--foreground"])
.current_dir(root)
.exec();
if let Ok(log) = OperationalLog::open(&daemon_log) {
log.emit_with_fields(
LogLevel::Error,
"sloop::daemon",
"restart_exec_failed",
json!({"path": executable, "error": error.to_string()}),
);
}
}
}
}
}
enum ServeExit {
Stopped,
Restart { root: PathBuf, daemon_log: PathBuf },
}
fn serve_current_repository_once() -> Result<ServeExit, DaemonError> {
let cwd = std::env::current_dir().map_err(DaemonError::CurrentDirectory)?;
let repository = Repository::discover(&cwd)?;
let config = Config::load(&repository)?;
let classifier = Arc::new(VendorErrorClassifier::built_in().map_err(DaemonError::Catalog)?);
fs::create_dir_all(&repository.state_dir).map_err(|source| DaemonError::Io {
path: repository.state_dir.clone(),
source,
})?;
fs::set_permissions(&repository.state_dir, fs::Permissions::from_mode(0o700)).map_err(
|source| DaemonError::Io {
path: repository.state_dir.clone(),
source,
},
)?;
let runtime_root = repository
.runtime_dir
.parent()
.expect("repository runtime directories have a parent");
fs::create_dir_all(runtime_root).map_err(|source| DaemonError::Io {
path: runtime_root.to_path_buf(),
source,
})?;
fs::set_permissions(runtime_root, fs::Permissions::from_mode(0o700)).map_err(|source| {
DaemonError::Io {
path: runtime_root.to_path_buf(),
source,
}
})?;
fs::create_dir(&repository.runtime_dir)
.or_else(|source| {
if source.kind() == io::ErrorKind::AlreadyExists {
Ok(())
} else {
Err(source)
}
})
.map_err(|source| DaemonError::Io {
path: repository.runtime_dir.clone(),
source,
})?;
fs::set_permissions(&repository.runtime_dir, fs::Permissions::from_mode(0o700)).map_err(
|source| DaemonError::Io {
path: repository.runtime_dir.clone(),
source,
},
)?;
let lock = acquire_daemon_lock(&repository.lock_path)?;
let legacy_lock_path = repository.runtime_dir.join("daemon.lock");
let legacy_lock = acquire_daemon_lock(&legacy_lock_path)?;
let identity = json!({
"pid": std::process::id(),
"started_at_ms": process_start_time(std::process::id()),
"socket": repository.operator_socket,
});
let _ = lock.set_len(0);
let _ = {
use std::io::Write as _;
writeln!(&lock, "{identity}")
};
let clock: Arc<dyn Clock> = match std::env::var_os("SLOOP_TEST_CLOCK_PATH") {
Some(path) => Arc::new(FileClock::new(path.into())),
None => Arc::new(SystemClock),
};
let db = Db::open(&repository.db_path, clock.now_ms())
.map_err(StoreError::from)
.map_err(DaemonError::Store)?;
let local_work_state = LocalSqlite::from_db(db.clone());
let run_store = RunStore::from_db(db.clone());
run_store
.clear_restart_draining(clock.now_ms())
.map_err(DaemonError::Store)?;
run_store
.trim_events(EVENT_RETENTION)
.map_err(DaemonError::Store)?;
if let Some(agent) = &config.agent {
local_work_state
.backfill_ticket_targets(&agent.default_target, clock.now_ms())
.map_err(DaemonError::Store)?;
}
let _ = index_projects(
&repository.root,
&config.project_dir,
&local_work_state,
clock.now_ms(),
&config.project_prefix,
)?;
reconcile_tickets(
&repository.root,
&local_work_state,
&run_store,
clock.now_ms(),
config.delete_missing_after_ms,
)?;
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.map_err(DaemonError::Runtime)?;
let root = repository.root.clone();
let daemon_log = repository.daemon_log.clone();
let control = runtime.block_on(serve(
repository,
config,
db,
lock,
legacy_lock,
clock,
classifier,
))?;
drop(runtime);
Ok(match control {
DaemonControl::Stop => ServeExit::Stopped,
DaemonControl::Restart => ServeExit::Restart { root, daemon_log },
})
}
async fn serve(
repository: Repository,
config: Config,
db: Db,
_lock: fs::File,
_legacy_lock: fs::File,
clock: Arc<dyn Clock>,
classifier: Arc<VendorErrorClassifier>,
) -> Result<DaemonControl, DaemonError> {
if repository.operator_socket.exists() {
fs::remove_file(&repository.operator_socket).map_err(|source| DaemonError::Io {
path: repository.operator_socket.clone(),
source,
})?;
}
let listener =
UnixListener::bind(&repository.operator_socket).map_err(|source| DaemonError::Io {
path: repository.operator_socket.clone(),
source,
})?;
fs::set_permissions(
&repository.operator_socket,
fs::Permissions::from_mode(0o600),
)
.map_err(|source| DaemonError::Io {
path: repository.operator_socket.clone(),
source,
})?;
let log = OperationalLog::open(&repository.daemon_log).map_err(|source| DaemonError::Io {
path: repository.daemon_log.clone(),
source,
})?;
log.emit(LogLevel::Info, "sloop::daemon", "daemon_started");
let run_store = RunStore::from_db(db.clone());
let paused = run_store.paused().map_err(DaemonError::Store)?;
let (dispatcher_tx, dispatcher_rx) = mpsc::channel(DISPATCH_CHANNEL_CAPACITY);
let (events_tx, events_rx) = mpsc::channel(DISPATCH_CHANNEL_CAPACITY);
let (shutdown_tx, mut shutdown_rx) = mpsc::channel::<DaemonControl>(1);
let shutdown_flag = Arc::new(AtomicBool::new(false));
let ticket_source = match &config.ticket_source {
TicketSourceConfig::Markdown => {
TicketFeeder::markdown(&repository.root, &config.ticket_dir)
}
TicketSourceConfig::Exec(argv) => TicketFeeder::exec(&repository.root, argv.clone()),
};
let work_state_author_enabled = ticket_source.supports_authoring();
let local_work_state = LocalSqlite::from_db(db.clone());
let work_state = Arc::new(LocalSqlite::from_db_with_clock_and_reporter(
db,
clock.clone(),
ticket_source.exec_reporter(),
)) as Arc<dyn crate::work_state::WorkState>;
let mut state = DispatcherState {
pid: std::process::id(),
paused,
draining: false,
restart_acknowledged: false,
restart_signalled: false,
max_agents: config.max_parallel_tasks,
stall_report_after_ms: config.stall_report_after_ms,
stall_after_ms: config.stall_after_ms,
ticket_prefix: config.ticket_prefix.clone(),
project_prefix: config.project_prefix.clone(),
running_hours: config.running_hours.clone(),
agent: config.agent.clone(),
flows: config.flows.clone(),
default_flow: config.default_flow.clone(),
flow_test_cmd: config.flow_test_cmd.clone(),
root: repository.root.clone(),
project_dir: config.project_dir.clone(),
ticket_dir: config.ticket_dir.clone(),
ticket_source,
work_state_author_enabled,
worktree_dir: repository.root.join(&config.worktree_dir),
worktree_retention_ms: config.worktree_retention_ms,
state_dir: repository.state_dir.clone(),
runtime_dir: repository.runtime_dir.clone(),
socket: repository.operator_socket.clone(),
daemon_log: repository.daemon_log.clone(),
local_work_state,
work_state,
run_store,
storage_full: Cell::new(false),
reconciliation_blocked: false,
active: HashSet::new(),
supervised: HashSet::new(),
suspected_dead: HashSet::new(),
recovering: HashSet::new(),
cancelling: HashSet::new(),
stalling: HashSet::new(),
worker_tokens: HashMap::new(),
worker_listeners: HashMap::new(),
worker_socket_paths: HashMap::new(),
pending_exits: HashMap::new(),
reported_stalls: HashMap::new(),
output_notify: Arc::new(Notify::new()),
requests_tx: dispatcher_tx.clone(),
log: log.clone(),
clock,
run_ids: Arc::new(RandomRunIds) as Arc<dyn RunIdSource>,
classifier,
shutdown: shutdown_tx.clone(),
shutdown_flag: shutdown_flag.clone(),
};
restore_reported_output_stalls(&mut state);
reconcile_merged_ticket_triggers(&state, &log);
recover_inflight_runs(&mut state, &events_tx, &log).await?;
let dispatcher_task = tokio::spawn(run_dispatcher(
state,
dispatcher_rx,
events_rx,
events_tx,
log.clone(),
));
loop {
tokio::select! {
accepted = listener.accept() => {
let (stream, _) = accepted.map_err(|source| DaemonError::Io {
path: repository.operator_socket.clone(),
source,
})?;
let dispatcher_tx = dispatcher_tx.clone();
let log = log.clone();
let shutdown = shutdown_tx.clone();
tokio::spawn(async move {
if let Err(error) = handle_connection(stream, dispatcher_tx, shutdown).await {
log.emit_with_fields(
LogLevel::Error,
"sloop::socket",
"connection_failed",
json!({"error": error.to_string()}),
);
}
});
}
control = shutdown_rx.recv() => {
let control = control.unwrap_or(DaemonControl::Stop);
shutdown_flag.store(true, Ordering::Release);
if control == DaemonControl::Stop {
log.emit(LogLevel::Info, "sloop::daemon", "daemon_stopped");
}
let _ = fs::remove_file(&repository.operator_socket);
dispatcher_task.abort();
let _ = dispatcher_task.await;
return Ok(control);
}
}
}
}
async fn handle_connection(
stream: UnixStream,
dispatcher: mpsc::Sender<DispatcherMessage>,
shutdown: mpsc::Sender<DaemonControl>,
) -> io::Result<()> {
let reader = AsyncBufReader::new(stream);
let mut limited = reader.take(MAX_ENVELOPE_BYTES + 1);
let mut bytes = Vec::new();
let read = limited.read_until(b'\n', &mut bytes).await?;
if read == 0 {
return Ok(());
}
let mut stream = limited.into_inner().into_inner();
let envelope = if bytes.len() as u64 > MAX_ENVELOPE_BYTES {
Err(protocol_error("request envelope is too large"))
} else {
std::str::from_utf8(&bytes)
.map_err(|_| protocol_error("request envelope must be UTF-8"))
.and_then(|line| RequestEnvelope::decode(line.trim_end()).map_err(|error| error.body))
};
let is_stop = matches!(
&envelope,
Ok(envelope) if matches!(envelope.request, Request::Stop(_))
);
let is_restart = matches!(
&envelope,
Ok(envelope) if matches!(envelope.request, Request::Restart(_))
);
let response = match envelope {
Ok(envelope) if envelope.token.is_some() => ResponseEnvelope::failure(
Some(envelope.id),
unauthorized("operator socket does not accept worker tokens"),
),
Ok(envelope)
if !matches!(
envelope.request.capability(),
Capability::Operator | Capability::Both
) =>
{
ResponseEnvelope::failure(
Some(envelope.id),
unauthorized(
"worker verbs are not available on the operator socket; \
run `sloop` or `sloop show <ticket>` to inspect tickets from here",
),
)
}
Ok(envelope) => dispatch_envelope(envelope, RequestOrigin::Operator, &dispatcher).await,
Err(error) => ResponseEnvelope::failure(None, error),
};
let stopping = is_stop && response.ok;
let encoded = serde_json::to_vec(&response).map_err(io::Error::other)?;
stream.write_all(&encoded).await?;
stream.write_all(b"\n").await?;
stream.shutdown().await?;
if stopping {
let _ = shutdown.send(DaemonControl::Stop).await;
} else if is_restart && response.ok {
let _ = dispatcher
.send(DispatcherMessage::RestartAcknowledged)
.await;
}
Ok(())
}
async fn handle_worker_connection(
stream: UnixStream,
run_id: String,
dispatcher: mpsc::Sender<DispatcherMessage>,
) -> io::Result<()> {
let reader = AsyncBufReader::new(stream);
let mut limited = reader.take(MAX_ENVELOPE_BYTES + 1);
let mut bytes = Vec::new();
let read = limited.read_until(b'\n', &mut bytes).await?;
if read == 0 {
return Ok(());
}
let mut stream = limited.into_inner().into_inner();
let envelope = if bytes.len() as u64 > MAX_ENVELOPE_BYTES {
Err(protocol_error("request envelope is too large"))
} else {
std::str::from_utf8(&bytes)
.map_err(|_| protocol_error("request envelope must be UTF-8"))
.and_then(|line| RequestEnvelope::decode(line.trim_end()).map_err(|error| error.body))
};
let response = match envelope {
Ok(envelope)
if !matches!(
envelope.request.capability(),
Capability::Worker | Capability::Both
) =>
{
ResponseEnvelope::failure(
Some(envelope.id),
unauthorized("operator verbs are not available on a worker socket"),
)
}
Ok(envelope) => {
let token = envelope.token.clone();
dispatch_envelope(
envelope,
RequestOrigin::Worker { run_id, token },
&dispatcher,
)
.await
}
Err(error) => ResponseEnvelope::failure(None, error),
};
let encoded = serde_json::to_vec(&response).map_err(io::Error::other)?;
stream.write_all(&encoded).await?;
stream.write_all(b"\n").await?;
stream.shutdown().await
}
async fn dispatch_envelope(
envelope: RequestEnvelope,
origin: RequestOrigin,
dispatcher: &mpsc::Sender<DispatcherMessage>,
) -> ResponseEnvelope {
let (reply_tx, reply_rx) = oneshot::channel();
let id = envelope.id;
if dispatcher
.send(DispatcherMessage::Request {
id: id.clone(),
request: envelope.request,
origin,
reply: reply_tx,
})
.await
.is_err()
{
ResponseEnvelope::failure(Some(id), internal("dispatcher is unavailable"))
} else {
reply_rx.await.unwrap_or_else(|_| {
ResponseEnvelope::failure(Some(id), internal("dispatcher dropped response"))
})
}
}
pub(super) async fn serve_worker_socket(
listener: UnixListener,
run_id: String,
dispatcher: mpsc::Sender<DispatcherMessage>,
log: OperationalLog,
) {
loop {
let Ok((stream, _)) = listener.accept().await else {
return;
};
let run_id = run_id.clone();
let dispatcher = dispatcher.clone();
let log = log.clone();
tokio::spawn(async move {
if let Err(error) = handle_worker_connection(stream, run_id.clone(), dispatcher).await {
log.emit_with_fields(
LogLevel::Error,
"sloop::socket",
"worker_connection_failed",
json!({"run_id": run_id, "error": error.to_string()}),
);
}
});
}
}
fn acquire_daemon_lock(path: &Path) -> Result<File, DaemonError> {
let lock = OpenOptions::new()
.create(true)
.truncate(false)
.read(true)
.write(true)
.open(path)
.map_err(|source| DaemonError::Io {
path: path.to_path_buf(),
source,
})?;
let deadline = Instant::now() + LOCK_GRACE;
loop {
match lock.try_lock_exclusive() {
Ok(()) => return Ok(lock),
Err(source) if source.kind() == io::ErrorKind::WouldBlock => {
if Instant::now() >= deadline {
return Err(DaemonError::AlreadyRunning);
}
std::thread::sleep(LOCK_POLL);
}
Err(source) => {
return Err(DaemonError::Io {
path: path.to_path_buf(),
source,
});
}
}
}
}
fn spawn_daemon(repository: &Repository) -> Result<(), DaemonError> {
let executable = std::env::current_exe().map_err(DaemonError::CurrentExecutable)?;
let mut command = Command::new(executable);
command
.args(["daemon", "--foreground"])
.current_dir(&repository.root)
.stdin(Stdio::null())
.stdout(Stdio::null())
.stderr(Stdio::null());
unsafe {
command.pre_exec(|| {
if libc::setsid() == -1 {
return Err(io::Error::last_os_error());
}
Ok(())
});
}
command.spawn().map(|_| ()).map_err(DaemonError::Spawn)
}
fn send_existing(
repository: &Repository,
request: Request,
) -> Result<ResponseEnvelope, DaemonError> {
match send(&repository.operator_socket, request.clone()) {
Ok(response) => Ok(response),
Err(current_error) => {
let Some(identity) = read_lock_identity(&repository.lock_path) else {
return Err(current_error);
};
if !process_identity_matches(identity.pid, identity.started_at_ms) {
return Err(current_error);
}
let Some(socket) = identity.socket else {
return Err(current_error);
};
if socket == repository.operator_socket {
return Err(current_error);
}
send(&socket, request)
}
}
}
fn send(socket: &Path, request: Request) -> Result<ResponseEnvelope, DaemonError> {
let mut stream = StdUnixStream::connect(socket).map_err(DaemonError::Connect)?;
stream
.set_read_timeout(Some(CLIENT_TIMEOUT))
.map_err(DaemonError::Connect)?;
stream
.set_write_timeout(Some(CLIENT_TIMEOUT))
.map_err(DaemonError::Connect)?;
let sequence = NEXT_REQUEST_ID.fetch_add(1, Ordering::Relaxed);
let envelope = RequestEnvelope::new(
RequestId::new(format!("req-{}-{sequence}", std::process::id())),
request,
None,
);
serde_json::to_writer(&mut stream, &envelope).map_err(DaemonError::Encode)?;
stream.write_all(b"\n").map_err(DaemonError::Write)?;
let mut line = String::new();
let mut reader = BufReader::new(stream).take(MAX_ENVELOPE_BYTES + 1);
reader.read_line(&mut line).map_err(DaemonError::Read)?;
if line.len() as u64 > MAX_ENVELOPE_BYTES {
return Err(DaemonError::InvalidResponse(
"response envelope is too large".into(),
));
}
serde_json::from_str(line.trim_end()).map_err(DaemonError::Decode)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct LockIdentity {
pub pid: u32,
pub started_at_ms: Option<i64>,
pub socket: Option<PathBuf>,
}
pub fn read_lock_identity(path: &Path) -> Option<LockIdentity> {
let content = fs::read_to_string(path).ok()?;
let value: serde_json::Value = serde_json::from_str(content.trim()).ok()?;
Some(LockIdentity {
pid: u32::try_from(value["pid"].as_u64()?).ok()?,
started_at_ms: value["started_at_ms"].as_i64(),
socket: value["socket"].as_str().map(PathBuf::from),
})
}
#[derive(Debug)]
pub enum DaemonError {
Config(ConfigError),
Catalog(CatalogError),
Store(StoreError),
WorkState(crate::work_state::SourceError),
CurrentDirectory(io::Error),
CurrentExecutable(io::Error),
Io {
path: PathBuf,
source: io::Error,
},
AlreadyRunning,
Runtime(io::Error),
Spawn(io::Error),
Connect(io::Error),
Write(io::Error),
Read(io::Error),
Encode(serde_json::Error),
Decode(serde_json::Error),
InvalidResponse(String),
Frontmatter {
path: PathBuf,
error: FrontmatterError,
},
IdAllocation(IdError),
}
impl DaemonError {
pub fn error_body(&self) -> ErrorBody {
let code = match self {
Self::Config(_) => ErrorCode::InvalidArguments,
_ => ErrorCode::DaemonUnavailable,
};
ErrorBody {
code,
message: self.to_string(),
details: json!({}),
}
}
}
impl From<ConfigError> for DaemonError {
fn from(error: ConfigError) -> Self {
Self::Config(error)
}
}
impl From<IdError> for DaemonError {
fn from(error: IdError) -> Self {
Self::IdAllocation(error)
}
}
impl std::fmt::Display for DaemonError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Config(error) => error.fmt(formatter),
Self::Catalog(error) => error.fmt(formatter),
Self::Store(error) => error.fmt(formatter),
Self::WorkState(error) => write!(formatter, "work source error: {error:?}"),
Self::CurrentDirectory(error) => {
write!(formatter, "cannot read current directory: {error}")
}
Self::CurrentExecutable(error) => {
write!(formatter, "cannot locate sloop executable: {error}")
}
Self::Io { path, source } => write!(formatter, "{}: {source}", path.display()),
Self::AlreadyRunning => formatter.write_str("another sloop daemon holds the lock"),
Self::Runtime(error) => write!(formatter, "cannot start async runtime: {error}"),
Self::Spawn(error) => write!(formatter, "cannot spawn daemon: {error}"),
Self::Connect(error) => write!(formatter, "cannot connect to daemon: {error}"),
Self::Write(error) => write!(formatter, "cannot write daemon request: {error}"),
Self::Read(error) => write!(formatter, "cannot read daemon response: {error}"),
Self::Encode(error) => write!(formatter, "cannot encode daemon request: {error}"),
Self::Decode(error) => write!(formatter, "cannot decode daemon response: {error}"),
Self::InvalidResponse(message) => formatter.write_str(message),
Self::Frontmatter { path, error } => write!(formatter, "{}: {error}", path.display()),
Self::IdAllocation(error) => error.fmt(formatter),
}
}
}
impl std::error::Error for DaemonError {}