using Microsoft.Extensions.Hosting; using WxAgent.Core; namespace WxAgent.Service; public interface IAgentEventSource { IAsyncEnumerable 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 { await foreach (var item in source.ListenAsync(stoppingToken)) { if (remoteQueue is not null && options.Remote is { IsConfigured: true } remote && item.Kind == "message" && item.ChatId is { Length: > 0 } chatId && item.ChatType is { } chatType) { _ = remoteQueue.Enqueue(options.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); } } } }