git-refs: compact the push queue on Drain · Entire
git-refs: compact the push queue on Drain
40cded1→main·
Soph·2w ago·2 files·+109 added/-16 removed
Enqueue only appends, so a long-lived session that re-enqueues the same checkpoint ref across many writes without pushing would grow the queue file unboundedly — Remove was the only point that rewrote it. Drain already de-duplicates in memory; now it also rewrites the file to that de-duplicated set when the on-disk queue held redundant lines (duplicate refs or malformed/blank records), bounding the file to one line per distinct queued ref.
readLocked reports the raw non-empty line count so Drain compacts only when rawLines > len(refs) (i.e. there was actually something redundant). The rewrite is factored into rewriteLocked, shared with Remove, and stays atomic (temp + rename). Drain still returns the refs and does not clear them — they survive until a confirmed Remove.
Co-Authored-By: Claude Opus 4.8 (1M context) noreply@anthropic.com
Sessions
5608d53754d2View transcript
Changes
2
cmd/entire/cli/checkpoint
Mpushqueue.go+49/-16
Mpushqueue_test.go+60
83 unmodified lines
// Drain returns the de-duplicated refs currently queued, in first-seen order. It
// does NOT remove them; call Remove after a confirmed push so a failed push
// retries next time. A missing queue file yields no refs.
//
// It compacts the file in place: when the on-disk queue held redundant lines
// (duplicate enqueues of the same ref, or malformed/blank lines), Drain rewrites
// it to the de-duplicated set. Enqueue only ever appends, so without this the
// file would grow unboundedly between the Removes that are otherwise the sole
// compaction point (e.g. a long-lived session that keeps re-enqueuing the same
// checkpoint ref but never pushes).
func (q *PushQueue) Drain() ([]plumbing.ReferenceName, error) {
release, err := flock.Acquire(q.lockPath())
if err != nil {
return nil, fmt.Errorf("lock push queue: %w", err)
}
defer release()
return q.readLocked()
refs, rawLines, err := q.readLocked()
if err != nil {
return nil, err
}
if rawLines > len(refs) {
if err := q.rewriteLocked(refs); err != nil {
return nil, err
}
}
return refs, nil
}
// Remove deletes the given refs from the queue, preserving any entries appended
}
current, err := q.readLocked()
current, _, err := q.readLocked()
if err != nil {
return err
}
for _, r := range refs {
removed[r.String()] = struct{}{}
}
var buf bytes.Buffer
kept := make([]plumbing.ReferenceName, 0, len(current))
for _, r := range current {
if _, drop := removed[r.String()]; drop {
continue
}
kept = append(kept, r)
}
return q.rewriteLocked(kept)
}
// rewriteLocked replaces the queue file with exactly refs (de-duplicated, one
// line each), or removes the file when refs is empty so a clean repo has no
// stray queue. The caller must hold the lock. The write is atomic (temp file +
// rename) so a concurrent reader never sees a half-written queue.
func (q *PushQueue) rewriteLocked(refs []plumbing.ReferenceName) error {
if len(refs) == 0 {
if err := os.Remove(q.queuePath()); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove empty push queue: %w", err)
}
return nil
}
var buf bytes.Buffer
for _, r := range refs {
line, err := json.Marshal(pushQueueEntry{Ref: r.String()})
if err != nil {
return fmt.Errorf("encode push queue entry: %w", err)
}
buf.Write(line)
buf.WriteByte('\n')
}
if buf.Len() == 0 {
// Nothing left: drop the file so a clean repo has no stray queue.
if err := os.Remove(q.queuePath()); err != nil && !os.IsNotExist(err) {
return fmt.Errorf("remove empty push queue: %w", err)
}
return nil
}
if err := writeFileAtomicInDir(q.dir, q.queuePath(), buf.Bytes()); err != nil {
return fmt.Errorf("rewrite push queue: %w", err)
}
}
// readLocked parses the queue file into de-duplicated refs, preserving first-seen
// order. The caller must hold the lock. Malformed lines are skipped rather than
// failing the whole drain — a single bad record must not strand every queued ref.
func (q *PushQueue) readLocked() ([]plumbing.ReferenceName, error) {
// rawLines is the number of non-empty lines seen (including duplicates and
// malformed records), so callers can detect when the file holds more than the
// de-duplicated set and is worth compacting: rawLines > len(refs) exactly when
// there were redundant lines.
func (q *PushQueue) readLocked() (refs []plumbing.ReferenceName, rawLines int, err error) {
f, err := os.Open(q.queuePath())
if err != nil {
if os.IsNotExist(err) {
return nil, nil
return nil, 0, nil
}
return nil, fmt.Errorf("open push queue: %w", err)
return nil, 0, fmt.Errorf("open push queue: %w", err)
}
defer f.Close()
var refs []plumbing.ReferenceName
seen := make(map[string]struct{})
scanner := bufio.NewScanner(f)
scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024)
if len(line) == 0 {
continue
}
rawLines++
var entry pushQueueEntry
if err := json.Unmarshal(line, &entry); err != nil || entry.Ref == "" {
continue
}
refs = append(refs, plumbing.ReferenceName(entry.Ref))
}
if err := scanner.Err(); err != nil {
return nil, fmt.Errorf("read push queue: %w", err)
return nil, 0, fmt.Errorf("read push queue: %w", err)
}
return refs, nil
return refs, rawLines, nil
}
// writeFileAtomicInDir writes data to a temp file in dir and renames it over
Mcmd/entire/cli/checkpoint/pushqueue.go+49/-16
2 unmodified lines
3
4
5
6
7
8
9
53 unmodified lines
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
2 unmodified lines
import (
"os"
"path/filepath"
"strings"
"testing"
"github.com/go-git/go-git/v6/plumbing"
53 unmodified lines
assert.Equal(t, []plumbing.ReferenceName{a}, refs, "duplicates collapse to one")
// nonEmptyLineCount returns how many non-blank lines the queue file holds.
func nonEmptyLineCount(t *testing.T, q *PushQueue) int {
t.Helper()
data, err := os.ReadFile(q.queuePath())
if os.IsNotExist(err) {
return 0
}
require.NoError(t, err)
count := 0
for _, line := range strings.Split(string(data), "\n") {
if strings.TrimSpace(line) != "" {
count++
}
}
return count
}
func TestPushQueue_DrainCompactsRedundantEntries(t *testing.T) {
t.Parallel()
q := NewPushQueue(t.TempDir())
a := mustRefName(t, "a1b2c3d4e5f6")
b := mustRefName(t, "b2c3d4e5f6a1")
require.NoError(t, q.Enqueue(a))
require.NoError(t, q.Enqueue(a))
require.NoError(t, q.Enqueue(b))
require.NoError(t, q.Enqueue(a))
require.Equal(t, 4, nonEmptyLineCount(t, q), "enqueue only appends")
refs, err := q.Drain()
require.NoError(t, err)
assert.Equal(t, []plumbing.ReferenceName{a, b}, refs)
assert.Equal(t, 2, nonEmptyLineCount(t, q), "Drain compacts the file to the de-duplicated set")
// The refs still survive until Remove, and a re-drain does not rewrite again.
refs, err = q.Drain()
require.NoError(t, err)
assert.Equal(t, []plumbing.ReferenceName{a, b}, refs, "compaction preserves queued refs")
assert.Equal(t, 2, nonEmptyLineCount(t, q))
}
func TestPushQueue_DrainCompactsMalformedLines(t *testing.T) {
t.Parallel()
q := NewPushQueue(t.TempDir())
a := mustRefName(t, "a1b2c3d4e5f6")
require.NoError(t, q.Enqueue(a))
f, err := os.OpenFile(q.queuePath(), os.O_WRONLY|os.O_APPEND, 0o600)
require.NoError(t, err)
_, err = f.WriteString("not json\n\n")
require.NoError(t, err)
require.NoError(t, f.Close())
refs, err := q.Drain()
require.NoError(t, err)
assert.Equal(t, []plumbing.ReferenceName{a}, refs)
assert.Equal(t, 1, nonEmptyLineCount(t, q), "Drain drops malformed lines from disk")
}
func TestPushQueue_RemovePreservesLaterEntries(t *testing.T) {
t.Parallel()
q := NewPushQueue(t.TempDir())