Documentation
¶
Index ¶
- Variables
- type Client
- func (c *Client) AMQP() *amqp.Client
- func (c *Client) Close(ctx context.Context) error
- func (c *Client) Connect(ctx context.Context) error
- func (c *Client) IsConnected() bool
- func (c *Client) MQTT() *mqtt.Client
- func (c *Client) Publish(ctx context.Context, topic string, payload []byte, opts ...Option) error
- func (c *Client) PublishToQueue(ctx context.Context, queue string, payload []byte, opts ...Option) error
- func (c *Client) Subscribe(ctx context.Context, topic string, handler MessageHandler, opts ...Option) error
- func (c *Client) SubscribeToQueue(ctx context.Context, queue, group string, handler MessageHandler, ...) error
- func (c *Client) Unsubscribe(ctx context.Context, topic string, opts ...Option) error
- func (c *Client) UnsubscribeFromQueue(ctx context.Context, queue string, opts ...Option) error
- type Config
- type Message
- type MessageHandler
- type Option
- func WithAutoAck(autoAck bool) Option
- func WithExchange(exchange string) Option
- func WithImmediate(immediate bool) Option
- func WithMandatory(mandatory bool) Option
- func WithProperties(props map[string]string) Option
- func WithProtocol(protocol Protocol) Option
- func WithQoS(qos byte) Option
- func WithRetain(retain bool) Option
- func WithRoutingKey(key string) Option
- type Protocol
- type PublishOptions
- type SubscribeOptions
Constants ¶
This section is empty.
Variables ¶
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 (*Client) IsConnected ¶
IsConnected reports whether the client is connected.
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 ¶
Unsubscribe removes a topic subscription.
type Config ¶
type Config struct {
MQTT *mqttclient.Options
AMQP *amqpclient.Options
DefaultProtocol Protocol
}
Config configures the unified client transports.
func (*Config) SetAMQP ¶
func (c *Config) SetAMQP(opts *amqpclient.Options) *Config
SetAMQP sets AMQP transport options.
func (*Config) SetDefaultProtocol ¶
SetDefaultProtocol sets the protocol used when operation options do not specify one.
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.
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 ¶
WithAutoAck controls AMQP auto-ack for subscriptions (publish/subscribe only). It has no effect for MQTT subscriptions.
func WithExchange ¶
WithExchange sets the AMQP exchange (publish only).
func WithImmediate ¶
WithImmediate sets the AMQP immediate flag (publish only).
func WithMandatory ¶
WithMandatory sets the AMQP mandatory flag (publish only).
func WithProperties ¶
WithProperties sets unified properties.
func WithProtocol ¶
WithProtocol sets protocol routing for publish/subscribe operations.
func WithRetain ¶
WithRetain sets the MQTT retain flag (publish only).
func WithRoutingKey ¶
WithRoutingKey sets the AMQP routing key (publish only).
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 ¶
SubscribeOptions configure subscription behavior.