diff --git a/src/V2/Support/ParallelChildGroup.php b/src/V2/Support/ParallelChildGroup.php index ef5aa2ef..bb2322e5 100644 --- a/src/V2/Support/ParallelChildGroup.php +++ b/src/V2/Support/ParallelChildGroup.php @@ -799,7 +799,7 @@ private static function groupCompletedSuccessfully( $parentRun->setRelation('historyEvents', $parentRun->historyEvents() ->lockForUpdate() ->get()); } - $activitiesBySequence = collect(RunActivityView::activitiesForRun($parentRun)) + $activitiesBySequence = collect(RunActivityView::activitiesForRun($parentRun, decodePayloads: false)) ->filter(static fn (array $activity): bool => is_int($activity['sequence'] ?? null)) ->keyBy(static fn (array $activity): string => (string) $activity['sequence']); $timersBySequence = collect(RunTimerView::timersForRun($parentRun)) diff --git a/src/V2/Support/RunActivityView.php b/src/V2/Support/RunActivityView.php index e45412cd..c286f5a2 100644 --- a/src/V2/Support/RunActivityView.php +++ b/src/V2/Support/RunActivityView.php @@ -27,12 +27,16 @@ final class RunActivityView /** * Typed history is the default authority. Callers rendering bounded * projection-backed views may disable it to avoid loading durable history. + * Metadata-only callers can omit payload fields without fetching external objects. * * @return list> */ - public static function activitiesForRun(WorkflowRun $run, bool $useDurableHistory = true): array - { - return self::activityStates($run, $useDurableHistory); + public static function activitiesForRun( + WorkflowRun $run, + bool $useDurableHistory = true, + bool $decodePayloads = true, + ): array { + return self::activityStates($run, $useDurableHistory, $decodePayloads); } /** @@ -108,8 +112,11 @@ public static function isDiagnosticOnly(array $activity): bool /** * @return list> */ - private static function activityStates(WorkflowRun $run, bool $useDurableHistory = true): array - { + private static function activityStates( + WorkflowRun $run, + bool $useDurableHistory = true, + bool $decodePayloads = true, + ): array { $relations = ['activityExecutions.attempts']; if ($useDurableHistory) { @@ -182,6 +189,7 @@ private static function activityStates(WorkflowRun $run, bool $useDurableHistory $run, $execution, $attemptsByActivityId[$activityId] ?? [], + $decodePayloads, ); } @@ -235,6 +243,7 @@ private static function presentActivity( WorkflowRun $run, ?ActivityExecution $execution = null, array $attemptStates = [], + bool $decodePayloads = true, ): array { $attempts = self::presentAttempts($state, $execution, $attemptStates); $latestAttempt = $attempts === [] @@ -291,13 +300,15 @@ private static function presentActivity( 'created_at' => $state['created_at'] ?? null, 'started_at' => $state['started_at'] ?? ($latestAttempt['started_at'] ?? null), 'closed_at' => self::activityClosedAt($status, $state, $latestAttempt), - 'arguments' => self::publicTypedValue($state['arguments'] ?? null, $payloadCodec, $namespace, []), - 'result' => self::publicTypedValue( - $unsupportedReason === null ? ($state['result'] ?? null) : null, - $payloadCodec, - $namespace, - null, - ), + ...($decodePayloads ? [ + 'arguments' => self::publicTypedValue($state['arguments'] ?? null, $payloadCodec, $namespace, []), + 'result' => self::publicTypedValue( + $unsupportedReason === null ? ($state['result'] ?? null) : null, + $payloadCodec, + $namespace, + null, + ), + ] : []), 'attempts' => $attempts, ]; } diff --git a/src/V2/Support/RunSummaryProjector.php b/src/V2/Support/RunSummaryProjector.php index 175f68fb..26dd7161 100644 --- a/src/V2/Support/RunSummaryProjector.php +++ b/src/V2/Support/RunSummaryProjector.php @@ -58,7 +58,7 @@ public static function project(WorkflowRun $run): WorkflowRunSummary : CurrentRunResolver::forInstance($run->instance); $isTerminal = $run->status->isTerminal(); - $activities = RunActivityView::activitiesForRun($run); + $activities = RunActivityView::activitiesForRun($run, decodePayloads: false); $timers = RunTimerView::timersForRun($run); $openActivity = $isTerminal diff --git a/src/V2/Support/RunTaskView.php b/src/V2/Support/RunTaskView.php index 1ae5bc2b..b5736dd6 100644 --- a/src/V2/Support/RunTaskView.php +++ b/src/V2/Support/RunTaskView.php @@ -23,7 +23,7 @@ public static function forRun(WorkflowRun $run): array $run->loadMissing(['tasks', 'activityExecutions', 'timers', 'historyEvents']); /** @var Collection> $activities */ - $activities = collect(RunActivityView::activitiesForRun($run)) + $activities = collect(RunActivityView::activitiesForRun($run, decodePayloads: false)) ->filter(static fn (array $activity): bool => is_string($activity['id'] ?? null)) ->keyBy(static fn (array $activity): string => $activity['id']); /** @var Collection> $timers */ diff --git a/src/V2/Support/RunWaitView.php b/src/V2/Support/RunWaitView.php index 24b605a2..415793ac 100644 --- a/src/V2/Support/RunWaitView.php +++ b/src/V2/Support/RunWaitView.php @@ -51,7 +51,7 @@ public static function forRun(WorkflowRun $run): array $waits = []; - foreach (RunActivityView::activitiesForRun($run) as $activity) { + foreach (RunActivityView::activitiesForRun($run, decodePayloads: false) as $activity) { if (! is_string($activity['id'] ?? null)) { continue; } diff --git a/tests/Feature/V2/V2ActivityMetadataProjectionTest.php b/tests/Feature/V2/V2ActivityMetadataProjectionTest.php new file mode 100644 index 00000000..a416a18d --- /dev/null +++ b/tests/Feature/V2/V2ActivityMetadataProjectionTest.php @@ -0,0 +1,128 @@ +reads++; + + return $this->bytes; + } + + public function delete(string $uri): void + { + } + }; + $this->app->instance(ExternalPayloadStoragePolicy::class, new class( + $driver + ) implements ExternalPayloadStoragePolicy { + public function __construct( + private readonly ExternalPayloadStorageDriver $driver + ) { + } + + public function driverFor(?string $namespace): ?ExternalPayloadStorageDriver + { + return $this->driver; + } + + public function thresholdBytesFor(?string $namespace): ?int + { + return 1; + } + }); + $stored = ExternalPayloads::externalize($bytes, 'avro', $driver, 1); + $instance = WorkflowInstance::query()->create([ + 'id' => 'metadata-projection', + 'workflow_class' => 'MetadataWorkflow', + 'workflow_type' => 'metadata.workflow', + 'run_count' => 1, + ]); + $run = WorkflowRun::query()->create([ + 'workflow_instance_id' => $instance->id, + 'run_number' => 1, + 'workflow_class' => 'MetadataWorkflow', + 'workflow_type' => 'metadata.workflow', + 'status' => 'completed', + 'namespace' => 'default', + 'payload_codec' => 'avro', + 'started_at' => now(), + 'closed_at' => now(), + 'last_progress_at' => now(), + ]); + $instance->forceFill([ + 'current_run_id' => $run->id, + ])->save(); + $activity = ActivityExecution::query()->create([ + 'workflow_run_id' => $run->id, + 'sequence' => 1, + 'activity_class' => 'MetadataActivity', + 'activity_type' => 'metadata.activity', + 'status' => 'completed', + 'attempt_count' => 1, + 'payload_codec' => 'avro', + 'arguments' => $stored, + 'result' => $stored, + 'closed_at' => now(), + ]); + WorkflowHistoryEvent::record($run, HistoryEventType::ActivityCompleted, [ + 'activity_execution_id' => $activity->id, + 'activity' => ActivitySnapshot::fromExecution($activity), + ]); + $runId = $run->id; + unset($run, $activity); + $run = WorkflowRun::query()->findOrFail($runId); + + $summary = RunSummaryProjector::project($run); + $this->assertSame(0, $driver->reads, 'Summary projection must not download activity data.'); + $this->assertSame('completed', $summary->status); + + $metadata = RunActivityView::activitiesForRun($run->fresh(), decodePayloads: false)[0]; + $this->assertSame(0, $driver->reads); + $this->assertSame(RunActivityView::HISTORY_AUTHORITY_TYPED, $metadata['history_authority']); + $this->assertSame('completed', $metadata['status']); + $this->assertArrayNotHasKey('arguments', $metadata); + $this->assertArrayNotHasKey('result', $metadata); + + $decoded = RunActivityView::activitiesForRun($run->fresh())[0]; + $this->assertSame(['preserved value'], $decoded['arguments']); + $this->assertSame(['preserved value'], $decoded['result']); + $this->assertGreaterThan(0, $driver->reads); + unset($decoded['arguments'], $decoded['result']); + $this->assertSame($metadata, $decoded); + } +}