use std::sync::Arc;
use std::time::{Duration, Instant};
use code_system_graph_model::stable_id;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use thiserror::Error;
pub const DEFAULT_MAX_SCAN_WALL_TIME_MS: u64 = 21_600_000;
pub const DEFAULT_MAX_NO_PROGRESS_TIME_MS: u64 = 300_000;
pub const DEFAULT_MAX_CODEGRAPH_SYNC_WALL_TIME_MS_PER_REPO: u64 = 3_600_000;
pub const DEFAULT_MAX_WORKER_MEMORY_BYTES: u64 = 17_179_869_184;
pub const DEFAULT_GRACEFUL_TERMINATION_MS: u64 = 5_000;
pub const DEFAULT_WATCH_IDLE_TIMEOUT_MS: u64 = 28_800_000;
pub const DEFAULT_MAX_WATCH_SESSION_WALL_TIME_MS: u64 = 86_400_000;
pub const DEFAULT_MIN_WATCH_RESCAN_INTERVAL_MS: u64 = 10_000;
pub const DEFAULT_MAX_CHECKPOINT_CACHE_BYTES: u64 = 10_737_418_240;
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub struct ExecutionPolicyOverrides {
pub max_scan_wall_time_ms: Option<u64>,
pub max_no_progress_time_ms: Option<u64>,
#[serde(rename = "maxCodeGraphSyncWallTimeMsPerRepo")]
pub max_codegraph_sync_wall_time_ms_per_repo: Option<u64>,
pub max_worker_memory_bytes: Option<u64>,
pub graceful_termination_ms: Option<u64>,
pub watch_idle_timeout_ms: Option<u64>,
pub max_watch_session_wall_time_ms: Option<u64>,
pub min_watch_rescan_interval_ms: Option<u64>,
pub max_checkpoint_cache_bytes: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct ExecutionPolicy {
pub max_scan_wall_time_ms: u64,
pub max_no_progress_time_ms: u64,
#[serde(rename = "maxCodeGraphSyncWallTimeMsPerRepo")]
pub max_codegraph_sync_wall_time_ms_per_repo: u64,
pub max_worker_memory_bytes: u64,
pub graceful_termination_ms: u64,
pub watch_idle_timeout_ms: u64,
pub max_watch_session_wall_time_ms: u64,
pub min_watch_rescan_interval_ms: u64,
pub max_checkpoint_cache_bytes: u64,
}
impl Default for ExecutionPolicy {
fn default() -> Self {
Self {
max_scan_wall_time_ms: DEFAULT_MAX_SCAN_WALL_TIME_MS,
max_no_progress_time_ms: DEFAULT_MAX_NO_PROGRESS_TIME_MS,
max_codegraph_sync_wall_time_ms_per_repo:
DEFAULT_MAX_CODEGRAPH_SYNC_WALL_TIME_MS_PER_REPO,
max_worker_memory_bytes: DEFAULT_MAX_WORKER_MEMORY_BYTES,
graceful_termination_ms: DEFAULT_GRACEFUL_TERMINATION_MS,
watch_idle_timeout_ms: DEFAULT_WATCH_IDLE_TIMEOUT_MS,
max_watch_session_wall_time_ms: DEFAULT_MAX_WATCH_SESSION_WALL_TIME_MS,
min_watch_rescan_interval_ms: DEFAULT_MIN_WATCH_RESCAN_INTERVAL_MS,
max_checkpoint_cache_bytes: DEFAULT_MAX_CHECKPOINT_CACHE_BYTES,
}
}
}
impl ExecutionPolicy {
pub fn resolve(
overrides: Option<&ExecutionPolicyOverrides>,
) -> Result<Self, InvalidExecutionPolicy> {
let mut policy = Self::default();
if let Some(values) = overrides {
macro_rules! apply {
($field:ident) => {
if let Some(value) = values.$field {
policy.$field = value;
}
};
}
apply!(max_scan_wall_time_ms);
apply!(max_no_progress_time_ms);
apply!(max_codegraph_sync_wall_time_ms_per_repo);
apply!(max_worker_memory_bytes);
apply!(graceful_termination_ms);
apply!(watch_idle_timeout_ms);
apply!(max_watch_session_wall_time_ms);
apply!(min_watch_rescan_interval_ms);
apply!(max_checkpoint_cache_bytes);
}
policy.validate()?;
Ok(policy)
}
fn validate(&self) -> Result<(), InvalidExecutionPolicy> {
for (field, value) in self.canonical_values() {
let invalid_bytes = field.ends_with("Bytes") && usize::try_from(value).is_err();
let invalid_sqlite_quota =
field == "maxCheckpointCacheBytes" && i64::try_from(value).is_err();
let invalid_deadline = field.ends_with("Ms")
&& Instant::now()
.checked_add(Duration::from_millis(value))
.is_none();
if value == 0 || invalid_bytes || invalid_sqlite_quota || invalid_deadline {
return Err(InvalidExecutionPolicy::InvalidValue { field, value });
}
}
Self::require_not_greater(
"maxNoProgressTimeMs",
self.max_no_progress_time_ms,
"maxScanWallTimeMs",
self.max_scan_wall_time_ms,
)?;
Self::require_not_greater(
"maxCodeGraphSyncWallTimeMsPerRepo",
self.max_codegraph_sync_wall_time_ms_per_repo,
"maxScanWallTimeMs",
self.max_scan_wall_time_ms,
)?;
Self::require_not_greater(
"gracefulTerminationMs",
self.graceful_termination_ms,
"maxNoProgressTimeMs",
self.max_no_progress_time_ms,
)?;
Self::require_not_greater(
"watchIdleTimeoutMs",
self.watch_idle_timeout_ms,
"maxWatchSessionWallTimeMs",
self.max_watch_session_wall_time_ms,
)?;
Self::require_not_greater(
"minWatchRescanIntervalMs",
self.min_watch_rescan_interval_ms,
"watchIdleTimeoutMs",
self.watch_idle_timeout_ms,
)
}
fn require_not_greater(
field: &'static str,
value: u64,
maximum_field: &'static str,
maximum: u64,
) -> Result<(), InvalidExecutionPolicy> {
if value > maximum {
return Err(InvalidExecutionPolicy::InvalidRelationship {
field,
value,
maximum_field,
maximum,
});
}
Ok(())
}
fn canonical_values(&self) -> [(&'static str, u64); 9] {
[
("maxScanWallTimeMs", self.max_scan_wall_time_ms),
("maxNoProgressTimeMs", self.max_no_progress_time_ms),
(
"maxCodeGraphSyncWallTimeMsPerRepo",
self.max_codegraph_sync_wall_time_ms_per_repo,
),
("maxWorkerMemoryBytes", self.max_worker_memory_bytes),
("gracefulTerminationMs", self.graceful_termination_ms),
("watchIdleTimeoutMs", self.watch_idle_timeout_ms),
(
"maxWatchSessionWallTimeMs",
self.max_watch_session_wall_time_ms,
),
(
"minWatchRescanIntervalMs",
self.min_watch_rescan_interval_ms,
),
("maxCheckpointCacheBytes", self.max_checkpoint_cache_bytes),
]
}
#[must_use]
pub fn fingerprint(&self) -> String {
let canonical = self
.canonical_values()
.into_iter()
.map(|(name, value)| format!("{name}={value}"))
.collect::<Vec<_>>()
.join(";");
stable_id("execution-policy", &canonical)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Error)]
pub enum InvalidExecutionPolicy {
#[error("execution policy `{field}` must be positive and representable; received {value}")]
InvalidValue {
field: &'static str,
value: u64,
},
#[error("execution policy `{field}` ({value}) must not exceed `{maximum_field}` ({maximum})")]
InvalidRelationship {
field: &'static str,
value: u64,
maximum_field: &'static str,
maximum: u64,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum JobPhase {
Configuration,
Discovery,
Fingerprinting,
Extraction,
GraphAssembly,
Communities,
Publication,
CodeGraphSync,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecutionResource {
WallTimeMs,
NoProgressTimeMs,
WorkerMemoryBytes,
WorkUnits,
WorkerProcess,
WorkerProtocolBytes,
Cancellation,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "camelCase")]
pub struct ExecutionSummary {
pub run_id: String,
pub duration_ms: u64,
pub peak_worker_memory_bytes: u64,
pub completed_work_units: u64,
pub checkpoint_hits: u64,
pub checkpoints_written: u64,
pub measured_artifacts: u64,
pub artifact_duration_p50_ms: u64,
pub artifact_duration_p95_ms: u64,
pub artifact_duration_p99_ms: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Error)]
#[error(
"execution `{run_id}` exceeded {resource:?} during {phase:?}: observed {observed}, maximum {maximum}, completed {completed_units} units"
)]
pub struct ExecutionLimitExceeded {
pub run_id: String,
pub phase: JobPhase,
pub resource: ExecutionResource,
pub observed: u64,
pub maximum: u64,
pub completed_units: u64,
}
pub trait MonotonicClock: std::fmt::Debug + Send + Sync {
fn now(&self) -> Duration;
}
#[derive(Debug)]
struct SystemMonotonicClock {
origin: Instant,
}
impl SystemMonotonicClock {
fn new() -> Self {
Self {
origin: Instant::now(),
}
}
}
impl MonotonicClock for SystemMonotonicClock {
fn now(&self) -> Duration {
self.origin.elapsed()
}
}
#[derive(Debug)]
pub struct ScanJobTracker {
run_id: String,
policy: ExecutionPolicy,
clock: Arc<dyn MonotonicClock>,
started: Duration,
last_progress: Duration,
phase: JobPhase,
completed_units: u64,
}
impl ScanJobTracker {
#[must_use]
pub fn new(run_id: impl Into<String>, policy: ExecutionPolicy) -> Self {
Self::with_clock(run_id, policy, Arc::new(SystemMonotonicClock::new()))
}
#[must_use]
pub fn with_clock(
run_id: impl Into<String>,
policy: ExecutionPolicy,
clock: Arc<dyn MonotonicClock>,
) -> Self {
let now = clock.now();
Self {
run_id: run_id.into(),
policy,
clock,
started: now,
last_progress: now,
phase: JobPhase::Configuration,
completed_units: 0,
}
}
pub fn enter_phase(&mut self, phase: JobPhase) -> Result<(), ExecutionLimitExceeded> {
self.check_time()?;
self.phase = phase;
Ok(())
}
pub fn progress(&mut self, amount: u64) -> Result<(), ExecutionLimitExceeded> {
let completed_units = self
.completed_units
.checked_add(amount)
.ok_or_else(|| self.exceeded(ExecutionResource::WorkUnits, u64::MAX, u64::MAX - 1))?;
self.check_time()?;
self.completed_units = completed_units;
self.last_progress = self.clock.now();
Ok(())
}
pub fn check_time(&self) -> Result<(), ExecutionLimitExceeded> {
let now = self.clock.now();
self.check_duration(
ExecutionResource::WallTimeMs,
now.saturating_sub(self.started),
self.policy.max_scan_wall_time_ms,
)?;
self.check_duration(
ExecutionResource::NoProgressTimeMs,
now.saturating_sub(self.last_progress),
self.policy.max_no_progress_time_ms,
)
}
fn check_duration(
&self,
resource: ExecutionResource,
observed: Duration,
maximum: u64,
) -> Result<(), ExecutionLimitExceeded> {
let observed = u64::try_from(observed.as_millis()).unwrap_or(u64::MAX);
if observed > maximum {
return Err(self.exceeded(resource, observed, maximum));
}
Ok(())
}
fn exceeded(
&self,
resource: ExecutionResource,
observed: u64,
maximum: u64,
) -> ExecutionLimitExceeded {
ExecutionLimitExceeded {
run_id: self.run_id.clone(),
phase: self.phase,
resource,
observed,
maximum,
completed_units: self.completed_units,
}
}
}
#[cfg(test)]
mod tests {
use std::sync::atomic::{AtomicU64, Ordering};
use super::*;
#[derive(Debug, Default)]
struct FakeClock {
milliseconds: AtomicU64,
}
impl FakeClock {
fn advance(&self, milliseconds: u64) {
self.milliseconds.fetch_add(milliseconds, Ordering::Relaxed);
}
}
impl MonotonicClock for FakeClock {
fn now(&self) -> Duration {
Duration::from_millis(self.milliseconds.load(Ordering::Relaxed))
}
}
#[test]
fn defaults_should_be_generous_and_finite() {
let policy = ExecutionPolicy::default();
assert_eq!(policy.max_scan_wall_time_ms, 21_600_000);
assert_eq!(policy.max_worker_memory_bytes, 17_179_869_184);
assert_eq!(policy.max_checkpoint_cache_bytes, 10_737_418_240);
}
#[test]
fn partial_override_should_preserve_other_defaults() {
let policy = ExecutionPolicy::resolve(Some(&ExecutionPolicyOverrides {
max_scan_wall_time_ms: Some(28_800_000),
..ExecutionPolicyOverrides::default()
}))
.expect("valid override");
assert_eq!(policy.max_scan_wall_time_ms, 28_800_000);
assert_eq!(policy.max_no_progress_time_ms, 300_000);
}
#[test]
fn zero_should_be_rejected() {
let error = ExecutionPolicy::resolve(Some(&ExecutionPolicyOverrides {
max_worker_memory_bytes: Some(0),
..ExecutionPolicyOverrides::default()
}))
.expect_err("zero must fail");
assert!(matches!(
error,
InvalidExecutionPolicy::InvalidValue {
field: "maxWorkerMemoryBytes",
value: 0
}
));
}
#[test]
fn technically_unrepresentable_values_should_be_rejected() {
let quota = ExecutionPolicy::resolve(Some(&ExecutionPolicyOverrides {
max_checkpoint_cache_bytes: Some(u64::MAX),
..ExecutionPolicyOverrides::default()
}))
.expect_err("SQLite quota overflow must fail");
assert!(matches!(quota, InvalidExecutionPolicy::InvalidValue { .. }));
}
#[test]
fn subordinate_deadline_should_not_exceed_scan_deadline() {
let error = ExecutionPolicy::resolve(Some(&ExecutionPolicyOverrides {
max_scan_wall_time_ms: Some(1_000),
max_no_progress_time_ms: Some(1_001),
max_codegraph_sync_wall_time_ms_per_repo: Some(1_000),
graceful_termination_ms: Some(500),
..ExecutionPolicyOverrides::default()
}))
.expect_err("relationship must fail");
assert!(matches!(
error,
InvalidExecutionPolicy::InvalidRelationship {
field: "maxNoProgressTimeMs",
..
}
));
}
#[test]
fn fingerprint_should_ignore_yaml_field_order() {
let first = ExecutionPolicy::default();
let second = ExecutionPolicy::resolve(Some(&ExecutionPolicyOverrides::default()))
.expect("defaults valid");
assert_eq!(first.fingerprint(), second.fingerprint());
}
#[test]
fn injected_clock_should_accept_exact_deadline_and_reject_one_unit_over() {
let clock = Arc::new(FakeClock::default());
let policy = ExecutionPolicy {
max_scan_wall_time_ms: 10,
max_no_progress_time_ms: 10,
..ExecutionPolicy::default()
};
let tracker = ScanJobTracker::with_clock("run", policy, clock.clone());
clock.advance(10);
tracker.check_time().expect("exact deadline is inclusive");
clock.advance(1);
let error = tracker.check_time().expect_err("one over must fail");
assert_eq!(error.resource, ExecutionResource::WallTimeMs);
assert_eq!(error.observed, 11);
assert_eq!(error.maximum, 10);
}
#[test]
fn phase_changes_should_not_fake_progress() {
let clock = Arc::new(FakeClock::default());
let policy = ExecutionPolicy {
max_scan_wall_time_ms: 100,
max_no_progress_time_ms: 5,
..ExecutionPolicy::default()
};
let mut tracker = ScanJobTracker::with_clock("run", policy, clock.clone());
clock.advance(5);
tracker
.enter_phase(JobPhase::Discovery)
.expect("exact idle deadline is inclusive");
clock.advance(1);
let error = tracker
.enter_phase(JobPhase::Fingerprinting)
.expect_err("phase churn must not renew watchdog");
assert_eq!(error.resource, ExecutionResource::NoProgressTimeMs);
}
}