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())); } } }