Skip to content

Commit 70b69d2

Browse files
authored
Merge pull request #562 from hypervel/fix/pre-outbox-framework-correctness
Fix queue publication and worker lifecycle correctness
2 parents c343ec7 + 491fe20 commit 70b69d2

84 files changed

Lines changed: 5250 additions & 1139 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

‎docs/plans/2026-09-04-1527-components-pre-outbox-framework-correctness.md‎

Lines changed: 476 additions & 0 deletions
Large diffs are not rendered by default.

‎docs/todo.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121
## Framework-wide
2222

2323
- Resolve the production dependency cycle between `hypervel/foundation` and `hypervel/testing`, which currently causes test-only classes, PHPUnit integration, and Mockery to propagate through `hypervel/support` into every package that depends on Support. Decide whether Foundation's testing namespace moves to `hypervel/testing` or Foundation changes its dependency to development-only while preserving split-package installation and public namespaces.
24-
- Design a connection-owned service identity and capability API for external backends that packages can query without repeated hot-path probes. Redis/Valkey and database connections already expose fragments of this information in different forms; prefer lazy detection cached for the current connection or pool generation, with invalidation on reconnect and purge, over an eager process-global startup registry that performs unused I/O or survives a backend change. Start with concrete consumers and capability checks rather than a universal version-comparison abstraction. Database Queue's `SKIP LOCKED` selection is one current consumer: MySQL/MariaDB engine and version detection runs on every pop and must not be cached on the worker-lived queue object.
24+
- Design a connection-owned service identity and capability API for external backends that packages can query without repeated hot-path probes. Redis/Valkey and database connections already expose fragments of this information in different forms; prefer lazy detection cached for the current connection or pool generation, with invalidation on reconnect and purge, over an eager process-global startup registry that performs unused I/O or survives a backend change. Start with concrete consumers and capability checks rather than a universal version-comparison abstraction.
2525
- Find a clean, simple framework-wide solution for configuration-dependent services resolved before worker configuration reload. `server:reload` refreshes the existing configuration repository, but objects that have already copied configuration into their own state remain stale. For example, `SentryServiceProvider` eagerly resolves a worker-lifetime Hub and client during boot, so DSN, environment, and sampling changes are not applied until a full restart; resolving `Cache::store('some-store')` from a service provider populates `CacheManager`'s store cache before reload, so changes to that store's driver, connection, prefix, or other captured configuration are likewise not applied. Define the reload contract, audit framework-owned eager resolutions and manager caches, and solve the lifecycle at their shared owning boundary instead of adding package-specific refresh hooks or application workarounds.
2626
- Investigate where requiring and directly using a PHP extension would make framework code significantly faster than its current pure-PHP implementation. The framework already declares bundled extensions it depends on, so the question is which hot paths are doing in PHP what a C extension does natively. The worked example is `ext-gmp` for identifier encoding: UUID and ULID string conversion and any base32/base58/base62 short-id work exceed 64 bits, so `ramsey/uuid` and `symfony/uid` convert them digit by digit in PHP, while `gmp_init()`/`gmp_strval()` do arbitrary-base conversion natively — a hand-rolled base-36 UUID conversion measured 14.6 µs against 0.5 µs for the GMP equivalent with byte-identical output. Anything that fits in a 64-bit int (snowflakes, timestamps, counters) needs no extension, and hashing, encryption, and signatures are already C. Measure `Str::uuid()`/`Str::ulid()` and the other candidates before adding a requirement, and weigh each new extension against installation cost.
2727
- Convert untyped `$config->get()` calls across `src/` to the typed getters (`string()`, `integer()`, `float()`, `boolean()`, `array()`) without call-site defaults, for every key that isn't genuinely nullable. Defaults live in the merged config files — declare any key currently defaulted only at a call site in its package's config file as part of the conversion. Typed getters throw `InvalidArgumentException` naming the key on misconfiguration instead of letting a wrong type propagate silently, and give phpstan real return types. Bootstrap code that runs before config merging keeps its call-site defaults. Approved modernization per the Porting Packages policy in `AGENTS.md`; new code already follows the rule.

‎src/broadcasting/src/BroadcastManager.php‎

Lines changed: 15 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@
1212
use Hypervel\Broadcasting\Broadcasters\NullBroadcaster;
1313
use Hypervel\Broadcasting\Broadcasters\PusherBroadcaster;
1414
use Hypervel\Broadcasting\Broadcasters\RedisBroadcaster;
15+
use Hypervel\Bus\DispatchLockContext;
1516
use Hypervel\Bus\UniqueLock;
1617
use Hypervel\Contracts\Broadcasting\Broadcaster;
1718
use Hypervel\Contracts\Broadcasting\Factory as BroadcastingFactoryContract;
@@ -237,6 +238,17 @@ public function queue(mixed $event): void
237238
)
238239
->pushOn($queue, $broadcastEvent);
239240

241+
if ($event instanceof ShouldBeUnique) {
242+
$dispatch = $push;
243+
$push = static function () use ($broadcastEvent, $dispatch): mixed {
244+
try {
245+
return $dispatch();
246+
} finally {
247+
DispatchLockContext::release($broadcastEvent);
248+
}
249+
};
250+
}
251+
240252
$event instanceof ShouldRescue
241253
? $this->rescue($push)
242254
: $push();
@@ -245,13 +257,10 @@ public function queue(mixed $event): void
245257
/**
246258
* Determine if the broadcastable event must be unique and determine if we can acquire the necessary lock.
247259
*/
248-
protected function mustBeUniqueAndCannotAcquireLock(mixed $event): bool
260+
protected function mustBeUniqueAndCannotAcquireLock(object $event): bool
249261
{
250-
return ! (new UniqueLock(
251-
method_exists($event, 'uniqueVia')
252-
? $event->uniqueVia()
253-
: $this->app->make(Cache::class)
254-
))->acquire($event);
262+
return ! (new UniqueLock($this->app->make(Cache::class)))
263+
->acquireForDispatch($event);
255264
}
256265

257266
/**

‎src/broadcasting/src/UniqueBroadcastEvent.php‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ public function __construct(mixed $event)
4343
/**
4444
* Resolve the cache implementation that should manage the event's uniqueness.
4545
*/
46-
public function uniqueVia(): Repository
46+
public function uniqueVia(): ?Repository
4747
{
4848
return method_exists($this->event, 'uniqueVia')
4949
? $this->event->uniqueVia()

‎src/bus/README.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,12 @@ Bus for Hypervel
33

44
[![Ask DeepWiki](https://deepwiki.com/badge.svg)](https://deepwiki.com/hypervel/bus)
55

6-
Ported from: https://github.com/laravel/framework/tree/13.x/src/Illuminate/Bus
6+
Documentation: https://hypervel.org/docs/queues
77

88
## Differences From Laravel
99

1010
Hypervel does not include Laravel's DynamoDB batch repository because DynamoDB is not a supported database backend.
1111

1212
`DatabaseBatchRepository::setConnection()` is intentionally omitted. The repository is shared for the worker lifetime, so mutating its connection would race across coroutines. Configure `queue.batching.database` instead; each repository operation resolves that connection when it runs.
13+
14+
Ported from: https://github.com/laravel/framework/tree/13.x/src/Illuminate/Bus

‎src/bus/src/DebounceLock.php‎

Lines changed: 96 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@
99
use Hypervel\Queue\Attributes\ReadsQueueAttributes;
1010
use Hypervel\Support\CarbonImmutable;
1111
use Hypervel\Support\Str;
12+
use Throwable;
1213

1314
class DebounceLock
1415
{
@@ -31,20 +32,72 @@ public function __construct(
3132
*/
3233
public function acquire(mixed $job, ?int $debounceFor = null, ?int $maxWait = null): array
3334
{
34-
$cache = $this->resolveCache($job);
35+
[$cache, $key, $ttl, $maxWait] = $this->resolveLockValues($job, $debounceFor, $maxWait);
36+
37+
return $this->acquireResolvedLock($cache, $key, $ttl, $maxWait);
38+
}
39+
40+
/**
41+
* Store and register a debounce owner token for a dispatch operation.
42+
*
43+
* @return array{owner: string, maxWaitExceeded: bool}
44+
*/
45+
public function acquireForDispatch(object $job, ?int $debounceFor = null, ?int $maxWait = null): array
46+
{
47+
[$cache, $key, $ttl, $maxWait] = $this->resolveLockValues($job, $debounceFor, $maxWait);
48+
$result = $this->acquireResolvedLock($cache, $key, $ttl, $maxWait);
49+
50+
if (isset(class_uses_recursive($job)[Queueable::class])) {
51+
$job->debounceOwner = $result['owner'];
52+
}
53+
54+
DispatchLockContext::registerDebounce($job, $cache, $key, $result['owner']);
3555

56+
return $result;
57+
}
58+
59+
/**
60+
* Resolve the values needed to acquire a debounce lock.
61+
*
62+
* @return array{Cache, string, int, null|int}
63+
*/
64+
protected function resolveLockValues(mixed $job, ?int $debounceFor, ?int $maxWait): array
65+
{
66+
$cache = $this->resolveCache($job);
3667
$ttl = max(($debounceFor ?? $this->getDebounceDelay($job)) * 10, 300);
3768

38-
$cache->put($key = static::getKey($job), $owner = Str::random(40), $ttl);
69+
return [
70+
$cache,
71+
static::getKey($job),
72+
$ttl,
73+
$maxWait ?? $this->getMaxDebounceWait($job),
74+
];
75+
}
76+
77+
/**
78+
* Store a debounce owner token using resolved values.
79+
*
80+
* @return array{owner: string, maxWaitExceeded: bool}
81+
*/
82+
protected function acquireResolvedLock(Cache $cache, string $key, int $ttl, ?int $maxWait): array
83+
{
84+
$cache->put($key, $owner = Str::random(40), $ttl);
85+
86+
try {
87+
$maxWaitExceeded = $this->maxWaitExceeded($cache, $key, $ttl, $maxWait);
88+
} catch (Throwable $exception) {
89+
try {
90+
static::releaseOwned($cache, $key, $owner);
91+
} catch (Throwable) {
92+
// A cleanup failure must not replace the acquisition failure.
93+
}
94+
95+
throw $exception;
96+
}
3997

4098
return [
4199
'owner' => $owner,
42-
'maxWaitExceeded' => $this->maxWaitExceeded(
43-
$cache,
44-
$key,
45-
$ttl,
46-
$maxWait ?? $this->getMaxDebounceWait($job)
47-
),
100+
'maxWaitExceeded' => $maxWaitExceeded,
48101
];
49102
}
50103

@@ -93,16 +146,47 @@ public function getCurrentOwner(mixed $job): ?string
93146
*/
94147
public function release(mixed $job, string $owner = ''): void
95148
{
96-
$key = static::getKey($job);
97-
98149
$cache = $this->resolveCache($job);
99150

151+
static::releaseOwned($cache, static::getKey($job), $owner);
152+
}
153+
154+
/**
155+
* Remove the maximum wait timestamp for the given job.
156+
*/
157+
public function releaseMaxWait(mixed $job): void
158+
{
159+
$this->resolveCache($job)->forget(static::getKey($job) . ':first_dispatched_at');
160+
}
161+
162+
/**
163+
* Release a debounce lock from already-resolved provenance.
164+
*
165+
* @internal
166+
*/
167+
public static function releaseOwned(Cache $cache, string $key, string $owner = ''): void
168+
{
100169
if ($owner !== '' && $cache->get($key) !== $owner) {
101170
return;
102171
}
103172

104-
$cache->forget($key);
105-
$cache->forget($key . ':first_dispatched_at');
173+
$exception = null;
174+
175+
try {
176+
$cache->forget($key);
177+
} catch (Throwable $throwable) {
178+
$exception = $throwable;
179+
}
180+
181+
try {
182+
$cache->forget($key . ':first_dispatched_at');
183+
} catch (Throwable $throwable) {
184+
$exception ??= $throwable;
185+
}
186+
187+
if ($exception !== null) {
188+
throw $exception;
189+
}
106190
}
107191

108192
/**

0 commit comments

Comments
 (0)