diff --git a/DuckDB.NET.Benchmarks/AppenderBenchmark.cs b/DuckDB.NET.Benchmarks/AppenderBenchmark.cs new file mode 100644 index 00000000..162ed2e6 --- /dev/null +++ b/DuckDB.NET.Benchmarks/AppenderBenchmark.cs @@ -0,0 +1,67 @@ +using BenchmarkDotNet.Attributes; +using DuckDB.NET.Data; + +namespace DuckDB.NET.Benchmarks; + +[MemoryDiagnoser] +public class AppenderBenchmark +{ + private DuckDBConnection connection = null!; + + [Params(1_000_000)] + public int RowCount { get; set; } + + [GlobalSetup] + public void Setup() + { + connection = new DuckDBConnection("DataSource=:memory:"); + connection.Open(); + } + + [GlobalCleanup] + public void Cleanup() + { + connection.Dispose(); + } + + [IterationSetup] + public void IterationSetup() + { + using var command = connection.CreateCommand(); + command.CommandText = "DROP TABLE IF EXISTS bench; CREATE TABLE bench (a INTEGER, b BIGINT, c DOUBLE, d BOOLEAN);"; + command.ExecuteNonQuery(); + } + + [Benchmark(Baseline = true)] + public void AppendRowsWithCreateRow() + { + using var appender = connection.CreateAppender("bench"); + + for (var i = 0; i < RowCount; i++) + { + appender.CreateRow() + .AppendValue(i) + .AppendValue((long)i) + .AppendValue((double)i) + .AppendValue(i % 2 == 0) + .EndRow(); + } + } + + [Benchmark] + public void AppendRowsWithAppendRow() + { + using var appender = connection.CreateAppender("bench"); + + for (var i = 0; i < RowCount; i++) + { + appender.AppendRow(i, static (row, value) => + { + row.AppendValue(value) + .AppendValue((long)value) + .AppendValue((double)value) + .AppendValue(value % 2 == 0); + }); + } + } +} diff --git a/DuckDB.NET.Benchmarks/Benchmarks.csproj b/DuckDB.NET.Benchmarks/Benchmarks.csproj new file mode 100644 index 00000000..2de0564f --- /dev/null +++ b/DuckDB.NET.Benchmarks/Benchmarks.csproj @@ -0,0 +1,33 @@ + + + + Exe + net10.0 + enable + enable + true + Full + true + ..\keyPair.snk + + + + + + + + + + + + + + + false + PreserveNewest + runtimes\%(RecursiveDir)\%(FileName)%(Extension) + + + + diff --git a/DuckDB.NET.Benchmarks/MappedAppenderBenchmark.cs b/DuckDB.NET.Benchmarks/MappedAppenderBenchmark.cs new file mode 100644 index 00000000..1c1bec89 --- /dev/null +++ b/DuckDB.NET.Benchmarks/MappedAppenderBenchmark.cs @@ -0,0 +1,96 @@ +using BenchmarkDotNet.Attributes; +using DuckDB.NET.Data; +using DuckDB.NET.Data.Mapping; + +namespace DuckDB.NET.Benchmarks; + +[MemoryDiagnoser] +public class MappedAppenderBenchmark +{ + private DuckDBConnection connection = null!; + private BenchRow[] rows = null!; + private IPropertyMapping[] mappings = null!; + + [Params(1_000_000)] + public int RowCount { get; set; } + + public sealed class BenchRow + { + public int Id { get; init; } + public long Score { get; init; } + public double Value { get; init; } + public bool Active { get; init; } + } + + public sealed class BenchRowMap : DuckDBAppenderMap + { + public BenchRowMap() + { + Map(row => row.Id); + Map(row => row.Score); + Map(row => row.Value); + Map(row => row.Active); + } + } + + [GlobalSetup] + public void Setup() + { + connection = new DuckDBConnection("DataSource=:memory:"); + connection.Open(); + + rows = new BenchRow[RowCount]; + for (var i = 0; i < RowCount; i++) + { + rows[i] = new BenchRow { Id = i, Score = i, Value = i, Active = i % 2 == 0 }; + } + + mappings = new BenchRowMap().PropertyMappings.ToArray(); + } + + [GlobalCleanup] + public void Cleanup() + { + connection.Dispose(); + } + + [IterationSetup] + public void IterationSetup() + { + using var command = connection.CreateCommand(); + command.CommandText = "DROP TABLE IF EXISTS bench_mapped; CREATE TABLE bench_mapped (a INTEGER, b BIGINT, c DOUBLE, d BOOLEAN);"; + command.ExecuteNonQuery(); + } + + // Mirrors the mapped appender loop before it adopted AppendRow. + [Benchmark(Baseline = true)] + public void AppendMappedRowsWithCreateRow() + { + using var appender = connection.CreateAppender("bench_mapped"); + + foreach (var record in rows) + { + AppendRecordWithCreateRow(appender, record); + } + } + + [Benchmark] + public void AppendMappedRowsWithAppendRow() + { + using var appender = connection.CreateAppender("bench_mapped"); + appender.AppendRecords(rows); + } + + private void AppendRecordWithCreateRow(DuckDBAppender appender, BenchRow record) + { + ArgumentNullException.ThrowIfNull(record); + + var row = appender.CreateRow(); + foreach (var mapping in mappings) + { + mapping.AppendToRow(row, record); + } + + row.EndRow(); + } +} diff --git a/DuckDB.NET.Benchmarks/NativeLibraryLoader.cs b/DuckDB.NET.Benchmarks/NativeLibraryLoader.cs new file mode 100644 index 00000000..429de245 --- /dev/null +++ b/DuckDB.NET.Benchmarks/NativeLibraryLoader.cs @@ -0,0 +1,48 @@ +using System.Reflection; +using System.Runtime.CompilerServices; +using System.Runtime.InteropServices; + +namespace DuckDB.NET.Benchmarks; + +/// +/// Loads the platform-native DuckDB library for the benchmark process. +/// +internal static class NativeLibraryLoader +{ + [ModuleInitializer] + public static void Init() + { + if (GetRid() is not { } rid) + { + return; + } + + _ = NativeLibrary.TryLoad(Path.Join("runtimes", rid, "native", "duckdb"), Assembly.GetExecutingAssembly(), DllImportSearchPath.AssemblyDirectory, out _) || + NativeLibrary.TryLoad(Path.Join("runtimes", rid, "native", "libduckdb"), Assembly.GetExecutingAssembly(), DllImportSearchPath.AssemblyDirectory, out _); + } + + private static string? GetRid() + { + if (RuntimeInformation.IsOSPlatform(OSPlatform.Windows)) + { + return Environment.Is64BitProcess ? "win-x64" : "win-x86"; + } + + if (RuntimeInformation.IsOSPlatform(OSPlatform.Linux)) + { + return RuntimeInformation.ProcessArchitecture switch + { + Architecture.X64 => "linux-x64", + Architecture.Arm64 => "linux-arm64", + _ => null, + }; + } + + if (RuntimeInformation.IsOSPlatform(OSPlatform.OSX)) + { + return "osx"; + } + + return null; + } +} diff --git a/DuckDB.NET.Benchmarks/Program.cs b/DuckDB.NET.Benchmarks/Program.cs new file mode 100644 index 00000000..00251844 --- /dev/null +++ b/DuckDB.NET.Benchmarks/Program.cs @@ -0,0 +1,15 @@ +using BenchmarkDotNet.Configs; +using BenchmarkDotNet.Jobs; +using BenchmarkDotNet.Running; +using BenchmarkDotNet.Toolchains.InProcess.Emit; +using DuckDB.NET.Benchmarks; + +// The repo's Directory.Build.props renames the output assembly to DuckDB.NET.Benchmarks +// while the project file stays Benchmarks.csproj, so BenchmarkDotNet's default toolchain +// can't locate the csproj. Run in-process to avoid the separate build/spawn step. +var config = DefaultConfig.Instance + .AddJob(Job.Default.WithToolchain(InProcessEmitToolchain.Instance)); + +BenchmarkSwitcher + .FromTypes([typeof(AppenderBenchmark), typeof(MappedAppenderBenchmark)]) + .Run(args, config); diff --git a/DuckDB.NET.Data/Data.csproj b/DuckDB.NET.Data/Data.csproj index 4afd3811..184a03d9 100644 --- a/DuckDB.NET.Data/Data.csproj +++ b/DuckDB.NET.Data/Data.csproj @@ -30,6 +30,8 @@ Fixes: + + diff --git a/DuckDB.NET.Data/DuckDBAppender.cs b/DuckDB.NET.Data/DuckDBAppender.cs index d32e3ca7..ef9216e9 100644 --- a/DuckDB.NET.Data/DuckDBAppender.cs +++ b/DuckDB.NET.Data/DuckDBAppender.cs @@ -1,14 +1,20 @@ using DuckDB.NET.Data.Common; using DuckDB.NET.Data.DataChunk.Writer; using DuckDB.NET.Data.Extensions; -using System.Diagnostics; -using System.Diagnostics.CodeAnalysis; - namespace DuckDB.NET.Data; +/// +/// Appends rows to a DuckDB table. +/// +/// +/// Instances are not thread-safe. Do not call other methods on the same appender from an +/// callback. +/// public class DuckDBAppender : IDisposable { private bool closed; + private bool isAppendingRow; + private bool isFaulted; private readonly Native.DuckDBAppender nativeAppender; private readonly string qualifiedTableName; @@ -17,6 +23,7 @@ public class DuckDBAppender : IDisposable private readonly DuckDBLogicalType[] logicalTypes; private readonly DuckDBDataChunk dataChunk; private readonly VectorDataWriterBase[] vectorWriters; + private DuckDBAppenderRow? reusableRow; internal DuckDBAppender(Native.DuckDBAppender appender, string qualifiedTableName) { @@ -43,13 +50,110 @@ internal DuckDBAppender(Native.DuckDBAppender appender, string qualifiedTableNam /// internal IReadOnlyList LogicalTypes => logicalTypes; + /// + /// Creates an independent row. The caller must append every column and call + /// . + /// public IDuckDBAppenderRow CreateRow() { - if (closed) + EnsureUsable(); + return new DuckDBAppenderRow(qualifiedTableName, vectorWriters, PrepareRow(), dataChunk, nativeAppender); + } + + /// + /// Appends a complete row using a reusable row instance. + /// + /// A callback that appends every column value. This method calls + /// after the callback returns. + /// + /// The row is valid only during the callback and must not be retained. The callback must not + /// call other methods on this appender. If the callback or automatic + /// fails, all previously completed rows are flushed, + /// the failed row is discarded, and the appender cannot be reused. + /// + public void AppendRow(Action writeRow) + { + ArgumentNullException.ThrowIfNull(writeRow); + + // Pass the callback as state so the adapter remains static and allocation-free. + AppendRow(writeRow, static (row, callback) => callback(row)); + } + + /// + /// Appends a complete row using a reusable row instance without exposing that instance as a + /// return value. The row passed to is only valid for the duration of + /// the callback and must not be retained or used after the callback returns. + /// + /// The type of value used to populate the row. + /// The value used to populate the row. + /// A callback that appends every column value. This method calls + /// after the callback returns. + /// + /// The callback must not call other methods on this appender. If the callback or automatic + /// fails, all previously completed rows are flushed, + /// the failed row is discarded, and the appender cannot be reused. + /// + public void AppendRow(TState state, Action writeRow) + { + ArgumentNullException.ThrowIfNull(writeRow); + EnsureUsable(); + + DuckDBAppenderRow? row = null; + isAppendingRow = true; + + try { - throw new InvalidOperationException("Appender is already closed"); + row = CreateReusableRow(); + writeRow(row, state); + row.EndRow(); + } + catch (Exception appendException) + { + if (row is not null) + { + try + { + FinalizeFailedAppendRow(row); + } + catch (Exception finalizationException) + { + throw new AggregateException( + "Appending the row failed and the previously completed rows could not be finalized", + appendException, + finalizationException); + } + } + + throw; + } + finally + { + isAppendingRow = false; + } + } + + /// + /// Creates a row whose instance may be reused by the next call. This is only safe for internal + /// callers that create, populate and end each row without exposing the row reference. + /// + internal DuckDBAppenderRow CreateReusableRow() + { + var rowIndex = PrepareRow(); + + if (reusableRow is null) + { + reusableRow = new DuckDBAppenderRow(qualifiedTableName, vectorWriters, rowIndex, dataChunk, nativeAppender); + } + else + { + reusableRow.Reset(rowIndex); } + return reusableRow; + } + + private ulong PrepareRow() + { if (rowCount % DuckDBGlobalData.VectorSize == 0) { AppendDataChunk(); @@ -60,16 +164,13 @@ public IDuckDBAppenderRow CreateRow() } rowCount++; - return new DuckDBAppenderRow(qualifiedTableName, vectorWriters, rowCount - 1, dataChunk, nativeAppender); + return rowCount - 1; } public void Clear() { - if (closed) - { - throw new InvalidOperationException("Appender is already closed"); - } - + EnsureUsable(); + var state = NativeMethods.Appender.DuckDBAppenderClear(nativeAppender); if (!state.IsSuccess()) { @@ -78,9 +179,16 @@ public void Clear() rowCount = 0; NativeMethods.DataChunks.DuckDBDataChunkReset(dataChunk); + InitVectorWriters(); } public void Close() + { + EnsureUsable(); + CloseCore(); + } + + private void CloseCore() { closed = true; @@ -88,18 +196,6 @@ public void Close() { AppendDataChunk(); - foreach (var logicalType in logicalTypes) - { - logicalType.Dispose(); - } - - foreach (var writer in vectorWriters) - { - writer?.Dispose(); - } - - dataChunk.Dispose(); - var state = NativeMethods.Appender.DuckDBAppenderClose(nativeAppender); if (!state.IsSuccess()) { @@ -108,8 +204,30 @@ public void Close() } finally { - nativeAppender.Close(); + try + { + DisposeManagedResources(); + } + finally + { + nativeAppender.Close(); + } + } + } + + private void DisposeManagedResources() + { + foreach (var logicalType in logicalTypes) + { + logicalType.Dispose(); + } + + foreach (var writer in vectorWriters) + { + writer?.Dispose(); } + + dataChunk.Dispose(); } public void Dispose() @@ -143,4 +261,37 @@ private void AppendDataChunk() NativeMethods.DataChunks.DuckDBDataChunkReset(dataChunk); } + + private void FinalizeFailedAppendRow(DuckDBAppenderRow row) + { + // The row index is also the number of completed rows before the failed row in this chunk. + rowCount = row.ChunkRowIndex; + row.Invalidate(); + isFaulted = true; + + CloseCore(); + } + + private void EnsureNotAppendingRow() + { + if (isAppendingRow) + { + throw new InvalidOperationException("The appender cannot be used from inside an AppendRow callback"); + } + } + + private void EnsureUsable() + { + EnsureNotAppendingRow(); + + if (isFaulted) + { + throw new InvalidOperationException("The appender cannot be reused after an AppendRow callback failed"); + } + + if (closed) + { + throw new InvalidOperationException("Appender is already closed"); + } + } } diff --git a/DuckDB.NET.Data/DuckDBAppenderRow.cs b/DuckDB.NET.Data/DuckDBAppenderRow.cs index 186ab69a..cd697b2a 100644 --- a/DuckDB.NET.Data/DuckDBAppenderRow.cs +++ b/DuckDB.NET.Data/DuckDBAppenderRow.cs @@ -7,10 +7,12 @@ public class DuckDBAppenderRow : IDuckDBAppenderRow private int columnIndex = 0; private readonly string qualifiedTableName; private readonly VectorDataWriterBase[] vectorWriters; - private readonly ulong rowIndex; + private ulong rowIndex; private readonly DuckDBDataChunk dataChunk; private readonly Native.DuckDBAppender nativeAppender; + internal ulong ChunkRowIndex => rowIndex; + internal DuckDBAppenderRow(string qualifiedTableName, VectorDataWriterBase[] vectorWriters, ulong rowIndex, DuckDBDataChunk dataChunk, Native.DuckDBAppender nativeAppender) { @@ -21,6 +23,23 @@ internal DuckDBAppenderRow(string qualifiedTableName, VectorDataWriterBase[] vec this.nativeAppender = nativeAppender; } + /// + /// Re-targets this row instance at a new row index so the appender can reuse a single + /// instead of allocating one per row. The table name, vector + /// writers, data chunk and native appender are stable for the lifetime of the appender, so only + /// the row index and column cursor need to be reset. + /// + internal void Reset(ulong rowIndex) + { + this.rowIndex = rowIndex; + columnIndex = 0; + } + + internal void Invalidate() + { + columnIndex = vectorWriters.Length; + } + public void EndRow() { if (columnIndex < vectorWriters.Length) diff --git a/DuckDB.NET.Data/DuckDBMappedAppender.cs b/DuckDB.NET.Data/DuckDBMappedAppender.cs index ee87f303..25a3444f 100644 --- a/DuckDB.NET.Data/DuckDBMappedAppender.cs +++ b/DuckDB.NET.Data/DuckDBMappedAppender.cs @@ -76,14 +76,14 @@ private void AppendRecord(T record) throw new ArgumentNullException(nameof(record)); } - var row = appender.CreateRow(); - - foreach (var mapping in mappings) + // Pass both values as state so the callback does not capture per record. + appender.AppendRow((Record: record, Mappings: mappings), static (row, state) => { - row = mapping.AppendToRow(row, record); - } - - row.EndRow(); + foreach (var mapping in state.Mappings) + { + mapping.AppendToRow(row, state.Record); + } + }); } private static DuckDBType GetExpectedDuckDBType(Type type) diff --git a/DuckDB.NET.Data/Mapping/DuckDBAppenderMap.cs b/DuckDB.NET.Data/Mapping/DuckDBAppenderMap.cs index 0f8d576f..ff60dcf6 100644 --- a/DuckDB.NET.Data/Mapping/DuckDBAppenderMap.cs +++ b/DuckDB.NET.Data/Mapping/DuckDBAppenderMap.cs @@ -68,7 +68,7 @@ internal interface IPropertyMapping { Type PropertyType { get; } PropertyMappingType MappingType { get; } - IDuckDBAppenderRow AppendToRow(IDuckDBAppenderRow row, T record); + void AppendToRow(IDuckDBAppenderRow row, T record); } internal sealed class PropertyMapping : IPropertyMapping @@ -77,16 +77,17 @@ internal sealed class PropertyMapping : IPropertyMapping public Func Getter { get; set; } = _ => default!; public PropertyMappingType MappingType { get; set; } - public IDuckDBAppenderRow AppendToRow(IDuckDBAppenderRow row, T record) + public void AppendToRow(IDuckDBAppenderRow row, T record) { var value = Getter(record); if (value is null) { - return row.AppendNullValue(); + row.AppendNullValue(); + return; } - return value switch + _ = value switch { // Reference types string v => row.AppendValue(v), @@ -125,9 +126,9 @@ internal sealed class DefaultValueMapping : IPropertyMapping public Type PropertyType { get; set; } = typeof(object); public PropertyMappingType MappingType { get; set; } - public IDuckDBAppenderRow AppendToRow(IDuckDBAppenderRow row, T record) + public void AppendToRow(IDuckDBAppenderRow row, T record) { - return row.AppendDefault(); + row.AppendDefault(); } } @@ -136,8 +137,8 @@ internal sealed class NullValueMapping : IPropertyMapping public Type PropertyType { get; set; } = typeof(object); public PropertyMappingType MappingType { get; set; } - public IDuckDBAppenderRow AppendToRow(IDuckDBAppenderRow row, T record) + public void AppendToRow(IDuckDBAppenderRow row, T record) { - return row.AppendNullValue(); + row.AppendNullValue(); } } diff --git a/DuckDB.NET.Test/DuckDBManagedAppenderTests.cs b/DuckDB.NET.Test/DuckDBManagedAppenderTests.cs index e81d8618..f570f49b 100644 --- a/DuckDB.NET.Test/DuckDBManagedAppenderTests.cs +++ b/DuckDB.NET.Test/DuckDBManagedAppenderTests.cs @@ -1,4 +1,5 @@ using System.Globalization; +using DuckDB.NET.Data.Common; using FluentAssertions.Common; namespace DuckDB.NET.Test; @@ -483,23 +484,28 @@ public void WrongTypesThrowException() } [Fact] - public void ClosedAdapterThrowException() + public void ClosedAppenderRejectsFurtherOperations() { var table = "CREATE TABLE managedAppenderClosedAdapterTest(a BOOLEAN, c Date, b TINYINT);"; Command.CommandText = table; Command.ExecuteNonQuery(); - Connection.Invoking(dbConnection => - { - using var appender = dbConnection.CreateAppender("managedAppenderClosedAdapterTest"); - appender.Close(); - var row = appender.CreateRow(); - row - .AppendValue(false) - .AppendValue((byte)1) - .AppendValue((short?)1) - .EndRow(); - }).Should().Throw(); + using var appender = Connection.CreateAppender("managedAppenderClosedAdapterTest"); + appender.Close(); + + appender.Invoking(value => value.CreateRow()) + .Should().Throw() + .WithMessage("Appender is already closed"); + appender.Invoking(value => value.AppendRow(row => row.AppendValue(false).AppendNullValue().AppendNullValue())) + .Should().Throw() + .WithMessage("Appender is already closed"); + appender.Invoking(value => value.Clear()) + .Should().Throw() + .WithMessage("Appender is already closed"); + appender.Invoking(value => value.Close()) + .Should().Throw() + .WithMessage("Appender is already closed"); + appender.Invoking(value => value.Dispose()).Should().NotThrow(); } [Fact] @@ -715,6 +721,322 @@ public void AppendDefault() reader.GetInt32(2).Should().Be(30); } + [Fact] + public void CreateRowReturnsIndependentRows() + { + Command.CommandText = "CREATE TABLE managedAppenderRowLifetime(a INTEGER, b INTEGER)"; + Command.ExecuteNonQuery(); + + using (var appender = Connection.CreateAppender("managedAppenderRowLifetime")) + { + var completed = appender.CreateRow(); + completed.AppendValue((int?)10).AppendValue((int?)11).EndRow(); + + var current = appender.CreateRow(); + + completed.Should().NotBeSameAs(current); + completed.Invoking(row => row.AppendValue((int?)99)).Should().Throw(); + current.AppendValue((int?)20).AppendValue((int?)21).EndRow(); + } + + Command.CommandText = "SELECT a, b FROM managedAppenderRowLifetime ORDER BY a"; + using var reader = Command.ExecuteReader(); + reader.Read().Should().BeTrue(); + reader.GetInt32(0).Should().Be(10); + reader.GetInt32(1).Should().Be(11); + reader.Read().Should().BeTrue(); + reader.GetInt32(0).Should().Be(20); + reader.GetInt32(1).Should().Be(21); + reader.Read().Should().BeFalse(); + } + + [Fact] + public void AppendRowStateOverloadWritesRows() + { + Command.CommandText = "CREATE TABLE managedAppenderScopedRow(a INTEGER, b VARCHAR)"; + Command.ExecuteNonQuery(); + + using (var appender = Connection.CreateAppender("managedAppenderScopedRow")) + { + for (var i = 0; i < 3; i++) + { + appender.AppendRow((Id: i, Name: $"row-{i}"), static (row, value) => + { + row.AppendValue(value.Id).AppendValue(value.Name); + }); + } + } + + Command.CommandText = "SELECT a, b FROM managedAppenderScopedRow ORDER BY a"; + using var reader = Command.ExecuteReader(); + for (var i = 0; i < 3; i++) + { + reader.Read().Should().BeTrue(); + reader.GetInt32(0).Should().Be(i); + reader.GetString(1).Should().Be($"row-{i}"); + } + + reader.Read().Should().BeFalse(); + } + + [Fact] + public void AppendRowActionOverloadWritesCompleteRow() + { + Command.CommandText = "CREATE TABLE managedAppenderActionRow(a INTEGER, b VARCHAR)"; + Command.ExecuteNonQuery(); + + using (var appender = Connection.CreateAppender("managedAppenderActionRow")) + { + appender.AppendRow(row => row.AppendValue((int?)42).AppendValue("answer")); + } + + Command.CommandText = "SELECT a, b FROM managedAppenderActionRow"; + using var reader = Command.ExecuteReader(); + reader.Read().Should().BeTrue(); + reader.GetInt32(0).Should().Be(42); + reader.GetString(1).Should().Be("answer"); + reader.Read().Should().BeFalse(); + } + + [Fact] + public void AppendRowActionOverloadRejectsNullCallback() + { + Command.CommandText = "CREATE TABLE managedAppenderNullAction(a INTEGER)"; + Command.ExecuteNonQuery(); + + using var appender = Connection.CreateAppender("managedAppenderNullAction"); + appender.Invoking(value => value.AppendRow((Action)null!)) + .Should().Throw() + .WithParameterName("writeRow"); + } + + [Fact] + public void AppendRowStateOverloadRejectsNullCallback() + { + Command.CommandText = "CREATE TABLE managedAppenderNullStateAction(a INTEGER)"; + Command.ExecuteNonQuery(); + + using var appender = Connection.CreateAppender("managedAppenderNullStateAction"); + appender.Invoking(value => value.AppendRow(1, (Action)null!)) + .Should().Throw() + .WithParameterName("writeRow"); + } + + [Fact] + public void IncompleteAppendRowDiscardsFailedRowAndFaultsAppender() + { + Command.CommandText = "CREATE TABLE managedAppenderIncompleteScopedRow(a INTEGER, b INTEGER)"; + Command.ExecuteNonQuery(); + + using (var appender = Connection.CreateAppender("managedAppenderIncompleteScopedRow")) + { + appender.Invoking(value => value.AppendRow(1, static (row, state) => row.AppendValue(state))) + .Should().Throw() + .WithMessage("*specified only 1 values"); + + appender.Invoking(value => value.AppendRow((2, 3), static (row, state) => + row.AppendValue(state.Item1).AppendValue(state.Item2))) + .Should().Throw() + .WithMessage("*cannot be reused*"); + } + + Command.CommandText = "SELECT a, b FROM managedAppenderIncompleteScopedRow"; + using var reader = Command.ExecuteReader(); + reader.Read().Should().BeFalse(); + } + + [Fact] + public void AppendRowFailureFlushesCompletedRowsAndFaultsAppender() + { + Command.CommandText = "CREATE TABLE managedAppenderThrownScopedRow(a INTEGER, b INTEGER[])"; + Command.ExecuteNonQuery(); + + IDuckDBAppenderRow failedRow = null!; + + using (var appender = Connection.CreateAppender("managedAppenderThrownScopedRow")) + { + appender.AppendRow(row => row.AppendValue((int?)1).AppendValue(new[] { 1, 2 })); + + appender.Invoking(value => value.AppendRow(row => + { + failedRow = row; + row.AppendValue((int?)2).AppendValue(new[] { 3, 4 }); + throw new InvalidOperationException("callback failed"); + })) + .Should().Throw() + .WithMessage("callback failed"); + + failedRow.Should().NotBeNull(); + failedRow.Invoking(row => row.AppendValue((int?)99)).Should().Throw(); + + appender.Invoking(value => value.AppendRow(row => row.AppendValue((int?)3).AppendValue(new[] { 5, 6 }))) + .Should().Throw() + .WithMessage("*cannot be reused*"); + + appender.Invoking(value => value.CreateRow()) + .Should().Throw() + .WithMessage("*cannot be reused*"); + + appender.Invoking(value => value.Clear()) + .Should().Throw() + .WithMessage("*cannot be reused*"); + + appender.Invoking(value => value.Close()) + .Should().Throw() + .WithMessage("*cannot be reused*"); + } + + Command.CommandText = "SELECT a, b FROM managedAppenderThrownScopedRow ORDER BY a"; + using var reader = Command.ExecuteReader(); + reader.Read().Should().BeTrue(); + reader.GetInt32(0).Should().Be(1); + reader.GetFieldValue>(1).Should().Equal(1, 2); + reader.Read().Should().BeFalse(); + } + + [Fact] + public void AppendRowFailureAggregatesFinalizationFailureAndCanBeDisposed() + { + Command.CommandText = "CREATE TABLE managedAppenderFailedFinalization(a INTEGER UNIQUE)"; + Command.ExecuteNonQuery(); + Command.CommandText = "INSERT INTO managedAppenderFailedFinalization VALUES (1)"; + Command.ExecuteNonQuery(); + + using var appender = Connection.CreateAppender("managedAppenderFailedFinalization"); + appender.AppendRow(row => row.AppendValue((int?)1)); + + var exception = appender.Invoking(value => value.AppendRow(row => + { + row.AppendValue((int?)2); + throw new InvalidOperationException("callback failed"); + })) + .Should().Throw().Which; + + exception.InnerExceptions.Should().HaveCount(2); + exception.InnerExceptions[0].Should().BeOfType() + .Which.Message.Should().Be("callback failed"); + exception.InnerExceptions[1].Should().BeOfType() + .Which.ErrorType.Should().Be(DuckDBErrorType.Constraint); + + appender.Invoking(value => value.Dispose()).Should().NotThrow(); + appender.Invoking(value => value.Dispose()).Should().NotThrow(); + + Command.CommandText = "SELECT count(*) FROM managedAppenderFailedFinalization"; + Command.ExecuteScalar().Should().Be(1); + } + + [Fact] + public void AppendRowFailureFlushesCompletedRowsAcrossDataChunks() + { + Command.CommandText = "CREATE TABLE managedAppenderFailedAcrossChunks(a INTEGER)"; + Command.ExecuteNonQuery(); + + var completedRowCount = checked((int)DuckDBGlobalData.VectorSize + 3); + + using (var appender = Connection.CreateAppender("managedAppenderFailedAcrossChunks")) + { + for (var i = 0; i < completedRowCount; i++) + { + appender.AppendRow(i, static (row, value) => row.AppendValue(value)); + } + + appender.Invoking(value => value.AppendRow(completedRowCount, static (row, failedValue) => + { + row.AppendValue(failedValue); + throw new InvalidOperationException("callback failed"); + })) + .Should().Throw() + .WithMessage("callback failed"); + } + + Command.CommandText = """ + SELECT count(*)::BIGINT, min(a), max(a), sum(a)::BIGINT + FROM managedAppenderFailedAcrossChunks + """; + using var reader = Command.ExecuteReader(); + reader.Read().Should().BeTrue(); + reader.GetInt64(0).Should().Be(completedRowCount); + reader.GetInt32(1).Should().Be(0); + reader.GetInt32(2).Should().Be(completedRowCount - 1); + reader.GetInt64(3).Should().Be((long)(completedRowCount - 1) * completedRowCount / 2); + } + + [Fact] + public void AppendRowFailureLeavesCompletedRowsInCallerTransaction() + { + Command.CommandText = "CREATE TABLE managedAppenderFailedTransaction(a INTEGER)"; + Command.ExecuteNonQuery(); + + using (var transaction = Connection.BeginTransaction()) + { + using (var appender = Connection.CreateAppender("managedAppenderFailedTransaction")) + { + appender.AppendRow(row => row.AppendValue((int?)1)); + + appender.Invoking(value => value.AppendRow(row => + { + row.AppendValue((int?)2); + throw new InvalidOperationException("callback failed"); + })) + .Should().Throw() + .WithMessage("callback failed"); + } + + Command.CommandText = "SELECT count(*) FROM managedAppenderFailedTransaction"; + Command.ExecuteScalar().Should().Be(1); + + transaction.Rollback(); + } + + Command.CommandText = "SELECT count(*) FROM managedAppenderFailedTransaction"; + Command.ExecuteScalar().Should().Be(0); + } + + [Fact] + public void AppendRowWritesListValue() + { + Command.CommandText = "CREATE TABLE managedAppenderScopedListRow(a INTEGER, b INTEGER[])"; + Command.ExecuteNonQuery(); + + using (var appender = Connection.CreateAppender("managedAppenderScopedListRow")) + { + appender.AppendRow(row => row.AppendValue((int?)1).AppendValue(new[] { 1, 2 })); + } + + Command.CommandText = "SELECT a, b FROM managedAppenderScopedListRow"; + using var reader = Command.ExecuteReader(); + reader.Read().Should().BeTrue(); + reader.GetInt32(0).Should().Be(1); + reader.GetFieldValue>(1).Should().Equal(1, 2); + reader.Read().Should().BeFalse(); + } + + [Fact] + public void AppendRowRejectsReentrantAppenderUseAndFaultsAppender() + { + Command.CommandText = "CREATE TABLE managedAppenderReentrantScopedRow(a INTEGER)"; + Command.ExecuteNonQuery(); + + using (var appender = Connection.CreateAppender("managedAppenderReentrantScopedRow")) + { + appender.Invoking(value => value.AppendRow(row => + { + row.AppendValue((int?)1); + value.AppendRow(nested => nested.AppendValue((int?)2)); + })) + .Should().Throw() + .WithMessage("*inside an AppendRow callback"); + + appender.Invoking(value => value.AppendRow(row => row.AppendValue((int?)3))) + .Should().Throw() + .WithMessage("*cannot be reused*"); + } + + Command.CommandText = "SELECT a FROM managedAppenderReentrantScopedRow"; + using var reader = Command.ExecuteReader(); + reader.Read().Should().BeFalse(); + } + [Fact] public void ClearAppender() { @@ -851,4 +1173,4 @@ private enum EnumNotValidValueTestEnum { NotValid = 12345, } -} \ No newline at end of file +} diff --git a/DuckDB.NET.slnx b/DuckDB.NET.slnx index 715e640c..b40a6d65 100644 --- a/DuckDB.NET.slnx +++ b/DuckDB.NET.slnx @@ -1,4 +1,5 @@ +