From 513f569f77bbd0490b220daff7952ec0ee9c19c7 Mon Sep 17 00:00:00 2001 From: BL19 Date: Tue, 25 Aug 2026 09:51:33 +0200 Subject: [PATCH] feat: add durable condition waits --- .../Flows/IFlowContext.cs | 10 + .../Nodes/ConditionNode.cs | 24 + src/Spindle.Abstractions/Nodes/NodeKind.cs | 5 + src/Spindle.Abstractions/README.md | 29 + .../Waiting/ConditionCallback.cs | 11 + .../Conditions/ConditionWaitRecord.cs | 24 + .../Conditions/IConditionWaitStore.cs | 20 + .../ISpindleStore.cs | 2 + .../ISpindleStoreSession.cs | 3 + .../Nodes/INodeStore.cs | 9 + ...260825072322_AddConditionWaits.Designer.cs | 740 +++++++++++++++++ .../20260825072322_AddConditionWaits.cs | 48 ++ .../SpindleDbContextModelSnapshot.cs | 35 + ...260825072319_AddConditionWaits.Designer.cs | 748 +++++++++++++++++ .../20260825072319_AddConditionWaits.cs | 48 ++ .../SpindleDbContextModelSnapshot.cs | 35 + ...260825072326_AddConditionWaits.Designer.cs | 749 ++++++++++++++++++ .../20260825072326_AddConditionWaits.cs | 48 ++ .../SpindleDbContextModelSnapshot.cs | 35 + ...260825072315_AddConditionWaits.Designer.cs | 738 +++++++++++++++++ .../20260825072315_AddConditionWaits.cs | 47 ++ .../SpindleDbContextModelSnapshot.cs | 35 + .../EFCoreSpindleStore.cs | 5 + .../Entities/ConditionWaitEntity.cs | 21 + .../FactoryStoreSession.cs | 19 +- .../SpindleDbContext.cs | 10 + .../Stores/EFCoreConditionWaitStore.cs | 63 ++ .../Stores/EFCoreNodeStore.cs | 57 +- .../InMemorySpindleStore.cs | 6 + .../Stores/InMemoryConditionWaitStore.cs | 52 ++ .../Stores/InMemoryNodeStore.cs | 39 +- .../Execution/LocalStepExecutor.cs | 133 +++- .../Execution/StepExecutionRegistration.cs | 1 + .../Execution/StepExecutionResult.cs | 18 +- src/Spindle.Runtime/Execution/StepExecutor.cs | 14 +- src/Spindle.Runtime/FlowExecutionSession.cs | 40 + src/Spindle.Runtime/FlowExecutor.cs | 4 + .../NamespacedFlowContextWrapper.cs | 3 + .../Nodes/ConditionNodeInitialization.cs | 5 + .../Nodes/RuntimeConditionNode.cs | 40 + src/Spindle.Runtime/Nodes/RuntimeNodeState.cs | 4 +- src/Spindle.Runtime/RuntimeFlowContext.cs | 58 ++ src/Spindle.Runtime/RuntimeSpindleRuntime.cs | 40 + .../Nodes/FlowContextConditionExtensions.cs | 176 ++++ .../EntityFrameworkPersistenceTests.cs | 62 ++ .../ConditionWaitTests.cs | 268 +++++++ .../Stores/CountingNodeStore.cs | 12 +- .../Stores/CountingSpindleStore.cs | 3 + .../Stores/CoutingStoreSession.cs | 3 + 49 files changed, 4577 insertions(+), 22 deletions(-) create mode 100644 src/Spindle.Abstractions/Nodes/ConditionNode.cs create mode 100644 src/Spindle.Abstractions/Waiting/ConditionCallback.cs create mode 100644 src/Spindle.Persistence.Abstractions/Conditions/ConditionWaitRecord.cs create mode 100644 src/Spindle.Persistence.Abstractions/Conditions/IConditionWaitStore.cs create mode 100644 src/Spindle.Persistence.EFCore.MySql/Migrations/20260825072322_AddConditionWaits.Designer.cs create mode 100644 src/Spindle.Persistence.EFCore.MySql/Migrations/20260825072322_AddConditionWaits.cs create mode 100644 src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/20260825072319_AddConditionWaits.Designer.cs create mode 100644 src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/20260825072319_AddConditionWaits.cs create mode 100644 src/Spindle.Persistence.EFCore.SqlServer/Migrations/20260825072326_AddConditionWaits.Designer.cs create mode 100644 src/Spindle.Persistence.EFCore.SqlServer/Migrations/20260825072326_AddConditionWaits.cs create mode 100644 src/Spindle.Persistence.EFCore.Sqlite/Migrations/20260825072315_AddConditionWaits.Designer.cs create mode 100644 src/Spindle.Persistence.EFCore.Sqlite/Migrations/20260825072315_AddConditionWaits.cs create mode 100644 src/Spindle.Persistence.EFCore/Entities/ConditionWaitEntity.cs create mode 100644 src/Spindle.Persistence.EFCore/Stores/EFCoreConditionWaitStore.cs create mode 100644 src/Spindle.Persistence.InMemory/Stores/InMemoryConditionWaitStore.cs create mode 100644 src/Spindle.Runtime/Nodes/ConditionNodeInitialization.cs create mode 100644 src/Spindle.Runtime/Nodes/RuntimeConditionNode.cs create mode 100644 src/Spindle/Nodes/FlowContextConditionExtensions.cs create mode 100644 tests/Spindle.Runtime.Tests/ConditionWaitTests.cs diff --git a/src/Spindle.Abstractions/Flows/IFlowContext.cs b/src/Spindle.Abstractions/Flows/IFlowContext.cs index 8762d8a..abe8167 100644 --- a/src/Spindle.Abstractions/Flows/IFlowContext.cs +++ b/src/Spindle.Abstractions/Flows/IFlowContext.cs @@ -101,6 +101,16 @@ SignalNode WaitForSignal( CorrelationKey correlationKey, SignalWaitOptions? options = null); + /// + /// Declares a durable condition that is checked immediately and then at the specified interval. + /// + ConditionNode WaitForCondition( + string id, + string name, + TimeSpan pollingInterval, + IReadOnlyList dependencies, + ConditionCallback condition); + /// /// Forks the execution into a namespace. This allows awaits inside the fork without affecting the entire flow. /// diff --git a/src/Spindle.Abstractions/Nodes/ConditionNode.cs b/src/Spindle.Abstractions/Nodes/ConditionNode.cs new file mode 100644 index 0000000..d1d837b --- /dev/null +++ b/src/Spindle.Abstractions/Nodes/ConditionNode.cs @@ -0,0 +1,24 @@ +using Spindle.Abstractions.Core; + +namespace Spindle.Abstractions.Nodes; + +/// +/// Represents a durable wait that polls an asynchronous condition until it succeeds. +/// +public abstract class ConditionNode : WaitNode +{ + /// + /// Gets the delay between unsuccessful condition checks. + /// + public abstract TimeSpan PollingInterval { get; } + + /// + /// Gets the optional timeout measured from the node's first declaration. + /// + public abstract TimeSpan? Timeout { get; } + + /// + /// Configures the maximum duration in which the condition may become true. + /// + public abstract ConditionNode WithTimeout(TimeSpan timeout); +} diff --git a/src/Spindle.Abstractions/Nodes/NodeKind.cs b/src/Spindle.Abstractions/Nodes/NodeKind.cs index a9aa433..9402e7e 100644 --- a/src/Spindle.Abstractions/Nodes/NodeKind.cs +++ b/src/Spindle.Abstractions/Nodes/NodeKind.cs @@ -41,4 +41,9 @@ public enum NodeKind /// A fork node that namespaces its children and allows async/await inside of a segment of a flow. /// Fork, + + /// + /// A condition that is evaluated repeatedly until it succeeds or times out. + /// + ConditionWait, } diff --git a/src/Spindle.Abstractions/README.md b/src/Spindle.Abstractions/README.md index 5b08c1b..506786c 100644 --- a/src/Spindle.Abstractions/README.md +++ b/src/Spindle.Abstractions/README.md @@ -14,6 +14,8 @@ It can also be used to only reference the abstractions without needing to includ - `DelayNode` represents a durable timer. - `SignalNode` represents a durable signal wait and exposes its signal name and correlation key. +- `ConditionNode` represents a durable polling wait. It checks immediately once + its inputs are ready, then schedules another check after each false result. - `WaitAllNode` and `WaitAnyNode` are durable barriers whose results identify the terminal input outcomes or winning input node. @@ -41,3 +43,30 @@ Barriers default to terminal-outcome semantics. Pass satisfy the barrier. `WaitAny` does not cancel inputs that did not win. When no display name is supplied, a node uses its explicit ID as its persisted name; explicit names remain available for operator-facing descriptions. + +## Waiting for a condition + +Use `WaitForCondition` when an external system cannot push a signal. Condition +callbacks run locally and may consume zero, one, or two typed node inputs. A +false result is not a failure: the node persists a timer and is checked again +after the polling interval. + +```csharp +var delivery = ctx.Step( + "delivery", + "Create delivery", + async () => await deliveryService.CreateAsync()); + +await ctx.WaitForCondition( + "wait-for-delivery", + "Wait for delivery", + TimeSpan.FromMinutes(5), + delivery, + details => deliveryService.CheckIsDeliveredAsync(details.Id)) + .WithTimeout(TimeSpan.FromDays(31)); +``` + +Timeouts are measured from the node's original declaration and survive flow +replay. A timed-out condition throws `TimeoutException` when awaited. Exceptions +from the condition callback fail the node immediately; only a returned `false` +schedules another poll. diff --git a/src/Spindle.Abstractions/Waiting/ConditionCallback.cs b/src/Spindle.Abstractions/Waiting/ConditionCallback.cs new file mode 100644 index 0000000..da80eb5 --- /dev/null +++ b/src/Spindle.Abstractions/Waiting/ConditionCallback.cs @@ -0,0 +1,11 @@ +using Spindle.Abstractions.Nodes; +using Spindle.Abstractions.Steps; + +namespace Spindle.Abstractions.Waiting; + +/// +/// Evaluates whether a durable condition wait has completed. +/// +public delegate ValueTask ConditionCallback( + NodeInputs inputs, + IStepExecutionContext context); diff --git a/src/Spindle.Persistence.Abstractions/Conditions/ConditionWaitRecord.cs b/src/Spindle.Persistence.Abstractions/Conditions/ConditionWaitRecord.cs new file mode 100644 index 0000000..0292d1f --- /dev/null +++ b/src/Spindle.Persistence.Abstractions/Conditions/ConditionWaitRecord.cs @@ -0,0 +1,24 @@ +using Spindle.Abstractions.Core; + +namespace Spindle.Persistence.Conditions; + +/// +/// Describes the durable scheduling metadata for a condition wait. +/// +public sealed record ConditionWaitRecord +{ + /// Gets the owning flow instance. + public required FlowInstanceId FlowInstanceId { get; init; } + + /// Gets the condition node identifier. + public required NodeId NodeId { get; init; } + + /// Gets the delay between unsuccessful checks. + public required TimeSpan PollingInterval { get; init; } + + /// Gets the optional absolute timeout deadline. + public DateTimeOffset? ExpiresAt { get; init; } + + /// Gets the time at which the condition was first declared. + public required DateTimeOffset CreatedAt { get; init; } +} diff --git a/src/Spindle.Persistence.Abstractions/Conditions/IConditionWaitStore.cs b/src/Spindle.Persistence.Abstractions/Conditions/IConditionWaitStore.cs new file mode 100644 index 0000000..2c8874f --- /dev/null +++ b/src/Spindle.Persistence.Abstractions/Conditions/IConditionWaitStore.cs @@ -0,0 +1,20 @@ +using Spindle.Abstractions.Core; + +namespace Spindle.Persistence.Conditions; + +/// +/// Persists scheduling metadata for durable condition waits. +/// +public interface IConditionWaitStore +{ + /// Creates condition scheduling metadata. + ValueTask CreateAsync( + ConditionWaitRecord wait, + CancellationToken cancellationToken = default); + + /// Gets condition scheduling metadata. + ValueTask GetAsync( + FlowInstanceId flowInstanceId, + NodeId nodeId, + CancellationToken cancellationToken = default); +} diff --git a/src/Spindle.Persistence.Abstractions/ISpindleStore.cs b/src/Spindle.Persistence.Abstractions/ISpindleStore.cs index b6ba431..4155495 100644 --- a/src/Spindle.Persistence.Abstractions/ISpindleStore.cs +++ b/src/Spindle.Persistence.Abstractions/ISpindleStore.cs @@ -10,6 +10,8 @@ public interface ISpindleStore Timers.ITimerStore Timers { get; } + Conditions.IConditionWaitStore Conditions { get; } + Signals.ISignalStore Signals { get; } Messaging.IOutboxStore Outbox { get; } diff --git a/src/Spindle.Persistence.Abstractions/ISpindleStoreSession.cs b/src/Spindle.Persistence.Abstractions/ISpindleStoreSession.cs index 380def9..2e0d94f 100644 --- a/src/Spindle.Persistence.Abstractions/ISpindleStoreSession.cs +++ b/src/Spindle.Persistence.Abstractions/ISpindleStoreSession.cs @@ -1,4 +1,5 @@ using Spindle.Persistence.FlowDefinitions; +using Spindle.Persistence.Conditions; using Spindle.Persistence.FlowInstances; using Spindle.Persistence.History; using Spindle.Persistence.Leases; @@ -19,6 +20,8 @@ public interface ISpindleStoreSession ITimerStore Timers { get; } + IConditionWaitStore Conditions { get; } + ISignalStore Signals { get; } IOutboxStore Outbox { get; } diff --git a/src/Spindle.Persistence.Abstractions/Nodes/INodeStore.cs b/src/Spindle.Persistence.Abstractions/Nodes/INodeStore.cs index cede2fd..c7fe957 100644 --- a/src/Spindle.Persistence.Abstractions/Nodes/INodeStore.cs +++ b/src/Spindle.Persistence.Abstractions/Nodes/INodeStore.cs @@ -50,9 +50,18 @@ ValueTask MarkRunningAsync( ValueTask MarkWaitingAsync( FlowInstanceId flowInstanceId, NodeId nodeId, + int attempt, DateTimeOffset updatedAt, CancellationToken cancellationToken = default); + ValueTask MarkTimedOutAsync( + FlowInstanceId flowInstanceId, + NodeId nodeId, + int attempt, + string error, + DateTimeOffset timedOutAt, + CancellationToken cancellationToken = default); + ValueTask MarkCompletedAsync( FlowInstanceId flowInstanceId, NodeId nodeId, diff --git a/src/Spindle.Persistence.EFCore.MySql/Migrations/20260825072322_AddConditionWaits.Designer.cs b/src/Spindle.Persistence.EFCore.MySql/Migrations/20260825072322_AddConditionWaits.Designer.cs new file mode 100644 index 0000000..4bd0bd0 --- /dev/null +++ b/src/Spindle.Persistence.EFCore.MySql/Migrations/20260825072322_AddConditionWaits.Designer.cs @@ -0,0 +1,740 @@ +// +using System.Collections.Generic; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Spindle.Persistence.EFCore; + +#nullable disable + +namespace Spindle.Persistence.EFCore.MySql.Migrations +{ + [DbContext(typeof(SpindleDbContext))] + [Migration("20260825072322_AddConditionWaits")] + partial class AddConditionWaits + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.10") + .HasAnnotation("Relational:MaxIdentifierLength", 64); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("ExpiresAt") + .HasColumnType("bigint"); + + b.Property("PollingIntervalTicks") + .HasColumnType("bigint"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("ExpiresAt"); + + b.ToTable("ConditionWaits", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("bigint"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("EventType") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("FlowInstanceId") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("NodeId") + .HasColumnType("longtext"); + + b.HasKey("Id"); + + b.ToTable("ExecutionHistories", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowDefinitionEntity", b => + { + b.Property("FlowName") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("FlowVersion") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("DefinitionHash") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("FlowTypeName") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("UpdatedAt") + .HasColumnType("bigint"); + + b.HasKey("FlowName", "FlowVersion"); + + b.ToTable("FlowDefinitions", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowInstanceEntity", b => + { + b.Property("InstanceId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("CompletedAt") + .HasColumnType("bigint"); + + b.Property("CorrelationKey") + .HasColumnType("longtext"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("DefinitionHash") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("Error") + .HasColumnType("longtext"); + + b.Property("FlowName") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("FlowVersion") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("IdempotencyKey") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("Status") + .HasColumnType("int"); + + b.Property("UpdatedAt") + .HasColumnType("bigint"); + + b.ComplexProperty(typeof(Dictionary), "Input", "Spindle.Persistence.EFCore.Entities.FlowInstanceEntity.Input#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Input_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Input_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Input_TypeName"); + }); + + b.HasKey("InstanceId"); + + b.HasIndex("FlowName", "IdempotencyKey") + .IsUnique(); + + b.HasIndex("Status", "UpdatedAt"); + + b.ToTable("FlowInstances", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.InboxMessageEntity", b => + { + b.Property("MessageId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("Kind") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("ProcessedAt") + .HasColumnType("bigint"); + + b.Property("ReceivedAt") + .HasColumnType("bigint"); + + b.ComplexProperty(typeof(Dictionary), "Payload", "Spindle.Persistence.EFCore.Entities.InboxMessageEntity.Payload#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Payload_TypeName"); + }); + + b.HasKey("MessageId"); + + b.ToTable("InboxMessages", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeDependencyEntity", b => + { + b.Property("FlowInstanceId") + .HasColumnType("varchar(255)"); + + b.Property("NodeId") + .HasColumnType("varchar(255)"); + + b.Property("DependsOnId") + .HasColumnType("varchar(255)"); + + b.Property("Position") + .HasColumnType("int"); + + b.HasKey("FlowInstanceId", "NodeId", "DependsOnId"); + + b.HasIndex("FlowInstanceId", "DependsOnId"); + + b.HasIndex("FlowInstanceId", "NodeId"); + + b.ToTable("NodeDependencies"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("Attempt") + .HasColumnType("int"); + + b.Property("CompletedAt") + .HasColumnType("bigint"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("DependencyMode") + .HasColumnType("int"); + + b.Property("DispatchMode") + .HasColumnType("int"); + + b.Property("Error") + .HasColumnType("longtext"); + + b.Property("HandlerId") + .HasColumnType("longtext"); + + b.Property("Kind") + .HasColumnType("int"); + + b.Property("Name") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("Queue") + .HasColumnType("longtext"); + + b.Property("RetryAt") + .HasColumnType("bigint"); + + b.Property("StartedAt") + .HasColumnType("bigint"); + + b.Property("Status") + .HasColumnType("int"); + + b.Property("UpdatedAt") + .HasColumnType("bigint"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("FlowInstanceId", "Status"); + + b.HasIndex("Status", "CreatedAt"); + + b.ToTable("NodeInstances", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.OutboxMessageEntity", b => + { + b.Property("MessageId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("Headers") + .IsRequired() + .HasColumnType("json"); + + b.Property("Kind") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("PublishedAt") + .HasColumnType("bigint"); + + b.ComplexProperty(typeof(Dictionary), "Payload", "Spindle.Persistence.EFCore.Entities.OutboxMessageEntity.Payload#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Payload_TypeName"); + }); + + b.HasKey("MessageId"); + + b.HasIndex("PublishedAt", "CreatedAt"); + + b.ToTable("OutboxMessages", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("int"); + + b.Property("CorrelationKey") + .HasColumnType("longtext"); + + b.Property("FlowInstanceId") + .HasColumnType("longtext"); + + b.Property("RaisedAt") + .HasColumnType("bigint"); + + b.Property("SignalName") + .IsRequired() + .HasColumnType("longtext"); + + b.HasKey("Id"); + + b.ToTable("Signals", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("CompletedAt") + .HasColumnType("bigint"); + + b.Property("CorrelationKey") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("ExpiresAt") + .HasColumnType("bigint"); + + b.Property("SignalName") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("SignalName", "CorrelationKey", "CompletedAt"); + + b.ToTable("SignalWaits", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.StepAttemptEntity", b => + { + b.Property("AttemptId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("Attempt") + .HasColumnType("int"); + + b.Property("CompletedAt") + .HasColumnType("bigint"); + + b.Property("Error") + .HasColumnType("longtext"); + + b.Property("FlowInstanceId") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("NodeId") + .IsRequired() + .HasColumnType("longtext"); + + b.Property("StartedAt") + .HasColumnType("bigint"); + + b.Property("Status") + .HasColumnType("int"); + + b.Property("WorkerId") + .IsRequired() + .HasColumnType("longtext"); + + b.HasKey("AttemptId"); + + b.ToTable("StepAttempts", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.StepLeaseEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("AcquiredAt") + .HasColumnType("bigint"); + + b.Property("ExpiresAt") + .HasColumnType("bigint"); + + b.Property("Owner") + .IsRequired() + .HasColumnType("longtext"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.ToTable("StepLeases", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.TimerEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("DueAt") + .HasColumnType("bigint"); + + b.Property("FiredAt") + .HasColumnType("bigint"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("FiredAt", "DueAt"); + + b.ToTable("Timers", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", null) + .WithOne() + .HasForeignKey("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", "FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => + { + b1.Property("ExecutionHistoryEntityId") + .HasColumnType("bigint"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Payload_TypeName"); + + b1.HasKey("ExecutionHistoryEntityId"); + + b1.ToTable("ExecutionHistories"); + + b1.WithOwner() + .HasForeignKey("ExecutionHistoryEntityId"); + }); + + b.Navigation("Payload"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowDefinitionEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Definition", b1 => + { + b1.Property("FlowDefinitionEntityFlowName") + .HasColumnType("varchar(255)"); + + b1.Property("FlowDefinitionEntityFlowVersion") + .HasColumnType("varchar(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Definition_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Definition_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Definition_TypeName"); + + b1.HasKey("FlowDefinitionEntityFlowName", "FlowDefinitionEntityFlowVersion"); + + b1.ToTable("FlowDefinitions"); + + b1.WithOwner() + .HasForeignKey("FlowDefinitionEntityFlowName", "FlowDefinitionEntityFlowVersion"); + }); + + b.Navigation("Definition"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowInstanceEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Result", b1 => + { + b1.Property("FlowInstanceEntityInstanceId") + .HasColumnType("varchar(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Result_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Result_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Result_TypeName"); + + b1.HasKey("FlowInstanceEntityInstanceId"); + + b1.ToTable("FlowInstances"); + + b1.WithOwner() + .HasForeignKey("FlowInstanceEntityInstanceId"); + }); + + b.Navigation("Result"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeDependencyEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", "DependsOn") + .WithMany("Dependents") + .HasForeignKey("FlowInstanceId", "DependsOnId") + .OnDelete(DeleteBehavior.Restrict) + .IsRequired(); + + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", "Node") + .WithMany("Dependencies") + .HasForeignKey("FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Restrict) + .IsRequired(); + + b.Navigation("DependsOn"); + + b.Navigation("Node"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Input", b1 => + { + b1.Property("NodeInstanceEntityFlowInstanceId") + .HasColumnType("varchar(255)"); + + b1.Property("NodeInstanceEntityNodeId") + .HasColumnType("varchar(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Input_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Input_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Input_TypeName"); + + b1.HasKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + + b1.ToTable("NodeInstances"); + + b1.WithOwner() + .HasForeignKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + }); + + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Result", b1 => + { + b1.Property("NodeInstanceEntityFlowInstanceId") + .HasColumnType("varchar(255)"); + + b1.Property("NodeInstanceEntityNodeId") + .HasColumnType("varchar(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Result_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Result_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Result_TypeName"); + + b1.HasKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + + b1.ToTable("NodeInstances"); + + b1.WithOwner() + .HasForeignKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + }); + + b.Navigation("Input"); + + b.Navigation("Result"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => + { + b1.Property("SignalEntityId") + .HasColumnType("int"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("longblob") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("longtext") + .HasColumnName("Payload_TypeName"); + + b1.HasKey("SignalEntityId"); + + b1.ToTable("Signals"); + + b1.WithOwner() + .HasForeignKey("SignalEntityId"); + }); + + b.Navigation("Payload"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.Navigation("Dependencies"); + + b.Navigation("Dependents"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/src/Spindle.Persistence.EFCore.MySql/Migrations/20260825072322_AddConditionWaits.cs b/src/Spindle.Persistence.EFCore.MySql/Migrations/20260825072322_AddConditionWaits.cs new file mode 100644 index 0000000..89f7ad6 --- /dev/null +++ b/src/Spindle.Persistence.EFCore.MySql/Migrations/20260825072322_AddConditionWaits.cs @@ -0,0 +1,48 @@ +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace Spindle.Persistence.EFCore.MySql.Migrations +{ + /// + public partial class AddConditionWaits : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "ConditionWaits", + columns: table => new + { + FlowInstanceId = table.Column(type: "varchar(255)", maxLength: 255, nullable: false), + NodeId = table.Column(type: "varchar(255)", maxLength: 255, nullable: false), + PollingIntervalTicks = table.Column(type: "bigint", nullable: false), + ExpiresAt = table.Column(type: "bigint", nullable: true), + CreatedAt = table.Column(type: "bigint", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_ConditionWaits", x => new { x.FlowInstanceId, x.NodeId }); + table.ForeignKey( + name: "FK_ConditionWaits_NodeInstances_FlowInstanceId_NodeId", + columns: x => new { x.FlowInstanceId, x.NodeId }, + principalTable: "NodeInstances", + principalColumns: new[] { "FlowInstanceId", "NodeId" }, + onDelete: ReferentialAction.Cascade); + }) + .Annotation("MySQL:Charset", "utf8mb4"); + + migrationBuilder.CreateIndex( + name: "IX_ConditionWaits_ExpiresAt", + table: "ConditionWaits", + column: "ExpiresAt"); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "ConditionWaits"); + } + } +} diff --git a/src/Spindle.Persistence.EFCore.MySql/Migrations/SpindleDbContextModelSnapshot.cs b/src/Spindle.Persistence.EFCore.MySql/Migrations/SpindleDbContextModelSnapshot.cs index eef1b6f..aefecaa 100644 --- a/src/Spindle.Persistence.EFCore.MySql/Migrations/SpindleDbContextModelSnapshot.cs +++ b/src/Spindle.Persistence.EFCore.MySql/Migrations/SpindleDbContextModelSnapshot.cs @@ -19,6 +19,32 @@ protected override void BuildModel(ModelBuilder modelBuilder) .HasAnnotation("ProductVersion", "10.0.10") .HasAnnotation("Relational:MaxIdentifierLength", 64); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("varchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("bigint"); + + b.Property("ExpiresAt") + .HasColumnType("bigint"); + + b.Property("PollingIntervalTicks") + .HasColumnType("bigint"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("ExpiresAt"); + + b.ToTable("ConditionWaits", (string)null); + }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => { b.Property("Id") @@ -467,6 +493,15 @@ protected override void BuildModel(ModelBuilder modelBuilder) b.ToTable("Timers", (string)null); }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", null) + .WithOne() + .HasForeignKey("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", "FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => { b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => diff --git a/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/20260825072319_AddConditionWaits.Designer.cs b/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/20260825072319_AddConditionWaits.Designer.cs new file mode 100644 index 0000000..f4551fb --- /dev/null +++ b/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/20260825072319_AddConditionWaits.Designer.cs @@ -0,0 +1,748 @@ +// +using System; +using System.Collections.Generic; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Npgsql.EntityFrameworkCore.PostgreSQL.Metadata; +using Spindle.Persistence.EFCore; + +#nullable disable + +namespace Spindle.Persistence.EFCore.PostgreSQL.Migrations +{ + [DbContext(typeof(SpindleDbContext))] + [Migration("20260825072319_AddConditionWaits")] + partial class AddConditionWaits + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.10") + .HasAnnotation("Relational:MaxIdentifierLength", 63); + + NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("ExpiresAt") + .HasColumnType("timestamp with time zone"); + + b.Property("PollingIntervalTicks") + .HasColumnType("bigint"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("ExpiresAt"); + + b.ToTable("ConditionWaits", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("bigint"); + + NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id")); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("EventType") + .IsRequired() + .HasColumnType("text"); + + b.Property("FlowInstanceId") + .IsRequired() + .HasColumnType("text"); + + b.Property("NodeId") + .HasColumnType("text"); + + b.HasKey("Id"); + + b.ToTable("ExecutionHistories", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowDefinitionEntity", b => + { + b.Property("FlowName") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("FlowVersion") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("DefinitionHash") + .IsRequired() + .HasColumnType("text"); + + b.Property("FlowTypeName") + .IsRequired() + .HasColumnType("text"); + + b.Property("UpdatedAt") + .HasColumnType("timestamp with time zone"); + + b.HasKey("FlowName", "FlowVersion"); + + b.ToTable("FlowDefinitions", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowInstanceEntity", b => + { + b.Property("InstanceId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("CompletedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("CorrelationKey") + .HasColumnType("text"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("DefinitionHash") + .IsRequired() + .HasColumnType("text"); + + b.Property("Error") + .HasColumnType("text"); + + b.Property("FlowName") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("FlowVersion") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("IdempotencyKey") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("Status") + .HasColumnType("integer"); + + b.Property("UpdatedAt") + .HasColumnType("timestamp with time zone"); + + b.ComplexProperty(typeof(Dictionary), "Input", "Spindle.Persistence.EFCore.Entities.FlowInstanceEntity.Input#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Input_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Input_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Input_TypeName"); + }); + + b.HasKey("InstanceId"); + + b.HasIndex("FlowName", "IdempotencyKey") + .IsUnique(); + + b.HasIndex("Status", "UpdatedAt"); + + b.ToTable("FlowInstances", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.InboxMessageEntity", b => + { + b.Property("MessageId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("Kind") + .IsRequired() + .HasColumnType("text"); + + b.Property("ProcessedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("ReceivedAt") + .HasColumnType("timestamp with time zone"); + + b.ComplexProperty(typeof(Dictionary), "Payload", "Spindle.Persistence.EFCore.Entities.InboxMessageEntity.Payload#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Payload_TypeName"); + }); + + b.HasKey("MessageId"); + + b.ToTable("InboxMessages", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeDependencyEntity", b => + { + b.Property("FlowInstanceId") + .HasColumnType("character varying(255)"); + + b.Property("NodeId") + .HasColumnType("character varying(255)"); + + b.Property("DependsOnId") + .HasColumnType("character varying(255)"); + + b.Property("Position") + .HasColumnType("integer"); + + b.HasKey("FlowInstanceId", "NodeId", "DependsOnId"); + + b.HasIndex("FlowInstanceId", "DependsOnId"); + + b.HasIndex("FlowInstanceId", "NodeId"); + + b.ToTable("NodeDependencies"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("Attempt") + .HasColumnType("integer"); + + b.Property("CompletedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("DependencyMode") + .HasColumnType("integer"); + + b.Property("DispatchMode") + .HasColumnType("integer"); + + b.Property("Error") + .HasColumnType("text"); + + b.Property("HandlerId") + .HasColumnType("text"); + + b.Property("Kind") + .HasColumnType("integer"); + + b.Property("Name") + .IsRequired() + .HasColumnType("text"); + + b.Property("Queue") + .HasColumnType("text"); + + b.Property("RetryAt") + .HasColumnType("timestamp with time zone"); + + b.Property("StartedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("Status") + .HasColumnType("integer"); + + b.Property("UpdatedAt") + .HasColumnType("timestamp with time zone"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("FlowInstanceId", "Status"); + + b.HasIndex("Status", "CreatedAt"); + + b.ToTable("NodeInstances", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.OutboxMessageEntity", b => + { + b.Property("MessageId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("Headers") + .IsRequired() + .HasColumnType("jsonb"); + + b.Property("Kind") + .IsRequired() + .HasColumnType("text"); + + b.Property("PublishedAt") + .HasColumnType("timestamp with time zone"); + + b.ComplexProperty(typeof(Dictionary), "Payload", "Spindle.Persistence.EFCore.Entities.OutboxMessageEntity.Payload#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Payload_TypeName"); + }); + + b.HasKey("MessageId"); + + b.HasIndex("PublishedAt", "CreatedAt"); + + b.ToTable("OutboxMessages", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("integer"); + + NpgsqlPropertyBuilderExtensions.UseIdentityByDefaultColumn(b.Property("Id")); + + b.Property("CorrelationKey") + .HasColumnType("text"); + + b.Property("FlowInstanceId") + .HasColumnType("text"); + + b.Property("RaisedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("SignalName") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("Id"); + + b.ToTable("Signals", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("CompletedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("CorrelationKey") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("ExpiresAt") + .HasColumnType("timestamp with time zone"); + + b.Property("SignalName") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("SignalName", "CorrelationKey", "CompletedAt"); + + b.ToTable("SignalWaits", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.StepAttemptEntity", b => + { + b.Property("AttemptId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("Attempt") + .HasColumnType("integer"); + + b.Property("CompletedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("Error") + .HasColumnType("text"); + + b.Property("FlowInstanceId") + .IsRequired() + .HasColumnType("text"); + + b.Property("NodeId") + .IsRequired() + .HasColumnType("text"); + + b.Property("StartedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("Status") + .HasColumnType("integer"); + + b.Property("WorkerId") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("AttemptId"); + + b.ToTable("StepAttempts", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.StepLeaseEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("AcquiredAt") + .HasColumnType("timestamp with time zone"); + + b.Property("ExpiresAt") + .HasColumnType("timestamp with time zone"); + + b.Property("Owner") + .IsRequired() + .HasColumnType("text"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.ToTable("StepLeases", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.TimerEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("DueAt") + .HasColumnType("timestamp with time zone"); + + b.Property("FiredAt") + .HasColumnType("timestamp with time zone"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("FiredAt", "DueAt"); + + b.ToTable("Timers", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", null) + .WithOne() + .HasForeignKey("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", "FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => + { + b1.Property("ExecutionHistoryEntityId") + .HasColumnType("bigint"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Payload_TypeName"); + + b1.HasKey("ExecutionHistoryEntityId"); + + b1.ToTable("ExecutionHistories"); + + b1.WithOwner() + .HasForeignKey("ExecutionHistoryEntityId"); + }); + + b.Navigation("Payload"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowDefinitionEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Definition", b1 => + { + b1.Property("FlowDefinitionEntityFlowName") + .HasColumnType("character varying(255)"); + + b1.Property("FlowDefinitionEntityFlowVersion") + .HasColumnType("character varying(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Definition_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Definition_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Definition_TypeName"); + + b1.HasKey("FlowDefinitionEntityFlowName", "FlowDefinitionEntityFlowVersion"); + + b1.ToTable("FlowDefinitions"); + + b1.WithOwner() + .HasForeignKey("FlowDefinitionEntityFlowName", "FlowDefinitionEntityFlowVersion"); + }); + + b.Navigation("Definition"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowInstanceEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Result", b1 => + { + b1.Property("FlowInstanceEntityInstanceId") + .HasColumnType("character varying(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Result_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Result_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Result_TypeName"); + + b1.HasKey("FlowInstanceEntityInstanceId"); + + b1.ToTable("FlowInstances"); + + b1.WithOwner() + .HasForeignKey("FlowInstanceEntityInstanceId"); + }); + + b.Navigation("Result"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeDependencyEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", "DependsOn") + .WithMany("Dependents") + .HasForeignKey("FlowInstanceId", "DependsOnId") + .OnDelete(DeleteBehavior.Restrict) + .IsRequired(); + + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", "Node") + .WithMany("Dependencies") + .HasForeignKey("FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Restrict) + .IsRequired(); + + b.Navigation("DependsOn"); + + b.Navigation("Node"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Input", b1 => + { + b1.Property("NodeInstanceEntityFlowInstanceId") + .HasColumnType("character varying(255)"); + + b1.Property("NodeInstanceEntityNodeId") + .HasColumnType("character varying(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Input_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Input_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Input_TypeName"); + + b1.HasKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + + b1.ToTable("NodeInstances"); + + b1.WithOwner() + .HasForeignKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + }); + + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Result", b1 => + { + b1.Property("NodeInstanceEntityFlowInstanceId") + .HasColumnType("character varying(255)"); + + b1.Property("NodeInstanceEntityNodeId") + .HasColumnType("character varying(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Result_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Result_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Result_TypeName"); + + b1.HasKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + + b1.ToTable("NodeInstances"); + + b1.WithOwner() + .HasForeignKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + }); + + b.Navigation("Input"); + + b.Navigation("Result"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => + { + b1.Property("SignalEntityId") + .HasColumnType("integer"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("bytea") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("text") + .HasColumnName("Payload_TypeName"); + + b1.HasKey("SignalEntityId"); + + b1.ToTable("Signals"); + + b1.WithOwner() + .HasForeignKey("SignalEntityId"); + }); + + b.Navigation("Payload"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.Navigation("Dependencies"); + + b.Navigation("Dependents"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/20260825072319_AddConditionWaits.cs b/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/20260825072319_AddConditionWaits.cs new file mode 100644 index 0000000..fe820cd --- /dev/null +++ b/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/20260825072319_AddConditionWaits.cs @@ -0,0 +1,48 @@ +using System; +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace Spindle.Persistence.EFCore.PostgreSQL.Migrations +{ + /// + public partial class AddConditionWaits : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "ConditionWaits", + columns: table => new + { + FlowInstanceId = table.Column(type: "character varying(255)", maxLength: 255, nullable: false), + NodeId = table.Column(type: "character varying(255)", maxLength: 255, nullable: false), + PollingIntervalTicks = table.Column(type: "bigint", nullable: false), + ExpiresAt = table.Column(type: "timestamp with time zone", nullable: true), + CreatedAt = table.Column(type: "timestamp with time zone", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_ConditionWaits", x => new { x.FlowInstanceId, x.NodeId }); + table.ForeignKey( + name: "FK_ConditionWaits_NodeInstances_FlowInstanceId_NodeId", + columns: x => new { x.FlowInstanceId, x.NodeId }, + principalTable: "NodeInstances", + principalColumns: new[] { "FlowInstanceId", "NodeId" }, + onDelete: ReferentialAction.Cascade); + }); + + migrationBuilder.CreateIndex( + name: "IX_ConditionWaits_ExpiresAt", + table: "ConditionWaits", + column: "ExpiresAt"); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "ConditionWaits"); + } + } +} diff --git a/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/SpindleDbContextModelSnapshot.cs b/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/SpindleDbContextModelSnapshot.cs index c7bc48c..199a413 100644 --- a/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/SpindleDbContextModelSnapshot.cs +++ b/src/Spindle.Persistence.EFCore.PostgreSQL/Migrations/SpindleDbContextModelSnapshot.cs @@ -23,6 +23,32 @@ protected override void BuildModel(ModelBuilder modelBuilder) NpgsqlModelBuilderExtensions.UseIdentityByDefaultColumns(modelBuilder); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("character varying(255)"); + + b.Property("CreatedAt") + .HasColumnType("timestamp with time zone"); + + b.Property("ExpiresAt") + .HasColumnType("timestamp with time zone"); + + b.Property("PollingIntervalTicks") + .HasColumnType("bigint"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("ExpiresAt"); + + b.ToTable("ConditionWaits", (string)null); + }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => { b.Property("Id") @@ -475,6 +501,15 @@ protected override void BuildModel(ModelBuilder modelBuilder) b.ToTable("Timers", (string)null); }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", null) + .WithOne() + .HasForeignKey("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", "FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => { b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => diff --git a/src/Spindle.Persistence.EFCore.SqlServer/Migrations/20260825072326_AddConditionWaits.Designer.cs b/src/Spindle.Persistence.EFCore.SqlServer/Migrations/20260825072326_AddConditionWaits.Designer.cs new file mode 100644 index 0000000..8938bf5 --- /dev/null +++ b/src/Spindle.Persistence.EFCore.SqlServer/Migrations/20260825072326_AddConditionWaits.Designer.cs @@ -0,0 +1,749 @@ +// +using System; +using System.Collections.Generic; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Metadata; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Spindle.Persistence.EFCore; + +#nullable disable + +namespace Spindle.Persistence.EFCore.SqlServer.Migrations +{ + [DbContext(typeof(SpindleDbContext))] + [Migration("20260825072326_AddConditionWaits")] + partial class AddConditionWaits + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder + .HasAnnotation("ProductVersion", "10.0.10") + .HasAnnotation("Relational:MaxIdentifierLength", 128); + + SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("ExpiresAt") + .HasColumnType("datetimeoffset"); + + b.Property("PollingIntervalTicks") + .HasColumnType("bigint"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("ExpiresAt"); + + b.ToTable("ConditionWaits", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("bigint"); + + SqlServerPropertyBuilderExtensions.UseIdentityColumn(b.Property("Id")); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("EventType") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("FlowInstanceId") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("NodeId") + .HasColumnType("nvarchar(max)"); + + b.HasKey("Id"); + + b.ToTable("ExecutionHistories", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowDefinitionEntity", b => + { + b.Property("FlowName") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("FlowVersion") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("DefinitionHash") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("FlowTypeName") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("UpdatedAt") + .HasColumnType("datetimeoffset"); + + b.HasKey("FlowName", "FlowVersion"); + + b.ToTable("FlowDefinitions", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowInstanceEntity", b => + { + b.Property("InstanceId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("CompletedAt") + .HasColumnType("datetimeoffset"); + + b.Property("CorrelationKey") + .HasColumnType("nvarchar(max)"); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("DefinitionHash") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("Error") + .HasColumnType("nvarchar(max)"); + + b.Property("FlowName") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("FlowVersion") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("IdempotencyKey") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("Status") + .HasColumnType("int"); + + b.Property("UpdatedAt") + .HasColumnType("datetimeoffset"); + + b.ComplexProperty(typeof(Dictionary), "Input", "Spindle.Persistence.EFCore.Entities.FlowInstanceEntity.Input#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Input_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Input_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Input_TypeName"); + }); + + b.HasKey("InstanceId"); + + b.HasIndex("FlowName", "IdempotencyKey") + .IsUnique() + .HasFilter("[IdempotencyKey] IS NOT NULL"); + + b.HasIndex("Status", "UpdatedAt"); + + b.ToTable("FlowInstances", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.InboxMessageEntity", b => + { + b.Property("MessageId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("Kind") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("ProcessedAt") + .HasColumnType("datetimeoffset"); + + b.Property("ReceivedAt") + .HasColumnType("datetimeoffset"); + + b.ComplexProperty(typeof(Dictionary), "Payload", "Spindle.Persistence.EFCore.Entities.InboxMessageEntity.Payload#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Payload_TypeName"); + }); + + b.HasKey("MessageId"); + + b.ToTable("InboxMessages", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeDependencyEntity", b => + { + b.Property("FlowInstanceId") + .HasColumnType("nvarchar(255)"); + + b.Property("NodeId") + .HasColumnType("nvarchar(255)"); + + b.Property("DependsOnId") + .HasColumnType("nvarchar(255)"); + + b.Property("Position") + .HasColumnType("int"); + + b.HasKey("FlowInstanceId", "NodeId", "DependsOnId"); + + b.HasIndex("FlowInstanceId", "DependsOnId"); + + b.HasIndex("FlowInstanceId", "NodeId"); + + b.ToTable("NodeDependencies"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("Attempt") + .HasColumnType("int"); + + b.Property("CompletedAt") + .HasColumnType("datetimeoffset"); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("DependencyMode") + .HasColumnType("int"); + + b.Property("DispatchMode") + .HasColumnType("int"); + + b.Property("Error") + .HasColumnType("nvarchar(max)"); + + b.Property("HandlerId") + .HasColumnType("nvarchar(max)"); + + b.Property("Kind") + .HasColumnType("int"); + + b.Property("Name") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("Queue") + .HasColumnType("nvarchar(max)"); + + b.Property("RetryAt") + .HasColumnType("datetimeoffset"); + + b.Property("StartedAt") + .HasColumnType("datetimeoffset"); + + b.Property("Status") + .HasColumnType("int"); + + b.Property("UpdatedAt") + .HasColumnType("datetimeoffset"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("FlowInstanceId", "Status"); + + b.HasIndex("Status", "CreatedAt"); + + b.ToTable("NodeInstances", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.OutboxMessageEntity", b => + { + b.Property("MessageId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("Headers") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("Kind") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("PublishedAt") + .HasColumnType("datetimeoffset"); + + b.ComplexProperty(typeof(Dictionary), "Payload", "Spindle.Persistence.EFCore.Entities.OutboxMessageEntity.Payload#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Payload_TypeName"); + }); + + b.HasKey("MessageId"); + + b.HasIndex("PublishedAt", "CreatedAt"); + + b.ToTable("OutboxMessages", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("int"); + + SqlServerPropertyBuilderExtensions.UseIdentityColumn(b.Property("Id")); + + b.Property("CorrelationKey") + .HasColumnType("nvarchar(max)"); + + b.Property("FlowInstanceId") + .HasColumnType("nvarchar(max)"); + + b.Property("RaisedAt") + .HasColumnType("datetimeoffset"); + + b.Property("SignalName") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.HasKey("Id"); + + b.ToTable("Signals", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("CompletedAt") + .HasColumnType("datetimeoffset"); + + b.Property("CorrelationKey") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("ExpiresAt") + .HasColumnType("datetimeoffset"); + + b.Property("SignalName") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("SignalName", "CorrelationKey", "CompletedAt"); + + b.ToTable("SignalWaits", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.StepAttemptEntity", b => + { + b.Property("AttemptId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("Attempt") + .HasColumnType("int"); + + b.Property("CompletedAt") + .HasColumnType("datetimeoffset"); + + b.Property("Error") + .HasColumnType("nvarchar(max)"); + + b.Property("FlowInstanceId") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("NodeId") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.Property("StartedAt") + .HasColumnType("datetimeoffset"); + + b.Property("Status") + .HasColumnType("int"); + + b.Property("WorkerId") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.HasKey("AttemptId"); + + b.ToTable("StepAttempts", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.StepLeaseEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("AcquiredAt") + .HasColumnType("datetimeoffset"); + + b.Property("ExpiresAt") + .HasColumnType("datetimeoffset"); + + b.Property("Owner") + .IsRequired() + .HasColumnType("nvarchar(max)"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.ToTable("StepLeases", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.TimerEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("DueAt") + .HasColumnType("datetimeoffset"); + + b.Property("FiredAt") + .HasColumnType("datetimeoffset"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("FiredAt", "DueAt"); + + b.ToTable("Timers", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", null) + .WithOne() + .HasForeignKey("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", "FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => + { + b1.Property("ExecutionHistoryEntityId") + .HasColumnType("bigint"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Payload_TypeName"); + + b1.HasKey("ExecutionHistoryEntityId"); + + b1.ToTable("ExecutionHistories"); + + b1.WithOwner() + .HasForeignKey("ExecutionHistoryEntityId"); + }); + + b.Navigation("Payload"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowDefinitionEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Definition", b1 => + { + b1.Property("FlowDefinitionEntityFlowName") + .HasColumnType("nvarchar(255)"); + + b1.Property("FlowDefinitionEntityFlowVersion") + .HasColumnType("nvarchar(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Definition_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Definition_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Definition_TypeName"); + + b1.HasKey("FlowDefinitionEntityFlowName", "FlowDefinitionEntityFlowVersion"); + + b1.ToTable("FlowDefinitions"); + + b1.WithOwner() + .HasForeignKey("FlowDefinitionEntityFlowName", "FlowDefinitionEntityFlowVersion"); + }); + + b.Navigation("Definition"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowInstanceEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Result", b1 => + { + b1.Property("FlowInstanceEntityInstanceId") + .HasColumnType("nvarchar(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Result_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Result_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Result_TypeName"); + + b1.HasKey("FlowInstanceEntityInstanceId"); + + b1.ToTable("FlowInstances"); + + b1.WithOwner() + .HasForeignKey("FlowInstanceEntityInstanceId"); + }); + + b.Navigation("Result"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeDependencyEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", "DependsOn") + .WithMany("Dependents") + .HasForeignKey("FlowInstanceId", "DependsOnId") + .OnDelete(DeleteBehavior.Restrict) + .IsRequired(); + + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", "Node") + .WithMany("Dependencies") + .HasForeignKey("FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Restrict) + .IsRequired(); + + b.Navigation("DependsOn"); + + b.Navigation("Node"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Input", b1 => + { + b1.Property("NodeInstanceEntityFlowInstanceId") + .HasColumnType("nvarchar(255)"); + + b1.Property("NodeInstanceEntityNodeId") + .HasColumnType("nvarchar(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Input_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Input_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Input_TypeName"); + + b1.HasKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + + b1.ToTable("NodeInstances"); + + b1.WithOwner() + .HasForeignKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + }); + + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Result", b1 => + { + b1.Property("NodeInstanceEntityFlowInstanceId") + .HasColumnType("nvarchar(255)"); + + b1.Property("NodeInstanceEntityNodeId") + .HasColumnType("nvarchar(255)"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Result_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Result_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Result_TypeName"); + + b1.HasKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + + b1.ToTable("NodeInstances"); + + b1.WithOwner() + .HasForeignKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + }); + + b.Navigation("Input"); + + b.Navigation("Result"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => + { + b1.Property("SignalEntityId") + .HasColumnType("int"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("varbinary(max)") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("nvarchar(max)") + .HasColumnName("Payload_TypeName"); + + b1.HasKey("SignalEntityId"); + + b1.ToTable("Signals"); + + b1.WithOwner() + .HasForeignKey("SignalEntityId"); + }); + + b.Navigation("Payload"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.Navigation("Dependencies"); + + b.Navigation("Dependents"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/src/Spindle.Persistence.EFCore.SqlServer/Migrations/20260825072326_AddConditionWaits.cs b/src/Spindle.Persistence.EFCore.SqlServer/Migrations/20260825072326_AddConditionWaits.cs new file mode 100644 index 0000000..afb208e --- /dev/null +++ b/src/Spindle.Persistence.EFCore.SqlServer/Migrations/20260825072326_AddConditionWaits.cs @@ -0,0 +1,48 @@ +using System; +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace Spindle.Persistence.EFCore.SqlServer.Migrations +{ + /// + public partial class AddConditionWaits : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "ConditionWaits", + columns: table => new + { + FlowInstanceId = table.Column(type: "nvarchar(255)", maxLength: 255, nullable: false), + NodeId = table.Column(type: "nvarchar(255)", maxLength: 255, nullable: false), + PollingIntervalTicks = table.Column(type: "bigint", nullable: false), + ExpiresAt = table.Column(type: "datetimeoffset", nullable: true), + CreatedAt = table.Column(type: "datetimeoffset", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_ConditionWaits", x => new { x.FlowInstanceId, x.NodeId }); + table.ForeignKey( + name: "FK_ConditionWaits_NodeInstances_FlowInstanceId_NodeId", + columns: x => new { x.FlowInstanceId, x.NodeId }, + principalTable: "NodeInstances", + principalColumns: new[] { "FlowInstanceId", "NodeId" }, + onDelete: ReferentialAction.Cascade); + }); + + migrationBuilder.CreateIndex( + name: "IX_ConditionWaits_ExpiresAt", + table: "ConditionWaits", + column: "ExpiresAt"); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "ConditionWaits"); + } + } +} diff --git a/src/Spindle.Persistence.EFCore.SqlServer/Migrations/SpindleDbContextModelSnapshot.cs b/src/Spindle.Persistence.EFCore.SqlServer/Migrations/SpindleDbContextModelSnapshot.cs index 7c75ad0..38b9026 100644 --- a/src/Spindle.Persistence.EFCore.SqlServer/Migrations/SpindleDbContextModelSnapshot.cs +++ b/src/Spindle.Persistence.EFCore.SqlServer/Migrations/SpindleDbContextModelSnapshot.cs @@ -23,6 +23,32 @@ protected override void BuildModel(ModelBuilder modelBuilder) SqlServerModelBuilderExtensions.UseIdentityColumns(modelBuilder); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("nvarchar(255)"); + + b.Property("CreatedAt") + .HasColumnType("datetimeoffset"); + + b.Property("ExpiresAt") + .HasColumnType("datetimeoffset"); + + b.Property("PollingIntervalTicks") + .HasColumnType("bigint"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("ExpiresAt"); + + b.ToTable("ConditionWaits", (string)null); + }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => { b.Property("Id") @@ -476,6 +502,15 @@ protected override void BuildModel(ModelBuilder modelBuilder) b.ToTable("Timers", (string)null); }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", null) + .WithOne() + .HasForeignKey("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", "FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => { b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => diff --git a/src/Spindle.Persistence.EFCore.Sqlite/Migrations/20260825072315_AddConditionWaits.Designer.cs b/src/Spindle.Persistence.EFCore.Sqlite/Migrations/20260825072315_AddConditionWaits.Designer.cs new file mode 100644 index 0000000..01480ac --- /dev/null +++ b/src/Spindle.Persistence.EFCore.Sqlite/Migrations/20260825072315_AddConditionWaits.Designer.cs @@ -0,0 +1,738 @@ +// +using System.Collections.Generic; +using Microsoft.EntityFrameworkCore; +using Microsoft.EntityFrameworkCore.Infrastructure; +using Microsoft.EntityFrameworkCore.Migrations; +using Microsoft.EntityFrameworkCore.Storage.ValueConversion; +using Spindle.Persistence.EFCore; + +#nullable disable + +namespace Spindle.Persistence.EFCore.Sqlite.Migrations +{ + [DbContext(typeof(SpindleDbContext))] + [Migration("20260825072315_AddConditionWaits")] + partial class AddConditionWaits + { + /// + protected override void BuildTargetModel(ModelBuilder modelBuilder) + { +#pragma warning disable 612, 618 + modelBuilder.HasAnnotation("ProductVersion", "10.0.10"); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("ExpiresAt") + .HasColumnType("INTEGER"); + + b.Property("PollingIntervalTicks") + .HasColumnType("INTEGER"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("ExpiresAt"); + + b.ToTable("ConditionWaits", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("INTEGER"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("EventType") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("FlowInstanceId") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("NodeId") + .HasColumnType("TEXT"); + + b.HasKey("Id"); + + b.ToTable("ExecutionHistories", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowDefinitionEntity", b => + { + b.Property("FlowName") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("FlowVersion") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("DefinitionHash") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("FlowTypeName") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("UpdatedAt") + .HasColumnType("INTEGER"); + + b.HasKey("FlowName", "FlowVersion"); + + b.ToTable("FlowDefinitions", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowInstanceEntity", b => + { + b.Property("InstanceId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("CompletedAt") + .HasColumnType("INTEGER"); + + b.Property("CorrelationKey") + .HasColumnType("TEXT"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("DefinitionHash") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("Error") + .HasColumnType("TEXT"); + + b.Property("FlowName") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("FlowVersion") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("IdempotencyKey") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("Status") + .HasColumnType("INTEGER"); + + b.Property("UpdatedAt") + .HasColumnType("INTEGER"); + + b.ComplexProperty(typeof(Dictionary), "Input", "Spindle.Persistence.EFCore.Entities.FlowInstanceEntity.Input#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Input_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Input_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Input_TypeName"); + }); + + b.HasKey("InstanceId"); + + b.HasIndex("FlowName", "IdempotencyKey") + .IsUnique(); + + b.HasIndex("Status", "UpdatedAt"); + + b.ToTable("FlowInstances", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.InboxMessageEntity", b => + { + b.Property("MessageId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("Kind") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("ProcessedAt") + .HasColumnType("INTEGER"); + + b.Property("ReceivedAt") + .HasColumnType("INTEGER"); + + b.ComplexProperty(typeof(Dictionary), "Payload", "Spindle.Persistence.EFCore.Entities.InboxMessageEntity.Payload#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Payload_TypeName"); + }); + + b.HasKey("MessageId"); + + b.ToTable("InboxMessages", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeDependencyEntity", b => + { + b.Property("FlowInstanceId") + .HasColumnType("TEXT"); + + b.Property("NodeId") + .HasColumnType("TEXT"); + + b.Property("DependsOnId") + .HasColumnType("TEXT"); + + b.Property("Position") + .HasColumnType("INTEGER"); + + b.HasKey("FlowInstanceId", "NodeId", "DependsOnId"); + + b.HasIndex("FlowInstanceId", "DependsOnId"); + + b.HasIndex("FlowInstanceId", "NodeId"); + + b.ToTable("NodeDependencies"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("Attempt") + .HasColumnType("INTEGER"); + + b.Property("CompletedAt") + .HasColumnType("INTEGER"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("DependencyMode") + .HasColumnType("INTEGER"); + + b.Property("DispatchMode") + .HasColumnType("INTEGER"); + + b.Property("Error") + .HasColumnType("TEXT"); + + b.Property("HandlerId") + .HasColumnType("TEXT"); + + b.Property("Kind") + .HasColumnType("INTEGER"); + + b.Property("Name") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("Queue") + .HasColumnType("TEXT"); + + b.Property("RetryAt") + .HasColumnType("INTEGER"); + + b.Property("StartedAt") + .HasColumnType("INTEGER"); + + b.Property("Status") + .HasColumnType("INTEGER"); + + b.Property("UpdatedAt") + .HasColumnType("INTEGER"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("FlowInstanceId", "Status"); + + b.HasIndex("Status", "CreatedAt"); + + b.ToTable("NodeInstances", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.OutboxMessageEntity", b => + { + b.Property("MessageId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("Headers") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("Kind") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("PublishedAt") + .HasColumnType("INTEGER"); + + b.ComplexProperty(typeof(Dictionary), "Payload", "Spindle.Persistence.EFCore.Entities.OutboxMessageEntity.Payload#SerializedPayload", b1 => + { + b1.IsRequired(); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Payload_TypeName"); + }); + + b.HasKey("MessageId"); + + b.HasIndex("PublishedAt", "CreatedAt"); + + b.ToTable("OutboxMessages", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalEntity", b => + { + b.Property("Id") + .ValueGeneratedOnAdd() + .HasColumnType("INTEGER"); + + b.Property("CorrelationKey") + .HasColumnType("TEXT"); + + b.Property("FlowInstanceId") + .HasColumnType("TEXT"); + + b.Property("RaisedAt") + .HasColumnType("INTEGER"); + + b.Property("SignalName") + .IsRequired() + .HasColumnType("TEXT"); + + b.HasKey("Id"); + + b.ToTable("Signals", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("CompletedAt") + .HasColumnType("INTEGER"); + + b.Property("CorrelationKey") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("ExpiresAt") + .HasColumnType("INTEGER"); + + b.Property("SignalName") + .IsRequired() + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("SignalName", "CorrelationKey", "CompletedAt"); + + b.ToTable("SignalWaits", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.StepAttemptEntity", b => + { + b.Property("AttemptId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("Attempt") + .HasColumnType("INTEGER"); + + b.Property("CompletedAt") + .HasColumnType("INTEGER"); + + b.Property("Error") + .HasColumnType("TEXT"); + + b.Property("FlowInstanceId") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("NodeId") + .IsRequired() + .HasColumnType("TEXT"); + + b.Property("StartedAt") + .HasColumnType("INTEGER"); + + b.Property("Status") + .HasColumnType("INTEGER"); + + b.Property("WorkerId") + .IsRequired() + .HasColumnType("TEXT"); + + b.HasKey("AttemptId"); + + b.ToTable("StepAttempts", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.StepLeaseEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("AcquiredAt") + .HasColumnType("INTEGER"); + + b.Property("ExpiresAt") + .HasColumnType("INTEGER"); + + b.Property("Owner") + .IsRequired() + .HasColumnType("TEXT"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.ToTable("StepLeases", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.TimerEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("DueAt") + .HasColumnType("INTEGER"); + + b.Property("FiredAt") + .HasColumnType("INTEGER"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("FiredAt", "DueAt"); + + b.ToTable("Timers", (string)null); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", null) + .WithOne() + .HasForeignKey("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", "FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => + { + b1.Property("ExecutionHistoryEntityId") + .HasColumnType("INTEGER"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Payload_TypeName"); + + b1.HasKey("ExecutionHistoryEntityId"); + + b1.ToTable("ExecutionHistories"); + + b1.WithOwner() + .HasForeignKey("ExecutionHistoryEntityId"); + }); + + b.Navigation("Payload"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowDefinitionEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Definition", b1 => + { + b1.Property("FlowDefinitionEntityFlowName") + .HasColumnType("TEXT"); + + b1.Property("FlowDefinitionEntityFlowVersion") + .HasColumnType("TEXT"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Definition_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Definition_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Definition_TypeName"); + + b1.HasKey("FlowDefinitionEntityFlowName", "FlowDefinitionEntityFlowVersion"); + + b1.ToTable("FlowDefinitions"); + + b1.WithOwner() + .HasForeignKey("FlowDefinitionEntityFlowName", "FlowDefinitionEntityFlowVersion"); + }); + + b.Navigation("Definition"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.FlowInstanceEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Result", b1 => + { + b1.Property("FlowInstanceEntityInstanceId") + .HasColumnType("TEXT"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Result_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Result_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Result_TypeName"); + + b1.HasKey("FlowInstanceEntityInstanceId"); + + b1.ToTable("FlowInstances"); + + b1.WithOwner() + .HasForeignKey("FlowInstanceEntityInstanceId"); + }); + + b.Navigation("Result"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeDependencyEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", "DependsOn") + .WithMany("Dependents") + .HasForeignKey("FlowInstanceId", "DependsOnId") + .OnDelete(DeleteBehavior.Restrict) + .IsRequired(); + + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", "Node") + .WithMany("Dependencies") + .HasForeignKey("FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Restrict) + .IsRequired(); + + b.Navigation("DependsOn"); + + b.Navigation("Node"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Input", b1 => + { + b1.Property("NodeInstanceEntityFlowInstanceId") + .HasColumnType("TEXT"); + + b1.Property("NodeInstanceEntityNodeId") + .HasColumnType("TEXT"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Input_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Input_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Input_TypeName"); + + b1.HasKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + + b1.ToTable("NodeInstances"); + + b1.WithOwner() + .HasForeignKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + }); + + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Result", b1 => + { + b1.Property("NodeInstanceEntityFlowInstanceId") + .HasColumnType("TEXT"); + + b1.Property("NodeInstanceEntityNodeId") + .HasColumnType("TEXT"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Result_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Result_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Result_TypeName"); + + b1.HasKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + + b1.ToTable("NodeInstances"); + + b1.WithOwner() + .HasForeignKey("NodeInstanceEntityFlowInstanceId", "NodeInstanceEntityNodeId"); + }); + + b.Navigation("Input"); + + b.Navigation("Result"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.SignalEntity", b => + { + b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => + { + b1.Property("SignalEntityId") + .HasColumnType("INTEGER"); + + b1.Property("ContentType") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Payload_ContentType"); + + b1.Property("Data") + .IsRequired() + .HasColumnType("BLOB") + .HasColumnName("Payload_Data"); + + b1.Property("TypeName") + .IsRequired() + .HasColumnType("TEXT") + .HasColumnName("Payload_TypeName"); + + b1.HasKey("SignalEntityId"); + + b1.ToTable("Signals"); + + b1.WithOwner() + .HasForeignKey("SignalEntityId"); + }); + + b.Navigation("Payload"); + }); + + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", b => + { + b.Navigation("Dependencies"); + + b.Navigation("Dependents"); + }); +#pragma warning restore 612, 618 + } + } +} diff --git a/src/Spindle.Persistence.EFCore.Sqlite/Migrations/20260825072315_AddConditionWaits.cs b/src/Spindle.Persistence.EFCore.Sqlite/Migrations/20260825072315_AddConditionWaits.cs new file mode 100644 index 0000000..591e19d --- /dev/null +++ b/src/Spindle.Persistence.EFCore.Sqlite/Migrations/20260825072315_AddConditionWaits.cs @@ -0,0 +1,47 @@ +using Microsoft.EntityFrameworkCore.Migrations; + +#nullable disable + +namespace Spindle.Persistence.EFCore.Sqlite.Migrations +{ + /// + public partial class AddConditionWaits : Migration + { + /// + protected override void Up(MigrationBuilder migrationBuilder) + { + migrationBuilder.CreateTable( + name: "ConditionWaits", + columns: table => new + { + FlowInstanceId = table.Column(type: "TEXT", maxLength: 255, nullable: false), + NodeId = table.Column(type: "TEXT", maxLength: 255, nullable: false), + PollingIntervalTicks = table.Column(type: "INTEGER", nullable: false), + ExpiresAt = table.Column(type: "INTEGER", nullable: true), + CreatedAt = table.Column(type: "INTEGER", nullable: false) + }, + constraints: table => + { + table.PrimaryKey("PK_ConditionWaits", x => new { x.FlowInstanceId, x.NodeId }); + table.ForeignKey( + name: "FK_ConditionWaits_NodeInstances_FlowInstanceId_NodeId", + columns: x => new { x.FlowInstanceId, x.NodeId }, + principalTable: "NodeInstances", + principalColumns: new[] { "FlowInstanceId", "NodeId" }, + onDelete: ReferentialAction.Cascade); + }); + + migrationBuilder.CreateIndex( + name: "IX_ConditionWaits_ExpiresAt", + table: "ConditionWaits", + column: "ExpiresAt"); + } + + /// + protected override void Down(MigrationBuilder migrationBuilder) + { + migrationBuilder.DropTable( + name: "ConditionWaits"); + } + } +} diff --git a/src/Spindle.Persistence.EFCore.Sqlite/Migrations/SpindleDbContextModelSnapshot.cs b/src/Spindle.Persistence.EFCore.Sqlite/Migrations/SpindleDbContextModelSnapshot.cs index b77fe10..0a15e7a 100644 --- a/src/Spindle.Persistence.EFCore.Sqlite/Migrations/SpindleDbContextModelSnapshot.cs +++ b/src/Spindle.Persistence.EFCore.Sqlite/Migrations/SpindleDbContextModelSnapshot.cs @@ -17,6 +17,32 @@ protected override void BuildModel(ModelBuilder modelBuilder) #pragma warning disable 612, 618 modelBuilder.HasAnnotation("ProductVersion", "10.0.10"); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.Property("FlowInstanceId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("NodeId") + .HasMaxLength(255) + .HasColumnType("TEXT"); + + b.Property("CreatedAt") + .HasColumnType("INTEGER"); + + b.Property("ExpiresAt") + .HasColumnType("INTEGER"); + + b.Property("PollingIntervalTicks") + .HasColumnType("INTEGER"); + + b.HasKey("FlowInstanceId", "NodeId"); + + b.HasIndex("ExpiresAt"); + + b.ToTable("ConditionWaits", (string)null); + }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => { b.Property("Id") @@ -465,6 +491,15 @@ protected override void BuildModel(ModelBuilder modelBuilder) b.ToTable("Timers", (string)null); }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", b => + { + b.HasOne("Spindle.Persistence.EFCore.Entities.NodeInstanceEntity", null) + .WithOne() + .HasForeignKey("Spindle.Persistence.EFCore.Entities.ConditionWaitEntity", "FlowInstanceId", "NodeId") + .OnDelete(DeleteBehavior.Cascade) + .IsRequired(); + }); + modelBuilder.Entity("Spindle.Persistence.EFCore.Entities.ExecutionHistoryEntity", b => { b.OwnsOne("Spindle.Abstractions.Snapshot.SerializedPayload", "Payload", b1 => diff --git a/src/Spindle.Persistence.EFCore/EFCoreSpindleStore.cs b/src/Spindle.Persistence.EFCore/EFCoreSpindleStore.cs index 5c59865..3df7b9a 100644 --- a/src/Spindle.Persistence.EFCore/EFCoreSpindleStore.cs +++ b/src/Spindle.Persistence.EFCore/EFCoreSpindleStore.cs @@ -1,5 +1,6 @@ using Microsoft.EntityFrameworkCore; using Spindle.Persistence.FlowDefinitions; +using Spindle.Persistence.Conditions; using Spindle.Persistence.FlowInstances; using Spindle.Persistence.History; using Spindle.Persistence.Leases; @@ -31,6 +32,8 @@ public EFCoreSpindleStore( public ITimerStore Timers => _rootSession.Timers; + public IConditionWaitStore Conditions => _rootSession.Conditions; + public ISignalStore Signals => _rootSession.Signals; public IOutboxStore Outbox => _rootSession.Outbox; @@ -89,6 +92,7 @@ public ContextStoreSession(SpindleDbContext context) FlowInstances = new EFCoreFlowInstanceStore(context); Nodes = new EFCoreNodeStore(context); Timers = new EFCoreTimerStore(context); + Conditions = new EFCoreConditionWaitStore(context); Signals = new EFCoreSignalStore(context); Outbox = new EFCoreOutboxStore(context); Inbox = new EFCoreInboxStore(context); @@ -100,6 +104,7 @@ public ContextStoreSession(SpindleDbContext context) public IFlowInstanceStore FlowInstances { get; } public INodeStore Nodes { get; } public ITimerStore Timers { get; } + public IConditionWaitStore Conditions { get; } public ISignalStore Signals { get; } public IOutboxStore Outbox { get; } public IInboxStore Inbox { get; } diff --git a/src/Spindle.Persistence.EFCore/Entities/ConditionWaitEntity.cs b/src/Spindle.Persistence.EFCore/Entities/ConditionWaitEntity.cs new file mode 100644 index 0000000..a200d44 --- /dev/null +++ b/src/Spindle.Persistence.EFCore/Entities/ConditionWaitEntity.cs @@ -0,0 +1,21 @@ +using Microsoft.EntityFrameworkCore; +using System.ComponentModel.DataAnnotations; + +namespace Spindle.Persistence.EFCore.Entities; + +[PrimaryKey(nameof(FlowInstanceId), nameof(NodeId))] +[Index(nameof(ExpiresAt))] +internal sealed class ConditionWaitEntity +{ + [MaxLength(255)] + public required string FlowInstanceId { get; init; } + + [MaxLength(255)] + public required string NodeId { get; init; } + + public required long PollingIntervalTicks { get; init; } + + public DateTimeOffset? ExpiresAt { get; init; } + + public required DateTimeOffset CreatedAt { get; init; } +} diff --git a/src/Spindle.Persistence.EFCore/FactoryStoreSession.cs b/src/Spindle.Persistence.EFCore/FactoryStoreSession.cs index b8b17c7..7585869 100644 --- a/src/Spindle.Persistence.EFCore/FactoryStoreSession.cs +++ b/src/Spindle.Persistence.EFCore/FactoryStoreSession.cs @@ -2,6 +2,7 @@ using Spindle.Abstractions.Core; using Spindle.Abstractions.Snapshot; using Spindle.Persistence.FlowDefinitions; +using Spindle.Persistence.Conditions; using Spindle.Persistence.FlowInstances; using Spindle.Persistence.History; using Spindle.Persistence.Leases; @@ -20,6 +21,7 @@ public FactoryStoreSession(EFCoreSpindleStore store) FlowInstances = new FlowInstanceStore(store); Nodes = new NodeStore(store); Timers = new TimerStore(store); + Conditions = new ConditionWaitStore(store); Signals = new SignalStore(store); Outbox = new OutboxStore(store); Inbox = new InboxStore(store); @@ -31,6 +33,7 @@ public FactoryStoreSession(EFCoreSpindleStore store) public IFlowInstanceStore FlowInstances { get; } public INodeStore Nodes { get; } public ITimerStore Timers { get; } + public IConditionWaitStore Conditions { get; } public ISignalStore Signals { get; } public IOutboxStore Outbox { get; } public IInboxStore Inbox { get; } @@ -99,8 +102,11 @@ public ValueTask MarkReadyAsync(FlowInstanceId flowInstanceId, NodeId nodeId, Da public ValueTask MarkRunningAsync(FlowInstanceId flowInstanceId, NodeId nodeId, StepAttemptId attemptId, string workerId, DateTimeOffset startedAt, CancellationToken cancellationToken = default) => store.ExecuteDirectAsync((session, token) => session.Nodes.MarkRunningAsync(flowInstanceId, nodeId, attemptId, workerId, startedAt, token), cancellationToken); - public ValueTask MarkWaitingAsync(FlowInstanceId flowInstanceId, NodeId nodeId, DateTimeOffset updatedAt, CancellationToken cancellationToken = default) => - store.ExecuteDirectAsync((session, token) => session.Nodes.MarkWaitingAsync(flowInstanceId, nodeId, updatedAt, token), cancellationToken); + public ValueTask MarkWaitingAsync(FlowInstanceId flowInstanceId, NodeId nodeId, int attempt, DateTimeOffset updatedAt, CancellationToken cancellationToken = default) => + store.ExecuteDirectAsync((session, token) => session.Nodes.MarkWaitingAsync(flowInstanceId, nodeId, attempt, updatedAt, token), cancellationToken); + + public ValueTask MarkTimedOutAsync(FlowInstanceId flowInstanceId, NodeId nodeId, int attempt, string error, DateTimeOffset timedOutAt, CancellationToken cancellationToken = default) => + store.ExecuteDirectAsync((session, token) => session.Nodes.MarkTimedOutAsync(flowInstanceId, nodeId, attempt, error, timedOutAt, token), cancellationToken); public ValueTask MarkCompletedAsync(FlowInstanceId flowInstanceId, NodeId nodeId, int attempt, SerializedPayload? result, DateTimeOffset completedAt, CancellationToken cancellationToken = default) => store.ExecuteDirectAsync((session, token) => session.Nodes.MarkCompletedAsync(flowInstanceId, nodeId, attempt, result, completedAt, token), cancellationToken); @@ -127,6 +133,15 @@ public ValueTask MarkFiredAsync(FlowInstanceId flowInstanceId, NodeId nodeId, Da store.ExecuteDirectAsync((session, token) => session.Timers.MarkFiredAsync(flowInstanceId, nodeId, firedAt, token), cancellationToken); } + private sealed class ConditionWaitStore(EFCoreSpindleStore store) : IConditionWaitStore + { + public ValueTask CreateAsync(ConditionWaitRecord wait, CancellationToken cancellationToken = default) => + store.ExecuteDirectAsync((session, token) => session.Conditions.CreateAsync(wait, token), cancellationToken); + + public ValueTask GetAsync(FlowInstanceId flowInstanceId, NodeId nodeId, CancellationToken cancellationToken = default) => + store.ExecuteAsync((session, token) => session.Conditions.GetAsync(flowInstanceId, nodeId, token), cancellationToken); + } + private sealed class SignalStore(EFCoreSpindleStore store) : ISignalStore { public ValueTask CreateWaitAsync(SignalWaitRecord wait, CancellationToken cancellationToken = default) => diff --git a/src/Spindle.Persistence.EFCore/SpindleDbContext.cs b/src/Spindle.Persistence.EFCore/SpindleDbContext.cs index 03a630e..9617811 100644 --- a/src/Spindle.Persistence.EFCore/SpindleDbContext.cs +++ b/src/Spindle.Persistence.EFCore/SpindleDbContext.cs @@ -11,6 +11,7 @@ public sealed class SpindleDbContext( : DbContext(options) { internal DbSet ExecutionHistories => Set(); + internal DbSet ConditionWaits => Set(); internal DbSet FlowDefinitions => Set(); internal DbSet FlowInstances => Set(); internal DbSet InboxMessages => Set(); @@ -62,6 +63,15 @@ protected override void OnModelCreating(ModelBuilder modelBuilder) payload => ConfigurePayload(payload, "Payload")); }); + modelBuilder.Entity(entity => + { + entity.ToTable("ConditionWaits"); + entity.HasOne() + .WithOne() + .HasForeignKey(wait => new { wait.FlowInstanceId, wait.NodeId }) + .OnDelete(DeleteBehavior.Cascade); + }); + modelBuilder.Entity(entity => { entity.ToTable("FlowDefinitions"); diff --git a/src/Spindle.Persistence.EFCore/Stores/EFCoreConditionWaitStore.cs b/src/Spindle.Persistence.EFCore/Stores/EFCoreConditionWaitStore.cs new file mode 100644 index 0000000..832061b --- /dev/null +++ b/src/Spindle.Persistence.EFCore/Stores/EFCoreConditionWaitStore.cs @@ -0,0 +1,63 @@ +using Microsoft.EntityFrameworkCore; +using Spindle.Abstractions.Core; +using Spindle.Persistence.Conditions; +using Spindle.Persistence.EFCore.Entities; + +namespace Spindle.Persistence.EFCore.Stores; + +internal sealed class EFCoreConditionWaitStore(SpindleDbContext context) : IConditionWaitStore +{ + public async ValueTask CreateAsync( + ConditionWaitRecord wait, + CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + if (await context.ConditionWaits.AnyAsync( + existing => existing.FlowInstanceId == wait.FlowInstanceId.Value && + existing.NodeId == wait.NodeId.Value, + cancellationToken)) + { + throw new InvalidOperationException( + $"Condition wait '{wait.NodeId}' already exists for flow instance '{wait.FlowInstanceId}'."); + } + + await context.ConditionWaits.AddAsync( + new ConditionWaitEntity + { + FlowInstanceId = wait.FlowInstanceId.Value, + NodeId = wait.NodeId.Value, + PollingIntervalTicks = wait.PollingInterval.Ticks, + ExpiresAt = wait.ExpiresAt, + CreatedAt = wait.CreatedAt + }, + cancellationToken); + await context.SaveChangesAsync(cancellationToken); + } + + public async ValueTask GetAsync( + FlowInstanceId flowInstanceId, + NodeId nodeId, + CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + var entity = await context.ConditionWaits + .AsNoTracking() + .FirstOrDefaultAsync( + wait => wait.FlowInstanceId == flowInstanceId.Value && + wait.NodeId == nodeId.Value, + cancellationToken); + + return entity is null + ? null + : new ConditionWaitRecord + { + FlowInstanceId = new FlowInstanceId(entity.FlowInstanceId), + NodeId = new NodeId(entity.NodeId), + PollingInterval = TimeSpan.FromTicks(entity.PollingIntervalTicks), + ExpiresAt = entity.ExpiresAt, + CreatedAt = entity.CreatedAt + }; + } +} diff --git a/src/Spindle.Persistence.EFCore/Stores/EFCoreNodeStore.cs b/src/Spindle.Persistence.EFCore/Stores/EFCoreNodeStore.cs index 3e9383e..5ff0f9e 100644 --- a/src/Spindle.Persistence.EFCore/Stores/EFCoreNodeStore.cs +++ b/src/Spindle.Persistence.EFCore/Stores/EFCoreNodeStore.cs @@ -222,7 +222,11 @@ public async ValueTask MarkReadyAsync( if (currentStatus == null) throw new InvalidOperationException( $"Node '{nodeId}' does not exist for flow instance '{flowInstanceId}'."); - var isTerminal = currentStatus is NodeStatus.Completed or NodeStatus.Failed or NodeStatus.Cancelled; + var isTerminal = currentStatus is NodeStatus.Completed + or NodeStatus.Failed + or NodeStatus.Cancelled + or NodeStatus.TimedOut + or NodeStatus.Skipped; if (isTerminal) return; // It does exist and is not terminal, so we can set it as ready @@ -277,6 +281,7 @@ await context.StepAttempts.AddAsync(new StepAttemptEntity public async ValueTask MarkWaitingAsync( FlowInstanceId flowInstanceId, NodeId nodeId, + int attempt, DateTimeOffset updatedAt, CancellationToken cancellationToken = default) { @@ -295,6 +300,52 @@ await context.NodeInstances .SetProperty(x => x.Status, _ => NodeStatus.Waiting) .SetProperty(x => x.UpdatedAt, _ => updatedAt) , cancellationToken); + + await CompleteAttemptAsync( + flowInstanceId, + nodeId, + attempt, + NodeStatus.Waiting, + updatedAt, + null, + cancellationToken); + await context.SaveChangesAsync(cancellationToken); + } + + public async ValueTask MarkTimedOutAsync( + FlowInstanceId flowInstanceId, + NodeId nodeId, + int attempt, + string error, + DateTimeOffset timedOutAt, + CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + var updated = await context.NodeInstances + .Where(x => x.FlowInstanceId == flowInstanceId.Value && x.NodeId == nodeId.Value) + .ExecuteUpdateAsync(update => update + .SetProperty(x => x.Status, _ => NodeStatus.TimedOut) + .SetProperty(x => x.Error, _ => error) + .SetProperty(x => x.CompletedAt, _ => timedOutAt) + .SetProperty(x => x.UpdatedAt, _ => timedOutAt), + cancellationToken); + + if (updated == 0) + { + throw new InvalidOperationException( + $"Node '{nodeId}' does not exist for flow instance '{flowInstanceId}'."); + } + + await CompleteAttemptAsync( + flowInstanceId, + nodeId, + attempt, + NodeStatus.TimedOut, + timedOutAt, + error, + cancellationToken); + await context.SaveChangesAsync(cancellationToken); } public async ValueTask MarkCompletedAsync( @@ -368,7 +419,9 @@ public async ValueTask MarkDependentsReadyAsync( var q = context.NodeInstances .Where(x => x.FlowInstanceId == flowInstanceId.Value) - .Where(x => x.Kind == NodeKind.Step && x.Status == NodeStatus.Pending); + .Where(x => + (x.Kind == NodeKind.Step || x.Kind == NodeKind.ConditionWait) && + x.Status == NodeStatus.Pending); if (updatedNodes is { Count: > 0 }) { diff --git a/src/Spindle.Persistence.InMemory/InMemorySpindleStore.cs b/src/Spindle.Persistence.InMemory/InMemorySpindleStore.cs index d5b1cc8..ac86c70 100644 --- a/src/Spindle.Persistence.InMemory/InMemorySpindleStore.cs +++ b/src/Spindle.Persistence.InMemory/InMemorySpindleStore.cs @@ -1,5 +1,6 @@ using Spindle.Persistence; using Spindle.Persistence.FlowDefinitions; +using Spindle.Persistence.Conditions; using Spindle.Persistence.FlowInstances; using Spindle.Persistence.History; using Spindle.Persistence.Leases; @@ -29,6 +30,8 @@ public InMemorySpindleStore() public ITimerStore Timers => _session.Timers; + public IConditionWaitStore Conditions => _session.Conditions; + public ISignalStore Signals => _session.Signals; public IOutboxStore Outbox => _session.Outbox; @@ -84,6 +87,7 @@ public InMemorySpindleStoreSession() FlowInstances = new InMemoryFlowInstanceStore(); Nodes = new InMemoryNodeStore(); Timers = new InMemoryTimerStore(); + Conditions = new InMemoryConditionWaitStore(); Signals = new InMemorySignalStore(); Outbox = new InMemoryOutboxStore(); Inbox = new InMemoryInboxStore(); @@ -99,6 +103,8 @@ public InMemorySpindleStoreSession() public ITimerStore Timers { get; } + public IConditionWaitStore Conditions { get; } + public ISignalStore Signals { get; } public IOutboxStore Outbox { get; } diff --git a/src/Spindle.Persistence.InMemory/Stores/InMemoryConditionWaitStore.cs b/src/Spindle.Persistence.InMemory/Stores/InMemoryConditionWaitStore.cs new file mode 100644 index 0000000..b72bd34 --- /dev/null +++ b/src/Spindle.Persistence.InMemory/Stores/InMemoryConditionWaitStore.cs @@ -0,0 +1,52 @@ +using Spindle.Abstractions.Core; +using Spindle.Persistence.Conditions; + +namespace Spindle.Persistence.InMemory.Stores; + +/// +/// Stores durable condition metadata in memory. +/// +public sealed class InMemoryConditionWaitStore : IConditionWaitStore +{ + private readonly object _gate = new(); + private readonly Dictionary<(FlowInstanceId FlowInstanceId, NodeId NodeId), ConditionWaitRecord> _waits = []; + + /// + public ValueTask CreateAsync( + ConditionWaitRecord wait, + CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + lock (_gate) + { + var key = (wait.FlowInstanceId, wait.NodeId); + if (_waits.ContainsKey(key)) + { + throw new InvalidOperationException( + $"Condition wait '{wait.NodeId}' already exists for flow instance '{wait.FlowInstanceId}'."); + } + + _waits.Add(key, wait with { }); + } + + return ValueTask.CompletedTask; + } + + /// + public ValueTask GetAsync( + FlowInstanceId flowInstanceId, + NodeId nodeId, + CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + lock (_gate) + { + return ValueTask.FromResult( + _waits.TryGetValue((flowInstanceId, nodeId), out var wait) + ? wait with { } + : null); + } + } +} diff --git a/src/Spindle.Persistence.InMemory/Stores/InMemoryNodeStore.cs b/src/Spindle.Persistence.InMemory/Stores/InMemoryNodeStore.cs index f87eb79..5883182 100644 --- a/src/Spindle.Persistence.InMemory/Stores/InMemoryNodeStore.cs +++ b/src/Spindle.Persistence.InMemory/Stores/InMemoryNodeStore.cs @@ -155,7 +155,11 @@ public ValueTask MarkReadyAsync( { var node = GetRequired(flowInstanceId, nodeId); - if (node.Status is NodeStatus.Completed or NodeStatus.Failed or NodeStatus.Cancelled) + if (node.Status is NodeStatus.Completed + or NodeStatus.Failed + or NodeStatus.Cancelled + or NodeStatus.TimedOut + or NodeStatus.Skipped) { return ValueTask.CompletedTask; } @@ -211,6 +215,7 @@ public ValueTask MarkRunningAsync( public ValueTask MarkWaitingAsync( FlowInstanceId flowInstanceId, NodeId nodeId, + int attempt, DateTimeOffset updatedAt, CancellationToken cancellationToken = default) { @@ -224,6 +229,35 @@ public ValueTask MarkWaitingAsync( Status = NodeStatus.Waiting, UpdatedAt = updatedAt }; + + CompleteLatestAttempt(flowInstanceId, nodeId, attempt, NodeStatus.Waiting, updatedAt, null); + } + + return ValueTask.CompletedTask; + } + + public ValueTask MarkTimedOutAsync( + FlowInstanceId flowInstanceId, + NodeId nodeId, + int attempt, + string error, + DateTimeOffset timedOutAt, + CancellationToken cancellationToken = default) + { + cancellationToken.ThrowIfCancellationRequested(); + + lock (_gate) + { + var node = GetRequired(flowInstanceId, nodeId); + _nodes[(flowInstanceId, nodeId)] = node with + { + Status = NodeStatus.TimedOut, + Error = error, + CompletedAt = timedOutAt, + UpdatedAt = timedOutAt + }; + + CompleteLatestAttempt(flowInstanceId, nodeId, attempt, NodeStatus.TimedOut, timedOutAt, error); } return ValueTask.CompletedTask; @@ -309,7 +343,8 @@ public ValueTask MarkDependentsReadyAsync( foreach (var node in nodes) { - if (node.Kind != NodeKind.Step || node.Status != NodeStatus.Pending) + if (node.Kind is not (NodeKind.Step or NodeKind.ConditionWait) || + node.Status != NodeStatus.Pending) { continue; } diff --git a/src/Spindle.Runtime/Execution/LocalStepExecutor.cs b/src/Spindle.Runtime/Execution/LocalStepExecutor.cs index da7344a..2245a53 100644 --- a/src/Spindle.Runtime/Execution/LocalStepExecutor.cs +++ b/src/Spindle.Runtime/Execution/LocalStepExecutor.cs @@ -6,6 +6,7 @@ using Spindle.Persistence; using Spindle.Persistence.Leases; using Spindle.Persistence.Nodes; +using Spindle.Persistence.Timers; using Spindle.Runtime; using System.Diagnostics; @@ -162,6 +163,22 @@ await storeSession.Nodes .ConfigureAwait(false); } + if (registration.Kind == NodeKind.ConditionWait) + { + if (result is not bool conditionSatisfied) + { + throw new InvalidOperationException( + $"Condition node '{step.NodeId}' did not return a Boolean result."); + } + + return await CompleteConditionCheckAsync( + session, + running, + conditionSatisfied, + cancellationToken) + .ConfigureAwait(false); + } + await ConcurrencyHelper.AquireLock(step.FlowInstanceId); await store .ExecuteAsync( @@ -171,7 +188,7 @@ await storeSession.Nodes .MarkCompletedAsync( step.FlowInstanceId, step.NodeId, - step.Attempt, + running.Attempt, SerializerReflection.Serialize(serializer, result, registration.ResultType), timeProvider.GetUtcNow(), storeCancellationToken) @@ -213,7 +230,7 @@ await store storeSession.Nodes.MarkFailedAsync( step.FlowInstanceId, step.NodeId, - step.Attempt, + running.Attempt, exception.Message, timeProvider.GetUtcNow(), retryAt: null, @@ -241,6 +258,118 @@ await store return StepExecutionResult.Failed; } + private async Task CompleteConditionCheckAsync( + FlowExecutionSession session, + NodeInstanceRecord condition, + bool satisfied, + CancellationToken cancellationToken) + { + var wait = await store.Conditions + .GetAsync(condition.FlowInstanceId, condition.NodeId, cancellationToken) + .ConfigureAwait(false) + ?? throw new InvalidOperationException( + $"Condition metadata for node '{condition.NodeId}' does not exist."); + var now = timeProvider.GetUtcNow(); + + if (wait.ExpiresAt is { } expiresAt && now >= expiresAt) + { + var error = $"Condition node '{condition.NodeId}' timed out at {expiresAt:O}."; + await store.ExecuteAsync( + async (storeSession, storeCancellationToken) => + { + await storeSession.Nodes.MarkTimedOutAsync( + condition.FlowInstanceId, + condition.NodeId, + condition.Attempt, + error, + now, + storeCancellationToken) + .ConfigureAwait(false); + await storeSession.Leases.ReleaseStepLeaseAsync( + condition.FlowInstanceId, + condition.NodeId, + workerId, + storeCancellationToken) + .ConfigureAwait(false); + }, + cancellationToken) + .ConfigureAwait(false); + + return StepExecutionResult.TimedOut; + } + + if (satisfied) + { + await store.ExecuteAsync( + async (storeSession, storeCancellationToken) => + { + await storeSession.Nodes.MarkCompletedAsync( + condition.FlowInstanceId, + condition.NodeId, + condition.Attempt, + result: null, + now, + storeCancellationToken) + .ConfigureAwait(false); + await storeSession.Nodes.MarkDependentsReadyAsync( + condition.FlowInstanceId, + [condition.NodeId], + now, + storeCancellationToken) + .ConfigureAwait(false); + await storeSession.Leases.ReleaseStepLeaseAsync( + condition.FlowInstanceId, + condition.NodeId, + workerId, + storeCancellationToken) + .ConfigureAwait(false); + }, + cancellationToken) + .ConfigureAwait(false); + + session.SetResult(condition.NodeId, Unit.Value); + return StepExecutionResult.Succeeded; + } + + var dueAt = now.Add(wait.PollingInterval); + if (wait.ExpiresAt is { } deadline && deadline < dueAt) + { + dueAt = deadline; + } + + await store.ExecuteAsync( + async (storeSession, storeCancellationToken) => + { + await storeSession.Nodes.MarkWaitingAsync( + condition.FlowInstanceId, + condition.NodeId, + condition.Attempt, + now, + storeCancellationToken) + .ConfigureAwait(false); + await storeSession.Timers.CreateAsync( + new TimerRecord + { + FlowInstanceId = condition.FlowInstanceId, + NodeId = condition.NodeId, + DueAt = dueAt, + CreatedAt = now + }, + storeCancellationToken) + .ConfigureAwait(false); + await storeSession.Leases.ReleaseStepLeaseAsync( + condition.FlowInstanceId, + condition.NodeId, + workerId, + storeCancellationToken) + .ConfigureAwait(false); + }, + cancellationToken) + .ConfigureAwait(false); + + return StepExecutionResult.Waiting; + } + private async ValueTask BuildInputsAsync( FlowExecutionSession session, Persistence.Nodes.NodeInstanceRecord step, diff --git a/src/Spindle.Runtime/Execution/StepExecutionRegistration.cs b/src/Spindle.Runtime/Execution/StepExecutionRegistration.cs index 11ecb09..173b0cd 100644 --- a/src/Spindle.Runtime/Execution/StepExecutionRegistration.cs +++ b/src/Spindle.Runtime/Execution/StepExecutionRegistration.cs @@ -6,6 +6,7 @@ namespace Spindle; internal sealed record StepExecutionRegistration( NodeId NodeId, + NodeKind Kind, Type ResultType, IReadOnlyList DependencyResultTypes, Func> Execute); diff --git a/src/Spindle.Runtime/Execution/StepExecutionResult.cs b/src/Spindle.Runtime/Execution/StepExecutionResult.cs index ba84cbc..a19f882 100644 --- a/src/Spindle.Runtime/Execution/StepExecutionResult.cs +++ b/src/Spindle.Runtime/Execution/StepExecutionResult.cs @@ -1,18 +1,28 @@ +using Spindle.Abstractions.Snapshot; + namespace Spindle; internal readonly record struct StepExecutionResult( bool Executed, - bool Completed) + NodeStatus Status) { public static StepExecutionResult NotExecuted { get; } = new( Executed: false, - Completed: false); + NodeStatus.Ready); public static StepExecutionResult Failed { get; } = new( Executed: true, - Completed: false); + NodeStatus.Failed); public static StepExecutionResult Succeeded { get; } = new( Executed: true, - Completed: true); + NodeStatus.Completed); + + public static StepExecutionResult Waiting { get; } = new( + Executed: true, + NodeStatus.Waiting); + + public static StepExecutionResult TimedOut { get; } = new( + Executed: true, + NodeStatus.TimedOut); } diff --git a/src/Spindle.Runtime/Execution/StepExecutor.cs b/src/Spindle.Runtime/Execution/StepExecutor.cs index 0cdc1e8..13a4e7d 100644 --- a/src/Spindle.Runtime/Execution/StepExecutor.cs +++ b/src/Spindle.Runtime/Execution/StepExecutor.cs @@ -52,7 +52,7 @@ public async ValueTask ExecuteReadyStepsAsync( var remaining = maxCount - attempted.Count; var readySteps = steps.Values .Where(step => - step.Kind == NodeKind.Step && + (step.Kind is NodeKind.Step or NodeKind.ConditionWait) && step.Status == NodeStatus.Ready && !attempted.Contains(step.NodeId)) .OrderBy(step => step.CreatedAt) @@ -93,19 +93,17 @@ public async ValueTask ExecuteReadyStepsAsync( executed++; - if (!result.Completed) + steps[step.NodeId] = step with { Status = result.Status }; + session.UpsertNode(steps[step.NodeId]); + + if (result.Status != NodeStatus.Completed) { - steps[step.NodeId] = step with { Status = NodeStatus.Failed }; - session.UpsertNode(steps[step.NodeId]); continue; } - steps[step.NodeId] = step with { Status = NodeStatus.Completed }; - session.UpsertNode(steps[step.NodeId]); - foreach (var dependent in steps.Values .Where(candidate => - candidate.Kind == NodeKind.Step && + (candidate.Kind is NodeKind.Step or NodeKind.ConditionWait) && candidate.Status == NodeStatus.Pending && candidate.Dependencies.All(dependency => steps.TryGetValue(dependency, out var prerequisite) && diff --git a/src/Spindle.Runtime/FlowExecutionSession.cs b/src/Spindle.Runtime/FlowExecutionSession.cs index c26d725..61071d6 100644 --- a/src/Spindle.Runtime/FlowExecutionSession.cs +++ b/src/Spindle.Runtime/FlowExecutionSession.cs @@ -2,6 +2,7 @@ using Spindle.Abstractions.Snapshot; using Spindle.Abstractions.Nodes; using Spindle.Abstractions.Steps; +using Spindle.Abstractions.Waiting; using Spindle.Persistence.Nodes; namespace Spindle; @@ -49,11 +50,30 @@ public void Register( _registrations[nodeId] = new StepExecutionRegistration( nodeId, + NodeKind.Step, typeof(TResult), dependencyResultTypes.ToArray(), Execute); } + public void RegisterCondition( + NodeId nodeId, + IReadOnlyList dependencyResultTypes, + ConditionCallback callback) + { + async ValueTask Execute(NodeInputs inputs, IStepExecutionContext context) + { + return await callback(inputs, context).ConfigureAwait(false); + } + + _registrations[nodeId] = new StepExecutionRegistration( + nodeId, + NodeKind.ConditionWait, + typeof(bool), + dependencyResultTypes.ToArray(), + Execute); + } + public bool TryGet( NodeId nodeId, out StepExecutionRegistration registration) @@ -146,6 +166,26 @@ public IReadOnlyList GetPendingNodeDeclarations() public IReadOnlyList GetPendingNodeInitializations() => _pendingInitializations.Values.ToArray(); + public bool TryConfigureConditionTimeout( + NodeId nodeId, + TimeSpan timeout) + { + if (!_pendingInitializations.TryGetValue(nodeId, out var initialization) || + initialization is not ConditionNodeInitialization condition) + { + return false; + } + + _pendingInitializations[nodeId] = condition with + { + ConditionWait = condition.ConditionWait with + { + ExpiresAt = condition.ConditionWait.CreatedAt.Add(timeout) + } + }; + return true; + } + public void MarkNodeDeclarationsFlushed() { _pendingNodeDeclarations.Clear(); diff --git a/src/Spindle.Runtime/FlowExecutor.cs b/src/Spindle.Runtime/FlowExecutor.cs index 077e571..9d21b8a 100644 --- a/src/Spindle.Runtime/FlowExecutor.cs +++ b/src/Spindle.Runtime/FlowExecutor.cs @@ -167,6 +167,10 @@ await storeSession.Timers.CreateAsync(timer.Timer, cancellationToken) await storeSession.Signals.CreateWaitAsync(signal.SignalWait, cancellationToken) .ConfigureAwait(false); break; + case ConditionNodeInitialization condition: + await storeSession.Conditions.CreateAsync(condition.ConditionWait, cancellationToken) + .ConfigureAwait(false); + break; } } diff --git a/src/Spindle.Runtime/NamespacedFlowContextWrapper.cs b/src/Spindle.Runtime/NamespacedFlowContextWrapper.cs index 1e13bb8..00d4973 100644 --- a/src/Spindle.Runtime/NamespacedFlowContextWrapper.cs +++ b/src/Spindle.Runtime/NamespacedFlowContextWrapper.cs @@ -74,4 +74,7 @@ public SignalNode WaitForSignal(string id, string name, Signal public SignalNode WaitForSignal(string id, SignalName signalName, CorrelationKey correlationKey, SignalWaitOptions? options = null) => parent.WaitForSignal(WrapId(id), signalName, correlationKey, options); + public ConditionNode WaitForCondition(string id, string name, TimeSpan pollingInterval, IReadOnlyList dependencies, ConditionCallback condition) + => parent.WaitForCondition(WrapId(id), name, pollingInterval, dependencies, condition); + } diff --git a/src/Spindle.Runtime/Nodes/ConditionNodeInitialization.cs b/src/Spindle.Runtime/Nodes/ConditionNodeInitialization.cs new file mode 100644 index 0000000..bf06525 --- /dev/null +++ b/src/Spindle.Runtime/Nodes/ConditionNodeInitialization.cs @@ -0,0 +1,5 @@ +using Spindle.Persistence.Conditions; + +namespace Spindle; + +internal sealed record ConditionNodeInitialization(ConditionWaitRecord ConditionWait) : NodeInitialization; diff --git a/src/Spindle.Runtime/Nodes/RuntimeConditionNode.cs b/src/Spindle.Runtime/Nodes/RuntimeConditionNode.cs new file mode 100644 index 0000000..09f7e13 --- /dev/null +++ b/src/Spindle.Runtime/Nodes/RuntimeConditionNode.cs @@ -0,0 +1,40 @@ +using Spindle.Abstractions.Core; +using Spindle.Abstractions.Nodes; + +namespace Spindle; + +internal sealed class RuntimeConditionNode( + RuntimeNodeState state, + TimeSpan pollingInterval, + TimeSpan? timeout, + Func configureTimeout) + : ConditionNode, IRuntimeNode +{ + public override NodeId Id => state.Id; + + public override string Name => state.Name; + + public override NodeKind Kind => state.Kind; + + public override TimeSpan PollingInterval => pollingInterval; + + public override TimeSpan? Timeout => timeout; + + public Type ResultType => typeof(Unit); + + public FlowInstanceId FlowInstanceId => state.FlowInstanceId; + + public override ValueTask GetResultAsync(CancellationToken cancellationToken = default) + => state.GetResultAsync(cancellationToken); + + public override ConditionNode WithTimeout(TimeSpan value) + { + if (value <= TimeSpan.Zero) + { + throw new ArgumentOutOfRangeException(nameof(value)); + } + + configureTimeout(value); + return new RuntimeConditionNode(state, pollingInterval, value, configureTimeout); + } +} diff --git a/src/Spindle.Runtime/Nodes/RuntimeNodeState.cs b/src/Spindle.Runtime/Nodes/RuntimeNodeState.cs index 9d52bec..7fa1ef4 100644 --- a/src/Spindle.Runtime/Nodes/RuntimeNodeState.cs +++ b/src/Spindle.Runtime/Nodes/RuntimeNodeState.cs @@ -46,7 +46,9 @@ record = await Store.Nodes : Serializer.Deserialize(record.Result), NodeStatus.Failed => throw new InvalidOperationException( $"Node '{Id}' failed: {record.Error}"), - NodeStatus.TimedOut or NodeStatus.Cancelled => throw new TaskCanceledException(), + NodeStatus.TimedOut => throw new TimeoutException( + $"Node '{Id}' timed out: {record.Error}"), + NodeStatus.Cancelled => throw new TaskCanceledException(), _ => throw new FlowSuspendedException() }; } diff --git a/src/Spindle.Runtime/RuntimeFlowContext.cs b/src/Spindle.Runtime/RuntimeFlowContext.cs index 74fb56e..c5ca1f1 100644 --- a/src/Spindle.Runtime/RuntimeFlowContext.cs +++ b/src/Spindle.Runtime/RuntimeFlowContext.cs @@ -5,6 +5,7 @@ using Spindle.Abstractions.Snapshot; using Spindle.Abstractions.Waiting; using Spindle.Persistence; +using Spindle.Persistence.Conditions; using Spindle.Persistence.Nodes; using Spindle.Persistence.Signals; using Spindle.Persistence.Timers; @@ -248,6 +249,63 @@ public SignalNode WaitForSignal( correlationKey); } + public ConditionNode WaitForCondition( + string id, + string name, + TimeSpan pollingInterval, + IReadOnlyList dependencies, + ConditionCallback condition) + { + cancellationToken.ThrowIfCancellationRequested(); + ArgumentNullException.ThrowIfNull(dependencies); + ArgumentNullException.ThrowIfNull(condition); + if (pollingInterval <= TimeSpan.Zero) + { + throw new ArgumentOutOfRangeException(nameof(pollingInterval)); + } + + ValidateDependencies(dependencies); + var nodeId = new NodeId(id); + var dependencyIds = dependencies.Select(dependency => dependency.Id).ToArray(); + var dependencyResultTypes = dependencies.Select(GetDependencyResultType).ToArray(); + session.RegisterCondition(nodeId, dependencyResultTypes, condition); + + if (!session.TryGetNode(nodeId, out _)) + { + var now = timeProvider.GetUtcNow(); + session.TryDeclareNode( + new NodeInstanceRecord + { + FlowInstanceId = InstanceId, + NodeId = nodeId, + Name = name, + Kind = NodeKind.ConditionWait, + Status = DependenciesSucceeded(dependencyIds) + ? NodeStatus.Ready + : NodeStatus.Pending, + DispatchMode = StepDispatchMode.LocalWorker, + DependencyMode = DependencySatisfactionMode.AllSucceeded, + Dependencies = dependencyIds, + CreatedAt = now, + UpdatedAt = now + }, + new ConditionNodeInitialization( + new ConditionWaitRecord + { + FlowInstanceId = InstanceId, + NodeId = nodeId, + PollingInterval = pollingInterval, + CreatedAt = now + })); + } + + return new RuntimeConditionNode( + CreateState(nodeId, name, NodeKind.ConditionWait), + pollingInterval, + timeout: null, + timeout => session.TryConfigureConditionTimeout(nodeId, timeout)); + } + public ForkNode Fork(string id, Func> descriptor) { var nsWrapper = new NamespacedFlowContextWrapper(id, this); diff --git a/src/Spindle.Runtime/RuntimeSpindleRuntime.cs b/src/Spindle.Runtime/RuntimeSpindleRuntime.cs index 01981c4..4f004c6 100644 --- a/src/Spindle.Runtime/RuntimeSpindleRuntime.cs +++ b/src/Spindle.Runtime/RuntimeSpindleRuntime.cs @@ -4,6 +4,7 @@ using Spindle.Abstractions.Core; using Spindle.Abstractions.Exceptions; using Spindle.Abstractions.Flows; +using Spindle.Abstractions.Nodes; using Spindle.Abstractions.Snapshot; using Spindle.Persistence; using Spindle.Persistence.FlowDefinitions; @@ -339,6 +340,45 @@ await storeSession.Timers continue; } + if (step.Kind == NodeKind.ConditionWait) + { + var condition = await storeSession.Conditions + .GetAsync(timer.FlowInstanceId, timer.NodeId, storeCancellationToken) + .ConfigureAwait(false) + ?? throw new InvalidOperationException( + $"Condition metadata for node '{timer.NodeId}' does not exist."); + + if (condition.ExpiresAt is { } expiresAt && now >= expiresAt) + { + await storeSession.Nodes.MarkTimedOutAsync( + timer.FlowInstanceId, + timer.NodeId, + step.Attempt, + $"Condition node '{timer.NodeId}' timed out at {expiresAt:O}.", + now, + storeCancellationToken) + .ConfigureAwait(false); + } + else + { + await storeSession.Nodes.MarkReadyAsync( + timer.FlowInstanceId, + timer.NodeId, + now, + storeCancellationToken) + .ConfigureAwait(false); + } + + await storeSession.Timers.MarkFiredAsync( + timer.FlowInstanceId, + timer.NodeId, + now, + storeCancellationToken) + .ConfigureAwait(false); + fired++; + continue; + } + await storeSession.Nodes .MarkCompletedAsync(timer.FlowInstanceId, timer.NodeId, -1, result: null, now, storeCancellationToken) .ConfigureAwait(false); diff --git a/src/Spindle/Nodes/FlowContextConditionExtensions.cs b/src/Spindle/Nodes/FlowContextConditionExtensions.cs new file mode 100644 index 0000000..225cc8f --- /dev/null +++ b/src/Spindle/Nodes/FlowContextConditionExtensions.cs @@ -0,0 +1,176 @@ +using Spindle.Abstractions.Flows; +using Spindle.Abstractions.Nodes; +using Spindle.Abstractions.Steps; + +namespace Spindle; + +/// +/// Provides typed convenience overloads for declaring durable condition waits. +/// +public static class FlowContextConditionExtensions +{ + /// Declares an unnamed condition without node inputs. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + TimeSpan pollingInterval, + Func> condition) + => WaitForCondition(context, id, id, pollingInterval, condition); + + /// Declares a named condition without node inputs. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + string name, + TimeSpan pollingInterval, + Func> condition) + { + ArgumentNullException.ThrowIfNull(condition); + return context.WaitForCondition( + id, + name, + pollingInterval, + [], + (_, _) => condition()); + } + + /// Declares an unnamed condition that receives an execution context. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + TimeSpan pollingInterval, + Func> condition) + => WaitForCondition(context, id, id, pollingInterval, condition); + + /// Declares a named condition that receives an execution context. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + string name, + TimeSpan pollingInterval, + Func> condition) + { + ArgumentNullException.ThrowIfNull(condition); + return context.WaitForCondition( + id, + name, + pollingInterval, + [], + (_, executionContext) => condition(executionContext)); + } + + /// Declares an unnamed condition with one typed node input. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + TimeSpan pollingInterval, + Node input1, + Func> condition) + => WaitForCondition(context, id, id, pollingInterval, input1, condition); + + /// Declares a named condition with one typed node input. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + string name, + TimeSpan pollingInterval, + Node input1, + Func> condition) + { + ArgumentNullException.ThrowIfNull(condition); + return context.WaitForCondition( + id, + name, + pollingInterval, + [input1], + (inputs, _) => condition(inputs.Get(0))); + } + + /// Declares an unnamed condition with one typed input and an execution context. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + TimeSpan pollingInterval, + Node input1, + Func> condition) + => WaitForCondition(context, id, id, pollingInterval, input1, condition); + + /// Declares a named condition with one typed input and an execution context. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + string name, + TimeSpan pollingInterval, + Node input1, + Func> condition) + { + ArgumentNullException.ThrowIfNull(condition); + return context.WaitForCondition( + id, + name, + pollingInterval, + [input1], + (inputs, executionContext) => condition(inputs.Get(0), executionContext)); + } + + /// Declares an unnamed condition with two typed node inputs. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + TimeSpan pollingInterval, + Node input1, + Node input2, + Func> condition) + => WaitForCondition(context, id, id, pollingInterval, input1, input2, condition); + + /// Declares a named condition with two typed node inputs. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + string name, + TimeSpan pollingInterval, + Node input1, + Node input2, + Func> condition) + { + ArgumentNullException.ThrowIfNull(condition); + return context.WaitForCondition( + id, + name, + pollingInterval, + [input1, input2], + (inputs, _) => condition(inputs.Get(0), inputs.Get(1))); + } + + /// Declares an unnamed condition with two typed inputs and an execution context. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + TimeSpan pollingInterval, + Node input1, + Node input2, + Func> condition) + => WaitForCondition(context, id, id, pollingInterval, input1, input2, condition); + + /// Declares a named condition with two typed inputs and an execution context. + public static ConditionNode WaitForCondition( + this IFlowContext context, + string id, + string name, + TimeSpan pollingInterval, + Node input1, + Node input2, + Func> condition) + { + ArgumentNullException.ThrowIfNull(condition); + return context.WaitForCondition( + id, + name, + pollingInterval, + [input1, input2], + (inputs, executionContext) => condition( + inputs.Get(0), + inputs.Get(1), + executionContext)); + } +} diff --git a/tests/Spindle.Persistence.EFCore.Tests/EntityFrameworkPersistenceTests.cs b/tests/Spindle.Persistence.EFCore.Tests/EntityFrameworkPersistenceTests.cs index c3d5de3..85ed198 100644 --- a/tests/Spindle.Persistence.EFCore.Tests/EntityFrameworkPersistenceTests.cs +++ b/tests/Spindle.Persistence.EFCore.Tests/EntityFrameworkPersistenceTests.cs @@ -8,6 +8,7 @@ using Spindle.Abstractions.Nodes; using Spindle.Abstractions.Steps; using Spindle.Persistence.EFCore; +using Spindle.Persistence.Conditions; using Spindle.Persistence.EFCore.Sqlite; using Spindle.Persistence.FlowDefinitions; using Spindle.Persistence.FlowInstances; @@ -23,6 +24,67 @@ namespace Spindle.Persistence.EFCore.Tests; public sealed class EntityFrameworkPersistenceTests { + [Fact] + public async Task ConditionWaits_PersistPollingIntervalAndDeadline() + { + var connectionString = $"Data Source=SpindleConditionWaitTests-{Guid.NewGuid():N};Mode=Memory;Cache=Shared"; + await using var databaseAnchor = new SqliteConnection(connectionString); + await databaseAnchor.OpenAsync(); + + var services = new ServiceCollection(); + services.AddSpindleSqlite(connectionString); + await using var provider = services.BuildServiceProvider(); + var database = provider.GetRequiredService(); + await database.Database.MigrateAsync(); + + var store = provider.GetRequiredService(); + var now = DateTimeOffset.Parse("2026-08-25T12:00:00Z"); + var instanceId = new FlowInstanceId("condition-instance"); + var nodeId = new NodeId("condition"); + + await store.FlowInstances.CreateAsync(new FlowInstanceRecord + { + InstanceId = instanceId, + FlowName = new FlowName("condition-flow"), + FlowVersion = new FlowVersion("1"), + DefinitionHash = "hash", + Status = FlowInstanceStatus.Waiting, + Input = new SerializedPayload + { + ContentType = "application/json", + TypeName = typeof(object).FullName!, + Data = [] + }, + CreatedAt = now, + UpdatedAt = now + }); + await store.Nodes.CreateAsync(new NodeInstanceRecord + { + FlowInstanceId = instanceId, + NodeId = nodeId, + Name = "Condition", + Kind = NodeKind.ConditionWait, + Status = NodeStatus.Waiting, + DispatchMode = StepDispatchMode.LocalWorker, + CreatedAt = now, + UpdatedAt = now + }); + await store.Conditions.CreateAsync(new ConditionWaitRecord + { + FlowInstanceId = instanceId, + NodeId = nodeId, + PollingInterval = TimeSpan.FromMinutes(5), + ExpiresAt = now.AddDays(31), + CreatedAt = now + }); + + var persisted = await store.Conditions.GetAsync(instanceId, nodeId); + + Assert.NotNull(persisted); + Assert.Equal(TimeSpan.FromMinutes(5), persisted.PollingInterval); + Assert.Equal(now.AddDays(31), persisted.ExpiresAt); + } + [Fact] public async Task GenericNodesMigration_PreservesExistingNodesAndDependencies() { diff --git a/tests/Spindle.Runtime.Tests/ConditionWaitTests.cs b/tests/Spindle.Runtime.Tests/ConditionWaitTests.cs new file mode 100644 index 0000000..3f44d82 --- /dev/null +++ b/tests/Spindle.Runtime.Tests/ConditionWaitTests.cs @@ -0,0 +1,268 @@ +using Spindle.Abstractions.Core; +using Spindle.Abstractions.Flows; +using Spindle.Abstractions.Nodes; +using Spindle.Abstractions.Snapshot; +using Xunit; + +namespace Spindle.Runtime.Tests; + +public sealed class ConditionWaitTests : TestBase +{ + [Fact] + public async Task Condition_IsCheckedImmediatelyAndCompletesFlow() + { + var (runtime, store, _) = CreateRuntime(); + var flowName = new FlowName("condition-immediate"); + var options = new StartFlowOptions { IdempotencyKey = "condition-immediate" }; + var checks = 0; + + runtime.RegisterFlow( + flowName, + async (context, _) => + { + await context.WaitForCondition( + "ready", + "Ready", + TimeSpan.FromMinutes(5), + () => ValueTask.FromResult(++checks == 1)); + return new TestResult(42); + }); + + var handle = await runtime.StartAsync( + flowName, + new TestRequest(0), + options); + var result = await runtime.RunAsync( + flowName, + new TestRequest(0), + options); + var node = Assert.Single(await store.Nodes.GetByFlowInstanceAsync(handle.InstanceId)); + + Assert.Equal(new TestResult(42), result); + Assert.Equal(1, checks); + Assert.Equal(NodeKind.ConditionWait, node.Kind); + Assert.Equal(NodeStatus.Completed, node.Status); + } + + [Fact] + public async Task FalseCondition_SchedulesNextCheckAndSurvivesReplay() + { + var initial = DateTimeOffset.Parse("2026-08-25T10:00:00Z"); + var (runtime, store, serializer, clock) = CreateRuntime(initial); + var flowName = new FlowName("condition-polls"); + var options = new StartFlowOptions { IdempotencyKey = "condition-polls" }; + var checks = 0; + + runtime.RegisterFlow( + flowName, + async (context, _) => + { + await context.WaitForCondition( + "ready", + TimeSpan.FromMinutes(5), + () => ValueTask.FromResult(++checks >= 2)); + return new TestResult(checks); + }); + + var handle = await runtime.StartAsync( + flowName, + new TestRequest(0), + options); + + await Assert.ThrowsAsync(() => runtime + .RunAsync(flowName, new TestRequest(0), options) + .AsTask()); + + var timer = await store.Timers.GetAsync(handle.InstanceId, new NodeId("ready")); + var waiting = await store.Nodes.GetAsync(handle.InstanceId, new NodeId("ready")); + + Assert.Equal(1, checks); + Assert.NotNull(timer); + Assert.Equal(initial.AddMinutes(5), timer.DueAt); + Assert.Null(timer.FiredAt); + Assert.Equal(NodeStatus.Waiting, waiting?.Status); + + clock.SetUtcNow(initial.AddMinutes(5)); + var restartedRuntime = new RuntimeSpindleRuntime( + store, + options: new RuntimeSpindleOptions + { + TimeProvider = clock, + Serializer = serializer + }); + restartedRuntime.RegisterFlow( + flowName, + async (context, _) => + { + await context.WaitForCondition( + "ready", + TimeSpan.FromMinutes(5), + () => ValueTask.FromResult(++checks >= 2)); + return new TestResult(checks); + }); + + var result = await restartedRuntime.RunAsync( + flowName, + new TestRequest(0), + options); + + Assert.Equal(new TestResult(2), result); + Assert.Equal(2, checks); + } + + [Fact] + public async Task ConditionTimeout_TimesOutNodeAndFailsFlow() + { + var initial = DateTimeOffset.Parse("2026-08-25T10:00:00Z"); + var (runtime, store, _, clock) = CreateRuntime(initial); + var flowName = new FlowName("condition-timeout"); + var options = new StartFlowOptions { IdempotencyKey = "condition-timeout" }; + var reachedAfterWait = false; + + runtime.RegisterFlow( + flowName, + async (context, _) => + { + await context.WaitForCondition( + "ready", + TimeSpan.FromMinutes(5), + () => ValueTask.FromResult(false)) + .WithTimeout(TimeSpan.FromMinutes(2)); + reachedAfterWait = true; + return new TestResult(1); + }); + + var handle = await runtime.StartAsync( + flowName, + new TestRequest(0), + options); + + await Assert.ThrowsAsync(() => runtime + .RunAsync(flowName, new TestRequest(0), options) + .AsTask()); + + var timer = await store.Timers.GetAsync(handle.InstanceId, new NodeId("ready")); + Assert.Equal(initial.AddMinutes(2), timer?.DueAt); + + clock.SetUtcNow(initial.AddMinutes(2)); + await Assert.ThrowsAsync(() => runtime + .RunAsync(flowName, new TestRequest(0), options) + .AsTask()); + + var node = await store.Nodes.GetAsync(handle.InstanceId, new NodeId("ready")); + var instance = await store.FlowInstances.GetAsync(handle.InstanceId); + Assert.Equal(NodeStatus.TimedOut, node?.Status); + Assert.Equal(FlowInstanceStatus.Failed, instance?.Status); + Assert.False(reachedAfterWait); + } + + [Fact] + public async Task ConditionCallbackException_FailsNodeWithoutPollingAgain() + { + var (runtime, store, _) = CreateRuntime(); + var flowName = new FlowName("condition-error"); + var options = new StartFlowOptions { IdempotencyKey = "condition-error" }; + + runtime.RegisterFlow( + flowName, + async (context, _) => + { + await context.WaitForCondition( + "ready", + TimeSpan.FromMinutes(5), + () => ValueTask.FromException(new InvalidOperationException("check failed"))); + return new TestResult(1); + }); + + var handle = await runtime.StartAsync( + flowName, + new TestRequest(0), + options); + + await Assert.ThrowsAsync(() => runtime + .RunAsync(flowName, new TestRequest(0), options) + .AsTask()); + + var node = await store.Nodes.GetAsync(handle.InstanceId, new NodeId("ready")); + Assert.Equal(NodeStatus.Failed, node?.Status); + Assert.Null(await store.Timers.GetAsync(handle.InstanceId, new NodeId("ready"))); + } + + [Fact] + public async Task TypedConditionInputs_AreMaterializedInDeclarationOrder() + { + var (runtime, _, _) = CreateRuntime(); + var flowName = new FlowName("condition-inputs"); + + runtime.RegisterFlow( + flowName, + async (context, _) => + { + var first = context.Step( + "first", + "First", + () => ValueTask.FromResult(20)); + var second = context.Step( + "second", + "Second", + () => ValueTask.FromResult(22)); + var condition = context.WaitForCondition( + "sum-ready", + TimeSpan.FromMinutes(1), + first, + second, + (left, right) => ValueTask.FromResult(left + right == 42)); + var result = context.Step( + "result", + "Result", + condition, + _ => ValueTask.FromResult(new TestResult(42))); + return await result; + }); + + var result = await runtime.RunAsync( + flowName, + new TestRequest(0)); + + Assert.Equal(new TestResult(42), result); + } + + [Fact] + public async Task ConditionInsideFork_UsesNamespacedNodeId() + { + var (runtime, store, _) = CreateRuntime(); + var flowName = new FlowName("condition-fork"); + var options = new StartFlowOptions { IdempotencyKey = "condition-fork" }; + + runtime.RegisterFlow( + flowName, + async (context, _) => + { + await context.Fork( + "branch", + async branch => + { + await branch.WaitForCondition( + "ready", + TimeSpan.FromMinutes(1), + () => ValueTask.FromResult(true)); + }); + return new TestResult(42); + }); + + var handle = await runtime.StartAsync( + flowName, + new TestRequest(0), + options); + var result = await runtime.RunAsync( + flowName, + new TestRequest(0), + options); + var node = Assert.Single(await store.Nodes.GetByFlowInstanceAsync(handle.InstanceId)); + + Assert.Equal(new TestResult(42), result); + Assert.Equal(new NodeId("branch/ready"), node.NodeId); + Assert.Equal(NodeKind.ConditionWait, node.Kind); + Assert.Equal(NodeStatus.Completed, node.Status); + } +} diff --git a/tests/Spindle.Runtime.Tests/Stores/CountingNodeStore.cs b/tests/Spindle.Runtime.Tests/Stores/CountingNodeStore.cs index 3897aac..eacc99d 100644 --- a/tests/Spindle.Runtime.Tests/Stores/CountingNodeStore.cs +++ b/tests/Spindle.Runtime.Tests/Stores/CountingNodeStore.cs @@ -127,13 +127,23 @@ public ValueTask MarkRunningAsync( public ValueTask MarkWaitingAsync( FlowInstanceId flowInstanceId, NodeId nodeId, + int attempt, DateTimeOffset updatedAt, CancellationToken cancellationToken = default) { MarkWaitingCalls++; - return inner.MarkWaitingAsync(flowInstanceId, nodeId, updatedAt, cancellationToken); + return inner.MarkWaitingAsync(flowInstanceId, nodeId, attempt, updatedAt, cancellationToken); } + public ValueTask MarkTimedOutAsync( + FlowInstanceId flowInstanceId, + NodeId nodeId, + int attempt, + string error, + DateTimeOffset timedOutAt, + CancellationToken cancellationToken = default) + => inner.MarkTimedOutAsync(flowInstanceId, nodeId, attempt, error, timedOutAt, cancellationToken); + public ValueTask MarkCompletedAsync( FlowInstanceId flowInstanceId, NodeId nodeId, diff --git a/tests/Spindle.Runtime.Tests/Stores/CountingSpindleStore.cs b/tests/Spindle.Runtime.Tests/Stores/CountingSpindleStore.cs index b2f2c8b..bac87cc 100644 --- a/tests/Spindle.Runtime.Tests/Stores/CountingSpindleStore.cs +++ b/tests/Spindle.Runtime.Tests/Stores/CountingSpindleStore.cs @@ -1,5 +1,6 @@ using Spindle.Persistence; using Spindle.Persistence.FlowDefinitions; +using Spindle.Persistence.Conditions; using Spindle.Persistence.FlowInstances; using Spindle.Persistence.History; using Spindle.Persistence.Leases; @@ -31,6 +32,8 @@ public CountingSpindleStore( public ITimerStore Timers => _inner.Timers; + public IConditionWaitStore Conditions => _inner.Conditions; + public ISignalStore Signals => _inner.Signals; public IOutboxStore Outbox => _inner.Outbox; diff --git a/tests/Spindle.Runtime.Tests/Stores/CoutingStoreSession.cs b/tests/Spindle.Runtime.Tests/Stores/CoutingStoreSession.cs index 90af066..f22cb22 100644 --- a/tests/Spindle.Runtime.Tests/Stores/CoutingStoreSession.cs +++ b/tests/Spindle.Runtime.Tests/Stores/CoutingStoreSession.cs @@ -1,5 +1,6 @@ using Spindle.Persistence; using Spindle.Persistence.FlowDefinitions; +using Spindle.Persistence.Conditions; using Spindle.Persistence.FlowInstances; using Spindle.Persistence.History; using Spindle.Persistence.Leases; @@ -23,6 +24,8 @@ internal sealed class CountingStoreSession( public ITimerStore Timers => inner.Timers; + public IConditionWaitStore Conditions => inner.Conditions; + public ISignalStore Signals => inner.Signals; public IOutboxStore Outbox => inner.Outbox;