Refactor realtime message handling to send subscriber keys and add comprehensive tests for subscription message upsert behavior

This commit is contained in:
ArnabChatterjee20k
2026-04-06 17:10:57 +05:30
parent 97d46c6273
commit 6bc9adece8
2 changed files with 148 additions and 1 deletions
+1 -1
View File
@@ -441,7 +441,7 @@ $server->onWorkerStart(function (int $workerId) use ($server, $register, $stats,
]
];
$server->send($realtime->getSubscribers($event), json_encode([
$server->send(array_keys($realtime->getSubscribers($event)), json_encode([
'type' => 'event',
'data' => $event['data']
]));
@@ -146,6 +146,153 @@ class RealtimeCustomClientQueryTestWithMessage extends Scope
);
}
/**
* @param array<int, array<string, mixed>> $payloadEntries
* @return array<string, mixed>
*/
private function sendSubscribeMessage(WebSocketClient $client, array $payloadEntries): array
{
$client->send(\json_encode([
'type' => 'subscribe',
'data' => $payloadEntries,
]));
$response = \json_decode($client->receive(), true);
$this->assertEquals('response', $response['type'] ?? null);
$this->assertEquals('subscribe', $response['data']['to'] ?? null);
$this->assertTrue($response['data']['success'] ?? false);
$this->assertArrayHasKey('subscriptions', $response['data']);
$this->assertIsArray($response['data']['subscriptions']);
return $response;
}
/**
* subscriptionId: update with id from connected, create by omitting id, explicit new id,
* duplicate id in one bulk (last wins), mixed bulk, idempotent repeat, empty queries → select-all.
*/
public function testSubscribeMessageUpsertCreateAndEdgeCases(): void
{
$user = $this->getUser();
$session = $user['session'] ?? '';
$projectId = $this->getProject()['$id'];
$headers = [
'origin' => 'http://localhost',
'cookie' => 'a_session_' . $projectId . '=' . $session,
];
$queryString = \http_build_query([
'project' => $projectId,
'channels' => ['documents'],
]);
$client = new WebSocketClient(
'ws://appwrite.test/v1/realtime?' . $queryString,
[
'headers' => $headers,
'timeout' => 30,
]
);
$connected = \json_decode($client->receive(), true);
$this->assertEquals('connected', $connected['type'] ?? null);
$mapping = $connected['data']['subscriptions'] ?? [];
$this->assertNotEmpty($mapping);
$initialSubscriptionId = $mapping[\array_key_first($mapping)];
$q1 = [Query::equal('status', ['q1'])->toString()];
$r1 = $this->sendSubscribeMessage($client, [[
'subscriptionId' => $initialSubscriptionId,
'channels' => ['documents'],
'queries' => $q1,
]]);
$this->assertCount(1, $r1['data']['subscriptions']);
$this->assertSame($initialSubscriptionId, $r1['data']['subscriptions'][0]['subscriptionId']);
$this->assertSame($q1, $r1['data']['subscriptions'][0]['queries']);
$q2 = [Query::equal('status', ['q2'])->toString()];
$r2 = $this->sendSubscribeMessage($client, [[
'subscriptionId' => $initialSubscriptionId,
'channels' => ['documents'],
'queries' => $q2,
]]);
$this->assertSame($initialSubscriptionId, $r2['data']['subscriptions'][0]['subscriptionId']);
$this->assertSame($q2, $r2['data']['subscriptions'][0]['queries']);
$rOmit = $this->sendSubscribeMessage($client, [[
'channels' => ['documents'],
'queries' => [Query::equal('status', ['omitted-slot'])->toString()],
]]);
$mintedId = $rOmit['data']['subscriptions'][0]['subscriptionId'];
$this->assertNotSame($initialSubscriptionId, $mintedId);
$this->assertNotEmpty($mintedId);
$explicitNewId = ID::unique();
$qExplicit = [Query::equal('status', ['explicit'])->toString()];
$rExplicit = $this->sendSubscribeMessage($client, [[
'subscriptionId' => $explicitNewId,
'channels' => ['documents'],
'queries' => $qExplicit,
]]);
$this->assertSame($explicitNewId, $rExplicit['data']['subscriptions'][0]['subscriptionId']);
$this->assertSame($qExplicit, $rExplicit['data']['subscriptions'][0]['queries']);
$qFirst = [Query::equal('status', ['dup-a'])->toString()];
$qSecond = [Query::equal('status', ['dup-b'])->toString()];
$rDup = $this->sendSubscribeMessage($client, [
[
'subscriptionId' => $initialSubscriptionId,
'channels' => ['documents'],
'queries' => $qFirst,
],
[
'subscriptionId' => $initialSubscriptionId,
'channels' => ['documents'],
'queries' => $qSecond,
],
]);
$this->assertCount(2, $rDup['data']['subscriptions']);
$this->assertSame($initialSubscriptionId, $rDup['data']['subscriptions'][0]['subscriptionId']);
$this->assertSame($initialSubscriptionId, $rDup['data']['subscriptions'][1]['subscriptionId']);
$this->assertSame($qSecond, $rDup['data']['subscriptions'][1]['queries']);
$rMixed = $this->sendSubscribeMessage($client, [
[
'subscriptionId' => $initialSubscriptionId,
'channels' => ['documents'],
'queries' => [Query::equal('status', ['mixed-update'])->toString()],
],
[
'channels' => ['documents'],
'queries' => [Query::equal('status', ['mixed-new'])->toString()],
],
]);
$this->assertCount(2, $rMixed['data']['subscriptions']);
$this->assertSame($initialSubscriptionId, $rMixed['data']['subscriptions'][0]['subscriptionId']);
$mixedSecondId = $rMixed['data']['subscriptions'][1]['subscriptionId'];
$this->assertNotSame($initialSubscriptionId, $mixedSecondId);
$this->assertNotEmpty($mixedSecondId);
$rSame = $this->sendSubscribeMessage($client, [[
'subscriptionId' => $initialSubscriptionId,
'channels' => ['documents'],
'queries' => [Query::equal('status', ['idempotent'])->toString()],
]]);
$rSameAgain = $this->sendSubscribeMessage($client, [[
'subscriptionId' => $initialSubscriptionId,
'channels' => ['documents'],
'queries' => [Query::equal('status', ['idempotent'])->toString()],
]]);
$this->assertSame($rSame['data']['subscriptions'][0]['queries'], $rSameAgain['data']['subscriptions'][0]['queries']);
$rEmpty = $this->sendSubscribeMessage($client, [[
'subscriptionId' => $initialSubscriptionId,
'channels' => ['documents'],
'queries' => [],
]]);
$this->assertCount(1, $rEmpty['data']['subscriptions']);
$this->assertSame($initialSubscriptionId, $rEmpty['data']['subscriptions'][0]['subscriptionId']);
$client->close();
}
public function testInvalidQueryShouldNotSubscribe(): void
{
$user = $this->getUser();