Versions in this module Expand all Collapse all v0 v0.5.2 Jan 31, 2023 Changes in this version + var ErrClosed = errors.New("queue has been closed") + var ErrDropExceedsMaxPacketSize = errors.New("maximum packet size exceeded") + var ErrDropExpired = errors.New("the message is expired") + var ErrDropExpiredInflight = errors.New("the inflight message is expired") + var ErrDropQueueFull = errors.New("the message queue is full") + func ElemExpiry(now time.Time, elem *Elem) bool + type Elem struct + At time.Time + Expiry time.Time + func (e *Elem) Decode(b []byte) (err error) + func (e *Elem) Encode() []byte + type InitOptions struct + CleanStart bool + Notifier Notifier + ReadBytesLimit uint32 + Version packets.Version + type InternalError struct + Err error + func (i *InternalError) Error() string + type MessageWithID interface + ID func() packets.PacketID + SetID func(id packets.PacketID) + type MockMessageWithID struct + func NewMockMessageWithID(ctrl *gomock.Controller) *MockMessageWithID + func (m *MockMessageWithID) EXPECT() *MockMessageWithIDMockRecorder + func (m *MockMessageWithID) ID() packets.PacketID + func (m *MockMessageWithID) SetID(id packets.PacketID) + type MockMessageWithIDMockRecorder struct + func (mr *MockMessageWithIDMockRecorder) ID() *gomock.Call + func (mr *MockMessageWithIDMockRecorder) SetID(id interface{}) *gomock.Call + type MockNotifier struct + func NewMockNotifier(ctrl *gomock.Controller) *MockNotifier + func (m *MockNotifier) EXPECT() *MockNotifierMockRecorder + func (m *MockNotifier) NotifyDropped(elem *Elem, err error) + func (m *MockNotifier) NotifyInflightAdded(delta int) + func (m *MockNotifier) NotifyMsgQueueAdded(delta int) + type MockNotifierMockRecorder struct + func (mr *MockNotifierMockRecorder) NotifyDropped(elem, err interface{}) *gomock.Call + func (mr *MockNotifierMockRecorder) NotifyInflightAdded(delta interface{}) *gomock.Call + func (mr *MockNotifierMockRecorder) NotifyMsgQueueAdded(delta interface{}) *gomock.Call + type MockStore struct + func NewMockStore(ctrl *gomock.Controller) *MockStore + func (m *MockStore) Add(elem *Elem) error + func (m *MockStore) Clean() error + func (m *MockStore) Close() error + func (m *MockStore) EXPECT() *MockStoreMockRecorder + func (m *MockStore) Init(opts *InitOptions) error + func (m *MockStore) Read(pids []packets.PacketID) ([]*Elem, error) + func (m *MockStore) ReadInflight(maxSize uint) ([]*Elem, error) + func (m *MockStore) Remove(pid packets.PacketID) error + func (m *MockStore) Replace(elem *Elem) (bool, error) + type MockStoreMockRecorder struct + func (mr *MockStoreMockRecorder) Add(elem interface{}) *gomock.Call + func (mr *MockStoreMockRecorder) Clean() *gomock.Call + func (mr *MockStoreMockRecorder) Close() *gomock.Call + func (mr *MockStoreMockRecorder) Init(opts interface{}) *gomock.Call + func (mr *MockStoreMockRecorder) Read(pids interface{}) *gomock.Call + func (mr *MockStoreMockRecorder) ReadInflight(maxSize interface{}) *gomock.Call + func (mr *MockStoreMockRecorder) Remove(pid interface{}) *gomock.Call + func (mr *MockStoreMockRecorder) Replace(elem interface{}) *gomock.Call + type Notifier interface + NotifyDropped func(elem *Elem, err error) + NotifyInflightAdded func(delta int) + NotifyMsgQueueAdded func(delta int) + type Publish struct + func (p *Publish) Decode(b *bytes.Buffer) (err error) + func (p *Publish) Encode(b *bytes.Buffer) + func (p *Publish) ID() packets.PacketID + func (p *Publish) SetID(id packets.PacketID) + type Pubrel struct + PacketID packets.PacketID + func (p *Pubrel) Decode(b *bytes.Buffer) (err error) + func (p *Pubrel) Encode(b *bytes.Buffer) + func (p *Pubrel) ID() packets.PacketID + func (p *Pubrel) SetID(id packets.PacketID) + type Store interface + Add func(elem *Elem) error + Clean func() error + Close func() error + Init func(opts *InitOptions) error + Read func(pids []packets.PacketID) ([]*Elem, error) + ReadInflight func(maxSize uint) (elems []*Elem, err error) + Remove func(pid packets.PacketID) error + Replace func(elem *Elem) (replaced bool, err error)