132 lines
4.6 KiB
C#
132 lines
4.6 KiB
C#
using System.Text.Json.Nodes;
|
|
using Nexus.Api.Repositories;
|
|
|
|
namespace Nexus.Api.Services;
|
|
|
|
/// <summary>
|
|
/// Re-establishes the official OpenClaw session subscriptions after every
|
|
/// Gateway reconnect and subscribes active durable Nexus runs to sanitized
|
|
/// message/tool/approval events.
|
|
/// </summary>
|
|
public sealed class OpenClawEventSubscriptionCoordinator(
|
|
IGatewayConnector connector,
|
|
IServiceScopeFactory scopeFactory,
|
|
ILogger<OpenClawEventSubscriptionCoordinator> logger) : BackgroundService
|
|
{
|
|
private static readonly TimeSpan SyncInterval = TimeSpan.FromSeconds(2);
|
|
private readonly Dictionary<string, OpenClawRunSubscription> _messageSubscriptions =
|
|
new(StringComparer.Ordinal);
|
|
private DateTimeOffset? _connectionEpoch;
|
|
private bool _sessionCatalogSubscribed;
|
|
|
|
protected override async Task ExecuteAsync(CancellationToken stoppingToken)
|
|
{
|
|
while (!stoppingToken.IsCancellationRequested)
|
|
{
|
|
try
|
|
{
|
|
await SynchronizeOnceAsync(stoppingToken);
|
|
}
|
|
catch (OperationCanceledException) when (stoppingToken.IsCancellationRequested)
|
|
{
|
|
break;
|
|
}
|
|
catch (Exception exception)
|
|
{
|
|
logger.LogWarning(
|
|
exception,
|
|
"Could not synchronize OpenClaw event subscriptions; retrying.");
|
|
}
|
|
|
|
await Task.Delay(SyncInterval, stoppingToken);
|
|
}
|
|
}
|
|
|
|
public async Task SynchronizeOnceAsync(CancellationToken cancellationToken = default)
|
|
{
|
|
if (connector.ConnectionState != GatewayConnectionState.Connected)
|
|
{
|
|
ResetConnectionState();
|
|
return;
|
|
}
|
|
|
|
var connectionEpoch = connector.LastConnectedAt;
|
|
if (_connectionEpoch != connectionEpoch)
|
|
{
|
|
ResetConnectionState();
|
|
_connectionEpoch = connectionEpoch;
|
|
}
|
|
|
|
if (!_sessionCatalogSubscribed && connector.Supports("sessions.subscribe"))
|
|
{
|
|
await connector.InvokeAsync(
|
|
"sessions.subscribe",
|
|
new JsonObject(),
|
|
cancellationToken: cancellationToken);
|
|
_sessionCatalogSubscribed = true;
|
|
}
|
|
|
|
if (!connector.Supports("sessions.messages.subscribe"))
|
|
return;
|
|
|
|
await using var scope = scopeFactory.CreateAsyncScope();
|
|
var repository = scope.ServiceProvider.GetRequiredService<IOpenClawRunRepository>();
|
|
var activeRuns = await repository.GetActiveSubscriptionsAsync(cancellationToken);
|
|
var desiredSubscriptions = activeRuns.ToDictionary(
|
|
BuildSubscriptionKey,
|
|
subscription => subscription,
|
|
StringComparer.Ordinal);
|
|
var includeApprovals = connector.GrantedScopes.Contains("operator.admin")
|
|
|| connector.GrantedScopes.Contains("operator.approvals");
|
|
|
|
foreach (var activeRun in activeRuns)
|
|
{
|
|
var subscriptionKey = BuildSubscriptionKey(activeRun);
|
|
if (_messageSubscriptions.ContainsKey(subscriptionKey))
|
|
continue;
|
|
|
|
var parameters = new JsonObject
|
|
{
|
|
["key"] = activeRun.SessionKey,
|
|
["agentId"] = activeRun.AgentId
|
|
};
|
|
if (includeApprovals)
|
|
parameters["includeApprovals"] = true;
|
|
|
|
await connector.InvokeAsync(
|
|
"sessions.messages.subscribe",
|
|
parameters,
|
|
cancellationToken: cancellationToken);
|
|
_messageSubscriptions[subscriptionKey] = activeRun;
|
|
}
|
|
|
|
if (connector.Supports("sessions.messages.unsubscribe"))
|
|
{
|
|
foreach (var obsolete in _messageSubscriptions
|
|
.Where(item => !desiredSubscriptions.ContainsKey(item.Key))
|
|
.ToList())
|
|
{
|
|
await connector.InvokeAsync(
|
|
"sessions.messages.unsubscribe",
|
|
new JsonObject
|
|
{
|
|
["key"] = obsolete.Value.SessionKey,
|
|
["agentId"] = obsolete.Value.AgentId
|
|
},
|
|
cancellationToken: cancellationToken);
|
|
_messageSubscriptions.Remove(obsolete.Key);
|
|
}
|
|
}
|
|
}
|
|
|
|
private void ResetConnectionState()
|
|
{
|
|
_connectionEpoch = null;
|
|
_sessionCatalogSubscribed = false;
|
|
_messageSubscriptions.Clear();
|
|
}
|
|
|
|
private static string BuildSubscriptionKey(OpenClawRunSubscription subscription)
|
|
=> $"{subscription.SessionKey}\u001f{subscription.AgentId}";
|
|
}
|