diff --git a/harmony.sln b/harmony.sln index a44dc0d..f9d4080 100644 --- a/harmony.sln +++ b/harmony.sln @@ -22,6 +22,8 @@ Project("{2150E333-8FDC-42A3-9474-1A3956D46DE8}") = "Solution Items", "Solution src\.editorconfig = src\.editorconfig EndProjectSection EndProject +Project("{FAE04EC0-301F-11D3-BF4B-00C04F79EFBC}") = "SIL.Harmony.Benchmarks", "src\SIL.Harmony.Benchmarks\SIL.Harmony.Benchmarks.csproj", "{E5C038C9-AB5B-456C-AFC9-037A8E399486}" +EndProject Global GlobalSection(SolutionConfigurationPlatforms) = preSolution Debug|Any CPU = Debug|Any CPU @@ -55,5 +57,9 @@ Global {8EB807F6-C548-4016-856E-2ACCA5603036}.Debug|Any CPU.Build.0 = Debug|Any CPU {8EB807F6-C548-4016-856E-2ACCA5603036}.Release|Any CPU.ActiveCfg = Release|Any CPU {8EB807F6-C548-4016-856E-2ACCA5603036}.Release|Any CPU.Build.0 = Release|Any CPU + {E5C038C9-AB5B-456C-AFC9-037A8E399486}.Debug|Any CPU.ActiveCfg = Debug|Any CPU + {E5C038C9-AB5B-456C-AFC9-037A8E399486}.Debug|Any CPU.Build.0 = Debug|Any CPU + {E5C038C9-AB5B-456C-AFC9-037A8E399486}.Release|Any CPU.ActiveCfg = Release|Any CPU + {E5C038C9-AB5B-456C-AFC9-037A8E399486}.Release|Any CPU.Build.0 = Release|Any CPU EndGlobalSection EndGlobal diff --git a/src/SIL.Harmony.Benchmarks/AddSnapshotsBenchmarks.cs b/src/SIL.Harmony.Benchmarks/AddSnapshotsBenchmarks.cs new file mode 100644 index 0000000..3a9d503 --- /dev/null +++ b/src/SIL.Harmony.Benchmarks/AddSnapshotsBenchmarks.cs @@ -0,0 +1,160 @@ +using System.Diagnostics.CodeAnalysis; +using BenchmarkDotNet.Attributes; +using Microsoft.EntityFrameworkCore; +using SIL.Harmony.Db; +using SIL.Harmony.Tests; + +namespace SIL.Harmony.Benchmarks; + +public enum AddSnapshotsWorkload +{ + /// Create many distinct entities (root snapshots, all inserts, no FindAsync hits). + CreateNew, + + /// Word + Definition per commit — two projected table types. + MultiTypeCreate, + + /// Word + Tag + WordTag per commit — three types plus a reference graph. + ReferencedCreate, + + /// Seed created words, then update each once — updates project onto existing rows (FindAsync per snapshot). + UpdateExisting, + + /// Modify one entity many times — large snapshot count (many intermediates), dedup, few projected rows. + ModifySameEntity, + + /// Create then delete each entity — delete group + insert group, two SaveChanges. + CreateThenDelete, + + /// Create, delete, then apply-on-deleted — mixed EntityIsDeleted projection skips. + CreateDeleteModify, +} + +// Isolates CrdtRepository.AddSnapshots (the slow, non-FAST path) from the rest of the sync pipeline. +// Expensive DB seeding happens once in GlobalSetup; each iteration gets a clean copy via ForkDatabase() and +// recomputes the snapshot batch so no EF-tracked state leaks across iterations. +// disable warning about waiting for sync code, benchmarkdotnet does not support async code, and it doesn't deadlock when waiting. +[SuppressMessage("Usage", "VSTHRD002:Avoid problematic synchronous waits")] +public class AddSnapshotsBenchmarks +{ + private DataModelTestBase _template = null!; + private HashSet _measuredCommitIds = null!; + + private DataModelTestBase _local = null!; + private CrdtRepository _repository = null!; + private ObjectSnapshot[] _snapshotsToAdd = null!; + + [Params( + 1000 + // , 10_000 + )] + public int ChangeCount { get; set; } + + [ParamsAllValues] + public AddSnapshotsWorkload Workload { get; set; } + + [GlobalSetup] + public void GlobalSetup() + { + _template = new DataModelTestBase(alwaysValidate: false, performanceTest: true); + var clientId = Guid.NewGuid(); + + List seed; + List measured; + switch (Workload) + { + case AddSnapshotsWorkload.CreateNew: + seed = []; + measured = BenchmarkWorkloadBuilders.BuildCreateWords(_template, clientId, ChangeCount); + break; + case AddSnapshotsWorkload.MultiTypeCreate: + seed = []; + measured = BenchmarkWorkloadBuilders.BuildWordsWithDefinitions(_template, clientId, ChangeCount); + break; + case AddSnapshotsWorkload.ReferencedCreate: + seed = []; + measured = BenchmarkWorkloadBuilders.BuildWordsWithTags(_template, clientId, ChangeCount); + break; + case AddSnapshotsWorkload.UpdateExisting: + (seed, measured) = BenchmarkWorkloadBuilders.BuildUpdateExisting(_template, clientId, ChangeCount); + break; + case AddSnapshotsWorkload.ModifySameEntity: + seed = []; + measured = BenchmarkWorkloadBuilders.BuildModifySameWord(_template, clientId, ChangeCount); + break; + case AddSnapshotsWorkload.CreateThenDelete: + seed = []; + measured = BenchmarkWorkloadBuilders.BuildCreateThenDelete(_template, clientId, ChangeCount); + break; + case AddSnapshotsWorkload.CreateDeleteModify: + seed = []; + measured = BenchmarkWorkloadBuilders.BuildCreateDeleteModify(_template, clientId, ChangeCount); + break; + default: + throw new ArgumentOutOfRangeException(); + } + + // Seed commits go through the full pipeline so their snapshots and projected rows already exist. + if (seed.Count > 0) + ((ISyncable)_template.DataModel).AddRangeFromSync(seed).Wait(); + + // The measured commits are present in the database but their snapshots are NOT yet persisted; that's the + // work AddSnapshots performs. Adding only the commits mirrors the state right before UpdateSnapshots runs. + var repository = _template.CreateRepository(); + repository.AddCommits(measured).GetAwaiter().GetResult(); + _measuredCommitIds = measured.Select(c => c.Id).ToHashSet(); + } + + [IterationSetup] + public void IterationSetup() + { + _local = _template.ForkDatabase(alwaysValidate: false); + _repository = _local.CreateRepository(); + + // Load the measured commits fresh from the fork (tracked) so AddSnapshots resolves their Commit navigation + // from the change tracker instead of trying to re-insert them. + var measuredCommits = _local.DbContext.Commits + .Include(c => c.ChangeEntities) + .Where(c => EF.Parameter(_measuredCommitIds).Contains(c.Id)) + .ToArray() + .ToSortedSet(); + + // Prepopulate the snapshot lookup the same way DataModel.UpdateSnapshots does: existing snapshots (untracked) + // plus null for entities without one, so SnapshotWorker doesn't issue a per-entity query while computing. + var entityIds = measuredCommits + .SelectMany(c => c.ChangeEntities.Select(ce => ce.EntityId)) + .ToHashSet(); + var snapshotLookup = _repository.CurrentSnapshots() + .Include(s => s.Commit) + .Where(s => EF.Parameter(entityIds).Contains(s.EntityId)) + .ToDictionary(s => s.EntityId, s => (ObjectSnapshot?)s); + foreach (var entityId in entityIds) + snapshotLookup.TryAdd(entityId, null); + + var worker = new SnapshotWorker(snapshotLookup, _repository, _local.CrdtConfig); + _snapshotsToAdd = worker.ComputeSnapshotsToPersist(measuredCommits).GetAwaiter().GetResult().ToArray(); + } + + [Benchmark] + public void AddSnapshots() + { + _repository.AddSnapshots(_snapshotsToAdd).Wait(); + } + + [IterationCleanup] + public void IterationCleanup() + { + _repository.DisposeAsync().AsTask().Wait(); + _local.DisposeAsync().AsTask().Wait(); + _repository = null!; + _local = null!; + _snapshotsToAdd = null!; + } + + [GlobalCleanup] + public void GlobalCleanup() + { + _template.DisposeAsync().AsTask().Wait(); + _template = null!; + } +} diff --git a/src/SIL.Harmony.Benchmarks/BenchmarkWorkloadBuilders.cs b/src/SIL.Harmony.Benchmarks/BenchmarkWorkloadBuilders.cs new file mode 100644 index 0000000..15a4475 --- /dev/null +++ b/src/SIL.Harmony.Benchmarks/BenchmarkWorkloadBuilders.cs @@ -0,0 +1,168 @@ +using SIL.Harmony.Changes; +using SIL.Harmony.Tests; + +namespace SIL.Harmony.Benchmarks; + +/// +/// Shared commit-building helpers for the benchmark suites so (full sync +/// pipeline) and (isolated snapshot persist) exercise the same domain scenarios. +/// Builders are pure: they only use the source for its change factories and +/// counter and never touch the database. +/// +public static class BenchmarkWorkloadBuilders +{ + public static Commit NewCommit(Guid clientId, DateTimeOffset dateTime, params IChange[] changes) + { + var commit = new Commit(Guid.NewGuid()) + { + ClientId = clientId, + HybridDateTime = new HybridDateTime(dateTime, 0) + }; + for (var i = 0; i < changes.Length; i++) + { + commit.ChangeEntities.Add(DataModel.ToChangeEntity(changes[i], i, commit.Id)); + } + return commit; + } + + /// Create many distinct words (one SetWord per commit). + public static List BuildCreateWords(DataModelTestBase src, Guid clientId, int count) + { + var commits = new List(count); + for (var i = 0; i < count; i++) + { + commits.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(Guid.NewGuid(), $"entity {i}"))); + } + return commits; + } + + /// Create words each with a new definition in the same commit. + public static List BuildWordsWithDefinitions(DataModelTestBase src, Guid clientId, int count) + { + var commits = new List(count); + for (var i = 0; i < count; i++) + { + var wordId = Guid.NewGuid(); + commits.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(wordId, $"entity {i}"), + src.NewDefinition(wordId, $"definition {i}", "noun"))); + } + return commits; + } + + /// Create words each with a tag and WordTag link in the same commit. + public static List BuildWordsWithTags(DataModelTestBase src, Guid clientId, int count) + { + var commits = new List(count); + for (var i = 0; i < count; i++) + { + var wordId = Guid.NewGuid(); + var tagId = Guid.NewGuid(); + commits.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(wordId, $"entity {i}"), + src.SetTag(tagId, $"tag {i}"), + src.TagWord(wordId, tagId))); + } + return commits; + } + + /// Create one word, then modify that same word repeatedly. + public static List BuildModifySameWord(DataModelTestBase src, Guid clientId, int count) + { + var wordId = Guid.NewGuid(); + var commits = new List(count); + commits.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(wordId, "entity 0"))); + for (var i = 1; i < count; i++) + { + commits.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(wordId, $"entity {i}"))); + } + return commits; + } + + /// Create words then delete them. + public static List BuildCreateThenDelete(DataModelTestBase src, Guid clientId, int count) + { + var commits = new List(count * 2); + for (var i = 0; i < count; i++) + { + var wordId = Guid.NewGuid(); + commits.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(wordId, $"entity {i}"))); + commits.Add(NewCommit(clientId, src.NextDate(), + src.DeleteWord(wordId))); + } + return commits; + } + + /// Create, delete, then modify (apply-after-delete) for each word. + public static List BuildCreateDeleteModify(DataModelTestBase src, Guid clientId, int count) + { + var commits = new List(count * 3); + for (var i = 0; i < count; i++) + { + var wordId = Guid.NewGuid(); + commits.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(wordId, $"entity {i}"))); + commits.Add(NewCommit(clientId, src.NextDate(), + src.DeleteWord(wordId))); + // SetWordNote supports apply-on-existing (including deleted) without undeleting + commits.Add(NewCommit(clientId, src.NextDate(), + src.SetWordNote(wordId, $"note {i}"))); + } + return commits; + } + + /// + /// Local already has create + late modify; sync inserts a mid-history modify for each word + /// (forces stale snapshot deletion / rebuild). Returns the seed commits and the mid-history commits to sync. + /// + public static (List seed, List toSync) BuildOutOfOrderInsert(DataModelTestBase src, Guid clientId, int count) + { + var seed = new List(count * 2); + var toSync = new List(count); + for (var i = 0; i < count; i++) + { + var wordId = Guid.NewGuid(); + var createTime = src.NextDate(); + var midTime = src.NextDate(); + var lateTime = src.NextDate(); + + seed.Add(NewCommit(clientId, createTime, + src.SetWord(wordId, $"entity {i}"))); + // Mid-history change is what gets synced after local already has create + late + toSync.Add(NewCommit(clientId, midTime, + src.SetWordNote(wordId, $"note {i}"))); + seed.Add(NewCommit(clientId, lateTime, + src.SetWord(wordId, $"entity {i} late"))); + } + + return (seed, toSync); + } + + /// + /// Seed a set of created words, then modify each one exactly once. The updates are the measured batch; + /// their snapshots must update the already-projected rows (FindAsync per snapshot in the slow path). + /// + public static (List seed, List measured) BuildUpdateExisting(DataModelTestBase src, Guid clientId, int count) + { + var wordIds = new Guid[count]; + var seed = new List(count); + for (var i = 0; i < count; i++) + { + wordIds[i] = Guid.NewGuid(); + seed.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(wordIds[i], $"entity {i}"))); + } + + var measured = new List(count); + for (var i = 0; i < count; i++) + { + measured.Add(NewCommit(clientId, src.NextDate(), + src.SetWord(wordIds[i], $"entity {i} updated"))); + } + return (seed, measured); + } +} diff --git a/src/SIL.Harmony.Benchmarks/DataModelSyncBenchmarks.cs b/src/SIL.Harmony.Benchmarks/DataModelSyncBenchmarks.cs new file mode 100644 index 0000000..3830b72 --- /dev/null +++ b/src/SIL.Harmony.Benchmarks/DataModelSyncBenchmarks.cs @@ -0,0 +1,135 @@ +using System.Diagnostics.CodeAnalysis; +using BenchmarkDotNet.Attributes; +using BenchmarkDotNet.Engines; +using SIL.Harmony.Changes; +using SIL.Harmony.Tests; + +namespace SIL.Harmony.Benchmarks; + +public enum SyncWorkload +{ + /// Create many distinct words (one SetWord per commit). + CreateWords, + + /// Create words each with a new definition in the same commit. + WordsWithDefinitions, + + /// Create words each with a tag and WordTag link in the same commit. + WordsWithTags, + + /// Create one word, then modify that same word repeatedly. + ModifySameWord, + + /// Create words then delete them. + CreateThenDelete, + + /// Create, delete, then modify (apply-after-delete) for each word. + CreateDeleteModify, + + /// + /// Local already has create + late modify; sync inserts a mid-history modify + /// for each word (forces stale snapshot deletion / rebuild). + /// + OutOfOrderInsert, +} + +// [SimpleJob(RunStrategy.Monitoring)] +[MemoryDiagnoser] +[SuppressMessage("Usage", "VSTHRD002:Avoid problematic synchronous waits")] +public class DataModelSyncBenchmarks +{ + private DataModelTestBase remote = null!; + private DataModelTestBase local = null!; + + [Params( + 1000 + // , 10_000 + )] + public int ChangeCount { get; set; } + + [ParamsAllValues] + public SyncWorkload Workload { get; set; } + + private Commit[] _commits = null!; + private HashSet? _syncCommitIds; + + [GlobalSetup] + public void GlobalSetup() + { + _syncCommitIds = null; + remote = new DataModelTestBase(alwaysValidate: false, performanceTest: true); + var clientId = Guid.NewGuid(); + List commits; + switch (Workload) + { + case SyncWorkload.CreateWords: + commits = BenchmarkWorkloadBuilders.BuildCreateWords(remote, clientId, ChangeCount); + break; + case SyncWorkload.WordsWithDefinitions: + commits = BenchmarkWorkloadBuilders.BuildWordsWithDefinitions(remote, clientId, ChangeCount); + break; + case SyncWorkload.WordsWithTags: + commits = BenchmarkWorkloadBuilders.BuildWordsWithTags(remote, clientId, ChangeCount); + break; + case SyncWorkload.ModifySameWord: + commits = BenchmarkWorkloadBuilders.BuildModifySameWord(remote, clientId, ChangeCount); + break; + case SyncWorkload.CreateThenDelete: + commits = BenchmarkWorkloadBuilders.BuildCreateThenDelete(remote, clientId, ChangeCount); + break; + case SyncWorkload.CreateDeleteModify: + commits = BenchmarkWorkloadBuilders.BuildCreateDeleteModify(remote, clientId, ChangeCount); + break; + case SyncWorkload.OutOfOrderInsert: + var (seed, toSync) = BenchmarkWorkloadBuilders.BuildOutOfOrderInsert(remote, clientId, ChangeCount); + _syncCommitIds = toSync.Select(c => c.Id).ToHashSet(); + commits = [.. seed, .. toSync]; + break; + default: + throw new ArgumentOutOfRangeException(); + } + ((ISyncable)remote.DataModel).AddRangeFromSync(commits).Wait(); + } + + [IterationSetup] + public void IterationSetup() + { + local = new DataModelTestBase(alwaysValidate: false, performanceTest: true); + _ = local.WriteNextChange(local.SetWord(Guid.NewGuid(), "entity1")).Result; + //cant share commits between iterations, because EF modifies them + var allCommits = remote.DataModel.GetChanges(new SyncState([])).Result.MissingFromClient; + + if (_syncCommitIds is null) + { + _commits = allCommits; + return; + } + + var seed = new List(allCommits.Length); + var toSync = new List(_syncCommitIds.Count); + foreach (var commit in allCommits) + { + if (_syncCommitIds.Contains(commit.Id)) + toSync.Add(commit); + else + seed.Add(commit); + } + + ((ISyncable)local.DataModel).AddRangeFromSync(seed).Wait(); + _commits = toSync.ToArray(); + } + + [IterationCleanup] + public void IterationCleanup() + { + local.DisposeAsync().AsTask().Wait(); + local = null!; + } + + [Benchmark] + public void SyncCommits() + { + ((ISyncable)local.DataModel).AddRangeFromSync(_commits) + .Wait(); + } +} diff --git a/src/SIL.Harmony.Benchmarks/Program.cs b/src/SIL.Harmony.Benchmarks/Program.cs new file mode 100644 index 0000000..23e438d --- /dev/null +++ b/src/SIL.Harmony.Benchmarks/Program.cs @@ -0,0 +1,12 @@ +using BenchmarkDotNet.Configs; +using BenchmarkDotNet.Engines; +using BenchmarkDotNet.Jobs; +using BenchmarkDotNet.Running; +using SIL.Harmony.Benchmarks; + + +var config = DefaultConfig.Instance + .AddJob(Job.MediumRun.WithStrategy(RunStrategy.Monitoring)); +BenchmarkSwitcher + .FromTypes([typeof(DataModelSyncBenchmarks), typeof(AddSnapshotsBenchmarks)]) + .Run(args, config); diff --git a/src/SIL.Harmony.Benchmarks/SIL.Harmony.Benchmarks.csproj b/src/SIL.Harmony.Benchmarks/SIL.Harmony.Benchmarks.csproj new file mode 100644 index 0000000..90275bd --- /dev/null +++ b/src/SIL.Harmony.Benchmarks/SIL.Harmony.Benchmarks.csproj @@ -0,0 +1,26 @@ + + + Exe + + + AnyCPU + pdbonly + true + true + true + Release + false + + + + + + + + + + + + + + \ No newline at end of file diff --git a/src/SIL.Harmony.Tests/DataModelTestBase.cs b/src/SIL.Harmony.Tests/DataModelTestBase.cs index 8dcf44f..2f69660 100644 --- a/src/SIL.Harmony.Tests/DataModelTestBase.cs +++ b/src/SIL.Harmony.Tests/DataModelTestBase.cs @@ -2,6 +2,7 @@ using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.DependencyInjection.Extensions; +using Microsoft.Extensions.Options; using SIL.Harmony.Changes; using SIL.Harmony.Config; using SIL.Harmony.Db; @@ -65,6 +66,15 @@ public void SetCurrentDate(DateTime dateTime) currentDate = dateTime; } + /// + /// Creates a repository over this instance's DbContext. Exposed so benchmarks can drive + /// methods (e.g. AddSnapshots) directly without going through the sync pipeline. + /// + internal CrdtRepository CreateRepository() => + _services.GetRequiredService().CreateRepositorySync(); + + internal HarmonyConfig CrdtConfig => _services.GetRequiredService>().Value; + private static int _instanceCount = 0; private DateTimeOffset currentDate = new(new DateTime(2000, 1, 1, 0, 0, 0).AddHours(_instanceCount++)); public DateTimeOffset NextDate() => currentDate = currentDate.AddDays(1); @@ -132,6 +142,11 @@ public IChange SetWord(Guid entityId, string value) return new SetWordTextChange(entityId, value); } + public IChange SetWordNote(Guid entityId, string note) + { + return new SetWordNoteChange(entityId, note); + } + public IChange DeleteWord(Guid entityId) { return new DeleteChange(entityId); diff --git a/src/SIL.Harmony.Tests/SIL.Harmony.Tests.csproj b/src/SIL.Harmony.Tests/SIL.Harmony.Tests.csproj index f542263..24f3074 100644 --- a/src/SIL.Harmony.Tests/SIL.Harmony.Tests.csproj +++ b/src/SIL.Harmony.Tests/SIL.Harmony.Tests.csproj @@ -41,4 +41,8 @@ + + + + diff --git a/src/SIL.Harmony/Config/HarmonyConfig.cs b/src/SIL.Harmony/Config/HarmonyConfig.cs index 65556aa..9bd4bf5 100644 --- a/src/SIL.Harmony/Config/HarmonyConfig.cs +++ b/src/SIL.Harmony/Config/HarmonyConfig.cs @@ -1,3 +1,4 @@ +using System.Collections.Concurrent; using System.Text.Json; using System.Text.Json.Serialization.Metadata; using Microsoft.EntityFrameworkCore; @@ -30,6 +31,13 @@ public class HarmonyConfig private readonly Lazy _lazyJsonSerializerOptions; private readonly Lazy _lazyChangeDiscriminatorMaps; + /// + /// Cache of derived projected-table SQL metadata (keyed by projected CLR type), used by + /// . Stored on the config so it's shared across repositories and + /// db contexts and only built once per type. + /// + internal ConcurrentDictionary ProjectedTableInfoCache { get; } = new(); + public HarmonyConfig() { _lazyChangeDiscriminatorMaps = new Lazy(BuildChangeDiscriminatorMaps); diff --git a/src/SIL.Harmony/CrdtKernel.cs b/src/SIL.Harmony/CrdtKernel.cs index e1e6302..fef511a 100644 --- a/src/SIL.Harmony/CrdtKernel.cs +++ b/src/SIL.Harmony/CrdtKernel.cs @@ -52,6 +52,7 @@ public static IServiceCollection AddCrdtDataCore(this IServiceCollection service services.AddSingleton(sp => sp.GetRequiredService>().Value.JsonSerializerOptions); services.AddSingleton(TimeProvider.System); services.AddScoped(NewTimeProvider); + services.AddSingleton(); services.AddScoped(); //must use factory method because DataModel constructor is internal services.AddScoped(provider => new DataModel( diff --git a/src/SIL.Harmony/DataModel.cs b/src/SIL.Harmony/DataModel.cs index 01e5219..8cc2a25 100644 --- a/src/SIL.Harmony/DataModel.cs +++ b/src/SIL.Harmony/DataModel.cs @@ -124,7 +124,7 @@ public ValueTask DisposeAsync() return ValueTask.CompletedTask; } - private static ChangeEntity ToChangeEntity(IChange change, int index, Guid commitId) + internal static ChangeEntity ToChangeEntity(IChange change, int index, Guid commitId) { return new ChangeEntity() { @@ -192,19 +192,25 @@ private async Task UpdateSnapshots(CrdtRepository repo, SortedSet commit if (commitsToApply.Count == 0) return; var oldestAddedCommit = commitsToApply.First(); await repo.DeleteStaleSnapshots(oldestAddedCommit); - Dictionary snapshotLookup = []; + Dictionary snapshotLookup = []; if (commitsToApply.Count > 10) { // Bulk-load relevant snapshots to minimize DB queries var entityIds = commitsToApply .SelectMany(c => c.ChangeEntities.Select(ce => ce.EntityId)) - .Distinct(); + .ToHashSet(); //EF.Parameter forces a single JSON parameter; without it EF 10+ emits one parameter per id and overflows SQLite's parameter limit snapshotLookup = await repo.CurrentSnapshots() + .Include(s => s.Commit) .Where(s => EF.Parameter(entityIds).Contains(s.EntityId)) - .Select(s => new KeyValuePair(s.EntityId, s.Id)) - .ToDictionaryAsync(s => s.Key, s => s.Value); + .ToDictionaryAsync(s => s.EntityId, s => (ObjectSnapshot?)s); + entityIds.ExceptWith(snapshotLookup.Keys); + foreach (Guid entityId in entityIds) + { + //snapshot does not exist, store null to tell SnapshotWorker NOT to attempt to fetch it from the database + snapshotLookup[entityId] = null; + } } var snapshotWorker = new SnapshotWorker(snapshotLookup, repo, _crdtConfig.Value); diff --git a/src/SIL.Harmony/Db/CrdtDbContextFactory.cs b/src/SIL.Harmony/Db/CrdtDbContextFactory.cs index 890b1ee..d50a1b1 100644 --- a/src/SIL.Harmony/Db/CrdtDbContextFactory.cs +++ b/src/SIL.Harmony/Db/CrdtDbContextFactory.cs @@ -1,6 +1,7 @@ using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.ChangeTracking; using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Metadata; namespace SIL.Harmony.Db; @@ -65,6 +66,7 @@ public DbSet Set() where TEntity : class return context.Set(); } + public IModel Model => context.Model; public DatabaseFacade Database => context.Database; public ChangeTracker ChangeTracker => context.ChangeTracker; diff --git a/src/SIL.Harmony/Db/CrdtRepository.cs b/src/SIL.Harmony/Db/CrdtRepository.cs index d7a8a10..b0e2822 100644 --- a/src/SIL.Harmony/Db/CrdtRepository.cs +++ b/src/SIL.Harmony/Db/CrdtRepository.cs @@ -1,8 +1,10 @@ using System.Collections.Concurrent; +using System.Collections.Frozen; using System.Reflection; using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.ChangeTracking; using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Metadata; using Microsoft.EntityFrameworkCore.Storage; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging; @@ -52,14 +54,17 @@ internal class CrdtRepository : IDisposable, IAsyncDisposable private readonly ICrdtDbContext _dbContext; private readonly IOptions _crdtConfig; private readonly ILogger _logger; + private readonly FastProjection _fastProjection; public CrdtRepository(ICrdtDbContext dbContext, IOptions crdtConfig, ILogger logger, + FastProjection fastProjection, Commit? ignoreChangesAfter = null) { _crdtConfig = crdtConfig; _dbContext = ignoreChangesAfter is not null ? new ScopedDbContext(dbContext, ignoreChangesAfter) : dbContext; _logger = logger; + _fastProjection = fastProjection; //we can't use the scoped db context is it prevents access to the DbSet for the Snapshots, //but since we're using a custom query, we can use it directly and apply the scoped filters manually _currentSnapshotsQueryable = MakeCurrentSnapshotsQuery(dbContext, ignoreChangesAfter); @@ -300,85 +305,15 @@ public async Task> GetChanges(SyncState remoteState) return await _dbContext.Commits.GetChanges(remoteState); } - public async Task AddSnapshots(IEnumerable snapshots) + public Task AddSnapshots(IEnumerable snapshots) { - var latestProjectByEntityId = new Dictionary(); - foreach (var grouping in snapshots.GroupBy(s => s.EntityIsDeleted).OrderByDescending(g => g.Key))//execute deletes first - { - foreach (var snapshot in grouping.DefaultOrderDescending()) - { - _dbContext.Add(snapshot); - if (latestProjectByEntityId.TryGetValue(snapshot.EntityId, out var latestProjected)) - { - // there might be a deleted and un-deleted snapshot for the same entity in the same batch - // in that case there's only a 50% chance that they're in the right order, so we need to explicitly only project the latest one - if (snapshot.Commit.CompareKey.CompareTo(latestProjected) < 0) - { - continue; - } - } - latestProjectByEntityId[snapshot.EntityId] = snapshot.Commit.CompareKey; - - await ProjectSnapshot(snapshot); - } - - try - { - await _dbContext.SaveChangesAsync(); - } - catch (DbUpdateException e) - { - var entries = string.Join(Environment.NewLine, e.Entries.Select(entry => entry.ToString())); - var message = $"Error saving snapshots (deleted: {grouping.Key}): {e.Message}{Environment.NewLine}{entries}"; - _logger.LogError(e, message); - throw new DbUpdateException(message, e); - } - } - } - - private async ValueTask ProjectSnapshot(ObjectSnapshot objectSnapshot) - { - if (!_crdtConfig.Value.EnableProjectedTables) return; - - //need to check if an entry exists already, even if this is the root commit it may have already been added to the db - var existingEntry = await GetEntityEntry(objectSnapshot.Entity.DbObject.GetType(), objectSnapshot.EntityId); - if (existingEntry is null && objectSnapshot.EntityIsDeleted) return; - - if (existingEntry is null) // add - { - // this is a new entity even though it might not be a root snapshot, because we only project the latest snapshot of each entity per sync - - //if we don't make a copy first then the entity will be tracked by the context and be modified - //by future changes in the same session - var entity = objectSnapshot.Entity.Copy().DbObject; - - var newEntry = _dbContext.Entry(entity); - // only mark this single entry as added, rather than the whole graph (this matches the update behaviour below) - newEntry.State = EntityState.Added; - newEntry.Property(ObjectSnapshot.ShadowRefName).CurrentValue = objectSnapshot.Id; - } - else if (objectSnapshot.EntityIsDeleted) // delete - { - _dbContext.Remove(existingEntry.Entity); - } - else // update - { - var entity = objectSnapshot.Entity.DbObject; - existingEntry.CurrentValues.SetValues(entity); - existingEntry.Property(ObjectSnapshot.ShadowRefName).CurrentValue = objectSnapshot.Id; - } - } - - private async ValueTask GetEntityEntry(Type entityType, Guid entityId) - { - if (!_crdtConfig.Value.EnableProjectedTables) return null; - var entity = await _dbContext.FindAsync(entityType, entityId); - return entity is not null ? _dbContext.Entry(entity) : null; + var snapshotList = snapshots as IReadOnlyCollection ?? snapshots.ToArray(); + return _fastProjection.AddSnapshotsRawAsync(_dbContext, snapshotList); } public CrdtRepository GetScopedRepository(Commit excludeChangesAfterCommit) { - return new CrdtRepository(_dbContext, _crdtConfig, _logger, excludeChangesAfterCommit); + return new CrdtRepository(_dbContext, _crdtConfig, _logger, _fastProjection, excludeChangesAfterCommit); } /// @@ -506,6 +441,7 @@ public DbSet Set() where TEntity : class throw new NotSupportedException("can not support Set when using scoped db context"); } + public IModel Model => inner.Model; public DatabaseFacade Database => inner.Database; public ChangeTracker ChangeTracker => inner.ChangeTracker; diff --git a/src/SIL.Harmony/Db/FastProjection.cs b/src/SIL.Harmony/Db/FastProjection.cs new file mode 100644 index 0000000..c9d2ad0 --- /dev/null +++ b/src/SIL.Harmony/Db/FastProjection.cs @@ -0,0 +1,254 @@ +using System.Data.Common; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Metadata; +using Microsoft.EntityFrameworkCore.Storage; +using Microsoft.Extensions.Options; +using SIL.Harmony.Config; + +namespace SIL.Harmony.Db; + +/// +/// Projects snapshots into the projected tables. Snapshots are still inserted through EF (unchanged), +/// but the projected tables are populated with hand-written raw SQL +/// `INSERT ... ON CONFLICT(pk) DO UPDATE` (one upsert command per entity row) instead of going +/// through EF's change tracker. Everything the SQL needs (table/column names, primary key, the +/// SnapshotId shadow FK, value converters) is derived from the EF model, so no per-entity code is +/// required. +/// +internal class FastProjection +{ + private readonly HarmonyConfig _crdtConfig; + + public FastProjection(IOptions crdtConfig) + { + _crdtConfig = crdtConfig.Value; + } + + public async Task AddSnapshotsRawAsync( + ICrdtDbContext dbContext, + IReadOnlyCollection snapshots) + { + // AddSnapshots is normally called inside a caller-managed transaction; reuse it so the raw + // SQL runs on the same connection/transaction as the snapshot insert. When called outside a + // transaction (e.g. direct repository tests) we open and commit our own. + var ownTransaction = dbContext.Database.CurrentTransaction is null + ? await dbContext.Database.BeginTransactionAsync() + : null; + try + { + // 1. persist the snapshot rows exactly as before + dbContext.AddRange(snapshots); + await dbContext.SaveChangesAsync(); + + if (_crdtConfig.EnableProjectedTables) + { + await ProjectAsync(dbContext, snapshots); + } + + if (ownTransaction is not null) await ownTransaction.CommitAsync(); + } + finally + { + if (ownTransaction is not null) await ownTransaction.DisposeAsync(); + } + } + + private async Task ProjectAsync( + ICrdtDbContext dbContext, + IReadOnlyCollection snapshots) + { + // 2. dedup to the latest snapshot per entity (a batch can contain several snapshots for the + // same entity - intermediate + latest - and possibly a deleted and undeleted one). + var latest = new Dictionary(); + foreach (var snapshot in snapshots) + { + if (latest.TryGetValue(snapshot.EntityId, out var existing) && + existing.Commit.CompareKey.CompareTo(snapshot.Commit.CompareKey) >= 0) + { + continue; + } + latest[snapshot.EntityId] = snapshot; + } + + var connection = dbContext.Database.GetDbConnection(); + var transaction = dbContext.Database.CurrentTransaction!.GetDbTransaction(); + var sqlHelper = dbContext.Database.GetService(); + + // 3. group by projected CLR type and order the types so FK parents are written first + var byType = latest.Values + .GroupBy(s => s.Entity.DbObject.GetType()) + .ToDictionary(g => g.Key, g => g.ToList()); + var orderedTypes = OrderTypesByDependency(dbContext.Model, byType.Keys); + + // deletes first (like the slow path) so a unique value can be freed and re-inserted within the + // same batch; children before parents (reverse dependency order) to satisfy FK constraints. + for (var i = orderedTypes.Count - 1; i >= 0; i--) + { + var deleted = byType[orderedTypes[i]].Where(s => s.EntityIsDeleted).ToList(); + if (deleted.Count == 0) continue; + var info = GetTableInfo(dbContext, orderedTypes[i], sqlHelper); + await DeletePerQueryAsync(connection, transaction, info, deleted); + } + + // upserts: parents first + foreach (var type in orderedTypes) + { + var live = byType[type].Where(s => !s.EntityIsDeleted).ToList(); + if (live.Count == 0) continue; + var info = GetTableInfo(dbContext, type, sqlHelper); + await UpsertPerQueryAsync(connection, transaction, info, live); + } + } + + private static async Task UpsertPerQueryAsync( + DbConnection connection, DbTransaction transaction, ProjectedTableInfo info, List rows) + { + await using var command = connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = info.InsertSql; + var parameters = new DbParameter[info.Columns.Count]; + for (var i = 0; i < info.Columns.Count; i++) + { + var p = command.CreateParameter(); + p.ParameterName = "@p" + i; + command.Parameters.Add(p); + parameters[i] = p; + } + command.Prepare(); + + foreach (var snapshot in rows) + { + var dbObject = snapshot.Entity.DbObject; + for (var i = 0; i < info.Columns.Count; i++) + { + parameters[i].Value = ToParameterValue(GetProviderValue(info.Columns[i], snapshot, dbObject)); + } + await command.ExecuteNonQueryAsync(); + } + } + + private static async Task DeletePerQueryAsync( + DbConnection connection, DbTransaction transaction, ProjectedTableInfo info, List rows) + { + await using var command = connection.CreateCommand(); + command.Transaction = transaction; + command.CommandText = info.DeleteSql; + var p = command.CreateParameter(); + p.ParameterName = "@p0"; + command.Parameters.Add(p); + command.Prepare(); + + foreach (var snapshot in rows) + { + p.Value = ToParameterValue(GetProviderValue(info.PrimaryKey, snapshot, snapshot.Entity.DbObject)); + await command.ExecuteNonQueryAsync(); + } + } + + private static object? GetProviderValue(ColumnInfo column, ObjectSnapshot snapshot, object dbObject) + { + if (column.IsShadowSnapshotId) return snapshot.Id; + var raw = column.Property.PropertyInfo is { } pi + ? pi.GetValue(dbObject) + : column.Property.FieldInfo?.GetValue(dbObject); + var converter = column.Property.GetValueConverter(); + return converter is null ? raw : converter.ConvertToProvider(raw); + } + + private static object ToParameterValue(object? value) => value ?? DBNull.Value; + + // ---- model metadata --------------------------------------------------------------------- + + private ProjectedTableInfo GetTableInfo(ICrdtDbContext dbContext, Type clrType, ISqlGenerationHelper sqlHelper) + { + return _crdtConfig.ProjectedTableInfoCache.GetOrAdd(clrType, t => BuildTableInfo(dbContext, t, sqlHelper)); + } + + private static ProjectedTableInfo BuildTableInfo(ICrdtDbContext dbContext, Type clrType, ISqlGenerationHelper sqlHelper) + { + var entityType = dbContext.Model.FindEntityType(clrType) + ?? throw new InvalidOperationException($"No EF entity type found for projected type {clrType.Name}"); + var tableName = entityType.GetTableName() + ?? throw new InvalidOperationException($"No table name found for projected type {clrType.Name}"); + var schema = entityType.GetSchema(); + var storeObject = StoreObjectIdentifier.Table(tableName, schema); + var pkPropertyNames = (entityType.FindPrimaryKey()?.Properties + ?? throw new InvalidOperationException($"No primary key found for projected type {clrType.Name}")) + .Select(p => p.Name) + .ToHashSet(); + + var columns = new List(); + foreach (var property in entityType.GetProperties()) + { + var columnName = property.GetColumnName(storeObject); + if (columnName is null) continue; // not mapped to this table + columns.Add(new ColumnInfo( + sqlHelper.DelimitIdentifier(columnName), + property, + property.Name == ObjectSnapshot.ShadowRefName, + pkPropertyNames.Contains(property.Name))); + } + + var pkColumns = columns.Where(c => c.IsPrimaryKey).ToList(); + if (pkColumns.Count != 1) + throw new NotSupportedException($"Fast projection requires a single-column primary key for {clrType.Name}"); + var pk = pkColumns[0]; + + var delimitedTable = sqlHelper.DelimitIdentifier(tableName, schema); + var columnList = string.Join(",", columns.Select(c => c.DelimitedName)); + var parameterList = string.Join(",", columns.Select((_, i) => "@p" + i)); + var setClause = string.Join(",", columns.Where(c => !c.IsPrimaryKey) + .Select(c => $"{c.DelimitedName}=excluded.{c.DelimitedName}")); + var onConflict = string.IsNullOrEmpty(setClause) + ? $"ON CONFLICT ({pk.DelimitedName}) DO NOTHING" + : $"ON CONFLICT ({pk.DelimitedName}) DO UPDATE SET {setClause}"; + + return new ProjectedTableInfo( + columns, + pk, + InsertSql: $"INSERT INTO {delimitedTable} ({columnList}) VALUES ({parameterList}) {onConflict};", + DeleteSql: $"DELETE FROM {delimitedTable} WHERE {pk.DelimitedName}=@p0;"); + } + + /// + /// Post-order DFS over FK edges so principal (parent) types come before dependents. Self + /// references and FKs to non-projected types (e.g. the SnapshotId FK to Snapshots) are ignored. + /// + private static List OrderTypesByDependency(IModel model, IEnumerable types) + { + var typeSet = types.ToHashSet(); + var ordered = new List(typeSet.Count); + var visited = new HashSet(); + + void Visit(Type type) + { + if (!visited.Add(type)) return; + var entityType = model.FindEntityType(type); + if (entityType is not null) + { + foreach (var fk in entityType.GetForeignKeys()) + { + var principal = fk.PrincipalEntityType.ClrType; + if (principal != type && typeSet.Contains(principal)) Visit(principal); + } + } + ordered.Add(type); + } + + foreach (var type in typeSet) Visit(type); + return ordered; + } + + internal sealed record ColumnInfo( + string DelimitedName, + IProperty Property, + bool IsShadowSnapshotId, + bool IsPrimaryKey); + + internal sealed record ProjectedTableInfo( + IReadOnlyList Columns, + ColumnInfo PrimaryKey, + string InsertSql, + string DeleteSql); +} diff --git a/src/SIL.Harmony/Db/ICrdtDbContext.cs b/src/SIL.Harmony/Db/ICrdtDbContext.cs index 438a225..59f52cc 100644 --- a/src/SIL.Harmony/Db/ICrdtDbContext.cs +++ b/src/SIL.Harmony/Db/ICrdtDbContext.cs @@ -1,6 +1,7 @@ using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.ChangeTracking; using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Metadata; namespace SIL.Harmony.Db; @@ -11,6 +12,7 @@ public interface ICrdtDbContext : IDisposable, IAsyncDisposable Task SaveChangesAsync(CancellationToken cancellationToken = default); ValueTask FindAsync(Type entityType, params object?[]? keyValues); DbSet Set() where TEntity : class; + IModel Model { get; } DatabaseFacade Database { get; } ChangeTracker ChangeTracker { get; } EntityEntry Entry(TEntity entity) where TEntity : class; diff --git a/src/SIL.Harmony/SIL.Harmony.csproj b/src/SIL.Harmony/SIL.Harmony.csproj index 094fe40..4d202a5 100644 --- a/src/SIL.Harmony/SIL.Harmony.csproj +++ b/src/SIL.Harmony/SIL.Harmony.csproj @@ -8,6 +8,7 @@ + diff --git a/src/SIL.Harmony/SnapshotWorker.cs b/src/SIL.Harmony/SnapshotWorker.cs index 37a8c1e..3eaf243 100644 --- a/src/SIL.Harmony/SnapshotWorker.cs +++ b/src/SIL.Harmony/SnapshotWorker.cs @@ -11,7 +11,7 @@ namespace SIL.Harmony; /// internal class SnapshotWorker { - private readonly Dictionary _snapshotLookup; + private readonly Dictionary _snapshotCache; private readonly CrdtRepository _crdtRepository; private readonly HarmonyConfig _crdtConfig; private readonly Dictionary _pendingSnapshots = []; @@ -19,13 +19,13 @@ internal class SnapshotWorker private readonly List _newIntermediateSnapshots = []; private SnapshotWorker(Dictionary snapshots, - Dictionary snapshotLookup, + Dictionary snapshotCache, CrdtRepository crdtRepository, HarmonyConfig crdtConfig) { _pendingSnapshots = snapshots; _crdtRepository = crdtRepository; - _snapshotLookup = snapshotLookup; + _snapshotCache = snapshotCache; _crdtConfig = crdtConfig; } @@ -41,12 +41,12 @@ internal static async Task> ApplyCommitsToSnaps return snapshots; } - /// a dictionary of entity id to latest snapshot id + /// a dictionary of entity id to latest snapshot id /// /// - internal SnapshotWorker(Dictionary snapshotLookup, + internal SnapshotWorker(Dictionary snapshotCache, CrdtRepository crdtRepository, - HarmonyConfig crdtConfig) : this([], snapshotLookup, crdtRepository, crdtConfig) + HarmonyConfig crdtConfig) : this([], snapshotCache, crdtRepository, crdtConfig) { } @@ -60,6 +60,17 @@ await _crdtRepository.AddSnapshots([ ]); } + /// + /// Applies the commits to snapshots the same way does, but returns the full list + /// of snapshots that would be persisted instead of writing them. Used by benchmarks to isolate the + /// step from commit application. + /// + internal async Task> ComputeSnapshotsToPersist(SortedSet commits) + { + await ApplyCommitChanges(commits); + return [.. _rootSnapshots.Values, .. _newIntermediateSnapshots, .. _pendingSnapshots.Values]; + } + private async ValueTask ApplyCommitChanges(SortedSet commits) { var intermediateSnapshots = new Dictionary(); @@ -159,14 +170,13 @@ private async ValueTask MarkDeleted(Guid deletedEntityId, ChangeContext context) return rootSnapshot; } - if (_snapshotLookup.TryGetValue(entityId, out var snapshotId)) + if (_snapshotCache.TryGetValue(entityId, out snapshot)) { - if (snapshotId is null) return null; - return await _crdtRepository.FindSnapshot(snapshotId.Value, true); + return snapshot; } snapshot = await _crdtRepository.GetCurrentSnapshotByObjectId(entityId, true); - _snapshotLookup[entityId] = snapshot?.Id; + _snapshotCache[entityId] = snapshot; return snapshot; }