diff --git a/InterprocessMemory.TestWorker/Program.cs b/InterprocessMemory.TestWorker/Program.cs index e5e48f4..91d9614 100644 --- a/InterprocessMemory.TestWorker/Program.cs +++ b/InterprocessMemory.TestWorker/Program.cs @@ -18,6 +18,8 @@ /// concurrent_producer <name> <producerId> — enqueue 1000 unique integers /// try_write_lock <name> — try the cross-process write lock for 250 ms /// orphan_write_lock <name> — acquire a write lock and exit without releasing it +/// hold_write_lock <name> — acquire a write lock, print "holding", and keep it until killed +/// hold_read_lock <name> — acquire a read lock, print "holding", and keep it until killed /// if (args.Length < 2) { @@ -43,6 +45,8 @@ args.Length >= 3 ? int.Parse(args[2]) : 0), "try_write_lock" => TryWriteLock(bufferName), "orphan_write_lock" => OrphanWriteLock(bufferName), + "hold_write_lock" => HoldWriteLock(bufferName), + "hold_read_lock" => HoldReadLock(bufferName), _ => Error($"Unknown role: {role}") }; @@ -197,6 +201,32 @@ static int OrphanWriteLock(string name) return 0; // Deliberately skip Dispose/Release; process teardown closes only the mapping handle. } +static int HoldWriteLock(string name) +{ + using var region = MemoryRegion.OpenExisting(name); + if (!region.TryAcquireWriteLock(TimeSpan.FromSeconds(5))) + return Error("failed to acquire held test lock"); + + // The parent reads this line, probes the lock while we are alive, and then kills us + // to simulate a crash while the lock is held. + Console.WriteLine("holding"); + Console.Out.Flush(); + Thread.Sleep(Timeout.Infinite); + return 0; +} + +static int HoldReadLock(string name) +{ + using var region = MemoryRegion.OpenExisting(name); + if (!region.TryAcquireReadLock(TimeSpan.FromSeconds(5))) + return Error("failed to acquire held test read lock"); + + Console.WriteLine("holding"); + Console.Out.Flush(); + Thread.Sleep(Timeout.Infinite); + return 0; +} + static int Error(string msg) { Console.Error.WriteLine(msg); diff --git a/InterprocessMemory.Tests/ChangeVerificationTests.cs b/InterprocessMemory.Tests/ChangeVerificationTests.cs index 803c093..222eddc 100644 --- a/InterprocessMemory.Tests/ChangeVerificationTests.cs +++ b/InterprocessMemory.Tests/ChangeVerificationTests.cs @@ -1041,18 +1041,19 @@ public void Audit4_FieldDefinition_NullOrEmptyName_AllFactoriesReject() // ── AUDIT-5: ReleaseWriteLock CAS-by-owner ─────────────────────────────── [Test] - public void Audit5_ReleaseWriteLock_WithoutAcquire_IsNoOp() + public void Audit5_ReleaseWriteLock_WithoutAcquire_Throws() { - // Caller bug: releasing a lock that wasn't acquired by this process. Without the - // owner-CAS guard, ReleaseWriteLock would zero ownership metadata and free the lock — - // dangerous when another process legitimately holds it. With the fix it's a logged no-op. + // Caller bug: releasing a lock that wasn't acquired by this thread. Without the owner + // guard, ReleaseWriteLock would zero ownership metadata and free the lock — dangerous + // when another process legitimately holds it. The guard leaves the lock state untouched + // and reports the misuse instead of hiding it. using var buf = new MemoryRegion(N("Audit5_NoAcquire"), new MemoryRegionOptions { Capacity = 4096 }); - // No acquire here. Release should be a safe no-op (logs a warning). - Assert.DoesNotThrow(() => buf.ReleaseWriteLock()); + // No acquire here. + Assert.Throws(() => buf.ReleaseWriteLock()); - // Now actually acquire — must still work normally after the no-op release. + // Now actually acquire — must still work normally after the rejected release. Assert.That(buf.TryAcquireWriteLock(TimeSpan.FromSeconds(1)), Is.True); buf.ReleaseWriteLock(); } diff --git a/InterprocessMemory.Tests/ConcurrentMessageQueueTests.cs b/InterprocessMemory.Tests/ConcurrentMessageQueueTests.cs index a475f97..20405ef 100644 --- a/InterprocessMemory.Tests/ConcurrentMessageQueueTests.cs +++ b/InterprocessMemory.Tests/ConcurrentMessageQueueTests.cs @@ -370,8 +370,11 @@ public async Task MPMC_HighContention_ShouldMaintainIntegrity() var producersDone = 0; var errors = new ConcurrentBag(); + // The producers and consumers spin, so each gets a dedicated thread. On the shared thread pool + // the four spinning producers occupy every worker of a small machine, the consumers cannot start + // for seconds (the pool adds threads slowly), and the producers give up on the full queue. // Producers - var producers = Enumerable.Range(0, producerCount).Select(producerId => Task.Run(() => + var producers = Enumerable.Range(0, producerCount).Select(producerId => Task.Factory.StartNew(() => { try { @@ -397,36 +400,40 @@ public async Task MPMC_HighContention_ShouldMaintainIntegrity() { Interlocked.Increment(ref producersDone); } - })).ToArray(); + }, TaskCreationOptions.LongRunning)).ToArray(); // Consumers - var consumers = Enumerable.Range(0, consumerCount).Select(consumerId => Task.Run(() => + var consumers = Enumerable.Range(0, consumerCount).Select(consumerId => Task.Factory.StartNew(() => { var readBuffer = new byte[64]; - int emptyReads = 0; + var deadline = System.Diagnostics.Stopwatch.StartNew(); - while (received.Count < totalExpected && emptyReads < 1000) + // Wait for the full count instead of giving up after a number of empty polls: in a fresh + // process the consumers can burn through any such budget before the producers have sent + // anything (JIT, scheduling), and the producers then block forever on a full queue. + while (received.Count < totalExpected) { + if (deadline.Elapsed > TimeSpan.FromSeconds(30)) + { + errors.Add($"Consumer {consumerId} timed out with {received.Count}/{totalExpected} received"); + return; + } + var bytesRead = buffer.TryRead(readBuffer); if (bytesRead >= 4) { received.Add(BitConverter.ToInt32(readBuffer, 0)); - emptyReads = 0; + } + else if (Volatile.Read(ref producersDone) == producerCount && buffer.ApproximateCount == 0) + { + Thread.Sleep(1); } else { - emptyReads++; - if (producersDone == producerCount && buffer.ApproximateCount == 0) - { - Thread.Sleep(1); - } - else - { - Thread.SpinWait(10); - } + Thread.SpinWait(10); } } - })).ToArray(); + }, TaskCreationOptions.LongRunning)).ToArray(); await Task.WhenAll(producers.Concat(consumers)); diff --git a/InterprocessMemory.Tests/CrossProcessTests.cs b/InterprocessMemory.Tests/CrossProcessTests.cs index 60e05f5..1caaf73 100644 --- a/InterprocessMemory.Tests/CrossProcessTests.cs +++ b/InterprocessMemory.Tests/CrossProcessTests.cs @@ -293,6 +293,147 @@ public void CrossProcess_OrphanWriteLock_IsRecovered() region.ReleaseWriteLock(); } + /// + /// Starts a child that holds a lock until it is killed, and returns once the child has reported + /// that it holds the lock. is hold_write_lock or hold_read_lock. + /// + private static Process StartLockHolder(string bufferName, string role = "hold_write_lock") + { + Process process = Process.Start(CreateHelperStartInfo(role, bufferName))!; + string? line = process.StandardOutput.ReadLine(); + if (line != "holding") + { + try + { process.Kill(); } + catch (InvalidOperationException) { /* already exited */ } + process.Dispose(); + Assert.Fail($"lock holder did not report 'holding' (got '{line}')"); + } + + return process; + } + + /// Kills the process (a no-op when it has already exited) and waits for it to be gone. + private static void KillAndWait(Process process) + { + try + { process.Kill(); } + catch (InvalidOperationException) { /* already exited */ } + process.WaitForExit(); + } + + [Test, Timeout(30000)] + public void CrossProcess_LiveWriteLockOwner_IsNotReportedAsOrphan() + { + // Process.StartTime is derived per observing process from a wall-clock boot-time snapshot on + // Linux, so comparing the owner's recorded value with the observer's value reported every + // live owner as an impostor and let waiters steal its lock. + string name = GetUniqueName("LiveOwner"); + using var region = MemoryRegion.CreateOrOpen(name, 256); + using Process holder = StartLockHolder(name); + try + { + for (int i = 0; i < 50; i++) + { + Assert.That(region.IsWriteLockOrphaned(), Is.False, $"probe {i}"); + Thread.Sleep(10); + } + + Assert.That(region.TryAcquireWriteLock(TimeSpan.FromMilliseconds(300)), Is.False, + "a live owner's lock must not be taken over"); + } + finally + { + KillAndWait(holder); + } + } + + [Test, Timeout(30000)] + public void CrossProcess_InfiniteWait_RecoversWhenOwnerProcessDies() + { + string name = GetUniqueName("InfiniteOrphan"); + using var region = MemoryRegion.CreateOrOpen(name, 256); + using Process holder = StartLockHolder(name); + try + { + // Release on the acquiring thread: the write lock is thread-affine. + Task waiter = Task.Run(() => + { + bool acquired = region.TryAcquireWriteLock(Timeout.InfiniteTimeSpan); + if (acquired) + region.ReleaseWriteLock(); + return acquired; + }); + + // Give the waiter time to run its first orphan probe while the owner is still alive. + Thread.Sleep(500); + Assert.That(waiter.IsCompleted, Is.False, "the owner is alive, so the waiter must still be waiting"); + + KillAndWait(holder); + + Assert.That(waiter.Wait(TimeSpan.FromSeconds(10)), Is.True, + "an infinite wait must keep probing the owner and recover once it has died"); + Assert.That(waiter.Result, Is.True); + } + finally + { + KillAndWait(holder); + } + } + + [Test, Timeout(40000)] + public void CrossProcess_FiniteWait_RecoversPromptlyWhenOwnerDies() + { + // The orphan probe used to run only at the start of the wait and at 75% of the timeout, so + // with a 30 s timeout an owner that died after one second was noticed after ~22 s. + string name = GetUniqueName("FiniteOrphan"); + using var region = MemoryRegion.CreateOrOpen(name, 256); + using Process holder = StartLockHolder(name); + try + { + Task waiter = Task.Run(() => + { + bool acquired = region.TryAcquireWriteLock(TimeSpan.FromSeconds(30)); + if (acquired) + region.ReleaseWriteLock(); + return acquired; + }); + + Thread.Sleep(500); + KillAndWait(holder); + + Assert.That(waiter.Wait(TimeSpan.FromSeconds(5)), Is.True, + "recovery must not wait for most of the lock timeout"); + Assert.That(waiter.Result, Is.True); + } + finally + { + KillAndWait(holder); + } + } + + [Test, Timeout(30000)] + public void CrossProcess_DeadReaderProcess_NeedsForceResetLocks() + { + // Read locks are not attributed to an owner, so a reader that dies leaves the shared count + // above zero and nothing can tell that it is stale. ForceResetLocks is the operator's way out. + string name = GetUniqueName("DeadReader"); + using var region = MemoryRegion.CreateOrOpen(name, 256); + using Process holder = StartLockHolder(name, "hold_read_lock"); + Assert.That(region.GetLockOwnerInfo().ReaderCount, Is.EqualTo(1)); + KillAndWait(holder); + + Assert.That(region.TryAcquireWriteLock(TimeSpan.FromMilliseconds(500)), Is.False, + "the dead reader's lock is still counted"); + Assert.That(region.GetLockOwnerInfo().ReaderCount, Is.EqualTo(1)); + + region.ForceResetLocks(); + + Assert.That(region.GetLockOwnerInfo().ReaderCount, Is.EqualTo(0)); + Assert.That(region.TryAcquireWriteLock(TimeSpan.FromSeconds(1)), Is.True); + region.ReleaseWriteLock(); + } + // ── Schema ─────────────────────────────────────────────────────────────── public struct IpcTestSchema : IMemorySchema diff --git a/InterprocessMemory.Tests/ExtremeStressTests.cs b/InterprocessMemory.Tests/ExtremeStressTests.cs index 61795f1..ec77332 100644 --- a/InterprocessMemory.Tests/ExtremeStressTests.cs +++ b/InterprocessMemory.Tests/ExtremeStressTests.cs @@ -127,7 +127,7 @@ public async Task MPMC_BurstTraffic_ShouldHandleSpikes() for (int burst = 0; burst < burstCount; burst++) { var received = new ConcurrentBag(); - var producerDone = false; + var deadline = System.Diagnostics.Stopwatch.StartNew(); // Burst producer var producer = Task.Run(() => @@ -137,30 +137,34 @@ public async Task MPMC_BurstTraffic_ShouldHandleSpikes() { BitConverter.TryWriteBytes(data, burst * messagesPerBurst + i); while (!buffer.TryWrite(data)) + { + if (deadline.Elapsed > TimeSpan.FromSeconds(30)) + { + errors.Add($"Burst {burst}: producer blocked at message {i}"); + return; + } + Thread.SpinWait(1); + } } - producerDone = true; }); - // Consumer + // Consumer: wait for the full burst instead of giving up after a number of empty polls. + // That budget used to run out in microseconds when the consumer started before the + // producer, after which the producer spun forever on a full queue and the test hung. var consumer = Task.Run(() => { var readBuf = new byte[64]; - int emptyCount = 0; - while (received.Count < messagesPerBurst && emptyCount < 10000) + while (received.Count < messagesPerBurst) { + if (deadline.Elapsed > TimeSpan.FromSeconds(30)) + return; + var bytesRead = buffer.TryRead(readBuf); if (bytesRead > 0) - { received.Add(BitConverter.ToInt32(readBuf, 0)); - emptyCount = 0; - } else - { - emptyCount++; - if (producerDone && buffer.ApproximateCount == 0) - break; - } + Thread.SpinWait(1); } }); diff --git a/InterprocessMemory.Tests/LibraryHardeningTests.cs b/InterprocessMemory.Tests/LibraryHardeningTests.cs index a7d075c..b05d307 100644 --- a/InterprocessMemory.Tests/LibraryHardeningTests.cs +++ b/InterprocessMemory.Tests/LibraryHardeningTests.cs @@ -82,10 +82,21 @@ public void WriteLock_CannotBeReleasedByDifferentThread() Assert.That(buffer.TryAcquireWriteLock(TimeSpan.FromSeconds(1)), Is.True); - var invalidRelease = new Thread(buffer.ReleaseWriteLock); + // An unhandled exception on a raw Thread would take the test host down, so capture it. + Exception? releaseError = null; + var invalidRelease = new Thread(() => + { + try + { buffer.ReleaseWriteLock(); } + catch (Exception ex) { releaseError = ex; } + }); invalidRelease.Start(); invalidRelease.Join(); + // The misuse is reported instead of being silently ignored (a silently ignored release is + // what turned `await` inside a lock scope into a lock nobody could ever release). + Assert.That(releaseError, Is.InstanceOf()); + var contender = Task.Run(() => buffer.TryAcquireWriteLock(TimeSpan.FromMilliseconds(50))); Assert.That(contender.Result, Is.False); @@ -94,6 +105,378 @@ public void WriteLock_CannotBeReleasedByDifferentThread() buffer.ReleaseWriteLock(); } + [Test] + public void Dispose_WhileThreadWaitsForWriteLock_ReleasesWaiterInsteadOfCrashing() + { + // Dispose used to unmap the view while a thread was still spinning on the header, which is an + // AccessViolationException: the whole process dies and the exception cannot be caught. + var buffer = new MemoryRegion( + N("DisposeWriteWaiter"), + new MemoryRegionOptions { Capacity = 256 }); + + Assert.That(buffer.TryAcquireWriteLock(TimeSpan.FromSeconds(1)), Is.True); + + using var started = new ManualResetEventSlim(false); + var waiter = Task.Run(() => + { + started.Set(); + return buffer.TryAcquireWriteLock(Timeout.InfiniteTimeSpan); + }); + Assert.That(started.Wait(TimeSpan.FromSeconds(5)), Is.True); + Thread.Sleep(100); // let the waiter reach its spin loop + + buffer.Dispose(); + + var error = Assert.Throws(() => waiter.Wait(TimeSpan.FromSeconds(10))); + Assert.That(error!.InnerException, Is.InstanceOf()); + } + + [Test] + public void Dispose_WhileThreadWaitsForReadLock_ReleasesWaiterInsteadOfCrashing() + { + var buffer = new MemoryRegion( + N("DisposeReadWaiter"), + new MemoryRegionOptions { Capacity = 256 }); + + Assert.That(buffer.TryAcquireWriteLock(TimeSpan.FromSeconds(1)), Is.True); + + using var started = new ManualResetEventSlim(false); + var waiter = Task.Run(() => + { + started.Set(); + return buffer.TryAcquireReadLock(Timeout.InfiniteTimeSpan); + }); + Assert.That(started.Wait(TimeSpan.FromSeconds(5)), Is.True); + Thread.Sleep(100); + + buffer.Dispose(); + + var error = Assert.Throws(() => waiter.Wait(TimeSpan.FromSeconds(10))); + Assert.That(error!.InnerException, Is.InstanceOf()); + } + + [Test] + public void StructuredMemory_WriteLockGuard_DisposedOnAnotherThread_ThrowsAndStaysHeld() + { + string name = N("WriteGuardThread"); + using var memory = StructuredMemory.CreateOrOpen(name, new SimpleSchema()); + using var peer = StructuredMemory.OpenExisting(name, new SimpleSchema()); + + // What `await` inside a lock scope does: the guard is disposed on a different thread. + var guard = memory.AcquireWriteLock(); + Exception? error = null; + Task.Run(() => + { + try + { guard.Dispose(); } + catch (Exception ex) { error = ex; } + }).Wait(); + + Assert.That(error, Is.InstanceOf()); + + // Nothing was released or unbalanced: the lock still belongs to the acquiring thread... + Assert.Throws(() => peer.AcquireWriteLock(TimeSpan.FromMilliseconds(100))); + + // ...which can still release it properly, after which other holders get in. + guard.Dispose(); + using (peer.AcquireWriteLock(TimeSpan.FromSeconds(1))) + { + } + } + + [Test] + public void StructuredMemory_ReadLockGuard_DisposedOnAnotherThread_ThrowsAndStaysHeld() + { + string name = N("ReadGuardThread"); + using var memory = StructuredMemory.CreateOrOpen(name, new SimpleSchema()); + using var peer = StructuredMemory.OpenExisting(name, new SimpleSchema()); + + var guard = memory.AcquireReadLock(); + Exception? error = null; + Task.Run(() => + { + try + { guard.Dispose(); } + catch (Exception ex) { error = ex; } + }).Wait(); + + Assert.That(error, Is.InstanceOf()); + Assert.Throws(() => peer.AcquireWriteLock(TimeSpan.FromMilliseconds(100))); + + guard.Dispose(); + using (peer.AcquireWriteLock(TimeSpan.FromSeconds(1))) + { + } + } + + [Test] + public void DisposeGracePeriod_DefaultsToTenMilliseconds_AndRejectsNegativeValues() + { + Assert.That(MemoryRegion.DisposeGracePeriod, Is.EqualTo(TimeSpan.FromMilliseconds(10))); + Assert.Throws( + () => MemoryRegion.DisposeGracePeriod = TimeSpan.FromMilliseconds(-1)); + } + + [Test] + public void Dispose_KeepsTheMappingForTheGracePeriod() + { + TimeSpan original = MemoryRegion.DisposeGracePeriod; + try + { + MemoryRegion.DisposeGracePeriod = TimeSpan.FromMilliseconds(100); + var withGrace = new MemoryRegion(N("Grace"), new MemoryRegionOptions { Capacity = 256 }); + var stopwatch = System.Diagnostics.Stopwatch.StartNew(); + withGrace.Dispose(); + Assert.That(stopwatch.ElapsedMilliseconds, Is.GreaterThanOrEqualTo(90)); + + MemoryRegion.DisposeGracePeriod = TimeSpan.Zero; + var withoutGrace = new MemoryRegion(N("NoGrace"), new MemoryRegionOptions { Capacity = 256 }); + stopwatch.Restart(); + withoutGrace.Dispose(); + Assert.That(stopwatch.ElapsedMilliseconds, Is.LessThan(90)); + } + finally + { + MemoryRegion.DisposeGracePeriod = original; + } + } + + [Test, Timeout(120000)] + public void Dispose_WhileThreadsBusyPollTheQueue_DoesNotCrash() + { + // The common shutdown pattern: consumers spin on TryDequeue while another thread disposes. + // Without DisposeGracePeriod this killed the process with an AccessViolationException within a + // few hundred rounds on a 4-core machine. This is a mitigation test, not a proof: a thread that + // is descheduled for longer than the grace period at the wrong moment can still fail. + var random = new Random(42); + int pollerCount = Math.Max(3, Environment.ProcessorCount - 1); + + for (int round = 0; round < 200; round++) + { + string name = N($"BusyPoll{round}"); + var queue = InterprocessMemory.ConcurrentQueue.CreateOrOpen(name, 64); + var pollers = new Task[pollerCount]; + for (int i = 0; i < pollers.Length; i++) + { + pollers[i] = Task.Run(() => + { + try + { + while (true) + { + queue.TryDequeue(out _); + queue.TryEnqueue(1); + } + } + catch (ObjectDisposedException) + { + // Expected: the queue was disposed under the poller. + } + }); + } + + Thread.Sleep(random.Next(0, 3)); + queue.Dispose(); + + Assert.That(Task.WaitAll(pollers, TimeSpan.FromSeconds(10)), Is.True, $"round {round}"); + MemoryRegion.Remove(name); + } + } + + private static int CountOpenFilesMatching(string name) + { + return Directory.GetFiles("/proc/self/fd").Count(path => + { + try + { return new FileInfo(path).LinkTarget?.Contains(name) == true; } + catch (IOException) { return false; } + }); + } + + [Test] + public void SharedArray_FailedOpen_ReleasesTheMappingImmediately() + { + if (!OperatingSystem.IsLinux()) + Assert.Ignore("Counts the process's open file descriptors through /proc."); + + // Opening with another element type fails the header check after the region is mapped. + // Without an explicit dispose the file descriptor stayed open until the finalizer ran. + string name = N("ArrayLeak"); + using var owner = SharedArray.CreateOrOpen(name, 4); + int before = CountOpenFilesMatching(name); + + Assert.Throws(() => SharedArray.OpenExisting(name)); + + Assert.That(CountOpenFilesMatching(name), Is.EqualTo(before)); + } + + public struct WideSchema : IMemorySchema + { + public IEnumerable GetFields() + { + yield return FieldDefinition.Scalar("Id"); + yield return FieldDefinition.String("Name", 8); + } + } + + [Test] + public void StructuredMemory_AutoLockedAccess_DoesNotAllocatePerCall() + { + // Values wider than eight bytes and strings take the region lock automatically. Each of those + // calls used to allocate a delegate (about 104 bytes) for the lock guard. + using var memory = StructuredMemory.CreateOrOpen(N("NoAlloc"), new WideSchema()); + Guid id = Guid.NewGuid(); + + for (int i = 0; i < 2000; i++) + { + memory.Write("Id", id); + memory.Read("Id"); + memory.WriteString("Name", "ab"); + } + + const int Rounds = 10_000; + long before = GC.GetAllocatedBytesForCurrentThread(); + for (int i = 0; i < Rounds; i++) + { + memory.Write("Id", id); + memory.Read("Id"); + memory.WriteString("Name", "ab"); + } + long bytesPerCall = (GC.GetAllocatedBytesForCurrentThread() - before) / (Rounds * 3); + + Assert.That(bytesPerCall, Is.LessThan(8)); + } + + [Test] + public void OrphanLockTimeout_IsDisabledByDefault() + { + // A time limit takes the lock away from a healthy owner that merely holds it for long + // (a long transaction, a paused debugger), so it is opt-in. Dead owners are still recovered. + var options = new MemoryRegionOptions(); + + Assert.That(options.OrphanLockTimeout, Is.EqualTo(TimeSpan.Zero)); + Assert.That(options.EnableOrphanLockDetection, Is.True); + } + + [Test] + public void ForceResetLocks_ClearsReadersAndHeldWriteLock() + { + using var buffer = new MemoryRegion( + N("ForceReset"), + new MemoryRegionOptions { Capacity = 256 }); + + Assert.That(buffer.TryAcquireReadLock(TimeSpan.FromSeconds(1)), Is.True); + Assert.That(buffer.GetLockOwnerInfo().ReaderCount, Is.EqualTo(1)); + Assert.That(buffer.TryAcquireWriteLock(TimeSpan.FromMilliseconds(100)), Is.False); + + buffer.ForceResetLocks(); + + Assert.That(buffer.GetLockOwnerInfo().ReaderCount, Is.EqualTo(0)); + Assert.That(buffer.TryAcquireWriteLock(TimeSpan.FromSeconds(1)), Is.True); + + // Resetting also frees a write lock that is still "held"; its holder learns that on release. + buffer.ForceResetLocks(); + Assert.Throws(() => buffer.ReleaseWriteLock()); + + Assert.That(buffer.TryAcquireWriteLock(TimeSpan.FromSeconds(1)), Is.True); + buffer.ReleaseWriteLock(); + } + + [Test] + public void StructuredMemory_ForceResetLocks_UnblocksWriter() + { + string name = N("StructuredReset"); + using var memory = StructuredMemory.CreateOrOpen(name, new SimpleSchema()); + using var peer = StructuredMemory.OpenExisting(name, new SimpleSchema()); + + // A reader that never comes back, as after a crash. + var staleReader = peer.AcquireReadLock(); + Assert.Throws(() => memory.AcquireWriteLock(TimeSpan.FromMilliseconds(100))); + + memory.ForceResetLocks(); + + using (memory.AcquireWriteLock(TimeSpan.FromSeconds(1))) + { + } + + staleReader.Dispose(); + } + + [Test] + public void Remove_DeletesLinuxRegion_SoItCanBeRecreatedWithDifferentCapacity() + { + if (!OperatingSystem.IsLinux()) + Assert.Ignore("Only Linux keeps regions in /dev/shm after their last user is gone."); + + string name = N("Remove"); + using (MemoryRegion.CreateOrOpen(name, 256)) + { + } + + // The region outlives its users, so a different capacity is rejected... + Assert.Throws(() => MemoryRegion.CreateOrOpen(name, 512)); + + // ...until it is removed. + Assert.That(MemoryRegion.Remove(name), Is.True); + Assert.That(MemoryRegion.Remove(name), Is.False); + + using var fresh = MemoryRegion.CreateOrOpen(name, 512); + Assert.That(fresh.IsOwner, Is.True); + Assert.That(fresh.Capacity, Is.EqualTo(512)); + } + + [Test] + public void Remove_ClearsRegionLeftInitializingByCrashedCreator() + { + if (!OperatingSystem.IsLinux()) + Assert.Ignore("Only Linux keeps regions in /dev/shm after their last user is gone."); + + // "IPMI": the creator died after claiming initialization and before publishing the header, + // which makes every opener wait and then time out. + string name = N("StuckInit"); + using (var file = new FileStream("/dev/shm/" + name, FileMode.CreateNew)) + { + file.SetLength(128 + 256); + file.Write(BitConverter.GetBytes(0x494D5049u)); + } + + Assert.That(MemoryRegion.Remove(name), Is.True); + + using var region = MemoryRegion.CreateOrOpen(name, 256); + Assert.That(region.IsOwner, Is.True); + } + + [Test] + public void Remove_FileBackedRegion_DeletesTheFile() + { + string name = N("RemoveFile"); + string path = Path.Combine(Path.GetTempPath(), name + ".bin"); + var options = new MemoryRegionOptions { FilePath = path }; + try + { + using (MemoryRegion.CreateOrOpen(name, 256, options)) + { + } + + Assert.That(File.Exists(path), Is.True); + Assert.That(MemoryRegion.Remove(name, options), Is.True); + Assert.That(File.Exists(path), Is.False); + Assert.That(MemoryRegion.Remove(name, options), Is.False); + } + finally + { + File.Delete(path); + } + } + + [Test] + public void Remove_UnknownRegion_ReturnsFalse_AndBadNamesAreRejected() + { + Assert.That(MemoryRegion.Remove(N("Nothing")), Is.False); + Assert.Throws(() => MemoryRegion.Remove("a/b")); + Assert.Throws(() => MemoryRegion.Remove("")); + } + [Test] public void ReadLock_DoubleRelease_DoesNotBreakWriterExclusion() { diff --git a/InterprocessMemory.Tests/TypeLayoutFingerprintTests.cs b/InterprocessMemory.Tests/TypeLayoutFingerprintTests.cs new file mode 100644 index 0000000..d475f3d --- /dev/null +++ b/InterprocessMemory.Tests/TypeLayoutFingerprintTests.cs @@ -0,0 +1,92 @@ +using System.Collections.Generic; +using System.IO; +using NUnit.Framework; +using InterprocessMemory; + +namespace InterprocessMemory.Tests; + +[TestFixture] +public class TypeLayoutFingerprintTests +{ + private static string N(string prefix) => $"Fingerprint_{prefix}_{Guid.NewGuid():N}"; + + [Test] + public void MarshallableTypes_KeepTheirRelease300Fingerprint() + { + // These values were produced by 3.0.0. A different value would make a process built from + // this version reject regions created by a 3.0.0 process (and vice versa) with + // "different format or element type", so they must only ever change together with the + // shared-memory format version. + Assert.That(TypeLayoutFingerprint.Create(), + Is.EqualTo(new TypeLayoutFingerprint(0xD4616E4290E3AB20, 0xEB0DEB962578DCD5))); + Assert.That(TypeLayoutFingerprint.Create(), + Is.EqualTo(new TypeLayoutFingerprint(0xD52D441606CE7C3F, 0x4D50B1F77C44F34C))); + Assert.That(TypeLayoutFingerprint.Create(), + Is.EqualTo(new TypeLayoutFingerprint(0x7EE1FD251B406A46, 0x9427C9B56B8CC942))); + Assert.That(TypeLayoutFingerprint.Create(), + Is.EqualTo(new TypeLayoutFingerprint(0xB7337E9F0D12BB73, 0x2C0F045AB00D96C8))); + } + + [Test] + public void GenericUnmanagedStructs_HaveStableDistinctFingerprints() + { + // Marshal.SizeOf rejects generic types, so (int, int) used to fail with ArgumentException + // in every typed container even though it satisfies the unmanaged constraint. + var tupleA = TypeLayoutFingerprint.Create<(int, int)>(); + var tupleB = TypeLayoutFingerprint.Create<(int, int)>(); + var differentArgument = TypeLayoutFingerprint.Create<(int, uint)>(); + var differentShape = TypeLayoutFingerprint.Create<(int, int, int)>(); + var pair = TypeLayoutFingerprint.Create>(); + + Assert.That(tupleA, Is.EqualTo(tupleB)); + Assert.That(tupleA, Is.Not.EqualTo(differentArgument)); + Assert.That(tupleA, Is.Not.EqualTo(differentShape)); + Assert.That(tupleA, Is.Not.EqualTo(pair)); + } + + [Test] + public void SingleProducerQueue_AcceptsGenericStruct() + { + string name = N("SpscTuple"); + using var producer = SingleProducerQueue<(int, long)>.CreateOrOpen(name, 8); + Assert.That(producer.TryEnqueue((7, 9L)), Is.True); + + using var consumer = SingleProducerQueue<(int, long)>.OpenExisting(name); + Assert.That(consumer.TryDequeue(out (int, long) item), Is.True); + Assert.That(item, Is.EqualTo((7, 9L))); + } + + [Test] + public void ConcurrentQueue_AcceptsGenericStruct() + { + string name = N("MpmcPair"); + using var queue = InterprocessMemory.ConcurrentQueue>.CreateOrOpen(name, 8); + Assert.That(queue.TryEnqueue(new KeyValuePair(3, 4)), Is.True); + + using var reader = InterprocessMemory.ConcurrentQueue>.OpenExisting(name); + Assert.That(reader.TryDequeue(out KeyValuePair item), Is.True); + Assert.That(item.Key, Is.EqualTo(3)); + Assert.That(item.Value, Is.EqualTo(4)); + } + + [Test] + public void SharedArray_AcceptsGenericStruct() + { + string name = N("ArrayTuple"); + using var owner = SharedArray<(int, int)>.CreateOrOpen(name, 4); + owner[2] = (5, 6); + + using var reader = SharedArray<(int, int)>.OpenExisting(name); + Assert.That(reader[2], Is.EqualTo((5, 6))); + } + + [Test] + public void OpeningGenericStructQueueWithDifferentLayout_Throws() + { + // Same size (8 bytes), different element type: only the fingerprint tells them apart. + string name = N("Mismatch"); + using var owner = SingleProducerQueue<(int, int)>.CreateOrOpen(name, 4); + + Assert.Throws(() => SingleProducerQueue<(int, uint)>.OpenExisting(name)); + } +} diff --git a/InterprocessMemory/ConcurrentMessageQueue.cs b/InterprocessMemory/ConcurrentMessageQueue.cs index e8a8d46..5d79794 100644 --- a/InterprocessMemory/ConcurrentMessageQueue.cs +++ b/InterprocessMemory/ConcurrentMessageQueue.cs @@ -194,7 +194,7 @@ private ConcurrentMessageQueue( if (createOrOpen) { - _slotCount = RoundUpToPowerOf2(capacity!.Value); + _slotCount = PowerOfTwo.RoundUp(capacity!.Value, "capacity"); _maxMessageSize = maxMessageSize!.Value; _slotTotalSize = RoundUpToMultiple( checked(SlotHeaderSize + _maxMessageSize), 8); @@ -309,10 +309,13 @@ private void ValidateAndLoadBuffer(int? requestedCapacity, int? requestedMaxMess storedMaxMessageSize <= 0 || storedSlotStride != expectedSlotStride) throw new InvalidDataException("The message queue header contains invalid sizing metadata."); - if (requestedCapacity.HasValue && - RoundUpToPowerOf2(requestedCapacity.Value) != storedSlotCount) - throw new InvalidOperationException( - $"Capacity mismatch: expected {RoundUpToPowerOf2(requestedCapacity.Value)}, found {storedSlotCount}"); + if (requestedCapacity.HasValue) + { + int expectedCapacity = PowerOfTwo.RoundUp(requestedCapacity.Value, "capacity"); + if (expectedCapacity != storedSlotCount) + throw new InvalidOperationException( + $"Capacity mismatch: expected {expectedCapacity}, found {storedSlotCount}"); + } if (requestedMaxMessageSize.HasValue && requestedMaxMessageSize.Value != storedMaxMessageSize) throw new InvalidOperationException( @@ -633,36 +636,10 @@ public void Dispose() if (Interlocked.Exchange(ref _disposed, 1) != 0) return; + // The MemoryHandle wraps an unmanaged pointer and owns nothing, so there is no finalizer + // here: if Dispose is never called, the MemoryRegion's own finalizer unmaps the memory. _memoryHandle.Dispose(); - // See SingleProducerByteStream: only dispose the managed _buffer on the deterministic - // path. The finalizer below skips it to avoid touching peer objects whose own - // finalizers may have already run. _buffer?.Dispose(); - GC.SuppressFinalize(this); - } - - /// - /// Releases unmanaged resources if Dispose was not called. - /// - ~ConcurrentMessageQueue() - { - if (Interlocked.Exchange(ref _disposed, 1) != 0) - return; - try - { _memoryHandle.Dispose(); } - catch { /* best-effort */ } - } - - [MethodImpl(MethodImplOptions.AggressiveInlining)] - private static int RoundUpToPowerOf2(int value) - { - value--; - value |= value >> 1; - value |= value >> 2; - value |= value >> 4; - value |= value >> 8; - value |= value >> 16; - return value + 1; } [MethodImpl(MethodImplOptions.AggressiveInlining)] diff --git a/InterprocessMemory/ConcurrentQueue.cs b/InterprocessMemory/ConcurrentQueue.cs index 0cceb89..8ac8f79 100644 --- a/InterprocessMemory/ConcurrentQueue.cs +++ b/InterprocessMemory/ConcurrentQueue.cs @@ -99,7 +99,7 @@ private ConcurrentQueue( if (createOrOpen) { - _capacity = RoundUpToPowerOf2(capacity!.Value); + _capacity = PowerOfTwo.RoundUp(capacity!.Value, "capacity"); _capacityMask = _capacity - 1; _slotStride = RoundUpToMultiple(checked(SlotHeaderSize + _elementSize), 8); long regionCapacity = checked(HeaderSize + (long)_capacity * _slotStride); @@ -179,10 +179,13 @@ private void ValidateAndLoad(int? requestedCapacity) _header->FingerprintHigh != _fingerprint.High) throw new InvalidDataException("The queue has a different format or element type."); - if (requestedCapacity.HasValue && - RoundUpToPowerOf2(requestedCapacity.Value) != storedCapacity) - throw new InvalidOperationException( - $"Capacity mismatch: expected {RoundUpToPowerOf2(requestedCapacity.Value)}, found {storedCapacity}."); + if (requestedCapacity.HasValue) + { + int expectedCapacity = PowerOfTwo.RoundUp(requestedCapacity.Value, "capacity"); + if (expectedCapacity != storedCapacity) + throw new InvalidOperationException( + $"Capacity mismatch: expected {expectedCapacity}, found {storedCapacity}."); + } long expectedRegionSize = checked(HeaderSize + (long)storedCapacity * storedStride); if (_region.Capacity != expectedRegionSize) @@ -314,19 +317,6 @@ public bool TryDequeue( Volatile.Read(ref _header->FailedDequeues)); } - private static int RoundUpToPowerOf2(int value) - { - if (value > 1 << 30) - throw new ArgumentOutOfRangeException(nameof(value)); - value--; - value |= value >> 1; - value |= value >> 2; - value |= value >> 4; - value |= value >> 8; - value |= value >> 16; - return value + 1; - } - private static int RoundUpToMultiple(int value, int multiple) => checked((value + multiple - 1) / multiple * multiple); @@ -343,7 +333,6 @@ public void Dispose() return; _memoryHandle.Dispose(); _region.Dispose(); - GC.SuppressFinalize(this); } } } diff --git a/InterprocessMemory/IMemoryRegion.cs b/InterprocessMemory/IMemoryRegion.cs index ef36713..22bbdc3 100644 --- a/InterprocessMemory/IMemoryRegion.cs +++ b/InterprocessMemory/IMemoryRegion.cs @@ -74,6 +74,13 @@ public readonly struct LockOwnerInfo /// Gets whether the lock is orphaned (owner process died) /// public bool IsOrphan { get; init; } + + /// + /// Gets the number of read locks currently held across all processes. Read locks are not + /// attributed to an owner, so a process that dies while holding one leaves this count + /// permanently above zero (writers then time out); see . + /// + public int ReaderCount { get; init; } } /// @@ -131,13 +138,21 @@ public interface IMemoryRegion : IDisposable Memory GetMemory(long offset, int length); /// - /// Tries to acquire an exclusive write lock with timeout + /// Tries to acquire an exclusive write lock with timeout. + /// The wait is released with if the region is disposed + /// by another thread while waiting. /// bool TryAcquireWriteLock(TimeSpan timeout); /// - /// Releases the write lock + /// Releases the write lock. The lock is owned by the acquiring thread of the acquiring + /// process, so it must be released on that same thread: do not await between + /// acquiring and releasing it. /// + /// + /// The calling thread does not own the write lock, or the lock was taken over (for example by + /// orphan-lock recovery) before it was released. + /// void ReleaseWriteLock(); /// diff --git a/InterprocessMemory/MemoryRegion.cs b/InterprocessMemory/MemoryRegion.cs index e0845b8..91393c4 100644 --- a/InterprocessMemory/MemoryRegion.cs +++ b/InterprocessMemory/MemoryRegion.cs @@ -59,15 +59,67 @@ public sealed unsafe class MemoryRegion : IMemoryRegion private volatile int _disposed; + // Threads currently spinning inside TryAcquireWriteLock/TryAcquireReadLock. Those waits can be + // unbounded, so they are the one place where another thread can realistically call Dispose() + // while the shared header is being touched; Dispose() waits for this count to drain before it + // unmaps the view. Read/Write are not tracked: they are short, and two interlocked operations + // per call would cost more than they protect. + private int _activeWaiters; + + // How often a lock waiter re-probes whether the owning process is still alive. + private const long OrphanCheckIntervalMs = 250; + + // Upper bound on how long Dispose() waits for lock waiters to notice the disposed flag. + private static readonly TimeSpan s_waiterDrainTimeout = TimeSpan.FromSeconds(5); + + // Ticks. See DisposeGracePeriod. + private static long s_disposeGraceTicks = TimeSpan.FromMilliseconds(10).Ticks; + + /// + /// How long keeps the memory mapped after the region has been marked + /// disposed, so that operations on other threads that are already past their disposed check can + /// finish before the view is unmapped (default: 10 ms; disables it). + /// This is a process-wide setting. + /// + /// It is a best-effort mitigation, not a guarantee. The lock-free members (Read, + /// Write, the typed queues, SharedArray) deliberately do no per-call bookkeeping, so + /// cannot know whether another thread is still inside one. A thread that is + /// descheduled for longer than this period at exactly that moment would touch unmapped memory, + /// which terminates the process with an . The only complete + /// protection is to stop and join every thread that uses an instance before disposing it. + /// + /// + /// The value is negative. + public static TimeSpan DisposeGracePeriod + { + get => TimeSpan.FromTicks(Interlocked.Read(ref s_disposeGraceTicks)); + set + { + if (value < TimeSpan.Zero) + throw new ArgumentOutOfRangeException(nameof(value), "DisposeGracePeriod must not be negative."); + + Interlocked.Exchange(ref s_disposeGraceTicks, value.Ticks); + } + } + // Cached once per process: stamped into the header at lock acquire so an orphan check // can distinguish "same PID, same process" from "same PID, recycled by the OS for an - // unrelated process". Process.StartTime can throw under restricted permissions (Linux - // containers without /proc, certain Windows ACLs) — in that case we store 0 and the - // orphan check silently falls back to PID-only matching. + // unrelated process". The start time can be unreadable under restricted permissions + // (Linux containers without /proc, certain Windows ACLs) — in that case we store 0 and + // the orphan check silently falls back to PID-only matching. + // + // Windows: Process.StartTime.ToBinary() (the kernel creation time, identical for every observer). + // Linux: the kernel start tick count from /proc//stat. Process.StartTime must NOT be used + // there: every process derives it from its own wall-clock boot-time snapshot, so the value the + // owner records and the value another process computes for the same owner differ by + // milliseconds and would make every live owner look like an impostor. private static readonly long s_processStartTimeBinary = TryCaptureProcessStartTime(); private static long TryCaptureProcessStartTime() { + if (OperatingSystem.IsLinux()) + return TryReadLinuxStartTicks(Environment.ProcessId); + try { using var p = Process.GetCurrentProcess(); @@ -79,6 +131,53 @@ private static long TryCaptureProcessStartTime() } } + /// + /// Reads field 22 (starttime, clock ticks since boot) of /proc/<pid>/stat. + /// Returns 0 when it cannot be read. + /// + private static long TryReadLinuxStartTicks(int pid) + { + try + { + string stat = File.ReadAllText("/proc/" + pid + "/stat"); + + // The command name (field 2) is parenthesised and may itself contain spaces or + // parentheses, so split only what follows the LAST ')'. The first token after it + // is field 3, which makes field 22 index 19. + int commEnd = stat.LastIndexOf(')'); + if (commEnd < 0) + return 0; + + string[] fields = stat.Substring(commEnd + 2).Split(' '); + return fields.Length > 19 && long.TryParse(fields[19], out long ticks) ? ticks : 0; + } + catch + { + return 0; + } + } + + /// + /// True when the process currently using is provably not the one + /// that recorded (PID reuse). + /// + private static bool IsOwnerStartTimeMismatch(Process process, int ownerPid, long storedStartTime) + { + if (OperatingSystem.IsLinux()) + { + // 3.0.0 recorded Process.StartTime.ToBinary() here, which is negative for a local + // DateTime and not comparable across processes. Tick counts are positive. Treat the + // legacy format as "unknown" and keep the PID-only decision. + if (storedStartTime < 0) + return false; + + long currentTicks = TryReadLinuxStartTicks(ownerPid); + return currentTicks != 0 && currentTicks != storedStartTime; + } + + return process.StartTime.ToBinary() != storedStartTime; + } + /// /// Extended header structure with orphan lock detection support. /// WriterLockState and ReaderCount are placed in separate 64-byte cache lines @@ -178,6 +277,54 @@ internal static MemoryRegion OpenExisting( RegionKind regionKind) => new(name, capacityBytes: null, options, createOrOpen: false, regionKind); + /// + /// Deletes the backing storage of a named region so that the next + /// starts from scratch. + /// Use it to get rid of a region that a crash left unusable, for example one whose creator died + /// during initialization, or to change the capacity or element type of an existing region. + /// + /// Linux keeps the region as a file in /dev/shm that outlives its users, so it has to be + /// removed explicitly. Windows named sections are reference counted by the kernel and vanish when + /// the last handle closes; for them (without ) this + /// returns false. + /// + /// + /// Stop every process that uses the region first. Processes that still have it mapped keep + /// working on the removed storage, while later openers get a new, independent region. + /// This applies to every region kind (typed queues, arrays, structured memory), all of which are + /// addressed by the same name. + /// + /// + /// The region name that was passed to CreateOrOpen. + /// + /// Pass the same options (in particular ) that the region + /// was created with when it is file backed. + /// + /// true when backing storage was deleted; false when there was nothing to delete. + public static bool Remove(string name, MemoryRegionOptions? options = null) + { + ValidateFlatName(name); + + string? path = options?.FilePath; + if (string.IsNullOrEmpty(path)) + { + if (OperatingSystem.IsWindows()) + return false; + + if (!OperatingSystem.IsLinux()) + throw new PlatformNotSupportedException( + "MemoryRegion requires Windows or Linux unless MemoryRegionOptions.FilePath is used."); + + path = "/dev/shm/" + name; + } + + if (!File.Exists(path)) + return false; + + File.Delete(path); + return true; + } + internal MemoryRegion(string name, MemoryRegionOptions? options = null) : this( name, @@ -466,7 +613,9 @@ private bool InitializeOrOpen() } if (sw.Elapsed > TimeSpan.FromSeconds(5)) throw new TimeoutException( - "Timed out waiting for shared memory to be initialized by another process"); + "Timed out waiting for shared memory to be initialized by another process. " + + "If that process crashed during initialization, stop every user of the region " + + "and call MemoryRegion.Remove(name)."); Thread.SpinWait(100); } @@ -715,14 +864,30 @@ public bool TryAcquireWriteLock(TimeSpan timeout) ThrowIfDisposed(); TimeoutHelper.Validate(timeout, nameof(timeout)); + EnterWait(); + try + { + return TryAcquireWriteLockCore(timeout); + } + finally + { + ExitWait(); + } + } + + private bool TryAcquireWriteLockCore(TimeSpan timeout) + { var header = (SharedHeader*)_basePtr; - var sw = Stopwatch.StartNew(); + long start = Stopwatch.GetTimestamp(); // not a Stopwatch instance: that would allocate on every lock var spinner = new SpinWait(); - bool orphanCheckDone = false; - bool orphanCheckNearTimeout = false; + long nextOrphanCheckMs = 0; while (true) { + // Dispose() unmaps the header while we may still be spinning on it. It waits for + // registered waiters (EnterWait) to leave, and we leave as soon as we see the flag. + ThrowIfDisposed(); + if (Interlocked.CompareExchange(ref header->WriterLockState, 1, 0) == 0) { bool success = false; @@ -739,7 +904,9 @@ public bool TryAcquireWriteLock(TimeSpan timeout) var readerSpinner = new SpinWait(); while (Volatile.Read(ref header->ReaderCount) > 0) { - if (TimeoutHelper.HasExpired(sw, timeout)) + ThrowIfDisposed(); + + if (TimeoutHelper.HasExpired(start, timeout)) { return false; // Will release lock in finally } @@ -765,28 +932,21 @@ public bool TryAcquireWriteLock(TimeSpan timeout) } } - if (TimeoutHelper.HasExpired(sw, timeout)) + if (TimeoutHelper.HasExpired(start, timeout)) return false; - if (_options.EnableOrphanLockDetection) + if (_options.EnableOrphanLockDetection && ElapsedMilliseconds(start) >= nextOrphanCheckMs) { - // Check on first CAS failure; re-check when nearing timeout (≥75% elapsed) - // so a lock that becomes orphaned mid-wait is still recovered before giving up. - bool nearTimeout = TimeoutHelper.IsNearExpiry(sw, timeout, 0.75); + // Check on the first CAS failure and then periodically. The owner may die at any + // point while we wait, and a wait with Timeout.InfiniteTimeSpan has no deadline to + // key a one-off re-check on. The probe costs a process lookup, hence the interval. + nextOrphanCheckMs = ElapsedMilliseconds(start) + OrphanCheckIntervalMs; - if (!orphanCheckDone || (nearTimeout && !orphanCheckNearTimeout)) + if (IsWriteLockOrphaned()) { - if (!orphanCheckDone) - orphanCheckDone = true; - else - orphanCheckNearTimeout = true; - - if (IsWriteLockOrphaned()) - { - _logger?.LogWarning("Detected orphan write lock, attempting recovery"); - TryForceReleaseWriteLock(); - continue; - } + _logger?.LogWarning("Detected orphan write lock, attempting recovery"); + TryForceReleaseWriteLock(); + continue; } } @@ -805,21 +965,28 @@ public void ReleaseWriteLock() long currentThreadId = Environment.CurrentManagedThreadId; int ownerPid = Volatile.Read(ref header->LockOwnerProcessId); long ownerThreadId = Volatile.Read(ref header->LockOwnerThreadId); + // Releasing a lock this thread does not own must be loud. Silently ignoring it (the + // previous behaviour) turned `await` inside a lock scope — the continuation resumes on + // another thread — into a write lock that no process could ever release again. if (ownerPid != currentPid || ownerThreadId != currentThreadId) { _logger?.LogWarning( - "ReleaseWriteLock called from PID {Pid}/thread {ThreadId} but lock owner is PID {OwnerPid}/thread {OwnerThreadId} — ignored", + "ReleaseWriteLock called from PID {Pid}/thread {ThreadId} but lock owner is PID {OwnerPid}/thread {OwnerThreadId}", currentPid, currentThreadId, ownerPid, ownerThreadId); - return; + throw new SynchronizationLockException( + $"The write lock must be released by the thread that acquired it " + + $"(caller PID {currentPid}/thread {currentThreadId}, lock owner PID {ownerPid}/thread {ownerThreadId}). " + + "Do not await inside a lock scope."); } int prev = Interlocked.CompareExchange(ref header->LockOwnerProcessId, 0, currentPid); if (prev != currentPid) { _logger?.LogWarning( - "ReleaseWriteLock called from PID {Pid} but lock owner is {OwnerPid} — ignored", + "ReleaseWriteLock called from PID {Pid} but lock owner is {OwnerPid}", currentPid, prev); - return; + throw new SynchronizationLockException( + "The write lock was taken over (for example by orphan-lock recovery) before it was released."); } header->LockOwnerThreadId = 0; @@ -838,19 +1005,35 @@ public bool TryAcquireReadLock(TimeSpan timeout) ThrowIfDisposed(); TimeoutHelper.Validate(timeout, nameof(timeout)); + EnterWait(); + try + { + return TryAcquireReadLockCore(timeout); + } + finally + { + ExitWait(); + } + } + + private bool TryAcquireReadLockCore(TimeSpan timeout) + { var header = (SharedHeader*)_basePtr; - var sw = Stopwatch.StartNew(); + long start = Stopwatch.GetTimestamp(); var spinner = new SpinWait(); while (true) { + // See TryAcquireWriteLockCore: leave promptly once Dispose() has started. + ThrowIfDisposed(); + // Fast path: peek the writer flag without any atomic. If a writer is active, // wait — touching ReaderCount unnecessarily would create cache-line traffic on // the reader-side line and prolong the writer's release-then-drain phase. int writerState = Volatile.Read(ref header->WriterLockState); if (writerState != 0) { - if (TimeoutHelper.HasExpired(sw, timeout)) + if (TimeoutHelper.HasExpired(start, timeout)) return false; spinner.SpinOnce(); continue; @@ -872,7 +1055,7 @@ public bool TryAcquireReadLock(TimeSpan timeout) // twice extra — same penalty as the previous design's CAS-rollback path. Interlocked.Decrement(ref header->ReaderCount); - if (TimeoutHelper.HasExpired(sw, timeout)) + if (TimeoutHelper.HasExpired(start, timeout)) return false; spinner.SpinOnce(); @@ -924,14 +1107,13 @@ public bool IsWriteLockOrphaned() // PID-reuse defense: even when a process with this PID exists, it might be an // unrelated process that the OS recycled the PID for after the real owner died. - // Compare the captured StartTime; mismatch ⇒ impostor ⇒ orphan. + // Compare the captured start time; mismatch ⇒ impostor ⇒ orphan. long storedStartTime = header->LockOwnerProcessStartTime; if (storedStartTime != 0) { try { - long currentStartTime = process.StartTime.ToBinary(); - if (currentStartTime != storedStartTime) + if (IsOwnerStartTimeMismatch(process, ownerPid, storedStartTime)) { _logger?.LogWarning( "Lock owner PID {Pid} still exists but its StartTime differs (orphan from PID reuse)", @@ -1025,6 +1207,40 @@ public bool TryForceReleaseWriteLock() return true; } + /// + /// Unconditionally clears the shared write lock and the read-lock count. + /// + /// This is the recovery tool for a lock state that cannot be recovered automatically. Read + /// locks are not attributed to an owner, so when a process dies while holding one the + /// shared reader count stays above zero forever (see ) + /// and every writer times out, including after the region is reopened on Linux, where the + /// backing file outlives its users. + /// + /// + /// Call it only when no process is inside a critical section of this region; resetting locks + /// that are legitimately held lets a writer run concurrently with them. A thread that held + /// a lock when it was reset gets from its release. + /// + /// + public void ForceResetLocks() + { + ThrowIfDisposed(); + + var header = (SharedHeader*)_basePtr; + + _logger?.LogWarning( + "Force resetting locks of '{Name}' (writer state {WriterState}, readers {Readers})", + _name, Volatile.Read(ref header->WriterLockState), Volatile.Read(ref header->ReaderCount)); + + Volatile.Write(ref header->LockOwnerProcessId, 0); + header->LockOwnerThreadId = 0; + header->LockOwnerProcessStartTime = 0; + header->LockAcquiredTimestamp = 0; + Thread.MemoryBarrier(); + Volatile.Write(ref header->WriterLockState, 0); + Interlocked.Exchange(ref header->ReaderCount, 0); + } + /// public LockOwnerInfo GetLockOwnerInfo() { @@ -1037,7 +1253,8 @@ public LockOwnerInfo GetLockOwnerInfo() ProcessId = header->LockOwnerProcessId, ThreadId = header->LockOwnerThreadId, AcquiredTimestamp = header->LockAcquiredTimestamp, - IsOrphan = IsWriteLockOrphaned() + IsOrphan = IsWriteLockOrphaned(), + ReaderCount = Volatile.Read(ref header->ReaderCount) }; } @@ -1105,7 +1322,35 @@ private void ThrowIfDisposed() } /// - /// Releases all resources used by this buffer + /// Registers the calling thread as a lock waiter. Increment first, then check the flag: with + /// the full fences of the two interlocked operations, either this thread sees the disposed + /// flag or sees this thread in . + /// + private void EnterWait() + { + Interlocked.Increment(ref _activeWaiters); + if (_disposed != 0) + { + Interlocked.Decrement(ref _activeWaiters); + throw new ObjectDisposedException(nameof(MemoryRegion)); + } + } + + private void ExitWait() => Interlocked.Decrement(ref _activeWaiters); + + private static long ElapsedMilliseconds(long startTimestamp) => + (long)Stopwatch.GetElapsedTime(startTimestamp).TotalMilliseconds; + + /// + /// Releases all resources used by this buffer. Threads blocked in + /// or are released with an + /// before the memory is unmapped. + /// + /// Stop and join every thread that uses this instance before disposing it. Other members check + /// the disposed flag but are not tracked, so a thread that is inside one of them while the + /// memory is unmapped crashes the process; narrows that window + /// but does not close it. + /// /// public void Dispose() { @@ -1113,6 +1358,23 @@ public void Dispose() return; _logger?.LogDebug("Disposing shared buffer '{Name}'", _name); + + // Waiters observe _disposed on every spin iteration and leave within a few milliseconds. + // Never unmap underneath a thread that is still dereferencing the header: that is an + // AccessViolationException, which terminates the process and cannot be caught. + var drain = Stopwatch.StartNew(); + var drainSpinner = new SpinWait(); + while (Volatile.Read(ref _activeWaiters) > 0 && drain.Elapsed < s_waiterDrainTimeout) + drainSpinner.SpinOnce(); + + // Calls that passed their disposed check just before the flag was set are still running + // against the mapping and are not tracked (that would cost every call two interlocked + // operations; measured at roughly 9x on a queue round trip). New calls fail fast with + // ObjectDisposedException, so a short pause lets those in-flight calls complete. + long graceTicks = Interlocked.Read(ref s_disposeGraceTicks); + if (graceTicks > 0) + Thread.Sleep(TimeSpan.FromTicks(graceTicks)); + Cleanup(disposing: true); GC.SuppressFinalize(this); } diff --git a/InterprocessMemory/MemoryRegionOptions.cs b/InterprocessMemory/MemoryRegionOptions.cs index 9d685b3..7aa01bb 100644 --- a/InterprocessMemory/MemoryRegionOptions.cs +++ b/InterprocessMemory/MemoryRegionOptions.cs @@ -20,7 +20,9 @@ public sealed class MemoryRegionOptions public static readonly TimeSpan DefaultLockTimeout = TimeSpan.FromSeconds(5); /// - /// Default orphan lock timeout: 30 seconds + /// A reasonable value to pass as when you opt in to + /// time-based lock takeover (30 seconds). It is NOT the default: + /// is (disabled) unless you set it. /// public static readonly TimeSpan DefaultOrphanLockTimeout = TimeSpan.FromSeconds(30); @@ -45,14 +47,24 @@ public sealed class MemoryRegionOptions public string? FilePath { get; set; } /// - /// Gets or sets whether to enable orphan lock detection and recovery + /// Gets or sets whether a waiting process recovers a write lock whose owner process has exited + /// (or whose PID was reused by another process) /// public bool EnableOrphanLockDetection { get; set; } = true; /// - /// Gets or sets the timeout after which a lock is considered orphaned + /// Gets or sets the time after which a write lock held by a process that is still ALIVE is + /// also treated as orphaned, so that a waiter may take it over. + /// (the default) disables this. + /// + /// A lock whose owner process has exited is always recovered while + /// is on. A time limit additionally breaks mutual + /// exclusion for a healthy owner that simply holds the lock for longer than the limit (a long + /// transaction, a paused debugger), so enable it only when every critical section is known to + /// be shorter than the value you choose. + /// /// - public TimeSpan OrphanLockTimeout { get; set; } = DefaultOrphanLockTimeout; + public TimeSpan OrphanLockTimeout { get; set; } = TimeSpan.Zero; /// /// Gets or sets whether to enable checksum verification diff --git a/InterprocessMemory/PowerOfTwo.cs b/InterprocessMemory/PowerOfTwo.cs new file mode 100644 index 0000000..5a5f392 --- /dev/null +++ b/InterprocessMemory/PowerOfTwo.cs @@ -0,0 +1,26 @@ +using System; +using System.Numerics; + +namespace InterprocessMemory +{ + internal static class PowerOfTwo + { + /// The largest power of two that fits in a positive . + public const int MaxInt32 = 1 << 30; + + /// + /// Rounds a requested item count up to the next power of two, which the ring buffers need so + /// they can wrap with a mask instead of a modulo. + /// + public static int RoundUp(int value, string paramName) + { + if (value <= 0 || value > MaxInt32) + throw new ArgumentOutOfRangeException( + paramName, + value, + $"The value must be between 1 and {MaxInt32} so it can be rounded up to a power of two."); + + return (int)BitOperations.RoundUpToPowerOf2((uint)value); + } + } +} diff --git a/InterprocessMemory/SharedArray.cs b/InterprocessMemory/SharedArray.cs index 87ade64..f52f12f 100644 --- a/InterprocessMemory/SharedArray.cs +++ b/InterprocessMemory/SharedArray.cs @@ -64,24 +64,39 @@ private SharedArray(string name, int? length, bool createOrOpen) _buffer = MemoryRegion.CreateOrOpen( name, checked(ArrayHeaderSize + dataSize), - options: null, + CreateRegionOptions(), RegionKind.SharedArray); - - if (_buffer.IsOwner) - InitializeHeader(); - else - ValidateAndLoadHeader(expectedLength: _length); } else { _buffer = MemoryRegion.OpenExisting( name, - options: null, + CreateRegionOptions(), RegionKind.SharedArray); - ValidateAndLoadHeader(expectedLength: null); + } + + // The region is already mapped here. Opening an array of another element type or length + // throws from the header check, and without this the mapping and (on Linux) its file + // descriptor would stay open until the finalizer runs. + try + { + if (createOrOpen && _buffer.IsOwner) + InitializeHeader(); + else + ValidateAndLoadHeader(expectedLength: createOrOpen ? _length : null); + } + catch + { + _buffer.Dispose(); + throw; } } + // The array exposes no statistics, so the region's per-call counters would only cost an + // interlocked operation on every element access. + private static MemoryRegionOptions CreateRegionOptions() => + new() { EnableStatistics = false }; + private void InitializeHeader() { Span header = stackalloc byte[ArrayHeaderSize]; @@ -92,6 +107,10 @@ private void InitializeHeader() BinaryPrimitives.WriteUInt64LittleEndian(header.Slice(16), _fingerprint.Low); BinaryPrimitives.WriteUInt64LittleEndian(header.Slice(24), _fingerprint.High); _buffer.Write(header, 0); + + // Publish the magic last. Region.Write does no fencing, so on a weakly ordered CPU another + // process could otherwise see the magic before the fields it announces. + Thread.MemoryBarrier(); BinaryPrimitives.WriteUInt32LittleEndian(header, ArrayMagic); _buffer.Write(header.Slice(0, sizeof(uint)), 0); } @@ -99,17 +118,23 @@ private void InitializeHeader() private void ValidateAndLoadHeader(int? expectedLength) { Span header = stackalloc byte[ArrayHeaderSize]; + Span magic = stackalloc byte[sizeof(uint)]; var sw = Stopwatch.StartNew(); while (true) { - _buffer.Read(header, 0); - if (BinaryPrimitives.ReadUInt32LittleEndian(header) == ArrayMagic) + _buffer.Read(magic, 0); + if (BinaryPrimitives.ReadUInt32LittleEndian(magic) == ArrayMagic) break; if (sw.Elapsed > TimeSpan.FromSeconds(5)) throw new InvalidDataException("Timed out waiting for the shared-array header."); Thread.SpinWait(100); } + // Read the fields only after the magic has been seen, never in the same copy: a single + // copy gives no ordering between the magic and the bytes that follow it. + Thread.MemoryBarrier(); + _buffer.Read(header, 0); + int version = BinaryPrimitives.ReadInt32LittleEndian(header.Slice(4)); int storedLength = BinaryPrimitives.ReadInt32LittleEndian(header.Slice(8)); int storedElementSize = BinaryPrimitives.ReadInt32LittleEndian(header.Slice(12)); @@ -287,20 +312,8 @@ public void Dispose() if (Interlocked.Exchange(ref _disposed, 1) != 0) return; + // No finalizer: if Dispose is never called, the MemoryRegion's own finalizer unmaps the memory. _buffer?.Dispose(); - GC.SuppressFinalize(this); - } - - /// - /// Releases unmanaged resources if Dispose was not called. Does NOT proactively dispose - /// the inner — that has its own finalizer and - /// touching it from here risks running against an already-finalized peer (finalizer - /// order is undefined). The peer's finalizer reclaims its unmanaged handles directly. - /// - ~SharedArray() - { - // Just mark disposed so a racing manual Dispose is a no-op. No managed work here. - Interlocked.Exchange(ref _disposed, 1); } } } diff --git a/InterprocessMemory/SingleProducerByteStream.cs b/InterprocessMemory/SingleProducerByteStream.cs index 2c9b93e..846b488 100644 --- a/InterprocessMemory/SingleProducerByteStream.cs +++ b/InterprocessMemory/SingleProducerByteStream.cs @@ -431,28 +431,10 @@ public void Dispose() if (Interlocked.Exchange(ref _disposed, 1) != 0) return; - // MemoryHandle wraps an unmanaged pointer — safe to dispose anywhere. + // The MemoryHandle wraps an unmanaged pointer and owns nothing, so there is no finalizer + // here: if Dispose is never called, the MemoryRegion's own finalizer unmaps the memory. _memoryHandle.Dispose(); - // _buffer (MemoryRegion) is a managed object with its OWN finalizer. - // We only proactively dispose it on the deterministic path; from our finalizer we - // let the GC handle it to avoid touching a possibly-already-finalized peer. _buffer?.Dispose(); - GC.SuppressFinalize(this); - } - - /// - /// Releases unmanaged resources if Dispose was not called. - /// Skips the managed _buffer.Dispose() — its own finalizer reclaims it. - /// - ~SingleProducerByteStream() - { - // Guard against double-dispose if Dispose() already ran. _disposed is volatile so - // we observe its current value here. - if (Interlocked.Exchange(ref _disposed, 1) != 0) - return; - try - { _memoryHandle.Dispose(); } - catch { /* unmanaged release; best-effort */ } } [MethodImpl(MethodImplOptions.AggressiveInlining)] diff --git a/InterprocessMemory/SingleProducerQueue.cs b/InterprocessMemory/SingleProducerQueue.cs index a8e0284..0b8faf4 100644 --- a/InterprocessMemory/SingleProducerQueue.cs +++ b/InterprocessMemory/SingleProducerQueue.cs @@ -71,7 +71,7 @@ private SingleProducerQueue(string name, int? capacity, bool createOrOpen) if (createOrOpen) { - _capacity = RoundUpToPowerOf2(capacity!.Value); + _capacity = PowerOfTwo.RoundUp(capacity!.Value, "capacity"); _capacityMask = _capacity - 1; long regionCapacity = checked(HeaderSize + (long)_capacity * _elementSize); if (regionCapacity > int.MaxValue) @@ -132,10 +132,13 @@ private void ValidateAndLoad(int? requestedCapacity) _header->FingerprintHigh != _fingerprint.High) throw new InvalidDataException("The queue has a different format or element type."); - if (requestedCapacity.HasValue && - RoundUpToPowerOf2(requestedCapacity.Value) != storedCapacity) - throw new InvalidOperationException( - $"Capacity mismatch: expected {RoundUpToPowerOf2(requestedCapacity.Value)}, found {storedCapacity}."); + if (requestedCapacity.HasValue) + { + int expectedCapacity = PowerOfTwo.RoundUp(requestedCapacity.Value, "capacity"); + if (expectedCapacity != storedCapacity) + throw new InvalidOperationException( + $"Capacity mismatch: expected {expectedCapacity}, found {storedCapacity}."); + } long expectedRegionSize = checked(HeaderSize + (long)storedCapacity * _elementSize); if (_region.Capacity != expectedRegionSize) @@ -224,19 +227,6 @@ public bool TryDequeue( return true; } - private static int RoundUpToPowerOf2(int value) - { - if (value > 1 << 30) - throw new ArgumentOutOfRangeException(nameof(value)); - value--; - value |= value >> 1; - value |= value >> 2; - value |= value >> 4; - value |= value >> 8; - value |= value >> 16; - return value + 1; - } - [MethodImpl(MethodImplOptions.AggressiveInlining)] private void ThrowIfDisposed() { @@ -250,7 +240,6 @@ public void Dispose() return; _memoryHandle.Dispose(); _region.Dispose(); - GC.SuppressFinalize(this); } } } diff --git a/InterprocessMemory/StructuredMemory.cs b/InterprocessMemory/StructuredMemory.cs index 6d4a65a..39e6b7a 100644 --- a/InterprocessMemory/StructuredMemory.cs +++ b/InterprocessMemory/StructuredMemory.cs @@ -40,6 +40,11 @@ public sealed class StructuredMemory : IDisposable where TSchema : stru private readonly ThreadLocal _writeLockDepth = new(() => 0); private readonly ThreadLocal _readLockDepth = new(() => 0); + // A method group is converted to a new delegate on every use, which allocated on each + // automatically locked read or write. Convert once per instance instead. + private readonly Action _decrementWriteLockDepth; + private readonly Action _decrementReadLockDepth; + /// /// Gets the schema instance defining the memory layout /// @@ -94,6 +99,8 @@ internal StructuredMemory( _schema = schema; _compatibility = compatibility; + _decrementWriteLockDepth = DecrementWriteLockDepth; + _decrementReadLockDepth = DecrementReadLockDepth; _fields = BuildFieldMetadata(schema); if (_fields.Count == 0) @@ -107,15 +114,18 @@ internal StructuredMemory( long totalSize = SchemaHeaderSize + CalculateTotalSize(_fields); + // No statistics are exposed here, so skip the region's per-call counters. + var regionOptions = new MemoryRegionOptions { EnableStatistics = false }; + if (create) { _buffer = MemoryRegion.CreateOrOpen( - name, totalSize, options: null, RegionKind.StructuredMemory); + name, totalSize, regionOptions, RegionKind.StructuredMemory); } else { _buffer = MemoryRegion.OpenExisting( - name, options: null, RegionKind.StructuredMemory); + name, regionOptions, RegionKind.StructuredMemory); if (_buffer.Capacity != totalSize) { _buffer.Dispose(); @@ -720,6 +730,11 @@ public WriteLock AcquireWriteLock() /// Lock acquisition timeout /// A disposable lock guard that releases the lock on dispose /// Thrown when the lock cannot be acquired within the timeout + /// + /// The lock is owned by the calling thread. Dispose the guard on that same thread; do not + /// await inside the guarded scope, or throws + /// . + /// public WriteLock AcquireWriteLock(TimeSpan timeout) { ThrowIfDisposed(); @@ -729,14 +744,14 @@ public WriteLock AcquireWriteLock(TimeSpan timeout) if (_writeLockDepth.Value > 0) { IncrementWriteLockDepth(); - return new WriteLock(null, DecrementWriteLockDepth); + return new WriteLock(null, _decrementWriteLockDepth); } if (!_buffer.TryAcquireWriteLock(timeout)) throw new TimeoutException($"Failed to acquire write lock within {timeout}"); IncrementWriteLockDepth(); - return new WriteLock(_buffer, DecrementWriteLockDepth); + return new WriteLock(_buffer, _decrementWriteLockDepth); } /// @@ -759,6 +774,10 @@ public ReadLock AcquireReadLock() /// Lock acquisition timeout /// A disposable lock guard that releases the lock on dispose /// Thrown when the lock cannot be acquired within the timeout + /// + /// Dispose the guard on the thread that acquired it; do not await inside the guarded + /// scope, or throws . + /// public ReadLock AcquireReadLock(TimeSpan timeout) { ThrowIfDisposed(); @@ -768,14 +787,27 @@ public ReadLock AcquireReadLock(TimeSpan timeout) if (_readLockDepth.Value > 0 || _writeLockDepth.Value > 0) { IncrementReadLockDepth(); - return new ReadLock(null, DecrementReadLockDepth); + return new ReadLock(null, _decrementReadLockDepth); } if (!_buffer.TryAcquireReadLock(timeout)) throw new TimeoutException($"Failed to acquire read lock within {timeout}"); IncrementReadLockDepth(); - return new ReadLock(_buffer, DecrementReadLockDepth); + return new ReadLock(_buffer, _decrementReadLockDepth); + } + + /// + /// Unconditionally clears the cross-process write lock and read-lock count of this region. + /// A process that dies while holding a read lock leaves the shared reader count above zero + /// forever, which makes every later writer time out and, on Linux, survives reopening the region. + /// Call this only when no process is inside a critical section of the region. + /// See . + /// + public void ForceResetLocks() + { + ThrowIfDisposed(); + ((MemoryRegion)_buffer).ForceResetLocks(); } /// @@ -807,6 +839,9 @@ private void WriteSchemaHeader() BitConverter.TryWriteBytes(header.Slice(12), _schemaHash); _buffer.Write(header, 0); + + // Publish the magic last (see SharedArray.InitializeHeader). + Thread.MemoryBarrier(); BitConverter.TryWriteBytes(header, SchemaMagic); _buffer.Write(header.Slice(0, sizeof(uint)), 0); } @@ -814,11 +849,12 @@ private void WriteSchemaHeader() private void ValidateSchemaCompatibility() { Span header = stackalloc byte[SchemaHeaderSize]; + Span magic = stackalloc byte[sizeof(uint)]; var sw = Stopwatch.StartNew(); while (true) { - _buffer.Read(header, 0); - if (BitConverter.ToUInt32(header) == SchemaMagic) + _buffer.Read(magic, 0); + if (BitConverter.ToUInt32(magic) == SchemaMagic) break; if (sw.Elapsed > TimeSpan.FromSeconds(5)) throw new InvalidDataException( @@ -826,6 +862,10 @@ private void ValidateSchemaCompatibility() Thread.SpinWait(100); } + // Read the fields only after the magic has been seen (see SharedArray.ValidateAndLoadHeader). + Thread.MemoryBarrier(); + _buffer.Read(header, 0); + StoredSchemaVersion = BitConverter.ToInt32(header.Slice(4)); int storedFieldCount = BitConverter.ToInt32(header.Slice(8)); int storedHash = BitConverter.ToInt32(header.Slice(12)); @@ -1115,17 +1155,6 @@ public void Dispose() _buffer?.Dispose(); _writeLockDepth.Dispose(); _readLockDepth.Dispose(); - GC.SuppressFinalize(this); - } - - /// - /// Releases unmanaged resources if Dispose was not called - /// - ~StructuredMemory() - { - Interlocked.Exchange(ref _disposed, 1); - // MemoryRegion owns its unmanaged finalizer path. Avoid invoking - // managed Dispose logic from this finalizer. } private static void EnsureScalarField(FieldMetadata metadata) @@ -1152,24 +1181,46 @@ public struct WriteLock : IDisposable { private IMemoryRegion? _buffer; private Action? _onDispose; + private readonly int _ownerThreadId; internal WriteLock(IMemoryRegion? buffer, Action? onDispose = null) { _buffer = buffer; _onDispose = onDispose; + _ownerThreadId = Environment.CurrentManagedThreadId; } /// - /// Releases the write lock if not already released + /// Releases the write lock if not already released. /// + /// + /// Called on a different thread than the one that acquired the lock, which is what happens + /// when the guarded scope contains an await. Nothing is released in that case. + /// public void Dispose() { + if (_onDispose == null) + return; + + // The lock and its reentrancy depth belong to the acquiring thread. Releasing from + // another thread would leave that thread believing it still holds the lock. + if (Environment.CurrentManagedThreadId != _ownerThreadId) + throw new SynchronizationLockException( + "A write lock guard must be disposed on the thread that acquired it. " + + "Do not await inside a lock scope."); + var onDispose = Interlocked.Exchange(ref _onDispose, null); if (onDispose != null) { - _buffer?.ReleaseWriteLock(); - _buffer = null; - onDispose.Invoke(); + try + { + _buffer?.ReleaseWriteLock(); + } + finally + { + _buffer = null; + onDispose.Invoke(); + } } } } @@ -1186,24 +1237,44 @@ public struct ReadLock : IDisposable { private IMemoryRegion? _buffer; private Action? _onDispose; + private readonly int _ownerThreadId; internal ReadLock(IMemoryRegion? buffer, Action? onDispose = null) { _buffer = buffer; _onDispose = onDispose; + _ownerThreadId = Environment.CurrentManagedThreadId; } /// - /// Releases the read lock if not already released + /// Releases the read lock if not already released. /// + /// + /// Called on a different thread than the one that acquired the lock. Nothing is released + /// in that case. See . + /// public void Dispose() { + if (_onDispose == null) + return; + + if (Environment.CurrentManagedThreadId != _ownerThreadId) + throw new SynchronizationLockException( + "A read lock guard must be disposed on the thread that acquired it. " + + "Do not await inside a lock scope."); + var onDispose = Interlocked.Exchange(ref _onDispose, null); if (onDispose != null) { - _buffer?.ReleaseReadLock(); - _buffer = null; - onDispose.Invoke(); + try + { + _buffer?.ReleaseReadLock(); + } + finally + { + _buffer = null; + onDispose.Invoke(); + } } } } diff --git a/InterprocessMemory/TimeoutHelper.cs b/InterprocessMemory/TimeoutHelper.cs index 6f533be..05f4ed2 100644 --- a/InterprocessMemory/TimeoutHelper.cs +++ b/InterprocessMemory/TimeoutHelper.cs @@ -21,11 +21,13 @@ public static bool HasExpired(Stopwatch stopwatch, TimeSpan timeout) return timeout != Timeout.InfiniteTimeSpan && stopwatch.Elapsed > timeout; } - public static bool IsNearExpiry(Stopwatch stopwatch, TimeSpan timeout, double fraction) + /// + /// Same check against a start value. Unlike a + /// instance this does not allocate, which matters on paths that run per call. + /// + public static bool HasExpired(long startTimestamp, TimeSpan timeout) { - return timeout != Timeout.InfiniteTimeSpan - && timeout > TimeSpan.Zero - && stopwatch.Elapsed.TotalMilliseconds >= timeout.TotalMilliseconds * fraction; + return timeout != Timeout.InfiniteTimeSpan && Stopwatch.GetElapsedTime(startTimestamp) > timeout; } } } diff --git a/InterprocessMemory/TypeLayoutFingerprint.cs b/InterprocessMemory/TypeLayoutFingerprint.cs index 5f29fbd..1a11496 100644 --- a/InterprocessMemory/TypeLayoutFingerprint.cs +++ b/InterprocessMemory/TypeLayoutFingerprint.cs @@ -3,6 +3,7 @@ using System.Collections.Generic; using System.Linq; using System.Reflection; +using System.Runtime.CompilerServices; using System.Runtime.InteropServices; using System.Security.Cryptography; using System.Text; @@ -14,7 +15,21 @@ internal readonly record struct TypeLayoutFingerprint(ulong Low, ulong High) public static TypeLayoutFingerprint Create() where T : unmanaged { var descriptor = new StringBuilder(256); - AppendType(descriptor, typeof(T), new HashSet()); + try + { + AppendType(descriptor, typeof(T), new HashSet()); + } + catch (ArgumentException) + { + // Marshal.SizeOf/OffsetOf reject generic types (ValueTuple, KeyValuePair, ...), at any + // nesting depth, although they satisfy the unmanaged constraint. Describe those with + // the managed layout instead. Types the marshaller accepts keep the descriptor above, + // so their fingerprint does not change and processes built from different releases + // can still open the same region. + descriptor.Clear(); + AppendManagedType(descriptor, typeof(T), new HashSet()); + } + byte[] hash = SHA256.HashData(Encoding.UTF8.GetBytes(descriptor.ToString())); return new TypeLayoutFingerprint( BinaryPrimitives.ReadUInt64LittleEndian(hash), @@ -55,5 +70,65 @@ private static void AppendType(StringBuilder target, Type type, HashSet ac active.Remove(type); } + + // Unsafe.SizeOf() is the size the containers actually copy (bool = 1 byte), and it works for + // generic types. Unsafe.SizeOf has no generic constraint, so it can be bound to any Type. + private static readonly MethodInfo s_sizeOf = typeof(Unsafe).GetMethod(nameof(Unsafe.SizeOf))!; + + private static int ManagedSizeOf(Type type) => + (int)s_sizeOf.MakeGenericMethod(type).Invoke(null, null)!; + + // Version-independent type name. Type.FullName of a constructed generic type embeds the + // assembly-qualified name (including Version=) of every type argument, which differs between + // .NET releases and would make processes on different runtimes disagree. + private static void AppendTypeName(StringBuilder target, Type type) + { + target.Append(type.Assembly.GetName().Name).Append(':'); + if (!type.IsConstructedGenericType) + { + target.Append(type.FullName); + return; + } + + target.Append(type.GetGenericTypeDefinition().FullName).Append('['); + Type[] arguments = type.GetGenericArguments(); + for (int i = 0; i < arguments.Length; i++) + { + if (i > 0) + target.Append(','); + AppendTypeName(target, arguments[i]); + } + target.Append(']'); + } + + private static void AppendManagedType(StringBuilder target, Type type, HashSet active) + { + AppendTypeName(target, type); + target.Append('|').Append(ManagedSizeOf(type)); + + var layout = type.StructLayoutAttribute; + target.Append('|').Append((int)(layout?.Value ?? LayoutKind.Auto)) + .Append('|').Append(layout?.Pack ?? 0) + .Append('|').Append(layout?.Size ?? 0); + + if (type.IsPrimitive || type.IsEnum || !active.Add(type)) + return; + + // Declaration order (metadata token), explicit FieldOffset values, the total size and the + // layout kind together pin down the memory layout without Marshal.OffsetOf. + var fields = type + .GetFields(BindingFlags.Instance | BindingFlags.Public | BindingFlags.NonPublic) + .OrderBy(field => field.MetadataToken); + + foreach (var field in fields) + { + int? explicitOffset = field.GetCustomAttribute()?.Value; + target.Append(";f:").Append(field.Name).Append('@') + .Append(explicitOffset?.ToString() ?? "-").Append(':'); + AppendManagedType(target, field.FieldType, active); + } + + active.Remove(type); + } } } diff --git a/README.md b/README.md index 7a79557..236292a 100644 --- a/README.md +++ b/README.md @@ -59,7 +59,18 @@ if (memory.TryAcquireWriteLock(TimeSpan.FromSeconds(1))) ``` Lock ownership includes the process ID, managed thread ID, and process start time. A process -waiting for a write lock can recover a lock left behind by a terminated process. +waiting for a write lock keeps checking whether the owner is still alive (also when it waits +with `Timeout.InfiniteTimeSpan`) and recovers a lock left behind by a terminated process. + +Recovery is based on the owner process being gone. A write lock held by a process that is still +alive is never taken over unless you opt in with `MemoryRegionOptions.OrphanLockTimeout`, which +trades mutual exclusion for liveness when a critical section may run longer than the timeout. + +The write lock belongs to the thread that acquired it. Release it on that same thread; releasing +from another thread throws `SynchronizationLockException` and leaves the lock held. In particular, +do not `await` between acquiring and releasing a lock, because the continuation may resume on a +different thread. Disposing a region while another thread is waiting for one of its locks releases +that thread with an `ObjectDisposedException`. ## Choosing a data structure @@ -204,6 +215,10 @@ using (memory.AcquireWriteLock()) UTF-8 strings, and arrays to prevent torn reads and writes. Use an explicit lock when several fields form one transaction. +The lock guards returned by `AcquireWriteLock()` and `AcquireReadLock()` must be disposed on the +thread that acquired them. Keep the guarded scope synchronous: an `await` inside it lets the guard +be disposed on another thread, which throws `SynchronizationLockException`. + ## Shared arrays ```csharp @@ -226,6 +241,40 @@ modifying the existing bytes. Version 3 does not migrate live 2.x regions. Stop every 2.x process, remove the named/file-backed region, and recreate it with version 3. See [MIGRATION.md](MIGRATION.md). +## Disposing while other threads are running + +Stop and join every thread that uses an instance before you dispose it. The lock-free members +(`Read`, `Write`, the typed queues, `SharedArray`) read and write through a raw pointer and do no +per-call bookkeeping, because tracking calls costs every one of them two interlocked operations +(roughly 9 times slower for a queue round trip, and about 17 times slower for two threads +exchanging items). A call that is still running when the memory is unmapped terminates the process +with an `AccessViolationException`, which cannot be caught. + +What the library does to reduce the risk: + +- Calls made after `Dispose()` started fail with `ObjectDisposedException`. +- Threads waiting for a lock are released with `ObjectDisposedException`, and `Dispose()` waits for + them before it unmaps anything. +- `Dispose()` keeps the memory mapped for `MemoryRegion.DisposeGracePeriod` (10 ms by default, + process-wide, `TimeSpan.Zero` disables it) so that calls already in flight can finish. This is + best effort: a thread that is descheduled for longer than that at exactly the wrong moment still + hits unmapped memory. + +## Recovering after a crash + +Windows named sections disappear when their last handle closes. On Linux a region is a file in +`/dev/shm` that outlives its users, so a crash can leave state behind: + +- A process that dies while holding a **read lock** leaves the shared reader count above zero and + every writer times out. Read locks have no owner, so this cannot be detected automatically. + `GetLockOwnerInfo().ReaderCount` shows the stale count; once no process is inside a critical + section, call `MemoryRegion.ForceResetLocks()` (or `StructuredMemory.ForceResetLocks()`). +- A creator that dies during initialization makes every opener time out. A region with a different + capacity or element type than the one you now want is rejected as well. In both cases stop all + users and call `MemoryRegion.Remove(name)` (pass the same `MemoryRegionOptions` for file-backed + regions); the next `CreateOrOpen` starts from scratch. It works for every data structure because + they are all addressed by the same name. + ## Build and test ```shell