Skip to content
Merged
380 changes: 348 additions & 32 deletions docs/cli.md

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion docs/message.md
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,7 @@ $message->header(EventIdHeader::class)->eventId; // 'a4a4a4a4-4a4a-...'
```
:::warning
The `PlayheadHeader` is only added if the stream is playhead based.
Streams that are written without a playhead, for example custom streams,
Streams that are written without a playhead, for example [custom streams](store.md#custom-streams),
do not have this header. Use `hasHeader` before you access it.
:::

Expand Down
85 changes: 85 additions & 0 deletions docs/store.md
Original file line number Diff line number Diff line change
Expand Up @@ -494,6 +494,91 @@ The store locks the table for writing during each save by default.
Use the transactional method if you want to call multiple save methods in one transaction.
:::

#### Custom streams

The `StreamDoctrineDbalStore` is not limited to aggregates.
You can save any event into your own stream, for example integration events
that are derived from your domain events and consumed by other bounded contexts.
For this, the message needs a `StreamNameHeader` with the name of the stream.

```php
use Patchlevel\EventSourcing\Attribute\Event;

#[Event('hotel.integration.guest_checked_in')]
final class GuestCheckedInIntegrationEvent
{
public function __construct(
public readonly string $hotelId,
public readonly string $guestName,
) {
}
}
```
```php
use Patchlevel\EventSourcing\Attribute\Processor;
use Patchlevel\EventSourcing\Attribute\Subscribe;
use Patchlevel\EventSourcing\Message\Message;
use Patchlevel\EventSourcing\Store\Header\StreamNameHeader;
use Patchlevel\EventSourcing\Store\Store;

#[Processor('hotel_integration_publisher')]
final class HotelIntegrationPublisher
{
public function __construct(
private readonly Store $store,
) {
}

#[Subscribe(GuestIsCheckedIn::class)]
public function onGuestIsCheckedIn(GuestIsCheckedIn $event): void
{
$this->store->save(
Message::create(new GuestCheckedInIntegrationEvent(
$event->hotelId->toString(),
$event->guestName,
))->withHeader(new StreamNameHeader('integration-hotel')),
);
}
}
```
These messages are part of the store like any other message.
Every [subscriber](subscription.md) can subscribe to them,
and you can load them with the `StreamCriterion`.

```php
use Patchlevel\EventSourcing\Store\Criteria\Criteria;
use Patchlevel\EventSourcing\Store\Criteria\StreamCriterion;
use Patchlevel\EventSourcing\Store\Store;

/** @var Store $store */
$stream = $store->load(
new Criteria(new StreamCriterion('integration-hotel')),
);
```
Only the `StreamNameHeader` is required, the other headers are optional:

| Header | If missing |
|--------------------|-------------------------------------------------|
| `StreamNameHeader` | A `MissingDataForStorage` exception is thrown |
| `PlayheadHeader` | The playhead is saved as `null` |
| `EventIdHeader` | A new UUIDv7 is generated |
| `RecordedOnHeader` | The current time of the clock is used |

:::warning
Choose a stream name that cannot collide with your aggregate streams, which are named
like `profile-e3e3e3e3-3e3e-3e3e-3e3e-3e3e3e3e3e3e` by default. A prefix like `integration-` helps.
:::

:::note
Without a `PlayheadHeader` there is no optimistic locking for this stream,
because the unique index on `stream` and `playhead` does not apply to `null` values.
:::

:::note
Custom streams are only supported by the `StreamDoctrineDbalStore`.
The event class still needs the `#[Event]` attribute and must be known by the event serializer.
:::

### Update

It is not possible to update events.
Expand Down
9 changes: 7 additions & 2 deletions src/Console/Command/SubscriptionStatusCommand.php
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@

use Patchlevel\EventSourcing\Console\InputHelper;
use Patchlevel\EventSourcing\Console\OutputStyle;
use Patchlevel\EventSourcing\Subscription\Engine\SubscriptionEngineCriteria;
use Patchlevel\EventSourcing\Subscription\Store\SubscriptionNotFound;
use Patchlevel\EventSourcing\Subscription\Subscription;
use Patchlevel\EventSourcing\Subscription\SubscriptionError;
Expand Down Expand Up @@ -41,9 +42,10 @@ protected function execute(InputInterface $input, OutputInterface $output): int
$io = new OutputStyle($input, $output);

$id = InputHelper::nullableString($input->getArgument('id'));
$subscriptions = $this->engine->subscriptions();

if ($id === null) {
$subscriptions = $this->engine->subscriptions($this->subscriptionEngineCriteria($input));

$io->table(
[
'id',
Expand All @@ -69,7 +71,10 @@ protected function execute(InputInterface $input, OutputInterface $output): int
return 0;
}

$subscription = $this->findSubscription($subscriptions, $id);
$subscription = $this->findSubscription(
$this->engine->subscriptions(new SubscriptionEngineCriteria([$id])),
$id,
);

$io->horizontalTable(
[
Expand Down
2 changes: 1 addition & 1 deletion src/Subscription/Engine/GapResolverStoreMessageLoader.php
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ private function generator(int $startIndex, int $retry = 0): Generator
}

$retry = 0;
$currentIndex++;
$currentIndex = $stream->index() + 1;

yield $stream->index() => $message;
}
Expand Down
Loading
Loading