From d902202fe905a9be2d84551e02e639ef5e9b3a16 Mon Sep 17 00:00:00 2001 From: Jeppe Bundgaard Date: Tue, 30 Jun 2026 12:24:47 +0200 Subject: [PATCH] Fallback relay commands after broker dispatch failures --- .../classes/edge_gateway_manager.php | 49 ++++++++++++++++++- .../EdgeGatewayManagerCommandQueueTest.php | 7 +++ 2 files changed, 55 insertions(+), 1 deletion(-) diff --git a/services/nginx/app/modules/edgegateway/classes/edge_gateway_manager.php b/services/nginx/app/modules/edgegateway/classes/edge_gateway_manager.php index 3ab43edc..753ea80a 100644 --- a/services/nginx/app/modules/edgegateway/classes/edge_gateway_manager.php +++ b/services/nginx/app/modules/edgegateway/classes/edge_gateway_manager.php @@ -5638,6 +5638,7 @@ BASH; $delivery = (array)($job->delivery_json->value() ?? []); $preferredChannel = (string)($delivery['preferred_channel'] ?? self::DELIVERY_CHANNEL_API); $requireFastPath = !empty($delivery['require_fast_path']); + $waitTimeoutSeconds = self::COMMAND_WAIT_TIMEOUT_SECONDS; $effectiveStatus = self::resolveGatewayStatus( $gateway->status->value() === null ? null : (string)$gateway->status->value(), @@ -5673,6 +5674,24 @@ BASH; ); throw $exception; } + + if ($this->commandDisallowsAutomaticRetry($job)) { + $this->finalizeCommandJob( + $job, + false, + [], + $exception->getMessage(), + $gateway, + $this->isCommandTimeoutError($exception->getMessage()) ? 'TIMED_OUT' : 'FAILED' + ); + throw $exception; + } + + $this->requeueCommandForApiPolling($job, $exception->getMessage()); + $waitTimeoutSeconds = max( + self::COMMAND_WAIT_TIMEOUT_SECONDS, + self::COMMAND_POLL_TIMEOUT_SECONDS + 5 + ); } } } @@ -5682,7 +5701,35 @@ BASH; throw new Exception('Edge broker fast path is unavailable'); } - return $this->waitForCommandResult((int)$job->id); + return $this->waitForCommandResult((int)$job->id, $waitTimeoutSeconds); + } + + private function commandDisallowsAutomaticRetry(edge_gateway_command_jobs_o $job): bool + { + $delivery = (array)($job->delivery_json->value() ?? []); + return filter_var($delivery['no_auto_retry'] ?? false, FILTER_VALIDATE_BOOLEAN); + } + + private function requeueCommandForApiPolling(edge_gateway_command_jobs_o $job, string $errorMessage): void + { + $delivery = $this->buildDeliveryMetadata( + array_merge( + (array)($job->delivery_json->value() ?? []), + [ + 'preferred_channel' => self::DELIVERY_CHANNEL_API, + 'delivery_channel' => null, + 'fallback_reason' => 'broker_dispatch_failed', + 'last_dispatch_error' => $errorMessage, + ] + ), + self::COMMAND_EXPIRES_AFTER_SECONDS + ); + + $job->status->set('PENDING'); + $job->response_json->set([]); + $job->completed_at->set(null); + $job->error_message->set(null); + $job->delivery_json->set($delivery); } private function isCommandTimeoutError(string $errorMessage): bool diff --git a/services/nginx/app/tests/Unit/Selfserve/EdgeGatewayManagerCommandQueueTest.php b/services/nginx/app/tests/Unit/Selfserve/EdgeGatewayManagerCommandQueueTest.php index df0750d5..0fa059c2 100644 --- a/services/nginx/app/tests/Unit/Selfserve/EdgeGatewayManagerCommandQueueTest.php +++ b/services/nginx/app/tests/Unit/Selfserve/EdgeGatewayManagerCommandQueueTest.php @@ -42,6 +42,13 @@ it('keeps relay dispatch and discovery queueing on the edge gateway manager', fu expect($managerSource)->toContain("'batch_dedupe_key' => \$batchDedupeKey"); expect($managerSource)->toContain("'idempotency_key' => \$batchId . ':' . (string)\$command['target']"); expect($managerSource)->toContain("'no_auto_retry' => \$isGatePulse"); + expect($managerSource)->toContain('private function commandDisallowsAutomaticRetry'); + expect($managerSource)->toContain('private function requeueCommandForApiPolling'); + expect($managerSource)->toContain("'preferred_channel' => self::DELIVERY_CHANNEL_API"); + expect($managerSource)->toContain("'fallback_reason' => 'broker_dispatch_failed'"); + expect($managerSource)->toContain("\$this->requeueCommandForApiPolling(\$job, \$exception->getMessage());"); + expect($managerSource)->toContain('self::COMMAND_POLL_TIMEOUT_SECONDS + 5'); + expect($managerSource)->toContain('return $this->waitForCommandResult((int)$job->id, $waitTimeoutSeconds);'); expect($managerSource)->toContain('private function buildRelayBatchDedupeKey'); expect($managerSource)->toContain('private function findRecentRelayBatchByDedupeKey'); expect($managerSource)->toContain('private function enforceRelayCommandRateLimit');