mod check_cache;
mod collect;
mod commands;
mod config_supervisor;
mod env_gate;
mod finalize;
mod groups;
mod heartbeat;
mod host_perf;
mod idle_sampler;
mod job_object;
mod job_tail;
mod live_tail;
mod log_tail;
mod logs;
mod ping;
mod process;
mod process_perf;
mod self_update;
#[cfg(target_os = "windows")]
mod klp;
mod command_replay;
mod events_outbox;
mod local_scheduler;
mod nats_retry;
mod obs_outbox;
mod outbox;
mod script_cache;
mod staleness;
mod startup_event;
mod winlog;
#[cfg(target_os = "windows")]
mod client_shortcut;
#[cfg(target_os = "windows")]
mod cwd_expand;
#[cfg(target_os = "windows")]
mod process_as_user;
#[cfg(target_os = "windows")]
mod service;
#[cfg(target_os = "windows")]
mod session_supervisor;
#[cfg(target_os = "windows")]
mod capture_encode;
#[cfg(target_os = "windows")]
mod capture_frame_io;
#[cfg(target_os = "windows")]
mod capture_probe;
#[cfg(target_os = "windows")]
mod remote_session;
#[cfg(target_os = "windows")]
mod screen_capture;
use std::path::{Path, PathBuf};
use anyhow::{Context, Result};
use clap::Parser;
use kanade_shared::config::{LogSection, load_agent_config};
use kanade_shared::{default_paths, subject};
use tracing::info;
use tracing_subscriber::layer::SubscriberExt;
use tracing_subscriber::util::SubscriberInitExt;
const AGENT_VERSION: &str = env!("CARGO_PKG_VERSION");
#[derive(Parser, Debug)]
#[command(
name = "kanade-agent",
about = "Windows endpoint management agent (kanade)",
version
)]
struct Cli {
#[arg(long)]
config: Option<PathBuf>,
#[arg(long, hide = true)]
session_agent: bool,
#[arg(long, hide = true)]
capture_probe: bool,
#[arg(long, hide = true, default_value_t = 10)]
capture_probe_secs: u64,
#[arg(long, hide = true, default_value_t = 75)]
capture_probe_quality: u8,
#[arg(long, hide = true)]
capture_probe_save: Option<PathBuf>,
#[arg(long, hide = true)]
session_capture: bool,
#[arg(long, hide = true, default_value_t = 75)]
session_capture_quality: u8,
#[arg(long, hide = true, default_value_t = 10)]
session_capture_max_fps: u8,
#[arg(long, hide = true, default_value_t = 0)]
session_capture_output: u32,
#[arg(long, hide = true, default_value_t = 0)]
session_capture_secs: u64,
#[arg(long, hide = true)]
capture_decode: Option<PathBuf>,
}
fn main() -> Result<()> {
let cli = Cli::parse();
if cli.session_agent {
return run_session_agent();
}
if cli.capture_probe {
return run_capture_probe(&cli);
}
if cli.session_capture {
return run_session_capture(&cli);
}
if cli.capture_decode.is_some() {
return run_capture_decode(&cli);
}
if let Ok(exe) = std::env::current_exe() {
use kanade_shared::boot_sentinel::{BootDecision, BootSentinel, DEFAULT_MAX_ATTEMPTS};
let sentinel = BootSentinel::new(&default_paths::data_dir(), exe, AGENT_VERSION);
if let BootDecision::RolledBack { from } = sentinel.check_on_boot(DEFAULT_MAX_ATTEMPTS) {
eprintln!(
"boot sentinel: {from} crash-looped on boot — rolled back to last-good; \
exiting (64) for restart"
);
std::process::exit(64);
}
}
#[cfg(target_os = "windows")]
{
match service::try_run_as_service() {
Ok(()) => return Ok(()),
Err(e) if service::is_not_under_scm(&e) => {
}
Err(e) => return Err(anyhow::anyhow!("service dispatcher failed: {e}")),
}
}
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.context("build tokio runtime")?;
runtime.block_on(run_agent())
}
fn run_session_capture(cli: &Cli) -> Result<()> {
#[cfg(target_os = "windows")]
{
use std::io::Write;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use kanade_shared::wire::{FrameMeta, TileEncoding};
use crate::capture_encode::encode_tiles;
use crate::capture_frame_io::{FrameHeader, encode_frame};
use crate::screen_capture::{Capture, CaptureSession};
let mut session = CaptureSession::new(cli.session_capture_output).context(
"attach capture to display — this must run in the interactive desktop session",
)?;
let fps = cli.session_capture_max_fps.max(1);
let min_interval = Duration::from_millis(1000 / u64::from(fps));
let deadline = (cli.session_capture_secs > 0)
.then(|| Instant::now() + Duration::from_secs(cli.session_capture_secs));
let mut out = std::io::stdout().lock();
let mut scratch = Vec::new();
let mut frame_seq: u64 = 0;
let mut in_gap = false;
loop {
if deadline.is_some_and(|d| Instant::now() >= d) {
break;
}
let started = Instant::now();
match session.next_frame(100)? {
Capture::Idle => {
if in_gap {
in_gap = false;
let bytes = encode_frame(&FrameHeader::Resumed, &[])?;
if out.write_all(&bytes).is_err() || out.flush().is_err() {
return Ok(());
}
}
}
Capture::Unavailable(reason) => {
if !in_gap {
in_gap = true;
let bytes = encode_frame(&FrameHeader::Gap, reason.as_bytes())?;
if out.write_all(&bytes).is_err() || out.flush().is_err() {
return Ok(());
}
}
}
Capture::Frame(frame) => {
in_gap = false;
let tiles = encode_tiles(&frame, cli.session_capture_quality, &mut scratch)?;
let tile_count = u16::try_from(tiles.len()).unwrap_or(u16::MAX);
let captured_at_ms = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
for (i, tile) in tiles.iter().enumerate() {
let header = FrameHeader::Tile {
meta: FrameMeta {
frame_seq,
tile_index: u16::try_from(i).unwrap_or(u16::MAX),
tile_count,
x: tile.x,
y: tile.y,
w: tile.w,
h: tile.h,
screen_w: frame.width,
screen_h: frame.height,
captured_at_ms,
},
encoding: TileEncoding::Jpeg,
};
let bytes = encode_frame(&header, &tile.jpeg)?;
if out.write_all(&bytes).is_err() {
return Ok(());
}
}
if out.flush().is_err() {
return Ok(());
}
frame_seq += 1;
}
}
let elapsed = started.elapsed();
if elapsed < min_interval {
std::thread::sleep(min_interval - elapsed);
}
}
Ok(())
}
#[cfg(not(target_os = "windows"))]
{
let _ = cli;
anyhow::bail!("--session-capture is Windows-only")
}
}
fn run_capture_decode(cli: &Cli) -> Result<()> {
#[cfg(target_os = "windows")]
{
use std::collections::BTreeSet;
use std::io::BufReader;
use crate::capture_frame_io::read_frame;
let path = cli
.capture_decode
.as_ref()
.ok_or_else(|| anyhow::anyhow!("--capture-decode needs a path"))?;
let file = std::fs::File::open(path)
.with_context(|| format!("open frame dump {}", path.display()))?;
let mut reader = BufReader::new(file);
let mut tiles = 0u64;
let mut bytes = 0u64;
let mut frames = BTreeSet::new();
let mut screens = BTreeSet::new();
let mut largest = 0usize;
let mut gaps: Vec<String> = Vec::new();
let mut resumes = 0u64;
loop {
match read_frame(&mut reader) {
Ok(msg) => {
if let Some(reason) = msg.as_gap() {
gaps.push(reason);
continue;
}
if msg.is_resumed() {
resumes += 1;
continue;
}
let (meta, _enc) = msg
.as_tile()
.ok_or_else(|| anyhow::anyhow!("message is neither tile nor gap"))?;
tiles += 1;
bytes += msg.payload.len() as u64;
largest = largest.max(msg.payload.len());
frames.insert(meta.frame_seq);
screens.insert((meta.screen_w, meta.screen_h));
anyhow::ensure!(
msg.payload.starts_with(&[0xFF, 0xD8]),
"tile {} of frame {} is not a JPEG",
meta.tile_index,
meta.frame_seq,
);
meta.validate()
.map_err(|e| anyhow::anyhow!("frame {} geometry: {e}", meta.frame_seq))?;
anyhow::ensure!(
msg.payload.len() <= kanade_shared::wire::MAX_TILE_BYTES,
"tile {} of frame {} is {} bytes, over the {} wire budget",
meta.tile_index,
meta.frame_seq,
msg.payload.len(),
kanade_shared::wire::MAX_TILE_BYTES,
);
}
Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => break,
Err(e) => return Err(anyhow::anyhow!("decode frame {}: {e}", tiles + 1)),
}
}
println!("frames {}", frames.len());
println!("tiles {tiles}");
println!(
"tiles/frame {:.1}",
if frames.is_empty() {
0.0
} else {
tiles as f64 / frames.len() as f64
}
);
println!(
"payload {:.0} KB total | {:.0} KB mean | {:.0} KB largest",
bytes as f64 / 1024.0,
if tiles == 0 {
0.0
} else {
bytes as f64 / tiles as f64 / 1024.0
},
largest as f64 / 1024.0,
);
for (w, h) in &screens {
println!("screen {w}x{h}");
}
if !gaps.is_empty() || resumes > 0 {
println!("gaps {} ({resumes} resumed)", gaps.len());
for g in &gaps {
println!(" - {g}");
}
}
anyhow::ensure!(tiles > 0, "dump contained no frames");
println!(
"\nevery tile round-tripped, is a JPEG, fits its screen, and is under \
the {} KB wire budget",
kanade_shared::wire::MAX_TILE_BYTES / 1024,
);
Ok(())
}
#[cfg(not(target_os = "windows"))]
{
let _ = cli;
anyhow::bail!("--capture-decode is Windows-only")
}
}
fn run_capture_probe(cli: &Cli) -> Result<()> {
#[cfg(target_os = "windows")]
{
capture_probe::run(
cli.capture_probe_secs,
cli.capture_probe_quality,
cli.capture_probe_save.clone(),
)
}
#[cfg(not(target_os = "windows"))]
{
let _ = cli;
anyhow::bail!(
"--capture-probe is Windows-only: it measures DXGI Desktop Duplication, \
which has no non-Windows counterpart in this agent"
)
}
}
fn run_session_agent() -> Result<()> {
#[cfg(target_os = "windows")]
{
use std::io::Write;
use windows::Win32::System::SystemInformation::GetTickCount;
use windows::Win32::UI::Input::KeyboardAndMouse::{GetLastInputInfo, LASTINPUTINFO};
let mut out = std::io::stdout();
loop {
let line = unsafe {
let mut lii = LASTINPUTINFO {
cbSize: std::mem::size_of::<LASTINPUTINFO>() as u32,
dwTime: 0,
};
if GetLastInputInfo(&mut lii).as_bool() {
let idle_ms = GetTickCount().wrapping_sub(lii.dwTime);
format!("{{\"idle_ms\":{idle_ms}}}\n")
} else {
"{\"idle_ms\":null}\n".to_string()
}
};
if out.write_all(line.as_bytes()).is_err() || out.flush().is_err() {
break;
}
std::thread::sleep(env_gate::SESSION_IDLE_SAMPLE_INTERVAL);
}
}
#[cfg(not(target_os = "windows"))]
{
}
Ok(())
}
pub(crate) async fn run_agent() -> Result<()> {
let cli = Cli::parse();
let cfg_path =
default_paths::find_config(cli.config.as_deref(), "KANADE_AGENT_CONFIG", "agent.toml")?;
let cfg =
load_agent_config(&cfg_path).with_context(|| format!("load config from {cfg_path:?}"))?;
let _log_guard = init_tracing(&cfg.log)
.with_context(|| format!("init tracing from [log] in {cfg_path:?}"))?;
cleanup_stale_upgrade_artifacts();
info!(
pc_id = %cfg.agent.id,
nats_url = %cfg.agent.nats_url,
version = AGENT_VERSION,
log_path = %cfg.log.path,
log_keep_days = cfg.log.keep_days,
"starting kanade-agent",
);
let staleness_tracker = staleness::Tracker::new();
let client = kanade_shared::nats_client::connect_with_event_callback(
kanade_shared::nats_client::NatsRole::Agent,
&cfg.agent.nats_url,
staleness_tracker.on_event(),
)
.await?;
info!("connected to NATS");
let cmd_all = client.subscribe(subject::COMMANDS_ALL).await?;
let cmd_self = client
.subscribe(subject::commands_pc(&cfg.agent.id))
.await?;
info!(
commands_all = subject::COMMANDS_ALL,
commands_self = %subject::commands_pc(&cfg.agent.id),
"subscribed",
);
let pc_id = cfg.agent.id.clone();
let cfg_rx = config_supervisor::spawn(client.clone(), pc_id.clone(), staleness_tracker.clone());
tokio::spawn(heartbeat::heartbeat_loop(
client.clone(),
pc_id.clone(),
AGENT_VERSION.to_string(),
cfg_rx.clone(),
));
tokio::spawn(host_perf::host_perf_loop(
client.clone(),
pc_id.clone(),
cfg_rx.clone(),
));
tokio::spawn(process_perf::process_perf_loop(
client.clone(),
pc_id.clone(),
cfg_rx.clone(),
));
#[cfg(target_os = "windows")]
client_shortcut::spawn(cfg_rx.clone());
tokio::spawn(self_update::run(
client.clone(),
pc_id.clone(),
AGENT_VERSION.to_string(),
cfg_rx.clone(),
staleness_tracker.clone(),
));
tokio::spawn(async {
tokio::time::sleep(std::time::Duration::from_secs(30)).await;
if let Ok(exe) = std::env::current_exe() {
let sentinel = kanade_shared::boot_sentinel::BootSentinel::new(
&default_paths::data_dir(),
exe,
AGENT_VERSION,
);
if let Err(e) = sentinel.confirm_healthy() {
tracing::warn!(error = %e, "boot sentinel: confirm_healthy failed");
}
}
});
tokio::spawn(logs::serve(
client.clone(),
pc_id.clone(),
std::path::PathBuf::from(&cfg.log.path),
staleness_tracker.clone(),
));
tokio::spawn(job_tail::serve(
client.clone(),
pc_id.clone(),
staleness_tracker.clone(),
));
tokio::spawn(ping::serve(
client.clone(),
pc_id.clone(),
AGENT_VERSION.to_string(),
std::env::var("COMPUTERNAME")
.ok()
.or_else(|| std::env::var("HOSTNAME").ok()),
Some(std::env::consts::OS.to_string()),
staleness_tracker.clone(),
));
let check_sink =
check_cache::CheckSink::load(default_paths::data_dir().join("check_results.json"));
#[cfg(target_os = "windows")]
let notif_tx =
tokio::sync::broadcast::channel::<kanade_shared::ipc::notifications::Notification>(
klp::notify_bus::BROADCAST_CAPACITY,
)
.0;
#[cfg(target_os = "windows")]
let amend_tx = tokio::sync::broadcast::channel::<
kanade_shared::ipc::notifications::NotificationAmend,
>(klp::notify_bus::BROADCAST_CAPACITY)
.0;
#[cfg(target_os = "windows")]
{
let initial_snapshot = klp::state::eval_once(
&pc_id,
AGENT_VERSION,
&cfg_rx.borrow(),
klp::state::client_online(&client),
&check_sink.checks(),
);
let (state_tx, state_rx) = tokio::sync::watch::channel(initial_snapshot);
tokio::spawn(klp::state::eval_loop(
state_tx,
cfg_rx.clone(),
pc_id.clone(),
AGENT_VERSION.to_string(),
client.clone(),
check_sink.clone(),
));
let _klp_handle = klp::server::spawn(klp::server::ListenerContext {
pc_id: std::sync::Arc::from(pc_id.as_str()),
agent_version: std::sync::Arc::from(AGENT_VERSION),
config_rx: cfg_rx.clone(),
state_rx,
log_path: std::path::PathBuf::from(&cfg.log.path),
nats: client.clone(),
notif_tx: notif_tx.clone(),
amend_tx: amend_tx.clone(),
});
}
if !cfg.agent.groups.is_empty() {
tracing::warn!(
local_groups = ?cfg.agent.groups,
"agent.toml::[agent] groups is deprecated; use `kanade agent groups set` instead — local value is ignored",
);
}
let dedup = commands::shared_dedup_cache();
let script_cache = script_cache::ScriptCache::new(
async_nats::jetstream::new(client.clone()),
default_paths::data_dir().join("script_cache"),
);
let (groups_rx, _groups_handle) = groups::spawn(
client.clone(),
pc_id.clone(),
dedup.clone(),
staleness_tracker.clone(),
script_cache.clone(),
check_sink.clone(),
);
#[cfg(target_os = "windows")]
klp::notify_bus::spawn(
client.clone(),
pc_id.clone(),
groups_rx.clone(),
notif_tx,
amend_tx,
);
command_replay::spawn(
client.clone(),
pc_id.clone(),
dedup.clone(),
staleness_tracker.clone(),
script_cache.clone(),
check_sink.clone(),
groups_rx.clone(),
);
let outbox_dir = default_paths::data_dir().join("outbox");
let _outbox_handle = outbox::spawn_drain(client.clone(), outbox_dir.clone());
let events_outbox_dir = default_paths::data_dir().join("events-outbox");
let _events_outbox_handle =
events_outbox::spawn_drain(client.clone(), events_outbox_dir.clone());
let obs_outbox_dir = default_paths::data_dir().join("obs-outbox");
let _obs_outbox_handle = obs_outbox::spawn_drain(client.clone(), obs_outbox_dir.clone());
startup_event::emit(&pc_id, AGENT_VERSION, &obs_outbox_dir);
tokio::spawn(idle_sampler::run(pc_id.clone(), obs_outbox_dir.clone()));
#[cfg(target_os = "windows")]
if let Ok(self_exe) = std::env::current_exe() {
tokio::spawn(session_supervisor::run(self_exe.clone()));
tokio::spawn(remote_session::serve(
client.clone(),
pc_id.clone(),
self_exe,
));
}
tokio::spawn(winlog::run(pc_id.clone(), obs_outbox_dir.clone()));
let completions_path = default_paths::data_dir().join("local_completions.json");
local_scheduler::spawn(
client.clone(),
pc_id.clone(),
completions_path,
groups_rx,
staleness_tracker.clone(),
script_cache.clone(),
check_sink.clone(),
);
let _ = tokio::join!(
commands::command_loop(
client.clone(),
pc_id.clone(),
dedup.clone(),
staleness_tracker.clone(),
cmd_all,
script_cache.clone(),
check_sink.clone(),
),
commands::command_loop(
client.clone(),
pc_id.clone(),
dedup.clone(),
staleness_tracker.clone(),
cmd_self,
script_cache.clone(),
check_sink.clone(),
),
);
Ok(())
}
fn init_tracing(log: &LogSection) -> Result<Option<tracing_appender::non_blocking::WorkerGuard>> {
let env_filter = tracing_subscriber::EnvFilter::try_from_default_env()
.unwrap_or_else(|_| log.level.clone().into());
if log.keep_days == 0 {
let _ = tracing_subscriber::registry()
.with(env_filter)
.with(tracing_subscriber::fmt::layer().with_writer(std::io::stdout))
.try_init();
return Ok(None);
}
let path = Path::new(&log.path);
let dir = path
.parent()
.with_context(|| format!("[log] path '{}' has no parent dir", log.path))?;
let stem = path.file_stem().and_then(|s| s.to_str()).unwrap_or("agent");
let ext = path.extension().and_then(|s| s.to_str()).unwrap_or("log");
std::fs::create_dir_all(dir).with_context(|| format!("create log dir {dir:?}"))?;
let appender = tracing_appender::rolling::Builder::new()
.filename_prefix(stem)
.filename_suffix(ext)
.rotation(tracing_appender::rolling::Rotation::DAILY)
.max_log_files(log.keep_days)
.build(dir)
.context("build rolling file appender")?;
let (file_writer, guard) = tracing_appender::non_blocking(appender);
let _ = tracing_subscriber::registry()
.with(env_filter)
.with(tracing_subscriber::fmt::layer().with_writer(std::io::stdout))
.with(
tracing_subscriber::fmt::layer()
.with_writer(file_writer)
.with_ansi(false),
)
.try_init();
Ok(Some(guard))
}
fn cleanup_stale_upgrade_artifacts() {
let Ok(current) = std::env::current_exe() else {
return;
};
let Some(exe_dir) = current.parent() else {
return;
};
let Some(exe_name) = current.file_name().and_then(|n| n.to_str()) else {
return;
};
for suffix in ["old", "new"] {
let path = exe_dir.join(format!("{exe_name}.{suffix}"));
if !path.exists() {
continue;
}
match std::fs::remove_file(&path) {
Ok(_) => tracing::info!(?path, suffix, "removed stale upgrade artifact"),
Err(e) => {
tracing::warn!(?path, suffix, error = %e, "couldn't remove stale upgrade artifact")
}
}
}
}