Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions src/CodeyBox.Api/ReloadableSandboxProvider.cs
Original file line number Diff line number Diff line change
Expand Up @@ -207,6 +207,17 @@ public IReadOnlyList<ActiveSandboxProgress> SnapshotActiveSandboxProgress() =>
.SelectMany(static provider => provider.Progress.SnapshotActiveSandboxProgress())
.ToArray();

public async ValueTask<IReadOnlyList<ActiveSandboxProgress>> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default)
{
// 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<DiskGuardSample> SampleDiskGuardState() =>
ActivatedProviders
.SelectMany(static provider => provider.DiskGuard.SampleDiskGuardState())
Expand Down
27 changes: 24 additions & 3 deletions src/CodeyBox.Core/SandboxAbstractions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -818,7 +818,15 @@ public interface IActiveSandboxProvider
/// signal for detached VM-local work and use <see cref="Status"/> to report a
/// richer reason when providers can expose changing activity.
/// </summary>
public sealed record ActiveSandboxProgress(WorkItemId WorkItemId, string SandboxId, string? Status = null);
/// <param name="WorkItemId">Work item that owns the sandbox.</param>
/// <param name="SandboxId">Provider-side sandbox identifier.</param>
/// <param name="Status">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.</param>
/// <param name="CpuFraction">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.</param>
public sealed record ActiveSandboxProgress(WorkItemId WorkItemId, string SandboxId, string? Status = null, double? CpuFraction = null);

/// <summary>
/// Optional provider capability for reporting active sandbox ownership.
Expand All @@ -827,10 +835,23 @@ public sealed record ActiveSandboxProgress(WorkItemId WorkItemId, string Sandbox
public interface IActiveSandboxProgressProvider
{
/// <summary>
/// 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
/// <see cref="SnapshotActiveSandboxProgressAsync"/>, so both members report
/// the same projection.
/// </summary>
IReadOnlyList<ActiveSandboxProgress> SnapshotActiveSandboxProgress();

/// <summary>
/// Asynchronously samples and snapshots currently-active sandboxes, refreshing
/// activity projections such as guest CPU metrics where supported. The
/// default delegates to the non-refreshing
/// <see cref="SnapshotActiveSandboxProgress"/>.
/// </summary>
ValueTask<IReadOnlyList<ActiveSandboxProgress>> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default)
=> ValueTask.FromResult(SnapshotActiveSandboxProgress());
}

/// <summary>
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -312,6 +312,9 @@ public async Task DisposeLeakedAsync(ManagedSandboxInfo sandbox, CancellationTok
public IReadOnlyList<ActiveSandboxProgress> SnapshotActiveSandboxProgress() =>
_progressProvider?.SnapshotActiveSandboxProgress() ?? [];

public ValueTask<IReadOnlyList<ActiveSandboxProgress>> SnapshotActiveSandboxProgressAsync(CancellationToken ct = default) =>
_progressProvider?.SnapshotActiveSandboxProgressAsync(ct) ?? ValueTask.FromResult<IReadOnlyList<ActiveSandboxProgress>>([]);

public IReadOnlyList<SandboxHostPoolEntry> SnapshotHostPool() =>
_hostPoolSnapshot?.SnapshotHostPool() ?? [];

Expand Down
53 changes: 32 additions & 21 deletions src/CodeyBox.Orchestrator/WorkerProgressActivitySource.cs
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,12 @@ 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.
/// </summary>
public sealed record WorkerProgressActivity(string Reason);
/// <param name="Reason">Short machine-readable label naming the signal that
/// fired (e.g. <c>process-cpu</c>, <c>active-sandbox-change</c>).</param>
/// <param name="CpuFraction">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.</param>
public sealed record WorkerProgressActivity(string Reason, double? CpuFraction = null);

/// <summary>
/// Narrow activity probe settings resolved by the watchdog for the current
Expand Down Expand Up @@ -92,77 +97,83 @@ private DefaultWorkerProgressActivitySource(
_initialCpuSampleAttempts = Math.Max(1, initialCpuSampleAttempts);
}

public ValueTask<WorkerProgressActivity?> ObserveAsync(
public async ValueTask<WorkerProgressActivity?> ObserveAsync(
WorkerRegistration worker,
WorkItemId itemId,
WorkerProgressActivityProbe probe,
CancellationToken ct)
{
ct.ThrowIfCancellationRequested();
if (!string.Equals(worker.CurrentWorkItemId, itemId.ToString(), StringComparison.Ordinal))
return ValueTask.FromResult<WorkerProgressActivity?>(null);
return null;

if (probe.ProcessCpuProgressSignalEnabled
&& TryObserveProcessCpu(itemId, out var cpuReason))
{
return ValueTask.FromResult<WorkerProgressActivity?>(
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<WorkerProgressActivity?>(
new WorkerProgressActivity(sandboxReason));
return sandboxActivity;
}

return ValueTask.FromResult<WorkerProgressActivity?>(null);
return null;
}

private bool TryObserveActiveSandbox(WorkItemId itemId, out string reason)
private async ValueTask<WorkerProgressActivity?> TryObserveActiveSandboxAsync(WorkItemId itemId, CancellationToken ct)
{
reason = "";
if (_activeSandboxProvider is not { } activeProvider)
return false;
return null;

IReadOnlyList<ActiveSandboxProgress> snapshot;
try
{
snapshot = activeProvider.SnapshotActiveSandboxProgress();
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 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();

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)
Expand Down
8 changes: 6 additions & 2 deletions src/CodeyBox.Orchestrator/WorkerProgressWatchdog.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using System.Collections.Concurrent;
using System.Globalization;
using Microsoft.Extensions.Hosting;
using Microsoft.Extensions.Logging;
using CodeyBox.Core;
Expand Down Expand Up @@ -194,9 +195,12 @@ public async Task RunOnceAsync(CancellationToken ct)
if (workerActivity is not null)
{
_workerActivityProgress[activityKey] = new WorkerActivityProgress(now, 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}; treating as progress",
worker.WorkerId, itemId, workerActivity.Reason);
"Watchdog: worker {WorkerId} for item {ItemId} has live activity signal {Reason} (cpu={CpuFraction}); treating as progress",
worker.WorkerId, itemId, workerActivity.Reason, cpuText);
continue;
}

Expand Down
53 changes: 53 additions & 0 deletions src/CodeyBox.Sandbox.Incus/IIncusInstanceStateReader.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
namespace CodeyBox.Sandbox.Incus;

/// <summary>
/// Seam for reading instance guest state (CPU usage) from Incus.
/// </summary>
internal interface IIncusInstanceStateReader
{
Task<IncusGuestCpuSample?> ReadStateAsync(
IncusSandboxOptions options,
string instanceName,
CancellationToken ct);
}

/// <summary>
/// Default implementation that queries Incus instance state via the CLI/API runner.
/// </summary>
internal sealed class DefaultIncusInstanceStateReader(
IncusCliRunner cli,
TimeProvider timeProvider) : IIncusInstanceStateReader
{
public async Task<IncusGuestCpuSample?> ReadStateAsync(
IncusSandboxOptions options,
string instanceName,
CancellationToken ct)
{
// 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,
[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;
}
}
57 changes: 57 additions & 0 deletions src/CodeyBox.Sandbox.Incus/IncusCpuActivityEvaluator.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,57 @@
namespace CodeyBox.Sandbox.Incus;

/// <summary>
/// Pure evaluation of whether an Incus guest consumed enough CPU over an
/// interval to count as watchdog progress.
/// </summary>
public static class IncusCpuActivityEvaluator
{
/// <summary>
/// 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 <paramref name="elapsed"/>, reaches
/// <paramref name="thresholdPercent"/>. 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.
/// </summary>
/// <param name="previousSample">The earlier sample, or null when this is
/// the first observation (which always establishes a baseline and reports
/// inactive).</param>
/// <param name="currentSample">The latest sample.</param>
/// <param name="elapsed">Wall time between the two samples.</param>
/// <param name="thresholdPercent">Activity threshold as a percentage of one
/// CPU core (5 = 5% of a core, i.e. 0.05 CPU fraction).</param>
public static IncusCpuActivityEvaluation Evaluate(
IncusGuestCpuSample? previousSample,
IncusGuestCpuSample currentSample,
TimeSpan elapsed,
double thresholdPercent)
{
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);
}

/// <summary>
/// Convenience wrapper over <see cref="Evaluate"/> reporting only the
/// active/inactive decision.
/// </summary>
public static bool IsActive(
IncusGuestCpuSample? previousSample,
IncusGuestCpuSample currentSample,
TimeSpan elapsed,
double thresholdPercent) =>
Evaluate(previousSample, currentSample, elapsed, thresholdPercent).IsActive;
}
25 changes: 25 additions & 0 deletions src/CodeyBox.Sandbox.Incus/IncusGuestCpuSample.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
namespace CodeyBox.Sandbox.Incus;

/// <summary>
/// A point-in-time sample of an Incus guest's cumulative CPU consumption — the
/// monotonic <c>cpu.usage</c> nanoseconds counter reported by
/// <c>incus query /1.0/instances/&lt;name&gt;/state</c>.
/// </summary>
/// <param name="CpuUsageNanoseconds">Cumulative guest CPU time in nanoseconds
/// since the instance started; resets when the instance restarts.</param>
/// <param name="Timestamp">When the sample was read. The sampling interval is
/// derived from consecutive sample timestamps.</param>
public readonly record struct IncusGuestCpuSample(
long CpuUsageNanoseconds,
DateTimeOffset Timestamp = default);

/// <summary>
/// Evaluation result of guest CPU activity over an elapsed interval.
/// </summary>
/// <param name="IsActive">Whether the guest met the configured activity
/// threshold.</param>
/// <param name="CpuFraction">Measured guest CPU as a fraction of one core
/// averaged over the sample interval (1.0 = one fully busy core).</param>
public readonly record struct IncusCpuActivityEvaluation(
bool IsActive,
double CpuFraction);
21 changes: 21 additions & 0 deletions src/CodeyBox.Sandbox.Incus/IncusSandboxOptions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -230,6 +230,23 @@ 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);
/// <summary>
/// Guest-CPU activity threshold for the watchdog's active-sandbox signal,
/// as a percentage of one CPU core averaged over
/// <see cref="ActivitySampleInterval"/> (5 = 5% of one core). The emitted
/// progress signature changes only while the guest meets this threshold.
/// </summary>
public double ActivityCpuThresholdPercent { get; init; } = 5.0;
/// <summary>
/// Minimum wall-clock interval between guest-CPU state queries for a given
/// sandbox. Snapshots requested sooner reuse the last sampled projection.
/// </summary>
public TimeSpan ActivitySampleInterval { get; init; } = TimeSpan.FromSeconds(5);
/// <summary>
/// Per-call deadline for each <c>incus query /1.0/instances/&lt;name&gt;/state</c>
/// read. A timed-out query contributes no activity signal.
/// </summary>
public TimeSpan ActivityQueryTimeout { get; init; } = TimeSpan.FromSeconds(5);
public IncusDiskGuardOptions? DiskGuard { get; init; } = new();

/// <summary>
Expand Down Expand Up @@ -353,6 +370,10 @@ public static IReadOnlyList<string> 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)
Expand Down
Loading
Loading