Skip to content
30 changes: 30 additions & 0 deletions InterprocessMemory.TestWorker/Program.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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
/// </summary>
if (args.Length < 2)
{
Expand All @@ -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}")
};

Expand Down Expand Up @@ -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);
Expand Down
15 changes: 8 additions & 7 deletions InterprocessMemory.Tests/ChangeVerificationTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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<SynchronizationLockException>(() => 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();
}
Expand Down
39 changes: 23 additions & 16 deletions InterprocessMemory.Tests/ConcurrentMessageQueueTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -370,8 +370,11 @@ public async Task MPMC_HighContention_ShouldMaintainIntegrity()
var producersDone = 0;
var errors = new ConcurrentBag<string>();

// 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
{
Expand All @@ -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));

Expand Down
141 changes: 141 additions & 0 deletions InterprocessMemory.Tests/CrossProcessTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -293,6 +293,147 @@ public void CrossProcess_OrphanWriteLock_IsRecovered()
region.ReleaseWriteLock();
}

/// <summary>
/// Starts a child that holds a lock until it is killed, and returns once the child has reported
/// that it holds the lock. <paramref name="role"/> is <c>hold_write_lock</c> or <c>hold_read_lock</c>.
/// </summary>
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;
}

/// <summary>Kills the process (a no-op when it has already exited) and waits for it to be gone.</summary>
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<bool> 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<bool> 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
Expand Down
30 changes: 17 additions & 13 deletions InterprocessMemory.Tests/ExtremeStressTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ public async Task MPMC_BurstTraffic_ShouldHandleSpikes()
for (int burst = 0; burst < burstCount; burst++)
{
var received = new ConcurrentBag<int>();
var producerDone = false;
var deadline = System.Diagnostics.Stopwatch.StartNew();

// Burst producer
var producer = Task.Run(() =>
Expand All @@ -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);
}
});

Expand Down
Loading