From a12d22ffc4b8f0239f85de0908358db67d4b2ca5 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Thu, 11 Sep 2025 14:10:43 +0100 Subject: [PATCH 01/22] Support withTransaction using Futures --- src/di/container.ts | 2 +- src/lib/postgres.ts | 23 ++++++++++++++++++++++- 2 files changed, 23 insertions(+), 2 deletions(-) diff --git a/src/di/container.ts b/src/di/container.ts index 7d597fb..244655a 100644 --- a/src/di/container.ts +++ b/src/di/container.ts @@ -118,7 +118,7 @@ export async function configureDependencies(): Promise { poolSettings: defaultPoolSettings, }); - await postgres.withTransaction((transaction) => + await postgres.withTransactionP((transaction) => postgresEventStore.initialize({ transaction, database: env.EVENT_STORE_DATABASE_NAME, diff --git a/src/lib/postgres.ts b/src/lib/postgres.ts index ac973c0..72fa437 100644 --- a/src/lib/postgres.ts +++ b/src/lib/postgres.ts @@ -6,6 +6,7 @@ export { }; import { Pool, PoolConfig, PoolClient, QueryConfig, QueryResult } from 'pg'; +import { Future } from '@/lib/Future'; class PostgresTransaction { public closed: boolean = false; @@ -96,7 +97,7 @@ class Postgres { } // Execute an action with a transaction that will be automatically committed at the end. - async withTransaction( + async withTransactionP( f: (t: PostgresTransaction) => Promise, ): Promise { const connection = await this.pool.connect(); @@ -105,4 +106,24 @@ class Postgres { if (!transaction.closed) await transaction.commit(); return result; } + + // Execute an action with a transaction that will be automatically committed at the end. + withTransaction( + onConnectionError: (e: Error) => E, + f: (t: PostgresTransaction) => Future, + ): Future { + return Future.attemptP(() => this.pool.connect()) + .mapRej(onConnectionError) + .chain((connection) => { + const transaction = new PostgresTransaction(connection); + return f(transaction).chain((res) => { + if (!transaction.closed) { + Future.attemptP(transaction.commit) + .mapRej(onConnectionError) + .map(() => res); + } + return Future.resolve(res); + }); + }); + } } From 5f60b68da2c0153bc2284fca2068be6d8ceca506 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Thu, 11 Sep 2025 14:11:06 +0100 Subject: [PATCH 02/22] Implement commandHandler --- src/app/commandHandler.ts | 89 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 89 insertions(+) create mode 100644 src/app/commandHandler.ts diff --git a/src/app/commandHandler.ts b/src/app/commandHandler.ts new file mode 100644 index 0000000..7772656 --- /dev/null +++ b/src/app/commandHandler.ts @@ -0,0 +1,89 @@ +export { handleCommand }; + +import { Response } from '@/lib/router'; +import { EventStore } from '@/lib/eventSourcing/eventStore'; +import { Decoder, decode } from '@/lib/json/decoder'; +import * as express from 'express'; +import * as router from '@/lib/router'; +import { Future } from '@/lib/Future'; +import { Result, Failure } from '@/lib/Result'; +import { Postgres } from '@/lib/postgres'; +import { PostgresEventStore } from '@/app/postgresEventStore'; +import { Serializer } from '@/common/serializedEvent/Serializer'; +import { Deserializer } from '@/common/serializedEvent/Deserializer'; + +type Projections = {}; +type Services = {}; + +type CommandController = { + decoder: Decoder; + handler: (v: { + command: Command; + store: EventStore; + projections: Projections; + services: Services; + }) => Future; +}; + +function handleCommand( + serializer: Serializer, + deserializer: Deserializer, + eventStoreTable: string, + postgres: Postgres, + services: Services, + projections: Projections, + { decoder, handler }: CommandController, +): express.Handler { + return router.route((req) => + decodeCommand(decoder, req).chain((command) => + withEventStore( + postgres, + serializer, + deserializer, + eventStoreTable, + (store) => + handler({ + command, + store, + projections, + services, + }), + ), + ), + ); +} + +function decodeCommand( + decoder: Decoder, + req: express.Request, +): Future { + const decoded: Result = decode(decoder, req.body); + if (decoded instanceof Failure) { + return Future.reject( + router.json({ + status: 400, + content: { message: `Unable to decode command: ${decoded.error}` }, + }), + ); + } + + return Future.resolve(decoded.value); +} + +function withEventStore( + postgres: Postgres, + serializer: Serializer, + deserializer: Deserializer, + eventStoreTable: string, + f: (s: EventStore) => Future, +): Future { + const onError = (_: Error) => + router.json({ + status: 500, + content: { message: 'Internal Server Error' }, + }); + + return postgres.withTransaction(onError, (t) => + f(new PostgresEventStore(t, serializer, deserializer, eventStoreTable)), + ); +} From 4d884d60307f67e938ca4de3b2b1a80c6c6fc379 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Tue, 23 Sep 2025 13:59:35 +0100 Subject: [PATCH 03/22] Introduce time module --- package-lock.json | 10 +++ package.json | 1 + src/lib/time.ts | 180 ++++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 191 insertions(+) create mode 100644 src/lib/time.ts diff --git a/package-lock.json b/package-lock.json index 39307c5..ef43731 100644 --- a/package-lock.json +++ b/package-lock.json @@ -12,6 +12,7 @@ "class-validator": "^0.14.2", "express": "^4.18.2", "fluture": "^14.0.0", + "luxon": "^3.7.2", "minio": "^8.0.5", "mongodb": "^5.4.0", "nodemailer": "^7.0.5", @@ -3477,6 +3478,15 @@ "integrity": "sha512-6FlzubTLZG3J2a/NVCAleEhjzq5oxgHyaCU9yYXvcLsvoVaHJq/s5xXI6/XXP6tz7R9xAOtHnSO/tXtF3WRTlA==", "license": "MIT" }, + "node_modules/luxon": { + "version": "3.7.2", + "resolved": "https://registry.npmjs.org/luxon/-/luxon-3.7.2.tgz", + "integrity": "sha512-vtEhXh/gNjI9Yg1u4jX/0YVPMvxzHuGgCm6tC5kZyb08yjGWGnqAjGJvcXbqQR2P3MyMEFnRbpcdFS6PBcLqew==", + "license": "MIT", + "engines": { + "node": ">=12" + } + }, "node_modules/math-intrinsics": { "version": "1.1.0", "resolved": "https://registry.npmjs.org/math-intrinsics/-/math-intrinsics-1.1.0.tgz", diff --git a/package.json b/package.json index 670f21b..748a6cc 100644 --- a/package.json +++ b/package.json @@ -14,6 +14,7 @@ "class-validator": "^0.14.2", "express": "^4.18.2", "fluture": "^14.0.0", + "luxon": "^3.7.2", "minio": "^8.0.5", "mongodb": "^5.4.0", "nodemailer": "^7.0.5", diff --git a/src/lib/time.ts b/src/lib/time.ts new file mode 100644 index 0000000..bf1b02e --- /dev/null +++ b/src/lib/time.ts @@ -0,0 +1,180 @@ +// Sane utilities for dealing with time. +export { DateOnly, TimeOfDay, POSIX }; + +import { DateTime } from 'luxon'; +import * as s from '@/lib/json/schema'; +import { Failure, Success } from '@/lib/Result'; + +// POSIX time is the nominal time since 1970-01-01 00:00 UTC. +// Like DateTime, but without timezone confusion. +class POSIX { + static fromDate(d: Date): POSIX { + return new POSIX(d.valueOf()); + } + + static now(): POSIX { + return new POSIX(Date.now()); + } + + // The number of milliseconds for this date since midnight at the beginning + // of January 1, 1970, UTC. + value: number; + constructor(n: number) { + this.value = n; + } + + toDate(): Date { + return new Date(this.value); + } + + greaterThan(other: POSIX) { + return this.value > other.value; + } + + compare(other: POSIX): number { + return this.value > other.value ? 1 : this.value < other.value ? -1 : 0; + } + + static fromUTCDateAndTime(date: DateOnly, time: TimeOfDay): POSIX { + const s = `${date.pretty()}T${time.pretty()}Z`; + const luxonDate = DateTime.fromISO(s, { zone: 'UTC' }); + return POSIX.fromDate(luxonDate.toJSDate()); + } + + toUTCDateAndTime(): { date: DateOnly; time: TimeOfDay } { + const dt = DateTime.fromMillis(this.value, { zone: 'UTC' }); + const date = new DateOnly(dt.year, dt.month, dt.day); + const time = TimeOfDay.fromParts({ + hours: dt.hour, + minutes: dt.minute, + seconds: dt.second, + }); + return { date, time }; + } + + toLocalDateAndTime(): { date: DateOnly; time: TimeOfDay } { + const dt = DateTime.fromMillis(this.value, { zone: 'UTC' }).toLocal(); + const date = new DateOnly(dt.year, dt.month, dt.day); + const time = TimeOfDay.fromParts({ + hours: dt.hour, + minutes: dt.minute, + seconds: dt.second, + }); + return { date, time }; + } + + static schema: s.Schema = s.number.dimap( + (n) => new POSIX(n), + (p) => p.value, + ); +} + +// pad to two digits +const padded = (v: number) => v.toString().padStart(2, '0'); + +// Year, month and day. +class DateOnly { + readonly year: number; + readonly month: number; // 1-12 + readonly day: number; // 1-30ish + + constructor(year: number, month: number, day: number) { + this.year = year; + this.month = month; + this.day = day; + } + + static today(): DateOnly { + return DateOnly.fromDate(new Date()); + } + + static fromDate(date: Date): DateOnly { + return new DateOnly( + date.getFullYear(), + date.getMonth() + 1, + date.getDate(), + ); + } + + pretty() { + return `${this.year}-${padded(this.month)}-${padded(this.day)}`; + } + + static schema: s.Schema = s.string.then( + (s) => { + const parts = s.split('-'); + if (parts.length !== 3) { + return Failure('Invalid Date'); + } + const year = parseInt(parts[0] as string, 10); + const month = parseInt(parts[1] as string, 10); + const day = parseInt(parts[2] as string, 10); + + if (isNaN(year) || isNaN(month) || isNaN(day)) { + return Failure('Invalid Date'); + } + + return Success(new DateOnly(year, month, day)); + }, + (date) => date.pretty(), + ); + + greaterThan(other: DateOnly) { + return this.compare(other) == 1; + } + + compare(other: DateOnly): number { + return this.year > other.year + ? 1 + : this.year < other.year + ? -1 + : this.month > other.month + ? 1 + : this.month < other.month + ? -1 + : this.day > other.day + ? 1 + : this.day < other.day + ? -1 + : 0; + } + + addMonths(months: number): DateOnly { + const luxonDate = DateTime.fromObject({ + year: this.year, + month: this.month, + day: this.day, + }); + const newLuxonDate = luxonDate.plus({ months }); + + return new DateOnly( + newLuxonDate.year, + newLuxonDate.month, + newLuxonDate.day, + ); + } +} + +class TimeOfDay { + constructor(readonly seconds: number) {} + + static fromParts({ + hours, + minutes, + seconds, + }: { + hours: number; + minutes: number; + seconds: number; + }): TimeOfDay { + return new TimeOfDay(hours * 60 * 60 + minutes * 60 + seconds); + } + + // HH:MM:SS + pretty() { + const hours = padded(Math.floor(this.seconds / (60 * 60))); + const minutes = padded(Math.floor(this.seconds / 60) % 60); + const seconds = padded(this.seconds % 60); + return `${hours}:${minutes}:${seconds}`; + } +} From c15baa21b9e9102edfb1ff3ec722c70edce539c5 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Tue, 23 Sep 2025 13:59:59 +0100 Subject: [PATCH 04/22] Introduce eventSourcing/event --- src/lib/eventSourcing/event.ts | 55 ++++++++++++++++++++++++++++++++++ 1 file changed, 55 insertions(+) create mode 100644 src/lib/eventSourcing/event.ts diff --git a/src/lib/eventSourcing/event.ts b/src/lib/eventSourcing/event.ts new file mode 100644 index 0000000..a0099fa --- /dev/null +++ b/src/lib/eventSourcing/event.ts @@ -0,0 +1,55 @@ +export { + Event, + Aggregate, + TransformationEvent, + CreationEvent, + type EventInfo, + EventInfo_schema, +}; + +import * as s from '@/lib/json/schema'; +import { Schema } from '@/lib/json/schema'; +import { POSIX } from '@/lib/time'; + +// @ts-ignore +class Id { + // @ts-expect-error _tag's existence prevents structural comparison + private readonly _tag: null = null; + + constructor(public value: string) {} + + static schema(): Schema> { + return s.string.dimap( + (v) => new Id(v), + (id) => id.value, + ); + } +} + +// Class which all events derive from. Used for type constraints. +abstract class Aggregate {} + +// Class which all events derive from. Used for type constraints. +abstract class Event<_T extends Aggregate> {} + +// The first event for an aggregate. +abstract class CreationEvent extends Event { + abstract createAggregate(): T; +} + +// Any event that is not the first one for an aggregate. +abstract class TransformationEvent extends Event { + abstract transformAggregate(aggregate: T): T; +} + +// Information about an event. Not the event payload. +type EventInfo = s.Infer; + +const EventInfo_schema = s.object({ + event_id: Id.schema>(), + aggregate_id: Id.schema(), + aggregate_version: s.number, + correlation_id: Id.schema>(), + causation_id: Id.schema>(), + recorded_on: POSIX.schema, +}); From 8e39f2e3568ae32764dc8b194014963ccc193c7f Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Tue, 23 Sep 2025 15:30:34 +0100 Subject: [PATCH 05/22] Add watch command --- package.json | 1 + 1 file changed, 1 insertion(+) diff --git a/package.json b/package.json index 748a6cc..7ff2d50 100644 --- a/package.json +++ b/package.json @@ -4,6 +4,7 @@ "main": "dist/index.js", "scripts": { "build": "tsc", + "watch": "tsc --watch", "start": "tsx src/index.ts", "test": "tsx tests/unit/main.ts", "format": "prettier --write .", From f7d6ada12f01ccc70e374a3204fa1da6f565bc69 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Tue, 23 Sep 2025 15:31:14 +0100 Subject: [PATCH 06/22] Add luxon types --- package-lock.json | 8 ++++++++ package.json | 1 + 2 files changed, 9 insertions(+) diff --git a/package-lock.json b/package-lock.json index ef43731..11c8d06 100644 --- a/package-lock.json +++ b/package-lock.json @@ -24,6 +24,7 @@ }, "devDependencies": { "@types/express": "^4.17.21", + "@types/luxon": "^3.7.1", "@types/node": "^24.3.0", "@types/nodemailer": "^7.0.1", "@types/pg": "^8.6.6", @@ -2358,6 +2359,13 @@ "dev": true, "license": "MIT" }, + "node_modules/@types/luxon": { + "version": "3.7.1", + "resolved": "https://registry.npmjs.org/@types/luxon/-/luxon-3.7.1.tgz", + "integrity": "sha512-H3iskjFIAn5SlJU7OuxUmTEpebK6TKB8rxZShDslBMZJ5u9S//KM1sbdAisiSrqwLQncVjnpi2OK2J51h+4lsg==", + "dev": true, + "license": "MIT" + }, "node_modules/@types/mime": { "version": "1.3.5", "resolved": "https://registry.npmjs.org/@types/mime/-/mime-1.3.5.tgz", diff --git a/package.json b/package.json index 7ff2d50..8ef9749 100644 --- a/package.json +++ b/package.json @@ -31,6 +31,7 @@ }, "devDependencies": { "@types/express": "^4.17.21", + "@types/luxon": "^3.7.1", "@types/node": "^24.3.0", "@types/nodemailer": "^7.0.1", "@types/pg": "^8.6.6", From 80a1667e0bda9f267e1581b5e399e61efb342860 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 13:56:29 +0100 Subject: [PATCH 07/22] Change signature of Decoder.then --- src/lib/json/decoder.ts | 51 +++++++++++++++++++++++------------------ src/lib/json/schema.ts | 8 ++----- src/lib/time.ts | 12 +++++----- 3 files changed, 37 insertions(+), 34 deletions(-) diff --git a/src/lib/json/decoder.ts b/src/lib/json/decoder.ts index aca6e8d..080939c 100644 --- a/src/lib/json/decoder.ts +++ b/src/lib/json/decoder.ts @@ -49,7 +49,9 @@ export { triple, always, fail, + failure, optional, + succeed, }; import { Result, Success, Failure, traverse } from '@/lib/Result'; @@ -71,8 +73,8 @@ class Decoder { this.run = run; } - then(f: (v: T) => DecodeResult): Decoder { - return new Decoder((v) => this.run(v).then(f)); + then(f: (v: T) => Decoder): Decoder { + return new Decoder((u) => this.run(u).then((v) => f(v).run(u))); } map(f: (v: T) => W): Decoder { @@ -91,39 +93,45 @@ function showPath([path, error]: [Path, string]): string { type DecodeResult = Result<[Path, string], T>; type Path = List; -const fail = (msg: string): DecodeResult => Failure([List.empty(), msg]); +const failure = (msg: string): DecodeResult => + Failure([List.empty(), msg]); + +const fail = (msg: string): Decoder => + new Decoder((_) => Failure([List.empty(), msg])); const always = (v: T): Decoder => new Decoder((_) => Success(v)); +const succeed = always; + const any: Decoder = new Decoder((v) => Success(v)); const string: Decoder = new Decoder((v) => typeof v === 'string' ? Success(v) - : fail('expected string but found ' + typeof v), + : failure('expected string but found ' + typeof v), ); const number: Decoder = new Decoder((v) => typeof v === 'number' ? Success(v) - : fail('expected number but found ' + typeof v), + : failure('expected number but found ' + typeof v), ); const stringNumber: Decoder = string.then((s) => { const v = parseInt(s, 10); - return isNaN(v) ? fail('not a valid number: ' + s) : Success(v); + return isNaN(v) ? fail('not a valid number: ' + s) : succeed(v); }); const boolean: Decoder = new Decoder((v) => typeof v === 'boolean' ? Success(v) - : fail('expected boolean but found ' + typeof v), + : failure('expected boolean but found ' + typeof v), ); const array = (decodeValue: Decoder): Decoder> => new Decoder((input) => { if (!Array.isArray(input)) { - return fail('expected array but found ' + typeof input); + return failure('expected array but found ' + typeof input); } return traverse(List.from(input), decodeValue.run).map((list) => @@ -139,7 +147,7 @@ type DecoderDef = { const object = (decoders: DecoderDef): Decoder => new Decoder((input) => { if (typeof input !== 'object' || input === null) { - return fail('expected object but found ' + typeof input); + return failure('expected object but found ' + typeof input); } const obj = input as { [P in keyof A]: unknown }; @@ -168,7 +176,7 @@ type ObjectMap = { [x: string]: A }; const objectMap = (decoder: Decoder): Decoder> => new Decoder((input) => { if (typeof input !== 'object' || input === null) { - return fail('expected object but found ' + typeof input); + return failure('expected object but found ' + typeof input); } const result = {} as ObjectMap; @@ -197,10 +205,10 @@ const pair = ( ): Decoder<[L, R]> => new Decoder((input) => { if (!Array.isArray(input)) { - return fail('expected array but found ' + typeof input); + return failure('expected array but found ' + typeof input); } if (input.length !== 2) { - return fail( + return failure( 'expected array with 2 elements but it found ' + input.length, ); } @@ -218,10 +226,10 @@ const triple = ( ): Decoder<[A, B, C]> => new Decoder((input) => { if (!Array.isArray(input)) { - return fail('expected array but found ' + typeof input); + return failure('expected array but found ' + typeof input); } if (input.length !== 3) { - return fail( + return failure( 'expected array with 3 elements but it found ' + input.length, ); } @@ -236,7 +244,7 @@ const triple = ( const oneOf = (decoders: Array>): Decoder => new Decoder((input) => { - let decoded: DecodeResult = fail('no decoders'); + let decoded: DecodeResult = failure('no decoders'); const errors: Array<[Path, string]> = []; @@ -248,11 +256,10 @@ const oneOf = (decoders: Array>): Decoder => errors.push(decoded.error); } - const failure = Failure<[Path, string], V>([ + return Failure<[Path, string], V>([ List.empty(), errors.map(showPath).join('\n'), ]); - return failure; }); const maybe = (decoder: Decoder): Decoder> => @@ -262,19 +269,19 @@ const nullable = (decoder: Decoder): Decoder> => oneOf([nullP, decoder]); const nullP: Decoder = new Decoder((v) => - v === null ? Success(null) : fail('expected null but found ' + typeof v), + v === null ? Success(null) : failure('expected null but found ' + typeof v), ); const undefinedP: Decoder = new Decoder((v) => v === undefined ? Success(undefined) - : fail('expected `undefined` ' + typeof v), + : failure('expected `undefined` ' + typeof v), ); // Useful for parsing tag names in discriminated unions. const stringLiteral = (str: T): Decoder => new Decoder((v) => - v === str ? Success(v as T) : fail(`expected '${str}' but found '${v}'`), + v === str ? Success(v as T) : failure(`expected '${str}' but found '${v}'`), ); // An object field that may be absent. @@ -283,8 +290,8 @@ const optional = (decoder: Decoder): Decoder> => // Define a recursive decoder function rec(f: (p: Decoder) => Decoder): Decoder { - const base: Decoder = new Decoder((_) => - fail('A recursive decoder cannot immediately call itself.'), + const base: Decoder = fail( + 'A recursive decoder cannot immediately call itself.', ); const top = f(base); // @ts-expect-error will complain that 'run' is readonly. But we are doing this on purpose here. diff --git a/src/lib/json/schema.ts b/src/lib/json/schema.ts index 7714521..90b9d12 100644 --- a/src/lib/json/schema.ts +++ b/src/lib/json/schema.ts @@ -34,7 +34,6 @@ import * as D from '@/lib/json/decoder'; import { Encoder, EncoderDef } from '@/lib/json/encoder'; import { Json } from '@/lib/json/types'; import * as E from '@/lib/json/encoder'; -import { List } from '@/lib/List'; import { Maybe, Nullable } from '@/lib/Maybe'; // Infer the type from a schema definition @@ -54,11 +53,8 @@ class Schema { return new Schema(this.decoder.map(p), this.encoder.rmap(s)); } - then(p: (v: A) => Result, s: (v: W) => A): Schema { - return new Schema( - this.decoder.then((v) => p(v).mapFailure((e) => [List.empty(), e])), - this.encoder.rmap(s), - ); + then(p: (v: A) => Decoder, s: (v: W) => A): Schema { + return new Schema(this.decoder.then(p), this.encoder.rmap(s)); } } diff --git a/src/lib/time.ts b/src/lib/time.ts index bf1b02e..aeb4845 100644 --- a/src/lib/time.ts +++ b/src/lib/time.ts @@ -3,7 +3,7 @@ export { DateOnly, TimeOfDay, POSIX }; import { DateTime } from 'luxon'; import * as s from '@/lib/json/schema'; -import { Failure, Success } from '@/lib/Result'; +import * as d from '@/lib/json/decoder'; // POSIX time is the nominal time since 1970-01-01 00:00 UTC. // Like DateTime, but without timezone confusion. @@ -101,20 +101,20 @@ class DateOnly { } static schema: s.Schema = s.string.then( - (s) => { - const parts = s.split('-'); + (str) => { + const parts = str.split('-'); if (parts.length !== 3) { - return Failure('Invalid Date'); + return d.fail('Invalid Date'); } const year = parseInt(parts[0] as string, 10); const month = parseInt(parts[1] as string, 10); const day = parseInt(parts[2] as string, 10); if (isNaN(year) || isNaN(month) || isNaN(day)) { - return Failure('Invalid Date'); + return d.fail('Invalid Date'); } - return Success(new DateOnly(year, month, day)); + return d.succeed(new DateOnly(year, month, day)); }, (date) => date.pretty(), ); From dc648e16598818c94c4776559b0f9779a253d693 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 19:28:48 +0100 Subject: [PATCH 08/22] Add Infer to Maybe --- src/lib/Maybe.ts | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/src/lib/Maybe.ts b/src/lib/Maybe.ts index fc2447b..f8be9f7 100644 --- a/src/lib/Maybe.ts +++ b/src/lib/Maybe.ts @@ -17,6 +17,7 @@ Values can be extracted using `instsanceof` tests. export { type Maybe, type Nullable, + type Infer, CallableJust as Just, CallableNothing as Nothing, from, @@ -29,6 +30,9 @@ import Callable from '@/lib/Callable'; type Maybe = Just | Nothing; type Nullable = T | null; +// Infer the type from a Maybe definition +type Infer> = A extends Maybe ? B : never; + // prettier-ignore export interface IMaybe { isJust() : boolean; From 1cc084c3e5e84c49eeb670ad7afca2ae6f14d726 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 19:29:17 +0100 Subject: [PATCH 09/22] Parameterise Event --- src/lib/eventSourcing/event.ts | 22 +++++++++++++--------- src/lib/eventSourcing/eventStore.ts | 4 ++-- 2 files changed, 15 insertions(+), 11 deletions(-) diff --git a/src/lib/eventSourcing/event.ts b/src/lib/eventSourcing/event.ts index a0099fa..056e4cb 100644 --- a/src/lib/eventSourcing/event.ts +++ b/src/lib/eventSourcing/event.ts @@ -1,10 +1,11 @@ export { Event, - Aggregate, + type Aggregate, TransformationEvent, CreationEvent, type EventInfo, EventInfo_schema, + Id, }; import * as s from '@/lib/json/schema'; @@ -27,18 +28,21 @@ class Id { } // Class which all events derive from. Used for type constraints. -abstract class Aggregate {} +interface Aggregate { + readonly aggregateId: Id>; + readonly aggregateVersion: number; +} // Class which all events derive from. Used for type constraints. -abstract class Event<_T extends Aggregate> {} +abstract class Event<_T extends Aggregate<_T>> {} // The first event for an aggregate. -abstract class CreationEvent extends Event { +abstract class CreationEvent> extends Event { abstract createAggregate(): T; } // Any event that is not the first one for an aggregate. -abstract class TransformationEvent extends Event { +abstract class TransformationEvent> extends Event { abstract transformAggregate(aggregate: T): T; } @@ -46,10 +50,10 @@ abstract class TransformationEvent extends Event { type EventInfo = s.Infer; const EventInfo_schema = s.object({ - event_id: Id.schema>(), - aggregate_id: Id.schema(), + event_id: Id.schema>>(), + aggregate_id: Id.schema>(), aggregate_version: s.number, - correlation_id: Id.schema>(), - causation_id: Id.schema>(), + correlation_id: Id.schema>>(), + causation_id: Id.schema>>(), recorded_on: POSIX.schema, }); diff --git a/src/lib/eventSourcing/eventStore.ts b/src/lib/eventSourcing/eventStore.ts index ecb7a90..53a5b0d 100644 --- a/src/lib/eventSourcing/eventStore.ts +++ b/src/lib/eventSourcing/eventStore.ts @@ -1,5 +1,5 @@ export { type EventStore }; -import { Event } from '@/common/event/Event'; +import { Event } from '@/lib/eventSourcing/event'; import { Aggregate } from '@/common/aggregate/Aggregate'; import { AggregateAndEventIdsInLastEvent } from '@/common/eventStore/AggregateAndEventIdsInLastEvent'; @@ -23,7 +23,7 @@ interface EventStore { aggregateId: string, ): Promise>; - saveEvent(event: Event): Promise; + saveEvent(event: Event): Promise; doesEventAlreadyExist(eventId: string): Promise; } From 068a06af956f63476f1861ff35431c1e0968a411 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 19:53:33 +0100 Subject: [PATCH 10/22] Implement event decoding --- src/app/event.ts | 123 +++++++++++++++++++++++++++++++++++++++++++++++ 1 file changed, 123 insertions(+) create mode 100644 src/app/event.ts diff --git a/src/app/event.ts b/src/app/event.ts new file mode 100644 index 0000000..b446d43 --- /dev/null +++ b/src/app/event.ts @@ -0,0 +1,123 @@ +import { + Aggregate, + Id, + CreationEvent, + TransformationEvent, +} from '@/lib/eventSourcing/event'; + +import { Decoder } from '@/lib/json/decoder'; +import * as s from '@/lib/json/schema'; +import * as d from '@/lib/json/decoder'; +import { Schema } from '@/lib/json/schema'; +import { Maybe, Nothing, Just } from '@/lib/Maybe'; +import * as m from '@/lib/Maybe'; + +class User implements Aggregate { + constructor( + readonly aggregateId: Id, + readonly aggregateVersion: number, + readonly name: string, + ) {} +} + +class CreateUser implements CreationEvent { + static type: 'CreateUserr' = 'CreateUserr'; + constructor(readonly values: s.Infer) {} + static schemaArgs = s.object({ + type: s.stringLiteral(CreateUser.type), + name: s.string, + }); + + static schema = CreateUser.schemaArgs.dimap( + (v) => new CreateUser(v), + (v) => v.values, + ); + + createAggregate() { + return new User(new Id('wat'), 0, this.values.name); + } +} + +class AddName implements TransformationEvent { + static type: 'AddName' = 'AddName'; + constructor(readonly values: s.Infer) {} + + static schemaArgs = s.object({ + type: s.stringLiteral(AddName.type), + name: s.string, + }); + + static schema = AddName.schemaArgs.dimap( + (v) => new AddName(v), + (v) => v.values, + ); + + transformAggregate(agg: User): User { + const u = new User( + agg.aggregateId, + agg.aggregateVersion + 1, + this.values.name, + ); + return u; + } +} + +function decodeEvent( + ts: Array<{ schemaS: Schema; type: string }>, +): Decoder> { + return d.object({ type: d.string }).then(({ type: ty }) => { + const found = ts.find((t) => t.type == ty); + if (found === undefined) { + return d.succeed(Nothing()); + } + + return found.schemaS.decoder.map(Just) as Decoder>; + }); +} + +type Accepted }> = d.Infer; + +type W = d.Infer; + +const dd = accept({ + [AddName.type]: AddName.schema.decoder, + [CreateUser.type]: CreateUser.schema.decoder, +}); + +function accept }>( + ds: T, +): Decoder>> { + return d + .object({ type: d.string }) + .then(({ type: ty }): Decoder>> => { + const decoder: undefined | Decoder> = ds[ty as keyof T]; + return (decoder ? decoder.map(Just) : d.succeed(Nothing())) as Decoder< + Maybe> + >; + }); +} + +// ---------- + +type Accepted2 = T[number]; + +const vv = accept([AddName, CreateUser]); + +type WW = m.Infer>; + +type EventClass = { type: string; schema: Schema }; + +// Given some event classes, creates a decoder for those classes. +// Makes sure to error if decoding those class object fail, but +// succeeds if the encoded event was of another class. +function accept( + ts: T, +): Decoder>> { + type Ty = s.Infer; + return d.object({ type: d.string }).then(({ type: ty }) => { + const c: undefined | EventClass = ts.find((t) => t.type === ty); + return c + ? (c.schema.decoder.map(Just) as Decoder>) + : d.succeed(Nothing()); + }); +} From 175b66a7e2d401cd864da935fc2c6bce139c2606 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 19:53:40 +0100 Subject: [PATCH 11/22] Implement projection handler --- src/app/projectionHandler.ts | 80 ++++++++++++++++++++++++++++++++++++ 1 file changed, 80 insertions(+) create mode 100644 src/app/projectionHandler.ts diff --git a/src/app/projectionHandler.ts b/src/app/projectionHandler.ts new file mode 100644 index 0000000..d3c1fb1 --- /dev/null +++ b/src/app/projectionHandler.ts @@ -0,0 +1,80 @@ +export { handleProjection }; + +import { Response } from '@/lib/router'; +import { Event } from '@/lib/eventSourcing/event'; +import { Decoder, decode } from '@/lib/json/decoder'; +import * as express from 'express'; +import * as router from '@/lib/router'; +import * as d from '@/lib/json/decoder'; +import { Future } from '@/lib/Future'; +import { Result, Failure } from '@/lib/Result'; +import { Maybe, Nothing } from '@/lib/Maybe'; + +type Projections = {}; +type ProjectionStore = {}; +type Mongo = {}; + +type ProjectionController> = { + decoder: Decoder>; + handler: (v: { + event: E; + projections: Projections; + store: ProjectionStore; + }) => Future; +}; + +function handleProjection>( + projections: Projections, + mongo: Mongo, + { decoder, handler }: ProjectionController, +): express.Handler { + return router.route((req) => + decodeEvent(decoder, req).chain((event) => + withProjectionStore(mongo, (store) => + handler({ + event, + projections, + store, + }), + ), + ), + ); +} + +function decodeEvent( + decoder: Decoder>, + req: express.Request, +): Future { + const bodyDecoder: Decoder> = d + .object({ payload: decoder }) + .map((r) => r.payload); + + const decoded: Result> = decode(bodyDecoder, req.body); + + if (decoded instanceof Failure) { + return Future.reject( + router.json({ + status: 400, + content: { message: `Unable to decode command: ${decoded.error}` }, + }), + ); + } + + if (decoded.value instanceof Nothing) { + return Future.reject( + router.json({ + status: 200, + content: { message: 'Ignored' }, + }), + ); + } + + return Future.resolve(decoded.value.value); +} + +function withProjectionStore( + _mongo: Mongo, + _f: (s: ProjectionStore) => Future, +): Future { + throw new Error('TODO'); +} From 959c789f11ad38edc1ada2fe2481da130000d570 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 20:01:15 +0100 Subject: [PATCH 12/22] Implement query handler --- src/app/queryHandler.ts | 70 +++++++++++++++++++++++++++++++++++++++++ 1 file changed, 70 insertions(+) create mode 100644 src/app/queryHandler.ts diff --git a/src/app/queryHandler.ts b/src/app/queryHandler.ts new file mode 100644 index 0000000..e1ecb68 --- /dev/null +++ b/src/app/queryHandler.ts @@ -0,0 +1,70 @@ +export { handleQuery }; + +import { Response } from '@/lib/router'; +import { Event } from '@/lib/eventSourcing/event'; +import { Decoder, decode } from '@/lib/json/decoder'; +import * as express from 'express'; +import * as router from '@/lib/router'; +import * as d from '@/lib/json/decoder'; +import { Future } from '@/lib/Future'; +import { Result, Failure } from '@/lib/Result'; +import { Maybe, Nothing } from '@/lib/Maybe'; + +type Projections = {}; +type Services = {}; + +type QueryController> = { + decoder: Decoder>; + handler: (v: { + event: E; + projections: Projections; + services: Services; + }) => Future; +}; + +function handleQuery>( + projections: Projections, + services: Services, + { decoder, handler }: QueryController, +): express.Handler { + return router.route((req) => + decodeEvent(decoder, req).chain((event) => + handler({ + event, + projections, + services, + }), + ), + ); +} + +function decodeEvent( + decoder: Decoder>, + req: express.Request, +): Future { + const bodyDecoder: Decoder> = d + .object({ payload: decoder }) + .map((r) => r.payload); + + const decoded: Result> = decode(bodyDecoder, req.body); + + if (decoded instanceof Failure) { + return Future.reject( + router.json({ + status: 400, + content: { message: `Unable to decode command: ${decoded.error}` }, + }), + ); + } + + if (decoded.value instanceof Nothing) { + return Future.reject( + router.json({ + status: 200, + content: { message: 'Ignored' }, + }), + ); + } + + return Future.resolve(decoded.value.value); +} From 943a045bf0e8eaa187a9cecce51fb2107a9a8af2 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 20:01:25 +0100 Subject: [PATCH 13/22] Implement reaction handler --- src/app/reactionHandler.ts | 106 +++++++++++++++++++++++++++++++++++++ 1 file changed, 106 insertions(+) create mode 100644 src/app/reactionHandler.ts diff --git a/src/app/reactionHandler.ts b/src/app/reactionHandler.ts new file mode 100644 index 0000000..3d3ae48 --- /dev/null +++ b/src/app/reactionHandler.ts @@ -0,0 +1,106 @@ +export { handleReaction }; + +import { Response } from '@/lib/router'; +import { EventStore } from '@/lib/eventSourcing/eventStore'; +import { Event } from '@/lib/eventSourcing/event'; +import { Decoder, decode } from '@/lib/json/decoder'; +import * as express from 'express'; +import * as router from '@/lib/router'; +import * as d from '@/lib/json/decoder'; +import { Future } from '@/lib/Future'; +import { Result, Failure } from '@/lib/Result'; +import { Maybe, Nothing } from '@/lib/Maybe'; +import { Postgres } from '@/lib/postgres'; +import { PostgresEventStore } from '@/app/postgresEventStore'; +import { Serializer } from '@/common/serializedEvent/Serializer'; +import { Deserializer } from '@/common/serializedEvent/Deserializer'; + +type Projections = {}; +type Services = {}; + +type ReactionController> = { + decoder: Decoder>; + handler: (v: { + event: E; + projections: Projections; + services: Services; + store: EventStore; + }) => Future; +}; + +function handleReaction>( + projections: Projections, + services: Services, + postgres: Postgres, + serializer: Serializer, + deserializer: Deserializer, + eventStoreTable: string, + { decoder, handler }: ReactionController, +): express.Handler { + return router.route((req) => + decodeEvent(decoder, req).chain((event) => + withEventStore( + postgres, + serializer, + deserializer, + eventStoreTable, + (store) => + handler({ + event, + projections, + services, + store, + }), + ), + ), + ); +} + +function decodeEvent( + decoder: Decoder>, + req: express.Request, +): Future { + const bodyDecoder: Decoder> = d + .object({ payload: decoder }) + .map((r) => r.payload); + + const decoded: Result> = decode(bodyDecoder, req.body); + + if (decoded instanceof Failure) { + return Future.reject( + router.json({ + status: 400, + content: { message: `Unable to decode command: ${decoded.error}` }, + }), + ); + } + + if (decoded.value instanceof Nothing) { + return Future.reject( + router.json({ + status: 200, + content: { message: 'Ignored' }, + }), + ); + } + + return Future.resolve(decoded.value.value); +} + +function withEventStore( + postgres: Postgres, + serializer: Serializer, + deserializer: Deserializer, + eventStoreTable: string, + f: (s: EventStore) => Future, +): Future { + const onError = (_: Error) => + router.json({ + status: 500, + content: { message: 'Internal Server Error' }, + }); + + return postgres.withTransaction(onError, (t) => + f(new PostgresEventStore(t, serializer, deserializer, eventStoreTable)), + ); +} From 6584bf67af5a660762351916a941285c0d971742 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 22:54:07 +0100 Subject: [PATCH 14/22] Consolidate user of decodeEvent --- src/app/event.ts | 39 ++---------------------------------- src/app/projectionHandler.ts | 23 +++++++++++++++++++-- src/app/queryHandler.ts | 38 ++++++++++------------------------- src/app/reactionHandler.ts | 38 +++-------------------------------- 4 files changed, 37 insertions(+), 101 deletions(-) diff --git a/src/app/event.ts b/src/app/event.ts index b446d43..85abeec 100644 --- a/src/app/event.ts +++ b/src/app/event.ts @@ -1,3 +1,5 @@ +export { accept, type EventClass }; + import { Aggregate, Id, @@ -62,45 +64,8 @@ class AddName implements TransformationEvent { } } -function decodeEvent( - ts: Array<{ schemaS: Schema; type: string }>, -): Decoder> { - return d.object({ type: d.string }).then(({ type: ty }) => { - const found = ts.find((t) => t.type == ty); - if (found === undefined) { - return d.succeed(Nothing()); - } - - return found.schemaS.decoder.map(Just) as Decoder>; - }); -} - -type Accepted }> = d.Infer; - -type W = d.Infer; - -const dd = accept({ - [AddName.type]: AddName.schema.decoder, - [CreateUser.type]: CreateUser.schema.decoder, -}); - -function accept }>( - ds: T, -): Decoder>> { - return d - .object({ type: d.string }) - .then(({ type: ty }): Decoder>> => { - const decoder: undefined | Decoder> = ds[ty as keyof T]; - return (decoder ? decoder.map(Just) : d.succeed(Nothing())) as Decoder< - Maybe> - >; - }); -} - // ---------- -type Accepted2 = T[number]; - const vv = accept([AddName, CreateUser]); type WW = m.Infer>; diff --git a/src/app/projectionHandler.ts b/src/app/projectionHandler.ts index d3c1fb1..ec95c6f 100644 --- a/src/app/projectionHandler.ts +++ b/src/app/projectionHandler.ts @@ -1,4 +1,4 @@ -export { handleProjection }; +export { handleProjection, decodeEvent, accept }; import { Response } from '@/lib/router'; import { Event } from '@/lib/eventSourcing/event'; @@ -8,7 +8,9 @@ import * as router from '@/lib/router'; import * as d from '@/lib/json/decoder'; import { Future } from '@/lib/Future'; import { Result, Failure } from '@/lib/Result'; -import { Maybe, Nothing } from '@/lib/Maybe'; +import { Maybe, Nothing, Just } from '@/lib/Maybe'; +import { Schema } from '@/lib/json/schema'; +import * as s from '@/lib/json/schema'; type Projections = {}; type ProjectionStore = {}; @@ -78,3 +80,20 @@ function withProjectionStore( ): Future { throw new Error('TODO'); } + +type EventClass = { type: string; schema: Schema }; + +// Given some event classes, creates a decoder for those classes. +// Makes sure to error if decoding those class object fail, but +// succeeds if the encoded event was of another class. +function accept( + ts: T, +): Decoder>> { + type Ty = s.Infer; + return d.object({ type: d.string }).then(({ type: ty }) => { + const c: undefined | EventClass = ts.find((t) => t.type === ty); + return c + ? (c.schema.decoder.map(Just) as Decoder>) + : d.succeed(Nothing()); + }); +} diff --git a/src/app/queryHandler.ts b/src/app/queryHandler.ts index e1ecb68..057db14 100644 --- a/src/app/queryHandler.ts +++ b/src/app/queryHandler.ts @@ -5,18 +5,16 @@ import { Event } from '@/lib/eventSourcing/event'; import { Decoder, decode } from '@/lib/json/decoder'; import * as express from 'express'; import * as router from '@/lib/router'; -import * as d from '@/lib/json/decoder'; import { Future } from '@/lib/Future'; import { Result, Failure } from '@/lib/Result'; -import { Maybe, Nothing } from '@/lib/Maybe'; type Projections = {}; type Services = {}; -type QueryController> = { - decoder: Decoder>; +type QueryController = { + decoder: Decoder; handler: (v: { - event: E; + query: Query; projections: Projections; services: Services; }) => Future; @@ -28,9 +26,9 @@ function handleQuery>( { decoder, handler }: QueryController, ): express.Handler { return router.route((req) => - decodeEvent(decoder, req).chain((event) => + decodeQuery(decoder, req).chain((query) => handler({ - event, + query, projections, services, }), @@ -38,33 +36,19 @@ function handleQuery>( ); } -function decodeEvent( - decoder: Decoder>, +function decodeQuery( + decoder: Decoder, req: express.Request, -): Future { - const bodyDecoder: Decoder> = d - .object({ payload: decoder }) - .map((r) => r.payload); - - const decoded: Result> = decode(bodyDecoder, req.body); - +): Future { + const decoded: Result = decode(decoder, req.body); if (decoded instanceof Failure) { return Future.reject( router.json({ status: 400, - content: { message: `Unable to decode command: ${decoded.error}` }, - }), - ); - } - - if (decoded.value instanceof Nothing) { - return Future.reject( - router.json({ - status: 200, - content: { message: 'Ignored' }, + content: { message: `Unable to decode request: ${decoded.error}` }, }), ); } - return Future.resolve(decoded.value.value); + return Future.resolve(decoded.value); } diff --git a/src/app/reactionHandler.ts b/src/app/reactionHandler.ts index 3d3ae48..b612c84 100644 --- a/src/app/reactionHandler.ts +++ b/src/app/reactionHandler.ts @@ -3,17 +3,16 @@ export { handleReaction }; import { Response } from '@/lib/router'; import { EventStore } from '@/lib/eventSourcing/eventStore'; import { Event } from '@/lib/eventSourcing/event'; -import { Decoder, decode } from '@/lib/json/decoder'; +import { Decoder } from '@/lib/json/decoder'; import * as express from 'express'; import * as router from '@/lib/router'; -import * as d from '@/lib/json/decoder'; import { Future } from '@/lib/Future'; -import { Result, Failure } from '@/lib/Result'; -import { Maybe, Nothing } from '@/lib/Maybe'; +import { Maybe } from '@/lib/Maybe'; import { Postgres } from '@/lib/postgres'; import { PostgresEventStore } from '@/app/postgresEventStore'; import { Serializer } from '@/common/serializedEvent/Serializer'; import { Deserializer } from '@/common/serializedEvent/Deserializer'; +import { decodeEvent } from '@/app/projectionHandler'; type Projections = {}; type Services = {}; @@ -56,37 +55,6 @@ function handleReaction>( ); } -function decodeEvent( - decoder: Decoder>, - req: express.Request, -): Future { - const bodyDecoder: Decoder> = d - .object({ payload: decoder }) - .map((r) => r.payload); - - const decoded: Result> = decode(bodyDecoder, req.body); - - if (decoded instanceof Failure) { - return Future.reject( - router.json({ - status: 400, - content: { message: `Unable to decode command: ${decoded.error}` }, - }), - ); - } - - if (decoded.value instanceof Nothing) { - return Future.reject( - router.json({ - status: 200, - content: { message: 'Ignored' }, - }), - ); - } - - return Future.resolve(decoded.value.value); -} - function withEventStore( postgres: Postgres, serializer: Serializer, From b93c7b50f1c6f11e46aeb0840e5095435e58874b Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Wed, 24 Sep 2025 23:58:40 +0100 Subject: [PATCH 15/22] Create better Mongo abstraction --- src/lib/mongo.ts | 156 ++++++++++++++++++++++++++++++++++++++++++++ src/lib/postgres.ts | 11 ++++ 2 files changed, 167 insertions(+) create mode 100644 src/lib/mongo.ts diff --git a/src/lib/mongo.ts b/src/lib/mongo.ts new file mode 100644 index 0000000..2922ca0 --- /dev/null +++ b/src/lib/mongo.ts @@ -0,0 +1,156 @@ +export { Mongo, MongoTransaction }; + +import { + ClientSession, + Filter, + FindOptions, + Document, + ReplaceOptions, + InsertOneOptions, + CountOptions, + Db, + ReadConcern, + WriteConcern, + ReadPreference, + TransactionOptions, + OptionalUnlessRequiredId, + WithId, + MongoClientOptions, + MongoClient, +} from 'mongodb'; +import { Future } from '@/lib/Future'; + +class MongoTransaction { + public closed: boolean = false; + constructor( + public readonly session: ClientSession, + public readonly database: Db, + ) {} + + async commit() { + if (this.closed) { + throw new Error('Committing a closed transaction'); + } + + try { + await this.session.commitTransaction(); + this.closed = true; + } catch (error) { + this.closed = true; + throw new Error(`Failed to commit transaction: ${error}`); + } + } + + async abort() { + if (this.closed) { + throw new Error('Aborting a closed transaction'); + } + + try { + await this.session.abortTransaction(); + } catch (error) { + console.error('Failed to abort MongoDB transaction', error as Error); + } + this.closed = true; + } + + async find( + collectionName: string, + filter: Filter, + options?: FindOptions, + ): Promise[]> { + this.checkOpen(); + return this.database + .collection(collectionName) + .find(filter, { ...options, session: this.session }) + .toArray(); + } + + async replaceOne( + collectionName: string, + filter: Filter, + replacement: T, + options?: ReplaceOptions, + ): Promise { + this.checkOpen(); + return this.database + .collection(collectionName) + .replaceOne(filter, replacement, { + ...options, + session: this.session, + }); + } + + async insertOne( + collectionName: string, + document: T & OptionalUnlessRequiredId, + options?: InsertOneOptions, + ): Promise { + this.checkOpen(); + await this.database + .collection(collectionName) + .insertOne(document, { ...options, session: this.session }); + } + + async countDocuments( + collectionName: string, + filter: Filter, + options?: CountOptions, + ): Promise { + this.checkOpen(); + return this.database + .collection(collectionName) + .countDocuments(filter, { ...options, session: this.session }); + } + + private checkOpen() { + if (this.closed) { + throw new Error('Session must be active to read or write to MongoDB!'); + } + } +} + +const transactionOptions: TransactionOptions = { + readConcern: new ReadConcern('snapshot'), + writeConcern: new WriteConcern('majority'), + readPreference: ReadPreference.primary, +}; + +class Mongo { + client: MongoClient; + + constructor( + public values: { + username: string; + password: string; + host: string; + port: number; + database: string; + settings: MongoClientOptions; + }, + ) { + const connectionString = + `mongodb://${values.username}:${values.password}@${values.host}` + + `:${values.port.toString()}/${values.database}` + + '?serverSelectionTimeoutMS=10000&connectTimeoutMS=10000&authSource=admin'; + this.client = new MongoClient(connectionString, values.settings); + } + + // Execute an action with a transaction that will be automatically committed at the end. + withTransaction( + f: (t: MongoTransaction) => Future, + ): Future { + const session = this.client.startSession(); + + session.startTransaction(transactionOptions); + const database = this.client.db(this.values.database); + const transaction = new MongoTransaction(session, database); + return f(transaction).finally( + Future.create((_, res) => { + if (!transaction.closed) transaction.abort(); + session.endSession(); + return res(); + }), + ); + } +} diff --git a/src/lib/postgres.ts b/src/lib/postgres.ts index 72fa437..4f1ec0b 100644 --- a/src/lib/postgres.ts +++ b/src/lib/postgres.ts @@ -17,6 +17,7 @@ class PostgresTransaction { if (this.closed) { throw new Error('Committing a closed transaction'); } + try { await this.connection.query('COMMIT'); this.closed = true; @@ -24,6 +25,8 @@ class PostgresTransaction { this.closed = true; throw new Error(`Failed to commit transaction: ${error}`); } + + await this.release(); } async abort() { @@ -38,6 +41,14 @@ class PostgresTransaction { } this.closed = true; + await this.release(); + } + + async release() { + if (!this.closed) { + throw new Error('Releasing an active transaction'); + } + try { this.connection.release(); } catch (error) { From 2e2c0418b54007d9d515eb56d9fa1b7c0c3d60a3 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Thu, 25 Sep 2025 09:45:06 +0100 Subject: [PATCH 16/22] Initialise Mongo in configureDependencies --- src/di/container.ts | 26 +++++++++++++++++++++++++- src/lib/mongo.ts | 4 ++-- 2 files changed, 27 insertions(+), 3 deletions(-) diff --git a/src/di/container.ts b/src/di/container.ts index 244655a..e8f90b2 100644 --- a/src/di/container.ts +++ b/src/di/container.ts @@ -21,6 +21,8 @@ import { MembershipApplicationRepository } from '@/domain/cookingClub/membership import { CuisineRepository } from '@/domain/cookingClub/membership/projection/membersByCuisine/CuisineRepository'; import env from '@/app/environment'; import { Postgres, defaultPoolSettings } from '@/lib/postgres'; +import { Mongo } from '@/lib/mongo'; +import { ServerApiVersion } from 'mongodb'; import * as postgresEventStore from '@/app/postgresEventStore'; function registerEnvironmentVariables() { @@ -102,6 +104,7 @@ function registerScopedServices() { type Dependencies = { postgres: Postgres; + mongo: Mongo; }; export async function configureDependencies(): Promise { @@ -118,6 +121,27 @@ export async function configureDependencies(): Promise { poolSettings: defaultPoolSettings, }); + const mongo = new Mongo({ + user: env.MONGODB_PROJECTION_DATABASE_USERNAME, + password: env.MONGODB_PROJECTION_DATABASE_PASSWORD, + host: env.MONGODB_PROJECTION_HOST, + port: env.MONGODB_PROJECTION_PORT, + database: env.MONGODB_PROJECTION_DATABASE_NAME, + settings: { + maxPoolSize: 20, + minPoolSize: 5, + maxIdleTimeMS: 10 * 60 * 1000, // 10 minutes + maxConnecting: 30, + waitQueueTimeoutMS: 2000, + replicaSet: 'rs0', + serverApi: { + version: ServerApiVersion.v1, + strict: true, + deprecationErrors: true, + }, + }, + }); + await postgres.withTransactionP((transaction) => postgresEventStore.initialize({ transaction, @@ -131,5 +155,5 @@ export async function configureDependencies(): Promise { }), ); - return { postgres }; + return { postgres, mongo }; } diff --git a/src/lib/mongo.ts b/src/lib/mongo.ts index 2922ca0..e2c3f97 100644 --- a/src/lib/mongo.ts +++ b/src/lib/mongo.ts @@ -121,7 +121,7 @@ class Mongo { constructor( public values: { - username: string; + user: string; password: string; host: string; port: number; @@ -130,7 +130,7 @@ class Mongo { }, ) { const connectionString = - `mongodb://${values.username}:${values.password}@${values.host}` + + `mongodb://${values.user}:${values.password}@${values.host}` + `:${values.port.toString()}/${values.database}` + '?serverSelectionTimeoutMS=10000&connectTimeoutMS=10000&authSource=admin'; this.client = new MongoClient(connectionString, values.settings); From 2abd8c36af17cead028316ed879256f89b4b5214 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Thu, 25 Sep 2025 13:50:17 +0100 Subject: [PATCH 17/22] Implement Hydrator --- src/lib/eventSourcing/event.ts | 5 +- src/lib/eventSourcing/eventStore.ts | 143 +++++++++++++++++++++++++++- 2 files changed, 142 insertions(+), 6 deletions(-) diff --git a/src/lib/eventSourcing/event.ts b/src/lib/eventSourcing/event.ts index 056e4cb..2fe6f47 100644 --- a/src/lib/eventSourcing/event.ts +++ b/src/lib/eventSourcing/event.ts @@ -10,6 +10,7 @@ export { import * as s from '@/lib/json/schema'; import { Schema } from '@/lib/json/schema'; +import { Encoder } from '@/lib/json/encoder'; import { POSIX } from '@/lib/time'; // @ts-ignore @@ -34,7 +35,9 @@ interface Aggregate { } // Class which all events derive from. Used for type constraints. -abstract class Event<_T extends Aggregate<_T>> {} +abstract class Event> { + abstract encoder: Encoder>; +} // The first event for an aggregate. abstract class CreationEvent> extends Event { diff --git a/src/lib/eventSourcing/eventStore.ts b/src/lib/eventSourcing/eventStore.ts index 53a5b0d..499e2c1 100644 --- a/src/lib/eventSourcing/eventStore.ts +++ b/src/lib/eventSourcing/eventStore.ts @@ -1,7 +1,19 @@ -export { type EventStore }; -import { Event } from '@/lib/eventSourcing/event'; -import { Aggregate } from '@/common/aggregate/Aggregate'; -import { AggregateAndEventIdsInLastEvent } from '@/common/eventStore/AggregateAndEventIdsInLastEvent'; +export { type EventStore, type AggregateAndEventIdsInLastEvent, Hydrator }; + +import { + CreationEvent, + TransformationEvent, + EventInfo, + Event, + Aggregate, + Id, +} from '@/lib/eventSourcing/event'; +import { Json } from '@/lib/json/types'; +import { Schema } from '@/lib/json/schema'; +import * as s from '@/lib/json/schema'; +import * as d from '@/lib/json/decoder'; +import { POSIX } from '@/lib/time'; +import { Result, Failure } from '@/lib/Result'; /* Note [Event Store] @@ -18,8 +30,14 @@ import { AggregateAndEventIdsInLastEvent } from '@/common/eventStore/AggregateAn */ +interface AggregateAndEventIdsInLastEvent> { + aggregate: T; + eventIdOfLastEvent: Id>; + correlationIdOfLastEvent: Id>; +} + interface EventStore { - findAggregate( + findAggregate>( aggregateId: string, ): Promise>; @@ -27,3 +45,118 @@ interface EventStore { doesEventAlreadyExist(eventId: string): Promise; } + +// ----------------------------------------------------------------------- + +// Serialized representation of an event. +type Serialized

= s.Infer>>; + +const schema_Serialized = (payload: Schema) => + s.object({ + event_id: Id.schema>>(), + aggregate_id: Id.schema>(), + aggregate_version: s.number, + correlation_id: Id.schema>>(), + causation_id: Id.schema>>(), + recorded_on: POSIX.schema, + payload, + }); + +// ----------------------------------------------------------------------- + +// Convenient runtime representation of data in a serialized event. +type EventData = { info: EventInfo; payload: E }; + +function toSerialized

({ + info, + payload, +}: { + info: EventInfo; + payload: P; +}): Serialized

{ + return { ...info, payload }; +} + +function fromSerialized

(serialized: Serialized

): { + info: EventInfo; + payload: P; +} { + const { payload, ...info } = serialized; + return { info, payload }; +} + +const schema_EventData = (s: Schema): Schema> => + schema_Serialized(s).dimap(fromSerialized, toSerialized); + +// ----------------------------------------------------------------------- + +type Constructor = new (...args: any[]) => T; + +type Schemas> = { + creation: Schema>>; + transformation: Schema>>; +}; + +/* Note [Hydrator] + + We need some type-safe way to decode events for an aggregate. That is, without casting. + We perform type-directed decoding, where we specify the type of the aggregate, + then use the decoders we have for creation and transformation events for that aggregate. + + This ensures that we will never apply an incorrect aggregate transformation or create + an aggregate of the incorrect type. +*/ +class Hydrator { + private tmap = new Map, Schemas>(); + + constructor() {} + + // add support for deserializing an aggregate's events. + add>({ + aggregate, + creation, + transformation, + }: { + aggregate: Constructor; + creation: Schema>; + transformation: Schema>; + }): void { + this.tmap.set(aggregate, { + creation: schema_EventData(creation), + transformation: schema_EventData(transformation), + }); + } + + // Build an aggregate from all its serialized events. + hydrate>( + aggregate: Constructor, + serialized: Json[], + ): Result { + const schemas = this.tmap.get(aggregate) as undefined | Schemas; + if (schemas == undefined) { + throw new Error(`Unknown aggregate ${aggregate.name}`); + } + + if (serialized.length === 0) { + return Failure('No events'); + } + + return d + .decode(serialized[0], schemas.creation.decoder) + .then(({ payload: first, info }) => + d + .decode(serialized.slice(1), d.array(schemas.transformation.decoder)) + .map((es) => { + let aggregate = first.createAggregate(); + let lastEvent = info; + + for (const t of es) { + aggregate = t.payload.transformAggregate(aggregate); + lastEvent = t.info; + } + + return { aggregate, lastEvent }; + }), + ); + } +} From 7f54e5c10371b31d66b9079baa4c6bc391c608af Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Thu, 25 Sep 2025 18:17:34 +0100 Subject: [PATCH 18/22] Implement find and insert in postgresEventStore --- src/app/event.ts | 189 +++++++++++++++++++++++++++- src/app/postgresEventStore.ts | 180 +++++++++++++++----------- src/lib/eventSourcing/event.ts | 24 ++-- src/lib/eventSourcing/eventStore.ts | 27 ++-- 4 files changed, 326 insertions(+), 94 deletions(-) diff --git a/src/app/event.ts b/src/app/event.ts index 85abeec..a00c90f 100644 --- a/src/app/event.ts +++ b/src/app/event.ts @@ -1,18 +1,123 @@ -export { accept, type EventClass }; +export { accept }; import { Aggregate, Id, CreationEvent, TransformationEvent, + Event, + EventInfo, } from '@/lib/eventSourcing/event'; import { Decoder } from '@/lib/json/decoder'; +import { Encoder } from '@/lib/json/encoder'; import * as s from '@/lib/json/schema'; import * as d from '@/lib/json/decoder'; import { Schema } from '@/lib/json/schema'; +import { Json } from '@/lib/json/types'; import { Maybe, Nothing, Just } from '@/lib/Maybe'; +import { Result, Failure } from '@/lib/Result'; import * as m from '@/lib/Maybe'; +import { POSIX } from '@/lib/time'; + +type Serialized

= s.Infer>>; + +const schema_Serialized = (payload: Schema) => + s.object({ + event_id: Id.schema>>(), + aggregate_id: Id.schema>(), + aggregate_version: s.number, + correlation_id: Id.schema>>(), + causation_id: Id.schema>>(), + recorded_on: POSIX.schema, + payload, + }); + +function toSerialized

({ + info, + payload, +}: { + info: EventInfo; + payload: P; +}): Serialized

{ + return { ...info, payload }; +} + +function fromSerialized

(serialized: Serialized

): { + info: EventInfo; + payload: P; +} { + const { payload, ...info } = serialized; + return { info, payload }; +} + +type EventData = { info: EventInfo; payload: E }; + +const schema_EventData = (s: Schema): Schema> => + schema_Serialized(s).dimap(fromSerialized, toSerialized); + +type Constructor = new (...args: any[]) => T; + +type Schemas> = { + creation: Schema>>; + transformation: Schema>>; +}; + +type EventName = string; + +class Deserializer { + private tmap = new Map, Schemas>(); + + constructor() {} + + // add support for deserializing an aggregate's events. + add>({ + aggregate, + creation, + transformation, + }: { + aggregate: Constructor; + creation: Schema>; + transformation: Schema>; + }): void { + this.tmap.set(aggregate, { + creation: schema_EventData(creation), + transformation: schema_EventData(transformation), + }); + } + + deserialize>( + aggregate: Constructor, + serialized: Json[], + ): Result { + const schemas = this.tmap.get(aggregate) as undefined | Schemas; + if (schemas == undefined) { + throw new Error(`Unknown aggregate ${aggregate.name}`); + } + + if (serialized.length === 0) { + return Failure('No events'); + } + + return d + .decode(serialized[0], schemas.creation.decoder) + .then(({ payload: first, info }) => + d + .decode(serialized.slice(1), d.array(schemas.transformation.decoder)) + .map((es) => { + let aggregate = first.createAggregate(); + let lastEvent = info; + + for (const t of es) { + aggregate = t.payload.transformAggregate(aggregate); + lastEvent = t.info; + } + + return { aggregate, lastEvent }; + }), + ); + } +} class User implements Aggregate { constructor( @@ -22,37 +127,76 @@ class User implements Aggregate { ) {} } -class CreateUser implements CreationEvent { +class Other implements Aggregate { + constructor( + readonly aggregateId: Id, + readonly aggregateVersion: number, + ) {} +} + +class CreateUser implements CreationEvent { static type: 'CreateUserr' = 'CreateUserr'; - constructor(readonly values: s.Infer) {} static schemaArgs = s.object({ type: s.stringLiteral(CreateUser.type), + aggregateId: Id.schema(), name: s.string, }); - static schema = CreateUser.schemaArgs.dimap( (v) => new CreateUser(v), (v) => v.values, ); + schema = CreateUser.schema; + constructor(readonly values: s.Infer) {} createAggregate() { return new User(new Id('wat'), 0, this.values.name); } } -class AddName implements TransformationEvent { +class AddName implements TransformationEvent { static type: 'AddName' = 'AddName'; constructor(readonly values: s.Infer) {} static schemaArgs = s.object({ type: s.stringLiteral(AddName.type), + aggregateId: Id.schema(), + name: s.string, + }); + + encoder = AddName.schema.encoder as Encoder>; + + static schema = AddName.schemaArgs.dimap( + (v) => new AddName(v), + (v) => v.values, + ); + readonly schema = AddName.schema; + + transformAggregate(agg: User): User { + const u = new User( + agg.aggregateId, + agg.aggregateVersion + 1, + this.values.name, + ); + return u; + } +} + +class RemoveName implements TransformationEvent { + static type: 'RemoveName' = 'RemoveName'; + constructor(readonly values: s.Infer) {} + + static schemaArgs = s.object({ + type: s.stringLiteral(AddName.type), + aggregateId: Id.schema(), name: s.string, }); + encoder = RemoveName.schema.encoder as Encoder>; static schema = AddName.schemaArgs.dimap( (v) => new AddName(v), (v) => v.values, ); + readonly schema = AddName.schema; transformAggregate(agg: User): User { const u = new User( @@ -86,3 +230,38 @@ function accept( : d.succeed(Nothing()); }); } + +// Create a schema given a list of event classes +function makeSchema( + ts: T, +): Schema> { + type Ty = s.Infer; + + const decoder: Decoder = d + .object({ type: d.string }) + .then(({ type: ty }) => { + const c: undefined | EventClass = ts.find((t) => t.type === ty); + + if (c === undefined) return d.fail(`Unknown event type: ${ty}`); + + return c.schema.decoder as Decoder; + }); + + const encoder: Encoder = new Encoder((v: Ty) => { + const ty = ts.find((t) => t.type === v.type); + if (ty === undefined) { + throw new Error(`Unable to encode unknown event type: ${v.type}`); + } + + return ty.schema.encoder.run(v); + }); + + return new Schema(decoder, encoder); +} + +const ww = new Deserializer(); +ww.add({ + aggregate: User, + creation: makeSchema([CreateUser]).decoder, + transformation: makeSchema([AddName, RemoveName]).decoder, +}); diff --git a/src/app/postgresEventStore.ts b/src/app/postgresEventStore.ts index 6b3f81c..d8b8fc4 100644 --- a/src/app/postgresEventStore.ts +++ b/src/app/postgresEventStore.ts @@ -1,66 +1,91 @@ export { initialize, PostgresEventStore }; -import { Serializer } from '@/common/serializedEvent/Serializer'; -import { Deserializer } from '@/common/serializedEvent/Deserializer'; +import { Json } from '@/lib/json/types'; import { SerializedEvent } from '@/common/serializedEvent/SerializedEvent'; -import { Event } from '@/common/event/Event'; -import { CreationEvent } from '@/common/event/CreationEvent'; -import { TransformationEvent } from '@/common/event/TransformationEvent'; -import { Aggregate } from '@/common/aggregate/Aggregate'; -import { AggregateAndEventIdsInLastEvent } from '@/common/eventStore/AggregateAndEventIdsInLastEvent'; -import { EventStore } from '@/lib/eventSourcing/eventStore'; +import { + Id, + Event, + EventClass, + Aggregate, + CreationEvent, + TransformationEvent, + EventInfo, +} from '@/lib/eventSourcing/event'; +import { + EventStore, + Hydrator, + Constructor, + EventData, +} from '@/lib/eventSourcing/eventStore'; import { PostgresTransaction } from '@/lib/postgres'; import { log } from '@/common/util/Logger'; +import { IdGenerator } from '@/common/util/IdGenerator'; +import { encode } from '@/lib/json/schema'; +import { POSIX } from '@/lib/time'; +import { DateTime } from 'luxon'; class PostgresEventStore implements EventStore { constructor( private transaction: PostgresTransaction, - private readonly serializer: Serializer, - private readonly deserializer: Deserializer, + private readonly hydrator: Hydrator, private readonly eventStoreTable: string, ) {} - async findAggregate( - aggregateId: string, - ): Promise> { - const serializedEvents = - await this.findAllSerializedEventsByAggregateId(aggregateId); - const events = serializedEvents.map((e) => - this.deserializer.deserialize(e), - ); - - const firstEvent: Event | undefined = events[0]; - if (firstEvent == undefined) { - throw new Error(`No events found for aggregateId: ${aggregateId}`); - } - const creationEvent: Event = firstEvent; - if (!this.isCreationEventForAggregate(creationEvent)) { - throw new Error('First event is not a creation event'); - } - const transformationEvents = events.slice(1); + async find>( + cls: Constructor, + aggregateId: Id, + ): Promise<{ aggregate: T; lastEvent: EventInfo }> { + const events = await this.findAll(aggregateId); + const { lastEvent, aggregate } = this.hydrator + .hydrate(cls, events) + .unwrap((e) => e); - let aggregate = creationEvent.createAggregate(); - let eventIdOfLastEvent = creationEvent.eventId; - let correlationIdOfLastEvent = creationEvent.correlationId; + return { aggregate, lastEvent }; + } - for (const transformationEvent of transformationEvents) { - if (!this.isTransformationEventForAggregate(transformationEvent)) { - throw new Error('Event is not a transformation event'); + async save, T extends Aggregate>(args: { + aggregate: Constructor; + event: CreationEvent | TransformationEvent; + event_id?: Id>; + correlation_id?: Id>; + causation_id?: Id>; + }): Promise { + const event = args.event; + const event_id = args.event_id || new Id(IdGenerator.generateRandomId()); + let info: EventInfo; + switch (true) { + case event instanceof CreationEvent: { + const aggregate: T = event.createAggregate(); + info = { + event_id, + aggregate_id: aggregate.aggregateId, + aggregate_version: 0, + correlation_id: args.correlation_id || event_id, + causation_id: args.causation_id || event_id, + recorded_on: POSIX.now(), + }; + break; } - aggregate = transformationEvent.transformAggregate(aggregate); - eventIdOfLastEvent = transformationEvent.eventId; - correlationIdOfLastEvent = transformationEvent.correlationId; + case event instanceof TransformationEvent: { + const { aggregate, lastEvent } = await this.find( + args.aggregate, + event.values.aggregateId, + ); + info = { + event_id, + aggregate_id: aggregate.aggregateId, + aggregate_version: aggregate.aggregateVersion + 1, + correlation_id: lastEvent.correlation_id, + causation_id: lastEvent.causation_id, + recorded_on: POSIX.now(), + }; + break; + } + default: + return event satisfies never; } - return { - aggregate, - eventIdOfLastEvent, - correlationIdOfLastEvent, - }; - } - - async saveEvent(event: Event): Promise { - await this.saveSerializedEvent(this.serializer.serialize(event)); + await this.insert({ info, event }); } async doesEventAlreadyExist(eventId: string): Promise { @@ -68,9 +93,9 @@ class PostgresEventStore implements EventStore { return event !== null; } - private async findAllSerializedEventsByAggregateId( - aggregateId: string, - ): Promise { + private async findAll>( + aggregateId: Id, + ): Promise { const sql = ` SELECT id, event_id, aggregate_id, causation_id, correlation_id, aggregate_version, json_payload, json_metadata, recorded_on, event_name @@ -80,8 +105,8 @@ class PostgresEventStore implements EventStore { `; try { - const result = await this.transaction.query(sql, [aggregateId]); - return result.rows.map(this.mapRowToSerializedEvent); + const result = await this.transaction.query(sql, [aggregateId.value]); + return result.rows; } catch (error) { throw new Error( `Failed to fetch events for aggregate: ${aggregateId}: ${error}`, @@ -89,34 +114,31 @@ class PostgresEventStore implements EventStore { } } - private async saveSerializedEvent( - serializedEvent: SerializedEvent, - ): Promise { + private async insert>({ info, event }: EventData) { const sql = ` - INSERT INTO ${this.eventStoreTable} ( - event_id, aggregate_id, causation_id, correlation_id, - aggregate_version, json_payload, json_metadata, recorded_on, event_name - ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) - `; + INSERT INTO ${this.eventStoreTable} ( + event_id, aggregate_id, causation_id, correlation_id, + aggregate_version, json_payload, json_metadata, recorded_on, event_name + ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) + `; + const payload = encode(event.schema, event); const values = [ - serializedEvent.event_id, - serializedEvent.aggregate_id, - serializedEvent.causation_id, - serializedEvent.correlation_id, - serializedEvent.aggregate_version.toString(), - serializedEvent.json_payload, - serializedEvent.json_metadata, - serializedEvent.recorded_on, - serializedEvent.event_name, + info.event_id.value, + info.aggregate_id.value, + info.causation_id.value, + info.correlation_id.value, + info.aggregate_version.toString(), + JSON.stringify(payload), + '{}', + encodePOSIX(info.recorded_on), + event.type, ]; try { await this.transaction.query(sql, values); } catch (error) { - throw new Error( - `Failed to save event: ${serializedEvent.event_id}: ${error}`, - ); + throw new Error(`Failed to save event: ${info.event_id}: ${error}`); } } @@ -168,6 +190,20 @@ class PostgresEventStore implements EventStore { } } +function encodePOSIX(value: POSIX): string { + const { date, time } = value.toUTCDateAndTime(); + return `${date.pretty()}T${time.pretty()}Z`; +} + +function decodePOSIX(str: string): POSIX { + const date = DateTime.fromISO(str, { zone: 'UTC' }); + if (!date.isValid) { + throw new Error(`Invalid ISO date: ${str}`); + } + + return new POSIX(date.toMillis()); +} + // Prepare the database to be used as an event store. async function initialize({ transaction, @@ -199,7 +235,7 @@ async function initialize({ aggregate_version BIGINT NOT NULL, causation_id TEXT NOT NULL, correlation_id TEXT NOT NULL, - recorded_on TEXT NOT NULL, + recorded_on TIMESTAMPTZ NOT NULL, event_name TEXT NOT NULL, json_payload TEXT NOT NULL, json_metadata TEXT NOT NULL, diff --git a/src/lib/eventSourcing/event.ts b/src/lib/eventSourcing/event.ts index 2fe6f47..b736ffd 100644 --- a/src/lib/eventSourcing/event.ts +++ b/src/lib/eventSourcing/event.ts @@ -1,6 +1,7 @@ export { - Event, + type Event, type Aggregate, + EventClass, TransformationEvent, CreationEvent, type EventInfo, @@ -10,7 +11,6 @@ export { import * as s from '@/lib/json/schema'; import { Schema } from '@/lib/json/schema'; -import { Encoder } from '@/lib/json/encoder'; import { POSIX } from '@/lib/time'; // @ts-ignore @@ -31,21 +31,31 @@ class Id { // Class which all events derive from. Used for type constraints. interface Aggregate { readonly aggregateId: Id>; - readonly aggregateVersion: number; + aggregateVersion: number; } +type Event> = EventClass; + // Class which all events derive from. Used for type constraints. -abstract class Event> { - abstract encoder: Encoder>; +abstract class EventClass> { + abstract type: string; + abstract values: { aggregateId: Id }; + abstract schema: Schema; } // The first event for an aggregate. -abstract class CreationEvent> extends Event { +abstract class CreationEvent> extends EventClass< + Self, + T +> { abstract createAggregate(): T; } // Any event that is not the first one for an aggregate. -abstract class TransformationEvent> extends Event { +abstract class TransformationEvent< + Self, + T extends Aggregate, +> extends EventClass { abstract transformAggregate(aggregate: T): T; } diff --git a/src/lib/eventSourcing/eventStore.ts b/src/lib/eventSourcing/eventStore.ts index 499e2c1..d880efa 100644 --- a/src/lib/eventSourcing/eventStore.ts +++ b/src/lib/eventSourcing/eventStore.ts @@ -1,4 +1,11 @@ -export { type EventStore, type AggregateAndEventIdsInLastEvent, Hydrator }; +export { + type EventStore, + type AggregateAndEventIdsInLastEvent, + Hydrator, + type Constructor, + type EventData, + schema_EventData, +}; import { CreationEvent, @@ -65,24 +72,24 @@ const schema_Serialized = (payload: Schema) => // ----------------------------------------------------------------------- // Convenient runtime representation of data in a serialized event. -type EventData = { info: EventInfo; payload: E }; +type EventData = { info: EventInfo; event: E }; function toSerialized

({ info, - payload, + event, }: { info: EventInfo; - payload: P; + event: P; }): Serialized

{ - return { ...info, payload }; + return { ...info, payload: event }; } function fromSerialized

(serialized: Serialized

): { info: EventInfo; - payload: P; + event: P; } { const { payload, ...info } = serialized; - return { info, payload }; + return { info, event: payload }; } const schema_EventData = (s: Schema): Schema> => @@ -129,12 +136,12 @@ class Hydrator { // Build an aggregate from all its serialized events. hydrate>( - aggregate: Constructor, + cls: Constructor, serialized: Json[], ): Result { - const schemas = this.tmap.get(aggregate) as undefined | Schemas; + const schemas = this.tmap.get(cls) as undefined | Schemas; if (schemas == undefined) { - throw new Error(`Unknown aggregate ${aggregate.name}`); + throw new Error(`Unknown aggregate ${cls.name}`); } if (serialized.length === 0) { From 026186e1a2f5fbee7ad2ba00b86d4515264c5d4a Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Thu, 25 Sep 2025 18:49:16 +0100 Subject: [PATCH 19/22] Type-safe event serialization --- src/app/postgresEventStore.ts | 129 +++++++++------------------- src/lib/eventSourcing/eventStore.ts | 47 ++++++++-- 2 files changed, 79 insertions(+), 97 deletions(-) diff --git a/src/app/postgresEventStore.ts b/src/app/postgresEventStore.ts index d8b8fc4..0cb8b75 100644 --- a/src/app/postgresEventStore.ts +++ b/src/app/postgresEventStore.ts @@ -1,14 +1,13 @@ export { initialize, PostgresEventStore }; import { Json } from '@/lib/json/types'; -import { SerializedEvent } from '@/common/serializedEvent/SerializedEvent'; import { Id, Event, - EventClass, Aggregate, CreationEvent, TransformationEvent, + EventClass, EventInfo, } from '@/lib/eventSourcing/event'; import { @@ -16,13 +15,13 @@ import { Hydrator, Constructor, EventData, + schema_EventData, } from '@/lib/eventSourcing/eventStore'; import { PostgresTransaction } from '@/lib/postgres'; import { log } from '@/common/util/Logger'; import { IdGenerator } from '@/common/util/IdGenerator'; import { encode } from '@/lib/json/schema'; import { POSIX } from '@/lib/time'; -import { DateTime } from 'luxon'; class PostgresEventStore implements EventStore { constructor( @@ -43,7 +42,7 @@ class PostgresEventStore implements EventStore { return { aggregate, lastEvent }; } - async save, T extends Aggregate>(args: { + async save, T extends Aggregate>(args: { aggregate: Constructor; event: CreationEvent | TransformationEvent; event_id?: Id>; @@ -85,24 +84,32 @@ class PostgresEventStore implements EventStore { return event satisfies never; } - await this.insert({ info, event }); + await this.insert>({ info, event }); } - async doesEventAlreadyExist(eventId: string): Promise { - const event = await this.findSerializedEventByEventId(eventId); - return event !== null; + async doesEventAlreadyExist(eventId: Id>): Promise { + const sql = ` + SELECT 1 + FROM ${this.eventStoreTable} + WHERE event_id = $1`; + + try { + const result = await this.transaction.query(sql, [eventId.value]); + return result.rows.length > 0; + } catch (error) { + throw new Error(`Failed to fetch event: ${eventId}: ${error}`); + } } private async findAll>( aggregateId: Id, ): Promise { const sql = ` - SELECT id, event_id, aggregate_id, causation_id, correlation_id, - aggregate_version, json_payload, json_metadata, recorded_on, event_name - FROM ${this.eventStoreTable} - WHERE aggregate_id = $1 - ORDER BY aggregate_version ASC - `; + SELECT id, event_id, aggregate_id, causation_id, correlation_id, + aggregate_version, json_payload, json_metadata, recorded_on, event_name + FROM ${this.eventStoreTable} + WHERE aggregate_id = $1 + ORDER BY aggregate_version ASC`; try { const result = await this.transaction.query(sql, [aggregateId.value]); @@ -114,94 +121,40 @@ class PostgresEventStore implements EventStore { } } - private async insert>({ info, event }: EventData) { + private async insert>(edata: EventData) { const sql = ` INSERT INTO ${this.eventStoreTable} ( event_id, aggregate_id, causation_id, correlation_id, aggregate_version, json_payload, json_metadata, recorded_on, event_name - ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) - `; + ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9)`; + + const serialized = encode(schema_EventData(edata.event.schema), edata); - const payload = encode(event.schema, event); const values = [ - info.event_id.value, - info.aggregate_id.value, - info.causation_id.value, - info.correlation_id.value, - info.aggregate_version.toString(), - JSON.stringify(payload), + // @ts-ignore + serialized.event_id, + // @ts-ignore + serialized.aggregate_id, + // @ts-ignore + serialized.causation_id, + // @ts-ignore + serialized.correlation_id, + // @ts-ignore + serialized.aggregate_version, + // @ts-ignore + serialized.payload, '{}', - encodePOSIX(info.recorded_on), - event.type, + // @ts-ignore + serialized.recorded_on, + edata.event.type, ]; try { await this.transaction.query(sql, values); } catch (error) { - throw new Error(`Failed to save event: ${info.event_id}: ${error}`); + throw new Error(`Failed to save event: ${edata.info.event_id}: ${error}`); } } - - private async findSerializedEventByEventId( - eventId: string, - ): Promise { - const sql = ` - SELECT id, event_id, aggregate_id, causation_id, correlation_id, - aggregate_version, json_payload, json_metadata, recorded_on, event_name - FROM ${this.eventStoreTable} - WHERE event_id = $1 - `; - - try { - const result = await this.transaction.query(sql, [eventId]); - return result.rows.length > 0 - ? this.mapRowToSerializedEvent(result.rows[0]) - : null; - } catch (error) { - throw new Error(`Failed to fetch event: ${eventId}: ${error}`); - } - } - - private mapRowToSerializedEvent(row: any): SerializedEvent { - return { - id: row.id, - event_id: row.event_id, - aggregate_id: row.aggregate_id, - causation_id: row.causation_id, - correlation_id: row.correlation_id, - aggregate_version: row.aggregate_version, - json_payload: row.json_payload, - json_metadata: row.json_metadata, - recorded_on: row.recorded_on, - event_name: row.event_name, - }; - } - - private isCreationEventForAggregate( - event: Event, - ): event is CreationEvent { - return event instanceof CreationEvent; - } - - private isTransformationEventForAggregate( - event: Event, - ): event is TransformationEvent { - return event instanceof TransformationEvent; - } -} - -function encodePOSIX(value: POSIX): string { - const { date, time } = value.toUTCDateAndTime(); - return `${date.pretty()}T${time.pretty()}Z`; -} - -function decodePOSIX(str: string): POSIX { - const date = DateTime.fromISO(str, { zone: 'UTC' }); - if (!date.isValid) { - throw new Error(`Invalid ISO date: ${str}`); - } - - return new POSIX(date.toMillis()); } // Prepare the database to be used as an event store. diff --git a/src/lib/eventSourcing/eventStore.ts b/src/lib/eventSourcing/eventStore.ts index d880efa..90b5087 100644 --- a/src/lib/eventSourcing/eventStore.ts +++ b/src/lib/eventSourcing/eventStore.ts @@ -17,10 +17,12 @@ import { } from '@/lib/eventSourcing/event'; import { Json } from '@/lib/json/types'; import { Schema } from '@/lib/json/schema'; +import { Decoder } from '@/lib/json/schema'; import * as s from '@/lib/json/schema'; import * as d from '@/lib/json/decoder'; import { POSIX } from '@/lib/time'; import { Result, Failure } from '@/lib/Result'; +import { DateTime } from 'luxon'; /* Note [Event Store] @@ -44,13 +46,20 @@ interface AggregateAndEventIdsInLastEvent> { } interface EventStore { - findAggregate>( - aggregateId: string, - ): Promise>; - - saveEvent(event: Event): Promise; - - doesEventAlreadyExist(eventId: string): Promise; + find>( + cls: Constructor, + aggregateId: Id, + ): Promise<{ aggregate: T; lastEvent: EventInfo }>; + + save, T extends Aggregate>(args: { + aggregate: Constructor; + event: CreationEvent | TransformationEvent; + event_id?: Id>; + correlation_id?: Id>; + causation_id?: Id>; + }): Promise; + + doesEventAlreadyExist(eventId: Id>): Promise; } // ----------------------------------------------------------------------- @@ -65,10 +74,30 @@ const schema_Serialized = (payload: Schema) => aggregate_version: s.number, correlation_id: Id.schema>>(), causation_id: Id.schema>>(), - recorded_on: POSIX.schema, - payload, + recorded_on: schema_UTC, + payload: schema_StringifiedJSON(payload), }); +const schema_UTC: s.Schema = s.string.then( + (s) => { + const date = DateTime.fromISO(s, { zone: 'UTC' }); + return date.isValid + ? d.succeed(new POSIX(date.toMillis())) + : d.fail(`Invalid ISO date: ${s}`); + }, + (s) => { + const { date, time } = s.toUTCDateAndTime(); + return `${date.pretty()}T${time.pretty()}Z`; + }, +); + +const schema_StringifiedJSON = (inner: s.Schema): s.Schema => + s.string.then( + (str: string): Decoder => + new Decoder((_: unknown) => inner.decoder.run(JSON.parse(str))), + (t: T): string => JSON.stringify(inner.encoder.run(t)), + ); + // ----------------------------------------------------------------------- // Convenient runtime representation of data in a serialized event. From 51ccd2b9dd8edf08ab8ad6d7755fe7f2b16e5c07 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Thu, 25 Sep 2025 18:55:36 +0100 Subject: [PATCH 20/22] Fix types in eventStore --- src/lib/eventSourcing/eventStore.ts | 14 +++++++------- src/lib/json/decoder.ts | 2 +- 2 files changed, 8 insertions(+), 8 deletions(-) diff --git a/src/lib/eventSourcing/eventStore.ts b/src/lib/eventSourcing/eventStore.ts index 90b5087..89dc4cd 100644 --- a/src/lib/eventSourcing/eventStore.ts +++ b/src/lib/eventSourcing/eventStore.ts @@ -17,7 +17,7 @@ import { } from '@/lib/eventSourcing/event'; import { Json } from '@/lib/json/types'; import { Schema } from '@/lib/json/schema'; -import { Decoder } from '@/lib/json/schema'; +import { Decoder } from '@/lib/json/decoder'; import * as s from '@/lib/json/schema'; import * as d from '@/lib/json/decoder'; import { POSIX } from '@/lib/time'; @@ -129,8 +129,8 @@ const schema_EventData = (s: Schema): Schema> => type Constructor = new (...args: any[]) => T; type Schemas> = { - creation: Schema>>; - transformation: Schema>>; + creation: Schema>>; + transformation: Schema>>; }; /* Note [Hydrator] @@ -154,8 +154,8 @@ class Hydrator { transformation, }: { aggregate: Constructor; - creation: Schema>; - transformation: Schema>; + creation: Schema>; + transformation: Schema>; }): void { this.tmap.set(aggregate, { creation: schema_EventData(creation), @@ -179,7 +179,7 @@ class Hydrator { return d .decode(serialized[0], schemas.creation.decoder) - .then(({ payload: first, info }) => + .then(({ event: first, info }) => d .decode(serialized.slice(1), d.array(schemas.transformation.decoder)) .map((es) => { @@ -187,7 +187,7 @@ class Hydrator { let lastEvent = info; for (const t of es) { - aggregate = t.payload.transformAggregate(aggregate); + aggregate = t.event.transformAggregate(aggregate); lastEvent = t.info; } diff --git a/src/lib/json/decoder.ts b/src/lib/json/decoder.ts index 080939c..fd375c3 100644 --- a/src/lib/json/decoder.ts +++ b/src/lib/json/decoder.ts @@ -26,7 +26,7 @@ export { type FromJSON, type Infer, - type Decoder, // export only the abstract type, not constructors. + Decoder, type DecoderDef, type DecodeResult, decode, From baa0ec3ed45e08612c560ac879337ddf9264ea19 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Fri, 26 Sep 2025 14:54:44 +0100 Subject: [PATCH 21/22] Abstract withEventStore --- src/app/commandHandler.ts | 46 +++++++------------------------------- src/app/reactionHandler.ts | 46 +++++++------------------------------- src/di/container.ts | 3 +++ 3 files changed, 19 insertions(+), 76 deletions(-) diff --git a/src/app/commandHandler.ts b/src/app/commandHandler.ts index 7772656..3d00628 100644 --- a/src/app/commandHandler.ts +++ b/src/app/commandHandler.ts @@ -7,10 +7,6 @@ import * as express from 'express'; import * as router from '@/lib/router'; import { Future } from '@/lib/Future'; import { Result, Failure } from '@/lib/Result'; -import { Postgres } from '@/lib/postgres'; -import { PostgresEventStore } from '@/app/postgresEventStore'; -import { Serializer } from '@/common/serializedEvent/Serializer'; -import { Deserializer } from '@/common/serializedEvent/Deserializer'; type Projections = {}; type Services = {}; @@ -26,28 +22,20 @@ type CommandController = { }; function handleCommand( - serializer: Serializer, - deserializer: Deserializer, - eventStoreTable: string, - postgres: Postgres, + withEventStore: (f: (store: EventStore) => T) => T, services: Services, projections: Projections, { decoder, handler }: CommandController, ): express.Handler { return router.route((req) => decodeCommand(decoder, req).chain((command) => - withEventStore( - postgres, - serializer, - deserializer, - eventStoreTable, - (store) => - handler({ - command, - store, - projections, - services, - }), + withEventStore((store) => + handler({ + command, + store, + projections, + services, + }), ), ), ); @@ -69,21 +57,3 @@ function decodeCommand( return Future.resolve(decoded.value); } - -function withEventStore( - postgres: Postgres, - serializer: Serializer, - deserializer: Deserializer, - eventStoreTable: string, - f: (s: EventStore) => Future, -): Future { - const onError = (_: Error) => - router.json({ - status: 500, - content: { message: 'Internal Server Error' }, - }); - - return postgres.withTransaction(onError, (t) => - f(new PostgresEventStore(t, serializer, deserializer, eventStoreTable)), - ); -} diff --git a/src/app/reactionHandler.ts b/src/app/reactionHandler.ts index b612c84..1791196 100644 --- a/src/app/reactionHandler.ts +++ b/src/app/reactionHandler.ts @@ -8,10 +8,6 @@ import * as express from 'express'; import * as router from '@/lib/router'; import { Future } from '@/lib/Future'; import { Maybe } from '@/lib/Maybe'; -import { Postgres } from '@/lib/postgres'; -import { PostgresEventStore } from '@/app/postgresEventStore'; -import { Serializer } from '@/common/serializedEvent/Serializer'; -import { Deserializer } from '@/common/serializedEvent/Deserializer'; import { decodeEvent } from '@/app/projectionHandler'; type Projections = {}; @@ -28,47 +24,21 @@ type ReactionController> = { }; function handleReaction>( + withEventStore: (f: (s: EventStore) => T) => T, projections: Projections, services: Services, - postgres: Postgres, - serializer: Serializer, - deserializer: Deserializer, - eventStoreTable: string, { decoder, handler }: ReactionController, ): express.Handler { return router.route((req) => decodeEvent(decoder, req).chain((event) => - withEventStore( - postgres, - serializer, - deserializer, - eventStoreTable, - (store) => - handler({ - event, - projections, - services, - store, - }), + withEventStore((store) => + handler({ + event, + projections, + services, + store, + }), ), ), ); } - -function withEventStore( - postgres: Postgres, - serializer: Serializer, - deserializer: Deserializer, - eventStoreTable: string, - f: (s: EventStore) => Future, -): Future { - const onError = (_: Error) => - router.json({ - status: 500, - content: { message: 'Internal Server Error' }, - }); - - return postgres.withTransaction(onError, (t) => - f(new PostgresEventStore(t, serializer, deserializer, eventStoreTable)), - ); -} diff --git a/src/di/container.ts b/src/di/container.ts index e8f90b2..9449339 100644 --- a/src/di/container.ts +++ b/src/di/container.ts @@ -21,6 +21,7 @@ import { MembershipApplicationRepository } from '@/domain/cookingClub/membership import { CuisineRepository } from '@/domain/cookingClub/membership/projection/membersByCuisine/CuisineRepository'; import env from '@/app/environment'; import { Postgres, defaultPoolSettings } from '@/lib/postgres'; +import { Hydrator } from '@/lib/eventSourcing/eventStore'; import { Mongo } from '@/lib/mongo'; import { ServerApiVersion } from 'mongodb'; import * as postgresEventStore from '@/app/postgresEventStore'; @@ -112,6 +113,8 @@ export async function configureDependencies(): Promise { registerSingletons(); registerScopedServices(); + const hydrator = new Hydrator(); + const postgres = new Postgres({ user: env.EVENT_STORE_USER, password: env.EVENT_STORE_PASSWORD, From f177531b966f9d7968c786bb8f4410f52d7d0064 Mon Sep 17 00:00:00 2001 From: Marcelo Lazaroni Date: Fri, 26 Sep 2025 15:15:35 +0100 Subject: [PATCH 22/22] Fix compilation errors --- src/app/event.ts | 196 ++-------------------------- src/app/postgresEventStore.ts | 2 +- src/di/container.ts | 3 - src/lib/eventSourcing/event.ts | 6 +- src/lib/eventSourcing/eventStore.ts | 36 +++++ src/lib/eventSourcing/projection.ts | 26 ++++ 6 files changed, 75 insertions(+), 194 deletions(-) create mode 100644 src/lib/eventSourcing/projection.ts diff --git a/src/app/event.ts b/src/app/event.ts index a00c90f..5c078c3 100644 --- a/src/app/event.ts +++ b/src/app/event.ts @@ -1,123 +1,11 @@ -export { accept }; - import { Aggregate, Id, CreationEvent, TransformationEvent, - Event, - EventInfo, } from '@/lib/eventSourcing/event'; -import { Decoder } from '@/lib/json/decoder'; -import { Encoder } from '@/lib/json/encoder'; import * as s from '@/lib/json/schema'; -import * as d from '@/lib/json/decoder'; -import { Schema } from '@/lib/json/schema'; -import { Json } from '@/lib/json/types'; -import { Maybe, Nothing, Just } from '@/lib/Maybe'; -import { Result, Failure } from '@/lib/Result'; -import * as m from '@/lib/Maybe'; -import { POSIX } from '@/lib/time'; - -type Serialized

= s.Infer>>; - -const schema_Serialized = (payload: Schema) => - s.object({ - event_id: Id.schema>>(), - aggregate_id: Id.schema>(), - aggregate_version: s.number, - correlation_id: Id.schema>>(), - causation_id: Id.schema>>(), - recorded_on: POSIX.schema, - payload, - }); - -function toSerialized

({ - info, - payload, -}: { - info: EventInfo; - payload: P; -}): Serialized

{ - return { ...info, payload }; -} - -function fromSerialized

(serialized: Serialized

): { - info: EventInfo; - payload: P; -} { - const { payload, ...info } = serialized; - return { info, payload }; -} - -type EventData = { info: EventInfo; payload: E }; - -const schema_EventData = (s: Schema): Schema> => - schema_Serialized(s).dimap(fromSerialized, toSerialized); - -type Constructor = new (...args: any[]) => T; - -type Schemas> = { - creation: Schema>>; - transformation: Schema>>; -}; - -type EventName = string; - -class Deserializer { - private tmap = new Map, Schemas>(); - - constructor() {} - - // add support for deserializing an aggregate's events. - add>({ - aggregate, - creation, - transformation, - }: { - aggregate: Constructor; - creation: Schema>; - transformation: Schema>; - }): void { - this.tmap.set(aggregate, { - creation: schema_EventData(creation), - transformation: schema_EventData(transformation), - }); - } - - deserialize>( - aggregate: Constructor, - serialized: Json[], - ): Result { - const schemas = this.tmap.get(aggregate) as undefined | Schemas; - if (schemas == undefined) { - throw new Error(`Unknown aggregate ${aggregate.name}`); - } - - if (serialized.length === 0) { - return Failure('No events'); - } - - return d - .decode(serialized[0], schemas.creation.decoder) - .then(({ payload: first, info }) => - d - .decode(serialized.slice(1), d.array(schemas.transformation.decoder)) - .map((es) => { - let aggregate = first.createAggregate(); - let lastEvent = info; - - for (const t of es) { - aggregate = t.payload.transformAggregate(aggregate); - lastEvent = t.info; - } - - return { aggregate, lastEvent }; - }), - ); - } -} class User implements Aggregate { constructor( @@ -127,14 +15,7 @@ class User implements Aggregate { ) {} } -class Other implements Aggregate { - constructor( - readonly aggregateId: Id, - readonly aggregateVersion: number, - ) {} -} - -class CreateUser implements CreationEvent { +export class CreateUser implements CreationEvent { static type: 'CreateUserr' = 'CreateUserr'; static schemaArgs = s.object({ type: s.stringLiteral(CreateUser.type), @@ -153,7 +34,7 @@ class CreateUser implements CreationEvent { } } -class AddName implements TransformationEvent { +export class AddName implements TransformationEvent { static type: 'AddName' = 'AddName'; constructor(readonly values: s.Infer) {} @@ -163,8 +44,6 @@ class AddName implements TransformationEvent { name: s.string, }); - encoder = AddName.schema.encoder as Encoder>; - static schema = AddName.schemaArgs.dimap( (v) => new AddName(v), (v) => v.values, @@ -181,22 +60,21 @@ class AddName implements TransformationEvent { } } -class RemoveName implements TransformationEvent { +export class RemoveName implements TransformationEvent { static type: 'RemoveName' = 'RemoveName'; - constructor(readonly values: s.Infer) {} + constructor(readonly values: s.Infer) {} static schemaArgs = s.object({ - type: s.stringLiteral(AddName.type), + type: s.stringLiteral(RemoveName.type), aggregateId: Id.schema(), name: s.string, }); - encoder = RemoveName.schema.encoder as Encoder>; - static schema = AddName.schemaArgs.dimap( - (v) => new AddName(v), + static schema = RemoveName.schemaArgs.dimap( + (v) => new RemoveName(v), (v) => v.values, ); - readonly schema = AddName.schema; + readonly schema = RemoveName.schema; transformAggregate(agg: User): User { const u = new User( @@ -207,61 +85,3 @@ class RemoveName implements TransformationEvent { return u; } } - -// ---------- - -const vv = accept([AddName, CreateUser]); - -type WW = m.Infer>; - -type EventClass = { type: string; schema: Schema }; - -// Given some event classes, creates a decoder for those classes. -// Makes sure to error if decoding those class object fail, but -// succeeds if the encoded event was of another class. -function accept( - ts: T, -): Decoder>> { - type Ty = s.Infer; - return d.object({ type: d.string }).then(({ type: ty }) => { - const c: undefined | EventClass = ts.find((t) => t.type === ty); - return c - ? (c.schema.decoder.map(Just) as Decoder>) - : d.succeed(Nothing()); - }); -} - -// Create a schema given a list of event classes -function makeSchema( - ts: T, -): Schema> { - type Ty = s.Infer; - - const decoder: Decoder = d - .object({ type: d.string }) - .then(({ type: ty }) => { - const c: undefined | EventClass = ts.find((t) => t.type === ty); - - if (c === undefined) return d.fail(`Unknown event type: ${ty}`); - - return c.schema.decoder as Decoder; - }); - - const encoder: Encoder = new Encoder((v: Ty) => { - const ty = ts.find((t) => t.type === v.type); - if (ty === undefined) { - throw new Error(`Unable to encode unknown event type: ${v.type}`); - } - - return ty.schema.encoder.run(v); - }); - - return new Schema(decoder, encoder); -} - -const ww = new Deserializer(); -ww.add({ - aggregate: User, - creation: makeSchema([CreateUser]).decoder, - transformation: makeSchema([AddName, RemoveName]).decoder, -}); diff --git a/src/app/postgresEventStore.ts b/src/app/postgresEventStore.ts index 0cb8b75..6d42f1f 100644 --- a/src/app/postgresEventStore.ts +++ b/src/app/postgresEventStore.ts @@ -146,7 +146,7 @@ class PostgresEventStore implements EventStore { '{}', // @ts-ignore serialized.recorded_on, - edata.event.type, + edata.event.values.type, ]; try { diff --git a/src/di/container.ts b/src/di/container.ts index 9449339..e8f90b2 100644 --- a/src/di/container.ts +++ b/src/di/container.ts @@ -21,7 +21,6 @@ import { MembershipApplicationRepository } from '@/domain/cookingClub/membership import { CuisineRepository } from '@/domain/cookingClub/membership/projection/membersByCuisine/CuisineRepository'; import env from '@/app/environment'; import { Postgres, defaultPoolSettings } from '@/lib/postgres'; -import { Hydrator } from '@/lib/eventSourcing/eventStore'; import { Mongo } from '@/lib/mongo'; import { ServerApiVersion } from 'mongodb'; import * as postgresEventStore from '@/app/postgresEventStore'; @@ -113,8 +112,6 @@ export async function configureDependencies(): Promise { registerSingletons(); registerScopedServices(); - const hydrator = new Hydrator(); - const postgres = new Postgres({ user: env.EVENT_STORE_USER, password: env.EVENT_STORE_PASSWORD, diff --git a/src/lib/eventSourcing/event.ts b/src/lib/eventSourcing/event.ts index b736ffd..e014725 100644 --- a/src/lib/eventSourcing/event.ts +++ b/src/lib/eventSourcing/event.ts @@ -38,8 +38,10 @@ type Event> = EventClass; // Class which all events derive from. Used for type constraints. abstract class EventClass> { - abstract type: string; - abstract values: { aggregateId: Id }; + abstract values: { + type: string; + aggregateId: Id; + }; abstract schema: Schema; } diff --git a/src/lib/eventSourcing/eventStore.ts b/src/lib/eventSourcing/eventStore.ts index 89dc4cd..f4dba1b 100644 --- a/src/lib/eventSourcing/eventStore.ts +++ b/src/lib/eventSourcing/eventStore.ts @@ -5,6 +5,7 @@ export { type Constructor, type EventData, schema_EventData, + makeSchema, }; import { @@ -17,6 +18,7 @@ import { } from '@/lib/eventSourcing/event'; import { Json } from '@/lib/json/types'; import { Schema } from '@/lib/json/schema'; +import { Encoder } from '@/lib/json/encoder'; import { Decoder } from '@/lib/json/decoder'; import * as s from '@/lib/json/schema'; import * as d from '@/lib/json/decoder'; @@ -196,3 +198,37 @@ class Hydrator { ); } } + +// ----------------------------------------------------------------------- + +type EventConstructor = { type: string; schema: Schema }; + +// Create an efficient schema given a list of event classes +// +// To be used when joining schemas for the Hydrator +function makeSchema( + ts: T, +): Schema> { + type Ty = s.Infer; + + const decoder: Decoder = d + .object({ type: d.string }) + .then(({ type: ty }) => { + const c: undefined | EventConstructor = ts.find((t) => t.type === ty); + + if (c === undefined) return d.fail(`Unknown event type: ${ty}`); + + return c.schema.decoder as Decoder; + }); + + const encoder: Encoder = new Encoder((v: Ty) => { + const ty = ts.find((t) => t.type === v.type); + if (ty === undefined) { + throw new Error(`Unable to encode unknown event type: ${v.type}`); + } + + return ty.schema.encoder.run(v); + }); + + return new Schema(decoder, encoder); +} diff --git a/src/lib/eventSourcing/projection.ts b/src/lib/eventSourcing/projection.ts new file mode 100644 index 0000000..76e8426 --- /dev/null +++ b/src/lib/eventSourcing/projection.ts @@ -0,0 +1,26 @@ +export { accept }; + +import { Maybe, Nothing, Just } from '@/lib/Maybe'; +import { Decoder } from '@/lib/json/decoder'; +import * as d from '@/lib/json/decoder'; +import * as s from '@/lib/json/schema'; +import { Schema } from '@/lib/json/schema'; + +type EventConstructor = { type: string; schema: Schema }; + +// Given some event classes, creates a decoder for those classes. +// Makes sure to error if decoding those class object fail, but +// succeeds if the encoded event was of another class. +// +// To be used in decoding events for projections and reactions. +function accept( + ts: T, +): Decoder>> { + type Ty = s.Infer; + return d.object({ type: d.string }).then(({ type: ty }) => { + const c = ts.find((t) => t.type === ty); + return c + ? (c.schema.decoder.map(Just) as Decoder>) + : d.succeed(Nothing()); + }); +}