diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 6113ef876..95801af97 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -394,10 +394,10 @@ jobs: (github.event_name == 'push' && (github.ref == 'refs/heads/main' || github.ref == 'refs/heads/dev')) || needs.changes.outputs.docs_only == 'false' runs-on: ubuntu-latest - # Full coverage takes 34-35 minutes after a roughly 10-minute restore/build, - # then still needs time to upload the report. Keep the fail-fast guard while - # leaving enough headroom for runner variance and cleanup. - timeout-minutes: 55 + # Full coverage takes roughly 40 minutes after restore/build, + # then still needs time to generate and upload the report. Keep the fail-fast + # guard while leaving enough headroom for runner variance and cleanup. + timeout-minutes: 75 services: agent-tool-admission-redis: image: redis:7.2.3 diff --git a/docs/contracts/nyxid-assistant-conformance/v1/sources.json b/docs/contracts/nyxid-assistant-conformance/v1/sources.json index 2b2a1f536..bed374e9f 100644 --- a/docs/contracts/nyxid-assistant-conformance/v1/sources.json +++ b/docs/contracts/nyxid-assistant-conformance/v1/sources.json @@ -2,8 +2,8 @@ "schema_version": 1, "aevatar": { "repository": "https://github.com/AevatarAI/aevatar.git", - "revision": "d0d49a256ed15d24457649703b97618b6b64891c", - "contract_files_sha256": "09cf89aedd527a9a34a950708f13bfd4f1f40bd51ad200f92e40c5065b211720", + "revision": "20d1cd104ce24fc6942dfaf4bc23b6bc40ad5432", + "contract_files_sha256": "7668e545af09bab23c99f9f698f9cbae032ec7d6ddee56dd5b5551cfb6b3674f", "files": { "agents/Aevatar.GAgents.NyxidChat/NyxIdActionPostconditionPort.cs": "7791de469b567dcde70a0f8e2a88cc818972ca557617a2538294e8ccabd5bda0", "agents/Aevatar.GAgents.NyxidChat/NyxIdAssistantActionRegistry.cs": "60e6f67c94ae11b1bf0dac036ad8ac0c35901e31787b1f0c8173964f6a12d263", @@ -12,12 +12,12 @@ "agents/Aevatar.GAgents.NyxidChat/protos/nyxid_chat_recovery_secret.proto": "07dbc449a732df6c7a6d0a97054ddbe35a517f4af4c5670b05ed0bea0bf2011a", "agents/Aevatar.GAgents.NyxidChat/protos/nyxid_chat_task.proto": "5d698e6d75b90605eb40092899a775ecd45a455aaa17656ccebaf247fc71d5c4", "docs/adr/0048-nyxid-assistant-operation-class-boundary.md": "884aca09774e773e68154c923fec8078610b2cf8e97f581fedc36e10451ccec3", - "src/Aevatar.AI.Abstractions/ai_messages.proto": "50e334e9fdbc1c11e0f70345095de2c7d9f3ea5f4b9f84a30f729377621576e2", + "src/Aevatar.AI.Abstractions/ai_messages.proto": "6831b61c1016d1584654c0edf61788e9c448a7260449936e0e121fdfa789475e", "src/Aevatar.AI.ToolProviders.NyxId/NyxIdApiAccessContracts.cs": "a2e526a0a227f868304f122e9a65796164089fa69b0c27926f4c27c78b3ab0de", "src/Aevatar.AI.ToolProviders.NyxId/NyxIdAssistantToolSource.cs": "16b25f5bdd5004bc0c5402ae0adf2d294c270ae83135c95f0f342089ea80da05", "src/Aevatar.AI.ToolProviders.NyxId/Tools/NyxIdRequestKeyCreateTool.cs": "2c4f2cda99154f2e667c6cfd291497e697ef11df17f081f96ec70070a8af8b8c", "src/Aevatar.AI.ToolProviders.NyxId/Tools/NyxIdRequestKeyRotateTool.cs": "18212bb64644cfbca401065bccce439ea5fa00316deff57d730a0d9ac2650e53", - "src/Aevatar.Mainnet.Host.Api/Hosting/MainnetHostBuilderExtensions.cs": "88cf9ee20ea024c84768316850f92262003b7a55a2a853066445257bdd7a5d87" + "src/Aevatar.Mainnet.Host.Api/Hosting/MainnetHostBuilderExtensions.cs": "8baedc2ee9ea022c0a8da8c5b2e0c9f50dcaaf983ebc12f288a335a6213a793f" } }, "nyxid": { diff --git a/src/platform/Aevatar.GAgentService.Hosting/Endpoints/ScopeWorkflowScheduleEndpoints.cs b/src/platform/Aevatar.GAgentService.Hosting/Endpoints/ScopeWorkflowScheduleEndpoints.cs index d4cbed7cd..52e8509eb 100644 --- a/src/platform/Aevatar.GAgentService.Hosting/Endpoints/ScopeWorkflowScheduleEndpoints.cs +++ b/src/platform/Aevatar.GAgentService.Hosting/Endpoints/ScopeWorkflowScheduleEndpoints.cs @@ -8,6 +8,7 @@ using Aevatar.GAgentService.Abstractions.Schedules; using Aevatar.GAgentService.Abstractions.Services; using Aevatar.GAgentService.Hosting.Endpoints.Schedules; +using Aevatar.Workflow.Infrastructure.CapabilityApi; using Google.Protobuf.WellKnownTypes; using Microsoft.AspNetCore.Builder; using Microsoft.AspNetCore.Http; @@ -20,9 +21,30 @@ internal static class ScopeWorkflowScheduleEndpoints { private const string ChatEndpointId = "chat"; private const string DefaultWorkflowScheduleNyxIdScope = "proxy"; + private const string DefaultExternalTriggerDeliveryIdJsonPath = "event_id"; + private const string DefaultExternalTriggerDeliveryIdHeader = "X-NyxID-Delivery-Id"; + private const string DefaultExternalTriggerHmacSignatureHeader = "X-NyxID-Signature"; + private const string DefaultExternalTriggerHmacTimestampHeader = "X-NyxID-Timestamp"; public static RouteGroupBuilder MapScopeWorkflowScheduleEndpoints(this RouteGroupBuilder group) { + group.MapPost("/{scopeId}/workflows/{workflowId}/external-trigger", UpsertExternalTrigger) + .Produces(StatusCodes.Status200OK) + .Produces(StatusCodes.Status400BadRequest) + .Produces(StatusCodes.Status403Forbidden) + .Produces(StatusCodes.Status404NotFound) + .Produces(StatusCodes.Status409Conflict); + group.MapGet("/{scopeId}/workflows/{workflowId}/external-trigger", GetExternalTrigger) + .Produces(StatusCodes.Status200OK) + .Produces(StatusCodes.Status400BadRequest) + .Produces(StatusCodes.Status403Forbidden) + .Produces(StatusCodes.Status404NotFound); + group.MapDelete("/{scopeId}/workflows/{workflowId}/external-trigger", DeleteExternalTrigger) + .Produces(StatusCodes.Status204NoContent) + .Produces(StatusCodes.Status400BadRequest) + .Produces(StatusCodes.Status403Forbidden) + .Produces(StatusCodes.Status404NotFound) + .Produces(StatusCodes.Status409Conflict); group.MapGet("/{scopeId}/workflows/{workflowId}/schedules", List) .Produces(StatusCodes.Status200OK) .Produces(StatusCodes.Status400BadRequest) @@ -80,6 +102,80 @@ public static RouteGroupBuilder MapScopeWorkflowScheduleEndpoints(this RouteGrou return group; } + internal static async Task UpsertExternalTrigger( + HttpContext http, + string scopeId, + string workflowId, + WorkflowExternalTriggerConfigurationHttpRequest input, + [FromServices] IScopeWorkflowQueryPort workflowQueryPort, + CancellationToken ct = default) + { + var resolved = await ResolveWorkflowAsync(http, scopeId, workflowId, workflowQueryPort, ct); + if (resolved.Result != null) + return resolved.Result; + + var routeKey = await ResolveExternalTriggerRouteKeyAsync(http, resolved.Workflow!, ct) + ?? GenerateExternalTriggerRouteKey(); + var put = await WorkflowWebhookBindingEndpoints.HandlePutAsync( + http, + scopeId, + routeKey, + BuildExternalTriggerBindingRequest(resolved.Workflow!, input), + ct); + if (!IsSuccessStatus(put)) + return put; + + var record = await ResolveExternalTriggerBindingAsync(http, resolved.Workflow!, routeKey, ct); + return record == null + ? Results.Json( + new + { + code = "WORKFLOW_EXTERNAL_TRIGGER_BINDING_NOT_FOUND", + message = "External trigger binding was not found after upsert.", + }, + statusCode: StatusCodes.Status409Conflict) + : Results.Ok(WorkflowExternalTriggerHttpResult.FromBinding(resolved.Workflow!, record)); + } + + internal static async Task GetExternalTrigger( + HttpContext http, + string scopeId, + string workflowId, + [FromServices] IScopeWorkflowQueryPort workflowQueryPort, + CancellationToken ct = default) + { + var resolved = await ResolveWorkflowAsync(http, scopeId, workflowId, workflowQueryPort, ct); + if (resolved.Result != null) + return resolved.Result; + + var record = await ResolveExternalTriggerBindingAsync(http, resolved.Workflow!, routeKey: null, ct); + return Results.Ok(record == null + ? WorkflowExternalTriggerHttpResult.NotConfigured(resolved.Workflow!) + : WorkflowExternalTriggerHttpResult.FromBinding(resolved.Workflow!, record)); + } + + internal static async Task DeleteExternalTrigger( + HttpContext http, + string scopeId, + string workflowId, + [FromServices] IScopeWorkflowQueryPort workflowQueryPort, + CancellationToken ct = default) + { + var resolved = await ResolveWorkflowAsync(http, scopeId, workflowId, workflowQueryPort, ct); + if (resolved.Result != null) + return resolved.Result; + + var record = await ResolveExternalTriggerBindingAsync(http, resolved.Workflow!, routeKey: null, ct); + if (record == null) + return Results.NotFound(); + + return await WorkflowWebhookBindingEndpoints.HandleDeleteAsync( + http, + scopeId, + record.RouteKey, + ct); + } + internal static async Task Create( HttpContext http, string scopeId, @@ -472,6 +568,7 @@ private static bool BelongsToWorkflow(ScheduledDispatchSummary schedule, ScopeWo if (schedule.ScheduleKind != ScheduledDispatchScheduleKind.Workflow || schedule.TargetKind != ScheduledDispatchTargetKind.ServiceInvocation || !string.Equals(schedule.ServiceEndpointId, ChatEndpointId, StringComparison.Ordinal) || + !string.Equals(schedule.ServiceIdentity.TenantId, workflow.ScopeId, StringComparison.Ordinal) || !string.Equals(schedule.ServiceId, workflow.PublishedServiceId, StringComparison.Ordinal)) { return false; @@ -481,6 +578,71 @@ private static bool BelongsToWorkflow(ScheduledDispatchSummary schedule, ScopeWo string.Equals(schedule.ServiceKey, workflow.ServiceKey, StringComparison.Ordinal); } + private static WorkflowWebhookBindingEndpoints.PutWorkflowWebhookBindingRequest BuildExternalTriggerBindingRequest( + ScopeWorkflowSummary workflow, + WorkflowExternalTriggerConfigurationHttpRequest input) => + new( + WorkflowName: workflow.WorkflowName, + SourceId: input.SourceId, + PromptTemplate: input.PromptTemplate, + PromptJsonPath: input.PromptJsonPath, + DeliveryIdHeader: input.DeliveryIdHeader ?? DefaultExternalTriggerDeliveryIdHeader, + DeliveryIdJsonPath: input.DeliveryIdJsonPath ?? DefaultExternalTriggerDeliveryIdJsonPath, + HmacSecret: input.HmacSecret, + HmacSignatureHeader: input.HmacSignatureHeader ?? DefaultExternalTriggerHmacSignatureHeader, + HmacTimestampHeader: input.HmacTimestampHeader ?? DefaultExternalTriggerHmacTimestampHeader, + MaxTimestampSkewSeconds: input.MaxTimestampSkewSeconds, + DefinitionActorId: workflow.ActorId, + TargetRevisionId: workflow.ActiveRevisionId, + PreviousHmacSecret: input.PreviousHmacSecret, + TimeZoneId: input.TimeZoneId, + EnableUnattendedEffects: input.EnableUnattendedEffects); + + private static async Task ResolveExternalTriggerRouteKeyAsync( + HttpContext http, + ScopeWorkflowSummary workflow, + CancellationToken ct) + { + var record = await ResolveExternalTriggerBindingAsync(http, workflow, routeKey: null, ct); + return record?.RouteKey; + } + + private static async Task ResolveExternalTriggerBindingAsync( + HttpContext http, + ScopeWorkflowSummary workflow, + string? routeKey, + CancellationToken ct) + { + var store = http.RequestServices.GetService(typeof(IWorkflowWebhookBindingStore)) as IWorkflowWebhookBindingStore; + if (store == null) + return null; + + if (!string.IsNullOrWhiteSpace(routeKey)) + { + var record = await store.GetAsync(routeKey.Trim(), ct); + return IsExternalTriggerBinding(record, workflow) ? record : null; + } + + var records = await store.ListByScopeAsync(workflow.ScopeId, ct); + return records.FirstOrDefault(record => IsExternalTriggerBinding(record, workflow)); + } + + private static bool IsExternalTriggerBinding( + WorkflowWebhookBindingRecord? record, + ScopeWorkflowSummary workflow) => + record != null && + string.Equals(record.ScopeId, workflow.ScopeId, StringComparison.Ordinal) && + string.Equals(record.DefinitionActorId, workflow.ActorId, StringComparison.Ordinal) && + string.Equals(record.TargetRevisionId, workflow.ActiveRevisionId, StringComparison.Ordinal) && + string.Equals(record.WorkflowName, workflow.WorkflowName, StringComparison.OrdinalIgnoreCase); + + private static string GenerateExternalTriggerRouteKey() => + $"workflow-external-trigger-{Guid.NewGuid():N}"; + + private static bool IsSuccessStatus(IResult result) => + result is IStatusCodeHttpResult { StatusCode: >= 200 and < 300 or null } || + result is not IStatusCodeHttpResult; + private static ScheduledDispatchConfiguration BuildConfiguration( ScopeWorkflowSummary workflow, WorkflowScheduleConfigurationHttpRequest input, @@ -602,6 +764,11 @@ private static bool TryCreateInvalidScheduleIdResult(string? scheduleId, out IRe return false; } + private static string NormalizeRequired(string? value, string paramName) => + string.IsNullOrWhiteSpace(value) + ? throw new ArgumentException($"{paramName} is required.", paramName) + : value.Trim(); + private static ScheduledServiceInvocationNyxIdSubjectRef? ResolveAuthenticatedNyxIdOwnerSubject(HttpContext http) { var ownerUserId = ReadFirstClaim( @@ -646,6 +813,85 @@ private sealed record WorkflowScheduleOwnershipResult( IResult? Result); } +[JsonUnmappedMemberHandling(JsonUnmappedMemberHandling.Disallow)] +public sealed record WorkflowExternalTriggerConfigurationHttpRequest +{ + public string? SourceId { get; init; } + public string? PromptTemplate { get; init; } + public string? PromptJsonPath { get; init; } + public string? DeliveryIdHeader { get; init; } + public string? DeliveryIdJsonPath { get; init; } + public string? HmacSecret { get; init; } + public string? PreviousHmacSecret { get; init; } + public string? HmacSignatureHeader { get; init; } + public string? HmacTimestampHeader { get; init; } + public int? MaxTimestampSkewSeconds { get; init; } + public string? TimeZoneId { get; init; } + public bool EnableUnattendedEffects { get; init; } = true; +} + +public sealed record WorkflowExternalTriggerHttpResult +{ + public bool Configured { get; init; } + public required string WorkflowId { get; init; } + public required string ServiceId { get; init; } + public string RevisionId { get; init; } = string.Empty; + public string Status { get; init; } = string.Empty; + public string RouteKey { get; init; } = string.Empty; + public string FireUrl { get; init; } = string.Empty; + public string SourceId { get; init; } = string.Empty; + public string DefinitionActorId { get; init; } = string.Empty; + public string TargetRevisionId { get; init; } = string.Empty; + public string DeliveryIdHeader { get; init; } = string.Empty; + public string DeliveryIdJsonPath { get; init; } = string.Empty; + public string HmacSignatureHeader { get; init; } = string.Empty; + public string HmacTimestampHeader { get; init; } = string.Empty; + public bool HmacSecretSet { get; init; } + public bool PreviousHmacSecretSet { get; init; } + public bool AgentKeyReady { get; init; } + public bool UnattendedEffectsEnabled { get; init; } + public long UpdatedAtUnixMs { get; init; } + + public static WorkflowExternalTriggerHttpResult FromBinding( + ScopeWorkflowSummary workflow, + WorkflowWebhookBindingRecord record) => + new() + { + Configured = true, + WorkflowId = workflow.WorkflowId, + ServiceId = workflow.PublishedServiceId, + RevisionId = record.TargetRevisionId ?? workflow.ActiveRevisionId, + Status = record.CallerDurableCredential != null ? "active" : "configured", + RouteKey = record.RouteKey, + FireUrl = BuildWebhookFireUrl(record.RouteKey), + SourceId = record.SourceId ?? string.Empty, + DefinitionActorId = record.DefinitionActorId ?? string.Empty, + TargetRevisionId = record.TargetRevisionId ?? string.Empty, + DeliveryIdHeader = record.DeliveryIdHeader ?? string.Empty, + DeliveryIdJsonPath = record.DeliveryIdJsonPath ?? string.Empty, + HmacSignatureHeader = record.HmacSignatureHeader ?? string.Empty, + HmacTimestampHeader = record.HmacTimestampHeader ?? string.Empty, + HmacSecretSet = !string.IsNullOrWhiteSpace(record.HmacSecret), + PreviousHmacSecretSet = !string.IsNullOrWhiteSpace(record.PreviousHmacSecret), + AgentKeyReady = record.CallerDurableCredential != null, + UnattendedEffectsEnabled = record.CallerAuthority != null && record.UnattendedEffectAuthorization != null, + UpdatedAtUnixMs = record.UpdatedAtUnixMs, + }; + + public static WorkflowExternalTriggerHttpResult NotConfigured(ScopeWorkflowSummary workflow) => + new() + { + Configured = false, + WorkflowId = workflow.WorkflowId, + ServiceId = workflow.PublishedServiceId, + RevisionId = workflow.ActiveRevisionId, + Status = "not_configured", + }; + + private static string BuildWebhookFireUrl(string routeKey) => + $"/api/workflow-webhooks/{Uri.EscapeDataString(routeKey)}"; +} + [JsonUnmappedMemberHandling(JsonUnmappedMemberHandling.Disallow)] public sealed record WorkflowScheduleConfigurationHttpRequest { diff --git a/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/WorkflowWebhookAgentKeyMaterializer.cs b/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/WorkflowWebhookAgentKeyMaterializer.cs index e753e8ab5..336c6db32 100644 --- a/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/WorkflowWebhookAgentKeyMaterializer.cs +++ b/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/WorkflowWebhookAgentKeyMaterializer.cs @@ -11,7 +11,7 @@ namespace Aevatar.Workflow.Infrastructure.CapabilityApi; -internal interface IWorkflowWebhookAgentKeyMaterializer +public interface IWorkflowWebhookAgentKeyMaterializer { Task MaterializeAsync( WorkflowCallerNyxIdAuthority callerAuthority, @@ -27,7 +27,7 @@ Task RevokeAsync( CancellationToken ct); } -internal sealed record WorkflowWebhookAgentKeyMaterializationResult( +public sealed record WorkflowWebhookAgentKeyMaterializationResult( DurableCallerCredentialRef? Credential, int StatusCode, string ErrorCode) diff --git a/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/WorkflowWebhookBindingEndpoints.cs b/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/WorkflowWebhookBindingEndpoints.cs index 13ccd21cc..c5426edce 100644 --- a/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/WorkflowWebhookBindingEndpoints.cs +++ b/src/workflow/Aevatar.Workflow.Infrastructure/CapabilityApi/WorkflowWebhookBindingEndpoints.cs @@ -22,7 +22,7 @@ namespace Aevatar.Workflow.Infrastructure.CapabilityApi; /// ingress at /api/workflow-webhooks/{routeKey} resolves it dynamically — /// no host configuration change or redeploy per workflow. /// -internal static class WorkflowWebhookBindingEndpoints +public static class WorkflowWebhookBindingEndpoints { public static void Map(IEndpointRouteBuilder group) { @@ -51,7 +51,7 @@ public sealed record PutWorkflowWebhookBindingRequest( string? TimeZoneId = null, bool EnableUnattendedEffects = false); - internal static async Task HandlePutAsync( + public static async Task HandlePutAsync( HttpContext http, string scopeId, string routeKey, @@ -359,7 +359,7 @@ internal static async Task HandleListAsync( return Results.Ok(new { bindings = records.Select(ToView).ToArray() }); } - internal static async Task HandleDeleteAsync( + public static async Task HandleDeleteAsync( HttpContext http, string scopeId, string routeKey, diff --git a/test/Aevatar.GAgentService.Integration.Tests/ScopeWorkflowEndpointsTests.cs b/test/Aevatar.GAgentService.Integration.Tests/ScopeWorkflowEndpointsTests.cs index b490adfc8..9e94d16cb 100644 --- a/test/Aevatar.GAgentService.Integration.Tests/ScopeWorkflowEndpointsTests.cs +++ b/test/Aevatar.GAgentService.Integration.Tests/ScopeWorkflowEndpointsTests.cs @@ -1,3 +1,5 @@ +using System.Reflection; +using System.Security.Cryptography; using System.Text; using Aevatar.AI.Abstractions; using Aevatar.CQRS.Core.Abstractions.Interactions; @@ -6,20 +8,30 @@ using Aevatar.GAgentService.Abstractions.Ports; using Aevatar.GAgentService.Abstractions.Queries; using Aevatar.GAgentService.Abstractions.Schedules; +using Aevatar.GAgentService.Abstractions.Schedules.Authorization; using Aevatar.GAgentService.Abstractions.Services; using Aevatar.GAgentService.Application.Workflows; using Aevatar.GAgentService.Governance.Abstractions; using Aevatar.GAgentService.Governance.Abstractions.Ports; using Aevatar.GAgentService.Governance.Abstractions.Queries; +using Aevatar.GAgents.Channel.Abstractions; +using Aevatar.GAgents.Channel.Identity.Abstractions; using Aevatar.GAgentService.Hosting.Endpoints; +using Aevatar.GAgentService.Hosting.Endpoints.Schedules; using Aevatar.Studio.Application; +using Aevatar.Studio.Application.Provisioning; using Aevatar.Studio.Application.Studio.Abstractions; using Aevatar.Studio.Application.Studio.Contracts; +using Aevatar.Foundation.Abstractions; using Aevatar.Foundation.Abstractions.Connectors; +using Aevatar.Foundation.Abstractions.Credentials; using Aevatar.CQRS.Core.Abstractions.Commands; using Aevatar.Workflow.Abstractions; +using Aevatar.Workflow.Abstractions.Credentials; using Aevatar.Workflow.Application.Abstractions.ExternalCapabilities; +using CredentialWorkflowCallerAuthority = Aevatar.Workflow.Abstractions.WorkflowCallerNyxIdAuthority; using Aevatar.Workflow.Application.Abstractions.Runs; +using Aevatar.Workflow.Infrastructure.CapabilityApi; using Google.Protobuf.WellKnownTypes; using Microsoft.AspNetCore.Http; using Microsoft.Extensions.Configuration; @@ -1667,6 +1679,211 @@ public async Task HandleUpsertWorkflowAsync_ShouldReturnAccepted_WithLocation_Wh body.Should().NotContain("\"workflow\""); } + [Fact] + public async Task WorkflowExternalTriggerUpsert_ShouldCreateWebhookBindingWithAgentKeyAndFireUrl() + { + var bindingStore = new RecordingWorkflowWebhookBindingStore(); + var materializer = new RecordingWorkflowWebhookAgentKeyMaterializer(); + var tokenProvider = new RecordingWorkflowCallerAccessTokenProvider(); + var http = CreateExternalTriggerHttpContext(bindingStore, materializer, tokenProvider); + http.Request.Headers["X-NyxID-Delegation-Token"] = "proxy-delegation-token"; + var workflowQueryPort = new RecordingScopeWorkflowQueryPort + { + LookupResult = RunnableWorkflow(), + }; + var input = new WorkflowExternalTriggerConfigurationHttpRequest + { + SourceId = "nyxid-external-trigger", + PromptTemplate = "\"Run {{event_id}}\"", + HmacSecret = "delivery-signing-secret-at-least-32-bytes", + }; + + var result = await ScopeWorkflowScheduleEndpoints.UpsertExternalTrigger( + http, + "scope-alpha", + "wf-alpha", + input, + workflowQueryPort, + CancellationToken.None); + + await result.ExecuteAsync(http); + var body = await ReadBodyAsync(http.Response); + + http.Response.StatusCode.Should().Be(StatusCodes.Status200OK); + body.Should().Contain("\"configured\":true"); + body.Should().Contain("\"status\":\"active\""); + body.Should().Contain("\"fireUrl\":\"/api/workflow-webhooks/workflow-external-trigger-"); + body.Should().Contain("\"routeKey\":\"workflow-external-trigger-"); + body.Should().Contain("\"agentKeyReady\":true"); + body.Should().Contain("\"hmacSecretSet\":true"); + body.Should().NotContain("delivery-signing-secret-at-least-32-bytes"); + body.Should().NotContain("workflow-triggers"); + bindingStore.Records.Should().ContainSingle(); + var record = bindingStore.Records.Values.Single(); + record.ScopeId.Should().Be("scope-alpha"); + record.WorkflowName.Should().Be("workflow-alpha"); + record.DefinitionActorId.Should().Be("definition-actor-alpha"); + record.TargetRevisionId.Should().Be("rev-alpha"); + record.CallerDurableCredential.Should().NotBeNull(); + tokenProvider.Authorities.Should().ContainSingle().Which.BindingId.Should().Be("binding-caller-alpha"); + materializer.Materialized.Should().ContainSingle().Which.RouteKey.Should().Be(record.RouteKey); + + var secondHttp = CreateExternalTriggerHttpContext(bindingStore, materializer, tokenProvider); + secondHttp.Request.Headers["X-NyxID-Delegation-Token"] = "proxy-delegation-token"; + var secondResult = await ScopeWorkflowScheduleEndpoints.UpsertExternalTrigger( + secondHttp, + "scope-alpha", + "wf-alpha", + input, + workflowQueryPort, + CancellationToken.None); + + await secondResult.ExecuteAsync(secondHttp); + var secondBody = await ReadBodyAsync(secondHttp.Response); + + secondHttp.Response.StatusCode.Should().Be(StatusCodes.Status200OK); + secondBody.Should().Contain($"\"routeKey\":\"{record.RouteKey}\""); + bindingStore.Records.Should().ContainSingle(); + materializer.Revoked.Should().ContainSingle().Which.Credential.Ref.Should().Be("webhook-agent-key-ref-1"); + } + + [Fact] + public async Task WorkflowExternalTriggerUpsert_ShouldRejectForeignRouteBindingWithoutMutation() + { + var bindingStore = new RecordingWorkflowWebhookBindingStore(); + bindingStore.Records["workflow-external-trigger-foreign"] = ExternalTriggerBindingRecord( + routeKey: "workflow-external-trigger-foreign", + scopeId: "scope-other"); + var http = CreateExternalTriggerHttpContext(bindingStore); + var workflowQueryPort = new RecordingScopeWorkflowQueryPort + { + LookupResult = RunnableWorkflow(), + }; + + var result = await ScopeWorkflowScheduleEndpoints.UpsertExternalTrigger( + http, + "scope-alpha", + "wf-alpha", + new WorkflowExternalTriggerConfigurationHttpRequest + { + PromptTemplate = "\"Run {{event_id}}\"", + HmacSecret = "delivery-signing-secret-at-least-32-bytes", + EnableUnattendedEffects = false, + }, + workflowQueryPort, + CancellationToken.None); + + await result.ExecuteAsync(http); + + http.Response.StatusCode.Should().Be(StatusCodes.Status200OK); + bindingStore.Records.Should().HaveCount(2); + bindingStore.Records["workflow-external-trigger-foreign"].ScopeId.Should().Be("scope-other"); + } + + [Fact] + public async Task WorkflowExternalTriggerGet_ShouldReturnNotConfiguredWhenBindingIsAbsent() + { + var bindingStore = new RecordingWorkflowWebhookBindingStore(); + var http = CreateExternalTriggerHttpContext(bindingStore); + var workflowQueryPort = new RecordingScopeWorkflowQueryPort + { + LookupResult = RunnableWorkflow(), + }; + + var result = await ScopeWorkflowScheduleEndpoints.GetExternalTrigger( + http, + "scope-alpha", + "wf-alpha", + workflowQueryPort, + CancellationToken.None); + + await result.ExecuteAsync(http); + var body = await ReadBodyAsync(http.Response); + + http.Response.StatusCode.Should().Be(StatusCodes.Status200OK); + body.Should().Contain("\"configured\":false"); + body.Should().Contain("\"status\":\"not_configured\""); + body.Should().Contain("\"fireUrl\":\"\""); + body.Should().NotContain("workflow-triggers"); + } + + [Fact] + public async Task WorkflowExternalTriggerGet_ShouldReturnConfiguredWebhookBinding() + { + var bindingStore = new RecordingWorkflowWebhookBindingStore(); + bindingStore.Records["route-alpha"] = ExternalTriggerBindingRecord( + routeKey: "route-alpha", + callerDurableCredential: new DurableCallerCredentialRef + { + Ref = "credential-ref-alpha", + Purpose = CredentialSecretPurposes.WorkflowWebhookBindingAgentKey, + OwnerScopeKey = "scope-alpha", + SubjectId = "caller-alpha", + SourceKind = DurableCallerCredentialSourceKind.WebhookBinding, + ProviderCredentialId = "provider-key-alpha", + }); + var http = CreateExternalTriggerHttpContext(bindingStore); + var workflowQueryPort = new RecordingScopeWorkflowQueryPort + { + LookupResult = RunnableWorkflow(), + }; + + var result = await ScopeWorkflowScheduleEndpoints.GetExternalTrigger( + http, + "scope-alpha", + "wf-alpha", + workflowQueryPort, + CancellationToken.None); + + await result.ExecuteAsync(http); + var body = await ReadBodyAsync(http.Response); + + http.Response.StatusCode.Should().Be(StatusCodes.Status200OK); + body.Should().Contain("\"configured\":true"); + body.Should().Contain("\"routeKey\":\"route-alpha\""); + body.Should().Contain("\"fireUrl\":\"/api/workflow-webhooks/route-alpha\""); + body.Should().Contain("\"agentKeyReady\":true"); + body.Should().Contain("\"hmacSecretSet\":true"); + body.Should().NotContain("delivery-signing-secret-at-least-32-bytes"); + body.Should().NotContain("workflow-triggers"); + } + + [Fact] + public async Task WorkflowExternalTriggerDelete_ShouldRemoveBindingAndRevokeCredential() + { + var bindingStore = new RecordingWorkflowWebhookBindingStore(); + bindingStore.Records["route-alpha"] = ExternalTriggerBindingRecord( + routeKey: "route-alpha", + callerDurableCredential: new DurableCallerCredentialRef + { + Ref = "credential-ref-alpha", + Purpose = CredentialSecretPurposes.WorkflowWebhookBindingAgentKey, + OwnerScopeKey = "scope-alpha", + SubjectId = "caller-alpha", + SourceKind = DurableCallerCredentialSourceKind.WebhookBinding, + ProviderCredentialId = "provider-key-alpha", + }); + var materializer = new RecordingWorkflowWebhookAgentKeyMaterializer(); + var http = CreateExternalTriggerHttpContext(bindingStore, materializer); + var workflowQueryPort = new RecordingScopeWorkflowQueryPort + { + LookupResult = RunnableWorkflow(), + }; + + var result = await ScopeWorkflowScheduleEndpoints.DeleteExternalTrigger( + http, + "scope-alpha", + "wf-alpha", + workflowQueryPort, + CancellationToken.None); + + await result.ExecuteAsync(http); + + http.Response.StatusCode.Should().Be(StatusCodes.Status204NoContent); + bindingStore.Records.Should().BeEmpty(); + materializer.Revoked.Should().ContainSingle().Which.Credential.Ref.Should().Be("credential-ref-alpha"); + } + [Fact] public async Task WorkflowScheduleCreate_ShouldResolvePublishedServiceTargetWithoutTeamOwner() { @@ -2332,11 +2549,12 @@ private static ScopeWorkflowQueryApplicationService BuildQueryApplicationService private static DefaultHttpContext CreateHttpContext( string scopeId = "user-1", - IUserConfigQueryPort? userConfigQueryPort = null) + IUserConfigQueryPort? userConfigQueryPort = null, + Action? configureServices = null) { var http = new DefaultHttpContext { - RequestServices = BuildRequestServices(userConfigQueryPort), + RequestServices = BuildRequestServices(userConfigQueryPort, configureServices), }; http.Response.Body = new MemoryStream(); http.User = new ClaimsPrincipal( @@ -2377,15 +2595,23 @@ private static DefaultHttpContext CreateAnonymousHttpContext() return http; } - private static ServiceProvider BuildRequestServices(IUserConfigQueryPort? userConfigQueryPort = null) + private static ServiceProvider BuildRequestServices( + IUserConfigQueryPort? userConfigQueryPort = null, + Action? configureServices = null) { var services = new ServiceCollection() .AddLogging() .AddOptions() - .AddSingleton(new ConfigurationBuilder().Build()) + .AddSingleton(new ConfigurationBuilder() + .AddInMemoryCollection(new Dictionary + { + ["Aevatar:Authentication:Enabled"] = "true", + }) + .Build()) .AddSingleton(new TestHostEnvironment()); if (userConfigQueryPort != null) services.AddSingleton(userConfigQueryPort); + configureServices?.Invoke(services); return services.BuildServiceProvider(); } @@ -2453,33 +2679,44 @@ private static WorkflowRunEventEnvelope BuildRawObservedWorkflowExecutionStarted }, string.Empty); - private static ScheduledDispatchSummary WorkflowScheduleSummary(string scheduleId) => new( - scheduleId, - "Daily run", - ScheduledDispatchTargetKind.ServiceInvocation, - "target-actor-alpha", - Any.Pack(new ChatRequestEvent()).TypeUrl, - "svc-key-alpha", - "svc-alpha", - "chat", - "0 9 * * *", - "UTC", - true, - DateTimeOffset.UtcNow, - DateTimeOffset.UtcNow, - null, - null, - string.Empty, - string.Empty, - string.Empty, - string.Empty, - string.Empty, - 0, - 0, - new Dictionary(StringComparer.Ordinal), - $"actor:{scheduleId}", - "run workflow", - ScheduledDispatchScheduleKind.Workflow); + private static ScheduledDispatchSummary WorkflowScheduleSummary(string scheduleId) => + new( + scheduleId, + "Daily run", + ScheduledDispatchTargetKind.ServiceInvocation, + "target-actor-alpha", + Any.Pack(new ChatRequestEvent()).TypeUrl, + "svc-key-alpha", + "svc-alpha", + "chat", + "0 9 * * *", + "UTC", + true, + DateTimeOffset.UtcNow, + DateTimeOffset.UtcNow, + null, + null, + string.Empty, + string.Empty, + string.Empty, + string.Empty, + string.Empty, + 0, + 0, + new Dictionary(StringComparer.Ordinal), + $"actor:{scheduleId}", + "run workflow", + ScheduledDispatchScheduleKind.Workflow) + { + ServiceIdentity = new ServiceIdentity + { + TenantId = "scope-alpha", + AppId = "workflow-app", + Namespace = "workflow-ns", + ServiceId = "svc-alpha", + }, + ServiceRevisionId = "rev-alpha", + }; private sealed class RecordingScopeWorkflowQueryPort : IScopeWorkflowQueryPort, @@ -2538,10 +2775,258 @@ public Task LookupCatalogueByWorkflowIdAsync } } + private sealed class RecordingWorkflowCallerAccessTokenProvider : IWorkflowCallerAccessTokenProvider + { + public List Authorities { get; } = []; + + public Task IssueAsync( + CredentialWorkflowCallerAuthority authority, + CancellationToken ct = default) + { + Authorities.Add(authority.Clone()); + return Task.FromResult("issued-provisioning-token-alpha"); + } + } + + private static DefaultHttpContext CreateExternalTriggerHttpContext( + RecordingWorkflowWebhookBindingStore bindingStore, + RecordingWorkflowWebhookAgentKeyMaterializer? materializer = null, + RecordingWorkflowCallerAccessTokenProvider? tokenProvider = null) + { + var bindingReader = new FakeWorkflowActorBindingReader(); + bindingReader.Bindings["definition-actor-alpha"] = ExternalTriggerWorkflowBinding(); + return CreateHttpContext( + "scope-alpha", + configureServices: services => + { + services.AddSingleton(bindingStore); + services.AddSingleton(bindingReader); + services.AddSingleton(new FakeExternalIdentityBindingQueryPort()); + services.AddSingleton(tokenProvider ?? new RecordingWorkflowCallerAccessTokenProvider()); + services.AddSingleton(materializer ?? new RecordingWorkflowWebhookAgentKeyMaterializer()); + }); + } + + private static WorkflowActorBinding ExternalTriggerWorkflowBinding() + { + const string workflowYaml = "name: workflow-alpha\nsteps: []\n"; + var plan = ExternalTriggerDurableWritePlan(); + return new WorkflowActorBinding( + WorkflowActorKind.Definition, + "definition-actor-alpha", + "definition-actor-alpha", + string.Empty, + "workflow-alpha", + workflowYaml, + new Dictionary(), + ExternalCapabilityExecutionMode.Durable, + ScopeId: "scope-alpha", + SourceVersion: 1, + CapabilityAdmissionPlan: plan, + WorkflowId: "wf-alpha", + RevisionId: "rev-alpha"); + } + + private static WorkflowCapabilityAdmissionPlan ExternalTriggerDurableWritePlan() + { + var request = new NyxIdRequestSelector + { + UserServiceId = "service-alpha", + Method = NyxIdRequestMethod.Post, + PathTemplate = "/v1/resources", + BodyMode = NyxIdRequestBodyMode.Json, + BodyRequired = true, + ResponseMode = NyxIdRequestResponseMode.Text, + }; + var policy = new NyxIdOperationExecutionPolicy + { + Risk = NyxIdOperationRisk.Write, + Approval = NyxIdOperationApproval.Required, + EnforcementOwner = NyxIdOperationEnforcementOwner.Aevatar, + AllowedExecutionModes = + { + ExternalCapabilityExecutionMode.Interactive, + ExternalCapabilityExecutionMode.Durable, + }, + }; + var requestDigest = WorkflowCapabilityAdmissionPlanIntegrity.ComputeNyxIdRequestContractDigest(request); + var grant = new NyxIdExplicitRequestGrant + { + WorkflowId = "wf-alpha", + RevisionId = "rev-alpha", + CallSiteId = "workflow-alpha/update_resource", + RequestContractDigest = requestDigest, + GrantorAuthority = NyxIdExplicitRequestGrantorAuthority.AevatarWorkflowBinder, + GrantorOwnerKind = ExternalCapabilityAuthorizationOwnerKind.Personal, + GrantorOwnerSubject = "caller-alpha", + Risk = NyxIdOperationRisk.Write, + AllowedExecutionModes = + { + ExternalCapabilityExecutionMode.Interactive, + ExternalCapabilityExecutionMode.Durable, + }, + }; + var capability = new NyxIdUserRequestCapabilityRef + { + Request = request, + ServiceSlugSnapshot = "api-resource-service", + ContractDigest = WorkflowCapabilityAdmissionPlanIntegrity + .ComputeNyxIdExplicitRequestProofDigest(requestDigest, "api-resource-service"), + ExplicitRequestGrantDigest = WorkflowCapabilityAdmissionPlanIntegrity + .ComputeNyxIdExplicitRequestGrantDigest(grant), + ExecutionPolicy = policy, + }; + var plan = new WorkflowCapabilityAdmissionPlan + { + SchemaVersion = WorkflowCapabilityAdmissionPlanIntegrity.SchemaVersion, + DefinitionDigest = "sha256:definition", + ExecutionMode = ExternalCapabilityExecutionMode.Durable, + DurableAuthorizationOwner = new ExternalCapabilityAuthorizationOwner + { + Authority = WorkflowCapabilityAdmissionPlanIntegrity.NyxIdAuthority, + OwnerKind = ExternalCapabilityAuthorizationOwnerKind.Personal, + OwnerSubject = "caller-alpha", + }, + }; + plan.InvocationAdmissions.Add(new WorkflowCapabilityInvocationAdmission + { + CallSiteId = grant.CallSiteId, + Capability = new ExternalWorkflowCapabilityRef { NyxIdUserRequest = capability }, + NyxIdExplicitRequestGrant = grant, + }); + plan.AdmissionDigest = WorkflowCapabilityAdmissionPlanIntegrity.ComputeAdmissionDigest(plan); + return plan; + } + + private static WorkflowWebhookBindingRecord ExternalTriggerBindingRecord( + string routeKey, + string scopeId = "scope-alpha", + DurableCallerCredentialRef? callerDurableCredential = null) => new( + RouteKey: routeKey, + ScopeId: scopeId, + WorkflowName: "workflow-alpha", + SourceId: "nyxid-external-trigger", + PromptTemplate: "Run {{event_id}}", + PromptJsonPath: null, + DeliveryIdHeader: "X-NyxID-Delivery-Id", + DeliveryIdJsonPath: "event_id", + HmacSecret: "delivery-signing-secret-at-least-32-bytes", + HmacSignatureHeader: "X-NyxID-Signature", + HmacTimestampHeader: "X-NyxID-Timestamp", + MaxTimestampSkewSeconds: 300, + UpdatedAtUnixMs: DateTimeOffset.UtcNow.ToUnixTimeMilliseconds(), + DefinitionActorId: "definition-actor-alpha", + TargetRevisionId: "rev-alpha", + CallerAuthority: callerDurableCredential == null ? null : new CredentialWorkflowCallerAuthority + { + Platform = OwnerScope.NyxIdPlatform, + ExternalUserId = "caller-alpha", + Scope = "proxy", + BindingId = "binding-caller-alpha", + }, + CallerDurableCredential: callerDurableCredential); + + private sealed class RecordingWorkflowWebhookBindingStore : IWorkflowWebhookBindingStore + { + public Dictionary Records { get; } = new(StringComparer.Ordinal); + + public Task GetAsync(string routeKey, CancellationToken ct = default) => + Task.FromResult(Records.GetValueOrDefault(routeKey)); + + public async Task TryPutOwnedAsync(WorkflowWebhookBindingRecord record, CancellationToken ct = default) => + (await PutOwnedAsync(record, ct)).Succeeded; + + public Task PutOwnedAsync( + WorkflowWebhookBindingRecord record, + CancellationToken ct = default) + { + if (Records.TryGetValue(record.RouteKey, out var existing) && + !string.Equals(existing.ScopeId, record.ScopeId, StringComparison.Ordinal)) + { + return Task.FromResult(new WorkflowWebhookBindingPutResult(false, null)); + } + + Records[record.RouteKey] = record; + return Task.FromResult(new WorkflowWebhookBindingPutResult(true, existing)); + } + + public async Task TryDeleteOwnedAsync( + string routeKey, + string scopeId, + CancellationToken ct = default) => + (await DeleteOwnedAsync(routeKey, scopeId, ct)).Succeeded; + + public Task DeleteOwnedAsync( + string routeKey, + string scopeId, + CancellationToken ct = default) + { + if (!Records.TryGetValue(routeKey, out var existing) || + !string.Equals(existing.ScopeId, scopeId, StringComparison.Ordinal)) + { + return Task.FromResult(new WorkflowWebhookBindingDeleteResult(false, null)); + } + + Records.Remove(routeKey); + return Task.FromResult(new WorkflowWebhookBindingDeleteResult(true, existing)); + } + + public Task> ListByScopeAsync( + string scopeId, + CancellationToken ct = default) => + Task.FromResult>( + Records.Values.Where(record => string.Equals(record.ScopeId, scopeId, StringComparison.Ordinal)).ToArray()); + } + + private sealed class RecordingWorkflowWebhookAgentKeyMaterializer : IWorkflowWebhookAgentKeyMaterializer + { + public List<(CredentialWorkflowCallerAuthority Authority, string ScopeId, string RouteKey)> Materialized { get; } = []; + public List<(CredentialWorkflowCallerAuthority? Authority, DurableCallerCredentialRef Credential, string AuditReason)> Revoked { get; } = []; + + public Task MaterializeAsync( + CredentialWorkflowCallerAuthority callerAuthority, + WorkflowCapabilityAdmissionPlan admissionPlan, + string scopeId, + string routeKey, + CancellationToken ct) + { + Materialized.Add((callerAuthority.Clone(), scopeId, routeKey)); + return Task.FromResult(WorkflowWebhookAgentKeyMaterializationResult.Success(new DurableCallerCredentialRef + { + Ref = $"webhook-agent-key-ref-{Materialized.Count}", + Purpose = CredentialSecretPurposes.WorkflowWebhookBindingAgentKey, + OwnerScopeKey = scopeId, + SubjectId = callerAuthority.ExternalUserId, + SourceKind = DurableCallerCredentialSourceKind.WebhookBinding, + ProviderCredentialId = $"provider-key-{Materialized.Count}", + })); + } + + public Task RevokeAsync( + CredentialWorkflowCallerAuthority? callerAuthority, + DurableCallerCredentialRef credential, + string auditReason, + CancellationToken ct) + { + Revoked.Add((callerAuthority?.Clone(), credential.Clone(), auditReason)); + return Task.FromResult(true); + } + } + + private sealed class FakeExternalIdentityBindingQueryPort : IExternalIdentityBindingQueryPort + { + public Task ResolveAsync( + ExternalSubjectRef externalSubject, + CancellationToken ct = default) => + Task.FromResult(new BindingId { Value = $"binding-{externalSubject.ExternalUserId}" }); + } + private sealed class RecordingWorkflowScheduledDispatchService : IScheduledDispatchApplicationService { public List Created { get; } = []; public List CreateContexts { get; } = []; + public List Ensured { get; } = []; + public List EnsureContexts { get; } = []; public List<(string ScheduleId, ScheduledDispatchConfiguration Configuration)> Updated { get; } = []; public List UpdateContexts { get; } = []; public List EnableContexts { get; } = []; @@ -2555,7 +3040,93 @@ private sealed class RecordingWorkflowScheduledDispatchService : IScheduledDispa public (string CronExpression, string? Timezone, int Count, DateTimeOffset? FromUtc)? LastPreview { get; private set; } public ArgumentException? PreviewError { get; init; } public List RunNowScheduleIds { get; } = []; + public List RunNowTeamOwners { get; } = []; + public TeamMemberAutomationOwner? LastTeamAutomationGet { get; private set; } + public (string ScheduleId, string ScopeId, string? TeamId, string? MemberId)? LastTeamScheduleGet { get; private set; } public ScheduledDispatchDetail? Detail { get; init; } + public List BeginCredentialOperations { get; } = []; + public List<(string ScheduleId, string OperationId, string ErrorCode)> FailedCredentialOperations { get; } = []; + public List<(string ScheduleId, string OperationId)> RecordedCredentialCandidates { get; } = []; + public List<(string ScheduleId, string OperationId)> CompletedCredentialOperations { get; } = []; + + public Task BeginTeamAutomationCredentialOperationAsync( + TeamAutomationCredentialOperation operation, + CancellationToken ct = default) + { + BeginCredentialOperations.Add(operation); + return Task.FromResult(Committed( + operation.ScheduleId, + operation.OperationId, + operation.IdempotencyKey, + TeamAutomationOperationObservationStages.Begin, + ownsEffectAttempt: true, + "cmd-team-begin", + effectAttemptId: "attempt-alpha", + credentialEffectLocator: operation.CredentialEffectLocator, + newOperationCommitted: true)); + } + + public Task RecordTeamAutomationCredentialCandidateAsync( + string scheduleId, + TeamMemberAutomationOwner owner, + string operationId, + string idempotencyKey, + string effectAttemptId, + ScheduledInvocationAgentKeyCredentialReference credential, + ScheduledInvocationAuthorizationOwner credentialOwner, + CancellationToken ct = default) + { + RecordedCredentialCandidates.Add((scheduleId, operationId)); + return Task.FromResult(Committed( + scheduleId, + operationId, + idempotencyKey, + TeamAutomationOperationObservationStages.Candidate, + ownsEffectAttempt: false, + "cmd-team-candidate", + candidateCredential: credential, + candidateOwner: credentialOwner)); + } + + public Task CompleteTeamAutomationCredentialOperationAsync( + string scheduleId, + TeamMemberAutomationOwner owner, + string operationId, + string idempotencyKey, + string effectAttemptId, + ScheduledInvocationAgentKeyCredentialReference credential, + ScheduledDispatchConfiguration configuration, + CancellationToken ct = default) + { + CompletedCredentialOperations.Add((scheduleId, operationId)); + return Task.FromResult(Committed( + scheduleId, + operationId, + idempotencyKey, + TeamAutomationOperationObservationStages.Complete, + ownsEffectAttempt: false, + "cmd-team-complete")); + } + + public Task FailTeamAutomationCredentialOperationAsync( + string scheduleId, + TeamMemberAutomationOwner owner, + string operationId, + string idempotencyKey, + string effectAttemptId, + string errorCode, + CancellationToken ct = default) + { + FailedCredentialOperations.Add((scheduleId, operationId, errorCode)); + return Task.FromResult(Committed( + scheduleId, + operationId, + idempotencyKey, + TeamAutomationOperationObservationStages.Fail, + ownsEffectAttempt: false, + "cmd-team-fail", + errorCode)); + } public Task CreateAsync( ScheduledDispatchConfiguration configuration, @@ -2570,8 +3141,12 @@ public Task CreateAsync( public Task EnsureAsync( ScheduledDispatchConfiguration configuration, ScheduledDispatchMutationContext? context = null, - CancellationToken ct = default) => - Task.FromResult(MutationReceipt(configuration.ScheduleId)); + CancellationToken ct = default) + { + Ensured.Add(configuration); + EnsureContexts.Add(context); + return Task.FromResult(MutationReceipt(configuration.ScheduleId)); + } public Task UpdateAsync( string scheduleId, @@ -2624,6 +3199,39 @@ public Task DeleteAsync( return Task.FromResult(Detail?.Schedule.ScheduleId == scheduleId ? Detail : null); } + public Task GetTeamAutomationAsync( + string scheduleId, + TeamMemberAutomationOwner owner, + CancellationToken ct = default) + { + LastScheduleGet = scheduleId; + LastTeamAutomationGet = owner; + return Task.FromResult(IsTeamSchedule(scheduleId, owner) ? Detail : null); + } + + public Task GetTeamScheduleAsync( + string scheduleId, + string scopeId, + string? teamId = null, + string? memberId = null, + CancellationToken ct = default) + { + LastTeamScheduleGet = (scheduleId, scopeId, teamId, memberId); + if (Detail?.Schedule is not { } schedule || schedule.ScheduleId != scheduleId) + return Task.FromResult(null); + + if (!schedule.TeamOwned || !string.Equals(schedule.TeamOwnerScopeId, scopeId, StringComparison.Ordinal)) + return Task.FromResult(null); + + if (!string.IsNullOrWhiteSpace(teamId) && !string.Equals(schedule.TeamId, teamId, StringComparison.Ordinal)) + return Task.FromResult(null); + + if (!string.IsNullOrWhiteSpace(memberId) && !string.Equals(schedule.TeamOwnerMemberId, memberId, StringComparison.Ordinal)) + return Task.FromResult(null); + + return Task.FromResult(Detail); + } + public Task ListAsync( int take = 50, string? cursor = null, @@ -2660,7 +3268,29 @@ public Task RunNowAsync( { RunNowScheduleIds.Add(scheduleId); RunNowContexts.Add(context); - return Task.FromResult(new ScheduledDispatchRunNowReceipt( + return Task.FromResult(RunNowReceipt(scheduleId)); + } + + public Task RunTeamAutomationNowAsync( + string scheduleId, + TeamMemberAutomationOwner owner, + CancellationToken ct = default) + { + RunNowScheduleIds.Add(scheduleId); + RunNowTeamOwners.Add(owner); + return Task.FromResult(RunNowReceipt(scheduleId)); + } + + private bool IsTeamSchedule(string scheduleId, TeamMemberAutomationOwner owner) => + Detail?.Schedule is { } schedule && + schedule.ScheduleId == scheduleId && + schedule.TeamOwned && + string.Equals(schedule.TeamOwnerScopeId, owner.ScopeId, StringComparison.Ordinal) && + string.Equals(schedule.TeamOwnerMemberId, owner.MemberId, StringComparison.Ordinal) && + string.Equals(schedule.TeamId, owner.TeamId, StringComparison.Ordinal); + + private static ScheduledDispatchRunNowReceipt RunNowReceipt(string scheduleId) => + new( scheduleId, $"actor:{scheduleId}", DateTimeOffset.UtcNow, @@ -2669,8 +3299,44 @@ public Task RunNowAsync( "cmd-run-now", "corr-run-now", DateTimeOffset.UtcNow, - "accepted")); - } + "accepted"); + + private static TeamAutomationCommittedMutationReceipt Committed( + string scheduleId, + string operationId, + string idempotencyKey, + string stage, + bool ownsEffectAttempt, + string commandId, + string errorCode = "", + string effectAttemptId = "", + ScheduledInvocationAgentKeyCredentialReference? candidateCredential = null, + ScheduledInvocationAuthorizationOwner? candidateOwner = null, + ScheduledCredentialEffectLocator? credentialEffectLocator = null, + bool newOperationCommitted = false) => + new( + MutationReceipt(scheduleId) with { CommandId = commandId }, + new TeamAutomationOperationCommittedOutcome( + scheduleId, + operationId, + idempotencyKey, + stage, + ownsEffectAttempt, + StateVersion: 1, + errorCode, + ErrorMessage: string.Empty, + ObservedAtUtc: DateTimeOffset.UtcNow, + PendingRevocationCredential: null, + PendingRevocationOwner: null, + NyxIdRevocationPending: false, + VaultRevocationPending: false, + EffectAttemptId: effectAttemptId, + EffectAttemptGeneration: ownsEffectAttempt ? 1 : 0, + EffectAttemptExpiresAtUtc: ownsEffectAttempt ? DateTimeOffset.UtcNow.AddMinutes(5) : null, + CandidateCredential: candidateCredential, + CandidateOwner: candidateOwner, + CredentialEffectLocator: credentialEffectLocator, + NewOperationCommitted: newOperationCommitted)); private static ScheduledDispatchMutationReceipt MutationReceipt(string scheduleId) => new( scheduleId, diff --git a/test/Aevatar.GAgents.ChannelRuntime.Tests/AgentRunReplyGenerationExecutorTests.cs b/test/Aevatar.GAgents.ChannelRuntime.Tests/AgentRunReplyGenerationExecutorTests.cs index 28e82af8b..dd7b56449 100644 --- a/test/Aevatar.GAgents.ChannelRuntime.Tests/AgentRunReplyGenerationExecutorTests.cs +++ b/test/Aevatar.GAgents.ChannelRuntime.Tests/AgentRunReplyGenerationExecutorTests.cs @@ -520,7 +520,7 @@ await fixture.Executor.BuildLlmStepExecutionAsync( } [Fact] - public async Task BuildInitialStepState_WhenRegistrationAgentKeyCannotResolve_ShouldPreserveFailClosedRegistrationAuthority() + public async Task BuildInitialStepState_WhenRegistrationAgentKeyCannotResolve_ShouldClearAllNyxIdCredentials() { var fixture = CreateProfiledChannelExecutor(); var request = fixture.Request.Clone();