client

package
v0.51.0 Latest Latest
Warning

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

Go to latest
Published: Aug 5, 2026 License: Apache-2.0 Imports: 10 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	ErrNoTransport     = errors.New("client has no transport configured")
	ErrInvalidProtocol = errors.New("invalid protocol")
	ErrNoRouteProtocol = errors.New("requested protocol transport not configured")
	ErrConnectFailed   = errors.New("connect failed")
	ErrCloseFailed     = errors.New("close failed")
	ErrPublishFailed   = errors.New("publish failed")
	ErrSubscribeFailed = errors.New("subscribe failed")
	ErrUnsubFailed     = errors.New("unsubscribe failed")

	ErrQueuePublishFailed   = errors.New("queue publish failed")
	ErrQueueSubscribeFailed = errors.New("queue subscribe failed")
	ErrQueueUnsubFailed     = errors.New("queue unsubscribe failed")
)

Functions

This section is empty.

Types

type Client

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

Client provides a unified messaging API over MQTT or AMQP transports.

func New

func New(cfg *Config) (*Client, error)

New creates a unified client configured with one or both transports.

func NewAMQP

func NewAMQP(opts *amqp.Options) (*Client, error)

NewAMQP creates a unified client backed by the AMQP client.

func NewMQTT

func NewMQTT(opts *mqtt.Options) (*Client, error)

NewMQTT creates a unified client backed by the MQTT client.

func (*Client) AMQP

func (c *Client) AMQP() *amqp.Client

AMQP returns the underlying AMQP client, if configured.

func (*Client) Close

func (c *Client) Close(ctx context.Context) error

Close terminates the connection.

func (*Client) Connect

func (c *Client) Connect(ctx context.Context) error

Connect establishes a connection to the broker.

func (*Client) IsConnected

func (c *Client) IsConnected() bool

IsConnected reports whether the client is connected.

func (*Client) MQTT

func (c *Client) MQTT() *mqtt.Client

MQTT returns the underlying MQTT client, if configured.

func (*Client) Publish

func (c *Client) Publish(ctx context.Context, topic string, payload []byte, opts ...Option) error

Publish publishes a message to a topic.

func (*Client) PublishToQueue

func (c *Client) PublishToQueue(ctx context.Context, queue string, payload []byte, opts ...Option) error

PublishToQueue publishes a message to a durable queue.

func (*Client) Subscribe

func (c *Client) Subscribe(ctx context.Context, topic string, handler MessageHandler, opts ...Option) error

Subscribe subscribes to a topic and routes matching messages to handler.

func (*Client) SubscribeToQueue

func (c *Client) SubscribeToQueue(ctx context.Context, queue, group string, handler MessageHandler, opts ...Option) error

SubscribeToQueue subscribes to a queue with a consumer group.

func (*Client) Unsubscribe

func (c *Client) Unsubscribe(ctx context.Context, topic string, opts ...Option) error

Unsubscribe removes a topic subscription.

func (*Client) UnsubscribeFromQueue

func (c *Client) UnsubscribeFromQueue(ctx context.Context, queue string, opts ...Option) error

UnsubscribeFromQueue removes a queue subscription.

type Config

type Config struct {
	MQTT            *mqttclient.Options
	AMQP            *amqpclient.Options
	DefaultProtocol Protocol
}

Config configures the unified client transports.

func NewConfig

func NewConfig() *Config

NewConfig creates Config with sensible defaults.

func (*Config) SetAMQP

func (c *Config) SetAMQP(opts *amqpclient.Options) *Config

SetAMQP sets AMQP transport options.

func (*Config) SetDefaultProtocol

func (c *Config) SetDefaultProtocol(protocol Protocol) *Config

SetDefaultProtocol sets the protocol used when operation options do not specify one.

func (*Config) SetMQTT

func (c *Config) SetMQTT(opts *mqttclient.Options) *Config

SetMQTT sets MQTT transport options.

type Message

type Message struct {
	Topic      string
	Payload    []byte
	Properties map[string]string
	Timestamp  time.Time

	Queue  string
	Offset uint64
	// contains filtered or unexported fields
}

Message represents a unified message across protocols.

func (*Message) Ack

func (m *Message) Ack() error

Ack acknowledges a queue message when supported.

func (*Message) Nack

func (m *Message) Nack() error

Nack negatively acknowledges a queue message when supported.

func (*Message) Reject

func (m *Message) Reject() error

Reject rejects a queue message when supported.

type MessageHandler

type MessageHandler func(msg *Message)

MessageHandler handles an incoming unified message.

type Option

type Option interface {
	// contains filtered or unexported methods
}

Option applies to publish and/or subscribe operations.

func WithAutoAck

func WithAutoAck(autoAck bool) Option

WithAutoAck controls AMQP auto-ack for subscriptions (publish/subscribe only). It has no effect for MQTT subscriptions.

func WithExchange

func WithExchange(exchange string) Option

WithExchange sets the AMQP exchange (publish only).

func WithImmediate

func WithImmediate(immediate bool) Option

WithImmediate sets the AMQP immediate flag (publish only).

func WithMandatory

func WithMandatory(mandatory bool) Option

WithMandatory sets the AMQP mandatory flag (publish only).

func WithProperties

func WithProperties(props map[string]string) Option

WithProperties sets unified properties.

func WithProtocol

func WithProtocol(protocol Protocol) Option

WithProtocol sets protocol routing for publish/subscribe operations.

func WithQoS

func WithQoS(qos byte) Option

WithQoS sets the MQTT QoS for publish/subscribe.

func WithRetain

func WithRetain(retain bool) Option

WithRetain sets the MQTT retain flag (publish only).

func WithRoutingKey

func WithRoutingKey(key string) Option

WithRoutingKey sets the AMQP routing key (publish only).

type Protocol

type Protocol string

Protocol identifies the messaging transport.

const (
	// ProtocolMQTT routes operations over MQTT.
	ProtocolMQTT Protocol = "mqtt"
	// ProtocolAMQP routes operations over AMQP 0.9.1.
	ProtocolAMQP Protocol = "amqp"
)

type PublishOptions

type PublishOptions struct {
	Protocol   Protocol
	QoS        *byte
	Retain     *bool
	Exchange   string
	RoutingKey string
	Mandatory  *bool
	Immediate  *bool

	Properties map[string]string
}

PublishOptions configure publishing behavior.

type SubscribeOptions

type SubscribeOptions struct {
	Protocol Protocol
	QoS      *byte
	AutoAck  *bool
}

SubscribeOptions configure subscription behavior.

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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