diff --git a/src/Chaptarr.Core.Test/Books/ProviderAliasRepositoryRetryFixture.cs b/src/Chaptarr.Core.Test/Books/ProviderAliasRepositoryRetryFixture.cs new file mode 100644 index 00000000..300a2d46 --- /dev/null +++ b/src/Chaptarr.Core.Test/Books/ProviderAliasRepositoryRetryFixture.cs @@ -0,0 +1,262 @@ +using System; +using System.Collections.Generic; +using System.IO; +using System.Linq; +using Dapper; +using Microsoft.Data.Sqlite; +using Npgsql; +using NUnit.Framework; +using NzbDrone.Common.Messaging; +using NzbDrone.Core.Books; +using NzbDrone.Core.Datastore; +using NzbDrone.Core.Messaging.Events; + +namespace Chaptarr.Core.Test.Books +{ + [TestFixture] + public class ProviderAliasRepositoryRetryFixture + { + private static PostgresException PgError(string sqlState) => + new PostgresException("boom", "ERROR", "ERROR", sqlState); + + [Test] + public void should_retry_after_a_postgres_unique_violation_and_then_succeed() + { + var calls = 0; + + ProviderAliasRepository.RetryOnUniqueViolation( + () => + { + calls++; + if (calls == 1) + { + throw PgError(PostgresErrorCodes.UniqueViolation); + } + }, + maxAttempts: 3); + + Assert.That(calls, Is.EqualTo(2)); + } + + [Test] + public void should_retry_after_a_sqlite_unique_violation() + { + var calls = 0; + + ProviderAliasRepository.RetryOnUniqueViolation( + () => + { + calls++; + if (calls == 1) + { + throw new SqliteException("SQLite Error 19: 'UNIQUE constraint failed: ProviderAliasIndex.EntityType'.", 19); + } + }, + maxAttempts: 3); + + Assert.That(calls, Is.EqualTo(2)); + } + + [Test] + public void should_find_the_violation_when_it_is_wrapped() + { + Assert.That( + ProviderAliasRepository.IsUniqueViolation(new InvalidOperationException("wrapper", PgError(PostgresErrorCodes.UniqueViolation))), + Is.True); + } + + [Test] + public void should_not_retry_other_database_errors() + { + var calls = 0; + + Assert.Throws(() => ProviderAliasRepository.RetryOnUniqueViolation( + () => + { + calls++; + throw PgError(PostgresErrorCodes.UndefinedTable); + }, + maxAttempts: 3)); + + Assert.That(calls, Is.EqualTo(1), "only unique violations are retried"); + } + + [Test] + public void should_not_retry_sqlite_constraint_errors_that_are_not_unique_violations() + { + Assert.That(ProviderAliasRepository.IsUniqueViolation(new SqliteException("SQLite Error 19: 'NOT NULL constraint failed: X.Y'.", 19)), Is.False); + } + + [Test] + public void should_give_up_and_rethrow_after_the_last_attempt() + { + var calls = 0; + + Assert.Throws(() => ProviderAliasRepository.RetryOnUniqueViolation( + () => + { + calls++; + throw PgError(PostgresErrorCodes.UniqueViolation); + }, + maxAttempts: 3)); + + Assert.That(calls, Is.EqualTo(3), "a persistent violation is surfaced, not swallowed"); + } + + private sealed class NoopEventAggregator : IEventAggregator + { + public void PublishEvent(TEvent @event) + where TEvent : class, IEvent + { + } + } + + private sealed class FailingFirstAttemptRepository : ProviderAliasRepository + { + private readonly Func _failure; + + public FailingFirstAttemptRepository(IMainDatabase database, Func failure) + : base(database, new NoopEventAggregator()) + { + _failure = failure; + } + + public int Attempts { get; private set; } + public List> ItemsPerAttempt { get; } = new List>(); + + internal override void ReplaceAliasesOnce(string entityType, int entityId, string scope, List items) + { + Attempts++; + ItemsPerAttempt.Add(items); + + var failure = _failure(Attempts); + if (failure != null) + { + throw failure; + } + } + } + + // The repository constructor probes Database.DatabaseType, which opens a connection; hand it an + // in-memory SQLite one. The subclass overrides ReplaceAliasesOnce, so no query ever runs on it. + private static MainDatabase UnusedDatabase() => new MainDatabase(new Database("main", () => + { + var conn = new SqliteConnection("Data Source=:memory:"); + conn.Open(); + return conn; + })); + + [OneTimeSetUp] + public void OneTimeSetUp() + { + if (TableMapping.Mapper.TableMap.Count == 0) + { + TableMapping.Map(); + } + } + + [Test] + public void replace_aliases_should_go_through_the_retry_and_pass_the_same_aliases_each_attempt() + { + var sut = new FailingFirstAttemptRepository( + UnusedDatabase(), + attempt => attempt == 1 ? PgError(PostgresErrorCodes.UniqueViolation) : null); + + var aliases = new List + { + new ProviderAlias { EntityType = "Author", EntityId = 7, Scope = "author", Provider = "hc", NormalizedProviderId = "123" } + }; + + Assert.DoesNotThrow(() => sut.ReplaceAliases("Author", 7, "author", aliases)); + + Assert.That(sut.Attempts, Is.EqualTo(2), "the wiring must retry a unique violation"); + Assert.That(sut.ItemsPerAttempt.Select(x => x.Count), Is.EqualTo(new[] { 1, 1 })); + } + + [Test] + public void replace_aliases_should_not_retry_other_errors() + { + var sut = new FailingFirstAttemptRepository(UnusedDatabase(), attempt => PgError(PostgresErrorCodes.UndefinedTable)); + + Assert.Throws(() => sut.ReplaceAliases("Author", 7, "author", new List())); + + Assert.That(sut.Attempts, Is.EqualTo(1)); + } + + [Test] + public void replace_aliases_should_replace_the_entitys_rows_in_a_real_database_and_leave_other_entities_alone() + { + var databasePath = Path.Combine(TestContext.CurrentContext.WorkDirectory, $"provider_alias_{Guid.NewGuid():N}.db"); + var connectionString = new SqliteConnectionStringBuilder + { + DataSource = databasePath, + Mode = SqliteOpenMode.ReadWriteCreate, + Pooling = false + }.ToString(); + + try + { + using (var setup = new SqliteConnection(connectionString)) + { + setup.Open(); + setup.Execute(@" + CREATE TABLE ""ProviderAliasIndex"" ( + ""Id"" INTEGER PRIMARY KEY, + ""EntityType"" TEXT NOT NULL, + ""EntityId"" INTEGER NOT NULL, + ""Scope"" TEXT NOT NULL, + ""Provider"" TEXT NOT NULL, + ""NormalizedProviderId"" TEXT NOT NULL, + ""CreatedAt"" TEXT NOT NULL, + ""UpdatedAt"" TEXT NOT NULL + ); + CREATE UNIQUE INDEX ""IX_ProviderAliasIndex_Unique"" ON ""ProviderAliasIndex"" (""EntityType"", ""EntityId"", ""Scope"", ""Provider"", ""NormalizedProviderId""); + INSERT INTO ""ProviderAliasIndex"" (""EntityType"", ""EntityId"", ""Scope"", ""Provider"", ""NormalizedProviderId"", ""CreatedAt"", ""UpdatedAt"") + VALUES ('Author', 7, 'author', 'hc', 'old', '2026-01-01', '2026-01-01'), + ('Author', 8, 'author', 'hc', 'other-author', '2026-01-01', '2026-01-01');"); + } + + var database = new MainDatabase(new Database("main", () => + { + var conn = new SqliteConnection(connectionString); + conn.Open(); + return conn; + })); + var sut = new ProviderAliasRepository(database, new NoopEventAggregator()); + var now = DateTime.UtcNow; + + var aliases = new List + { + new ProviderAlias { EntityType = "Author", EntityId = 7, Scope = "author", Provider = "hc", NormalizedProviderId = "123", CreatedAt = now, UpdatedAt = now }, + new ProviderAlias { EntityType = "Author", EntityId = 7, Scope = "author", Provider = "gr", NormalizedProviderId = "456", CreatedAt = now, UpdatedAt = now } + }; + + sut.ReplaceAliases("Author", 7, "author", aliases); + sut.ReplaceAliases("Author", 7, "author", aliases); // idempotent: replacing twice must not collide with itself + + using (var verify = new SqliteConnection(connectionString)) + { + verify.Open(); + var author7 = verify.Query(@"SELECT ""Provider"" || ':' || ""NormalizedProviderId"" FROM ""ProviderAliasIndex"" WHERE ""EntityId"" = 7 ORDER BY 1;").ToList(); + var author8 = verify.Query(@"SELECT ""NormalizedProviderId"" FROM ""ProviderAliasIndex"" WHERE ""EntityId"" = 8;").ToList(); + + Assert.That(author7, Is.EqualTo(new[] { "gr:456", "hc:123" }), "old alias replaced, new ones inserted once each"); + Assert.That(author8, Is.EqualTo(new[] { "other-author" }), "another entity's aliases are untouched"); + } + } + finally + { + try + { + if (File.Exists(databasePath)) + { + File.Delete(databasePath); + } + } + catch + { + } + } + } + } +} diff --git a/src/NzbDrone.Core/Books/Repositories/ProviderAliasRepository.cs b/src/NzbDrone.Core/Books/Repositories/ProviderAliasRepository.cs index 4763769a..55de4c38 100644 --- a/src/NzbDrone.Core/Books/Repositories/ProviderAliasRepository.cs +++ b/src/NzbDrone.Core/Books/Repositories/ProviderAliasRepository.cs @@ -2,7 +2,10 @@ using System.Collections.Generic; using System.Data; using System.Linq; +using System.Threading; using Dapper; +using Microsoft.Data.Sqlite; +using Npgsql; using NzbDrone.Core.Datastore; using NzbDrone.Core.Messaging.Events; @@ -22,7 +25,25 @@ public ProviderAliasRepository(IMainDatabase database, IEventAggregator eventAgg { } + private const int MaxReplaceAttempts = 3; + public void ReplaceAliases(string entityType, int entityId, string scope, IEnumerable aliases) + { + var items = aliases?.ToList() ?? new List(); + + // Replacing an entity's aliases is DELETE-then-INSERT in a READ COMMITTED transaction. Two writers + // replacing the SAME entity at once (a bulk author edit updates thousands of authors while async + // handlers of each update's event refresh the same aliases) both delete nothing, one inserts and + // commits, and the other's INSERT then violates IX_ProviderAliasIndex_Unique - which failed a whole + // author-editor save. Retrying lets the loser's DELETE see the winner's committed rows and replace + // them, so the last writer wins instead of the request failing. + RetryOnUniqueViolation( + () => ReplaceAliasesOnce(entityType, entityId, scope, items), + MaxReplaceAttempts); + } + + // internal virtual so tests can make an attempt fail and prove ReplaceAliases retries it. + internal virtual void ReplaceAliasesOnce(string entityType, int entityId, string scope, List items) { using (var conn = _database.OpenConnection()) using (var transaction = conn.BeginTransaction(IsolationLevel.ReadCommitted)) @@ -34,16 +55,64 @@ public void ReplaceAliases(string entityType, int entityId, string scope, IEnume new { entityType, entityId, scope }, transaction); - var items = aliases?.ToList() ?? new List(); if (items.Count > 0) { - InsertMany(items, conn, transaction); + // Fresh instances per attempt: a failed InsertMany must not leave ids on the caller's objects. + var toInsert = items.Select(item => new ProviderAlias + { + EntityType = item.EntityType, + EntityId = item.EntityId, + Scope = item.Scope, + Provider = item.Provider, + NormalizedProviderId = item.NormalizedProviderId, + CreatedAt = item.CreatedAt, + UpdatedAt = item.UpdatedAt + }).ToList(); + + InsertMany(toInsert, conn, transaction); } transaction.Commit(); } } + internal static void RetryOnUniqueViolation(Action action, int maxAttempts) + { + for (var attempt = 1; ; attempt++) + { + try + { + action(); + return; + } + catch (Exception ex) when (attempt < maxAttempts && IsUniqueViolation(ex)) + { + // Small growing pause so two colliding writers do not immediately collide again. + Thread.Sleep(attempt * 15); + } + } + } + + internal static bool IsUniqueViolation(Exception exception) + { + for (var ex = exception; ex != null; ex = ex.InnerException) + { + if (ex is PostgresException pg && pg.SqlState == PostgresErrorCodes.UniqueViolation) + { + return true; + } + + if (ex is SqliteException sqlite && + sqlite.SqliteErrorCode == 19 && + sqlite.Message.IndexOf("UNIQUE constraint failed", StringComparison.OrdinalIgnoreCase) >= 0) + { + return true; + } + } + + return false; + } + public void DeleteAliases(string entityType, int entityId) { Delete(a => a.EntityType == entityType && a.EntityId == entityId);