diff --git a/NOTICE.html b/NOTICE.html
index 164ee3ebcb..b94072c048 100644
--- a/NOTICE.html
+++ b/NOTICE.html
@@ -532,13 +532,6 @@
| NETStandard.Library | 2.0.3 |
diff --git a/qodana.yaml b/qodana.yaml
index 150e8cc1a7..a7e5733599 100644
--- a/qodana.yaml
+++ b/qodana.yaml
@@ -36,13 +36,6 @@ dependencyOverrides:
- key: "MIT"
url: "https://www.nuget.org/packages/Microsoft.Diagnostics.Runtime/2.2.332302"
- # the referenced license file is not trivial but appears to imply that this library is MIT
- - name: "Mono.Posix.NETStandard"
- version: "1.0.0"
- licenses:
- - key: "MIT"
- url: "https://github.com/mono/mono/blob/main/LICENSE"
-
# not automatically detectable by Qodana
- name: "OpenTelemetry"
version: "1.16.0"
diff --git a/src/Directory.Packages.props b/src/Directory.Packages.props
index 38755be9f3..cf693dd4db 100644
--- a/src/Directory.Packages.props
+++ b/src/Directory.Packages.props
@@ -50,7 +50,6 @@
-
diff --git a/src/EventStore.Core.Tests/Index/IndexV1/when_a_ptable_is_corrupt_on_disk.cs b/src/EventStore.Core.Tests/Index/IndexV1/when_a_ptable_is_corrupt_on_disk.cs
index f70460783f..aeb8e9d912 100644
--- a/src/EventStore.Core.Tests/Index/IndexV1/when_a_ptable_is_corrupt_on_disk.cs
+++ b/src/EventStore.Core.Tests/Index/IndexV1/when_a_ptable_is_corrupt_on_disk.cs
@@ -58,4 +58,15 @@ public void the_hash_is_invalid()
var exc = Assert.Throws (() => PTable.FromFile(_copiedfilename, Constants.PTableInitialReaderCount, Constants.PTableMaxReaderCountDefault, 16, false));
Assert.IsInstanceOf(exc.InnerException);
}
+
+ [Test]
+ public void failed_open_releases_the_file_handle()
+ {
+ Assert.Throws(() => PTable.FromFile(_copiedfilename, 1, 1, 16, false));
+
+ Assert.DoesNotThrow(() =>
+ {
+ using var stream = new FileStream(_copiedfilename, FileMode.Open, FileAccess.ReadWrite, FileShare.None);
+ });
+ }
}
diff --git a/src/EventStore.Core.Tests/Index/IndexV1/when_opening_a_ptable_with_one_reader.cs b/src/EventStore.Core.Tests/Index/IndexV1/when_opening_a_ptable_with_one_reader.cs
new file mode 100644
index 0000000000..c93e6b536a
--- /dev/null
+++ b/src/EventStore.Core.Tests/Index/IndexV1/when_opening_a_ptable_with_one_reader.cs
@@ -0,0 +1,52 @@
+using System.Threading.Tasks;
+using EventStore.Core.Index;
+using NUnit.Framework;
+
+namespace EventStore.Core.Tests.Index.IndexV1;
+
+[TestFixture(PTableVersions.IndexV2)]
+[TestFixture(PTableVersions.IndexV3)]
+[TestFixture(PTableVersions.IndexV4)]
+public class when_opening_a_ptable_with_one_reader : SpecificationWithFile
+{
+ private readonly byte _version;
+ private PTable _ptable;
+
+ public when_opening_a_ptable_with_one_reader(byte version)
+ {
+ _version = version;
+ }
+
+ [SetUp]
+ public override async Task SetUp()
+ {
+ await base.SetUp();
+
+ var memTable = new HashListMemTable(_version, maxSize: 1);
+ memTable.Add(0x010100000000, 1, 42);
+ _ptable = PTable.FromMemtable(
+ memTable,
+ Filename,
+ initialReaders: 1,
+ maxReaders: 1,
+ cacheDepth: 16,
+ skipIndexVerify: false,
+ useBloomFilter: false,
+ lruCacheSize: 0);
+ }
+
+ [TearDown]
+ public override void TearDown()
+ {
+ _ptable.MarkForDestruction();
+ _ptable.WaitForDisposal(1_000);
+ base.TearDown();
+ }
+
+ [Test]
+ public void reader_is_reusable_after_initialization()
+ {
+ Assert.That(_ptable.TryGetOneValue(0x010100000000, 1, out var position), Is.True);
+ Assert.That(position, Is.EqualTo(42));
+ }
+}
diff --git a/src/EventStore.Core.Tests/Index/MemTableTests.cs b/src/EventStore.Core.Tests/Index/MemTableTests.cs
index 6de1ee0c59..c583783475 100644
--- a/src/EventStore.Core.Tests/Index/MemTableTests.cs
+++ b/src/EventStore.Core.Tests/Index/MemTableTests.cs
@@ -14,6 +14,45 @@ public HashListMemTableTests()
: base(() => new HashListMemTable(PTableVersions.IndexV2, maxSize: 20))
{
}
+
+ [Test]
+ public void ordered_iteration_does_not_allocate_a_key_view_for_each_stream()
+ {
+ const int streamCount = 1_024;
+ var warmupTable = CreateTable(streamCount: 2);
+ _ = MeasureOrderedIterationAllocations(warmupTable);
+
+ var table = CreateTable(streamCount);
+ var firstIteration = MeasureOrderedIterationAllocations(table);
+ var repeatedIteration = MeasureOrderedIterationAllocations(table);
+
+ Assert.That(firstIteration - repeatedIteration, Is.LessThanOrEqualTo(streamCount * IntPtr.Size));
+ }
+
+ private static HashListMemTable CreateTable(int streamCount)
+ {
+ var table = new HashListMemTable(PTableVersions.IndexV2, maxSize: streamCount);
+ for (var stream = 0; stream < streamCount; stream++)
+ {
+ table.Add((ulong)stream, version: 0, position: stream);
+ }
+
+ return table;
+ }
+
+ private static long MeasureOrderedIterationAllocations(HashListMemTable table)
+ {
+ var allocatedBefore = GC.GetAllocatedBytesForCurrentThread();
+ var count = 0;
+ foreach (var _ in table.IterateAllInOrder())
+ {
+ count++;
+ }
+ var allocated = GC.GetAllocatedBytesForCurrentThread() - allocatedBefore;
+
+ GC.KeepAlive(count);
+ return allocated;
+ }
}
[TestFixture]
diff --git a/src/EventStore.Core.Tests/Index/ReverseComparerTests.cs b/src/EventStore.Core.Tests/Index/ReverseComparerTests.cs
index f38eb18aac..23be27b579 100644
--- a/src/EventStore.Core.Tests/Index/ReverseComparerTests.cs
+++ b/src/EventStore.Core.Tests/Index/ReverseComparerTests.cs
@@ -1,3 +1,4 @@
+using System;
using EventStore.Core.Index;
using NUnit.Framework;
@@ -23,4 +24,23 @@ public void same_values_are_equal()
{
Assert.AreEqual(0, new ReverseComparer().Compare(5, 5));
}
+
+ [Test]
+ public void comparing_value_types_does_not_allocate_per_comparison()
+ {
+ const int maximumOneTimeAllocation = 24;
+ var comparer = new ReverseComparer();
+ _ = comparer.Compare(2, 1);
+
+ var allocatedBefore = GC.GetAllocatedBytesForCurrentThread();
+ var result = 0;
+ for (ulong i = 0; i < 1_000; i++)
+ {
+ result += comparer.Compare(i + 1, i);
+ }
+ var allocated = GC.GetAllocatedBytesForCurrentThread() - allocatedBefore;
+
+ GC.KeepAlive(result);
+ Assert.That(allocated, Is.LessThanOrEqualTo(maximumOneTimeAllocation));
+ }
}
diff --git a/src/EventStore.Core.Tests/TransactionLog/Unbuffered/UnbufferedTests.cs b/src/EventStore.Core.Tests/TransactionLog/Unbuffered/UnbufferedTests.cs
deleted file mode 100644
index afa339d9de..0000000000
--- a/src/EventStore.Core.Tests/TransactionLog/Unbuffered/UnbufferedTests.cs
+++ /dev/null
@@ -1,807 +0,0 @@
-using System;
-using System.IO;
-using EventStore.Core.TransactionLog.Unbuffered;
-using NUnit.Framework;
-
-namespace EventStore.Core.Tests.TransactionLog.Unbuffered;
-
-[TestFixture]
-public class UnbufferedTests : SpecificationWithDirectory
-{
- [Test]
- public void when_resizing_a_file()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
- stream.SetLength(4096 * 1024);
- stream.Close();
- Assert.AreEqual(4096 * 1024, new FileInfo(filename).Length);
- }
-
- [Test]
- public void when_expanding_an_aligned_file_by_one_page()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(0, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize, stream.Position); //verify position
- stream.SetLength(initialFileSize + 4096); //expand file by 4KB
- Assert.AreEqual(initialFileSize, stream.Position); //position should not change
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize + 4096, new FileInfo(filename).Length); //file size should increase by 4KB
- }
-
- [Test]
- public void when_expanding_an_aligned_file_by_one_byte_less_than_one_page()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(0, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize, stream.Position); //verify position
- stream.SetLength(initialFileSize + 4095); //expand file by 4KB - 1
- Assert.AreEqual(initialFileSize, stream.Position); //position should not change
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize + 4096, new FileInfo(filename).Length); //file size should increase by 4KB
- }
-
- [Test]
- public void when_expanding_an_aligned_file_by_one_byte_more_than_one_page()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(0, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize, stream.Position); //verify position
- stream.SetLength(initialFileSize + 4097); //expand file by 4KB + 1
- Assert.AreEqual(initialFileSize, stream.Position); //position should not change
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize + 4096 * 2, new FileInfo(filename).Length); //file size should increase by 4KB x 2
- }
-
- [Test]
- public void when_expanding_an_aligned_file_by_one_byte()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(0, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize, stream.Position); //verify position
- stream.SetLength(initialFileSize + 1); //expand file by 1 byte
- Assert.AreEqual(initialFileSize, stream.Position); //position should not change
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize + 4096, new FileInfo(filename).Length); //file size should increase by 4KB
- }
-
- [Test]
- public void when_truncating_an_aligned_file_by_one_page()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(0, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize, stream.Position); //verify position
- stream.SetLength(initialFileSize - 4096); //truncate file by 4KB
- Assert.AreEqual(initialFileSize - 4096, stream.Position); //position should decrease by 4KB
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize - 4096, new FileInfo(filename).Length); //file size should decrease by 4KB
- }
-
- [Test]
- public void when_truncating_an_aligned_file_by_one_page_and_position_one_page_from_eof()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(-4096, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize - 4096, stream.Position); //verify position
- stream.SetLength(initialFileSize - 4096); //truncate file by 4KB
- Assert.AreEqual(initialFileSize - 4096, stream.Position); //position should not change
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize - 4096, new FileInfo(filename).Length); //file size should decrease by 4KB
- }
-
- [Test]
- public void when_truncating_an_aligned_file_by_one_byte_less_than_a_page()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(0, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize, stream.Position); //verify position
- stream.SetLength(initialFileSize - 4095); //truncate file by 4KB - 1
- Assert.AreEqual(initialFileSize, stream.Position); //position should not change
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize, new FileInfo(filename).Length); //file size should not change
- }
-
- [Test]
- public void when_truncating_an_aligned_file_by_one_byte_more_than_a_page()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(0, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize, stream.Position); //verify position
- stream.SetLength(initialFileSize - 4097); //truncate file by 4KB + 1
- Assert.AreEqual(initialFileSize - 4096, stream.Position); //position should decrease by 4KB
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize - 4096, new FileInfo(filename).Length); //file size should decrease by 4KB
- }
-
- [Test]
- public void when_truncating_an_aligned_file_by_one_byte()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096);
-
- var initialFileSize = 4096 * 1024;
- stream.SetLength(initialFileSize); //initial size of 4MB
-
- stream.Seek(0, SeekOrigin.End);
-
- Assert.AreEqual(initialFileSize, stream.Position); //verify position
- stream.SetLength(initialFileSize - 1); //truncate file by 1 byte
- Assert.AreEqual(initialFileSize, stream.Position); //position should not change
-
- stream.Close();
-
- Assert.AreEqual(initialFileSize, new FileInfo(filename).Length); //file size should not change
- }
-
- [Test]
- public void when_writing_less_than_buffer()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = GetBytes(255);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Write(bytes, 0, bytes.Length);
- Assert.AreEqual(bytes.Length, stream.Position);
- Assert.AreEqual(0, new FileInfo(filename).Length);
- }
- }
-
- [Test]
- public void when_writing_more_than_buffer()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = GetBytes(9000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Write(bytes, 0, bytes.Length);
- Assert.AreEqual(4096 * 2, new FileInfo(filename).Length);
- var read = ReadAllBytesShared(filename);
- for (var i = 0; i < 4096 * 2; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void when_writing_less_than_buffer_and_closing()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = GetBytes(255);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Write(bytes, 0, bytes.Length);
- }
-
- Assert.AreEqual(4096, new FileInfo(filename).Length);
- var read = ReadAllBytesShared(filename);
-
- for (var i = 0; i < 255; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
-
- [Test, Ignore("Requires a 4gb file")]
- public void when_seeking_greater_than_2gb()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var GIGABYTE = 1024L * 1024L * 1024L;
- try
- {
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.SetLength(4L * GIGABYTE);
- stream.Seek(3L * GIGABYTE, SeekOrigin.Begin);
- }
- }
- finally
- {
- File.Delete(filename);
- }
- }
-
- [Test]
- public void when_writing_less_than_buffer_and_seeking()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = GetBytes(255);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Write(bytes, 0, bytes.Length);
- stream.Seek(0, SeekOrigin.Begin);
- Assert.AreEqual(0, stream.Position);
- Assert.AreEqual(4096, new FileInfo(filename).Length);
- var read = ReadAllBytesShared(filename);
-
- for (var i = 0; i < 255; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void when_writing_exact_to_alignment_and_writing_again()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = GetBytes(4096);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Write(bytes, 0, bytes.Length);
- Assert.AreEqual(4096, stream.Position);
- bytes = GetBytes(15);
- stream.Write(bytes, 0, bytes.Length);
- Assert.AreEqual(4111, stream.Position);
- stream.Flush();
- Assert.AreEqual(4111, stream.Position);
- Assert.AreEqual(8192, new FileInfo(filename).Length);
- var read = ReadAllBytesShared(filename);
-
- for (var i = 0; i < 255; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void when_writing_then_seeking_exact_to_alignment_and_writing_again()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = GetBytes(8192);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Write(bytes, 0, 5012);
- Assert.AreEqual(5012, stream.Position);
- stream.Seek(4096, SeekOrigin.Begin);
- Assert.AreEqual(4096, stream.Position);
- bytes = GetBytes(15);
- stream.Write(bytes, 0, bytes.Length);
- Assert.AreEqual(4111, stream.Position);
- stream.Flush();
- Assert.AreEqual(4111, stream.Position);
- Assert.AreEqual(8192, new FileInfo(filename).Length);
- var read = ReadAllBytesShared(filename);
-
- for (var i = 0; i < 255; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
-
- [Test]
- public void when_seeking_non_exact_to_zero_block_and_writing()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 4096 * 64);
- var bytes = GetBytes(512);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Seek(128, SeekOrigin.Begin);
- stream.Write(bytes, 0, bytes.Length);
- stream.Flush();
- }
-
- using (var stream = new FileStream(filename, FileMode.Open))
- {
- var read = new byte[128];
- stream.Read(read, 0, 128);
- for (var i = 0; i < read.Length; i++)
- {
- Assert.AreEqual(i, read[i]);
- }
- }
- }
-
- [Test]
- public void when_writing_multiple_times()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = GetBytes(256);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Write(bytes, 0, bytes.Length);
- Assert.AreEqual(256, stream.Position);
- stream.Flush();
- Assert.AreEqual(256, stream.Position);
- stream.Write(bytes, 0, bytes.Length);
- Assert.AreEqual(512, stream.Position);
- stream.Flush();
- Assert.AreEqual(512, stream.Position);
- Assert.AreEqual(4096, new FileInfo(filename).Length);
- var read = ReadAllBytesShared(filename);
-
- for (var i = 0; i < 512; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
-
- [Test]
- public void when_reading_multiple_times()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 20000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Seek(4096 + 15, SeekOrigin.Begin);
- var read = new byte[1000];
- stream.Read(read, 0, 500);
- Assert.AreEqual(4096 + 15 + 500, stream.Position);
- stream.Read(read, 500, 500);
- Assert.AreEqual(4096 + 15 + 1000, stream.Position);
- for (var i = 0; i < read.Length; i++)
- {
- Assert.AreEqual((i + 15) % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void when_reading_multiple_times_no_seek()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 20000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- var read = new byte[1000];
- stream.Read(read, 0, 500);
- Assert.AreEqual(500, stream.Position);
- stream.Read(read, 500, 500);
- Assert.AreEqual(1000, stream.Position);
- for (var i = 0; i < read.Length; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void when_reading_multiple_times_over_page_size()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 20000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- var read = new byte[6000];
- stream.Read(read, 0, 3000);
- Assert.AreEqual(3000, stream.Position);
- var total = stream.Read(read, 3000, 3000);
- Assert.AreEqual(3000, total);
- Assert.AreEqual(6000, stream.Position);
- for (var i = 0; i < read.Length; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void when_reading_multiple_times_on_page_size()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 20000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- var read = new byte[6000];
- stream.Read(read, 0, 3000);
- Assert.AreEqual(3000, stream.Position);
- var total = stream.Read(read, 3000, 1096);
- Assert.AreEqual(1096, total);
- total = stream.Read(read, 4096, read.Length - 4096);
- Assert.AreEqual(read.Length - 4096, total);
- Assert.AreEqual(6000, stream.Position);
- for (var i = 0; i < read.Length; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
-
- [Test]
- public void when_reading_multiple_times_exact_page_size()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 4096 * 100 + 50);
- Span expected = GetBytes(4096);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- var read = new byte[4096];
- for (var i = 0; i < 100; i++)
- {
- var total = stream.Read(read, 0, 4096);
- Assert.AreEqual(4096 * (i + 1), stream.Position);
- Assert.AreEqual(4096, total);
- if (!expected.SequenceEqual(read))
- {
- for (var j = 0; j < read.Length; j++)
- {
- Assert.AreEqual(j % 256, read[j]);
- }
- }
- }
-
- expected = expected.Slice(0, 50);
- var total2 = stream.Read(read, 0, 50);
- Assert.AreEqual(409600 + 50, stream.Position);
- Assert.AreEqual(50, total2);
- if (!expected.SequenceEqual(read))
- {
- for (var j = 0; j < 50; j++)
- {
- Assert.AreEqual(j % 256, read[j]);
- }
- }
- }
- }
-
- [Test]
- public void when_reading_multiple_times_offset_page_size()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 4096 * 100 + 50);
- Span expected = GetBytes(4096 + 50);
- expected = expected[50..];
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Seek(50, SeekOrigin.Begin);
- var read = new byte[4096];
- for (var i = 0; i < 100; i++)
- {
- var total = stream.Read(read, 0, 4096);
- Assert.AreEqual(4096 * (i + 1) + 50, stream.Position);
- Assert.AreEqual(4096, total);
- if (!expected.SequenceEqual(read))
- {
- for (var j = 0; j < read.Length; j++)
- {
- Assert.AreEqual((j + 50) % 256, read[j]);
- }
- }
- }
-
- Assert.AreEqual(4096 * 100 + 50, stream.Position);
- }
- }
-
- [Test]
- public void when_writing_more_than_buffer_and_closing()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = GetBytes(9000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Write(bytes, 0, bytes.Length);
- stream.Close();
- Assert.AreEqual(4096 * 3, new FileInfo(filename).Length);
- var read = File.ReadAllBytes(filename);
- for (var i = 0; i < 9000; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void when_reading_on_aligned_buffer()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 20000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- var read = new byte[4096];
- stream.Read(read, 0, 4096);
- for (var i = 0; i < 4096; i++)
- {
- Assert.AreEqual(i % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void when_reading_on_unaligned_buffer()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 20000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Seek(15, SeekOrigin.Begin);
- var read = new byte[999];
- stream.Read(read, 0, read.Length);
- for (var i = 0; i < read.Length; i++)
- {
- Assert.AreEqual((i + 15) % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void seek_and_read_on_unaligned_buffer()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- MakeFile(filename, 20000);
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Seek(4096 + 15, SeekOrigin.Begin);
- Assert.AreEqual(4096 + 15, stream.Position);
- var read = new byte[999];
- stream.Read(read, 0, read.Length);
- Assert.AreEqual(4096 + 15 + 999, stream.Position);
- for (var i = 0; i < read.Length; i++)
- {
- Assert.AreEqual((i + 15) % 256, read[i]);
- }
- }
- }
-
- [Test]
- public void seek_end_of_file()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Seek(0, SeekOrigin.End);
- Assert.AreEqual(stream.Length, stream.Position);
- }
- }
-
- [Test]
- public void seek_origin_end_to_mid_of_file()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- stream.Seek(-30, SeekOrigin.End);
- Assert.AreEqual(stream.Length - 30, stream.Position);
- }
- }
-
-
- [Test]
- public void seek_current_unimplemented()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
-
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- Assert.Throws(() => stream.Seek(0, SeekOrigin.Current));
- }
- }
-
-
- [Test]
- public void seek_write_seek_read_in_buffer()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- using (var stream = UnbufferedFileStream.Create(filename, FileMode.CreateNew, FileAccess.ReadWrite,
- FileShare.ReadWrite, 4096, 4096, false, 4096))
- {
- var buffer = GetBytes(255);
- stream.Seek(4096 + 15, SeekOrigin.Begin);
- stream.Write(buffer, 0, buffer.Length);
- stream.Seek(4096 + 15, SeekOrigin.Begin);
- var read = new byte[255];
- stream.Read(read, 0, read.Length);
- for (var i = 0; i < read.Length; i++)
- {
- Assert.AreEqual(i % 255, read[i]);
- }
- }
- }
-
- [Test]
- public void same_as_file_stream_on_reads()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = BuildBytes(4096 * 128);
- File.WriteAllBytes(filename, bytes);
- using (var f = new FileStream(filename, FileMode.Open, FileAccess.Read))
- {
- using (var b = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.Read, FileShare.Read,
- 4096, 4096, false, 4096))
- {
- Span readf = new byte[4096];
- var readb = new byte[4096];
- for (var i = 0; i < 128; i++)
- {
- var totalf = f.Read(readf);
- var totalb = b.Read(readb, 0, 4096);
-
- if (!readf[..totalf].SequenceEqual(readb[..totalb]))
- {
- for (var j = 0; j < 4096; j++)
- {
- Assert.AreEqual(totalf, totalb);
- Assert.AreEqual(readf[j], readb[j]);
- }
- }
- }
- }
- }
- }
-
- [Test]
- public void same_as_file_stream_on_reads_with_bigger_buffer()
- {
- var filename = GetFilePathFor(Guid.NewGuid().ToString());
- var bytes = BuildBytes(4096 * 128);
- File.WriteAllBytes(filename, bytes);
- using (var f = new FileStream(filename, FileMode.Open, FileAccess.Read))
- {
- using (var b = UnbufferedFileStream.Create(filename, FileMode.Open, FileAccess.Read, FileShare.Read,
- 4096, 4096 * 4, false, 4096))
- {
- Span readf = new byte[4096];
- var readb = new byte[4096];
- for (var i = 0; i < 128; i++)
- {
- var totalf = f.Read(readf);
- var totalb = b.Read(readb, 0, 4096);
-
- if (!readf[..totalf].SequenceEqual(readb[..totalb]))
- {
- for (var j = 0; j < 4096; j++)
- {
- Assert.AreEqual(totalf, totalb);
- Assert.AreEqual(readf[j], readb[j]);
- }
- }
- }
- }
- }
- }
-
- private byte[] BuildBytes(int count)
- {
- using var ret = new MemoryStream();
- for (int i = 0; i < count / 16; i++)
- {
- ret.Write(Guid.NewGuid().ToByteArray(), 0, 16);
- }
-
- return ret.ToArray();
- }
-
-
- private byte[] ReadAllBytesShared(string filename)
- {
- using (var fs = File.Open(filename, FileMode.Open, FileAccess.Read, FileShare.ReadWrite))
- {
- var ret = new byte[fs.Length];
- fs.Read(ret, 0, (int)fs.Length);
- return ret;
- }
- }
-
- private void MakeFile(string filename, int size)
- {
- var bytes = GetBytes(size);
- File.WriteAllBytes(filename, bytes);
- }
-
- private byte[] GetBytes(int size)
- {
- var ret = new byte[size];
- for (var i = 0; i < size; i++)
- {
- ret[i] = (byte)(i % 256);
- }
-
- return ret;
- }
-}
diff --git a/src/EventStore.Core/EventStore.Core.csproj b/src/EventStore.Core/EventStore.Core.csproj
index b8195ca616..bb9c03ca7e 100644
--- a/src/EventStore.Core/EventStore.Core.csproj
+++ b/src/EventStore.Core/EventStore.Core.csproj
@@ -43,7 +43,6 @@
-
diff --git a/src/EventStore.Core/Index/HashListMemTable.cs b/src/EventStore.Core/Index/HashListMemTable.cs
index 55b3d273e6..192b19f94c 100644
--- a/src/EventStore.Core/Index/HashListMemTable.cs
+++ b/src/EventStore.Core/Index/HashListMemTable.cs
@@ -14,6 +14,7 @@ public class HashListMemTable : IMemTable
{
private static readonly IComparer MemTableComparer = new EventNumberComparer();
private static readonly IComparer LogPosComparer = new LogPositionComparer();
+ private static readonly IComparer DescendingStreamHashComparer = new ReverseComparer();
private static readonly TimeSpan DefaultLockTimeout = TimeSpan.FromMilliseconds(10_000);
public long Count
@@ -128,7 +129,7 @@ public bool TryGetOneValue(ulong stream, long number, out long position)
return false;
}
- var key = list.Keys[endIdx];
+ var key = list.GetKeyAtIndex(endIdx);
if (key.EvNum == number)
{
position = key.LogPos;
@@ -158,7 +159,7 @@ public bool TryGetLatestEntry(ulong stream, out IndexEntry entry)
try
{
- var latest = list.Keys[list.Count - 1];
+ var latest = list.GetKeyAtIndex(list.Count - 1);
entry = new IndexEntry(hash, latest.EvNum, latest.LogPos);
return true;
}
@@ -203,7 +204,7 @@ public bool TryGetLatestEntry(ulong stream, out IndexEntry entry)
return null;
}
- var latestBeforePosition = list.Keys[endIdx];
+ var latestBeforePosition = list.GetKeyAtIndex(endIdx);
return new(hash, latestBeforePosition.EvNum, latestBeforePosition.LogPos);
}
catch (SearchStoppedException)
@@ -219,7 +220,7 @@ await isForThisStream(new IndexEntry(hash, e.EvNum, e.LogPos), token),
return null;
}
- var latestBeforePosition = list.Keys[maxIdx];
+ var latestBeforePosition = list.GetKeyAtIndex(maxIdx);
return new(hash, latestBeforePosition.EvNum, latestBeforePosition.LogPos);
}
finally
@@ -242,7 +243,7 @@ public bool TryGetOldestEntry(ulong stream, out IndexEntry entry)
try
{
- var oldest = list.Keys[0];
+ var oldest = list.GetKeyAtIndex(0);
entry = new IndexEntry(hash, oldest.EvNum, oldest.LogPos);
return true;
}
@@ -282,7 +283,7 @@ public bool TryGetNextEntry(ulong stream, long afterNumber, out IndexEntry entry
return false;
}
- var e = list.Keys[endIdx];
+ var e = list.GetKeyAtIndex(endIdx);
entry = new IndexEntry(hash, e.EvNum, e.LogPos);
return true;
}
@@ -322,7 +323,7 @@ public bool TryGetPreviousEntry(ulong stream, long beforeNumber, out IndexEntry
return false;
}
- var e = list.Keys[endIdx];
+ var e = list.GetKeyAtIndex(endIdx);
entry = new IndexEntry(hash, e.EvNum, e.LogPos);
return true;
}
@@ -340,14 +341,14 @@ public IEnumerable IterateAllInOrder()
//Log.Trace("Sorting array in HashListMemTable.IterateAllInOrder...");
var keys = _hash.Keys.ToArray();
- Array.Sort(keys, new ReverseComparer());
+ Array.Sort(keys, DescendingStreamHashComparer);
foreach (var key in keys)
{
var list = _hash[key];
for (int i = list.Count - 1; i >= 0; --i)
{
- var x = list.Keys[i];
+ var x = list.GetKeyAtIndex(i);
yield return new IndexEntry(key, x.EvNum, x.LogPos);
}
}
@@ -377,7 +378,7 @@ public IReadOnlyList GetRange(ulong stream, long startNumber, long e
var endIdx = list.UpperBound(new Entry(endNumber, long.MaxValue));
for (int i = endIdx; i >= 0; i--)
{
- var key = list.Keys[i];
+ var key = list.GetKeyAtIndex(i);
if (key.EvNum < startNumber || ret.Count == limit)
{
break;
@@ -458,6 +459,6 @@ public class ReverseComparer : IComparer where T : IComparable
{
public int Compare(T x, T y)
{
- return -x.CompareTo(y);
+ return Comparer.Default.Compare(y, x);
}
}
diff --git a/src/EventStore.Core/Index/PTable.cs b/src/EventStore.Core/Index/PTable.cs
index d2e21e6cde..0567644637 100644
--- a/src/EventStore.Core/Index/PTable.cs
+++ b/src/EventStore.Core/Index/PTable.cs
@@ -11,11 +11,9 @@
using EventStore.Core.DataStructures;
using EventStore.Core.DataStructures.ProbabilisticFilter;
using EventStore.Core.Exceptions;
-using EventStore.Core.TransactionLog.Unbuffered;
using ILogger = Serilog.ILogger;
using MD5 = EventStore.Core.Hashing.MD5;
using Range = EventStore.Core.Data.Range;
-using RuntimeInformation = System.Runtime.RuntimeInformation;
namespace EventStore.Core.Index;
@@ -355,19 +353,8 @@ internal UnmanagedMemoryAppendOnlyList CacheMidpointsAndVerifyHash(int
Log.Debug("Disabling Verification of PTable");
}
- WorkItem workItem = null;
-
- Stream stream;
- if (RuntimeInformation.IsUnix)
- {
- workItem = GetWorkItem();
- stream = workItem.Stream;
- }
- else
- {
- stream = UnbufferedFileStream.Create(_filename, FileMode.Open, FileAccess.Read, FileShare.Read,
- 4096, 4096, false, 4096);
- }
+ var workItem = GetWorkItem();
+ var stream = workItem.Stream;
UnmanagedMemoryAppendOnlyList midpoints = null;
@@ -539,17 +526,7 @@ internal UnmanagedMemoryAppendOnlyList CacheMidpointsAndVerifyHash(int
}
finally
{
- if (RuntimeInformation.IsUnix)
- {
- if (workItem != null)
- {
- ReturnWorkItem(workItem);
- }
- }
- else
- {
- stream?.Dispose();
- }
+ ReturnWorkItem(workItem);
}
}
diff --git a/src/EventStore.Core/TransactionLog/Unbuffered/UnbufferedFileStream.cs b/src/EventStore.Core/TransactionLog/Unbuffered/UnbufferedFileStream.cs
deleted file mode 100644
index efea76548c..0000000000
--- a/src/EventStore.Core/TransactionLog/Unbuffered/UnbufferedFileStream.cs
+++ /dev/null
@@ -1,376 +0,0 @@
-using System;
-using System.IO;
-using System.Runtime.InteropServices;
-using Microsoft.Win32.SafeHandles;
-using RuntimeInformation = System.Runtime.RuntimeInformation;
-
-namespace EventStore.Core.TransactionLog.Unbuffered
-{
- //NOTE THIS DOES NOT SUPPORT ALL STREAM OPERATIONS AS YOU MIGHT EXPECT IT SUPPORTS WHAT WE USE!
- public unsafe class UnbufferedFileStream : Stream
- {
- private static readonly INativeFile NativeFile;
-
- private byte* _writeBuffer;
- private byte* _readBuffer;
- private readonly int _writeBufferSize;
- private readonly int _readBufferSize;
- private readonly IntPtr _writeBufferOriginal;
- private readonly IntPtr _readBufferOriginal;
- private readonly uint _blockSize;
- private long _bufferedCount;
- private bool _aligned;
- private long _lastPosition;
- private bool _needsFlush;
- private SafeFileHandle _handle;
- private long _readLocation = -1;
- private bool _needsRead;
-
- static UnbufferedFileStream() =>
- NativeFile = RuntimeInformation.IsWindows ? new NativeFileWindows() : new NativeFileUnix();
-
- private UnbufferedFileStream(SafeFileHandle handle, uint blockSize, int internalWriteBufferSize,
- int internalReadBufferSize)
- {
- _handle = handle;
- _readBufferSize = internalReadBufferSize;
- _writeBufferSize = internalWriteBufferSize;
- _writeBufferOriginal = Marshal.AllocHGlobal((int)(internalWriteBufferSize + blockSize));
- _readBufferOriginal = Marshal.AllocHGlobal((int)(internalReadBufferSize + blockSize));
- _readBuffer = Align(_readBufferOriginal, blockSize);
- _writeBuffer = Align(_writeBufferOriginal, blockSize);
- _blockSize = blockSize;
- }
-
- private byte* Align(IntPtr buf, uint alignTo)
- {
- //This makes an aligned buffer linux needs this.
- //The buffer must originally be at least one alignment bigger!
- var diff = alignTo - (buf.ToInt64() % alignTo);
- var aligned = (IntPtr)(buf.ToInt64() + diff);
- return (byte*)aligned;
- }
-
- public static UnbufferedFileStream Create(string path,
- FileMode mode,
- FileAccess acc,
- FileShare share,
- int internalWriteBufferSize,
- int internalReadBufferSize,
- bool writeThrough,
- uint minBlockSize)
- {
- var blockSize = NativeFile.GetDriveSectorSize(path);
- blockSize = blockSize > minBlockSize ? blockSize : minBlockSize;
- if (internalWriteBufferSize % blockSize != 0)
- {
- throw new Exception("write buffer size must be aligned to block size of " + blockSize + " bytes");
- }
-
- if (internalReadBufferSize % blockSize != 0)
- {
- throw new Exception("read buffer size must be aligned to block size of " + blockSize + " bytes");
- }
-
- var handle = NativeFile.CreateUnbufferedRW(path, acc, share, mode, writeThrough);
- return new UnbufferedFileStream(handle, blockSize, internalWriteBufferSize, internalReadBufferSize);
- }
-
- public override void Flush()
- {
- CheckDisposed();
- if (!_needsFlush)
- {
- return;
- }
-
- var alignedbuffer = (int)GetLowestAlignment(_bufferedCount);
- var positionAligned = GetLowestAlignment(_lastPosition);
- if (!_aligned)
- {
- SeekInternal(positionAligned, SeekOrigin.Begin);
- }
-
- if (_bufferedCount == alignedbuffer)
- {
- InternalWrite(_writeBuffer, (uint)_bufferedCount);
- _lastPosition = positionAligned + _bufferedCount;
- _bufferedCount = 0;
- _aligned = true;
- }
- else
- {
- var left = _bufferedCount - alignedbuffer;
- InternalWrite(_writeBuffer, (uint)(alignedbuffer + _blockSize));
- _lastPosition = positionAligned + alignedbuffer;
- SetBuffer(alignedbuffer, left);
- _bufferedCount = left;
- _aligned = false;
- }
-
- _needsFlush = false;
- }
-
- private static void MemCopy(byte[] src, long srcOffset, byte* dest, long destOffset, long count)
- {
- fixed (byte* p = src)
- {
- MemCopy(p, srcOffset, dest, destOffset, count);
- }
- }
-
- private static void MemCopy(byte* src, long srcOffset, byte[] dest, long destOffset, long count)
- {
- fixed (byte* p = dest)
- {
- MemCopy(src, srcOffset, p, destOffset, count);
- }
- }
-
- private static void MemCopy(byte* src, long srcOffset, byte* dest, long destOffset, long count)
- {
- byte* psrc = src + srcOffset;
- byte* pdest = dest + destOffset;
-
- for (var i = 0; i < count; i++)
- {
- *pdest = *psrc;
- pdest++;
- psrc++;
- }
- }
-
- private void SeekInternal(long positionAligned, SeekOrigin origin)
- {
- NativeFile.Seek(_handle, positionAligned, origin);
- }
-
- private void InternalWrite(byte* buffer, uint count)
- {
- var written = 0;
- NativeFile.Write(_handle, buffer, count, ref written);
- }
-
- public override long Seek(long offset, SeekOrigin origin)
- {
- long mungedOffset = offset;
- CheckDisposed();
- if (origin == SeekOrigin.Current)
- {
- throw new NotImplementedException("only supports seek origin begin/end");
- }
-
- if (origin == SeekOrigin.End)
- {
- mungedOffset = Length + offset;
- }
-
- var aligned = GetLowestAlignment(mungedOffset);
- var left = (int)(mungedOffset - aligned);
- Flush();
- _bufferedCount = left;
- _aligned = aligned == left;
- _lastPosition = aligned;
- //TODO cant do two seeks + a read here.
- SeekInternal(aligned, SeekOrigin.Begin);
- _needsRead = true;
- return offset;
- }
-
- private long GetLowestAlignment(long offset)
- {
- return offset - (offset % _blockSize);
- }
-
- public override void SetLength(long value)
- {
- CheckDisposed();
- var aligned = GetLowestAlignment(value);
- aligned = aligned == value ? aligned : aligned + _blockSize;
- NativeFile.SetFileSize(_handle, aligned);
- if (Position > aligned)
- {
- Seek(aligned, SeekOrigin.Begin);
- }
- }
-
- public override int Read(byte[] buffer, int offset, int count)
- {
- CheckDisposed();
- if (offset < 0 || buffer.Length < offset)
- {
- throw new ArgumentException("offset");
- }
-
- if (count < 0 || buffer.Length < count)
- {
- throw new ArgumentException("offset");
- }
-
- if (offset + count > buffer.Length)
- {
- throw new ArgumentException("offset + count must be less than size of array");
- }
-
- var position = GetLowestAlignment(Position);
- var roffset = (int)(Position - position);
-
- var bytesRead = _readBufferSize;
- if (_readLocation + _readBufferSize <= position || _readLocation > position || _readLocation == -1)
- {
- SeekInternal(position, SeekOrigin.Begin);
- bytesRead = NativeFile.Read(_handle, _readBuffer, 0, _readBufferSize);
- _readLocation = position;
- }
- else if (_readLocation != position)
- {
- roffset += (int)(position - _readLocation);
- }
-
- var bytesAvailable = bytesRead - roffset;
- if (bytesAvailable <= 0)
- {
- return 0;
- }
-
- var toCopy = count > bytesAvailable ? bytesAvailable : count;
-
- MemCopy(_readBuffer, roffset, buffer, offset, toCopy);
- _bufferedCount += toCopy;
- if (count - toCopy == 0)
- {
- return toCopy;
- }
-
- return toCopy + Read(buffer, offset + toCopy, count - toCopy);
- }
-
- public override void Write(byte[] buffer, int offset, int count)
- {
- CheckDisposed();
- var done = false;
- long left = count;
- long current = offset;
- if (_needsRead)
- {
- SeekInternal(_lastPosition, SeekOrigin.Begin);
- NativeFile.Read(_handle, _writeBuffer, 0, (int)_blockSize);
- SeekInternal(_lastPosition, SeekOrigin.Begin);
- _needsRead = false;
- }
-
- while (!done)
- {
- _needsFlush = true;
- if (_bufferedCount + left < _writeBufferSize)
- {
- CopyBuffer(buffer, current, left);
- done = true;
- current += left;
- }
- else
- {
- var toFill = _writeBufferSize - _bufferedCount;
- CopyBuffer(buffer, current, toFill);
- Flush();
- left -= toFill;
- current += toFill;
- done = left == 0;
- }
- }
- }
-
- private void CopyBuffer(byte[] buffer, long offset, long count)
- {
- MemCopy(buffer, offset, _writeBuffer, _bufferedCount, count);
- _bufferedCount += count;
- }
-
- public override bool CanRead
- {
- get
- {
- CheckDisposed();
- return true;
- }
- }
-
- public override bool CanSeek
- {
- get
- {
- CheckDisposed();
- return true;
- }
- }
-
- public override bool CanWrite
- {
- get
- {
- CheckDisposed();
- return true;
- }
- }
-
- public override long Length
- {
- get
- {
- CheckDisposed();
- return NativeFile.GetFileSize(_handle);
- }
- }
-
- public override long Position
- {
- get
- {
- CheckDisposed();
- if (_aligned)
- {
- return _lastPosition + _bufferedCount;
- }
-
- return GetLowestAlignment(_lastPosition) + _bufferedCount;
- }
- set
- {
- CheckDisposed();
- Seek(value, SeekOrigin.Begin);
- }
- }
-
- private void SetBuffer(long alignedbuffer, long left)
- {
- MemCopy(_writeBuffer, alignedbuffer, _writeBuffer, 0, left);
- }
-
- [System.Diagnostics.Conditional("DEBUG")]
- private void CheckDisposed()
- {
- //only check in debug
- if (_handle == null)
- {
- throw new ObjectDisposedException("object is disposed.");
- }
- }
-
- protected override void Dispose(bool disposing)
- {
- if (_handle == null)
- {
- return;
- }
-
- Flush();
- _handle.Close();
- _handle = null;
- _readBuffer = (byte*)IntPtr.Zero;
- _writeBuffer = (byte*)IntPtr.Zero;
- Marshal.FreeHGlobal(_readBufferOriginal);
- Marshal.FreeHGlobal(_writeBufferOriginal);
- GC.SuppressFinalize(this);
- }
- }
-}
diff --git a/src/EventStore.Native/EventStore.Native.csproj b/src/EventStore.Native/EventStore.Native.csproj
deleted file mode 100644
index 5dca5f3d32..0000000000
--- a/src/EventStore.Native/EventStore.Native.csproj
+++ /dev/null
@@ -1,12 +0,0 @@
-
-
- true
- true
-
-
-
-
-
-
-
-
diff --git a/src/EventStore.Native/TransactionLog/Unbuffered/ExtendedFileOptions.cs b/src/EventStore.Native/TransactionLog/Unbuffered/ExtendedFileOptions.cs
deleted file mode 100644
index 23d53f3711..0000000000
--- a/src/EventStore.Native/TransactionLog/Unbuffered/ExtendedFileOptions.cs
+++ /dev/null
@@ -1,12 +0,0 @@
-using System;
-
-namespace EventStore.Core.TransactionLog.Unbuffered;
-
-[Flags]
-public enum ExtendedFileOptions
-{
- NoBuffering = unchecked((int)0x20000000),
- Overlapped = unchecked((int)0x40000000),
- SequentialScan = unchecked((int)0x08000000),
- WriteThrough = unchecked((int)0x80000000)
-}
diff --git a/src/EventStore.Native/TransactionLog/Unbuffered/INativeFile.cs b/src/EventStore.Native/TransactionLog/Unbuffered/INativeFile.cs
deleted file mode 100644
index 4153809ac8..0000000000
--- a/src/EventStore.Native/TransactionLog/Unbuffered/INativeFile.cs
+++ /dev/null
@@ -1,20 +0,0 @@
-using System.IO;
-using Microsoft.Win32.SafeHandles;
-
-namespace EventStore.Core.TransactionLog.Unbuffered;
-
-public interface INativeFile
-{
- uint GetDriveSectorSize(string path);
- long GetPageSize(string path);
- void SetFileSize(SafeFileHandle handle, long count);
- unsafe void Write(SafeFileHandle handle, byte* buffer, uint count, ref int written);
- unsafe int Read(SafeFileHandle handle, byte* buffer, int offset, int count);
- long GetFileSize(SafeFileHandle handle);
- SafeFileHandle Create(string path, FileAccess acc, FileShare readWrite, FileMode mode, int flags);
-
- SafeFileHandle CreateUnbufferedRW(string path, FileAccess acc, FileShare share, FileMode mode,
- bool writeThrough);
-
- void Seek(SafeFileHandle handle, long position, SeekOrigin origin);
-}
diff --git a/src/EventStore.Native/TransactionLog/Unbuffered/MacCaching.cs b/src/EventStore.Native/TransactionLog/Unbuffered/MacCaching.cs
deleted file mode 100644
index 2f90e8ed60..0000000000
--- a/src/EventStore.Native/TransactionLog/Unbuffered/MacCaching.cs
+++ /dev/null
@@ -1,34 +0,0 @@
-using System.Runtime.InteropServices;
-using Microsoft.Win32.SafeHandles;
-using Mono.Unix;
-using RuntimeInformation = System.Runtime.RuntimeInformation;
-
-namespace EventStore.Core.TransactionLog.Unbuffered;
-
-internal static class MacCaching
-{
- // ReSharper disable once InconsistentNaming
- private const uint MAC_F_NOCACHE = 48;
-
- [DllImport("libc")]
- static extern int fcntl(int fd, uint command, int arg);
-
- public static void Disable(SafeFileHandle handle)
- {
- if (!RuntimeInformation.IsOSX)
- {
- return;
- }
-
- long r;
- do
- {
- r = fcntl(handle.DangerousGetHandle().ToInt32(), MAC_F_NOCACHE, 1);
- } while (UnixMarshal.ShouldRetrySyscall((int)r));
-
- if (r == -1)
- {
- UnixMarshal.ThrowExceptionForLastError();
- }
- }
-}
diff --git a/src/EventStore.Native/TransactionLog/Unbuffered/NativeFileUnix.cs b/src/EventStore.Native/TransactionLog/Unbuffered/NativeFileUnix.cs
deleted file mode 100644
index ed59bd9ad5..0000000000
--- a/src/EventStore.Native/TransactionLog/Unbuffered/NativeFileUnix.cs
+++ /dev/null
@@ -1,166 +0,0 @@
-using System;
-using System.ComponentModel;
-using System.IO;
-using System.Runtime;
-using Microsoft.Win32.SafeHandles;
-using Mono.Unix;
-using Mono.Unix.Native;
-
-namespace EventStore.Core.TransactionLog.Unbuffered;
-
-public unsafe class NativeFileUnix : INativeFile
-{
- public uint GetDriveSectorSize(string path)
- {
- return 0;
- }
-
- public long GetPageSize(string path)
- {
- int r;
- do
- {
- r = (int)Syscall.sysconf(SysconfName._SC_PAGESIZE);
- } while (UnixMarshal.ShouldRetrySyscall(r));
-
- UnixMarshal.ThrowExceptionForLastErrorIf(r);
- return r;
- }
-
- public void SetFileSize(SafeFileHandle handle, long count)
- {
- int r;
- do
- {
- r = Syscall.ftruncate(handle.DangerousGetHandle().ToInt32(), count);
- } while (UnixMarshal.ShouldRetrySyscall(r));
-
- UnixMarshal.ThrowExceptionForLastErrorIf(r);
- FSync(handle);
- }
-
- private static void FSync(SafeFileHandle handle)
- {
- Syscall.fsync(handle.DangerousGetHandle().ToInt32());
- }
-
- public void Write(SafeFileHandle handle, byte* buffer, uint count, ref int written)
- {
- int ret;
- do
- {
- ret = (int)Syscall.write(handle.DangerousGetHandle().ToInt32(), buffer, count);
- } while (UnixMarshal.ShouldRetrySyscall(ret));
-
- if (ret == -1)
- {
- UnixMarshal.ThrowExceptionForLastErrorIf(ret);
- }
-
- written = (int)count;
- }
-
- public int Read(SafeFileHandle handle, byte* buffer, int offset, int count)
- {
- int r;
- do
- {
- r = (int)Syscall.read(handle.DangerousGetHandle().ToInt32(), buffer, (ulong)count);
- } while (UnixMarshal.ShouldRetrySyscall(r));
-
- if (r == -1)
- {
- UnixMarshal.ThrowExceptionForLastError();
- }
-
- return count;
- }
-
- public long GetFileSize(SafeFileHandle handle)
- {
- Stat s;
- int r;
- do
- {
- r = Syscall.fstat(handle.DangerousGetHandle().ToInt32(), out s);
- } while (UnixMarshal.ShouldRetrySyscall(r));
-
- UnixMarshal.ThrowExceptionForLastErrorIf(r);
- return s.st_size;
- }
-
- //TODO UNBUFF use FileAccess etc or do custom?
- public SafeFileHandle Create(string path, FileAccess acc, FileShare readWrite, FileMode mode, int flags)
- {
- //TODO convert flags or separate methods?
- return new SafeFileHandle((IntPtr)0, true);
- }
-
-
- public SafeFileHandle CreateUnbufferedRW(string path, FileAccess acc, FileShare share, FileMode mode,
- bool writeThrough)
- {
- //O_RDONLY is 0
- var direct = RuntimeInformation.IsOSX ? OpenFlags.O_RDONLY : OpenFlags.O_DIRECT;
- var flags = GetFlags(acc, mode) | direct;
- var han = Syscall.open(path, flags, FilePermissions.S_IRWXU);
- if (han < 0)
- {
- throw new Win32Exception();
- }
-
- var handle = new SafeFileHandle((IntPtr)han, true);
- if (handle.IsInvalid)
- {
- throw new Exception("Invalid handle");
- }
-
- MacCaching.Disable(handle);
-
- return handle;
- }
-
- private static OpenFlags GetFlags(FileAccess acc, FileMode mode)
- {
- var flags = OpenFlags.O_RDONLY; //RDONLY is 0
- switch (acc)
- {
- case FileAccess.Read:
- flags |= OpenFlags.O_RDONLY;
- break;
- case FileAccess.Write:
- flags |= OpenFlags.O_WRONLY;
- break;
- case FileAccess.ReadWrite:
- flags |= OpenFlags.O_RDWR;
- break;
- }
-
- switch (mode)
- {
- case FileMode.Append:
- flags |= OpenFlags.O_APPEND;
- break;
- case FileMode.Create:
- case FileMode.CreateNew:
- flags |= OpenFlags.O_CREAT;
- break;
- case FileMode.Truncate:
- flags |= OpenFlags.O_TRUNC;
- break;
- }
-
- return flags;
- }
-
- public void Seek(SafeFileHandle handle, long position, SeekOrigin origin)
- {
- int r;
- do
- {
- r = (int)Syscall.lseek(handle.DangerousGetHandle().ToInt32(), position, SeekFlags.SEEK_SET);
- } while (UnixMarshal.ShouldRetrySyscall(r));
-
- UnixMarshal.ThrowExceptionForLastErrorIf(r);
- }
-}
diff --git a/src/EventStore.Native/TransactionLog/Unbuffered/NativeFileWindows.cs b/src/EventStore.Native/TransactionLog/Unbuffered/NativeFileWindows.cs
deleted file mode 100644
index c87997d1e1..0000000000
--- a/src/EventStore.Native/TransactionLog/Unbuffered/NativeFileWindows.cs
+++ /dev/null
@@ -1,128 +0,0 @@
-using System;
-using System.ComponentModel;
-using System.IO;
-using Microsoft.Win32.SafeHandles;
-
-namespace EventStore.Core.TransactionLog.Unbuffered;
-
-public unsafe class NativeFileWindows : INativeFile
-{
- public uint GetDriveSectorSize(string path)
- {
- WinNative.GetDiskFreeSpace(Path.GetPathRoot(path), out _, out var size, out _, out _);
- return size;
- }
-
- public long GetPageSize(string path)
- {
- return GetDriveSectorSize(path);
- }
-
- public void SetFileSize(SafeFileHandle handle, long count)
- {
- var low = (int)(count & 0xffffffff);
- var high = (int)(count >> 32);
- WinNative.SetFilePointer(handle, low, ref high, WinNative.EMoveMethod.Begin);
- if (!WinNative.SetEndOfFile(handle))
- {
- throw new Win32Exception();
- }
-
- FSync(handle);
- }
-
- private static void FSync(SafeFileHandle handle)
- {
- WinNative.FlushFileBuffers(handle);
- }
-
- public void Write(SafeFileHandle handle, byte* buffer, uint count, ref int written)
- {
- if (!WinNative.WriteFile(handle, buffer, count, ref written, IntPtr.Zero))
- {
- throw new Win32Exception();
- }
- }
-
- public int Read(SafeFileHandle handle, byte* buffer, int offset, int count)
- {
- var read = 0;
-
- if (!WinNative.ReadFile(handle, buffer, count, ref read, 0))
- {
- throw new Win32Exception();
- }
-
- return read;
- }
-
- public long GetFileSize(SafeFileHandle handle)
- {
- if (!WinNative.GetFileSizeEx(handle, out var size))
- {
- throw new Win32Exception();
- }
-
- return size;
- }
-
- //TODO UNBUFF use FileAccess etc or do custom?
- public SafeFileHandle Create(string path, FileAccess acc, FileShare readWrite, FileMode mode, int flags)
- {
- var handle = WinNative.CreateFile(path,
- acc,
- FileShare.ReadWrite,
- IntPtr.Zero,
- mode,
- flags,
- IntPtr.Zero);
- if (handle.IsInvalid)
- {
- throw new Win32Exception();
- }
-
- return handle;
- }
-
-
- public SafeFileHandle CreateUnbufferedRW(string path, FileAccess acc, FileShare share, FileMode mode,
- bool writeThrough)
- {
- var flags = ExtendedFileOptions.NoBuffering;
- if (writeThrough)
- {
- flags = flags | ExtendedFileOptions.WriteThrough;
- }
-
- var handle = WinNative.CreateFile(path,
- acc,
- share,
- IntPtr.Zero,
- FileMode.OpenOrCreate,
- (int)flags,
- IntPtr.Zero);
- if (handle.IsInvalid)
- {
- throw new Win32Exception();
- }
-
- return handle;
- }
-
- public void Seek(SafeFileHandle handle, long position, SeekOrigin origin)
- {
- var low = (int)(position & 0xffffffff);
- var high = (int)(position >> 32);
- var moveMethod = origin switch
- {
- SeekOrigin.Current => WinNative.EMoveMethod.Current,
- SeekOrigin.End => WinNative.EMoveMethod.End,
- _ => WinNative.EMoveMethod.Begin
- };
- var f = WinNative.SetFilePointer(handle, low, ref high, moveMethod);
- if (f == WinNative.INVALID_SET_FILE_POINTER)
- {
- throw new Win32Exception();
- }
- }
-}
diff --git a/src/EventStore.Native/TransactionLog/Unbuffered/WinNative.cs b/src/EventStore.Native/TransactionLog/Unbuffered/WinNative.cs
deleted file mode 100644
index 20d1996d87..0000000000
--- a/src/EventStore.Native/TransactionLog/Unbuffered/WinNative.cs
+++ /dev/null
@@ -1,78 +0,0 @@
-using System;
-using System.IO;
-using System.Runtime.InteropServices;
-using Microsoft.Win32.SafeHandles;
-
-namespace EventStore.Core.TransactionLog.Unbuffered;
-
-internal static unsafe class WinNative
-{
- [DllImport("KERNEL32", SetLastError = true, CharSet = CharSet.Auto, BestFitMapping = false)]
- public static extern bool GetDiskFreeSpace(string path,
- out uint sectorsPerCluster,
- out uint bytesPerSector,
- out uint numberOfFreeClusters,
- out uint totalNumberOfClusters);
-
- [DllImport("kernel32.dll", SetLastError = true)]
- internal static extern bool WriteFile(
- SafeFileHandle hFile,
- byte* aBuffer,
- UInt32 cbToWrite,
- ref int cbThatWereWritten,
- IntPtr pOverlapped);
-
- [DllImport("kernel32", SetLastError = true)]
- public static extern bool ReadFile
- (
- SafeFileHandle hFile,
- byte* pBuffer,
- int numberOfBytesToRead,
- ref int pNumberOfBytesRead,
- int overlapped
- );
-
- [DllImport("kernel32.dll")]
- public static extern bool GetFileSizeEx(SafeFileHandle hFile, out long lpFileSize);
-
- [DllImport("kernel32.dll", SetLastError = true)]
- internal static extern UInt32 SetFilePointer(
- SafeFileHandle hFile,
- Int32 cbDistanceToMove,
- IntPtr pDistanceToMoveHigh,
- EMoveMethod fMoveMethod);
-
- [DllImport("KERNEL32", SetLastError = true, CharSet = CharSet.Auto, BestFitMapping = false)]
- public static extern SafeFileHandle CreateFile(String fileName,
- FileAccess desiredAccess,
- FileShare shareMode,
- IntPtr securityAttrs,
- FileMode creationDisposition,
- int flagsAndAttributes,
- IntPtr templateFile);
-
- public enum EMoveMethod : uint
- {
- Begin = 0,
- Current = 1,
- End = 2
- }
-
- [DllImport("kernel32.dll", SetLastError = true)]
- internal static extern bool SetEndOfFile(
- SafeFileHandle hFile);
-
- [DllImport("kernel32.dll", SetLastError = true)]
- internal static extern bool FlushFileBuffers(SafeFileHandle filehandle);
-
- [DllImport("kernel32.dll", SetLastError = true)]
- [return: MarshalAs(UnmanagedType.Bool)]
- public static extern bool CloseHandle(IntPtr hObject);
-
- [DllImport("Kernel32.dll", SetLastError = true, CharSet = CharSet.Auto)]
- public static extern int SetFilePointer(SafeFileHandle handle, int lDistanceToMove,
- ref int lpDistanceToMoveHigh, EMoveMethod dwMoveMethod);
-
- // ReSharper disable once InconsistentNaming
- public const int INVALID_SET_FILE_POINTER = -1;
-}
diff --git a/src/EventStore.sln b/src/EventStore.sln
index ef6553395d..7262d86925 100644
--- a/src/EventStore.sln
+++ b/src/EventStore.sln
@@ -23,8 +23,6 @@ Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "EventStore.Transport.Tcp",
EndProject
Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "EventStore.ClusterNode", "EventStore.ClusterNode\EventStore.ClusterNode.csproj", "{FDFCE363-84A4-4A6F-B97C-D9F57D8D124F}"
EndProject
-Project("{9A19103F-16F7-4668-BE54-9A1E7A4F7556}") = "EventStore.Native", "EventStore.Native\EventStore.Native.csproj", "{03FA0763-6D41-44C0-911D-662AE2B10084}"
-EndProject
Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Solution Items", "Solution Items", "{A19631BB-65CF-4726-8732-4EFD977E9B32}"
ProjectSection(SolutionItems) = preProject
.editorconfig = .editorconfig
@@ -195,18 +193,6 @@ Global
{FDFCE363-84A4-4A6F-B97C-D9F57D8D124F}.Release|ARM64.Build.0 = Release|ARM64
{FDFCE363-84A4-4A6F-B97C-D9F57D8D124F}.Release|x64.ActiveCfg = Release|x64
{FDFCE363-84A4-4A6F-B97C-D9F57D8D124F}.Release|x64.Build.0 = Release|x64
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Debug|Any CPU.Build.0 = Debug|Any CPU
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Debug|ARM64.ActiveCfg = Debug|ARM64
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Debug|ARM64.Build.0 = Debug|ARM64
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Debug|x64.ActiveCfg = Debug|x64
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Debug|x64.Build.0 = Debug|x64
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Release|Any CPU.ActiveCfg = Release|Any CPU
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Release|Any CPU.Build.0 = Release|Any CPU
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Release|ARM64.ActiveCfg = Release|ARM64
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Release|ARM64.Build.0 = Release|ARM64
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Release|x64.ActiveCfg = Release|x64
- {03FA0763-6D41-44C0-911D-662AE2B10084}.Release|x64.Build.0 = Release|x64
{066B6514-86A0-41D8-BB6A-D9B05141D026}.Debug|Any CPU.ActiveCfg = Debug|Any CPU
{066B6514-86A0-41D8-BB6A-D9B05141D026}.Debug|Any CPU.Build.0 = Debug|Any CPU
{066B6514-86A0-41D8-BB6A-D9B05141D026}.Debug|ARM64.ActiveCfg = Debug|ARM64
|