fix(siteruntime): injectable ScriptExecutionScheduler + un-leakable expression-fault test

Gitea #18: the process-wide ScriptExecutionScheduler was reached via a static
singleton (Shared) directly inside ScriptActor, ScriptExecutionActor, AlarmActor,
AlarmExecutionActor, and ScriptSchedulerStatsReporter, so it could not be
substituted — a shared mutable global for anything running more than one logical
site in a process, most visibly the test assembly. Add an optional scheduler
injection seam to all five: the Host injects nothing and gets the shared pool
(byte-for-byte unchanged), while a test (or a future multi-site host) hands an
actor its own instance. Spawning actors thread their scheduler down to the
execution children they create. Shared() now also recreates a disposed cached
instance rather than returning it (new IsDisposed guard), so a disposed scheduler
can never silently poison every later script/alarm execution.

Gitea #16: ScriptActorTests now runs on its OWN scheduler (disposed at class
teardown), so the faulted-task test can no longer strand a worker on the
process-wide pool and starve unrelated test classes (one prior run failed 22
tests). EvalGate's wait is bounded to 30 s (was unbounded — the actual leak
mechanism) and reset per test; the teardown releases one permit per worker
instead of Release(Entries), which under-released in exactly the failure case.
Full SiteRuntime suite: 6/6 runs green (524/524); was 3/6 failing on clean main.

Fixes: #18
Fixes: #16
This commit was merged in pull request #21.
This commit is contained in:
Joseph Doherty
2026-07-17 15:36:39 -04:00
parent 3e84eee195
commit b1b4874090
9 changed files with 205 additions and 37 deletions
@@ -216,6 +216,7 @@ When the Instance Actor is stopped (due to disable, delete, or redeployment), Ak
- Executes the script in the Akka actor context. - Executes the script in the Akka actor context.
- Has access to the full Script Runtime API (see below). - Has access to the full Script Runtime API (see below).
- Returns the script's return value (if defined) to the caller, then stops. - Returns the script's return value (if defined) to the caller, then stops.
- The script body itself runs on the **dedicated `ScriptExecutionScheduler`** (a bounded set of dedicated threads), not the shared .NET thread pool, so blocking script I/O cannot starve the global pool or stall Akka dispatchers. The scheduler is **process-wide by default** (one pool per host, sized from `SiteRuntimeOptions.ScriptExecutionThreadCount`), but each script/alarm actor takes it through an **optional injection seam** rather than reaching for the static directly: the Host injects nothing and gets the shared pool, while tests (or a future multi-site host) can hand an actor its own instance. The shared accessor also recreates a disposed pool rather than returning it, so a disposed scheduler can never silently poison later executions.
### Handling `Instance.CallScript` ### Handling `Instance.CallScript`
- When an external caller (another Script Execution Actor, an Alarm Execution Actor, or a routed call from the Inbound API) sends a `CallScript` message to the Script Actor, it spawns a Script Execution Actor to handle the call. - When an external caller (another Script Execution Actor, an Alarm Execution Actor, or a routed call from the Inbound API) sends a `CallScript` message to the Script Actor, it spawns a Script Execution Actor to handle the call.
@@ -41,6 +41,15 @@ public class AlarmActor : ReceiveActor
private readonly ISiteHealthCollector? _healthCollector; private readonly ISiteHealthCollector? _healthCollector;
private readonly IServiceProvider? _serviceProvider; private readonly IServiceProvider? _serviceProvider;
/// <summary>
/// Script-execution scheduler seam (#18): the process-wide
/// <see cref="ScriptExecutionScheduler"/> when null, or an injected instance so this
/// alarm's trigger-expression evaluation and spawned on-trigger scripts run on a
/// caller-owned pool. Resolved lazily at each use so the null (host) path is
/// unchanged.
/// </summary>
private readonly ScriptExecutionScheduler? _scheduler;
/// <summary> /// <summary>
/// The optional site operational-event log, resolved once from /// The optional site operational-event log, resolved once from
/// <see cref="_serviceProvider"/> at construction and cached. The /// <see cref="_serviceProvider"/> at construction and cached. The
@@ -125,6 +134,7 @@ public class AlarmActor : ReceiveActor
/// <param name="onTriggerExecutionTimeoutSeconds">The on-trigger script's per-script /// <param name="onTriggerExecutionTimeoutSeconds">The on-trigger script's per-script
/// execution timeout in seconds (from its <see cref="ResolvedScript.ExecutionTimeoutSeconds"/>), /// execution timeout in seconds (from its <see cref="ResolvedScript.ExecutionTimeoutSeconds"/>),
/// or null/non-positive to use the global default.</param> /// or null/non-positive to use the global default.</param>
/// <param name="scheduler">Optional script-execution scheduler override (#18); null uses the process-wide shared scheduler.</param>
public AlarmActor( public AlarmActor(
string alarmName, string alarmName,
string instanceName, string instanceName,
@@ -139,7 +149,9 @@ public class AlarmActor : ReceiveActor
ISiteHealthCollector? healthCollector = null, ISiteHealthCollector? healthCollector = null,
IServiceProvider? serviceProvider = null, IServiceProvider? serviceProvider = null,
// Per-script timeout for the on-trigger script (null = global). // Per-script timeout for the on-trigger script (null = global).
int? onTriggerExecutionTimeoutSeconds = null) int? onTriggerExecutionTimeoutSeconds = null,
// Script-execution scheduler seam (#18); null uses the process-wide shared scheduler.
ScriptExecutionScheduler? scheduler = null)
{ {
_alarmName = alarmName; _alarmName = alarmName;
_instanceName = instanceName; _instanceName = instanceName;
@@ -149,6 +161,7 @@ public class AlarmActor : ReceiveActor
_logger = logger; _logger = logger;
_healthCollector = healthCollector; _healthCollector = healthCollector;
_serviceProvider = serviceProvider; _serviceProvider = serviceProvider;
_scheduler = scheduler;
// Resolve the optional site event logger once and cache it, // Resolve the optional site event logger once and cache it,
// rather than calling GetService on every alarm transition. // rather than calling GetService on every alarm transition.
_siteEventLogger = serviceProvider?.GetService<ISiteEventLogger>(); _siteEventLogger = serviceProvider?.GetService<ISiteEventLogger>();
@@ -554,7 +567,7 @@ public class AlarmActor : ReceiveActor
return false; return false;
} }
}, CancellationToken.None, TaskCreationOptions.DenyChildAttach, }, CancellationToken.None, TaskCreationOptions.DenyChildAttach,
ScriptExecutionScheduler.Shared(_options)).Unwrap().PipeTo(self, _scheduler ?? ScriptExecutionScheduler.Shared(_options)).Unwrap().PipeTo(self,
success: r => new ExpressionEvalResult(r), success: r => new ExpressionEvalResult(r),
failure: ex => new ExpressionEvalFailed(ex)); failure: ex => new ExpressionEvalFailed(ex));
} }
@@ -685,7 +698,10 @@ public class AlarmActor : ReceiveActor
// Per-script timeout from the on-trigger script (null = global). // Per-script timeout from the on-trigger script (null = global).
_onTriggerExecutionTimeoutSeconds, _onTriggerExecutionTimeoutSeconds,
// The firing context's execution id (null today). // The firing context's execution id (null today).
parentExecutionId)); parentExecutionId,
// Scheduler seam (#18): share this alarm's scheduler override with the
// spawned on-trigger script body (null = process-wide shared).
_scheduler));
Context.ActorOf(props, executionId); Context.ActorOf(props, executionId);
} }
@@ -52,7 +52,10 @@ public class AlarmExecutionActor : ReceiveActor
// alarm on-trigger script. Null or non-positive falls back to the global. // alarm on-trigger script. Null or non-positive falls back to the global.
int? executionTimeoutSeconds = null, int? executionTimeoutSeconds = null,
// The firing context's execution id (null today). // The firing context's execution id (null today).
Guid? parentExecutionId = null) Guid? parentExecutionId = null,
// Script-execution scheduler seam (#18): the process-wide scheduler by
// default; null selects the shared default.
ScriptExecutionScheduler? scheduler = null)
{ {
var self = Self; var self = Self;
var parent = Context.Parent; var parent = Context.Parent;
@@ -61,7 +64,7 @@ public class AlarmExecutionActor : ReceiveActor
alarmName, instanceName, level, priority, message, alarmName, instanceName, level, priority, message,
compiledScript, instanceActor, compiledScript, instanceActor,
sharedScriptLibrary, options, self, parent, logger, sharedScriptLibrary, options, self, parent, logger,
executionTimeoutSeconds, parentExecutionId); executionTimeoutSeconds, parentExecutionId, scheduler);
} }
private static void ExecuteAlarmScript( private static void ExecuteAlarmScript(
@@ -78,7 +81,8 @@ public class AlarmExecutionActor : ReceiveActor
IActorRef parent, IActorRef parent,
ILogger logger, ILogger logger,
int? executionTimeoutSeconds, int? executionTimeoutSeconds,
Guid? parentExecutionId) Guid? parentExecutionId,
ScriptExecutionScheduler? scheduler)
{ {
// Per-script timeout overrides the global default. A null or // Per-script timeout overrides the global default. A null or
// non-positive per-script value (≤ 0) falls back to the global. // non-positive per-script value (≤ 0) falls back to the global.
@@ -88,8 +92,9 @@ public class AlarmExecutionActor : ReceiveActor
: options.ScriptExecutionTimeoutSeconds); : options.ScriptExecutionTimeoutSeconds);
// Run the alarm on-trigger body on the dedicated // Run the alarm on-trigger body on the dedicated
// script-execution scheduler, not the shared .NET thread pool. // script-execution scheduler, not the shared .NET thread pool. An injected
var scheduler = ScriptExecutionScheduler.Shared(options); // scheduler (#18) overrides the process-wide default.
var executionScheduler = scheduler ?? ScriptExecutionScheduler.Shared(options);
_ = Task.Factory.StartNew(async () => _ = Task.Factory.StartNew(async () =>
{ {
@@ -157,6 +162,6 @@ public class AlarmExecutionActor : ReceiveActor
{ {
self.Tell(PoisonPill.Instance); self.Tell(PoisonPill.Instance);
} }
}, CancellationToken.None, TaskCreationOptions.DenyChildAttach, scheduler).Unwrap(); }, CancellationToken.None, TaskCreationOptions.DenyChildAttach, executionScheduler).Unwrap();
} }
} }
@@ -40,6 +40,15 @@ public class ScriptActor : ReceiveActor, IWithTimers
private readonly ISiteHealthCollector? _healthCollector; private readonly ISiteHealthCollector? _healthCollector;
private readonly IServiceProvider? _serviceProvider; private readonly IServiceProvider? _serviceProvider;
/// <summary>
/// Script-execution scheduler seam (#18): the process-wide
/// <see cref="ScriptExecutionScheduler"/> when null, or an injected instance so
/// this actor's trigger-expression evaluation and spawned script bodies run on a
/// caller-owned pool instead of the shared one. Resolved lazily at each use so the
/// null (host) path stays byte-for-byte identical to the previous static call.
/// </summary>
private readonly ScriptExecutionScheduler? _scheduler;
private Script<object?>? _compiledScript; private Script<object?>? _compiledScript;
private ScriptTriggerConfig? _triggerConfig; private ScriptTriggerConfig? _triggerConfig;
private TimeSpan? _minTimeBetweenRuns; private TimeSpan? _minTimeBetweenRuns;
@@ -101,6 +110,7 @@ public class ScriptActor : ReceiveActor, IWithTimers
/// <param name="initialAttributes">Initial attribute snapshot used to seed expression trigger evaluation state.</param> /// <param name="initialAttributes">Initial attribute snapshot used to seed expression trigger evaluation state.</param>
/// <param name="healthCollector">Optional health metrics collector.</param> /// <param name="healthCollector">Optional health metrics collector.</param>
/// <param name="serviceProvider">Optional DI service provider for script execution context services.</param> /// <param name="serviceProvider">Optional DI service provider for script execution context services.</param>
/// <param name="scheduler">Optional script-execution scheduler override (#18); null uses the process-wide shared scheduler.</param>
public ScriptActor( public ScriptActor(
string scriptName, string scriptName,
string instanceName, string instanceName,
@@ -113,7 +123,8 @@ public class ScriptActor : ReceiveActor, IWithTimers
Script<object?>? compiledTriggerExpression = null, Script<object?>? compiledTriggerExpression = null,
IReadOnlyDictionary<string, object?>? initialAttributes = null, IReadOnlyDictionary<string, object?>? initialAttributes = null,
ISiteHealthCollector? healthCollector = null, ISiteHealthCollector? healthCollector = null,
IServiceProvider? serviceProvider = null) IServiceProvider? serviceProvider = null,
ScriptExecutionScheduler? scheduler = null)
{ {
_scriptName = scriptName; _scriptName = scriptName;
_instanceName = instanceName; _instanceName = instanceName;
@@ -124,6 +135,7 @@ public class ScriptActor : ReceiveActor, IWithTimers
_logger = logger; _logger = logger;
_healthCollector = healthCollector; _healthCollector = healthCollector;
_serviceProvider = serviceProvider; _serviceProvider = serviceProvider;
_scheduler = scheduler;
_minTimeBetweenRuns = scriptConfig.MinTimeBetweenRuns; _minTimeBetweenRuns = scriptConfig.MinTimeBetweenRuns;
_executionTimeoutSeconds = scriptConfig.ExecutionTimeoutSeconds; _executionTimeoutSeconds = scriptConfig.ExecutionTimeoutSeconds;
_scope = scriptConfig.Scope; _scope = scriptConfig.Scope;
@@ -313,7 +325,7 @@ public class ScriptActor : ReceiveActor, IWithTimers
return false; return false;
} }
}, CancellationToken.None, TaskCreationOptions.DenyChildAttach, }, CancellationToken.None, TaskCreationOptions.DenyChildAttach,
ScriptExecutionScheduler.Shared(_options)).Unwrap().PipeTo(self, _scheduler ?? ScriptExecutionScheduler.Shared(_options)).Unwrap().PipeTo(self,
success: r => new ExpressionEvalResult(r), success: r => new ExpressionEvalResult(r),
failure: ex => new ExpressionEvalFailed(ex)); failure: ex => new ExpressionEvalFailed(ex));
} }
@@ -480,7 +492,10 @@ public class ScriptActor : ReceiveActor, IWithTimers
// an inbound-API-routed call supplies the inbound request's id. // an inbound-API-routed call supplies the inbound request's id.
parentExecutionId, parentExecutionId,
// Per-script timeout override (null = use global). // Per-script timeout override (null = use global).
_executionTimeoutSeconds)); _executionTimeoutSeconds,
// Scheduler seam (#18): thread this actor's scheduler override down so
// spawned script bodies share the same pool (null = process-wide shared).
_scheduler));
Context.ActorOf(props, executionId); Context.ActorOf(props, executionId);
} }
@@ -69,7 +69,12 @@ public class ScriptExecutionActor : ReceiveActor
Guid? parentExecutionId = null, Guid? parentExecutionId = null,
// Per-script execution timeout override (seconds). Null or // Per-script execution timeout override (seconds). Null or
// non-positive falls back to the global ScriptExecutionTimeoutSeconds. // non-positive falls back to the global ScriptExecutionTimeoutSeconds.
int? executionTimeoutSeconds = null) int? executionTimeoutSeconds = null,
// Script-execution scheduler seam (#18): the process-wide
// ScriptExecutionScheduler by default; a test (or a future multi-site
// host) can inject its own instance so script bodies never run on the
// shared process-wide pool. Null selects the shared default.
ScriptExecutionScheduler? scheduler = null)
{ {
// Immediately begin execution // Immediately begin execution
var self = Self; var self = Self;
@@ -79,7 +84,7 @@ public class ScriptExecutionActor : ReceiveActor
scriptName, instanceName, compiledScript, parameters, callDepth, scriptName, instanceName, compiledScript, parameters, callDepth,
instanceActor, sharedScriptLibrary, options, replyTo, correlationId, instanceActor, sharedScriptLibrary, options, replyTo, correlationId,
self, parent, logger, scope, healthCollector, serviceProvider, self, parent, logger, scope, healthCollector, serviceProvider,
parentExecutionId, executionTimeoutSeconds); parentExecutionId, executionTimeoutSeconds, scheduler);
} }
private static void ExecuteScript( private static void ExecuteScript(
@@ -100,7 +105,8 @@ public class ScriptExecutionActor : ReceiveActor
ISiteHealthCollector? healthCollector, ISiteHealthCollector? healthCollector,
IServiceProvider? serviceProvider, IServiceProvider? serviceProvider,
Guid? parentExecutionId, Guid? parentExecutionId,
int? executionTimeoutSeconds) int? executionTimeoutSeconds,
ScriptExecutionScheduler? scheduler)
{ {
// Per-script timeout overrides the global default. A null or // Per-script timeout overrides the global default. A null or
// non-positive per-script value (≤ 0) falls back to the global. // non-positive per-script value (≤ 0) falls back to the global.
@@ -111,8 +117,9 @@ public class ScriptExecutionActor : ReceiveActor
// Run the script body on the dedicated script-execution // Run the script body on the dedicated script-execution
// scheduler, not the shared .NET thread pool, so blocking script I/O cannot // scheduler, not the shared .NET thread pool, so blocking script I/O cannot
// starve the global pool and stall Akka dispatchers / HTTP handling. // starve the global pool and stall Akka dispatchers / HTTP handling. An
var scheduler = ScriptExecutionScheduler.Shared(options); // injected scheduler (#18) overrides the process-wide default.
var executionScheduler = scheduler ?? ScriptExecutionScheduler.Shared(options);
// Notification Outbox: the site communication actor that Notify.Status queries // Notification Outbox: the site communication actor that Notify.Status queries
// central through. Resolved by actor path so the Notify helper does not need an // central through. Resolved by actor path so the Notify helper does not need an
@@ -329,6 +336,6 @@ public class ScriptExecutionActor : ReceiveActor
// Stop self after execution completes // Stop self after execution completes
self.Tell(PoisonPill.Instance); self.Tell(PoisonPill.Instance);
} }
}, CancellationToken.None, TaskCreationOptions.DenyChildAttach, scheduler).Unwrap(); }, CancellationToken.None, TaskCreationOptions.DenyChildAttach, executionScheduler).Unwrap();
} }
} }
@@ -32,20 +32,30 @@ public sealed class ScriptExecutionScheduler : TaskScheduler, IDisposable
private static readonly object SharedLock = new(); private static readonly object SharedLock = new();
/// <summary> /// <summary>
/// The process-wide script-execution scheduler. Lazily created on first use with the /// The process-wide script-execution scheduler, used as the default when no scheduler
/// thread count from <see cref="SiteRuntimeOptions.ScriptExecutionThreadCount"/>; the /// is injected. Lazily created on first use with the thread count from
/// first caller wins, subsequent calls reuse the existing instance. /// <see cref="SiteRuntimeOptions.ScriptExecutionThreadCount"/>; the first caller wins,
/// subsequent calls reuse the existing instance.
///
/// If the cached instance has been disposed it is recreated rather than handed back:
/// a disposed scheduler can execute no work, so returning it would silently poison
/// every subsequent script and alarm execution for the rest of the process. Nothing
/// disposes the shared instance today, but the guard keeps that latent failure mode
/// closed regardless of a future shutdown path or a multi-site host.
/// </summary> /// </summary>
/// <param name="options">Site runtime options supplying the thread count for the scheduler.</param> /// <param name="options">Site runtime options supplying the thread count for the scheduler.</param>
/// <returns>The process-wide <see cref="ScriptExecutionScheduler"/> instance, creating it on first call.</returns> /// <returns>A live process-wide <see cref="ScriptExecutionScheduler"/> instance, creating it on first call.</returns>
public static ScriptExecutionScheduler Shared(SiteRuntimeOptions options) public static ScriptExecutionScheduler Shared(SiteRuntimeOptions options)
{ {
if (_shared != null) var existing = _shared;
return _shared; if (existing is { IsDisposed: false })
return existing;
lock (SharedLock) lock (SharedLock)
{ {
return _shared ??= new ScriptExecutionScheduler(options.ScriptExecutionThreadCount); if (_shared is null || _shared.IsDisposed)
_shared = new ScriptExecutionScheduler(options.ScriptExecutionThreadCount);
return _shared;
} }
} }
@@ -73,6 +83,10 @@ public sealed class ScriptExecutionScheduler : TaskScheduler, IDisposable
} }
} }
/// <summary><see langword="true"/> once <see cref="Dispose"/> has run; a disposed scheduler
/// can execute no further work and must not be handed out by <see cref="Shared"/>.</summary>
public bool IsDisposed => Volatile.Read(ref _disposed) != 0;
/// <inheritdoc /> /// <inheritdoc />
public override int MaximumConcurrencyLevel => _threads.Count; public override int MaximumConcurrencyLevel => _threads.Count;
@@ -24,21 +24,30 @@ public sealed class ScriptSchedulerStatsReporter : BackgroundService
private readonly ILogger<ScriptSchedulerStatsReporter> _logger; private readonly ILogger<ScriptSchedulerStatsReporter> _logger;
private readonly TimeSpan _pollInterval; private readonly TimeSpan _pollInterval;
/// <summary>
/// Script-execution scheduler seam (#18): the process-wide scheduler when null, or
/// an injected instance so the reporter samples the same pool the actors run on.
/// </summary>
private readonly ScriptExecutionScheduler? _scheduler;
/// <summary>Initializes a new instance of <see cref="ScriptSchedulerStatsReporter"/>.</summary> /// <summary>Initializes a new instance of <see cref="ScriptSchedulerStatsReporter"/>.</summary>
/// <param name="collector">The site health collector that receives the scheduler gauges.</param> /// <param name="collector">The site health collector that receives the scheduler gauges.</param>
/// <param name="options">Site runtime options supplying the shared scheduler's thread count.</param> /// <param name="options">Site runtime options supplying the shared scheduler's thread count.</param>
/// <param name="logger">Logger instance.</param> /// <param name="logger">Logger instance.</param>
/// <param name="pollInterval">Poll interval override; defaults to <see cref="DefaultPollInterval"/> (10 s).</param> /// <param name="pollInterval">Poll interval override; defaults to <see cref="DefaultPollInterval"/> (10 s).</param>
/// <param name="scheduler">Optional script-execution scheduler override (#18); null samples the process-wide shared scheduler.</param>
public ScriptSchedulerStatsReporter( public ScriptSchedulerStatsReporter(
ISiteHealthCollector collector, ISiteHealthCollector collector,
SiteRuntimeOptions options, SiteRuntimeOptions options,
ILogger<ScriptSchedulerStatsReporter> logger, ILogger<ScriptSchedulerStatsReporter> logger,
TimeSpan? pollInterval = null) TimeSpan? pollInterval = null,
ScriptExecutionScheduler? scheduler = null)
{ {
_collector = collector ?? throw new ArgumentNullException(nameof(collector)); _collector = collector ?? throw new ArgumentNullException(nameof(collector));
_options = options ?? throw new ArgumentNullException(nameof(options)); _options = options ?? throw new ArgumentNullException(nameof(options));
_logger = logger ?? throw new ArgumentNullException(nameof(logger)); _logger = logger ?? throw new ArgumentNullException(nameof(logger));
_pollInterval = pollInterval ?? DefaultPollInterval; _pollInterval = pollInterval ?? DefaultPollInterval;
_scheduler = scheduler;
} }
/// <inheritdoc /> /// <inheritdoc />
@@ -67,7 +76,7 @@ public sealed class ScriptSchedulerStatsReporter : BackgroundService
{ {
try try
{ {
var scheduler = ScriptExecutionScheduler.Shared(_options); var scheduler = _scheduler ?? ScriptExecutionScheduler.Shared(_options);
_collector.SetScriptSchedulerStats( _collector.SetScriptSchedulerStats(
scheduler.QueueDepth, scheduler.QueueDepth,
scheduler.BusyThreadCount, scheduler.BusyThreadCount,
@@ -26,6 +26,14 @@ public class ScriptActorTests : TestKit, IDisposable
private readonly SiteRuntimeOptions _options; private readonly SiteRuntimeOptions _options;
private readonly ScriptCompilationService _compilationService; private readonly ScriptCompilationService _compilationService;
// Gitea #16/#18: every ScriptActor built here runs its trigger-expression
// evaluation and spawned script bodies on THIS class's own scheduler, never the
// process-wide singleton. That is what makes the faulted-task test below
// un-leakable across the assembly: any worker a test strands (or the 30 s bounded
// EvalGate wait leaves parked) lives on this instance, which is disposed at class
// teardown, so it can never starve unrelated test classes' script executions.
private readonly ScriptExecutionScheduler _scheduler;
public ScriptActorTests() public ScriptActorTests()
{ {
_compilationService = new ScriptCompilationService( _compilationService = new ScriptCompilationService(
@@ -37,11 +45,13 @@ public class ScriptActorTests : TestKit, IDisposable
MaxScriptCallDepth = 10, MaxScriptCallDepth = 10,
ScriptExecutionTimeoutSeconds = 30 ScriptExecutionTimeoutSeconds = 30
}; };
_scheduler = new ScriptExecutionScheduler(_options.ScriptExecutionThreadCount);
} }
void IDisposable.Dispose() void IDisposable.Dispose()
{ {
Shutdown(); Shutdown();
_scheduler.Dispose();
} }
private Script<object?> CompileScript(string code) private Script<object?> CompileScript(string code)
@@ -74,7 +84,12 @@ public class ScriptActorTests : TestKit, IDisposable
scriptConfig, scriptConfig,
_sharedLibrary, _sharedLibrary,
_options, _options,
NullLogger<ScriptActor>.Instance))); NullLogger<ScriptActor>.Instance,
null, // compiledTriggerExpression
null, // initialAttributes
null, // healthCollector
null, // serviceProvider
_scheduler)));
// Ask pattern (WP-22) for CallScript // Ask pattern (WP-22) for CallScript
scriptActor.Tell(new ScriptCallRequest("GetAnswer", null, 0, "corr-1")); scriptActor.Tell(new ScriptCallRequest("GetAnswer", null, 0, "corr-1"));
@@ -103,7 +118,12 @@ public class ScriptActorTests : TestKit, IDisposable
scriptConfig, scriptConfig,
_sharedLibrary, _sharedLibrary,
_options, _options,
NullLogger<ScriptActor>.Instance))); NullLogger<ScriptActor>.Instance,
null, // compiledTriggerExpression
null, // initialAttributes
null, // healthCollector
null, // serviceProvider
_scheduler)));
var parameters = new Dictionary<string, object?> { ["x"] = 3, ["y"] = 4 }; var parameters = new Dictionary<string, object?> { ["x"] = 3, ["y"] = 4 };
scriptActor.Tell(new ScriptCallRequest("Add", parameters, 0, "corr-2")); scriptActor.Tell(new ScriptCallRequest("Add", parameters, 0, "corr-2"));
@@ -131,7 +151,12 @@ public class ScriptActorTests : TestKit, IDisposable
scriptConfig, scriptConfig,
_sharedLibrary, _sharedLibrary,
_options, _options,
NullLogger<ScriptActor>.Instance))); NullLogger<ScriptActor>.Instance,
null, // compiledTriggerExpression
null, // initialAttributes
null, // healthCollector
null, // serviceProvider
_scheduler)));
scriptActor.Tell(new ScriptCallRequest("Broken", null, 0, "corr-3")); scriptActor.Tell(new ScriptCallRequest("Broken", null, 0, "corr-3"));
@@ -161,7 +186,12 @@ public class ScriptActorTests : TestKit, IDisposable
scriptConfig, scriptConfig,
_sharedLibrary, _sharedLibrary,
_options, _options,
NullLogger<ScriptActor>.Instance))); NullLogger<ScriptActor>.Instance,
null, // compiledTriggerExpression
null, // initialAttributes
null, // healthCollector
null, // serviceProvider
_scheduler)));
// Send an attribute change that matches the trigger // Send an attribute change that matches the trigger
scriptActor.Tell(new AttributeValueChanged( scriptActor.Tell(new AttributeValueChanged(
@@ -194,7 +224,12 @@ public class ScriptActorTests : TestKit, IDisposable
scriptConfig, scriptConfig,
_sharedLibrary, _sharedLibrary,
_options, _options,
NullLogger<ScriptActor>.Instance))); NullLogger<ScriptActor>.Instance,
null, // compiledTriggerExpression
null, // initialAttributes
null, // healthCollector
null, // serviceProvider
_scheduler)));
// First trigger -- should execute // First trigger -- should execute
scriptActor.Tell(new AttributeValueChanged( scriptActor.Tell(new AttributeValueChanged(
@@ -228,7 +263,12 @@ public class ScriptActorTests : TestKit, IDisposable
scriptConfig, scriptConfig,
_sharedLibrary, _sharedLibrary,
_options, _options,
NullLogger<ScriptActor>.Instance))); NullLogger<ScriptActor>.Instance,
null, // compiledTriggerExpression
null, // initialAttributes
null, // healthCollector
null, // serviceProvider
_scheduler)));
// First call -- fails // First call -- fails
scriptActor.Tell(new ScriptCallRequest("Failing", null, 0, "corr-fail-1")); scriptActor.Tell(new ScriptCallRequest("Failing", null, 0, "corr-fail-1"));
@@ -293,7 +333,8 @@ public class ScriptActorTests : TestKit, IDisposable
triggerExpression, triggerExpression,
null, null,
null, null,
null))); null,
_scheduler)));
return (actor, instance); return (actor, instance);
} }
@@ -355,6 +396,7 @@ public class ScriptActorTests : TestKit, IDisposable
[Fact] [Fact]
public void ExpressionEvalTaskFault_ClearsInFlight_AndDrainsPendingEvaluation() public void ExpressionEvalTaskFault_ClearsInFlight_AndDrainsPendingEvaluation()
{ {
EvalGate.Reset(); // fresh gate + zeroed counter; never inherit a prior run's state
var expr = CompileRawTriggerExpression( var expr = CompileRawTriggerExpression(
"ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests.Actors.EvalGate.Block()"); "ZB.MOM.WW.ScadaBridge.SiteRuntime.Tests.Actors.EvalGate.Block()");
var (actor, _) = CreateTriggeredActor( var (actor, _) = CreateTriggeredActor(
@@ -363,6 +405,12 @@ public class ScriptActorTests : TestKit, IDisposable
{ {
actor.Tell(Change("A", "1")); // eval starts on the scheduler and BLOCKS → _evalInFlight = true actor.Tell(Change("A", "1")); // eval starts on the scheduler and BLOCKS → _evalInFlight = true
AwaitAssert(() => Assert.Equal(1, EvalGate.Entries), TimeSpan.FromSeconds(10)); AwaitAssert(() => Assert.Equal(1, EvalGate.Entries), TimeSpan.FromSeconds(10));
// #18 seam: the blocked evaluation is running on THIS class's injected
// scheduler — not the process-wide singleton — so a worker it strands can
// never starve another test class.
Assert.Equal(1, _scheduler.BusyThreadCount);
actor.Tell(Change("A", "2")); // coalesces → _evalPending = true actor.Tell(Change("A", "2")); // coalesces → _evalPending = true
// Simulate the faulted scheduler task the PipeTo failure mapping now surfaces // Simulate the faulted scheduler task the PipeTo failure mapping now surfaces
@@ -374,7 +422,12 @@ public class ScriptActorTests : TestKit, IDisposable
} }
finally finally
{ {
EvalGate.Gate.Release(EvalGate.Entries); // ALWAYS free the shared scheduler threads // Release generously (one per worker) so every parked/queued Block returns
// at once and this class's scheduler drains promptly at teardown. The old
// Release(Entries) under-released in exactly the failure case (Entries read
// as 1 while a second eval was still queued), stranding a worker; the
// bounded 30 s wait in Block is the ultimate backstop.
EvalGate.Gate.Release(_options.ScriptExecutionThreadCount);
} }
} }
@@ -541,8 +594,38 @@ public class ScriptActorTests : TestKit, IDisposable
/// </summary> /// </summary>
public static class EvalGate public static class EvalGate
{ {
public static readonly SemaphoreSlim Gate = new(0); private static SemaphoreSlim _gate = new(0);
private static int _entries; private static int _entries;
/// <summary>The semaphore that script-scheduler threads park on inside <see cref="Block"/>.</summary>
public static SemaphoreSlim Gate => _gate;
/// <summary>Number of evaluations that have entered <see cref="Block"/> since the last <see cref="Reset"/>.</summary>
public static int Entries => Volatile.Read(ref _entries); public static int Entries => Volatile.Read(ref _entries);
public static bool Block() { Interlocked.Increment(ref _entries); Gate.Wait(); return false; }
/// <summary>
/// Resets the gate to a fresh, permit-free state and zeroes the entry counter.
/// Called at the start of the test so counts never carry across runs — the old
/// static leaked <c>_entries</c> from whatever compiled an expression before it
/// (Gitea #16).
/// </summary>
public static void Reset()
{
_gate = new SemaphoreSlim(0);
Volatile.Write(ref _entries, 0);
}
/// <summary>
/// Counts the entry, then blocks the calling script-scheduler thread until a permit
/// is released. The wait is BOUNDED (30 s): even if a test releases too few permits,
/// a stranded evaluation frees its worker instead of parking it for the rest of the
/// process — the leak mechanism behind Gitea #16. Returns false so a trigger
/// expression built on it evaluates to "not fired".
/// </summary>
public static bool Block()
{
Interlocked.Increment(ref _entries);
_gate.Wait(TimeSpan.FromSeconds(30));
return false;
}
} }
@@ -45,6 +45,24 @@ public class ScriptExecutionSchedulerTests
Assert.Same(a, b); Assert.Same(a, b);
} }
// Gitea #18: IsDisposed backs the Shared() recreate guard so a disposed scheduler
// can never be handed back and silently poison every later execution. Verified on a
// local instance — deliberately NOT by disposing the process-wide singleton, which
// would strand any script another test class queued on it (the very static hazard
// #18 is about).
[Fact]
public void IsDisposed_FlipsOnDispose()
{
var scheduler = new ScriptExecutionScheduler(1);
Assert.False(scheduler.IsDisposed);
scheduler.Dispose();
Assert.True(scheduler.IsDisposed);
scheduler.Dispose(); // idempotent — stays disposed, does not throw
Assert.True(scheduler.IsDisposed);
}
// SiteRuntime-S2/UA5: the scheduler exposes observability gauges so a stuck / // SiteRuntime-S2/UA5: the scheduler exposes observability gauges so a stuck /
// saturated script-execution pool is visible to operators via site health. // saturated script-execution pool is visible to operators via site health.
[Fact] [Fact]