From c8e6efe203001689ad874e96489ae5e1c9757e85 Mon Sep 17 00:00:00 2001 From: Sebastian Krupinski Date: Sat, 8 Aug 2026 00:20:32 -0400 Subject: [PATCH] feat: improve deferred event processing Signed-off-by: Sebastian Krupinski --- .../Execution/TerminationReport.php | 3 + core/lib/Event/DeferredProcessingResult.php | 3 + core/lib/Event/EventDispatcher.php | 55 +++++++++++++++---- core/lib/Kernel.php | 9 +++ tests/php/Unit/Event/EventDispatcherTest.php | 48 ++++++++++++++++ 5 files changed, 107 insertions(+), 11 deletions(-) diff --git a/core/lib/Application/Execution/TerminationReport.php b/core/lib/Application/Execution/TerminationReport.php index 60a6e33..e401544 100644 --- a/core/lib/Application/Execution/TerminationReport.php +++ b/core/lib/Application/Execution/TerminationReport.php @@ -15,6 +15,9 @@ final readonly class TerminationReport public array $failures = [], public bool $deadlineExceeded = false, public bool $limitExceeded = false, + public int $deferredListenerInvocations = 0, + public bool $deferredEventLimitExceeded = false, + public bool $deferredListenerInvocationLimitExceeded = false, ) { } } diff --git a/core/lib/Event/DeferredProcessingResult.php b/core/lib/Event/DeferredProcessingResult.php index 3a313c7..e9ba31b 100644 --- a/core/lib/Event/DeferredProcessingResult.php +++ b/core/lib/Event/DeferredProcessingResult.php @@ -11,6 +11,9 @@ final readonly class DeferredProcessingResult public int $remaining, public bool $deadlineExceeded, public bool $limitExceeded = false, + public int $listenerInvocations = 0, + public bool $eventLimitExceeded = false, + public bool $listenerInvocationLimitExceeded = false, ) { } } diff --git a/core/lib/Event/EventDispatcher.php b/core/lib/Event/EventDispatcher.php index 3871a04..32e0772 100644 --- a/core/lib/Event/EventDispatcher.php +++ b/core/lib/Event/EventDispatcher.php @@ -13,6 +13,10 @@ use Psr\Log\LoggerInterface; final class EventDispatcher implements EventDispatcherInterface, DeferredEventProcessorInterface { + private const DEFAULT_DEFERRED_PROCESSING_TIMEOUT_SECONDS = 300.0; + private const DEFAULT_MAX_DEFERRED_EVENTS = 1000; + private const DEFAULT_MAX_DEFERRED_LISTENER_INVOCATIONS = 50000; + /** @var array> */ private array $deferred = []; private ?string $activeExecution = null; @@ -22,7 +26,19 @@ final class EventDispatcher implements EventDispatcherInterface, DeferredEventPr private readonly EventListenerRegistry $registry, private readonly ContainerInterface $container, private readonly LoggerInterface $logger, + private readonly float $deferredProcessingTimeoutSeconds = self::DEFAULT_DEFERRED_PROCESSING_TIMEOUT_SECONDS, + private readonly int $maxDeferredEvents = self::DEFAULT_MAX_DEFERRED_EVENTS, + private readonly int $maxDeferredListenerInvocations = self::DEFAULT_MAX_DEFERRED_LISTENER_INVOCATIONS, ) { + if ($this->deferredProcessingTimeoutSeconds <= 0) { + throw new \InvalidArgumentException('The deferred processing timeout must be greater than zero.'); + } + if ($this->maxDeferredEvents <= 0) { + throw new \InvalidArgumentException('The deferred event limit must be greater than zero.'); + } + if ($this->maxDeferredListenerInvocations <= 0) { + throw new \InvalidArgumentException('The deferred listener invocation limit must be greater than zero.'); + } } public function dispatch(Event $event): void @@ -61,13 +77,15 @@ final class EventDispatcher implements EventDispatcherInterface, DeferredEventPr } try { - $processed = 0; - $deadline = microtime(true) + 1.0; + $processedEvents = 0; + $listenerInvocations = 0; + $deadline = microtime(true) + $this->deferredProcessingTimeoutSeconds; $deadlineExceeded = false; - $limitExceeded = false; + $eventLimitExceeded = false; + $listenerInvocationLimitExceeded = false; while (($event = array_shift($this->deferred[$executionId])) !== null) { - if ($processed >= 1000) { - $limitExceeded = true; + if ($processedEvents >= $this->maxDeferredEvents) { + $eventLimitExceeded = true; array_unshift($this->deferred[$executionId], $event); break; } @@ -76,14 +94,29 @@ final class EventDispatcher implements EventDispatcherInterface, DeferredEventPr array_unshift($this->deferred[$executionId], $event); break; } - $processed += $this->invoke($event, DeliveryMode::Deferred); + + $eventListenerCount = count($this->registry->listeners( + $event->getName(), + DeliveryMode::Deferred, + )); + if ($listenerInvocations + $eventListenerCount > $this->maxDeferredListenerInvocations) { + $listenerInvocationLimitExceeded = true; + array_unshift($this->deferred[$executionId], $event); + break; + } + + $listenerInvocations += $this->invoke($event, DeliveryMode::Deferred); + $processedEvents++; } return new DeferredProcessingResult( - $processed, - count($this->deferred[$executionId]), - $deadlineExceeded, - $limitExceeded, + processed: $processedEvents, + remaining: count($this->deferred[$executionId]), + deadlineExceeded: $deadlineExceeded, + limitExceeded: $eventLimitExceeded || $listenerInvocationLimitExceeded, + listenerInvocations: $listenerInvocations, + eventLimitExceeded: $eventLimitExceeded, + listenerInvocationLimitExceeded: $listenerInvocationLimitExceeded, ); } finally { $this->discardDeferred($executionId); @@ -106,10 +139,10 @@ final class EventDispatcher implements EventDispatcherInterface, DeferredEventPr break; } + $processed++; try { $service = $this->container->get($listener->service); $service->{$listener->method}($event); - $processed++; } catch (\Throwable $error) { $this->logger->error('Event listener failed.', [ 'event' => $event->getName(), diff --git a/core/lib/Kernel.php b/core/lib/Kernel.php index 3635bc9..dd4f2c0 100644 --- a/core/lib/Kernel.php +++ b/core/lib/Kernel.php @@ -238,6 +238,9 @@ class Kernel implements KernelInterface $remaining = 0; $deadlineExceeded = false; $limitExceeded = false; + $listenerInvocations = 0; + $eventLimitExceeded = false; + $listenerInvocationLimitExceeded = false; $failures = []; try { @@ -249,6 +252,9 @@ class Kernel implements KernelInterface $remaining = $result->remaining; $deadlineExceeded = $result->deadlineExceeded; $limitExceeded = $result->limitExceeded; + $listenerInvocations = $result->listenerInvocations; + $eventLimitExceeded = $result->eventLimitExceeded; + $listenerInvocationLimitExceeded = $result->listenerInvocationLimitExceeded; } } catch (\Throwable $e) { $failures[] = $e; @@ -287,6 +293,9 @@ class Kernel implements KernelInterface failures: $failures, deadlineExceeded: $deadlineExceeded, limitExceeded: $limitExceeded, + deferredListenerInvocations: $listenerInvocations, + deferredEventLimitExceeded: $eventLimitExceeded, + deferredListenerInvocationLimitExceeded: $listenerInvocationLimitExceeded, ); } diff --git a/tests/php/Unit/Event/EventDispatcherTest.php b/tests/php/Unit/Event/EventDispatcherTest.php index ba4310c..138162b 100644 --- a/tests/php/Unit/Event/EventDispatcherTest.php +++ b/tests/php/Unit/Event/EventDispatcherTest.php @@ -230,9 +230,57 @@ final class EventDispatcherTest extends TestCase $result = $dispatcher->processDeferred('test'); self::assertSame(1000, $result->processed); + self::assertSame(1000, $result->listenerInvocations); self::assertSame(1, $result->remaining); self::assertFalse($result->deadlineExceeded); self::assertTrue($result->limitExceeded); + self::assertTrue($result->eventLimitExceeded); + self::assertFalse($result->listenerInvocationLimitExceeded); + } + + #[Test] + #[TestDox('Deferred event and listener invocation limits are tracked separately')] + public function boundsDeferredListenerInvocations(): void + { + $recursive = new RecursiveListener(); + $recording = new RecordingListener(); + $registry = new EventListenerRegistry(); + $registry->listen( + 'test', + 'test.event', + RecursiveListener::class, + 'deferred', + DeliveryMode::Deferred, + ); + $registry->listen( + 'test', + 'test.event', + RecordingListener::class, + 'deferred', + DeliveryMode::Deferred, + ); + $registry->freeze(); + $dispatcher = new EventDispatcher( + $registry, + new RecordingContainer([ + RecursiveListener::class => $recursive, + RecordingListener::class => $recording, + ]), + new NullLogger(), + maxDeferredListenerInvocations: 3, + ); + $recursive->dispatcher = $dispatcher; + $dispatcher->beginExecution('test'); + $dispatcher->dispatch(new Event('test.event')); + + $result = $dispatcher->processDeferred('test'); + + self::assertSame(1, $result->processed); + self::assertSame(2, $result->listenerInvocations); + self::assertSame(1, $result->remaining); + self::assertTrue($result->limitExceeded); + self::assertFalse($result->eventLimitExceeded); + self::assertTrue($result->listenerInvocationLimitExceeded); } }