Versions in this module Expand all Collapse all v1 v1.1.0 Jul 21, 2026 Changes in this version + const EventAgentSuppressionAdded + const EventDomainSendingFailed + const EventDomainSendingVerified + const EventDomainSuppressionAdded + const EventEmailBlocked + const EventEmailBounced + const EventEmailComplained + const EventEmailDelivered + const EventEmailFailed + const EventEmailFlagged + const EventEmailReceived + const EventEmailReviewApproved + const EventEmailReviewRejected + const EventEmailReviewRequested + const EventEmailSent + const SchemaVersion + var AllEventTypes = []string + var ExperimentalEventTypes = []string + func DeterministicEventID(parts ...string) string + func IsValidEventType(name string) bool + type DeliveryEnqueuer interface + EnqueueDeliveryTx func(ctx context.Context, tx pgx.Tx, deliveryID string) (int64, error) + type Envelope struct + CreatedAt time.Time + Data any + ID string + SchemaVersion string + Type string + type Event struct + AgentID string + ConversationID string + CreatedAt time.Time + Data any + ID string + Labels []string + MessageID string + Type string + UserID string + func NewEvent(eventType, userID string, data any) Event + func (e Event) AsEnvelope() Envelope + type FanOutArgs struct + EventID string + func (FanOutArgs) Kind() string + type FanOutEnqueuer interface + EnqueueFanOutTx func(ctx context.Context, tx pgx.Tx, eventID string) (int64, error) + type FanOutJobs struct + func NewFanOutJobs(pool *pgxpool.Pool, identityStore identityReader, deliveryEnq DeliveryEnqueuer, ...) *FanOutJobs + func (j *FanOutJobs) EnqueueFanOutTx(ctx context.Context, tx pgx.Tx, eventID string) (int64, error) + func (j *FanOutJobs) ReconcilePending(ctx context.Context, pool *pgxpool.Pool) (int, error) + func (j *FanOutJobs) RegisterJobs(w *river.Workers) []*river.PeriodicJob + func (j *FanOutJobs) SetEnqueuer(e jobs.Enqueuer) + type FanOutReconcileArgs struct + func (FanOutReconcileArgs) Kind() string + type FanOutReconcileWorker struct + func (w *FanOutReconcileWorker) Work(ctx context.Context, _ *river.Job[FanOutReconcileArgs]) error + type FanOutWorker struct + func NewFanOutWorker(pool *pgxpool.Pool, identityStore identityReader, deliveryEnq DeliveryEnqueuer, ...) *FanOutWorker + func (w *FanOutWorker) Timeout(*river.Job[FanOutArgs]) time.Duration + func (w *FanOutWorker) Work(ctx context.Context, job *river.Job[FanOutArgs]) error + type FeatureFlag interface + Enabled func() bool + type Outbox interface + DeleteExpiredWebhookEvents func(ctx context.Context) (int, error) + Enabled func() bool + PublishBestEffortTx func(ctx context.Context, tx pgx.Tx, e Event) (wrote bool) + PublishTx func(ctx context.Context, tx pgx.Tx, e Event) error + SetFanOutEnqueuer func(e FanOutEnqueuer) + func NewOutbox(pool *pgxpool.Pool, flag FeatureFlag) Outbox + type OutboxPublisher struct + func NewOutboxPublisher(outbox Outbox, pool *pgxpool.Pool) *OutboxPublisher + func (p *OutboxPublisher) Publish(ctx context.Context, e Event) + type OutboxWorker struct + func NewOutboxWorker(pool *pgxpool.Pool, store identityReader) *OutboxWorker + func (w *OutboxWorker) Start(ctx context.Context) + func (w *OutboxWorker) Tick(ctx context.Context) + func (w *OutboxWorker) WithDeliveryEnqueuer(e DeliveryEnqueuer) *OutboxWorker + func (w *OutboxWorker) WithMetrics(m telemetry.Metrics) *OutboxWorker + type Publisher interface + Publish func(ctx context.Context, e Event) + type StaticFlag bool + func (f StaticFlag) Enabled() bool