feat: add account-scoped data synchronization

This commit is contained in:
2026-09-22 09:58:58 +08:00
parent ab9ff389f2
commit 72e040546a
37 changed files with 4252 additions and 159 deletions
@@ -10,7 +10,8 @@ public sealed class RemoteAgentHostedService(
ServiceOptions options,
IAgentBackend backend,
ILogger<RemoteAgentHostedService> logger,
RemoteEventQueue? remoteQueue = null) : BackgroundService
RemoteEventQueue? remoteQueue = null,
IRemoteDataCollector? dataCollector = null) : BackgroundService
{
private static readonly TimeSpan HeartbeatInterval = TimeSpan.FromSeconds(10);
private static readonly TimeSpan RetryDelay = TimeSpan.FromSeconds(2);
@@ -41,6 +42,17 @@ public sealed class RemoteAgentHostedService(
using var client = new RemoteControlClient(remote);
var ledger = new RemoteTaskLedger(Path.Combine(options.DataDirectory, "remote-task-ledger.json"));
var eventQueue = remoteQueue ?? new RemoteEventQueue(Path.Combine(options.DataDirectory, "remote-event-queue.json"));
var dataQueue = options.EnableDataSync
? new RemoteDataBatchQueue(Path.Combine(options.DataDirectory, "remote-data-queue.json"), options.DataSyncQueueMaxItems, options.DataSyncQueueMaxBytes)
: null;
var syncState = options.EnableDataSync
? new RemoteDataSyncStateStore(Path.Combine(options.DataDirectory, "remote-data-sync-state.json"))
: null;
var lastDataSyncAt = DateTimeOffset.MinValue;
string? registeredActiveAccountId = null;
bool? registeredActiveAccountVerified = null;
long registeredReportingConfigVersion = -1;
ReportingConfig? activeReporting = null;
var retry = RetryDelay;
while (!stoppingToken.IsCancellationRequested)
{
@@ -54,16 +66,43 @@ public sealed class RemoteAgentHostedService(
continue;
}
var reporting = runtime.Reporting;
activeReporting = reporting;
var snapshot = await ReadSnapshotAsync(remote, stoppingToken);
if (client.AuthState != RemoteAuthState.Authenticated)
var registrationNeedsRefresh = client.AuthState != RemoteAuthState.Authenticated
|| !string.Equals(registeredActiveAccountId, snapshot.ActiveAccountId, StringComparison.Ordinal)
|| registeredActiveAccountVerified != snapshot.ActiveAccountVerified
|| registeredReportingConfigVersion != reporting.ConfigVersion;
if (registrationNeedsRefresh)
{
await client.RegisterAsync(CreateRegistration(remote, reporting, snapshot), stoppingToken);
registeredActiveAccountId = snapshot.ActiveAccountId;
registeredActiveAccountVerified = snapshot.ActiveAccountVerified;
registeredReportingConfigVersion = reporting.ConfigVersion;
retry = RetryDelay;
}
await client.HeartbeatAsync(CreateHeartbeat(remote, reporting, snapshot, null), stoppingToken);
await ReplayUnreportedResultsAsync(client, ledger, reporting, stoppingToken);
await client.FlushEventsAsync(eventQueue, reporting, stoppingToken);
if (dataQueue is not null && syncState is not null && dataCollector is not null
&& DateTimeOffset.UtcNow - lastDataSyncAt >= TimeSpan.FromSeconds(options.DataSyncIntervalSeconds))
{
await ReconcileDataStatusAsync(client, reporting, dataQueue, syncState, stoppingToken);
try
{
var confirmedBatches = await client.FlushDataBatchesAsync(dataQueue, reporting, stoppingToken);
foreach (var confirmedBatch in confirmedBatches)
syncState.MarkConfirmed(confirmedBatch);
}
catch (Exception exception) when (exception is HttpRequestException or RemoteClientException)
{
logger.LogWarning("Data sync delivery deferred; type={ExceptionType}.", exception.GetType().Name);
}
// Collection must continue while the platform is unreachable; the durable queue is
// the offline buffer and will be flushed on the next successful connection.
await CollectDataBatchesAsync(remote, reporting, dataCollector, dataQueue, syncState, stoppingToken);
lastDataSyncAt = DateTimeOffset.UtcNow;
}
var pollAccountIds = reporting.Accounts
.Where(account => account.Enabled)
.Select(account => account.AccountId)
@@ -94,18 +133,120 @@ public sealed class RemoteAgentHostedService(
{
logger.LogWarning("Remote control-plane request failed; code={Code}; status={StatusCode}; correlationId={CorrelationId}.",
exception.Code, exception.StatusCode, exception.CorrelationId);
await TryCollectOfflineDataBatchesAsync(remote, activeReporting, dataCollector, dataQueue, syncState, stoppingToken);
await DelayAsync(retry, stoppingToken);
retry = TimeSpan.FromSeconds(Math.Min(retry.TotalSeconds * 2, 30));
}
catch (Exception exception)
{
logger.LogWarning("Remote agent cycle failed; type={ExceptionType}.", exception.GetType().Name);
await TryCollectOfflineDataBatchesAsync(remote, activeReporting, dataCollector, dataQueue, syncState, stoppingToken);
await DelayAsync(retry, stoppingToken);
retry = TimeSpan.FromSeconds(Math.Min(retry.TotalSeconds * 2, 30));
}
}
}
private async Task CollectDataBatchesAsync(
RemoteAgentOptions remote,
ReportingConfig reporting,
IRemoteDataCollector collector,
RemoteDataBatchQueue queue,
RemoteDataSyncStateStore state,
CancellationToken cancellationToken)
{
foreach (var account in reporting.Accounts.Where(item => item.Enabled))
{
cancellationToken.ThrowIfCancellationRequested();
var scopes = AuthorizedChatScopes(reporting, account.AccountId);
if (scopes.Count == 0 || queue.Pending().Any(batch => string.Equals(batch.AccountId, account.AccountId, StringComparison.Ordinal)
&& string.Equals(batch.StreamKey, "messages", StringComparison.Ordinal)))
continue;
var checkpoint = state.GetLatest(account.AccountId, "messages");
var collection = await collector.CollectAsync(account.AccountId, scopes, checkpoint, cancellationToken).ConfigureAwait(false);
var sourceChanged = !string.Equals(checkpoint.SourceGeneration, collection.SourceGeneration, StringComparison.Ordinal);
var cursorStart = sourceChanged ? new Dictionary<string, long>(StringComparer.Ordinal) : checkpoint.ConfirmedCursors;
var cursorChanged = !cursorStart.OrderBy(item => item.Key).SequenceEqual(collection.NextCursors.OrderBy(item => item.Key));
if (!sourceChanged && !cursorChanged && !collection.HasNewItems && checkpoint.ConfirmedSequence > 0)
continue;
var sequence = sourceChanged ? 1 : checkpoint.ConfirmedSequence + 1;
var batch = new RemoteSyncBatch(
remote.NodeId!, account.AccountId, Guid.NewGuid().ToString("N"), collection.SourceGeneration,
"messages", sequence, SerializeCursors(cursorStart), SerializeCursors(collection.NextCursors),
string.Empty, collection.IsComplete ? "complete" : "partial", collection.Conversations, collection.Messages);
batch = batch with { PayloadHash = RemoteDataBatchAuthorization.ComputePayloadHash(batch) };
var queued = queue.Enqueue(reporting, batch);
if (queued.Accepted && queued.Batch is not null)
state.MarkCollected(queued.Batch, collection.UnavailableSources);
if (!collection.IsComplete)
logger.LogWarning("Data sync is partial; accountId={AccountId}; unavailableSources={UnavailableSourceCount}.", account.AccountId, collection.UnavailableSources.Count);
}
}
private async Task TryCollectOfflineDataBatchesAsync(
RemoteAgentOptions remote,
ReportingConfig? reporting,
IRemoteDataCollector? collector,
RemoteDataBatchQueue? queue,
RemoteDataSyncStateStore? state,
CancellationToken cancellationToken)
{
if (reporting is null || collector is null || queue is null || state is null)
return;
try
{
await CollectDataBatchesAsync(remote, reporting, collector, queue, state, cancellationToken).ConfigureAwait(false);
}
catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested)
{
throw;
}
catch (Exception exception)
{
logger.LogWarning("Offline data collection deferred; type={ExceptionType}.", exception.GetType().Name);
}
}
private async Task ReconcileDataStatusAsync(
RemoteControlClient client,
ReportingConfig reporting,
RemoteDataBatchQueue queue,
RemoteDataSyncStateStore state,
CancellationToken cancellationToken)
{
foreach (var account in reporting.Accounts.Where(item => item.Enabled))
{
var local = state.GetLatest(account.AccountId, "messages");
RemoteDataSyncStatus remote;
try
{
remote = await client.GetDataSyncStatusAsync(account.AccountId, cancellationToken).ConfigureAwait(false);
}
catch (RemoteClientException exception) when (exception.StatusCode == 404 && exception.Code == "NotFound")
{
// Older control planes do not expose the reconciliation endpoint; retain the legacy ACK path.
continue;
}
catch (HttpRequestException)
{
// A disconnected platform must not prevent local DB collection.
return;
}
if (!state.Reconcile(remote, account.AccountId, "messages"))
continue;
var dropped = queue.DropForAccount(account.AccountId, "messages");
logger.LogWarning(
"Data sync checkpoint reconciled with the platform; accountId={AccountId}; localSequence={LocalSequence}; platformSequence={PlatformSequence}; droppedPendingBatches={DroppedPendingBatches}.",
account.AccountId, local.ConfirmedSequence, remote.ConfirmedSequence, dropped);
}
}
private static string SerializeCursors(IReadOnlyDictionary<string, long> cursors) =>
JsonSerializer.Serialize(cursors.OrderBy(item => item.Key).ToDictionary(item => item.Key, item => item.Value), RemoteJson.Options);
private async Task ProcessTaskAsync(
RemoteControlClient client,
RemoteAgentOptions remote,
@@ -178,6 +319,15 @@ public sealed class RemoteAgentHostedService(
}
var renewed = await client.RenewTaskAsync(started, cancellationToken);
if (task.Kind == "sync-data")
{
var result = new RemoteTaskResult(task.TaskId, task.AccountId, renewed.LeaseGeneration,
RemoteTaskStatus.Succeeded, null, "后台数据同步已受理;采集器将在下一轮同步周期执行。", false, null, Guid.NewGuid().ToString("N"));
ledger.Complete(result);
await ReportResultAsync(client, ledger, result, reporting, cancellationToken);
return;
}
if (task.Kind == "send-text")
{
if (!TryReadSendTextPayload(task.Payload, out var targetId, out var text))
@@ -271,7 +421,17 @@ public sealed class RemoteAgentHostedService(
var visibleSessions = sessions
.Where(session => scopes.Any(scope => SessionMatchesScope(session, scope, contacts)))
.ToArray();
var page = visibleSessions.ToPage(payload.Limit, payload.Offset);
var matchedScopeCount = scopes.Count(scope => sessions.Any(session => SessionMatchesScope(session, scope, contacts)));
var coverage = CreateCoverage(
matchedScopeCount == scopes.Count
? ReadCoverageStates.Complete
: matchedScopeCount == 0 ? ReadCoverageStates.Unknown : ReadCoverageStates.Partial,
"uia-session-list",
sessions.Count,
scopes.Count,
matchedScopeCount,
matchedScopeCount == 0 ? "ScopeIdentityNotObserved" : null);
var page = visibleSessions.ToPage(payload.Limit, payload.Offset) with { Coverage = coverage };
return new RemoteReadExecution(JsonSerializer.SerializeToElement(page, RemoteJson.Options), scopes);
}
case "read-contacts":
@@ -294,7 +454,17 @@ public sealed class RemoteAgentHostedService(
}
var allowed = scopes.Select(scope => scope.ChatId).ToHashSet(StringComparer.Ordinal);
var items = all.Where(contact => allowed.Contains(contact.Id)).ToArray();
var resultPage = items.ToPage(payload.Limit, payload.Offset);
var matchedScopeCount = scopes.Count(scope => all.Any(contact => string.Equals(contact.Id, scope.ChatId, StringComparison.Ordinal)));
var coverage = CreateCoverage(
matchedScopeCount == scopes.Count
? ReadCoverageStates.Complete
: matchedScopeCount == 0 ? ReadCoverageStates.Unknown : ReadCoverageStates.Partial,
"uia-contact-list",
all.Count,
scopes.Count,
matchedScopeCount,
matchedScopeCount == 0 ? "ScopeIdentityNotObserved" : null);
var resultPage = items.ToPage(payload.Limit, payload.Offset) with { Coverage = coverage };
return new RemoteReadExecution(JsonSerializer.SerializeToElement(resultPage, RemoteJson.Options), scopes);
}
case "read-messages":
@@ -326,7 +496,8 @@ public sealed class RemoteAgentHostedService(
if (candidateSessions.Length != 1)
throw new ServiceException("ChatIdentityUnconfirmed", 409, "The requested chat identity is not uniquely visible.");
var messages = await backend.MessagesAsync(task.AccountId, candidateSessions[0].AutomationId, payload.IncludeContent, cancellationToken);
var page = messages.ToPage(payload.Limit, payload.Offset);
var coverage = CreateCoverage(ReadCoverageStates.Complete, "uia-session-messages", messages.Count, 1, 1);
var page = messages.ToPage(payload.Limit, payload.Offset) with { Coverage = coverage };
return new RemoteReadExecution(JsonSerializer.SerializeToElement(page, RemoteJson.Options), scopes);
}
default:
@@ -405,6 +576,15 @@ public sealed class RemoteAgentHostedService(
}
}
private static ReadCoverage CreateCoverage(
string state,
string source,
int observedCount,
int authorizedScopeCount,
int matchedScopeCount,
string? errorCode = null) =>
new(state, source, observedCount, authorizedScopeCount, matchedScopeCount, DateTimeOffset.UtcNow, errorCode);
private static void ValidatePage(int limit, int offset)
{
if (limit is < 1 or > 200 || offset < 0)
@@ -527,12 +707,19 @@ public sealed class RemoteAgentHostedService(
.Select(identity =>
{
var account = reporting.FindAccount(identity.AccountId);
return new RemoteAccountSummary(
var summary = new RemoteAccountSummary(
identity.AccountId,
string.Equals(identity.AccountId, snapshot.ActiveAccountId, StringComparison.Ordinal),
identity.Verified,
account?.AllowedChats.Count(chat => chat.Type == ReportingChatType.Group && chat.Enabled && chat.IdentityVerified) ?? 0,
account?.AllowedChats.Count(chat => chat.Type == ReportingChatType.Private && chat.Enabled && chat.IdentityVerified) ?? 0);
return summary with
{
AllowedChats = account?.AllowedChats
.Where(chat => chat.Enabled && chat.IdentityVerified)
.Select(chat => new RemoteAllowedChatSummary(chat.ChatId, chat.Type))
.ToArray() ?? []
};
}).ToArray());
private static RemoteHeartbeat CreateHeartbeat(RemoteAgentOptions remote, ReportingConfig reporting, BackendSnapshot snapshot, string? errorCode) =>