DefaultConcurrencyHandler
in package
implements
ConcurrencyHandler
Table of Contents
Interfaces
Properties
- $concurrentSlots : array<string, int>
- $config : array<string|int, mixed>
- $state : StateAdapter
Methods
- __construct() : mixed
- process() : void
- Process an incoming message, applying the concurrency strategy.
- applyStrategy() : void
- dequeueAll() : array<string|int, Message>
- drainAllQueued() : void
- processConcurrent() : void
- processDebounce() : void
- processDrop() : void
- processQueue() : void
Properties
$concurrentSlots
private
array<string, int>
$concurrentSlots
= []
$config read-only
private
array<string|int, mixed>
$config
= []
$state read-only
private
StateAdapter
$state
Methods
__construct()
public
__construct(StateAdapter $state[, array<string|int, mixed> $config = [] ]) : mixed
Parameters
- $state : StateAdapter
- $config : array<string|int, mixed> = []
process()
Process an incoming message, applying the concurrency strategy.
public
process(Adapter $adapter, string $threadId, Message $message, callable $processCallback[, ServerRequestInterface|null $request = null ]) : void
Parameters
- $adapter : Adapter
-
The platform adapter
- $threadId : string
-
The canonical thread ID
- $message : Message
-
The incoming message (post-dedup, post-middleware)
- $processCallback : callable
-
fn(Adapter, string $threadId, Message, array $skippedMessages, int $totalSinceLastHandler): void
- $request : ServerRequestInterface|null = null
-
The original PSR-7 request (for job serialization)
applyStrategy()
private
applyStrategy(Strategy $strategy, Adapter $adapter, string $threadId, string $lockKey, Message $message, Handler $handler, int $debounceMs, int $maxConcurrent, int $maxQueueSize, callable $processCallback) : void
Parameters
dequeueAll()
private
dequeueAll(string $threadId, Handler $handler) : array<string|int, Message>
Parameters
- $threadId : string
- $handler : Handler
Return values
array<string|int, Message>drainAllQueued()
private
drainAllQueued(Adapter $adapter, string $threadId, Handler $handler, callable $processCallback) : void
Parameters
processConcurrent()
private
processConcurrent(Adapter $adapter, string $threadId, Message $message, int $maxConcurrent, callable $processCallback) : void
Parameters
processDebounce()
private
processDebounce(Adapter $adapter, string $threadId, string $lockKey, Message $message, Handler $handler, int $debounceMs, int $maxQueueSize, callable $processCallback) : void
Parameters
processDrop()
private
processDrop(Adapter $adapter, string $threadId, string $lockKey, Message $message, Handler $handler, callable $processCallback) : void
Parameters
processQueue()
private
processQueue(Adapter $adapter, string $threadId, string $lockKey, Message $message, Handler $handler, int $maxQueueSize, callable $processCallback) : void