Fix Telnyx bookmark resumption

This commit is contained in:
Sipke Schoorstra 2023-08-16 19:49:45 +02:00
parent 8c3f2658e8
commit 1cc7a91fc5
16 changed files with 89 additions and 37 deletions

View file

@ -166,7 +166,7 @@ internal class ProtoActorWorkflowRuntime : IWorkflowRuntime
/// <inheritdoc />
public async Task<ICollection<WorkflowExecutionResult>> 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 };

View file

@ -46,7 +46,7 @@ public abstract class AnswerCallBase : Activity<CallAnsweredPayload>
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)
{

View file

@ -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.
/// </summary>
[Activity(Constants.Namespace, "Bridge two calls.", Kind = ActivityKind.Task)]
[WebhookDriven(WebhookEventTypes.CallBridged)]
[PublicAPI]
public abstract class BridgeCallsBase : Activity<BridgedCallsOutput>
{
@ -44,7 +42,7 @@ public abstract class BridgeCallsBase : Activity<BridgedCallsOutput>
{
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<ITelnyxClient>();
try
@ -53,7 +51,7 @@ public abstract class BridgeCallsBase : Activity<BridgedCallsOutput>
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)
{

View file

@ -34,7 +34,7 @@ public class CallAnswered : Activity<CallAnsweredPayload>
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));
}
}

View file

@ -34,7 +34,7 @@ public class CallHangup : Activity<CallHangupPayload>
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));
}
}

View file

@ -82,8 +82,8 @@ public class DialAndWait : Activity<CallPayload>
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)

View file

@ -170,7 +170,7 @@ public class GatherUsingSpeak : Activity<CallGatherEndedPayload>
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)
{

View file

@ -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);
}
/// <summary>

View file

@ -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)
{

View file

@ -79,7 +79,7 @@ public abstract class StartRecordingBase : Activity<CallRecordingSavedPayload>
{
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)
{

View file

@ -108,8 +108,8 @@ public class TransferCall : Activity<CallPayload>
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;
}

View file

@ -48,7 +48,7 @@ public class WebhookEvent : Activity<Payload>
var eventType = EventType;
var payload = new WebhookEventBookmarkPayload(eventType);
context.CreateBookmark(new BookmarkOptions(payload, Resume, Type));
context.CreateBookmark(new BookmarkOptions(payload, Resume, Type, false));
}
}

View file

@ -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;
/// <summary>
/// Triggers all workflows starting with or blocked on a <see cref="CallAnswered"/> activity.
/// </summary>
[PublicAPI]
internal class TriggerCallBridgedActivities : INotificationHandler<TelnyxWebhookReceived>
{
private readonly IWorkflowInbox _workflowInbox;
private readonly ILogger _logger;
public TriggerCallBridgedActivities(IWorkflowInbox workflowInbox, ILogger<TriggerCallBridgedActivities> 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<string, object>().AddInput(callBridgedPayload);
var callControlId = callBridgedPayload.CallControlId;
var activityTypeNames = new[]
{
ActivityTypeNameHelper.GenerateTypeName<BridgeCalls>(),
ActivityTypeNameHelper.GenerateTypeName<FlowBridgeCalls>(),
};
foreach (var activityTypeName in activityTypeNames)
{
await _workflowInbox.SubmitAsync(new NewWorkflowInboxMessage
{
ActivityTypeName = activityTypeName,
BookmarkPayload = new WebhookEventBookmarkPayload(WebhookEventTypes.CallBridged, callControlId),
CorrelationId = correlationId,
Input = input
}, cancellationToken);
}
}
}

View file

@ -46,13 +46,13 @@ internal class TriggerWebhookActivities : INotificationHandler<TelnyxWebhookRece
var bookmarkPayload = new WebhookEventBookmarkPayload(eventType);
var input = new Dictionary<string, object>().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
{

View file

@ -39,27 +39,19 @@ internal class TriggerWebhookDrivenActivities : INotificationHandler<TelnyxWebho
var input = new Dictionary<string, object>().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);
}

View file

@ -263,10 +263,11 @@ public class ActivityExecutionContext : IExecutionContext
/// </summary>
/// <param name="payloads">The payloads to create bookmarks for.</param>
/// <param name="callback">An optional callback that is invoked when the bookmark is resumed.</param>
public void CreateBookmarks(IEnumerable<object> payloads, ExecuteActivityDelegate? callback = default)
/// <param name="includeActivityInstanceId">Whether or not the activity instance ID should be included in the bookmark payload.</param>
public void CreateBookmarks(IEnumerable<object> payloads, ExecuteActivityDelegate? callback = default, bool includeActivityInstanceId = true)
{
foreach (var payload in payloads)
CreateBookmark(new BookmarkOptions(payload, callback));
CreateBookmark(new BookmarkOptions(payload, callback, IncludeActivityInstanceId: includeActivityInstanceId));
}
/// <summary>
@ -293,8 +294,9 @@ public class ActivityExecutionContext : IExecutionContext
/// </summary>
/// <param name="payload">The payload to associate with the bookmark.</param>
/// <param name="callback">An optional callback that is invoked when the bookmark is resumed.</param>
/// <param name="includeActivityInstanceId">Whether or not the activity instance ID should be included in the bookmark payload.</param>
/// <returns>The created bookmark.</returns>
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));
/// <summary>
/// Creates a bookmark so that this activity can be resumed at a later time.