Update API endpoint for Polling Observer (#5787)

* Rename and refactor journal update endpoint

Replaced `/workflow-instances/{id}/journal/has-updates` endpoint with `/workflow-instances/{id}/updated-at` to simplify API responses. Deleted `HasUpdates` related classes and introduced `GetUpdatedAtResponse` for consistency and clarity. Updated client contracts accordingly.

* Remove HasUpdates endpoint and refactor workflow observer

Deleted the HasUpdates endpoint and refactored related code to use an updated timestamp approach instead. Improved nullable handling in WorkflowInstanceDesigner and ensured proper observer disposal to avoid memory leaks. Updated workflow observer factory and observer implementations to support observer names and enhanced logging.

* Rename updated workflow instance endpoint and handle execution state

Renamed the endpoint from "/updated-at" to "/execution-state" to better reflect its purpose. Updated related response models and documentation to capture workflow execution state details such as status, sub-status, and last updated timestamp.

* Enable SignalR for real-time workflows

Add a flag to use SignalR and activate real-time workflows when enabled. Refactor code to wrap SignalR setup in conditional checks based on the new flag. This enhances the application's interactivity through real-time capabilities.

* Remove obsolete endpoints and rename execution state paths

Deleted the outdated Api1 and DynamicWorkflows endpoints under Elsa.Server.Web. Also, renamed paths related to execution state models and endpoint to remove "Journal" from the namespace for better clarity and organization.
This commit is contained in:
Sipke Schoorstra 2024-07-18 07:51:46 +02:00 committed by GitHub
parent f448e9520a
commit 149f91e538
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
10 changed files with 73 additions and 139 deletions

View file

@ -1,33 +0,0 @@
using Elsa.Abstractions;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Parameters;
using JetBrains.Annotations;
namespace Elsa.Server.Web.Endpoints.Api1.Get;
/// <summary>
/// Returns a message.
/// </summary>
[UsedImplicitly]
public class Get : ElsaEndpointWithoutRequest
{
/// <inheritdoc />
public override void Configure()
{
Get("/api-1");
AllowAnonymous();
}
/// <inheritdoc />
public override async Task HandleAsync(CancellationToken ct)
{
await Task.Delay(1000, ct);
var response = new
{
Message = "OK"
};
await SendOkAsync(response, ct);
}
}

View file

@ -1,37 +0,0 @@
using Elsa.Abstractions;
using Elsa.Workflows.Activities;
using Elsa.Workflows.Models;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Options;
using Elsa.Workflows.Runtime.Parameters;
namespace Elsa.Server.Web.Endpoints.DynamicWorkflows.Post;
public class Post(IWorkflowRegistry workflowRegistry, IWorkflowRuntime workflowRuntime) : ElsaEndpointWithoutRequest
{
public override void Configure()
{
Post("/dynamic-workflows");
AllowAnonymous();
}
public override async Task HandleAsync(CancellationToken ct)
{
var workflow = new Workflow
{
Identity = new WorkflowIdentity("DynamicWorkflow1", 1, "DynamicWorkflow1:v1"),
Root = new Sequence
{
Activities =
{
new WriteLine("Step 1"),
new WriteLine("Step 2"),
new WriteLine("Step 3")
}
}
};
await workflowRegistry.RegisterAsync(workflow, ct);
await workflowRuntime.StartWorkflowAsync("DynamicWorkflow1", new StartWorkflowRuntimeParams());
}
}

View file

@ -48,6 +48,7 @@ const bool runEFCoreMigrations = true;
const bool useMemoryStores = false;
const bool useCaching = true;
const bool useReadOnlyMode = false;
const bool useSignalR = true;
const DistributedCachingTransport distributedCachingTransport = DistributedCachingTransport.MassTransit;
const MassTransitBroker useMassTransitBroker = MassTransitBroker.Memory;
@ -250,7 +251,6 @@ services
{
api.AddFastEndpointsAssembly<Program>();
})
.UseRealTimeWorkflows()
.UseCSharp(options =>
{
options.AppendScript("string Greet(string name) => $\"Hello {name}!\";");
@ -322,6 +322,11 @@ services
elsa.UseQuartz(quartz => { quartz.UseSqlite(sqliteConnectionString); });
}
if (useSignalR)
{
elsa.UseRealTimeWorkflows();
}
if (useMassTransit)
{
elsa.UseMassTransit(massTransit =>
@ -418,19 +423,18 @@ if (app.Environment.IsDevelopment())
}
// SignalR.
app.UseWorkflowsSignalRHubs();
if (useSignalR)
{
app.UseWorkflowsSignalRHubs();
}
// Run.
app.Run();
/// <summary>
/// The main entry point for the application made public for end to end testing.
/// </summary>
[UsedImplicitly]
public partial class Program
{
/// <summary>
/// Set by the test runner to configure the module for testing.
/// </summary>
public static Action<IModule>? ConfigureForTest { get; set; }
}
}

View file

@ -1,6 +1,7 @@
using Elsa.Api.Client.Resources.WorkflowDefinitions.Models;
using Elsa.Api.Client.Resources.WorkflowInstances.Models;
using Elsa.Api.Client.Resources.WorkflowInstances.Requests;
using Elsa.Api.Client.Resources.WorkflowInstances.Responses;
using Elsa.Api.Client.Shared.Models;
using Refit;
@ -47,14 +48,13 @@ public interface IWorkflowInstancesApi
Task<PagedListResponse<WorkflowExecutionLogRecord>> GetFilteredJournalAsync(string workflowInstanceId, GetFilteredJournalRequest? filter, int? skip = default, int? take = default, CancellationToken cancellationToken = default);
/// <summary>
/// Checks if there are updates in the journal for a specific workflow instance.
/// Returns the execution state of the specified workflow instance.
/// </summary>
/// <param name="workflowInstanceId">The ID of the workflow instance for which to check for updates.</param>
/// <param name="request">The request containing the ID and time from which to check for updates.</param>
/// <param name="workflowInstanceId">The ID of the workflow instance for which to return its execution state.</param>
/// <param name="cancellationToken">The cancellation token.</param>
/// <returns>Returns whether updates are available for the journal.</returns>
[Get("/workflow-instances/{workflowInstanceId}/journal/has-updates")]
Task<bool> HasJournalUpdates(string workflowInstanceId, [Query]HasJournalUpdateRequest request, CancellationToken cancellationToken = default);
/// <returns>Returns a response containing the execution state.</returns>
[Get("/workflow-instances/{workflowInstanceId}/execution-state")]
Task<WorkflowInstanceExecutionStateResponse> GetExecutionStateAsync(string workflowInstanceId, CancellationToken cancellationToken = default);
/// <summary>
/// Deletes a workflow instance.

View file

@ -1,11 +0,0 @@
namespace Elsa.Api.Client.Resources.WorkflowInstances.Models;
/// A request to update a journal for a workflow instance.
public class HasJournalUpdateRequest
{
/// The unique identifier of a workflow instance.
public string WorkflowInstanceId { get; set; } = default!;
/// The start date for checking for updates in the workflow instance journal.
public DateTimeOffset UpdatesSince { get; set; }
}

View file

@ -0,0 +1,6 @@
using Elsa.Api.Client.Resources.WorkflowInstances.Enums;
namespace Elsa.Api.Client.Resources.WorkflowInstances.Responses;
/// Represents the response containing the last updated timestamp of a workflow instance.
public record WorkflowInstanceExecutionStateResponse(WorkflowStatus Status, WorkflowSubStatus WorkflowSubStatus, DateTimeOffset UpdatedAt);

View file

@ -0,0 +1,37 @@
using Elsa.Abstractions;
using Elsa.Extensions;
using Elsa.Workflows.Management.Contracts;
using JetBrains.Annotations;
namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.ExecutionState;
/// Returns the execution state of the specified workflow instance.
[PublicAPI]
internal class ExecutionState(IWorkflowInstanceStore store) : ElsaEndpoint<Request, Response>
{
/// <inheritdoc />
public override void Configure()
{
Get("/workflow-instances/{id}/execution-state");
ConfigurePermissions("read:workflow-instances");
}
/// <inheritdoc />
public override async Task HandleAsync(Request request, CancellationToken cancellationToken)
{
var workflowInstance = await store.FindAsync(request.WorkflowInstanceId, cancellationToken);
if (workflowInstance == null)
{
await SendNotFoundAsync(cancellationToken);
return;
}
var response = new Response(
workflowInstance.Status,
workflowInstance.SubStatus,
workflowInstance.UpdatedAt);
await SendOkAsync(response, cancellationToken);
}
}

View file

@ -0,0 +1,13 @@
using FastEndpoints;
namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.ExecutionState;
/// The request to check the execution state of the workflow instance.
public class Request
{
/// The unique identifier of a workflow instance.
[BindFrom("id")] public string WorkflowInstanceId { get; set; } = default!;
}
/// Represents the response containing the last updated timestamp of a workflow instance.
public record Response(WorkflowStatus Status, WorkflowSubStatus SubStatus, DateTimeOffset UpdatedAt);

View file

@ -1,32 +0,0 @@
using Elsa.Abstractions;
using Elsa.Common.Entities;
using Elsa.Common.Models;
using Elsa.Workflows.Runtime.Contracts;
using Elsa.Workflows.Runtime.Filters;
using Elsa.Workflows.Runtime.OrderDefinitions;
using JetBrains.Annotations;
namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal.HasUpdates;
/// Endpoint that checks if there are updates for a workflow instance.
[PublicAPI]
internal class HasUpdates(IWorkflowExecutionLogStore store) : ElsaEndpoint<Request, bool>
{
/// <inheritdoc />
public override void Configure()
{
Get("/workflow-instances/{id}/journal/has-updates");
ConfigurePermissions("read:workflow-instances");
}
/// <inheritdoc />
public override async Task<bool> ExecuteAsync(Request request, CancellationToken cancellationToken)
{
var pageArgs = PageArgs.From(1, 1, 0, 1);
var filter = new WorkflowExecutionLogRecordFilter { WorkflowInstanceId = request.WorkflowInstanceId };
var order = new WorkflowExecutionLogRecordOrder<long>(x => x.Sequence, OrderDirection.Descending);
var pageOfRecords = await store.FindManyAsync(filter, pageArgs, order, cancellationToken);
return pageOfRecords.Items.Any(item => item.Timestamp >= request.UpdatesSince);
}
}

View file

@ -1,13 +0,0 @@
using FastEndpoints;
namespace Elsa.Workflows.Api.Endpoints.WorkflowInstances.Journal.HasUpdates;
/// The request to check if there are updates for a workflow instance journal.
public class Request
{
/// The unique identifier of a workflow instance.
[BindFrom("id")] public string WorkflowInstanceId { get; set; } = default!;
/// The start date for checking for updates in the workflow instance journal.
public DateTimeOffset UpdatesSince { get; set; }
}