Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions .changeset/effect-3-22-update.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
---
'@codeforbreakfast/eventsourcing-store-postgres': patch
'@codeforbreakfast/eventsourcing-transport-websocket': patch
'@codeforbreakfast/eventsourcing-testing-contracts': patch
---

Build and test against Effect 3.22.

`eventsourcing-store-postgres` now depends on `@effect/sql` 0.52, `@effect/sql-pg` 0.53 and `@effect/experimental` 0.61. `eventsourcing-transport-websocket` now depends on `@effect/platform` 0.97.

`eventsourcing-store-postgres` fixes lost events on new subscriptions. `subscribe` and `subscribeAll` used to return before Postgres `LISTEN` was active, so an event committed in that gap never reached the subscriber. They now wait until `LISTEN` is active. To detect that, the store sends probe notifications to the channel it is listening on, repeating every 10ms until one comes back. A probe payload starts with `eventstore_listen_probe:`, and every listener on that channel receives it. If your own code listens on the `eventstore_events_*` channels, ignore payloads with that prefix. Older versions of this package log a parse error for each probe they receive.

In `eventsourcing-testing-contracts`, the test helpers fail with Effect's tagged errors instead of a plain `Error`. `expectError` fails with `NoSuchElementException` when the error does not match the predicate. `waitForConnectionState` and `collectMessages` fail with `TimeoutException` when they time out. The declared error type is still `Error`, so existing code compiles unchanged, and you can now match these failures by `_tag`.
104 changes: 47 additions & 57 deletions bun.lock

Large diffs are not rendered by default.

6 changes: 3 additions & 3 deletions examples/todo-app/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -21,9 +21,9 @@
"@codeforbreakfast/eventsourcing-projections": "workspace:*",
"@codeforbreakfast/eventsourcing-store": "workspace:*",
"@codeforbreakfast/eventsourcing-store-filesystem": "workspace:*",
"@effect/cli": "0.72.1",
"@effect/platform-bun": "0.86.0",
"effect": "3.19.9"
"@effect/cli": "0.77.2",
"@effect/platform-bun": "0.91.2",
"effect": "3.22.2"
},
"devDependencies": {
"@types/bun": "latest"
Expand Down
16 changes: 11 additions & 5 deletions examples/todo-app/src/domain/todoAggregate.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { Effect, Match, Option, ParseResult, Schema, pipe } from 'effect';
import { Data, Effect, Match, Option, ParseResult, Schema, pipe } from 'effect';
import {
makeAggregateRoot,
defineAggregateEventStore,
Expand Down Expand Up @@ -91,20 +91,26 @@ const applyEvent =
Match.orElse(() => handleNonCreatedEvent(state, event))
);

class TodoCommandError extends Data.TaggedError('TodoCommandError')<{
readonly message: string;
}> {}

const requireExistingTodo = <A, E, R>(
operation: string,
onSome: (state: TodoState) => Effect.Effect<A, E, R>
): ((state: Readonly<Option.Option<TodoState>>) => Effect.Effect<A, E | Error, R>) =>
): ((state: Readonly<Option.Option<TodoState>>) => Effect.Effect<A, E | TodoCommandError, R>) =>
Option.match({
onNone: () => Effect.fail(new Error(`Cannot ${operation} non-existent TODO`)),
onNone: () =>
Effect.fail(new TodoCommandError({ message: `Cannot ${operation} non-existent TODO` })),
onSome,
});

const failIfDeletedTodo =
(operation: string) =>
(state: TodoState): Effect.Effect<TodoState, Error> =>
(state: TodoState): Effect.Effect<TodoState, TodoCommandError> =>
Effect.if(state.deleted, {
onTrue: () => Effect.fail(new Error(`Cannot ${operation} deleted TODO`)),
onTrue: () =>
Effect.fail(new TodoCommandError({ message: `Cannot ${operation} deleted TODO` })),
onFalse: () => Effect.succeed(state),
});

Expand Down
14 changes: 2 additions & 12 deletions examples/todo-app/src/infrastructure/processManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,12 +16,7 @@ const isTodoCreated = (event: unknown): event is EventRecord<TodoCreated, UserId
const isTodoDeleted = (event: unknown): event is EventRecord<TodoDeleted, UserId> =>
typeof event === 'object' && event !== null && 'type' in event && event.type === 'TodoDeleted';

const parseUserId = (userId: unknown): Effect.Effect<UserId, Error> =>
pipe(
userId,
Schema.decodeUnknown(UserIdSchema),
Effect.mapError(() => new Error(`Invalid UserId: ${String(userId)}`))
);
const parseUserId = Schema.decodeUnknown(UserIdSchema);

const provideInitiatorFromUserId =
<TResult, TError>(effect: Effect.Effect<TResult, TError>) =>
Expand Down Expand Up @@ -49,12 +44,7 @@ const executeAndCommit = <TEvent extends EventRecord<unknown, UserId>, TError>(
) =>
pipe(command, withCommandInitiator(event), Effect.flatMap(commitEvents(state.nextEventNumber)));

const parseTodoId = (todoId: string): Effect.Effect<TodoId, Error> =>
pipe(
todoId,
Schema.decode(TodoIdSchema),
Effect.mapError(() => new Error(`Invalid TodoId: ${todoId}`))
);
const parseTodoId = Schema.decode(TodoIdSchema);

const handleCommand = <TEvent extends EventRecord<unknown, UserId>, TError>(
event: TEvent,
Expand Down
14 changes: 7 additions & 7 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -47,17 +47,17 @@
"@changesets/cli": "2.31.1",
"@commitlint/cli": "20.5.3",
"@commitlint/config-conventional": "20.5.3",
"@effect/cluster": "0.55.0",
"@effect/language-service": "0.60.0",
"@effect/platform": "0.93.6",
"@effect/platform-bun": "0.86.0",
"@effect/rpc": "0.72.2",
"@effect/workflow": "0.15.0",
"@effect/cluster": "0.60.2",
"@effect/language-service": "0.87.3",
"@effect/platform": "0.97.2",
"@effect/platform-bun": "0.91.2",
"@effect/rpc": "0.76.2",
"@effect/workflow": "0.19.1",
"@types/bun": "1.4.2",
"@typescript-eslint/eslint-plugin": "8.70.1",
"@typescript-eslint/parser": "8.70.1",
"dependency-cruiser": "17.4.3",
"effect": "3.19.9",
"effect": "3.22.2",
"eslint": "9.39.5",
"eslint-config-prettier": "10.1.8",
"eslint-plugin-compat": "6.2.1",
Expand Down
2 changes: 1 addition & 1 deletion packages/bun-test-effect/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,6 @@
"effect": ">=3.0.0"
},
"devDependencies": {
"effect": "3.19.9"
"effect": "3.22.2"
}
}
2 changes: 1 addition & 1 deletion packages/eslint-effect/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,6 @@
"eslint": "9.39.5",
"@typescript-eslint/eslint-plugin": "8.70.1",
"typescript": "5.9.3",
"effect": "3.19.9"
"effect": "3.22.2"
}
}
4 changes: 2 additions & 2 deletions packages/eslint-effect/test/no-if-statement.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { Effect, Either, Option } from 'effect';
import { Cause, Effect, Either, Option } from 'effect';

const either = Either.right(42);
const option = Option.some(42);
Expand All @@ -23,7 +23,7 @@ if (event.type === 'TodoCreated') {

// eslint-disable-next-line effect/no-if-statement
if (state.deleted) {
const _unused1 = Effect.fail(new Error('Cannot complete deleted TODO'));
const _unused1 = Effect.fail(new Cause.RuntimeException('Cannot complete deleted TODO'));
} else {
const _unused2 = Effect.succeed(state);
}
Expand Down
8 changes: 5 additions & 3 deletions packages/eslint-effect/test/prefer-match-over-ternary.test.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { pipe, Effect, Match, Option } from 'effect';
import { pipe, Cause, Effect, Match, Option } from 'effect';

// Should fail - ternary with Effect calls in return statement
const ternaryWithEffect = (condition: boolean) => {
Expand Down Expand Up @@ -49,8 +49,10 @@ const ternaryInEffectSucceed = (condition: boolean) => Effect.succeed(condition

// Should fail - ternary inside Effect.fail (complex condition)
const ternaryInEffectFail = (hasError: boolean) =>
// eslint-disable-next-line effect/prefer-match-over-ternary
Effect.fail(hasError ? new Error('Critical') : new Error('Warning'));
Effect.fail(
// eslint-disable-next-line effect/prefer-match-over-ternary
hasError ? new Cause.RuntimeException('Critical') : new Cause.RuntimeException('Warning')
);

// Should NOT fail - simple literal equality INSIDE Effect.succeed (no duplication)
const simpleLiteralInEffect = (id: string) => Effect.succeed(id === 'user-1' ? 'John' : 'Guest');
Expand Down
2 changes: 1 addition & 1 deletion packages/eventsourcing-aggregates/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@
"@types/node": "24.19.0",
"@codeforbreakfast/bun-test-effect": "workspace:*",
"@codeforbreakfast/eventsourcing-store-inmemory": "workspace:*",
"effect": "3.19.9",
"effect": "3.22.2",
"typescript": "5.9.3",
"eslint": "9.39.5"
},
Expand Down
2 changes: 1 addition & 1 deletion packages/eventsourcing-commands/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@
"devDependencies": {
"@codeforbreakfast/bun-test-effect": "workspace:*",
"@types/node": "24.19.0",
"effect": "3.19.9",
"effect": "3.22.2",
"typescript": "5.9.3",
"eslint": "9.39.5"
},
Expand Down
2 changes: 1 addition & 1 deletion packages/eventsourcing-projections/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@
},
"devDependencies": {
"@types/node": "24.19.0",
"effect": "3.19.9",
"effect": "3.22.2",
"typescript": "5.9.3",
"eslint": "9.39.5"
},
Expand Down
2 changes: 1 addition & 1 deletion packages/eventsourcing-protocol/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -77,7 +77,7 @@
"@codeforbreakfast/bun-test-effect": "workspace:*",
"@codeforbreakfast/eventsourcing-transport-inmemory": "workspace:*",
"@types/node": "24.19.0",
"effect": "3.19.9",
"effect": "3.22.2",
"typescript": "5.9.3",
"eslint": "9.39.5"
},
Expand Down
2 changes: 1 addition & 1 deletion packages/eventsourcing-server/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@
"@codeforbreakfast/eventsourcing-store-inmemory": "workspace:*",
"@codeforbreakfast/eventsourcing-testing-contracts": "workspace:*",
"@codeforbreakfast/bun-test-effect": "workspace:*",
"effect": "3.19.9",
"effect": "3.22.2",
"typescript": "5.9.3",
"eslint": "9.39.5"
},
Expand Down
4 changes: 2 additions & 2 deletions packages/eventsourcing-store-filesystem/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -68,9 +68,9 @@
"devDependencies": {
"@codeforbreakfast/bun-test-effect": "workspace:*",
"@codeforbreakfast/eventsourcing-testing-contracts": "workspace:*",
"@effect/platform-bun": "0.86.0",
"@effect/platform-bun": "0.91.2",
"@types/node": "24.19.0",
"effect": "3.19.9",
"effect": "3.22.2",
"eslint": "9.39.5",
"typescript": "5.9.3"
},
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import {
Cause,
Chunk,
Effect,
Stream,
Expand Down Expand Up @@ -318,11 +319,8 @@ const appendEventsToStream =
);
};

const parseJsonContent = <V>(content: string): Effect.Effect<V, Error, never> =>
Effect.try({
try: () => JSON.parse(content) as V,
catch: () => new Error('Failed to parse event'),
});
const parseJsonContent = <V>(content: string): Effect.Effect<V, Cause.UnknownException, never> =>
Effect.try(() => JSON.parse(content) as V);

const readEventFromFile = <V>(
eventPath: string,
Expand Down
2 changes: 1 addition & 1 deletion packages/eventsourcing-store-inmemory/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -69,7 +69,7 @@
"@codeforbreakfast/bun-test-effect": "workspace:*",
"@codeforbreakfast/eventsourcing-testing-contracts": "workspace:*",
"@types/node": "24.19.0",
"effect": "3.19.9",
"effect": "3.22.2",
"typescript": "5.9.3",
"eslint": "9.39.5"
},
Expand Down
10 changes: 5 additions & 5 deletions packages/eventsourcing-store-postgres/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -67,17 +67,17 @@
},
"dependencies": {
"@codeforbreakfast/eventsourcing-store": "workspace:*",
"@effect/experimental": "0.57.11",
"@effect/sql": "0.48.6",
"@effect/sql-pg": "0.49.7",
"@effect/experimental": "0.61.1",
"@effect/sql": "0.52.1",
"@effect/sql-pg": "0.53.0",
"type-fest": "5.10.0"
},
"devDependencies": {
"@codeforbreakfast/bun-test-effect": "workspace:*",
"@codeforbreakfast/eventsourcing-testing-contracts": "workspace:*",
"@effect/platform-bun": "0.86.0",
"@effect/platform-bun": "0.91.2",
"@types/node": "24.19.0",
"effect": "3.19.9",
"effect": "3.22.2",
"eslint": "9.39.5",
"typescript": "5.9.3"
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,11 +42,15 @@
events: readonly string[]
) => pipe(events, Stream.fromIterable, Stream.run(store.append({ streamId, eventNumber: 0 })));

const forkPerStreamSubscription = (store: EventStore<string>, streamId: EventStreamId) => {
const forkPerStreamSubscription = (
store: EventStore<string>,
streamId: EventStreamId,
count: number
) => {
const subscription = store.subscribe({ streamId, eventNumber: 0 });
return pipe(
subscription,
Effect.flatMap((stream) => collectEvents(stream, 2))
Effect.flatMap((stream) => collectEvents(stream, count))
);
};

Expand All @@ -60,7 +64,7 @@

const setupDualSubscriptions = (store: EventStore<string>, streamId: EventStreamId) =>
Effect.all({
perStreamFiber: forkPerStreamSubscription(store, streamId),
perStreamFiber: forkPerStreamSubscription(store, streamId, 2),
allEventsFiber: forkAllEventsSubscription(store, 2),
});

Expand Down Expand Up @@ -209,6 +213,63 @@
})
);

const appendImmediatelyAndJoin = <A>(
store: EventStore<string>,
streamId: EventStreamId,
fiber: Fiber.Fiber<Chunk.Chunk<A>, unknown>
) =>
pipe(
appendEvents(store, streamId, ['immediate-event']),
Effect.andThen(Fiber.join(fiber)),
Effect.timeout('2 seconds')
);

const runImmediatePerStreamTest = (streamId: EventStreamId) =>
pipe(
StringEventStore,
Effect.flatMap((store) =>
Effect.flatMap(forkPerStreamSubscription(store, streamId, 1), (fiber) =>
appendImmediatelyAndJoin(store, streamId, fiber)
)
),
Effect.map((events) => {
expect(Array.from(events)).toEqual(['immediate-event']);
})
);

const runImmediateAllEventsTest = (streamId: EventStreamId) =>
pipe(
StringEventStore,
Effect.flatMap((store) =>
Effect.flatMap(forkAllEventsSubscription(store, 1), (fiber) =>
appendImmediatelyAndJoin(store, streamId, fiber)
)
),
Effect.map((events) => {
expect(Array.from(Chunk.map(events, (e) => e.event))).toEqual(['immediate-event']);
})
);

describe('Subscriptions are live as soon as they return', () => {
it('should deliver an event appended as soon as subscribe returns', () =>

Check failure on line 254 in packages/eventsourcing-store-postgres/src/bridge.integration.test.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Add at least one assertion to this test case.

See more on https://sonarcloud.io/project/issues?id=CodeForBreakfast_eventsourcing&issues=AaDxYJgZZu8ixgGxqO-A&open=AaDxYJgZZu8ixgGxqO-A&pullRequest=399
pipe(
randomId(),
decodeStreamId,
Effect.flatMap(runImmediatePerStreamTest),
Effect.provide(TestLayer),
Effect.runPromise
));

it('should deliver an event appended as soon as subscribeAll returns', () =>

Check failure on line 263 in packages/eventsourcing-store-postgres/src/bridge.integration.test.ts

View check run for this annotation

SonarQubeCloud / SonarCloud Code Analysis

Add at least one assertion to this test case.

See more on https://sonarcloud.io/project/issues?id=CodeForBreakfast_eventsourcing&issues=AaDxYJgaZu8ixgGxqO-B&open=AaDxYJgaZu8ixgGxqO-B&pullRequest=399
pipe(
randomId(),
decodeStreamId,
Effect.flatMap(runImmediateAllEventsTest),
Effect.provide(TestLayer),
Effect.runPromise
));
});

describe('Bridge notification publishing', () => {
it('should publish events to BOTH per-stream subscribers AND all-events subscribers', () => {
const streamId = decodeStreamId(randomId());
Expand Down
Loading
Loading