gcs

package
v0.10.68 Latest Latest
Warning

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

Go to latest
Published: Sep 9, 2026 License: Apache-2.0 Imports: 21 Imported by: 0

README

GCS Workqueue Implementation

This package implements a Google Cloud Storage (GCS) backed workqueue that provides reliable, persistent task processing with state management.

Bucket Organization

The GCS workqueue uses object prefixes to organize tasks by their state within a single bucket:

Prefixes
  • queued/ - Tasks waiting to be processed
  • in-progress/ - Tasks currently being processed by a worker
  • dead-letter/ - Tasks that have failed after exceeding maximum retry attempts
State Transitions
queued/{key} → in-progress/{key} → [completed] (deleted)
                      ↓
                 dead-letter/{key} (on failure)
                      ↑
                 queued/{key} (on requeue)

Object Metadata

Each object stores metadata to track task state:

  • priority - Zero-padded 8-digit priority for lexicographic ordering (higher = processed first)
  • attempts - Number of processing attempts
  • lease-expiration - When the current lease expires (for in-progress tasks)
  • not-before - Earliest time the task should be processed (RFC3339 format)
  • failed-time - When the task was moved to dead letter queue (RFC3339 format)
  • last-attempted - Unix timestamp of last processing attempt

Key Features

  • Priority-based processing - Higher priority tasks processed first
  • Lease-based ownership - In-progress tasks have renewable leases to prevent multiple workers processing the same task
  • Automatic retry with backoff - Failed tasks automatically requeued with exponential backoff
  • Dead letter handling - Tasks exceeding retry limits moved to dead letter queue
  • Orphan detection - Detects and handles tasks with expired leases
  • Deduplication - Duplicate queue requests update priority/timing instead of creating duplicates

Producer backpressure

QueuedDepth(ctx, client, queueName, limit) counts the keys under queued/, counting no further than limit. A result equal to limit means there are at least that many; any smaller result is exact. An error returns a count of zero, never a partial one — a caller that logged the error and carried on would otherwise read the keys counted before the failure as headroom.

Use it instead of Enumerate for that question. Enumerate lists every object in the bucket, in every state, and sorts as it goes — the dispatcher can afford that because it runs once per dispatch and needs the whole picture, but a bulk producer polling it would pay for the size of its own backlog on every check.

What a call costs. One list request per 1,000 keys counted, rounded up: GCS caps a page at 1,000 however large maxResults is, so a limit of 2,000 is two round trips and a limit of 200,000 is two hundred, sequentially, on every poll. Size a threshold knowing that. Only the object name is selected, so the bytes stay small even when the round trips do not.

It emits workqueue_queued_depth_latency_seconds and workqueue_queued_depth_errors_total, labelled with the queueName passed in — which should be the name the queue was built with, or "" for a queue with none. The pair exists because this read is polled: without it a stalled producer looks the same whether the queue is full and it is throttling correctly, or the depth read is failing and the caller is swallowing it.

Three things it does not do, all deliberate:

  • It counts every queued key, including one whose not-before has not arrived: a delayed key is work the queue still owes, and reading not-before would mean fetching each object's metadata.
  • It counts nothing in progress and nothing dead-lettered — those are bounded by the queue's concurrency and by attention, not by how fast a producer runs.
  • It counts one bucket. A sharded queue (see hyperqueue) gives each shard its own workqueue module and so its own bucket, so a count against one of them sees 1/N of the depth and reads as headroom. A producer in front of a sharded queue has to sum the shards.

A limit below 1 is an error. No count satisfies the contract above, and a threshold misconfigured to zero should be heard about rather than read as either "always full" or "always room" depending on its sign.

Metrics

The implementation exports Prometheus metrics for:

  • Queue sizes (queued, in-progress, dead-lettered)
  • Processing latency and wait times
  • Retry attempts and completion rates
  • Task priorities and attempt distributions

Documentation

Overview

Package gcs provides a Google Cloud Storage-backed workqueue implementation.

Keys are stored as GCS objects under queued/, in-progress/, and dead-letter/ prefixes. The implementation supports priority ordering, not-before delays, and automatic lease refresh to detect orphaned work.

Index

Examples

Constants

This section is empty.

Variables

View Source
var RefreshInterval = 5 * time.Minute

RefreshInterval is the period on which we refresh the lease of owned objects It is surfaced as a global, so that it can be mutated by tests and exposed as a flag by binaries wrapping this library. However, binary authors should use caution to pass consistent values to the key ingress and dispatchers, or they may see unexpected behavior. TODO(mattmoor): What's the right balance here?

View Source
var TrackWorkAttemptMinThreshold = 20

The minimum number of attempts before tracking work attempts. This is to minimize the cardinality of the metric.

Functions

func NewWorkQueue

func NewWorkQueue(client ClientInterface, limit int, opts ...Option) workqueue.Interface

NewWorkQueue creates a new GCS-backed workqueue.

Example

ExampleNewWorkQueue demonstrates constructing a GCS-backed workqueue.

package main

import (
	"fmt"
)

func main() {
	// In production, pass a real *storage.BucketHandle obtained from a
	// cloud.google.com/go/storage client.
	//
	//   client, err := storage.NewClient(ctx)
	//   bucket := client.Bucket("my-workqueue-bucket")
	//   wq := gcs.NewWorkQueue(bucket, 10)
	//
	// The limit parameter controls the maximum number of keys dequeued per
	// Enumerate call.
	fmt.Println("GCS workqueue limit:", 10)
}
Output:
GCS workqueue limit: 10

func QueuedDepth added in v0.10.41

func QueuedDepth(ctx context.Context, client ClientInterface, queueName string, limit int) (int, error)

QueuedDepth counts the keys waiting in the queue, counting no further than limit. A result equal to limit means there are AT LEAST that many; any smaller result is exact. An error returns a count of zero, never a partial one.

It exists for producers that have to decide whether to enqueue more. The question they are really asking is "is the queue deeper than N", not "how deep is it", and the difference matters: a bulk producer polls this, so the cost of the answer has to be a function of N rather than of the size of the backlog.

Enumerate is the wrong call for that. It lists every object in the bucket, in every state, and sorts as it goes — the dispatcher can afford that because it runs once per dispatch and needs the whole picture. A producer holding several million keys back cannot.

queueName labels the metrics this emits, and should be the name the queue was built with (see WithName). Pass "" for a queue that has none.

What a call costs

One list request per 1,000 keys counted, rounded up: GCS caps a page at 1,000 however large maxResults is, so a limit of 2,000 is two round trips and a limit of 200,000 is two hundred, sequentially, on every poll. Size a threshold knowing that. Only the object name is selected, so the bytes are small even when the round trips are not, and the last page may over-fetch — the iterator re-reads the page size on every fetch rather than narrowing it to what is left to count, so the final request can return up to 999 names the count discards.

Three things it does NOT do, all deliberate:

  • It counts every queued key, including one whose not-before time has not arrived. A delayed key is work the queue still owes, which is what backpressure is about; and reading not-before means fetching each object's metadata, which would cost far more than the answer is worth.
  • It counts nothing in progress and nothing dead-lettered. Those are bounded by the queue's concurrency and by attention respectively, not by how fast a producer runs.
  • It counts ONE bucket. A sharded queue (see the hyperqueue package) gives each shard its own workqueue module and so its own bucket, so a count against one of them sees 1/N of the depth and reads as headroom. A producer in front of a sharded queue has to sum the shards.

A limit below 1 is an error rather than a defined answer. There is no count that satisfies the contract above, and a producer whose threshold has been misconfigured to zero should hear about it: the alternative is a bound that reads as either "always full" or "always room" depending on the sign.

Types

type ClientInterface

type ClientInterface interface {
	Object(name string) *storage.ObjectHandle
	Objects(ctx context.Context, q *storage.Query) *storage.ObjectIterator
}

ClientInterface is an interface that abstracts the GCS client.

Objects must forward q, including its attribute selection, to the underlying listing. Enumerate selects exactly the attributes it reads (see keyAttrs). An implementation that substitutes its own query lists the full projection. One that narrows the selection further lists observed keys with a zero Generation and Metageneration, and their requeue then skips orphan recovery.

type Option added in v0.7.1

type Option func(*wq)

Option configures a GCS-backed workqueue created by NewWorkQueue.

func WithIdentity added in v0.10.19

func WithIdentity(identity string) Option

WithIdentity sets the identity recorded on each in-progress object this queue starts, e.g. the region of the dispatcher claiming the key. Because the owner rides on the object, every process enumerating the queue sees who holds what, surfaced via the workqueue_in_progress_keys_by_owner metric. Defaults to "", which records nothing.

func WithName added in v0.7.1

func WithName(name string) Option

WithName sets the queue_name label applied to every Prometheus metric this workqueue emits. Use it to disambiguate multiple workqueues running in the same Cloud Run service (which share K_SERVICE / K_REVISION). Defaults to "".

func WithScheduledWaitWarningThreshold added in v0.10.28

func WithScheduledWaitWarningThreshold(threshold time.Duration) Option

WithScheduledWaitWarningThreshold emits a bounded structured warning when a key is successfully claimed after waiting at least threshold from its scheduled eligibility time. Non-positive values disable the warning.

Example

ExampleWithScheduledWaitWarningThreshold demonstrates enabling the optional structured warning for keys that wait too long after becoming eligible.

package main

import (
	"fmt"
	"time"

	"chainguard.dev/driftlessaf/workqueue/gcs"
)

func main() {
	const threshold = 30 * time.Minute
	_ = gcs.NewWorkQueue(nil, 10, gcs.WithScheduledWaitWarningThreshold(threshold))

	fmt.Println("Scheduled wait warning threshold:", threshold)
}
Output:
Scheduled wait warning threshold: 30m0s

Jump to

Keyboard shortcuts

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