From dae3f97704571e53e60c0ef583f5b893123923b1 Mon Sep 17 00:00:00 2001 From: Adam Frisby Date: Thu, 24 Sep 2026 18:17:33 +0000 Subject: [PATCH 1/3] codeybox: Watchdog is blind inside Incus VMs: add guest-CPU activity signal so busy-but-silent agents aren't killed CodeyBox-WorkItem: 2a07eb4dc99247e092b95c0bee4b605c CodeyBox-Agent: antigravity/gemini-3.8-flash-high CodeyBox-Prompt-Revision: 1 Co-Authored-By: CodeyBox --- src/CodeyBox.Api/ReloadableSandboxProvider.cs | 12 +++ src/CodeyBox.Core/SandboxAbstractions.cs | 9 +- .../SandboxAdmissionControlledProvider.cs | 3 + .../WorkerProgressActivitySource.cs | 42 ++++----- .../WorkerProgressWatchdog.cs | 15 +++- .../IIncusInstanceStateReader.cs | 49 +++++++++++ .../IncusCpuActivityEvaluator.cs | 87 +++++++++++++++++++ .../IncusGuestCpuSample.cs | 19 ++++ .../IncusSandboxOptions.cs | 7 ++ .../IncusStateParser.cs | 50 +++++++++++ 10 files changed, 268 insertions(+), 25 deletions(-) create mode 100644 src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs create mode 100644 src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs create mode 100644 src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs create mode 100644 src/CodeyBox.Sandbox.Incus/IncusStateParser.cs diff --git a/src/CodeyBox.Api/ReloadableSandboxProvider.cs b/src/CodeyBox.Api/ReloadableSandboxProvider.cs index 264551b95..414b84cd8 100644 --- a/src/CodeyBox.Api/ReloadableSandboxProvider.cs +++ b/src/CodeyBox.Api/ReloadableSandboxProvider.cs @@ -207,6 +207,18 @@ public IReadOnlyList SnapshotActiveSandboxProgress() => .SelectMany(static provider => provider.Progress.SnapshotActiveSandboxProgress()) .ToArray(); + public async ValueTask> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default) + { + var activated = ActivatedProviders; + var results = new List(); + foreach (var provider in activated) + { + var progress = await provider.Progress.SnapshotActiveSandboxProgressAsync(ct).ConfigureAwait(false); + results.AddRange(progress); + } + return results; + } + public IReadOnlyList SampleDiskGuardState() => ActivatedProviders .SelectMany(static provider => provider.DiskGuard.SampleDiskGuardState()) diff --git a/src/CodeyBox.Core/SandboxAbstractions.cs b/src/CodeyBox.Core/SandboxAbstractions.cs index 73c7af22e..f9e54aa0b 100644 --- a/src/CodeyBox.Core/SandboxAbstractions.cs +++ b/src/CodeyBox.Core/SandboxAbstractions.cs @@ -818,7 +818,7 @@ public interface IActiveSandboxProvider /// signal for detached VM-local work and use to report a /// richer reason when providers can expose changing activity. /// -public sealed record ActiveSandboxProgress(WorkItemId WorkItemId, string SandboxId, string? Status = null); +public sealed record ActiveSandboxProgress(WorkItemId WorkItemId, string SandboxId, string? Status = null, double? CpuFraction = null); /// /// Optional provider capability for reporting active sandbox ownership. @@ -831,6 +831,13 @@ public interface IActiveSandboxProgressProvider /// monitoring needs. /// IReadOnlyList SnapshotActiveSandboxProgress(); + + /// + /// Asynchronously samples and snapshots currently-active sandboxes, refreshing + /// activity projections such as guest CPU metrics where supported. + /// + ValueTask> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default) + => ValueTask.FromResult(SnapshotActiveSandboxProgress()); } /// diff --git a/src/CodeyBox.Orchestrator/SandboxAdmissionControlledProvider.cs b/src/CodeyBox.Orchestrator/SandboxAdmissionControlledProvider.cs index 13d489339..1df58927a 100644 --- a/src/CodeyBox.Orchestrator/SandboxAdmissionControlledProvider.cs +++ b/src/CodeyBox.Orchestrator/SandboxAdmissionControlledProvider.cs @@ -312,6 +312,9 @@ public async Task DisposeLeakedAsync(ManagedSandboxInfo sandbox, CancellationTok public IReadOnlyList SnapshotActiveSandboxProgress() => _progressProvider?.SnapshotActiveSandboxProgress() ?? []; + public ValueTask> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default) => + _progressProvider?.SnapshotActiveSandboxProgressAsync(ct) ?? ValueTask.FromResult>([]); + public IReadOnlyList SnapshotHostPool() => _hostPoolSnapshot?.SnapshotHostPool() ?? []; diff --git a/src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs b/src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs index df1c04c88..4d9943a32 100644 --- a/src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs +++ b/src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs @@ -8,7 +8,7 @@ namespace CodeyBox.Orchestrator; /// A live worker-side signal that should count as watchdog progress even when /// the work item row and agent stream files are quiet. /// -public sealed record WorkerProgressActivity(string Reason); +public sealed record WorkerProgressActivity(string Reason, double? CpuFraction = null); /// /// Narrow activity probe settings resolved by the watchdog for the current @@ -92,7 +92,7 @@ private DefaultWorkerProgressActivitySource( _initialCpuSampleAttempts = Math.Max(1, initialCpuSampleAttempts); } - public ValueTask ObserveAsync( + public async ValueTask ObserveAsync( WorkerRegistration worker, WorkItemId itemId, WorkerProgressActivityProbe probe, @@ -100,45 +100,45 @@ private DefaultWorkerProgressActivitySource( { ct.ThrowIfCancellationRequested(); if (!string.Equals(worker.CurrentWorkItemId, itemId.ToString(), StringComparison.Ordinal)) - return ValueTask.FromResult(null); + return null; if (probe.ProcessCpuProgressSignalEnabled && TryObserveProcessCpu(itemId, out var cpuReason)) { - return ValueTask.FromResult( - new WorkerProgressActivity(cpuReason)); + return new WorkerProgressActivity(cpuReason); } if (probe.ActiveSandboxProgressSignalEnabled - && TryObserveActiveSandbox(itemId, out var sandboxReason)) + && await TryObserveActiveSandboxAsync(itemId, ct).ConfigureAwait(false) is { } sandboxActivity) { - return ValueTask.FromResult( - new WorkerProgressActivity(sandboxReason)); + return sandboxActivity; } - return ValueTask.FromResult(null); + return null; } - private bool TryObserveActiveSandbox(WorkItemId itemId, out string reason) + private async ValueTask TryObserveActiveSandboxAsync(WorkItemId itemId, CancellationToken ct) { - reason = ""; if (_activeSandboxProvider is not { } activeProvider) - return false; + return null; IReadOnlyList snapshot; try { - snapshot = activeProvider.SnapshotActiveSandboxProgress(); + snapshot = await activeProvider.SnapshotActiveSandboxProgressAsync(ct).ConfigureAwait(false); } catch { - return false; + return null; } - var signatureParts = snapshot + var matchingEntries = snapshot .Where(entry => entry.WorkItemId == itemId && !string.IsNullOrWhiteSpace(entry.SandboxId)) + .ToArray(); + + var signatureParts = matchingEntries .Select(entry => $"{entry.SandboxId}\0{entry.Status ?? ""}") .Order(StringComparer.Ordinal) .ToArray(); @@ -146,23 +146,23 @@ private bool TryObserveActiveSandbox(WorkItemId itemId, out string reason) if (signatureParts.Length == 0) { _activeSandboxSignatures.TryRemove(itemId, out _); - return false; + return null; } var signature = string.Join("\0\0", signatureParts); + var cpuFraction = matchingEntries.Select(e => e.CpuFraction).FirstOrDefault(f => f.HasValue); + if (!_activeSandboxSignatures.TryGetValue(itemId, out var previous)) { _activeSandboxSignatures[itemId] = signature; - reason = "active-sandbox"; - return true; + return new WorkerProgressActivity("active-sandbox", cpuFraction); } if (string.Equals(signature, previous, StringComparison.Ordinal)) - return false; + return null; _activeSandboxSignatures[itemId] = signature; - reason = "active-sandbox-change"; - return true; + return new WorkerProgressActivity("active-sandbox-change", cpuFraction); } private bool TryObserveProcessCpu(WorkItemId itemId, out string reason) diff --git a/src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs b/src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs index 3c69b4ad9..7d68b7566 100644 --- a/src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs +++ b/src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs @@ -194,9 +194,18 @@ public async Task RunOnceAsync(CancellationToken ct) if (workerActivity is not null) { _workerActivityProgress[activityKey] = new WorkerActivityProgress(now, workerActivity.Reason); - _log.LogDebug( - "Watchdog: worker {WorkerId} for item {ItemId} has live activity signal {Reason}; treating as progress", - worker.WorkerId, itemId, workerActivity.Reason); + if (workerActivity.CpuFraction is { } cpuFraction) + { + _log.LogDebug( + "Watchdog: worker {WorkerId} for item {ItemId} has live activity signal {Reason} (cpu={CpuFraction:P1}); treating as progress", + worker.WorkerId, itemId, workerActivity.Reason, cpuFraction); + } + else + { + _log.LogDebug( + "Watchdog: worker {WorkerId} for item {ItemId} has live activity signal {Reason}; treating as progress", + worker.WorkerId, itemId, workerActivity.Reason); + } continue; } diff --git a/src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs b/src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs new file mode 100644 index 000000000..717bf09b7 --- /dev/null +++ b/src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs @@ -0,0 +1,49 @@ +namespace CodeyBox.Sandbox.Incus; + +/// +/// Seam for reading instance guest state (CPU usage) from Incus. +/// +internal interface IIncusInstanceStateReader +{ + Task ReadStateAsync( + IncusSandboxOptions options, + string instanceName, + CancellationToken ct); +} + +/// +/// Default implementation that queries Incus instance state via the CLI/API runner. +/// +internal sealed class DefaultIncusInstanceStateReader( + IncusCliRunner cli, + TimeProvider timeProvider) : IIncusInstanceStateReader +{ + public async Task ReadStateAsync( + IncusSandboxOptions options, + string instanceName, + CancellationToken ct) + { + ArgumentNullException.ThrowIfNull(options); + ArgumentException.ThrowIfNullOrWhiteSpace(instanceName); + + var result = await cli.RunAllowFailureAsync( + options, + [options.BinaryPath, "query", $"/1.0/instances/{instanceName}/state?project={options.ProjectName}"], + stdin: null, + timeout: options.ActivityQueryTimeout, + ct: ct, + heavyOperation: false, + maxStdoutBytes: 64 * 1024, + maxStderrBytes: 4096).ConfigureAwait(false); + + if (result.ExitCode != 0 || string.IsNullOrWhiteSpace(result.Stdout)) + return null; + + if (IncusStateParser.TryParseGuestCpuUsage(result.Stdout, out var cpuUsageNs)) + { + return new IncusGuestCpuSample(cpuUsageNs, timeProvider.GetUtcNow()); + } + + return null; + } +} diff --git a/src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs b/src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs new file mode 100644 index 000000000..07afeae3e --- /dev/null +++ b/src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs @@ -0,0 +1,87 @@ +namespace CodeyBox.Sandbox.Incus; + +/// +/// Pure evaluation functions determining if an Incus guest consumed sufficient +/// CPU over an interval to count as active progress. +/// +public static class IncusCpuActivityEvaluator +{ + public static bool IsActive( + IncusGuestCpuSample? previousSample, + IncusGuestCpuSample currentSample, + TimeSpan elapsed, + double thresholdPercent) => + Evaluate(previousSample, currentSample, elapsed, thresholdPercent).IsActive; + + public static bool IsActive( + long? previousCpuNanoseconds, + long currentCpuNanoseconds, + TimeSpan elapsed, + double thresholdPercent) => + Evaluate( + previousCpuNanoseconds is { } prev ? new IncusGuestCpuSample(prev) : null, + new IncusGuestCpuSample(currentCpuNanoseconds), + elapsed, + thresholdPercent).IsActive; + + public static bool IsActiveFraction( + IncusGuestCpuSample? previousSample, + IncusGuestCpuSample currentSample, + TimeSpan elapsed, + double thresholdFraction) => + EvaluateFraction(previousSample, currentSample, elapsed, thresholdFraction).IsActive; + + public static bool IsActiveFraction( + long? previousCpuNanoseconds, + long currentCpuNanoseconds, + TimeSpan elapsed, + double thresholdFraction) => + EvaluateFraction( + previousCpuNanoseconds is { } prev ? new IncusGuestCpuSample(prev) : null, + new IncusGuestCpuSample(currentCpuNanoseconds), + elapsed, + thresholdFraction).IsActive; + + public static IncusCpuActivityEvaluation Evaluate( + IncusGuestCpuSample? previousSample, + IncusGuestCpuSample currentSample, + TimeSpan elapsed, + double thresholdPercent) + { + if (previousSample is null) + return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); + + if (elapsed <= TimeSpan.Zero) + return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); + + if (currentSample.CpuUsageNanoseconds < previousSample.Value.CpuUsageNanoseconds) + return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); + + if (currentSample.CpuUsageNanoseconds < 0 || previousSample.Value.CpuUsageNanoseconds < 0) + return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); + + long deltaNs = currentSample.CpuUsageNanoseconds - previousSample.Value.CpuUsageNanoseconds; + if (deltaNs == 0) + return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); + + if (!double.IsFinite(thresholdPercent) || thresholdPercent <= 0) + return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); + + double elapsedSeconds = elapsed.TotalSeconds; + if (elapsedSeconds <= 0) + return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); + + double cpuFraction = (double)deltaNs / (elapsedSeconds * 1_000_000_000.0); + double cpuPercent = cpuFraction * 100.0; + + bool isActive = cpuPercent >= thresholdPercent; + return new IncusCpuActivityEvaluation(isActive, cpuFraction, cpuPercent); + } + + public static IncusCpuActivityEvaluation EvaluateFraction( + IncusGuestCpuSample? previousSample, + IncusGuestCpuSample currentSample, + TimeSpan elapsed, + double thresholdFraction) => + Evaluate(previousSample, currentSample, elapsed, thresholdFraction * 100.0); +} diff --git a/src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs b/src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs new file mode 100644 index 000000000..70ce92715 --- /dev/null +++ b/src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs @@ -0,0 +1,19 @@ +namespace CodeyBox.Sandbox.Incus; + +/// +/// A point-in-time sample of an Incus guest's cumulative CPU nanoseconds. +/// +public readonly record struct IncusGuestCpuSample( + long CpuUsageNanoseconds, + DateTimeOffset Timestamp = default); + +/// +/// Evaluation result of guest CPU activity over an elapsed interval. +/// +public readonly record struct IncusCpuActivityEvaluation( + bool IsActive, + double CpuFraction, + double CpuPercent) +{ + public static implicit operator bool(IncusCpuActivityEvaluation evaluation) => evaluation.IsActive; +} diff --git a/src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs b/src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs index 52ff2b7d9..1cf33823e 100644 --- a/src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs +++ b/src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs @@ -230,6 +230,9 @@ public sealed record IncusSandboxOptions public bool CaptureResourceMetrics { get; init; } public TimeSpan ResourceMetricsCaptureTimeout { get; init; } = TimeSpan.FromSeconds(5); public TimeSpan ResourceMetricsSampleInterval { get; init; } = TimeSpan.FromSeconds(10); + public double ActivityCpuThresholdPercent { get; init; } = 5.0; + public TimeSpan ActivitySampleInterval { get; init; } = TimeSpan.FromSeconds(5); + public TimeSpan ActivityQueryTimeout { get; init; } = TimeSpan.FromSeconds(5); public IncusDiskGuardOptions? DiskGuard { get; init; } = new(); /// @@ -353,6 +356,10 @@ public static IReadOnlyList Validate(IncusSandboxOptions options) } RequirePositiveDuration(options.ResourceMetricsCaptureTimeout, nameof(ResourceMetricsCaptureTimeout), errors); RequirePositiveDuration(options.ResourceMetricsSampleInterval, nameof(ResourceMetricsSampleInterval), errors); + if (!double.IsFinite(options.ActivityCpuThresholdPercent) || options.ActivityCpuThresholdPercent <= 0 || options.ActivityCpuThresholdPercent > 10000) + errors.Add($"{nameof(ActivityCpuThresholdPercent)} must be greater than zero and at most 10000."); + RequirePositiveDuration(options.ActivitySampleInterval, nameof(ActivitySampleInterval), errors); + RequirePositiveDuration(options.ActivityQueryTimeout, nameof(ActivityQueryTimeout), errors); if (options.MaxConcurrentOperations is < 1 or > 64) errors.Add($"{nameof(MaxConcurrentOperations)} must be between 1 and 64."); if (options.MaxConcurrentBoots is < 1 or > 64) diff --git a/src/CodeyBox.Sandbox.Incus/IncusStateParser.cs b/src/CodeyBox.Sandbox.Incus/IncusStateParser.cs new file mode 100644 index 000000000..cfbb6c359 --- /dev/null +++ b/src/CodeyBox.Sandbox.Incus/IncusStateParser.cs @@ -0,0 +1,50 @@ +using System.Text.Json; + +namespace CodeyBox.Sandbox.Incus; + +/// +/// Safe JSON parser for Incus instance state responses. +/// +internal static class IncusStateParser +{ + public static bool TryParseGuestCpuUsage(string? json, out long cpuUsageNanoseconds) + { + cpuUsageNanoseconds = 0; + if (string.IsNullOrWhiteSpace(json)) + return false; + + try + { + using var doc = JsonDocument.Parse(json); + var root = doc.RootElement; + var target = root.TryGetProperty("metadata", out var meta) && meta.ValueKind == JsonValueKind.Object + ? meta + : root; + + if (target.ValueKind != JsonValueKind.Object) + return false; + + if (target.TryGetProperty("status", out var statusProp) + && statusProp.ValueKind == JsonValueKind.String + && !string.Equals(statusProp.GetString(), "Running", StringComparison.OrdinalIgnoreCase)) + { + return false; + } + + if (target.TryGetProperty("cpu", out var cpuProp) && cpuProp.ValueKind == JsonValueKind.Object) + { + if (cpuProp.TryGetProperty("usage", out var usageProp) && usageProp.TryGetInt64(out var ns) && ns >= 0) + { + cpuUsageNanoseconds = ns; + return true; + } + } + + return false; + } + catch (JsonException) + { + return false; + } + } +} From ea6d03934febe5b942bcd9f012da0967f4e3ee1b Mon Sep 17 00:00:00 2001 From: Adam Frisby Date: Thu, 24 Sep 2026 19:57:00 +0000 Subject: [PATCH 2/3] codeybox: preempt checkpoint Watchdog is blind inside Incus VMs: add guest-CPU activity signal so busy-but-silent agents aren't killed CodeyBox-WorkItem: 2a07eb4dc99247e092b95c0bee4b605c CodeyBox-Agent: antigravity/gemini-3.8-flash-high Co-Authored-By: CodeyBox From fb08ff9319d5b4fe4628120f520a961b07e44225 Mon Sep 17 00:00:00 2001 From: Adam Frisby Date: Thu, 24 Sep 2026 20:24:26 +0000 Subject: [PATCH 3/3] codeybox: wire the Incus guest-CPU activity signal into the provider and add the required tests Audit rework: the previous commit added the sampling machinery but never connected it, so the watchdog signal could never fire. - IncusSandboxProvider now overrides SnapshotActiveSandboxProgressAsync: per-owned-sandbox cpu.usage sampling via an injected IIncusInstanceStateReader, a per-sandbox cache gated by ActivitySampleInterval, serialized refreshes, bounded fan-out, and an epoch-countered status that changes only while the guest exceeds ActivityCpuThresholdPercent. Sync SnapshotActiveSandboxProgress reports the same last-known projection. - DefaultIncusInstanceStateReader now validates the instance name and options identity before interpolating them into the incus query path. - IncusCpuActivityEvaluator reduced to Evaluate + IsActive; dropped the lossy implicit bool, the redundant fraction/primitive overloads, and the unused CpuPercent; sample Timestamp is now consumed for elapsed. - WorkerProgressActivitySource rethrows caller cancellation instead of downgrading it to no-signal; watchdog debug log collapsed to one template carrying the measured CPU fraction; CpuFraction units documented on both contracts; ReloadableSandboxProvider fans provider snapshots out in parallel. - Tests: pure-function coverage of the evaluator (first sample, counter reset/regression, threshold boundary, bad elapsed/threshold), the state parser, and the reader's input guards; provider tests with a fake state reader (busy guest changes the signature, idle/failed/within-interval keep it stable); FakeTimeProvider watchdog integration proving a silent busy Incus guest is not recovered past ProgressTimeout while a silent idle guest is. CodeyBox-WorkItem: 2a07eb4dc99247e092b95c0bee4b605c CodeyBox-Agent: antigravity/gemini-3.8-flash-high CodeyBox-Prompt-Revision: 1 Co-Authored-By: CodeyBox --- src/CodeyBox.Api/ReloadableSandboxProvider.cs | 15 +- src/CodeyBox.Core/SandboxAbstractions.cs | 20 +- .../WorkerProgressActivitySource.cs | 11 + .../WorkerProgressWatchdog.cs | 19 +- .../IIncusInstanceStateReader.cs | 8 +- .../IncusCpuActivityEvaluator.cs | 108 ++--- .../IncusGuestCpuSample.cs | 18 +- .../IncusSandboxOptions.cs | 14 + .../IncusSandboxProvider.cs | 164 ++++++- .../IncusGuestCpuActivityTests.cs | 309 ++++++++++++ .../IncusSandboxLifecycleTests.cs | 442 ++++++++++++++++++ 11 files changed, 1025 insertions(+), 103 deletions(-) create mode 100644 tests/CodeyBox.Tests/IncusGuestCpuActivityTests.cs diff --git a/src/CodeyBox.Api/ReloadableSandboxProvider.cs b/src/CodeyBox.Api/ReloadableSandboxProvider.cs index 414b84cd8..5b0d7f2a1 100644 --- a/src/CodeyBox.Api/ReloadableSandboxProvider.cs +++ b/src/CodeyBox.Api/ReloadableSandboxProvider.cs @@ -209,14 +209,13 @@ public IReadOnlyList SnapshotActiveSandboxProgress() => public async ValueTask> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default) { - var activated = ActivatedProviders; - var results = new List(); - foreach (var provider in activated) - { - var progress = await provider.Progress.SnapshotActiveSandboxProgressAsync(ct).ConfigureAwait(false); - results.AddRange(progress); - } - return results; + // Fan out in parallel: providers that sample live signals (guest CPU + // queries) would otherwise serialize their per-call timeouts. + var snapshots = ActivatedProviders + .Select(provider => provider.Progress.SnapshotActiveSandboxProgressAsync(ct).AsTask()) + .ToArray(); + var results = await Task.WhenAll(snapshots).ConfigureAwait(false); + return results.SelectMany(static snapshot => snapshot).ToArray(); } public IReadOnlyList SampleDiskGuardState() => diff --git a/src/CodeyBox.Core/SandboxAbstractions.cs b/src/CodeyBox.Core/SandboxAbstractions.cs index f9e54aa0b..fd402ebd6 100644 --- a/src/CodeyBox.Core/SandboxAbstractions.cs +++ b/src/CodeyBox.Core/SandboxAbstractions.cs @@ -818,6 +818,14 @@ public interface IActiveSandboxProvider /// signal for detached VM-local work and use to report a /// richer reason when providers can expose changing activity. /// +/// Work item that owns the sandbox. +/// Provider-side sandbox identifier. +/// Provider-defined activity marker. Watchdog signature +/// tracking treats a changed value as progress, so a provider that emits a +/// live signal must keep it stable while the sandbox is idle. +/// Latest measured CPU activity as a fraction of one +/// core averaged over the provider's sample interval (1.0 = one fully busy +/// core); null when the provider does not measure CPU. public sealed record ActiveSandboxProgress(WorkItemId WorkItemId, string SandboxId, string? Status = null, double? CpuFraction = null); /// @@ -827,14 +835,20 @@ public sealed record ActiveSandboxProgress(WorkItemId WorkItemId, string Sandbox public interface IActiveSandboxProgressProvider { /// - /// Snapshot of currently-active sandboxes, projected to the fields progress - /// monitoring needs. + /// Last-known snapshot of currently-active sandboxes, projected to the + /// fields progress monitoring needs. This member never refreshes activity + /// projections: providers that sample live signals (e.g. guest CPU) return + /// the cached values produced by the most recent + /// , so both members report + /// the same projection. /// IReadOnlyList SnapshotActiveSandboxProgress(); /// /// Asynchronously samples and snapshots currently-active sandboxes, refreshing - /// activity projections such as guest CPU metrics where supported. + /// activity projections such as guest CPU metrics where supported. The + /// default delegates to the non-refreshing + /// . /// ValueTask> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default) => ValueTask.FromResult(SnapshotActiveSandboxProgress()); diff --git a/src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs b/src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs index 4d9943a32..47eee88ab 100644 --- a/src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs +++ b/src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs @@ -8,6 +8,11 @@ namespace CodeyBox.Orchestrator; /// A live worker-side signal that should count as watchdog progress even when /// the work item row and agent stream files are quiet. /// +/// Short machine-readable label naming the signal that +/// fired (e.g. process-cpu, active-sandbox-change). +/// Measured CPU activity as a fraction of one core +/// averaged over the provider's sample interval (1.0 = one fully busy core); +/// null when the signal carries no CPU measurement. public sealed record WorkerProgressActivity(string Reason, double? CpuFraction = null); /// @@ -127,6 +132,12 @@ private DefaultWorkerProgressActivitySource( { snapshot = await activeProvider.SnapshotActiveSandboxProgressAsync(ct).ConfigureAwait(false); } + catch (OperationCanceledException) when (ct.IsCancellationRequested) + { + // Caller-driven cancellation aborts the sweep; it is not the + // best-effort "no signal" the blanket catch below encodes. + throw; + } catch { return null; diff --git a/src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs b/src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs index 7d68b7566..52397f94d 100644 --- a/src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs +++ b/src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs @@ -1,4 +1,5 @@ using System.Collections.Concurrent; +using System.Globalization; using Microsoft.Extensions.Hosting; using Microsoft.Extensions.Logging; using CodeyBox.Core; @@ -194,18 +195,12 @@ public async Task RunOnceAsync(CancellationToken ct) if (workerActivity is not null) { _workerActivityProgress[activityKey] = new WorkerActivityProgress(now, workerActivity.Reason); - if (workerActivity.CpuFraction is { } cpuFraction) - { - _log.LogDebug( - "Watchdog: worker {WorkerId} for item {ItemId} has live activity signal {Reason} (cpu={CpuFraction:P1}); treating as progress", - worker.WorkerId, itemId, workerActivity.Reason, cpuFraction); - } - else - { - _log.LogDebug( - "Watchdog: worker {WorkerId} for item {ItemId} has live activity signal {Reason}; treating as progress", - worker.WorkerId, itemId, workerActivity.Reason); - } + var cpuText = workerActivity.CpuFraction is { } cpuFraction + ? cpuFraction.ToString("P1", CultureInfo.InvariantCulture) + : "n/a"; + _log.LogDebug( + "Watchdog: worker {WorkerId} for item {ItemId} has live activity signal {Reason} (cpu={CpuFraction}); treating as progress", + worker.WorkerId, itemId, workerActivity.Reason, cpuText); continue; } diff --git a/src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs b/src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs index 717bf09b7..6e09b230f 100644 --- a/src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs +++ b/src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs @@ -23,8 +23,12 @@ internal sealed class DefaultIncusInstanceStateReader( string instanceName, CancellationToken ct) { - ArgumentNullException.ThrowIfNull(options); - ArgumentException.ThrowIfNullOrWhiteSpace(instanceName); + // The instance name and project are interpolated into the Incus REST + // path below, so they must carry the same identifier guard every other + // Incus call site applies — an unvalidated name could smuggle '/', + // '?', or '%' escapes and redirect the query at a different endpoint. + IncusInputValidation.ValidateOptionsIdentity(options); + IncusInputValidation.ValidateInstanceName(instanceName, nameof(instanceName)); var result = await cli.RunAllowFailureAsync( options, diff --git a/src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs b/src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs index 07afeae3e..86535f644 100644 --- a/src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs +++ b/src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs @@ -1,87 +1,57 @@ namespace CodeyBox.Sandbox.Incus; /// -/// Pure evaluation functions determining if an Incus guest consumed sufficient -/// CPU over an interval to count as active progress. +/// Pure evaluation of whether an Incus guest consumed enough CPU over an +/// interval to count as watchdog progress. /// public static class IncusCpuActivityEvaluator { - public static bool IsActive( - IncusGuestCpuSample? previousSample, - IncusGuestCpuSample currentSample, - TimeSpan elapsed, - double thresholdPercent) => - Evaluate(previousSample, currentSample, elapsed, thresholdPercent).IsActive; - - public static bool IsActive( - long? previousCpuNanoseconds, - long currentCpuNanoseconds, - TimeSpan elapsed, - double thresholdPercent) => - Evaluate( - previousCpuNanoseconds is { } prev ? new IncusGuestCpuSample(prev) : null, - new IncusGuestCpuSample(currentCpuNanoseconds), - elapsed, - thresholdPercent).IsActive; - - public static bool IsActiveFraction( - IncusGuestCpuSample? previousSample, - IncusGuestCpuSample currentSample, - TimeSpan elapsed, - double thresholdFraction) => - EvaluateFraction(previousSample, currentSample, elapsed, thresholdFraction).IsActive; - - public static bool IsActiveFraction( - long? previousCpuNanoseconds, - long currentCpuNanoseconds, - TimeSpan elapsed, - double thresholdFraction) => - EvaluateFraction( - previousCpuNanoseconds is { } prev ? new IncusGuestCpuSample(prev) : null, - new IncusGuestCpuSample(currentCpuNanoseconds), - elapsed, - thresholdFraction).IsActive; - + /// + /// Compares the cumulative guest-CPU counters of two consecutive samples. + /// The guest counts as active only when the counter delta, expressed as a + /// percentage of one CPU core over , reaches + /// . The first sample (no previous), a + /// counter reset or regression (current below previous), negative counters, + /// a non-positive elapsed interval, and a non-positive or non-finite + /// threshold all evaluate to inactive. + /// + /// The earlier sample, or null when this is + /// the first observation (which always establishes a baseline and reports + /// inactive). + /// The latest sample. + /// Wall time between the two samples. + /// Activity threshold as a percentage of one + /// CPU core (5 = 5% of a core, i.e. 0.05 CPU fraction). public static IncusCpuActivityEvaluation Evaluate( IncusGuestCpuSample? previousSample, IncusGuestCpuSample currentSample, TimeSpan elapsed, double thresholdPercent) { - if (previousSample is null) - return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); - - if (elapsed <= TimeSpan.Zero) - return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); - - if (currentSample.CpuUsageNanoseconds < previousSample.Value.CpuUsageNanoseconds) - return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); - - if (currentSample.CpuUsageNanoseconds < 0 || previousSample.Value.CpuUsageNanoseconds < 0) - return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); - - long deltaNs = currentSample.CpuUsageNanoseconds - previousSample.Value.CpuUsageNanoseconds; - if (deltaNs == 0) - return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); - - if (!double.IsFinite(thresholdPercent) || thresholdPercent <= 0) - return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); - - double elapsedSeconds = elapsed.TotalSeconds; - if (elapsedSeconds <= 0) - return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0, CpuPercent: 0.0); - - double cpuFraction = (double)deltaNs / (elapsedSeconds * 1_000_000_000.0); - double cpuPercent = cpuFraction * 100.0; - - bool isActive = cpuPercent >= thresholdPercent; - return new IncusCpuActivityEvaluation(isActive, cpuFraction, cpuPercent); + if (previousSample is not { } previous + || elapsed <= TimeSpan.Zero + || currentSample.CpuUsageNanoseconds < 0 + || previous.CpuUsageNanoseconds < 0 + || currentSample.CpuUsageNanoseconds < previous.CpuUsageNanoseconds + || !double.IsFinite(thresholdPercent) + || thresholdPercent <= 0) + { + return new IncusCpuActivityEvaluation(IsActive: false, CpuFraction: 0.0); + } + + var deltaNs = currentSample.CpuUsageNanoseconds - previous.CpuUsageNanoseconds; + var cpuFraction = deltaNs / (elapsed.TotalSeconds * 1_000_000_000.0); + return new IncusCpuActivityEvaluation(cpuFraction * 100.0 >= thresholdPercent, cpuFraction); } - public static IncusCpuActivityEvaluation EvaluateFraction( + /// + /// Convenience wrapper over reporting only the + /// active/inactive decision. + /// + public static bool IsActive( IncusGuestCpuSample? previousSample, IncusGuestCpuSample currentSample, TimeSpan elapsed, - double thresholdFraction) => - Evaluate(previousSample, currentSample, elapsed, thresholdFraction * 100.0); + double thresholdPercent) => + Evaluate(previousSample, currentSample, elapsed, thresholdPercent).IsActive; } diff --git a/src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs b/src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs index 70ce92715..364d2ed44 100644 --- a/src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs +++ b/src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs @@ -1,8 +1,14 @@ namespace CodeyBox.Sandbox.Incus; /// -/// A point-in-time sample of an Incus guest's cumulative CPU nanoseconds. +/// A point-in-time sample of an Incus guest's cumulative CPU consumption — the +/// monotonic cpu.usage nanoseconds counter reported by +/// incus query /1.0/instances/<name>/state. /// +/// Cumulative guest CPU time in nanoseconds +/// since the instance started; resets when the instance restarts. +/// When the sample was read. The sampling interval is +/// derived from consecutive sample timestamps. public readonly record struct IncusGuestCpuSample( long CpuUsageNanoseconds, DateTimeOffset Timestamp = default); @@ -10,10 +16,10 @@ public readonly record struct IncusGuestCpuSample( /// /// Evaluation result of guest CPU activity over an elapsed interval. /// +/// Whether the guest met the configured activity +/// threshold. +/// Measured guest CPU as a fraction of one core +/// averaged over the sample interval (1.0 = one fully busy core). public readonly record struct IncusCpuActivityEvaluation( bool IsActive, - double CpuFraction, - double CpuPercent) -{ - public static implicit operator bool(IncusCpuActivityEvaluation evaluation) => evaluation.IsActive; -} + double CpuFraction); diff --git a/src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs b/src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs index 1cf33823e..0e3a24e95 100644 --- a/src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs +++ b/src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs @@ -230,8 +230,22 @@ public sealed record IncusSandboxOptions public bool CaptureResourceMetrics { get; init; } public TimeSpan ResourceMetricsCaptureTimeout { get; init; } = TimeSpan.FromSeconds(5); public TimeSpan ResourceMetricsSampleInterval { get; init; } = TimeSpan.FromSeconds(10); + /// + /// Guest-CPU activity threshold for the watchdog's active-sandbox signal, + /// as a percentage of one CPU core averaged over + /// (5 = 5% of one core). The emitted + /// progress signature changes only while the guest meets this threshold. + /// public double ActivityCpuThresholdPercent { get; init; } = 5.0; + /// + /// Minimum wall-clock interval between guest-CPU state queries for a given + /// sandbox. Snapshots requested sooner reuse the last sampled projection. + /// public TimeSpan ActivitySampleInterval { get; init; } = TimeSpan.FromSeconds(5); + /// + /// Per-call deadline for each incus query /1.0/instances/<name>/state + /// read. A timed-out query contributes no activity signal. + /// public TimeSpan ActivityQueryTimeout { get; init; } = TimeSpan.FromSeconds(5); public IncusDiskGuardOptions? DiskGuard { get; init; } = new(); diff --git a/src/CodeyBox.Sandbox.Incus/IncusSandboxProvider.cs b/src/CodeyBox.Sandbox.Incus/IncusSandboxProvider.cs index ab0dfb284..56c3b4069 100644 --- a/src/CodeyBox.Sandbox.Incus/IncusSandboxProvider.cs +++ b/src/CodeyBox.Sandbox.Incus/IncusSandboxProvider.cs @@ -92,6 +92,7 @@ public sealed class IncusSandboxProvider : private readonly Func _optionsAccessor; private readonly ILogger _log; private readonly IncusCliRunner _cli; + private readonly IIncusInstanceStateReader _stateReader; private readonly ITimingStore? _timings; private readonly ISandboxResourceUsageStore? _resourceUsageStore; private readonly IDiskSpaceProbe _diskProbe; @@ -102,6 +103,10 @@ public sealed class IncusSandboxProvider : private readonly string _lifecycleStagingRootPath; private readonly SemaphoreSlim _hostPreflightLock = new(1, 1); private readonly SemaphoreSlim _hostProvisioningInputGate = new(1, 1); + // Fan-out cap for the lightweight per-sandbox `incus query .../state` + // reads behind SnapshotActiveSandboxProgressAsync — deliberately outside + // the heavy-operation gate so watchdog probes never starve lifecycle ops. + private const int MaxConcurrentGuestStateQueries = 4; private readonly ConcurrentDictionary _baselineLocks = new(StringComparer.Ordinal); private readonly IncusSharedPackageArchiveCache _sharedPackageArchives = new(); // Boot gate: staggers concurrent VM boots (incus start + guest-agent wait) @@ -112,6 +117,8 @@ public sealed class IncusSandboxProvider : private int _bootGateCapacity; private readonly ConcurrentDictionary _activeNames = new(StringComparer.Ordinal); private readonly ConcurrentDictionary _activeOwners = new(StringComparer.Ordinal); + private readonly ConcurrentDictionary _guestActivity = new(StringComparer.Ordinal); + private readonly SemaphoreSlim _guestActivityRefreshLock = new(1, 1); private readonly ConcurrentDictionary _uncertainBaselines = new(StringComparer.Ordinal); private long _lastPoolFreeBytes = -1; private string? _lastPoolName; @@ -145,7 +152,8 @@ internal IncusSandboxProvider( IDiskSpaceProbe? diskProbe = null, TimeProvider? timeProvider = null, Func? newGuid = null, - Func? environmentVariableReader = null) + Func? environmentVariableReader = null, + IIncusInstanceStateReader? stateReader = null) { _optionsAccessor = optionsAccessor ?? throw new ArgumentNullException(nameof(optionsAccessor)); _log = log ?? throw new ArgumentNullException(nameof(log)); @@ -156,6 +164,7 @@ internal IncusSandboxProvider( _newGuid = newGuid ?? Guid.NewGuid; _environmentVariableReader = environmentVariableReader ?? Environment.GetEnvironmentVariable; _cli = new IncusCliRunner(runner, _timeProvider); + _stateReader = stateReader ?? new DefaultIncusInstanceStateReader(_cli, _timeProvider); var initialOptions = ReadValidatedOptions(); _lifecycleProjectName = initialOptions.ProjectName; _lifecycleStagingRootPath = ResolveStagingRootPath(initialOptions); @@ -1046,10 +1055,158 @@ await DeleteVerifiedOwnedInstanceAsync( .Select(owner => (owner.WorkItemId, (IShutdownTeardownSandbox)owner.Sandbox)) .ToArray(); + /// + /// Last-known projection without refreshing guest state — the same values + /// the most recent + /// produced, or the never-active baseline before the first refresh. + /// public IReadOnlyList SnapshotActiveSandboxProgress() => - _activeOwners.Values - .Select(owner => new ActiveSandboxProgress(owner.WorkItemId, owner.Sandbox.Id, "incus-running")) + _activeOwners + .Select(pair => ProjectActiveSandboxProgress(pair.Key, pair.Value.WorkItemId)) + .ToArray(); + + /// + /// Refreshes the per-sandbox guest-CPU projection before snapshotting so + /// the watchdog sees a signature that changes while the guest is busy and + /// stays stable while it is idle or unreachable. Each owned instance is + /// queried at most once per + /// ; a failed or + /// timed-out query contributes no signal and never throws into the caller. + /// + public async ValueTask> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default) + { + var options = ReadOptions(); + var owners = _activeOwners + .Select(static pair => (pair.Value.WorkItemId, InstanceName: pair.Key)) .ToArray(); + await RefreshGuestActivityAsync(options, owners, ct).ConfigureAwait(false); + return SnapshotActiveSandboxProgress(); + } + + private ActiveSandboxProgress ProjectActiveSandboxProgress(string instanceName, WorkItemId workItemId) + { + if (_guestActivity.TryGetValue(instanceName, out var state)) + { + lock (state) + return new ActiveSandboxProgress(workItemId, instanceName, state.Status, state.CpuFraction); + } + return new ActiveSandboxProgress(workItemId, instanceName, FormatGuestActivityStatus(0), CpuFraction: null); + } + + private async Task RefreshGuestActivityAsync( + IncusSandboxOptions options, + IReadOnlyList<(WorkItemId WorkItemId, string InstanceName)> owners, + CancellationToken ct) + { + // Serialize refreshes: two concurrent snapshots must not double-sample + // or double-count the same guest-CPU interval. + await _guestActivityRefreshLock.WaitAsync(ct).ConfigureAwait(false); + try + { + var live = new HashSet(StringComparer.Ordinal); + foreach (var owner in owners) + live.Add(owner.InstanceName); + foreach (var cached in _guestActivity.Keys) + { + if (!live.Contains(cached)) + _guestActivity.TryRemove(cached, out _); + } + if (owners.Count == 0) + return; + + var now = _timeProvider.GetUtcNow(); + var due = new List(owners.Count); + foreach (var owner in owners) + { + var state = _guestActivity.GetOrAdd(owner.InstanceName, static _ => new GuestActivityState()); + lock (state) + { + if (now - state.LastQueryAt >= options.ActivitySampleInterval) + due.Add(owner.InstanceName); + } + } + if (due.Count == 0) + return; + + using var gate = new SemaphoreSlim(MaxConcurrentGuestStateQueries); + await Task.WhenAll(due.Select( + instanceName => SampleGuestActivityAsync(options, instanceName, gate, ct))).ConfigureAwait(false); + } + finally + { + _guestActivityRefreshLock.Release(); + } + } + + private async Task SampleGuestActivityAsync( + IncusSandboxOptions options, + string instanceName, + SemaphoreSlim gate, + CancellationToken ct) + { + await gate.WaitAsync(ct).ConfigureAwait(false); + try + { + IncusGuestCpuSample? sample; + try + { + sample = await _stateReader.ReadStateAsync(options, instanceName, ct).ConfigureAwait(false); + } + catch (OperationCanceledException) when (ct.IsCancellationRequested) + { + throw; + } + catch (Exception ex) + { + // Best-effort probe: a failed query yields no fresh sample and + // must never surface to the watchdog. + _log.LogDebug(ex, "Incus guest CPU query failed for {InstanceName}; no activity signal", instanceName); + sample = null; + } + + var queryAt = _timeProvider.GetUtcNow(); + if (!_guestActivity.TryGetValue(instanceName, out var state)) + return; + + lock (state) + { + // Rate-limit failed queries too, not just successful reads. + state.LastQueryAt = queryAt; + if (sample is not { } current) + return; + + var previous = state.LastSample; + var elapsed = previous is { } prev ? current.Timestamp - prev.Timestamp : TimeSpan.Zero; + var evaluation = IncusCpuActivityEvaluator.Evaluate( + previous, current, elapsed, options.ActivityCpuThresholdPercent); + state.LastSample = current; + state.CpuFraction = previous is null ? null : evaluation.CpuFraction; + if (evaluation.IsActive) + state.ActiveEpochs++; + } + } + finally + { + gate.Release(); + } + } + + private static string FormatGuestActivityStatus(long activeEpochs) => + string.Create(CultureInfo.InvariantCulture, $"incus-cpu-{activeEpochs}"); + + /// + /// Per-sandbox guest-CPU sampling state. counts + /// the number of intervals that met the activity threshold, so the emitted + /// changes only while the guest is measurably busy. + /// + private sealed class GuestActivityState + { + internal DateTimeOffset LastQueryAt; + internal IncusGuestCpuSample? LastSample; + internal long ActiveEpochs; + internal double? CpuFraction; + internal string Status => FormatGuestActivityStatus(ActiveEpochs); + } private async Task ResolveOrEnsureBaselineAsync( IncusSandboxOptions options, @@ -3023,6 +3180,7 @@ private void MarkInactive(string name) { _activeNames.TryRemove(name, out _); _activeOwners.TryRemove(name, out _); + _guestActivity.TryRemove(name, out _); } private static bool ProjectListContains(string json, string projectName) diff --git a/tests/CodeyBox.Tests/IncusGuestCpuActivityTests.cs b/tests/CodeyBox.Tests/IncusGuestCpuActivityTests.cs new file mode 100644 index 000000000..3ee019a9e --- /dev/null +++ b/tests/CodeyBox.Tests/IncusGuestCpuActivityTests.cs @@ -0,0 +1,309 @@ +using CodeyBox.HostProcess; +using CodeyBox.Sandbox.Incus; +using ControllableTimeProvider = Microsoft.Extensions.Time.Testing.FakeTimeProvider; + +namespace CodeyBox.Tests; + +/// +/// Pure-function coverage for , the +/// payload handling, and the input guards on +/// . +/// +public sealed class IncusGuestCpuActivityTests +{ + private const double ThresholdPercent = 5.0; + private static readonly TimeSpan Interval = TimeSpan.FromSeconds(10); + + private static IncusGuestCpuSample Sample(long cpuNanoseconds, long timestampSeconds = 0) => + new(cpuNanoseconds, DateTimeOffset.UnixEpoch.AddSeconds(timestampSeconds)); + + [Fact] + public void Evaluate_FirstSample_IsInactiveBaseline() + { + var evaluation = IncusCpuActivityEvaluator.Evaluate( + previousSample: null, + Sample(1_000_000_000), + Interval, + ThresholdPercent); + + Assert.False(evaluation.IsActive); + Assert.Equal(0.0, evaluation.CpuFraction); + } + + [Fact] + public void Evaluate_BusyDelta_IsActiveAndReportsFraction() + { + // 4 CPU-seconds consumed over a 10-second interval = 0.4 cores = 40%. + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(1_000_000_000, timestampSeconds: 100), + Sample(5_000_000_000, timestampSeconds: 110), + Interval, + ThresholdPercent); + + Assert.True(evaluation.IsActive); + Assert.Equal(0.4, evaluation.CpuFraction, precision: 6); + } + + [Fact] + public void Evaluate_IdleDelta_IsInactive() + { + // ~8.8 ms of CPU over 10 s — the observed near-idle VM rate. + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(1_000_000_000, timestampSeconds: 100), + Sample(1_008_800_000, timestampSeconds: 110), + Interval, + ThresholdPercent); + + Assert.False(evaluation.IsActive); + Assert.Equal(0.00088, evaluation.CpuFraction, precision: 6); + } + + [Fact] + public void Evaluate_ExactlyAtThreshold_IsActive() + { + // 0.5 CPU-seconds over 10 s = exactly 5% of one core. + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(1_000_000_000), + Sample(1_500_000_000), + Interval, + ThresholdPercent); + + Assert.True(evaluation.IsActive); + } + + [Fact] + public void Evaluate_JustBelowThreshold_IsInactive() + { + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(1_000_000_000), + Sample(1_499_999_999), + Interval, + ThresholdPercent); + + Assert.False(evaluation.IsActive); + } + + [Fact] + public void Evaluate_CounterResetOrRegression_IsInactive() + { + // A VM restart resets cpu.usage; the regression must not look like a + // giant positive delta (and must not wrap to a bogus fraction). + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(50_000_000_000), + Sample(1_000_000), + Interval, + ThresholdPercent); + + Assert.False(evaluation.IsActive); + Assert.Equal(0.0, evaluation.CpuFraction); + } + + [Fact] + public void Evaluate_UnchangedCounter_IsInactive() + { + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(7_000_000_000), + Sample(7_000_000_000), + Interval, + ThresholdPercent); + + Assert.False(evaluation.IsActive); + } + + [Theory] + [InlineData(0)] + [InlineData(-1)] + public void Evaluate_NonPositiveElapsed_IsInactive(long elapsedTicks) + { + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(1_000_000_000), + Sample(9_000_000_000), + TimeSpan.FromTicks(elapsedTicks), + ThresholdPercent); + + Assert.False(evaluation.IsActive); + } + + [Theory] + [InlineData(0.0)] + [InlineData(-5.0)] + [InlineData(double.NaN)] + [InlineData(double.PositiveInfinity)] + public void Evaluate_InvalidThreshold_IsInactive(double thresholdPercent) + { + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(1_000_000_000), + Sample(9_000_000_000), + Interval, + thresholdPercent); + + Assert.False(evaluation.IsActive); + } + + [Theory] + [InlineData(-5, 1_000_000_000)] + [InlineData(1_000_000_000, -5)] + public void Evaluate_NegativeCounter_IsInactive(long previousNs, long currentNs) + { + var evaluation = IncusCpuActivityEvaluator.Evaluate( + Sample(previousNs), + Sample(currentNs), + Interval, + ThresholdPercent); + + Assert.False(evaluation.IsActive); + } + + [Fact] + public void IsActive_MatchesEvaluateDecision() + { + Assert.True(IncusCpuActivityEvaluator.IsActive( + Sample(1_000_000_000), Sample(5_000_000_000), Interval, ThresholdPercent)); + Assert.False(IncusCpuActivityEvaluator.IsActive( + Sample(1_000_000_000), Sample(1_001_000_000), Interval, ThresholdPercent)); + } + + [Fact] + public void TryParseGuestCpuUsage_MetadataWrappedPayload_ReadsUsage() + { + const string json = """ + {"metadata": {"status": "Running", "cpu": {"usage": 123456789}}} + """; + + Assert.True(IncusStateParser.TryParseGuestCpuUsage(json, out var usage)); + Assert.Equal(123456789, usage); + } + + [Fact] + public void TryParseGuestCpuUsage_DirectPayload_ReadsUsage() + { + Assert.True(IncusStateParser.TryParseGuestCpuUsage( + """{"status": "Running", "cpu": {"usage": 42}}""", out var usage)); + Assert.Equal(42, usage); + } + + [Fact] + public void TryParseGuestCpuUsage_NonRunningInstance_ReturnsFalse() + { + Assert.False(IncusStateParser.TryParseGuestCpuUsage( + """{"status": "Stopped", "cpu": {"usage": 42}}""", out _)); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData(" ")] + [InlineData("not json")] + [InlineData("{}")] + [InlineData("""{"cpu": {"usage": -5}}""")] + [InlineData("""{"cpu": {}}""")] + public void TryParseGuestCpuUsage_UnusablePayload_ReturnsFalse(string? json) + { + Assert.False(IncusStateParser.TryParseGuestCpuUsage(json, out var usage)); + Assert.Equal(0, usage); + } + + [Theory] + [InlineData("bad/name")] + [InlineData("bad?name")] + [InlineData("bad%2fname")] + [InlineData("bad name")] + public async Task ReadStateAsync_InvalidInstanceName_ThrowsBeforeQuerying(string instanceName) + { + var runner = new NeverInvokedProcessRunner(); + var reader = new DefaultIncusInstanceStateReader( + new IncusCliRunner(runner), TimeProvider.System); + + await Assert.ThrowsAsync(() => + reader.ReadStateAsync(new IncusSandboxOptions(), instanceName, CancellationToken.None)); + Assert.Equal(0, runner.Calls); + } + + [Fact] + public async Task ReadStateAsync_InvalidProjectIdentity_ThrowsBeforeQuerying() + { + var runner = new NeverInvokedProcessRunner(); + var reader = new DefaultIncusInstanceStateReader( + new IncusCliRunner(runner), TimeProvider.System); + + await Assert.ThrowsAsync(() => + reader.ReadStateAsync( + new IncusSandboxOptions { ProjectName = "bad/project" }, + "codeybox-vm", + CancellationToken.None)); + Assert.Equal(0, runner.Calls); + } + + [Fact] + public async Task ReadStateAsync_RunningInstance_ReturnsTimestampedSample() + { + var time = new ControllableTimeProvider(DateTimeOffset.UtcNow); + var runner = new ScriptedStateRunner( + """{"metadata": {"status": "Running", "cpu": {"usage": 987654321}}}"""); + var reader = new DefaultIncusInstanceStateReader(new IncusCliRunner(runner, time), time); + + var sample = await reader.ReadStateAsync( + new IncusSandboxOptions(), "codeybox-vm-1", CancellationToken.None); + + Assert.NotNull(sample); + Assert.Equal(987654321, sample!.Value.CpuUsageNanoseconds); + Assert.Equal(time.GetUtcNow(), sample.Value.Timestamp); + var argv = Assert.Single(runner.Arguments); + Assert.Equal(["incus", "query", "/1.0/instances/codeybox-vm-1/state?project=codeybox"], argv); + } + + [Fact] + public async Task ReadStateAsync_FailedQuery_ReturnsNull() + { + var runner = new ScriptedStateRunner(null); + var reader = new DefaultIncusInstanceStateReader( + new IncusCliRunner(runner), TimeProvider.System); + + var sample = await reader.ReadStateAsync( + new IncusSandboxOptions(), "codeybox-vm-1", CancellationToken.None); + + Assert.Null(sample); + } + + private sealed class NeverInvokedProcessRunner : IProcessRunner + { + internal int Calls { get; private set; } + + public Task RunAsync( + IReadOnlyList argv, + string? stdin, + CancellationToken ct, + Action? stdoutChunkCallback = null, + Action? stderrChunkCallback = null, + int? maxStdoutBytes = null, + int? maxStderrBytes = null, + IReadOnlyDictionary? environment = null, + bool killOnOutputLimit = true) + { + Calls++; + throw new InvalidOperationException("The Incus CLI must not run for rejected input."); + } + } + + private sealed class ScriptedStateRunner(string? stdout) : IProcessRunner + { + internal List> Arguments { get; } = []; + + public Task RunAsync( + IReadOnlyList argv, + string? stdin, + CancellationToken ct, + Action? stdoutChunkCallback = null, + Action? stderrChunkCallback = null, + int? maxStdoutBytes = null, + int? maxStderrBytes = null, + IReadOnlyDictionary? environment = null, + bool killOnOutputLimit = true) + { + Arguments.Add(argv.ToArray()); + return Task.FromResult(stdout is null + ? new ProcessRunResult(1, string.Empty, "query failed") + : new ProcessRunResult(0, stdout, string.Empty)); + } + } +} diff --git a/tests/CodeyBox.Tests/IncusSandboxLifecycleTests.cs b/tests/CodeyBox.Tests/IncusSandboxLifecycleTests.cs index 7f4590f9f..2ccff4c82 100644 --- a/tests/CodeyBox.Tests/IncusSandboxLifecycleTests.cs +++ b/tests/CodeyBox.Tests/IncusSandboxLifecycleTests.cs @@ -1,5 +1,6 @@ using CodeyBox.Core; using CodeyBox.HostProcess; +using CodeyBox.Orchestrator; using CodeyBox.Sandbox; using CodeyBox.Sandbox.Incus; using Microsoft.Extensions.Logging.Abstractions; @@ -3155,6 +3156,447 @@ public async Task StopAndPreserve_ForceStopStillRunning_LaterDisposeForceDeletes } } + // ── Guest-CPU activity signal ──────────────────────────────────────────── + + [Fact] + public async Task ActiveSandboxProgress_BusyGuest_ChangesSignatureBetweenSnapshots() + { + var fixture = PrepareRetainedAdoptionFixture("codeybox-cpu-busy"); + var time = new ControllableTimeProvider(DateTimeOffset.UtcNow); + var reader = new ScriptedInstanceStateReader(time) + { + NextCpuUsage = BusyCpuCounter(time), + }; + var runner = new RetainedAdoptionRunner( + fixture.Options.StagingDirectory!, + fixture.SandboxName, + fixture.Manifest.LeaseTokenSha256, + fixture.ManifestHash); + var options = fixture.Options with + { + // The boot stagger is awaited through the injected clock; a fake + // clock would hang CreateAsync until the test advances it. + BootLaunchDelay = TimeSpan.Zero, + }; + var provider = new IncusSandboxProvider( + () => options, + NullLogger.Instance, + timings: null, + runner, + timeProvider: time, + stateReader: reader); + + try + { + var adopted = await provider.CreateAsync(fixture.RequestSpec); + var workItemId = fixture.RequestSpec.TimingWorkItemId!.Value; + + var first = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + Assert.Equal(workItemId, first.WorkItemId); + Assert.Equal(fixture.SandboxName, first.SandboxId); + Assert.Null(first.CpuFraction); + + // ~80% of one core over the 5-second sample interval exceeds the + // default 5% threshold, so the emitted status advances. + time.Advance(TimeSpan.FromSeconds(5)); + var second = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + Assert.NotEqual(first.Status, second.Status); + Assert.Equal(0.8, second.CpuFraction!.Value, precision: 3); + + // The non-refreshing view exposes the same projection. + Assert.Equal(second.Status, Assert.Single(provider.SnapshotActiveSandboxProgress()).Status); + + time.Advance(TimeSpan.FromSeconds(5)); + var third = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + Assert.NotEqual(second.Status, third.Status); + + await adopted.DisposeAsync(); + Assert.Empty(provider.SnapshotActiveSandboxProgress()); + } + finally + { + if (Directory.Exists(fixture.Options.StagingDirectory)) + Directory.Delete(fixture.Options.StagingDirectory, recursive: true); + } + } + + [Fact] + public async Task ActiveSandboxProgress_IdleGuest_KeepsStableSignature() + { + var fixture = PrepareRetainedAdoptionFixture("codeybox-cpu-idle"); + var time = new ControllableTimeProvider(DateTimeOffset.UtcNow); + var reader = new ScriptedInstanceStateReader(time) + { + NextCpuUsage = () => 1_000_000_000, + }; + var runner = new RetainedAdoptionRunner( + fixture.Options.StagingDirectory!, + fixture.SandboxName, + fixture.Manifest.LeaseTokenSha256, + fixture.ManifestHash); + var options = fixture.Options with + { + // The boot stagger is awaited through the injected clock; a fake + // clock would hang CreateAsync until the test advances it. + BootLaunchDelay = TimeSpan.Zero, + }; + var provider = new IncusSandboxProvider( + () => options, + NullLogger.Instance, + timings: null, + runner, + timeProvider: time, + stateReader: reader); + + try + { + await provider.CreateAsync(fixture.RequestSpec); + + var first = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + time.Advance(TimeSpan.FromSeconds(5)); + var second = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + time.Advance(TimeSpan.FromSeconds(5)); + var third = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + + Assert.Equal(first.Status, second.Status); + Assert.Equal(first.Status, third.Status); + Assert.Equal(3, reader.Calls); + } + finally + { + if (Directory.Exists(fixture.Options.StagingDirectory)) + Directory.Delete(fixture.Options.StagingDirectory, recursive: true); + } + } + + [Fact] + public async Task ActiveSandboxProgress_FailedQuery_YieldsNoSignalAndKeepsStableSignature() + { + var fixture = PrepareRetainedAdoptionFixture("codeybox-cpu-failure"); + var time = new ControllableTimeProvider(DateTimeOffset.UtcNow); + var reader = new ScriptedInstanceStateReader(time); + var runner = new RetainedAdoptionRunner( + fixture.Options.StagingDirectory!, + fixture.SandboxName, + fixture.Manifest.LeaseTokenSha256, + fixture.ManifestHash); + var options = fixture.Options with + { + // The boot stagger is awaited through the injected clock; a fake + // clock would hang CreateAsync until the test advances it. + BootLaunchDelay = TimeSpan.Zero, + }; + var provider = new IncusSandboxProvider( + () => options, + NullLogger.Instance, + timings: null, + runner, + timeProvider: time, + stateReader: reader); + + try + { + await provider.CreateAsync(fixture.RequestSpec); + + reader.NextCpuUsage = () => 1_000_000_000; + var first = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + + // A throwing reader must not propagate into the watchdog probe and + // must not alter the signature. + reader.Failure = new InvalidOperationException("incus query exploded"); + time.Advance(TimeSpan.FromSeconds(5)); + var second = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + Assert.Equal(first.Status, second.Status); + + // A reader that returns no usable sample behaves identically. + reader.Failure = null; + reader.NextCpuUsage = () => null; + time.Advance(TimeSpan.FromSeconds(5)); + var third = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + Assert.Equal(first.Status, third.Status); + } + finally + { + if (Directory.Exists(fixture.Options.StagingDirectory)) + Directory.Delete(fixture.Options.StagingDirectory, recursive: true); + } + } + + [Fact] + public async Task ActiveSandboxProgress_WithinSampleInterval_ReusesCachedProjection() + { + var fixture = PrepareRetainedAdoptionFixture("codeybox-cpu-gated"); + var time = new ControllableTimeProvider(DateTimeOffset.UtcNow); + var reader = new ScriptedInstanceStateReader(time) + { + NextCpuUsage = () => 1_000_000_000, + }; + var runner = new RetainedAdoptionRunner( + fixture.Options.StagingDirectory!, + fixture.SandboxName, + fixture.Manifest.LeaseTokenSha256, + fixture.ManifestHash); + var options = fixture.Options with + { + // The boot stagger is awaited through the injected clock; a fake + // clock would hang CreateAsync until the test advances it. + BootLaunchDelay = TimeSpan.Zero, + }; + var provider = new IncusSandboxProvider( + () => options, + NullLogger.Instance, + timings: null, + runner, + timeProvider: time, + stateReader: reader); + + try + { + await provider.CreateAsync(fixture.RequestSpec); + + var first = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + // Inside ActivitySampleInterval the cached projection is reused — + // no second query, no signature change. + var second = Assert.Single(await provider.SnapshotActiveSandboxProgressAsync()); + Assert.Equal(1, reader.Calls); + Assert.Equal(first.Status, second.Status); + + time.Advance(fixture.Options.ActivitySampleInterval); + _ = await provider.SnapshotActiveSandboxProgressAsync(); + Assert.Equal(2, reader.Calls); + } + finally + { + if (Directory.Exists(fixture.Options.StagingDirectory)) + Directory.Delete(fixture.Options.StagingDirectory, recursive: true); + } + } + + [Fact] + public async Task Watchdog_SilentStreamBusyIncusGuest_KeepsWorkerAlivePastProgressTimeout() + { + var fixture = PrepareRetainedAdoptionFixture("codeybox-cpu-watchdog-busy"); + var time = new ControllableTimeProvider(DateTimeOffset.UtcNow); + var reader = new ScriptedInstanceStateReader(time) + { + // Each query reports ~80% of one core consumed since the previous + // sample, whatever the interval between watchdog sweeps. + NextCpuUsage = BusyCpuCounter(time), + }; + var runner = new RetainedAdoptionRunner( + fixture.Options.StagingDirectory!, + fixture.SandboxName, + fixture.Manifest.LeaseTokenSha256, + fixture.ManifestHash); + var options = fixture.Options with + { + // The boot stagger is awaited through the injected clock; a fake + // clock would hang CreateAsync until the test advances it. + BootLaunchDelay = TimeSpan.Zero, + }; + var provider = new IncusSandboxProvider( + () => options, + NullLogger.Instance, + timings: null, + runner, + timeProvider: time, + stateReader: reader); + + using var scratch = TestScratchDirectory.Create("codeybox-incus-watchdog-"); + var dbPath = scratch.DbPath("watchdog.db"); + using var store = new SqliteWorkItemStore(dbPath); + using var registry = new SqliteWorkerRegistry(dbPath); + var queue = new InMemoryTaskQueue(); + try + { + var adopted = await provider.CreateAsync(fixture.RequestSpec); + var itemId = fixture.RequestSpec.TimingWorkItemId!.Value; + var item = WatchdogItem(itemId, time.GetUtcNow() - TimeSpan.FromMinutes(45)); + await store.CreateAsync(item); + await PlantWatchdogWorkerAsync(registry, itemId, time.GetUtcNow()); + + // streams: null — a completely silent agent stream, the exact + // scenario that used to get busy Incus guests killed. + var watchdog = new WorkerProgressWatchdog( + registry, store, queue, + new WorkerProgressWatchdogOptions + { + ProgressTimeout = TimeSpan.FromMinutes(30), + CheckInterval = TimeSpan.FromMinutes(1), + ProcessCpuProgressSignalEnabled = false, + ActiveSandboxProgressSignalEnabled = true, + }, + NullLogger.Instance, + streams: null, + activitySource: new DefaultWorkerProgressActivitySource(provider), + timeProvider: time); + + // Three sweeps each a full ProgressTimeout apart: the changing + // guest-CPU signature keeps registering fresh progress. + for (var sweep = 0; sweep < 3; sweep++) + { + await watchdog.RunOnceAsync(CancellationToken.None); + time.Advance(TimeSpan.FromMinutes(31)); + } + + var after = await store.GetAsync(item.Id); + Assert.Equal(WorkItemState.Working, after!.State); + Assert.Equal(0, after.RecoveryAttempts); + Assert.Equal(0, queue.Count); + Assert.True(reader.Calls >= 3, $"expected at least 3 guest-state queries, saw {reader.Calls}"); + + await adopted.DisposeAsync(); + } + finally + { + if (Directory.Exists(fixture.Options.StagingDirectory)) + Directory.Delete(fixture.Options.StagingDirectory, recursive: true); + } + } + + [Fact] + public async Task Watchdog_SilentStreamIdleIncusGuest_RecoversWorkerPastProgressTimeout() + { + var fixture = PrepareRetainedAdoptionFixture("codeybox-cpu-watchdog-idle"); + var time = new ControllableTimeProvider(DateTimeOffset.UtcNow); + var reader = new ScriptedInstanceStateReader(time) + { + // A genuinely hung guest: the cpu.usage counter never advances. + NextCpuUsage = () => 1_000_000_000, + }; + var runner = new RetainedAdoptionRunner( + fixture.Options.StagingDirectory!, + fixture.SandboxName, + fixture.Manifest.LeaseTokenSha256, + fixture.ManifestHash); + var options = fixture.Options with + { + // The boot stagger is awaited through the injected clock; a fake + // clock would hang CreateAsync until the test advances it. + BootLaunchDelay = TimeSpan.Zero, + }; + var provider = new IncusSandboxProvider( + () => options, + NullLogger.Instance, + timings: null, + runner, + timeProvider: time, + stateReader: reader); + + using var scratch = TestScratchDirectory.Create("codeybox-incus-watchdog-"); + var dbPath = scratch.DbPath("watchdog.db"); + using var store = new SqliteWorkItemStore(dbPath); + using var registry = new SqliteWorkerRegistry(dbPath); + var queue = new InMemoryTaskQueue(); + try + { + await provider.CreateAsync(fixture.RequestSpec); + var itemId = fixture.RequestSpec.TimingWorkItemId!.Value; + var item = WatchdogItem(itemId, time.GetUtcNow() - TimeSpan.FromMinutes(45)); + await store.CreateAsync(item); + await PlantWatchdogWorkerAsync(registry, itemId, time.GetUtcNow()); + + var watchdog = new WorkerProgressWatchdog( + registry, store, queue, + new WorkerProgressWatchdogOptions + { + ProgressTimeout = TimeSpan.FromMinutes(30), + CheckInterval = TimeSpan.FromMinutes(1), + ProcessCpuProgressSignalEnabled = false, + ActiveSandboxProgressSignalEnabled = true, + }, + NullLogger.Instance, + streams: null, + activitySource: new DefaultWorkerProgressActivitySource(provider), + timeProvider: time); + + // First sweep records the initial signature (first sighting counts + // as progress); the second, a full timeout later, sees the + // unchanged idle signature and recovers the worker. + await watchdog.RunOnceAsync(CancellationToken.None); + time.Advance(TimeSpan.FromMinutes(31)); + await watchdog.RunOnceAsync(CancellationToken.None); + + var after = await store.GetAsync(item.Id); + Assert.Equal(WorkItemState.Queued, after!.State); + Assert.Equal(1, after.RecoveryAttempts); + Assert.Equal(1, queue.Count); + } + finally + { + if (Directory.Exists(fixture.Options.StagingDirectory)) + Directory.Delete(fixture.Options.StagingDirectory, recursive: true); + } + } + + private static WorkItem WatchdogItem(WorkItemId itemId, DateTimeOffset updatedAt) => new() + { + Id = itemId, + ProjectId = new ProjectId("test"), + Title = "t", + Prompt = "p", + State = WorkItemState.Working, + UpdatedAt = updatedAt, + DependsOn = [], + }; + + private static async Task PlantWatchdogWorkerAsync( + SqliteWorkerRegistry registry, WorkItemId itemId, DateTimeOffset now) + { + // A fresh heartbeat — the watchdog must NOT rely on heartbeat staleness; + // progress is decided from item.UpdatedAt + stream/activity signals. + await registry.RegisterAsync(new WorkerRegistration + { + WorkerId = Guid.NewGuid().ToString(), + HostName = "host", + ProcessId = 1, + StartedAt = now.AddHours(-1), + LastHeartbeatAt = now, + CurrentWorkItemId = itemId.ToString(), + }); + } + + /// + /// A cpu.usage counter that grows at a fixed fraction of one core against + /// the injected clock, so the sampled fraction is identical regardless of + /// the interval between queries. + /// + private static Func BusyCpuCounter(TimeProvider time, double fractionOfCore = 0.8) + { + var cpu = 1_000_000_000L; + var lastAt = time.GetUtcNow(); + return () => + { + var now = time.GetUtcNow(); + cpu += (long)((now - lastAt).TotalSeconds * fractionOfCore * 1_000_000_000.0); + lastAt = now; + return cpu; + }; + } + + private sealed class ScriptedInstanceStateReader(TimeProvider timeProvider) : IIncusInstanceStateReader + { + internal int Calls { get; private set; } + internal List QueriedInstances { get; } = []; + internal Exception? Failure { get; set; } + internal Func NextCpuUsage { get; set; } = static () => null; + + public Task ReadStateAsync( + IncusSandboxOptions options, + string instanceName, + CancellationToken ct) + { + Calls++; + QueriedInstances.Add(instanceName); + if (Failure is { } failure) + throw failure; + var usage = NextCpuUsage(); + return Task.FromResult(usage is { } nanoseconds + ? new IncusGuestCpuSample(nanoseconds, timeProvider.GetUtcNow()) + : (IncusGuestCpuSample?)null); + } + } + private static IncusSandboxOptions FastLifecycleOptions() => new() { CaptureResourceMetrics = false,