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 ¶
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?
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
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
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
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