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
4 changes: 2 additions & 2 deletions composer.json
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,6 @@
"symfony/property-access": "^6.0.0|^7.0.0|^8.0.0",
"symfony/property-info": "^6.0.0|^7.0.0|^8.0.0",
"symfony/runtime": "^6.0.0|^7.0.0|^8.0.0",
"symfony/serializer": "^6.0.0|^7.0.0|^8.0.0",
"symfony/uid": "^6.0.0|^7.0.0|^8.0.0",
"symfony/yaml": "^6.0.0|^7.0.0|^8.0.0"
},
Expand All @@ -45,7 +44,8 @@
"ringcentral/psr7": "^1.3",
"squizlabs/php_codesniffer": "^3.6",
"symfony/http-client": "^6.0.0|^7.0.0|^8.0.0",
"symfony/process": "^6.0.0|^7.0.0|^8.0.0"
"symfony/process": "^6.0.0|^7.0.0|^8.0.0",
"symfony/serializer": "^6.0.0|^7.0.0|^8.0.0"
},
"config": {
"optimize-autoloader": true,
Expand Down
58 changes: 47 additions & 11 deletions src/Hub/Transport/Redis/RedisSerializer.php
Original file line number Diff line number Diff line change
Expand Up @@ -4,26 +4,62 @@

namespace Freddie\Hub\Transport\Redis;

use Freddie\Message\Message;
use Freddie\Message\Update;
use Symfony\Component\Serializer\Encoder\JsonEncoder;
use Symfony\Component\Serializer\Normalizer\ObjectNormalizer;
use Symfony\Component\Serializer\Serializer;
use Symfony\Component\Serializer\SerializerInterface;
use UnexpectedValueException;

use function json_decode;
use function json_encode;

use const JSON_THROW_ON_ERROR;

/**
* Hand-rolled (de)serializer for the Redis transport.
*
* The wire shape is intentionally identical to the previous Symfony
* ObjectNormalizer output so mixed hub versions stay interoperable during a
* rolling deploy; only the reflection-heavy normalizer is dropped, since this
* runs on every message on every worker.
*/
final readonly class RedisSerializer
{
public function __construct(
private SerializerInterface $serializer = new Serializer([new ObjectNormalizer()], [new JsonEncoder()]),
) {
}

public function serialize(Update $update): string
{
return $this->serializer->serialize($update, 'json');
$message = $update->message;

return json_encode([
'topics' => $update->topics,
'message' => [
'id' => $message->id,
'data' => $message->data,
'private' => $message->private,
'event' => $message->event,
'retry' => $message->retry,
],
], JSON_THROW_ON_ERROR);
}

public function deserialize(string $payload): Update
{
return $this->serializer->deserialize($payload, Update::class, 'json');
/** @var array{topics?: string[], message?: array{id?: string|null, data?: string|null, private?: bool, event?: string|null, retry?: int|null}} $data */
$data = json_decode($payload, true, flags: JSON_THROW_ON_ERROR);

// Required keys throw (as the previous ObjectNormalizer did); optional
// message fields fall back to their defaults (a missing id is regenerated).
$topics = $data['topics']
?? throw new UnexpectedValueException('Malformed Mercure update: missing "topics".');
$message = $data['message']
?? throw new UnexpectedValueException('Malformed Mercure update: missing "message".');

return new Update(
$topics,
new Message(
id: $message['id'] ?? null,
data: $message['data'] ?? null,
private: $message['private'] ?? false,
event: $message['event'] ?? null,
retry: $message['retry'] ?? null,
),
);
}
}
14 changes: 14 additions & 0 deletions tests/Pest.php
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,9 @@
use Psr\Http\Message\ResponseInterface;
use Psr\Http\Message\ServerRequestInterface;
use ReflectionClass;
use Symfony\Component\Serializer\Encoder\JsonEncoder;
use Symfony\Component\Serializer\Normalizer\ObjectNormalizer;
use Symfony\Component\Serializer\Serializer;

function handle(App $app, ServerRequestInterface $request): ResponseInterface
{
Expand Down Expand Up @@ -46,6 +49,17 @@ function jwt_config(): Configuration
);
}

/**
* The Symfony-based serializer previously used for the Redis transport,
* kept as a dev dependency to assert wire-format compatibility.
*/
function legacy_redis_serializer(): Serializer
{
static $serializer;

return $serializer ??= new Serializer([new ObjectNormalizer()], [new JsonEncoder()]);
}

function create_jwt(array $claims): string
{
$builder = jwt_config()->builder();
Expand Down
102 changes: 102 additions & 0 deletions tests/Unit/Hub/Transport/Redis/RedisSerializerTest.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,102 @@
<?php

declare(strict_types=1);

namespace Freddie\Tests\Unit\Hub\Transport\Redis;

use ErrorException;
use Freddie\Hub\Transport\Redis\RedisSerializer;
use Freddie\Message\Message;
use Freddie\Message\Update;
use UnexpectedValueException;

use function Freddie\Tests\legacy_redis_serializer;

it('round-trips an update with all message fields', function () {
$serializer = new RedisSerializer();
$update = new Update(
['/foo', '/bar'],
new Message(id: '01ARZ3', data: "line1\nline2", private: true, event: 'ping', retry: 3000),
);

$result = $serializer->deserialize($serializer->serialize($update));

expect($result->topics)->toBe(['/foo', '/bar']);
expect($result->message->id)->toBe('01ARZ3');
expect($result->message->data)->toBe("line1\nline2");
expect($result->message->private)->toBeTrue();
expect($result->message->event)->toBe('ping');
expect($result->message->retry)->toBe(3000);
});

it('round-trips an update with null message fields and preserves the id', function () {
$serializer = new RedisSerializer();
$update = new Update('/foo', new Message(id: '01BX5Z'));

$result = $serializer->deserialize($serializer->serialize($update));

expect($result->topics)->toBe(['/foo']);
expect($result->message->id)->toBe('01BX5Z');
expect($result->message->data)->toBeNull();
expect($result->message->private)->toBeFalse();
expect($result->message->event)->toBeNull();
expect($result->message->retry)->toBeNull();
});

// Guarantees mixed old/new hubs interoperate during a rolling deploy (Redis
// pub/sub has no format versioning), by cross-checking against the previous
// Symfony ObjectNormalizer wire format.
it('stays wire-compatible with the Symfony ObjectNormalizer format', function () {
$serializer = new RedisSerializer();
$objectNormalizer = legacy_redis_serializer();
$update = new Update(['/foo'], new Message(id: '01CX', data: 'hi', private: true, event: 'e', retry: 1));

// new serialize -> old deserialize
/** @var Update $fromNew */
$fromNew = $objectNormalizer->deserialize($serializer->serialize($update), Update::class, 'json');
expect($fromNew->topics)->toBe(['/foo']);
expect($fromNew->message->id)->toBe('01CX');
expect($fromNew->message->data)->toBe('hi');
expect($fromNew->message->private)->toBeTrue();
expect($fromNew->message->event)->toBe('e');
expect($fromNew->message->retry)->toBe(1);

// old serialize -> new deserialize
$fromOld = $serializer->deserialize($objectNormalizer->serialize($update, 'json'));
expect($fromOld->topics)->toBe(['/foo']);
expect($fromOld->message->id)->toBe('01CX');
expect($fromOld->message->data)->toBe('hi');
expect($fromOld->message->private)->toBeTrue();
expect($fromOld->message->event)->toBe('e');
expect($fromOld->message->retry)->toBe(1);
});

// Matches the previous ObjectNormalizer contract: optional message fields fall
// back to their defaults (a missing id is regenerated) WITHOUT emitting a warning.
it('tolerates missing optional message fields without emitting a warning', function () {
set_error_handler(
static fn (int $errno, string $errstr) => throw new ErrorException($errstr),
E_WARNING | E_NOTICE,
);

try {
$update = (new RedisSerializer())->deserialize('{"topics":["/foo"],"message":{"data":"hi"}}');
} finally {
restore_error_handler();
}

expect(strlen($update->message->id))->toBe(26); // a fresh Ulid was generated
expect($update->message->data)->toBe('hi');
expect($update->message->private)->toBeFalse();
expect($update->message->event)->toBeNull();
expect($update->message->retry)->toBeNull();
});

// Matches the previous ObjectNormalizer contract: required keys throw when absent.
it('throws when the payload is missing required topics', function () {
(new RedisSerializer())->deserialize('{"message":{"id":"01CX","data":"hi"}}');
})->throws(UnexpectedValueException::class);

it('throws when the payload is missing the message', function () {
(new RedisSerializer())->deserialize('{"topics":["/foo"]}');
})->throws(UnexpectedValueException::class);
Loading