queuedproto

package module
v0.0.13 Latest Latest
Warning

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

Go to latest
Published: Jul 21, 2026 License: MIT Imports: 8 Imported by: 0

Documentation

Overview

Package queuedproto - поддержка протокола queued.

Кроме структур протокола самого queued, поддерживаются форматы перловых очередей. Perl addons (используются во всех современных проектах на perl) поддерживаются, perl extensions (используются только в одном крупном старом проекте) - нет.

Типы/константы названы так же, как в исходниках queued. В перле используется другая терминология.

Index

Constants

View Source
const (
	CmdAddItem      = Cmd(20)
	CmdGetActive    = Cmd(21)
	CmdDeleteItems  = Cmd(22)
	CmdUpdateItems  = Cmd(24)
	CmdGetItems     = Cmd(28)
	CmdFullUpdate   = Cmd(30)
	CmdAddData      = Cmd(34)
	CmdGetQueueStat = Cmd(36)
)

Коды команд queued

View Source
const (
	// errors
	RcOK                  = uint8(0)
	RcLogicError          = uint8(1)
	RcWrongItem           = uint8(2)
	RcMemAllocationFailed = uint8(3)
	RcUnknownQueueType    = uint8(4)
	RcReqAlreadyProcessed = uint8(8)
	RcItemLocked          = uint8(9)
	RcWrongRequest        = uint8(10)
	RcBadRequestLength    = uint8(11)
	RcWrongVersion        = uint8(12)

	// warnings
	RcNoActiveItems           = uint8(5)
	RcInsufficientActiveCount = uint8(6)
	RcWrongID                 = uint8(7)
)

Коды ответов queued

View Source
const (
	FlagEnableExtensions = uint8(0x80)
	FlagEncodeInUTF8     = uint8(0x40)
	FlagEnableAddons     = uint8(0x20)

	VersionMask = ^(FlagEnableExtensions | FlagEncodeInUTF8 | FlagEnableAddons)
)

Флаги расширений, указываются в поле Version

View Source
const (
	TimestampAbsolute      = uint8(0) // указанный UnixTime - абсолютный (секунды с 1970-01-01T00:00:00Z)
	TimestampRelativePlus  = uint8(1) // в Unixtime относительное время, которое необходимо прибавить к текущему серверному времени
	TimestampRelativeMinus = uint8(2) // в Unixtime относительное время, которое необходимо отнять от текущего серверного времени
)

Типы времени для команды CmdUpdateItems

View Source
const FlagIgnoreWrongItem = uint32(1)

FlagIgnoreWrongItem - флаг команды CmdAddItem

View Source
const MaxIProtoLen = 1 << 20

MaxIProtoLen caps how many elements a generated UnmarshalIProto will pre-allocate for a length-prefixed slice or map, bounding memory amplification when decoding untrusted input: a short message can encode a huge element count, and make([]T, n) / make(map, n) reserves n*sizeof(elem) before any element is read. Decoding a slice or map longer than this fails with iproto.ErrOverflow. Declare a constant of this name in the package to override the default.

View Source
const SkipTimeUpdate = 0xFFFFFFFF

SkipTimeUpdate - для запроса ReqFullUpdate - magic value для UnixTime, в случае передачи этого значения время не изменяется.

Variables

View Source
var (
	ErrQueued              = errors.New("queued error")
	ErrLogic               = fmt.Errorf("%w: logic error", ErrQueued)
	ErrWrongItem           = fmt.Errorf("%w: wrong item", ErrQueued)
	ErrMemAllocationFailed = fmt.Errorf("%w: queued memory allocation failure", ErrQueued)
	ErrUnknownQueueType    = fmt.Errorf("%w: unknown queue type", ErrQueued)
	ErrReqAlreadyProcessed = fmt.Errorf("%w: request has been already processed (duplicate ReqID)", ErrQueued)
	ErrItemLocked          = fmt.Errorf("%w: item locked", ErrQueued)
	ErrWrongRequest        = fmt.Errorf("%w: wrong request", ErrQueued)
	ErrBadRequestLength    = fmt.Errorf("%w: bad request length", ErrQueued)
	ErrWrongVersion        = fmt.Errorf("%w: wrong version", ErrQueued)
)

Протокольные ошибки queued

Functions

func ErrorByRetCode

func ErrorByRetCode(rc uint8) error

ErrorByRetCode - возвращает ошибку по коду, для неизвестных ошибок в тексте сообщается код. Для RcOK и всех ворнингов возвращается nil, если ворнинги нужны - их нужно обрабатывать вручную.

func IsSoftError

func IsSoftError(err error) bool

IsSoftError сообщает, является ли ошибка логической - т.е. не связанной с низкоуровневыми проблемами (разрывы соединения, таймауты, блокировки, сбои сервера).

Все неизвестные протокольные ошибки считаются "жёсткими".

Запросы с логическими ошибками повторять не надо, с низкоуровневыми - можно/нужно.

Можно передавать также обёрнутые ошибки.

Types

type Cmd added in v0.0.12

type Cmd uint32

Cmd - команда queued.

func (Cmd) String added in v0.0.12

func (i Cmd) String() string

type EventID

type EventID struct {
	Shard uint `iproto:"-"`
	ID    uint64
}

EventID - составной ID с номером шарда. Нужен, т.к. протокол (и сам queued) не поддерживает шардинг. Связка Shard:ID уникальна, тогда как сам по себе ID в списке событий может быть неуникальным, если события пришли из разных шардов. Такое возможно только при автогенерации ID событий на стороне queued.

В структурах запросах этот тип используется только для слайсов (чтобы не перелопачивать слайсы целиком из параметров метода, которому номер шарда должен как-то поступать). Скалярные айдишники надо будет скопировать (только ID, передающийся по прококолу).

func (EventID) MarshalIProto

func (recvEventID EventID) MarshalIProto(buf []byte) ([]byte, error)

func (EventID) MarshalText

func (eid EventID) MarshalText() ([]byte, error)

MarshalText - для упрощения логгирования списков событий при помощи zerolog.Event.Interface (не нужно переваливать слайс айдишников событий в слайс строк/стрингеров).

func (EventID) String

func (eid EventID) String() string

String возвращает EventID в перловом текстовом формате

func (*EventID) UnmarshalIProto

func (recv_EventID *EventID) UnmarshalIProto(buf []byte) ([]byte, error)

type ItemList

type ItemList struct {
	Items []QueueItem `iproto:"u32"`
}

ItemList - список событий

func (ItemList) MarshalIProto

func (recvItemList ItemList) MarshalIProto(buf []byte) ([]byte, error)

func (*ItemList) UnmarshalIProto

func (recv_ItemList *ItemList) UnmarshalIProto(buf []byte) ([]byte, error)

type PerlAddonData

type PerlAddonData struct {
	ID   uint16 // addons.CreationTimeID, addons.RetryID, addons.Retry2ID или addons.ProducerID
	Data []byte `iproto:"u8"`
}

PerlAddonData - данные аддона. Константы ID и структуры формата данных оперделены в пакете addons

type PerlAddons

type PerlAddons struct {
	List []PerlAddonData `iproto:"u8"`
}

PerlAddons - список аддонов

func (PerlAddons) MarshalIProto

func (recvPerlAddons PerlAddons) MarshalIProto(buf []byte) ([]byte, error)

func (*PerlAddons) UnmarshalIProto

func (recv_PerlAddons *PerlAddons) UnmarshalIProto(buf []byte) ([]byte, error)

type PerlData

type PerlData struct {
	Version    uint8          // версия (содержит флаги FlagEnableExtensions, FlagEncodeInUTF8, FlagEnableAddons)
	Addons     PerlAddons     // список данных аддонов  (если установлен FlagEnableAddons - иначе поле не передаётся!)
	Extensions PerlExtensions // данные расширения (если установлен FlagEnableExtensions - иначе поле не передаётся!)
	Data       []byte         // данные в формате perl pack
}

PerlData - событие очереди в перловом формате.

Аддоны декодируются, только если установлен флаг FlagEnableAddons, иначе Addons.List устанавливается в nil. При кодировании, если len(Addons.List) != 0, они кодируются, и автоматически устанавливается флаг FlagEnableAddons.

Extensions - аналогично, однако их внутренности (пока) не обрабатываются (просто декодируются, как слайс байтов).

При кодировании флаг FlagEncodeInUTF8 устанавливается автоматически всегда.

func (PerlData) MarshalIProto

func (pd PerlData) MarshalIProto(buf []byte) ([]byte, error)

MarshalIProto кодирует данные в перловом формате, передавая необязательные поля, только если они заданы. В этом случае устанавливаются соответствующие им флаги (FlagEnableAddons - для Addons, FlagEnableExtensions - для Extensions).

func (*PerlData) UnmarshalIProto

func (pd *PerlData) UnmarshalIProto(buf []byte) ([]byte, error)

UnmarshalIProto декодирует данные в перловом формате, обрабатывая необязательные поля, только если установлены соответствующие им флаги (FlagEnableAddons - для Addons, FlagEnableExtensions - для Extensions).

Если флаг для необязательного поля не установлен, оно сбрасывается в дефлотное значение, так сделано, чтобы не прилетели данные от предыдущей записи в случае переиспользования объекта.

type PerlExtensions

type PerlExtensions struct {
	Flags uint32
	Data  []byte `iproto:"u32"`
}

PerlExtensions - данные расширений (пока не обрабатываются, де/кодируются как слайс байтов)

func (PerlExtensions) MarshalIProto

func (recvPerlExtensions PerlExtensions) MarshalIProto(buf []byte) ([]byte, error)

func (*PerlExtensions) UnmarshalIProto

func (recv_PerlExtensions *PerlExtensions) UnmarshalIProto(buf []byte) ([]byte, error)

type QueueItem

type QueueItem struct {
	EventID EventID
	Data    []byte `iproto:"u16"` // Данные (до 4кБ)
}

QueueItem - событие в очереди

type ReqAddData

type ReqAddData struct {
	StorageType uint16 // номер очереди
	ID          uint64 // ID изменяемого события
	Data        []byte `iproto:"u16"` // дописываемые данные (в сумме с существующими должно быть до 4кБ)
}

ReqAddData - формат tuple запроса на добавление данных в конец события

func (ReqAddData) Cmd

func (ReqAddData) Cmd() Cmd

Cmd - команда queued: CmdAddData (34)

func (ReqAddData) MarshalIProto

func (recvReqAddData ReqAddData) MarshalIProto(buf []byte) ([]byte, error)

func (ReqAddData) QueueID added in v0.0.12

func (req ReqAddData) QueueID() uint16

func (*ReqAddData) UnmarshalIProto

func (recv_ReqAddData *ReqAddData) UnmarshalIProto(buf []byte) ([]byte, error)

type ReqAddItem

type ReqAddItem struct {
	Pid         uint32 // pid клиента
	ReqID       uint32 // номер запроса, в случае таймаута на сервере при перепосылке запроса ReqID 2й попытки должен совпадать с ReqID первой попытки
	StorageType uint16 // номер очереди
	ID          uint64 // ID события, 0 - queued сам назначает ID (автоинкремент)
	UnixTime    uint32 // время активации события
	Data        []byte `iproto:"u16"` // данные события, до 4096 байт
	Flags       uint32 `iproto:"-"`   // FlagIgnoreWrongItem. Поле опциональное (передаётся, только если не 0)
}

ReqAddItem - формат tuple запроса на создание новой записи в очереди.

func (ReqAddItem) Cmd

func (ReqAddItem) Cmd() Cmd

Cmd - команда queued: CmdAddItem (20)

func (ReqAddItem) MarshalIProto

func (req ReqAddItem) MarshalIProto(buf []byte) ([]byte, error)

MarshalIProto кодирует запрос на создание новой записи в очереди.

Кастомный маршалер необходим из-за опциональности поля Flags.

func (ReqAddItem) QueueID added in v0.0.12

func (req ReqAddItem) QueueID() uint16

func (*ReqAddItem) UnmarshalIProto

func (req *ReqAddItem) UnmarshalIProto(buf []byte) ([]byte, error)

UnmarshalIProto декодирует запрос на создание новой записи в очереди.

Метод объявлен чисто для симметрии маршалеров, для проектов клиентов/обработчиков очередей он не нужен. Пригодится для тестов, прокси iproto/queued->grpc, ну и для гошной версии queued :)

type ReqDeleteItems

type ReqDeleteItems struct {
	StorageType uint16   // номер очереди
	IDs         []uint64 `iproto:"u32"`
}

ReqDeleteItems - формат tuple запроса на удаление событий по списку ID

func (ReqDeleteItems) Cmd

func (ReqDeleteItems) Cmd() Cmd

Cmd - команда queued: CmdDeleteItems (22)

func (ReqDeleteItems) MarshalIProto

func (recvReqDeleteItems ReqDeleteItems) MarshalIProto(buf []byte) ([]byte, error)

func (ReqDeleteItems) QueueID added in v0.0.12

func (req ReqDeleteItems) QueueID() uint16

func (*ReqDeleteItems) UnmarshalIProto

func (recv_ReqDeleteItems *ReqDeleteItems) UnmarshalIProto(buf []byte) ([]byte, error)

type ReqFullUpdate

type ReqFullUpdate struct {
	StorageType uint16         // номер очереди
	UnixTime    uint32         // новое время активации обновлённых событий, если [SkipTimeUpdate] (0xFFFFFFFF), время не обновляется
	Items       []UpdQueueItem `iproto:"u32"`
}

ReqFullUpdate - формат tuple запроса на обновление нескольких событий в очереди

func (ReqFullUpdate) Cmd

func (ReqFullUpdate) Cmd() Cmd

Cmd - команда queued: CmdFullUpdate (30)

func (ReqFullUpdate) MarshalIProto

func (recvReqFullUpdate ReqFullUpdate) MarshalIProto(buf []byte) ([]byte, error)

func (ReqFullUpdate) QueueID added in v0.0.12

func (req ReqFullUpdate) QueueID() uint16

func (*ReqFullUpdate) UnmarshalIProto

func (recv_ReqFullUpdate *ReqFullUpdate) UnmarshalIProto(buf []byte) ([]byte, error)

type ReqGetActive

type ReqGetActive struct {
	StorageType uint16 // номер очереди
	Count       uint32 // количество событий в ответе
}

ReqGetActive - формат tuple запроса на получение Count активных событий из очереди

func (ReqGetActive) Cmd

func (ReqGetActive) Cmd() Cmd

Cmd - команда queued: CmdGetActive (21)

func (ReqGetActive) MarshalIProto

func (recvReqGetActive ReqGetActive) MarshalIProto(buf []byte) ([]byte, error)

func (ReqGetActive) QueueID added in v0.0.12

func (req ReqGetActive) QueueID() uint16

func (*ReqGetActive) UnmarshalIProto

func (recv_ReqGetActive *ReqGetActive) UnmarshalIProto(buf []byte) ([]byte, error)

type ReqGetItems

type ReqGetItems struct {
	StorageType uint16   // номер очереди
	IDs         []uint64 `iproto:"u32"`
}

ReqGetItems - формат tuple запроса на получение событий по ID

func (ReqGetItems) Cmd

func (ReqGetItems) Cmd() Cmd

Cmd - команда queued: CmdGetItems (28)

func (ReqGetItems) MarshalIProto

func (recvReqGetItems ReqGetItems) MarshalIProto(buf []byte) ([]byte, error)

func (ReqGetItems) QueueID added in v0.0.12

func (req ReqGetItems) QueueID() uint16

func (*ReqGetItems) UnmarshalIProto

func (recv_ReqGetItems *ReqGetItems) UnmarshalIProto(buf []byte) ([]byte, error)

type ReqQueueStat

type ReqQueueStat struct {
	Version     uint8
	StorageType uint16
}

ReqQueueStat - получение статистики

func (ReqQueueStat) Cmd

func (ReqQueueStat) Cmd() Cmd

Cmd - команда queued: CmdGetQueueStat (36)

func (ReqQueueStat) MarshalIProto

func (recvReqQueueStat ReqQueueStat) MarshalIProto(buf []byte) ([]byte, error)

func (ReqQueueStat) QueueID added in v0.0.12

func (req ReqQueueStat) QueueID() uint16

func (*ReqQueueStat) UnmarshalIProto

func (recv_ReqQueueStat *ReqQueueStat) UnmarshalIProto(buf []byte) ([]byte, error)

type ReqUpdateItems

type ReqUpdateItems struct {
	StorageType   uint16   // номер очереди
	UnixTime      uint32   // новое время активации для всех указанных событий
	IDs           []uint64 `iproto:"u32"`
	TimestampType uint8    // TimestampAbsolute, TimestampRelativePlus или TimestampRelativeMinus
}

ReqUpdateItems - формат tuple запроса на обновление времени срабатывания события

func (ReqUpdateItems) Cmd

func (ReqUpdateItems) Cmd() Cmd

Cmd - команда queued: CmdUpdateItems(24)

func (ReqUpdateItems) MarshalIProto

func (recvReqUpdateItems ReqUpdateItems) MarshalIProto(buf []byte) ([]byte, error)

func (ReqUpdateItems) QueueID added in v0.0.12

func (req ReqUpdateItems) QueueID() uint16

func (*ReqUpdateItems) UnmarshalIProto

func (recv_ReqUpdateItems *ReqUpdateItems) UnmarshalIProto(buf []byte) ([]byte, error)

type Request

type Request interface {
	iproto.Marshaler
	Cmd() Cmd
}

Request - структуры команд протокола должны поддерживать этот интерфейс (возвращать код команды методом Cmd)

type RespQueueStat

type RespQueueStat struct {
	ItemsCount  uint32 // кол-во событий в очереди
	ActiveCount uint32 // кол-во активных событий в очереди
	LockedCount uint32 // кол-во заблокированных событий в очереди
}

RespQueueStat - ответ на команду запроса статистики

func (RespQueueStat) MarshalIProto

func (recvRespQueueStat RespQueueStat) MarshalIProto(buf []byte) ([]byte, error)

func (RespQueueStat) String

func (i RespQueueStat) String() string

String возвращает статистику в формате строки

func (*RespQueueStat) UnmarshalIProto

func (recv_RespQueueStat *RespQueueStat) UnmarshalIProto(buf []byte) ([]byte, error)

type UpdQueueItem

type UpdQueueItem struct {
	EventID EventID
	Data    []byte `iproto:"u16"` // Данные (до 4кБ)
}

UpdQueueItem - аналог структуры QueueItem для запроса ReqFullUpdate. Если Data == nil, данные не обновляются - только время активации. Чтобы передать пустой массив данных, надо указать []byte{}.

func (UpdQueueItem) MarshalIProto

func (uqi UpdQueueItem) MarshalIProto(buf []byte) ([]byte, error)

MarshalIProto кодирует структуру UpdQueueItem. Если Data == nil, передаётся длина 0xFFFF.

func (*UpdQueueItem) UnmarshalIProto

func (uqi *UpdQueueItem) UnmarshalIProto(buf []byte) ([]byte, error)

UnmarshalIProto декодирует структуру UpdQueueItem. Если переданæ длина данных 0xFFFF, поле Data будет равняться nil.

Directories

Path Synopsis
Package addons - поддержка перловых аддонов.
Package addons - поддержка перловых аддонов.

Jump to

Keyboard shortcuts

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