66
77use Closure ;
88use Patchlevel \EventSourcing \Aggregate \AggregateHeader ;
9+ use Patchlevel \EventSourcing \Clock \SystemClock ;
910use Patchlevel \EventSourcing \Message \HeaderNotFound ;
1011use Patchlevel \EventSourcing \Message \Message ;
12+ use Patchlevel \EventSourcing \Metadata \Event \EventRegistry ;
1113use Patchlevel \EventSourcing \Store \Criteria \AggregateIdCriterion ;
1214use Patchlevel \EventSourcing \Store \Criteria \AggregateNameCriterion ;
1315use Patchlevel \EventSourcing \Store \Criteria \ArchivedCriterion ;
1416use Patchlevel \EventSourcing \Store \Criteria \Criteria ;
17+ use Patchlevel \EventSourcing \Store \Criteria \EventsCriterion ;
1518use Patchlevel \EventSourcing \Store \Criteria \FromIndexCriterion ;
1619use Patchlevel \EventSourcing \Store \Criteria \FromPlayheadCriterion ;
1720use Patchlevel \EventSourcing \Store \Criteria \StreamCriterion ;
21+ use Patchlevel \EventSourcing \Store \Criteria \ToIndexCriterion ;
22+ use Patchlevel \EventSourcing \Store \Header \EventIdHeader ;
23+ use Patchlevel \EventSourcing \Store \Header \IndexHeader ;
1824use Patchlevel \EventSourcing \Store \Header \PlayheadHeader ;
25+ use Patchlevel \EventSourcing \Store \Header \RecordedOnHeader ;
1926use Patchlevel \EventSourcing \Store \Header \StreamNameHeader ;
27+ use Psr \Clock \ClockInterface ;
28+ use Ramsey \Uuid \Uuid ;
29+ use Throwable ;
2030
2131use function array_filter ;
2232use function array_map ;
23- use function array_push ;
2433use function array_reverse ;
2534use function array_slice ;
2635use function array_unique ;
3544
3645final class InMemoryStore implements StreamStore
3746{
38- /** @param array<positive-int|0, Message> $messages */
47+ /** @var array<0|positive-int, Message> */
48+ private array $ messages = [];
49+
50+ /** @param list<Message> $messages */
3951 public function __construct (
40- private array $ messages = [],
52+ array $ messages = [],
53+ private readonly EventRegistry |null $ eventRegistry = null ,
54+ private readonly ClockInterface $ clock = new SystemClock (),
4155 ) {
56+ $ this ->save (...$ messages );
4257 }
4358
4459 public function load (
@@ -71,7 +86,25 @@ public function count(Criteria|null $criteria = null): int
7186
7287 public function save (Message ...$ messages ): void
7388 {
74- array_push ($ this ->messages , ...$ messages );
89+ $ count = count ($ this ->messages );
90+
91+ foreach ($ messages as $ message ) {
92+ $ count ++;
93+
94+ if (!$ message ->hasHeader (IndexHeader::class)) {
95+ $ message = $ message ->withHeader (new IndexHeader ($ count ));
96+ }
97+
98+ if (!$ message ->hasHeader (EventIdHeader::class)) {
99+ $ message = $ message ->withHeader (new EventIdHeader (Uuid::uuid7 ()->toString ()));
100+ }
101+
102+ if (!$ message ->hasHeader (RecordedOnHeader::class)) {
103+ $ message = $ message ->withHeader (new RecordedOnHeader ($ this ->clock ->now ()));
104+ }
105+
106+ $ this ->messages [] = $ message ;
107+ }
75108 }
76109
77110 /**
@@ -81,7 +114,14 @@ public function save(Message ...$messages): void
81114 */
82115 public function transactional (Closure $ function ): void
83116 {
84- $ function ();
117+ $ messages = $ this ->messages ;
118+ try {
119+ $ function ();
120+ } catch (Throwable $ e ) {
121+ $ this ->messages = $ messages ;
122+
123+ throw $ e ;
124+ }
85125 }
86126
87127 /** @return list<string> */
@@ -134,9 +174,11 @@ private function filter(Criteria|null $criteria): array
134174 return $ this ->messages ;
135175 }
136176
177+ $ eventRegistry = $ this ->eventRegistry ;
178+
137179 return array_filter (
138180 $ this ->messages ,
139- static function (Message $ message, int $ index ) use ($ criteria ): bool {
181+ static function (Message $ message ) use ($ criteria, $ eventRegistry ): bool {
140182 foreach ($ criteria ->all () as $ criterion ) {
141183 switch ($ criterion ::class) {
142184 case AggregateIdCriterion::class:
@@ -222,7 +264,35 @@ static function (Message $message, int $index) use ($criteria): bool {
222264
223265 break ;
224266 case FromIndexCriterion::class:
225- if ($ index < $ criterion ->fromIndex ) {
267+ try {
268+ $ index = $ message ->header (IndexHeader::class)->index ;
269+ } catch (HeaderNotFound ) {
270+ return false ;
271+ }
272+
273+ if ($ index <= $ criterion ->fromIndex ) {
274+ return false ;
275+ }
276+
277+ break ;
278+ case ToIndexCriterion::class:
279+ try {
280+ $ index = $ message ->header (IndexHeader::class)->index ;
281+ } catch (HeaderNotFound ) {
282+ return false ;
283+ }
284+
285+ if ($ index >= $ criterion ->toIndex ) {
286+ return false ;
287+ }
288+
289+ break ;
290+ case EventsCriterion::class:
291+ if ($ eventRegistry === null ) {
292+ throw new MissingEventRegistry ($ criterion ::class);
293+ }
294+
295+ if (!in_array ($ eventRegistry ->eventName ($ message ->event ()::class), $ criterion ->events )) {
226296 return false ;
227297 }
228298
0 commit comments