From 58a459bb15a736f245f6f2ccf1e9fcc8a6375568 Mon Sep 17 00:00:00 2001 From: Rogee Date: Tue, 22 Sep 2026 17:32:16 +0800 Subject: [PATCH] Split large connection-authorized batches by payload size --- .../RemoteAgentHostedService.cs | 61 ++++++++++++++----- 1 file changed, 47 insertions(+), 14 deletions(-) diff --git a/node-agent/WxAgent.Service/RemoteAgentHostedService.cs b/node-agent/WxAgent.Service/RemoteAgentHostedService.cs index 21e4613..922fb94 100644 --- a/node-agent/WxAgent.Service/RemoteAgentHostedService.cs +++ b/node-agent/WxAgent.Service/RemoteAgentHostedService.cs @@ -192,25 +192,41 @@ public sealed class RemoteAgentHostedService( continue; var sequence = sourceChanged ? 1 : checkpoint.ConfirmedSequence + 1; - var remainingConversations = collection.Conversations.ToList(); - var remainingMessages = collection.Messages.ToList(); + var conversationIndex = 0; + var messageIndex = 0; var cursorStartJson = SerializeCursors(cursorStart); var cursorEndJson = SerializeCursors(collection.NextCursors); var allQueued = true; RemoteSyncBatch? lastQueuedBatch = null; - while (remainingConversations.Count > 0 || remainingMessages.Count > 0 || lastQueuedBatch is null) + while (conversationIndex < collection.Conversations.Count || messageIndex < collection.Messages.Count || lastQueuedBatch is null) { - var conversationCount = Math.Min(remainingConversations.Count, RemoteDataProtocol.MaxBatchMessages); - var conversations = remainingConversations.Take(conversationCount).ToArray(); - remainingConversations.RemoveRange(0, conversationCount); - var messageCount = Math.Min(remainingMessages.Count, RemoteDataProtocol.MaxBatchMessages - conversations.Length); - var messages = remainingMessages.Take(messageCount).ToArray(); - remainingMessages.RemoveRange(0, messageCount); - var batch = new RemoteSyncBatch( - remote.NodeId!, account.AccountId, Guid.NewGuid().ToString("N"), collection.SourceGeneration, - "messages", sequence++, cursorStartJson, cursorEndJson, string.Empty, - collection.IsComplete ? "complete" : "partial", conversations, messages); - batch = batch with { PayloadHash = RemoteDataBatchAuthorization.ComputePayloadHash(batch) }; + var conversations = new List(); + var messages = new List(); + while (conversationIndex < collection.Conversations.Count && conversations.Count + messages.Count < RemoteDataProtocol.MaxBatchMessages) + { + conversations.Add(collection.Conversations[conversationIndex]); + var candidate = CreateSyncBatch(remote, account.AccountId, collection.SourceGeneration, sequence, cursorStartJson, cursorEndJson, collection.IsComplete, conversations, messages); + if (JsonSerializer.SerializeToUtf8Bytes(candidate, RemoteJson.Options).LongLength > RemoteDataProtocol.MaxBatchBytes) + { + conversations.RemoveAt(conversations.Count - 1); + break; + } + conversationIndex++; + } + while (messageIndex < collection.Messages.Count && conversations.Count + messages.Count < RemoteDataProtocol.MaxBatchMessages) + { + messages.Add(collection.Messages[messageIndex]); + var candidate = CreateSyncBatch(remote, account.AccountId, collection.SourceGeneration, sequence, cursorStartJson, cursorEndJson, collection.IsComplete, conversations, messages); + if (JsonSerializer.SerializeToUtf8Bytes(candidate, RemoteJson.Options).LongLength > RemoteDataProtocol.MaxBatchBytes) + { + messages.RemoveAt(messages.Count - 1); + break; + } + messageIndex++; + } + if (conversations.Count == 0 && messages.Count == 0) + throw new WxAgentException(WxAgentErrorCode.InvalidArgument, "A single sync item exceeds the maximum batch size."); + var batch = CreateSyncBatch(remote, account.AccountId, collection.SourceGeneration, sequence++, cursorStartJson, cursorEndJson, collection.IsComplete, conversations, messages); var queued = queue.Enqueue(reporting, batch); if (!queued.Accepted || queued.Batch is null) { @@ -303,6 +319,23 @@ public sealed class RemoteAgentHostedService( } } + private static RemoteSyncBatch CreateSyncBatch( + RemoteAgentOptions remote, + string accountId, + string sourceGeneration, + long sequence, + string cursorStart, + string cursorEnd, + bool complete, + IReadOnlyList conversations, + IReadOnlyList messages) + { + var batch = new RemoteSyncBatch( + remote.NodeId!, accountId, Guid.NewGuid().ToString("N"), sourceGeneration, "messages", sequence, + cursorStart, cursorEnd, string.Empty, complete ? "complete" : "partial", conversations, messages); + return batch with { PayloadHash = RemoteDataBatchAuthorization.ComputePayloadHash(batch) }; + } + private static string SerializeCursors(IReadOnlyDictionary cursors) => JsonSerializer.Serialize(cursors.OrderBy(item => item.Key).ToDictionary(item => item.Key, item => item.Value), RemoteJson.Options);