use std::fs;
use std::path::PathBuf;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Instant;
use tauri::{AppHandle, Emitter, Runtime};
use crate::collector::artifacts::{self, CollectorContext};
use crate::collector::manifest;
use crate::collector::types::*;
pub const COLLECTION_PROGRESS_EVENT: &str = "collection-progress";
pub fn run_collection<R: Runtime>(
request_id: String,
output_root: Option<String>,
enabled_families: Option<Vec<String>>,
app: AppHandle<R>,
) -> Result<CollectionResult, crate::error::AppError> {
let start = Instant::now();
let mut profile = CollectionProfile::embedded();
if let Some(ref families) = enabled_families {
profile.filter_by_families(families);
}
let total_items = profile.total_items();
let bundle_id = generate_bundle_id();
let bundle_root = resolve_bundle_root(output_root.as_deref(), &bundle_id)?;
let evidence_root = bundle_root.join("evidence");
fs::create_dir(&evidence_root).map_err(crate::error::AppError::Io)?;
for subdir in &[
"logs",
"registry",
"event-logs",
"exports",
"command-output",
] {
fs::create_dir(evidence_root.join(subdir)).map_err(crate::error::AppError::Io)?;
}
let completed = Arc::new(AtomicUsize::new(0));
let results: Arc<Mutex<Vec<ArtifactResult>>> =
Arc::new(Mutex::new(Vec::with_capacity(total_items)));
emit_progress(
&app,
&request_id,
0,
total_items,
"Starting collection...",
None,
);
let ctx_logs = CollectorContext {
bundle_evidence_root: evidence_root.clone(),
completed: Arc::clone(&completed),
results: Arc::clone(&results),
};
let ctx_registry = CollectorContext {
bundle_evidence_root: evidence_root.clone(),
completed: Arc::clone(&completed),
results: Arc::clone(&results),
};
let ctx_evtx = CollectorContext {
bundle_evidence_root: evidence_root.clone(),
completed: Arc::clone(&completed),
results: Arc::clone(&results),
};
let ctx_exports = CollectorContext {
bundle_evidence_root: evidence_root.clone(),
completed: Arc::clone(&completed),
results: Arc::clone(&results),
};
let ctx_commands = CollectorContext {
bundle_evidence_root: evidence_root.clone(),
completed: Arc::clone(&completed),
results: Arc::clone(&results),
};
let progress_completed = Arc::clone(&completed);
let progress_app = app.clone();
let progress_request_id = request_id.clone();
let progress_handle = std::thread::spawn(move || {
let mut last_reported = 0usize;
loop {
std::thread::sleep(std::time::Duration::from_millis(250));
let current = progress_completed.load(Ordering::Relaxed);
if current != last_reported {
emit_progress(
&progress_app,
&progress_request_id,
current,
total_items,
"Collecting diagnostics...",
None,
);
last_reported = current;
}
if current >= total_items {
break;
}
}
});
std::thread::scope(|s| {
s.spawn(|| artifacts::collect_logs(&profile.logs, &ctx_logs));
s.spawn(|| artifacts::export_registry_keys(&profile.registry, &ctx_registry));
s.spawn(|| artifacts::copy_event_logs(&profile.event_logs, &ctx_evtx));
s.spawn(|| artifacts::copy_exports(&profile.exports, &ctx_exports));
s.spawn(|| artifacts::run_commands(&profile.commands, &ctx_commands));
});
let _ = progress_handle.join();
let all_results = results.lock().expect("collector results mutex poisoned");
let counts = compute_counts(&all_results);
let gaps = compute_gaps(&all_results);
let duration_ms = start.elapsed().as_millis() as u64;
manifest::write_manifest(
&bundle_root,
&bundle_id,
&profile,
&all_results,
&counts,
duration_ms,
)?;
manifest::write_notes(&bundle_root, &profile, &counts, duration_ms)?;
emit_progress(
&app,
&request_id,
total_items,
total_items,
"Collection complete.",
None,
);
Ok(CollectionResult {
bundle_path: bundle_root.to_string_lossy().into_owned(),
bundle_id,
artifact_counts: counts,
duration_ms,
gaps,
})
}
const MAX_BUNDLE_HOSTNAME_CHARS: usize = 64;
fn sanitize_bundle_hostname(hostname: &str) -> String {
let sanitized = hostname
.chars()
.take(MAX_BUNDLE_HOSTNAME_CHARS)
.map(|character| {
if character.is_ascii_alphanumeric() || matches!(character, '-' | '_') {
character
} else {
'_'
}
})
.collect::<String>();
if sanitized.is_empty() {
"unknown".to_string()
} else {
sanitized
}
}
fn generate_bundle_id() -> String {
let now = chrono::Utc::now();
let hostname = std::env::var("COMPUTERNAME")
.or_else(|_| std::env::var("HOSTNAME"))
.unwrap_or_else(|_| "unknown".to_string());
let hostname = sanitize_bundle_hostname(&hostname);
let nonce = uuid::Uuid::new_v4().simple();
format!(
"CMTRACE-{}-{}-{nonce}",
now.format("%Y%m%d-%H%M%S"),
hostname
)
}
fn resolve_bundle_root(
output_root: Option<&str>,
bundle_id: &str,
) -> Result<PathBuf, crate::error::AppError> {
let base = match output_root {
Some(root) => PathBuf::from(root),
None => {
let program_data = std::env::var("ProgramData")
.or_else(|_| std::env::var("PROGRAMDATA"))
.unwrap_or_else(|_| {
std::env::temp_dir().to_string_lossy().into_owned()
});
PathBuf::from(program_data)
.join("CmtraceOpen")
.join("Evidence")
}
};
fs::create_dir_all(&base).map_err(crate::error::AppError::Io)?;
let canonical_base = base.canonicalize().map_err(crate::error::AppError::Io)?;
let bundle_root = canonical_base.join(bundle_id);
fs::create_dir(&bundle_root).map_err(crate::error::AppError::Io)?;
Ok(bundle_root)
}
fn compute_counts(results: &[ArtifactResult]) -> ArtifactCounts {
let mut collected = 0u32;
let mut missing = 0u32;
let mut failed = 0u32;
for r in results {
match r.status {
ArtifactStatus::Collected => collected += 1,
ArtifactStatus::Missing => missing += 1,
ArtifactStatus::Failed => failed += 1,
}
}
ArtifactCounts {
collected,
missing,
failed,
total: results.len() as u32,
}
}
fn compute_gaps(results: &[ArtifactResult]) -> Vec<CollectionGap> {
let mut gaps: Vec<CollectionGap> = results
.iter()
.filter(|r| !matches!(r.status, ArtifactStatus::Collected))
.map(|r| CollectionGap {
artifact_id: r.id.clone(),
category: r.category.clone(),
reason: r.error.clone().unwrap_or_else(|| format!("{:?}", r.status)),
})
.collect();
gaps.sort_by(|left, right| {
left.artifact_id
.cmp(&right.artifact_id)
.then_with(|| left.category.cmp(&right.category))
});
gaps
}
fn emit_progress<R: Runtime>(
app: &AppHandle<R>,
request_id: &str,
completed_items: usize,
total_items: usize,
message: &str,
current_item: Option<&str>,
) {
let payload = CollectionProgressPayload {
request_id: request_id.to_string(),
message: message.to_string(),
current_item: current_item.map(|s| s.to_string()),
completed_items,
total_items,
};
let _ = app.emit(COLLECTION_PROGRESS_EVENT, payload);
}
#[cfg(test)]
mod tests {
use super::{generate_bundle_id, resolve_bundle_root, sanitize_bundle_hostname};
use crate::error::AppError;
use std::fs;
use std::path::PathBuf;
use std::time::{SystemTime, UNIX_EPOCH};
#[test]
fn generated_bundle_ids_include_unpredictable_nonce() {
let first = generate_bundle_id();
let second = generate_bundle_id();
assert_ne!(
first, second,
"bundle IDs must not collide within one second"
);
let nonce = first.rsplit('-').next().expect("bundle nonce");
uuid::Uuid::parse_str(nonce).expect("bundle ID ends with a UUID nonce");
}
#[test]
fn bundle_hostname_is_a_bounded_safe_path_component() {
let sanitized = sanitize_bundle_hostname(r"..\..\outside/host name:?");
assert_eq!(sanitized, "______outside_host_name__");
assert!(sanitized
.chars()
.all(|character| character.is_ascii_alphanumeric() || matches!(character, '-' | '_')));
assert_eq!(sanitize_bundle_hostname("HOST-01_lab"), "HOST-01_lab");
assert_eq!(sanitize_bundle_hostname(""), "unknown");
assert_eq!(sanitize_bundle_hostname(&"a".repeat(1_000)).len(), 64);
}
#[test]
fn bundle_root_creation_is_exclusive() {
let base = create_temp_dir("collector-exclusive-bundle");
let bundle_id = "CMTRACE-existing";
fs::create_dir(base.join(bundle_id)).expect("create pre-existing bundle root");
let error = resolve_bundle_root(Some(base.to_string_lossy().as_ref()), bundle_id)
.expect_err("pre-existing bundle root must be rejected");
match error {
AppError::Io(error) => {
assert_eq!(error.kind(), std::io::ErrorKind::AlreadyExists)
}
other => panic!("expected AlreadyExists I/O error, got {other}"),
}
fs::remove_dir_all(&base).expect("remove temp root");
}
fn create_temp_dir(prefix: &str) -> PathBuf {
let unique = SystemTime::now()
.duration_since(UNIX_EPOCH)
.expect("system time before unix epoch")
.as_nanos();
let path = std::env::temp_dir().join(format!("{prefix}-{}-{unique}", std::process::id()));
fs::create_dir_all(&path).expect("create temp dir");
path
}
}