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
Original file line number Diff line number Diff line change
Expand Up @@ -5,12 +5,16 @@
namespace Hypervel\Reverb\Protocols\Pusher\Channels\Concerns;

use Hypervel\Reverb\Contracts\Connection;
use Hypervel\Reverb\FailureReporter;
use Hypervel\Reverb\Protocols\Pusher\Channels\ChannelConnection;
use Hypervel\Reverb\Protocols\Pusher\EventDispatcher;
use Hypervel\Reverb\Protocols\Pusher\MetricsHandler;
use Hypervel\Reverb\Protocols\Pusher\MetricType;
use Hypervel\Reverb\Servers\Hypervel\Contracts\SharedState;
use Hypervel\Reverb\Webhooks\Contracts\WebhookDispatcher;
use Hypervel\Reverb\Webhooks\DeferredWebhookManager;
use Swoole\Coroutine\CanceledException;
use Throwable;

trait InteractsWithPresenceChannels
{
Expand Down Expand Up @@ -153,11 +157,28 @@ public function data(): array
];
}

$snapshot = app(MetricsHandler::class)->gather(
$connection->app(),
MetricType::Presence->value,
['channel' => $this->name()],
);
try {
$snapshot = app(MetricsHandler::class)->gather(
$connection->app(),
MetricType::Presence->value,
['channel' => $this->name()],
);
} catch (CanceledException $exception) {
throw $exception;
} catch (Throwable $exception) {
// The subscription is already committed, so still answer it with
// this worker's members when other workers or servers can't be reached.
FailureReporter::report($exception);

$snapshot = [
'users' => collect($this->connections->all())
->map(fn (ChannelConnection $member): array => $member->data())
->unique('user_id')
->values()
->all(),
];
}

$connections = collect($snapshot['users']);

if ($connections->contains(fn ($connection) => ! isset($connection['user_id']))) {
Expand All @@ -174,7 +195,7 @@ public function data(): array
'presence' => [
'count' => $connections->count(),
'ids' => $connections->map(fn ($connection) => $connection['user_id'])->values()->all(),
'hash' => $connections->keyBy('user_id')->map->user_info->toArray(),
'hash' => $connections->pluck('user_info', 'user_id')->map(fn ($info) => $info ?: (object) [])->all(),
],
];
}
Expand Down
21 changes: 21 additions & 0 deletions src/reverb/src/Protocols/Pusher/Http/Controllers/Controller.php
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,11 @@

abstract class Controller
{
/**
* The number of seconds either side of the current time a request signature remains valid.
*/
protected const int SIGNATURE_TOLERANCE = 600;

/**
* Verify that the incoming request is valid.
*
Expand Down Expand Up @@ -88,6 +93,8 @@ protected function verifySignature(Request $request, Application $application, s
if (! is_string($authSignature) || ! hash_equals($signature, $authSignature)) {
throw new HttpException(401, 'Authentication signature invalid.');
}

$this->verifySignatureTimestamp($query);
}

/**
Expand All @@ -103,4 +110,18 @@ protected static function formatQueryParametersForVerification(array $params): s
return "{$key}={$value}";
})->implode('&');
}

/**
* Verify that the request signature has not expired.
*
* @throws HttpException
*/
protected function verifySignatureTimestamp(array $query): void
{
$timestamp = $query['auth_timestamp'] ?? null;

if (! is_numeric($timestamp) || abs(time() - (int) $timestamp) > static::SIGNATURE_TOLERANCE) {
throw new HttpException(401, 'Authentication signature invalid.');
}
}
}
17 changes: 8 additions & 9 deletions src/scout/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,16 +3,15 @@ Scout for Hypervel

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

Ported from: https://github.com/laravel/scout
Documentation: https://hypervel.org/docs/scout

Differences From Laravel
---
## Differences From Laravel

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P3: The differences rewrite drops the documented Searchable::removeAllFromSearch() force-flag difference, but the method still deviates from Laravel's parameterless removeAllFromSearch(): src/scout/src/Searchable.php:316 is removeAllFromSearch(bool $force = false), and the scout:flush command enables it via Scout::guardModelFlush. That is a deliberate lasting public-contract difference that the README should keep (or point to the docs section covering it) so porters don't miss the force behavior.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At src/scout/README.md, line 8:

<comment>The differences rewrite drops the documented `Searchable::removeAllFromSearch()` force-flag difference, but the method still deviates from Laravel's parameterless `removeAllFromSearch()`: `src/scout/src/Searchable.php:316` is `removeAllFromSearch(bool $force = false)`, and the `scout:flush` command enables it via `Scout::guardModelFlush`. That is a deliberate lasting public-contract difference that the README should keep (or point to the docs section covering it) so porters don't miss the force behavior.</comment>

<file context>
@@ -3,16 +3,15 @@ Scout for Hypervel
 
-Differences From Laravel
----
+## Differences From Laravel
 
 - Algolia 4 is the only supported Algolia client.
</file context>

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Keeping this out of the README. Its differences section only lists differences that ported Laravel code has to account for. removeAllFromSearch() still works with no arguments, as in Laravel. The optional force argument only matters to an application that registers a model-flush guard, and the Scout documentation covers it under Removing Records and Customizing Scout Lifecycles.


- Algolia 4 is the only supported Algolia client.
- Numeric values passed to Algolia `where`, `whereIn`, and `whereNotIn` compile as numeric comparisons; numeric-looking strings remain facet values.
- Queue mode supports dedicated connection and queue selection; nonqueued indexing is deferred until after HTTP responses and runs immediately without an active request.
- Command imports use bounded coroutine concurrency.
- Meilisearch requests use bounded retries and sign tenant tokens from an explicit parent-key UID and secret.
- Destructive index deletion requires the configured Scout prefix.
- Boot-time lifecycle callbacks can prepare builders, documents, settings, and model flushes; external engines also support completion-aware filtered deletion.
- `Searchable::removeAllFromSearch()` accepts an optional force flag, which the explicit `scout:flush` command enables.
- Without a queue, indexing is deferred until after the HTTP response is sent, and runs immediately outside a request.
- Pausing search syncing with `withoutSyncingToSearch()` or `disableSearchSyncing()` applies only to the current coroutine, so other requests keep indexing. See [Pausing Indexing](https://hypervel.org/docs/scout#pausing-indexing).
- `MeilisearchEngine::generateTenantToken()` takes the search rules, the parent key's UID, the key itself and an optional expiry. Laravel's engine forwards the call to the Meilisearch client, whose method takes the UID, the search rules and an options array. See [Tenant Tokens](https://hypervel.org/docs/scout#meilisearch-tenant-tokens).
- `scout:delete-all-indexes` refuses to run without a configured Scout prefix unless you pass `--force`.

Ported from: https://github.com/laravel/scout
4 changes: 2 additions & 2 deletions src/scout/src/Builder.php
Original file line number Diff line number Diff line change
Expand Up @@ -114,13 +114,13 @@ class Builder
*/
public function __construct(
Model $model,
string $query,
?string $query,
?Closure $callback = null,
bool $softDelete = false
) {
/** @var SearchableInterface&TModel $model */
$this->model = $model;
$this->query = $query;
$this->query = $query ?? '';
$this->callback = $callback;

if ($softDelete) {
Expand Down
2 changes: 1 addition & 1 deletion src/scout/src/Contracts/SearchableInterface.php
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ interface SearchableInterface
*
* @return Builder<Model&static>
*/
public static function search(string $query = '', ?Closure $callback = null): Builder;
public static function search(?string $query = '', ?Closure $callback = null): Builder;
Comment thread
qodo-free-for-open-source-projects[bot] marked this conversation as resolved.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: Normalize $query to '' before resolving $scoutBuilder; search(null) otherwise passes null to custom builders whose constructor requires string and throws a TypeError.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At src/scout/src/Contracts/SearchableInterface.php, line 30:

<comment>Normalize `$query` to `''` before resolving `$scoutBuilder`; `search(null)` otherwise passes null to custom builders whose constructor requires `string` and throws a `TypeError`.</comment>

<file context>
@@ -27,7 +27,7 @@ interface SearchableInterface
      * @return Builder<Model&static>
      */
-    public static function search(string $query = '', ?Closure $callback = null): Builder;
+    public static function search(?string $query = '', ?Closure $callback = null): Builder;
 
     /**
</file context>

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Keeping this as is. Laravel's search() also passes the raw query to static::$scoutBuilder. Since search() returns a Builder, a custom builder is a Builder subclass, and Builder's constructor accepts null and stores it as an empty query. Normalizing the query in search() as well would duplicate that.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2: Existing SearchableInterface implementations with string $query now fail class loading. Add an upgrade note requiring those implementations to widen the parameter to ?string.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At src/scout/src/Contracts/SearchableInterface.php, line 30:

<comment>Existing `SearchableInterface` implementations with `string $query` now fail class loading. Add an upgrade note requiring those implementations to widen the parameter to `?string`.</comment>

<file context>
@@ -27,7 +27,7 @@ interface SearchableInterface
      * @return Builder<Model&static>
      */
-    public static function search(string $query = '', ?Closure $callback = null): Builder;
+    public static function search(?string $query = '', ?Closure $callback = null): Builder;
 
     /**
</file context>

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

No note needed. Laravel's search($query = '', $callback = null) has no parameter type, so overrides ported from Laravel don't declare string $query, and an untyped $query is compatible with ?string. Only an override copied from Hypervel's earlier signature is affected, and the upgrade guide already recommends moving 0.3 code into a fresh 0.4 application and updating it to the Laravel-style APIs rather than listing individual signature changes.


/**
* Get the requested models from an array of object IDs.
Expand Down
2 changes: 1 addition & 1 deletion src/scout/src/Searchable.php
Original file line number Diff line number Diff line change
Expand Up @@ -240,7 +240,7 @@ public function searchIndexShouldBeUpdated(): bool
*
* @return Builder<static>
*/
public static function search(string $query = '', ?Closure $callback = null): Builder
public static function search(?string $query = '', ?Closure $callback = null): Builder
{
// @phpstan-ignore staticProperty.notFound (models may define the documented custom builder property)
$builder = static::$scoutBuilder ?? Builder::class;
Expand Down
2 changes: 2 additions & 0 deletions src/websocket-server/src/Security.php
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,8 @@ class Security

public const string SEC_WEBSOCKET_KEY = 'sec-websocket-key';

public const string SEC_WEBSOCKET_VERSION = 'sec-websocket-version';

public const string SEC_WEBSOCKET_PROTOCOL = 'sec-websocket-protocol';

/**
Expand Down
16 changes: 13 additions & 3 deletions src/websocket-server/src/Server.php
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
use Swoole\WebSocket\Frame;
use Swoole\WebSocket\Server as WebSocketServer;
use Symfony\Component\HttpFoundation\Response;
use Symfony\Component\HttpKernel\Exception\HttpException;
use Throwable;

class Server implements BootstrapsForServer, OnHandshakeInterface, OnCloseInterface, OnMessageInterface
Expand Down Expand Up @@ -98,8 +99,9 @@ public function bootstrapForServer(string $serverName): void
* Handle the WebSocket handshake request.
*
* Converts the Swoole request to HttpFoundation, validates the WebSocket
* security key, dispatches through the Router for route matching and
* middleware execution, then builds the 101 Switching Protocols response.
* security key and version, dispatches through the Router for route
* matching and middleware execution, then builds the 101 Switching
* Protocols response.
*/
public function onHandshake(Request $request, SwooleResponse $response): void
{
Expand Down Expand Up @@ -135,13 +137,21 @@ public function onHandshake(Request $request, SwooleResponse $response): void

$this->logger->debug(sprintf('WebSocket: fd[%d] start a handshake request.', $fd));

// Validate sec-websocket-key before routing
// Validate sec-websocket-key and sec-websocket-version before routing
$key = $httpRequest->headers->get(Security::SEC_WEBSOCKET_KEY);
$security = $this->container->make(Security::class);
if (! $key || $security->isInvalidSecurityKey($key)) {
throw new WebSocketHandshakeException('sec-websocket-key is invalid!');
}

if ($httpRequest->headers->get(Security::SEC_WEBSOCKET_VERSION) !== Security::VERSION) {
throw new HttpException(Response::HTTP_UPGRADE_REQUIRED, 'sec-websocket-version is unsupported!', headers: [
'Upgrade' => 'websocket',
'Connection' => 'Upgrade',
'Sec-WebSocket-Version' => Security::VERSION,
]);
}

// Route matching + middleware via Router.
// dispatchToCallback() performs the full Router context lifecycle
// (findRoute, context setup, RouteMatched event, middleware pipeline)
Expand Down
5 changes: 4 additions & 1 deletion tests/Integration/Cache/CacheFunnelTestCase.php
Original file line number Diff line number Diff line change
Expand Up @@ -240,7 +240,10 @@ public function testFunnelLeaseRefreshExtendsLifetime(): void
$this->markTestSkipped('This cache store does not return refreshable funnel leases.');
}

usleep(1_100_000);
// Whole-second stores can report the same lifetime before and after
// the refresh when a second ends between the reads, unless at least
// two seconds have passed since acquisition.
usleep(2_100_000);

$decayedLifetime = $lease->getRemainingLifetime();
$this->assertNotNull($decayedLifetime);
Expand Down
15 changes: 15 additions & 0 deletions tests/Integration/Redis/DurationLimiterIntegrationTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,8 @@ public function testItFailsImmediatelyOrRetriesForAWhileBasedOnAGivenTimeout():
{
$store = [];

$this->waitForNextSecond();

(new DurationLimiter($this->redis(), 'key', 1, 1))->block(2, function () use (&$store) {
$store[] = 1;
});
Expand Down Expand Up @@ -193,6 +195,8 @@ public function testAcquireResetsAfterDecay(): void
{
$limiter = new DurationLimiter($this->redis(), 'reset-after-decay-key', 1, 1);

$this->waitForNextSecond();

$this->assertTrue($limiter->acquire());
$this->assertFalse($limiter->acquire());

Expand Down Expand Up @@ -224,6 +228,17 @@ public function testAcquireUsesTheSelectedConnectionPrefix(): void
}
}

/**
* Wait until just after the next whole second.
*
* One-second windows end on a whole second, so starting just after one
* keeps an immediate second attempt inside the first window.
*/
private function waitForNextSecond(): void
{
usleep((int) ((1.05 - fmod(microtime(true), 1)) * 1_000_000));
}

/**
* Get the Redis connection for testing.
*/
Expand Down
62 changes: 62 additions & 0 deletions tests/Integration/Reverb/RedisServerTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -164,8 +164,70 @@ public function testPresenceMemberNotificationsWithRedisScaling(): void
$this->assertNotNull($message, 'Client one did not receive member_added for user 2');
$this->assertSame('User 2', $this->decodeEventData($message)['user_info']['name']);

$this->disconnect($clientTwo);

$message = $this->receiveMemberRemoved($clientOne, 2);
$this->assertNotNull($message, 'Client one did not receive member_removed for user 2');

$this->disconnect($clientOne);
}

public function testIncludesExistingMembersInSubscriptionSucceededWhenScaling(): void
{
['client' => $clientOne, 'socketId' => $socketIdOne] = $this->connect();
$this->subscribe($clientOne, $socketIdOne, 'presence-redis-existing-members-channel', [
'user_id' => 1,
'user_info' => ['name' => 'User 1'],
]);

['client' => $clientTwo, 'socketId' => $socketIdTwo] = $this->connect();
$response = $this->subscribe($clientTwo, $socketIdTwo, 'presence-redis-existing-members-channel', [
'user_id' => 2,
'user_info' => ['name' => 'User 2'],
]);

$message = $this->decodeEventMessage($response);
$this->assertSame('pusher_internal:subscription_succeeded', $message['event']);
$this->assertSame([
'presence' => [
'count' => 2,
'ids' => [1, 2],
'hash' => [1 => ['name' => 'User 1'], 2 => ['name' => 'User 2']],
],
], $this->decodeEventData($message));

$this->disconnect($clientOne);
$this->disconnect($clientTwo);
}

public function testDoesNotCacheInternalEventsOnAPresenceCacheChannelWhenScaling(): void
{
['client' => $clientOne, 'socketId' => $socketIdOne] = $this->connect();
$this->subscribe($clientOne, $socketIdOne, 'presence-cache-redis-internal-channel', [
'user_id' => 1,
'user_info' => ['name' => 'User 1'],
]);

['client' => $clientTwo, 'socketId' => $socketIdTwo] = $this->connect();
$this->subscribe($clientTwo, $socketIdTwo, 'presence-cache-redis-internal-channel', [
'user_id' => 2,
'user_info' => ['name' => 'User 2'],
]);

$this->assertNotNull($this->receiveMemberAdded($clientOne, 2), 'Client one did not receive member_added for user 2');

// A cached member event would be replayed to the next subscriber instead of a cache miss.
['client' => $clientThree, 'socketId' => $socketIdThree] = $this->connect();
$this->subscribe($clientThree, $socketIdThree, 'presence-cache-redis-internal-channel', [
'user_id' => 3,
'user_info' => ['name' => 'User 3'],
]);

$this->assertNotNull($this->receiveEvent($clientThree, 'pusher:cache_miss'), 'Client three did not receive a cache miss');

$this->disconnect($clientOne);
$this->disconnect($clientTwo);
$this->disconnect($clientThree);
}

// ── HTTP API with Redis scaling ────────────────────────────────────
Expand Down
37 changes: 37 additions & 0 deletions tests/Integration/Reverb/ServerTest.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,8 @@

namespace Hypervel\Tests\Integration\Reverb;

use Swoole\Coroutine\Http\Client;

/**
* End-to-end integration tests for the Reverb WebSocket server.
*
Expand Down Expand Up @@ -41,6 +43,41 @@ public function testFailsToConnectWithInvalidAppKey(): void
$client->close();
}

public function testRejectsAHandshakeWithAnInvalidWebsocketKey(): void
{
$client = new Client($this->getServerHost(), $this->getServerPort());
$client->set(['timeout' => 5]);
$client->setHeaders([
'Connection' => 'Upgrade',
'Upgrade' => 'websocket',
'Sec-WebSocket-Key' => 'invalid-key',
'Sec-WebSocket-Version' => '13',
]);
$client->get('/app/' . $this->appKey);

$this->assertSame(400, $client->getStatusCode());

$client->close();
}

public function testRejectsAHandshakeRequestingAnUnsupportedWebsocketVersion(): void
{
$client = new Client($this->getServerHost(), $this->getServerPort());
$client->set(['timeout' => 5]);
$client->setHeaders([
'Connection' => 'Upgrade',
'Upgrade' => 'websocket',
'Sec-WebSocket-Key' => 'dGhlIHNhbXBsZSBub25jZQ==',
'Sec-WebSocket-Version' => '8',
]);
$client->get('/app/' . $this->appKey);

$this->assertSame(426, $client->getStatusCode());

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P3: This test verifies only the 426 status code and not the Sec-WebSocket-Version response header, which RFC 6455 §4.2.2 requires on a version-rejection handshake and which the server does send (src/websocket-server/src/Server.php throws the 426 with Sec-WebSocket-Version => Security::VERSION). Since the PR highlights the RFC 6455 426 behavior, assert the header (plus the Upgrade: websocket) so a regression that drops it is caught.

Prompt for AI agents
Check if this issue is valid — if so, understand the root cause and fix it. When an issue isn't valid or won't be fixed in this PR, reply in its thread with the reason and then resolve the thread. At tests/Integration/Reverb/ServerTest.php, line 75:

<comment>This test verifies only the 426 status code and not the `Sec-WebSocket-Version` response header, which RFC 6455 §4.2.2 requires on a version-rejection handshake and which the server does send (src/websocket-server/src/Server.php throws the 426 with `Sec-WebSocket-Version => Security::VERSION`). Since the PR highlights the RFC 6455 426 behavior, assert the header (plus the `Upgrade: websocket`) so a regression that drops it is caught.</comment>

<file context>
@@ -41,6 +43,40 @@ public function testFailsToConnectWithInvalidAppKey(): void
+        ]);
+        $client->get('/app/' . $this->appKey);
+
+        $this->assertSame(426, $client->getStatusCode());
+
+        $client->close();
</file context>
Suggested change
$this->assertSame(426, $client->getStatusCode());
$this->assertSame(426, $client->getStatusCode());
$this->assertSame('13', $client->getHeaders()['sec-websocket-version'] ?? null);

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 4f1f7aa. The live test now also asserts Sec-WebSocket-Version: 13. A running Reverb server renders handshake errors through the application's exception handler, while the unit tests cover the WebSocket server's own handler, so this checks the header on the production path.

I didn't add a separate Upgrade assertion. The server sets Upgrade, Connection and Sec-WebSocket-Version together on the same exception, ServerHandshakeTest already asserts all three, and the version header check shows those headers reach the client through the application's handler.

$this->assertSame('13', $client->getHeaders()['sec-websocket-version'] ?? null);

$client->close();
}

// ── Channel subscriptions ──────────────────────────────────────────

public function testCanSubscribeToAPublicChannel(): void
Expand Down
Loading
Loading