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.ready_output.is_none() && opts.on_output_hook.is_none() && !opts.pty.unwrap_or(false)
}
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()
}
pub(crate) fn is_some(&self) -> bool {
self.0.is_some()
}
}
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,
}
impl SinkPipe {
pub(crate) fn new(log_format: String) -> Result<(Self, std::io::PipeWriter)> {
let (reader, writer) = std::io::pipe().into_diagnostic()?;
Ok((Self { reader, log_format }, 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()?;
tokio::process::Command::new(&*crate::env::PITCHFORK_BIN)
.arg("log-sink")
.arg("--daemon-id")
.arg(id.qualified())
.arg("--log-format")
.arg(&self.log_format)
.stdin(std::process::Stdio::from(reader))
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.into_diagnostic()
}
}