Skip to content

Repository files navigation

logicblocks.event.store

PyPI - Version Python - Version Documentation Status CircleCI

Eventing infrastructure for event-sourced architectures.

Table of Contents

Installation

pip install logicblocks-event-store

Usage

Basic Example

import asyncio

from logicblocks.event.store import EventStore, adapters
from logicblocks.event.types import NewEvent, StreamIdentifier
from logicblocks.event.projection import Projector


class ProfileProjector(
    Projector[StreamIdentifier, dict[str, str], dict[str, str]]
):
    def initial_state_factory(self) -> dict[str, str]:
        return {}

    def initial_metadata_factory(self) -> dict[str, str]:
        return {}

    def id_factory(self, state, source: StreamIdentifier) -> str:
        return source.stream

    def profile_created(self, state, event):
        state['name'] = event.payload['name']
        state['email'] = event.payload['email']
        return state

    def date_of_birth_set(self, state, event):
        state['dob'] = event.payload['dob']
        return state


async def main():
    adapter = adapters.InMemoryEventStorageAdapter()
    store = EventStore(adapter)

    stream = store.stream(category="profiles", stream="joe.bloggs")
    # metadata is required; pass metadata=None when the event has no metadata
    profile_created_event = NewEvent(name="profile-created",
                                     payload={"name": "Joe Bloggs", "email": "joe.bloggs@example.com"},
                                     metadata=None)
    date_of_birth_set_event = NewEvent(name="date-of-birth-set", payload={"dob": "1992-07-10"},
                                       metadata={"actor": "user-123"})

    await stream.publish(
        events=[
            profile_created_event
        ])
    await stream.publish(
        events=[
            date_of_birth_set_event
        ]
    )

    projector = ProfileProjector()
    projection = await projector.project(source=stream)
    profile = projection.state


asyncio.run(main())


# profile == {
#   "name": "Joe Bloggs", 
#   "email": "joe.bloggs@example.com", 
#   "dob": "1992-07-10"
# }

Finalising State

Override finalise_state to do work once per projection rather than in every event handler, for example validation or normalisation. Handlers can then make cheap, unvalidated changes and leave the expensive work to the end:

class ValidatedProfileProjector(ProfileProjector):
    def finalise_state(self, state: dict[str, str]) -> dict[str, str]:
        if "email" not in state:
            raise ValueError("profile has no email")
        return {**state, "email": state["email"].strip()}

finalise_state is called once at the end of each project() call, after all events are applied and before id_factory derives the projection id, including when the source has no events. Keep in mind that:

  • with ProjectionEventProcessor, which projects one event at a time, it runs once per processed event, so the saving comes from multi-event folds such as rebuilds;
  • its output is passed back in as the starting state when a projection is resumed (via state=, and always by ProjectionEventProcessor), so it must not change what later handlers, update_metadata or id_factory compute. Validation and recomputing derived fields are fine; lossy normalisation, such as clamping or truncation, isn't;
  • its output must survive a round trip through the projection store, and it must not change fields id_factory uses for projections that are already stored, otherwise those projections need rebuilding;
  • update_metadata and apply() see unfinalised state;
  • handlers mustn't rely on coercion or defaults that validation would apply;
  • finalise_state is a reserved name, so no event may be named finalise-state;
  • subclasses that override project() must call finalise_state themselves.

Features

  • Event modelling:
    • Log / category / stream based: events are grouped into logs of categories of streams.
    • Arbitrary payloads and metadata: events can have arbitrary payloads and metadata limited only by what the underlying storage backend can support.
    • Bi-temporality support: events included timestamps for both the time the event occurred and the time the event was recorded in the log.
  • Event storage:
    • Immutable and append only: the event store is modelled as an append-only log of immutable events.
    • Consistency guarantees: concurrent stream updates can optionally be handled with optimistic concurrency control.
    • Write conditions: an extensible write condition system allows pre-conditions to be evaluated before publish.
    • Ordering guarantees: event writes are serialised (at log level by default, but customisable) to guarantee consistent ordering at scan time.
    • asyncio support: the event store is implemented using asyncio and can be used in cooperative multitasking applications.
  • Storage adapters:
    • Storage adapter abstraction: adapters are provided for different storage backends, currently including:
      • an in-memory implementation for testing and experimentation; and
      • a PostgreSQL backed implementation for production use.
    • Extensible to other backends: the storage adapter abstract base class is designed to be relatively easily implemented to support other storage backends.
  • Projections:
    • Reduction: event sequences can be reduced to a single value, a projection, using a projector.
    • Metadata: projections have metadata for keeping track of things like update timestamps, versions, etc.
    • Storage: a general purpose projection store allows easy management of projections for the majority of use cases, utilising the same adapter architecture as the event store, with a rich and customisable query language providing store search.
    • Snapshotting: coming soon.
  • Types:
    • Type hints: includes type hints for all public classes and functions.
    • Value types: includes serialisable value types for identifiers, events and projections.
    • Pydantic support: coming soon.
  • Testing utilities:
    • Builders: includes builders for events to simplify testing.
    • Data generators: includes random data generators for events and event attributes.
    • Storage adapter tests: includes tests for storage adapters to ensure consistency across implementations.

Documentation

Development

This project uses mise for tool management. To get started:

mise install
mise run

See CONTRIBUTING.md for detailed development instructions.

Contributing

Bug reports and pull requests are welcome on GitHub at https://github.com/logicblocks/event.store.

See CONTRIBUTING.md for guidelines on:

  • Reporting bugs and requesting features
  • Setting up your development environment
  • Running tests and code quality checks
  • Submitting pull requests

This project is intended to be a safe, welcoming space for collaboration, and contributors are expected to adhere to the code of conduct.

License

Copyright © 2025 LogicBlocks Maintainers

Distributed under the terms of the MIT License.