use crate::Result;
use crate::daemon::RunOptions;
use crate::daemon_id::DaemonId;
use crate::log_store::LogStore;
use crate::supervisor::SUPERVISOR;
use miette::IntoDiagnostic;
use std::io::PipeReader;
pub(crate) fn is_supported(opts: &RunOptions) -> bool {
!opts.pty.unwrap_or(false)
}
#[derive(Default)]
pub(crate) struct WatchFor {
ready_pattern: Option<String>,
hook: Option<crate::config_types::OnOutputHook>,
}
impl WatchFor {
pub(crate) fn from_opts(id: &DaemonId, opts: &RunOptions) -> Self {
Self {
ready_pattern: opts.ready_output.as_ref().map(|o| o.pattern.clone()),
hook: opts
.on_output_hook
.as_ref()
.filter(|hook| {
hook.validate(id.name())
.inspect_err(|e| error!("{e}"))
.is_ok()
})
.cloned(),
}
}
pub(crate) fn is_empty(&self) -> bool {
self.ready_pattern.is_none() && self.hook.is_none()
}
fn args(&self, relay_token: u64) -> Vec<String> {
let mut args = Vec::new();
if self.is_empty() {
return args;
}
args.push("--relay-token".into());
args.push(relay_token.to_string());
if let Some(ref pattern) = self.ready_pattern {
args.push("--ready-pattern".into());
args.push(pattern.clone());
}
if let Some(ref hook) = self.hook {
args.push("--report-output".into());
if let Some(ref filter) = hook.filter {
args.push("--output-filter".into());
args.push(filter.clone());
}
if let Some(ref regex) = hook.regex {
args.push("--output-regex".into());
args.push(regex.clone());
}
args.push("--output-debounce-ms".into());
args.push(hook.debounce_duration().as_millis().to_string());
}
args
}
}
static NEXT_RELAY_TOKEN: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(1);
pub(crate) struct Relay {
token: u64,
tx: tokio::sync::mpsc::Sender<super::OutputLine>,
}
pub(crate) struct OutputRelay {
id: DaemonId,
token: u64,
}
impl OutputRelay {
pub(crate) fn register(
id: &DaemonId,
tx: tokio::sync::mpsc::Sender<super::OutputLine>,
) -> Self {
let token = NEXT_RELAY_TOKEN.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
SUPERVISOR
.sink_output
.lock()
.expect("sink_output lock poisoned")
.insert(id.clone(), Relay { token, tx });
Self {
id: id.clone(),
token,
}
}
pub(crate) fn token(&self) -> u64 {
self.token
}
}
impl Drop for OutputRelay {
fn drop(&mut self) {
let mut relays = SUPERVISOR
.sink_output
.lock()
.expect("sink_output lock poisoned");
if relays
.get(&self.id)
.is_some_and(|current| current.token == self.token)
{
relays.remove(&self.id);
}
}
}
pub(crate) async fn deliver_reported_line(
id: &DaemonId,
token: u64,
fires_hook: bool,
text: String,
) {
let tx = SUPERVISOR
.sink_output
.lock()
.expect("sink_output lock poisoned")
.get(id)
.filter(|relay| relay.token == token)
.map(|relay| relay.tx.clone());
let Some(tx) = tx else {
debug!("dropping output reported for {id} by a sink from a finished start attempt");
return;
};
if tx
.send(super::OutputLine {
text,
source: super::OutputSource::Sink { fires_hook },
})
.await
.is_err()
{
debug!("monitoring task for {id} stopped before its sink's output arrived");
}
}
const SPAWN_RETRY_DELAY: std::time::Duration = std::time::Duration::from_millis(200);
const SPAWN_RETRY_MAX_DELAY: std::time::Duration = std::time::Duration::from_secs(5);
pub(crate) async fn wait_for_output(
id: &DaemonId,
since: chrono::DateTime<chrono::Local>,
timeout: std::time::Duration,
) {
let _ = tokio::time::timeout(timeout, settle(id, since)).await;
}
async fn settle(id: &DaemonId, since: chrono::DateTime<chrono::Local>) {
const POLL: std::time::Duration = std::time::Duration::from_millis(20);
const SETTLED_FOR: std::time::Duration = std::time::Duration::from_millis(80);
let mut newest: Option<i64> = None;
let mut last_change = tokio::time::Instant::now();
loop {
match newest_entry_id(id, since).await {
Ok(latest) if latest != newest && latest.is_some() => {
newest = latest;
last_change = tokio::time::Instant::now();
}
Ok(_) => {
if newest.is_some() && last_change.elapsed() >= SETTLED_FOR {
return;
}
}
Err(e) => {
debug!("could not check {id}'s captured output: {e}");
}
}
tokio::time::sleep(POLL).await;
}
}
async fn newest_entry_id(
id: &DaemonId,
since: chrono::DateTime<chrono::Local>,
) -> Result<Option<i64>> {
let query = crate::log_store::LogQuery {
daemon_ids: vec![id.qualified()],
from: Some(since),
limit: Some(1),
order_desc: true,
..Default::default()
};
tokio::task::spawn_blocking(move || {
crate::log_store::sqlite::LOG_STORE
.query(&query)
.map(|entries| entries.first().map(|entry| entry.id))
})
.await
.map_err(|e| miette::miette!("log store check did not run: {e}"))?
}
pub(crate) struct PendingSink(Option<tokio::process::Child>);
impl PendingSink {
pub(crate) fn new(child: tokio::process::Child) -> Self {
Self(Some(child))
}
pub(crate) fn take(&mut self) -> Option<tokio::process::Child> {
self.0.take()
}
}
impl Drop for PendingSink {
fn drop(&mut self) {
if let Some(mut child) = self.0.take() {
let _ = child.start_kill();
}
}
}
pub(crate) struct SinkPipe {
reader: PipeReader,
log_format: String,
watch_for: WatchFor,
relay_token: u64,
}
impl SinkPipe {
pub(crate) fn new(
log_format: String,
watch_for: WatchFor,
relay_token: u64,
) -> Result<(Self, std::io::PipeWriter)> {
let (reader, writer) = std::io::pipe().into_diagnostic()?;
Ok((
Self {
reader,
log_format,
watch_for,
relay_token,
},
writer,
))
}
pub(crate) fn start(&self, id: &DaemonId) -> Result<tokio::process::Child> {
self.spawn(id)
}
pub(crate) fn supervise(self, id: DaemonId, token: u64, first: tokio::process::Child) {
tokio::spawn(async move {
let mut child = first;
loop {
match child.wait_with_output().await {
Ok(out) if out.status.success() => {
debug!("log sink for {id} finished");
break;
}
Ok(out) => {
if !SUPERVISOR.monitor_token_valid(&id, token) {
break;
}
warn!(
"log sink for {id} exited unexpectedly ({}); restarting it",
out.status
);
}
Err(e) => {
if !SUPERVISOR.monitor_token_valid(&id, token) {
break;
}
warn!("lost track of the log sink for {id}: {e}; restarting it");
}
}
let mut delay = SPAWN_RETRY_DELAY;
child = loop {
match self.spawn(&id) {
Ok(child) => break child,
Err(e) => {
error!("failed to start log sink for {id}: {e}; retrying in {delay:?}");
tokio::time::sleep(delay).await;
delay = (delay * 2).min(SPAWN_RETRY_MAX_DELAY);
if !SUPERVISOR.monitor_token_valid(&id, token) {
return;
}
}
}
};
}
});
}
fn spawn(&self, id: &DaemonId) -> Result<tokio::process::Child> {
let reader = self.reader.try_clone().into_diagnostic()?;
let mut cmd = tokio::process::Command::new(&*crate::env::PITCHFORK_BIN);
cmd.arg("log-sink")
.arg("--daemon-id")
.arg(id.qualified())
.arg("--log-format")
.arg(&self.log_format);
cmd.args(self.watch_for.args(self.relay_token));
cmd.stdin(std::process::Stdio::from(reader))
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.into_diagnostic()
}
}