managementhttp

package
v1.0.1 Latest Latest
Warning

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

Go to latest
Published: Aug 28, 2026 License: MIT Imports: 16 Imported by: 0

Documentation

Overview

Package managementhttp transports queue management contracts over a bounded authenticated HTTP boundary.

Index

Examples

Constants

View Source
const MaxFleetEndpoints = 100

MaxFleetEndpoints bounds one resolved management fleet and all operation fan-out.

Variables

View Source
var (
	// ErrInvalidFleetConfiguration reports an unusable dynamic fleet client.
	ErrInvalidFleetConfiguration = errors.New("managementhttp: invalid fleet configuration")
	// ErrInvalidFleetEndpoints reports malformed, duplicate, or unbounded resolved endpoints.
	ErrInvalidFleetEndpoints = errors.New("managementhttp: invalid fleet endpoints")
	// ErrFleetTargetUnavailable reports a worker target absent from the current fleet snapshot.
	ErrFleetTargetUnavailable = errors.New("managementhttp: fleet target unavailable")
	// ErrFleetUnavailable reports a fleet operation with no trustworthy remote result.
	ErrFleetUnavailable = errors.New("managementhttp: fleet unavailable")
)
View Source
var (
	ErrInvalidConfiguration = errors.New("management HTTP: invalid configuration")
	ErrInvalidRequest       = errors.New("management HTTP: invalid request")
	ErrInvalidResponse      = errors.New("management HTTP: invalid response")
	ErrResponseTooLarge     = errors.New("management HTTP: response too large")
	ErrRemoteFailure        = errors.New("management HTTP: remote failure")
)

Functions

func NewHandler

func NewHandler(config HandlerConfig) (http.Handler, error)

NewHandler creates a worker-side management HTTP handler.

Types

type Client

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

Client implements remote management contracts.

func NewClient

func NewClient(config ClientConfig) (*Client, error)

NewClient creates a bounded management HTTP client.

func (*Client) Execute

func (c *Client) Execute(
	ctx context.Context,
	request management.Command,
) (management.CommandResult, error)

Execute sends one validated management command to the configured worker endpoint and returns its validated acknowledgement.

Example
package main

import (
	"context"
	"fmt"
	"net/http/httptest"
	"time"

	"github.com/faustbrian/go-queue/management"
	"github.com/faustbrian/go-queue/managementhttp"
)

func main() {
	controller := exampleController{}
	handler, err := managementhttp.NewHandler(managementhttp.HandlerConfig{
		Token: "replace-with-a-secret", Controller: controller,
	})
	if err != nil {
		panic(err)
	}
	server := httptest.NewTLSServer(handler)
	defer server.Close()
	client, err := managementhttp.NewClient(managementhttp.ClientConfig{
		BaseURL: server.URL, Token: "replace-with-a-secret",
		HTTPClient: server.Client(),
	})
	if err != nil {
		panic(err)
	}
	requestedAt := time.Date(2026, 7, 16, 10, 0, 0, 0, time.UTC)
	result, err := client.Execute(context.Background(), management.Command{
		ID: "command-1", IdempotencyKey: "deployment-1", Actor: "operator-1",
		Reason:      "drain before deployment",
		Protocol:    management.ProtocolVersion{Major: 1},
		Action:      management.CommandDrain,
		Target:      management.Target{Kind: management.TargetWorker, Name: "worker-1"},
		RequestedAt: requestedAt, Deadline: requestedAt.Add(time.Minute),
	})
	if err != nil {
		panic(err)
	}
	fmt.Println(result.Status)
}

type exampleController struct{}

func (exampleController) Execute(
	_ context.Context,
	command management.Command,
) (management.CommandResult, error) {
	return management.CommandResult{
		CommandID: command.ID, IdempotencyKey: command.IdempotencyKey,
		WorkerID: "worker-1", Protocol: command.Protocol,
		Status:      management.CommandAcknowledged,
		CompletedAt: command.RequestedAt.Add(time.Second),
	}, nil
}
Output:
acknowledged

func (*Client) Inspect

func (c *Client) Inspect(
	ctx context.Context,
	request management.InspectRequest,
) (management.JobRecord, error)

Inspect returns one remote record at explicit payload visibility.

func (*Client) ListDeadLetters

func (c *Client) ListDeadLetters(
	ctx context.Context,
	request management.PageRequest,
) (management.RecordPage, error)

ListDeadLetters returns one remote bounded dead-letter page.

func (*Client) ListFailures

func (c *Client) ListFailures(
	ctx context.Context,
	request management.PageRequest,
) (management.RecordPage, error)

ListFailures returns one remote bounded failure page.

Example
package main

import (
	"context"
	"fmt"
	"net/http/httptest"
	"time"

	"github.com/faustbrian/go-queue/management"
	"github.com/faustbrian/go-queue/managementhttp"
)

func main() {
	reader := exampleRecordReader{}
	handler, err := managementhttp.NewHandler(managementhttp.HandlerConfig{
		Token: "replace-with-a-secret", Records: reader,
	})
	if err != nil {
		panic(err)
	}
	server := httptest.NewTLSServer(handler)
	defer server.Close()
	client, err := managementhttp.NewClient(managementhttp.ClientConfig{
		BaseURL: server.URL, Token: "replace-with-a-secret",
		HTTPClient: server.Client(),
	})
	if err != nil {
		panic(err)
	}
	page, err := client.ListFailures(context.Background(), management.PageRequest{
		Limit: 25, Sort: management.SortOccurredAt,
		Direction: management.SortDescending,
	})
	if err != nil {
		panic(err)
	}
	fmt.Println(page.Items[0].ID, page.Items[0].Payload.Visibility == management.PayloadHidden)
}

type exampleRecordReader struct{}

func (exampleRecordReader) ListFailures(
	context.Context,
	management.PageRequest,
) (management.RecordPage, error) {
	return management.RecordPage{Items: []management.JobRecord{{
		Kind: management.RecordFailure, ID: "failure-1", Backend: "example",
		Queue: "critical", OccurredAt: time.Unix(2, 0).UTC(), Attempts: 1,
		FailureCode: "handler_failed", Payload: management.Payload{Size: 7},
	}}}, nil
}

func (exampleRecordReader) ListDeadLetters(
	context.Context,
	management.PageRequest,
) (management.RecordPage, error) {
	return management.RecordPage{}, nil
}

func (exampleRecordReader) Inspect(
	context.Context,
	management.InspectRequest,
) (management.JobRecord, error) {
	return management.JobRecord{}, nil
}
Output:
failure-1 true

func (*Client) ListQueues

ListQueues returns one remote bounded queue-status page.

func (*Client) ListWorkers

ListWorkers returns one remote bounded worker-status page.

type ClientConfig

type ClientConfig struct {
	BaseURL          string
	Token            string
	HTTPClient       *http.Client
	MaxResponseBytes int64
}

ClientConfig provides the worker endpoint, credential, and response bound.

type Endpoint

type Endpoint struct {
	ID      string
	BaseURL string
}

Endpoint is one stable worker-management transport target.

type EndpointResolver

type EndpointResolver interface {
	ResolveEndpoints(context.Context) ([]Endpoint, error)
}

EndpointResolver returns the current bounded worker-management fleet.

type EndpointResolverFunc

type EndpointResolverFunc func(context.Context) ([]Endpoint, error)

EndpointResolverFunc adapts a function into an EndpointResolver.

func (EndpointResolverFunc) ResolveEndpoints

func (resolve EndpointResolverFunc) ResolveEndpoints(ctx context.Context) ([]Endpoint, error)

ResolveEndpoints implements EndpointResolver.

type FleetClient

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

FleetClient aggregates bounded worker status and routes control to current endpoints.

func NewFleetClient

func NewFleetClient(config FleetClientConfig) (*FleetClient, error)

NewFleetClient creates a dynamic multi-endpoint management client.

func (*FleetClient) Execute

func (client *FleetClient) Execute(
	ctx context.Context,
	command management.Command,
) (management.CommandResult, error)

func (*FleetClient) Inspect

func (client *FleetClient) Inspect(
	ctx context.Context,
	request management.InspectRequest,
) (management.JobRecord, error)

func (*FleetClient) ListDeadLetters

func (client *FleetClient) ListDeadLetters(
	ctx context.Context,
	request management.PageRequest,
) (management.RecordPage, error)

func (*FleetClient) ListFailures

func (client *FleetClient) ListFailures(
	ctx context.Context,
	request management.PageRequest,
) (management.RecordPage, error)

func (*FleetClient) ListQueues

func (client *FleetClient) ListQueues(
	ctx context.Context,
	request management.StatusPageRequest,
) (management.QueueStatusPage, error)

func (*FleetClient) ListWorkers

func (client *FleetClient) ListWorkers(
	ctx context.Context,
	request management.StatusPageRequest,
) (management.WorkerStatusPage, error)

type FleetClientConfig

type FleetClientConfig struct {
	Resolver         EndpointResolver
	Token            string
	HTTPClient       *http.Client
	MaxResponseBytes int64
}

FleetClientConfig configures dynamic authenticated worker-management access.

type HandlerConfig

type HandlerConfig struct {
	Token      string
	Status     management.StatusReader
	Controller management.Controller
	Records    management.RecordReader
}

HandlerConfig provides the worker-side management services and shared transport credential.

Jump to

Keyboard shortcuts

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