From 714a0d0c2cf83c213c38a22b6717988aaeeddce9 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Mon, 10 Aug 2026 14:22:15 -0400 Subject: [PATCH 1/2] chore(storage): reduce platform-specific index risk Signed-off-by: Yordis Prieto --- NOTICE.html | 7 - qodana.yaml | 7 - src/Directory.Packages.props | 1 - .../when_a_ptable_is_corrupt_on_disk.cs | 11 + .../when_opening_a_ptable_with_one_reader.cs | 52 ++ .../Unbuffered/UnbufferedTests.cs | 807 ------------------ src/EventStore.Core/EventStore.Core.csproj | 1 - src/EventStore.Core/Index/PTable.cs | 29 +- .../Unbuffered/UnbufferedFileStream.cs | 376 -------- .../EventStore.Native.csproj | 12 - .../Unbuffered/ExtendedFileOptions.cs | 12 - .../TransactionLog/Unbuffered/INativeFile.cs | 20 - .../TransactionLog/Unbuffered/MacCaching.cs | 34 - .../Unbuffered/NativeFileUnix.cs | 166 ---- .../Unbuffered/NativeFileWindows.cs | 128 --- .../TransactionLog/Unbuffered/WinNative.cs | 78 -- src/EventStore.sln | 14 - 17 files changed, 66 insertions(+), 1689 deletions(-) create mode 100644 src/EventStore.Core.Tests/Index/IndexV1/when_opening_a_ptable_with_one_reader.cs delete mode 100644 src/EventStore.Core.Tests/TransactionLog/Unbuffered/UnbufferedTests.cs delete mode 100644 src/EventStore.Core/TransactionLog/Unbuffered/UnbufferedFileStream.cs delete mode 100644 src/EventStore.Native/EventStore.Native.csproj delete mode 100644 src/EventStore.Native/TransactionLog/Unbuffered/ExtendedFileOptions.cs delete mode 100644 src/EventStore.Native/TransactionLog/Unbuffered/INativeFile.cs delete mode 100644 src/EventStore.Native/TransactionLog/Unbuffered/MacCaching.cs delete mode 100644 src/EventStore.Native/TransactionLog/Unbuffered/NativeFileUnix.cs delete mode 100644 src/EventStore.Native/TransactionLog/Unbuffered/NativeFileWindows.cs delete mode 100644 src/EventStore.Native/TransactionLog/Unbuffered/WinNative.cs diff --git a/NOTICE.html b/NOTICE.html index 164ee3ebcb..b94072c048 100644 --- a/NOTICE.html +++ b/NOTICE.html @@ -532,13 +532,6 @@

Third-party software list

© Microsoft Corporation. All rights reserved. - - Mono.Posix.NETStandard1.0.0 -
- MIT -
- - NETStandard.Library2.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/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/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 From bd4752612146297f282c12c230a4b45946f965f2 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Mon, 10 Aug 2026 13:58:52 -0400 Subject: [PATCH 2/2] fix(index): keep PTable construction allocation-efficient Signed-off-by: Yordis Prieto --- .../Index/MemTableTests.cs | 39 +++++++++++++++++++ .../Index/ReverseComparerTests.cs | 20 ++++++++++ src/EventStore.Core/Index/HashListMemTable.cs | 23 +++++------ 3 files changed, 71 insertions(+), 11 deletions(-) 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/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); } }