postgres

package
v0.4.4 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 2, 2026 License: MIT Imports: 19 Imported by: 0

Documentation

Overview

Package postgres contains implementations of go-eventually interfaces specific to PostgreSQL, such as Aggregate Repository, Event Store, etc.

Index

Constants

View Source
const (
	// DefaultAggregateTableName is the default aggregate snapshot table.
	DefaultAggregateTableName = "aggregates"
	// DefaultEventsTableName is the default domain events table.
	DefaultEventsTableName = "events"
	// DefaultStreamsTableName is the default event streams table.
	DefaultStreamsTableName = "event_streams"
)

Variables

View Source
var ErrTransactionRequired = errors.New("postgres.TransactionAwareAggregateRepository: transaction required but not found in context")

ErrTransactionRequired reports a missing transaction in context.

Functions

func RunMigrations

func RunMigrations(db *sql.DB) error

RunMigrations runs the latest migrations for the postgres integration.

Make sure to run these in the entrypoint of your application, ideally before building a postgres interface implementation.

Types

type AggregateRepository

type AggregateRepository[ID aggregate.ID, T aggregate.Root[ID]] struct {
	// contains filtered or unexported fields
}

AggregateRepository implements aggregate.Repository for PostgreSQL. It reads aggregate snapshots from the pool and saves snapshots and recorded domain events together in a self-managed Serializable transaction. Table names can be configured using the available functional options.

func NewAggregateRepository

func NewAggregateRepository[ID aggregate.ID, T aggregate.Root[ID]](
	conn *pgxpool.Pool,
	aggregateType aggregate.Type[ID, T],
	aggregateSerde serde.Bytes[T],
	messageSerde serde.Bytes[message.Message],
	options ...Option[ID, T],
) AggregateRepository[ID, T]

NewAggregateRepository returns a new AggregateRepository instance.

func (AggregateRepository[ID, T]) Get

func (repo AggregateRepository[ID, T]) Get(ctx context.Context, id ID) (T, error)

Get returns the aggregate.Root instance specified by id. Returns aggregate.ErrRootNotFound if the aggregate does not exist.

func (AggregateRepository[ID, T]) Save

func (repo AggregateRepository[ID, T]) Save(ctx context.Context, root T) error

Save saves the snapshot and recorded events of root in its own transaction.

type EventStore

type EventStore struct {
	// contains filtered or unexported fields
}

EventStore is an event.Store implementation targeted to PostgreSQL databases.

The implementation uses "event_streams" and "events" as their operational tables. Updates to these tables are transactional.

func NewEventStore

func NewEventStore(conn *pgxpool.Pool, messageSerde serde.Bytes[message.Message]) EventStore

NewEventStore returns a new EventStore instance.

func (EventStore) Append

func (es EventStore) Append(
	ctx context.Context,
	id event.StreamID,
	expected version.Check,
	events ...event.Envelope,
) (version.Version, error)

Append implements event.Store.

func (EventStore) Stream

func (es EventStore) Stream(
	ctx context.Context,
	id event.StreamID,
	selector version.Selector,
) *event.Stream

Stream implements the event.Streamer interface.

type Option

type Option[ID aggregate.ID, T aggregate.Root[ID]] interface {
	// contains filtered or unexported methods
}

Option configures the shared implementation of PostgreSQL aggregate repositories.

func WithAggregateTableName

func WithAggregateTableName[ID aggregate.ID, T aggregate.Root[ID]](tableName string) Option[ID, T]

WithAggregateTableName configures the aggregate snapshot table.

func WithEventsTableName

func WithEventsTableName[ID aggregate.ID, T aggregate.Root[ID]](tableName string) Option[ID, T]

WithEventsTableName configures the domain events table.

func WithStreamsTableName

func WithStreamsTableName[ID aggregate.ID, T aggregate.Root[ID]](tableName string) Option[ID, T]

WithStreamsTableName configures the event streams table.

type TransactionAwareAggregateRepository added in v0.4.4

type TransactionAwareAggregateRepository[ID aggregate.ID, T aggregate.Root[ID]] struct {
	// contains filtered or unexported fields
}

TransactionAwareAggregateRepository implements aggregate.Repository using caller-owned transactions for both reads and writes. It never begins, commits, or rolls back a transaction. Transactions must use Serializable isolation. Save success means writes are staged; the caller must commit the transaction. After a failed or rolled-back save, discard the aggregate and reload before retrying.

func NewTransactionAwareAggregateRepository added in v0.4.4

func NewTransactionAwareAggregateRepository[ID aggregate.ID, T aggregate.Root[ID]](
	retrieveTx TxRetriever,
	aggregateType aggregate.Type[ID, T],
	aggregateSerde serde.Bytes[T],
	messageSerde serde.Bytes[message.Message],
	options ...Option[ID, T],
) TransactionAwareAggregateRepository[ID, T]

NewTransactionAwareAggregateRepository returns a repository requiring a transaction from retrieveTx for every operation.

func (TransactionAwareAggregateRepository[ID, T]) Get added in v0.4.4

func (repo TransactionAwareAggregateRepository[ID, T]) Get(ctx context.Context, id ID) (T, error)

Get reads an aggregate snapshot within the transaction retrieved from ctx. Returns ErrTransactionRequired if no transaction is available, or aggregate.ErrRootNotFound if the aggregate does not exist.

func (TransactionAwareAggregateRepository[ID, T]) Save added in v0.4.4

func (repo TransactionAwareAggregateRepository[ID, T]) Save(ctx context.Context, root T) error

Save writes the snapshot and recorded events within the transaction retrieved from ctx. Returns ErrTransactionRequired without flushing events if no transaction is available. The caller must roll back on error.

type TxRetriever added in v0.4.4

type TxRetriever func(context.Context) (pgx.Tx, bool)

TxRetriever retrieves a caller-owned transaction from context.

Directories

Path Synopsis
Package internal contains utilities and helper functions useful in the scope of eventually's postgres implementation.
Package internal contains utilities and helper functions useful in the scope of eventually's postgres implementation.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL