Versions in this module Expand all Collapse all v1 v1.5.0 Aug 1, 2026 Changes in this version + const EventContactDue v1.4.1 Jul 26, 2026 v1.4.0 Jul 26, 2026 v1.3.0 Jul 25, 2026 v1.2.2 Jul 24, 2026 v1.2.1 Jul 23, 2026 v1.2.0 Jul 22, 2026 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