diff --git a/app/Http/Controllers/Api/WorkerController.php b/app/Http/Controllers/Api/WorkerController.php index 66dc4624..822e6915 100644 --- a/app/Http/Controllers/Api/WorkerController.php +++ b/app/Http/Controllers/Api/WorkerController.php @@ -1821,16 +1821,16 @@ function () use ( $this->authorizeServiceOperationReplays($request, (string) $namespace, $taskId, $commands); $commands = $this->canonicalizeWorkflowStreamPayloadCodecs($commands); - $commands = app(WorkflowStreamCommandProcessor::class)->process( - $taskId, - (string) $namespace, - $commands, - ); + $streamCommands = $commands; + $commands = app(WorkflowStreamCommandProcessor::class)->withoutDirectives($commands); $commands = WorkflowCommandNormalizer::normalize( $commands, WorkerProtocol::requestVersion($request), ); $outcome = $bridge->complete($taskId, $commands); + if (($outcome['completed'] ?? false) === true) { + app(WorkflowStreamCommandProcessor::class)->process($taskId, (string) $namespace, $streamCommands); + } $this->applyStickyCacheClaim( $taskId, $validated['lease_owner'], diff --git a/app/Support/WorkflowStreamCommandProcessor.php b/app/Support/WorkflowStreamCommandProcessor.php index 9c4d88ed..bf10a0c9 100644 --- a/app/Support/WorkflowStreamCommandProcessor.php +++ b/app/Support/WorkflowStreamCommandProcessor.php @@ -7,11 +7,12 @@ use Workflow\V2\Models\WorkflowTask; /** - * Applies replay-safe Workflow Stream directives before task completion. + * Applies replay-safe Workflow Stream directives after successful admission. * * The directive rides a record_side_effect command so the append/close and * SideEffectRecorded event commit in one outer database transaction. The - * directive is stripped before the package command normalizer sees it. + * directive is stripped before the package command normalizer sees it. Effects + * are applied only after admission succeeds, inside the same outer transaction. */ final class WorkflowStreamCommandProcessor { @@ -19,6 +20,19 @@ public function __construct( private readonly WorkflowStreamService $streams, ) {} + /** + * @param list> $commands + * @return list> + */ + public function withoutDirectives(array $commands): array + { + foreach ($commands as &$command) { + unset($command['workflow_stream']); + } + + return $commands; + } + /** * @param list> $commands * @return list> diff --git a/composer.json b/composer.json index a91607af..dee537a5 100644 --- a/composer.json +++ b/composer.json @@ -48,7 +48,7 @@ }, "extra": { "durable-workflow": { - "product-train": "2.4.38" + "product-train": "2.4.39" }, "laravel": { "dont-discover": [] diff --git a/composer.lock b/composer.lock index 156ce704..3a613e1c 100644 --- a/composer.lock +++ b/composer.lock @@ -4,7 +4,7 @@ "Read more about it at https://getcomposer.org/doc/01-basic-usage.md#installing-dependencies", "This file is @generated automatically" ], - "content-hash": "df708d2799646bb6864288d08677c22e", + "content-hash": "cbd2cce87f07478179460c70eb2e728c", "packages": [ { "name": "apache/avro", diff --git a/docker-compose.dedicated-matching.yml b/docker-compose.dedicated-matching.yml index f58517c7..65151065 100644 --- a/docker-compose.dedicated-matching.yml +++ b/docker-compose.dedicated-matching.yml @@ -32,13 +32,13 @@ name: durable-workflow-server # daemon reports `shape: dedicated`. # Generated by scripts/ci/sync-source-release.mjs. Do not edit the fallback. -x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.4.38}} +x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.4.39}} x-server-environment: &server-environment APP_NAME: "Durable Workflow Server" APP_ENV: ${APP_ENV:-local} DW_SERVER_KEY: ${DW_SERVER_KEY:-} - APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.4.38}} + APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.4.39}} APP_DEBUG: ${APP_DEBUG:-false} DB_CONNECTION: mysql DB_HOST: mysql diff --git a/docker-compose.memo-rolling.yml b/docker-compose.memo-rolling.yml index 9777f5a0..424c5730 100644 --- a/docker-compose.memo-rolling.yml +++ b/docker-compose.memo-rolling.yml @@ -49,14 +49,14 @@ services: command: ["server-bootstrap"] environment: <<: *runtime-environment - APP_VERSION: ${APP_VERSION:-2.4.38} + APP_VERSION: ${APP_VERSION:-2.4.39} successor: image: ${DW_MEMO_SUCCESSOR_IMAGE:-durable-workflow/server-memo-rolling:local} ports: !override [] environment: <<: *runtime-environment - APP_VERSION: ${APP_VERSION:-2.4.38} + APP_VERSION: ${APP_VERSION:-2.4.39} DW_SERVER_ID: memo-successor DW_SERVER_TOPOLOGY_SHAPE: standalone_server DW_SERVER_PROCESS_CLASS: server_http_node diff --git a/docker-compose.published.yml b/docker-compose.published.yml index 7da3de88..25711eaf 100644 --- a/docker-compose.published.yml +++ b/docker-compose.published.yml @@ -1,13 +1,13 @@ name: durable-workflow-server # Generated by scripts/ci/sync-source-release.mjs. Do not edit the fallback. -x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.4.38}} +x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.4.39}} x-server-environment: &server-environment APP_NAME: "Durable Workflow Server" APP_ENV: ${APP_ENV:-local} DW_SERVER_KEY: ${DW_SERVER_KEY:-} - APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.4.38}} + APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.4.39}} APP_DEBUG: ${APP_DEBUG:-false} LOG_CHANNEL: ${LOG_CHANNEL:-stderr} LOG_LEVEL: ${LOG_LEVEL:-info} diff --git a/docker-compose.small-cluster.yml b/docker-compose.small-cluster.yml index dcde5e55..a12b02ce 100644 --- a/docker-compose.small-cluster.yml +++ b/docker-compose.small-cluster.yml @@ -12,7 +12,7 @@ x-server-build: &server-build x-server-environment: &server-environment APP_NAME: "Durable Workflow Server" APP_ENV: testing - APP_VERSION: ${APP_VERSION:-2.4.38} + APP_VERSION: ${APP_VERSION:-2.4.39} APP_DEBUG: "false" DW_SERVER_KEY: ${DW_SERVER_KEY:-base64:5Zt4nUhlCm3DD0nLXZJQdHiwPfb56yGo9gNV/g3jYbY=} DB_CONNECTION: ${DW_SMALL_CLUSTER_DB:-mysql} diff --git a/docker-compose.yml b/docker-compose.yml index d0fd5589..1b8f1a44 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -15,7 +15,7 @@ services: DW_SERVER_KEY: "${DW_SERVER_KEY:-}" DW_SERVER_TOPOLOGY_SHAPE: standalone_server DW_SERVER_PROCESS_CLASS: server_http_node - APP_VERSION: "${APP_VERSION:-2.4.38}" + APP_VERSION: "${APP_VERSION:-2.4.39}" APP_DEBUG: "false" DB_CONNECTION: mysql DB_HOST: mysql @@ -62,7 +62,7 @@ services: APP_NAME: "Durable Workflow Server" APP_ENV: local DW_SERVER_KEY: "${DW_SERVER_KEY:-}" - APP_VERSION: "${APP_VERSION:-2.4.38}" + APP_VERSION: "${APP_VERSION:-2.4.39}" APP_DEBUG: "false" DB_CONNECTION: mysql DB_HOST: mysql @@ -124,7 +124,7 @@ services: DW_SERVER_KEY: "${DW_SERVER_KEY:-}" DW_SERVER_TOPOLOGY_SHAPE: standalone_server DW_SERVER_PROCESS_CLASS: worker_node - APP_VERSION: "${APP_VERSION:-2.4.38}" + APP_VERSION: "${APP_VERSION:-2.4.39}" DB_CONNECTION: mysql DB_HOST: mysql DB_PORT: 3306 @@ -178,7 +178,7 @@ services: DW_SERVER_KEY: "${DW_SERVER_KEY:-}" DW_SERVER_TOPOLOGY_SHAPE: standalone_server DW_SERVER_PROCESS_CLASS: scheduler_node - APP_VERSION: "${APP_VERSION:-2.4.38}" + APP_VERSION: "${APP_VERSION:-2.4.39}" DB_CONNECTION: mysql DB_HOST: mysql DB_PORT: 3306 diff --git a/k8s/README.md b/k8s/README.md index 2d324146..54e3c7d9 100644 --- a/k8s/README.md +++ b/k8s/README.md @@ -13,7 +13,7 @@ The checked-in manifests are synchronized with the repository's stable source release and pin its Docker Hub tag: ```text -durableworkflow/server:2.4.38 +durableworkflow/server:2.4.39 ``` Before production use, patch every workload image to the exact published tag or @@ -21,15 +21,15 @@ digest you intend to run: ```bash kubectl set image -n durable-workflow deploy/durable-workflow-server \ - server=durableworkflow/server:2.4.38 + server=durableworkflow/server:2.4.39 kubectl set image -n durable-workflow deploy/durable-workflow-worker \ - worker=durableworkflow/server:2.4.38 + worker=durableworkflow/server:2.4.39 kubectl set image -n durable-workflow cronjob/durable-workflow-scheduler \ - scheduler=durableworkflow/server:2.4.38 + scheduler=durableworkflow/server:2.4.39 ``` GitHub Container Registry publishes the same release line at -`ghcr.io/durable-workflow/server:2.4.38`. Digest pinning is preferred for strict +`ghcr.io/durable-workflow/server:2.4.39`. Digest pinning is preferred for strict change control. The manifests expect you to provide: diff --git a/k8s/helm/durable-workflow/Chart.yaml b/k8s/helm/durable-workflow/Chart.yaml index 89cb22d0..e3a5f436 100644 --- a/k8s/helm/durable-workflow/Chart.yaml +++ b/k8s/helm/durable-workflow/Chart.yaml @@ -5,11 +5,11 @@ type: application # The chart's own semver version. Bumped on every chart release; treated as # independent of the server image version (appVersion). Breaking-change rules # for this version live in docs/helm-upgrading.md alongside the chart. -version: 0.1.134 +version: 0.1.135 # The immutable Durable Workflow Server identity this chart release packages. # The onboarding default in values.yaml and appVersion are generated from the # checked-in source release record. -appVersion: "2.4.38" +appVersion: "2.4.39" kubeVersion: ">=1.27.0-0" home: https://durable-workflow.github.io/docs/2.0/deployment sources: @@ -30,7 +30,7 @@ annotations: # exact commit that most recently changed the packaged chart. org.opencontainers.image.source: https://github.com/durable-workflow/server dev.durable-workflow.source-revision: "unreleased" - dev.durable-workflow.image-reference: "docker.io/durableworkflow/server:2.4.38" + dev.durable-workflow.image-reference: "docker.io/durableworkflow/server:2.4.39" artifacthub.io/license: MIT artifacthub.io/category: integration-delivery # Free-form changelog for the current chart release shown by Artifact Hub. diff --git a/k8s/helm/durable-workflow/README.md b/k8s/helm/durable-workflow/README.md index 02711428..7b6ae984 100644 --- a/k8s/helm/durable-workflow/README.md +++ b/k8s/helm/durable-workflow/README.md @@ -63,7 +63,7 @@ helm install durable-workflow ./k8s/helm/durable-workflow \ ```yaml image: - tag: "2.4.38" + tag: "2.4.39" # Pin a digest in production: # digest: "sha256:abc123..." # memoPayloadStorage: "raw-json-v1" # Required for a digest or custom image. diff --git a/k8s/helm/durable-workflow/ci/existing-secrets-values.yaml b/k8s/helm/durable-workflow/ci/existing-secrets-values.yaml index 528a9b37..25a33826 100644 --- a/k8s/helm/durable-workflow/ci/existing-secrets-values.yaml +++ b/k8s/helm/durable-workflow/ci/existing-secrets-values.yaml @@ -1,7 +1,7 @@ # CI fixture: GitOps / externally-managed-secret path. The chart consumes # existing Secrets and renders no Secret resources of its own. image: - tag: "2.4.38" + tag: "2.4.39" externalDatabase: connection: pgsql diff --git a/k8s/helm/durable-workflow/ci/ingress-and-hpa-values.yaml b/k8s/helm/durable-workflow/ci/ingress-and-hpa-values.yaml index 360fe5ae..53dd12a1 100644 --- a/k8s/helm/durable-workflow/ci/ingress-and-hpa-values.yaml +++ b/k8s/helm/durable-workflow/ci/ingress-and-hpa-values.yaml @@ -1,6 +1,6 @@ # CI fixture: ingress + autoscaling enabled. Exercises optional templates. image: - tag: "2.4.38" + tag: "2.4.39" externalDatabase: connection: mysql diff --git a/k8s/helm/durable-workflow/ci/inline-secrets-values.yaml b/k8s/helm/durable-workflow/ci/inline-secrets-values.yaml index 7994ab13..292fe26b 100644 --- a/k8s/helm/durable-workflow/ci/inline-secrets-values.yaml +++ b/k8s/helm/durable-workflow/ci/inline-secrets-values.yaml @@ -2,7 +2,7 @@ # chart's render path is exercised end-to-end. Real deployments should use # existingSecret instead. image: - tag: "2.4.38" + tag: "2.4.39" externalDatabase: connection: mysql diff --git a/k8s/helm/durable-workflow/templates/_helpers.tpl b/k8s/helm/durable-workflow/templates/_helpers.tpl index 8024c962..a01a2ed7 100644 --- a/k8s/helm/durable-workflow/templates/_helpers.tpl +++ b/k8s/helm/durable-workflow/templates/_helpers.tpl @@ -88,7 +88,7 @@ resolved by an explicit capability declaration or an existing workload marker. {{- define "durable-workflow.memoPayloadStorageForImage" -}} {{- $image := toString . -}} {{- $normalized := regexReplaceAll "^index\\.docker\\.io/" $image "docker.io/" -}} -{{- if eq $normalized "docker.io/durableworkflow/server:2.4.38" -}} +{{- if eq $normalized "docker.io/durableworkflow/server:2.4.39" -}} dual-v1 {{- else if regexMatch "^docker\\.io/durableworkflow/server:2\\.0\\.0-rc\\.[0-9]+$" $normalized -}} {{- $releaseCandidate := atoi (regexFind "[0-9]+$" $normalized) -}} diff --git a/k8s/helm/durable-workflow/values.yaml b/k8s/helm/durable-workflow/values.yaml index 6c3fb3a0..3b71081a 100644 --- a/k8s/helm/durable-workflow/values.yaml +++ b/k8s/helm/durable-workflow/values.yaml @@ -21,7 +21,7 @@ image: registry: docker.io repository: durableworkflow/server # Generated by scripts/ci/sync-source-release.mjs. Do not edit this default. - tag: "2.4.38" + tag: "2.4.39" # Optional digest pin. When set, takes precedence over tag for change control. # Example: "sha256:abc123..." digest: "" diff --git a/k8s/helm/examples/values-dev.yaml b/k8s/helm/examples/values-dev.yaml index 6d26c78e..12ae8f9c 100644 --- a/k8s/helm/examples/values-dev.yaml +++ b/k8s/helm/examples/values-dev.yaml @@ -3,7 +3,7 @@ # shape in production. image: - tag: "2.4.38" + tag: "2.4.39" externalDatabase: connection: mysql diff --git a/k8s/helm/examples/values-external-secrets-operator.yaml b/k8s/helm/examples/values-external-secrets-operator.yaml index b97320d4..622177fe 100644 --- a/k8s/helm/examples/values-external-secrets-operator.yaml +++ b/k8s/helm/examples/values-external-secrets-operator.yaml @@ -5,7 +5,7 @@ # concern. image: - tag: "2.4.38" + tag: "2.4.39" externalDatabase: connection: pgsql diff --git a/k8s/helm/examples/values-production-existing-secrets.yaml b/k8s/helm/examples/values-production-existing-secrets.yaml index 4d6cdc84..d43bd08f 100644 --- a/k8s/helm/examples/values-production-existing-secrets.yaml +++ b/k8s/helm/examples/values-production-existing-secrets.yaml @@ -10,7 +10,7 @@ image: repository: durable-workflow/server # Pin a digest in production for change-control auditability. digest: "" # e.g. "sha256:abc123..." - tag: "2.4.38" + tag: "2.4.39" externalDatabase: connection: pgsql diff --git a/k8s/migration-job.yaml b/k8s/migration-job.yaml index 9ecd0e0f..53df94a3 100644 --- a/k8s/migration-job.yaml +++ b/k8s/migration-job.yaml @@ -13,7 +13,7 @@ spec: restartPolicy: OnFailure containers: - name: migrate - image: durableworkflow/server:2.4.38 + image: durableworkflow/server:2.4.39 command: ["server-entrypoint"] args: ["server-bootstrap"] envFrom: diff --git a/k8s/scheduler-cronjob.yaml b/k8s/scheduler-cronjob.yaml index 52c5cbe4..279e92ef 100644 --- a/k8s/scheduler-cronjob.yaml +++ b/k8s/scheduler-cronjob.yaml @@ -24,7 +24,7 @@ spec: restartPolicy: Never containers: - name: scheduler - image: durableworkflow/server:2.4.38 + image: durableworkflow/server:2.4.39 command: ["server-entrypoint"] args: - sh diff --git a/k8s/secret.yaml b/k8s/secret.yaml index cf52c517..2184bb2a 100644 --- a/k8s/secret.yaml +++ b/k8s/secret.yaml @@ -12,7 +12,7 @@ metadata: app.kubernetes.io/name: durable-workflow data: APP_NAME: "Durable Workflow Server" - APP_VERSION: "2.4.38" + APP_VERSION: "2.4.39" APP_ENV: production APP_DEBUG: "false" DB_CONNECTION: mysql diff --git a/k8s/server-deployment.yaml b/k8s/server-deployment.yaml index 1b732e76..a6b0f6a4 100644 --- a/k8s/server-deployment.yaml +++ b/k8s/server-deployment.yaml @@ -23,7 +23,7 @@ spec: spec: containers: - name: server - image: durableworkflow/server:2.4.38 + image: durableworkflow/server:2.4.39 ports: - containerPort: 8080 name: http diff --git a/k8s/worker-deployment.yaml b/k8s/worker-deployment.yaml index 2447474d..5e1eee08 100644 --- a/k8s/worker-deployment.yaml +++ b/k8s/worker-deployment.yaml @@ -19,7 +19,7 @@ spec: spec: containers: - name: worker - image: durableworkflow/server:2.4.38 + image: durableworkflow/server:2.4.39 command: ["server-entrypoint"] args: ["php", "artisan", "queue:work", "--sleep=1", "--tries=3", "--max-time=3600"] envFrom: diff --git a/resources/release/source-release.json b/resources/release/source-release.json index b73bc014..c118fdca 100644 --- a/resources/release/source-release.json +++ b/resources/release/source-release.json @@ -1,9 +1,9 @@ { "schema": "durable-workflow.server.source-release/v1", "server": { - "version": "2.4.38" + "version": "2.4.39" }, "helm_chart": { - "version": "0.1.134" + "version": "0.1.135" } } diff --git a/scripts/k8s-kind-smoke.sh b/scripts/k8s-kind-smoke.sh index caff3005..00b56221 100755 --- a/scripts/k8s-kind-smoke.sh +++ b/scripts/k8s-kind-smoke.sh @@ -7,7 +7,7 @@ cluster="${K8S_SMOKE_CLUSTER:-durable-workflow-server-smoke}" image="${K8S_SMOKE_IMAGE:-durableworkflow/server:k8s-smoke}" # Generated by scripts/ci/sync-source-release.mjs so the smoke replaces the # same default shipped by the public manifests. -manifest_image="durableworkflow/server:2.4.38" +manifest_image="durableworkflow/server:2.4.39" kind_node_image="${K8S_SMOKE_KIND_NODE_IMAGE:-kindest/node:v1.29.4}" artifact_dir="${K8S_SMOKE_ARTIFACT_DIR:-/tmp/durable-workflow-k8s-kind-smoke-artifacts}" rendered_dir="${artifact_dir}/rendered-manifests" diff --git a/tests/Feature/WorkflowStreamsTest.php b/tests/Feature/WorkflowStreamsTest.php index 531d77aa..55d45c54 100644 --- a/tests/Feature/WorkflowStreamsTest.php +++ b/tests/Feature/WorkflowStreamsTest.php @@ -553,6 +553,80 @@ public function test_worker_completion_commits_stream_append_with_side_effect_hi ->count()); } + public function test_timed_out_completion_records_timeout_without_committing_stream_output(): void + { + $task = $this->claimStreamCompletionTask(1); + $this->travel(2)->seconds(); + $this->withHeaders($this->workerHeaders())->postJson("/api/worker/workflow-tasks/{$task['task_id']}/complete", [ + 'lease_owner' => $task['lease_owner'], 'workflow_task_attempt' => $task['workflow_task_attempt'], + 'commands' => [$this->streamCompletionCommand($task)], + ])->assertConflict()->assertJsonPath('recorded', false)->assertJsonPath('reason', 'run_timed_out'); + $this->assertDatabaseCount('workflow_durable_stream_items', 0); + $this->assertDatabaseCount('workflow_durable_streams', 0); + $this->assertDatabaseHas('workflow_history_events', ['workflow_run_id' => $task['run_id'], 'event_type' => 'WorkflowTimedOut']); + $this->assertDatabaseMissing('workflow_history_events', ['workflow_run_id' => $task['run_id'], 'event_type' => 'SideEffectRecorded']); + $this->assertSame('failed', WorkflowRun::query()->findOrFail($task['run_id'])->status->value); + } + + public function test_terminal_completion_commits_its_last_stream_item_and_close(): void + { + $task = $this->claimStreamCompletionTask(); + $close = $this->streamCompletionCommand($task); + $close['workflow_stream'] = [...$close['workflow_stream'], 'operation' => 'close', 'command_ordinal' => 1]; + unset($close['workflow_stream']['items']); + $this->withHeaders($this->workerHeaders())->postJson("/api/worker/workflow-tasks/{$task['task_id']}/complete", [ + 'lease_owner' => $task['lease_owner'], 'workflow_task_attempt' => $task['workflow_task_attempt'], + 'commands' => [$this->streamCompletionCommand($task), $close, ['type' => 'complete_workflow', 'result' => Serializer::serializeWithCodec('avro', 'done')]], + ])->assertOk()->assertJsonPath('recorded', true); + $this->assertDatabaseCount('workflow_durable_stream_items', 1); + $this->assertSame('closed', WorkflowDurableStream::query()->sole()->status); + $this->assertSame('completed', WorkflowRun::query()->findOrFail($task['run_id'])->status->value); + } + + public function test_stream_failure_rolls_back_successful_native_admission(): void + { + $task = $this->claimStreamCompletionTask(); + $before = WorkflowRun::query()->findOrFail($task['run_id'])->status; + $command = $this->streamCompletionCommand($task); + $command['workflow_stream']['max_pending_items'] = 1; + $second = $command['workflow_stream']['items'][0]; + $second['idempotency_key'] = substr($second['idempotency_key'], 0, -1).'1'; + $command['workflow_stream']['items'][] = $second; + $this->withHeaders($this->workerHeaders())->postJson("/api/worker/workflow-tasks/{$task['task_id']}/complete", [ + 'lease_owner' => $task['lease_owner'], 'workflow_task_attempt' => $task['workflow_task_attempt'], + 'commands' => [$command, ['type' => 'complete_workflow', 'result' => Serializer::serializeWithCodec('avro', 'done')]], + ])->assertStatus(429); + $this->assertDatabaseCount('workflow_durable_stream_items', 0); + $this->assertDatabaseMissing('workflow_history_events', ['workflow_run_id' => $task['run_id'], 'event_type' => 'SideEffectRecorded']); + $this->assertSame($before, WorkflowRun::query()->findOrFail($task['run_id'])->status); + $this->assertSame(TaskStatus::Leased, WorkflowTask::query()->findOrFail($task['task_id'])->status); + } + + private function claimStreamCompletionTask(?int $timeout = null): array + { + config()->set('workflows.v2.types.workflows', ['tests.external-greeting-workflow' => ExternalGreetingWorkflow::class]); + $this->withHeaders($this->apiHeaders())->postJson('/api/workflows', [ + 'workflow_type' => 'tests.external-greeting-workflow', 'task_queue' => 'stream-command-queue', + 'input' => ['codec' => 'avro', 'blob' => Serializer::serializeWithCodec('avro', ['Ada'])], + ...($timeout === null ? [] : ['run_timeout_seconds' => $timeout]), + ])->assertCreated(); + $this->registerWorker('stream-command-worker', 'stream-command-queue', supportedWorkflowTypes: ['tests.external-greeting-workflow']); + + return $this->withHeaders($this->workerHeaders())->postJson('/api/worker/workflow-tasks/poll', [ + 'worker_id' => 'stream-command-worker', 'task_queue' => 'stream-command-queue', + ])->assertOk()->json('task'); + } + + private function streamCompletionCommand(array $task): array + { + $identity = $task['workflow_command_id'] ?: $task['task_id']; + + return ['type' => 'record_side_effect', 'result' => Serializer::serializeWithCodec('avro', null), + 'workflow_stream' => ['operation' => 'append', 'stream_name' => 'tokens', 'command_identity' => $identity, + 'command_ordinal' => 0, 'items' => [['payload' => Serializer::serializeWithCodec('avro', ['token']), + 'payload_codec' => 'avro', 'idempotency_key' => "dw-stream:{$identity}:0:0"]]]]; + } + /** * @param list> $items */