48 lines
2.3 KiB
C#
48 lines
2.3 KiB
C#
using Microsoft.Extensions.Hosting;
|
|
using WxAgent.Core;
|
|
|
|
namespace WxAgent.Service;
|
|
|
|
public interface IAgentEventSource
|
|
{
|
|
IAsyncEnumerable<AgentEvent> ListenAsync(CancellationToken cancellationToken);
|
|
}
|
|
|
|
public sealed class EventPump(IAgentBackend backend, EventHub hub, ServiceOptions options, RemoteEventQueue? remoteQueue = null) : BackgroundService
|
|
{
|
|
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
|
{
|
|
if (backend is not IAgentEventSource source || backend.Capabilities.All(c => c.Operation != "listener-events" || !c.Enabled)) return;
|
|
while (!stoppingToken.IsCancellationRequested)
|
|
{
|
|
try
|
|
{
|
|
var remote = options.Remote;
|
|
var reporting = options.Reporting;
|
|
if (remoteQueue is not null && remote is { IsConfigured: true })
|
|
{
|
|
var accounts = await backend.AccountsAsync(stoppingToken, refreshUiIdentity: false);
|
|
reporting = ReportingAuthorization.ForVerifiedAccounts(options.Reporting,
|
|
accounts.Where(RemoteAgentHostedService.HasLiveWeChatBinding).Select(account => account.AccountId));
|
|
}
|
|
|
|
await foreach (var item in source.ListenAsync(stoppingToken))
|
|
{
|
|
if (remoteQueue is not null && remote is { IsConfigured: true } && item.Kind == "message"
|
|
&& item.ChatId is { Length: > 0 } chatId && item.ChatType is { } chatType)
|
|
{
|
|
_ = remoteQueue.Enqueue(reporting, remote.NodeId!, item.AccountId, chatId, chatType,
|
|
item.Kind, item.At, item.Content);
|
|
}
|
|
var localItem = item with { Content = null };
|
|
foreach (var credential in options.ReadCredentials())
|
|
if (credential.AccountIds.Length == 0 || credential.AccountIds.Contains(item.AccountId, StringComparer.Ordinal))
|
|
hub.Publish(credential.PrincipalId, localItem);
|
|
}
|
|
}
|
|
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested) { return; }
|
|
catch (Exception) { await Task.Delay(TimeSpan.FromSeconds(2), stoppingToken); }
|
|
}
|
|
}
|
|
}
|