using System.Collections.Concurrent; using System.Data; using System.Globalization; using System.Security.Cryptography; using System.Text; using System.Text.Json; using Microsoft.Data.SqlClient; using Microsoft.Extensions.Options; var builder = WebApplication.CreateBuilder(args); builder.WebHost.ConfigureKestrel(options => options.Limits.MaxRequestBodySize = 1_048_576); builder.Services.AddOptions() .Bind(builder.Configuration.GetSection("XiaoshuSync")) .Validate(options => !string.IsNullOrWhiteSpace(options.KeyId), "XiaoshuSync KeyId 未配置") .Validate(options => options.Secret.Length >= 32, "XiaoshuSync Secret 必须至少 32 位") .Validate(options => options.MaxClockSkewSeconds is >= 30 and <= 900, "签名时间窗口必须在 30 至 900 秒之间") .Validate(options => options.DefaultPageSize > 0 && options.MaxPageSize >= options.DefaultPageSize && options.MaxPageSize <= 5000, "同步分页配置无效") .Validate(options => options.MaxCommandBodyBytes is >= 1024 and <= 1_048_576, "写入请求大小上限无效") .Validate(options => options.PathBase.StartsWith('/') && !options.PathBase.EndsWith('/'), "PathBase 必须以 / 开头且不能以 / 结尾") .ValidateOnStart(); builder.Services.AddSingleton(); builder.Services.AddSingleton(); var app = builder.Build(); app.UsePathBase(app.Services.GetRequiredService>().Value.PathBase); app.UseExceptionHandler(exceptionApp => exceptionApp.Run(async context => { context.Response.StatusCode = StatusCodes.Status500InternalServerError; context.Response.ContentType = "application/json"; await context.Response.WriteAsJsonAsync(new { error = "sync_bridge_error", message = "同步桥执行失败,请查看服务器日志" }); })); app.Use(async (context, next) => { if (!context.Request.Path.StartsWithSegments("/v1") || context.Request.Path.Equals("/v1/health")) { await next(); return; } context.Response.Headers.CacheControl = "no-store"; var configured = context.RequestServices.GetRequiredService>().Value; if (context.Request.Path.StartsWithSegments("/v1/commands") && context.Request.ContentLength > configured.MaxCommandBodyBytes) { context.Response.StatusCode = StatusCodes.Status413PayloadTooLarge; await context.Response.WriteAsJsonAsync(new { error = "payload_too_large" }); return; } var rejection = await RequestAuthenticator.ValidateAsync(context, configured, context.RequestServices.GetRequiredService()); if (rejection is not null) { context.Response.StatusCode = rejection.Value.Status; await context.Response.WriteAsJsonAsync(new { error = rejection.Value.Message }); return; } await next(); }); app.MapGet("/v1/health", () => Results.Ok(new { service = "xiaoshu-legacy-sync", status = "ok", version = "v1", time = DateTimeOffset.UtcNow })); app.MapGet("/v1/manifest", async (LegacyRepository repository, CancellationToken cancellationToken) => Results.Ok(await repository.ManifestAsync(cancellationToken))); app.MapGet("/v1/changes", async (string dataset, string? cursor, int? limit, LegacyRepository repository, CancellationToken cancellationToken) => { if (!DatasetCatalog.All.TryGetValue(dataset, out var definition)) return Results.BadRequest(new { error = "unsupported_dataset", message = $"不支持的数据集:{dataset}" }); return Results.Ok(await repository.ChangesAsync(definition, cursor, limit, cancellationToken)); }); app.MapPost("/v1/commands/{operation}", async (string operation, HttpContext context, LegacyRepository repository, IOptions configured, CancellationToken cancellationToken) => { if (!CommandCatalog.Allowed.Contains(operation)) return Results.BadRequest(new { error = "unsupported_operation", message = $"不支持的写入操作:{operation}" }); if (context.Request.ContentLength is > 0 && context.Request.ContentLength > configured.Value.MaxCommandBodyBytes) return Results.StatusCode(StatusCodes.Status413PayloadTooLarge); var command = await context.Request.ReadFromJsonAsync(cancellationToken: cancellationToken); if (command is null || string.IsNullOrWhiteSpace(command.IdempotencyKey) || command.IdempotencyKey.Length > 128 || string.IsNullOrWhiteSpace(command.Reason) || command.Reason.Length > 500 || string.IsNullOrWhiteSpace(command.Actor) || command.Actor.Length > 200 || command.Payload.ValueKind is not JsonValueKind.Object) return Results.BadRequest(new { error = "invalid_command", message = "actor、reason、idempotencyKey 和对象 payload 必填且长度必须有效" }); return Results.Ok(await repository.ExecuteCommandAsync(operation, command, cancellationToken)); }); app.MapGet("/health", () => Results.Ok(new { service = "xiaoshu-legacy-sync", status = "ok", time = DateTimeOffset.UtcNow })); app.Run(); sealed class SyncOptions { public string KeyId { get; init; } = ""; public string Secret { get; init; } = ""; public string[] AllowedIps { get; init; } = []; public int MaxClockSkewSeconds { get; init; } = 300; public int DefaultPageSize { get; init; } = 500; public int MaxPageSize { get; init; } = 1000; public int MaxCommandBodyBytes { get; init; } = 1_048_576; public string LegacyTimeZoneId { get; init; } = "China Standard Time"; public string PathBase { get; init; } = "/xiaoshu-sync"; } readonly record struct AuthRejection(int Status, string Message); static class RequestAuthenticator { public static async Task ValidateAsync(HttpContext context, SyncOptions options, ReplayGuard replayGuard) { if (string.IsNullOrWhiteSpace(options.KeyId) || string.IsNullOrWhiteSpace(options.Secret)) return new(503, "同步桥密钥尚未配置"); var remoteIp = context.Connection.RemoteIpAddress?.ToString() ?? ""; if (options.AllowedIps.Length > 0 && !options.AllowedIps.Contains(remoteIp, StringComparer.OrdinalIgnoreCase)) return new(403, "来源 IP 不在白名单"); var keyId = context.Request.Headers["X-Xiaoshu-Key-Id"].ToString(); var timestampText = context.Request.Headers["X-Xiaoshu-Timestamp"].ToString(); var nonce = context.Request.Headers["X-Xiaoshu-Nonce"].ToString(); var signature = context.Request.Headers["X-Xiaoshu-Signature"].ToString(); if (!long.TryParse(timestampText, out var timestamp) || string.IsNullOrWhiteSpace(nonce) || string.IsNullOrWhiteSpace(signature) || !FixedText(keyId, options.KeyId)) return new(401, "签名头缺失或无效"); if (Math.Abs(DateTimeOffset.UtcNow.ToUnixTimeSeconds() - timestamp) > options.MaxClockSkewSeconds) return new(401, "请求时间戳已过期"); if (!replayGuard.TryUse(nonce, timestamp, options.MaxClockSkewSeconds)) return new(409, "nonce 已使用,拒绝重放"); context.Request.EnableBuffering(); using var reader = new StreamReader(context.Request.Body, Encoding.UTF8, leaveOpen: true); var body = await reader.ReadToEndAsync(); context.Request.Body.Position = 0; var bodyHash = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(body))).ToLowerInvariant(); var canonical = $"{context.Request.Method.ToUpperInvariant()}\n{context.Request.PathBase}{context.Request.Path}{context.Request.QueryString}\n{timestampText}\n{nonce}\n{bodyHash}"; using var hmac = new HMACSHA256(Encoding.UTF8.GetBytes(options.Secret)); var expected = Convert.ToHexString(hmac.ComputeHash(Encoding.UTF8.GetBytes(canonical))).ToLowerInvariant(); return FixedText(signature.ToLowerInvariant(), expected) ? null : new(401, "请求签名无效"); } private static bool FixedText(string left, string right) => left.Length == right.Length && CryptographicOperations.FixedTimeEquals(Encoding.UTF8.GetBytes(left), Encoding.UTF8.GetBytes(right)); } sealed class ReplayGuard { private readonly ConcurrentDictionary _nonces = new(); public bool TryUse(string nonce, long timestamp, int lifetimeSeconds) { var floor = DateTimeOffset.UtcNow.ToUnixTimeSeconds() - lifetimeSeconds; foreach (var stale in _nonces.Where(item => item.Value < floor).Select(item => item.Key)) _nonces.TryRemove(stale, out _); return nonce.Length is >= 16 and <= 128 && _nonces.TryAdd(nonce, timestamp); } } sealed record LegacyCommand(string Actor, string Reason, string IdempotencyKey, string? ExpectedUpdatedAt, JsonElement Payload); sealed record ProjectionRecord(string ClassName, string KeyField, string LegacyKey, Dictionary Fields); sealed record ChangeItem(string Key, string Operation, DateTimeOffset? ChangedAt, IReadOnlyList Records); sealed record ChangePage(string Dataset, string Cursor, bool HasMore, int Count, DateTimeOffset ServerTime, string Phase, long Watermark, IReadOnlyList Items); sealed record DatasetDefinition(string Key, string Label, string CountSql, string KeysSql, string SnapshotSql, string LookupSql, string LegacyKeyColumn, IReadOnlyList Records); sealed record RecordMapping(string ClassName, string KeyField, string Prefix, string SourcePrefix); sealed record CursorState(bool Incremental, long Position, long SnapshotWatermark); sealed record ChangeLogRow(long Id, string LegacyKey, string Operation, DateTimeOffset? ChangedAt, string? ItemKey); static class DatasetCatalog { private const string ActiveCommon = "COALESCE(c.Status,99)<>-2"; public static readonly IReadOnlyDictionary All = new Dictionary(StringComparer.OrdinalIgnoreCase) { ["users"] = new("users", "全部业务账号", "SELECT COUNT_BIG(*) FROM ZL_User", "SELECT CAST(UserID AS bigint) legacyKey FROM ZL_User ORDER BY UserID", """SELECT TOP (@limit) CAST(u.UserID AS bigint) legacyKey,COALESCE(u.LastLoginTimes,u.RegTime) changedAt,u.UserID [user__legacyUserId],u.UserName [user__username],u.HoneyName [user__nickname],u.Email [user__email],u.Mobile [user__mobile],u.GroupID [user__legacyGroupId],u.ParentUserID [user__parentUserId],u.RegTime [user__registeredAt],u.LastLoginTimes [user__lastLoginAt],u.LoginTimes [user__loginCount],u.State [user__legacyState],u.Purse [user__purse],u.SilverCoin [user__silverCoin],u.UserExp [user__userExp],u.UserPoint [user__userPoint],u.DummyPurse [user__dummyPurse],u.UserCreit [user__credit] FROM ZL_User u WHERE u.UserID>@after ORDER BY u.UserID""", """SELECT TOP (1) CAST(u.UserID AS bigint) legacyKey,COALESCE(u.LastLoginTimes,u.RegTime) changedAt,u.UserID [user__legacyUserId],u.UserName [user__username],u.HoneyName [user__nickname],u.Email [user__email],u.Mobile [user__mobile],u.GroupID [user__legacyGroupId],u.ParentUserID [user__parentUserId],u.RegTime [user__registeredAt],u.LastLoginTimes [user__lastLoginAt],u.LoginTimes [user__loginCount],u.State [user__legacyState],u.Purse [user__purse],u.SilverCoin [user__silverCoin],u.UserExp [user__userExp],u.UserPoint [user__userPoint],u.DummyPurse [user__dummyPurse],u.UserCreit [user__credit] FROM ZL_User u WHERE u.UserID=@legacyKey""", "legacyKey", [new("_User", "legacyUserId", "user__", "legacy-sync:user:")]), ["nodes"] = new("nodes", "课程与词库目录", "SELECT COUNT_BIG(*) FROM ZL_Node", "SELECT CAST(NodeID AS bigint) legacyKey FROM ZL_Node ORDER BY NodeID", """SELECT TOP (@limit) CAST(n.NodeID AS bigint) legacyKey,NULL changedAt,n.NodeID [node__nodeId],n.ParentID [node__parentId],n.NodeName [node__nodeName],n.NodeType [node__nodeType],n.NodeDir [node__nodeDir],n.NodePic [node__nodePicUrl],n.Description [node__description],n.OrderID [node__orderId],n.CUser [node__cuser],n.CUName [node__cuname],n.EditDate [node__editDate] FROM ZL_Node n WHERE n.NodeID>@after ORDER BY n.NodeID""", """SELECT TOP (1) CAST(n.NodeID AS bigint) legacyKey,NULL changedAt,n.NodeID [node__nodeId],n.ParentID [node__parentId],n.NodeName [node__nodeName],n.NodeType [node__nodeType],n.NodeDir [node__nodeDir],n.NodePic [node__nodePicUrl],n.Description [node__description],n.OrderID [node__orderId],n.CUser [node__cuser],n.CUName [node__cuname],n.EditDate [node__editDate] FROM ZL_Node n WHERE n.NodeID=@legacyKey""", "legacyKey", [new("Node", "nodeId", "node__", "legacy-sync:node:")]), ["course-bindings"] = Content("course-bindings", "课程绑定", 58, "ZL_C_kcbd", "b", ["yhid","kcid","yxx","cksl","syjd"]), ["appointments"] = Content("appointments", "预约排课", 54, "ZL_C_order", "a", ["szyh","pl","fxpl","plxm","kcid","yykcid","dslx","jffs","bxrq","sdsd","yysj","dszt","kssj","jssj","scsj","szmd"]), ["lessons"] = Content("lessons", "上课记录", 59, "ZL_C_skjl", "l", ["xymz","jsmz","kcid","kclx","kcmc","pf","pjnr","pldp","kzsj","yyds","plpjsj","szmdid"]), ["learning-records"] = Content("learning-records", "每日学习记录", 56, "ZL_C_ss", "d", ["userId","pl","dqrq","dsid","learned","ygg","djq","xxqs","fxrl","szmdid"]), ["practice-records"] = Content("practice-records", "练习记录", 53, "ZL_C_lxjl", "p", ["yhid","kcid","scid","xxcs","jrscb"]), ["memory-records"] = Content("memory-records", "抗遗忘记录", 60, "ZL_C_gywjl", "m", ["yhid","plid","kcid","xxjlid","fxzt","wcsj","kywrq","kywsj","OrderID"]), ["assessments"] = Content("assessments", "测评档案", 61, "ZL_C_cpda", "e", ["userId","df","askid","wrong","answerid","dontKnow","prev_score:prevScore","totalScore"]), ["vocabulary"] = Content("vocabulary", "词库", 52, "ZL_C_ck", "v", ["sy","yb","lj","yp"]), ["special-training"] = Content("special-training", "专项训练", 2, "ZL_C_Article", "r", ["ico","mp4","author","source","content","synopsis"], " AND c.NodeID=3") }; private static DatasetDefinition Content(string key, string label, int modelId, string addonTable, string alias, string[] addonFields, string extraWhere = "") { var addonSelect = string.Join(',', addonFields.Select(field => { var mapping = field.Split(':', 2); var sourceField = mapping[0]; var targetField = mapping.Length == 2 ? mapping[1] : ToCamel(sourceField); return $"{alias}.[{sourceField}] [addon__{targetField}]"; })); var select = $"CAST(c.GeneralID AS bigint) legacyKey,COALESCE(c.UpDateTime,c.CreateTime) changedAt,c.GeneralID [common__generalId],c.ItemID [common__itemId],c.ModelID [common__modelId],c.NodeID [common__nodeId],c.TableName [common__tableName],c.Title [common__title],c.Subtitle [common__subtitle],c.Inputer [common__inputer],c.CreateTime [common__createTime],c.UpDateTime [common__upDateTime],c.Status [common__status],c.OrderID [common__orderId],c.TopImg [common__topImg],{alias}.ID [addon__id],{addonSelect}"; var from = $" FROM ZL_CommonModel c INNER JOIN {addonTable} {alias} ON {alias}.ID=c.ItemID WHERE c.ModelID={modelId}{extraWhere} AND {ActiveCommon}"; var snapshotSql = $"SELECT TOP (@limit) {select}{from} AND c.GeneralID>@after ORDER BY c.GeneralID"; var lookupSql = $"SELECT TOP (1) {select}{from} AND c.GeneralID=@legacyKey"; var countSql = $"SELECT COUNT_BIG(*){from}"; var keysSql = $"SELECT CAST(c.GeneralID AS bigint) legacyKey{from} ORDER BY c.GeneralID"; return new(key, label, countSql, keysSql, snapshotSql, lookupSql, "legacyKey", [new("CommonModel", "generalId", "common__", $"legacy-sync:{key}:common:"), new(ClassFor(modelId), "id", "addon__", $"legacy-sync:{key}:addon:")]); } private static string ClassFor(int modelId) => modelId switch { 2 => "ContentArticle", 52 => "VocabularyWord", 53 => "PracticeRecord", 54 => "CourseAppointment", 56 => "DailyStudyRecord", 58 => "CourseBinding", 59 => "LessonRecord", 60 => "MemoryPracticeRecord", 61 => "AssessmentProfile", _ => throw new InvalidOperationException() }; private static string ToCamel(string value) => value.Length == 0 ? value : char.ToLowerInvariant(value[0]) + value[1..]; } static class CommandCatalog { public static readonly HashSet Allowed = new(StringComparer.OrdinalIgnoreCase) { "member.create", "member.update-profile", "member.reset-password", "member.change-status", "member.change-group", "member.adjust-balance", "member.bind-agent", "coach.create", "coach.update", "coach.delete", "appointment.create", "appointment.update", "appointment.cancel", "appointment.restore", "appointment.complete", "appointment.transition", "course-binding.save", "course-binding.recycle", "course-binding.restore", "vocabulary.save", "vocabulary.move", "vocabulary.recycle", "vocabulary.restore", "vocabulary.import", "special-training.save", "special-training.publish", "special-training.unpublish", "special-training.recycle", "special-training.restore" }; } sealed class LegacyRepository(IConfiguration configuration, IOptions options) { private readonly string _connectionString = configuration.GetConnectionString("LegacySqlServer") ?? throw new InvalidOperationException("LegacySqlServer connection string is missing"); private readonly SyncOptions _options = options.Value; public async Task ManifestAsync(CancellationToken cancellationToken) { await using var connection = new SqlConnection(_connectionString); await connection.OpenAsync(cancellationToken); var watermarkStart = await GlobalWatermarkAsync(connection, cancellationToken); var datasets = new List(); foreach (var definition in DatasetCatalog.All.Values) { await using var command = new SqlCommand(definition.CountSql, connection) { CommandTimeout = 120 }; var count = Convert.ToInt64(await command.ExecuteScalarAsync(cancellationToken), CultureInfo.InvariantCulture); var checksum = await KeyChecksumAsync(connection, definition, cancellationToken); await using var watermarkCommand = new SqlCommand("SELECT COALESCE(MAX(Id),0) FROM dbo.XiaoshuSyncChangeLog WHERE Dataset=@dataset", connection); watermarkCommand.Parameters.AddWithValue("@dataset", definition.Key); var watermark = Convert.ToInt64(await watermarkCommand.ExecuteScalarAsync(cancellationToken), CultureInfo.InvariantCulture); datasets.Add(new { key = definition.Key, label = definition.Label, count, watermark, checksum }); } var watermarkEnd = await GlobalWatermarkAsync(connection, cancellationToken); var now = DateTimeOffset.UtcNow; return new { manifestTime = now, serverTime = now, consistent = watermarkStart == watermarkEnd, watermarkStart, watermarkEnd, datasets }; } private static async Task GlobalWatermarkAsync(SqlConnection connection, CancellationToken cancellationToken) { await using var command = new SqlCommand("SELECT COALESCE(MAX(Id),0) FROM dbo.XiaoshuSyncChangeLog", connection); return Convert.ToInt64(await command.ExecuteScalarAsync(cancellationToken), CultureInfo.InvariantCulture); } private static async Task KeyChecksumAsync(SqlConnection connection, DatasetDefinition definition, CancellationToken cancellationToken) { using var hash = IncrementalHash.CreateHash(HashAlgorithmName.SHA256); await using var command = new SqlCommand(definition.KeysSql, connection) { CommandTimeout = 180 }; await using var reader = await command.ExecuteReaderAsync(CommandBehavior.SequentialAccess, cancellationToken); while (await reader.ReadAsync(cancellationToken)) { var key = Convert.ToString(reader[0], CultureInfo.InvariantCulture) ?? ""; hash.AppendData(Encoding.UTF8.GetBytes(key)); hash.AppendData("\n"u8); } return Convert.ToHexString(hash.GetHashAndReset()).ToLowerInvariant(); } public async Task ChangesAsync(DatasetDefinition definition, string? cursor, int? requestedLimit, CancellationToken cancellationToken) { var state = CursorCodec.Decode(cursor); var limit = Math.Clamp(requestedLimit ?? _options.DefaultPageSize, 1, _options.MaxPageSize); await using var connection = new SqlConnection(_connectionString); await connection.OpenAsync(cancellationToken); if (state.Incremental) return await IncrementalChangesAsync(connection, definition, state, limit, cancellationToken); var snapshotWatermark = state.SnapshotWatermark; if (snapshotWatermark == 0) { await using var watermarkCommand = new SqlCommand("SELECT COALESCE(MAX(Id),0) FROM dbo.XiaoshuSyncChangeLog WHERE Dataset=@dataset", connection); watermarkCommand.Parameters.AddWithValue("@dataset", definition.Key); snapshotWatermark = Convert.ToInt64(await watermarkCommand.ExecuteScalarAsync(cancellationToken), CultureInfo.InvariantCulture); } await using var command = new SqlCommand(definition.SnapshotSql, connection) { CommandTimeout = 180 }; command.Parameters.Add(new SqlParameter("@after", SqlDbType.BigInt) { Value = state.Position }); command.Parameters.Add(new SqlParameter("@limit", SqlDbType.Int) { Value = limit + 1 }); var items = await ReadProjectionItemsAsync(command, definition, cancellationToken); var hasMore = items.Count > limit; if (hasMore) items.RemoveAt(items.Count - 1); var next = items.Count == 0 ? state.Position : long.Parse(items[^1].Key, CultureInfo.InvariantCulture); var nextCursor = hasMore ? CursorCodec.EncodeSnapshot(next, snapshotWatermark) : CursorCodec.EncodeIncremental(snapshotWatermark); return new ChangePage(definition.Key, nextCursor, hasMore, items.Count, DateTimeOffset.UtcNow, "snapshot", snapshotWatermark, items); } private async Task> ReadProjectionItemsAsync(SqlCommand command, DatasetDefinition definition, CancellationToken cancellationToken) { var items = new List(); await using var reader = await command.ExecuteReaderAsync(cancellationToken); while (await reader.ReadAsync(cancellationToken)) { var key = Convert.ToString(reader[definition.LegacyKeyColumn], CultureInfo.InvariantCulture) ?? ""; var records = new List(); foreach (var mapping in definition.Records) { var fields = new Dictionary(StringComparer.Ordinal); for (var index = 0; index < reader.FieldCount; index++) { var name = reader.GetName(index); if (!name.StartsWith(mapping.Prefix, StringComparison.Ordinal)) continue; var field = name[mapping.Prefix.Length..]; var value = reader.IsDBNull(index) ? null : reader.GetValue(index); if (value is DateTime date) value = LegacyDateTime.ToUtc(date, _options.LegacyTimeZoneId); fields[field] = value; } var recordKey = Convert.ToString(fields.GetValueOrDefault(mapping.KeyField), CultureInfo.InvariantCulture) ?? key; fields["sourceKey"] = mapping.SourcePrefix + recordKey; records.Add(new ProjectionRecord(mapping.ClassName, mapping.KeyField, recordKey, fields)); } DateTimeOffset? changedAt = reader["changedAt"] is DateTime changed ? LegacyDateTime.ToUtc(changed, _options.LegacyTimeZoneId) : null; items.Add(new ChangeItem(key, "upsert", changedAt, records)); } return items; } private async Task IncrementalChangesAsync(SqlConnection connection, DatasetDefinition definition, CursorState state, int limit, CancellationToken cancellationToken) { const string logSql = "SELECT TOP (@limit) Id,LegacyKey,Operation,ChangedAt,ItemKey FROM dbo.XiaoshuSyncChangeLog WHERE Dataset=@dataset AND Id>@after ORDER BY Id"; await using var logCommand = new SqlCommand(logSql, connection) { CommandTimeout = 120 }; logCommand.Parameters.AddWithValue("@dataset", definition.Key); logCommand.Parameters.AddWithValue("@after", state.Position); logCommand.Parameters.AddWithValue("@limit", limit + 1); var changes = new List(); await using (var reader = await logCommand.ExecuteReaderAsync(cancellationToken)) { while (await reader.ReadAsync(cancellationToken)) { DateTimeOffset? changedAt = reader["ChangedAt"] is DateTime changed ? new DateTimeOffset(DateTime.SpecifyKind(changed, DateTimeKind.Utc)) : null; changes.Add(new(Convert.ToInt64(reader["Id"], CultureInfo.InvariantCulture), Convert.ToString(reader["LegacyKey"], CultureInfo.InvariantCulture) ?? "", Convert.ToString(reader["Operation"], CultureInfo.InvariantCulture) ?? "U", changedAt, reader["ItemKey"] is DBNull ? null : Convert.ToString(reader["ItemKey"], CultureInfo.InvariantCulture))); } } var hasMore = changes.Count > limit; if (hasMore) changes.RemoveAt(changes.Count - 1); var items = new List(); foreach (var change in changes) { if (change.Operation.Equals("D", StringComparison.OrdinalIgnoreCase)) { items.Add(DeleteItem(definition, change)); continue; } await using var lookup = new SqlCommand(definition.LookupSql, connection) { CommandTimeout = 120 }; lookup.Parameters.Add(new SqlParameter("@legacyKey", SqlDbType.BigInt) { Value = long.Parse(change.LegacyKey, CultureInfo.InvariantCulture) }); var projected = await ReadProjectionItemsAsync(lookup, definition, cancellationToken); items.Add(projected.Count == 0 ? DeleteItem(definition, change) : projected[0] with { ChangedAt = change.ChangedAt ?? projected[0].ChangedAt }); } var next = changes.Count == 0 ? state.Position : changes[^1].Id; return new ChangePage(definition.Key, CursorCodec.EncodeIncremental(next), hasMore, items.Count, DateTimeOffset.UtcNow, "incremental", next, items); } private static ChangeItem DeleteItem(DatasetDefinition definition, ChangeLogRow change) { var records = definition.Records.Select((mapping, index) => new ProjectionRecord(mapping.ClassName, mapping.KeyField, index == 0 ? change.LegacyKey : change.ItemKey ?? change.LegacyKey, new Dictionary())).ToList(); return new ChangeItem(change.LegacyKey, "delete", change.ChangedAt, records); } public async Task ExecuteCommandAsync(string operation, LegacyCommand command, CancellationToken cancellationToken) { await using var connection = new SqlConnection(_connectionString); await connection.OpenAsync(cancellationToken); await using var sql = new SqlCommand("dbo.XiaoshuSync_ExecuteCommand", connection) { CommandType = CommandType.StoredProcedure, CommandTimeout = 120 }; sql.Parameters.AddWithValue("@Operation", operation); sql.Parameters.AddWithValue("@Actor", command.Actor); sql.Parameters.AddWithValue("@Reason", command.Reason); sql.Parameters.AddWithValue("@IdempotencyKey", command.IdempotencyKey); sql.Parameters.AddWithValue("@ExpectedUpdatedAt", (object?)command.ExpectedUpdatedAt ?? DBNull.Value); sql.Parameters.AddWithValue("@PayloadJson", command.Payload.GetRawText()); var output = sql.Parameters.Add("@ResultJson", SqlDbType.NVarChar, -1); output.Direction = ParameterDirection.Output; await sql.ExecuteNonQueryAsync(cancellationToken); using var document = JsonDocument.Parse(Convert.ToString(output.Value, CultureInfo.InvariantCulture) ?? "{}"); return new { operation, idempotencyKey = command.IdempotencyKey, result = document.RootElement.Clone(), serverTime = DateTimeOffset.UtcNow }; } } static class CursorCodec { public static string EncodeSnapshot(long position, long watermark) => Encode($"v2:s:{Math.Max(0, position)}:{Math.Max(0, watermark)}"); public static string EncodeIncremental(long position) => Encode($"v2:c:{Math.Max(0, position)}"); public static CursorState Decode(string? value) { if (string.IsNullOrWhiteSpace(value)) return new(false, 0, 0); try { var decoded = Encoding.UTF8.GetString(Convert.FromBase64String(value)); var parts = decoded.Split(':'); if (parts.Length == 2 && parts[0] == "v1" && long.TryParse(parts[1], out var legacy)) return new(false, Math.Max(0, legacy), 0); if (parts.Length >= 3 && parts[0] == "v2" && long.TryParse(parts[2], out var position)) { if (parts[1] == "c") return new(true, Math.Max(0, position), 0); if (parts[1] == "s") return new(false, Math.Max(0, position), parts.Length > 3 && long.TryParse(parts[3], out var watermark) ? Math.Max(0, watermark) : 0); } } catch { } return new(false, 0, 0); } private static string Encode(string value) => Convert.ToBase64String(Encoding.UTF8.GetBytes(value)); } static class LegacyDateTime { public static DateTimeOffset ToUtc(DateTime value, string timeZoneId) { var unspecified = DateTime.SpecifyKind(value, DateTimeKind.Unspecified); var zone = TimeZoneInfo.FindSystemTimeZoneById(timeZoneId); return new DateTimeOffset(TimeZoneInfo.ConvertTimeToUtc(unspecified, zone), TimeSpan.Zero); } }