Add git-refs per-checkpoint checkpoint backend · Entire
Add git-refs per-checkpoint checkpoint backend
989c1c0→main·
Soph·2w ago·11 files·+1,143 added/-14 removed
Second slice of #1471: a git-backed store that keeps one commit per
refs/entire/checkpoints/
- refs_store.go: gitRefsStore implements PersistentStore over per-checkpoint refs (orphan-then-parented history, stamps refs-v1, no Vercel merge, List enumerates local refs, optional AuthorReader). Reads resolve a ref → commit tree and use shared read helpers; writes build the subtree via the embedded *treeWriter.
- Share the tree-read helpers: the git-branch reader's Read/ReadSession* now delegate to readSummaryFromCheckpointTree / readSession*FromTree free functions (no behavior change) so both backends read the same way, just navigating to the tree differently.
- registry: BackendTypeGitRefs ("git-refs") registered built-in with gitBacked:true
- gitRefsBackendFactory; OpenEnv/OpenOptions gain RefFetcher; PrimaryIsRefs helper.
- pushqueue.go: flock-protected JSONL push-discovery queue in the git common dir (Enqueue/Drain/Remove); gitRefsStore.setRef enqueues best-effort. (Pre-push consumption lands in the next commit.)
- CheckpointVersionRefsV1 = "refs-v1".
Co-Authored-By: Claude Opus 4.8 (1M context) noreply@anthropic.com
Sessions
a4baf5553823View transcript
Changes
11
api/checkpoint
Merrors.go+6
cmd/entire/cli/checkpoint
Maliases.go+3
Mfetching_tree.go+6
Mopen.go+14/-1
Mpersistent.go+27/-11
Apushqueue.go+200
Apushqueue_test.go+109
Arefs_store.go+381
Arefs_store_seam_test.go+99
Arefs_store_test.go+269
Mregistry.go+29/-2
12 unmodified lines
13
14
15
16
17
18
19
20
21
12 unmodified lines
// CheckpointVersionBranchV1 identifies the branch-backed checkpoint metadata format.
const CheckpointVersionBranchV1 = "branch-v1"
// CheckpointVersionRefsV1 identifies the per-checkpoint-ref checkpoint metadata
// format (one ref per checkpoint at refs/entire/checkpoints/<shard>/<id>). The
// value follows the <family>-v<major> convention (cf. branch-v1) so
// checkpointpolicy.ParseFormat parses it.
const CheckpointVersionRefsV1 = "refs-v1"
Mapi/checkpoint/errors.go+6
53 unmodified lines
54
55
56
57
58
59
60
61
62
53 unmodified lines
// CheckpointVersionBranchV1 identifies the branch-backed checkpoint metadata format.
const CheckpointVersionBranchV1 = apicheckpoint.CheckpointVersionBranchV1
// CheckpointVersionRefsV1 identifies the per-checkpoint-ref checkpoint metadata format.
const CheckpointVersionRefsV1 = apicheckpoint.CheckpointVersionRefsV1
// Sentinel errors (re-exported so errors.Is keeps working across packages).
var (
ErrCheckpointNotFound = apicheckpoint.ErrCheckpointNotFound
)
Mcmd/entire/cli/checkpoint/aliases.go+3
16 unmodified lines
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
16 unmodified lines
// BlobFetchFunc fetches missing blob objects by hash from a remote.
type BlobFetchFunc func(ctx context.Context, hashes []plumbing.Hash) error
// RefFetchFunc fetches a single checkpoint ref from the remote into the local
// ref of the same name. The git-refs store uses it to resolve a checkpoint ref
// that is not present locally (e.g. written on another machine). The checkpoint
// package cannot resolve the remote target itself, so the CLI layer injects it.
type RefFetchFunc func(ctx context.Context, ref plumbing.ReferenceName) error
// FetchingTree wraps a git tree to automatically fetch missing blobs on demand.
// After a treeless fetch (--filter=blob:none), tree objects are available locally
// but blob objects are not. Each File() call checks whether the target blob
Mcmd/entire/cli/checkpoint/fetching_tree.go+6
18 unmodified lines
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
19 unmodified lines
62
63
64
52
65
66
67
68
18 unmodified lines
// fetching off.
BlobFetcher BlobFetchFunc
// RefFetcher is the CLI-level on-demand checkpoint-ref fetcher, used by the
// git-refs backend to resolve a checkpoint ref missing locally. nil leaves
// reads local-only; ignored by the git-branch backend.
RefFetcher RefFetchFunc
// Refs overrides the default committed-ref topology. A non-nil value wins,
// e.g. attach pins reads to Primary via PrimaryAsRead().
Refs *PersistentRefs
}
// PrimaryIsRefs reports whether the configured primary backend is the git-refs
// per-checkpoint store. It centralizes the topology check so push/pre-push code
// does not compare backend-type strings itself. A nil config (default) is the
// git-branch backend, so this returns false.
func PrimaryIsRefs(cfg *settings.CheckpointsConfig) bool {
return cfg != nil && cfg.Primary.Type == BackendTypeGitRefs
}
// Stores is the facade returned by Open: the persistent store plus the git-only
// ephemeral (shadow-branch) capability and resolved committed-ref topology.
type Stores struct {
19 unmodified lines
// default git-branch backend with no mirrors, preserving default behavior.
func Open(ctx context.Context, repo *git.Repository, opts OpenOptions) (*Stores, error) {
refs := resolveOpenRefs(ctx, opts)
env := OpenEnv{Repo: repo, BlobFetcher: opts.BlobFetcher, Refs: refs}
env := OpenEnv{Repo: repo, BlobFetcher: opts.BlobFetcher, RefFetcher: opts.RefFetcher, Refs: refs}
cfg, err := settings.LoadCheckpointsConfig(ctx)
if err != nil {
Mcmd/entire/cli/checkpoint/open.go+14/-1
1298 unmodified lines
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
46 unmodified lines
1361
1362
1363
1364
1365
1366
1367
1368
1369
1370
1371
1372
19 unmodified lines
1392
1393
1394
1395
1396
1397
1381
1398
1399
1400
1383
1384
1385
1386
1387
1388
1389
1390
1391
1401
1402
1403
1404
3 unmodified lines
1408
1409
1410
1401
1411
1412
1413
1414
1 unmodified line
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
15 unmodified lines
1441
1442
1443
1444
1445
1446
1447
1448
1449
1450
1298 unmodified lines
return nil, nil //nolint:nilnil,nilerr // Checkpoint directory not found
}
return readSummaryFromCheckpointTree(checkpointTree)
}
// readSummaryFromCheckpointTree reads the root CheckpointSummary from a checkpoint
// tree (the tree holding metadata.json plus the numbered session dirs). It is
// shared by the git-branch store (which descends to <shard>/<id>) and the
// git-refs store (whose ref tree is the checkpoint tree directly). It returns
// (nil, nil) when metadata.json is absent so callers normalize a missing
// checkpoint to ErrCheckpointNotFound via the contract.
func readSummaryFromCheckpointTree(checkpointTree *FetchingTree) (*CheckpointSummary, error) {
// Read root metadata.json as CheckpointSummary (auto-fetches blob if needed)
.metadataFile, err := checkpointTree.File(paths.MetadataFileName)
if err != nil {
46 unmodified lines
if err != nil {
return nil, err
}
return readSessionMetadataFromTree(sessionTree, sessionIndex)
}
metadataFile, err := sessionTree.File(paths.MetadataFileName)
Mcmd/entire/cli/checkpoint/persistent.go+27/-11
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
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
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
package checkpoint
import (
"bufio"
"bytes"
"context"
"encoding/json"
"fmt"
"os"
"path/filepath"
"github.com/go-git/go-git/v6"
"github.com/go-git/go-git/v6/plumbing"
"github.com/entireio/cli/cmd/entire/cli/internal/flock"
)
// Push-discovery queue file names, kept in the git common dir so every worktree
// sharing the object store enqueues into one queue. The git-refs backend cannot
// push every local checkpoint ref at pre-push time (reads fetch refs too, and
// deleting them after push would hurt local workflows), so each write records
// the ref it touched here and pre-push drains + batch-pushes exactly those.
const (
pushQueueFileName = "entire-checkpoint-push-queue.jsonl"
pushQueueLockName = "entire-checkpoint-push-queue.lock"
)
// pushQueueEntry is one JSONL record: a checkpoint ref awaiting push.
type pushQueueEntry struct {
Ref string `json:"ref"`
}
// PushQueue is a flock-protected JSONL list of checkpoint refs awaiting push,
// stored in the git common dir. Entries are removed only after a confirmed push
// (Remove), so an interrupted or failed push leaves them for the next pre-push.
// Duplicates are tolerated on disk and collapsed by Drain.
type PushQueue struct {
dir string
}
// NewPushQueue returns the push queue rooted at gitCommonDir.
func NewPushQueue(gitCommonDir string) *PushQueue {
return &PushQueue{dir: gitCommonDir}
}
// PushQueueForRepo resolves the git common dir for repo and returns its queue.
func PushQueueForRepo(ctx context.Context, repo *git.Repository) (*PushQueue, error) {
dir, err := resolveGitCommonDir(ctx, repo)
if err != nil {
return nil, err
}
return NewPushQueue(dir), nil
}
func (q *PushQueue) queuePath() string { return filepath.Join(q.dir, pushQueueFileName) }
func (q *PushQueue) lockPath() string { return filepath.Join(q.dir, pushQueueLockName) }
// Enqueue appends a ref to the queue. It is safe to enqueue a ref already
// present (or already pushed): Drain collapses duplicates and the batch push is
// idempotent. Enqueue takes the lock so concurrent writers never interleave a
// partial line.
func (q *PushQueue) Enqueue(ref plumbing.ReferenceName) error {
release, err := flock.Acquire(q.lockPath())
if err != nil {
return fmt.Errorf("lock push queue: %w", err)
}
defer release()
line, err := json.Marshal(pushQueueEntry{Ref: ref.String()})
if err != nil {
return fmt.Errorf("encode push queue entry: %w", err)
}
f, err := os.OpenFile(q.queuePath(), os.O_CREATE|os.O_WRONLY|os.O_APPEND, 0o600)
if err != nil {
return fmt.Errorf("open push queue: %w", err)
}
defer f.Close()
if _, err := f.Write(append(line, '\n')); err != nil {
return fmt.Errorf("append push queue entry: %w", err)
}
return nil
}
// 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.
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()
}
// Remove deletes the given refs from the queue, preserving any entries appended
// after a Drain (e.g. a write that landed during the push). Called after a
// confirmed push.
func (q *PushQueue) Remove(refs []plumbing.ReferenceName) error {
if len(refs) == 0 {
return nil
}
release, err := flock.Acquire(q.lockPath())
if err != nil {
return fmt.Errorf("lock push queue: %w", err)
}
defer release()
current, err := q.readLocked()
if err != nil {
return err
}
removed := make(map[string]struct{}, len(refs))
for _, r := range refs {
removed[r.String()] = struct{}{}
}
var buf bytes.Buffer
for _, r := range current {
if _, drop := removed[r.String()]; drop {
continue
}
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)
}
return nil
}
// 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) {
f, err := os.Open(q.queuePath())
if err != nil {
if os.IsNotExist(err) {
return nil, nil
}
return nil, 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)
for scanner.Scan() {
line := bytes.TrimSpace(scanner.Bytes())
if len(line) == 0 {
continue
}
var entry pushQueueEntry
if err := json.Unmarshal(line, &entry); err != nil || entry.Ref == "" {
continue
}
if _, dup := seen[entry.Ref]; dup {
continue
}
seen[entry.Ref] = struct{}{}
refs = append(refs, plumbing.ReferenceName(entry.Ref))
}
if err := scanner.Err(); err != nil {
return nil, fmt.Errorf("read push queue: %w", err)
}
return refs, nil
}
// writeFileAtomicInDir writes data to a temp file in dir and renames it over
// path, so a reader (under the lock) never sees a half-written queue.
func writeFileAtomicInDir(dir, path string, data []byte) error {
tmp, err := os.CreateTemp(dir, pushQueueFileName+".*")
if err != nil {
return fmt.Errorf("create temp push queue: %w", err)
}
tmpName := tmp.Name()
defer os.Remove(tmpName)
if _, err := tmp.Write(data); err != nil {
_ = tmp.Close()
return fmt.Errorf("write temp push queue: %w", err)
}
if err := tmp.Close(); err != nil {
return fmt.Errorf("close temp push queue: %w", err)
}
if err := os.Rename(tmpName, path); err != nil {
return fmt.Errorf("rename temp push queue: %w", err)
}
return nil
}
// Additional functions can be added here to further structure the content
// as needed.