Event Store Port¶
event_store_port
¶
Event store port for event-sourced aggregates.
Defines the EventStorePort contract for appending and retrieving
domain events. The interface is agnostic of storage backend — in-memory,
relational, or event-native implementations are all supported.
EventStorePort
¶
Bases: OutboundPort
Abstract base class for event stores that persist and retrieve domain events.
Responsibilities
- Append domain events to an aggregate's event stream.
- Retrieve events by version range from a stream.
- Track and report the current version of each stream.
- Enforce optimistic concurrency via
expected_version.
Non-Responsibilities
- Serialize or deserialize events — that belongs to infrastructure.
- Validate event schemas or enforce event versioning.
- Manage snapshots or compaction strategies.
Example
Source code in src/forging_blocks/application/ports/outbound/event_store_port.py
append_events(aggregate_id: UUID, events: Sequence[Event[EventPayloadType]], expected_version: int | None = None) -> Result[int, EventStoreError]
abstractmethod
async
¶
Append events to an aggregate's event stream.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
aggregate_id
|
UUID
|
The aggregate identifier. |
required |
events
|
Sequence[Event[EventPayloadType]]
|
The domain events to append. |
required |
expected_version
|
int | None
|
Expected current version for optimistic |
None
|
Returns:
| Type | Description |
|---|---|
Result[int, EventStoreError]
|
A |
Result[int, EventStoreError]
|
|
Source code in src/forging_blocks/application/ports/outbound/event_store_port.py
get_events(aggregate_id: UUID, from_version: int | None = None, to_version: int | None = None) -> Result[Sequence[Event[EventPayloadType]], EventStoreError]
abstractmethod
async
¶
Retrieve events from an aggregate's event stream.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
aggregate_id
|
UUID
|
The aggregate identifier. |
required |
from_version
|
int | None
|
Start version (inclusive). |
None
|
to_version
|
int | None
|
End version (inclusive). |
None
|
Returns:
| Type | Description |
|---|---|
Result[Sequence[Event[EventPayloadType]], EventStoreError]
|
A |
Result[Sequence[Event[EventPayloadType]], EventStoreError]
|
|
Source code in src/forging_blocks/application/ports/outbound/event_store_port.py
get_current_version(aggregate_id: UUID) -> Result[int, EventStoreError]
abstractmethod
async
¶
Retrieve the current version of an aggregate's event stream.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
aggregate_id
|
UUID
|
The aggregate identifier. |
required |
Returns:
| Type | Description |
|---|---|
Result[int, EventStoreError]
|
A |
Result[int, EventStoreError]
|
streams) on success, or an |