using System.Globalization;
using System.IdentityModel.Tokens.Jwt;
using System.Security.Claims;
using System.Text;
using System.Text.RegularExpressions;
using System.Text.Json;
using System.Text.Json.Nodes;
using Microsoft.AspNetCore.Http;
using Nexus.Api.Models;
namespace Nexus.Api.Services;
///
/// Stable, browser-safe facade over OpenClaw's protocol-v4 control plane.
/// Gateway payloads are intentionally normalized here so frontend components
/// never depend on OpenClaw's internal storage or raw protocol shapes.
///
public sealed class OpenClawControlService(
IGatewayConnector connector,
IConfiguration configuration,
IOpenClawWriteGate writeGate,
ILogger logger,
IHttpContextAccessor? httpContextAccessor = null,
IOpenClawOperationAuditStore? operationAuditStore = null,
IOpenClawManagementState? managementState = null) : IOpenClawControlService
{
private static readonly string[] ApprovalDecisions = ["allow-once", "allow-always", "deny"];
private static readonly HashSet ModelAuthStatuses =
[
"ok",
"expiring",
"expired",
"missing",
"static"
];
private static readonly HashSet ModelAuthProfileTypes =
[
"oauth",
"token",
"api_key"
];
private static readonly Regex SafeEnvironmentVariablePattern = new(
"^[A-Z][A-Z0-9_]{0,127}$",
RegexOptions.Compiled | RegexOptions.CultureInvariant);
private static readonly Regex SensitiveUsageTextPattern = new(
@"(?ix)
https?://
| \b[\w.+-]+@[\w.-]+\.[a-z]{2,}\b
| \b(?:sk|key|token|secret|password)[-_:=][^\s]+",
RegexOptions.Compiled | RegexOptions.CultureInvariant);
private static readonly HashSet CronMutationMethods =
[
"cron.add",
"cron.update",
"cron.remove",
"cron.run"
];
private static readonly HashSet AllowedCronPatchFields =
[
"name",
"displayName",
"agentId",
"sessionKey",
"description",
"enabled",
"deleteAfterRun",
"schedule",
"trigger",
"sessionTarget",
"wakeMode",
"payload",
"delivery",
"failureAlert"
];
private static readonly Regex DeliveryTargetPattern = new(
@"(?ix)
https?://[^\s]+
| \b[\w.+-]+@[\w.-]+\.[a-z]{2,}\b
| (? GetCapabilities()
{
var connected = connector.ConnectionState == GatewayConnectionState.Connected;
return CapabilityDefinitions
.Select(definition =>
{
var methodAvailable = connected && connector.Supports(definition.Method);
var scopeAvailable = connected && HasScope(definition.RequiredScope);
var managementAvailable =
!CronMutationMethods.Contains(definition.Method) ||
CronManagementEnabled;
var available = methodAvailable && scopeAvailable && managementAvailable;
var state = !connected
? "disconnected"
: !methodAvailable
? "unsupported"
: !scopeAvailable
? "forbidden"
: !managementAvailable
? "management_disabled"
: "ready";
var reason = state switch
{
"disconnected" => "Gateway nicht verbunden.",
"unsupported" => $"Gateway bietet {definition.Method} nicht an.",
"forbidden" => $"Scope {definition.RequiredScope} fehlt.",
"management_disabled" => "OpenClaw-Verwaltung ist in Nexus nicht freigegeben.",
_ => null
};
return new OpenClawCapabilityDto(
definition.Id,
definition.Label,
definition.Method,
definition.RequiredScope,
available,
state,
reason);
})
.ToArray();
}
public async Task GetOverviewAsync(
CancellationToken cancellationToken = default)
{
var tasksTask = GetTasksAsync(100, null, cancellationToken);
var sessionsTask = GetSessionsAsync(100, cancellationToken);
var cronTask = GetCronJobsAsync(100, cancellationToken);
var approvalsTask = GetApprovalsAsync(100, cancellationToken);
var activityTask = GetActivityAsync(100, null, cancellationToken);
var modelsTask = GetModelsAsync(cancellationToken);
var agentsTask = GetAgentsAsync(cancellationToken);
await Task.WhenAll(
tasksTask,
sessionsTask,
cronTask,
approvalsTask,
activityTask,
modelsTask,
agentsTask);
return new OpenClawOverviewDto(
GetConnection(),
GetCapabilities(),
await tasksTask,
await sessionsTask,
await cronTask,
await approvalsTask,
await activityTask,
await modelsTask,
await agentsTask,
DateTimeOffset.UtcNow);
}
public Task> GetTasksAsync(
int limit = 100,
string? cursor = null,
CancellationToken cancellationToken = default)
{
var parameters = new JsonObject
{
["limit"] = Math.Clamp(limit, 1, 500)
};
if (!string.IsNullOrWhiteSpace(cursor))
parameters["cursor"] = cursor.Trim();
return InvokeCollectionAsync(
"tasks.list",
"operator.read",
parameters,
response => ReadItems(response, "tasks", "items")
.Select(MapTask)
.Where(item => !string.IsNullOrWhiteSpace(item.Id))
.ToArray(),
response => ReadString(response, "nextCursor"),
cancellationToken);
}
public Task> GetSessionsAsync(
int limit = 100,
CancellationToken cancellationToken = default)
{
var canAbort = connector.Supports("sessions.abort") && HasScope("operator.write");
return InvokeCollectionAsync(
"sessions.list",
"operator.read",
new { limit = Math.Clamp(limit, 1, 500) },
response => ReadItems(response, "sessions", "items")
.Select(item => MapSession(item, canAbort))
.Where(item => !string.IsNullOrWhiteSpace(item.Key))
.ToArray(),
_ => null,
cancellationToken);
}
public Task> GetCronJobsAsync(
int limit = 100,
CancellationToken cancellationToken = default)
=> GetCronJobsAsync(
includeDisabled: true,
limit,
cursor: null,
cancellationToken);
public Task> GetCronJobsAsync(
bool includeDisabled,
int limit = 100,
string? cursor = null,
CancellationToken cancellationToken = default)
{
if (!TryDecodeCursor(cursor, out var offset))
{
return Task.FromResult(EmptyCollection(
"invalid",
"Der Cron-Cursor ist ungültig.",
"Liste ohne Cursor neu laden."));
}
var canRun =
CronManagementEnabled &&
connector.Supports("cron.run") &&
HasScope("operator.admin");
var boundedLimit = Math.Clamp(limit, 1, 200);
return InvokeCollectionAsync(
"cron.list",
"operator.read",
new
{
includeDisabled,
limit = boundedLimit,
offset
},
response => ReadItems(response, "jobs", "items")
.Select(item => MapCronJob(item, canRun))
.Where(item => !string.IsNullOrWhiteSpace(item.Id))
.ToArray(),
response => ReadOffsetPaginationCursor(
response,
offset,
boundedLimit,
"jobs",
"items"),
cancellationToken);
}
public async Task> GetCronJobAsync(
string jobId,
CancellationToken cancellationToken = default)
{
var normalizedJobId = jobId.Trim();
var unavailable = OperationUnavailable(
"cron.get",
"operator.read");
if (unavailable is not null)
return unavailable;
try
{
var response = await connector.InvokeAsync(
"cron.get",
new { id = normalizedJobId },
cancellationToken: cancellationToken);
var node = ReadNode(response, "job") ?? response;
if (node is null || string.IsNullOrWhiteSpace(ReadString(node, "id", "jobId")))
{
return new OpenClawOperationDto(
false,
"not_found",
"OpenClaw-Zeitplan wurde nicht gefunden.",
null,
"Zeitplanliste aktualisieren und die Job-ID prüfen.",
DateTimeOffset.UtcNow);
}
return new OpenClawOperationDto(
true,
"ready",
"OpenClaw-Zeitplan wurde geladen.",
MapCronJobDetail(node),
null,
DateTimeOffset.UtcNow);
}
catch (Exception exception)
{
logger.LogWarning(exception, "OpenClaw cron get failed for {JobId}", normalizedJobId);
return OperationFailure(exception);
}
}
public Task> GetCronRunsAsync(
string jobId,
int limit = 100,
string? cursor = null,
string? runId = null,
CancellationToken cancellationToken = default)
{
if (!TryDecodeCursor(cursor, out var offset))
{
return Task.FromResult(EmptyCollection(
"invalid",
"Der Cron-History-Cursor ist ungültig.",
"Historie ohne Cursor neu laden."));
}
var boundedLimit = Math.Clamp(limit, 1, 200);
var parameters = new JsonObject
{
["scope"] = "job",
["id"] = jobId.Trim(),
["limit"] = boundedLimit,
["offset"] = offset,
["sortDir"] = "desc"
};
if (!string.IsNullOrWhiteSpace(runId))
parameters["runId"] = runId.Trim();
return InvokeCollectionAsync(
"cron.runs",
"operator.read",
parameters,
response => ReadItems(response, "entries", "runs", "items")
.Select(MapCronRun)
.Where(item => !string.IsNullOrWhiteSpace(item.JobId))
.ToArray(),
response => ReadOffsetPaginationCursor(
response,
offset,
boundedLimit,
"entries",
"runs",
"items"),
cancellationToken);
}
public async Task> GetApprovalsAsync(
int limit = 100,
CancellationToken cancellationToken = default)
{
var disconnected = CollectionUnavailable(
"approval.history",
"operator.approvals");
if (disconnected is not null)
return disconnected;
var supportedMethods = new[]
{
"exec.approval.list",
"plugin.approval.list",
"approval.history"
}.Where(connector.Supports).ToArray();
if (supportedMethods.Length == 0)
{
return EmptyCollection(
"unsupported",
"Diese OpenClaw-Version bietet keine lesbare Approval-Liste an.",
"OpenClaw aktualisieren oder approval.history beziehungsweise *.approval.list aktivieren.");
}
try
{
var canResolve = connector.Supports("approval.resolve") &&
HasScope("operator.approvals");
var calls = supportedMethods.Select(method =>
connector.InvokeAsync(
method,
method == "approval.history"
? new { limit = Math.Clamp(limit, 1, 200) }
: new { },
cancellationToken: cancellationToken)).ToArray();
var responses = await Task.WhenAll(calls);
var unique = new Dictionary(StringComparer.Ordinal);
foreach (var response in responses)
{
foreach (var node in ReadItems(response, "approvals", "items", "history", "pending"))
{
var approval = MapApproval(node, canResolve);
if (!string.IsNullOrWhiteSpace(approval.Id))
unique[approval.Id] = approval;
}
}
var items = unique.Values
.OrderBy(item => string.Equals(item.Status, "pending", StringComparison.OrdinalIgnoreCase) ? 0 : 1)
.ThenByDescending(item => item.RequestedAt)
.Take(Math.Clamp(limit, 1, 200))
.ToArray();
return ReadyCollection(items);
}
catch (Exception exception)
{
logger.LogWarning(exception, "OpenClaw approval list failed");
return CollectionFailure(exception);
}
}
public async Task> GetActivityAsync(
int limit = 100,
string? cursor = null,
CancellationToken cancellationToken = default)
{
if (connector.ConnectionState != GatewayConnectionState.Connected)
return CollectionFailure(
new OpenClawGatewayRpcException(
"GATEWAY_DISCONNECTED",
"OpenClaw Gateway is not connected.",
retryable: true));
if (!connector.Supports("audit.activity.list"))
{
var eventItems = connector.GetRecentEvents(limit)
.Select(MapGatewayEvent)
.ToArray();
return new OpenClawCollectionDto(
eventItems.Length == 0 ? "unsupported" : "ready",
eventItems,
null,
eventItems.Length == 0
? "OpenClaw bietet audit.activity.list nicht an."
: "Audit-Ledger nicht verfügbar; echte Gateway-Events werden angezeigt.",
eventItems.Length == 0
? "OpenClaw aktualisieren oder Audit-Ledger aktivieren."
: null,
DateTimeOffset.UtcNow);
}
var parameters = new JsonObject
{
["limit"] = Math.Clamp(limit, 1, 500)
};
if (!string.IsNullOrWhiteSpace(cursor))
parameters["cursor"] = cursor.Trim();
return await InvokeCollectionAsync(
"audit.activity.list",
"operator.read",
parameters,
response => ReadItems(response, "events", "items")
.Select(MapActivity)
.ToArray(),
response => ReadString(response, "nextCursor"),
cancellationToken);
}
public Task> GetModelsAsync(
CancellationToken cancellationToken = default)
{
return InvokeCollectionAsync(
"models.list",
"operator.read",
new { view = "configured" },
response => ReadItems(response, "models", "items")
.Select(MapModel)
.Where(item => !string.IsNullOrWhiteSpace(item.Id))
.ToArray(),
_ => null,
cancellationToken);
}
public Task> GetModelAuthStatusAsync(
bool refresh = false,
CancellationToken cancellationToken = default)
{
return InvokeCollectionAsync(
"models.authStatus",
"operator.read",
new { refresh },
response => ReadItems(response, "providers", "items")
.Select(MapModelAuthProvider)
.Where(item => !string.IsNullOrWhiteSpace(item.Provider))
.OrderBy(item => item.DisplayName, StringComparer.OrdinalIgnoreCase)
.ToArray(),
_ => null,
cancellationToken);
}
public Task> GetAgentsAsync(
CancellationToken cancellationToken = default)
{
return InvokeCollectionAsync(
"agents.list",
"operator.read",
new { },
response => ReadItems(response, "agents", "items")
.Select(MapAgent)
.Where(item => !string.IsNullOrWhiteSpace(item.Id))
.ToArray(),
_ => null,
cancellationToken);
}
public async Task> CancelTaskAsync(
string taskId,
string? reason,
CancellationToken cancellationToken = default,
OpenClawInvocationContext? invocationContext = null)
{
var normalizedTaskId = taskId.Trim();
var taskRef = new EntityRefDto("openclaw-task", normalizedTaskId);
var writeUnavailable =
await WriteMutationUnavailableAsync(
"tasks.cancel",
cancellationToken,
"operator.write");
if (writeUnavailable is not null)
return EnrichUnavailableOrInvalid(
writeUnavailable,
invocationContext,
taskRef);
var normalizedReason = string.IsNullOrWhiteSpace(reason) ? null : reason.Trim();
var unavailable = OperationUnavailable("tasks.cancel", "operator.write");
if (unavailable is not null)
return EnrichUnavailableOrInvalid(unavailable, invocationContext, taskRef);
return await ExecuteMutationAsync(
"tasks.cancel",
"task",
normalizedTaskId,
[normalizedReason ?? string.Empty],
invocationContext,
cancellationToken,
async context =>
{
try
{
var parameters = new JsonObject { ["taskId"] = normalizedTaskId };
if (normalizedReason is not null)
parameters["reason"] = normalizedReason;
var response = await connector.InvokeAsync(
"tasks.cancel",
parameters,
cancellationToken: cancellationToken,
invocationContext: context with { IncludeIdempotencyParameter = false });
var found = ReadBool(response, "found") ?? true;
var cancelled = ReadBool(response, "cancelled") ?? false;
var taskNode = ReadNode(response, "task");
var task = taskNode is null ? null : MapTask(taskNode);
var message = !found
? "OpenClaw-Task wurde nicht gefunden."
: cancelled
? "OpenClaw-Task wurde abgebrochen."
: ReadString(response, "reason") ?? "OpenClaw konnte den Task nicht abbrechen.";
return new OpenClawOperationDto(
found && cancelled,
found && cancelled ? "completed" : "rejected",
message,
task,
found ? null : "Task-Liste aktualisieren und die Task-ID prüfen.",
DateTimeOffset.UtcNow);
}
catch (Exception exception)
{
logger.LogWarning(exception, "OpenClaw task cancellation failed for {TaskId}", taskId);
return OperationFailure(exception);
}
});
}
public async Task> AbortSessionAsync(
string sessionKey,
string? runId,
bool clearQueued,
CancellationToken cancellationToken = default,
OpenClawInvocationContext? invocationContext = null)
{
var normalizedSessionKey = sessionKey.Trim();
var sessionRef = new EntityRefDto("session", normalizedSessionKey);
var writeUnavailable =
await WriteMutationUnavailableAsync