messaging

package
v0.7.0 Latest Latest
Warning

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

Go to latest
Published: Jul 23, 2026 License: MIT Imports: 3 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Consumer

type Consumer interface {
	// Poll blocks until at least one event is available, ctx is cancelled, or
	// the consumer is closed. A non-empty batch may be returned alongside a nil
	// error even when some partitions failed.
	Poll(ctx context.Context) ([]*event.Event, error)

	// Commit acknowledges every event returned by Poll so far. Call it only
	// after the events have been processed — committing first turns a crash
	// mid-processing into silent data loss.
	Commit(ctx context.Context) error

	// Ping reports broker reachability, for the consuming service's health probe.
	Ping(ctx context.Context) error
}

Consumer reads events from a stream.

Delivery is at-least-once: an event may be redelivered after a crash or a rebalance, so implementations of the read loop must deduplicate on event.Event.ID rather than assume each event arrives once.

type Publisher

type Publisher interface {
	io.Closer
	Publish(ctx context.Context, event *event.Event) error
	Ping(ctx context.Context) error
}

Producer is an event producer

Directories

Path Synopsis
Package kafka implements the messaging interfaces against a Kafka broker.
Package kafka implements the messaging interfaces against a Kafka broker.
Package mock is a generated GoMock package.
Package mock is a generated GoMock package.

Jump to

Keyboard shortcuts

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