Skip to main content

eventsource-py

PyPI version Python Version CI License: MIT

Stop losing data. Start capturing history.

Traditional databases overwrite state on every update. Event sourcing captures what happened as a sequence of immutable events, giving you:

  • Complete audit trail - Know exactly what changed, when, and why
  • Time travel - Reconstruct state at any point in history
  • Multiple views - Build different read models from the same events
  • Reliable debugging - Replay events to reproduce any bug

eventsource-py makes this practical for Python applications with a clean, async-first API.

pip install eventsource-py

Quick Start

import asyncio
from uuid import UUID, uuid4
from pydantic import BaseModel
from eventsource import (
    DomainEvent, register_event, AggregateRoot, AggregateRepository,
    InMemoryEventStore, InMemoryEventBus, InMemoryCheckpointRepository,
)
from eventsource.application.subscriptions import SubscriptionManager, SubscriptionConfig

# 1. Define events - immutable facts that capture what happened
@register_event
class OrderPlaced(DomainEvent):
    event_type: str = "OrderPlaced"
    aggregate_type: str = "Order"
    customer_id: UUID
    total: float

@register_event
class OrderShipped(DomainEvent):
    event_type: str = "OrderShipped"
    aggregate_type: str = "Order"
    tracking_number: str

# 2. Define aggregate state and business logic
class OrderState(BaseModel):
    order_id: UUID
    customer_id: UUID | None = None
    total: float = 0.0
    status: str = "draft"

class Order(AggregateRoot[OrderState]):
    aggregate_type = "Order"

    def _get_initial_state(self) -> OrderState:
        return OrderState(order_id=self.aggregate_id)

    def _apply(self, event: DomainEvent) -> None:
        match event:
            case OrderPlaced():
                self._state = OrderState(
                    order_id=self.aggregate_id,
                    customer_id=event.customer_id,
                    total=event.total,
                    status="placed",
                )
            case OrderShipped():
                self._state = self._state.model_copy(update={"status": "shipped"})

    def place(self, customer_id: UUID, total: float) -> None:
        if self.version > 0:
            raise ValueError("Order already placed")
        self.apply_event(OrderPlaced(
            aggregate_id=self.aggregate_id,
            customer_id=customer_id,
            total=total,
            aggregate_version=self.get_next_version(),
        ))

    def ship(self, tracking_number: str) -> None:
        if self.state.status != "placed":
            raise ValueError("Order must be placed before shipping")
        self.apply_event(OrderShipped(
            aggregate_id=self.aggregate_id,
            tracking_number=tracking_number,
            aggregate_version=self.get_next_version(),
        ))

# 3. Define a projection - build read models from the event stream
class SalesReport:
    """Read model that tracks sales metrics from events."""
    def __init__(self):
        self.total_revenue = 0.0
        self.orders_placed = 0
        self.orders_shipped = 0

    def subscribed_to(self) -> list[type[DomainEvent]]:
        return [OrderPlaced, OrderShipped]

    async def handle(self, event: DomainEvent) -> None:
        match event:
            case OrderPlaced():
                self.total_revenue += event.total
                self.orders_placed += 1
            case OrderShipped():
                self.orders_shipped += 1

# 4. Wire it together
async def main():
    # Infrastructure
    store = InMemoryEventStore()
    bus = InMemoryEventBus()
    repo = AggregateRepository(
        event_store=store,
        aggregate_factory=Order,
        event_publisher=bus,  # Publishes events to the bus after saving
    )

    # Set up subscription manager with our projection
    manager = SubscriptionManager(store, bus, InMemoryCheckpointRepository())
    report = SalesReport()
    await manager.subscribe(report, SubscriptionConfig(start_from="beginning"), name="SalesReport")
    await manager.start()

    # Create some orders
    for i in range(3):
        order = repo.create_new(uuid4())
        order.place(customer_id=uuid4(), total=100.0 * (i + 1))
        await repo.save(order)

        if i == 0:  # Ship the first order
            shipped_order_id = order.aggregate_id
            order.ship(tracking_number="TRACK-001")
            await repo.save(order)

    await asyncio.sleep(0.1)  # Let events propagate

    # The projection built a read model from the event stream
    print(f"Revenue: ${report.total_revenue}")      # Revenue: $600.0
    print(f"Orders placed: {report.orders_placed}")  # Orders placed: 3
    print(f"Orders shipped: {report.orders_shipped}")  # Orders shipped: 1

    # Events are the source of truth - reload aggregate from its event history
    order = await repo.load(shipped_order_id)
    print(f"Order status: {order.state.status}")  # Order status: shipped
    print(f"Order version: {order.version}")      # Order version: 2 (placed + shipped)

    await manager.stop()

asyncio.run(main())

Alternative: the decider style

The Order aggregate above is imperative — commands are methods, and business rules live inside them. If you prefer your domain as pure functions, the same aggregate can be written in the decider style: commands become values, and the whole domain becomes two functions you can unit-test with plain asserts — no aggregate instance, no event loop, no fixtures.

from eventsource import DeciderAggregate, DomainCommand, CommandRejectedError

# Commands - intents as values, may be rejected
class PlaceOrder(DomainCommand):
    order_id: UUID
    customer_id: UUID
    total: float

class ShipOrder(DomainCommand):
    order_id: UUID
    tracking_number: str

OrderCommand = PlaceOrder | ShipOrder

# The aggregate as three pure static methods - no I/O, no self, no versions
class Order(DeciderAggregate[OrderState]):
    aggregate_type = "Order"

    @staticmethod
    def initial_state() -> OrderState:
        return OrderState()

    @staticmethod
    def decide(command: OrderCommand, state: OrderState) -> list[DomainEvent]:
        """Command + current state -> new events (or a rejection)."""
        match command, state:
            case PlaceOrder(order_id=oid, customer_id=cid, total=total), OrderState(status="draft"):
                return [OrderPlaced(aggregate_id=oid, customer_id=cid, total=total)]
            case PlaceOrder(), _:
                raise CommandRejectedError("Order already placed", command)
            case ShipOrder(order_id=oid, tracking_number=tn), OrderState(status="placed"):
                return [OrderShipped(aggregate_id=oid, tracking_number=tn)]
            case ShipOrder(), _:
                raise CommandRejectedError("Order must be placed before shipping", command)

    @staticmethod
    def evolve(state: OrderState, event: DomainEvent) -> OrderState:
        """State + event -> next state."""
        match event:
            case OrderPlaced(customer_id=cid, total=total):
                return state.model_copy(update={"customer_id": cid, "total": total, "status": "placed"})
            case OrderShipped():
                return state.model_copy(update={"status": "shipped"})
            case _:
                return state

Callers issue commands as values instead of calling methods — execute() is inherited from DeciderAggregate, and everything else (events, state, projection, wiring) is identical to the example above, and so is the output:

order_id = uuid4()
order = repo.create_new(order_id)
order.execute(PlaceOrder(order_id=order_id, customer_id=uuid4(), total=100.0))
await repo.save(order)

order.execute(ShipOrder(order_id=order_id, tracking_number="TRACK-001"))
await repo.save(order)

See The Decider Pattern for the trade-offs and benchmarks (spoiler: identical on replay, a few microseconds per command — maintainability is the deciding factor, not speed).

Production Ready

Swap in production backends when you're ready to deploy:

Component Development Production
Event Store InMemoryEventStore PostgreSQLEventStore, SQLiteEventStore
Event Bus InMemoryEventBus RedisEventBus, RabbitMQEventBus, KafkaEventBus
Checkpoints InMemoryCheckpointRepository PostgreSQLCheckpointRepository
# Add PostgreSQL + Redis for production
pip install eventsource-py[postgresql,redis]
All installation options
pip install eventsource-py[postgresql]  # PostgreSQL event store
pip install eventsource-py[sqlite]      # SQLite event store
pip install eventsource-py[redis]       # Redis event bus
pip install eventsource-py[rabbitmq]    # RabbitMQ event bus
pip install eventsource-py[kafka]       # Kafka event bus
pip install eventsource-py[telemetry]   # OpenTelemetry tracing
pip install eventsource-py[all]         # Everything

Features

  • Event Stores - PostgreSQL, SQLite, In-Memory with optimistic concurrency
  • Event Bus - Redis Streams, RabbitMQ, Kafka, In-Memory with consumer groups
  • Subscriptions - Catch-up from history, live events, checkpointing, graceful shutdown
  • Projections - Declarative handlers, retry logic, dead letter queues
  • Snapshots - Optimize aggregate loading for long event streams
  • Multi-tenancy - Built-in tenant isolation
  • Observability - OpenTelemetry integration

Documentation

Full Documentation - Guides, examples, and API reference

License

MIT

Metadata

Release files for eventsource-py 0.12.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for eventsource-py 0.12.0
File Size Uploaded
eventsource_py-0.12.0.tar.gz 2.9 MB Details

Built distribution (wheel)

Table of built distributions (wheels) for eventsource-py 0.12.0
File Interpreter ABI Platform
eventsource_py-0.12.0-py3-none-any.whl Python 3 none any Details

Total release size: 3.6 MB

Release files / eventsource_py-0.12.0.tar.gz

Download URL eventsource_py-0.12.0.tar.gz
Size 2.9 MB
Tags Source
SHA-256 checksum
How to use checksums
d173a3304c3dad22588a8cd57f8e1b80fd578114e6f70db373fbcb4a60f390b4
BLAKE2b-256 checksum
How to use checksums
56245cd0d851e1dd0682c21b57bfed50ad6228a48bb3e66a97e1cfe77f0e6069
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Aug 6, 2026.

Transparency log

Release files / eventsource_py-0.12.0-py3-none-any.whl

Download URL eventsource_py-0.12.0-py3-none-any.whl
Size 720.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
0714c4390766ec683dbd8990b8326fca48ee7901e3381b522703e93024d13205
BLAKE2b-256 checksum
How to use checksums
fcb8ebe94477843b8f2c9f5d4b64ba045d4ded074c76949f47b616a6317c2f58
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
Yes
Uploaded via twine/7.0.0 CPython/3.13.14

Provenance

Provenance describes where a file came from. On PyPI, provenance is shared via attestations, which provide a verifiable record of the build or publishing details. View details, limitations and caveats.

PyPI Publish Attestation

PyPI verified that this artifact, at this checksum, originated from the publisher listed below.

Signed by GitHub Actions, verified by PyPI on Aug 6, 2026.

Transparency log

Release history Release notifications | RSS feed

0.14.0

2 release files

This release

0.12.0 This release

2 release files

0.9.1

2 release files

0.9.0

2 release files

0.8.1

2 release files

0.8.0

2 release files

0.5.0

2 release files

0.4.0

2 release files

0.3.1

2 release files

0.3.0

2 release files

0.2.0

2 release files

0.1.3

2 release files

0.1.2

2 release files

0.1.0

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page