using System; using System.Collections.Generic; using Microsoft.Data.SqlClient; using Newtonsoft.Json; using SyncEngine.Configuration; namespace SyncEngine.Metadata { public sealed class SyncMetadataRepository { private readonly SyncOptions _options; public SyncMetadataRepository(SyncOptions options) { _options = options; } public void EnsureSchema(SqlConnection connection) { var sql = @" IF SCHEMA_ID('sync_metadata') IS NULL EXEC('CREATE SCHEMA sync_metadata'); IF OBJECT_ID('sync_metadata.runs', 'U') IS NULL CREATE TABLE sync_metadata.runs ( run_id UNIQUEIDENTIFIER NOT NULL PRIMARY KEY, leg NVARCHAR(32) NOT NULL, started_at DATETIME2 NOT NULL, finished_at DATETIME2 NULL, status NVARCHAR(16) NOT NULL, tables_processed INT NOT NULL DEFAULT 0, tables_excluded INT NOT NULL DEFAULT 0, rows_inserted INT NOT NULL DEFAULT 0, rows_updated INT NOT NULL DEFAULT 0, rows_soft_deleted INT NOT NULL DEFAULT 0, rows_skipped INT NOT NULL DEFAULT 0, errors NVARCHAR(MAX) NULL ); IF OBJECT_ID('sync_metadata.checkpoints', 'U') IS NULL CREATE TABLE sync_metadata.checkpoints ( run_id UNIQUEIDENTIFIER NOT NULL, table_name NVARCHAR(256) NOT NULL, last_pk NVARCHAR(256) NULL, rows_inserted INT NOT NULL DEFAULT 0, rows_updated INT NOT NULL DEFAULT 0, rows_soft_deleted INT NOT NULL DEFAULT 0, updated_at DATETIME2 NOT NULL, PRIMARY KEY (run_id, table_name) ); IF OBJECT_ID('sync_metadata.schema_versions', 'U') IS NULL CREATE TABLE sync_metadata.schema_versions ( object_key NVARCHAR(512) NOT NULL PRIMARY KEY, snapshot NVARCHAR(MAX) NOT NULL, updated_at DATETIME2 NOT NULL );"; using (var cmd = new SqlCommand(sql, connection)) { cmd.CommandTimeout = _options.CommandTimeoutSeconds; cmd.ExecuteNonQuery(); } } public Guid StartRun(SqlConnection connection, string leg, int tablesExcluded) { var runId = Guid.NewGuid(); var sql = @" INSERT INTO sync_metadata.runs (run_id, leg, started_at, status, tables_excluded) VALUES (@id, @leg, @started, 'running', @excluded)"; using (var cmd = new SqlCommand(sql, connection)) { cmd.Parameters.AddWithValue("@id", runId); cmd.Parameters.AddWithValue("@leg", leg); cmd.Parameters.AddWithValue("@started", DateTime.UtcNow); cmd.Parameters.AddWithValue("@excluded", tablesExcluded); cmd.ExecuteNonQuery(); } return runId; } public void FinishRun(SqlConnection connection, Guid runId, string status, SyncRunStats stats, IList errors) { var sql = @" UPDATE sync_metadata.runs SET finished_at = @finished, status = @status, tables_processed = @tables, rows_inserted = @ins, rows_updated = @upd, rows_soft_deleted = @del, rows_skipped = @skip, errors = @errors WHERE run_id = @id"; using (var cmd = new SqlCommand(sql, connection)) { cmd.Parameters.AddWithValue("@id", runId); cmd.Parameters.AddWithValue("@finished", DateTime.UtcNow); cmd.Parameters.AddWithValue("@status", status); cmd.Parameters.AddWithValue("@tables", stats.TablesProcessed); cmd.Parameters.AddWithValue("@ins", stats.RowsInserted); cmd.Parameters.AddWithValue("@upd", stats.RowsUpdated); cmd.Parameters.AddWithValue("@del", stats.RowsSoftDeleted); cmd.Parameters.AddWithValue("@skip", stats.TablesSkipped); cmd.Parameters.AddWithValue("@errors", errors == null || errors.Count == 0 ? (object)DBNull.Value : JsonConvert.SerializeObject(errors)); cmd.ExecuteNonQuery(); } } public void SaveCheckpoint(SqlConnection connection, Guid runId, string tableName, object lastPk, int inserted, int updated, int softDeleted) { var sql = @" MERGE sync_metadata.checkpoints AS t USING (SELECT @run AS run_id, @table AS table_name) AS s ON t.run_id = s.run_id AND t.table_name = s.table_name WHEN MATCHED THEN UPDATE SET last_pk = @pk, rows_inserted = @ins, rows_updated = @upd, rows_soft_deleted = @del, updated_at = @now WHEN NOT MATCHED THEN INSERT (run_id, table_name, last_pk, rows_inserted, rows_updated, rows_soft_deleted, updated_at) VALUES (@run, @table, @pk, @ins, @upd, @del, @now);"; using (var cmd = new SqlCommand(sql, connection)) { cmd.Parameters.AddWithValue("@run", runId); cmd.Parameters.AddWithValue("@table", tableName); cmd.Parameters.AddWithValue("@pk", lastPk == null ? (object)DBNull.Value : Convert.ToString(lastPk)); cmd.Parameters.AddWithValue("@ins", inserted); cmd.Parameters.AddWithValue("@upd", updated); cmd.Parameters.AddWithValue("@del", softDeleted); cmd.Parameters.AddWithValue("@now", DateTime.UtcNow); cmd.ExecuteNonQuery(); } } public object GetCheckpoint(SqlConnection connection, Guid runId, string tableName) { var sql = "SELECT last_pk FROM sync_metadata.checkpoints WHERE run_id = @run AND table_name = @table"; using (var cmd = new SqlCommand(sql, connection)) { cmd.Parameters.AddWithValue("@run", runId); cmd.Parameters.AddWithValue("@table", tableName); var result = cmd.ExecuteScalar(); if (result == null || result == DBNull.Value) { return null; } return result; } } public void SaveSchemaVersion(SqlConnection connection, string objectKey, string snapshot) { var sql = @" MERGE sync_metadata.schema_versions AS t USING (SELECT @key AS object_key) AS s ON t.object_key = s.object_key WHEN MATCHED THEN UPDATE SET snapshot = @snap, updated_at = @now WHEN NOT MATCHED THEN INSERT (object_key, snapshot, updated_at) VALUES (@key, @snap, @now);"; using (var cmd = new SqlCommand(sql, connection)) { cmd.Parameters.AddWithValue("@key", objectKey); cmd.Parameters.AddWithValue("@snap", snapshot); cmd.Parameters.AddWithValue("@now", DateTime.UtcNow); cmd.ExecuteNonQuery(); } } } public sealed class SyncRunStats { public int TablesProcessed { get; set; } public int TablesSkipped { get; set; } public int RowsInserted { get; set; } public int RowsUpdated { get; set; } public int RowsSoftDeleted { get; set; } } }