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);