Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions src/Spindle.Abstractions/Flows/IFlowContext.cs
Original file line number Diff line number Diff line change
Expand Up @@ -101,6 +101,16 @@ SignalNode<TSignal> WaitForSignal<TSignal>(
CorrelationKey correlationKey,
SignalWaitOptions? options = null);

/// <summary>
/// Declares a durable condition that is checked immediately and then at the specified interval.
/// </summary>
ConditionNode WaitForCondition(
string id,
string name,
TimeSpan pollingInterval,
IReadOnlyList<Node> dependencies,
ConditionCallback condition);

/// <summary>
/// Forks the execution into a namespace. This allows awaits inside the fork without affecting the entire flow.
/// </summary>
Expand Down
24 changes: 24 additions & 0 deletions src/Spindle.Abstractions/Nodes/ConditionNode.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
using Spindle.Abstractions.Core;

namespace Spindle.Abstractions.Nodes;

/// <summary>
/// Represents a durable wait that polls an asynchronous condition until it succeeds.
/// </summary>
public abstract class ConditionNode : WaitNode<Unit>
{
/// <summary>
/// Gets the delay between unsuccessful condition checks.
/// </summary>
public abstract TimeSpan PollingInterval { get; }

/// <summary>
/// Gets the optional timeout measured from the node's first declaration.
/// </summary>
public abstract TimeSpan? Timeout { get; }

/// <summary>
/// Configures the maximum duration in which the condition may become true.
/// </summary>
public abstract ConditionNode WithTimeout(TimeSpan timeout);
}
5 changes: 5 additions & 0 deletions src/Spindle.Abstractions/Nodes/NodeKind.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
/// </summary>
Fork,

/// <summary>
/// A condition that is evaluated repeatedly until it succeeds or times out.
/// </summary>
ConditionWait,
}
29 changes: 29 additions & 0 deletions src/Spindle.Abstractions/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ It can also be used to only reference the abstractions without needing to includ
- `DelayNode` represents a durable timer.
- `SignalNode<TSignal>` 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.

Expand Down Expand Up @@ -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.
11 changes: 11 additions & 0 deletions src/Spindle.Abstractions/Waiting/ConditionCallback.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
using Spindle.Abstractions.Nodes;
using Spindle.Abstractions.Steps;

namespace Spindle.Abstractions.Waiting;

/// <summary>
/// Evaluates whether a durable condition wait has completed.
/// </summary>
public delegate ValueTask<bool> ConditionCallback(
NodeInputs inputs,
IStepExecutionContext context);
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
using Spindle.Abstractions.Core;

namespace Spindle.Persistence.Conditions;

/// <summary>
/// Describes the durable scheduling metadata for a condition wait.
/// </summary>
public sealed record ConditionWaitRecord
{
/// <summary>Gets the owning flow instance.</summary>
public required FlowInstanceId FlowInstanceId { get; init; }

/// <summary>Gets the condition node identifier.</summary>
public required NodeId NodeId { get; init; }

/// <summary>Gets the delay between unsuccessful checks.</summary>
public required TimeSpan PollingInterval { get; init; }

/// <summary>Gets the optional absolute timeout deadline.</summary>
public DateTimeOffset? ExpiresAt { get; init; }

/// <summary>Gets the time at which the condition was first declared.</summary>
public required DateTimeOffset CreatedAt { get; init; }
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
using Spindle.Abstractions.Core;

namespace Spindle.Persistence.Conditions;

/// <summary>
/// Persists scheduling metadata for durable condition waits.
/// </summary>
public interface IConditionWaitStore
{
/// <summary>Creates condition scheduling metadata.</summary>
ValueTask CreateAsync(
ConditionWaitRecord wait,
CancellationToken cancellationToken = default);

/// <summary>Gets condition scheduling metadata.</summary>
ValueTask<ConditionWaitRecord?> GetAsync(
FlowInstanceId flowInstanceId,
NodeId nodeId,
CancellationToken cancellationToken = default);
}
2 changes: 2 additions & 0 deletions src/Spindle.Persistence.Abstractions/ISpindleStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ public interface ISpindleStore

Timers.ITimerStore Timers { get; }

Conditions.IConditionWaitStore Conditions { get; }

Signals.ISignalStore Signals { get; }

Messaging.IOutboxStore Outbox { get; }
Expand Down
3 changes: 3 additions & 0 deletions src/Spindle.Persistence.Abstractions/ISpindleStoreSession.cs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
using Spindle.Persistence.FlowDefinitions;
using Spindle.Persistence.Conditions;
using Spindle.Persistence.FlowInstances;
using Spindle.Persistence.History;
using Spindle.Persistence.Leases;
Expand All @@ -19,6 +20,8 @@ public interface ISpindleStoreSession

ITimerStore Timers { get; }

IConditionWaitStore Conditions { get; }

ISignalStore Signals { get; }

IOutboxStore Outbox { get; }
Expand Down
9 changes: 9 additions & 0 deletions src/Spindle.Persistence.Abstractions/Nodes/INodeStore.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading