WorkflowContext
in package
Straight-line deterministic operations available while a workflow Fiber is replayed.
Table of Contents
Constants
- MAX_PARALLEL_OPERATIONS : mixed = 1000
- MESSAGE_STREAM_CURSOR_SCHEMA : mixed = 'durable-workflow.v2.message-stream.cursor'
- MESSAGE_STREAM_SCHEMA : mixed = 'durable-workflow.v2.message-stream.message'
- MESSAGE_STREAM_SIGNAL : mixed = '__durable_workflow_message_stream'
- MAX_VERSION : mixed = 2147483647
- MIN_VERSION : mixed = -2147483648
Properties
- $runId : string
- $workflowId : string
- $cancellationRequested : bool
- $captureFrames : array<int, array<int, DeferredWorkflowOperation|ParallelWorkflowCommand>>
- $codec : PayloadCodec
- $execution : Fiber<mixed, mixed>|null
- $history : array<string|int, mixed>
- $localActivityExecutor : Closure|null
- $messageStreamCursors : array<string, int>
- $messageStreamMessages : array<string, array<int, MessageStreamMessage>>
- $messageStreamWaits : array<string, int>
- $workflowCommandId : string|null
- $workflowStreamCommandOrdinal : int
Methods
- __construct() : mixed
- activity() : mixed
- all() : array<int, mixed>
- Schedule every deferred leaf, then return results in declaration order.
- appendWorkflowStream() : void
- Append a replay-safe batch to a named run-scoped Workflow Stream.
- childWorkflow() : mixed
- closeWorkflowStream() : void
- consumeMessageStreamMessages() : array<int, MessageStreamMessage>
- continueAsNew() : never
- deferActivity() : DeferredWorkflowOperation
- Prepare an activity without scheduling it until an all/parallel barrier is reached.
- deferChildWorkflow() : DeferredWorkflowOperation
- Prepare a child workflow without starting it until an all/parallel barrier is reached.
- deferCondition() : DeferredWorkflowOperation
- Prepare a deterministic condition or signal-derived wait for a durable group.
- deferTimer() : DeferredWorkflowOperation
- Prepare a durable timer for an all/parallel barrier.
- deprecatePatch() : void
- Keep a patch marker alive after the legacy branch has been removed.
- errorWorkflowStream() : void
- getVersion() : int
- Select the newest supported version for a change, or replay its recorded decision.
- hasPendingMessageStreamMessages() : bool
- isCancellationRequested() : bool
- localActivity() : mixed
- Execute a registered activity in this workflow worker and record its outcome durably.
- messageStream() : MessageStream
- messageStreamCursor() : int
- messageStreamCursorAcknowledgements() : array<int, array{stream_name: string, through_position: int}>
- messageStreamPendingWaits() : array<int, array{stream_name: string, after_position: int}>
- parallel() : array<int, mixed>
- Alias for {@see self::all()}.
- patched() : bool
- Record or replay the standard -1 (legacy) / 1 (patched) decision.
- recordMessageStreamWait() : void
- saga() : Saga
- Create an isolated deterministic saga for activity compensation.
- select() : SelectionResult
- Schedule every durable member and resume with the first committed winner.
- sideEffect() : mixed
- signals() : array<int, array<int, mixed>>
- sleep() : void
- throwIfCancellationRequested() : void
- updates() : array<int, array<int, mixed>>
- upsertMemo() : void
- Merge non-indexed workflow memo metadata. Null removes a key; all other Avro values replace that key while unrelated memo entries are preserved.
- upsertSearchAttributes() : void
- waitCondition() : bool
- Suspend until the deterministic predicate is satisfied or its durable timeout elapses.
- assertActiveFiber() : void
- assertDeferredOperation() : DeferredWorkflowOperation|ParallelWorkflowCommand
- capture() : void
- captureOperation() : DeferredWorkflowOperation|ParallelWorkflowCommand
- conditionKey() : string|null
- finishWorkflowStream() : void
- isCapturing() : bool
- loadMessageStreamMessages() : void
- suspend() : mixed
- version() : int|bool|null
Constants
MAX_PARALLEL_OPERATIONS
public
mixed
MAX_PARALLEL_OPERATIONS
= 1000
MESSAGE_STREAM_CURSOR_SCHEMA
public
mixed
MESSAGE_STREAM_CURSOR_SCHEMA
= 'durable-workflow.v2.message-stream.cursor'
MESSAGE_STREAM_SCHEMA
public
mixed
MESSAGE_STREAM_SCHEMA
= 'durable-workflow.v2.message-stream.message'
MESSAGE_STREAM_SIGNAL
public
mixed
MESSAGE_STREAM_SIGNAL
= '__durable_workflow_message_stream'
MAX_VERSION
private
mixed
MAX_VERSION
= 2147483647
MIN_VERSION
private
mixed
MIN_VERSION
= -2147483648
Properties
$runId read-only
public
string
$runId
$workflowId read-only
public
string
$workflowId
$cancellationRequested read-only
private
bool
$cancellationRequested
= false
$captureFrames
private
array<int, array<int, DeferredWorkflowOperation|ParallelWorkflowCommand>>
$captureFrames
= []
$codec read-only
private
PayloadCodec
$codec
$execution read-only
private
Fiber<mixed, mixed>|null
$execution
$history read-only
private
array<string|int, mixed>
$history
$localActivityExecutor read-only
private
Closure|null
$localActivityExecutor
= null
$messageStreamCursors
private
array<string, int>
$messageStreamCursors
= []
$messageStreamMessages
private
array<string, array<int, MessageStreamMessage>>
$messageStreamMessages
= []
$messageStreamWaits
private
array<string, int>
$messageStreamWaits
= []
$workflowCommandId read-only
private
string|null
$workflowCommandId
= null
$workflowStreamCommandOrdinal
private
int
$workflowStreamCommandOrdinal
= 0
Methods
__construct()
public
__construct(string $workflowId, string $runId, array<int, array<string, mixed>> $history, PayloadCodec $codec[, bool $cancellationRequested = false ][, Fiber<mixed, mixed>|null $execution = null ][, string|null $workflowCommandId = null ][, Closure|null $localActivityExecutor = null ]) : mixed
Parameters
- $workflowId : string
- $runId : string
- $history : array<int, array<string, mixed>>
- $codec : PayloadCodec
- $cancellationRequested : bool = false
- $execution : Fiber<mixed, mixed>|null = null
- $workflowCommandId : string|null = null
- $localActivityExecutor : Closure|null = null
activity()
public
activity(string $activityType[, array<int, mixed> $arguments = [] ][, array<string, mixed> $options = [] ]) : mixed
Parameters
- $activityType : string
- $arguments : array<int, mixed> = []
- $options : array<string, mixed> = []
all()
Schedule every deferred leaf, then return results in declaration order.
public
all(iterable<int, callable(): mixed|DeferredWorkflowOperation> $operations) : array<int, mixed>
Closures are captured without suspending, so ordinary activity(), childWorkflow(), and sleep() calls remain straight-line. Nested all()/parallel() calls preserve their result shape. The first durable failure is thrown at this barrier.
Parameters
- $operations : iterable<int, callable(): mixed|DeferredWorkflowOperation>
Return values
array<int, mixed>appendWorkflowStream()
Append a replay-safe batch to a named run-scoped Workflow Stream.
public
appendWorkflowStream(string $streamName, array<int, WorkflowStreamAppendItem> $items[, int|null $maxPendingItems = null ]) : void
Parameters
- $streamName : string
- $items : array<int, WorkflowStreamAppendItem>
- $maxPendingItems : int|null = null
childWorkflow()
public
childWorkflow(string $workflowType[, array<int, mixed> $arguments = [] ][, array<string, mixed> $options = [] ]) : mixed
Parameters
- $workflowType : string
- $arguments : array<int, mixed> = []
- $options : array<string, mixed> = []
closeWorkflowStream()
public
closeWorkflowStream(string $streamName[, int|null $retentionSeconds = null ]) : void
Parameters
- $streamName : string
- $retentionSeconds : int|null = null
consumeMessageStreamMessages()
public
consumeMessageStreamMessages(string $name, int $maxItems) : array<int, MessageStreamMessage>
Parameters
- $name : string
- $maxItems : int
Return values
array<int, MessageStreamMessage>continueAsNew()
public
continueAsNew([array<int, mixed> $arguments = [] ][, string|null $workflowType = null ][, string|null $taskQueue = null ]) : never
Parameters
- $arguments : array<int, mixed> = []
- $workflowType : string|null = null
- $taskQueue : string|null = null
Return values
neverdeferActivity()
Prepare an activity without scheduling it until an all/parallel barrier is reached.
public
deferActivity(string $activityType[, array<int, mixed> $arguments = [] ][, array<string, mixed> $options = [] ]) : DeferredWorkflowOperation
Parameters
- $activityType : string
- $arguments : array<int, mixed> = []
- $options : array<string, mixed> = []
Return values
DeferredWorkflowOperationdeferChildWorkflow()
Prepare a child workflow without starting it until an all/parallel barrier is reached.
public
deferChildWorkflow(string $workflowType[, array<int, mixed> $arguments = [] ][, array<string, mixed> $options = [] ]) : DeferredWorkflowOperation
Parameters
- $workflowType : string
- $arguments : array<int, mixed> = []
- $options : array<string, mixed> = []
Return values
DeferredWorkflowOperationdeferCondition()
Prepare a deterministic condition or signal-derived wait for a durable group.
public
deferCondition(callable $predicate[, string|null $key = null ][, int|float|null $timeout = null ]) : DeferredWorkflowOperation
Parameters
- $predicate : callable
- $key : string|null = null
- $timeout : int|float|null = null
Return values
DeferredWorkflowOperationdeferTimer()
Prepare a durable timer for an all/parallel barrier.
public
deferTimer(int|float $seconds) : DeferredWorkflowOperation
Parameters
- $seconds : int|float
Return values
DeferredWorkflowOperationdeprecatePatch()
Keep a patch marker alive after the legacy branch has been removed.
public
deprecatePatch(string $changeId) : void
Parameters
- $changeId : string
errorWorkflowStream()
public
errorWorkflowStream(string $streamName, string $errorReason[, int|null $retentionSeconds = null ]) : void
Parameters
- $streamName : string
- $errorReason : string
- $retentionSeconds : int|null = null
getVersion()
Select the newest supported version for a change, or replay its recorded decision.
public
getVersion(string $changeId, int $minSupported, int $maxSupported) : int
Parameters
- $changeId : string
- $minSupported : int
- $maxSupported : int
Return values
inthasPendingMessageStreamMessages()
public
hasPendingMessageStreamMessages(string $name) : bool
Parameters
- $name : string
Return values
boolisCancellationRequested()
public
isCancellationRequested() : bool
Return values
boollocalActivity()
Execute a registered activity in this workflow worker and record its outcome durably.
public
localActivity(string $activityType[, array<int, mixed> $arguments = [] ][, array<string, mixed> $options = [] ]) : mixed
Parameters
- $activityType : string
- $arguments : array<int, mixed> = []
- $options : array<string, mixed> = []
messageStream()
public
messageStream(string $name) : MessageStream
Parameters
- $name : string
Return values
MessageStreammessageStreamCursor()
public
messageStreamCursor(string $name) : int
Parameters
- $name : string
Return values
intmessageStreamCursorAcknowledgements()
public
messageStreamCursorAcknowledgements() : array<int, array{stream_name: string, through_position: int}>
Return values
array<int, array{stream_name: string, through_position: int}>messageStreamPendingWaits()
public
messageStreamPendingWaits() : array<int, array{stream_name: string, after_position: int}>
Return values
array<int, array{stream_name: string, after_position: int}>parallel()
Alias for {@see self::all()}.
public
parallel(iterable<int, callable(): mixed|DeferredWorkflowOperation> $operations) : array<int, mixed>
Parameters
- $operations : iterable<int, callable(): mixed|DeferredWorkflowOperation>
Return values
array<int, mixed>patched()
Record or replay the standard -1 (legacy) / 1 (patched) decision.
public
patched(string $changeId) : bool
Parameters
- $changeId : string
Return values
boolrecordMessageStreamWait()
public
recordMessageStreamWait(string $name, int $afterPosition) : void
Parameters
- $name : string
- $afterPosition : int
saga()
Create an isolated deterministic saga for activity compensation.
public
saga() : Saga
Return values
Sagaselect()
Schedule every durable member and resume with the first committed winner.
public
select(iterable<int|string, callable(): mixed|DeferredWorkflowOperation|ParallelWorkflowCommand> $operations) : SelectionResult
Non-winning operations continue and remain available through durable handles.
Parameters
- $operations : iterable<int|string, callable(): mixed|DeferredWorkflowOperation|ParallelWorkflowCommand>
Return values
SelectionResultsideEffect()
public
sideEffect(callable(): mixed $operation) : mixed
Parameters
- $operation : callable(): mixed
signals()
public
signals(string $signalName) : array<int, array<int, mixed>>
Parameters
- $signalName : string
Return values
array<int, array<int, mixed>>sleep()
public
sleep(int|float $seconds) : void
Parameters
- $seconds : int|float
throwIfCancellationRequested()
public
throwIfCancellationRequested() : void
updates()
public
updates(string $updateName) : array<int, array<int, mixed>>
Parameters
- $updateName : string
Return values
array<int, array<int, mixed>>upsertMemo()
Merge non-indexed workflow memo metadata. Null removes a key; all other Avro values replace that key while unrelated memo entries are preserved.
public
upsertMemo(array<string, mixed> $entries) : void
Parameters
- $entries : array<string, mixed>
upsertSearchAttributes()
public
upsertSearchAttributes(array<string, mixed> $attributes) : void
Parameters
- $attributes : array<string, mixed>
waitCondition()
Suspend until the deterministic predicate is satisfied or its durable timeout elapses.
public
waitCondition(callable(): bool $predicate[, string|null $key = null ][, int|float|null $timeout = null ]) : bool
The result is true only when the condition was satisfied and false only when it timed out. Give repeated or otherwise ambiguous waits a stable key so replay can identify them.
Parameters
- $predicate : callable(): bool
- $key : string|null = null
- $timeout : int|float|null = null
Return values
boolassertActiveFiber()
private
assertActiveFiber() : void
assertDeferredOperation()
private
assertDeferredOperation(mixed $operation) : DeferredWorkflowOperation|ParallelWorkflowCommand
Parameters
- $operation : mixed
Return values
DeferredWorkflowOperation|ParallelWorkflowCommandcapture()
private
capture(DeferredWorkflowOperation|ParallelWorkflowCommand $operation) : void
Parameters
- $operation : DeferredWorkflowOperation|ParallelWorkflowCommand
captureOperation()
private
captureOperation(callable(): mixed $callback) : DeferredWorkflowOperation|ParallelWorkflowCommand
Parameters
- $callback : callable(): mixed
Return values
DeferredWorkflowOperation|ParallelWorkflowCommandconditionKey()
private
static conditionKey(string|null $key) : string|null
Parameters
- $key : string|null
Return values
string|nullfinishWorkflowStream()
private
finishWorkflowStream(string $streamName, string|null $errorReason, int|null $retentionSeconds) : void
Parameters
- $streamName : string
- $errorReason : string|null
- $retentionSeconds : int|null
isCapturing()
private
isCapturing() : bool
Return values
boolloadMessageStreamMessages()
private
loadMessageStreamMessages() : void
suspend()
private
suspend(WorkflowCommand|ParallelWorkflowCommand $command) : mixed
Parameters
- $command : WorkflowCommand|ParallelWorkflowCommand
version()
private
version(string $changeId, int $minSupported, int $maxSupported, string $resultKind) : int|bool|null
Parameters
- $changeId : string
- $minSupported : int
- $maxSupported : int
- $resultKind : string