{"version":3,"sources":["/home/runner/work/emmett/emmett/src/packages/emmett/dist/index.cjs","../src/typing/command.ts","../src/typing/event.ts","../src/typing/message.ts","../src/typing/workflow.ts","../src/typing/index.ts","../src/eventStore/afterCommit/afterEventStoreCommitHandler.ts","../src/eventStore/afterCommit/forwardToMessageBus.ts","../src/eventStore/events/index.ts","../src/eventStore/eventStore.ts","../src/eventStore/expectedVersion.ts","../src/eventStore/inMemoryEventStore.ts","../src/database/inMemoryDatabase.ts","../src/utils/collections/duplicates.ts","../src/utils/collections/merge.ts","../src/utils/collections/index.ts","../src/utils/deepEquals.ts","../src/utils/iterators.ts","../src/taskProcessing/taskProcessor.ts","../src/utils/locking/index.ts","../src/utils/numbers/bigint.ts","../src/utils/promises.ts","../src/utils/retry.ts","../src/serialization/json/JSONParser.ts","../src/utils/shutdown/gracefulShutdown.ts","../src/utils/strings/hashText.ts","../src/database/utils.ts","../src/eventStore/projections/inMemory/inMemoryProjection.ts","../src/eventStore/projections/inMemory/inMemoryProjectionSpec.ts","../src/testing/assertions.ts","../src/testing/deciderSpecification.ts","../src/testing/wrapEventStore.ts","../src/eventStore/versioning/downcasting.ts","../src/eventStore/versioning/upcasting.ts","../src/commandHandling/handleCommand.ts","../src/commandHandling/handleCommandWithDecider.ts","../src/messageBus/index.ts","../src/processors/processors.ts","../src/processors/inMemoryProcessors.ts","../src/projections/index.ts"],"names":["message","error","uuid","result","sum","item","operationResult","projection","event","command","inlineProjections"],"mappings":"AAAA;AACE;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACF,wDAA6B;AAC7B;AACA;AC6BO,IAAM,QAAA,EAAU,CAAA,GAClB,IAAA,EAAA,GAWa;AAChB,EAAA,MAAM,CAAC,IAAA,EAAM,IAAA,EAAM,QAAQ,EAAA,EAAI,IAAA;AAE/B,EAAA,OAAO,SAAA,IAAa,KAAA,EAAA,EACf,EAAE,IAAA,EAAM,IAAA,EAAM,QAAA,EAAU,IAAA,EAAM,UAAU,EAAA,EACxC,EAAE,IAAA,EAAM,IAAA,EAAM,IAAA,EAAM,UAAU,CAAA;AACrC,CAAA;AD1CA;AACA;AEiCO,IAAM,MAAA,EAAQ,CAAA,GAChB,IAAA,EAAA,GAOW;AACd,EAAA,MAAM,CAAC,IAAA,EAAM,IAAA,EAAM,QAAQ,EAAA,EAAI,IAAA;AAE/B,EAAA,OAAO,SAAA,IAAa,KAAA,EAAA,EACf,EAAE,IAAA,EAAM,IAAA,EAAM,QAAA,EAAU,IAAA,EAAM,QAAQ,EAAA,EACtC,EAAE,IAAA,EAAM,IAAA,EAAM,IAAA,EAAM,QAAQ,CAAA;AACnC,CAAA;AF1CA;AACA;AGHO,IAAM,QAAA,EAAU,CAAA,GAClB,IAAA,EAAA,GAYa;AAChB,EAAA,MAAM,CAAC,IAAA,EAAM,IAAA,EAAM,IAAA,EAAM,QAAQ,EAAA,EAAI,IAAA;AAErC,EAAA,OAAO,SAAA,IAAa,KAAA,EAAA,EACf,EAAE,IAAA,EAAM,IAAA,EAAM,QAAA,EAAU,KAAK,EAAA,EAC7B,EAAE,IAAA,EAAM,IAAA,EAAM,KAAK,CAAA;AAC1B,CAAA;AHXA;AACA;AIAO,IAAM,MAAA,EAAQ,CACnBA,QAAAA,EAAAA,GAC4B;AAC5B,EAAA,OAAO;AAAA,IACL,MAAA,EAAQ,OAAA;AAAA,IACR,OAAA,EAAAA;AAAA,EACF,CAAA;AACF,CAAA;AAEO,IAAM,KAAA,EAAO,CAClBA,QAAAA,EAAAA,GAC4B;AAC5B,EAAA,OAAO;AAAA,IACL,MAAA,EAAQ,MAAA;AAAA,IACR,OAAA,EAAAA;AAAA,EACF,CAAA;AACF,CAAA;AAEO,IAAM,QAAA,EAAU,CACrBA,QAAAA,EAAAA,GAC4B;AAC5B,EAAA,OAAO;AAAA,IACL,MAAA,EAAQ,SAAA;AAAA,IACR,OAAA,EAAAA;AAAA,EACF,CAAA;AACF,CAAA;AAEO,IAAM,SAAA,EAAW,CACtBA,QAAAA,EACA,IAAA,EAAA,GAC4B;AAC5B,EAAA,OAAO;AAAA,IACL,MAAA,EAAQ,UAAA;AAAA,IACR,OAAA,EAAAA,QAAAA;AAAA,IACA;AAAA,EACF,CAAA;AACF,CAAA;AAEO,IAAM,SAAA,EAAW,CAAA,EAAA,GAEQ;AAC9B,EAAA,OAAO;AAAA,IACL,MAAA,EAAQ;AAAA,EACV,CAAA;AACF,CAAA;AAEO,IAAM,OAAA,EAAS,CACpB,MAAA,EAAA,GAC4B;AAC5B,EAAA,OAAO;AAAA,IACL,MAAA,EAAQ,QAAA;AAAA,IACR;AAAA,EACF,CAAA;AACF,CAAA;AAEO,IAAM,MAAA,EAAQ,CACnB,MAAA,EAAA,GAC4B;AAC5B,EAAA,OAAO;AAAA,IACL,MAAA,EAAQ,OAAA;AAAA,IACR;AAAA,EACF,CAAA;AACF,CAAA;AAEO,IAAM,OAAA,EAAS,CAAA,EAAA,GAEU;AAC9B,EAAA,OAAO,EAAE,MAAA,EAAQ,SAAS,CAAA;AAC5B,CAAA;AJtBA;AACA;AKlEO,IAAM,aAAA,EAAe,KAAA;AAErB,IAAM,UAAA,EAAY,QAAA;AAClB,IAAM,WAAA,EAAa,CAAA,EAAA;AACA;ALmER;AACA;AMlCI;AAUP,EAAA;AAET,EAAA;AACI,IAAA;AACC,IAAA;AACAC,EAAAA;AAEO,IAAA;AACP,IAAA;AACT,EAAA;AACF;ANyBkB;AACA;AOnGL;AASED,EAAAA;AACH,IAAA;AACR,EAAA;AACF;AP6FgB;AACA;AQvGL;AAQA;AAIA;AAWA;AAIO,EAAA;AACjB;AAEU;AAIA;AR8EK;AACA;ASnBL;AAIA;AAGG,EAAA;AACN,IAAA;AACJ,MAAA;AACO,MAAA;AACT,IAAA;AAEO,IAAA;AACT,EAAA;AACF;ATekB;AACA;AUtHL;AACA;AAEA;AAGA;AAKP,EAAA;AAEY,EAAA;AAEA,EAAA;AAET,EAAA;AACT;AAEa;AAOE,EAAA;AAER,EAAA;AACO,IAAA;AACd;AAEa;AAIT,EAAA;AAGM,IAAA;AAGC,IAAA;AACT,EAAA;AACF;AAEa;AAKTC,EAAAA;AACA,EAAA;AACF;AVyFgB;AACA;AW/JHC;AXiKG;AACA;AYlKH;AZoKG;AACA;AarKL;AAII,EAAA;AACT,EAAA;AAEC,EAAA;AACT;AAEa;AAIC,EAAA;AAEI,EAAA;AACD,IAAA;AACD,IAAA;AACC,IAAA;AACH,MAAA;AACV,IAAA;AACc,IAAA;AAChB,EAAA;AAEa,EAAA;AAGf;Ab2JkB;AACA;AcvLhB;AAMe,EAAA;AAEA,EAAA;AAGI,IAAA;AAEJ,IAAA;AACJ,IAAA;AAIA,EAAA;AAGO,IAAA;AAEP,IAAA;AACR,EAAA;AAIY,EAAA;AACPC,IAAAA;AAEFA,IAAAA;AACN,EAAA;AAEO,EAAA;AACT;AdqKkB;AACA;AerMQ;AACxB,EAAA;AACA,EAAA;AACA,EAAA;AACF;AfuMkB;AACA;AgBjNE;AACL,EAAA;AAEX,EAAA;AAQJ;AAEM;AACK,EAAA;AACA,IAAA;AACT,EAAA;AACgB,EAAA;AACR,IAAA;AACA,IAAA;AACF,IAAA;AACA,IAAA;AACN,EAAA;AACO,EAAA;AACT;AAEM;AACQ,EAAA;AACd;AAEM;AACQ,EAAA;AACd;AAEM;AACK,EAAA;AACA,IAAA;AACT,EAAA;AACM,EAAA;AACA,EAAA;AACO,EAAA;AACP,EAAA;AACK,EAAA;AACJ,IAAA;AAEA,IAAA;AACP,EAAA;AACO,EAAA;AACT;AAEoB;AAIT,EAAA;AAEG,EAAA;AACN,IAAA;AACS,MAAA;AACF,QAAA;AACT,MAAA;AACK,IAAA;AACO,MAAA;AACA,MAAA;AACN,QAAA;AACM,UAAA;AACR,UAAA;AACF,QAAA;AACF,MAAA;AACY,MAAA;AACd,IAAA;AACF,EAAA;AACO,EAAA;AACT;AAEoB;AACT,EAAA;AAEE,EAAA;AACL,IAAA;AACS,MAAA;AACN,IAAA;AACO,MAAA;AACD,MAAA;AACL,QAAA;AACM,UAAA;AACR,UAAA;AACF,QAAA;AACF,MAAA;AACY,MAAA;AACd,IAAA;AACF,EAAA;AACO,EAAA;AACT;AAEM;AAIK,EAAA;AACH,EAAA;AACA,EAAA;AACU,EAAA;AACA,IAAA;AAChB,EAAA;AACO,EAAA;AACT;AAEM;AAIK,EAAA;AACA,EAAA;AAEH,EAAA;AACC,IAAA;AACA,IAAA;AACA,IAAA;AACP,EAAA;AACM,EAAA;AACE,IAAA;AACA,IAAA;AACA,IAAA;AACR,EAAA;AAEgB,EAAA;AACA,IAAA;AAChB,EAAA;AACO,EAAA;AACT;AAEM;AAIU,EAAA;AACA,EAAA;AAEJ,EAAA;AACD,IAAA;AACT,EAAA;AAEW,EAAA;AACG,IAAA;AACV,MAAA;AACF,IAAA;AAEM,IAAA;AACQ,IAAA;AACL,MAAA;AACT,IAAA;AACF,EAAA;AAEO,EAAA;AACT;AAEiB;AACD,EAAA;AACA,EAAA;AAER,EAAA;AACF,EAAA;AAEM,EAAA;AACN,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AACA,EAAA;AAEY,EAAA;AAET,EAAA;AACT;AAE8B;AACf,EAAA;AAEG,EAAA;AACF,IAAA;AACd,EAAA;AAEM,EAAA;AACA,EAAA;AAEF,EAAA;AAEI,EAAA;AACD,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACI,MAAA;AAEJ,IAAA;AACI,MAAA;AAEJ,IAAA;AACI,MAAA;AAEJ,IAAA;AACI,MAAA;AAEJ,IAAA;AACI,MAAA;AAEJ,IAAA;AACI,MAAA;AACL,QAAA;AACA,QAAA;AACF,MAAA;AAEG,IAAA;AACI,MAAA;AAEJ,IAAA;AACI,MAAA;AAEJ,IAAA;AACA,IAAA;AACA,IAAA;AACI,MAAA;AAEJ,IAAA;AACI,MAAA;AACL,QAAA;AACA,QAAA;AACF,MAAA;AAEG,IAAA;AACsB,MAAA;AAEtB,IAAA;AACqB,MAAA;AAErB,IAAA;AACqB,MAAA;AAErB,IAAA;AACI,MAAA;AACL,QAAA;AACA,QAAA;AACF,MAAA;AAEF,IAAA;AACS,MAAA;AACX,EAAA;AACF;AAI2B;AAEd,EAAA;AAMb;AhB4IkB;AACA;AiB3ZhB;AAGE,EAAA;AAEC,EAAA;AAES,IAAA;AACH,IAAA;AACC,EAAA;AACHC,EAAAA;AACT;AjByZkB;AACA;AkB/YL;AAMS,EAAA;AAAA,IAAA;AAAgC,EAAA;AALxB,iBAAA;AACL,kBAAA;AACD,kBAAA;AACc,kBAAA;AAIV,EAAA;AACf,IAAA;AACA,MAAA;AACD,QAAA;AACF,UAAA;AACF,QAAA;AACF,MAAA;AACF,IAAA;AAEY,IAAA;AACd,EAAA;AAEA,EAAA;AACc,IAAA;AACd,EAAA;AAEmC,EAAA;AAC1B,IAAA;AACK,MAAA;AACF,QAAA;AACG,UAAA;AACC,YAAA;AACJ,cAAA;AACD,YAAA;AAED,YAAA;AAEE,cAAA;AACA,cAAA;AACD,YAAA;AACF,UAAA;AACH,QAAA;AAEK,QAAA;AACK,QAAA;AACH,UAAA;AACP,QAAA;AACF,MAAA;AACY,MAAA;AACd,IAAA;AACF,EAAA;AAEQ,EAAA;AACG,IAAA;AACJ,IAAA;AACA,IAAA;AACP,EAAA;AAE6B,EAAA;AACvB,IAAA;AAEK,MAAA;AAGC,QAAA;AAEF,QAAA;AAEE,QAAA;AAEF,QAAA;AAEG,UAAA;AACP,QAAA;AAEK,QAAA;AACK,QAAA;AACZ,MAAA;AACOH,IAAAA;AACC,MAAA;AACFA,MAAAA;AACN,IAAA;AACK,MAAA;AAEE,MAAA;AAGA,QAAA;AACP,MAAA;AACF,IAAA;AACF,EAAA;AAEc,EAAA;AACR,IAAA;AACS,MAAA;AACX,IAAA;AACK,MAAA;AAGD,MAAA;AACG,QAAA;AACP,MAAA;AAEK,MAAA;AACP,IAAA;AACF,EAAA;AAEQ,kBAAA;AACA,IAAA;AAEDI,MAAAA;AAEL,IAAA;AAEI,IAAA;AAEK,MAAA;AACT,IAAA;AAGW,IAAA;AAEJ,IAAA;AACT,EAAA;AAEQ,kBAAA;AAGD,IAAA;AAEC,EAAA;AACV;AAEM;AAEA;AAOO,EAAA;AACL,IAAA;AAEE,IAAA;AAEF,IAAA;AACG,MAAA;AACH,QAAA;AACM,UAAA;AACN,QAAA;AACF,MAAA;AACC,IAAA;AAEO,IAAA;AACR,MAAA;AACI,MAAA;AACF,QAAA;AACF,MAAA;AACY,MAAA;AACJ,MAAA;AACD,IAAA;AACV,EAAA;AACH;AlBmWkB;AACA;AmB5gBL;AACL,EAAA;AACJ,IAAA;AACc,IAAA;AACf,EAAA;AAGa,EAAA;AAEP,EAAA;AACS,IAAA;AAGF,MAAA;AACR,QAAA;AAEW,UAAA;AAGC,YAAA;AAEN,YAAA;AACA,YAAA;AACF,UAAA;AACE,UAAA;AAEG,QAAA;AACV,MAAA;AACH,IAAA;AAEM,IAAA;AAEM,MAAA;AACD,QAAA;AACT,MAAA;AAGW,MAAA;AAEJ,MAAA;AACT,IAAA;AAEU,IAAA;AACI,MAAA;AACA,MAAA;AACH,QAAA;AACT,MAAA;AACM,MAAA;AACF,MAAA;AACG,MAAA;AACT,IAAA;AAEM,IAAA;AAIG,MAAA;AACI,QAAA;AAGD,UAAA;AAGF,UAAA;AACF,YAAA;AACF,UAAA;AACQ,YAAA;AACF,YAAA;AACN,UAAA;AACF,QAAA;AACE,QAAA;AACJ,MAAA;AACF,IAAA;AACF,EAAA;AACF;AnBsfkB;AACA;AoBllBL;AAGS;AACpB,EAAA;AACF;ApBklBkB;AACA;AqBxlBI;AACT,EAAA;AACb;AAWa;AACsB,EAAA;AAEjB,EAAA;AACA,IAAA;AACL,MAAA;AACA,MAAA;AACR,IAAA;AACA,EAAA;AAEI,EAAA;AACT;ArB8kBkB;AACA;AsBvmBA;AtBymBA;AACA;AuB1mBX;AACO,EAAA;AACJ,IAAA;AACR,EAAA;AACF;AA0B0B;AAEtB,EAAA;AAGY,IAAA;AACD,sBAAA;AAAmD;AAAA;AAGjD,MAAA;AACb,IAAA;AACF,EAAA;AAGE,EAAA;AAEM,IAAA;AAEO,IAAA;AACD,MAAA;AAEL,IAAA;AAGT,EAAA;AACF;AvBykBkB;AACA;AsBxnB4B;AAEpB;AAIX,EAAA;AAEN,EAAA;AACE,IAAA;AACD,MAAA;AACI,QAAA;AAEI,QAAA;AACF,UAAA;AACJ,YAAA;AACF,UAAA;AACF,QAAA;AACO,QAAA;AACAJ,MAAAA;AACG,QAAA;AACHA,UAAAA;AACE,UAAA;AACT,QAAA;AACMA,QAAAA;AACR,MAAA;AACF,IAAA;AACU,qBAAA;AACZ,EAAA;AACF;AtBonBkB;AACA;AwBjpBS;AACT,EAAA;AAGL,EAAA;AACE,IAAA;AACE,MAAA;AACb,IAAA;AACa,IAAA;AACA,MAAA;AACD,QAAA;AACV,MAAA;AACF,IAAA;AACF,EAAA;AAIc,EAAA;AAEF,EAAA;AACC,IAAA;AAEJ,MAAA;AACP,IAAA;AACa,IAAA;AACA,MAAA;AAEJ,QAAA;AACP,MAAA;AACF,IAAA;AACF,EAAA;AAGa,EAAA;AAAC,EAAA;AAChB;AxB0oBkB;AACA;AyBvrBE;AAEI;AAChB,EAAA;AACJ,IAAA;AACY,IAAA;AACd,EAAA;AAGa,EAAA;AACA,EAAA;AACf;AzBsrBkB;AACA;A0BxrBL;AAIT,EAAA;AAIJ;AAEa;AAOA;AAQLK,EAAAA;AACD,IAAA;AACW,IAAA;AACF,IAAA;AACZ,IAAA;AACU,MAAA;AACA,MAAA;AAEH,MAAA;AACO,QAAA;AACR,2BAAA;AAEF,QAAA;AACJ,IAAA;AACF,EAAA;AAEY,EAAA;AACVA,IAAAA;AAEKA,EAAAA;AACT;A1BkqBkB;AACA;AY5qBL;AACK,EAAA;AAET,EAAA;AAEH,IAAA;AAKM,MAAA;AACC,QAAA;AACP,MAAA;AAEM,MAAA;AAEA,MAAA;AACJ,QAAA;AACA,QAAA;AAGE,UAAA;AAEM,UAAA;AACA,UAAA;AAEA,UAAA;AAEF,UAAA;AACF,YAAA;AACE,cAAA;AACE,gBAAA;AACA,gBAAA;AACA,gBAAA;AACF,cAAA;AACE,cAAA;AACJ,YAAA;AACF,UAAA;AAEM,UAAA;AACA,UAAA;AACA,UAAA;AACE,UAAA;AAED,UAAA;AACL,YAAA;AACE,cAAA;AACA,cAAA;AACA,cAAA;AACF,YAAA;AACE,YAAA;AACJ,UAAA;AACF,QAAA;AACU,QAAA;AACR,UAAA;AAEM,UAAA;AACA,UAAA;AAIA,UAAA;AAEC,UAAA;AACT,QAAA;AACO,QAAA;AACL,UAAA;AAEM,UAAA;AACA,UAAA;AAIC,UAAA;AACT,QAAA;AACA,QAAA;AACE,UAAA;AAEM,UAAA;AAEF,UAAA;AACI,YAAA;AAA8C,cAAA;AAEpD,YAAA;AAEI,YAAA;AACF,cAAA;AACE,gBAAA;AACE,kBAAA;AAAA,oBAAA;AACc,oBAAA;AACE,oBAAA;AAEhB,kBAAA;AACA,kBAAA;AACF,gBAAA;AACF,cAAA;AACF,YAAA;AACE,cAAA;AACE,gBAAA;AACA,gBAAA;AACF,cAAA;AAEA,cAAA;AAEA,cAAA;AACE,gBAAA;AACE,kBAAA;AAAA,oBAAA;AACc,oBAAA;AACE,oBAAA;AAEhB,kBAAA;AACA,kBAAA;AACF,gBAAA;AACF,cAAA;AACF,YAAA;AACF,UAAA;AAEM,UAAA;AAEE,UAAA;AAED,UAAA;AACL,YAAA;AACE,cAAA;AACE,gBAAA;AACA,gBAAA;AACA,gBAAA;AACF,cAAA;AACE,cAAA;AACJ,YAAA;AACF,UAAA;AACF,QAAA;AACA,QAAA;AAKE,UAAA;AAEM,UAAA;AAEA,UAAA;AAA8C,YAAA;AAEpD,UAAA;AAEI,UAAA;AACF,YAAA;AACE,cAAA;AACE,gBAAA;AACE,kBAAA;AACA,kBAAA;AACA,kBAAA;AACA,kBAAA;AACF,gBAAA;AACE,gBAAA;AACJ,cAAA;AACF,YAAA;AACF,UAAA;AAEM,UAAA;AAGJ,UAAA;AAGA,YAAA;AACE,cAAA;AACE,gBAAA;AACE,kBAAA;AACA,kBAAA;AACA,kBAAA;AACA,kBAAA;AACF,gBAAA;AACE,gBAAA;AACJ,cAAA;AACF,YAAA;AACF,UAAA;AAEM,UAAA;AAEA,UAAA;AACC,YAAA;AACF,YAAA;AACH,YAAA;AACD,UAAA;AAEO,UAAA;AAED,UAAA;AACL,YAAA;AACE,cAAA;AACE,gBAAA;AACA,gBAAA;AACA,gBAAA;AACA,gBAAA;AACF,cAAA;AACE,cAAA;AACJ,YAAA;AACF,UAAA;AACF,QAAA;AACQ,QAAA;AAKE,UAAA;AAER,UAAA;AACM,UAAA;AAEA,UAAA;AAGH,UAAA;AAOD,YAAA;AACE,cAAA;AACE,gBAAA;AACA,gBAAA;AACF,cAAA;AACE,cAAA;AACJ,YAAA;AACF,UAAA;AAEM,UAAA;AAEF,UAAA;AACF,YAAA;AACE,cAAA;AACE,gBAAA;AACA,gBAAA;AACF,cAAA;AACE,cAAA;AACJ,YAAA;AAEG,UAAA;AACG,YAAA;AACA,YAAA;AACD,cAAA;AACH,cAAA;AACwC,YAAA;AAC1C,YAAA;AACK,cAAA;AACH,cAAA;AACE,gBAAA;AACA,gBAAA;AACF,cAAA;AACF,YAAA;AACF,UAAA;AAEI,UAAA;AACI,YAAA;AACD,cAAA;AACL,YAAA;AACA,YAAA;AACF,UAAA;AAEI,UAAA;AACI,YAAA;AACD,cAAA;AACH,cAAA;AACA,cAAA;AACE,gBAAA;AACA,gBAAA;AACF,cAAA;AACF,YAAA;AACA,YAAA;AACK,cAAA;AACH,cAAA;AACE,gBAAA;AACA,gBAAA;AACF,cAAA;AACF,YAAA;AACF,UAAA;AAEO,UAAA;AACL,YAAA;AACE,cAAA;AACA,cAAA;AACF,YAAA;AACE,YAAA;AACJ,UAAA;AACF,QAAA;AACF,MAAA;AAEO,MAAA;AACT,IAAA;AACF,EAAA;AACF;AZymBkB;AACA;A2Bl7BL;AAsBA;AAKH,EAAA;AAGF,EAAA;AAGA,EAAA;AACF,IAAA;AACJ,EAAA;AAGWC,EAAAA;AACHA,IAAAA;AACJ,MAAA;AACA,MAAA;AACD,IAAA;AACH,EAAA;AACF;AAuCa;AACX,EAAA;AACA,EAAA;AACA,EAAA;AACoF;AACpF,EAAA;AACe,EAAA;AACA,IAAA;AACD,MAAA;AACZ,IAAA;AACa,IAAA;AACR,MAAA;AACO,MAAA;AACX,IAAA;AACH,EAAA;AACU,EAAA;AAES,IAAA;AACD,MAAA;AACZ,IAAA;AACO,IAAA;AACF,MAAA;AACO,MAAA;AACX,IAAA;AAEH,EAAA;AACN;AAyBa;AAMH,EAAA;AAED,EAAA;AACG,IAAA;AAIA,MAAA;AAEKC,MAAAA;AACH,QAAA;AACA,UAAA;AACF,YAAA;AACK,UAAA;AACL,YAAA;AACF,UAAA;AACD,QAAA;AACH,MAAA;AACF,IAAA;AACA,IAAA;AACU,IAAA;AACR,MAAA;AACuE,IAAA;AAGjE,MAAA;AACA,MAAA;AAEK,MAAA;AACL,QAAA;AACI,UAAA;AACA,UAAA;AACR,QAAA;AACF,MAAA;AACF,IAAA;AACD,EAAA;AACH;AAyBa;AAMJ,EAAA;AACF,IAAA;AACH,IAAA;AAED,EAAA;AACH;A3B0yBkB;AACA;A4B1gCHN;A5B4gCG;AACA;A6BzgCL;AACCF,EAAAA;AACG,IAAA;AACf,EAAA;AACF;AAEyB;AACX,EAAA;AACA,EAAA;AAEA,EAAA;AACA,EAAA;AAEE,EAAA;AACD,IAAA;AACF,MAAA;AACT,IAAA;AACc,IAAA;AACf,EAAA;AACH;AAE2B;AACf,EAAA;AACZ;AAEa;AAIP,EAAA;AACQ,IAAA;AACHC,EAAAA;AACD,IAAA;AACF,IAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACF,MAAA;AACO,MAAA;AACT,IAAA;AAEA,IAAA;AACa,MAAA;AACX,MAAA;AACF,IAAA;AAEO,IAAA;AACT,EAAA;AACU,EAAA;AACZ;AAEa;AAIP,EAAA;AACE,IAAA;AACGA,EAAAA;AACD,IAAA;AAEF,IAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACF,MAAA;AACS,IAAA;AACT,MAAA;AACE,QAAA;AACA,QAAA;AACF,MAAA;AACF,IAAA;AAEO,IAAA;AACT,EAAA;AACU,EAAA;AACZ;AAEa;AAIP,EAAA;AACE,IAAA;AACG,IAAA;AACAA,EAAAA;AACD,IAAA;AAEF,IAAA;AACF,MAAA;AACE,QAAA;AACA,QAAA;AACF,MAAA;AACK,IAAA;AACO,MAAA;AACd,IAAA;AAEO,IAAA;AACT,EAAA;AACF;AAEa;AAIP,EAAA;AACI,IAAA;AACI,IAAA;AACHA,EAAAA;AACF,IAAA;AAED,IAAA;AACC,IAAA;AACP,EAAA;AACF;AAEa;AAKG,EAAA;AACF,IAAA;AAEN,uBAAA;AAAuB;AAAmB;AAAkC;AAChF,IAAA;AACJ;AAEa;AAKK,EAAA;AACJ,IAAA;AAEN,uBAAA;AAAuB;AAAmB;AAAiC;AAC/E,IAAA;AACJ;AAEa;AAKI,EAAA;AACH,IAAA;AAEN,uBAAA;AAAuB;AAAmB;AAA8B;AAC5E,IAAA;AACJ;AAE8B;AACrB,EAAA;AACO,IAAA;AACd,EAAA;AACF;AAEa;AAIKD,EAAAA;AAClB;AAEgB;AAIV,EAAA;AACQ,IAAA;AACd;AAEgB;AAIV,EAAA;AACQ,IAAA;AACd;AAIE;AAGgB,EAAA;AAClB;AAEgB;AAKV,EAAA;AACQ,IAAA;AACLA,MAAAA;AAAkD,UAAA;AAA2C,QAAA;AAClG,IAAA;AACJ;AAEgB;AAKF,EAAA;AACA,IAAA;AACG,uBAAA;AACb,IAAA;AACJ;AAEgB;AAGC,EAAA;AACA,EAAA;AACjB;AAEgB;AAGF,EAAA;AACd;AAYM;AAKA;AAOU;AACP,EAAA;AACS,IAAA;AACA,MAAA;AACd,IAAA;AACW,IAAA;AACG,MAAA;AACd,IAAA;AACc,IAAA;AACZ,MAAA;AACW,wBAAA;AACX,MAAA;AACF,IAAA;AACY,IAAA;AACV,MAAA;AACW,wBAAA;AAGX,MAAA;AACF,IAAA;AACA,IAAA;AACE,MAAA;AACW,wBAAA;AAGX,MAAA;AACF,IAAA;AACA,IAAA;AACE,MAAA;AACW,wBAAA;AACX,MAAA;AACA,MAAA;AACW,wBAAA;AAGJ,UAAA;AAIH,QAAA;AACJ,MAAA;AACF,IAAA;AACA,IAAA;AACE,MAAA;AACW,wBAAA;AAIe,UAAA;AAEtB,QAAA;AACJ,MAAA;AACF,IAAA;AACF,EAAA;AACF;AAEa;AACJ,EAAA;AACI,IAAA;AAEC,MAAA;AACN,MAAA;AACA,MAAA;AACF,IAAA;AACU,IAAA;AACF,IAAA;AACV,IAAA;AACa,MAAA;AACb,IAAA;AACA,IAAA;AACa,MAAA;AACb,IAAA;AACA,IAAA;AACc,MAAA;AACD,MAAA;AACb,IAAA;AACA,IAAA;AACc,MAAA;AACD,MAAA;AACb,IAAA;AACA,IAAA;AACc,MAAA;AACD,MAAA;AACb,IAAA;AACA,IAAA;AACc,MAAA;AACH,MAAA;AACP,QAAA;AACF,MAAA;AACF,IAAA;AACA,IAAA;AACc,MAAA;AACD,MAAA;AACb,IAAA;AACW,IAAA;AACE,MAAA;AACb,IAAA;AACA,IAAA;AACE,MAAA;AAES,QAAA;AAET,MAAA;AACF,IAAA;AACA,IAAA;AACa,MAAA;AACb,IAAA;AACW,IAAA;AACE,MAAA;AACb,IAAA;AACa,IAAA;AACA,MAAA;AACb,IAAA;AACA,IAAA;AAGa,MAAA;AACT,QAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AACF;A7B+5BkB;AACA;A8B7uCL;AACN,EAAA;AACP;AAYS;AAUP,EAAA;AACU,IAAA;AACC,MAAA;AACES,QAAAA;AACC,UAAA;AACE,YAAA;AAIA,YAAA;AACJ,cAAA;AACA,cAAA;AACF,YAAA;AAEA,YAAA;AACF,UAAA;AAEO,UAAA;AACC,YAAA;AACJ,cAAA;AAEI,cAAA;AACF,gBAAA;AACE,kBAAA;AACD,gBAAA;AACH,cAAA;AAEA,cAAA;AACF,YAAA;AACA,YAAA;AACE,cAAA;AAEI,cAAA;AACF,gBAAA;AACE,kBAAA;AACD,gBAAA;AACH,cAAA;AAEA,cAAA;AACF,YAAA;AACA,YAAA;AAGM,cAAA;AACF,gBAAA;AACA,gBAAA;AACE,kBAAA;AAEI,oBAAA;AAAU,sBAAA;AACR,oBAAA;AAEJ,kBAAA;AAEE,oBAAA;AACF,kBAAA;AACJ,gBAAA;AACA,gBAAA;AACF,cAAA;AACE,gBAAA;AACF,cAAA;AACF,YAAA;AACF,UAAA;AACF,QAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AACF;AAES;AAID,EAAA;AAEA,EAAA;AAIU,EAAA;AACd,IAAA;AACF,EAAA;AACF;AAES;AACD,EAAA;AACU,EAAA;AAClB;AAES;AAIHR,EAAAA;AAEK,EAAA;AAEJ,EAAA;AACH,IAAA;AACUA,MAAAA;AACR,MAAA;AACF,IAAA;AACA,IAAA;AACF,EAAA;AAEA,EAAA;AACEA,IAAAA;AACA,IAAA;AACF,EAAA;AAEa,EAAA;AACX,IAAA;AACUA,MAAAA;AACR,MAAA;AACF,IAAA;AACF,EAAA;AACF;A9B6rCkB;AACA;A+Bp1CL;AAGL,EAAA;AAEU,EAAA;AACX,IAAA;AACH,IAAA;AAMS,MAAA;AACT,IAAA;AAEM,IAAA;AASI,MAAA;AACN,QAAA;AACA,QAAA;AACF,MAAA;AAIF,IAAA;AAEA,IAAA;AAKQ,MAAA;AACJ,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AAEM,MAAA;AAEN,MAAA;AACE,QAAA;AACI,QAAA;AACL,MAAA;AAEM,MAAA;AACT,IAAA;AAEA,IAAA;AAGE,IAAA;AAGO,MAAA;AACT,IAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAAA;AAQF,EAAA;AAEO,EAAA;AACT;A/BmzCkB;AACA;A4B31CL;AAET,EAAA;AAEQ,IAAA;AAEA,IAAA;AACC,MAAA;AAEH,QAAA;AAGM,UAAA;AAGA,UAAA;AACA,YAAA;AACE,YAAA;AAEN,YAAA;AACK,cAAA;AACA,cAAA;AACF,YAAA;AACD,cAAA;AACE,gBAAA;AACA,gBAAA;AACA,gBAAA;AACA,gBAAA;AACF,cAAA;AAEA,cAAA;AACE,gBAAA;AACA,gBAAA;AACA,gBAAA;AACE,kBAAA;AACA,kBAAA;AACF,gBAAA;AAID,cAAA;AACH,YAAA;AAGM,YAAA;AACJ,cAAA;AACA,cAAA;AACE,gBAAA;AACE,kBAAA;AACA,kBAAA;AACA,kBAAA;AACD,gBAAA;AACH,cAAA;AACA,cAAA;AACE,gBAAA;AACE,kBAAA;AACA,kBAAA;AACA,kBAAA;AACD,gBAAA;AACH,cAAA;AACA,cAAA;AACE,gBAAA;AACE,kBAAA;AACA,kBAAA;AACD,gBAAA;AACH,cAAA;AACA,cAAA;AACE,gBAAA;AACF,cAAA;AACF,YAAA;AAEM,YAAA;AACJ,cAAA;AACA,cAAA;AACA,cAAA;AACA,cAAA;AACD,YAAA;AACH,UAAA;AAEO,UAAA;AACC,YAAA;AAIJ,cAAA;AACA,cAAA;AAEA,cAAA;AAEI,cAAA;AACF,gBAAA;AACED,mCAAAA;AAEF,gBAAA;AACF,cAAA;AACF,YAAA;AACA,YAAA;AAGE,cAAA;AACI,cAAA;AACF,gBAAA;AACA,gBAAA;AACF,cAAA;AACE,gBAAA;AAEA,gBAAA;AAEA,gBAAA;AACE,kBAAA;AAAA,oBAAA;AAC4B,oBAAA;AAE5B,kBAAA;AACA,kBAAA;AACF,gBAAA;AAEA,gBAAA;AACEC,kBAAAA;AACA,kBAAA;AACF,gBAAA;AAEA,gBAAA;AACE,kBAAA;AAAA,oBAAA;AAC4B,oBAAA;AAE5B,kBAAA;AACF,gBAAA;AACF,cAAA;AACF,YAAA;AACF,UAAA;AACF,QAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AACF;AAGa;AAQJ,EAAA;AACFO,IAAAA;AACO,IAAA;AACE,MAAA;AACEA,MAAAA;AACd,IAAA;AACF,EAAA;AACF;AAEa;AAQG,EAAA;AAChB;AAEa;AAGG;AAIE,EAAA;AACR,IAAA;AAEA,IAAA;AAEE,MAAA;AACC,MAAA;AACR,IAAA;AAEI,IAAA;AACH,MAAA;AACE,QAAA;AACF,MAAA;AACO,MAAA;AACT,IAAA;AAGY,IAAA;AACJ,MAAA;AAEF,MAAA;AAGF,QAAA;AACO,QAAA;AACT,MAAA;AACF,IAAA;AAEO,IAAA;AACT,EAAA;AACF;AAGa;AACK,EAAA;AACL,IAAA;AACK,MAAA;AAER,QAAA;AACQ,QAAA;AACT,MAAA;AACL,IAAA;AACF,EAAA;AACF;A5B6xCkB;AACA;AgChiDL;AAiBG,EAAA;AACL,IAAA;AAKH,EAAA;AACJ,IAAA;AAIF,EAAA;AAEO,EAAA;AACF,IAAA;AAAA;AAEG,IAAA;AACF,IAAA;AAEY,MAAA;AACJ,QAAA;AAGA,QAAA;AAGN,MAAA;AAED,IAAA;AACP,EAAA;AAIF;AAEa;AAiBG,EAAA;AACL,IAAA;AAKF,EAAA;AAAsB,IAAA;AAE7B,EAAA;AACF;AhC4+CkB;AACA;AiC3jDL;AAiBG,EAAA;AACL,IAAA;AAKH,EAAA;AACJ,IAAA;AAIF,EAAA;AAEO,EAAA;AACF,IAAA;AAAA;AAEG,IAAA;AACF,IAAA;AAEY,MAAA;AACJ,QAAA;AAGA,QAAA;AACN,MAAA;AAED,IAAA;AACP,EAAA;AACF;AAEa;AAiBG,EAAA;AACL,IAAA;AAKF,EAAA;AAAsB,IAAA;AAE7B,EAAA;AACF;AjC4gDkB;AACA;AW1kDL;AA6BA;AAGK,EAAA;AAKV,EAAA;AACS,IAAA;AAGf,EAAA;AAGM,EAAA;AAGAE,EAAAA;AAKA,EAAA;AACJ,IAAA;AACM,IAAA;AAaI,MAAA;AAEF,MAAA;AACJ,QAAA;AACA,QAAA;AACF,MAAA;AAEM,MAAA;AAEA,MAAA;AAEC,MAAA;AACL,QAAA;AACA,QAAA;AACA,QAAA;AACF,MAAA;AACF,IAAA;AAME,IAAA;AASM,MAAA;AACA,MAAA;AAIN,MAAA;AACE,QAAA;AACS,wBAAA;AACT,QAAA;AACF,MAAA;AAEM,MAAA;AACK,MAAA;AACA,yCAAA;AAIX,MAAA;AAEM,MAAA;AAOS,QAAA;AAIE,wBAAA;AAEV,MAAA;AAED,MAAA;AAIJ,QAAA;AACQ,QAAA;AACR,QAAA;AACF,MAAA;AAEO,MAAA;AACT,IAAA;AAEA,IAAA;AAYQ,MAAA;AACA,MAAA;AAKN,MAAA;AACE,QAAA;AACS,wBAAA;AACT,QAAA;AACF,MAAA;AAEM,MAAA;AAIE,QAAA;AACJ,UAAA;AACA,UAAA;AACA,UAAA;AACA,UAAA;AACF,QAAA;AACO,QAAA;AACFF,UAAAA;AACGA,UAAAA;AACN,UAAA;AACM,YAAA;AACD,YAAA;AACL,UAAA;AAIF,QAAA;AACD,MAAA;AAEK,MAAA;AACM,QAAA;AACZ,MAAA;AAEY,MAAA;AACP,QAAA;AACA,QAAA;AACJ,MAAA;AAGGE,MAAAA;AACI,QAAA;AACJ,UAAA;AACQ,UAAA;AACR,UAAA;AACA,UAAA;AACD,QAAA;AACH,MAAA;AAEM,MAAA;AACJ,QAAA;AACA,QAAA;AAEF,MAAA;AAEM,MAAA;AACJ,QAAA;AACA,wBAAA;AACF,MAAA;AAEO,MAAA;AACT,IAAA;AAEc,IAAA;AACN,MAAA;AAEC,MAAA;AACT,IAAA;AACF,EAAA;AAEO,EAAA;AACT;AX08CkB;AACA;AkC9rDL;AAEA,EAAA;AACG,EAAA;AACJ,EAAA;AACR,EAAA;AACF;AAMI;AAGA,EAAA;AAEA,EAAA;AACS,IAAA;AACF,MAAA;AACA,IAAA;AACA,MAAA;AACF,QAAA;AACM,QAAA;AACX,MAAA;AACU,IAAA;AACd,EAAA;AAEO,EAAA;AACT;AA+Ca;AAiBK,EAAA;AACJ,IAAA;AAQI,MAAA;AACF,MAAA;AAEA,MAAA;AAGA,MAAA;AAKJ,QAAA;AACA,QAAA;AACM,QAAA;AACI,UAAA;AACJ,UAAA;AAAA;AAAA;AAAA;AAUJ,UAAA;AAEF,QAAA;AACD,MAAA;AAIK,MAAA;AAAA;AAEJ,QAAA;AACA,QAAA;AACG,QAAA;AACD,MAAA;AAEQ,MAAA;AAEN,MAAA;AACF,MAAA;AAGO,MAAA;AACHP,QAAAA;AAEA,QAAA;AAEF,QAAA;AACM,UAAA;AACV,QAAA;AAEA,QAAA;AACF,MAAA;AAII,MAAA;AACK,QAAA;AACF,UAAA;AACH,UAAA;AACA,UAAA;AAAU;AAEV,UAAA;AACA,UAAA;AACF,QAAA;AACF,MAAA;AAOM,MAAA;AAWA,MAAA;AACJ,QAAA;AACA,QAAA;AACA,QAAA;AACM,UAAA;AAOJ,UAAA;AACF,QAAA;AACF,MAAA;AAGO,MAAA;AACF,QAAA;AACH,QAAA;AACU,QAAA;AACZ,MAAA;AACD,IAAA;AAEM,IAAA;AACT,EAAA;AACA,EAAA;AACE,IAAA;AAGF,EAAA;AACF;AAGgB;AAIZ,EAAA;AAIC,EAAA;AACT;AlCmjDkB;AACA;AmC1xDL;AAUO,EAAA;AAEV,EAAA;AACS,IAAA;AACf,EAAA;AAEO,EAAA;AACL,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACF,EAAA;AACF;AnCixDgB;AACA;AoCvvDL;AAIL,EAAA;AAIF,EAAA;AAEG,EAAA;AAEHM,IAAAA;AAEM,MAAA;AAEF,MAAA;AACQ,QAAA;AACR,UAAA;AACF,QAAA;AAEI,MAAA;AAEA,MAAA;AACR,IAAA;AAES,IAAA;AAGD,MAAA;AAEK,MAAA;AACH,QAAA;AAEA,QAAA;AACR,MAAA;AACF,IAAA;AAGET,IAAAA;AAGA,MAAA;AACF,IAAA;AAGE,IAAA;AAGM,MAAA;AAAoD,QAAA;AAE1D,MAAA;AAEI,MAAA;AACQ,QAAA;AACR,UAAA;AACF,QAAA;AACS,MAAA;AACT,QAAA;AACE,UAAA;AACD,QAAA;AACH,MAAA;AACF,IAAA;AAGE,IAAA;AAGW,MAAA;AACJ,QAAA;AAEL,QAAA;AACM,UAAA;AACJ,UAAA;AACD,QAAA;AACH,MAAA;AACF,IAAA;AAES,IAAA;AACD,MAAA;AACN,MAAA;AACO,MAAA;AACT,IAAA;AACF,EAAA;AACF;ApCytDkB;AACA;AqCj3DHE;AAwCF;AASJ,EAAA;AAAwB;AAEnB,IAAA;AACR,EAAA;AAEWF,EAAAA;AAA+B;AAEhC,IAAA;AACR,EAAA;AAEWA,EAAAA;AAA+B;AAEhC,IAAA;AACR,EAAA;AACV;AAEa;AAUL,EAAA;AACA,EAAA;AAGJ,EAAA;AAMJ;AAQa;AACA,EAAA;AACF,EAAA;AACX;AAyBa;AACH,EAAA;AACC,IAAA;AACC,MAAA;AACF,MAAA;AACN,IAAA;AACO,IAAA;AAIC,MAAA;AACF,MAAA;AACN,IAAA;AACF,EAAA;AACF;AA4Ia;AAsDA;AACA;AAEA;AAGA;AAWX;AAaM,EAAA;AACJ,IAAA;AACA,IAAA;AACA,IAAA;AACO,IAAA;AACG,IAAA;AACE,IAAA;AACH,IAAA;AACT,IAAA;AACA,IAAA;AACA,IAAA;AACA,IAAA;AACE,EAAA;AAEE,EAAA;AASF,EAAA;AACW,EAAA;AAEX,EAAA;AACA,EAAA;AAES,EAAA;AACP,IAAA;AAEM,IAAA;AACR,MAAA;AACA,MAAA;AACF,IAAA;AAEa,IAAA;AACC,MAAA;AACZ,MAAA;AACY,IAAA;AAChB,EAAA;AAEc,EAAA;AAMD,IAAA;AAEP,IAAA;AACU,MAAA;AACZ,MAAA;AACF,IAAA;AAEU,IAAA;AACF,MAAA;AACR,IAAA;AACF,EAAA;AAEO,EAAA;AAAA;AAED,IAAA;AACJ,IAAA;AACA,IAAA;AACA,IAAA;AAEE,IAAA;AAEI,MAAA;AACM,QAAA;AACN,UAAA;AACF,QAAA;AACA,QAAA;AACF,MAAA;AAEQ,MAAA;AACN,QAAA;AACF,MAAA;AAEW,MAAA;AAEA,MAAA;AAEX,MAAA;AAEI,MAAA;AACM,QAAA;AACN,UAAA;AACF,QAAA;AACO,QAAA;AACL,UAAA;AACF,QAAA;AACF,MAAA;AAEO,MAAA;AACK,QAAA;AACA,UAAA;AACN,YAAA;AACF,UAAA;AACM,UAAA;AACR,QAAA;AAEI,QAAA;AACM,UAAA;AACN,YAAA;AACF,UAAA;AACO,UAAA;AACT,QAAA;AAEI,QAAA;AACI,UAAA;AACJ,YAAA;AACE,cAAA;AACA,cAAA;AACF,YAAA;AACK,YAAA;AACP,UAAA;AACA,UAAA;AACF,QAAA;AAEI,QAAA;AACM,UAAA;AACN,YAAA;AACF,UAAA;AACO,UAAA;AACT,QAAA;AACQ,QAAA;AACN,UAAA;AACF,QAAA;AAEO,QAAA;AACL,UAAA;AACF,QAAA;AACC,MAAA;AACL,IAAA;AACA,IAAA;AACI,IAAA;AACK,MAAA;AACT,IAAA;AACQ,IAAA;AAID,MAAA;AAED,MAAA;AACK,QAAA;AACD,UAAA;AAEJ,UAAA;AACM,YAAA;AAEE,YAAA;AAAW;AAEfA,cAAAA;AAIA,8BAAA;AACF,YAAA;AAEI,YAAA;AACF,cAAA;AAEI,YAAA;AACJ,cAAA;AACA,cAAA;AACF,YAAA;AAEI,YAAA;AACF,cAAA;AAEI,gBAAA;AACE,kBAAA;AACA,kBAAA;AACA,kBAAA;AACA,kBAAA;AACA,kBAAA;AACF,gBAAA;AACA,gBAAA;AACF,cAAA;AAEE,cAAA;AAEF,gBAAA;AACF,cAAA;AACF,YAAA;AAGE,YAAA;AAGA,cAAA;AACA,cAAA;AACA,cAAA;AACF,YAAA;AAEI,YAAA;AACF,cAAA;AACA,cAAA;AACA,cAAA;AACF,YAAA;AAGE,YAAA;AAGA,cAAA;AACJ,UAAA;AAEO,UAAA;AACN,QAAA;AACIC,MAAAA;AACC,QAAA;AACN,UAAA;AACAA,UAAAA;AACF,QAAA;AACA,QAAA;AACO,QAAA;AACC,UAAA;AACCA,UAAAA;AACC,UAAA;AACV,QAAA;AACF,MAAA;AACF,IAAA;AACF,EAAA;AACF;AAWE;AAaM,EAAA;AACJM,IAAAA;AACc,IAAA;AACZ,MAAA;AACD,IAAA;AACE,IAAA;AACD,EAAA;AAQF,EAAA;AACG,IAAA;AACG,IAAA;AACKA,IAAAA;AACX,IAAA;AACA,IAAA;AACO,IAAA;AACG,MAAA;AAEL,MAAA;AAGS,QAAA;AACI,UAAA;AAEJ,QAAA;AAEN,MAAA;AACG,MAAA;AACX,IAAA;AACa,IAAA;AAId,EAAA;AACH;ArCs9CkB;AACA;AsChiEL;AAGJ,EAAA;AACQ,IAAA;AACL,MAAA;AAIC,MAAA;AACL,QAAA;AACD,MAAA;AACH,IAAA;AACc,IAAA;AACJ,MAAA;AACF,MAAA;AACJ,QAAA;AACF,MAAA;AAEM,MAAA;AACK,QAAA;AACX,MAAA;AAEM,MAAA;AAEA,MAAA;AAGJ,MAAA;AAIO,QAAA;AACL,UAAA;AAEE,UAAA;AAKJ,QAAA;AACF,MAAA;AAEM,MAAA;AACA,QAAA;AACC,QAAA;AACL,QAAA;AACA,MAAA;AAEO,MAAA;AACX,IAAA;AACF,EAAA;AACF;AA6BM;AAIE,EAAA;AAEA,EAAA;AAQE,IAAA;AAED,IAAA;AACO,MAAA;AACR,QAAA;AACF,MAAA;AAEK,IAAA;AACT,EAAA;AAEO,EAAA;AACT;AAEa;AAGL,EAAA;AAEQ,EAAA;AACJ,IAAA;AACC,IAAA;AACA,IAAA;AAES,MAAA;AAEd,IAAA;AACN,EAAA;AAEM,EAAA;AAKD,IAAA;AACH,IAAA;AACA,IAAA;AACE,MAAA;AACA,MAAA;AAED,IAAA;AACY,IAAA;AACd,EAAA;AAEa,EAAA;AAChB;AAEa;AAGL,EAAA;AAEQ,EAAA;AACJ,IAAA;AACC,IAAA;AACA,IAAA;AACX,EAAA;AAEM,EAAA;AACD,IAAA;AACH,IAAA;AACA,IAAA;AACE,MAAA;AACA,MAAA;AACD,IAAA;AACY,IAAA;AACd,EAAA;AAEa,EAAA;AAChB;AtCi9DkB;AACA;AuC1mEL;AAWLG,EAAAA;AAIA,EAAA;AACJA,IAAAA;AACU,IAAA;AACZ,EAAA;AAEI,EAAA;AACQ,IAAA;AAAY;AAElB,MAAA;AAA4C,0BAAA;AAElD,EAAA;AAEOA,EAAAA;AACT;AAQE;AAaW;AAgBH,EAAA;AACM,EAAA;AACZ;AAES;AAeH,EAAA;AACM,EAAA;AACZ;AAEuB;AACjB,EAAA;AACD,EAAA;AACT;AvC2iEkB;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA;AACA","file":"/home/runner/work/emmett/emmett/src/packages/emmett/dist/index.cjs","sourcesContent":[null,"import type { DefaultRecord } from './';\n\nexport type Command<\n  CommandType extends string = string,\n  CommandData extends DefaultRecord = DefaultRecord,\n  CommandMetaData extends DefaultRecord | undefined = undefined,\n> = Readonly<\n  CommandMetaData extends undefined\n    ? {\n        type: CommandType;\n        data: Readonly<CommandData>;\n        metadata?: DefaultCommandMetadata | undefined;\n      }\n    : {\n        type: CommandType;\n        data: CommandData;\n        metadata: CommandMetaData;\n      }\n> & { readonly kind?: 'Command' };\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport type AnyCommand = Command<any, any, any>;\n\nexport type CommandTypeOf<T extends Command> = T['type'];\nexport type CommandDataOf<T extends Command> = T['data'];\nexport type CommandMetaDataOf<T extends Command> = T extends {\n  metadata: infer M;\n}\n  ? M\n  : undefined;\n\nexport type CreateCommandType<\n  CommandType extends string,\n  CommandData extends DefaultRecord,\n  CommandMetaData extends DefaultRecord | undefined = undefined,\n> = Readonly<\n  CommandMetaData extends undefined\n    ? {\n        type: CommandType;\n        data: CommandData;\n        metadata?: DefaultCommandMetadata | undefined;\n      }\n    : {\n        type: CommandType;\n        data: CommandData;\n        metadata: CommandMetaData;\n      }\n> & { readonly kind?: 'Command' };\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport const command = <CommandType extends Command<string, any, any>>(\n  ...args: CommandMetaDataOf<CommandType> extends undefined\n    ? [\n        type: CommandTypeOf<CommandType>,\n        data: CommandDataOf<CommandType>,\n        metadata?: DefaultCommandMetadata | undefined,\n      ]\n    : [\n        type: CommandTypeOf<CommandType>,\n        data: CommandDataOf<CommandType>,\n        metadata: CommandMetaDataOf<CommandType>,\n      ]\n): CommandType => {\n  const [type, data, metadata] = args;\n\n  return metadata !== undefined\n    ? ({ type, data, metadata, kind: 'Command' } as CommandType)\n    : ({ type, data, kind: 'Command' } as CommandType);\n};\n\nexport type DefaultCommandMetadata = { now: Date };\n","import type { DefaultRecord } from './';\nimport type {\n  AnyRecordedMessageMetadata,\n  CombinedMessageMetadata,\n  CommonRecordedMessageMetadata,\n  GlobalPositionTypeOfRecordedMessageMetadata,\n  RecordedMessage,\n  RecordedMessageMetadata,\n  RecordedMessageMetadataWithGlobalPosition,\n  RecordedMessageMetadataWithoutGlobalPosition,\n  StreamPositionTypeOfRecordedMessageMetadata,\n} from './message';\n\nexport type BigIntStreamPosition = bigint;\nexport type BigIntGlobalPosition = bigint;\n\nexport type Event<\n  EventType extends string = string,\n  EventData extends DefaultRecord = DefaultRecord,\n  EventMetaData extends DefaultRecord | undefined = undefined,\n> = Readonly<\n  EventMetaData extends undefined\n    ? {\n        type: EventType;\n        data: EventData;\n      }\n    : {\n        type: EventType;\n        data: EventData;\n        metadata: EventMetaData;\n      }\n> & { readonly kind?: 'Event' };\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport type AnyEvent = Event<any, any, any>;\n\nexport type EventTypeOf<T extends Event> = T['type'];\nexport type EventDataOf<T extends Event> = T['data'];\nexport type EventMetaDataOf<T extends Event> = T extends { metadata: infer M }\n  ? M\n  : undefined;\n\nexport type CreateEventType<\n  EventType extends string,\n  EventData extends DefaultRecord,\n  EventMetaData extends DefaultRecord | undefined = undefined,\n> = Readonly<\n  EventMetaData extends undefined\n    ? {\n        type: EventType;\n        data: EventData;\n      }\n    : {\n        type: EventType;\n        data: EventData;\n        metadata: EventMetaData;\n      }\n> & { readonly kind?: 'Event' };\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport const event = <EventType extends Event<string, any, any>>(\n  ...args: EventMetaDataOf<EventType> extends undefined\n    ? [type: EventTypeOf<EventType>, data: EventDataOf<EventType>]\n    : [\n        type: EventTypeOf<EventType>,\n        data: EventDataOf<EventType>,\n        metadata: EventMetaDataOf<EventType>,\n      ]\n): EventType => {\n  const [type, data, metadata] = args;\n\n  return metadata !== undefined\n    ? ({ type, data, metadata, kind: 'Event' } as EventType)\n    : ({ type, data, kind: 'Event' } as EventType);\n};\n\nexport type CombinedReadEventMetadata<\n  EventType extends Event = Event,\n  EventMetaDataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n> = CombinedMessageMetadata<EventType, EventMetaDataType>;\n\nexport type ReadEvent<\n  EventType extends Event = Event,\n  EventMetaDataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n> = RecordedMessage<EventType, EventMetaDataType>;\n\nexport type AnyReadEvent<\n  EventMetaDataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n> = ReadEvent<AnyEvent, EventMetaDataType>;\n\nexport type CommonReadEventMetadata<StreamPosition = BigIntStreamPosition> =\n  CommonRecordedMessageMetadata<StreamPosition>;\n\nexport type ReadEventMetadata<\n  GlobalPosition = undefined,\n  StreamPosition = BigIntStreamPosition,\n> = RecordedMessageMetadata<GlobalPosition, StreamPosition>;\n\nexport type AnyReadEventMetadata = AnyRecordedMessageMetadata;\n\nexport type ReadEventMetadataWithGlobalPosition<\n  GlobalPosition = BigIntGlobalPosition,\n> = RecordedMessageMetadataWithGlobalPosition<GlobalPosition>;\n\nexport type ReadEventMetadataWithoutGlobalPosition<\n  StreamPosition = BigIntStreamPosition,\n> = RecordedMessageMetadataWithoutGlobalPosition<StreamPosition>;\n\nexport type GlobalPositionTypeOfReadEventMetadata<ReadEventMetadataType> =\n  GlobalPositionTypeOfRecordedMessageMetadata<ReadEventMetadataType>;\n\nexport type StreamPositionTypeOfReadEventMetadata<ReadEventMetadataType> =\n  StreamPositionTypeOfRecordedMessageMetadata<ReadEventMetadataType>;\n","import type {\n  AnyCommand,\n  AnyEvent,\n  BigIntGlobalPosition,\n  BigIntStreamPosition,\n  Command,\n  DefaultRecord,\n  Event,\n} from '.';\n\nexport type Message<\n  Type extends string = string,\n  Data extends DefaultRecord = DefaultRecord,\n  MetaData extends DefaultRecord | undefined = undefined,\n> = Command<Type, Data, MetaData> | Event<Type, Data, MetaData>;\n\nexport type AnyMessage = AnyEvent | AnyCommand;\n\nexport type MessageKindOf<T extends Message> = T['kind'];\nexport type MessageTypeOf<T extends Message> = T['type'];\nexport type MessageDataOf<T extends Message> = T['data'];\nexport type MessageMetaDataOf<T extends Message> = T extends {\n  metadata: infer M;\n}\n  ? M\n  : undefined;\n\nexport type CanHandle<T extends Message> = MessageTypeOf<T>[];\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport const message = <MessageType extends Message<string, any, any>>(\n  ...args: MessageMetaDataOf<MessageType> extends undefined\n    ? [\n        kind: MessageKindOf<MessageType>,\n        type: MessageTypeOf<MessageType>,\n        data: MessageDataOf<MessageType>,\n      ]\n    : [\n        kind: MessageKindOf<MessageType>,\n        type: MessageTypeOf<MessageType>,\n        data: MessageDataOf<MessageType>,\n        metadata: MessageMetaDataOf<MessageType>,\n      ]\n): MessageType => {\n  const [kind, type, data, metadata] = args;\n\n  return metadata !== undefined\n    ? ({ type, data, metadata, kind } as MessageType)\n    : ({ type, data, kind } as MessageType);\n};\n\nexport type CombinedMessageMetadata<\n  MessageType extends Message = Message,\n  MessageMetaDataType extends DefaultRecord = DefaultRecord,\n> =\n  MessageMetaDataOf<MessageType> extends undefined\n    ? MessageMetaDataType\n    : MessageMetaDataOf<MessageType> & MessageMetaDataType;\n\nexport type CombineMetadata<\n  MessageType extends Message = Message,\n  MessageMetaDataType extends DefaultRecord = DefaultRecord,\n> = MessageType & {\n  metadata: CombinedMessageMetadata<MessageType, MessageMetaDataType>;\n};\n\nexport type RecordedMessage<\n  MessageType extends Message = Message,\n  MessageMetaDataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n> = CombineMetadata<MessageType, MessageMetaDataType> & {\n  kind: NonNullable<MessageKindOf<Message>>;\n};\n\nexport type CommonRecordedMessageMetadata<\n  StreamPosition = BigIntStreamPosition,\n> = Readonly<{\n  messageId: string;\n  streamPosition: StreamPosition;\n  streamName: string;\n}>;\n\nexport type WithGlobalPosition<GlobalPosition> = Readonly<{\n  globalPosition: GlobalPosition;\n}>;\n\nexport type RecordedMessageMetadata<\n  GlobalPosition = undefined,\n  StreamPosition = BigIntStreamPosition,\n> = CommonRecordedMessageMetadata<StreamPosition> &\n  // eslint-disable-next-line @typescript-eslint/no-empty-object-type\n  (GlobalPosition extends undefined ? {} : WithGlobalPosition<GlobalPosition>);\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport type AnyRecordedMessage = Message<any, any, any>;\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport type AnyRecordedMessageMetadata = RecordedMessageMetadata<any, any>;\n\nexport type RecordedMessageMetadataWithGlobalPosition<\n  GlobalPosition = BigIntGlobalPosition,\n> = RecordedMessageMetadata<GlobalPosition>;\n\nexport type RecordedMessageMetadataWithoutGlobalPosition<\n  StreamPosition = BigIntStreamPosition,\n> = RecordedMessageMetadata<undefined, StreamPosition>;\n\nexport type GlobalPositionTypeOfRecordedMessageMetadata<\n  RecordedMessageMetadataType,\n> =\n  // eslint-disable-next-line @typescript-eslint/no-explicit-any\n  RecordedMessageMetadataType extends RecordedMessageMetadata<infer GP, any>\n    ? GP\n    : never;\n\nexport type StreamPositionTypeOfRecordedMessageMetadata<\n  RecordedMessageMetadataType,\n> =\n  // eslint-disable-next-line @typescript-eslint/no-explicit-any\n  RecordedMessageMetadataType extends RecordedMessageMetadata<any, infer SV>\n    ? SV\n    : never;\n","import type { AnyCommand } from './command';\nimport type { AnyEvent } from './event';\n\n/// Inspired by https://blog.bittacklr.be/the-workflow-pattern.html\n\nexport type Workflow<\n  Input extends AnyEvent | AnyCommand,\n  State,\n  Output extends AnyEvent | AnyCommand,\n> = {\n  decide: (command: Input, state: State) => WorkflowOutput<Output>[];\n  evolve: (currentState: State, event: WorkflowEvent<Output>) => State;\n  initialState: () => State;\n};\n\nexport type WorkflowEvent<Output extends AnyEvent | AnyCommand> = Extract<\n  Output,\n  { kind?: 'Event' }\n>;\n\nexport type WorkflowCommand<Output extends AnyEvent | AnyCommand> = Extract<\n  Output,\n  { kind?: 'Command' }\n>;\n\nexport type WorkflowOutput<TOutput extends AnyEvent | AnyCommand> =\n  | { action: 'Reply'; message: TOutput }\n  | { action: 'Send'; message: WorkflowCommand<TOutput> }\n  | { action: 'Publish'; message: WorkflowEvent<TOutput> }\n  | {\n      action: 'Schedule';\n      message: TOutput;\n      when: { afterInMs: number } | { at: Date };\n    }\n  | { action: 'Complete' }\n  | { action: 'Accept' }\n  | { action: 'Ignore'; reason: string }\n  | { action: 'Error'; reason: string };\n\nexport const reply = <TOutput extends AnyEvent | AnyCommand>(\n  message: TOutput,\n): WorkflowOutput<TOutput> => {\n  return {\n    action: 'Reply',\n    message,\n  };\n};\n\nexport const send = <TOutput extends AnyEvent | AnyCommand>(\n  message: WorkflowCommand<TOutput>,\n): WorkflowOutput<TOutput> => {\n  return {\n    action: 'Send',\n    message,\n  };\n};\n\nexport const publish = <TOutput extends AnyEvent | AnyCommand>(\n  message: WorkflowEvent<TOutput>,\n): WorkflowOutput<TOutput> => {\n  return {\n    action: 'Publish',\n    message,\n  };\n};\n\nexport const schedule = <TOutput extends AnyEvent | AnyCommand>(\n  message: TOutput,\n  when: { afterInMs: number } | { at: Date },\n): WorkflowOutput<TOutput> => {\n  return {\n    action: 'Schedule',\n    message,\n    when,\n  };\n};\n\nexport const complete = <\n  TOutput extends AnyEvent | AnyCommand,\n>(): WorkflowOutput<TOutput> => {\n  return {\n    action: 'Complete',\n  };\n};\n\nexport const ignore = <TOutput extends AnyEvent | AnyCommand>(\n  reason: string,\n): WorkflowOutput<TOutput> => {\n  return {\n    action: 'Ignore',\n    reason,\n  };\n};\n\nexport const error = <TOutput extends AnyEvent | AnyCommand>(\n  reason: string,\n): WorkflowOutput<TOutput> => {\n  return {\n    action: 'Error',\n    reason,\n  };\n};\n\nexport const accept = <\n  TOutput extends AnyEvent | AnyCommand,\n>(): WorkflowOutput<TOutput> => {\n  return { action: 'Accept' };\n};\n","export * from './deepReadonly';\n\nexport * from './command';\nexport * from './event';\nexport * from './message';\nexport * from './messageHandling';\n\nexport * from './decider';\nexport * from './workflow';\n\nexport type Brand<K, T> = K & { readonly __brand: T };\nexport type Flavour<K, T> = K & { readonly __brand?: T };\n\nexport type DefaultRecord = Record<string, unknown>;\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport type AnyRecord = Record<string, any>;\n\nexport type NonNullable<T> = T extends null | undefined ? never : T;\n\nexport const emmettPrefix = 'emt';\n\nexport const globalTag = 'global';\nexport const defaultTag = `${emmettPrefix}:default`;\nexport const unknownTag = `${emmettPrefix}:unknown`;\n","import {\n  type BatchRecordedMessageHandlerWithContext,\n  type BatchRecordedMessageHandlerWithoutContext,\n  type DefaultRecord,\n  type Event,\n  type ReadEvent,\n} from '../../typing';\nimport type { EventStore, EventStoreReadEventMetadata } from '../eventStore';\n\nexport type AfterEventStoreCommitHandler<\n  Store extends EventStore,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n> = HandlerContext extends undefined\n  ? BatchRecordedMessageHandlerWithoutContext<\n      Event,\n      EventStoreReadEventMetadata<Store>\n    >\n  : BatchRecordedMessageHandlerWithContext<\n      Event,\n      EventStoreReadEventMetadata<Store>,\n      NonNullable<HandlerContext>\n    >;\n\nexport type BeforeEventStoreCommitHandler<\n  Store extends EventStore,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n> = HandlerContext extends undefined\n  ? BatchRecordedMessageHandlerWithoutContext<\n      Event,\n      EventStoreReadEventMetadata<Store>\n    >\n  : BatchRecordedMessageHandlerWithContext<\n      Event,\n      EventStoreReadEventMetadata<Store>,\n      NonNullable<HandlerContext>\n    >;\n\ntype TryPublishMessagesAfterCommitOptions<\n  Store extends EventStore,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n> = {\n  onAfterCommit?: AfterEventStoreCommitHandler<Store, HandlerContext>;\n};\n\nexport async function tryPublishMessagesAfterCommit<Store extends EventStore>(\n  messages: ReadEvent<Event, EventStoreReadEventMetadata<Store>>[],\n  options: TryPublishMessagesAfterCommitOptions<Store, undefined> | undefined,\n): Promise<boolean>;\nexport async function tryPublishMessagesAfterCommit<\n  Store extends EventStore,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n>(\n  messages: ReadEvent<Event, EventStoreReadEventMetadata<Store>>[],\n  options:\n    | TryPublishMessagesAfterCommitOptions<Store, HandlerContext>\n    | undefined,\n  context: HandlerContext,\n): Promise<boolean>;\nexport async function tryPublishMessagesAfterCommit<\n  Store extends EventStore,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n>(\n  messages: ReadEvent<Event, EventStoreReadEventMetadata<Store>>[],\n  options:\n    | TryPublishMessagesAfterCommitOptions<Store, HandlerContext>\n    | undefined,\n  context?: HandlerContext,\n): Promise<boolean> {\n  if (options?.onAfterCommit === undefined) return false;\n\n  try {\n    await options?.onAfterCommit(messages, context!);\n    return true;\n  } catch (error) {\n    // TODO: enhance with tracing\n    console.error(`Error in on after commit hook`, error);\n    return false;\n  }\n}\n","import type { EventsPublisher } from '../../messageBus';\nimport type { DefaultRecord, Event, ReadEvent } from '../../typing';\nimport type { EventStore, EventStoreReadEventMetadata } from '../eventStore';\nimport type { AfterEventStoreCommitHandler } from './afterEventStoreCommitHandler';\n\nexport const forwardToMessageBus = <\n  Store extends EventStore,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n>(\n  eventPublisher: EventsPublisher,\n): AfterEventStoreCommitHandler<Store, HandlerContext> =>\n  (async (\n    messages: ReadEvent<Event, EventStoreReadEventMetadata<Store>>[],\n  ): Promise<void> => {\n    for (const message of messages) {\n      await eventPublisher.publish(message);\n    }\n  }) as AfterEventStoreCommitHandler<Store, HandlerContext>;\n","import {\n  event,\n  type Event,\n  type EventDataOf,\n  type ReadEvent,\n  type ReadEventMetadataWithGlobalPosition,\n} from '../../typing';\n\nexport const GlobalStreamCaughtUpType = '__emt:GlobalStreamCaughtUp';\n\nexport type GlobalStreamCaughtUp = Event<\n  '__emt:GlobalStreamCaughtUp',\n  { globalPosition: bigint },\n  { globalPosition: bigint }\n>;\n\nexport const isGlobalStreamCaughtUp = (\n  event: Event,\n): event is GlobalStreamCaughtUp => event.type === GlobalStreamCaughtUpType;\n\nexport const caughtUpEventFrom =\n  (position: bigint) =>\n  (\n    event: ReadEvent<Event, ReadEventMetadataWithGlobalPosition>,\n  ): event is ReadEvent<\n    GlobalStreamCaughtUp,\n    ReadEventMetadataWithGlobalPosition\n  > =>\n    event.type === GlobalStreamCaughtUpType &&\n    event.metadata?.globalPosition >= position;\n\nexport const globalStreamCaughtUp = (\n  data: EventDataOf<GlobalStreamCaughtUp>,\n): GlobalStreamCaughtUp =>\n  event<GlobalStreamCaughtUp>(GlobalStreamCaughtUpType, data, {\n    globalPosition: data.globalPosition,\n  });\n\nexport const isSubscriptionEvent = (\n  event: Event,\n): event is GlobalSubscriptionEvent => isGlobalStreamCaughtUp(event);\n\nexport const isNotInternalEvent = (event: Event): boolean =>\n  !isGlobalStreamCaughtUp(event);\n\nexport type GlobalSubscriptionEvent = GlobalStreamCaughtUp;\n","import type {\n  AnyReadEventMetadata,\n  BigIntGlobalPosition,\n  BigIntStreamPosition,\n  CommonReadEventMetadata,\n  DefaultRecord,\n  Event,\n  GlobalPositionTypeOfReadEventMetadata,\n  ReadEvent,\n  ReadEventMetadata,\n  StreamPositionTypeOfReadEventMetadata,\n  WithGlobalPosition,\n} from '../typing';\nimport type { AfterEventStoreCommitHandler } from './afterCommit';\nimport type { ExpectedStreamVersion } from './expectedVersion';\n\n// #region event-store\nexport interface EventStore<\n  ReadEventMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n> {\n  aggregateStream<\n    State,\n    EventType extends Event,\n    EventPayloadType extends Event = EventType,\n  >(\n    streamName: string,\n    options: AggregateStreamOptions<\n      State,\n      EventType,\n      ReadEventMetadataType,\n      EventPayloadType\n    >,\n  ): Promise<\n    AggregateStreamResult<\n      State,\n      StreamPositionTypeOfReadEventMetadata<ReadEventMetadataType>\n    >\n  >;\n\n  readStream<\n    EventType extends Event,\n    EventPayloadType extends Event = EventType,\n  >(\n    streamName: string,\n    options?: ReadStreamOptions<\n      StreamPositionTypeOfReadEventMetadata<ReadEventMetadataType>,\n      EventType,\n      EventPayloadType\n    >,\n  ): Promise<ReadStreamResult<EventType, ReadEventMetadataType>>;\n\n  appendToStream<\n    EventType extends Event,\n    EventPayloadType extends Event = EventType,\n  >(\n    streamName: string,\n    events: EventType[],\n    options?: AppendToStreamOptions<\n      StreamPositionTypeOfReadEventMetadata<ReadEventMetadataType>,\n      EventType,\n      EventPayloadType\n    >,\n  ): Promise<\n    AppendToStreamResult<\n      StreamPositionTypeOfReadEventMetadata<ReadEventMetadataType>\n    >\n  >;\n\n  streamExists(streamName: string): Promise<StreamExistsResult>;\n\n  // streamEvents(): ReadableStream<\n  //   ReadEvent<Event, ReadEventMetadataType> | GlobalSubscriptionEvent\n  // >;\n}\n\nexport type EventStoreReadEventMetadata<Store extends EventStore> =\n  Store extends EventStore<infer T>\n    ? T extends CommonReadEventMetadata<infer SP>\n      ? T extends WithGlobalPosition<infer GP>\n        ? ReadEventMetadata<GP, SP> & T\n        : ReadEventMetadata<undefined, SP> & T\n      : never\n    : never;\n\nexport type GlobalPositionTypeOfEventStore<Store extends EventStore> =\n  GlobalPositionTypeOfReadEventMetadata<EventStoreReadEventMetadata<Store>>;\n\nexport type StreamPositionTypeOfEventStore<Store extends EventStore> =\n  StreamPositionTypeOfReadEventMetadata<EventStoreReadEventMetadata<Store>>;\n\nexport type EventStoreSession<EventStoreType extends EventStore> = {\n  eventStore: EventStoreType;\n  close: () => Promise<void>;\n};\n\nexport interface EventStoreSessionFactory<EventStoreType extends EventStore> {\n  withSession<T = unknown>(\n    callback: (session: EventStoreSession<EventStoreType>) => Promise<T>,\n  ): Promise<T>;\n}\n// #endregion event-store\n\nexport const canCreateEventStoreSession = <Store extends EventStore>(\n  eventStore: Store | EventStoreSessionFactory<Store>,\n): eventStore is EventStoreSessionFactory<Store> => 'withSession' in eventStore;\n\nexport const nulloSessionFactory = <EventStoreType extends EventStore>(\n  eventStore: EventStoreType,\n): EventStoreSessionFactory<EventStoreType> => ({\n  withSession: (callback) => {\n    const nulloSession: EventStoreSession<EventStoreType> = {\n      eventStore,\n      close: () => Promise.resolve(),\n    };\n\n    return callback(nulloSession);\n  },\n});\n\n////////////////////////////////////////////////////////////////////\n/// Schema Versioning types\n////////////////////////////////////////////////////////////////////\n\nexport type EventStoreReadSchemaOptions<\n  StreamEvent extends Event = Event,\n  StoredEvent extends Event = StreamEvent,\n> = {\n  versioning?: {\n    upcast?: (event: StoredEvent) => StreamEvent;\n  };\n};\n\nexport type EventStoreAppendSchemaOptions<\n  StreamEvent extends Event = Event,\n  StoredEvent extends Event = StreamEvent,\n> = {\n  versioning?: {\n    downcast?: (event: StreamEvent) => StoredEvent;\n  };\n};\n\nexport type EventStoreSchemaOptions<\n  StreamEvent extends Event = Event,\n  StoredEvent extends Event = StreamEvent,\n> = EventStoreReadSchemaOptions<StreamEvent, StoredEvent> &\n  EventStoreAppendSchemaOptions<StreamEvent, StoredEvent>;\n\n////////////////////////////////////////////////////////////////////\n/// ReadStream types\n////////////////////////////////////////////////////////////////////\n\nexport type ReadStreamOptions<\n  StreamVersion = BigIntStreamPosition,\n  EventType extends Event = Event,\n  EventPayloadType extends Event = EventType,\n> = {\n  from?: StreamVersion;\n  to?: StreamVersion;\n  maxCount?: bigint;\n  expectedStreamVersion?: ExpectedStreamVersion<StreamVersion>;\n  schema?: EventStoreReadSchemaOptions<EventType, EventPayloadType>;\n};\n\nexport type ReadStreamResult<\n  EventType extends Event,\n  ReadEventMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n> = {\n  currentStreamVersion: StreamPositionTypeOfReadEventMetadata<ReadEventMetadataType>;\n  events: ReadEvent<EventType, ReadEventMetadataType>[];\n  streamExists: boolean;\n};\n\n////////////////////////////////////////////////////////////////////\n/// AggregateStream types\n////////////////////////////////////////////////////////////////////\n\ntype Evolve<\n  State,\n  EventType extends Event,\n  ReadEventMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n> =\n  | ((currentState: State, event: EventType) => State)\n  | ((\n      currentState: State,\n      event: ReadEvent<EventType, ReadEventMetadataType>,\n    ) => State)\n  | ((currentState: State, event: ReadEvent<EventType>) => State);\n\nexport type AggregateStreamOptions<\n  State,\n  EventType extends Event,\n  ReadEventMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  EventPayloadType extends Event = EventType,\n> = {\n  evolve: Evolve<State, EventType, ReadEventMetadataType>;\n  initialState: () => State;\n  read?: ReadStreamOptions<\n    StreamPositionTypeOfReadEventMetadata<ReadEventMetadataType>,\n    EventType,\n    EventPayloadType\n  >;\n};\n\nexport type AggregateStreamResult<\n  State,\n  StreamPosition = BigIntStreamPosition,\n> = {\n  currentStreamVersion: StreamPosition;\n  state: State;\n  streamExists: boolean;\n};\n\nexport type AggregateStreamResultWithGlobalPosition<\n  State,\n  StreamPosition = BigIntStreamPosition,\n  GlobalPosition = BigIntGlobalPosition,\n> =\n  | (AggregateStreamResult<State, StreamPosition> & {\n      streamExists: true;\n      lastEventGlobalPosition: GlobalPosition;\n    })\n  | (AggregateStreamResult<State, StreamPosition> & {\n      streamExists: false;\n    });\n\nexport type AggregateStreamResultOfEventStore<Store extends EventStore> =\n  // eslint-disable-next-line @typescript-eslint/no-explicit-any\n  Store['aggregateStream'] extends (...args: any[]) => Promise<infer R>\n    ? R\n    : never;\n\n////////////////////////////////////////////////////////////////////\n/// AppendToStream types\n////////////////////////////////////////////////////////////////////\n\nexport type AppendToStreamOptions<\n  StreamVersion = BigIntStreamPosition,\n  EventType extends Event = Event,\n  EventPayloadType extends Event = EventType,\n> = {\n  expectedStreamVersion?: ExpectedStreamVersion<StreamVersion>;\n  schema?: EventStoreAppendSchemaOptions<EventType, EventPayloadType>;\n};\n\nexport type AppendToStreamResult<StreamVersion = BigIntStreamPosition> = {\n  nextExpectedStreamVersion: StreamVersion;\n  createdNewStream: boolean;\n};\n\nexport type AppendToStreamResultWithGlobalPosition<\n  StreamVersion = BigIntStreamPosition,\n  GlobalPosition = BigIntGlobalPosition,\n> = AppendToStreamResult<StreamVersion> & {\n  lastEventGlobalPosition: GlobalPosition;\n};\n\nexport type AppendStreamResultOfEventStore<Store extends EventStore> =\n  // eslint-disable-next-line @typescript-eslint/no-explicit-any\n  Store['appendToStream'] extends (...args: any[]) => Promise<infer R>\n    ? R\n    : never;\n\n////////////////////////////////////////////////////////////////////\n/// StreamExists types\n////////////////////////////////////////////////////////////////////\n\nexport type StreamExistsResult = boolean;\n\n////////////////////////////////////////////////////////////////////\n/// DefaultEventStoreOptions\n////////////////////////////////////////////////////////////////////\n\nexport type DefaultEventStoreOptions<\n  Store extends EventStore,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n> = {\n  /**\n   * Pluggable set of hooks informing about the event store internal behaviour.\n   */\n  hooks?: {\n    /**\n     * This hook will be called **AFTER** events were stored in the event store.\n     * It's designed to handle scenarios where delivery and ordering guarantees do not matter much.\n     *\n     * **WARNINGS:**\n     *\n     *  1. It will be called **EXACTLY ONCE** if append succeded.\n     *  2. If the hook fails, its append **will still silently succeed**, and no error will be thrown.\n     *  3. Wen process crashes after events were committed, but before the hook was called, delivery won't be retried.\n     * That can lead to state inconsistencies.\n     *  4. In the case of high concurrent traffic, **race conditions may cause ordering issues**.\n     * For instance, where the second hook takes longer to process than the first one, ordering won't be guaranteed.\n     *\n     * @type {AfterEventStoreCommitHandler<Store, HandlerContext>}\n     */\n    onAfterCommit?: AfterEventStoreCommitHandler<Store, HandlerContext>;\n  };\n};\n","import { ConcurrencyError, EmmettError } from '../errors';\nimport type { BigIntStreamPosition, Flavour } from '../typing';\n\nexport type ExpectedStreamVersion<VersionType = BigIntStreamPosition> =\n  | ExpectedStreamVersionWithValue<VersionType>\n  | ExpectedStreamVersionGeneral;\n\nexport type ExpectedStreamVersionWithValue<VersionType = BigIntStreamPosition> =\n  Flavour<VersionType, 'StreamVersion'>;\n\nexport type ExpectedStreamVersionGeneral = Flavour<\n  'STREAM_EXISTS' | 'STREAM_DOES_NOT_EXIST' | 'NO_CONCURRENCY_CHECK',\n  'StreamVersion'\n>;\n\nexport const STREAM_EXISTS = 'STREAM_EXISTS' as ExpectedStreamVersionGeneral;\nexport const STREAM_DOES_NOT_EXIST =\n  'STREAM_DOES_NOT_EXIST' as ExpectedStreamVersionGeneral;\nexport const NO_CONCURRENCY_CHECK =\n  'NO_CONCURRENCY_CHECK' as ExpectedStreamVersionGeneral;\n\nexport const matchesExpectedVersion = <StreamVersion = BigIntStreamPosition>(\n  current: StreamVersion | undefined,\n  expected: ExpectedStreamVersion<StreamVersion>,\n  defaultVersion: StreamVersion,\n): boolean => {\n  if (expected === NO_CONCURRENCY_CHECK) return true;\n\n  if (expected == STREAM_DOES_NOT_EXIST) return current === defaultVersion;\n\n  if (expected == STREAM_EXISTS) return current !== defaultVersion;\n\n  return current === expected;\n};\n\nexport const assertExpectedVersionMatchesCurrent = <\n  StreamVersion = BigIntStreamPosition,\n>(\n  current: StreamVersion,\n  expected: ExpectedStreamVersion<StreamVersion> | undefined,\n  defaultVersion: StreamVersion,\n): void => {\n  expected ??= NO_CONCURRENCY_CHECK;\n\n  if (!matchesExpectedVersion(current, expected, defaultVersion))\n    throw new ExpectedVersionConflictError(current, expected);\n};\n\nexport class ExpectedVersionConflictError<\n  VersionType = BigIntStreamPosition,\n> extends ConcurrencyError {\n  constructor(\n    current: VersionType,\n    expected: ExpectedStreamVersion<VersionType>,\n  ) {\n    super(current?.toString(), expected?.toString());\n\n    // 👇️ because we are extending a built-in class\n    Object.setPrototypeOf(this, ExpectedVersionConflictError.prototype);\n  }\n}\n\nexport const isExpectedVersionConflictError = (\n  error: unknown,\n): error is ExpectedVersionConflictError =>\n  error instanceof ExpectedVersionConflictError ||\n  EmmettError.isInstanceOf<ConcurrencyError>(\n    error,\n    ExpectedVersionConflictError.Codes.ConcurrencyError,\n  );\n","import { v4 as uuid } from 'uuid';\nimport {\n  getInMemoryDatabase,\n  type InMemoryDatabase,\n} from '../database/inMemoryDatabase';\nimport type { ProjectionRegistration } from '../projections';\nimport type {\n  BigIntStreamPosition,\n  CombinedReadEventMetadata,\n  Event,\n  ReadEvent,\n  ReadEventMetadataWithGlobalPosition,\n} from '../typing';\nimport { tryPublishMessagesAfterCommit } from './afterCommit';\nimport {\n  type AggregateStreamOptions,\n  type AggregateStreamResult,\n  type AppendToStreamOptions,\n  type AppendToStreamResult,\n  type DefaultEventStoreOptions,\n  type EventStore,\n  type ReadStreamOptions,\n  type ReadStreamResult,\n  type StreamExistsResult,\n} from './eventStore';\nimport { assertExpectedVersionMatchesCurrent } from './expectedVersion';\nimport { handleInMemoryProjections } from './projections/inMemory';\nimport { downcastRecordedMessages, upcastRecordedMessages } from './versioning';\n\nexport const InMemoryEventStoreDefaultStreamVersion = 0n;\n\nexport type InMemoryEventStore =\n  EventStore<ReadEventMetadataWithGlobalPosition> & {\n    database: InMemoryDatabase;\n  };\n\nexport type InMemoryReadEventMetadata = ReadEventMetadataWithGlobalPosition;\n\nexport type InMemoryProjectionHandlerContext = {\n  eventStore?: InMemoryEventStore;\n  database?: InMemoryDatabase;\n};\n\nexport type InMemoryEventStoreOptions =\n  DefaultEventStoreOptions<InMemoryEventStore> & {\n    projections?: ProjectionRegistration<\n      'inline',\n      InMemoryReadEventMetadata,\n      InMemoryProjectionHandlerContext\n    >[];\n    database?: InMemoryDatabase;\n  };\n\nexport type InMemoryReadEvent<EventType extends Event = Event> = ReadEvent<\n  EventType,\n  ReadEventMetadataWithGlobalPosition\n>;\n\nexport const getInMemoryEventStore = (\n  eventStoreOptions?: InMemoryEventStoreOptions,\n): InMemoryEventStore => {\n  const streams = new Map<\n    string,\n    ReadEvent<Event, ReadEventMetadataWithGlobalPosition>[]\n  >();\n\n  const getAllEventsCount = () => {\n    return Array.from<ReadEvent[]>(streams.values())\n      .map((s) => s.length)\n      .reduce((p, c) => p + c, 0);\n  };\n\n  // Get the database instance to be used for projections\n  const database = eventStoreOptions?.database || getInMemoryDatabase();\n\n  // Extract inline projections from options\n  const inlineProjections = (eventStoreOptions?.projections ?? [])\n    .filter(({ type }) => type === 'inline')\n    .map(({ projection }) => projection);\n\n  // Create the event store object\n  const eventStore: InMemoryEventStore = {\n    database,\n    async aggregateStream<\n      State,\n      EventType extends Event,\n      EventPayloadType extends Event = EventType,\n    >(\n      streamName: string,\n      options: AggregateStreamOptions<\n        State,\n        EventType,\n        ReadEventMetadataWithGlobalPosition,\n        EventPayloadType\n      >,\n    ): Promise<AggregateStreamResult<State>> {\n      const { evolve, initialState, read } = options;\n\n      const result = await this.readStream<EventType, EventPayloadType>(\n        streamName,\n        read,\n      );\n\n      const events = result?.events ?? [];\n\n      const state = events.reduce((s, e) => evolve(s, e), initialState());\n\n      return {\n        currentStreamVersion: BigInt(events.length),\n        state,\n        streamExists: result.streamExists,\n      };\n    },\n\n    readStream: <\n      EventType extends Event,\n      EventPayloadType extends Event = EventType,\n    >(\n      streamName: string,\n      options?: ReadStreamOptions<\n        BigIntStreamPosition,\n        EventType,\n        EventPayloadType\n      >,\n    ): Promise<\n      ReadStreamResult<EventType, ReadEventMetadataWithGlobalPosition>\n    > => {\n      const events = streams.get(streamName);\n      const currentStreamVersion = events\n        ? BigInt(events.length)\n        : InMemoryEventStoreDefaultStreamVersion;\n\n      assertExpectedVersionMatchesCurrent(\n        currentStreamVersion,\n        options?.expectedStreamVersion,\n        InMemoryEventStoreDefaultStreamVersion,\n      );\n\n      const from = Number(options?.from ?? 0);\n      const to = Number(\n        options?.to ??\n          (options?.maxCount\n            ? (options.from ?? 0n) + options.maxCount\n            : (events?.length ?? 1)),\n      );\n\n      const resultEvents =\n        events !== undefined && events.length > 0\n          ? upcastRecordedMessages<\n              EventType,\n              EventPayloadType,\n              ReadEventMetadataWithGlobalPosition\n            >(\n              events.slice(from, to) as ReadEvent<\n                EventPayloadType,\n                ReadEventMetadataWithGlobalPosition\n              >[],\n              options?.schema?.versioning,\n            )\n          : [];\n\n      const result: ReadStreamResult<\n        EventType,\n        ReadEventMetadataWithGlobalPosition\n      > = {\n        currentStreamVersion,\n        events: resultEvents,\n        streamExists: events !== undefined && events.length > 0,\n      };\n\n      return Promise.resolve(result);\n    },\n\n    appendToStream: async <\n      EventType extends Event,\n      EventPayloadType extends Event = EventType,\n    >(\n      streamName: string,\n      events: EventType[],\n      options?: AppendToStreamOptions<\n        BigIntStreamPosition,\n        EventType,\n        EventPayloadType\n      >,\n    ): Promise<AppendToStreamResult> => {\n      const currentEvents = streams.get(streamName) ?? [];\n      const currentStreamVersion =\n        currentEvents.length > 0\n          ? BigInt(currentEvents.length)\n          : InMemoryEventStoreDefaultStreamVersion;\n\n      assertExpectedVersionMatchesCurrent(\n        currentStreamVersion,\n        options?.expectedStreamVersion,\n        InMemoryEventStoreDefaultStreamVersion,\n      );\n\n      const newEvents: ReadEvent<\n        EventType,\n        ReadEventMetadataWithGlobalPosition\n      >[] = events.map((event, index) => {\n        const metadata: ReadEventMetadataWithGlobalPosition = {\n          streamName,\n          messageId: uuid(),\n          streamPosition: BigInt(currentEvents.length + index + 1),\n          globalPosition: BigInt(getAllEventsCount() + index + 1),\n        };\n        return {\n          ...event,\n          kind: event.kind ?? 'Event',\n          metadata: {\n            ...('metadata' in event ? (event.metadata ?? {}) : {}),\n            ...metadata,\n          } as CombinedReadEventMetadata<\n            EventType,\n            ReadEventMetadataWithGlobalPosition\n          >,\n        };\n      });\n\n      const positionOfLastEventInTheStream = BigInt(\n        newEvents.slice(-1)[0]!.metadata.streamPosition,\n      );\n\n      streams.set(streamName, [\n        ...currentEvents,\n        ...downcastRecordedMessages(newEvents, options?.schema?.versioning),\n      ]);\n\n      // Process projections if there are any registered\n      if (inlineProjections.length > 0) {\n        await handleInMemoryProjections({\n          projections: inlineProjections,\n          events: newEvents,\n          database: eventStore.database,\n          eventStore,\n        });\n      }\n\n      const result: AppendToStreamResult = {\n        nextExpectedStreamVersion: positionOfLastEventInTheStream,\n        createdNewStream:\n          currentStreamVersion === InMemoryEventStoreDefaultStreamVersion,\n      };\n\n      await tryPublishMessagesAfterCommit<InMemoryEventStore>(\n        newEvents,\n        eventStoreOptions?.hooks,\n      );\n\n      return result;\n    },\n\n    streamExists: (streamName): Promise<StreamExistsResult> => {\n      const events = streams.get(streamName);\n\n      return Promise.resolve(events !== undefined && events.length > 0);\n    },\n  };\n\n  return eventStore;\n};\n","import { v7 as uuid } from 'uuid';\nimport { deepEquals } from '../utils';\nimport {\n  type DatabaseHandleOptionErrors,\n  type DatabaseHandleOptions,\n  type DatabaseHandleResult,\n  type DeleteResult,\n  type Document,\n  type DocumentHandler,\n  type InsertOneResult,\n  type OptionalUnlessRequiredIdAndVersion,\n  type ReplaceOneOptions,\n  type UpdateResult,\n  type WithIdAndVersion,\n  type WithoutId,\n} from './types';\nimport { expectedVersionValue, operationResult } from './utils';\n\nexport interface InMemoryDocumentsCollection<T extends Document> {\n  handle: (\n    id: string,\n    handle: DocumentHandler<T>,\n    options?: DatabaseHandleOptions,\n  ) => Promise<DatabaseHandleResult<T>>;\n  findOne: (predicate?: Predicate<T>) => Promise<T | null>;\n  find: (predicate?: Predicate<T>) => Promise<T[]>;\n  insertOne: (\n    document: OptionalUnlessRequiredIdAndVersion<T>,\n  ) => Promise<InsertOneResult>;\n  deleteOne: (predicate?: Predicate<T>) => Promise<DeleteResult>;\n  replaceOne: (\n    predicate: Predicate<T>,\n    document: WithoutId<T>,\n    options?: ReplaceOneOptions,\n  ) => Promise<UpdateResult>;\n}\n\nexport interface InMemoryDatabase {\n  collection: <T extends Document>(\n    name: string,\n  ) => InMemoryDocumentsCollection<T>;\n}\n\ntype Predicate<T> = (item: T) => boolean;\ntype CollectionName = string;\n\nexport const getInMemoryDatabase = (): InMemoryDatabase => {\n  const storage = new Map<CollectionName, WithIdAndVersion<Document>[]>();\n\n  return {\n    collection: <T extends Document, CollectionName extends string>(\n      collectionName: CollectionName,\n      collectionOptions: {\n        errors?: DatabaseHandleOptionErrors;\n      } = {},\n    ): InMemoryDocumentsCollection<T> => {\n      const ensureCollectionCreated = () => {\n        if (!storage.has(collectionName)) storage.set(collectionName, []);\n      };\n\n      const errors = collectionOptions.errors;\n\n      const collection = {\n        collectionName,\n        insertOne: async (\n          document: OptionalUnlessRequiredIdAndVersion<T>,\n        ): Promise<InsertOneResult> => {\n          ensureCollectionCreated();\n\n          const _id = (document._id as string | undefined | null) ?? uuid();\n          const _version = document._version ?? 1n;\n\n          const existing = await collection.findOne((c) => c._id === _id);\n\n          if (existing) {\n            return operationResult<InsertOneResult>(\n              {\n                successful: false,\n                insertedId: null,\n                nextExpectedVersion: _version,\n              },\n              { operationName: 'insertOne', collectionName, errors },\n            );\n          }\n\n          const documentsInCollection = storage.get(collectionName)!;\n          const newDocument = { ...document, _id, _version };\n          const newCollection = [...documentsInCollection, newDocument];\n          storage.set(collectionName, newCollection);\n\n          return operationResult<InsertOneResult>(\n            {\n              successful: true,\n              insertedId: _id,\n              nextExpectedVersion: _version,\n            },\n            { operationName: 'insertOne', collectionName, errors },\n          );\n        },\n        findOne: (predicate?: Predicate<T>): Promise<T | null> => {\n          ensureCollectionCreated();\n\n          const documentsInCollection = storage.get(collectionName);\n          const filteredDocuments = predicate\n            ? documentsInCollection?.filter((doc) => predicate(doc as T))\n            : documentsInCollection;\n\n          const firstOne = filteredDocuments?.[0] ?? null;\n\n          return Promise.resolve(firstOne as T | null);\n        },\n        find: (predicate?: Predicate<T>): Promise<T[]> => {\n          ensureCollectionCreated();\n\n          const documentsInCollection = storage.get(collectionName);\n          const filteredDocuments = predicate\n            ? documentsInCollection?.filter((doc) => predicate(doc as T))\n            : documentsInCollection;\n\n          return Promise.resolve(filteredDocuments as T[]);\n        },\n        deleteOne: (predicate?: Predicate<T>): Promise<DeleteResult> => {\n          ensureCollectionCreated();\n\n          const documentsInCollection = storage.get(collectionName)!;\n\n          if (predicate) {\n            const foundIndex = documentsInCollection.findIndex((doc) =>\n              predicate(doc as T),\n            );\n\n            if (foundIndex === -1) {\n              return Promise.resolve(\n                operationResult<DeleteResult>(\n                  {\n                    successful: false,\n                    matchedCount: 0,\n                    deletedCount: 0,\n                  },\n                  { operationName: 'deleteOne', collectionName, errors },\n                ),\n              );\n            } else {\n              const newCollection = documentsInCollection.toSpliced(\n                foundIndex,\n                1,\n              );\n\n              storage.set(collectionName, newCollection);\n\n              return Promise.resolve(\n                operationResult<DeleteResult>(\n                  {\n                    successful: true,\n                    matchedCount: 1,\n                    deletedCount: 1,\n                  },\n                  { operationName: 'deleteOne', collectionName, errors },\n                ),\n              );\n            }\n          }\n\n          const newCollection = documentsInCollection.slice(1);\n\n          storage.set(collectionName, newCollection);\n\n          return Promise.resolve(\n            operationResult<DeleteResult>(\n              {\n                successful: true,\n                matchedCount: 1,\n                deletedCount: 1,\n              },\n              { operationName: 'deleteOne', collectionName, errors },\n            ),\n          );\n        },\n        replaceOne: (\n          predicate: Predicate<T>,\n          document: WithoutId<T>,\n          options?: ReplaceOneOptions,\n        ): Promise<UpdateResult> => {\n          ensureCollectionCreated();\n\n          const documentsInCollection = storage.get(collectionName)!;\n\n          const firstIndex = documentsInCollection.findIndex((doc) =>\n            predicate(doc as T),\n          );\n\n          if (firstIndex === undefined || firstIndex === -1) {\n            return Promise.resolve(\n              operationResult<UpdateResult>(\n                {\n                  successful: false,\n                  matchedCount: 0,\n                  modifiedCount: 0,\n                  nextExpectedVersion: 0n,\n                },\n                { operationName: 'replaceOne', collectionName, errors },\n              ),\n            );\n          }\n\n          const existing = documentsInCollection[firstIndex]!;\n\n          if (\n            typeof options?.expectedVersion === 'bigint' &&\n            existing._version !== options.expectedVersion\n          ) {\n            return Promise.resolve(\n              operationResult<UpdateResult>(\n                {\n                  successful: false,\n                  matchedCount: 1,\n                  modifiedCount: 0,\n                  nextExpectedVersion: existing._version,\n                },\n                { operationName: 'replaceOne', collectionName, errors },\n              ),\n            );\n          }\n\n          const newVersion = existing._version + 1n;\n\n          const newCollection = documentsInCollection.with(firstIndex, {\n            _id: existing._id,\n            ...document,\n            _version: newVersion,\n          });\n\n          storage.set(collectionName, newCollection);\n\n          return Promise.resolve(\n            operationResult<UpdateResult>(\n              {\n                successful: true,\n                modifiedCount: 1,\n                matchedCount: firstIndex,\n                nextExpectedVersion: newVersion,\n              },\n              { operationName: 'replaceOne', collectionName, errors },\n            ),\n          );\n        },\n        handle: async (\n          id: string,\n          handle: DocumentHandler<T>,\n          options?: DatabaseHandleOptions,\n        ): Promise<DatabaseHandleResult<T>> => {\n          const { expectedVersion: version, ...operationOptions } =\n            options ?? {};\n          ensureCollectionCreated();\n          const existing = await collection.findOne(({ _id }) => _id === id);\n\n          const expectedVersion = expectedVersionValue(version);\n\n          if (\n            (existing == null && version === 'DOCUMENT_EXISTS') ||\n            (existing == null && expectedVersion != null) ||\n            (existing != null && version === 'DOCUMENT_DOES_NOT_EXIST') ||\n            (existing != null &&\n              expectedVersion !== null &&\n              existing._version !== expectedVersion)\n          ) {\n            return operationResult<DatabaseHandleResult<T>>(\n              {\n                successful: false,\n                document: existing as WithIdAndVersion<T>,\n              },\n              { operationName: 'handle', collectionName, errors },\n            );\n          }\n\n          const result = handle(existing !== null ? { ...existing } : null);\n\n          if (deepEquals(existing, result))\n            return operationResult<DatabaseHandleResult<T>>(\n              {\n                successful: true,\n                document: existing as WithIdAndVersion<T>,\n              },\n              { operationName: 'handle', collectionName, errors },\n            );\n\n          if (!existing && result) {\n            const newDoc = { ...result, _id: id };\n            const insertResult = await collection.insertOne({\n              ...newDoc,\n              _id: id,\n            } as OptionalUnlessRequiredIdAndVersion<T>);\n            return {\n              ...insertResult,\n              document: {\n                ...newDoc,\n                _version: insertResult.nextExpectedVersion,\n              } as unknown as WithIdAndVersion<T>,\n            };\n          }\n\n          if (existing && !result) {\n            const deleteResult = await collection.deleteOne(\n              ({ _id }) => id === _id,\n            );\n            return { ...deleteResult, document: null };\n          }\n\n          if (existing && result) {\n            const replaceResult = await collection.replaceOne(\n              ({ _id }) => id === _id,\n              result,\n              {\n                ...operationOptions,\n                expectedVersion: expectedVersion ?? 'DOCUMENT_EXISTS',\n              },\n            );\n            return {\n              ...replaceResult,\n              document: {\n                ...result,\n                _version: replaceResult.nextExpectedVersion,\n              } as unknown as WithIdAndVersion<T>,\n            };\n          }\n\n          return operationResult<DatabaseHandleResult<T>>(\n            {\n              successful: true,\n              document: existing as WithIdAndVersion<T>,\n            },\n            { operationName: 'handle', collectionName, errors },\n          );\n        },\n      };\n\n      return collection;\n    },\n  };\n};\n","export const hasDuplicates = <ArrayItem, Mapped>(\n  array: ArrayItem[],\n  predicate: (value: ArrayItem, index: number, array: ArrayItem[]) => Mapped,\n) => {\n  const mapped = array.map(predicate);\n  const uniqueValues = new Set(mapped);\n\n  return uniqueValues.size < mapped.length;\n};\n\nexport const getDuplicates = <ArrayItem, Mapped>(\n  array: ArrayItem[],\n  predicate: (value: ArrayItem, index: number, array: ArrayItem[]) => Mapped,\n): ArrayItem[] => {\n  const map = new Map<Mapped, ArrayItem[]>();\n\n  for (let i = 0; i < array.length; i++) {\n    const item = array[i]!;\n    const key = predicate(item, i, array);\n    if (!map.has(key)) {\n      map.set(key, []);\n    }\n    map.get(key)!.push(item);\n  }\n\n  return Array.from(map.values())\n    .filter((group) => group.length > 1)\n    .flat();\n};\n","export const merge = <T>(\n  array: T[],\n  item: T,\n  where: (current: T) => boolean,\n  onExisting: (current: T) => T,\n  onNotFound: () => T | undefined = () => undefined,\n) => {\n  let wasFound = false;\n\n  const result = array\n    // merge the existing item if matches condition\n    .map((p: T) => {\n      if (!where(p)) return p;\n\n      wasFound = true;\n      return onExisting(p);\n    })\n    // filter out item if undefined was returned\n    // for cases of removal\n    .filter((p) => p !== undefined)\n    // make TypeScript happy\n    .map((p) => {\n      if (!p) throw Error('That should not happen');\n\n      return p;\n    });\n\n  // if item was not found and onNotFound action is defined\n  // try to generate new item\n  if (!wasFound) {\n    const result = onNotFound();\n\n    if (result !== undefined) return [...array, item];\n  }\n\n  return result;\n};\n","import { getDuplicates, hasDuplicates } from './duplicates';\nimport { merge } from './merge';\n\nexport * from './merge';\n\nexport const arrayUtils = {\n  merge,\n  hasDuplicates,\n  getDuplicates,\n};\n","const isPrimitive = (value: unknown): boolean => {\n  const type = typeof value;\n  return (\n    value === null ||\n    value === undefined ||\n    type === 'boolean' ||\n    type === 'number' ||\n    type === 'string' ||\n    type === 'symbol' ||\n    type === 'bigint'\n  );\n};\n\nconst compareArrays = <T>(left: T[], right: T[]): boolean => {\n  if (left.length !== right.length) {\n    return false;\n  }\n  for (let i = 0; i < left.length; i++) {\n    const leftHas = i in left;\n    const rightHas = i in right;\n    if (leftHas !== rightHas) return false;\n    if (leftHas && !deepEquals(left[i], right[i])) return false;\n  }\n  return true;\n};\n\nconst compareDates = (left: Date, right: Date): boolean => {\n  return left.getTime() === right.getTime();\n};\n\nconst compareRegExps = (left: RegExp, right: RegExp): boolean => {\n  return left.toString() === right.toString();\n};\n\nconst compareErrors = (left: Error, right: Error): boolean => {\n  if (left.message !== right.message || left.name !== right.name) {\n    return false;\n  }\n  const leftKeys = Object.keys(left);\n  const rightKeys = Object.keys(right);\n  if (leftKeys.length !== rightKeys.length) return false;\n  const rightKeySet = new Set(rightKeys);\n  for (const key of leftKeys) {\n    if (!rightKeySet.has(key)) return false;\n    // @ts-expect-error - accessing dynamic keys\n    if (!deepEquals(left[key], right[key])) return false;\n  }\n  return true;\n};\n\nconst compareMaps = (\n  left: Map<unknown, unknown>,\n  right: Map<unknown, unknown>,\n): boolean => {\n  if (left.size !== right.size) return false;\n\n  for (const [key, value] of left) {\n    if (isPrimitive(key)) {\n      if (!right.has(key) || !deepEquals(value, right.get(key))) {\n        return false;\n      }\n    } else {\n      let found = false;\n      for (const [rightKey, rightValue] of right) {\n        if (deepEquals(key, rightKey) && deepEquals(value, rightValue)) {\n          found = true;\n          break;\n        }\n      }\n      if (!found) return false;\n    }\n  }\n  return true;\n};\n\nconst compareSets = (left: Set<unknown>, right: Set<unknown>): boolean => {\n  if (left.size !== right.size) return false;\n\n  for (const leftItem of left) {\n    if (isPrimitive(leftItem)) {\n      if (!right.has(leftItem)) return false;\n    } else {\n      let found = false;\n      for (const rightItem of right) {\n        if (deepEquals(leftItem, rightItem)) {\n          found = true;\n          break;\n        }\n      }\n      if (!found) return false;\n    }\n  }\n  return true;\n};\n\nconst compareArrayBuffers = (\n  left: ArrayBuffer,\n  right: ArrayBuffer,\n): boolean => {\n  if (left.byteLength !== right.byteLength) return false;\n  const leftView = new Uint8Array(left);\n  const rightView = new Uint8Array(right);\n  for (let i = 0; i < leftView.length; i++) {\n    if (leftView[i] !== rightView[i]) return false;\n  }\n  return true;\n};\n\nconst compareTypedArrays = (\n  left: ArrayBufferView,\n  right: ArrayBufferView,\n): boolean => {\n  if (left.constructor !== right.constructor) return false;\n  if (left.byteLength !== right.byteLength) return false;\n\n  const leftArray = new Uint8Array(\n    left.buffer,\n    left.byteOffset,\n    left.byteLength,\n  );\n  const rightArray = new Uint8Array(\n    right.buffer,\n    right.byteOffset,\n    right.byteLength,\n  );\n\n  for (let i = 0; i < leftArray.length; i++) {\n    if (leftArray[i] !== rightArray[i]) return false;\n  }\n  return true;\n};\n\nconst compareObjects = (\n  left: Record<string, unknown>,\n  right: Record<string, unknown>,\n): boolean => {\n  const keys1 = Object.keys(left);\n  const keys2 = Object.keys(right);\n\n  if (keys1.length !== keys2.length) {\n    return false;\n  }\n\n  for (const key of keys1) {\n    if (left[key] instanceof Function && right[key] instanceof Function) {\n      continue;\n    }\n\n    const isEqual = deepEquals(left[key], right[key]);\n    if (!isEqual) {\n      return false;\n    }\n  }\n\n  return true;\n};\n\nconst getType = (value: unknown): string => {\n  if (value === null) return 'null';\n  if (value === undefined) return 'undefined';\n\n  const primitiveType = typeof value;\n  if (primitiveType !== 'object') return primitiveType;\n\n  if (Array.isArray(value)) return 'array';\n  if (value instanceof Boolean) return 'boxed-boolean';\n  if (value instanceof Number) return 'boxed-number';\n  if (value instanceof String) return 'boxed-string';\n  if (value instanceof Date) return 'date';\n  if (value instanceof RegExp) return 'regexp';\n  if (value instanceof Error) return 'error';\n  if (value instanceof Map) return 'map';\n  if (value instanceof Set) return 'set';\n  if (value instanceof ArrayBuffer) return 'arraybuffer';\n  if (value instanceof DataView) return 'dataview';\n  if (value instanceof WeakMap) return 'weakmap';\n  if (value instanceof WeakSet) return 'weakset';\n\n  if (ArrayBuffer.isView(value)) return 'typedarray';\n\n  return 'object';\n};\n\nexport const deepEquals = <T>(left: T, right: T): boolean => {\n  if (left === right) return true;\n\n  if (isEquatable(left)) {\n    return left.equals(right);\n  }\n\n  const leftType = getType(left);\n  const rightType = getType(right);\n\n  if (leftType !== rightType) return false;\n\n  switch (leftType) {\n    case 'null':\n    case 'undefined':\n    case 'boolean':\n    case 'number':\n    case 'bigint':\n    case 'string':\n    case 'symbol':\n    case 'function':\n      return left === right;\n\n    case 'array':\n      return compareArrays(left as unknown[], right as unknown[]);\n\n    case 'date':\n      return compareDates(left as Date, right as Date);\n\n    case 'regexp':\n      return compareRegExps(left as RegExp, right as RegExp);\n\n    case 'error':\n      return compareErrors(left as Error, right as Error);\n\n    case 'map':\n      return compareMaps(\n        left as Map<unknown, unknown>,\n        right as Map<unknown, unknown>,\n      );\n\n    case 'set':\n      return compareSets(left as Set<unknown>, right as Set<unknown>);\n\n    case 'arraybuffer':\n      return compareArrayBuffers(left as ArrayBuffer, right as ArrayBuffer);\n\n    case 'dataview':\n    case 'weakmap':\n    case 'weakset':\n      return false;\n\n    case 'typedarray':\n      return compareTypedArrays(\n        left as ArrayBufferView,\n        right as ArrayBufferView,\n      );\n\n    case 'boxed-boolean':\n      return (left as boolean).valueOf() === (right as boolean).valueOf();\n\n    case 'boxed-number':\n      return (left as number).valueOf() === (right as number).valueOf();\n\n    case 'boxed-string':\n      return (left as string).valueOf() === (right as string).valueOf();\n\n    case 'object':\n      return compareObjects(\n        left as Record<string, unknown>,\n        right as Record<string, unknown>,\n      );\n\n    default:\n      return false;\n  }\n};\n\nexport type Equatable<T> = { equals: (right: T) => boolean } & T;\n\nexport const isEquatable = <T>(left: T): left is Equatable<T> => {\n  return (\n    left !== null &&\n    left !== undefined &&\n    typeof left === 'object' &&\n    'equals' in left &&\n    typeof left['equals'] === 'function'\n  );\n};\n","export const sum = (\n  iterator: Iterator<number, number, number> | Iterator<number>,\n) => {\n  let value,\n    done: boolean | undefined,\n    sum = 0;\n  do {\n    // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment\n    ({ value, done } = iterator.next());\n    sum += value || 0;\n  } while (!done);\n  return sum;\n};\n","import { EmmettError } from '../errors';\n\nexport type TaskQueue = TaskQueueItem[];\n\nexport type TaskQueueItem = {\n  task: () => Promise<void>;\n  options?: EnqueueTaskOptions;\n};\n\nexport type TaskProcessorOptions = {\n  maxActiveTasks: number;\n  maxQueueSize: number;\n  maxTaskIdleTime?: number;\n};\n\nexport type Task<T> = (context: TaskContext) => Promise<T>;\n\nexport type TaskContext = {\n  ack: () => void;\n};\n\nexport type EnqueueTaskOptions = { taskGroupId?: string };\n\nexport class TaskProcessor {\n  private queue: TaskQueue = [];\n  private isProcessing = false;\n  private activeTasks = 0;\n  private activeGroups: Set<string> = new Set();\n\n  constructor(private options: TaskProcessorOptions) {}\n\n  enqueue<T>(task: Task<T>, options?: EnqueueTaskOptions): Promise<T> {\n    if (this.queue.length >= this.options.maxQueueSize) {\n      return Promise.reject(\n        new EmmettError(\n          'Too many pending connections. Please try again later.',\n        ),\n      );\n    }\n\n    return this.schedule(task, options);\n  }\n\n  waitForEndOfProcessing(): Promise<void> {\n    return this.schedule(({ ack }) => Promise.resolve(ack()));\n  }\n\n  private schedule<T>(task: Task<T>, options?: EnqueueTaskOptions): Promise<T> {\n    return promiseWithDeadline(\n      (resolve, reject) => {\n        const taskWithContext = () => {\n          return new Promise<void>((resolveTask, failTask) => {\n            const taskPromise = task({\n              ack: resolveTask,\n            });\n\n            taskPromise.then(resolve).catch((err) => {\n              // eslint-disable-next-line @typescript-eslint/prefer-promise-reject-errors\n              failTask(err);\n              reject(err);\n            });\n          });\n        };\n\n        this.queue.push({ task: taskWithContext, options });\n        if (!this.isProcessing) {\n          this.ensureProcessing();\n        }\n      },\n      { deadline: this.options.maxTaskIdleTime },\n    );\n  }\n\n  private ensureProcessing(): void {\n    if (this.isProcessing) return;\n    this.isProcessing = true;\n    this.processQueue();\n  }\n\n  private processQueue(): void {\n    try {\n      while (\n        this.activeTasks < this.options.maxActiveTasks &&\n        this.queue.length > 0\n      ) {\n        const item = this.takeFirstAvailableItem();\n\n        if (item === null) return;\n\n        const groupId = item.options?.taskGroupId;\n\n        if (groupId) {\n          // Mark the group as active\n          this.activeGroups.add(groupId);\n        }\n\n        this.activeTasks++;\n        void this.executeItem(item);\n      }\n    } catch (error) {\n      console.error(error);\n      throw error;\n    } finally {\n      this.isProcessing = false;\n      if (\n        this.hasItemsToProcess() &&\n        this.activeTasks < this.options.maxActiveTasks\n      ) {\n        this.ensureProcessing();\n      }\n    }\n  }\n\n  private async executeItem({ task, options }: TaskQueueItem): Promise<void> {\n    try {\n      await task();\n    } finally {\n      this.activeTasks--;\n\n      // Mark the group as inactive after task completion\n      if (options && options.taskGroupId) {\n        this.activeGroups.delete(options.taskGroupId);\n      }\n\n      this.ensureProcessing();\n    }\n  }\n\n  private takeFirstAvailableItem = (): TaskQueueItem | null => {\n    const taskIndex = this.queue.findIndex(\n      (item) =>\n        !item.options?.taskGroupId ||\n        !this.activeGroups.has(item.options.taskGroupId),\n    );\n\n    if (taskIndex === -1) {\n      // All remaining tasks are blocked by active groups\n      return null;\n    }\n\n    // Remove the task from the queue\n    const [item] = this.queue.splice(taskIndex, 1);\n\n    return item ?? null;\n  };\n\n  private hasItemsToProcess = (): boolean =>\n    this.queue.findIndex(\n      (item) =>\n        !item.options?.taskGroupId ||\n        !this.activeGroups.has(item.options.taskGroupId),\n    ) !== -1;\n}\n\nconst DEFAULT_PROMISE_DEADLINE = 2147483647;\n\nconst promiseWithDeadline = <T>(\n  executor: (\n    resolve: (value: T | PromiseLike<T>) => void,\n    reject: (reason?: unknown) => void,\n  ) => void,\n  options: { deadline?: number },\n) => {\n  return new Promise<T>((resolve, reject) => {\n    let taskStarted = false;\n\n    const maxWaitingTime = options.deadline || DEFAULT_PROMISE_DEADLINE;\n\n    let timeoutId: NodeJS.Timeout | null = setTimeout(() => {\n      if (!taskStarted) {\n        reject(\n          new Error('Task was not started within the maximum waiting time'),\n        );\n      }\n    }, maxWaitingTime);\n\n    executor((value) => {\n      taskStarted = true;\n      if (timeoutId) {\n        clearTimeout(timeoutId);\n      }\n      timeoutId = null;\n      resolve(value);\n    }, reject);\n  });\n};\n","import { TaskProcessor } from '../../taskProcessing';\n\nexport type LockOptions = { lockId: number };\n\nexport type AcquireLockOptions = { lockId: string };\nexport type ReleaseLockOptions = { lockId: string };\n\nexport type Lock = {\n  acquire(options: AcquireLockOptions): Promise<void>;\n  tryAcquire(options: AcquireLockOptions): Promise<boolean>;\n  release(options: ReleaseLockOptions): Promise<boolean>;\n  withAcquire: <Result = unknown>(\n    handle: () => Promise<Result>,\n    options: AcquireLockOptions,\n  ) => Promise<Result>;\n};\n\nexport const InProcessLock = (): Lock => {\n  const taskProcessor = new TaskProcessor({\n    maxActiveTasks: Number.MAX_VALUE,\n    maxQueueSize: Number.MAX_VALUE,\n  });\n\n  // Map to store ack functions of currently held locks: lockId -> ack()\n  const locks = new Map<string, () => void>();\n\n  return {\n    async acquire({ lockId }: AcquireLockOptions): Promise<void> {\n      // If the lock is already held, we just queue up another task in the same group.\n      // TaskProcessor ensures tasks in the same group run one at a time.\n      await new Promise<void>((resolve, reject) => {\n        taskProcessor\n          .enqueue(\n            ({ ack }) => {\n              // When this task starts, it means the previous lock (if any) was released\n              // and now we have exclusive access.\n              locks.set(lockId, ack);\n              // We do NOT call ack() here. We hold onto the lock.\n              resolve();\n              return Promise.resolve();\n            },\n            { taskGroupId: lockId },\n          )\n          .catch(reject);\n      });\n    },\n\n    async tryAcquire({ lockId }: AcquireLockOptions): Promise<boolean> {\n      // If lock is already held, fail immediately\n      if (locks.has(lockId)) {\n        return false;\n      }\n\n      // TODO: Check pending queue\n      await this.acquire({ lockId });\n\n      return true;\n    },\n\n    release({ lockId }: ReleaseLockOptions): Promise<boolean> {\n      const ack = locks.get(lockId);\n      if (ack === undefined) {\n        return Promise.resolve(true);\n      }\n      locks.delete(lockId);\n      ack();\n      return Promise.resolve(true);\n    },\n\n    async withAcquire<Result = unknown>(\n      handle: () => Promise<Result>,\n      { lockId }: AcquireLockOptions,\n    ): Promise<Result> {\n      return taskProcessor.enqueue(\n        async ({ ack }) => {\n          // When this task starts, it means the previous lock (if any) was released\n          // and now we have exclusive access.\n          locks.set(lockId, ack);\n\n          // We do NOT call ack() here. We hold onto the lock.\n          try {\n            return await handle();\n          } finally {\n            locks.delete(lockId);\n            ack();\n          }\n        },\n        { taskGroupId: lockId },\n      );\n    },\n  };\n};\n","export const toNormalizedString = (value: bigint): string =>\n  value.toString().padStart(19, '0');\n\nexport const bigInt = {\n  toNormalizedString,\n};\n","export const delay = (ms: number): Promise<void> => {\n  return new Promise((resolve) => setTimeout(resolve, ms));\n};\n\nexport type AsyncAwaiter<T = void> = {\n  wait: Promise<T>;\n  resolve: (value: T | PromiseLike<T>) => void;\n  // eslint-disable-next-line @typescript-eslint/no-explicit-any\n  reject: (reason?: any) => void;\n  reset: () => void;\n};\n\n// TODO: Remove this after migrating to Node 22\nexport const asyncAwaiter = <T = void>(): AsyncAwaiter<T> => {\n  const result: AsyncAwaiter<T> = {} as AsyncAwaiter<T>;\n\n  (result.reset = () => {\n    result.wait = new Promise<T>((res, rej) => {\n      result.resolve = res;\n      result.reject = rej;\n    });\n  })();\n\n  return result;\n};\n","import retry from 'async-retry';\nimport { EmmettError } from '../errors';\nimport { JSONParser } from '../serialization';\n\nexport type AsyncRetryOptions<T = unknown> = retry.Options & {\n  shouldRetryResult?: (result: T) => boolean;\n  shouldRetryError?: (error?: unknown) => boolean;\n};\n\nexport const NoRetries: AsyncRetryOptions = { retries: 0 };\n\nexport const asyncRetry = async <T>(\n  fn: () => Promise<T>,\n  opts?: AsyncRetryOptions<T>,\n): Promise<T> => {\n  if (opts === undefined || opts.retries === 0) return fn();\n\n  return retry(\n    async (bail) => {\n      try {\n        const result = await fn();\n\n        if (opts?.shouldRetryResult && opts.shouldRetryResult(result)) {\n          throw new EmmettError(\n            `Retrying because of result: ${JSONParser.stringify(result)}`,\n          );\n        }\n        return result;\n      } catch (error) {\n        if (opts?.shouldRetryError && !opts.shouldRetryError(error)) {\n          bail(error as Error);\n          return undefined as unknown as T;\n        }\n        throw error;\n      }\n    },\n    opts ?? { retries: 0 },\n  );\n};\n","export class ParseError extends Error {\n  constructor(text: string) {\n    super(`Cannot parse! ${text}`);\n  }\n}\n\nexport type Mapper<From, To = From> =\n  | ((value: unknown) => To)\n  | ((value: Partial<From>) => To)\n  | ((value: From) => To)\n  | ((value: Partial<To>) => To)\n  | ((value: To) => To)\n  | ((value: Partial<To | From>) => To)\n  | ((value: To | From) => To);\n\nexport type MapperArgs<From, To = From> = Partial<From> &\n  From &\n  Partial<To> &\n  To;\n\nexport type ParseOptions<From, To = From> = {\n  reviver?: (key: string, value: unknown) => unknown;\n  map?: Mapper<From, To>;\n  typeCheck?: <To>(value: unknown) => value is To;\n};\n\nexport type StringifyOptions<From, To = From> = {\n  map?: Mapper<From, To>;\n};\n\nexport const JSONParser = {\n  stringify: <From, To = From>(\n    value: From,\n    options?: StringifyOptions<From, To>,\n  ) => {\n    return JSON.stringify(\n      options?.map ? options.map(value as MapperArgs<From, To>) : value,\n      //TODO: Consider adding support to DateTime and adding specific format to mark that's a bigint\n      // eslint-disable-next-line @typescript-eslint/no-unsafe-return\n      (_, v) => (typeof v === 'bigint' ? v.toString() : v),\n    );\n  },\n  parse: <From, To = From>(\n    text: string,\n    options?: ParseOptions<From, To>,\n  ): To | undefined => {\n    const parsed: unknown = JSON.parse(text, options?.reviver);\n\n    if (options?.typeCheck && !options?.typeCheck<To>(parsed))\n      throw new ParseError(text);\n\n    return options?.map\n      ? options.map(parsed as MapperArgs<From, To>)\n      : (parsed as To | undefined);\n  },\n};\n","export type ShutdownHandler = () => void | Promise<void>;\n\n/**\n * Registers handlers for OS signals to enable graceful shutdown.\n * Handles SIGTERM and SIGINT by default.\n * Works in Node.js, Bun, and Deno. Safely no-ops in Browser/Cloudflare Workers.\n *\n * @param handler - Function to call when shutdown signal is received\n * @returns Cleanup function to unregister the handlers\n */\nexport const onShutdown = (handler: ShutdownHandler): (() => void) => {\n  const signals = ['SIGTERM', 'SIGINT'] as const;\n\n  // Node.js/Bun\n  if (typeof process !== 'undefined' && typeof process.on === 'function') {\n    for (const signal of signals) {\n      process.on(signal, handler);\n    }\n    return () => {\n      for (const signal of signals) {\n        process.off(signal, handler);\n      }\n    };\n  }\n\n  // Deno\n  // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment, @typescript-eslint/no-explicit-any, @typescript-eslint/no-unsafe-member-access\n  const deno = (globalThis as any).Deno;\n  // eslint-disable-next-line @typescript-eslint/no-unsafe-member-access\n  if (deno && typeof deno.addSignalListener === 'function') {\n    for (const signal of signals) {\n      // eslint-disable-next-line @typescript-eslint/no-unsafe-call, @typescript-eslint/no-unsafe-member-access\n      deno.addSignalListener(signal, handler);\n    }\n    return () => {\n      for (const signal of signals) {\n        // eslint-disable-next-line @typescript-eslint/no-unsafe-call, @typescript-eslint/no-unsafe-member-access\n        deno.removeSignalListener(signal, handler);\n      }\n    };\n  }\n\n  // Browser/Cloudflare Workers: no-op\n  return () => {};\n};\n","const textEncoder = new TextEncoder();\n\nexport const hashText = async (text: string): Promise<bigint> => {\n  const hashBuffer = await crypto.subtle.digest(\n    'SHA-256',\n    textEncoder.encode(text),\n  );\n  // Create an array with a single element that is a 64-bit signed integer\n  // We take the first 8 bytes (so 64 bits) of the SHA-256 hash\n  const view = new BigInt64Array(hashBuffer, 0, 1);\n  return view[0]!;\n};\n","import { ConcurrencyInMemoryDatabaseError } from '../errors';\nimport { JSONParser } from '../serialization';\nimport type {\n  DatabaseHandleOptionErrors,\n  ExpectedDocumentVersion,\n  ExpectedDocumentVersionGeneral,\n  ExpectedDocumentVersionValue,\n  OperationResult,\n} from './types';\n\nexport const isGeneralExpectedDocumentVersion = (\n  version: ExpectedDocumentVersion | undefined,\n): version is ExpectedDocumentVersionGeneral => {\n  return (\n    version === 'DOCUMENT_DOES_NOT_EXIST' ||\n    version === 'DOCUMENT_EXISTS' ||\n    version === 'NO_CONCURRENCY_CHECK'\n  );\n};\n\nexport const expectedVersionValue = (\n  version: ExpectedDocumentVersion | undefined,\n): ExpectedDocumentVersionValue | null =>\n  version === undefined || isGeneralExpectedDocumentVersion(version)\n    ? null\n    : version;\n\nexport const operationResult = <T extends OperationResult>(\n  result: Omit<T, 'assertSuccess' | 'acknowledged' | 'assertSuccessful'>,\n  options: {\n    operationName: string;\n    collectionName: string;\n    errors?: DatabaseHandleOptionErrors;\n  },\n): T => {\n  const operationResult: T = {\n    ...result,\n    acknowledged: true,\n    successful: result.successful,\n    assertSuccessful: (errorMessage?: string) => {\n      const { successful } = result;\n      const { operationName, collectionName } = options;\n\n      if (!successful)\n        throw new ConcurrencyInMemoryDatabaseError(\n          errorMessage ??\n            `${operationName} on ${collectionName} failed. Expected document state does not match current one! Result: ${JSONParser.stringify(result)}!`,\n        );\n    },\n  } as T;\n\n  if (options.errors?.throwOnOperationFailures)\n    operationResult.assertSuccessful();\n\n  return operationResult;\n};\n","import type { InMemoryDatabase } from '../../../database/inMemoryDatabase';\nimport type {\n  ProjectionDefinition,\n  TruncateProjection,\n} from '../../../projections';\nimport type { CanHandle, Event, ReadEvent } from '../../../typing';\nimport {\n  type InMemoryProjectionHandlerContext,\n  type InMemoryReadEventMetadata,\n} from '../../inMemoryEventStore';\n\nexport const DATABASE_REQUIRED_ERROR_MESSAGE =\n  'Database is required in context for InMemory projections';\n\nexport type InMemoryProjectionDefinition<EventType extends Event> =\n  ProjectionDefinition<\n    EventType,\n    InMemoryReadEventMetadata,\n    InMemoryProjectionHandlerContext\n  >;\n\nexport type InMemoryProjectionHandlerOptions<EventType extends Event = Event> =\n  {\n    projections: InMemoryProjectionDefinition<EventType>[];\n    events: ReadEvent<EventType, InMemoryReadEventMetadata>[];\n    database: InMemoryDatabase;\n    eventStore?: InMemoryProjectionHandlerContext['eventStore'];\n  };\n\n/**\n * Handles projections for the InMemoryEventStore\n * Similar to the PostgreSQL implementation, this processes events through projections\n */\nexport const handleInMemoryProjections = async <\n  EventType extends Event = Event,\n>(\n  options: InMemoryProjectionHandlerOptions<EventType>,\n): Promise<void> => {\n  const { projections, events, database, eventStore } = options;\n\n  // Get all event types from the events batch to filter projections\n  const eventTypes = events.map((e) => e.type);\n\n  // Filter projections that can handle these event types\n  const relevantProjections = projections.filter((p) =>\n    p.canHandle.some((type) => eventTypes.includes(type)),\n  );\n\n  // Process each projection\n  for (const projection of relevantProjections) {\n    await projection.handle(events, {\n      eventStore,\n      database,\n    });\n  }\n};\n\nexport type InMemoryWithNotNullDocumentEvolve<\n  DocumentType extends Record<string, unknown>,\n  EventType extends Event,\n> = (\n  document: DocumentType,\n  event: ReadEvent<EventType, InMemoryReadEventMetadata>,\n) => DocumentType | null;\n\nexport type InMemoryWithNullableDocumentEvolve<\n  DocumentType extends Record<string, unknown>,\n  EventType extends Event,\n> = (\n  document: DocumentType | null,\n  event: ReadEvent<EventType, InMemoryReadEventMetadata>,\n) => DocumentType | null;\n\nexport type InMemoryDocumentEvolve<\n  DocumentType extends Record<string, unknown>,\n  EventType extends Event,\n> =\n  | InMemoryWithNotNullDocumentEvolve<DocumentType, EventType>\n  | InMemoryWithNullableDocumentEvolve<DocumentType, EventType>;\n\nexport type InMemoryProjectionOptions<EventType extends Event> = {\n  handle: (\n    events: ReadEvent<EventType, InMemoryReadEventMetadata>[],\n    context: InMemoryProjectionHandlerContext & { database: InMemoryDatabase },\n  ) => Promise<void>;\n  canHandle: CanHandle<EventType>;\n  truncate?: TruncateProjection<\n    InMemoryProjectionHandlerContext & { database: InMemoryDatabase }\n  >;\n};\n\n/**\n * Creates an InMemory projection\n */\nexport const inMemoryProjection = <EventType extends Event>({\n  truncate,\n  handle,\n  canHandle,\n}: InMemoryProjectionOptions<EventType>): InMemoryProjectionDefinition<EventType> => ({\n  canHandle,\n  handle: async (events, context) => {\n    if (!context.database) {\n      throw new Error(DATABASE_REQUIRED_ERROR_MESSAGE);\n    }\n    await handle(events, {\n      ...context,\n      database: context.database,\n    });\n  },\n  truncate: truncate\n    ? (context) => {\n        if (!context.database) {\n          throw new Error(DATABASE_REQUIRED_ERROR_MESSAGE);\n        }\n        return truncate({\n          ...context,\n          database: context.database,\n        });\n      }\n    : undefined,\n});\n\n/**\n * Creates a multi-stream projection for InMemoryDatabase\n */\nexport type InMemoryMultiStreamProjectionOptions<\n  DocumentType extends Record<string, unknown>,\n  EventType extends Event,\n> = {\n  canHandle: CanHandle<EventType>;\n  collectionName: string;\n  getDocumentId: (event: ReadEvent<EventType>) => string;\n} & (\n  | {\n      evolve: InMemoryWithNullableDocumentEvolve<DocumentType, EventType>;\n    }\n  | {\n      evolve: InMemoryWithNotNullDocumentEvolve<DocumentType, EventType>;\n      initialState: () => DocumentType;\n    }\n);\n\n/**\n * Creates a projection that handles events across multiple streams\n */\nexport const inMemoryMultiStreamProjection = <\n  DocumentType extends Record<string, unknown>,\n  EventType extends Event,\n>(\n  options: InMemoryMultiStreamProjectionOptions<DocumentType, EventType>,\n): InMemoryProjectionDefinition<EventType> => {\n  const { collectionName, getDocumentId, canHandle } = options;\n\n  return inMemoryProjection({\n    handle: async (\n      events: ReadEvent<EventType, InMemoryReadEventMetadata>[],\n      { database },\n    ) => {\n      const collection = database.collection<DocumentType>(collectionName);\n\n      for (const event of events) {\n        await collection.handle(getDocumentId(event), (document) => {\n          if ('initialState' in options) {\n            return options.evolve(document ?? options.initialState(), event);\n          } else {\n            return options.evolve(document, event);\n          }\n        });\n      }\n    },\n    canHandle,\n    truncate: async ({\n      database,\n    }: InMemoryProjectionHandlerContext & { database: InMemoryDatabase }) => {\n      // For InMemory database, we can't directly truncate a collection\n      // So we'll delete all documents from the collection\n      const collection = database.collection<DocumentType>(collectionName);\n      const documents = await collection.find();\n\n      for (const doc of documents) {\n        if (doc && '_id' in doc) {\n          const id = doc._id;\n          await collection.deleteOne((d) => d._id === id);\n        }\n      }\n    },\n  });\n};\n\n/**\n * Creates a single-stream projection for InMemoryDatabase\n */\nexport type InMemorySingleStreamProjectionOptions<\n  DocumentType extends Record<string, unknown>,\n  EventType extends Event,\n> = {\n  canHandle: CanHandle<EventType>;\n  getDocumentId?: (event: ReadEvent<EventType>) => string;\n  collectionName: string;\n} & (\n  | {\n      evolve: InMemoryWithNullableDocumentEvolve<DocumentType, EventType>;\n    }\n  | {\n      evolve: InMemoryWithNotNullDocumentEvolve<DocumentType, EventType>;\n      initialState: () => DocumentType;\n    }\n);\n\n/**\n * Creates a projection that handles events from a single stream\n */\nexport const inMemorySingleStreamProjection = <\n  DocumentType extends Record<string, unknown>,\n  EventType extends Event,\n>(\n  options: InMemorySingleStreamProjectionOptions<DocumentType, EventType>,\n): InMemoryProjectionDefinition<EventType> => {\n  return inMemoryMultiStreamProjection<DocumentType, EventType>({\n    ...options,\n    getDocumentId:\n      options.getDocumentId ?? ((event) => event.metadata.streamName),\n  });\n};\n","import { v4 as uuid } from 'uuid';\nimport {\n  handleInMemoryProjections,\n  type InMemoryProjectionDefinition,\n} from '.';\nimport {\n  getInMemoryDatabase,\n  type Document,\n  type InMemoryDatabase,\n} from '../../../database';\nimport { isErrorConstructor } from '../../../errors';\nimport { JSONParser } from '../../../serialization';\nimport {\n  assertFails,\n  AssertionError,\n  assertTrue,\n  type ThenThrows,\n} from '../../../testing';\nimport type { CombinedReadEventMetadata, ReadEvent } from '../../../typing';\nimport { type Event } from '../../../typing';\nimport type {\n  InMemoryEventStore,\n  InMemoryReadEventMetadata,\n} from '../../inMemoryEventStore';\n\n// Define a more specific type for T that extends Document\ntype DocumentWithId = Document & { _id?: string | number };\n\nexport type InMemoryProjectionSpecEvent<\n  EventType extends Event,\n  EventMetaDataType extends InMemoryReadEventMetadata =\n    InMemoryReadEventMetadata,\n> = EventType & {\n  metadata?: Partial<EventMetaDataType>;\n};\n\nexport type InMemoryProjectionSpecWhenOptions = { numberOfTimes: number };\n\nexport type InMemoryProjectionSpec<EventType extends Event> = (\n  givenEvents: InMemoryProjectionSpecEvent<EventType>[],\n) => {\n  when: (\n    events: InMemoryProjectionSpecEvent<EventType>[],\n    options?: InMemoryProjectionSpecWhenOptions,\n  ) => {\n    then: (assert: InMemoryProjectionAssert, message?: string) => Promise<void>;\n    thenThrows: <ErrorType extends Error = Error>(\n      ...args: Parameters<ThenThrows<ErrorType>>\n    ) => Promise<void>;\n  };\n};\n\nexport type InMemoryProjectionAssert = (options: {\n  database: InMemoryDatabase;\n}) => Promise<void | boolean>;\n\nexport type InMemoryProjectionSpecOptions<EventType extends Event> = {\n  projection: InMemoryProjectionDefinition<EventType>;\n};\n\nexport const InMemoryProjectionSpec = {\n  for: <EventType extends Event>(\n    options: InMemoryProjectionSpecOptions<EventType>,\n  ): InMemoryProjectionSpec<EventType> => {\n    const { projection } = options;\n\n    return (givenEvents: InMemoryProjectionSpecEvent<EventType>[]) => {\n      return {\n        when: (\n          events: InMemoryProjectionSpecEvent<EventType>[],\n          options?: InMemoryProjectionSpecWhenOptions,\n        ) => {\n          const allEvents: ReadEvent<EventType, InMemoryReadEventMetadata>[] =\n            [];\n\n          const run = async (database: InMemoryDatabase) => {\n            let globalPosition = 0n;\n            const numberOfTimes = options?.numberOfTimes ?? 1;\n\n            for (const event of [\n              ...givenEvents,\n              ...Array.from({ length: numberOfTimes }).flatMap(() => events),\n            ]) {\n              const metadata: InMemoryReadEventMetadata = {\n                globalPosition: ++globalPosition,\n                streamPosition: globalPosition,\n                streamName: event.metadata?.streamName ?? `test-${uuid()}`,\n                messageId: uuid(),\n              };\n\n              allEvents.push({\n                ...event,\n                kind: 'Event',\n                metadata: {\n                  ...metadata,\n                  ...('metadata' in event ? (event.metadata ?? {}) : {}),\n                } as CombinedReadEventMetadata<\n                  EventType,\n                  InMemoryReadEventMetadata\n                >,\n              });\n            }\n\n            // Create a minimal mock EventStore implementation\n            const mockEventStore = {\n              database,\n              aggregateStream: async () => {\n                return Promise.resolve({\n                  state: {},\n                  currentStreamVersion: 0n,\n                  streamExists: false,\n                });\n              },\n              readStream: async () => {\n                return Promise.resolve({\n                  events: [],\n                  currentStreamVersion: 0n,\n                  streamExists: false,\n                });\n              },\n              appendToStream: async () => {\n                return Promise.resolve({\n                  nextExpectedStreamVersion: 0n,\n                  createdNewStream: false,\n                });\n              },\n              streamExists: async () => {\n                return Promise.resolve(false);\n              },\n            } as InMemoryEventStore;\n\n            await handleInMemoryProjections({\n              events: allEvents,\n              projections: [projection],\n              database,\n              eventStore: mockEventStore,\n            });\n          };\n\n          return {\n            then: async (\n              assertFn: InMemoryProjectionAssert,\n              message?: string,\n            ): Promise<void> => {\n              const database = getInMemoryDatabase();\n              await run(database);\n\n              const succeeded = await assertFn({ database });\n\n              if (succeeded !== undefined && succeeded === false) {\n                assertFails(\n                  message ??\n                    \"Projection specification didn't match the criteria\",\n                );\n              }\n            },\n            thenThrows: async <ErrorType extends Error = Error>(\n              ...args: Parameters<ThenThrows<ErrorType>>\n            ): Promise<void> => {\n              const database = getInMemoryDatabase();\n              try {\n                await run(database);\n                throw new AssertionError('Handler did not fail as expected');\n              } catch (error) {\n                if (error instanceof AssertionError) throw error;\n\n                if (args.length === 0) return;\n\n                if (!isErrorConstructor(args[0])) {\n                  assertTrue(\n                    args[0](error as ErrorType),\n                    `Error didn't match the error condition: ${error?.toString()}`,\n                  );\n                  return;\n                }\n\n                assertTrue(\n                  error instanceof args[0],\n                  `Caught error is not an instance of the expected type: ${error?.toString()}`,\n                );\n\n                if (args[1]) {\n                  assertTrue(\n                    args[1](error as ErrorType),\n                    `Error didn't match the error condition: ${error?.toString()}`,\n                  );\n                }\n              }\n            },\n          };\n        },\n      };\n    };\n  },\n};\n\n// Helper functions for creating events in stream\nexport const eventInStream = <\n  EventType extends Event = Event,\n  EventMetaDataType extends InMemoryReadEventMetadata =\n    InMemoryReadEventMetadata,\n>(\n  streamName: string,\n  event: InMemoryProjectionSpecEvent<EventType, EventMetaDataType>,\n): InMemoryProjectionSpecEvent<EventType, EventMetaDataType> => {\n  return {\n    ...event,\n    metadata: {\n      ...(event.metadata ?? {}),\n      streamName: event.metadata?.streamName ?? streamName,\n    } as Partial<EventMetaDataType>,\n  };\n};\n\nexport const eventsInStream = <\n  EventType extends Event = Event,\n  EventMetaDataType extends InMemoryReadEventMetadata =\n    InMemoryReadEventMetadata,\n>(\n  streamName: string,\n  events: InMemoryProjectionSpecEvent<EventType, EventMetaDataType>[],\n): InMemoryProjectionSpecEvent<EventType, EventMetaDataType>[] => {\n  return events.map((e) => eventInStream(streamName, e));\n};\n\nexport const newEventsInStream = eventsInStream;\n\n// Assertion helpers for checking documents\nexport function documentExists<T extends DocumentWithId>(\n  expected: Partial<T>,\n  options: { inCollection: string; withId: string | number },\n): InMemoryProjectionAssert {\n  return async ({ database }) => {\n    const collection = database.collection<T>(options.inCollection);\n\n    const document = await collection.findOne((doc) => {\n      // Handle both string IDs and numeric IDs in a type-safe way\n      const docId = '_id' in doc ? doc._id : undefined;\n      return docId === options.withId;\n    });\n\n    if (!document) {\n      assertFails(\n        `Document with ID ${options.withId} does not exist in collection ${options.inCollection}`,\n      );\n      return Promise.resolve(false);\n    }\n\n    // Check that all expected properties exist with expected values\n    for (const [key, value] of Object.entries(expected)) {\n      const propKey = key as keyof typeof document;\n      if (\n        !(key in document) ||\n        JSONParser.stringify(document[propKey]) !== JSONParser.stringify(value)\n      ) {\n        assertFails(`Property ${key} doesn't match the expected value`);\n        return Promise.resolve(false);\n      }\n    }\n\n    return Promise.resolve(true);\n  };\n}\n\n// Helper for checking document contents\nexport const expectInMemoryDocuments = {\n  fromCollection: <T extends DocumentWithId>(collectionName: string) => ({\n    withId: (id: string | number) => ({\n      toBeEqual: (expected: Partial<T>): InMemoryProjectionAssert =>\n        documentExists<T>(expected, {\n          inCollection: collectionName,\n          withId: id,\n        }),\n    }),\n  }),\n};\n","import { JSONParser } from '../serialization';\nimport type { DefaultRecord } from '../typing';\nimport { deepEquals } from '../utils';\n\nexport class AssertionError extends Error {\n  constructor(message: string) {\n    super(message);\n  }\n}\n\nexport const isSubset = (superObj: unknown, subObj: unknown): boolean => {\n  const sup = superObj as DefaultRecord;\n  const sub = subObj as DefaultRecord;\n\n  assertOk(sup);\n  assertOk(sub);\n\n  return Object.keys(sub).every((ele: string) => {\n    if (typeof sub[ele] == 'object') {\n      return isSubset(sup[ele], sub[ele]);\n    }\n    return sub[ele] === sup[ele];\n  });\n};\n\nexport const assertFails = (message?: string) => {\n  throw new AssertionError(message ?? 'That should not ever happened, right?');\n};\n\nexport const assertThrowsAsync = async <TError extends Error>(\n  fun: () => Promise<void>,\n  errorCheck?: (error: Error) => boolean,\n): Promise<TError> => {\n  try {\n    await fun();\n  } catch (error) {\n    const typedError = error as TError;\n    if (typedError instanceof AssertionError || !errorCheck) {\n      assertFalse(\n        typedError instanceof AssertionError,\n        \"Function didn't throw expected error\",\n      );\n      return typedError;\n    }\n\n    assertTrue(\n      errorCheck(typedError),\n      `Error doesn't match the expected condition: ${JSONParser.stringify(error)}`,\n    );\n\n    return typedError;\n  }\n  throw new AssertionError(\"Function didn't throw expected error\");\n};\n\nexport const assertThrows = <TError extends Error>(\n  fun: () => void,\n  errorCheck?: (error: Error) => boolean,\n): TError => {\n  try {\n    fun();\n  } catch (error) {\n    const typedError = error as TError;\n\n    if (errorCheck) {\n      assertTrue(\n        errorCheck(typedError),\n        `Error doesn't match the expected condition: ${JSONParser.stringify(error)}`,\n      );\n    } else if (typedError instanceof AssertionError) {\n      assertFalse(\n        typedError instanceof AssertionError,\n        \"Function didn't throw expected error\",\n      );\n    }\n\n    return typedError;\n  }\n  throw new AssertionError(\"Function didn't throw expected error\");\n};\n\nexport const assertDoesNotThrow = <TError extends Error>(\n  fun: () => void,\n  errorCheck?: (error: Error) => boolean,\n): TError | null => {\n  try {\n    fun();\n    return null;\n  } catch (error) {\n    const typedError = error as TError;\n\n    if (errorCheck) {\n      assertFalse(\n        errorCheck(typedError),\n        `Error matching the expected condition was thrown!: ${JSONParser.stringify(error)}`,\n      );\n    } else {\n      assertFails(`Function threw an error: ${JSONParser.stringify(error)}`);\n    }\n\n    return typedError;\n  }\n};\n\nexport const assertRejects = async <T, TError extends Error = Error>(\n  promise: Promise<T>,\n  errorCheck?: ((error: TError) => boolean) | TError,\n) => {\n  try {\n    await promise;\n    throw new AssertionError(\"Function didn't throw expected error\");\n  } catch (error) {\n    if (!errorCheck) return;\n\n    if (errorCheck instanceof Error) assertDeepEqual(error, errorCheck);\n    else assertTrue(errorCheck(error as TError));\n  }\n};\n\nexport const assertMatches = (\n  actual: unknown,\n  expected: unknown,\n  message?: string,\n) => {\n  if (!isSubset(actual, expected))\n    throw new AssertionError(\n      message ??\n        `subObj:\\n${JSONParser.stringify(expected)}\\nis not subset of\\n${JSONParser.stringify(actual)}`,\n    );\n};\n\nexport const assertDeepEqual = <T = unknown>(\n  actual: T,\n  expected: T,\n  message?: string,\n) => {\n  if (!deepEquals(actual, expected))\n    throw new AssertionError(\n      message ??\n        `subObj:\\n${JSONParser.stringify(expected)}\\nis not equal to\\n${JSONParser.stringify(actual)}`,\n    );\n};\n\nexport const assertNotDeepEqual = <T = unknown>(\n  actual: T,\n  expected: T,\n  message?: string,\n) => {\n  if (deepEquals(actual, expected))\n    throw new AssertionError(\n      message ??\n        `subObj:\\n${JSONParser.stringify(expected)}\\nis equals to\\n${JSONParser.stringify(actual)}`,\n    );\n};\n\nexport const assertThat = <T>(item: T) => {\n  return {\n    isEqualTo: (other: T) => assertTrue(deepEquals(item, other)),\n  };\n};\n\nexport const assertDefined = (\n  value: unknown,\n  message?: string | Error,\n): asserts value => {\n  assertOk(value, message instanceof Error ? message.message : message);\n};\n\nexport function assertFalse(\n  condition: boolean,\n  message?: string,\n): asserts condition is false {\n  if (condition !== false)\n    throw new AssertionError(message ?? `Condition is true`);\n}\n\nexport function assertTrue(\n  condition: boolean,\n  message?: string,\n): asserts condition is true {\n  if (condition !== true)\n    throw new AssertionError(message ?? `Condition is false`);\n}\n\n// TODO: replace with assertDefined\nexport function assertOk<T>(\n  obj: T | null | undefined,\n  message?: string,\n): asserts obj is T {\n  if (!obj) throw new AssertionError(message ?? `Condition is not truthy`);\n}\n\nexport function assertEqual<T>(\n  expected: T | null | undefined,\n  actual: T | null | undefined,\n  message?: string,\n): void {\n  if (expected !== actual)\n    throw new AssertionError(\n      `${message ?? 'Objects are not equal'}:\\nExpected: ${JSONParser.stringify(expected)}\\nActual: ${JSONParser.stringify(actual)}`,\n    );\n}\n\nexport function assertNotEqual<T>(\n  obj: T | null | undefined,\n  other: T | null | undefined,\n  message?: string,\n): void {\n  if (obj === other)\n    throw new AssertionError(\n      message ?? `Objects are equal: ${JSONParser.stringify(obj)}`,\n    );\n}\n\nexport function assertIsNotNull<T extends object | bigint>(\n  result: T | null,\n): asserts result is T {\n  assertNotEqual(result, null);\n  assertOk(result);\n}\n\nexport function assertIsNull<T extends object>(\n  result: T | null,\n): asserts result is null {\n  assertEqual(result, null);\n}\n\ntype Call = {\n  arguments: unknown[];\n  result: unknown;\n  target: unknown;\n  this: unknown;\n};\n\nexport type ArgumentMatcher = (arg: unknown) => boolean;\n\nexport const argValue =\n  <T>(value: T): ArgumentMatcher =>\n  (arg) =>\n    deepEquals(arg, value);\n\nexport const argMatches =\n  <T>(matches: (arg: T) => boolean): ArgumentMatcher =>\n  (arg) =>\n    matches(arg as T);\n\n// eslint-disable-next-line @typescript-eslint/no-unsafe-function-type\nexport type MockedFunction = Function & { mock?: { calls: Call[] } };\n\nexport function verifyThat(fn: MockedFunction) {\n  return {\n    calledTimes: (times: number) => {\n      assertEqual(fn.mock?.calls?.length, times);\n    },\n    notCalled: () => {\n      assertEqual(fn?.mock?.calls?.length, 0);\n    },\n    called: () => {\n      assertTrue(\n        fn.mock?.calls.length !== undefined && fn.mock.calls.length > 0,\n      );\n    },\n    calledWith: (...args: unknown[]) => {\n      assertTrue(\n        fn.mock?.calls.length !== undefined &&\n          fn.mock.calls.length >= 1 &&\n          fn.mock.calls.some((call) => deepEquals(call.arguments, args)),\n      );\n    },\n    calledOnceWith: (...args: unknown[]) => {\n      assertTrue(\n        fn.mock?.calls.length !== undefined &&\n          fn.mock.calls.length === 1 &&\n          fn.mock.calls.some((call) => deepEquals(call.arguments, args)),\n      );\n    },\n    calledWithArgumentMatching: (...matches: ArgumentMatcher[]) => {\n      assertTrue(\n        fn.mock?.calls.length !== undefined && fn.mock.calls.length >= 1,\n      );\n      assertTrue(\n        fn.mock?.calls.length !== undefined &&\n          fn.mock.calls.length >= 1 &&\n          fn.mock.calls.some(\n            (call) =>\n              call.arguments &&\n              call.arguments.length >= matches.length &&\n              matches.every((match, index) => match(call.arguments[index])),\n          ),\n      );\n    },\n    notCalledWithArgumentMatching: (...matches: ArgumentMatcher[]) => {\n      assertFalse(\n        fn.mock?.calls.length !== undefined &&\n          fn.mock.calls.length >= 1 &&\n          fn.mock.calls[0]!.arguments &&\n          fn.mock.calls[0]!.arguments.length >= matches.length &&\n          matches.every((match, index) =>\n            match(fn.mock!.calls[0]!.arguments[index]),\n          ),\n      );\n    },\n  };\n}\n\nexport const assertThatArray = <T>(array: T[]) => {\n  return {\n    isEmpty: () =>\n      assertEqual(\n        array.length,\n        0,\n        `Array is not empty ${JSONParser.stringify(array)}`,\n      ),\n    isNotEmpty: () => assertNotEqual(array.length, 0, `Array is empty`),\n    hasSize: (length: number) => assertEqual(array.length, length),\n    containsElements: (other: T[]) => {\n      assertTrue(other.every((ts) => array.some((o) => deepEquals(ts, o))));\n    },\n    containsElementsMatching: (other: T[]) => {\n      assertTrue(other.every((ts) => array.some((o) => isSubset(o, ts))));\n    },\n    containsOnlyElementsMatching: (other: T[]) => {\n      assertEqual(array.length, other.length, `Arrays lengths don't match`);\n      assertTrue(other.every((ts) => array.some((o) => isSubset(o, ts))));\n    },\n    containsExactlyInAnyOrder: (other: T[]) => {\n      assertEqual(array.length, other.length);\n      assertTrue(array.every((ts) => other.some((o) => deepEquals(ts, o))));\n    },\n    containsExactlyInAnyOrderElementsOf: (other: T[]) => {\n      assertEqual(array.length, other.length);\n      assertTrue(array.every((ts) => other.some((o) => deepEquals(ts, o))));\n    },\n    containsExactlyElementsOf: (other: T[]) => {\n      assertEqual(array.length, other.length);\n      for (let i = 0; i < array.length; i++) {\n        assertTrue(deepEquals(array[i], other[i]));\n      }\n    },\n    containsExactly: (elem: T) => {\n      assertEqual(array.length, 1);\n      assertTrue(deepEquals(array[0], elem));\n    },\n    contains: (elem: T) => {\n      assertTrue(array.some((a) => deepEquals(a, elem)));\n    },\n    containsOnlyOnceElementsOf: (other: T[]) => {\n      assertTrue(\n        other\n          .map((o) => array.filter((a) => deepEquals(a, o)).length)\n          .filter((a) => a === 1).length === other.length,\n      );\n    },\n    containsAnyOf: (other: T[]) => {\n      assertTrue(array.some((a) => other.some((o) => deepEquals(a, o))));\n    },\n    allMatch: (matches: (item: T) => boolean) => {\n      assertTrue(array.every(matches));\n    },\n    anyMatches: (matches: (item: T) => boolean) => {\n      assertTrue(array.some(matches));\n    },\n    allMatchAsync: async (\n      matches: (item: T) => Promise<boolean>,\n    ): Promise<void> => {\n      for (const item of array) {\n        assertTrue(await matches(item));\n      }\n    },\n  };\n};\n","import { isErrorConstructor, type ErrorConstructor } from '../errors';\nimport { AssertionError, assertThatArray, assertTrue } from './assertions';\n\ntype ErrorCheck<ErrorType> = (error: ErrorType) => boolean;\n\nexport type ThenThrows<ErrorType extends Error> =\n  | (() => void)\n  | ((errorConstructor: ErrorConstructor<ErrorType>) => void)\n  | ((errorCheck: ErrorCheck<ErrorType>) => void)\n  | ((\n      errorConstructor: ErrorConstructor<ErrorType>,\n      errorCheck?: ErrorCheck<ErrorType>,\n    ) => void);\n\nexport type DeciderSpecification<Command, Event> = (\n  givenEvents: Event | Event[],\n) => {\n  when: (command: Command) => {\n    then: (expectedEvents: Event | Event[]) => void;\n    thenNothingHappened: () => void;\n    thenThrows: <ErrorType extends Error = Error>(\n      ...args: Parameters<ThenThrows<ErrorType>>\n    ) => void;\n  };\n};\nexport type AsyncDeciderSpecification<Command, Event> = (\n  givenEvents: Event | Event[],\n) => {\n  when: (command: Command) => {\n    then: (expectedEvents: Event | Event[]) => Promise<void>;\n    thenNothingHappened: () => Promise<void>;\n    thenThrows: <ErrorType extends Error = Error>(\n      ...args: Parameters<ThenThrows<ErrorType>>\n    ) => Promise<void>;\n  };\n};\n\nexport const DeciderSpecification = {\n  for: deciderSpecificationFor,\n};\n\nfunction deciderSpecificationFor<Command, Event, State>(decider: {\n  decide: (command: Command, state: State) => Event | Event[];\n  evolve: (state: State, event: Event) => State;\n  initialState: () => State;\n}): DeciderSpecification<Command, Event>;\nfunction deciderSpecificationFor<Command, Event, State>(decider: {\n  decide: (command: Command, state: State) => Promise<Event | Event[]>;\n  evolve: (state: State, event: Event) => State;\n  initialState: () => State;\n}): AsyncDeciderSpecification<Command, Event>;\nfunction deciderSpecificationFor<Command, Event, State>(decider: {\n  decide: (\n    command: Command,\n    state: State,\n  ) => Event | Event[] | Promise<Event | Event[]>;\n  evolve: (state: State, event: Event) => State;\n  initialState: () => State;\n}):\n  | DeciderSpecification<Command, Event>\n  | AsyncDeciderSpecification<Command, Event> {\n  {\n    return (givenEvents: Event | Event[]) => {\n      return {\n        when: (command: Command) => {\n          const handle = () => {\n            const existingEvents = Array.isArray(givenEvents)\n              ? givenEvents\n              : [givenEvents];\n\n            const currentState = existingEvents.reduce<State>(\n              decider.evolve,\n              decider.initialState(),\n            );\n\n            return decider.decide(command, currentState);\n          };\n\n          return {\n            then: (expectedEvents: Event | Event[]): void | Promise<void> => {\n              const resultEvents = handle();\n\n              if (resultEvents instanceof Promise) {\n                return resultEvents.then((events) => {\n                  thenHandler(events, expectedEvents);\n                });\n              }\n\n              thenHandler(resultEvents, expectedEvents);\n            },\n            thenNothingHappened: (): void | Promise<void> => {\n              const resultEvents = handle();\n\n              if (resultEvents instanceof Promise) {\n                return resultEvents.then((events) => {\n                  thenNothingHappensHandler(events);\n                });\n              }\n\n              thenNothingHappensHandler(resultEvents);\n            },\n            thenThrows: <ErrorType extends Error>(\n              ...args: Parameters<ThenThrows<ErrorType>>\n            ): void | Promise<void> => {\n              try {\n                const result = handle();\n                if (result instanceof Promise) {\n                  return result\n                    .then(() => {\n                      throw new AssertionError(\n                        'Handler did not fail as expected',\n                      );\n                    })\n                    .catch((error) => {\n                      thenThrowsErrorHandler(error, args);\n                    });\n                }\n                throw new AssertionError('Handler did not fail as expected');\n              } catch (error) {\n                thenThrowsErrorHandler(error, args);\n              }\n            },\n          };\n        },\n      };\n    };\n  }\n}\n\nfunction thenHandler<Event>(\n  events: Event | Event[],\n  expectedEvents: Event | Event[],\n): void {\n  const resultEventsArray = Array.isArray(events) ? events : [events];\n\n  const expectedEventsArray = Array.isArray(expectedEvents)\n    ? expectedEvents\n    : [expectedEvents];\n\n  assertThatArray(resultEventsArray).containsOnlyElementsMatching(\n    expectedEventsArray,\n  );\n}\n\nfunction thenNothingHappensHandler<Event>(events: Event | Event[]): void {\n  const resultEventsArray = Array.isArray(events) ? events : [events];\n  assertThatArray(resultEventsArray).isEmpty();\n}\n\nfunction thenThrowsErrorHandler<ErrorType extends Error>(\n  error: unknown,\n  args: Parameters<ThenThrows<ErrorType>>,\n): void {\n  if (error instanceof AssertionError) throw error;\n\n  if (args.length === 0) return;\n\n  if (!isErrorConstructor(args[0])) {\n    assertTrue(\n      args[0](error as ErrorType),\n      `Error didn't match the error condition: ${error?.toString()}`,\n    );\n    return;\n  }\n\n  assertTrue(\n    error instanceof args[0],\n    `Caught error is not an instance of the expected type: ${error?.toString()}`,\n  );\n\n  if (args[1]) {\n    assertTrue(\n      args[1](error as ErrorType),\n      `Error didn't match the error condition: ${error?.toString()}`,\n    );\n  }\n}\n","import type {\n  AggregateStreamOptions,\n  AggregateStreamResult,\n  AppendToStreamOptions,\n  AppendToStreamResult,\n  EventStore,\n  EventStoreReadEventMetadata,\n  ReadStreamOptions,\n  ReadStreamResult,\n  StreamPositionTypeOfEventStore,\n} from '../eventStore';\nimport { type Event, type EventMetaDataOf } from '../typing';\n\nexport type TestEventStream<EventType extends Event = Event> = [\n  string,\n  EventType[],\n];\n\nexport type EventStoreWrapper<Store extends EventStore> = Store & {\n  appendedEvents: Map<string, TestEventStream>;\n  setup<EventType extends Event>(\n    streamName: string,\n    events: EventType[],\n  ): Promise<AppendToStreamResult<StreamPositionTypeOfEventStore<Store>>>;\n};\n\nexport const WrapEventStore = <Store extends EventStore>(\n  eventStore: Store,\n): EventStoreWrapper<Store> => {\n  const appendedEvents = new Map<string, TestEventStream>();\n\n  const wrapped = {\n    ...eventStore,\n    aggregateStream<State, EventType extends Event>(\n      streamName: string,\n      options: AggregateStreamOptions<State, EventType>,\n    ): Promise<\n      AggregateStreamResult<State, StreamPositionTypeOfEventStore<Store>>\n    > {\n      return eventStore.aggregateStream(streamName, options);\n    },\n\n    async readStream<EventType extends Event>(\n      streamName: string,\n      options?: ReadStreamOptions<StreamPositionTypeOfEventStore<Store>>,\n    ): Promise<\n      ReadStreamResult<\n        EventType,\n        EventStoreReadEventMetadata<Store> & EventMetaDataOf<EventType>\n      >\n    > {\n      return (await eventStore.readStream(\n        streamName,\n        options,\n      )) as ReadStreamResult<\n        EventType,\n        EventStoreReadEventMetadata<Store> & EventMetaDataOf<EventType>\n      >;\n    },\n\n    appendToStream: async <EventType extends Event>(\n      streamName: string,\n      events: EventType[],\n      options?: AppendToStreamOptions<StreamPositionTypeOfEventStore<Store>>,\n    ): Promise<AppendToStreamResult<StreamPositionTypeOfEventStore<Store>>> => {\n      const result = await eventStore.appendToStream(\n        streamName,\n        events,\n        options,\n      );\n\n      const currentStream = appendedEvents.get(streamName) ?? [streamName, []];\n\n      appendedEvents.set(streamName, [\n        streamName,\n        [...currentStream[1], ...events],\n      ]);\n\n      return result;\n    },\n\n    appendedEvents,\n\n    setup: async <EventType extends Event>(\n      streamName: string,\n      events: EventType[],\n    ): Promise<AppendToStreamResult<StreamPositionTypeOfEventStore<Store>>> => {\n      return eventStore.appendToStream(streamName, events);\n    },\n\n    // streamEvents: (): ReadableStream<\n    //   // eslint-disable-next-line @typescript-eslint/no-redundant-type-constituents\n    //   ReadEvent<Event, ReadEventMetadataType> | GlobalSubscriptionEvent\n    // > => {\n    //   return eventStore.streamEvents();\n    // },\n  };\n\n  return wrapped as EventStoreWrapper<Store>;\n};\n","import type {\n  AnyMessage,\n  AnyRecordedMessageMetadata,\n  RecordedMessage,\n} from '../../typing';\n\nexport type MessageDowncast<\n  MessageType extends AnyMessage,\n  MessagePayloadType extends AnyMessage = MessageType,\n  RecordedMessageMetadataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n> =\n  | ((\n      message: RecordedMessage<MessageType, RecordedMessageMetadataType>,\n    ) => RecordedMessage<MessagePayloadType, RecordedMessageMetadataType>)\n  | ((message: MessageType) => MessagePayloadType);\n\nexport const downcastRecordedMessage = <\n  MessageType extends AnyMessage,\n  MessagePayloadType extends AnyMessage = MessageType,\n  RecordedMessageMetadataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n>(\n  recordedMessage:\n    | RecordedMessage<MessageType, RecordedMessageMetadataType>\n    | MessageType,\n  options?: {\n    downcast?: MessageDowncast<\n      MessageType,\n      MessagePayloadType,\n      RecordedMessageMetadataType\n    >;\n  },\n): RecordedMessage<MessagePayloadType, RecordedMessageMetadataType> => {\n  if (!options?.downcast)\n    return recordedMessage as unknown as RecordedMessage<\n      MessagePayloadType,\n      RecordedMessageMetadataType\n    >;\n\n  const downcasted = options.downcast(\n    recordedMessage as RecordedMessage<\n      MessageType,\n      RecordedMessageMetadataType\n    >,\n  );\n\n  return {\n    ...recordedMessage,\n    // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment\n    data: downcasted.data,\n    ...('metadata' in recordedMessage || 'metadata' in downcasted\n      ? {\n          metadata: {\n            ...('metadata' in recordedMessage\n              ? (recordedMessage.metadata as object)\n              : {}),\n            ...('metadata' in downcasted\n              ? (downcasted.metadata as object)\n              : {}),\n          },\n        }\n      : {}),\n  } as unknown as RecordedMessage<\n    MessagePayloadType,\n    RecordedMessageMetadataType\n  >;\n};\n\nexport const downcastRecordedMessages = <\n  MessageType extends AnyMessage,\n  MessagePayloadType extends AnyMessage = MessageType,\n  RecordedMessageMetadataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n>(\n  recordedMessages:\n    | RecordedMessage<MessageType, RecordedMessageMetadataType>[]\n    | MessageType[],\n  options?: {\n    downcast?: MessageDowncast<\n      MessageType,\n      MessagePayloadType,\n      RecordedMessageMetadataType\n    >;\n  },\n): RecordedMessage<MessagePayloadType, RecordedMessageMetadataType>[] => {\n  if (!options?.downcast)\n    return recordedMessages as unknown as RecordedMessage<\n      MessagePayloadType,\n      RecordedMessageMetadataType\n    >[];\n\n  return recordedMessages.map((recordedMessage) =>\n    downcastRecordedMessage(recordedMessage, options),\n  );\n};\n","import type {\n  AnyMessage,\n  AnyRecordedMessageMetadata,\n  RecordedMessage,\n} from '../../typing';\n\nexport type MessageUpcast<\n  MessageType extends AnyMessage,\n  MessagePayloadType extends AnyMessage = MessageType,\n  RecordedMessageMetadataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n> =\n  | ((message: MessagePayloadType) => MessageType)\n  | ((\n      message: RecordedMessage<MessagePayloadType, RecordedMessageMetadataType>,\n    ) => RecordedMessage<MessageType, RecordedMessageMetadataType>);\n\nexport const upcastRecordedMessage = <\n  MessageType extends AnyMessage,\n  MessagePayloadType extends AnyMessage = MessageType,\n  RecordedMessageMetadataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n>(\n  recordedMessage:\n    | RecordedMessage<MessagePayloadType, RecordedMessageMetadataType>\n    | MessagePayloadType,\n  options?: {\n    upcast?: MessageUpcast<\n      MessageType,\n      MessagePayloadType,\n      RecordedMessageMetadataType\n    >;\n  },\n): RecordedMessage<MessageType, RecordedMessageMetadataType> => {\n  if (!options?.upcast)\n    return recordedMessage as unknown as RecordedMessage<\n      MessageType,\n      RecordedMessageMetadataType\n    >;\n\n  const upcasted = options.upcast(\n    recordedMessage as RecordedMessage<\n      MessagePayloadType,\n      RecordedMessageMetadataType\n    >,\n  );\n\n  return {\n    ...recordedMessage,\n    // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment\n    data: upcasted.data,\n    ...('metadata' in recordedMessage || 'metadata' in upcasted\n      ? {\n          metadata: {\n            ...('metadata' in recordedMessage\n              ? (recordedMessage.metadata as object)\n              : {}),\n            ...('metadata' in upcasted ? (upcasted.metadata as object) : {}),\n          },\n        }\n      : {}),\n  } as unknown as RecordedMessage<MessageType, RecordedMessageMetadataType>;\n};\n\nexport const upcastRecordedMessages = <\n  MessageType extends AnyMessage,\n  MessagePayloadType extends AnyMessage = MessageType,\n  RecordedMessageMetadataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n>(\n  recordedMessages:\n    | RecordedMessage<MessagePayloadType, RecordedMessageMetadataType>[]\n    | MessagePayloadType[],\n  options?: {\n    upcast?: MessageUpcast<\n      MessageType,\n      MessagePayloadType,\n      RecordedMessageMetadataType\n    >;\n  },\n): RecordedMessage<MessageType, RecordedMessageMetadataType>[] => {\n  if (!options?.upcast)\n    return recordedMessages as unknown as RecordedMessage<\n      MessageType,\n      RecordedMessageMetadataType\n    >[];\n\n  return recordedMessages.map((recordedMessage) =>\n    upcastRecordedMessage(recordedMessage, options),\n  );\n};\n","import {\n  canCreateEventStoreSession,\n  isExpectedVersionConflictError,\n  NO_CONCURRENCY_CHECK,\n  nulloSessionFactory,\n  STREAM_DOES_NOT_EXIST,\n  type AppendStreamResultOfEventStore,\n  type AppendToStreamOptions,\n  type EventStore,\n  type EventStoreReadEventMetadata,\n  type EventStoreSession,\n  type ExpectedStreamVersion,\n  type ReadStreamOptions,\n  type StreamPositionTypeOfEventStore,\n} from '../eventStore';\nimport type { Event, StreamPositionTypeOfReadEventMetadata } from '../typing';\nimport { asyncRetry, NoRetries, type AsyncRetryOptions } from '../utils';\n\nexport const CommandHandlerStreamVersionConflictRetryOptions: AsyncRetryOptions =\n  {\n    retries: 3,\n    minTimeout: 100,\n    factor: 1.5,\n    shouldRetryError: isExpectedVersionConflictError,\n  };\n\nexport type CommandHandlerRetryOptions =\n  | AsyncRetryOptions\n  | { onVersionConflict: true | number | AsyncRetryOptions };\n\nconst fromCommandHandlerRetryOptions = (\n  retryOptions: CommandHandlerRetryOptions | undefined,\n): AsyncRetryOptions => {\n  if (retryOptions === undefined) return NoRetries;\n\n  if ('onVersionConflict' in retryOptions) {\n    if (typeof retryOptions.onVersionConflict === 'boolean')\n      return CommandHandlerStreamVersionConflictRetryOptions;\n    else if (typeof retryOptions.onVersionConflict === 'number')\n      return {\n        ...CommandHandlerStreamVersionConflictRetryOptions,\n        retries: retryOptions.onVersionConflict,\n      };\n    else return retryOptions.onVersionConflict;\n  }\n\n  return retryOptions;\n};\n\n// #region command-handler\nexport type CommandHandlerResult<\n  State,\n  StreamEvent extends Event,\n  Store extends EventStore,\n> = AppendStreamResultOfEventStore<Store> & {\n  newState: State;\n  newEvents: StreamEvent[];\n};\n\nexport type CommandHandlerOptions<\n  State,\n  StreamEvent extends Event,\n  StoredEvent extends Event = StreamEvent,\n> = {\n  evolve: (state: State, event: StreamEvent) => State;\n  initialState: () => State;\n  mapToStreamId?: (id: string) => string;\n  retry?: CommandHandlerRetryOptions;\n  schema?: {\n    versioning?: {\n      upcast?: (event: StoredEvent) => StreamEvent;\n      downcast?: (event: StreamEvent) => StoredEvent;\n    };\n  };\n};\n\nexport type HandleOptions<Store extends EventStore> = Parameters<\n  Store['appendToStream']\n>[2] &\n  (\n    | {\n        expectedStreamVersion?: ExpectedStreamVersion<\n          StreamPositionTypeOfEventStore<Store>\n        >;\n      }\n    | {\n        retry?: CommandHandlerRetryOptions;\n      }\n  );\n\ntype CommandHandlerFunction<State, StreamEvent extends Event> = (\n  state: State,\n) => StreamEvent | StreamEvent[] | Promise<StreamEvent | StreamEvent[]>;\n\nexport const CommandHandler =\n  <\n    State,\n    StreamEvent extends Event,\n    EventPayloadType extends Event = StreamEvent,\n  >(\n    options: CommandHandlerOptions<State, StreamEvent, EventPayloadType>,\n  ) =>\n  async <Store extends EventStore>(\n    store: Store,\n    id: string,\n    handle:\n      | CommandHandlerFunction<State, StreamEvent>\n      | CommandHandlerFunction<State, StreamEvent>[],\n    handleOptions?: HandleOptions<Store>,\n  ): Promise<CommandHandlerResult<State, StreamEvent, Store>> =>\n    asyncRetry(\n      async () => {\n        const result = await withSession<\n          Store,\n          CommandHandlerResult<\n            State,\n            StreamEvent,\n            StreamPositionTypeOfEventStore<Store>\n          >\n        >(store, async ({ eventStore }) => {\n          const { evolve, initialState } = options;\n          const mapToStreamId = options.mapToStreamId ?? ((id) => id);\n\n          const streamName = mapToStreamId(id);\n\n          // 1. Aggregate the stream\n          const aggregationResult = await eventStore.aggregateStream<\n            State,\n            StreamEvent,\n            EventPayloadType\n          >(streamName, {\n            evolve,\n            initialState,\n            read: {\n              schema: options.schema,\n              ...(handleOptions as ReadStreamOptions<\n                StreamPositionTypeOfReadEventMetadata<\n                  EventStoreReadEventMetadata<Store>\n                >,\n                StreamEvent,\n                EventPayloadType\n              >),\n              // expected stream version is passed to fail fast\n              // if stream is in the wrong state\n              // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment\n              expectedStreamVersion:\n                handleOptions?.expectedStreamVersion ?? NO_CONCURRENCY_CHECK,\n            },\n          });\n\n          // 2. Use the aggregate state\n\n          const {\n            // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment\n            currentStreamVersion,\n            streamExists: _streamExists,\n            ...restOfAggregationResult\n          } = aggregationResult;\n\n          let state = aggregationResult.state;\n\n          const handlers = Array.isArray(handle) ? handle : [handle];\n          let eventsToAppend: StreamEvent[] = [];\n\n          // 3. Run business logic\n          for (const handler of handlers) {\n            const result = await handler(state);\n\n            const newEvents = Array.isArray(result) ? result : [result];\n\n            if (newEvents.length > 0) {\n              state = newEvents.reduce(evolve, state);\n            }\n\n            eventsToAppend = [...eventsToAppend, ...newEvents];\n          }\n\n          //const newEvents = Array.isArray(result) ? result : [result];\n\n          if (eventsToAppend.length === 0) {\n            return {\n              ...restOfAggregationResult,\n              newEvents: [],\n              newState: state,\n              // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment\n              nextExpectedStreamVersion: currentStreamVersion,\n              createdNewStream: false,\n            } as unknown as CommandHandlerResult<State, StreamEvent, Store>;\n          }\n\n          // Either use:\n          // - provided expected stream version,\n          // - current stream version got from stream aggregation,\n          // - or expect stream not to exists otherwise.\n          // eslint-disable-next-line @typescript-eslint/no-unsafe-assignment\n          const expectedStreamVersion: ExpectedStreamVersion<\n            StreamPositionTypeOfEventStore<Store>\n          > =\n            handleOptions?.expectedStreamVersion ??\n            (aggregationResult.streamExists\n              ? (currentStreamVersion as ExpectedStreamVersion<\n                  StreamPositionTypeOfEventStore<Store>\n                >)\n              : STREAM_DOES_NOT_EXIST);\n\n          // 4. Append result to the stream\n          const appendResult = await eventStore.appendToStream(\n            streamName,\n            eventsToAppend,\n            {\n              ...(handleOptions as AppendToStreamOptions<\n                StreamPositionTypeOfReadEventMetadata<\n                  EventStoreReadEventMetadata<Store>\n                >,\n                StreamEvent,\n                EventPayloadType\n              >),\n              expectedStreamVersion,\n            },\n          );\n\n          // 5. Return result with updated state\n          return {\n            ...appendResult,\n            newEvents: eventsToAppend,\n            newState: state,\n          } as unknown as CommandHandlerResult<State, StreamEvent, Store>;\n        });\n\n        return result;\n      },\n      fromCommandHandlerRetryOptions(\n        handleOptions && 'retry' in handleOptions\n          ? handleOptions.retry\n          : options.retry,\n      ),\n    );\n// #endregion command-handler\n\nconst withSession = <EventStoreType extends EventStore, T = unknown>(\n  eventStore: EventStoreType,\n  callback: (session: EventStoreSession<EventStoreType>) => Promise<T>,\n) => {\n  const sessionFactory = canCreateEventStoreSession<EventStoreType>(eventStore)\n    ? eventStore\n    : nulloSessionFactory<EventStoreType>(eventStore);\n\n  return sessionFactory.withSession(callback);\n};\n","import type { EventStore } from '../eventStore';\nimport { type Command, type Event } from '../typing';\nimport type { Decider } from '../typing/decider';\nimport {\n  CommandHandler,\n  type CommandHandlerOptions,\n  type HandleOptions,\n} from './handleCommand';\n\n// #region command-handler\n\nexport type DeciderCommandHandlerOptions<\n  State,\n  CommandType extends Command,\n  StreamEvent extends Event,\n> = CommandHandlerOptions<State, StreamEvent> &\n  Decider<State, CommandType, StreamEvent>;\n\nexport const DeciderCommandHandler =\n  <State, CommandType extends Command, StreamEvent extends Event>(\n    options: DeciderCommandHandlerOptions<State, CommandType, StreamEvent>,\n  ) =>\n  async <Store extends EventStore>(\n    eventStore: Store,\n    id: string,\n    commands: CommandType | CommandType[],\n    handleOptions?: HandleOptions<Store>,\n  ) => {\n    const { decide, ...rest } = options;\n\n    const deciders = (Array.isArray(commands) ? commands : [commands]).map(\n      (command) => (state: State) => decide(command, state),\n    );\n\n    return CommandHandler<State, StreamEvent>(rest)(\n      eventStore,\n      id,\n      deciders,\n      handleOptions,\n    );\n  };\n// #endregion command-handler\n","import { EmmettError } from '../errors';\nimport {\n  type AnyCommand,\n  type AnyMessage,\n  type Command,\n  type CommandTypeOf,\n  type Event,\n  type EventTypeOf,\n  type Message,\n  type SingleMessageHandler,\n  type SingleRawMessageHandlerWithoutContext,\n} from '../typing';\n\nexport interface CommandSender {\n  send<CommandType extends Command = Command>(\n    command: CommandType,\n  ): Promise<void>;\n}\n\nexport interface EventsPublisher {\n  publish<EventType extends Event = Event>(event: EventType): Promise<void>;\n}\n\nexport type ScheduleOptions = { afterInMs: number } | { at: Date };\n\nexport interface MessageScheduler<CommandOrEvent extends Command | Event> {\n  schedule<MessageType extends CommandOrEvent>(\n    message: MessageType,\n    when?: ScheduleOptions,\n  ): void;\n}\n\nexport interface CommandBus extends CommandSender, MessageScheduler<Command> {}\n\nexport interface EventBus extends EventsPublisher, MessageScheduler<Event> {}\n\nexport interface MessageBus extends CommandBus, EventBus {\n  schedule<MessageType extends Command | Event>(\n    message: MessageType,\n    when?: ScheduleOptions,\n  ): void;\n}\n\nexport interface CommandProcessor {\n  handle<CommandType extends Command>(\n    commandHandler: SingleMessageHandler<CommandType>,\n    ...commandTypes: CommandTypeOf<CommandType>[]\n  ): void;\n}\nexport interface EventSubscription {\n  subscribe<EventType extends Event>(\n    eventHandler: SingleMessageHandler<EventType>,\n    ...eventTypes: EventTypeOf<EventType>[]\n  ): void;\n}\n\nexport type ScheduledMessage = {\n  message: Message;\n  options?: ScheduleOptions;\n};\n\nexport interface ScheduledMessageProcessor {\n  dequeue(): ScheduledMessage[];\n}\n\nexport type MessageSubscription = EventSubscription | CommandProcessor;\n\nexport const getInMemoryMessageBus = (): MessageBus &\n  EventSubscription &\n  CommandProcessor &\n  ScheduledMessageProcessor => {\n  const allHandlers = new Map<\n    string,\n    SingleRawMessageHandlerWithoutContext<AnyMessage>[]\n  >();\n  let pendingMessages: ScheduledMessage[] = [];\n\n  return {\n    send: async <CommandType extends Command = AnyCommand>(\n      command: CommandType,\n    ): Promise<void> => {\n      const handlers = allHandlers.get(command.type);\n\n      if (handlers === undefined || handlers.length === 0)\n        throw new EmmettError(\n          `No handler registered for command ${command.type}!`,\n        );\n\n      const commandHandler = handlers[0]!;\n\n      await commandHandler(command);\n    },\n\n    publish: async <EventType extends Event = Event>(\n      event: EventType,\n    ): Promise<void> => {\n      const handlers = allHandlers.get(event.type) ?? [];\n\n      for (const handler of handlers) {\n        const eventHandler = handler;\n\n        await eventHandler(event);\n      }\n    },\n\n    schedule: <MessageType extends Message>(\n      message: MessageType,\n      when?: ScheduleOptions,\n    ): void => {\n      pendingMessages = [...pendingMessages, { message, options: when }];\n    },\n\n    handle: <CommandType extends Command>(\n      commandHandler: SingleMessageHandler<CommandType>,\n      ...commandTypes: CommandTypeOf<CommandType>[]\n    ): void => {\n      const alreadyRegistered = [...allHandlers.keys()].filter((registered) =>\n        commandTypes.includes(registered),\n      );\n\n      if (alreadyRegistered.length > 0)\n        throw new EmmettError(\n          `Cannot register handler for commands ${alreadyRegistered.join(', ')} as they're already registered!`,\n        );\n      for (const commandType of commandTypes) {\n        allHandlers.set(commandType, [\n          commandHandler as SingleRawMessageHandlerWithoutContext<AnyMessage>,\n        ]);\n      }\n    },\n\n    subscribe<EventType extends Event>(\n      eventHandler: SingleMessageHandler<EventType>,\n      ...eventTypes: EventTypeOf<EventType>[]\n    ): void {\n      for (const eventType of eventTypes) {\n        if (!allHandlers.has(eventType)) allHandlers.set(eventType, []);\n\n        allHandlers.set(eventType, [\n          ...(allHandlers.get(eventType) ?? []),\n          eventHandler as SingleRawMessageHandlerWithoutContext<AnyMessage>,\n        ]);\n      }\n    },\n\n    dequeue: (): ScheduledMessage[] => {\n      const pending = pendingMessages;\n      pendingMessages = [];\n      return pending;\n    },\n  };\n};\n","import { v7 as uuid } from 'uuid';\nimport type { EmmettError } from '../errors';\nimport { upcastRecordedMessage } from '../eventStore';\nimport type { ProjectionDefinition } from '../projections';\nimport { JSONParser } from '../serialization';\nimport {\n  defaultTag,\n  type AnyEvent,\n  type AnyMessage,\n  type AnyReadEventMetadata,\n  type AnyRecordedMessageMetadata,\n  type BatchRecordedMessageHandlerWithContext,\n  type CanHandle,\n  type DefaultRecord,\n  type Event,\n  type GlobalPositionTypeOfRecordedMessageMetadata,\n  type Message,\n  type MessageHandlerResult,\n  type RecordedMessage,\n  type SingleMessageHandlerWithContext,\n  type SingleRecordedMessageHandlerWithContext,\n} from '../typing';\nimport { onShutdown } from '../utils/shutdown';\nimport { isBigint } from '../validation';\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport type CurrentMessageProcessorPosition<CheckpointType = any> =\n  | { lastCheckpoint: CheckpointType }\n  | 'BEGINNING'\n  | 'END';\n\nexport type GetCheckpoint<\n  MessageType extends AnyMessage = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n> = (\n  message: RecordedMessage<MessageType, MessageMetadataType>,\n) => CheckpointType | null;\n\nexport const getCheckpoint = <\n  MessageType extends AnyMessage = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n>(\n  message: RecordedMessage<MessageType, MessageMetadataType>,\n): CheckpointType | null => {\n  // eslint-disable-next-line @typescript-eslint/no-unsafe-return\n  return 'checkpoint' in message.metadata\n    ? // eslint-disable-next-line @typescript-eslint/no-unsafe-member-access\n      message.metadata.checkpoint\n    : 'globalPosition' in message.metadata &&\n        // eslint-disable-next-line @typescript-eslint/no-unsafe-member-access\n        isBigint(message.metadata.globalPosition)\n      ? // eslint-disable-next-line @typescript-eslint/no-unsafe-member-access\n        message.metadata.globalPosition\n      : 'streamPosition' in message.metadata &&\n          // eslint-disable-next-line @typescript-eslint/no-unsafe-member-access\n          isBigint(message.metadata.streamPosition)\n        ? // eslint-disable-next-line @typescript-eslint/no-unsafe-member-access\n          message.metadata.streamPosition\n        : null;\n};\n\nexport const wasMessageHandled = <\n  MessageType extends AnyMessage = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n>(\n  message: RecordedMessage<MessageType, MessageMetadataType>,\n  checkpoint: CheckpointType | null,\n): boolean => {\n  //TODO Make it smarter\n  const messageCheckpoint = getCheckpoint(message);\n  const checkpointBigint = checkpoint as bigint | null;\n\n  return (\n    messageCheckpoint !== null &&\n    messageCheckpoint !== undefined &&\n    checkpointBigint !== null &&\n    checkpointBigint !== undefined &&\n    messageCheckpoint <= checkpointBigint\n  );\n};\n\n// eslint-disable-next-line @typescript-eslint/no-explicit-any\nexport type MessageProcessorStartFrom<CheckpointType = any> =\n  | CurrentMessageProcessorPosition<CheckpointType>\n  | 'CURRENT';\n\nexport type MessageProcessorType = 'projector' | 'reactor';\nexport const MessageProcessorType = {\n  PROJECTOR: 'projector' as MessageProcessorType,\n  REACTOR: 'reactor' as MessageProcessorType,\n};\n\nexport type MessageProcessor<\n  MessageType extends AnyMessage = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n> = {\n  id: string;\n  instanceId: string;\n  type: string;\n  init: (options: Partial<HandlerContext>) => Promise<void>;\n  start: (\n    options: Partial<HandlerContext>,\n  ) => Promise<CurrentMessageProcessorPosition<CheckpointType> | undefined>;\n  close: (closeOptions: Partial<HandlerContext>) => Promise<void>;\n  isActive: boolean;\n  handle: BatchRecordedMessageHandlerWithContext<\n    MessageType,\n    MessageMetadataType,\n    Partial<HandlerContext>\n  >;\n};\n\nexport const MessageProcessor = {\n  result: {\n    skip: (options?: { reason?: string }): MessageHandlerResult => ({\n      type: 'SKIP',\n      ...(options ?? {}),\n    }),\n    stop: (options?: {\n      reason?: string;\n      error?: EmmettError;\n    }): MessageHandlerResult => ({\n      type: 'STOP',\n      ...(options ?? {}),\n    }),\n  },\n};\n\nexport type MessageProcessingScope<\n  HandlerContext extends DefaultRecord | undefined = undefined,\n> = <Result = MessageHandlerResult>(\n  handler: (context: HandlerContext) => Result | Promise<Result>,\n  partialContext: Partial<HandlerContext>,\n) => Result | Promise<Result>;\n\nexport type Checkpointer<\n  MessageType extends AnyMessage = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  HandlerContext extends DefaultRecord = DefaultRecord,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n> = {\n  read: ReadProcessorCheckpoint<CheckpointType, HandlerContext>;\n  store: StoreProcessorCheckpoint<\n    MessageType,\n    MessageMetadataType,\n    CheckpointType,\n    HandlerContext\n  >;\n};\n\nexport type ProcessorHooks<\n  HandlerContext extends DefaultRecord = DefaultRecord,\n> = {\n  onInit?: OnReactorInitHook<HandlerContext>;\n  onStart?: OnReactorStartHook<HandlerContext>;\n  onClose?: OnReactorCloseHook<HandlerContext>;\n};\n\nexport type BaseMessageProcessorOptions<\n  MessageType extends AnyMessage = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  HandlerContext extends DefaultRecord = DefaultRecord,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n> = {\n  type?: string;\n  processorId: string;\n  processorInstanceId?: string;\n  version?: number;\n  partition?: string;\n  startFrom?: MessageProcessorStartFrom<CheckpointType>;\n  stopAfter?: (\n    message: RecordedMessage<MessageType, MessageMetadataType>,\n  ) => boolean;\n  processingScope?: MessageProcessingScope<HandlerContext>;\n  checkpoints?: Checkpointer<\n    MessageType,\n    MessageMetadataType,\n    HandlerContext,\n    CheckpointType\n  >;\n  canHandle?: CanHandle<MessageType>;\n  hooks?: ProcessorHooks<HandlerContext>;\n};\n\nexport type HandlerOptions<\n  MessageType extends AnyMessage = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  HandlerContext extends DefaultRecord = DefaultRecord,\n> =\n  | {\n      eachMessage: SingleRecordedMessageHandlerWithContext<\n        MessageType,\n        MessageMetadataType,\n        HandlerContext\n      >;\n      eachBatch?: never;\n    }\n  | {\n      eachMessage?: never;\n      eachBatch: BatchRecordedMessageHandlerWithContext<\n        MessageType,\n        MessageMetadataType,\n        HandlerContext\n      >;\n    };\n\nexport type OnReactorInitHook<\n  HandlerContext extends DefaultRecord = DefaultRecord,\n> = (context: HandlerContext) => Promise<void>;\n\nexport type OnReactorStartHook<\n  HandlerContext extends DefaultRecord = DefaultRecord,\n> = (context: HandlerContext) => Promise<void>;\n\nexport type OnReactorCloseHook<\n  HandlerContext extends DefaultRecord = DefaultRecord,\n> = (context: HandlerContext) => Promise<void>;\n\nexport type ReactorOptions<\n  MessageType extends AnyMessage = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  HandlerContext extends DefaultRecord = DefaultRecord,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n  MessagePayloadType extends AnyMessage = MessageType,\n> = BaseMessageProcessorOptions<\n  MessageType,\n  MessageMetadataType,\n  HandlerContext,\n  CheckpointType\n> &\n  HandlerOptions<MessageType, MessageMetadataType, HandlerContext> & {\n    messageOptions?: {\n      schema?: {\n        versioning?: { upcast?: (event: MessagePayloadType) => MessageType };\n      };\n    };\n  };\n\nexport type ProjectorOptions<\n  EventType extends AnyEvent = AnyEvent,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  HandlerContext extends DefaultRecord = DefaultRecord,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n  EventPayloadType extends Event = EventType,\n> = Omit<\n  BaseMessageProcessorOptions<\n    EventType,\n    MessageMetadataType,\n    HandlerContext,\n    CheckpointType\n  >,\n  'type' | 'processorId'\n> & { processorId?: string } & {\n  truncateOnStart?: boolean;\n  projection: ProjectionDefinition<\n    EventType,\n    MessageMetadataType,\n    HandlerContext,\n    EventPayloadType\n  >;\n};\n\nexport const defaultProcessingMessageProcessingScope = <\n  HandlerContext = never,\n  Result = MessageHandlerResult,\n>(\n  handler: (context: HandlerContext) => Result | Promise<Result>,\n  partialContext: Partial<HandlerContext>,\n) => handler(partialContext as HandlerContext);\n\nexport type ReadProcessorCheckpointResult<CheckpointType = unknown> = {\n  lastCheckpoint: CheckpointType | null;\n};\n\nexport type ReadProcessorCheckpoint<\n  CheckpointType = unknown,\n  HandlerContext extends DefaultRecord = DefaultRecord,\n> = (\n  options: { processorId: string; partition?: string },\n  context: HandlerContext,\n) => Promise<ReadProcessorCheckpointResult<CheckpointType>>;\n\nexport type StoreProcessorCheckpointResult<CheckpointType = unknown> =\n  | {\n      success: true;\n      newCheckpoint: CheckpointType;\n    }\n  | { success: false; reason: 'IGNORED' | 'MISMATCH' | 'CURRENT_AHEAD' };\n\nexport type StoreProcessorCheckpoint<\n  MessageType extends Message = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  CheckpointType = unknown,\n  HandlerContext extends DefaultRecord | undefined = undefined,\n> =\n  | ((\n      options: {\n        message: RecordedMessage<MessageType, MessageMetadataType>;\n        processorId: string;\n        version: number | undefined;\n        lastCheckpoint: CheckpointType | null;\n        partition?: string;\n      },\n      context: HandlerContext,\n    ) => Promise<StoreProcessorCheckpointResult<CheckpointType | null>>)\n  | ((\n      options: {\n        message: RecordedMessage<MessageType, MessageMetadataType>;\n        processorId: string;\n        version: number | undefined;\n        lastCheckpoint: CheckpointType | null;\n        partition?: string;\n      },\n      context: HandlerContext,\n    ) => Promise<StoreProcessorCheckpointResult<CheckpointType>>);\n\nexport const defaultProcessorVersion = 1;\nexport const defaultProcessorPartition = defaultTag;\n\nexport const getProcessorInstanceId = (processorId: string): string =>\n  `${processorId}:${uuid()}`;\n\nexport const getProjectorId = (options: { projectionName: string }): string =>\n  `emt:processor:projector:${options.projectionName}`;\n\nexport const reactor = <\n  MessageType extends Message = AnyMessage,\n  MessageMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  HandlerContext extends DefaultRecord = DefaultRecord,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<MessageMetadataType>,\n  MessagePayloadType extends Message = MessageType,\n>(\n  options: ReactorOptions<\n    MessageType,\n    MessageMetadataType,\n    HandlerContext,\n    CheckpointType,\n    MessagePayloadType\n  >,\n): MessageProcessor<\n  MessageType,\n  MessageMetadataType,\n  HandlerContext,\n  CheckpointType\n> => {\n  const {\n    checkpoints,\n    processorId,\n    processorInstanceId: instanceId = getProcessorInstanceId(processorId),\n    type = MessageProcessorType.REACTOR,\n    version = defaultProcessorVersion,\n    partition = defaultProcessorPartition,\n    hooks = {},\n    processingScope = defaultProcessingMessageProcessingScope,\n    startFrom,\n    canHandle,\n    stopAfter,\n  } = options;\n\n  const eachMessage: SingleMessageHandlerWithContext<\n    MessageType,\n    MessageMetadataType,\n    HandlerContext\n  > =\n    'eachMessage' in options && options.eachMessage\n      ? options.eachMessage\n      : () => Promise.resolve();\n\n  let isInitiated = false;\n  let isActive = false;\n\n  let lastCheckpoint: CheckpointType | null = null;\n  let closeSignal: (() => void) | null = null;\n\n  const init = async (initOptions: Partial<HandlerContext>): Promise<void> => {\n    if (isInitiated) return;\n\n    if (hooks.onInit === undefined) {\n      isInitiated = true;\n      return;\n    }\n\n    return await processingScope(async (context) => {\n      await hooks.onInit!(context);\n      isInitiated = true;\n    }, initOptions);\n  };\n\n  const close = async (\n    closeOptions: Partial<HandlerContext>,\n  ): Promise<void> => {\n    // TODO: Align when active is set to false\n    // if (!isActive) return;\n\n    isActive = false;\n\n    if (closeSignal) {\n      closeSignal();\n      closeSignal = null;\n    }\n\n    if (hooks.onClose) {\n      await processingScope(hooks.onClose, closeOptions);\n    }\n  };\n\n  return {\n    // TODO: Consider whether not make it optional or add URN prefix\n    id: processorId,\n    instanceId,\n    type,\n    init,\n    start: async (\n      startOptions: Partial<HandlerContext>,\n    ): Promise<CurrentMessageProcessorPosition<CheckpointType> | undefined> => {\n      if (isActive) {\n        console.log(\n          `Processor ${processorId} with instance id ${instanceId} is already active. Start request ignored.`,\n        );\n        return;\n      }\n\n      console.log(\n        `Starting processor ${processorId} with instance id ${instanceId}`,\n      );\n\n      await init(startOptions);\n\n      isActive = true;\n\n      closeSignal = onShutdown(() => close({}));\n\n      if (lastCheckpoint !== null) {\n        console.log(\n          `Processor ${processorId} started with instance id ${instanceId}, checkpoint: ${JSONParser.stringify(lastCheckpoint)}`,\n        );\n        return {\n          lastCheckpoint,\n        };\n      }\n\n      return await processingScope(async (context) => {\n        if (hooks.onStart) {\n          console.log(\n            `Executing onStart hook for processor ${processorId} with instance id ${instanceId}`,\n          );\n          await hooks.onStart(context);\n        }\n\n        if (startFrom && startFrom !== 'CURRENT') {\n          console.log(\n            `Processor ${processorId} with instance id ${instanceId} starting from: ${JSONParser.stringify(startFrom)}`,\n          );\n          return startFrom;\n        }\n\n        if (checkpoints) {\n          const readResult = await checkpoints?.read(\n            {\n              processorId: processorId,\n              partition,\n            },\n            { ...startOptions, ...context },\n          );\n          lastCheckpoint = readResult.lastCheckpoint;\n        }\n\n        if (lastCheckpoint === null) {\n          console.log(\n            `Processor ${processorId} with instance id ${instanceId} starting from: BEGINNING`,\n          );\n          return 'BEGINNING';\n        }\n        console.log(\n          `Checkpoint read for processor ${processorId} with instance id ${instanceId}: ${JSONParser.stringify(lastCheckpoint)}`,\n        );\n\n        return {\n          lastCheckpoint,\n        };\n      }, startOptions);\n    },\n    close,\n    get isActive() {\n      return isActive;\n    },\n    handle: async (\n      messages: RecordedMessage<MessageType, MessageMetadataType>[],\n      partialContext: Partial<HandlerContext>,\n    ): Promise<MessageHandlerResult> => {\n      if (!isActive) return Promise.resolve();\n\n      try {\n        return await processingScope(async (context) => {\n          let result: MessageHandlerResult = undefined;\n\n          for (const message of messages) {\n            if (wasMessageHandled(message, lastCheckpoint)) continue;\n\n            const upcasted = upcastRecordedMessage(\n              // TODO: Make it smarter\n              message as unknown as RecordedMessage<\n                MessagePayloadType,\n                MessageMetadataType\n              >,\n              options.messageOptions?.schema?.versioning,\n            );\n\n            if (canHandle !== undefined && !canHandle.includes(upcasted.type))\n              continue;\n\n            const messageProcessingResult = await eachMessage(\n              upcasted,\n              context,\n            );\n\n            if (checkpoints) {\n              const storeCheckpointResult: StoreProcessorCheckpointResult<CheckpointType | null> =\n                await checkpoints.store(\n                  {\n                    processorId,\n                    version,\n                    message: upcasted,\n                    lastCheckpoint,\n                    partition,\n                  },\n                  context,\n                );\n\n              if (storeCheckpointResult.success) {\n                // TODO: Add correct handling of the storing checkpoint\n                lastCheckpoint = storeCheckpointResult.newCheckpoint;\n              }\n            }\n\n            if (\n              messageProcessingResult &&\n              messageProcessingResult.type === 'STOP'\n            ) {\n              isActive = false;\n              result = messageProcessingResult;\n              break;\n            }\n\n            if (stopAfter && stopAfter(upcasted)) {\n              isActive = false;\n              result = { type: 'STOP', reason: 'Stop condition reached' };\n              break;\n            }\n\n            if (\n              messageProcessingResult &&\n              messageProcessingResult.type === 'SKIP'\n            )\n              continue;\n          }\n\n          return result;\n        }, partialContext);\n      } catch (error) {\n        console.log(\n          `Error during message processing for processor ${processorId} with instance id ${instanceId}. Stopping the processor.`,\n          error,\n        );\n        isActive = false;\n        return {\n          type: 'STOP',\n          error: error as EmmettError,\n          reason: 'Error during message processing',\n        };\n      }\n    },\n  };\n};\n\nexport const projector = <\n  EventType extends Event = Event,\n  EventMetaDataType extends AnyRecordedMessageMetadata =\n    AnyRecordedMessageMetadata,\n  HandlerContext extends DefaultRecord = DefaultRecord,\n  CheckpointType =\n    GlobalPositionTypeOfRecordedMessageMetadata<EventMetaDataType>,\n  EventPayloadType extends Event = EventType,\n>(\n  options: ProjectorOptions<\n    EventType,\n    EventMetaDataType,\n    HandlerContext,\n    CheckpointType,\n    EventPayloadType\n  >,\n): MessageProcessor<\n  EventType,\n  EventMetaDataType,\n  HandlerContext,\n  CheckpointType\n> => {\n  const {\n    projection,\n    processorId = getProjectorId({\n      projectionName: projection.name ?? 'unknown',\n    }),\n    ...rest\n  } = options;\n\n  return reactor<\n    EventType,\n    EventMetaDataType,\n    HandlerContext,\n    CheckpointType,\n    EventPayloadType\n  >({\n    ...rest,\n    type: MessageProcessorType.PROJECTOR,\n    canHandle: projection.canHandle,\n    processorId,\n    messageOptions: options.projection.eventsOptions,\n    hooks: {\n      onInit: options.hooks?.onInit,\n      onStart:\n        (options.truncateOnStart && options.projection.truncate) ||\n        options.hooks?.onStart\n          ? async (context: HandlerContext) => {\n              if (options.truncateOnStart && options.projection.truncate)\n                await options.projection.truncate(context);\n\n              if (options.hooks?.onStart) await options.hooks?.onStart(context);\n            }\n          : undefined,\n      onClose: options.hooks?.onClose,\n    },\n    eachMessage: async (\n      event: RecordedMessage<EventType, EventMetaDataType>,\n      context: HandlerContext,\n    ) => projection.handle([event], context),\n  });\n};\n","import { getInMemoryDatabase, type InMemoryDatabase } from '../database';\nimport { EmmettError } from '../errors';\nimport {\n  type AnyEvent,\n  type AnyMessage,\n  type BatchRecordedMessageHandlerWithContext,\n  type MessageHandlerResult,\n  type ReadEventMetadataWithGlobalPosition,\n  type SingleRecordedMessageHandlerWithContext,\n} from '../typing';\nimport {\n  getCheckpoint,\n  MessageProcessor,\n  projector,\n  reactor,\n  type Checkpointer,\n  type MessageProcessingScope,\n  type ProjectorOptions,\n  type ReactorOptions,\n} from './processors';\n\nexport type InMemoryProcessorHandlerContext = {\n  database: InMemoryDatabase;\n};\n\nexport type InMemoryProcessor<MessageType extends AnyMessage = AnyMessage> =\n  MessageProcessor<\n    MessageType,\n    // TODO: generalize this to support other metadata types\n    ReadEventMetadataWithGlobalPosition,\n    InMemoryProcessorHandlerContext\n  > & { database: InMemoryDatabase };\n\nexport type InMemoryProcessorEachMessageHandler<\n  MessageType extends AnyMessage = AnyMessage,\n> = SingleRecordedMessageHandlerWithContext<\n  MessageType,\n  ReadEventMetadataWithGlobalPosition,\n  InMemoryProcessorHandlerContext\n>;\n\nexport type InMemoryProcessorEachBatchHandler<\n  MessageType extends AnyMessage = AnyMessage,\n> = BatchRecordedMessageHandlerWithContext<\n  MessageType,\n  ReadEventMetadataWithGlobalPosition,\n  InMemoryProcessorHandlerContext\n>;\n\nexport type InMemoryProcessorConnectionOptions = {\n  database?: InMemoryDatabase;\n};\n\ntype CheckpointDocument = {\n  _id: string;\n  lastCheckpoint: bigint | null;\n};\n\nexport type InMemoryCheckpointer<MessageType extends AnyMessage = AnyMessage> =\n  Checkpointer<\n    MessageType,\n    ReadEventMetadataWithGlobalPosition,\n    InMemoryProcessorHandlerContext\n  >;\n\nexport const inMemoryCheckpointer = <\n  MessageType extends AnyMessage = AnyMessage,\n>(): InMemoryCheckpointer<MessageType> => {\n  return {\n    read: async ({ processorId }, { database }) => {\n      const checkpoint = await database\n        .collection<CheckpointDocument>('emt_processor_checkpoints')\n        .findOne((d) => d._id === processorId);\n\n      return Promise.resolve({\n        lastCheckpoint: checkpoint?.lastCheckpoint ?? null,\n      });\n    },\n    store: async (context, { database }) => {\n      const { message, processorId, lastCheckpoint } = context;\n      const checkpoints = database.collection<CheckpointDocument>(\n        'emt_processor_checkpoints',\n      );\n\n      const checkpoint = await checkpoints.findOne(\n        (d) => d._id === processorId,\n      );\n\n      const currentPosition = checkpoint?.lastCheckpoint ?? null;\n\n      const newCheckpoint: bigint | null = getCheckpoint(message);\n\n      if (\n        currentPosition &&\n        (currentPosition === newCheckpoint ||\n          currentPosition !== lastCheckpoint)\n      ) {\n        return {\n          success: false,\n          reason:\n            currentPosition === newCheckpoint\n              ? 'IGNORED'\n              : newCheckpoint !== null && currentPosition > newCheckpoint\n                ? 'CURRENT_AHEAD'\n                : 'MISMATCH',\n        };\n      }\n\n      await checkpoints.handle(processorId, (existing) => ({\n        ...(existing ?? {}),\n        _id: processorId,\n        lastCheckpoint: newCheckpoint,\n      }));\n\n      return { success: true, newCheckpoint };\n    },\n  };\n};\n\ntype InMemoryConnectionOptions = {\n  connectionOptions?: InMemoryProcessorConnectionOptions;\n};\n\nexport type InMemoryReactorOptions<\n  MessageType extends AnyMessage = AnyMessage,\n> = ReactorOptions<\n  MessageType,\n  ReadEventMetadataWithGlobalPosition,\n  InMemoryProcessorHandlerContext\n> &\n  InMemoryConnectionOptions;\n\nexport type InMemoryProjectorOptions<EventType extends AnyEvent = AnyEvent> =\n  ProjectorOptions<\n    EventType,\n    ReadEventMetadataWithGlobalPosition,\n    InMemoryProcessorHandlerContext\n  > &\n    InMemoryConnectionOptions;\n\nexport type InMemoryProcessorOptions<\n  MessageType extends AnyMessage = AnyMessage,\n> =\n  | InMemoryReactorOptions<MessageType>\n  | InMemoryProjectorOptions<MessageType & AnyEvent>;\n\nconst inMemoryProcessingScope = (options: {\n  database: InMemoryDatabase | null;\n  processorId: string;\n}): MessageProcessingScope<InMemoryProcessorHandlerContext> => {\n  const processorDatabase = options.database;\n\n  const processingScope: MessageProcessingScope<\n    InMemoryProcessorHandlerContext\n  > = <Result = MessageHandlerResult>(\n    handler: (\n      context: InMemoryProcessorHandlerContext,\n    ) => Result | Promise<Result>,\n    partialContext: Partial<InMemoryProcessorHandlerContext>,\n  ) => {\n    const database = processorDatabase ?? partialContext?.database;\n\n    if (!database)\n      throw new EmmettError(\n        `InMemory processor '${options.processorId}' is missing database. Ensure that you passed it through options`,\n      );\n\n    return handler({ ...partialContext, database });\n  };\n\n  return processingScope;\n};\n\nexport const inMemoryProjector = <EventType extends AnyEvent = AnyEvent>(\n  options: InMemoryProjectorOptions<EventType>,\n): InMemoryProcessor<EventType> => {\n  const database = options.connectionOptions?.database ?? getInMemoryDatabase();\n\n  const hooks = {\n    onInit: options.hooks?.onInit,\n    onStart: options.hooks?.onStart,\n    onClose: options.hooks?.onClose\n      ? async (context: InMemoryProcessorHandlerContext) => {\n          if (options.hooks?.onClose) await options.hooks?.onClose(context);\n        }\n      : undefined,\n  };\n\n  const processor = projector<\n    EventType,\n    ReadEventMetadataWithGlobalPosition,\n    InMemoryProcessorHandlerContext\n  >({\n    ...options,\n    hooks,\n    processingScope: inMemoryProcessingScope({\n      database,\n      processorId:\n        options.processorId ?? `projection:${options.projection.name}`,\n    }),\n    checkpoints: inMemoryCheckpointer<EventType>(),\n  });\n\n  return Object.assign(processor, { database });\n};\n\nexport const inMemoryReactor = <MessageType extends AnyMessage = AnyMessage>(\n  options: InMemoryReactorOptions<MessageType>,\n): InMemoryProcessor<MessageType> => {\n  const database = options.connectionOptions?.database ?? getInMemoryDatabase();\n\n  const hooks = {\n    onInit: options.hooks?.onInit,\n    onStart: options.hooks?.onStart,\n    onClose: options.hooks?.onClose,\n  };\n\n  const processor = reactor({\n    ...options,\n    hooks,\n    processingScope: inMemoryProcessingScope({\n      database,\n      processorId: options.processorId,\n    }),\n    checkpoints: inMemoryCheckpointer<MessageType>(),\n  });\n\n  return Object.assign(processor, { database });\n};\n","import { EmmettError } from '../errors';\nimport type { EventStoreReadSchemaOptions } from '../eventStore';\nimport { JSONParser } from '../serialization';\nimport type {\n  AnyEvent,\n  AnyReadEventMetadata,\n  BatchRecordedMessageHandlerWithContext,\n  CanHandle,\n  DefaultRecord,\n  Event,\n} from '../typing';\nimport { arrayUtils } from '../utils';\n\nexport type ProjectionHandlingType = 'inline' | 'async';\n\nexport type ProjectionHandler<\n  EventType extends Event = AnyEvent,\n  EventMetaDataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n> = BatchRecordedMessageHandlerWithContext<\n  EventType,\n  EventMetaDataType,\n  ProjectionHandlerContext\n>;\n\nexport type TruncateProjection<\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n> = (context: ProjectionHandlerContext) => Promise<void>;\n\nexport type ProjectionInitOptions<\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n> = {\n  version: number;\n  status?: 'active' | 'inactive';\n  registrationType: ProjectionHandlingType;\n  context: ProjectionHandlerContext;\n};\n\nexport interface ProjectionDefinition<\n  EventType extends Event = AnyEvent,\n  EventMetaDataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n  EventPayloadType extends Event = EventType,\n> {\n  name?: string;\n  version?: number;\n  kind?: string;\n  canHandle: CanHandle<EventType>;\n  handle: ProjectionHandler<\n    EventType,\n    EventMetaDataType,\n    ProjectionHandlerContext\n  >;\n  truncate?: TruncateProjection<ProjectionHandlerContext>;\n  init?: (\n    options: ProjectionInitOptions<ProjectionHandlerContext>,\n  ) => void | Promise<void>;\n  eventsOptions?: {\n    schema?: EventStoreReadSchemaOptions<EventType, EventPayloadType>;\n  };\n}\n\nexport type ProjectionRegistration<\n  HandlingType extends ProjectionHandlingType,\n  ReadEventMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n> = {\n  type: HandlingType;\n  projection: ProjectionDefinition<\n    AnyEvent,\n    ReadEventMetadataType,\n    ProjectionHandlerContext,\n    AnyEvent\n  >;\n};\n\nexport const filterProjections = <\n  ReadEventMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n>(\n  type: ProjectionHandlingType,\n  projections: ProjectionRegistration<\n    ProjectionHandlingType,\n    ReadEventMetadataType,\n    ProjectionHandlerContext\n  >[],\n) => {\n  const inlineProjections = projections\n    .filter((projection) => projection.type === type)\n    .map(({ projection }) => projection);\n\n  const duplicateRegistrations = arrayUtils.getDuplicates(\n    inlineProjections,\n    (proj) => proj.name,\n  );\n\n  if (duplicateRegistrations.length > 0) {\n    throw new EmmettError(`You cannot register multiple projections with the same name (or without the name).\n      Ensure that:\n      ${JSONParser.stringify(duplicateRegistrations)}\n      have different names`);\n  }\n\n  return inlineProjections;\n};\n\nexport const projection = <\n  EventType extends Event = Event,\n  EventMetaDataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n  EventPayloadType extends Event = EventType,\n>(\n  definition: ProjectionDefinition<\n    EventType,\n    EventMetaDataType,\n    ProjectionHandlerContext,\n    EventPayloadType\n  >,\n): ProjectionDefinition<\n  EventType,\n  EventMetaDataType,\n  ProjectionHandlerContext,\n  EventPayloadType\n> => definition;\n\nexport const inlineProjections = <\n  ReadEventMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n>(\n  definitions: ProjectionDefinition<\n    // eslint-disable-next-line @typescript-eslint/no-explicit-any\n    any,\n    ReadEventMetadataType,\n    ProjectionHandlerContext\n  >[],\n): ProjectionRegistration<\n  'inline',\n  ReadEventMetadataType,\n  ProjectionHandlerContext\n>[] =>\n  definitions.map((definition) => ({\n    type: 'inline',\n    projection: definition,\n  }));\n\nexport const asyncProjections = <\n  ReadEventMetadataType extends AnyReadEventMetadata = AnyReadEventMetadata,\n  ProjectionHandlerContext extends DefaultRecord = DefaultRecord,\n>(\n  definitions: ProjectionDefinition<\n    AnyEvent,\n    ReadEventMetadataType,\n    ProjectionHandlerContext\n  >[],\n): ProjectionRegistration<\n  'inline',\n  ReadEventMetadataType,\n  ProjectionHandlerContext\n>[] =>\n  definitions.map((definition) => ({\n    type: 'inline',\n    projection: definition,\n  }));\n\nexport const projections = {\n  inline: inlineProjections,\n  async: asyncProjections,\n};\n"]}