From 1cc7a91fc538deccc5f526b5de5f7ce640ad5688 Mon Sep 17 00:00:00 2001 From: Sipke Schoorstra Date: Wed, 16 Aug 2023 19:49:45 +0200 Subject: [PATCH] Fix Telnyx bookmark resumption --- .../Services/ProtoActorWorkflowRuntime.cs | 2 +- .../Elsa.Telnyx/Activities/AnswerCallBase.cs | 2 +- .../Elsa.Telnyx/Activities/BridgeCallsBase.cs | 6 +- .../Elsa.Telnyx/Activities/CallAnswered.cs | 2 +- .../Elsa.Telnyx/Activities/CallHangup.cs | 2 +- .../Elsa.Telnyx/Activities/DialAndWait.cs | 4 +- .../Activities/GatherUsingSpeak.cs | 2 +- .../Elsa.Telnyx/Activities/PlayAudioBase.cs | 2 +- .../Elsa.Telnyx/Activities/SpeakTextBase.cs | 2 +- .../Activities/StartRecordingBase.cs | 2 +- .../Elsa.Telnyx/Activities/TransferCall.cs | 4 +- .../Elsa.Telnyx/Activities/WebhookEvent.cs | 2 +- .../Handlers/TriggerCallBridgedActivities.cs | 60 +++++++++++++++++++ .../Handlers/TriggerWebhookActivities.cs | 14 ++--- .../TriggerWebhookDrivenActivities.cs | 12 +--- .../Contexts/ActivityExecutionContext.cs | 8 ++- 16 files changed, 89 insertions(+), 37 deletions(-) create mode 100644 src/modules/Elsa.Telnyx/Handlers/TriggerCallBridgedActivities.cs diff --git a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs index 31c0b7bda..399e87cb3 100644 --- a/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs +++ b/src/modules/Elsa.ProtoActor/Services/ProtoActorWorkflowRuntime.cs @@ -166,7 +166,7 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime /// public async Task> ResumeWorkflowsAsync(string activityTypeName, object bookmarkPayload, TriggerWorkflowsRuntimeOptions options, CancellationToken cancellationToken = default) { - var hash = _hasher.Hash(activityTypeName, bookmarkPayload); + var hash = _hasher.Hash(activityTypeName, bookmarkPayload, options.ActivityInstanceId); var correlationId = options.CorrelationId; var workflowInstanceId = options.WorkflowInstanceId; var filter = new BookmarkFilter { Hash = hash, CorrelationId = correlationId, WorkflowInstanceId = workflowInstanceId }; diff --git a/src/modules/Elsa.Telnyx/Activities/AnswerCallBase.cs b/src/modules/Elsa.Telnyx/Activities/AnswerCallBase.cs index e620e6540..2f54982b9 100644 --- a/src/modules/Elsa.Telnyx/Activities/AnswerCallBase.cs +++ b/src/modules/Elsa.Telnyx/Activities/AnswerCallBase.cs @@ -46,7 +46,7 @@ public abstract class AnswerCallBase : Activity await telnyxClient.Calls.AnswerCallAsync(callControlId, request, context.CancellationToken); // Create a bookmark so we can resume the workflow when the call is answered. - context.CreateBookmark(new AnswerCallBookmarkPayload(callControlId), ResumeAsync); + context.CreateBookmark(new AnswerCallBookmarkPayload(callControlId), ResumeAsync, true); } catch (ApiException e) { diff --git a/src/modules/Elsa.Telnyx/Activities/BridgeCallsBase.cs b/src/modules/Elsa.Telnyx/Activities/BridgeCallsBase.cs index 50abfb75e..7649e94a2 100644 --- a/src/modules/Elsa.Telnyx/Activities/BridgeCallsBase.cs +++ b/src/modules/Elsa.Telnyx/Activities/BridgeCallsBase.cs @@ -1,5 +1,4 @@ using Elsa.Extensions; -using Elsa.Telnyx.Attributes; using Elsa.Telnyx.Bookmarks; using Elsa.Telnyx.Client.Models; using Elsa.Telnyx.Client.Services; @@ -18,7 +17,6 @@ namespace Elsa.Telnyx.Activities; /// Bridge two calls. /// [Activity(Constants.Namespace, "Bridge two calls.", Kind = ActivityKind.Task)] -[WebhookDriven(WebhookEventTypes.CallBridged)] [PublicAPI] public abstract class BridgeCallsBase : Activity { @@ -44,7 +42,7 @@ public abstract class BridgeCallsBase : Activity { var callControlIdA = context.GetPrimaryCallControlId(CallControlIdA) ?? throw new Exception("CallControlA is required"); var callControlIdB = context.GetSecondaryCallControlId(CallControlIdB) ?? throw new Exception("CallControlB is required"); - var request = new BridgeCallsRequest(callControlIdB, ClientState: context.CreateCorrelatingClientState()); + var request = new BridgeCallsRequest(callControlIdB, ClientState: context.CreateCorrelatingClientState(context.Id)); var telnyxClient = context.GetRequiredService(); try @@ -53,7 +51,7 @@ public abstract class BridgeCallsBase : Activity var bookmarkA = new WebhookEventBookmarkPayload(WebhookEventTypes.CallBridged, callControlIdA); var bookmarkB = new WebhookEventBookmarkPayload(WebhookEventTypes.CallBridged, callControlIdB); - context.CreateBookmarks(new[] { bookmarkA, bookmarkB }, ResumeAsync); + context.CreateBookmarks(new[] { bookmarkA, bookmarkB }, ResumeAsync, false); } catch (ApiException e) { diff --git a/src/modules/Elsa.Telnyx/Activities/CallAnswered.cs b/src/modules/Elsa.Telnyx/Activities/CallAnswered.cs index 988cf9fbc..4d3d6c869 100644 --- a/src/modules/Elsa.Telnyx/Activities/CallAnswered.cs +++ b/src/modules/Elsa.Telnyx/Activities/CallAnswered.cs @@ -34,7 +34,7 @@ public class CallAnswered : Activity foreach (var callControlId in callControlIds) { var payload = new CallAnsweredBookmarkPayload(callControlId); - context.CreateBookmark(new BookmarkOptions(payload, Resume, Type)); + context.CreateBookmark(new BookmarkOptions(payload, Resume, Type, IncludeActivityInstanceId: false)); } } diff --git a/src/modules/Elsa.Telnyx/Activities/CallHangup.cs b/src/modules/Elsa.Telnyx/Activities/CallHangup.cs index b6088ba3e..babc59575 100644 --- a/src/modules/Elsa.Telnyx/Activities/CallHangup.cs +++ b/src/modules/Elsa.Telnyx/Activities/CallHangup.cs @@ -34,7 +34,7 @@ public class CallHangup : Activity foreach (var callControlId in callControlIds) { var payload = new CallHangupBookmarkPayload(callControlId); - context.CreateBookmark(new BookmarkOptions(payload, Resume, Type)); + context.CreateBookmark(new BookmarkOptions(payload, Resume, Type, IncludeActivityInstanceId: false)); } } diff --git a/src/modules/Elsa.Telnyx/Activities/DialAndWait.cs b/src/modules/Elsa.Telnyx/Activities/DialAndWait.cs index df7f50a2e..db531eb34 100644 --- a/src/modules/Elsa.Telnyx/Activities/DialAndWait.cs +++ b/src/modules/Elsa.Telnyx/Activities/DialAndWait.cs @@ -82,8 +82,8 @@ public class DialAndWait : Activity var answeredBookmark = new WebhookEventBookmarkPayload(WebhookEventTypes.CallAnswered, response.CallControlId); var hangupBookmark = new WebhookEventBookmarkPayload(WebhookEventTypes.CallHangup, response.CallControlId); - context.CreateBookmark(answeredBookmark, OnCallAnswered); - context.CreateBookmark(hangupBookmark, OnCallHangup); + context.CreateBookmark(answeredBookmark, OnCallAnswered, false); + context.CreateBookmark(hangupBookmark, OnCallHangup, false); } private async ValueTask OnCallAnswered(ActivityExecutionContext context) diff --git a/src/modules/Elsa.Telnyx/Activities/GatherUsingSpeak.cs b/src/modules/Elsa.Telnyx/Activities/GatherUsingSpeak.cs index 3bb847f67..a531c01c4 100644 --- a/src/modules/Elsa.Telnyx/Activities/GatherUsingSpeak.cs +++ b/src/modules/Elsa.Telnyx/Activities/GatherUsingSpeak.cs @@ -170,7 +170,7 @@ public class GatherUsingSpeak : Activity await telnyxClient.Calls.GatherUsingSpeakAsync(callControlId, request, context.CancellationToken); // Create a bookmark so we can resume this activity when the call.gather.ended webhook comes back. - context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallGatherEnded, callControlId), ResumeAsync); + context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallGatherEnded, callControlId), ResumeAsync, true); } catch (ApiException e) { diff --git a/src/modules/Elsa.Telnyx/Activities/PlayAudioBase.cs b/src/modules/Elsa.Telnyx/Activities/PlayAudioBase.cs index 76debca73..9ec0f5433 100644 --- a/src/modules/Elsa.Telnyx/Activities/PlayAudioBase.cs +++ b/src/modules/Elsa.Telnyx/Activities/PlayAudioBase.cs @@ -102,7 +102,7 @@ public abstract class PlayAudioBase : Activity await HandleDisconnectedAsync(context); } - context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallPlaybackStarted, callControlId), ResumeAsync); + context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallPlaybackStarted, callControlId), ResumeAsync, false); } /// diff --git a/src/modules/Elsa.Telnyx/Activities/SpeakTextBase.cs b/src/modules/Elsa.Telnyx/Activities/SpeakTextBase.cs index 8d57cfc37..a5f72bb71 100644 --- a/src/modules/Elsa.Telnyx/Activities/SpeakTextBase.cs +++ b/src/modules/Elsa.Telnyx/Activities/SpeakTextBase.cs @@ -106,7 +106,7 @@ public abstract class SpeakTextBase : Activity await telnyxClient.Calls.SpeakTextAsync(callControlId, request, context.CancellationToken); // Create bookmark to resume the workflow when speaking has finished. - context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallSpeakEnded, callControlId), HandleSpeakingHasFinished); + context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallSpeakEnded, callControlId), HandleSpeakingHasFinished, true); } catch (ApiException e) { diff --git a/src/modules/Elsa.Telnyx/Activities/StartRecordingBase.cs b/src/modules/Elsa.Telnyx/Activities/StartRecordingBase.cs index f7e8f86c2..335cdbe4c 100644 --- a/src/modules/Elsa.Telnyx/Activities/StartRecordingBase.cs +++ b/src/modules/Elsa.Telnyx/Activities/StartRecordingBase.cs @@ -79,7 +79,7 @@ public abstract class StartRecordingBase : Activity { await telnyxClient.Calls.StartRecordingAsync(callControlId, request, context.CancellationToken); - context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallRecordingSaved, callControlId), ResumeAsync); + context.CreateBookmark(new WebhookEventBookmarkPayload(WebhookEventTypes.CallRecordingSaved, callControlId), ResumeAsync, false); } catch (ApiException e) { diff --git a/src/modules/Elsa.Telnyx/Activities/TransferCall.cs b/src/modules/Elsa.Telnyx/Activities/TransferCall.cs index ece154003..698bdf4ff 100644 --- a/src/modules/Elsa.Telnyx/Activities/TransferCall.cs +++ b/src/modules/Elsa.Telnyx/Activities/TransferCall.cs @@ -108,8 +108,8 @@ public class TransferCall : Activity var callControlId = payload.CallControlId; var answeredBookmark = new WebhookEventBookmarkPayload(WebhookEventTypes.CallAnswered, callControlId); var hangupBookmark = new WebhookEventBookmarkPayload(WebhookEventTypes.CallHangup, callControlId); - context.CreateBookmark(answeredBookmark, AnsweredAsync); - context.CreateBookmark(hangupBookmark, HangupAsync); + context.CreateBookmark(answeredBookmark, AnsweredAsync, false); + context.CreateBookmark(hangupBookmark, HangupAsync, false); return default; } diff --git a/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs b/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs index ed5e1842d..d987f15cd 100644 --- a/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs +++ b/src/modules/Elsa.Telnyx/Activities/WebhookEvent.cs @@ -48,7 +48,7 @@ public class WebhookEvent : Activity var eventType = EventType; var payload = new WebhookEventBookmarkPayload(eventType); - context.CreateBookmark(new BookmarkOptions(payload, Resume, Type)); + context.CreateBookmark(new BookmarkOptions(payload, Resume, Type, false)); } } diff --git a/src/modules/Elsa.Telnyx/Handlers/TriggerCallBridgedActivities.cs b/src/modules/Elsa.Telnyx/Handlers/TriggerCallBridgedActivities.cs new file mode 100644 index 000000000..b1ab739fa --- /dev/null +++ b/src/modules/Elsa.Telnyx/Handlers/TriggerCallBridgedActivities.cs @@ -0,0 +1,60 @@ +using Elsa.Extensions; +using Elsa.Mediator.Contracts; +using Elsa.Telnyx.Activities; +using Elsa.Telnyx.Bookmarks; +using Elsa.Telnyx.Events; +using Elsa.Telnyx.Extensions; +using Elsa.Telnyx.Payloads.Call; +using Elsa.Workflows.Core.Helpers; +using Elsa.Workflows.Runtime.Contracts; +using Elsa.Workflows.Runtime.Models; +using JetBrains.Annotations; +using Microsoft.Extensions.Logging; + +namespace Elsa.Telnyx.Handlers; + +/// +/// Triggers all workflows starting with or blocked on a activity. +/// +[PublicAPI] +internal class TriggerCallBridgedActivities : INotificationHandler +{ + private readonly IWorkflowInbox _workflowInbox; + private readonly ILogger _logger; + + public TriggerCallBridgedActivities(IWorkflowInbox workflowInbox, ILogger logger) + { + _workflowInbox = workflowInbox; + _logger = logger; + } + + public async Task HandleAsync(TelnyxWebhookReceived notification, CancellationToken cancellationToken) + { + var webhook = notification.Webhook; + var payload = webhook.Data.Payload; + + if (payload is not CallBridgedPayload callBridgedPayload) + return; + + var clientStatePayload = callBridgedPayload.GetClientStatePayload(); + var correlationId = clientStatePayload.CorrelationId; + var input = new Dictionary().AddInput(callBridgedPayload); + var callControlId = callBridgedPayload.CallControlId; + var activityTypeNames = new[] + { + ActivityTypeNameHelper.GenerateTypeName(), + ActivityTypeNameHelper.GenerateTypeName(), + }; + + foreach (var activityTypeName in activityTypeNames) + { + await _workflowInbox.SubmitAsync(new NewWorkflowInboxMessage + { + ActivityTypeName = activityTypeName, + BookmarkPayload = new WebhookEventBookmarkPayload(WebhookEventTypes.CallBridged, callControlId), + CorrelationId = correlationId, + Input = input + }, cancellationToken); + } + } +} \ No newline at end of file diff --git a/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookActivities.cs b/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookActivities.cs index 88a66ce24..6948b686f 100644 --- a/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookActivities.cs +++ b/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookActivities.cs @@ -46,13 +46,13 @@ internal class TriggerWebhookActivities : INotificationHandler().AddInput(webhook); - await _workflowInbox.SubmitAsync(new NewWorkflowInboxMessage - { - ActivityTypeName = activityType, - BookmarkPayload = bookmarkPayload, - CorrelationId = correlationId, - Input = input - }, cancellationToken); + // await _workflowInbox.SubmitAsync(new NewWorkflowInboxMessage + // { + // ActivityTypeName = activityType, + // BookmarkPayload = bookmarkPayload, + // CorrelationId = correlationId, + // Input = input + // }, cancellationToken); await _workflowInbox.SubmitAsync(new NewWorkflowInboxMessage { diff --git a/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs b/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs index 97b51a622..4a9a1af5e 100644 --- a/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs +++ b/src/modules/Elsa.Telnyx/Handlers/TriggerWebhookDrivenActivities.cs @@ -39,27 +39,19 @@ internal class TriggerWebhookDrivenActivities : INotificationHandler().AddInput(eventPayload.GetType().Name, eventPayload); var activityDescriptors = FindActivityDescriptors(eventType).ToList(); var clientStatePayload = ((Payload)webhook.Data.Payload).GetClientStatePayload(); + var activityInstanceId = clientStatePayload.ActivityInstanceId; var correlationId = clientStatePayload.CorrelationId; var bookmarkPayload = new WebhookEventBookmarkPayload(eventType); var bookmarkPayloadWithCallControl = new WebhookEventBookmarkPayload(eventType, callControlId); foreach (var activityDescriptor in activityDescriptors) { - await _workflowInbox.SubmitAsync(new NewWorkflowInboxMessage - { - ActivityTypeName = activityDescriptor.TypeName, - BookmarkPayload = bookmarkPayload, - CorrelationId = correlationId, - ActivityInstanceId = clientStatePayload.ActivityInstanceId, - Input = input - }, cancellationToken); - await _workflowInbox.SubmitAsync(new NewWorkflowInboxMessage { ActivityTypeName = activityDescriptor.TypeName, BookmarkPayload = bookmarkPayloadWithCallControl, CorrelationId = correlationId, - ActivityInstanceId = clientStatePayload.ActivityInstanceId, + ActivityInstanceId = activityInstanceId, Input = input }, cancellationToken); } diff --git a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs index 0234ff069..5715448f2 100644 --- a/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs +++ b/src/modules/Elsa.Workflows.Core/Contexts/ActivityExecutionContext.cs @@ -263,10 +263,11 @@ public class ActivityExecutionContext : IExecutionContext /// /// The payloads to create bookmarks for. /// An optional callback that is invoked when the bookmark is resumed. - public void CreateBookmarks(IEnumerable payloads, ExecuteActivityDelegate? callback = default) + /// Whether or not the activity instance ID should be included in the bookmark payload. + public void CreateBookmarks(IEnumerable payloads, ExecuteActivityDelegate? callback = default, bool includeActivityInstanceId = true) { foreach (var payload in payloads) - CreateBookmark(new BookmarkOptions(payload, callback)); + CreateBookmark(new BookmarkOptions(payload, callback, IncludeActivityInstanceId: includeActivityInstanceId)); } /// @@ -293,8 +294,9 @@ public class ActivityExecutionContext : IExecutionContext /// /// The payload to associate with the bookmark. /// An optional callback that is invoked when the bookmark is resumed. + /// Whether or not the activity instance ID should be included in the bookmark payload. /// The created bookmark. - public Bookmark CreateBookmark(object payload, ExecuteActivityDelegate callback) => CreateBookmark(new BookmarkOptions(payload, callback)); + public Bookmark CreateBookmark(object payload, ExecuteActivityDelegate callback, bool includeActivityInstanceId = true) => CreateBookmark(new BookmarkOptions(payload, callback, IncludeActivityInstanceId: includeActivityInstanceId)); /// /// Creates a bookmark so that this activity can be resumed at a later time.