Split large connection-authorized batches by payload size

This commit is contained in:
2026-09-22 17:32:16 +08:00
parent e057c5cdb3
commit 58a459bb15
@@ -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<RemoteSyncConversation>();
var messages = new List<RemoteSyncMessage>();
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<RemoteSyncConversation> conversations,
IReadOnlyList<RemoteSyncMessage> 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<string, long> cursors) =>
JsonSerializer.Serialize(cursors.OrderBy(item => item.Key).ToDictionary(item => item.Key, item => item.Value), RemoteJson.Options);