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
8 changes: 8 additions & 0 deletions Services/Configuration/WorkerSettings.cs
Original file line number Diff line number Diff line change
Expand Up @@ -92,6 +92,14 @@ public class WorkerSettings
/// </summary>
public int BgTimeoutSeconds { get; set; }

/// <summary>
/// How long a timed-out or cancelled invocation gets to unwind after PowerShell is asked to stop.
/// A pipeline blocked inside a .NET call (a lock, a sync wait, a read with no timeout) never sees a
/// stop request; past this grace its worker is abandoned and replaced, so the pool does not shrink
/// to nothing. Values below 1 use the default of 30 seconds.
/// </summary>
public int StopGraceSeconds { get; set; } = 30;

/// <summary>
/// How long an incoming HTTP request waits for a free PowerShell runspace when every worker in the
/// HTTP pool is already busy, before it is shed with <c>503 "Server busy, please retry"</c>.
Expand Down
1 change: 1 addition & 0 deletions Services/Hosting/Endpoints/HealthEndpoint.cs
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@ public static WebApplication MapCraftHealthEndpoint(
status = (httpReady && bgReady && storageReady) ? "ready" : "starting",
roles = new { frontend = roles.Frontend, http = roles.Http, background = roles.Background },
ready = new { http = httpReady, background = bgReady, storage = storageReady },
workers = new { hung = pool.HungWorkers, hungTotal = pool.HungWorkersTotal },
});
});

Expand Down
8 changes: 8 additions & 0 deletions Services/Orchestration/OrchestratorService.cs
Original file line number Diff line number Diff line change
Expand Up @@ -527,6 +527,14 @@ void StepDone(string status, string? error, WorkStore.ClaimedTask? next) =>
break;
}

if (worker is { IsHung: true })
{
// A step timed out on a worker that could not be stopped. Hand it back for replacement.
ReclaimSequentialWorker(worker, faulted: true);
worker = null;
worker = CheckoutSequentialWorker(jobCt);
}

var parameters = nextPayload ?? await _store.GetPayloadAsync(header.RunKey, step.Seq, jobCt);
nextPayload = null;
WorkStore.Finish finish;
Expand Down
170 changes: 133 additions & 37 deletions Services/PowerShellHost/PowerShellWorker.cs
Original file line number Diff line number Diff line change
Expand Up @@ -315,12 +315,14 @@ public void RestoreExportedModuleNames(ExportedModuleState state)
/// When a cancellation token is provided and fires, the PowerShell pipeline is
/// stopped via <see cref="PowerShell.Stop"/>. The resulting <c>PipelineStoppedException</c>
/// is normalized to an <see cref="OperationCanceledException"/> so timeout callers can
/// distinguish a cancelled request from a genuine script failure.
/// distinguish a cancelled request from a genuine script failure. A pipeline that does not
/// stop within <see cref="StopGrace"/> throws <see cref="WorkerHungException"/> and leaves
/// this worker <see cref="IsHung"/>.
/// </summary>
public async Task<Collection<PSObject>> InvokeAsync(string functionName, Dictionary<string, object?> parameters,
CancellationToken ct = default)
{
CancellationTokenRegistration? registration = null;
ThrowIfHung(functionName);
var prof = DispatchProfiler.Enabled;
long buildTicks = 0, runTicks = 0, copyTicks = 0;
try
Expand All @@ -330,19 +332,10 @@ public async Task<Collection<PSObject>> InvokeAsync(string functionName, Diction
_pwsh.AddCommand(functionName);
foreach (var p in parameters)
_pwsh.AddParameter(p.Key, p.Value);

// Register cancellation callback to stop the PS pipeline
if (ct.CanBeCanceled)
registration = ct.Register(() => _pwsh.Stop());
if (prof) buildTicks = System.Diagnostics.Stopwatch.GetTimestamp() - bStart;

var rStart = prof ? System.Diagnostics.Stopwatch.GetTimestamp() : 0;
// A fresh output collection per invocation: after an invocation that threw, the PowerShell object's
// own output buffer comes back null from the next EndInvoke, silently dropping that call's output.
using var outputs = new PSDataCollection<PSObject>();
var asyncResult = _pwsh.BeginInvoke<PSObject, PSObject>(null, outputs);
await Task.Factory.FromAsync(asyncResult, _pwsh.EndInvoke);
ct.ThrowIfCancellationRequested();
using var outputs = await RunPipelineAsync(functionName, ct);
if (prof) runTicks = System.Diagnostics.Stopwatch.GetTimestamp() - rStart;

var cpStart = prof ? System.Diagnostics.Stopwatch.GetTimestamp() : 0;
Expand All @@ -352,18 +345,21 @@ public async Task<Collection<PSObject>> InvokeAsync(string functionName, Diction
}
catch (PipelineStoppedException) when (ct.IsCancellationRequested)
{
// ct.Register(_pwsh.Stop) stopped the pipeline mid-invoke, so EndInvoke threw
// PipelineStoppedException rather than reaching ThrowIfCancellationRequested below.
// Normalize to OperationCanceledException so timeout callers return 504, not 500.
// The stop landed mid-invoke, so EndInvoke threw PipelineStoppedException rather than
// reaching ThrowIfCancellationRequested. Normalize to OperationCanceledException so
// timeout callers return 504, not 500.
throw new OperationCanceledException(ct);
}
finally
{
registration?.Dispose();
var cStart = prof ? System.Diagnostics.Stopwatch.GetTimestamp() : 0;
Cleanup();
if (prof) DispatchProfiler.RecordInvokeDetail(buildTicks, runTicks, copyTicks,
System.Diagnostics.Stopwatch.GetTimestamp() - cStart);
// A hung worker's pipeline is still running, and the worker never runs again.
if (!IsHung)
{
var cStart = prof ? System.Diagnostics.Stopwatch.GetTimestamp() : 0;
Cleanup();
if (prof) DispatchProfiler.RecordInvokeDetail(buildTicks, runTicks, copyTicks,
System.Diagnostics.Stopwatch.GetTimestamp() - cStart);
}
}
}

Expand All @@ -373,7 +369,7 @@ public async Task<Collection<PSObject>> InvokeAsync(string functionName, Diction
public async Task<Collection<PSObject>> InvokeScriptAsync(ScriptBlock scriptBlock,
Dictionary<string, object?>? parameters = null, CancellationToken ct = default)
{
CancellationTokenRegistration? registration = null;
ThrowIfHung("script");
try
{
StampOperationContext();
Expand All @@ -382,30 +378,117 @@ public async Task<Collection<PSObject>> InvokeScriptAsync(ScriptBlock scriptBloc
foreach (var p in parameters)
_pwsh.AddParameter(p.Key, p.Value);

if (ct.CanBeCanceled)
registration = ct.Register(() => _pwsh.Stop());

// A fresh output collection per invocation; see InvokeAsync.
using var outputs = new PSDataCollection<PSObject>();
var asyncResult = _pwsh.BeginInvoke<PSObject, PSObject>(null, outputs);
await Task.Factory.FromAsync(asyncResult, _pwsh.EndInvoke);

ct.ThrowIfCancellationRequested();
using var outputs = await RunPipelineAsync("script", ct);
return new Collection<PSObject>(outputs.ReadAll());
}
catch (PipelineStoppedException) when (ct.IsCancellationRequested)
{
// See InvokeAsync: a cancellation-triggered _pwsh.Stop() surfaces as
// PipelineStoppedException; normalize it to OperationCanceledException.
// See InvokeAsync.
throw new OperationCanceledException(ct);
}
finally
{
registration?.Dispose();
Cleanup();
if (!IsHung) Cleanup();
}
}

/// <summary>
/// How long a cancelled invocation gets to unwind after the stop request before the worker is
/// abandoned. See <see cref="Configuration.WorkerSettings.StopGraceSeconds"/>.
/// </summary>
public TimeSpan StopGrace { get; set; } = TimeSpan.FromSeconds(30);

/// <summary>
/// True once an invocation failed to stop within <see cref="StopGrace"/>. Its pipeline thread is
/// blocked in a call PowerShell cannot interrupt; the worker never runs again and the pool replaces it.
/// </summary>
public bool IsHung => _hungRun != null;

/// <summary>The abandoned pipeline of a hung worker; completes if the blocked call ever returns.</summary>
internal Task? HungRun => _hungRun;

private volatile Task? _hungRun;
private int _disposed;

private void ThrowIfHung(string function)
{
if (IsHung)
throw new WorkerHungException(Id, $"Worker W{Id} is hung and cannot run {function}");
}

/// <summary>
/// Run the queued commands. On cancellation the stop runs on a thread of its own, because
/// <see cref="PowerShell.Stop"/> blocks until the pipeline unwinds — which a pipeline blocked in .NET
/// never does. Called from the cancellation callback, as it used to be, it never returned, so the
/// timeout was never reported and the worker never came back to the pool.
/// </summary>
private async Task<PSDataCollection<PSObject>> RunPipelineAsync(string label, CancellationToken ct)
{
// A fresh output collection per invocation: after an invocation that threw, the PowerShell object's
// own output buffer comes back null from the next EndInvoke, silently dropping that call's output.
var outputs = new PSDataCollection<PSObject>();
var run = Task.Factory.FromAsync(_pwsh.BeginInvoke<PSObject, PSObject>(null, outputs), _pwsh.EndInvoke);
Task? stop = null;
try
{
if (ct.CanBeCanceled && !run.IsCompleted)
{
var cancelled = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously);
using (ct.Register(() => cancelled.TrySetResult()))
await Task.WhenAny(run, cancelled.Task);

if (!run.IsCompleted)
{
stop = Task.Factory.StartNew(StopPipeline, CancellationToken.None,
TaskCreationOptions.LongRunning, TaskScheduler.Default);
if (await Task.WhenAny(run, Task.Delay(StopGrace, CancellationToken.None)) != run)
{
Abandon(run, outputs, label);
throw new WorkerHungException(Id,
$"Worker W{Id} running {label} did not stop within {StopGrace.TotalSeconds:0}s of being cancelled and was abandoned");
}
}
}

await run;
ct.ThrowIfCancellationRequested();
return outputs;
}
catch (Exception ex) when (ex is not WorkerHungException)
{
outputs.Dispose();
throw;
}
finally
{
// The pipeline has finished, so a stop still in flight returns at once. Wait for it, or it could
// land on the cleanup pipelines that run next.
if (stop != null && !IsHung) await stop;
}
}

private void StopPipeline()
{
try { _pwsh.Stop(); }
catch (Exception ex) { _logger.LogDebug(ex, "Worker{Id}: pipeline stop threw", Id); }
}

private void Abandon(Task run, PSDataCollection<PSObject> outputs, string label)
{
var abandonedAt = System.Diagnostics.Stopwatch.GetTimestamp();
_hungRun = run;
_logger.LogError(
"[Pool] W{Id} {Function} did not stop within {Grace:0}s of being cancelled: it is blocked in a call " +
"PowerShell cannot interrupt. Abandoning the worker.", Id, label, StopGrace.TotalSeconds);
_ = run.ContinueWith(t =>
{
_ = t.Exception;
outputs.Dispose();
_logger.LogWarning("[Pool] Abandoned W{Id} {Function} returned {Seconds:0}s after it was abandoned",
Id, label, System.Diagnostics.Stopwatch.GetElapsedTime(abandonedAt).TotalSeconds);
}, CancellationToken.None, TaskContinuationOptions.None, TaskScheduler.Default);
}

public PSDataStreams Streams => _pwsh.Streams;

/// <summary>
Expand Down Expand Up @@ -656,6 +739,22 @@ public ExportedModuleState ExportModuleState()

public void Dispose()
{
// Nothing here owns unmanaged resources directly, but suppressing finalization keeps a
// derived type that adds a finalizer from having to re-implement IDisposable to do it.
GC.SuppressFinalize(this);
if (_hungRun is { IsCompleted: false } hung)
{
// Closing the runspace would wait on the blocked pipeline, possibly forever.
_ = hung.ContinueWith(_ => DisposeCore(), CancellationToken.None,
TaskContinuationOptions.None, TaskScheduler.Default);
return;
}
DisposeCore();
}

private void DisposeCore()
{
if (Interlocked.Exchange(ref _disposed, 1) == 1) return;
// PowerShell.Create(iss) ASSIGNS the runspace rather than creating it lazily, and an assigned
// runspace is caller-owned — _pwsh.Dispose() does not close it. Left open, the runspace keeps
// its ReuseThread pipeline thread alive, and a live thread roots the entire session state
Expand All @@ -664,8 +763,5 @@ public void Dispose()
var runspace = _pwsh.Runspace;
_pwsh.Dispose();
runspace?.Dispose();
// Nothing else here owns unmanaged resources directly, but suppressing finalization keeps a
// derived type that adds a finalizer from having to re-implement IDisposable to do it.
GC.SuppressFinalize(this);
}
}
Loading
Loading