Skip to content
4 changes: 2 additions & 2 deletions cmd/entire/cli/checkpoint/migrate.go
Original file line number Diff line number Diff line change
Expand Up @@ -127,7 +127,7 @@ func migrateBranchToRefs(
return nil
}
if alreadyImported {
if err := queue.Enqueue(refName); err != nil {
if err := queue.EnqueueRef(repo, refName); err != nil {
return fmt.Errorf("enqueue checkpoint %s for push: %w", cid, err)
}
result.Skipped++
Expand Down Expand Up @@ -176,7 +176,7 @@ func migrateBranchToRefs(

// The migration's queued-for-push contract is guaranteed for both new
// and already-imported refs. Duplicates collapse on Drain.
if err := queue.Enqueue(refName); err != nil {
if err := queue.EnqueueRef(repo, refName); err != nil {
return fmt.Errorf("enqueue checkpoint %s for push: %w", cid, err)
}
if migrated {
Expand Down
180 changes: 137 additions & 43 deletions cmd/entire/cli/checkpoint/pushqueue.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,15 @@ func (q *PushQueue) lock() (*os.Root, func(), error) {

// pushQueueEntry is one JSONL record: a checkpoint ref awaiting push.
type pushQueueEntry struct {
Ref string `json:"ref"`
Ref string `json:"ref"`
Hash string `json:"hash,omitempty"`
}

// PushQueueEntry identifies one observed generation of a checkpoint ref. Hash
// is zero only for queue records written by older Entire versions.
type PushQueueEntry struct {
Ref plumbing.ReferenceName
Hash plumbing.Hash
}

// PushQueue is a flock-protected JSONL list of checkpoint refs awaiting push,
Expand Down Expand Up @@ -78,13 +86,24 @@ func PushQueueForRepo(_ context.Context, repo *git.Repository) (*PushQueue, erro
// idempotent. Enqueue takes the lock so concurrent writers never interleave a
// partial line.
func (q *PushQueue) Enqueue(ref plumbing.ReferenceName) error {
return q.EnqueueEntry(PushQueueEntry{Ref: ref})
}

// EnqueueEntry appends an observed ref generation. A later generation of the
// same ref supersedes it when the queue is read, while exact-generation removal
// cannot erase the later entry.
func (q *PushQueue) EnqueueEntry(entry PushQueueEntry) error {
root, release, err := q.lock()
if err != nil {
return err
}
defer release()

line, err := json.Marshal(pushQueueEntry{Ref: ref.String()})
diskEntry := pushQueueEntry{Ref: entry.Ref.String()}
if !entry.Hash.IsZero() {
diskEntry.Hash = entry.Hash.String()
}
line, err := json.Marshal(diskEntry)
if err != nil {
return fmt.Errorf("encode push queue entry: %w", err)
}
Expand All @@ -99,6 +118,18 @@ func (q *PushQueue) Enqueue(ref plumbing.ReferenceName) error {
return nil
}

// EnqueueRef records the ref's current observed generation. When the ref
// cannot be resolved it still queues the ref, with no generation: a ref that
// misses the queue is never pushed, while a generation-less entry is the legacy
// form delivery already re-resolves (and prunes if the ref is gone).
func (q *PushQueue) EnqueueRef(repo *git.Repository, refName plumbing.ReferenceName) error {
entry := PushQueueEntry{Ref: refName}
if ref, err := repo.Reference(refName, true); err == nil {
entry.Hash = ref.Hash()
}
return q.EnqueueEntry(entry)
}

// 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.
Expand All @@ -110,37 +141,50 @@ func (q *PushQueue) Enqueue(ref plumbing.ReferenceName) error {
// 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) {
entries, err := q.DrainEntries()
return pushQueueRefs(entries), err
}

// DrainEntries returns the latest observed generation of every queued ref in
// first-seen ref order. It compacts redundant, superseded, and malformed lines.
func (q *PushQueue) DrainEntries() ([]PushQueueEntry, error) {
root, release, err := q.lock()
if err != nil {
return nil, err
}
defer release()

refs, rawLines, err := q.readLocked(root)
entries, rawLines, err := q.readEntriesLocked(root)
if err != nil {
return nil, err
}
if rawLines > len(refs) {
if err := q.rewriteLocked(root, refs); err != nil {
if rawLines > len(entries) {
if err := q.rewriteEntriesLocked(root, entries); err != nil {
return nil, err
}
}
return refs, nil
return entries, nil
}

// Peek returns the de-duplicated refs currently queued, in first-seen order,
// without mutating the queue file. Read-only counterpart to Drain (which
// compacts redundant lines in place) — for counters/status displays that must
// observe the queue without owning a push. A missing queue file yields no refs.
func (q *PushQueue) Peek() ([]plumbing.ReferenceName, error) {
entries, err := q.PeekEntries()
return pushQueueRefs(entries), err
}

// PeekEntries is the read-only counterpart to DrainEntries.
func (q *PushQueue) PeekEntries() ([]PushQueueEntry, error) {
root, release, err := q.lock()
if err != nil {
return nil, err
}
defer release()

refs, _, err := q.readLocked(root)
return refs, err
entries, _, err := q.readEntriesLocked(root)
return entries, err
}

// Remove deletes the given refs from the queue, preserving any entries appended
Expand All @@ -156,28 +200,59 @@ func (q *PushQueue) Remove(refs []plumbing.ReferenceName) error {
}
defer release()

current, _, err := q.readLocked(root)
current, _, err := q.readEntriesLocked(root)
if err != nil {
return err
}
removed := make(map[string]struct{}, len(refs))
for _, r := range refs {
removed[r.String()] = struct{}{}
}
kept := make([]plumbing.ReferenceName, 0, len(current))
for _, r := range current {
if _, drop := removed[r.String()]; drop {
kept := make([]PushQueueEntry, 0, len(current))
for _, entry := range current {
if _, drop := removed[entry.Ref.String()]; drop {
continue
}
kept = append(kept, r)
kept = append(kept, entry)
}
return q.rewriteLocked(root, kept)
return q.rewriteEntriesLocked(root, kept)
}

// Rotate moves the given refs to the back of the queue, keeping the relative
// order within both the moved and the untouched group. Refs not currently queued
// are ignored, and entries appended after a Drain keep their place ahead of the
// rotated ones.
// RemoveEntries removes only exact ref generations. A newer generation of the
// same ref therefore survives cleanup for a push that delivered an older one.
func (q *PushQueue) RemoveEntries(entries []PushQueueEntry) error {
if len(entries) == 0 {
return nil
}
root, release, err := q.lock()
if err != nil {
return err
}
defer release()

current, _, err := q.readEntriesLocked(root)
if err != nil {
return err
}
removed := make(map[PushQueueEntry]struct{}, len(entries))
for _, entry := range entries {
removed[entry] = struct{}{}
}
kept := make([]PushQueueEntry, 0, len(current))
for _, entry := range current {
if _, drop := removed[entry]; drop {
continue
}
kept = append(kept, entry)
}
return q.rewriteEntriesLocked(root, kept)
}

// Rotate moves the given queue entries to the back of the queue, keeping the
// relative order within both the moved and the untouched group. Entries match
// by exact generation (ref and hash), so a newer generation of the same ref
// enqueued meanwhile keeps its place. Entries not currently queued are ignored,
// and entries appended after a Drain keep their place ahead of the rotated ones.
//
// This is what keeps a bounded flush fair. The queue is drained in first-seen
// order and a flush that stops early leaves the rest for the next push, so
Expand All @@ -189,8 +264,8 @@ func (q *PushQueue) Remove(refs []plumbing.ReferenceName) error {
// One lock for the whole reorder, rather than Remove followed by Enqueue: that
// pair leaves a window where the refs are in neither the queue nor the remote,
// and a crash inside it drops them.
func (q *PushQueue) Rotate(refs []plumbing.ReferenceName) error {
if len(refs) == 0 {
func (q *PushQueue) Rotate(entries []PushQueueEntry) error {
if len(entries) == 0 {
return nil
}
root, release, err := q.lock()
Expand All @@ -199,43 +274,47 @@ func (q *PushQueue) Rotate(refs []plumbing.ReferenceName) error {
}
defer release()

current, _, err := q.readLocked(root)
current, _, err := q.readEntriesLocked(root)
if err != nil {
return err
}
rotating := make(map[string]struct{}, len(refs))
for _, r := range refs {
rotating[r.String()] = struct{}{}
rotating := make(map[PushQueueEntry]struct{}, len(entries))
for _, entry := range entries {
rotating[entry] = struct{}{}
}
kept := make([]plumbing.ReferenceName, 0, len(current))
moved := make([]plumbing.ReferenceName, 0, len(refs))
for _, r := range current {
if _, rotate := rotating[r.String()]; rotate {
moved = append(moved, r)
kept := make([]PushQueueEntry, 0, len(current))
moved := make([]PushQueueEntry, 0, len(entries))
for _, entry := range current {
if _, rotate := rotating[entry]; rotate {
moved = append(moved, entry)
continue
}
kept = append(kept, r)
kept = append(kept, entry)
}
if len(moved) == 0 {
return nil
}
return q.rewriteLocked(root, append(kept, moved...))
return q.rewriteEntriesLocked(root, append(kept, moved...))
}

// rewriteLocked replaces the queue file with exactly refs (de-duplicated, one
// rewriteEntriesLocked replaces the queue file with exactly entries (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(root *os.Root, refs []plumbing.ReferenceName) error {
if len(refs) == 0 {
func (q *PushQueue) rewriteEntriesLocked(root *os.Root, entries []PushQueueEntry) error {
if len(entries) == 0 {
if err := root.Remove(pushQueueFileName); 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()})
for _, entry := range entries {
diskEntry := pushQueueEntry{Ref: entry.Ref.String()}
if !entry.Hash.IsZero() {
diskEntry.Hash = entry.Hash.String()
}
line, err := json.Marshal(diskEntry)
if err != nil {
return fmt.Errorf("encode push queue entry: %w", err)
}
Expand All @@ -248,15 +327,15 @@ func (q *PushQueue) rewriteLocked(root *os.Root, refs []plumbing.ReferenceName)
return nil
}

// readLocked parses the queue file into de-duplicated refs, preserving first-seen
// readEntriesLocked parses the queue file into de-duplicated entries, 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.
//
// 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(root *os.Root) (refs []plumbing.ReferenceName, rawLines int, err error) {
func (q *PushQueue) readEntriesLocked(root *os.Root) (entries []PushQueueEntry, rawLines int, err error) {
f, err := root.Open(pushQueueFileName)
if err != nil {
if os.IsNotExist(err) {
Expand All @@ -266,7 +345,7 @@ func (q *PushQueue) readLocked(root *os.Root) (refs []plumbing.ReferenceName, ra
}
defer f.Close()

seen := make(map[string]struct{})
seen := make(map[string]int)
scanner := bufio.NewScanner(f)
scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024)
for scanner.Scan() {
Expand All @@ -279,16 +358,31 @@ func (q *PushQueue) readLocked(root *os.Root) (refs []plumbing.ReferenceName, ra
if err := json.Unmarshal(line, &entry); err != nil || entry.Ref == "" {
continue
}
if _, dup := seen[entry.Ref]; dup {
parsed := PushQueueEntry{Ref: plumbing.ReferenceName(entry.Ref)}
// A corrupt generation keeps the ref, generation-less: dropping the line
// would compact it away, and the queue is the only push discovery.
if plumbing.IsHash(entry.Hash) {
parsed.Hash = plumbing.NewHash(entry.Hash)
}
if idx, dup := seen[entry.Ref]; dup {
entries[idx] = parsed
continue
}
seen[entry.Ref] = struct{}{}
refs = append(refs, plumbing.ReferenceName(entry.Ref))
seen[entry.Ref] = len(entries)
entries = append(entries, parsed)
}
if err := scanner.Err(); err != nil {
return nil, 0, fmt.Errorf("read push queue: %w", err)
}
return refs, rawLines, nil
return entries, rawLines, nil
}

func pushQueueRefs(entries []PushQueueEntry) []plumbing.ReferenceName {
refs := make([]plumbing.ReferenceName, 0, len(entries))
for _, entry := range entries {
refs = append(refs, entry.Ref)
}
return refs
}

// writeQueueAtomic writes data to a temp file inside root and renames it over
Expand Down
Loading
Loading