use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use anyhow::Result;
use kamaji::native::NativeRuntime;
use kamaji::{Kamaji, MeshIdent};
use tokio::sync::oneshot;
use workload_spec::{
EnvVar, ExposeSpec, ImageRef, MeshExpose, Millis, NamespaceId, ResourceLimits, RestartPolicy,
StopPolicy, TenantId, TierTag, WorkloadSpec,
};
use super::LogBuffer;
pub(crate) const NATIVE_IDENTITY_DIGEST: &str =
"sha256:0000000000000000000000000000000000000000000000000000000000000000";
pub(crate) fn sanitize_ident(raw: &str) -> String {
raw.chars()
.map(|ch| {
if ch.is_ascii_alphanumeric() {
ch.to_ascii_lowercase()
} else {
'-'
}
})
.collect()
}
pub(crate) fn native_spec(ident: &str, argv: Vec<String>, env: Vec<EnvVar>) -> WorkloadSpec {
WorkloadSpec {
name: ident.to_string(),
image: ImageRef {
registry: "localhost".to_string(),
repository: format!("native/{ident}"),
tag: "dev".to_string(),
digest: NATIVE_IDENTITY_DIGEST.to_string(),
},
tier: TierTag("dev".to_string()),
replicas: 1,
command: Some(argv),
entrypoint: None,
workdir: None,
user: None,
env,
secrets: vec![],
volumes: vec![],
resources: ResourceLimits {
memory_mb: 512,
cpu_millis: 512,
memory_request_mb: None,
cpu_limit_millis: None,
pids_max: None,
scratch_floor_mb: None,
},
depends_on: vec![],
requires: vec![],
healthcheck: None,
restart_policy: RestartPolicy::Never,
archetype: None,
stop_policy: StopPolicy {
signal: 15,
grace_period: Millis::from_secs(5),
},
expose: ExposeSpec {
mesh: MeshExpose {
identity: MeshIdent(ident.to_string()),
ports: vec![],
allow_from: vec![],
},
public: None,
operator: None,
},
tenant: TenantId::singleton(),
namespace: NamespaceId::singleton(),
labels: Default::default(),
durability: None,
db: Vec::new(),
capabilities: Vec::new(),
annotations: Default::default(),
files: Vec::new(),
}
}
pub(crate) fn capture_paths(state_dir: &std::path::Path, ident: &str) -> (PathBuf, PathBuf) {
let dir = state_dir.join(ident);
(dir.join("stdout.log"), dir.join("stderr.log"))
}
#[cfg(unix)]
pub(crate) fn spawn_native_log_supervisor(
runtime: Arc<NativeRuntime>,
ident: MeshIdent,
log_buf: LogBuffer,
stdout_path: PathBuf,
stderr_path: PathBuf,
mut shutdown_rx: oneshot::Receiver<()>,
) -> tokio::task::JoinHandle<Result<()>> {
tokio::spawn(async move {
let mut out_tail = FileTail::new(stdout_path);
let mut err_tail = FileTail::new(stderr_path);
loop {
tokio::select! {
_ = &mut shutdown_rx => {
out_tail.drain_into(&log_buf).await;
err_tail.drain_into(&log_buf).await;
runtime.teardown_workload(&ident).await.ok();
return Ok(());
}
_ = tokio::time::sleep(Duration::from_millis(200)) => {
out_tail.drain_into(&log_buf).await;
err_tail.drain_into(&log_buf).await;
match runtime.get_workload(&ident).await {
Ok(Some(state)) if !state.status.is_terminal() => {}
Ok(Some(_)) => {
out_tail.drain_into(&log_buf).await;
err_tail.drain_into(&log_buf).await;
return Ok(());
}
Ok(None) | Err(_) => return Ok(()),
}
}
}
}
})
}
#[cfg(any(unix, test))]
pub(crate) struct FileTail {
path: PathBuf,
offset: u64,
partial: String,
}
#[cfg(any(unix, test))]
impl FileTail {
pub(crate) fn new(path: PathBuf) -> Self {
Self {
path,
offset: 0,
partial: String::new(),
}
}
pub(crate) async fn drain_into(&mut self, log_buf: &LogBuffer) {
use tokio::io::{AsyncReadExt, AsyncSeekExt};
let Ok(mut file) = tokio::fs::File::open(&self.path).await else {
return;
};
if file
.seek(std::io::SeekFrom::Start(self.offset))
.await
.is_err()
{
return;
}
let mut buf = Vec::new();
let Ok(n) = file.read_to_end(&mut buf).await else {
return;
};
if n == 0 {
return;
}
self.offset += n as u64;
self.partial.push_str(&String::from_utf8_lossy(&buf));
while let Some(idx) = self.partial.find('\n') {
let line: String = self.partial.drain(..=idx).collect();
log_buf
.push(line.trim_end_matches(['\n', '\r']).to_string())
.await;
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ident_is_lowercased_and_dns_safe() {
assert_eq!(
sanitize_ident("yah-cloud-admin/Cloud_Admin"),
"yah-cloud-admin-cloud-admin"
);
}
#[test]
fn native_spec_carries_the_identity_only_digest() {
let spec = native_spec("svc-comp", vec!["/bin/true".into()], vec![]);
assert_eq!(spec.image.digest, NATIVE_IDENTITY_DIGEST);
assert_eq!(
spec.command.as_deref(),
Some(&["/bin/true".to_string()][..])
);
assert!(matches!(spec.restart_policy, RestartPolicy::Never));
}
#[tokio::test]
async fn new_starts_at_zero_and_reads_existing_content() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("stdout.log");
tokio::fs::write(&path, b"first line\n").await.unwrap();
let buf = LogBuffer::new();
FileTail::new(path).drain_into(&buf).await;
assert_eq!(buf.since(0).await.0, vec!["first line".to_string()]);
}
#[test]
fn capture_paths_match_native_runtime_layout() {
let (out, err) = capture_paths(std::path::Path::new("/s"), "id");
assert_eq!(out, PathBuf::from("/s/id/stdout.log"));
assert_eq!(err, PathBuf::from("/s/id/stderr.log"));
}
#[tokio::test]
async fn file_tail_emits_whole_lines_and_carries_partial() {
use tokio::io::AsyncWriteExt;
let tmp = tempfile::tempdir().unwrap();
let path = tmp.path().join("stdout.log");
let log = LogBuffer::new();
let mut tail = FileTail::new(path.clone());
tail.drain_into(&log).await;
let (l0, c0) = log.since(0).await;
assert!(l0.is_empty());
let mut f = tokio::fs::File::create(&path).await.unwrap();
f.write_all(b"alpha\nbeta\npar").await.unwrap();
f.flush().await.unwrap();
tail.drain_into(&log).await;
let (l1, c1) = log.since(c0).await;
assert_eq!(l1, vec!["alpha".to_string(), "beta".to_string()]);
f.write_all(b"tial\ngamma\n").await.unwrap();
f.flush().await.unwrap();
tail.drain_into(&log).await;
let (l2, _) = log.since(c1).await;
assert_eq!(l2, vec!["partial".to_string(), "gamma".to_string()]);
}
}