Program.cs 29 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326
  1. using System.Collections.Concurrent;
  2. using System.Data;
  3. using System.Globalization;
  4. using System.Security.Cryptography;
  5. using System.Text;
  6. using System.Text.Json;
  7. using Microsoft.Data.SqlClient;
  8. using Microsoft.Extensions.Options;
  9. var builder = WebApplication.CreateBuilder(args);
  10. builder.WebHost.ConfigureKestrel(options => options.Limits.MaxRequestBodySize = 1_048_576);
  11. builder.Services.AddOptions<SyncOptions>()
  12. .Bind(builder.Configuration.GetSection("XiaoshuSync"))
  13. .Validate(options => !string.IsNullOrWhiteSpace(options.KeyId), "XiaoshuSync KeyId 未配置")
  14. .Validate(options => options.Secret.Length >= 32, "XiaoshuSync Secret 必须至少 32 位")
  15. .Validate(options => options.MaxClockSkewSeconds is >= 30 and <= 900, "签名时间窗口必须在 30 至 900 秒之间")
  16. .Validate(options => options.DefaultPageSize > 0 && options.MaxPageSize >= options.DefaultPageSize && options.MaxPageSize <= 5000, "同步分页配置无效")
  17. .Validate(options => options.MaxCommandBodyBytes is >= 1024 and <= 1_048_576, "写入请求大小上限无效")
  18. .Validate(options => options.PathBase.StartsWith('/') && !options.PathBase.EndsWith('/'), "PathBase 必须以 / 开头且不能以 / 结尾")
  19. .ValidateOnStart();
  20. builder.Services.AddSingleton<ReplayGuard>();
  21. builder.Services.AddSingleton<LegacyRepository>();
  22. var app = builder.Build();
  23. app.UsePathBase(app.Services.GetRequiredService<IOptions<SyncOptions>>().Value.PathBase);
  24. app.UseExceptionHandler(exceptionApp => exceptionApp.Run(async context =>
  25. {
  26. context.Response.StatusCode = StatusCodes.Status500InternalServerError;
  27. context.Response.ContentType = "application/json";
  28. await context.Response.WriteAsJsonAsync(new { error = "sync_bridge_error", message = "同步桥执行失败,请查看服务器日志" });
  29. }));
  30. app.Use(async (context, next) =>
  31. {
  32. if (!context.Request.Path.StartsWithSegments("/v1") || context.Request.Path.Equals("/v1/health")) { await next(); return; }
  33. context.Response.Headers.CacheControl = "no-store";
  34. var configured = context.RequestServices.GetRequiredService<IOptions<SyncOptions>>().Value;
  35. 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; }
  36. var rejection = await RequestAuthenticator.ValidateAsync(context, configured, context.RequestServices.GetRequiredService<ReplayGuard>());
  37. if (rejection is not null) { context.Response.StatusCode = rejection.Value.Status; await context.Response.WriteAsJsonAsync(new { error = rejection.Value.Message }); return; }
  38. await next();
  39. });
  40. app.MapGet("/v1/health", () => Results.Ok(new { service = "xiaoshu-legacy-sync", status = "ok", version = "v1", time = DateTimeOffset.UtcNow }));
  41. app.MapGet("/v1/manifest", async (LegacyRepository repository, CancellationToken cancellationToken) => Results.Ok(await repository.ManifestAsync(cancellationToken)));
  42. app.MapGet("/v1/changes", async (string dataset, string? cursor, int? limit, LegacyRepository repository, CancellationToken cancellationToken) =>
  43. {
  44. if (!DatasetCatalog.All.TryGetValue(dataset, out var definition)) return Results.BadRequest(new { error = "unsupported_dataset", message = $"不支持的数据集:{dataset}" });
  45. return Results.Ok(await repository.ChangesAsync(definition, cursor, limit, cancellationToken));
  46. });
  47. app.MapPost("/v1/commands/{operation}", async (string operation, HttpContext context, LegacyRepository repository, IOptions<SyncOptions> configured, CancellationToken cancellationToken) =>
  48. {
  49. if (!CommandCatalog.Allowed.Contains(operation)) return Results.BadRequest(new { error = "unsupported_operation", message = $"不支持的写入操作:{operation}" });
  50. if (context.Request.ContentLength is > 0 && context.Request.ContentLength > configured.Value.MaxCommandBodyBytes) return Results.StatusCode(StatusCodes.Status413PayloadTooLarge);
  51. var command = await context.Request.ReadFromJsonAsync<LegacyCommand>(cancellationToken: cancellationToken);
  52. 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 必填且长度必须有效" });
  53. return Results.Ok(await repository.ExecuteCommandAsync(operation, command, cancellationToken));
  54. });
  55. app.MapGet("/health", () => Results.Ok(new { service = "xiaoshu-legacy-sync", status = "ok", time = DateTimeOffset.UtcNow }));
  56. app.Run();
  57. sealed class SyncOptions
  58. {
  59. public string KeyId { get; init; } = "";
  60. public string Secret { get; init; } = "";
  61. public string[] AllowedIps { get; init; } = [];
  62. public int MaxClockSkewSeconds { get; init; } = 300;
  63. public int DefaultPageSize { get; init; } = 500;
  64. public int MaxPageSize { get; init; } = 1000;
  65. public int MaxCommandBodyBytes { get; init; } = 1_048_576;
  66. public string LegacyTimeZoneId { get; init; } = "China Standard Time";
  67. public string PathBase { get; init; } = "/xiaoshu-sync";
  68. }
  69. readonly record struct AuthRejection(int Status, string Message);
  70. static class RequestAuthenticator
  71. {
  72. public static async Task<AuthRejection?> ValidateAsync(HttpContext context, SyncOptions options, ReplayGuard replayGuard)
  73. {
  74. if (string.IsNullOrWhiteSpace(options.KeyId) || string.IsNullOrWhiteSpace(options.Secret)) return new(503, "同步桥密钥尚未配置");
  75. var remoteIp = context.Connection.RemoteIpAddress?.ToString() ?? "";
  76. if (options.AllowedIps.Length > 0 && !options.AllowedIps.Contains(remoteIp, StringComparer.OrdinalIgnoreCase)) return new(403, "来源 IP 不在白名单");
  77. var keyId = context.Request.Headers["X-Xiaoshu-Key-Id"].ToString();
  78. var timestampText = context.Request.Headers["X-Xiaoshu-Timestamp"].ToString();
  79. var nonce = context.Request.Headers["X-Xiaoshu-Nonce"].ToString();
  80. var signature = context.Request.Headers["X-Xiaoshu-Signature"].ToString();
  81. if (!long.TryParse(timestampText, out var timestamp) || string.IsNullOrWhiteSpace(nonce) || string.IsNullOrWhiteSpace(signature) || !FixedText(keyId, options.KeyId)) return new(401, "签名头缺失或无效");
  82. if (Math.Abs(DateTimeOffset.UtcNow.ToUnixTimeSeconds() - timestamp) > options.MaxClockSkewSeconds) return new(401, "请求时间戳已过期");
  83. if (!replayGuard.TryUse(nonce, timestamp, options.MaxClockSkewSeconds)) return new(409, "nonce 已使用,拒绝重放");
  84. context.Request.EnableBuffering();
  85. using var reader = new StreamReader(context.Request.Body, Encoding.UTF8, leaveOpen: true);
  86. var body = await reader.ReadToEndAsync(); context.Request.Body.Position = 0;
  87. var bodyHash = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(body))).ToLowerInvariant();
  88. var canonical = $"{context.Request.Method.ToUpperInvariant()}\n{context.Request.PathBase}{context.Request.Path}{context.Request.QueryString}\n{timestampText}\n{nonce}\n{bodyHash}";
  89. using var hmac = new HMACSHA256(Encoding.UTF8.GetBytes(options.Secret));
  90. var expected = Convert.ToHexString(hmac.ComputeHash(Encoding.UTF8.GetBytes(canonical))).ToLowerInvariant();
  91. return FixedText(signature.ToLowerInvariant(), expected) ? null : new(401, "请求签名无效");
  92. }
  93. private static bool FixedText(string left, string right) => left.Length == right.Length && CryptographicOperations.FixedTimeEquals(Encoding.UTF8.GetBytes(left), Encoding.UTF8.GetBytes(right));
  94. }
  95. sealed class ReplayGuard
  96. {
  97. private readonly ConcurrentDictionary<string, long> _nonces = new();
  98. public bool TryUse(string nonce, long timestamp, int lifetimeSeconds)
  99. {
  100. var floor = DateTimeOffset.UtcNow.ToUnixTimeSeconds() - lifetimeSeconds;
  101. foreach (var stale in _nonces.Where(item => item.Value < floor).Select(item => item.Key)) _nonces.TryRemove(stale, out _);
  102. return nonce.Length is >= 16 and <= 128 && _nonces.TryAdd(nonce, timestamp);
  103. }
  104. }
  105. sealed record LegacyCommand(string Actor, string Reason, string IdempotencyKey, string? ExpectedUpdatedAt, JsonElement Payload);
  106. sealed record ProjectionRecord(string ClassName, string KeyField, string LegacyKey, Dictionary<string, object?> Fields);
  107. sealed record ChangeItem(string Key, string Operation, DateTimeOffset? ChangedAt, IReadOnlyList<ProjectionRecord> Records);
  108. sealed record ChangePage(string Dataset, string Cursor, bool HasMore, int Count, DateTimeOffset ServerTime, string Phase, long Watermark, IReadOnlyList<ChangeItem> Items);
  109. sealed record DatasetDefinition(string Key, string Label, string CountSql, string KeysSql, string SnapshotSql, string LookupSql, string LegacyKeyColumn, IReadOnlyList<RecordMapping> Records);
  110. sealed record RecordMapping(string ClassName, string KeyField, string Prefix, string SourcePrefix);
  111. sealed record CursorState(bool Incremental, long Position, long SnapshotWatermark);
  112. sealed record ChangeLogRow(long Id, string LegacyKey, string Operation, DateTimeOffset? ChangedAt, string? ItemKey);
  113. static class DatasetCatalog
  114. {
  115. private const string ActiveCommon = "COALESCE(c.Status,99)<>-2";
  116. public static readonly IReadOnlyDictionary<string, DatasetDefinition> All = new Dictionary<string, DatasetDefinition>(StringComparer.OrdinalIgnoreCase)
  117. {
  118. ["users"] = new("users", "全部业务账号", "SELECT COUNT_BIG(*) FROM ZL_User",
  119. "SELECT CAST(UserID AS bigint) legacyKey FROM ZL_User ORDER BY UserID",
  120. """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""",
  121. """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:")]),
  122. ["nodes"] = new("nodes", "课程与词库目录", "SELECT COUNT_BIG(*) FROM ZL_Node",
  123. "SELECT CAST(NodeID AS bigint) legacyKey FROM ZL_Node ORDER BY NodeID",
  124. """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""",
  125. """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:")]),
  126. ["course-bindings"] = Content("course-bindings", "课程绑定", 58, "ZL_C_kcbd", "b", ["yhid","kcid","yxx","cksl","syjd"]),
  127. ["appointments"] = Content("appointments", "预约排课", 54, "ZL_C_order", "a", ["szyh","pl","fxpl","plxm","kcid","yykcid","dslx","jffs","bxrq","sdsd","yysj","dszt","kssj","jssj","scsj","szmd"]),
  128. ["lessons"] = Content("lessons", "上课记录", 59, "ZL_C_skjl", "l", ["xymz","jsmz","kcid","kclx","kcmc","pf","pjnr","pldp","kzsj","yyds","plpjsj","szmdid"]),
  129. ["learning-records"] = Content("learning-records", "每日学习记录", 56, "ZL_C_ss", "d", ["userId","pl","dqrq","dsid","learned","ygg","djq","xxqs","fxrl","szmdid"]),
  130. ["practice-records"] = Content("practice-records", "练习记录", 53, "ZL_C_lxjl", "p", ["yhid","kcid","scid","xxcs","jrscb"]),
  131. ["memory-records"] = Content("memory-records", "抗遗忘记录", 60, "ZL_C_gywjl", "m", ["yhid","plid","kcid","xxjlid","fxzt","wcsj","kywrq","kywsj","OrderID"]),
  132. ["assessments"] = Content("assessments", "测评档案", 61, "ZL_C_cpda", "e", ["userId","df","askid","wrong","answerid","dontKnow","prev_score:prevScore","totalScore"]),
  133. ["vocabulary"] = Content("vocabulary", "词库", 52, "ZL_C_ck", "v", ["sy","yb","lj","yp"]),
  134. ["special-training"] = Content("special-training", "专项训练", 2, "ZL_C_Article", "r", ["ico","mp4","author","source","content","synopsis"], " AND c.NodeID=3")
  135. };
  136. private static DatasetDefinition Content(string key, string label, int modelId, string addonTable, string alias, string[] addonFields, string extraWhere = "")
  137. {
  138. var addonSelect = string.Join(',', addonFields.Select(field =>
  139. {
  140. var mapping = field.Split(':', 2);
  141. var sourceField = mapping[0];
  142. var targetField = mapping.Length == 2 ? mapping[1] : ToCamel(sourceField);
  143. return $"{alias}.[{sourceField}] [addon__{targetField}]";
  144. }));
  145. 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}";
  146. var from = $" FROM ZL_CommonModel c INNER JOIN {addonTable} {alias} ON {alias}.ID=c.ItemID WHERE c.ModelID={modelId}{extraWhere} AND {ActiveCommon}";
  147. var snapshotSql = $"SELECT TOP (@limit) {select}{from} AND c.GeneralID>@after ORDER BY c.GeneralID";
  148. var lookupSql = $"SELECT TOP (1) {select}{from} AND c.GeneralID=@legacyKey";
  149. var countSql = $"SELECT COUNT_BIG(*){from}";
  150. var keysSql = $"SELECT CAST(c.GeneralID AS bigint) legacyKey{from} ORDER BY c.GeneralID";
  151. 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:")]);
  152. }
  153. 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() };
  154. private static string ToCamel(string value) => value.Length == 0 ? value : char.ToLowerInvariant(value[0]) + value[1..];
  155. }
  156. static class CommandCatalog
  157. {
  158. public static readonly HashSet<string> 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" };
  159. }
  160. sealed class LegacyRepository(IConfiguration configuration, IOptions<SyncOptions> options)
  161. {
  162. private readonly string _connectionString = configuration.GetConnectionString("LegacySqlServer") ?? throw new InvalidOperationException("LegacySqlServer connection string is missing");
  163. private readonly SyncOptions _options = options.Value;
  164. public async Task<object> ManifestAsync(CancellationToken cancellationToken)
  165. {
  166. await using var connection = new SqlConnection(_connectionString); await connection.OpenAsync(cancellationToken);
  167. var watermarkStart = await GlobalWatermarkAsync(connection, cancellationToken);
  168. var datasets = new List<object>();
  169. foreach (var definition in DatasetCatalog.All.Values)
  170. {
  171. await using var command = new SqlCommand(definition.CountSql, connection) { CommandTimeout = 120 };
  172. var count = Convert.ToInt64(await command.ExecuteScalarAsync(cancellationToken), CultureInfo.InvariantCulture);
  173. var checksum = await KeyChecksumAsync(connection, definition, cancellationToken);
  174. await using var watermarkCommand = new SqlCommand("SELECT COALESCE(MAX(Id),0) FROM dbo.XiaoshuSyncChangeLog WHERE Dataset=@dataset", connection);
  175. watermarkCommand.Parameters.AddWithValue("@dataset", definition.Key);
  176. var watermark = Convert.ToInt64(await watermarkCommand.ExecuteScalarAsync(cancellationToken), CultureInfo.InvariantCulture);
  177. datasets.Add(new { key = definition.Key, label = definition.Label, count, watermark, checksum });
  178. }
  179. var watermarkEnd = await GlobalWatermarkAsync(connection, cancellationToken);
  180. var now = DateTimeOffset.UtcNow;
  181. return new { manifestTime = now, serverTime = now, consistent = watermarkStart == watermarkEnd, watermarkStart, watermarkEnd, datasets };
  182. }
  183. private static async Task<long> GlobalWatermarkAsync(SqlConnection connection, CancellationToken cancellationToken)
  184. {
  185. await using var command = new SqlCommand("SELECT COALESCE(MAX(Id),0) FROM dbo.XiaoshuSyncChangeLog", connection);
  186. return Convert.ToInt64(await command.ExecuteScalarAsync(cancellationToken), CultureInfo.InvariantCulture);
  187. }
  188. private static async Task<string> KeyChecksumAsync(SqlConnection connection, DatasetDefinition definition, CancellationToken cancellationToken)
  189. {
  190. using var hash = IncrementalHash.CreateHash(HashAlgorithmName.SHA256);
  191. await using var command = new SqlCommand(definition.KeysSql, connection) { CommandTimeout = 180 };
  192. await using var reader = await command.ExecuteReaderAsync(CommandBehavior.SequentialAccess, cancellationToken);
  193. while (await reader.ReadAsync(cancellationToken))
  194. {
  195. var key = Convert.ToString(reader[0], CultureInfo.InvariantCulture) ?? "";
  196. hash.AppendData(Encoding.UTF8.GetBytes(key)); hash.AppendData("\n"u8);
  197. }
  198. return Convert.ToHexString(hash.GetHashAndReset()).ToLowerInvariant();
  199. }
  200. public async Task<ChangePage> ChangesAsync(DatasetDefinition definition, string? cursor, int? requestedLimit, CancellationToken cancellationToken)
  201. {
  202. var state = CursorCodec.Decode(cursor); var limit = Math.Clamp(requestedLimit ?? _options.DefaultPageSize, 1, _options.MaxPageSize);
  203. await using var connection = new SqlConnection(_connectionString); await connection.OpenAsync(cancellationToken);
  204. if (state.Incremental) return await IncrementalChangesAsync(connection, definition, state, limit, cancellationToken);
  205. var snapshotWatermark = state.SnapshotWatermark;
  206. if (snapshotWatermark == 0)
  207. {
  208. await using var watermarkCommand = new SqlCommand("SELECT COALESCE(MAX(Id),0) FROM dbo.XiaoshuSyncChangeLog WHERE Dataset=@dataset", connection);
  209. watermarkCommand.Parameters.AddWithValue("@dataset", definition.Key);
  210. snapshotWatermark = Convert.ToInt64(await watermarkCommand.ExecuteScalarAsync(cancellationToken), CultureInfo.InvariantCulture);
  211. }
  212. await using var command = new SqlCommand(definition.SnapshotSql, connection) { CommandTimeout = 180 };
  213. command.Parameters.Add(new SqlParameter("@after", SqlDbType.BigInt) { Value = state.Position }); command.Parameters.Add(new SqlParameter("@limit", SqlDbType.Int) { Value = limit + 1 });
  214. var items = await ReadProjectionItemsAsync(command, definition, cancellationToken);
  215. 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);
  216. var nextCursor = hasMore ? CursorCodec.EncodeSnapshot(next, snapshotWatermark) : CursorCodec.EncodeIncremental(snapshotWatermark);
  217. return new ChangePage(definition.Key, nextCursor, hasMore, items.Count, DateTimeOffset.UtcNow, "snapshot", snapshotWatermark, items);
  218. }
  219. private async Task<List<ChangeItem>> ReadProjectionItemsAsync(SqlCommand command, DatasetDefinition definition, CancellationToken cancellationToken)
  220. {
  221. var items = new List<ChangeItem>(); await using var reader = await command.ExecuteReaderAsync(cancellationToken);
  222. while (await reader.ReadAsync(cancellationToken))
  223. {
  224. var key = Convert.ToString(reader[definition.LegacyKeyColumn], CultureInfo.InvariantCulture) ?? ""; var records = new List<ProjectionRecord>();
  225. foreach (var mapping in definition.Records)
  226. {
  227. var fields = new Dictionary<string, object?>(StringComparer.Ordinal);
  228. for (var index = 0; index < reader.FieldCount; index++)
  229. {
  230. var name = reader.GetName(index); if (!name.StartsWith(mapping.Prefix, StringComparison.Ordinal)) continue;
  231. var field = name[mapping.Prefix.Length..]; var value = reader.IsDBNull(index) ? null : reader.GetValue(index);
  232. if (value is DateTime date) value = LegacyDateTime.ToUtc(date, _options.LegacyTimeZoneId); fields[field] = value;
  233. }
  234. var recordKey = Convert.ToString(fields.GetValueOrDefault(mapping.KeyField), CultureInfo.InvariantCulture) ?? key; fields["sourceKey"] = mapping.SourcePrefix + recordKey;
  235. records.Add(new ProjectionRecord(mapping.ClassName, mapping.KeyField, recordKey, fields));
  236. }
  237. DateTimeOffset? changedAt = reader["changedAt"] is DateTime changed ? LegacyDateTime.ToUtc(changed, _options.LegacyTimeZoneId) : null;
  238. items.Add(new ChangeItem(key, "upsert", changedAt, records));
  239. }
  240. return items;
  241. }
  242. private async Task<ChangePage> IncrementalChangesAsync(SqlConnection connection, DatasetDefinition definition, CursorState state, int limit, CancellationToken cancellationToken)
  243. {
  244. const string logSql = "SELECT TOP (@limit) Id,LegacyKey,Operation,ChangedAt,ItemKey FROM dbo.XiaoshuSyncChangeLog WHERE Dataset=@dataset AND Id>@after ORDER BY Id";
  245. await using var logCommand = new SqlCommand(logSql, connection) { CommandTimeout = 120 };
  246. logCommand.Parameters.AddWithValue("@dataset", definition.Key); logCommand.Parameters.AddWithValue("@after", state.Position); logCommand.Parameters.AddWithValue("@limit", limit + 1);
  247. var changes = new List<ChangeLogRow>(); await using (var reader = await logCommand.ExecuteReaderAsync(cancellationToken))
  248. {
  249. while (await reader.ReadAsync(cancellationToken))
  250. {
  251. DateTimeOffset? changedAt = reader["ChangedAt"] is DateTime changed ? new DateTimeOffset(DateTime.SpecifyKind(changed, DateTimeKind.Utc)) : null;
  252. 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)));
  253. }
  254. }
  255. var hasMore = changes.Count > limit; if (hasMore) changes.RemoveAt(changes.Count - 1);
  256. var items = new List<ChangeItem>();
  257. foreach (var change in changes)
  258. {
  259. if (change.Operation.Equals("D", StringComparison.OrdinalIgnoreCase)) { items.Add(DeleteItem(definition, change)); continue; }
  260. await using var lookup = new SqlCommand(definition.LookupSql, connection) { CommandTimeout = 120 };
  261. lookup.Parameters.Add(new SqlParameter("@legacyKey", SqlDbType.BigInt) { Value = long.Parse(change.LegacyKey, CultureInfo.InvariantCulture) });
  262. var projected = await ReadProjectionItemsAsync(lookup, definition, cancellationToken);
  263. items.Add(projected.Count == 0 ? DeleteItem(definition, change) : projected[0] with { ChangedAt = change.ChangedAt ?? projected[0].ChangedAt });
  264. }
  265. var next = changes.Count == 0 ? state.Position : changes[^1].Id;
  266. return new ChangePage(definition.Key, CursorCodec.EncodeIncremental(next), hasMore, items.Count, DateTimeOffset.UtcNow, "incremental", next, items);
  267. }
  268. private static ChangeItem DeleteItem(DatasetDefinition definition, ChangeLogRow change)
  269. {
  270. var records = definition.Records.Select((mapping, index) => new ProjectionRecord(mapping.ClassName, mapping.KeyField, index == 0 ? change.LegacyKey : change.ItemKey ?? change.LegacyKey, new Dictionary<string, object?>())).ToList();
  271. return new ChangeItem(change.LegacyKey, "delete", change.ChangedAt, records);
  272. }
  273. public async Task<object> ExecuteCommandAsync(string operation, LegacyCommand command, CancellationToken cancellationToken)
  274. {
  275. await using var connection = new SqlConnection(_connectionString); await connection.OpenAsync(cancellationToken);
  276. await using var sql = new SqlCommand("dbo.XiaoshuSync_ExecuteCommand", connection) { CommandType = CommandType.StoredProcedure, CommandTimeout = 120 };
  277. 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());
  278. var output = sql.Parameters.Add("@ResultJson", SqlDbType.NVarChar, -1); output.Direction = ParameterDirection.Output; await sql.ExecuteNonQueryAsync(cancellationToken);
  279. using var document = JsonDocument.Parse(Convert.ToString(output.Value, CultureInfo.InvariantCulture) ?? "{}"); return new { operation, idempotencyKey = command.IdempotencyKey, result = document.RootElement.Clone(), serverTime = DateTimeOffset.UtcNow };
  280. }
  281. }
  282. static class CursorCodec
  283. {
  284. public static string EncodeSnapshot(long position, long watermark) => Encode($"v2:s:{Math.Max(0, position)}:{Math.Max(0, watermark)}");
  285. public static string EncodeIncremental(long position) => Encode($"v2:c:{Math.Max(0, position)}");
  286. public static CursorState Decode(string? value)
  287. {
  288. if (string.IsNullOrWhiteSpace(value)) return new(false, 0, 0);
  289. try
  290. {
  291. var decoded = Encoding.UTF8.GetString(Convert.FromBase64String(value)); var parts = decoded.Split(':');
  292. if (parts.Length == 2 && parts[0] == "v1" && long.TryParse(parts[1], out var legacy)) return new(false, Math.Max(0, legacy), 0);
  293. if (parts.Length >= 3 && parts[0] == "v2" && long.TryParse(parts[2], out var position))
  294. {
  295. if (parts[1] == "c") return new(true, Math.Max(0, position), 0);
  296. 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);
  297. }
  298. }
  299. catch { }
  300. return new(false, 0, 0);
  301. }
  302. private static string Encode(string value) => Convert.ToBase64String(Encoding.UTF8.GetBytes(value));
  303. }
  304. static class LegacyDateTime
  305. {
  306. public static DateTimeOffset ToUtc(DateTime value, string timeZoneId)
  307. {
  308. var unspecified = DateTime.SpecifyKind(value, DateTimeKind.Unspecified);
  309. var zone = TimeZoneInfo.FindSystemTimeZoneById(timeZoneId);
  310. return new DateTimeOffset(TimeZoneInfo.ConvertTimeToUtc(unspecified, zone), TimeSpan.Zero);
  311. }
  312. }