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
25 changes: 3 additions & 22 deletions app/Http/Controllers/Api/HistoryController.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@

use App\Support\ControlPlaneProtocol;
use App\Support\ExternalPayloadEnvelopeService;
use App\Support\HistoryPageToken;
use App\Support\LegacyV1Projection;
use App\Support\LongPoller;
use App\Support\LongPollSignalStore;
Expand Down Expand Up @@ -50,7 +51,7 @@ public function show(Request $request, string $workflowId, string $runId): JsonR
}

$pageSize = $validated['page_size'] ?? 100;
$afterSequence = $this->decodePageToken($validated['next_page_token'] ?? null);
$afterSequence = HistoryPageToken::decode($validated['next_page_token'] ?? null);
$waitNewEvent = (bool) ($validated['wait_new_event'] ?? false);

$events = $waitNewEvent
Expand Down Expand Up @@ -81,7 +82,7 @@ public function show(Request $request, string $workflowId, string $runId): JsonR
'payload' => $this->eventPayload($namespace, $run, $event),
])->all(),
'next_page_token' => $hasMore && $lastSequence !== null
? self::encodePageToken((int) $lastSequence)
? HistoryPageToken::encode((int) $lastSequence)
: null,
];

Expand Down Expand Up @@ -229,26 +230,6 @@ private function compatibilityFleetReason(string $namespace, WorkflowRun $run):
);
}

private function decodePageToken(?string $token): ?int
{
if (! is_string($token) || trim($token) === '') {
return null;
}

$decoded = base64_decode($token, true);

if (! is_string($decoded) || ! ctype_digit($decoded)) {
return null;
}

return (int) $decoded;
}

private static function encodePageToken(int $sequence): string
{
return base64_encode((string) $sequence);
}

/**
* Surface the server-derived principal recorded on the underlying
* command at the top of the event response so audit clients can
Expand Down
45 changes: 43 additions & 2 deletions app/Http/Controllers/Api/SystemController.php
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,9 @@
use App\Support\WorkflowTaskFailureMetrics;
use Illuminate\Http\JsonResponse;
use Illuminate\Http\Request;
use Illuminate\Validation\ValidationException;
use Workflow\V2\Contracts\MatchingRole;
use Workflow\V2\Contracts\OperatorObservabilityRepository;
use Workflow\V2\Support\HealthCheck;
use Workflow\V2\Support\OperatorDashboardSummary;
use Workflow\V2\Support\OperatorMetrics;
Expand Down Expand Up @@ -185,14 +187,53 @@ public function boundedOperatorDashboard(Request $request): JsonResponse
return $this->dashboardResponse($request, includeHistoryAudits: false);
}

private function dashboardResponse(Request $request, bool $includeHistoryAudits): JsonResponse
public function workflowTypeOperatorDashboard(Request $request): JsonResponse
{
if ($response = ControlPlaneProtocol::rejectUnsupported($request)) {
return $response;
}
$validated = $request->validate([
'workflow_types' => ['required', 'string', 'json', 'max:4096'],
]);
$encoded = $validated['workflow_types'];
$types = json_decode($encoded, flags: JSON_THROW_ON_ERROR);
if (strlen(rawurlencode($encoded)) > 4096 || ! is_array($types) || ! array_is_list($types)) {
throw ValidationException::withMessages([
'workflow_types' => 'Provide a JSON list of workflow types with at most 4096 encoded bytes.',
]);
}
foreach ($types as $type) {
if (! is_string($type) || $type === '' || mb_strlen($type) > 255) {
throw ValidationException::withMessages([
'workflow_types' => 'Workflow types must be nonempty strings of at most 255 characters.',
]);
}
}

return $this->dashboardResponse($request, includeHistoryAudits: false, workflowTypes: $types);
}

/** @param list<string>|null $workflowTypes */
private function dashboardResponse(Request $request, bool $includeHistoryAudits, ?array $workflowTypes = null): JsonResponse
{
if ($response = ControlPlaneProtocol::rejectUnsupported($request)) {
return $response;
}

$namespace = (string) $request->attributes->get('namespace');
$dashboard = OperatorDashboardSummary::snapshot(null, $namespace, $includeHistoryAudits);
if ($workflowTypes === null) {
$dashboard = OperatorDashboardSummary::snapshot(null, $namespace, $includeHistoryAudits);
} else {
$observer = app(OperatorObservabilityRepository::class);
if (! method_exists($observer, 'workflowTypeDashboardSummary')) {
return ControlPlaneProtocol::json([
'message' => 'The installed workflow observer cannot filter dashboard totals by workflow type.',
'reason' => 'backend_capability_unavailable',
'capability' => 'workflow_type_dashboard',
], 501);
}
$dashboard = $observer->workflowTypeDashboardSummary($workflowTypes, namespace: $namespace);
}
if (! $includeHistoryAudits) {
$dashboard['operator_metrics']['capacity_evidence'] = $this->capacityEvidence->snapshot($namespace);
}
Expand Down
24 changes: 24 additions & 0 deletions app/Support/HistoryPageToken.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,24 @@
<?php

namespace App\Support;

final class HistoryPageToken
{
public static function decode(?string $token): ?int
{
if (! is_string($token) || trim($token) === '') {
return null;
}

$decoded = base64_decode($token, true);

return is_string($decoded) && ctype_digit($decoded)
? (int) $decoded
: null;
}

public static function encode(int $sequence): string
{
return base64_encode((string) $sequence);
}
}
98 changes: 80 additions & 18 deletions app/Support/WorkflowRunDiagnostics.php
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,10 @@
use App\Models\WorkerRegistration;
use Carbon\CarbonInterface;
use Illuminate\Support\Carbon;
use Illuminate\Support\Collection;
use Workflow\V2\Enums\ActivityAttemptStatus;
use Workflow\V2\Enums\ActivityStatus;
use Workflow\V2\Enums\HistoryEventType;
use Workflow\V2\Enums\TaskStatus;
use Workflow\V2\Enums\TaskType;
use Workflow\V2\Models\ActivityAttempt;
Expand Down Expand Up @@ -56,7 +58,13 @@ public function forRun(string $namespace, WorkflowRun $run, bool $includeLastEve
$activityTaskQueues = $this->activityTaskQueues($namespace, $pendingActivities);
$lastEvent = $this->lastEvent($run, $includeLastEventPayload);
$nextScheduledEvent = $this->nextScheduledEvent($summary, $taskRows->all());
$recentFailures = $this->recentFailures($run);
$failureRows = WorkflowFailure::query()
->where('workflow_run_id', $run->id)
->latest('created_at')
->orderByDesc('id')
->limit(self::FAILURE_LIMIT + 1)
->get();
$recentFailures = $this->recentFailures($run, $failureRows->take(self::FAILURE_LIMIT));
$latestWorkflowTaskFailure = $this->latestWorkflowTaskFailure($run);

$payload = [
Expand All @@ -71,6 +79,7 @@ public function forRun(string $namespace, WorkflowRun $run, bool $includeLastEve
'task_queue' => $taskQueue,
'activity_task_queues' => $activityTaskQueues,
'recent_failures' => $recentFailures,
'recent_failures_truncated' => $failureRows->count() > self::FAILURE_LIMIT,
'latest_workflow_task_failure' => $latestWorkflowTaskFailure,
'compatibility' => $this->compatibility($namespace, $run, $summary, $taskQueue),
'cancellation_cascade_supported' => class_exists(CancellationCascadeView::class),
Expand Down Expand Up @@ -663,28 +672,81 @@ private function activityTaskQueues(string $namespace, array $pendingActivities)
}

/**
* @param Collection<int, WorkflowFailure> $failures
* @return list<array<string, mixed>>
*/
private function recentFailures(WorkflowRun $run): array
private function recentFailures(WorkflowRun $run, Collection $failures): array
{
return WorkflowFailure::query()
$events = $this->failureEvents($run, $failures->pluck('id')->all());

return $failures
->map(function (WorkflowFailure $failure) use ($run, $events): array {
$event = $events->get($failure->id);

return $this->compact([
'failure_id' => $failure->id,
'source_kind' => $failure->source_kind,
'source_id' => $failure->source_id,
'propagation_kind' => $failure->propagation_kind,
'failure_category' => $this->enumValue($failure->failure_category),
'exception_class' => $failure->exception_class,
'message' => $failure->message,
'non_retryable' => (bool) $failure->non_retryable,
'handled' => (bool) $failure->handled,
'created_at' => $this->timestamp($failure->created_at),
'supporting_event' => $event instanceof WorkflowHistoryEvent ? [
'state' => 'retained',
'sequence' => (int) $event->sequence,
'event_type' => $this->enumValue($event->event_type),
'recorded_at' => $this->timestamp($event->recorded_at),
'next_page_token' => HistoryPageToken::encode(max(0, (int) $event->sequence - 1)),
] : [
'state' => $run->details_pruned_at === null ? 'unavailable' : 'pruned',
],
]);
})
->all();
}

/**
* @param list<string> $failureIds
* @return Collection<string, WorkflowHistoryEvent>
*/
private function failureEvents(WorkflowRun $run, array $failureIds): Collection
{
if ($failureIds === []) {
return collect();
}

// Aggregate only the requested failure identities. The result contains
// at most ten scalar rows, even when the retained history is large.
$references = WorkflowHistoryEvent::query()->toBase()
->where('workflow_run_id', $run->id)
->latest('created_at')
->limit(self::FAILURE_LIMIT)
->whereIn('payload->failure_id', $failureIds)
->whereIn('event_type', [
HistoryEventType::ActivityFailed->value,
HistoryEventType::ActivityTimedOut->value,
HistoryEventType::ChildRunFailed->value,
HistoryEventType::ChildRunCancelled->value,
HistoryEventType::ChildRunTerminated->value,
HistoryEventType::WorkflowFailed->value,
HistoryEventType::WorkflowTimedOut->value,
HistoryEventType::WorkflowCancelled->value,
HistoryEventType::WorkflowTerminated->value,
HistoryEventType::UpdateCompleted->value,
])
->select('payload->failure_id as failure_id')
->selectRaw('MAX(sequence) as sequence')
->groupBy('payload->failure_id')
->get();
$failureBySequence = $references->pluck('failure_id', 'sequence');

return WorkflowHistoryEvent::query()
->where('workflow_run_id', $run->id)
->whereIn('sequence', $failureBySequence->keys()->all())
->select(['id', 'workflow_run_id', 'sequence', 'event_type', 'recorded_at'])
->get()
->map(fn (WorkflowFailure $failure): array => $this->compact([
'failure_id' => $failure->id,
'source_kind' => $failure->source_kind,
'source_id' => $failure->source_id,
'propagation_kind' => $failure->propagation_kind,
'failure_category' => $this->enumValue($failure->failure_category),
'exception_class' => $failure->exception_class,
'message' => $failure->message,
'non_retryable' => (bool) $failure->non_retryable,
'handled' => (bool) $failure->handled,
'created_at' => $this->timestamp($failure->created_at),
]))
->all();
->keyBy(fn (WorkflowHistoryEvent $event): string => (string) $failureBySequence->get($event->sequence));
}

/**
Expand Down
4 changes: 2 additions & 2 deletions composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@
"require": {
"php": "^8.2",
"apache/avro": "^1.12",
"durable-workflow/workflow": "2.4.3",
"durable-workflow/workflow": "2.4.4",
"laravel/framework": "^13.30",
"laravel/tinker": "^3.0",
"league/flysystem-aws-s3-v3": "^3.35.3"
Expand Down Expand Up @@ -48,7 +48,7 @@
},
"extra": {
"durable-workflow": {
"product-train": "2.5.3"
"product-train": "2.5.4"
},
"laravel": {
"dont-discover": []
Expand Down
16 changes: 8 additions & 8 deletions composer.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions docker-compose.dedicated-matching.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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.5.3}}
x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.5.4}}

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.5.3}}
APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.5.4}}
APP_DEBUG: ${APP_DEBUG:-false}
DB_CONNECTION: mysql
DB_HOST: mysql
Expand Down
4 changes: 2 additions & 2 deletions docker-compose.memo-rolling.yml
Original file line number Diff line number Diff line change
Expand Up @@ -49,14 +49,14 @@ services:
command: ["server-bootstrap"]
environment:
<<: *runtime-environment
APP_VERSION: ${APP_VERSION:-2.5.3}
APP_VERSION: ${APP_VERSION:-2.5.4}

successor:
image: ${DW_MEMO_SUCCESSOR_IMAGE:-durable-workflow/server-memo-rolling:local}
ports: !override []
environment:
<<: *runtime-environment
APP_VERSION: ${APP_VERSION:-2.5.3}
APP_VERSION: ${APP_VERSION:-2.5.4}
DW_SERVER_ID: memo-successor
DW_SERVER_TOPOLOGY_SHAPE: standalone_server
DW_SERVER_PROCESS_CLASS: server_http_node
Expand Down
4 changes: 2 additions & 2 deletions docker-compose.published.yml
Original file line number Diff line number Diff line change
@@ -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.5.3}}
x-server-image: &server-image ${DW_SERVER_IMAGE:-durableworkflow/server:${DW_SERVER_TAG:-2.5.4}}

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.5.3}}
APP_VERSION: ${APP_VERSION:-${DW_SERVER_TAG:-2.5.4}}
APP_DEBUG: ${APP_DEBUG:-false}
LOG_CHANNEL: ${LOG_CHANNEL:-stderr}
LOG_LEVEL: ${LOG_LEVEL:-info}
Expand Down
Loading
Loading