queue

package
v0.10.0 Latest Latest
Warning

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

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

Documentation

Overview

Package queue provides a generic queue resource plus an in-memory adapter for demos and tests. Optional NATS/Kafka wrappers stay thin and injection-based so core ShiftLock does not require broker SDKs.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Backend

type Backend interface {
	Ping(ctx context.Context) error
	Pause(ctx context.Context) error
	Resume(ctx context.Context) error
	Depth(ctx context.Context) (int, error)
}

Backend is the minimal queue control surface.

type Config

type Config struct {
	ID          resource.ResourceID
	DisplayName string
	Backend     Backend
	// MaxDepth soft capacity signal for health (0 = 10000).
	MaxDepth int
}

Config configures a queue resource.

type Kafka

type Kafka struct {
	Backend Backend
}

Kafka is a thin optional wrapper documenting injection of a pause/ping backend.

type Memory

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

Memory is an in-process queue backend for demos/tests.

func NewMemory

func NewMemory(maxMsgs int) *Memory

NewMemory creates a bounded memory queue (default max 1024).

func (*Memory) Consume

func (m *Memory) Consume() (string, bool)

Consume removes and returns the next message.

func (*Memory) Depth

func (m *Memory) Depth(context.Context) (int, error)

func (*Memory) Pause

func (m *Memory) Pause(context.Context) error

func (*Memory) Paused

func (m *Memory) Paused() bool

Paused reports pause state.

func (*Memory) Ping

func (m *Memory) Ping(context.Context) error

func (*Memory) Publish

func (m *Memory) Publish(msg string) error

Publish enqueues a message (rejects when paused or full).

func (*Memory) Resume

func (m *Memory) Resume(context.Context) error

type NATS

type NATS struct {
	Backend Backend
}

NATS is a thin optional wrapper documenting injection of a pause/ping backend. No NATS SDK dependency is pulled into the module.

type Resource

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

Resource implements resource.Resource for queues.

func New

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

New constructs a queue resource.

func (*Resource) Capabilities

func (r *Resource) Capabilities() resource.ResourceCapabilities

func (*Resource) Depth

func (r *Resource) Depth(ctx context.Context) (int, error)

Depth returns approximate backlog.

func (*Resource) Describe

func (r *Resource) Describe() resource.Description

func (*Resource) Health

func (*Resource) ID

func (r *Resource) ID() resource.ResourceID

func (*Resource) Kind

func (r *Resource) Kind() resource.Kind

func (*Resource) Pause

func (r *Resource) Pause(ctx context.Context) error

Pause pauses consumers (application-level).

func (*Resource) Resume

func (r *Resource) Resume(ctx context.Context) error

Resume resumes consumers.

func (*Resource) Snapshot

func (r *Resource) Snapshot(ctx context.Context) (map[string]string, error)

Snapshot is sanitized (no payloads).

Jump to

Keyboard shortcuts

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