Documentation
¶
Overview ¶
Package managementhttp transports queue management contracts over a bounded authenticated HTTP boundary.
Index ¶
- Constants
- Variables
- func NewHandler(config HandlerConfig) (http.Handler, error)
- type Client
- func (c *Client) Execute(ctx context.Context, request management.Command) (management.CommandResult, error)
- func (c *Client) Inspect(ctx context.Context, request management.InspectRequest) (management.JobRecord, error)
- func (c *Client) ListDeadLetters(ctx context.Context, request management.PageRequest) (management.RecordPage, error)
- func (c *Client) ListFailures(ctx context.Context, request management.PageRequest) (management.RecordPage, error)
- func (c *Client) ListQueues(ctx context.Context, request management.StatusPageRequest) (management.QueueStatusPage, error)
- func (c *Client) ListWorkers(ctx context.Context, request management.StatusPageRequest) (management.WorkerStatusPage, error)
- type ClientConfig
- type Endpoint
- type EndpointResolver
- type EndpointResolverFunc
- type FleetClient
- func (client *FleetClient) Execute(ctx context.Context, command management.Command) (management.CommandResult, error)
- func (client *FleetClient) Inspect(ctx context.Context, request management.InspectRequest) (management.JobRecord, error)
- func (client *FleetClient) ListDeadLetters(ctx context.Context, request management.PageRequest) (management.RecordPage, error)
- func (client *FleetClient) ListFailures(ctx context.Context, request management.PageRequest) (management.RecordPage, error)
- func (client *FleetClient) ListQueues(ctx context.Context, request management.StatusPageRequest) (management.QueueStatusPage, error)
- func (client *FleetClient) ListWorkers(ctx context.Context, request management.StatusPageRequest) (management.WorkerStatusPage, error)
- type FleetClientConfig
- type HandlerConfig
Examples ¶
Constants ¶
const MaxFleetEndpoints = 100
MaxFleetEndpoints bounds one resolved management fleet and all operation fan-out.
Variables ¶
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 = errors.New("managementhttp: fleet target unavailable") ErrFleetUnavailable = errors.New("managementhttp: fleet unavailable") )
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 ¶
func (c *Client) ListQueues( ctx context.Context, request management.StatusPageRequest, ) (management.QueueStatusPage, error)
ListQueues returns one remote bounded queue-status page.
func (*Client) ListWorkers ¶
func (c *Client) ListWorkers( ctx context.Context, request management.StatusPageRequest, ) (management.WorkerStatusPage, error)
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 EndpointResolver ¶
EndpointResolver returns the current bounded worker-management fleet.
type EndpointResolverFunc ¶
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.