From 5a1008ba43e175dfa916a69573e92737fad2df36 Mon Sep 17 00:00:00 2001 From: Widthdom Date: Mon, 25 May 2026 12:08:37 +0900 Subject: [PATCH 1/3] Fix DbWriter transaction scope races (#1733 #1741 #1783 #2676) --- changelog.d/unreleased/1733.fixed.md | 16 ++ changelog.d/unreleased/1741.fixed.md | 16 ++ changelog.d/unreleased/1783.fixed.md | 16 ++ changelog.d/unreleased/2676.fixed.md | 16 ++ src/CodeIndex/Database/DbWriter.cs | 264 ++++++++++++++-------- tests/CodeIndex.Tests/ConcurrencyTests.cs | 56 +++++ tests/CodeIndex.Tests/DatabaseTests.cs | 45 ++++ 7 files changed, 339 insertions(+), 90 deletions(-) create mode 100644 changelog.d/unreleased/1733.fixed.md create mode 100644 changelog.d/unreleased/1741.fixed.md create mode 100644 changelog.d/unreleased/1783.fixed.md create mode 100644 changelog.d/unreleased/2676.fixed.md diff --git a/changelog.d/unreleased/1733.fixed.md b/changelog.d/unreleased/1733.fixed.md new file mode 100644 index 0000000000..75eba83321 --- /dev/null +++ b/changelog.d/unreleased/1733.fixed.md @@ -0,0 +1,16 @@ +--- +category: fixed +issues: + - 1733 +affected: + - src/CodeIndex/Database/DbWriter.cs + - tests/CodeIndex.Tests/ConcurrencyTests.cs +--- + +## English + +- **DbWriter transaction depth is serialized for shared writers (#1733)** — transaction scopes on one writer now wait across threads while preserving same-thread nested savepoints, preventing depth races from leaking or mis-targeting savepoints. + +## 日本語 + +- **共有 DbWriter の transaction depth を直列化しました (#1733)** — 1 つの writer 上の transaction scope はスレッド間で待機し、同一スレッドの nested savepoint は維持するため、depth race による savepoint 漏れや誤対象化を防ぎます。 diff --git a/changelog.d/unreleased/1741.fixed.md b/changelog.d/unreleased/1741.fixed.md new file mode 100644 index 0000000000..f3a973eeec --- /dev/null +++ b/changelog.d/unreleased/1741.fixed.md @@ -0,0 +1,16 @@ +--- +category: fixed +issues: + - 1741 +affected: + - src/CodeIndex/Database/DbWriter.cs + - tests/CodeIndex.Tests/DatabaseTests.cs +--- + +## English + +- **Savepoint cleanup failures are now diagnosable (#1741)** — savepoint scopes now throw an explicit invalid-operation error when their SQLite connection is missing and log cleanup rollback failures before best-effort disposal continues. + +## 日本語 + +- **savepoint cleanup の失敗を診断できるようにしました (#1741)** — savepoint scope は SQLite connection が欠けている場合に明示的な invalid-operation error を投げ、cleanup rollback の失敗は best-effort dispose を続ける前にログへ残します。 diff --git a/changelog.d/unreleased/1783.fixed.md b/changelog.d/unreleased/1783.fixed.md new file mode 100644 index 0000000000..a75eaa2ca1 --- /dev/null +++ b/changelog.d/unreleased/1783.fixed.md @@ -0,0 +1,16 @@ +--- +category: fixed +issues: + - 1783 +affected: + - src/CodeIndex/Database/DbWriter.cs + - tests/CodeIndex.Tests/DatabaseTests.cs +--- + +## English + +- **DbWriter transaction depth now survives begin failures (#1783)** — failed SQLite transaction or savepoint starts no longer leave `_transactionDepth` advanced without a matching active transaction. + +## 日本語 + +- **DbWriter の transaction depth が begin 失敗後も整合するようになりました (#1783)** — SQLite transaction / savepoint の開始に失敗しても、対応する active transaction がないまま `_transactionDepth` だけが進む状態を残しません。 diff --git a/changelog.d/unreleased/2676.fixed.md b/changelog.d/unreleased/2676.fixed.md new file mode 100644 index 0000000000..fe99d9e988 --- /dev/null +++ b/changelog.d/unreleased/2676.fixed.md @@ -0,0 +1,16 @@ +--- +category: fixed +issues: + - 2676 +affected: + - src/CodeIndex/Database/DbWriter.cs + - tests/CodeIndex.Tests/ConcurrencyTests.cs +--- + +## English + +- **Ready-bit stamping no longer starts nested SQLite transactions on shared writers (#2676)** — `SetReadyBit` and fold-ready stamping now share the writer transaction gate, so concurrent calls wait for an active scope before issuing `BEGIN IMMEDIATE`. + +## 日本語 + +- **ready-bit stamp が共有 writer 上で nested SQLite transaction を開始しないようにしました (#2676)** — `SetReadyBit` と fold-ready stamp は writer transaction gate を共有し、active scope が終わるまで待ってから `BEGIN IMMEDIATE` を発行します。 diff --git a/src/CodeIndex/Database/DbWriter.cs b/src/CodeIndex/Database/DbWriter.cs index 36483f343b..d0513210fb 100644 --- a/src/CodeIndex/Database/DbWriter.cs +++ b/src/CodeIndex/Database/DbWriter.cs @@ -1,4 +1,5 @@ using Microsoft.Data.Sqlite; +using CodeIndex.Cli; using CodeIndex.Indexer; using CodeIndex.Models; using System.Text; @@ -7,7 +8,10 @@ namespace CodeIndex.Database; /// /// Handles INSERT/UPSERT operations to the database with batch commits. +/// Transaction scopes on a writer are serialized; nested scopes are supported on the +/// owning thread and other threads wait until the active scope is disposed. /// バッチコミットによるINSERT/UPSERT処理を担当する。 +/// writer 上の transaction scope は直列化され、同一所有スレッドのネストのみ許可する。 /// public class DbWriter { @@ -20,12 +24,15 @@ public class DbWriter private readonly Action? _markWriteWork; internal static Action? FoldBackfillRowUpdatedForTesting { get; set; } internal static Action? BatchRowSkipWarningForTesting { get; set; } + private readonly object _transactionStateLock = new(); + private readonly SemaphoreSlim _transactionGate = new(1, 1); private const int BatchSize = 500; private const int MaxSqlVariables = 999; private const int SqliteConstraintErrorCode = 19; private int _rowSkipSavepointCounter; private long _batchRowsSkipped; private int _transactionDepth; + private int _transactionOwnerThreadId; // Outermost SqliteTransaction currently held open by this writer (null when no // transaction is active OR after the outermost transaction has been committed / // rolled back). Tracked so cached prepared commands can be re-pointed at the live @@ -119,22 +126,66 @@ private void ReleaseCommand(SqliteCommand cmd) /// public TransactionScope BeginTransaction() { - if (_transactionDepth == 0) + var ownsTransactionGate = EnterTransactionGate(); + try { - _transactionDepth++; - var txn = _conn.BeginTransaction(); - _activeTransaction = txn; - return new TransactionScope(txn, this); + if (_transactionDepth == 0) + { + var txn = _conn.BeginTransaction(); + _transactionDepth = 1; + _activeTransaction = txn; + return new TransactionScope(txn, this, ownsTransactionGate); + } + else + { + // Nested: use SAVEPOINT instead of BEGIN TRANSACTION + // ネスト: BEGIN TRANSACTIONの代わりにSAVEPOINTを使用 + var name = $"sp_{_transactionDepth}"; + Execute($"SAVEPOINT {name}"); + _transactionDepth++; + return new TransactionScope(name, _conn, this, ownsTransactionGate); + } } - else + catch + { + if (ownsTransactionGate) + ExitTransactionGate(); + throw; + } + } + + private bool EnterTransactionGate() + { + while (true) + { + lock (_transactionStateLock) + { + if (_transactionDepth > 0 && _transactionOwnerThreadId == Environment.CurrentManagedThreadId) + return false; + } + + _transactionGate.Wait(); + lock (_transactionStateLock) + { + if (_transactionDepth == 0) + { + _transactionOwnerThreadId = Environment.CurrentManagedThreadId; + return true; + } + } + + _transactionGate.Release(); + Thread.Yield(); + } + } + + private void ExitTransactionGate() + { + lock (_transactionStateLock) { - // Nested: use SAVEPOINT instead of BEGIN TRANSACTION - // ネスト: BEGIN TRANSACTIONの代わりにSAVEPOINTを使用 - var name = $"sp_{_transactionDepth}"; - _transactionDepth++; - Execute($"SAVEPOINT {name}"); - return new TransactionScope(name, _conn, this); + _transactionOwnerThreadId = 0; } + _transactionGate.Release(); } /// @@ -149,6 +200,7 @@ public sealed class TransactionScope : IDisposable private readonly string? _savepointName; private readonly SqliteConnection? _conn; private readonly DbWriter _writer; + private readonly bool _ownsTransactionGate; private const int StateActive = 0; private const int StateCommitting = 1; private const int StateCommitted = 2; @@ -158,18 +210,20 @@ public sealed class TransactionScope : IDisposable private int _disposeStarted; // Real transaction / 実トランザクション - internal TransactionScope(SqliteTransaction transaction, DbWriter writer) + internal TransactionScope(SqliteTransaction transaction, DbWriter writer, bool ownsTransactionGate = false) { _transaction = transaction; _writer = writer; + _ownsTransactionGate = ownsTransactionGate; } // Savepoint / セーブポイント - internal TransactionScope(string savepointName, SqliteConnection conn, DbWriter writer) + internal TransactionScope(string savepointName, SqliteConnection conn, DbWriter writer, bool ownsTransactionGate = false) { _savepointName = savepointName; _conn = conn; _writer = writer; + _ownsTransactionGate = ownsTransactionGate; } public void Commit() @@ -303,9 +357,10 @@ public void Dispose() ExecuteSql($"ROLLBACK TO SAVEPOINT {_savepointName}"); Volatile.Write(ref _state, StateRolledBack); } - catch + catch (Exception ex) { // Best effort during dispose / Dispose中はベストエフォート + GlobalToolLog.Error($"transaction_scope_dispose_rollback_failed {GlobalToolLog.FormatExceptionChain(ex)}"); Volatile.Write(ref _state, StateRolledBack); } break; @@ -313,20 +368,31 @@ public void Dispose() } finally { - _transaction?.Dispose(); - if (_writer._transactionDepth > 0) _writer._transactionDepth--; - // Safety net: even if Commit/Rollback was bypassed (e.g. uncommitted scope - // disposed after an exception), make sure the outer-transaction reference is - // cleared before the next RentCommand sees it. - // 安全弁: Commit/Rollback を経由せず Dispose された場合でも active reference を解除。 - if (_writer._transactionDepth == 0) - _writer._activeTransaction = null; + try + { + _transaction?.Dispose(); + if (_writer._transactionDepth > 0) _writer._transactionDepth--; + // Safety net: even if Commit/Rollback was bypassed (e.g. uncommitted scope + // disposed after an exception), make sure the outer-transaction reference is + // cleared before the next RentCommand sees it. + // 安全弁: Commit/Rollback を経由せず Dispose された場合でも active reference を解除。 + if (_writer._transactionDepth == 0) + _writer._activeTransaction = null; + } + finally + { + if (_ownsTransactionGate) + _writer.ExitTransactionGate(); + } } } private void ExecuteSql(string sql) { - using var cmd = _conn!.CreateCommand(); + if (_conn is null) + throw new InvalidOperationException("Savepoint transaction scope is missing its SQLite connection."); + + using var cmd = _conn.CreateCommand(); cmd.CommandText = sql; cmd.ExecuteNonQuery(); _writer._markWriteWork?.Invoke(); @@ -1692,61 +1758,70 @@ public bool OptimizeFtsIfIncrementalWriteThresholdReached(int threshold = Defaul /// True when the bit was actually stamped; false when re-verification failed. public bool MarkFoldReady(bool stampCurrentSymbolExtractorVersions = false) { - bool ownTransaction = !IsInTransaction(); - if (ownTransaction) - Execute("BEGIN IMMEDIATE"); + var ownsTransactionGate = EnterTransactionGate(); try { - if (stampCurrentSymbolExtractorVersions) + bool ownTransaction = !IsInTransaction(); + if (ownTransaction) + Execute("BEGIN IMMEDIATE"); + try + { + if (stampCurrentSymbolExtractorVersions) + StampSymbolExtractorVersions(); + + if (!AllFoldedColumnsBackfilledCore( + requireCurrentSymbolExtractorVersions: false, + requireCurrentFoldKeys: true)) + { + if (ownTransaction) + { + Execute("COMMIT"); + ownTransaction = false; + } + return false; + } + + // Inline the SetReadyBit body. SetReadyBit opens its own BEGIN IMMEDIATE + // when not already in a DbWriter-tracked transaction, but our raw + // BEGIN IMMEDIATE above is not tracked in _transactionDepth, so a direct + // SetReadyBit call would attempt a nested BEGIN IMMEDIATE and fail. + // SetReadyBit は _transactionDepth ベースでしか外側 transaction を見ないため、 + // 生 BEGIN IMMEDIATE 内では呼べない。内容を inline 展開する。 + int current; + using (var read = _conn.CreateCommand()) + { + read.CommandText = "PRAGMA user_version"; + var raw = read.ExecuteScalar(); + current = raw is long l ? (int)l : (raw is int i ? i : 0); + } + int next = current | DbContext.FoldReadyFlag; + if (next != current) + Execute($"PRAGMA user_version = {next}"); + + SetMeta("fold_key_version", NameFold.Version.ToString(System.Globalization.CultureInfo.InvariantCulture)); + SetMeta("fold_key_fingerprint", NameFold.Fingerprint()); StampSymbolExtractorVersions(); - if (!AllFoldedColumnsBackfilledCore( - requireCurrentSymbolExtractorVersions: false, - requireCurrentFoldKeys: true)) - { if (ownTransaction) { Execute("COMMIT"); ownTransaction = false; } - return false; - } - - // Inline the SetReadyBit body. SetReadyBit opens its own BEGIN IMMEDIATE - // when not already in a DbWriter-tracked transaction, but our raw - // BEGIN IMMEDIATE above is not tracked in _transactionDepth, so a direct - // SetReadyBit call would attempt a nested BEGIN IMMEDIATE and fail. - // SetReadyBit は _transactionDepth ベースでしか外側 transaction を見ないため、 - // 生 BEGIN IMMEDIATE 内では呼べない。内容を inline 展開する。 - int current; - using (var read = _conn.CreateCommand()) - { - read.CommandText = "PRAGMA user_version"; - var raw = read.ExecuteScalar(); - current = raw is long l ? (int)l : (raw is int i ? i : 0); + return true; } - int next = current | DbContext.FoldReadyFlag; - if (next != current) - Execute($"PRAGMA user_version = {next}"); - - SetMeta("fold_key_version", NameFold.Version.ToString(System.Globalization.CultureInfo.InvariantCulture)); - SetMeta("fold_key_fingerprint", NameFold.Fingerprint()); - StampSymbolExtractorVersions(); - - if (ownTransaction) + catch (Exception) { - Execute("COMMIT"); - ownTransaction = false; + if (ownTransaction) + { + try { Execute("ROLLBACK"); } catch (SqliteException) { /* best effort */ } + } + throw; } - return true; } - catch (Exception) + finally { - if (ownTransaction) - { - try { Execute("ROLLBACK"); } catch (SqliteException) { /* best effort */ } - } - throw; + if (ownsTransactionGate) + ExitTransactionGate(); } } @@ -3164,39 +3239,48 @@ private static int GetRowsPerInsertStatement(int columnCount) private void SetReadyBit(int flag) { - // The ready bits share a single PRAGMA user_version word, so two parallel - // cdidx writers (e.g. CI + a local rebuild) can each read the same prior - // value, OR in their own flag, and the slower writer's PRAGMA write clobbers - // the faster writer's flag. Wrap the read-modify-write in BEGIN IMMEDIATE so - // SQLite's reserved write lock serialises it across processes (issue #1513). - bool ownTransaction = !IsInTransaction(); - if (ownTransaction) - Execute("BEGIN IMMEDIATE"); + var ownsTransactionGate = EnterTransactionGate(); try { - int current; - using (var read = _conn.CreateCommand()) + // The ready bits share a single PRAGMA user_version word, so two parallel + // cdidx writers (e.g. CI + a local rebuild) can each read the same prior + // value, OR in their own flag, and the slower writer's PRAGMA write clobbers + // the faster writer's flag. Wrap the read-modify-write in BEGIN IMMEDIATE so + // SQLite's reserved write lock serialises it across processes (issue #1513). + bool ownTransaction = !IsInTransaction(); + if (ownTransaction) + Execute("BEGIN IMMEDIATE"); + try { - read.CommandText = "PRAGMA user_version"; - var raw = read.ExecuteScalar(); - current = raw is long l ? (int)l : (raw is int i ? i : 0); + int current; + using (var read = _conn.CreateCommand()) + { + read.CommandText = "PRAGMA user_version"; + var raw = read.ExecuteScalar(); + current = raw is long l ? (int)l : (raw is int i ? i : 0); + } + int next = current | flag; + if (next != current) + Execute($"PRAGMA user_version = {next}"); + if (ownTransaction) + { + Execute("COMMIT"); + ownTransaction = false; + } } - int next = current | flag; - if (next != current) - Execute($"PRAGMA user_version = {next}"); - if (ownTransaction) + catch (Exception) { - Execute("COMMIT"); - ownTransaction = false; + if (ownTransaction) + { + try { Execute("ROLLBACK"); } catch (SqliteException) { /* best effort */ } + } + throw; } } - catch (Exception) + finally { - if (ownTransaction) - { - try { Execute("ROLLBACK"); } catch (SqliteException) { /* best effort */ } - } - throw; + if (ownsTransactionGate) + ExitTransactionGate(); } } diff --git a/tests/CodeIndex.Tests/ConcurrencyTests.cs b/tests/CodeIndex.Tests/ConcurrencyTests.cs index 995ddcca77..ee9ab6238c 100644 --- a/tests/CodeIndex.Tests/ConcurrencyTests.cs +++ b/tests/CodeIndex.Tests/ConcurrencyTests.cs @@ -522,6 +522,62 @@ public async Task SetReadyBit_ConcurrentWriters_DoNotLoseFlags() $"Sample: {string.Join(", ", lostFlagIterations.Take(3).Select(v => $"i={v.iteration} user_version={v.finalUserVersion}"))}"); } + [Fact] + public async Task BeginTransaction_SharedWriterSerializesConcurrentScopes() + { + var writer = new DbWriter(_db.Connection); + using var start = new ManualResetEventSlim(false); + var errors = new ConcurrentBag(); + + var tasks = Enumerable.Range(0, 8).Select(worker => Task.Run(() => + { + start.Wait(); + try + { + for (var i = 0; i < 20; i++) + { + using var txn = writer.BeginTransaction(); + writer.UpsertFile(new FileRecord + { + Path = $"src/worker{worker}_{i}.cs", + Lang = "csharp", + Size = 10, + Lines = 1, + Modified = new DateTime(2025, 1, 1, 0, 0, 0, DateTimeKind.Utc), + Checksum = $"{worker}_{i}", + }); + txn.Commit(); + } + } + catch (Exception ex) + { + errors.Add(ex); + } + })).ToArray(); + + start.Set(); + await Task.WhenAll(tasks); + + Assert.True(errors.IsEmpty, string.Join(Environment.NewLine, errors.Select(e => e.ToString()))); + } + + [Fact] + public async Task SetReadyBit_SharedWriterWaitsForActiveScope() + { + var writer = new DbWriter(_db.Connection); + using var txn = writer.BeginTransaction(); + + var markTask = Task.Run(() => writer.MarkGraphReady()); + await Task.Delay(100); + Assert.False(markTask.IsCompleted); + + txn.Commit(); + txn.Dispose(); + await markTask; + + Assert.True((_db.GetUserVersion() & DbContext.GraphReadyFlag) != 0); + } + private static List BuildReferenceBatch(long fileId, string label, int count) { var refs = new List(count); diff --git a/tests/CodeIndex.Tests/DatabaseTests.cs b/tests/CodeIndex.Tests/DatabaseTests.cs index 21a447b928..0cdb151b38 100644 --- a/tests/CodeIndex.Tests/DatabaseTests.cs +++ b/tests/CodeIndex.Tests/DatabaseTests.cs @@ -189,6 +189,44 @@ public void BatchInProgress_ClearInsideRolledBackTransaction_LeavesRecoveryWarni Assert.Contains("Last batch did not complete", stderr); } + [Fact] + public void BeginTransaction_WhenBeginFails_RestoresTransactionDepth() + { + using var connection = new SqliteConnection(new SqliteConnectionStringBuilder { DataSource = _dbPath }.ConnectionString); + var writer = new DbWriter(connection); + + var ex = Assert.Throws(() => writer.BeginTransaction()); + + Assert.Contains("connection", ex.Message, StringComparison.OrdinalIgnoreCase); + Assert.Equal(0, GetTransactionDepth(writer)); + } + + [Fact] + public void TransactionScope_SavepointWithoutConnection_ThrowsExplicitInvalidOperation() + { + var scopeType = typeof(DbWriter).GetNestedType("TransactionScope") + ?? throw new InvalidOperationException("TransactionScope type was not found."); + var scope = Activator.CreateInstance( + scopeType, + System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic, + binder: null, + args: ["sp_missing_conn", null!, _writer, false], + culture: null) + ?? throw new InvalidOperationException("TransactionScope instance was not created."); + + using var disposable = (IDisposable)scope; + var commit = scopeType.GetMethod("Commit") + ?? throw new InvalidOperationException("Commit method was not found."); + + var ex = Assert.ThrowsAny(() => commit.Invoke(scope, null)); + var actual = ex is System.Reflection.TargetInvocationException { InnerException: { } inner } + ? inner + : ex; + + var invalidOperation = Assert.IsType(actual); + Assert.Contains("SQLite connection", invalidOperation.Message); + } + [Fact] public void Constructor_NewDatabaseEnablesIncrementalAutoVacuum() { @@ -239,6 +277,13 @@ INSERT INTO vacuum_payload (payload) } } + private static int GetTransactionDepth(DbWriter writer) + { + var field = typeof(DbWriter).GetField("_transactionDepth", System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic) + ?? throw new InvalidOperationException("_transactionDepth field was not found."); + return (int)field.GetValue(writer)!; + } + [Fact] public void RunIncrementalVacuum_ConvertsLegacyNoAutoVacuumDatabase() { From 6afc53d80644986dc8b42d6624f414f57baabc50 Mon Sep 17 00:00:00 2001 From: Widthdom Date: Mon, 25 May 2026 12:38:51 +0900 Subject: [PATCH 2/3] Address DbWriter logical transaction ownership review (#1733 #2676) --- src/CodeIndex/Database/DbWriter.cs | 86 ++++++++++++++++------- tests/CodeIndex.Tests/ConcurrencyTests.cs | 31 +++++++- tests/CodeIndex.Tests/DatabaseTests.cs | 2 +- 3 files changed, 91 insertions(+), 28 deletions(-) diff --git a/src/CodeIndex/Database/DbWriter.cs b/src/CodeIndex/Database/DbWriter.cs index c77d7782c6..6b1b85c447 100644 --- a/src/CodeIndex/Database/DbWriter.cs +++ b/src/CodeIndex/Database/DbWriter.cs @@ -26,13 +26,14 @@ public class DbWriter internal static Action? BatchRowSkipWarningForTesting { get; set; } private readonly object _transactionStateLock = new(); private readonly SemaphoreSlim _transactionGate = new(1, 1); + private readonly AsyncLocal _currentTransactionGateToken = new(); private const int BatchSize = 500; private const int MaxSqlVariables = 999; private const int SqliteConstraintErrorCode = 19; private int _rowSkipSavepointCounter; private long _batchRowsSkipped; private int _transactionDepth; - private int _transactionOwnerThreadId; + private Guid _transactionOwnerToken; // Outermost SqliteTransaction currently held open by this writer (null when no // transaction is active OR after the outermost transaction has been committed / // rolled back). Tracked so cached prepared commands can be re-pointed at the live @@ -126,7 +127,7 @@ private void ReleaseCommand(SqliteCommand cmd) /// public TransactionScope BeginTransaction() { - var ownsTransactionGate = EnterTransactionGate(); + var gateLease = EnterTransactionGate(); try { if (_transactionDepth == 0) @@ -134,7 +135,7 @@ public TransactionScope BeginTransaction() var txn = _conn.BeginTransaction(); _transactionDepth = 1; _activeTransaction = txn; - return new TransactionScope(txn, this, ownsTransactionGate); + return new TransactionScope(txn, this, gateLease); } else { @@ -143,25 +144,24 @@ public TransactionScope BeginTransaction() var name = $"sp_{_transactionDepth}"; Execute($"SAVEPOINT {name}"); _transactionDepth++; - return new TransactionScope(name, _conn, this, ownsTransactionGate); + return new TransactionScope(name, _conn, this, gateLease); } } catch { - if (ownsTransactionGate) - ExitTransactionGate(); + gateLease.Dispose(); throw; } } - private bool EnterTransactionGate() + private TransactionGateLease EnterTransactionGate() { while (true) { lock (_transactionStateLock) { - if (_transactionDepth > 0 && _transactionOwnerThreadId == Environment.CurrentManagedThreadId) - return false; + if (_transactionDepth > 0 && _transactionOwnerToken != Guid.Empty && _currentTransactionGateToken.Value == _transactionOwnerToken) + return TransactionGateLease.None; } _transactionGate.Wait(); @@ -169,8 +169,11 @@ private bool EnterTransactionGate() { if (_transactionDepth == 0) { - _transactionOwnerThreadId = Environment.CurrentManagedThreadId; - return true; + var previousToken = _currentTransactionGateToken.Value; + var token = Guid.NewGuid(); + _transactionOwnerToken = token; + _currentTransactionGateToken.Value = token; + return new TransactionGateLease(this, token, previousToken); } } @@ -179,15 +182,39 @@ private bool EnterTransactionGate() } } - private void ExitTransactionGate() + private void ExitTransactionGate(Guid token, Guid? previousToken) { lock (_transactionStateLock) { - _transactionOwnerThreadId = 0; + if (_transactionOwnerToken == token) + _transactionOwnerToken = Guid.Empty; } + if (_currentTransactionGateToken.Value == token) + _currentTransactionGateToken.Value = previousToken; _transactionGate.Release(); } + internal readonly struct TransactionGateLease + { + public static readonly TransactionGateLease None = new(null, Guid.Empty, null); + + private readonly DbWriter? _writer; + private readonly Guid _token; + private readonly Guid? _previousToken; + + public TransactionGateLease(DbWriter? writer, Guid token, Guid? previousToken) + { + _writer = writer; + _token = token; + _previousToken = previousToken; + } + + public void Dispose() + { + _writer?.ExitTransactionGate(_token, _previousToken); + } + } + /// /// RAII wrapper for transactions and savepoints. /// Ensures _transactionDepth is decremented and uncommitted changes are rolled back on Dispose. @@ -200,7 +227,7 @@ public sealed class TransactionScope : IDisposable private readonly string? _savepointName; private readonly SqliteConnection? _conn; private readonly DbWriter _writer; - private readonly bool _ownsTransactionGate; + private readonly TransactionGateLease _transactionGateLease; private const int StateActive = 0; private const int StateCommitting = 1; private const int StateCommitted = 2; @@ -210,20 +237,30 @@ public sealed class TransactionScope : IDisposable private int _disposeStarted; // Real transaction / 実トランザクション - internal TransactionScope(SqliteTransaction transaction, DbWriter writer, bool ownsTransactionGate = false) + internal TransactionScope(SqliteTransaction transaction, DbWriter writer) + : this(transaction, writer, TransactionGateLease.None) + { + } + + internal TransactionScope(SqliteTransaction transaction, DbWriter writer, TransactionGateLease transactionGateLease = default) { _transaction = transaction; _writer = writer; - _ownsTransactionGate = ownsTransactionGate; + _transactionGateLease = transactionGateLease; } // Savepoint / セーブポイント - internal TransactionScope(string savepointName, SqliteConnection conn, DbWriter writer, bool ownsTransactionGate = false) + internal TransactionScope(string savepointName, SqliteConnection conn, DbWriter writer) + : this(savepointName, conn, writer, TransactionGateLease.None) + { + } + + internal TransactionScope(string savepointName, SqliteConnection conn, DbWriter writer, TransactionGateLease transactionGateLease = default) { _savepointName = savepointName; _conn = conn; _writer = writer; - _ownsTransactionGate = ownsTransactionGate; + _transactionGateLease = transactionGateLease; } public void Commit() @@ -381,8 +418,7 @@ public void Dispose() } finally { - if (_ownsTransactionGate) - _writer.ExitTransactionGate(); + _transactionGateLease.Dispose(); } } } @@ -1766,7 +1802,7 @@ public bool OptimizeFtsIfIncrementalWriteThresholdReached(int threshold = Defaul /// True when the bit was actually stamped; false when re-verification failed. public bool MarkFoldReady(bool stampCurrentSymbolExtractorVersions = false) { - var ownsTransactionGate = EnterTransactionGate(); + var gateLease = EnterTransactionGate(); try { bool ownTransaction = !IsInTransaction(); @@ -1828,8 +1864,7 @@ public bool MarkFoldReady(bool stampCurrentSymbolExtractorVersions = false) } finally { - if (ownsTransactionGate) - ExitTransactionGate(); + gateLease.Dispose(); } } @@ -3252,7 +3287,7 @@ private static int GetRowsPerInsertStatement(int columnCount) private void SetReadyBit(int flag) { - var ownsTransactionGate = EnterTransactionGate(); + var gateLease = EnterTransactionGate(); try { // The ready bits share a single PRAGMA user_version word, so two parallel @@ -3300,8 +3335,7 @@ private void SetReadyBit(int flag) } finally { - if (ownsTransactionGate) - ExitTransactionGate(); + gateLease.Dispose(); } } diff --git a/tests/CodeIndex.Tests/ConcurrencyTests.cs b/tests/CodeIndex.Tests/ConcurrencyTests.cs index ee9ab6238c..112399dfd2 100644 --- a/tests/CodeIndex.Tests/ConcurrencyTests.cs +++ b/tests/CodeIndex.Tests/ConcurrencyTests.cs @@ -567,7 +567,11 @@ public async Task SetReadyBit_SharedWriterWaitsForActiveScope() var writer = new DbWriter(_db.Connection); using var txn = writer.BeginTransaction(); - var markTask = Task.Run(() => writer.MarkGraphReady()); + Task markTask; + using (ExecutionContext.SuppressFlow()) + { + markTask = Task.Run(() => writer.MarkGraphReady()); + } await Task.Delay(100); Assert.False(markTask.IsCompleted); @@ -578,6 +582,31 @@ public async Task SetReadyBit_SharedWriterWaitsForActiveScope() Assert.True((_db.GetUserVersion() & DbContext.GraphReadyFlag) != 0); } + [Fact] + public async Task BeginTransaction_NestedScopeAfterAwaitKeepsLogicalOwnership() + { + var writer = new DbWriter(_db.Connection); + + using var outer = writer.BeginTransaction(); + await Task.Yield(); + using var nested = writer.BeginTransaction(); + + writer.UpsertFile(new FileRecord + { + Path = "src/after_await.cs", + Lang = "csharp", + Size = 10, + Lines = 1, + Modified = new DateTime(2025, 1, 1, 0, 0, 0, DateTimeKind.Utc), + Checksum = "after-await", + }); + nested.Commit(); + outer.Commit(); + + var reader = new DbReader(_db.Connection); + Assert.Equal(1, reader.GetStatus().Files); + } + private static List BuildReferenceBatch(long fileId, string label, int count) { var refs = new List(count); diff --git a/tests/CodeIndex.Tests/DatabaseTests.cs b/tests/CodeIndex.Tests/DatabaseTests.cs index 0cdb151b38..9728e83a27 100644 --- a/tests/CodeIndex.Tests/DatabaseTests.cs +++ b/tests/CodeIndex.Tests/DatabaseTests.cs @@ -210,7 +210,7 @@ public void TransactionScope_SavepointWithoutConnection_ThrowsExplicitInvalidOpe scopeType, System.Reflection.BindingFlags.Instance | System.Reflection.BindingFlags.NonPublic, binder: null, - args: ["sp_missing_conn", null!, _writer, false], + args: ["sp_missing_conn", null!, _writer], culture: null) ?? throw new InvalidOperationException("TransactionScope instance was not created."); From cf7522ab0d075b9964bc20f48dfd3ee3389b20d8 Mon Sep 17 00:00:00 2001 From: Widthdom Date: Mon, 25 May 2026 12:43:39 +0900 Subject: [PATCH 3/3] Tighten DbWriter transaction gate ownership (#1733 #2676) --- src/CodeIndex/Database/DbWriter.cs | 10 +++++++- tests/CodeIndex.Tests/ConcurrencyTests.cs | 30 ++++++++--------------- 2 files changed, 19 insertions(+), 21 deletions(-) diff --git a/src/CodeIndex/Database/DbWriter.cs b/src/CodeIndex/Database/DbWriter.cs index 6b1b85c447..84d21c0772 100644 --- a/src/CodeIndex/Database/DbWriter.cs +++ b/src/CodeIndex/Database/DbWriter.cs @@ -33,6 +33,7 @@ public class DbWriter private int _rowSkipSavepointCounter; private long _batchRowsSkipped; private int _transactionDepth; + private int _transactionOwnerThreadId; private Guid _transactionOwnerToken; // Outermost SqliteTransaction currently held open by this writer (null when no // transaction is active OR after the outermost transaction has been committed / @@ -160,7 +161,10 @@ private TransactionGateLease EnterTransactionGate() { lock (_transactionStateLock) { - if (_transactionDepth > 0 && _transactionOwnerToken != Guid.Empty && _currentTransactionGateToken.Value == _transactionOwnerToken) + if (_transactionDepth > 0 && + _transactionOwnerThreadId == Environment.CurrentManagedThreadId && + _transactionOwnerToken != Guid.Empty && + _currentTransactionGateToken.Value == _transactionOwnerToken) return TransactionGateLease.None; } @@ -171,6 +175,7 @@ private TransactionGateLease EnterTransactionGate() { var previousToken = _currentTransactionGateToken.Value; var token = Guid.NewGuid(); + _transactionOwnerThreadId = Environment.CurrentManagedThreadId; _transactionOwnerToken = token; _currentTransactionGateToken.Value = token; return new TransactionGateLease(this, token, previousToken); @@ -187,7 +192,10 @@ private void ExitTransactionGate(Guid token, Guid? previousToken) lock (_transactionStateLock) { if (_transactionOwnerToken == token) + { + _transactionOwnerThreadId = 0; _transactionOwnerToken = Guid.Empty; + } } if (_currentTransactionGateToken.Value == token) _currentTransactionGateToken.Value = previousToken; diff --git a/tests/CodeIndex.Tests/ConcurrencyTests.cs b/tests/CodeIndex.Tests/ConcurrencyTests.cs index 112399dfd2..e2863a93c9 100644 --- a/tests/CodeIndex.Tests/ConcurrencyTests.cs +++ b/tests/CodeIndex.Tests/ConcurrencyTests.cs @@ -567,11 +567,7 @@ public async Task SetReadyBit_SharedWriterWaitsForActiveScope() var writer = new DbWriter(_db.Connection); using var txn = writer.BeginTransaction(); - Task markTask; - using (ExecutionContext.SuppressFlow()) - { - markTask = Task.Run(() => writer.MarkGraphReady()); - } + var markTask = Task.Run(() => writer.MarkGraphReady()); await Task.Delay(100); Assert.False(markTask.IsCompleted); @@ -583,28 +579,22 @@ public async Task SetReadyBit_SharedWriterWaitsForActiveScope() } [Fact] - public async Task BeginTransaction_NestedScopeAfterAwaitKeepsLogicalOwnership() + public async Task BeginTransaction_TaskRunWithFlowedExecutionContextStillWaitsForActiveScope() { var writer = new DbWriter(_db.Connection); - using var outer = writer.BeginTransaction(); - await Task.Yield(); - using var nested = writer.BeginTransaction(); - writer.UpsertFile(new FileRecord + var nestedTask = Task.Run(() => { - Path = "src/after_await.cs", - Lang = "csharp", - Size = 10, - Lines = 1, - Modified = new DateTime(2025, 1, 1, 0, 0, 0, DateTimeKind.Utc), - Checksum = "after-await", + using var nested = writer.BeginTransaction(); + nested.Commit(); }); - nested.Commit(); - outer.Commit(); + await Task.Delay(100); + Assert.False(nestedTask.IsCompleted); - var reader = new DbReader(_db.Connection); - Assert.Equal(1, reader.GetStatus().Files); + outer.Commit(); + outer.Dispose(); + await nestedTask; } private static List BuildReferenceBatch(long fileId, string label, int count)