Repository navigation
feat: record what each job attempt consumed, and run the redis store fully on grove kv - #33
Open
juicycleff wants to merge 8 commits into
Open
juicycleff wants to merge 8 commits into
juicycleff wants to merge 8 commits into
Conversation
exec.Usage was populated by every rung and then discarded inside the terminal closure, so nothing accumulated. A resource estimator cannot be built from data that was never kept, and the history it needs cannot be backfilled. Usage pairs the prediction with the measurement on one row: Resources is what the enqueuing process computed, PeakRSS and the rest are what the attempt really used, and InputBytes is the feature the two correlate through. Failed attempts are recorded too -- a job that OOMs is the strongest evidence about how much memory it needed. UsageRecorder is an optional store capability rather than part of the composite Store, following the WakeNotifier precedent. Usage is telemetry, not correctness: a backend without it runs jobs identically and simply keeps no history, which is what lets backends gain the capability one at a time. Recording happens before handleFailure, which increments RetryCount on its way to scheduling the retry, and through a context.WithoutCancel so a timed-out attempt still leaves its measurements behind.
The production backend is the one that actually needs to accumulate history, so it implements UsageRecorder first. The predicted resource set is stored as JSONB rather than scalar columns: unlike the job row nothing queries it, the table is read in bulk by whatever computes estimates, and a blob keeps custom resource keys intact without a column per key. Two indexes, one per reader: (name, recorded_at DESC) for an estimator reading one definition over a window, and recorded_at alone for the retention sweep. Purge selects before deleting because Postgres will not take a LIMIT on DELETE, which is also what lets a sweep over a table this size work in batches instead of locking it whole.
Completes the capability across every backend, so a fleet on any store accumulates the history an estimator needs rather than only the Postgres ones. Each backend expresses the same contract in its own terms. SQLite selects victims before deleting because it is usually built without DELETE ... LIMIT, which also keeps the batching identical to Postgres rather than dialect-dependent. Mongo deletes in one statement when unbounded and selects first when limited, since DeleteMany takes no limit. Redis has no query engine at all, so the ordering and both filters come from sorted sets written on record: one global, scored by time, for the newest-first read and the retention sweep, and one per definition for the read an estimator actually makes. Durations are stored as integer nanoseconds everywhere. Neither BSON nor SQLite has a duration type, and a consistent representation keeps the numbers comparable across backends.
First slice of removing the go-redis coupling, now that grove kv has sorted sets, sets, hashes, and a Store-level publish/subscribe. Where a TxPipeline collapsed several index writes into one round trip, the calls are now sequential and the comment says what a crash between them costs. In every case converted here the answer is a dangling index entry that the next read skips, not a half-written record: the entity is authoritative and the indexes are derived. Two capabilities keep go-redis in the package for now. Lease renewal and reclamation run Lua for their compare-and-set, and event streams use XADD/XRANGE; grove kv has no scripting or stream interface, so those are untouched rather than reimplemented less safely.
The two paths that could not be converted before, now that grove kv has script and stream capabilities. Lease renewal and reclamation keep their Lua verbatim -- the compare-and-set is the correctness core of lease handoff and was never the problem; only the transport was. The scripts are plain strings passed to kv.Eval rather than go-redis Script values, and scriptInt reads the result, treating anything unexpected as a loss because the safe answer to whether this caller won a contested claim is no. Event streams move to kv.XAdd/XRange. The poll batch is now a named constant with the reason attached: a waiter wants the first match, so a larger batch costs latency on every miss and buys nothing on a hit.
No production file in store/redis imports go-redis any more. The queue's sorted sets, the id sets, the artifact link hashes, the lease scripts, and the event streams all go through kv.Store, so the backend is whichever driver the caller opened rather than Redis by construction. Two changes are behavioural rather than mechanical. The dequeue claim issues one ZRem per candidate instead of one pipelined batch. Pipeline was never providing atomicity -- it was a Pipeline, not a TxPipeline -- so a single ZRem returning 1 still means this worker took the job, and the claim semantics are identical. What it costs is a round trip per candidate, bounded by the dequeue limit. The batch entity read keeps its single round trip: kv MGetRaw pipelines its GETs rather than issuing MGET, which is also why it stays correct against a cluster, where one queue's job keys span slots. Everywhere a TxPipeline collapsed several writes, the comment now says what a crash between them costs. In every case it is an index entry pointing at a record that is gone, which reads already skip, because the entity is authoritative and the indexes are derived. New: MissingCapabilities reports at construction which capabilities a driver lacks, so a mis-chosen backend says so at startup instead of failing the first dequeue with ErrNotSupported.
The whole package, tests included, now goes through grove kv. The requeue index check just reads a sorted set, so it uses kv.ZRange instead of building a second Redis client alongside the store. The scan-cost guard needed grove to grow hooks on collection commands first, which it now has. It counts keys touched rather than commands sent: a batch read is one call but N keys of Redis CPU, and MGetRaw reports all of its keys, so a full scan cannot hide behind a single call. Counting at the kv layer also means the guard measures the same thing whatever driver is underneath. Mutation-verified, as the original was: restoring the unbounded scan still fails this test, at 29 keys passing versus a 2000-key backlog.
Every redis store operation now goes through kv.Store, so a namespace hook on the store you pass in has to reach sorted sets, sets, scripts and streams as well as plain gets and sets. Grove kv v1.6.3 applied the hook to some of those and not to others. Hand Dispatch a namespaced store on v1.6.3 and it writes the job under the namespace and the queue entry outside it, and then it can't dequeue the job it just enqueued. v1.7.0 resolves the rewritten key on every operation. The new integration test opens kv with middleware.NewNamespace, drives enqueue, dequeue, lease reclaim, events, workers, leadership, cron and usage, then lists Redis raw and fails on any key outside the namespace. It fails at the dequeue on v1.6.3 and passes on v1.7.0. go-redis drops to an indirect dependency, since nothing in store/redis imports it any more.
This branch has not been deployed
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Every job attempt now leaves a usage row behind: what the enqueuer predicted, what the attempt actually used (wall time, CPU time, peak RSS, disk written), which executor ran it, and whether it completed or failed. Failed attempts are recorded too, since a job that OOMs is the best evidence you'll get of how much memory it needed. All four backends persist it (postgres, sqlite, mongo, redis) and the memory store keeps it for tests. It's an optional store capability,
job.UsageRecorder, so a backend without it runs jobs exactly as before and just keeps no history.The redis store also stops talking to Redis directly. Every operation, including the queue's sorted sets, the lease scripts and the event streams, goes through
kv.Store, and nothing instore/redisimports go-redis.Where this came from
feat/job-usagesat unpushed for seven weeks while main moved 197 commits. 116 of its 125 commits had already landed on main by other routes, so pushing it as it stood would have opened a PR with 78 conflicts. We cherry-picked the 7 commits that hadn't landed onto current main and leftfeat/job-usageitself alone. Two of the nine unique commits are dropped because main already has them: the Trove backend adapter is9e375dbwith identical code, and the bun modules went with the grove migration.The last commit is new. It bumps grove to v1.7.0 and adds a test for it, explained below.
What to look at when you review
The conflicts were with work that landed on main after the branch was cut, so a few resolutions are judgement calls.
worker/runner.go: main added a step that commits sandbox outputs after a successful run. Usage capture happens before that step, so an attempt whose output commit fails still records what it consumed.20260813120000to20261007120000on postgres and sqlite. Main added20260817120000in the meantime. Grove would still apply the older stamp, but the new one keeps file order, version order and rollback order the same. The branch was never released, so no database has the old stamp.keystype thatWithKeyPrefix(b17fb95) introduced. As written on the branch they used the old free functions and would have ignored the tenant prefix. A unit test pins the prefixed shape now.lease.gowas converted by hand. Main's version gainedUpdateLeasedJoband the fencing work after the branch was cut, andUpdateLeasedJobnow runs its compare-and-set throughkv.Evallike renewal and reclaim do.Behaviour changes
These came with the conversion and are documented in the code, but you should know about them before merging.
TxPipelineare now sequential. A crash between them leaves an index entry pointing at a record that's gone, which every read already skips because the entity is authoritative.ZRemper candidate where they used one pipelined batch. Claim semantics are the same (the pipeline was never a transaction). It costs a round trip per candidate, bounded by the dequeue limit.EVALon every call, where go-redis usedEVALSHAwith a fallback, so the script body goes over the wire each time.StartWakeListenerno longer waits for Redis to confirm the subscription, andstopcancels the subscriber without waiting for its goroutine to exit. A wake sent in the first instant after startup can be missed, and one can arrive just after stop. Neither affects correctness, since polling is what actually picks up jobs andPool.Wakeis a non-blocking send, but it's a weaker guarantee than before. The wake tests passed 20 out of 20 runs under-count=20.Why grove v1.7.0
The extension takes its
*kv.Storefrom the DI container, so a host app can hand Dispatch a store with a namespace hook on it. Grove kv v1.6.3 applied that hook to plain gets and sets but not to sorted sets, sets, scripts or streams, and with everything on kv that splits one store across two keyspaces.store/redis/namespace_test.goopens kv withmiddleware.NewNamespace, drives enqueue, dequeue, lease reclaim, events, workers, leadership, cron and usage, then lists Redis raw. On v1.6.3 Dispatch can't dequeue the job it just enqueued. On v1.7.0 every key stays in the namespace. The two Lua scripts only touchKEYS[1], so the hook covers them too.WithKeyPrefixis unchanged and doesn't depend on hooks either way.Testing
go build ./...andgo test ./...: 41 packages passgo test -race ./worker/... ./engine/...: passgo test -tags integration ./store/...against throwaway containers: memory, mongo, postgres, redis, sqlite and storetest all passgolangci-lint run ./...with a fresh cache: 0 issuesDependabot
#26 (Go modules) touches
go.modandgo.sum, same as this PR. It's also stale: it was opened in May, and it would take go-redis back to 9.19 from main's 9.21 and re-bump bun modules main no longer has. Closing it and letting dependabot regenerate after this merges is probably simplest. #25 (docs npm), #19, #18 and #6 (Actions) don't touch anything here.