using System;
using System.Collections.Generic;
using System.Threading;
using System.Threading.Tasks;
using Junction.Core.Monitoring;
using Junction.Core.Plugins;
using Junction.Core.Polling;
using Junction.Domain;
using Junction.Domain.Models;
using Junction.Domain.Persistence;
using Junction.Domain.Protocols;
using Microsoft.Extensions.Logging.Abstractions;
using Moq;
using Xunit;
namespace Junction.Tests.Unit
{
///
/// Behavior tests for . Repository and plugin loader are mocked;
/// drivers/factories are hand-rolled fakes. The real drives the
/// loops. Timings are deliberately generous to avoid CI flake.
///
public class MachineMonitorTests
{
private const string PluginsDir = "/fake/plugins";
private static readonly TimeSpan FastPoll = TimeSpan.FromMilliseconds(25);
private static Machine MachineWith(string protocolId) =>
new Machine(Guid.NewGuid(), "M-" + protocolId, protocolId, null, FastPoll);
private static MachineSnapshot Snapshot(Guid machineId) =>
new MachineSnapshot(machineId, DateTimeOffset.UtcNow, ConnectionState.Connected, Array.Empty());
private static Func RealEngineFactory() =>
() => new PollingEngine(NullLogger.Instance);
private static MachineMonitor NewMonitor(IMachineRepository repo, IPluginLoader loader) =>
new MachineMonitor(repo, loader, RealEngineFactory(), NullLogger.Instance);
private static Mock LoaderReturning(params IProtocolDriverFactory[] factories)
{
var loaded = new List();
foreach (var f in factories)
{
var manifest = new PluginManifest(f.ProtocolId, f.ProtocolId, "x.dll", "X", "1.0");
loaded.Add(new LoadedPlugin(new PluginDescriptor(manifest, "/fake/x.dll"), f));
}
var mock = new Mock();
mock.Setup(l => l.LoadFrom(It.IsAny()))
.Returns(Result>.Ok(loaded));
return mock;
}
private static Mock RepoReturning(params Machine[] machines)
{
var mock = new Mock();
mock.Setup(r => r.GetAllAsync(It.IsAny()))
.ReturnsAsync(Result>.Ok(machines));
mock.Setup(r => r.SaveSnapshotAsync(It.IsAny(), It.IsAny()))
.ReturnsAsync(Result.Ok());
return mock;
}
[Fact]
public async Task StartAsync_PollsMachine_RaisesEvent_UpdatesCache_Persists()
{
var machine = MachineWith("test");
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(machine.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning(machine);
var monitor = NewMonitor(repo.Object, loader.Object);
var eventRaised = new ManualResetEventSlim(false);
monitor.SnapshotUpdated += (_, snap) =>
{
if (snap.MachineId == machine.Id) eventRaised.Set();
};
var start = await monitor.StartAsync(PluginsDir, CancellationToken.None);
Assert.True(start.IsSuccess);
Assert.True(eventRaised.Wait(TimeSpan.FromSeconds(2)), "SnapshotUpdated was not raised");
await monitor.StopAsync();
Assert.True(monitor.LatestSnapshots.ContainsKey(machine.Id));
Assert.Equal(machine.Id, monitor.LatestSnapshots[machine.Id].MachineId);
repo.Verify(r => r.SaveSnapshotAsync(
It.Is(s => s.MachineId == machine.Id),
It.IsAny()), Times.AtLeastOnce);
}
[Fact]
public async Task StartAsync_UnknownProtocol_SkipsMachine_OthersRun_StillOk()
{
var known = MachineWith("test");
var unknown = MachineWith("nope");
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(known.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning(unknown, known);
var monitor = NewMonitor(repo.Object, loader.Object);
var start = await monitor.StartAsync(PluginsDir, CancellationToken.None);
Assert.True(start.IsSuccess);
// Give the known machine time to produce.
var sw = System.Diagnostics.Stopwatch.StartNew();
while (!monitor.LatestSnapshots.ContainsKey(known.Id) && sw.Elapsed < TimeSpan.FromSeconds(2))
{
await Task.Delay(20);
}
await monitor.StopAsync();
Assert.True(monitor.LatestSnapshots.ContainsKey(known.Id), "known machine did not run");
Assert.False(monitor.LatestSnapshots.ContainsKey(unknown.Id), "unknown-protocol machine must be skipped");
}
[Fact]
public async Task StartAsync_DriverCreateFails_SkipsMachine_StillOk()
{
var machine = MachineWith("test");
var factory = new FakeFactory("test", m =>
Result.Fail(OperationError.Of("test", "bad config")));
var loader = LoaderReturning(factory);
var repo = RepoReturning(machine);
var monitor = NewMonitor(repo.Object, loader.Object);
var start = await monitor.StartAsync(PluginsDir, CancellationToken.None);
Assert.True(start.IsSuccess);
await monitor.StopAsync();
Assert.False(monitor.LatestSnapshots.ContainsKey(machine.Id));
}
[Fact]
public async Task StartAsync_LoaderFails_ReturnsFail()
{
var loader = new Mock();
loader.Setup(l => l.LoadFrom(It.IsAny()))
.Returns(Result>.Fail(
new OperationError("PLUGIN_DIR_MISSING", "PluginLoader", "no dir")));
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
var start = await monitor.StartAsync(PluginsDir, CancellationToken.None);
Assert.False(start.IsSuccess);
Assert.NotEmpty(start.Errors);
repo.Verify(r => r.GetAllAsync(It.IsAny()), Times.Never);
}
[Fact]
public async Task StartAsync_RepoGetAllFails_ReturnsFail()
{
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(Guid.NewGuid())))));
var loader = LoaderReturning(factory);
var repo = new Mock();
repo.Setup(r => r.GetAllAsync(It.IsAny()))
.ReturnsAsync(Result>.Fail(
new OperationError("STORE_UNAVAILABLE", "repo", "db down")));
var monitor = NewMonitor(repo.Object, loader.Object);
var start = await monitor.StartAsync(PluginsDir, CancellationToken.None);
Assert.False(start.IsSuccess);
Assert.NotEmpty(start.Errors);
}
[Fact]
public async Task StopAsync_CancelsLoops_Cleanly()
{
var machine = MachineWith("test");
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(machine.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning(machine);
var monitor = NewMonitor(repo.Object, loader.Object);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
await Task.Delay(80);
var stop = monitor.StopAsync();
var completed = await Task.WhenAny(stop, Task.Delay(2000)) == stop;
Assert.True(completed, "StopAsync did not complete promptly");
await stop; // rethrows if faulted
// Idempotent: second stop is a no-op and must not throw.
await monitor.StopAsync();
}
[Fact]
public async Task AddOrUpdateMachineAsync_StartsNewMachine_RaisesEvent_CachesSnapshot()
{
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(m.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning(); // start with no machines
var monitor = NewMonitor(repo.Object, loader.Object);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
var added = MachineWith("test");
var eventRaised = new ManualResetEventSlim(false);
monitor.SnapshotUpdated += (_, snap) =>
{
if (snap.MachineId == added.Id) eventRaised.Set();
};
var result = await monitor.AddOrUpdateMachineAsync(added, CancellationToken.None);
Assert.True(result.IsSuccess);
Assert.True(eventRaised.Wait(TimeSpan.FromSeconds(2)), "SnapshotUpdated not raised for added machine");
await monitor.StopAsync();
Assert.True(monitor.LatestSnapshots.ContainsKey(added.Id));
}
[Fact]
public async Task AddOrUpdateMachineAsync_CalledTwice_RestartsSingleLoop_NoDuplicateStream()
{
// Count concurrently-active drivers: a duplicate loop would push this above 1.
int active = 0;
int maxActive = 0;
var gate = new object();
IProtocolDriver MakeDriver(Guid id) => new CountingDriver(
() =>
{
lock (gate)
{
active++;
if (active > maxActive) maxActive = active;
}
return Result.Ok(Snapshot(id));
},
() => { lock (gate) { active--; } });
var machine = MachineWith("test");
var factory = new FakeFactory("test", m => Result.Ok(MakeDriver(m.Id)));
var loader = LoaderReturning(factory);
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
var r1 = await monitor.AddOrUpdateMachineAsync(machine, CancellationToken.None);
Assert.True(r1.IsSuccess);
await Task.Delay(80);
// Restart same id.
var r2 = await monitor.AddOrUpdateMachineAsync(machine, CancellationToken.None);
Assert.True(r2.IsSuccess);
await Task.Delay(120);
await monitor.StopAsync();
// At any instant at most one loop for this machine was polling.
Assert.True(maxActive <= 1, $"expected a single active loop, saw {maxActive} concurrent");
Assert.True(monitor.LatestSnapshots.ContainsKey(machine.Id));
}
[Fact]
public async Task RemoveMachineAsync_StopsLoop_DropsFromCache_NoFurtherUpdates()
{
var machine = MachineWith("test");
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(m.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning(machine);
var monitor = NewMonitor(repo.Object, loader.Object);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
// Wait until it has produced.
var sw = System.Diagnostics.Stopwatch.StartNew();
while (!monitor.LatestSnapshots.ContainsKey(machine.Id) && sw.Elapsed < TimeSpan.FromSeconds(2))
{
await Task.Delay(20);
}
Assert.True(monitor.LatestSnapshots.ContainsKey(machine.Id), "machine did not start");
int updatesAfterRemove = 0;
var remove = await monitor.RemoveMachineAsync(machine.Id, CancellationToken.None);
Assert.True(remove.IsSuccess);
Assert.False(monitor.LatestSnapshots.ContainsKey(machine.Id), "removed machine must drop from cache");
// Subscribe only now; if the loop truly stopped no further updates arrive.
monitor.SnapshotUpdated += (_, snap) =>
{
if (snap.MachineId == machine.Id) Interlocked.Increment(ref updatesAfterRemove);
};
await Task.Delay(150);
Assert.Equal(0, updatesAfterRemove);
await monitor.StopAsync();
}
[Fact]
public async Task RemoveMachineAsync_NotPresent_IsOkIdempotent()
{
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(m.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
var result = await monitor.RemoveMachineAsync(Guid.NewGuid(), CancellationToken.None);
Assert.True(result.IsSuccess);
}
[Fact]
public async Task AddOrUpdateMachineAsync_UnknownProtocol_ReturnsFail()
{
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(m.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
var bad = MachineWith("nope");
var result = await monitor.AddOrUpdateMachineAsync(bad, CancellationToken.None);
Assert.False(result.IsSuccess);
Assert.NotEmpty(result.Errors);
Assert.False(monitor.LatestSnapshots.ContainsKey(bad.Id));
await monitor.StopAsync();
}
[Fact]
public async Task AvailableProtocols_EmptyBeforeStart_ContainsPluginProtocolAfterStart()
{
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(m.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
// Before start: no factory map yet.
Assert.Empty(monitor.AvailableProtocols);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
Assert.Contains("test", monitor.AvailableProtocols);
await monitor.StopAsync();
}
[Fact]
public async Task AddOrUpdateMachineAsync_BeforeStart_ReturnsFail()
{
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(m.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
var machine = MachineWith("test");
var result = await monitor.AddOrUpdateMachineAsync(machine, CancellationToken.None);
Assert.False(result.IsSuccess);
Assert.NotEmpty(result.Errors);
Assert.Contains(result.Errors, e => e.Message.Contains("not started"));
}
[Fact]
public async Task ProbeAsync_ReturnsCatalog_FromFactoryDriver()
{
var catalog = new List
{
new DataItemDescriptor("x1_pos", "X", "POSITION", "SAMPLE", "MILLIMETER"),
new DataItemDescriptor("dev1_avail", "avail", "AVAILABILITY", "EVENT", ""),
};
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver(
"test",
() => Result.Ok(Snapshot(m.Id)),
() => Result>.Ok(catalog))));
var loader = LoaderReturning(factory);
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
var machine = MachineWith("test");
var result = await monitor.ProbeAsync(machine, CancellationToken.None);
Assert.True(result.IsSuccess);
Assert.Equal(2, result.Value.Count);
Assert.Contains(result.Value, d => d.Id == "x1_pos" && d.Type == "POSITION");
await monitor.StopAsync();
}
[Fact]
public async Task ProbeAsync_UnknownProtocol_ReturnsFail()
{
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(m.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
await monitor.StartAsync(PluginsDir, CancellationToken.None);
var bad = MachineWith("nope");
var result = await monitor.ProbeAsync(bad, CancellationToken.None);
Assert.False(result.IsSuccess);
Assert.NotEmpty(result.Errors);
await monitor.StopAsync();
}
[Fact]
public async Task ProbeAsync_BeforeStart_ReturnsFail()
{
var factory = new FakeFactory("test", m => Result.Ok(
new FakeDriver("test", () => Result.Ok(Snapshot(m.Id)))));
var loader = LoaderReturning(factory);
var repo = RepoReturning();
var monitor = NewMonitor(repo.Object, loader.Object);
var machine = MachineWith("test");
var result = await monitor.ProbeAsync(machine, CancellationToken.None);
Assert.False(result.IsSuccess);
Assert.Contains(result.Errors, e => e.Message.Contains("not started"));
}
/// Hand-rolled factory; behavior supplied by a delegate.
private sealed class FakeFactory : IProtocolDriverFactory
{
private readonly Func> _create;
public FakeFactory(string protocolId, Func> create)
{
ProtocolId = protocolId;
_create = create;
}
public string ProtocolId { get; }
public Result Create(Machine machine) => _create(machine);
}
/// Hand-rolled driver; behavior supplied by a delegate. Optional canned probe catalog.
private sealed class FakeDriver : IProtocolDriver
{
private readonly Func> _read;
private readonly Func>>? _probe;
public FakeDriver(
string protocolId,
Func> read,
Func>>? probe = null)
{
ProtocolId = protocolId;
_read = read;
_probe = probe;
}
public string ProtocolId { get; }
public Task> ReadCurrentAsync(CancellationToken cancellationToken)
{
cancellationToken.ThrowIfCancellationRequested();
return Task.FromResult(_read());
}
public Task>> ProbeAsync(CancellationToken cancellationToken) =>
Task.FromResult(_probe != null
? _probe()
: Result>.Ok(Array.Empty()));
}
///
/// Driver that holds an "active" window across an awaited delay so a test can detect two
/// loops polling the same machine concurrently. runs at read
/// start (increment + produce), in a finally (decrement) so the
/// count stays balanced even on cancellation.
///
private sealed class CountingDriver : IProtocolDriver
{
private readonly Func> _enter;
private readonly Action _exit;
public CountingDriver(Func> enter, Action exit)
{
_enter = enter;
_exit = exit;
}
public string ProtocolId => "test";
public async Task> ReadCurrentAsync(CancellationToken cancellationToken)
{
cancellationToken.ThrowIfCancellationRequested();
Result r = _enter();
try
{
await Task.Delay(30, cancellationToken).ConfigureAwait(false);
return r;
}
finally
{
_exit();
}
}
public Task>> ProbeAsync(CancellationToken cancellationToken) =>
Task.FromResult(Result>.Ok(Array.Empty()));
}
}
}