Documentation
¶
Overview ¶
Package workqueue contains an interface for a simple key workqueue abstraction.
Metrics ¶
The GCS implementation exports the following Prometheus metrics:
- workqueue_in_progress_keys: The number of keys currently being processed
- workqueue_queued_keys: The number of keys currently in the backlog
- workqueue_notbefore_keys: The number of keys waiting on a 'not before' time
- workqueue_max_attempts: The maximum number of attempts for any queued or in-progress task
- workqueue_task_max_attempts: The maximum number of attempts for a given task above 20
- workqueue_process_latency_seconds: The duration taken to process a key
- workqueue_wait_latency_seconds: The duration the key waited to start
- workqueue_added_keys: The total number of queue requests
- workqueue_deduped_keys: The total number of keys that were deduped
- workqueue_attempts_at_completion: The number of attempts for successfully completed tasks
- workqueue_dead_lettered_keys: The number of keys currently in the dead letter queue
- workqueue_time_to_completion_seconds: The time from first queue to final outcome (success or dead-letter). The metric captures the full lifecycle duration including all retry attempts and backoff delays.
All metrics include service_name and revision_name labels. Additional labels vary by metric.
Index ¶
- Constants
- Variables
- func DeadLetterError(err error, reason string) error
- func DrainInterceptor(drain context.Context) grpc.UnaryServerInterceptor
- func GetRequeueDelay(err error) (time.Duration, bool)
- func GetRequeueOptions(err error) (delay time.Duration, floor bool, ok bool)
- func HasInfrastructureMarker(err error) bool
- func InfrastructureCauses(err error) []error
- func InfrastructureError(err error, causes ...error) error
- func IsInfrastructureError(err error) bool
- func NonRetriableError(err error, reason string) error
- func QueueKeys(keys ...QueueKey) error
- func RegisterWorkqueueServiceServer(s grpc.ServiceRegistrar, srv WorkqueueServiceServer)
- func RequeueAfter(delay time.Duration) error
- func RequeueAfterWithJitter(delay, jitter time.Duration) error
- func RequeueNotBefore(delay time.Duration) error
- type Client
- type DeadLetteredKey
- type GetKeyStateRequest
- type InProgressKey
- type Interface
- type Key
- type KeyState
- func (*KeyState) Descriptor() ([]byte, []int)deprecated
- func (x *KeyState) GetAttempts() int32
- func (x *KeyState) GetKey() string
- func (x *KeyState) GetNotBeforeTime() int64
- func (x *KeyState) GetPriority() int64
- func (x *KeyState) GetQueuedTime() int64
- func (x *KeyState) GetStatus() KeyState_Status
- func (*KeyState) ProtoMessage()
- func (x *KeyState) ProtoReflect() protoreflect.Message
- func (x *KeyState) Reset()
- func (x *KeyState) String() string
- type KeyState_Status
- func (KeyState_Status) Descriptor() protoreflect.EnumDescriptor
- func (x KeyState_Status) Enum() *KeyState_Status
- func (KeyState_Status) EnumDescriptor() ([]byte, []int)deprecated
- func (x KeyState_Status) Number() protoreflect.EnumNumber
- func (x KeyState_Status) String() string
- func (KeyState_Status) Type() protoreflect.EnumType
- type NoRetryDetails
- type ObservedInProgressKey
- type Options
- type OwnedInProgressKey
- type ProcessRequest
- func (*ProcessRequest) Descriptor() ([]byte, []int)deprecated
- func (x *ProcessRequest) GetDelaySeconds() int64
- func (x *ProcessRequest) GetKey() string
- func (x *ProcessRequest) GetPriority() int64
- func (x *ProcessRequest) LogAttrs() []any
- func (*ProcessRequest) ProtoMessage()
- func (x *ProcessRequest) ProtoReflect() protoreflect.Message
- func (x *ProcessRequest) Reset()
- func (x *ProcessRequest) String() string
- type ProcessResponse
- func (*ProcessResponse) Descriptor() ([]byte, []int)deprecated
- func (x *ProcessResponse) GetQueueKeys() []*QueueKeyRequest
- func (x *ProcessResponse) GetRequeueAfterSeconds() int64
- func (x *ProcessResponse) GetRequeueFloor() bool
- func (*ProcessResponse) ProtoMessage()
- func (x *ProcessResponse) ProtoReflect() protoreflect.Message
- func (x *ProcessResponse) Reset()
- func (x *ProcessResponse) String() string
- type QueueKey
- type QueueKeyRequest
- func (*QueueKeyRequest) Descriptor() ([]byte, []int)deprecated
- func (x *QueueKeyRequest) GetDelaySeconds() int64
- func (x *QueueKeyRequest) GetKey() string
- func (x *QueueKeyRequest) GetPriority() int64
- func (*QueueKeyRequest) ProtoMessage()
- func (x *QueueKeyRequest) ProtoReflect() protoreflect.Message
- func (x *QueueKeyRequest) Reset()
- func (x *QueueKeyRequest) String() string
- type QueuedKey
- type UnimplementedWorkqueueServiceServer
- type UnsafeWorkqueueServiceServer
- type WorkqueueServiceClient
- type WorkqueueServiceServer
Examples ¶
Constants ¶
const ( WorkqueueService_Process_FullMethodName = "/chainguard.workqueue.WorkqueueService/Process" WorkqueueService_GetKeyState_FullMethodName = "/chainguard.workqueue.WorkqueueService/GetKeyState" )
Variables ¶
var ( // DrainBackstop bounds how long DrainInterceptor waits, once draining // begins, for an in-flight handler to unwind cooperatively before // abandoning it and answering on its behalf. It must comfortably fit // inside the platform's SIGTERM grace (10s on Cloud Run) together with // the response flush. DrainBackstop = 3 * time.Second // DrainRequeueDelay is the minimum requeue delay the interceptor // answers with while draining; DrainRequeueJitter adds a random extra // delay in [0, jitter) so a retiring instance's in-flight keys do not // all come back at once. DrainRequeueDelay = 15 * time.Second DrainRequeueJitter = 30 * time.Second )
Note that these are variables, so that they can be modified by tests and made flags in binary entrypoints.
var ( // BackoffPeriod is the first step of the dispatcher's failure-retry // backoff. Every retriable dispatch failure requeues on a jittered // doubling curve starting here: the fast first step keeps races and // transient blips cheap, and the widening keeps persistent failures // (infrastructure storms, deterministic errors awaiting a fix) from // burning the dead-letter budget in minutes. Failures still count // against that budget regardless of spacing. The storage // implementations also use this as the unit of their legacy linear // backoff on bare Requeue (e.g. orphan recovery). BackoffPeriod = 30 * time.Second // MaximumBackoffPeriod caps the failure-retry backoff. MaximumBackoffPeriod = 10 * time.Minute // MaximumRequeueFloor caps the delay of a floored requeue (RequeueNotBefore). // A floor cannot be undercut by events or resync, so an unbounded floor would // starve a key indefinitely; the dispatcher clamps floored requeue delays to // this value. It is a variable so binaries can adjust it (e.g. to a resync // period) and tests can shrink it. MaximumRequeueFloor = time.Hour )
Note that these are variables, so that they can be modified by tests and made flags in binary entrypoints.
var ( KeyState_Status_name = map[int32]string{ 0: "UNKNOWN", 1: "QUEUED", 2: "IN_PROGRESS", 3: "DEAD_LETTER", } KeyState_Status_value = map[string]int32{ "UNKNOWN": 0, "QUEUED": 1, "IN_PROGRESS": 2, "DEAD_LETTER": 3, } )
Enum value maps for KeyState_Status.
var File_workqueue_proto protoreflect.FileDescriptor
var WorkqueueService_ServiceDesc = grpc.ServiceDesc{ ServiceName: "chainguard.workqueue.WorkqueueService", HandlerType: (*WorkqueueServiceServer)(nil), Methods: []grpc.MethodDesc{ { MethodName: "Process", Handler: _WorkqueueService_Process_Handler, }, { MethodName: "GetKeyState", Handler: _WorkqueueService_GetKeyState_Handler, }, }, Streams: []grpc.StreamDesc{}, Metadata: "workqueue.proto", }
WorkqueueService_ServiceDesc is the grpc.ServiceDesc for WorkqueueService service. It's only intended for direct use with grpc.RegisterService, and not to be introspected or modified (even as a copy)
Functions ¶
func DeadLetterError ¶ added in v0.10.29
DeadLetterError marks the error as non-retriable AND directs the dispatcher to move the key to the dead-letter queue immediately, instead of completing (dropping) it the way a plain NonRetriableError does. Use it for permanent refusals that need an operator disposition: the dead-letter object is the durable, enumerable signal (Enumerate, drain audits, requeue tooling) where a dropped key leaves only a log line.
The error also carries the same NoRetryDetails a NonRetriableError carries, deliberately: a dispatcher predating the marker degrades to the drop path (today's behavior), never to a retry loop.
Example ¶
ExampleDeadLetterError demonstrates a permanent refusal that must surface durably: the dispatcher moves the key to the dead-letter queue immediately (where Enumerate and drain audits see it) instead of completing it with only a logged reason, the way a plain NonRetriableError is dropped.
package main
import (
"errors"
"fmt"
"chainguard.dev/driftlessaf/workqueue"
)
func main() {
err := workqueue.DeadLetterError(
errors.New("line records several coordinates; an operator must adjudicate"),
"permanent",
)
// The dispatcher checks the dead-letter marker BEFORE the plain
// non-retriable details, because a DeadLetterError carries both: the
// NoRetryDetails keep an older dispatcher from retrying it.
if d := workqueue.GetDeadLetterDetails(err); d != nil {
fmt.Printf("dead-letter reason: %s\n", d.GetMessage())
}
if d := workqueue.GetNonRetriableDetails(err); d != nil {
fmt.Println("also reads as non-retriable")
}
// A plain NonRetriableError carries no dead-letter marker: it drops.
dropped := workqueue.NonRetriableError(errors.New("stale key"), "retired spelling")
fmt.Printf("plain non-retriable dead-letters: %v\n", workqueue.GetDeadLetterDetails(dropped) != nil)
}
Output: dead-letter reason: permanent also reads as non-retriable plain non-retriable dead-letters: false
func DrainInterceptor ¶ added in v0.9.29
func DrainInterceptor(drain context.Context) grpc.UnaryServerInterceptor
DrainInterceptor returns a unary server interceptor that converts orderly shutdown into requeues instead of severed callbacks.
gRPC request contexts are created from the transport stream and deliberately do not derive from any application context, so a top-level cancellation (e.g. signal.NotifyContext observing SIGTERM) is invisible to in-flight reconciles: the process keeps working until the platform's grace expires and SIGKILL severs every connection, which the dispatcher records as a failed attempt against each in-flight key. This interceptor bridges the two cancellation domains per request:
- a Process call arriving while drain is already canceled is answered immediately with a requeue, without starting work;
- an in-flight Process call has its request context canceled when drain fires, and its cooperative unwind (or, past DrainBackstop, the interceptor answering on its behalf) is translated into a requeue response rather than an error;
- a handler that completes successfully during the drain window keeps its real response.
Because the requeue is delivered as RequeueAfterSeconds, the dispatcher resets the key's attempt count: announced shutdowns stop consuming dead-letter budget entirely, while unannounced deaths keep their existing (attempt-consuming) semantics.
Wire it with the recovery interceptor after (inside) it, so handler panics remain the recovery interceptor's to translate:
grpc.ChainUnaryInterceptor(
...,
workqueue.DrainInterceptor(ctx), // ctx canceled on SIGTERM
recovery.UnaryServerInterceptor(),
)
Methods other than WorkqueueService/Process pass through untouched.
func GetRequeueDelay ¶
GetRequeueDelay extracts the requeue delay from an error if it's a requeue error. Returns the delay and true if the error is a requeue error, or 0 and false otherwise. It is a convenience wrapper over GetRequeueOptions for callers that don't need the floor flag.
func GetRequeueOptions ¶ added in v0.7.2
GetRequeueOptions extracts the requeue delay and whether its NotBefore should be treated as a floor from an error, if it is a requeue error (RequeueAfter or RequeueNotBefore). floor is true only for RequeueNotBefore; ok is false for non-requeue errors.
func HasInfrastructureMarker ¶ added in v0.10.3
HasInfrastructureMarker reports whether err carries the marker that InfrastructureError put on it. It is narrower than IsInfrastructureError, which also reports true for every error carrying codes.Unavailable: HasInfrastructureMarker matches the marker alone and never the status-code class. Use it where only a deliberate classification at the call site counts, for example where a scorer exempts a provider outage but a bare codes.Unavailable from the run's own work must still count against it. The marker is matched by type, never by message text.
Example ¶
ExampleHasInfrastructureMarker demonstrates the narrower predicate: it matches only an error a call site marked with InfrastructureError, and never the codes.Unavailable class that IsInfrastructureError also accepts. A caller that must exempt a marked provider failure, but still count a bare codes.Unavailable, uses this one.
package main
import (
"errors"
"fmt"
"time"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"chainguard.dev/driftlessaf/workqueue"
)
func main() {
marked := workqueue.InfrastructureError(workqueue.RequeueAfter(5*time.Minute), errors.New("connection reset by peer"))
fmt.Println(workqueue.HasInfrastructureMarker(marked))
// An unmarked codes.Unavailable carries no marker.
unavailable := status.Error(codes.Unavailable, "connection termination")
fmt.Println(workqueue.HasInfrastructureMarker(unavailable))
fmt.Println(workqueue.IsInfrastructureError(unavailable))
}
Output: true false true
func InfrastructureCauses ¶ added in v0.10.3
InfrastructureCauses returns the causes InfrastructureError recorded on err, or nil when err carries no marker or when the marker carries no cause. The marker's message is the wrapped error's message, so a caller that reports the failure in text (a log line, an eval score reasoning) needs the causes to name the provider failure behind it. Nil causes are dropped, because a nil cause names nothing. The marker is matched by type, never by message text.
Example ¶
ExampleInfrastructureCauses demonstrates reading the provider causes back off a marked error. The marker's message is the wrapped requeue's message, so a caller that reports the failure in text needs the causes to name the provider failure behind it.
package main
import (
"errors"
"fmt"
"time"
"chainguard.dev/driftlessaf/workqueue"
)
func main() {
cause := errors.New("dial tcp 10.0.0.1:443: connect: connection refused")
err := workqueue.InfrastructureError(workqueue.RequeueAfter(5*time.Minute), cause)
for _, c := range workqueue.InfrastructureCauses(err) {
fmt.Println(c)
}
// An unmarked error carries no causes.
fmt.Println(workqueue.InfrastructureCauses(errors.New("reconcile failed")))
}
Output: dial tcp 10.0.0.1:443: connect: connection refused []
func InfrastructureError ¶ added in v0.10.3
InfrastructureError marks err as an infrastructure failure, so that IsInfrastructureError reports it as one whatever gRPC code err carries. Use InfrastructureError when only the call site can recognise a provider failure, for example when it wraps a requeue error returned for a transport fault:
InfrastructureError(RequeueAfter(5*time.Minute), cause)
The marker does not change the wrapped error's message and does not hide the wrapped error. The wrapped error and every cause stay reachable through errors.Is and errors.As, so GetRequeueOptions still finds the requeue delay and callers still match the original provider cause. Match the marker with IsInfrastructureError, never on message text. InfrastructureError returns nil when err is nil.
Example ¶
ExampleInfrastructureError demonstrates marking a requeue error as an infrastructure failure when only the call site can recognise the provider fault. The marker is transparent: the message, the requeue delay, and the original cause all stay reachable.
package main
import (
"errors"
"fmt"
"time"
"chainguard.dev/driftlessaf/workqueue"
)
func main() {
cause := errors.New("dial tcp 10.0.0.1:443: connect: connection refused")
err := workqueue.InfrastructureError(workqueue.RequeueAfter(5*time.Minute), cause)
fmt.Println(err)
fmt.Println(workqueue.IsInfrastructureError(err))
delay, _ := workqueue.GetRequeueDelay(err)
fmt.Println(delay)
fmt.Println(errors.Is(err, cause))
}
Output: requeue after 5m0s (floor=false) true 5m0s true
func IsInfrastructureError ¶ added in v0.9.23
IsInfrastructureError reports whether err is an infrastructure failure and not an application verdict. An error is one of two forms: InfrastructureError marked it, or it carries codes.Unavailable. gRPC transports set codes.Unavailable when the receiver was not reachable or the connection broke mid-call (for example, the receiving instance was killed), and reconcilers propagate codes.Unavailable when a downstream dependency is unavailable. The classification does not change scheduling, because every retriable failure requeues on the same widening backoff curve. The classification is for observability (dispatch error events and logs), where the split between infrastructure churn and application failures makes the failure-rate dashboards readable.
Example ¶
ExampleIsInfrastructureError demonstrates how dispatch errors are classified: transport-level failures (the receiver died mid-call, no healthy backend, a dependency reporting itself unavailable) are infrastructure errors. The classification is observability-only — it separates infrastructure churn from application failures on dispatch error events — while scheduling retries every failure on the same widening backoff curve.
package main
import (
"errors"
"fmt"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
"chainguard.dev/driftlessaf/workqueue"
)
func main() {
// The gRPC transport synthesizes codes.Unavailable when the receiving
// instance is killed mid-dispatch.
infra := status.Error(codes.Unavailable, "connection termination")
fmt.Println(workqueue.IsInfrastructureError(infra))
// An ordinary application failure is not infrastructure.
app := errors.New("reconcile failed")
fmt.Println(workqueue.IsInfrastructureError(app))
}
Output: true false
func NonRetriableError ¶
NoRetryDetails marks the error as non-retriable with the given reason. If this error is returned to the dispatcher, it will not requeue the key.
func QueueKeys ¶
QueueKeys returns a sentinel error indicating additional keys should be queued. This is used by reconcilers to signal that dependent work should be enqueued. The current key will be completed after all queued keys are successfully added. If the current key is included in the list, it will be requeued (enters "dual state").
func RegisterWorkqueueServiceServer ¶
func RegisterWorkqueueServiceServer(s grpc.ServiceRegistrar, srv WorkqueueServiceServer)
func RequeueAfter ¶
RequeueAfter returns an error that indicates the work item should be requeued after the specified delay. The resulting NotBefore is undercuttable: a fresh enqueue of the same key (e.g. from a new event) before the delay elapses will pull the key forward and process it immediately. Use this for retry/backoff where promptly reacting to new events is desirable.
Example ¶
ExampleRequeueAfter demonstrates how to use RequeueAfter in a callback.
package main
import (
"context"
"fmt"
"time"
"chainguard.dev/driftlessaf/workqueue"
)
func main() {
callback := func(_ context.Context, _ string, _ workqueue.Options) error {
// Do some work...
// Request requeue with a 30-second delay
return workqueue.RequeueAfter(30 * time.Second)
}
// This would be used in a dispatcher
err := callback(context.Background(), "example-key", workqueue.Options{})
delay, ok := workqueue.GetRequeueDelay(err)
fmt.Printf("Requeue requested: %v, delay: %v\n", ok, delay)
}
Output: Requeue requested: true, delay: 30s
func RequeueAfterWithJitter ¶ added in v0.7.51
RequeueAfterWithJitter is like RequeueAfter, but adds a random extra delay in [0, jitter) so keys that failed together don't all come back at once.
Example ¶
ExampleRequeueAfterWithJitter demonstrates spreading out retries of keys that failed together.
package main
import (
"fmt"
"time"
"chainguard.dev/driftlessaf/workqueue"
)
func main() {
err := workqueue.RequeueAfterWithJitter(10*time.Second, 50*time.Second)
delay, ok := workqueue.GetRequeueDelay(err)
fmt.Printf("Requeue requested: %v, delay in [10s, 60s): %v\n", ok, delay >= 10*time.Second && delay < 60*time.Second)
}
Output: Requeue requested: true, delay in [10s, 60s): true
func RequeueNotBefore ¶ added in v0.7.2
RequeueNotBefore is like RequeueAfter, but the resulting NotBefore is a floor: subsequent enqueues of the same key that arrive before the delay elapses are coalesced onto it rather than pulling the key forward. This debounces a burst of events for a key the reconciler has decided to revisit at a fixed cadence (e.g. polling for CI to settle) into roughly one reconcile per delay, while still observing the latest state when the key fires.
Types ¶
type Client ¶
type Client interface {
WorkqueueServiceClient
Close() error
}
func NewWorkqueueClient ¶
type DeadLetteredKey ¶
type DeadLetteredKey interface {
Key
// GetFailedTime returns the time when the key was dead-lettered.
GetFailedTime() time.Time
// GetAttempts returns the number of attempts before the key was dead-lettered.
GetAttempts() int
}
DeadLetteredKey is a key that has been moved to the dead-letter queue after exceeding the maximum retry attempts.
type GetKeyStateRequest ¶
type GetKeyStateRequest struct {
// The key to retrieve information for
Key string `protobuf:"bytes,1,opt,name=key,proto3" json:"key,omitempty"`
// contains filtered or unexported fields
}
func (*GetKeyStateRequest) Descriptor
deprecated
func (*GetKeyStateRequest) Descriptor() ([]byte, []int)
Deprecated: Use GetKeyStateRequest.ProtoReflect.Descriptor instead.
func (*GetKeyStateRequest) GetKey ¶
func (x *GetKeyStateRequest) GetKey() string
func (*GetKeyStateRequest) ProtoMessage ¶
func (*GetKeyStateRequest) ProtoMessage()
func (*GetKeyStateRequest) ProtoReflect ¶
func (x *GetKeyStateRequest) ProtoReflect() protoreflect.Message
func (*GetKeyStateRequest) Reset ¶
func (x *GetKeyStateRequest) Reset()
func (*GetKeyStateRequest) String ¶
func (x *GetKeyStateRequest) String() string
type InProgressKey ¶
type InProgressKey interface {
Key
// Requeue returns this key to the queue.
Requeue(context.Context) error
// RequeueWithOptions returns this key to the queue with custom options.
RequeueWithOptions(context.Context, Options) error
}
InProgressKey is a shared interface that all in-progress key types must implement.
type Interface ¶
type Interface interface {
// Identity returns the identity recorded as the owner of keys started by
// this queue, or an empty string if no identity was configured.
Identity() string
// Queue adds an item to the workqueue.
Queue(ctx context.Context, key string, opts Options) error
// Enumerate returns:
// - a list of all of the in-progress keys,
// - a list of the next "N" keys in the queue (according to its configured ordering),
// - a list of all dead-lettered keys, or
// - an error if the workqueue is unable to enumerate the keys.
Enumerate(ctx context.Context) ([]ObservedInProgressKey, []QueuedKey, []DeadLetteredKey, error)
// Get retrieves the current state and metadata for a specific key.
Get(ctx context.Context, key string) (*KeyState, error)
}
Interface is the interface that workqueue implementations must implement.
type Key ¶
type Key interface {
// Name is the name of the key.
Name() string
// Priority is the priority of the key.
Priority() int64
}
Key is a shared interface that all key types must implement.
type KeyState ¶
type KeyState struct {
// The key name
Key string `protobuf:"bytes,1,opt,name=key,proto3" json:"key,omitempty"`
Status KeyState_Status `protobuf:"varint,2,opt,name=status,proto3,enum=chainguard.workqueue.KeyState_Status" json:"status,omitempty"`
// Priority of the key (if queued or in progress)
Priority int64 `protobuf:"varint,3,opt,name=priority,proto3" json:"priority,omitempty"`
// Number of attempts made (if in progress)
Attempts int32 `protobuf:"varint,4,opt,name=attempts,proto3" json:"attempts,omitempty"`
// Time when the key was first queued
QueuedTime int64 `protobuf:"varint,5,opt,name=queued_time,json=queuedTime,proto3" json:"queued_time,omitempty"`
// Time when the key should be processed (NotBefore)
NotBeforeTime int64 `protobuf:"varint,6,opt,name=not_before_time,json=notBeforeTime,proto3" json:"not_before_time,omitempty"`
// contains filtered or unexported fields
}
func (*KeyState) Descriptor
deprecated
func (*KeyState) GetAttempts ¶
func (*KeyState) GetNotBeforeTime ¶
func (*KeyState) GetPriority ¶
func (*KeyState) GetQueuedTime ¶
func (*KeyState) GetStatus ¶
func (x *KeyState) GetStatus() KeyState_Status
func (*KeyState) ProtoMessage ¶
func (*KeyState) ProtoMessage()
func (*KeyState) ProtoReflect ¶
func (x *KeyState) ProtoReflect() protoreflect.Message
type KeyState_Status ¶
type KeyState_Status int32
Current status of the key
const ( KeyState_UNKNOWN KeyState_Status = 0 KeyState_QUEUED KeyState_Status = 1 KeyState_IN_PROGRESS KeyState_Status = 2 KeyState_DEAD_LETTER KeyState_Status = 3 )
func (KeyState_Status) Descriptor ¶
func (KeyState_Status) Descriptor() protoreflect.EnumDescriptor
func (KeyState_Status) Enum ¶
func (x KeyState_Status) Enum() *KeyState_Status
func (KeyState_Status) EnumDescriptor
deprecated
func (KeyState_Status) EnumDescriptor() ([]byte, []int)
Deprecated: Use KeyState_Status.Descriptor instead.
func (KeyState_Status) Number ¶
func (x KeyState_Status) Number() protoreflect.EnumNumber
func (KeyState_Status) String ¶
func (x KeyState_Status) String() string
func (KeyState_Status) Type ¶
func (KeyState_Status) Type() protoreflect.EnumType
type NoRetryDetails ¶
type NoRetryDetails struct {
Message string `protobuf:"bytes,1,opt,name=message,proto3" json:"message,omitempty"` // A message describing why the key should not be retried.
// contains filtered or unexported fields
}
NoRetryDetails is a marker message that indicates that the key should not be retried.
func GetDeadLetterDetails ¶ added in v0.10.29
func GetDeadLetterDetails(err error) *NoRetryDetails
GetDeadLetterDetails extracts the NoRetryDetails from an error carrying the immediate-dead-letter marker (DeadLetterError). It returns nil for a plain NonRetriableError and for errors carrying neither detail, so a dispatcher checks it BEFORE GetNonRetriableDetails (which also matches a DeadLetterError).
func GetNonRetriableDetails ¶
func GetNonRetriableDetails(err error) *NoRetryDetails
GetNonRetriableDetails extracts the NoRetryDetails from the error if it exists. If the error is nil or does not contain NoRetryDetails, it returns nil.
func (*NoRetryDetails) Descriptor
deprecated
func (*NoRetryDetails) Descriptor() ([]byte, []int)
Deprecated: Use NoRetryDetails.ProtoReflect.Descriptor instead.
func (*NoRetryDetails) GetMessage ¶
func (x *NoRetryDetails) GetMessage() string
func (*NoRetryDetails) ProtoMessage ¶
func (*NoRetryDetails) ProtoMessage()
func (*NoRetryDetails) ProtoReflect ¶
func (x *NoRetryDetails) ProtoReflect() protoreflect.Message
func (*NoRetryDetails) Reset ¶
func (x *NoRetryDetails) Reset()
func (*NoRetryDetails) String ¶
func (x *NoRetryDetails) String() string
type ObservedInProgressKey ¶
type ObservedInProgressKey interface {
InProgressKey
// Owner returns the identity recorded when the key was claimed, or an empty
// string if no owner was recorded.
Owner() string
// IsOrphaned checks whether the key has been orphaned by it's owner.
IsOrphaned() bool
}
ObservedInProgressKey is a key that we have observed to be in progress, but that we are not the owner of.
type Options ¶
type Options struct {
// Priority is the priority of the key.
// Higher values are processed first.
Priority int64
// NotBefore is the earliest time that the key should be processed.
// When deduplicating, the earliest time is used unless NotBeforeFloor is set
// on the queued entry (see NotBeforeFloor).
NotBefore time.Time
// Delay is an optional duration to wait before processing the key.
// This is used when requeueing with a custom delay.
Delay time.Duration
// BackoffDelay is an optional duration to wait before reprocessing a key
// on a failure retry. Unlike Delay it does NOT reset the attempt count,
// so the dispatcher's attempts >= maxRetry dead-letter cutoff keeps
// firing, and it is applied regardless of priority. When zero (the
// default) the requeue path is unchanged: callers that never set it get
// the existing linear/priority backoff behavior. A caller that wants
// custom failure-retry backoff (e.g. decorrelated exponential jitter)
// computes the delay itself and passes it here, leaving the attempt
// counter intact. Takes precedence over Delay and the priority backoff
// when set.
BackoffDelay time.Duration
// NotBeforeFloor, when true, marks NotBefore as a floor that a non-floor
// enqueue of the same key cannot dedup to an earlier time. Merge rules for a
// queued entry:
// - a floored enqueue over a non-floored entry replaces its NotBefore
// outright (in either direction) and marks it floored — the non-floored
// time was undercuttable anyway;
// - a floored enqueue over a floored entry takes the later NotBefore;
// - a non-floor enqueue over a floored entry leaves it untouched (neither
// earlier nor later);
// - between two non-floor enqueues the earliest NotBefore wins (the default).
// Set via RequeueNotBefore.
NotBeforeFloor bool
}
Options is a set of options that can be passed when queuing a key.
type OwnedInProgressKey ¶
type OwnedInProgressKey interface {
InProgressKey
// Complete marks the key as successfully completed, and removes it from
// the in-progress key set.
Complete(context.Context) error
// Deadletter permanently removes this key from the queue, indicating it has
// failed after exceeding the maximum retry attempts.
Deadletter(context.Context) error
// GetAttempts returns the current attempt count for the key.
GetAttempts() int
// Context is the context of the process heartbeating the key.
Context() context.Context
}
OwnedInProgressKey is an in-progress key where we have initiated the work, and own until it completes either successfully (Complete), or unsuccessfully (Requeue or Fail).
type ProcessRequest ¶
type ProcessRequest struct {
// The key of the work item
Key string `protobuf:"bytes,1,opt,name=key,proto3" json:"key,omitempty"`
// The (optional) priority of the work item, where higher numbers are processed first.
Priority int64 `protobuf:"varint,2,opt,name=priority,proto3" json:"priority,omitempty"`
// The (optional) delay in second to wait before processing the work item.
DelaySeconds int64 `protobuf:"varint,3,opt,name=delay_seconds,json=delaySeconds,proto3" json:"delay_seconds,omitempty"`
// contains filtered or unexported fields
}
func (*ProcessRequest) Descriptor
deprecated
func (*ProcessRequest) Descriptor() ([]byte, []int)
Deprecated: Use ProcessRequest.ProtoReflect.Descriptor instead.
func (*ProcessRequest) GetDelaySeconds ¶
func (x *ProcessRequest) GetDelaySeconds() int64
func (*ProcessRequest) GetKey ¶
func (x *ProcessRequest) GetKey() string
func (*ProcessRequest) GetPriority ¶
func (x *ProcessRequest) GetPriority() int64
func (*ProcessRequest) LogAttrs ¶
func (x *ProcessRequest) LogAttrs() []any
LogAttrs returns a slice of attributes for logging purposes.
func (*ProcessRequest) ProtoMessage ¶
func (*ProcessRequest) ProtoMessage()
func (*ProcessRequest) ProtoReflect ¶
func (x *ProcessRequest) ProtoReflect() protoreflect.Message
func (*ProcessRequest) Reset ¶
func (x *ProcessRequest) Reset()
func (*ProcessRequest) String ¶
func (x *ProcessRequest) String() string
type ProcessResponse ¶
type ProcessResponse struct {
// Optional: If set, indicates the work item should be requeued after the specified delay.
// If not set or 0, the default behavior applies (success if no error, or exponential backoff on error).
// This field is only honored when the RPC returns successfully (no error).
RequeueAfterSeconds int64 `protobuf:"varint,1,opt,name=requeue_after_seconds,json=requeueAfterSeconds,proto3" json:"requeue_after_seconds,omitempty"`
// Optional: When true, requeue_after_seconds is treated as a floor that a
// subsequent enqueue cannot pull earlier (RequeueNotBefore semantics), rather
// than an undercuttable delay (RequeueAfter). Only meaningful when
// requeue_after_seconds > 0.
RequeueFloor bool `protobuf:"varint,3,opt,name=requeue_floor,json=requeueFloor,proto3" json:"requeue_floor,omitempty"`
// Optional: Additional keys to queue as a result of processing this item.
// These keys are queued BEFORE the current key is completed, ensuring dependent
// work is guaranteed to be queued before we mark the triggering work as done.
// If any queue operation fails, the current key will be retried.
QueueKeys []*QueueKeyRequest `protobuf:"bytes,2,rep,name=queue_keys,json=queueKeys,proto3" json:"queue_keys,omitempty"`
// contains filtered or unexported fields
}
func (*ProcessResponse) Descriptor
deprecated
func (*ProcessResponse) Descriptor() ([]byte, []int)
Deprecated: Use ProcessResponse.ProtoReflect.Descriptor instead.
func (*ProcessResponse) GetQueueKeys ¶
func (x *ProcessResponse) GetQueueKeys() []*QueueKeyRequest
func (*ProcessResponse) GetRequeueAfterSeconds ¶
func (x *ProcessResponse) GetRequeueAfterSeconds() int64
func (*ProcessResponse) GetRequeueFloor ¶ added in v0.7.2
func (x *ProcessResponse) GetRequeueFloor() bool
func (*ProcessResponse) ProtoMessage ¶
func (*ProcessResponse) ProtoMessage()
func (*ProcessResponse) ProtoReflect ¶
func (x *ProcessResponse) ProtoReflect() protoreflect.Message
func (*ProcessResponse) Reset ¶
func (x *ProcessResponse) Reset()
func (*ProcessResponse) String ¶
func (x *ProcessResponse) String() string
type QueueKey ¶
QueueKey represents a key to be queued with optional priority and delay.
func GetQueueKeys ¶
GetQueueKeys extracts queued keys from an error, if present. Returns nil if the error doesn't contain queue keys.
type QueueKeyRequest ¶
type QueueKeyRequest struct {
// The key to queue.
Key string `protobuf:"bytes,1,opt,name=key,proto3" json:"key,omitempty"`
// Optional priority (default: 0). Higher values are processed first.
Priority int64 `protobuf:"varint,2,opt,name=priority,proto3" json:"priority,omitempty"`
// Optional delay in seconds before the key becomes eligible for processing.
DelaySeconds int64 `protobuf:"varint,3,opt,name=delay_seconds,json=delaySeconds,proto3" json:"delay_seconds,omitempty"`
// contains filtered or unexported fields
}
QueueKeyRequest represents a key to be queued with optional priority and delay.
func (*QueueKeyRequest) Descriptor
deprecated
func (*QueueKeyRequest) Descriptor() ([]byte, []int)
Deprecated: Use QueueKeyRequest.ProtoReflect.Descriptor instead.
func (*QueueKeyRequest) GetDelaySeconds ¶
func (x *QueueKeyRequest) GetDelaySeconds() int64
func (*QueueKeyRequest) GetKey ¶
func (x *QueueKeyRequest) GetKey() string
func (*QueueKeyRequest) GetPriority ¶
func (x *QueueKeyRequest) GetPriority() int64
func (*QueueKeyRequest) ProtoMessage ¶
func (*QueueKeyRequest) ProtoMessage()
func (*QueueKeyRequest) ProtoReflect ¶
func (x *QueueKeyRequest) ProtoReflect() protoreflect.Message
func (*QueueKeyRequest) Reset ¶
func (x *QueueKeyRequest) Reset()
func (*QueueKeyRequest) String ¶
func (x *QueueKeyRequest) String() string
type QueuedKey ¶
type QueuedKey interface {
Key
// Start initiates processing of the key, returning an OwnedInProgressKey
// on success and an error on failure.
Start(context.Context) (OwnedInProgressKey, error)
}
QueuedKey is a key that is in the queue, waiting to be processed.
type UnimplementedWorkqueueServiceServer ¶
type UnimplementedWorkqueueServiceServer struct{}
UnimplementedWorkqueueServiceServer must be embedded to have forward compatible implementations.
NOTE: this should be embedded by value instead of pointer to avoid a nil pointer dereference when methods are called.
func (UnimplementedWorkqueueServiceServer) GetKeyState ¶
func (UnimplementedWorkqueueServiceServer) GetKeyState(context.Context, *GetKeyStateRequest) (*KeyState, error)
func (UnimplementedWorkqueueServiceServer) Process ¶
func (UnimplementedWorkqueueServiceServer) Process(context.Context, *ProcessRequest) (*ProcessResponse, error)
type UnsafeWorkqueueServiceServer ¶
type UnsafeWorkqueueServiceServer interface {
// contains filtered or unexported methods
}
UnsafeWorkqueueServiceServer may be embedded to opt out of forward compatibility for this service. Use of this interface is not recommended, as added methods to WorkqueueServiceServer will result in compilation errors.
type WorkqueueServiceClient ¶
type WorkqueueServiceClient interface {
Process(ctx context.Context, in *ProcessRequest, opts ...grpc.CallOption) (*ProcessResponse, error)
GetKeyState(ctx context.Context, in *GetKeyStateRequest, opts ...grpc.CallOption) (*KeyState, error)
}
WorkqueueServiceClient is the client API for WorkqueueService service.
For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
func NewWorkqueueServiceClient ¶
func NewWorkqueueServiceClient(cc grpc.ClientConnInterface) WorkqueueServiceClient
type WorkqueueServiceServer ¶
type WorkqueueServiceServer interface {
Process(context.Context, *ProcessRequest) (*ProcessResponse, error)
GetKeyState(context.Context, *GetKeyStateRequest) (*KeyState, error)
// contains filtered or unexported methods
}
WorkqueueServiceServer is the server API for WorkqueueService service. All implementations must embed UnimplementedWorkqueueServiceServer for forward compatibility.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package conformance provides a suite of conformance tests for workqueue implementations.
|
Package conformance provides a suite of conformance tests for workqueue implementations. |
|
Package dispatcher provides a workqueue dispatcher that dequeues keys and invokes a callback for each one.
|
Package dispatcher provides a workqueue dispatcher that dequeues keys and invokes a callback for each one. |
|
Package gcs provides a Google Cloud Storage-backed workqueue implementation.
|
Package gcs provides a Google Cloud Storage-backed workqueue implementation. |
|
Package hyperqueue provides a sharded WorkqueueService implementation that consistently distributes keys across N backend workqueue services using consistent hashing.
|
Package hyperqueue provides a sharded WorkqueueService implementation that consistently distributes keys across N backend workqueue services using consistent hashing. |
|
Package inmem provides an in-memory workqueue implementation intended for testing.
|
Package inmem provides an in-memory workqueue implementation intended for testing. |
|
Package serve runs a workqueue.WorkqueueServiceServer on the standard duplex gRPC server: trace-parent restoration, OpenTelemetry stats, metrics and recovery interceptors, a gRPC health service, and the metrics listener.
|
Package serve runs a workqueue.WorkqueueServiceServer on the standard duplex gRPC server: trace-parent restoration, OpenTelemetry stats, metrics and recovery interceptors, a gRPC health service, and the metrics listener. |