diff options
Diffstat (limited to 'sidecar')
| -rw-r--r-- | sidecar/store.go | 329 | ||||
| -rw-r--r-- | sidecar/store_test.go | 164 |
2 files changed, 493 insertions, 0 deletions
diff --git a/sidecar/store.go b/sidecar/store.go new file mode 100644 index 0000000..76e626f --- /dev/null +++ b/sidecar/store.go @@ -0,0 +1,329 @@ +// Package sidecar stores vr's append-only JSONL log in a private jj change. +// +// The change is an anonymous direct child of root(), so it is visible and +// durable in the jj repository without becoming part of project history or the +// working copy. A repo-local revset alias names it and git.private-commits keeps +// accidental --change pushes from publishing it. +package sidecar + +import ( + "errors" + "fmt" + "os" + "os/exec" + "path/filepath" + "strconv" + "strings" + "syscall" +) + +const ( + Alias = "vr_log" + LogFile = "vr-log.jsonl" + Description = "private: vr log" +) + +var ErrNoLog = errors.New("vr log has not been initialized") + +// Store is a repository's vr log. Open is read-only with respect to jj; the +// sidecar change and its config are created lazily by Append. +type Store struct { + root string + lockPath string +} + +// Open finds the jj workspace containing dir. +func Open(dir string) (*Store, error) { + root, err := jjOutput(dir, "workspace", "root") + if err != nil { + return nil, err + } + root = strings.TrimSpace(root) + configPath, err := jjOutput(root, "config", "path", "--repo") + if err != nil { + return nil, err + } + return &Store{ + root: root, + lockPath: filepath.Join(filepath.Dir(strings.TrimSpace(configPath)), "vr-agent-logger.lock"), + }, nil +} + +func (s *Store) Root() string { return s.root } + +func (s *Store) String() string { + return fmt.Sprintf("jj sidecar %s:%s", Alias, LogFile) +} + +// Read returns the current log. Before the first sidecar write it falls back +// to a repo-root legacy log, which makes migration transparent. +func (s *Store) Read() ([]byte, error) { + var data []byte + err := s.withLock(syscall.LOCK_SH, func() error { + rev, found, _, err := s.resolveExisting() + if err != nil { + return err + } + if found { + data, err = s.readRevision(rev) + return err + } + data, err = os.ReadFile(filepath.Join(s.root, LogFile)) + if errors.Is(err, os.ErrNotExist) { + return ErrNoLog + } + return err + }) + return data, err +} + +// Generation is a cheap identity for the current contents, used by vrsite to +// notice appends. A sidecar rewrite gets a new commit id. +func (s *Store) Generation() (string, error) { + var generation string + err := s.withLock(syscall.LOCK_SH, func() error { + rev, found, _, err := s.resolveExisting() + if err != nil { + return err + } + if found { + generation, err = s.logValue(rev, "commit_id") + return err + } + info, err := os.Stat(filepath.Join(s.root, LogFile)) + if errors.Is(err, os.ErrNotExist) { + generation = "empty" + return nil + } + if err != nil { + return err + } + generation = fmt.Sprintf("legacy:%d:%d", info.Size(), info.ModTime().UnixNano()) + return nil + }) + return generation, err +} + +// Append serializes sidecar rewrites across vr processes. makePayload runs +// while holding that lock and receives the complete current log, allowing note +// numbering and reply validation to be atomic with the append. +func (s *Store) Append(makePayload func(current []byte) ([]byte, error)) error { + return s.withLock(syscall.LOCK_EX, func() error { + rev, err := s.ensure() + if err != nil { + return err + } + sidecarData, err := s.readRevision(rev) + if err != nil { + return err + } + current := sidecarData + if len(current) == 0 { + legacy, legacyErr := os.ReadFile(filepath.Join(s.root, LogFile)) + if legacyErr == nil { + current = legacy + } else if !errors.Is(legacyErr, os.ErrNotExist) { + return legacyErr + } + } + payload, err := makePayload(current) + if err != nil { + return err + } + if len(payload) == 0 { + return errors.New("refusing to append an empty vr log payload") + } + + var appendData []byte + if len(sidecarData) == 0 && len(current) > 0 { + appendData = append(appendData, current...) + } + if len(current) > 0 && current[len(current)-1] != '\n' { + appendData = append(appendData, '\n') + } + appendData = append(appendData, payload...) + if appendData[len(appendData)-1] != '\n' { + appendData = append(appendData, '\n') + } + return s.runAppend(appendData) + }) +} + +func (s *Store) withLock(kind int, fn func() error) error { + if err := os.MkdirAll(filepath.Dir(s.lockPath), 0o755); err != nil { + return err + } + f, err := os.OpenFile(s.lockPath, os.O_CREATE|os.O_RDWR, 0o600) + if err != nil { + return err + } + defer f.Close() + if err := syscall.Flock(int(f.Fd()), kind); err != nil { + return err + } + defer syscall.Flock(int(f.Fd()), syscall.LOCK_UN) //nolint:errcheck + return fn() +} + +func (s *Store) ensure() (string, error) { + rev, found, configured, err := s.resolveExisting() + if err != nil { + return "", err + } + if !found { + if _, err := s.jj("new", "--no-edit", "root()", "-m", Description); err != nil { + return "", err + } + ids, err := s.descriptionMatches() + if err != nil { + return "", err + } + if len(ids) != 1 { + return "", fmt.Errorf("created %q sidecar but found %d matching root children", Description, len(ids)) + } + rev = ids[0] + } + if !configured { + if _, err := s.jj("config", "set", "--repo", "revset-aliases."+Alias, strconv.Quote(rev)); err != nil { + return "", err + } + } + private, err := s.jj("config", "get", "git.private-commits") + if err != nil { + return "", err + } + if !containsRevsetName(private, Alias) { + expr := fmt.Sprintf("(%s) | %s", strings.TrimSpace(private), Alias) + if _, err := s.jj("config", "set", "--repo", "git.private-commits", strconv.Quote(expr)); err != nil { + return "", err + } + } + return rev, nil +} + +func (s *Store) resolveExisting() (rev string, found, configured bool, err error) { + configuredValue, configErr := s.jj("config", "get", "revset-aliases."+Alias) + if configErr == nil { + rev, err = s.logValue("exactly("+Alias+", 1)", "change_id") + if err != nil { + return "", false, true, fmt.Errorf("repo-local revset alias %s is invalid: %w", Alias, err) + } + valid := fmt.Sprintf("exactly(%s & root()+ & subject(exact:%s), 1)", Alias, strconv.Quote(Description)) + if _, err := s.logValue(valid, "change_id"); err != nil { + return "", false, true, fmt.Errorf("repo-local revset alias %s=%q does not name vr's sidecar: %w", + Alias, strings.TrimSpace(configuredValue), err) + } + return rev, true, true, nil + } + + ids, err := s.descriptionMatches() + if err != nil { + return "", false, false, err + } + switch len(ids) { + case 0: + return "", false, false, nil + case 1: + return ids[0], true, false, nil + default: + return "", false, false, fmt.Errorf("found %d root children described %q; cannot choose a vr sidecar", len(ids), Description) + } +} + +func (s *Store) descriptionMatches() ([]string, error) { + revset := fmt.Sprintf("root()+ & subject(exact:%s)", strconv.Quote(Description)) + out, err := s.jj("log", "-G", "-r", revset, "-T", "change_id ++ \"\\n\"") + if err != nil { + return nil, err + } + var ids []string + for _, id := range strings.Fields(out) { + ids = append(ids, id) + } + return ids, nil +} + +func (s *Store) logValue(revset, template string) (string, error) { + out, err := s.jj("log", "-G", "-r", revset, "-T", template) + if err != nil { + return "", err + } + value := strings.TrimSpace(out) + if value == "" { + return "", fmt.Errorf("revision %q produced no %s", revset, template) + } + return value, nil +} + +func (s *Store) readRevision(rev string) ([]byte, error) { + listed, err := s.jj("file", "list", "-r", rev, "--", LogFile) + if err != nil { + return nil, err + } + if strings.TrimSpace(listed) == "" { + return nil, nil + } + out, err := s.jj("file", "show", "-r", rev, "--", LogFile) + return []byte(out), err +} + +func (s *Store) runAppend(data []byte) error { + f, err := os.CreateTemp(filepath.Dir(s.lockPath), "vr-log-append-*") + if err != nil { + return err + } + tmp := f.Name() + defer os.Remove(tmp) //nolint:errcheck + if _, err := f.Write(data); err != nil { + f.Close() + return err + } + if err := f.Close(); err != nil { + return err + } + _, err = s.jj("--config", "snapshot.auto-track=\"all()\"", + "--config", "snapshot.max-new-file-size=\"1GiB\"", + "run", "--clean", "--root", "-r", "exactly("+Alias+", 1)", "--", + "sh", "-c", "cat \"$1\" >> "+LogFile, "sh", tmp) + return err +} + +func (s *Store) jj(args ...string) (string, error) { + return jjOutput(s.root, args...) +} + +func jjOutput(dir string, args ...string) (string, error) { + full := append([]string{"--ignore-working-copy", "--no-pager"}, args...) + cmd := exec.Command("jj", full...) + cmd.Dir = dir + out, err := cmd.Output() + if err != nil { + message := strings.TrimSpace(string(out)) + if ee, ok := err.(*exec.ExitError); ok { + message = strings.TrimSpace(string(ee.Stderr)) + } + if message != "" { + return "", fmt.Errorf("jj %s: %w: %s", strings.Join(args, " "), err, message) + } + return "", fmt.Errorf("jj %s: %w", strings.Join(args, " "), err) + } + return string(out), nil +} + +func containsRevsetName(expr, name string) bool { + for i := 0; i+len(name) <= len(expr); i++ { + if expr[i:i+len(name)] != name { + continue + } + before := i == 0 || !isNameByte(expr[i-1]) + after := i+len(name) == len(expr) || !isNameByte(expr[i+len(name)]) + if before && after { + return true + } + } + return false +} + +func isNameByte(b byte) bool { + return b == '_' || b == '-' || b >= 'a' && b <= 'z' || b >= 'A' && b <= 'Z' || b >= '0' && b <= '9' +} diff --git a/sidecar/store_test.go b/sidecar/store_test.go new file mode 100644 index 0000000..85f50e1 --- /dev/null +++ b/sidecar/store_test.go @@ -0,0 +1,164 @@ +package sidecar + +import ( + "fmt" + "os" + "os/exec" + "path/filepath" + "strings" + "sync" + "testing" +) + +func TestAppendCreatesPrivateUnrelatedSidecar(t *testing.T) { + repo := testRepo(t) + // Store must be independent of a user's restrictive auto-track setting. + testJJ(t, repo, "config", "set", "--repo", "snapshot.auto-track", `"none()"`) + testJJ(t, repo, "config", "set", "--repo", "git.private-commits", `"description(glob:'secret:*')"`) + before := testJJ(t, repo, "status") + store, err := Open(repo) + if err != nil { + t.Fatal(err) + } + + if err := store.Append(func(current []byte) ([]byte, error) { + if len(current) != 0 { + t.Fatalf("new store contained %q", current) + } + return []byte("one\n"), nil + }); err != nil { + t.Fatal(err) + } + if err := store.Append(func(current []byte) ([]byte, error) { + if string(current) != "one\n" { + t.Fatalf("second append saw %q", current) + } + return []byte("two\n"), nil + }); err != nil { + t.Fatal(err) + } + + got, err := store.Read() + if err != nil { + t.Fatal(err) + } + if string(got) != "one\ntwo\n" { + t.Fatalf("sidecar log = %q", got) + } + if after := testJJ(t, repo, "status"); after != before { + t.Fatalf("working copy changed:\n--- before\n%s--- after\n%s", before, after) + } + if _, err := os.Stat(filepath.Join(repo, LogFile)); !os.IsNotExist(err) { + t.Fatalf("log materialized in the project working copy: %v", err) + } + + alias := strings.TrimSpace(testJJ(t, repo, "config", "get", "revset-aliases."+Alias)) + if alias == "" { + t.Fatal("repo-local sidecar alias was not configured") + } + private := testJJ(t, repo, "config", "get", "git.private-commits") + if !containsRevsetName(private, Alias) { + t.Fatalf("private commits = %q, missing %s", private, Alias) + } + if !strings.Contains(private, "secret:*") { + t.Fatalf("private commits = %q, previous setting was overwritten", private) + } + parent := strings.TrimSpace(testJJ(t, repo, "log", "-G", "-r", "exactly("+Alias+", 1)", + "-T", "parents.map(|p| p.commit_id()).join(\"\")")) + if parent != strings.Repeat("0", 40) { + t.Fatalf("sidecar parent = %q, want virtual root", parent) + } +} + +func TestLegacyLogIsReadThenImportedOnFirstAppend(t *testing.T) { + repo := testRepo(t) + legacy := []byte("legacy\n") + if err := os.WriteFile(filepath.Join(repo, LogFile), legacy, 0o644); err != nil { + t.Fatal(err) + } + store, err := Open(repo) + if err != nil { + t.Fatal(err) + } + if got, err := store.Read(); err != nil || string(got) != string(legacy) { + t.Fatalf("pre-migration read = %q, %v", got, err) + } + if err := store.Append(func(current []byte) ([]byte, error) { + if string(current) != string(legacy) { + t.Fatalf("migration callback saw %q", current) + } + return []byte("sidecar\n"), nil + }); err != nil { + t.Fatal(err) + } + if got, err := store.Read(); err != nil || string(got) != "legacy\nsidecar\n" { + t.Fatalf("migrated log = %q, %v", got, err) + } +} + +func TestConcurrentAppendsAreSerialized(t *testing.T) { + repo := testRepo(t) + store, err := Open(repo) + if err != nil { + t.Fatal(err) + } + const writers = 6 + errCh := make(chan error, writers) + var wg sync.WaitGroup + for range writers { + wg.Add(1) + go func() { + defer wg.Done() + errCh <- store.Append(func(current []byte) ([]byte, error) { + id := len(strings.Fields(string(current))) + 1 + return []byte(fmt.Sprintf("%d\n", id)), nil + }) + }() + } + wg.Wait() + close(errCh) + for err := range errCh { + if err != nil { + t.Fatal(err) + } + } + got, err := store.Read() + if err != nil { + t.Fatal(err) + } + if string(got) != "1\n2\n3\n4\n5\n6\n" { + t.Fatalf("serialized appends = %q", got) + } +} + +func testRepo(t *testing.T) string { + t.Helper() + if _, err := exec.LookPath("jj"); err != nil { + t.Skip("jj not installed") + } + home := t.TempDir() + t.Setenv("HOME", home) + t.Setenv("JJ_USER", "test") + t.Setenv("JJ_EMAIL", "[email protected]") + repo := filepath.Join(home, "repo") + if err := os.Mkdir(repo, 0o755); err != nil { + t.Fatal(err) + } + cmd := exec.Command("jj", "git", "init") + cmd.Dir = repo + if out, err := cmd.CombinedOutput(); err != nil { + t.Fatalf("jj git init: %v\n%s", err, out) + } + return repo +} + +func testJJ(t *testing.T, repo string, args ...string) string { + t.Helper() + cmd := exec.Command("jj", append([]string{"--ignore-working-copy", "--no-pager"}, args...)...) + cmd.Dir = repo + out, err := cmd.CombinedOutput() + if err != nil { + t.Fatalf("jj %s: %v\n%s", strings.Join(args, " "), err, out) + } + return string(out) +} |
