PHP code example of thesis / nats

1. Go to this page and download the library: Download thesis/nats library. Choose the download type require.

2. Extract the ZIP file and open the index.php.

3. Add this code to the index.php.
    
        
<?php
require_once('vendor/autoload.php');

/* Start to develop here. Best regards https://php-download.com/ */

    

thesis / nats example snippets




declare(strict_types=1);

mp\delay;
use function Amp\trapSignal;

$nc = new Nats\Client(Nats\Config::default());

$nc->subscribe('foo.*', static function (Nats\Delivery $delivery): void {
    dump("Received message {$delivery->message->payload} for consumer#1");
});

$nc->subscribe('foo.>', static function (Nats\Delivery $delivery): void {
    dump("Received message {$delivery->message->payload} for consumer#2");
});

$subscription = $nc->subscribe('foo.bar', static function (Nats\Delivery $delivery): void {
    dump("Received message {$delivery->message->payload} for consumer#3");
});

$nc->publish('foo.bar', new Nats\Message('Hello World!')); // visible for all consumers
$nc->publish('foo.baz', new Nats\Message('Hello World!')); // visible only for 1-2 consumers
$nc->publish('foo.bar.baz', new Nats\Message('Hello World!')); // visible only for 2 consumer

$subscription->stop();
$nc->publish('foo.bar', new Nats\Message('Hello World!')); // visible for 1-2 consumers

trapSignal([\SIGTERM, \SIGINT]);

$nc->drain();



declare(strict_types=1);

mp\trapSignal;

$nc = new Nats\Client(Nats\Config::default());

$nc->subscribe(
    subject: 'foo.>',
    handler: static function (Nats\Delivery $delivery): void {
        dump("Received message {$delivery->message->payload} for consumer#1");
    },
    queueGroup: 'test',
);

$nc->subscribe(
    subject: 'foo.>',
    handler: static function (Nats\Delivery $delivery): void {
        dump("Received message {$delivery->message->payload} for consumer#2");
    },
    queueGroup: 'test',
);

$nc->subscribe(
    subject: 'foo.>',
    handler: static function (Nats\Delivery $delivery): void {
        dump("Received message {$delivery->message->payload} for consumer#3");
    },
    queueGroup: 'test',
);

$nc->publish('foo.bar', new Nats\Message('x'));
$nc->publish('foo.baz', new Nats\Message('y'));
$nc->publish('foo.bar.baz', new Nats\Message('z'));

trapSignal([\SIGTERM, \SIGINT]);

$nc->drain();



declare(strict_types=1);

s\Client(Nats\Config::default());

$nc->subscribe('foo.>', static function (Nats\Delivery $delivery): void {
    dump("Received request {$delivery->message->payload}");
    $delivery->reply(new Nats\Message(strrev($delivery->message->payload ?? '')));
});

$response = $nc->request('foo.bar', new Nats\Message('Hello World!'));
dump("Received response {$response->message->payload}");

$nc->drain();



declare(strict_types=1);

use Thesis\Nats;
use Thesis\Nats\JetStream;

$consumer = $stream->createOrUpdateConsumer(...);

$batch = $consumer
    ->pulling()
    ->fetch(JetStream\FetchConfig::batch(10));

foreach ($batch as $delivery) {
    // handle message
    $delivery->ack();
}



declare(strict_types=1);

use Thesis\Nats;
use Thesis\Nats\JetStream;

$consumer = $stream->createOrUpdateConsumer(...);

$batch = $consumer
    ->pulling()
    ->fetch(JetStream\FetchConfig::bytes(1024));

foreach ($batch as $delivery) {
    // handle message
    $delivery->ack();
}



declare(strict_types=1);

use Thesis\Nats;
use Thesis\Nats\JetStream;

$consumer = $stream->createOrUpdateConsumer(...);

$batch = $consumer
    ->pulling()
    ->fetch(JetStream\FetchConfig::immediate());

foreach ($batch as $delivery) {
    // handle message
    $delivery->ack();
}



declare(strict_types=1);

s\JetStream\Api;
use Thesis\Time\TimeSpan;

$consumer = $stream->createOrUpdateConsumer(new Api\ConsumerConfig(
    durableName: 'EventPushConsumer',
    deliverSubject: 'push-consumer-delivery',
    ackPolicy: Api\AckPolicy::Explicit,
    idleHeartbeat: TimeSpan::fromSeconds(5),
    maxAckPending: 1,
));

$subscription = $consumer->push(
    static function (Nats\JetStream\Delivery $delivery): void {
        dump($delivery->message->payload);
        $delivery->ack();
    },
);



declare(strict_types=1);

s\JetStream\Api;
use Thesis\Time\TimeSpan;

$consumer = $stream->createOrUpdateConsumer(new Api\ConsumerConfig(
    durableName: 'EventPushConsumer',
    deliverSubject: 'push-consumer-delivery',
    deliverGroup: 'testing',
    ackPolicy: Api\AckPolicy::Explicit,
    idleHeartbeat: TimeSpan::fromSeconds(5),
    maxAckPending: 1,
));



declare(strict_types=1);

s\JetStream\Api;
use Thesis\Time\TimeSpan;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$consumer = $stream->createOrUpdateConsumer(new Api\ConsumerConfig(
    durableName: 'EventPullConsumer',
    ackPolicy: Api\AckPolicy::Explicit,
));

$subscription = $consumer->pull(
    static function (Nats\JetStream\Delivery $delivery): void {
        dump($delivery->message->payload);
        $delivery->ack();
    },
);



declare(strict_types=1);

s\JetStream\Api;
use Thesis\Time\TimeSpan;
use function Amp\trapSignal;

$consumer = $stream->createOrUpdateConsumer(new Api\ConsumerConfig(
    durableName: 'EventPullConsumer',
    ackPolicy: Api\AckPolicy::Explicit,
));

$consumer = $consumer->pulling();

$consumer->consume(
    static function (Nats\JetStream\Delivery $delivery): void {
        dump("Consumer#1: {$delivery->message->payload}");
        $delivery->ack();
    },
);

$consumer->consume(
    static function (Nats\JetStream\Delivery $delivery): void {
        dump("Consumer#2: {$delivery->message->payload}");
        $delivery->ack();
    },
);

trapSignal([\SIGINT, \SIGTERM]);

$consumer->drain();



declare(strict_types=1);

use Thesis\Nats\JetStream;;

$consumer->consume(
    static function (Nats\JetStream\Delivery $delivery): void {},
    new JetStream\PullConsumeConfig(
        maxMessages: 1_000,
    ),
);



declare(strict_types=1);

s\JetStream\Api\StreamConfig;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$js->deleteStream('EventStream');

$stream = $js->createStream(new StreamConfig(
    name: 'EventStream',
    description: 'Application events',
    subjects: ['events.*'],
));

for ($i = 0; $i < 5; ++$i) {
    $js->publish(
        subject: 'events.payment_rejected',
        message: new Nats\Message(
            payload: "Message#{$i}",
            headers: (new Nats\Headers())
                ->with(Nats\Header\MsgId::header(), "id:{$i}"),
        ),
    );
}

dump($stream->getLastMessageForSubject('events.payment_rejected')?->payload);

$nc->drain();



declare(strict_types=1);

s\JetStream\KeyValue\BucketConfig;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$kv = $js->createOrUpdateKeyValue(new BucketConfig(
    bucket: 'configs',
));

$kv->put('app.env', 'prod');
$kv->put('database.dsn', 'mysql:host=127.0.0.1;port=3306');

dump(
    $kv->get('app.env')?->value,
    $kv->get('database.dsn')?->value,
);

$nc->drain();



declare(strict_types=1);

s\JetStream\KeyValue\BucketConfig;
use function Amp\trapSignal;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$js->deleteKeyValue('configs');

$kv = $js->createOrUpdateKeyValue(new BucketConfig(
    bucket: 'configs',
));

$subscription = $kv->watch(static function (Nats\JetStream\KeyValue\Entry $entry): void {
    dump("Config key {$entry->key} value changed to {$entry->value}");
});

$kv->put('app.env', 'prod');
$kv->put('database.dsn', 'mysql:host=127.0.0.1;port=3306');

trapSignal([\SIGTERM, \SIGINT]);

$subscription->stop();

$nc->drain();



declare(strict_types=1);

s\JetStream\ObjectStore\ObjectMeta;
use Thesis\Nats\JetStream\ObjectStore\ResourceReader;
use Thesis\Nats\JetStream\ObjectStore\StoreConfig;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$js->deleteObjectStore('code');

$store = $js->createOrUpdateObjectStore(new StoreConfig(
    store: 'code',
));

$handle = fopen(__DIR__.'/app.php', 'r') ?? throw new \RuntimeException('Failed to open file.');

$store->put(new ObjectMeta(name: 'app.php'), new ResourceReader($handle));

fclose($handle);

$store->put(new ObjectMeta('config.php'), ' return [];');

dump(
    (string) $store->get('app.php'),
    (string) $store->get('config.php'),
);

$nc->drain();



declare(strict_types=1);

s\JetStream\ObjectStore\ObjectInfo;
use Thesis\Nats\JetStream\ObjectStore\ObjectMeta;
use Thesis\Nats\JetStream\ObjectStore\StoreConfig;
use function Amp\delay;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$js->deleteObjectStore('code');

$store = $js->createOrUpdateObjectStore(new StoreConfig(
    store: 'code',
));

$subscription = $store->watch(static function (ObjectInfo $info): void {
    dump("New object {$info->name} in the bucket {$info->bucket} at size {$info->size} bytes");
});

$store->put(new ObjectMeta('config.php'), ' return [];');
$store->put(new ObjectMeta('snippet.php'), ' echo 1 + 1;');

delay(0.5);

$subscription->stop();

$nc->drain();



declare(strict_types=1);

s\JetStream\Counter\CounterConfig;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$counter = $js->createOrUpdateCounter(new CounterConfig(
    name: 'atomics',
));

dump($counter->add('x', 1)); // 1
dump($counter->add('x', 2)); // 3



declare(strict_types=1);

s\JetStream\Counter\CounterConfig;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$counter = $js->createOrUpdateCounter(new CounterConfig(
    name: 'atomics',
));

dump($counter->add('x', 1)); // 1
dump($counter->get('x')?->value); // 1



declare(strict_types=1);

s\JetStream\Counter\CounterConfig;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$jetstream->deleteCounter('atomics');

$counter = $js->createOrUpdateCounter(new CounterConfig(
    name: 'atomics',
));

$counter->add('x', 1);
$counter->add('y', 1);
$counter->add('z', 1);

foreach ($counter->getMultiple() as $entry) {
    echo "{$entry->subject}: {$entry->value}\n";
}



declare(strict_types=1);

s\Header;
use Thesis\Nats\JetStream\Api\AckPolicy;
use Thesis\Nats\JetStream\Api\ConsumerConfig;
use Thesis\Nats\JetStream\Api\StreamConfig;
use Thesis\Nats\JetStream\Api\DeliverPolicy;
use function Amp\trapSignal;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$stream = $js->createStream(new StreamConfig(
    name: 'RecurrentsStream',
    subjects: [
        'recurrents',
        'scheduler.recurrents.*',
    ],
    allowMsgSchedules: true,
));

$js->publish('scheduler.recurrents.1', new Nats\Message(
    payload: '{"id":1}',
    headers: (new Nats\Headers())
        ->with(Header\Schedule::Header, new \DateTimeImmutable('+5 seconds'))
        ->with(Header\ScheduleTarget::header(), 'recurrents'),
));

$subscription = $stream
    ->createOrUpdateConsumer(new ConsumerConfig(
        durableName: 'RecurrentsConsumer',
        deliverPolicy: DeliverPolicy::New,
        ackPolicy: AckPolicy::None,
        filterSubjects: ['recurrents'],
    ))
    ->pull(static function (Nats\JetStream\Delivery $delivery): void {
        dump([
            $delivery->message->payload,
            $delivery->message->headers?->get(Header\Scheduler::header()),
            $delivery->message->headers?->get(Header\ScheduleNext::header()),
        ]);
    });

trapSignal([\SIGINT, \SIGTERM]);

$subscription->awaitCompletion();



declare(strict_types=1);

s\JetStream\Api\StreamConfig;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$stream = $js->createStream(new StreamConfig(
    name: 'Batches',
    description: 'Batch Stream',
    subjects: ['batch.*'],
    allowAtomicPublish: true,
));

$batch = $js->createPublishBatch();

for ($i = 0; $i < 999; ++$i) {
    $batch->publish('batch.orders', new Nats\Message("Order#{$i}"));
}

$batch->publish('batch.orders', new Nats\Message('Order#1000'), new Nats\PublishBatchOptions(commit: true));



declare(strict_types=1);

s\JetStream\Api\StreamConfig;

$nc = new Nats\Client(Nats\Config::default());
$js = $nc->jetStream();

$stream = $js->createStream(new StreamConfig(
    name: 'Batches',
    description: 'Batch Stream',
    subjects: ['batch.*'],
    allowAtomicPublish: true,
));

$js->publishBatch('batch.orders', [
    new Nats\Message('Order#1'),
    new Nats\Message('Order#2'),
    new Nats\Message('Order#3'),
]);



declare(strict_types=1);

use Thesis\Nats;
use Thesis\Nats\Micro;

eService(new Micro\ServiceConfig('EchoService', '1.0.0'));



declare(strict_types=1);

use Thesis\Nats;
use Thesis\Nats\Micro;

eService(new Micro\ServiceConfig('EchoService', '1.0.0'));

$srv
    ->addEndpoint(new Micro\EndpointConfig('srv.echo'), static function (Micro\Request $request): void {
        $request->respond(new Micro\Response($request->data));
    });

dump($nc->request('srv.echo', new Nats\Message('ping'))->message->payload);



declare(strict_types=1);

use Thesis\Nats;
use Thesis\Nats\Micro;

eService(new Micro\ServiceConfig('EchoService', '1.0.0'));

$srv
    ->addEndpoint(new Micro\EndpointConfig('srv.echo', subject: 'srv.echo.*'), static function (Micro\Request $request): void {
        $request->respond(new Micro\Response($request->subject));
    });

dump($nc->request('srv.echo.x', new Nats\Message('ping'))->message->payload);



declare(strict_types=1);

use Thesis\Nats;
use Thesis\Nats\Micro;

:4222'),
);

$nc
    ->createService(new Micro\ServiceConfig('EchoService', '1.0.0'))
        ->addGroup(new Micro\GroupConfig('srv.api'))
            ->addGroup(new Micro\GroupConfig('v1'))
                ->addEndpoint(new Micro\EndpointConfig('echo'), static function (Micro\Request $request): void {
                    $request->respond(new Micro\Response($request->data));
                });

dump($nc->request('srv.api.v1.echo', new Nats\Message('ping'))->message->payload);