From 1e8d052f16a07defb8c07fc5be8a6fafab5ad2e0 Mon Sep 17 00:00:00 2001 From: Piyush Sachdeva Date: Tue, 23 Jun 2026 16:51:52 +0000 Subject: [PATCH 1/3] Upload Claude memory files in amikalog beta:push beta:push now uploads the Claude memory files (~/.claude/projects//memory/*.md) of every project amikalog has captured Claude sessions for, under the same "/" key prefix as that project's sessions ("/memory/.md"). Memory files are edited in place, unlike append-only session JSONL, so they can diverge between this machine and the cloud copy other machines pushed. A dedicated manifest records each file's last-synced content hash and a 3-way rule decides per file: no cloud copy or only-local change uploads, only-cloud change pulls the cloud copy down, identical skips, and a file changed on both sides is merged with the host claude CLI, written back locally, and uploaded. A failed or unavailable merge never clobbers either side; if claude is not installed, diverged files are skipped with a warning while the rest upload. Adds eventlog.Downloader, apiclient.GetObjectByKey for single-object reads, and a --skip-memories flag. --- go/cmd/amikalog/push.go | 59 +++- go/internal/apiclient/client.go | 29 ++ go/internal/eventlog/memory.go | 450 ++++++++++++++++++++++++++++ go/internal/eventlog/memory_test.go | 303 +++++++++++++++++++ go/internal/eventlog/push.go | 29 ++ 5 files changed, 863 insertions(+), 7 deletions(-) create mode 100644 go/internal/eventlog/memory.go create mode 100644 go/internal/eventlog/memory_test.go diff --git a/go/cmd/amikalog/push.go b/go/cmd/amikalog/push.go index d244e006..814fa7b0 100644 --- a/go/cmd/amikalog/push.go +++ b/go/cmd/amikalog/push.go @@ -10,12 +10,21 @@ import ( "github.com/spf13/cobra" ) +// skipMemories disables uploading Claude memory files alongside captured events. +var skipMemories bool + var pushCmd = &cobra.Command{ Use: "beta:push", Short: "Upload captured events to your organization", Long: `Upload captured events that have not been pushed yet. Repeated runs upload only events captured since the last push. +Claude memory files (~/.claude/projects//memory/*.md) for the projects +you have captured sessions for are uploaded too. Because memory files are edited +in place, a file that changed both locally and in the cloud is merged with the +local claude CLI rather than overwritten; pass --skip-memories to upload only +events. + Set AMIKA_API_KEY to authenticate.`, Args: cobra.NoArgs, SilenceUsage: true, @@ -31,18 +40,42 @@ Set AMIKA_API_KEY to authenticate.`, } client := apiclient.NewClientWithTokenSource(config.APIURL(), apiclient.NewStaticTokenSource(key)) - report, err := eventlog.Push(stateDir, apiUploader{client: client}) + uploader := apiUploader{client: client} + out := cmd.OutOrStdout() + errOut := cmd.ErrOrStderr() + + report, err := eventlog.Push(stateDir, uploader) if err != nil { return err } - - out := cmd.OutOrStdout() fmt.Fprintf(out, "uploaded %d, skipped %d, failed %d\n", report.Uploaded, report.Skipped, report.Failed) - if report.Failed > 0 { - for _, e := range report.Errors { - fmt.Fprintf(cmd.ErrOrStderr(), "amikalog: %v\n", e) + for _, e := range report.Errors { + fmt.Fprintf(errOut, "amikalog: %v\n", e) + } + failed := report.Failed + + if !skipMemories { + home, herr := os.UserHomeDir() + if herr != nil { + return fmt.Errorf("resolving home directory: %w", herr) + } + mreport, merr := eventlog.PushMemories(stateDir, home, uploader, apiDownloader{client: client}, eventlog.NewClaudeMerger()) + if merr != nil { + return fmt.Errorf("pushing memories: %w", merr) } - return fmt.Errorf("%d file(s) failed to upload", report.Failed) + fmt.Fprintf(out, "memories: uploaded %d, merged %d, pulled %d, skipped %d, failed %d\n", + mreport.Uploaded, mreport.Merged, mreport.Pulled, mreport.Skipped, mreport.Failed) + for _, w := range mreport.Warnings { + fmt.Fprintf(errOut, "amikalog: %s\n", w) + } + for _, e := range mreport.Errors { + fmt.Fprintf(errOut, "amikalog: %v\n", e) + } + failed += mreport.Failed + } + + if failed > 0 { + return fmt.Errorf("%d file(s) failed to upload", failed) } return nil }, @@ -72,6 +105,18 @@ func (a apiUploader) Upload(objectKey string, data []byte) error { return a.client.UploadToSignedURL(resp.Objects[0].UploadURL, data, "application/json") } +// apiDownloader adapts the Amika API client to eventlog.Downloader: it fetches a +// single object's current bytes by key, used to detect whether a memory file's +// cloud copy has diverged before overwriting it. +type apiDownloader struct { + client *apiclient.Client +} + +func (a apiDownloader) Fetch(objectKey string) ([]byte, bool, error) { + return a.client.GetObjectByKey(objectKey) +} + func init() { + pushCmd.Flags().BoolVar(&skipMemories, "skip-memories", false, "Do not upload Claude memory files, only captured events") rootCmd.AddCommand(pushCmd) } diff --git a/go/internal/apiclient/client.go b/go/internal/apiclient/client.go index 2eccb53c..4ebd1941 100644 --- a/go/internal/apiclient/client.go +++ b/go/internal/apiclient/client.go @@ -605,6 +605,35 @@ func (c *Client) DownloadFromSignedURL(signedURL string) ([]byte, error) { return body, nil } +// GetObjectByKey fetches the current bytes of a single object by its exact +// bucket key. There is no single-object endpoint, so it lists the subtree at +// key (a prefix listing returns the object under its own key, possibly +// alongside keys that share it as a prefix) and downloads the entry whose Key +// matches exactly. found is false when no object has that exact key, which the +// caller can treat as "no cloud copy yet". +func (c *Client) GetObjectByKey(key string) (data []byte, found bool, err error) { + cursor := "" + for { + resp, err := c.ListDownloads(key, cursor, 0) + if err != nil { + return nil, false, err + } + for _, o := range resp.Objects { + if o.Key == key { + b, err := c.DownloadFromSignedURL(o.DownloadURL) + if err != nil { + return nil, false, err + } + return b, true, nil + } + } + if resp.NextCursor == nil || *resp.NextCursor == "" { + return nil, false, nil + } + cursor = *resp.NextCursor + } +} + func (c *Client) doJSON(method, path string, body interface{}, out interface{}) error { var bodyReader io.Reader if body != nil { diff --git a/go/internal/eventlog/memory.go b/go/internal/eventlog/memory.go new file mode 100644 index 00000000..82bab98e --- /dev/null +++ b/go/internal/eventlog/memory.go @@ -0,0 +1,450 @@ +package eventlog + +import ( + "bytes" + "context" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "io/fs" + "os" + "os/exec" + "path" + "path/filepath" + "sort" + "strings" + "time" +) + +// memoryManifestName is the file under /events that records the +// last-synced content hash of each uploaded memory file, kept separate from the +// session push manifest because memory files are edited in place (so they are +// tracked by content hash, not by the monotonically growing size that the +// append-only session manifest assumes). +const memoryManifestName = ".amikalog-memory-push-state.json" + +// memorySegment is the object-key segment under a repository prefix that holds +// that repository's memory files, mirroring the "sessions" segment used for +// captured events. +const memorySegment = "memory" + +// mergeTimeout bounds a single claude merge invocation. +const mergeTimeout = 5 * time.Minute + +// ErrMergerUnavailable is returned by a Merger when the backing coding agent is +// not available (e.g. the claude CLI is not installed on this host). PushMemories +// treats it as a reason to skip the file with a warning rather than fail the run, +// so a machine without claude still uploads non-diverged memory files. +var ErrMergerUnavailable = errors.New("merger unavailable") + +// Merger reconciles a memory file that changed on both this machine and in the +// cloud into a single merged document. Implementations should preserve every +// distinct fact from both sides; PushMemories writes the result back to the +// local file and uploads it. +type Merger interface { + Merge(local, cloud []byte) ([]byte, error) +} + +// MemoryPushReport summarizes a PushMemories run. +type MemoryPushReport struct { + // Uploaded is the number of memory files sent to the bucket because they were + // new or changed only locally since the last sync. + Uploaded int + // Merged is the number of memory files that had diverged on both sides and + // were reconciled by the Merger, then written back locally and uploaded. + Merged int + // Pulled is the number of memory files that changed only in the cloud and + // were written back to the local file (not uploaded). + Pulled int + // Skipped is the number of memory files already in sync (or skipped because + // a merge was needed but the Merger was unavailable). + Skipped int + // Failed is the number of memory files whose reconciliation returned an error. + Failed int + // Warnings holds non-fatal messages (e.g. a merge skipped because claude is + // not installed). + Warnings []string + // Errors holds one error per failed file (parallel to Failed). + Errors []error +} + +// memoryEntry records what was last synced for one memory file. +type memoryEntry struct { + // ObjectKey is the destination bucket key, pinned so it never changes. + ObjectKey string `json:"object_key"` + // SyncedHash is the sha256 (hex) of the content as of the last successful + // sync. It is the base for 3-way divergence detection: if the cloud copy + // still hashes to this, the cloud is unchanged since we last synced and the + // local copy can be uploaded; if it differs, both sides changed. + SyncedHash string `json:"synced_hash"` +} + +// memoryManifest tracks synced memory files keyed by their object key. +type memoryManifest struct { + Synced map[string]memoryEntry `json:"synced"` +} + +// memoryUnit is one memory file PushMemories may reconcile. +type memoryUnit struct { + // filePath is the absolute on-disk path of the memory file. + filePath string + // objectKey is the destination bucket key, lowercased. + objectKey string +} + +// memoryOutcome is the result of reconciling one memory file. +type memoryOutcome int + +const ( + outcomeSkipped memoryOutcome = iota + outcomeUploaded + outcomePulled + outcomeMerged +) + +// claudeProjectDirName returns the directory name Claude Code uses for a project +// working directory under ~/.claude/projects. Claude Code derives it by +// replacing every character that is not an ASCII letter or digit with '-', so +// "/home/u/my.app" becomes "-home-u-my-app". The mapping is lossy (a path +// containing '-' or '.' cannot be reversed unambiguously), but we only ever go +// cwd -> name, which is deterministic. +func claudeProjectDirName(cwd string) string { + return strings.Map(func(r rune) rune { + switch { + case r >= 'a' && r <= 'z', r >= 'A' && r <= 'Z', r >= '0' && r <= '9': + return r + default: + return '-' + } + }, cwd) +} + +// PushMemories uploads the Claude memory files of every project amikalog has +// captured Claude sessions for, reconciling each against its cloud copy so an +// in-place edit on this machine never clobbers an edit made elsewhere. +// +// A memory file lives at ~/.claude/projects//memory/*.md, where +// is derived from the session's working directory (claudeProjectDirName). +// Its object key is "/memory/", sharing the repository prefix of +// that project's captured sessions. For each file PushMemories compares the +// local content, the cloud content, and the last-synced hash recorded in a +// dedicated manifest: +// - no cloud copy -> upload local +// - identical -> skip +// - only local changed -> upload local +// - only cloud changed -> pull (write cloud to the local file) +// - both changed -> merge with the Merger, write back, and upload +// +// The run-wide push lock is held for the duration (the same lock Push takes), so +// concurrent pushes cannot interleave overwrites. Per-file failures are recorded +// in the report and do not abort the run. +func PushMemories(stateDir, home string, up Uploader, down Downloader, merger Merger) (MemoryPushReport, error) { + eventsBase := filepath.Join(stateDir, "events") + if err := os.MkdirAll(eventsBase, 0o755); err != nil { + return MemoryPushReport{}, fmt.Errorf("creating events dir %s: %w", eventsBase, err) + } + + lock, err := acquireLock(filepath.Join(eventsBase, pushLockName)) + if err != nil { + return MemoryPushReport{}, err + } + defer lock.release() + + manifestPath := filepath.Join(eventsBase, memoryManifestName) + manifest, err := loadMemoryManifest(manifestPath) + if err != nil { + return MemoryPushReport{}, err + } + + units, err := collectMemoryUnits(stateDir, home) + if err != nil { + return MemoryPushReport{}, err + } + + var report MemoryPushReport + for _, u := range units { + outcome, rerr := reconcileMemory(u, manifest, up, down, merger) + if rerr != nil { + if errors.Is(rerr, ErrMergerUnavailable) { + report.Skipped++ + report.Warnings = append(report.Warnings, fmt.Sprintf("%s: %v", u.objectKey, rerr)) + continue + } + report.Failed++ + report.Errors = append(report.Errors, rerr) + continue + } + switch outcome { + case outcomeUploaded: + report.Uploaded++ + case outcomeMerged: + report.Merged++ + case outcomePulled: + report.Pulled++ + default: + report.Skipped++ + } + // Persist after every reconciled file so an interrupted run (e.g. a slow + // merge that is killed) resumes without re-doing completed work. + if err := saveMemoryManifest(manifestPath, manifest); err != nil { + return report, err + } + } + return report, nil +} + +// collectMemoryUnits lists the memory files to reconcile: one per *.md file in +// the memory directory of each project amikalog has captured a Claude session +// for. The repository prefix for each project is taken from its sessions so a +// repo's memory files land alongside its events under the same "/" key. +func collectMemoryUnits(stateDir, home string) ([]memoryUnit, error) { + sessionsRoot := EventsDir(stateDir, SourceClaude) + entries, err := os.ReadDir(sessionsRoot) + if err != nil { + if os.IsNotExist(err) { + return nil, nil + } + return nil, fmt.Errorf("reading claude sessions: %w", err) + } + + // Map each captured project working directory to its repository segment. + repoByCwd := map[string]string{} + for _, e := range entries { + if e.IsDir() || !strings.HasSuffix(e.Name(), ".jsonl") { + continue + } + data, err := os.ReadFile(filepath.Join(sessionsRoot, e.Name())) + if err != nil { + return nil, fmt.Errorf("reading %s: %w", e.Name(), err) + } + cwd := cwdFromJSONL(data) + if cwd == "" { + continue + } + if _, ok := repoByCwd[cwd]; !ok { + repoByCwd[cwd] = repoSegmentFromJSONL(data) + } + } + + cwds := make([]string, 0, len(repoByCwd)) + for cwd := range repoByCwd { + cwds = append(cwds, cwd) + } + sort.Strings(cwds) + + var units []memoryUnit + seen := map[string]bool{} + for _, cwd := range cwds { + repoSeg := repoByCwd[cwd] + memoryDir := filepath.Join(home, ".claude", "projects", claudeProjectDirName(cwd), "memory") + walkErr := filepath.WalkDir(memoryDir, func(p string, d fs.DirEntry, err error) error { + if err != nil { + // The memory directory simply may not exist for this project. + if os.IsNotExist(err) { + return fs.SkipAll + } + return err + } + if d.IsDir() || !strings.HasSuffix(strings.ToLower(d.Name()), ".md") { + return nil + } + rel, err := relSlash(memoryDir, p) + if err != nil { + return err + } + objectKey := strings.ToLower(path.Join(repoSeg, memorySegment, rel)) + // Distinct projects can map to the same repo segment (e.g. a repo and + // a subdirectory of it); the first occurrence of a key wins. + if seen[objectKey] { + return nil + } + seen[objectKey] = true + units = append(units, memoryUnit{filePath: p, objectKey: objectKey}) + return nil + }) + if walkErr != nil { + return nil, fmt.Errorf("scanning memory dir %s: %w", memoryDir, walkErr) + } + } + return units, nil +} + +// reconcileMemory reconciles one memory file against its cloud copy using the +// last-synced hash in manifest, mutating manifest on success. See PushMemories +// for the decision table. +func reconcileMemory(u memoryUnit, manifest *memoryManifest, up Uploader, down Downloader, merger Merger) (memoryOutcome, error) { + localBytes, err := os.ReadFile(u.filePath) + if err != nil { + return outcomeSkipped, fmt.Errorf("reading %s: %w", u.filePath, err) + } + localHash := hashBytes(localBytes) + entry, known := manifest.Synced[u.objectKey] + + cloudBytes, cloudExists, err := down.Fetch(u.objectKey) + if err != nil { + return outcomeSkipped, fmt.Errorf("fetching %s: %w", u.objectKey, err) + } + + // No cloud copy yet: upload local as the initial version. + if !cloudExists { + if err := up.Upload(u.objectKey, localBytes); err != nil { + return outcomeSkipped, fmt.Errorf("uploading %s: %w", u.objectKey, err) + } + manifest.Synced[u.objectKey] = memoryEntry{ObjectKey: u.objectKey, SyncedHash: localHash} + return outcomeUploaded, nil + } + + cloudHash := hashBytes(cloudBytes) + + // Identical: nothing to send. Refresh the manifest so a first sync after the + // content already matched still records the base for future runs. + if cloudHash == localHash { + if !known || entry.SyncedHash != localHash { + manifest.Synced[u.objectKey] = memoryEntry{ObjectKey: u.objectKey, SyncedHash: localHash} + } + return outcomeSkipped, nil + } + + // Cloud unchanged since our last sync: local is the only change -> upload it. + if known && entry.SyncedHash == cloudHash { + if err := up.Upload(u.objectKey, localBytes); err != nil { + return outcomeSkipped, fmt.Errorf("uploading %s: %w", u.objectKey, err) + } + manifest.Synced[u.objectKey] = memoryEntry{ObjectKey: u.objectKey, SyncedHash: localHash} + return outcomeUploaded, nil + } + + // Local unchanged since our last sync: only the cloud changed -> pull it down. + if known && entry.SyncedHash == localHash { + if err := writeFileAtomic(u.filePath, cloudBytes); err != nil { + return outcomeSkipped, fmt.Errorf("writing %s: %w", u.filePath, err) + } + manifest.Synced[u.objectKey] = memoryEntry{ObjectKey: u.objectKey, SyncedHash: cloudHash} + return outcomePulled, nil + } + + // Both sides changed (or this is the first sync against a pre-existing cloud + // copy): merge with the coding agent, then converge both sides. On any merge + // failure neither side is touched, so a bad merge can never clobber data. + merged, err := merger.Merge(localBytes, cloudBytes) + if err != nil { + return outcomeSkipped, fmt.Errorf("merging %s: %w", u.objectKey, err) + } + if err := writeFileAtomic(u.filePath, merged); err != nil { + return outcomeSkipped, fmt.Errorf("writing merged %s: %w", u.filePath, err) + } + if err := up.Upload(u.objectKey, merged); err != nil { + return outcomeSkipped, fmt.Errorf("uploading merged %s: %w", u.objectKey, err) + } + manifest.Synced[u.objectKey] = memoryEntry{ObjectKey: u.objectKey, SyncedHash: hashBytes(merged)} + return outcomeMerged, nil +} + +// hashBytes returns the hex-encoded sha256 of b, used as a memory file's content +// fingerprint. +func hashBytes(b []byte) string { + sum := sha256.Sum256(b) + return hex.EncodeToString(sum[:]) +} + +// loadMemoryManifest reads the memory push manifest, returning an empty one when +// the file does not yet exist. +func loadMemoryManifest(path string) (*memoryManifest, error) { + data, err := os.ReadFile(path) + if err != nil { + if os.IsNotExist(err) { + return &memoryManifest{Synced: map[string]memoryEntry{}}, nil + } + return nil, fmt.Errorf("reading memory push manifest %s: %w", path, err) + } + var m memoryManifest + if err := json.Unmarshal(data, &m); err != nil { + return nil, fmt.Errorf("parsing memory push manifest %s: %w", path, err) + } + if m.Synced == nil { + m.Synced = map[string]memoryEntry{} + } + return &m, nil +} + +// saveMemoryManifest writes the manifest atomically (write-then-rename) so an +// interrupted write cannot corrupt the existing manifest. +func saveMemoryManifest(path string, m *memoryManifest) error { + data, err := json.MarshalIndent(m, "", " ") + if err != nil { + return err + } + tmp := path + ".tmp" + if err := os.WriteFile(tmp, data, 0o644); err != nil { + return fmt.Errorf("writing memory push manifest: %w", err) + } + if err := os.Rename(tmp, path); err != nil { + return fmt.Errorf("replacing memory push manifest: %w", err) + } + return nil +} + +// claudeMerger merges two versions of a memory file by invoking the locally +// installed claude CLI in one-shot print mode. +type claudeMerger struct{} + +// NewClaudeMerger returns a Merger backed by the host claude CLI. It returns +// ErrMergerUnavailable from Merge when the claude binary is not on PATH. +func NewClaudeMerger() Merger { return claudeMerger{} } + +// Merge runs `claude -p --output-format json` with both versions +// embedded in the prompt and returns the reconciled file content. It requires +// the claude CLI on PATH (ErrMergerUnavailable otherwise) and asks the model to +// emit only the merged file content. +func (claudeMerger) Merge(local, cloud []byte) ([]byte, error) { + exe, err := exec.LookPath("claude") + if err != nil { + return nil, fmt.Errorf("%w: claude CLI not found on PATH", ErrMergerUnavailable) + } + + ctx, cancel := context.WithTimeout(context.Background(), mergeTimeout) + defer cancel() + + cmd := exec.CommandContext(ctx, exe, "-p", buildMergePrompt(local, cloud), + "--output-format", "json", "--dangerously-skip-permissions") + var stdout, stderr bytes.Buffer + cmd.Stdout = &stdout + cmd.Stderr = &stderr + if err := cmd.Run(); err != nil { + return nil, fmt.Errorf("running claude merge: %w (stderr: %s)", err, strings.TrimSpace(stderr.String())) + } + + var out struct { + Result string `json:"result"` + IsError bool `json:"is_error"` + } + if err := json.Unmarshal(stdout.Bytes(), &out); err != nil { + return nil, fmt.Errorf("parsing claude output: %w", err) + } + if out.IsError { + return nil, fmt.Errorf("claude reported a merge error: %s", strings.TrimSpace(out.Result)) + } + merged := strings.TrimSpace(out.Result) + if merged == "" { + return nil, fmt.Errorf("claude returned an empty merge result") + } + return []byte(merged + "\n"), nil +} + +// buildMergePrompt constructs the one-shot prompt instructing the agent to +// reconcile the two versions into a single markdown document. +func buildMergePrompt(local, cloud []byte) string { + var b strings.Builder + b.WriteString("You are merging two versions of a Markdown memory file that have diverged. ") + b.WriteString("Produce a single reconciled version that preserves every distinct fact from BOTH versions, ") + b.WriteString("removes exact duplicates, and keeps one valid YAML frontmatter block if either version has one. ") + b.WriteString("Do not invent new facts. Output ONLY the merged file content, with no commentary, no explanation, and no code fences.\n\n") + b.WriteString("===== VERSION A (local) =====\n") + b.Write(local) + b.WriteString("\n===== VERSION B (cloud) =====\n") + b.Write(cloud) + b.WriteString("\n===== END =====\n") + return b.String() +} diff --git a/go/internal/eventlog/memory_test.go b/go/internal/eventlog/memory_test.go new file mode 100644 index 00000000..5e661027 --- /dev/null +++ b/go/internal/eventlog/memory_test.go @@ -0,0 +1,303 @@ +package eventlog + +import ( + "encoding/json" + "errors" + "os" + "path/filepath" + "sort" + "testing" +) + +// fakeDownloader serves canned cloud bytes by object key for reconcile tests. +type fakeDownloader struct { + data map[string][]byte + err error +} + +func (f *fakeDownloader) Fetch(key string) ([]byte, bool, error) { + if f.err != nil { + return nil, false, f.err + } + b, ok := f.data[key] + return b, ok, nil +} + +// fakeMerger returns a canned merge result (or error) and counts invocations. +type fakeMerger struct { + result []byte + err error + calls int +} + +func (f *fakeMerger) Merge(_, _ []byte) ([]byte, error) { + f.calls++ + if f.err != nil { + return nil, f.err + } + return f.result, nil +} + +// eventLineCWD renders one event line carrying a working directory and git info, +// as a Claude session's first line would. +func eventLineCWD(t *testing.T, cwd string, git *GitInfo) string { + t.Helper() + b, err := json.Marshal(Event{ + Source: SourceClaude, + HookEvent: "SessionStart", + SessionID: "sess", + CWD: cwd, + Git: git, + Payload: json.RawMessage(`{}`), + }) + if err != nil { + t.Fatalf("marshal: %v", err) + } + return string(b) +} + +func TestClaudeProjectDirName(t *testing.T) { + cases := map[string]string{ + "/home/amika/workspace/amika": "-home-amika-workspace-amika", + "/home/u/my.app": "-home-u-my-app", + "/work/a_b/c-d": "-work-a-b-c-d", + "relative/path": "relative-path", + } + for in, want := range cases { + if got := claudeProjectDirName(in); got != want { + t.Errorf("claudeProjectDirName(%q) = %q, want %q", in, got, want) + } + } +} + +func TestCollectMemoryUnits(t *testing.T) { + stateDir := t.TempDir() + home := t.TempDir() + cwd := "/work/myrepo" + + writeSessionFile(t, stateDir, SourceClaude, "20240101T000000.000000000Z_sess-a.jsonl", + eventLineCWD(t, cwd, &GitInfo{RepoRoot: "/work/myrepo"})) + + memDir := filepath.Join(home, ".claude", "projects", "-work-myrepo", "memory") + if err := os.MkdirAll(memDir, 0o755); err != nil { + t.Fatalf("mkdir: %v", err) + } + for name, body := range map[string]string{ + "foo.md": "x", + "MEMORY.md": "y", // uppercase name -> lowercased key + "notes.txt": "z", // not markdown -> ignored + } { + if err := os.WriteFile(filepath.Join(memDir, name), []byte(body), 0o644); err != nil { + t.Fatalf("write %s: %v", name, err) + } + } + + units, err := collectMemoryUnits(stateDir, home) + if err != nil { + t.Fatalf("collectMemoryUnits: %v", err) + } + got := make([]string, 0, len(units)) + for _, u := range units { + got = append(got, u.objectKey) + } + sort.Strings(got) + want := []string{"myrepo/memory/foo.md", "myrepo/memory/memory.md"} + if len(got) != len(want) || got[0] != want[0] || got[1] != want[1] { + t.Fatalf("object keys = %v, want %v", got, want) + } +} + +// reconcileFixture sets up one local memory file plus the fakes for a single +// reconcileMemory call. +type reconcileFixture struct { + unit memoryUnit + manifest *memoryManifest + up *fakeUploader + down *fakeDownloader + merger *fakeMerger +} + +func newReconcileFixture(t *testing.T, key, local string) *reconcileFixture { + t.Helper() + path := filepath.Join(t.TempDir(), "fact.md") + if err := os.WriteFile(path, []byte(local), 0o644); err != nil { + t.Fatalf("write local: %v", err) + } + return &reconcileFixture{ + unit: memoryUnit{filePath: path, objectKey: key}, + manifest: &memoryManifest{Synced: map[string]memoryEntry{}}, + up: newFakeUploader(), + down: &fakeDownloader{data: map[string][]byte{}}, + merger: &fakeMerger{}, + } +} + +func (f *reconcileFixture) localContent(t *testing.T) string { + t.Helper() + b, err := os.ReadFile(f.unit.filePath) + if err != nil { + t.Fatalf("read local: %v", err) + } + return string(b) +} + +func TestReconcileMemory_NoCloudCopy_Uploads(t *testing.T) { + f := newReconcileFixture(t, "repo/memory/f.md", "local") + outcome, err := reconcileMemory(f.unit, f.manifest, f.up, f.down, f.merger) + if err != nil || outcome != outcomeUploaded { + t.Fatalf("outcome=%v err=%v, want uploaded", outcome, err) + } + if got := string(f.up.bytesFor("repo/memory/f.md")); got != "local" { + t.Fatalf("uploaded %q, want %q", got, "local") + } + if f.manifest.Synced["repo/memory/f.md"].SyncedHash != hashBytes([]byte("local")) { + t.Fatalf("manifest hash not set to local hash") + } +} + +func TestReconcileMemory_Identical_Skips(t *testing.T) { + f := newReconcileFixture(t, "repo/memory/f.md", "same") + f.down.data["repo/memory/f.md"] = []byte("same") + outcome, err := reconcileMemory(f.unit, f.manifest, f.up, f.down, f.merger) + if err != nil || outcome != outcomeSkipped { + t.Fatalf("outcome=%v err=%v, want skipped", outcome, err) + } + if len(f.up.keys()) != 0 { + t.Fatalf("unexpected upload: %v", f.up.keys()) + } +} + +func TestReconcileMemory_OnlyLocalChanged_Uploads(t *testing.T) { + f := newReconcileFixture(t, "repo/memory/f.md", "v2") + f.down.data["repo/memory/f.md"] = []byte("v1") + f.manifest.Synced["repo/memory/f.md"] = memoryEntry{ObjectKey: "repo/memory/f.md", SyncedHash: hashBytes([]byte("v1"))} + outcome, err := reconcileMemory(f.unit, f.manifest, f.up, f.down, f.merger) + if err != nil || outcome != outcomeUploaded { + t.Fatalf("outcome=%v err=%v, want uploaded", outcome, err) + } + if got := string(f.up.bytesFor("repo/memory/f.md")); got != "v2" { + t.Fatalf("uploaded %q, want v2", got) + } + if f.merger.calls != 0 { + t.Fatalf("merger called unexpectedly") + } +} + +func TestReconcileMemory_OnlyCloudChanged_Pulls(t *testing.T) { + f := newReconcileFixture(t, "repo/memory/f.md", "v1") + f.down.data["repo/memory/f.md"] = []byte("v2") + f.manifest.Synced["repo/memory/f.md"] = memoryEntry{ObjectKey: "repo/memory/f.md", SyncedHash: hashBytes([]byte("v1"))} + outcome, err := reconcileMemory(f.unit, f.manifest, f.up, f.down, f.merger) + if err != nil || outcome != outcomePulled { + t.Fatalf("outcome=%v err=%v, want pulled", outcome, err) + } + if got := f.localContent(t); got != "v2" { + t.Fatalf("local content %q, want v2 (pulled from cloud)", got) + } + if len(f.up.keys()) != 0 { + t.Fatalf("unexpected upload on pull: %v", f.up.keys()) + } + if f.manifest.Synced["repo/memory/f.md"].SyncedHash != hashBytes([]byte("v2")) { + t.Fatalf("manifest hash not advanced to cloud hash") + } +} + +func TestReconcileMemory_BothChanged_Merges(t *testing.T) { + f := newReconcileFixture(t, "repo/memory/f.md", "local-v2") + f.down.data["repo/memory/f.md"] = []byte("cloud-v2") + f.manifest.Synced["repo/memory/f.md"] = memoryEntry{ObjectKey: "repo/memory/f.md", SyncedHash: hashBytes([]byte("base-v1"))} + f.merger.result = []byte("merged\n") + + outcome, err := reconcileMemory(f.unit, f.manifest, f.up, f.down, f.merger) + if err != nil || outcome != outcomeMerged { + t.Fatalf("outcome=%v err=%v, want merged", outcome, err) + } + if f.merger.calls != 1 { + t.Fatalf("merger calls = %d, want 1", f.merger.calls) + } + if got := f.localContent(t); got != "merged\n" { + t.Fatalf("local content %q, want merged", got) + } + if got := string(f.up.bytesFor("repo/memory/f.md")); got != "merged\n" { + t.Fatalf("uploaded %q, want merged", got) + } + if f.manifest.Synced["repo/memory/f.md"].SyncedHash != hashBytes([]byte("merged\n")) { + t.Fatalf("manifest hash not set to merged hash") + } +} + +func TestReconcileMemory_MergerUnavailable_DoesNotClobber(t *testing.T) { + f := newReconcileFixture(t, "repo/memory/f.md", "local-v2") + f.down.data["repo/memory/f.md"] = []byte("cloud-v2") + f.manifest.Synced["repo/memory/f.md"] = memoryEntry{ObjectKey: "repo/memory/f.md", SyncedHash: hashBytes([]byte("base-v1"))} + f.merger.err = ErrMergerUnavailable + + _, err := reconcileMemory(f.unit, f.manifest, f.up, f.down, f.merger) + if !errors.Is(err, ErrMergerUnavailable) { + t.Fatalf("err = %v, want ErrMergerUnavailable", err) + } + if got := f.localContent(t); got != "local-v2" { + t.Fatalf("local content changed to %q on merge failure", got) + } + if len(f.up.keys()) != 0 { + t.Fatalf("unexpected upload on merge failure: %v", f.up.keys()) + } +} + +func TestReconcileMemory_MergeError_DoesNotClobber(t *testing.T) { + f := newReconcileFixture(t, "repo/memory/f.md", "local-v2") + f.down.data["repo/memory/f.md"] = []byte("cloud-v2") + f.manifest.Synced["repo/memory/f.md"] = memoryEntry{ObjectKey: "repo/memory/f.md", SyncedHash: hashBytes([]byte("base-v1"))} + f.merger.err = errors.New("boom") + + _, err := reconcileMemory(f.unit, f.manifest, f.up, f.down, f.merger) + if err == nil || errors.Is(err, ErrMergerUnavailable) { + t.Fatalf("err = %v, want a non-unavailable error", err) + } + if got := f.localContent(t); got != "local-v2" { + t.Fatalf("local content changed to %q on merge failure", got) + } + if len(f.up.keys()) != 0 { + t.Fatalf("unexpected upload on merge failure: %v", f.up.keys()) + } +} + +func TestPushMemories_UploadsThenSkipsAndWritesManifest(t *testing.T) { + stateDir := t.TempDir() + home := t.TempDir() + cwd := "/work/myrepo" + writeSessionFile(t, stateDir, SourceClaude, "20240101T000000.000000000Z_sess-a.jsonl", + eventLineCWD(t, cwd, &GitInfo{RepoRoot: "/work/myrepo"})) + memDir := filepath.Join(home, ".claude", "projects", "-work-myrepo", "memory") + if err := os.MkdirAll(memDir, 0o755); err != nil { + t.Fatalf("mkdir: %v", err) + } + if err := os.WriteFile(filepath.Join(memDir, "fact.md"), []byte("hello"), 0o644); err != nil { + t.Fatalf("write: %v", err) + } + + up := newFakeUploader() + down := &fakeDownloader{data: map[string][]byte{}} + merger := &fakeMerger{} + + rep, err := PushMemories(stateDir, home, up, down, merger) + if err != nil { + t.Fatalf("PushMemories: %v", err) + } + if rep.Uploaded != 1 || rep.Failed != 0 { + t.Fatalf("first run report = %+v, want 1 uploaded", rep) + } + if _, err := os.Stat(filepath.Join(stateDir, "events", memoryManifestName)); err != nil { + t.Fatalf("manifest not written: %v", err) + } + + // Second run: the cloud now holds what we uploaded, so it is a no-op skip. + down.data["myrepo/memory/fact.md"] = up.bytesFor("myrepo/memory/fact.md") + rep2, err := PushMemories(stateDir, home, up, down, merger) + if err != nil { + t.Fatalf("PushMemories second run: %v", err) + } + if rep2.Uploaded != 0 || rep2.Skipped != 1 { + t.Fatalf("second run report = %+v, want 1 skipped", rep2) + } +} diff --git a/go/internal/eventlog/push.go b/go/internal/eventlog/push.go index 3e734f89..a74c1175 100644 --- a/go/internal/eventlog/push.go +++ b/go/internal/eventlog/push.go @@ -20,6 +20,14 @@ type Uploader interface { Upload(objectKey string, data []byte) error } +// Downloader fetches the current bytes of one object from the destination +// bucket by its object key. exists is false when the bucket has no object at +// that key. It is used by PushMemories to detect whether a memory file's cloud +// copy has diverged from the local one before overwriting it. +type Downloader interface { + Fetch(objectKey string) (data []byte, exists bool, err error) +} + // PushReport summarizes a Push run. type PushReport struct { // Uploaded is the number of session files uploaded this run. @@ -438,6 +446,27 @@ func repoSegmentFromJSONL(data []byte) string { } } +// cwdFromJSONL returns the working directory recorded in the first line of a +// session's JSONL snapshot that carries one, or "" when no line records a cwd. +// Like repoSegmentFromJSONL the snapshot is taken under the lock, so this never +// sees a partial line. It is used to locate the Claude memory directory for the +// project the session ran in. +func cwdFromJSONL(data []byte) string { + r := bufio.NewReader(bytes.NewReader(data)) + for { + line, readErr := r.ReadBytes('\n') + if len(line) > 0 { + var ev Event + if json.Unmarshal(line, &ev) == nil && ev.CWD != "" { + return ev.CWD + } + } + if readErr != nil { + return "" + } + } +} + // resolveRepoSegment returns the sanitized repository basename for a legacy // session directory, read from the first event file that carries git context. // It returns "unknown-repo" when the session has no event with a git repo root. From 01a5ef09be5e7b4500300f05d561fdcb38048a99 Mon Sep 17 00:00:00 2001 From: Piyush Sachdeva Date: Tue, 23 Jun 2026 17:20:58 +0000 Subject: [PATCH 2/3] Make memory upload opt-in and add --all-projects Uploading Claude memory files is now off by default: beta:push uploads only captured events unless --memories is passed. This keeps memory upload behind an explicit flag while it is being tested, before considering it as a default. Add --all-projects (requires --memories) to also upload memory for projects that have no captured amikalog session. The repository prefix for such a project is recovered from its own Claude transcript working directory and that directory's git repo, falling back to "unknown-repo". Replaces the previous --skip-memories flag. --- go/cmd/amikalog/push.go | 28 +++-- go/internal/eventlog/memory.go | 173 +++++++++++++++++++++------- go/internal/eventlog/memory_test.go | 73 +++++++++++- 3 files changed, 218 insertions(+), 56 deletions(-) diff --git a/go/cmd/amikalog/push.go b/go/cmd/amikalog/push.go index 814fa7b0..c7f941af 100644 --- a/go/cmd/amikalog/push.go +++ b/go/cmd/amikalog/push.go @@ -10,8 +10,12 @@ import ( "github.com/spf13/cobra" ) -// skipMemories disables uploading Claude memory files alongside captured events. -var skipMemories bool +// uploadMemories opts into uploading Claude memory files alongside captured +// events. allMemoryProjects extends that to projects with no captured session. +var ( + uploadMemories bool + allMemoryProjects bool +) var pushCmd = &cobra.Command{ Use: "beta:push", @@ -19,17 +23,20 @@ var pushCmd = &cobra.Command{ Long: `Upload captured events that have not been pushed yet. Repeated runs upload only events captured since the last push. -Claude memory files (~/.claude/projects//memory/*.md) for the projects -you have captured sessions for are uploaded too. Because memory files are edited -in place, a file that changed both locally and in the cloud is merged with the -local claude CLI rather than overwritten; pass --skip-memories to upload only -events. +Pass --memories to also upload Claude memory files +(~/.claude/projects//memory/*.md) for the projects you have captured +sessions for. Because memory files are edited in place, a file that changed both +locally and in the cloud is merged with the local claude CLI rather than +overwritten. Add --all-projects to include projects with no captured session. Set AMIKA_API_KEY to authenticate.`, Args: cobra.NoArgs, SilenceUsage: true, SilenceErrors: true, RunE: func(cmd *cobra.Command, _ []string) error { + if allMemoryProjects && !uploadMemories { + return fmt.Errorf("--all-projects requires --memories") + } key := os.Getenv(config.EnvAPIKey) if key == "" { return fmt.Errorf("set %s to push; amikalog authenticates with an org API key only", config.EnvAPIKey) @@ -54,12 +61,12 @@ Set AMIKA_API_KEY to authenticate.`, } failed := report.Failed - if !skipMemories { + if uploadMemories { home, herr := os.UserHomeDir() if herr != nil { return fmt.Errorf("resolving home directory: %w", herr) } - mreport, merr := eventlog.PushMemories(stateDir, home, uploader, apiDownloader{client: client}, eventlog.NewClaudeMerger()) + mreport, merr := eventlog.PushMemories(stateDir, home, allMemoryProjects, uploader, apiDownloader{client: client}, eventlog.NewClaudeMerger()) if merr != nil { return fmt.Errorf("pushing memories: %w", merr) } @@ -117,6 +124,7 @@ func (a apiDownloader) Fetch(objectKey string) ([]byte, bool, error) { } func init() { - pushCmd.Flags().BoolVar(&skipMemories, "skip-memories", false, "Do not upload Claude memory files, only captured events") + pushCmd.Flags().BoolVar(&uploadMemories, "memories", false, "Also upload Claude memory files for projects with captured sessions") + pushCmd.Flags().BoolVar(&allMemoryProjects, "all-projects", false, "With --memories, include projects that have no captured session") rootCmd.AddCommand(pushCmd) } diff --git a/go/internal/eventlog/memory.go b/go/internal/eventlog/memory.go index 82bab98e..deadbf66 100644 --- a/go/internal/eventlog/memory.go +++ b/go/internal/eventlog/memory.go @@ -121,16 +121,20 @@ func claudeProjectDirName(cwd string) string { }, cwd) } -// PushMemories uploads the Claude memory files of every project amikalog has -// captured Claude sessions for, reconciling each against its cloud copy so an -// in-place edit on this machine never clobbers an edit made elsewhere. +// PushMemories uploads Claude memory files, reconciling each against its cloud +// copy so an in-place edit on this machine never clobbers an edit made elsewhere. // -// A memory file lives at ~/.claude/projects//memory/*.md, where -// is derived from the session's working directory (claudeProjectDirName). -// Its object key is "/memory/", sharing the repository prefix of -// that project's captured sessions. For each file PushMemories compares the -// local content, the cloud content, and the last-synced hash recorded in a -// dedicated manifest: +// By default only the projects amikalog has captured Claude sessions for are +// considered, so each file's repository prefix comes from its session. When +// allProjects is true every ~/.claude/projects/*/memory directory is scanned, +// including projects with no captured session; the repository prefix for those +// is recovered from the project's own Claude transcript working directory (and +// the git repo there), falling back to "unknown-repo". +// +// A memory file lives at ~/.claude/projects//memory/*.md and its object +// key is "/memory/". For each file PushMemories compares the local +// content, the cloud content, and the last-synced hash recorded in a dedicated +// manifest: // - no cloud copy -> upload local // - identical -> skip // - only local changed -> upload local @@ -140,7 +144,7 @@ func claudeProjectDirName(cwd string) string { // The run-wide push lock is held for the duration (the same lock Push takes), so // concurrent pushes cannot interleave overwrites. Per-file failures are recorded // in the report and do not abort the run. -func PushMemories(stateDir, home string, up Uploader, down Downloader, merger Merger) (MemoryPushReport, error) { +func PushMemories(stateDir, home string, allProjects bool, up Uploader, down Downloader, merger Merger) (MemoryPushReport, error) { eventsBase := filepath.Join(stateDir, "events") if err := os.MkdirAll(eventsBase, 0o755); err != nil { return MemoryPushReport{}, fmt.Errorf("creating events dir %s: %w", eventsBase, err) @@ -158,7 +162,7 @@ func PushMemories(stateDir, home string, up Uploader, down Downloader, merger Me return MemoryPushReport{}, err } - units, err := collectMemoryUnits(stateDir, home) + units, err := collectMemoryUnits(stateDir, home, allProjects) if err != nil { return MemoryPushReport{}, err } @@ -196,49 +200,55 @@ func PushMemories(stateDir, home string, up Uploader, down Downloader, merger Me } // collectMemoryUnits lists the memory files to reconcile: one per *.md file in -// the memory directory of each project amikalog has captured a Claude session -// for. The repository prefix for each project is taken from its sessions so a -// repo's memory files land alongside its events under the same "/" key. -func collectMemoryUnits(stateDir, home string) ([]memoryUnit, error) { - sessionsRoot := EventsDir(stateDir, SourceClaude) - entries, err := os.ReadDir(sessionsRoot) +// a project's memory directory. By default only projects amikalog has captured a +// Claude session for are scanned, each keyed under the repository prefix from its +// sessions so a repo's memory lands alongside its events under the same "/" +// key. When allProjects is true every ~/.claude/projects/*/memory directory is +// scanned; a project with no captured session has its repository prefix recovered +// from its own Claude transcript (see repoSegmentForProjectDir). +func collectMemoryUnits(stateDir, home string, allProjects bool) ([]memoryUnit, error) { + projectsRoot := filepath.Join(home, ".claude", "projects") + + // Seed each captured project directory with its repository segment, derived + // from amikalog's own session capture. + repoByProjectDir, err := trackedProjectRepoSegments(stateDir) if err != nil { - if os.IsNotExist(err) { - return nil, nil - } - return nil, fmt.Errorf("reading claude sessions: %w", err) + return nil, err } - // Map each captured project working directory to its repository segment. - repoByCwd := map[string]string{} - for _, e := range entries { - if e.IsDir() || !strings.HasSuffix(e.Name(), ".jsonl") { - continue - } - data, err := os.ReadFile(filepath.Join(sessionsRoot, e.Name())) + // Decide which project directories to scan. + var projectDirs []string + if allProjects { + entries, err := os.ReadDir(projectsRoot) if err != nil { - return nil, fmt.Errorf("reading %s: %w", e.Name(), err) + if os.IsNotExist(err) { + return nil, nil + } + return nil, fmt.Errorf("reading claude projects: %w", err) } - cwd := cwdFromJSONL(data) - if cwd == "" { - continue + for _, e := range entries { + if e.IsDir() { + projectDirs = append(projectDirs, e.Name()) + } } - if _, ok := repoByCwd[cwd]; !ok { - repoByCwd[cwd] = repoSegmentFromJSONL(data) + } else { + for dir := range repoByProjectDir { + projectDirs = append(projectDirs, dir) } } - - cwds := make([]string, 0, len(repoByCwd)) - for cwd := range repoByCwd { - cwds = append(cwds, cwd) - } - sort.Strings(cwds) + sort.Strings(projectDirs) var units []memoryUnit seen := map[string]bool{} - for _, cwd := range cwds { - repoSeg := repoByCwd[cwd] - memoryDir := filepath.Join(home, ".claude", "projects", claudeProjectDirName(cwd), "memory") + for _, dir := range projectDirs { + projectDir := filepath.Join(projectsRoot, dir) + repoSeg, ok := repoByProjectDir[dir] + if !ok { + // Untracked project (allProjects only): recover its repository from + // the project's own Claude transcript. + repoSeg = repoSegmentForProjectDir(projectDir) + } + memoryDir := filepath.Join(projectDir, "memory") walkErr := filepath.WalkDir(memoryDir, func(p string, d fs.DirEntry, err error) error { if err != nil { // The memory directory simply may not exist for this project. @@ -271,6 +281,83 @@ func collectMemoryUnits(stateDir, home string) ([]memoryUnit, error) { return units, nil } +// trackedProjectRepoSegments maps each Claude project directory name amikalog has +// captured a session for to that project's repository segment, read from the +// session's working directory and git context. +func trackedProjectRepoSegments(stateDir string) (map[string]string, error) { + repoByProjectDir := map[string]string{} + sessionsRoot := EventsDir(stateDir, SourceClaude) + entries, err := os.ReadDir(sessionsRoot) + if err != nil { + if os.IsNotExist(err) { + return repoByProjectDir, nil + } + return nil, fmt.Errorf("reading claude sessions: %w", err) + } + for _, e := range entries { + if e.IsDir() || !strings.HasSuffix(e.Name(), ".jsonl") { + continue + } + data, err := os.ReadFile(filepath.Join(sessionsRoot, e.Name())) + if err != nil { + return nil, fmt.Errorf("reading %s: %w", e.Name(), err) + } + cwd := cwdFromJSONL(data) + if cwd == "" { + continue + } + dir := claudeProjectDirName(cwd) + if _, ok := repoByProjectDir[dir]; !ok { + repoByProjectDir[dir] = repoSegmentFromJSONL(data) + } + } + return repoByProjectDir, nil +} + +// repoSegmentForProjectDir recovers a repository segment for an untracked Claude +// project directory: it reads the working directory from the project's own Claude +// transcript and inspects that directory's git repository. It returns +// "unknown-repo" when there is no transcript cwd or the directory is not a repo. +func repoSegmentForProjectDir(projectDir string) string { + cwd := cwdFromClaudeTranscripts(projectDir) + if cwd == "" { + return unknownRepoSegment + } + if git := GatherGit(cwd); git != nil && git.RepoRoot != "" { + return sanitizeRepoSegment(filepath.Base(git.RepoRoot)) + } + return unknownRepoSegment +} + +// cwdFromClaudeTranscripts returns the working directory recorded in a project's +// Claude Code transcript (the *.jsonl files Claude writes directly under the +// project directory), or "" when none is found. Claude's transcript lines carry +// a top-level "cwd" field, so cwdFromJSONL (which reads Event.CWD, tagged "cwd") +// extracts it without needing the full Event shape. +func cwdFromClaudeTranscripts(projectDir string) string { + entries, err := os.ReadDir(projectDir) + if err != nil { + return "" + } + names := make([]string, 0, len(entries)) + for _, e := range entries { + if !e.IsDir() && strings.HasSuffix(e.Name(), ".jsonl") { + names = append(names, e.Name()) + } + } + sort.Strings(names) + for _, name := range names { + data, err := os.ReadFile(filepath.Join(projectDir, name)) + if err != nil { + continue + } + if cwd := cwdFromJSONL(data); cwd != "" { + return cwd + } + } + return "" +} + // reconcileMemory reconciles one memory file against its cloud copy using the // last-synced hash in manifest, mutating manifest on success. See PushMemories // for the decision table. diff --git a/go/internal/eventlog/memory_test.go b/go/internal/eventlog/memory_test.go index 5e661027..d1253b53 100644 --- a/go/internal/eventlog/memory_test.go +++ b/go/internal/eventlog/memory_test.go @@ -92,7 +92,7 @@ func TestCollectMemoryUnits(t *testing.T) { } } - units, err := collectMemoryUnits(stateDir, home) + units, err := collectMemoryUnits(stateDir, home, false) if err != nil { t.Fatalf("collectMemoryUnits: %v", err) } @@ -280,7 +280,7 @@ func TestPushMemories_UploadsThenSkipsAndWritesManifest(t *testing.T) { down := &fakeDownloader{data: map[string][]byte{}} merger := &fakeMerger{} - rep, err := PushMemories(stateDir, home, up, down, merger) + rep, err := PushMemories(stateDir, home, false, up, down, merger) if err != nil { t.Fatalf("PushMemories: %v", err) } @@ -293,7 +293,7 @@ func TestPushMemories_UploadsThenSkipsAndWritesManifest(t *testing.T) { // Second run: the cloud now holds what we uploaded, so it is a no-op skip. down.data["myrepo/memory/fact.md"] = up.bytesFor("myrepo/memory/fact.md") - rep2, err := PushMemories(stateDir, home, up, down, merger) + rep2, err := PushMemories(stateDir, home, false, up, down, merger) if err != nil { t.Fatalf("PushMemories second run: %v", err) } @@ -301,3 +301,70 @@ func TestPushMemories_UploadsThenSkipsAndWritesManifest(t *testing.T) { t.Fatalf("second run report = %+v, want 1 skipped", rep2) } } + +func TestCollectMemoryUnits_AllProjects(t *testing.T) { + requireGit(t) + stateDir := t.TempDir() + home := t.TempDir() + + // Tracked project: amikalog captured a session recording its cwd + repo. + trackedCwd := "/work/tracked" + writeSessionFile(t, stateDir, SourceClaude, "20240101T000000.000000000Z_sess-a.jsonl", + eventLineCWD(t, trackedCwd, &GitInfo{RepoRoot: "/work/tracked"})) + trackedMem := filepath.Join(home, ".claude", "projects", "-work-tracked", "memory") + if err := os.MkdirAll(trackedMem, 0o755); err != nil { + t.Fatalf("mkdir: %v", err) + } + if err := os.WriteFile(filepath.Join(trackedMem, "a.md"), []byte("a"), 0o644); err != nil { + t.Fatalf("write: %v", err) + } + + // Untracked project: no amikalog session, but a Claude transcript points at a + // real git repo so its repo segment is recovered from git. + repoDir := filepath.Join(t.TempDir(), "untracked-repo") + if err := os.MkdirAll(repoDir, 0o755); err != nil { + t.Fatalf("mkdir repo: %v", err) + } + initRepo(t, repoDir) + untrackedProject := filepath.Join(home, ".claude", "projects", claudeProjectDirName(repoDir)) + if err := os.MkdirAll(filepath.Join(untrackedProject, "memory"), 0o755); err != nil { + t.Fatalf("mkdir: %v", err) + } + // Claude's own transcript carries the working directory. + if err := os.WriteFile(filepath.Join(untrackedProject, "transcript.jsonl"), + []byte(`{"cwd":"`+repoDir+`"}`+"\n"), 0o644); err != nil { + t.Fatalf("write transcript: %v", err) + } + if err := os.WriteFile(filepath.Join(untrackedProject, "memory", "b.md"), []byte("b"), 0o644); err != nil { + t.Fatalf("write: %v", err) + } + + // Default mode: only the tracked project's memory is collected. + got, err := collectMemoryUnits(stateDir, home, false) + if err != nil { + t.Fatalf("collectMemoryUnits(false): %v", err) + } + if len(got) != 1 || got[0].objectKey != "tracked/memory/a.md" { + t.Fatalf("default mode keys = %v, want [tracked/memory/a.md]", keysOf(got)) + } + + // All-projects mode: both, with the untracked repo segment recovered from git. + got, err = collectMemoryUnits(stateDir, home, true) + if err != nil { + t.Fatalf("collectMemoryUnits(true): %v", err) + } + wantUntracked := "untracked-repo/memory/b.md" + keys := keysOf(got) + sort.Strings(keys) + if len(keys) != 2 || keys[0] != "tracked/memory/a.md" || keys[1] != wantUntracked { + t.Fatalf("all-projects keys = %v, want [tracked/memory/a.md %s]", keys, wantUntracked) + } +} + +func keysOf(units []memoryUnit) []string { + ks := make([]string, 0, len(units)) + for _, u := range units { + ks = append(ks, u.objectKey) + } + return ks +} From a1753866d27074472800fc854a3caad12d29a72d Mon Sep 17 00:00:00 2001 From: Piyush Sachdeva Date: Wed, 24 Jun 2026 04:50:44 +0000 Subject: [PATCH 3/3] Fix GetObjectByKey to list by parent folder, not full key GetObjectByKey listed the downloads endpoint with the full object key as the `prefix`, then exact-matched. The storage backend treats a listing prefix as a folder path, so a prefix that includes the filename matches nothing and the lookup always reported found=false. That made `amikalog beta:push --memories` re-upload every memory file on every run: reconcileMemory always took the "no cloud copy yet" branch because the cloud read-back never found the object, regardless of the local push manifest. Session/event uploads were unaffected because they dedup via the byte-size manifest and never call GetObjectByKey. List the object's parent folder (up to and including the final "/") and exact-match within it. Add regression tests that fake the downloads endpoint with folder-prefix semantics; they fail against the old full-key listing and pass with the fix. --- go/internal/apiclient/client.go | 21 +++-- go/internal/apiclient/client_test.go | 122 +++++++++++++++++++++++++++ 2 files changed, 137 insertions(+), 6 deletions(-) diff --git a/go/internal/apiclient/client.go b/go/internal/apiclient/client.go index 4ebd1941..ba8fb56e 100644 --- a/go/internal/apiclient/client.go +++ b/go/internal/apiclient/client.go @@ -606,15 +606,24 @@ func (c *Client) DownloadFromSignedURL(signedURL string) ([]byte, error) { } // GetObjectByKey fetches the current bytes of a single object by its exact -// bucket key. There is no single-object endpoint, so it lists the subtree at -// key (a prefix listing returns the object under its own key, possibly -// alongside keys that share it as a prefix) and downloads the entry whose Key -// matches exactly. found is false when no object has that exact key, which the -// caller can treat as "no cloud copy yet". +// bucket key. There is no single-object endpoint, so it lists the object's +// parent folder and downloads the entry whose Key matches exactly. found is +// false when no object has that exact key, which the caller can treat as "no +// cloud copy yet". +// +// The listing is restricted to the parent prefix (everything up to and +// including the final "/") rather than the full key: the storage backend treats +// a listing prefix as a folder path, so a prefix equal to the full object key +// (filename included) matches nothing. Listing the folder and exact-matching +// within it is the only reliable way to find the object. func (c *Client) GetObjectByKey(key string) (data []byte, found bool, err error) { + prefix := "" + if i := strings.LastIndex(key, "/"); i >= 0 { + prefix = key[:i+1] + } cursor := "" for { - resp, err := c.ListDownloads(key, cursor, 0) + resp, err := c.ListDownloads(prefix, cursor, 0) if err != nil { return nil, false, err } diff --git a/go/internal/apiclient/client_test.go b/go/internal/apiclient/client_test.go index fb981068..3ee1746d 100644 --- a/go/internal/apiclient/client_test.go +++ b/go/internal/apiclient/client_test.go @@ -466,3 +466,125 @@ func TestDownloadFromSignedURL_Non2xxIsHTTPError(t *testing.T) { t.Errorf("status = %d, want 404", httpErr.StatusCode) } } + +// folderPrefixBucket is an httptest server that mimics the storage backend's +// listing semantics: it treats the `prefix` query as a FOLDER path, so it +// returns an object only when the prefix is empty or ends in "/". A prefix equal +// to a full object key (filename included, no trailing slash) matches nothing — +// the exact behavior that made GetObjectByKey re-upload every memory file when +// it listed by the full key. Each object's download_url points back at this same +// server so the returned bytes can be fetched. +func folderPrefixBucket(t *testing.T, objects map[string]string) *httptest.Server { + t.Helper() + var srv *httptest.Server + srv = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/api/v0beta1/storage/downloads": + prefix := r.URL.Query().Get("prefix") + out := []map[string]any{} + for key := range objects { + // Folder semantics: only an empty or "/"-terminated prefix lists + // objects; a full-key prefix matches nothing. + if prefix != "" && !strings.HasSuffix(prefix, "/") { + continue + } + if !strings.HasPrefix(key, prefix) { + continue + } + out = append(out, map[string]any{ + "key": key, + "size": len(objects[key]), + "last_modified": "2026-01-01T00:00:00Z", + "download_url": srv.URL + "/dl?key=" + url.QueryEscape(key), + }) + } + w.WriteHeader(http.StatusOK) + json.NewEncoder(w).Encode(map[string]any{ + "bucket": "org-123", + "prefix": prefix, + "objects": out, + "expires_in": 3600, + "next_cursor": nil, + }) + case "/dl": + body, ok := objects[r.URL.Query().Get("key")] + if !ok { + w.WriteHeader(http.StatusNotFound) + return + } + w.WriteHeader(http.StatusOK) + _, _ = w.Write([]byte(body)) + default: + w.WriteHeader(http.StatusNotFound) + } + })) + t.Cleanup(srv.Close) + return srv +} + +// TestGetObjectByKey_FindsObjectViaFolderPrefix is the regression test for the +// memory re-upload bug: against a backend that only honors folder prefixes, +// GetObjectByKey must still find an object by listing its parent folder and +// exact-matching the key, rather than listing by the full key (which returns +// nothing). It also asserts the prefix actually sent is the folder, not the key. +func TestGetObjectByKey_FindsObjectViaFolderPrefix(t *testing.T) { + objects := map[string]string{ + "decyph-ai-app/memory/memory.md": `{"a":1}`, + // A sibling that shares the folder, to ensure exact-match (not just + // prefix-match) selects the right object. + "decyph-ai-app/memory/memory.md.bak": `{"a":2}`, + } + srv := folderPrefixBucket(t, objects) + + var gotPrefix string + c := NewClient(srv.URL, "key-xyz") + c.HTTP = recordPrefixClient(&gotPrefix) + + data, found, err := c.GetObjectByKey("decyph-ai-app/memory/memory.md") + if err != nil { + t.Fatalf("GetObjectByKey: %v", err) + } + if !found { + t.Fatal("found = false, want true (object exists under the folder)") + } + if string(data) != `{"a":1}` { + t.Errorf("data = %q, want %q", string(data), `{"a":1}`) + } + if gotPrefix != "decyph-ai-app/memory/" { + t.Errorf("listing prefix = %q, want the parent folder %q", gotPrefix, "decyph-ai-app/memory/") + } +} + +// TestGetObjectByKey_NotFound confirms a genuinely absent object reports found +// = false (so callers treat it as "no cloud copy yet"), even though the folder +// listing returns sibling objects. +func TestGetObjectByKey_NotFound(t *testing.T) { + srv := folderPrefixBucket(t, map[string]string{ + "decyph-ai-app/memory/memory.md": `{"a":1}`, + }) + c := NewClient(srv.URL, "key-xyz") + + data, found, err := c.GetObjectByKey("decyph-ai-app/memory/absent.md") + if err != nil { + t.Fatalf("GetObjectByKey: %v", err) + } + if found { + t.Errorf("found = true, want false for an absent key (data %q)", string(data)) + } +} + +// recordPrefixClient returns an *http.Client whose transport records the latest +// `prefix` query parameter seen on a downloads listing request before forwarding +// it unchanged, letting a test assert which prefix GetObjectByKey listed by. +func recordPrefixClient(prefix *string) *http.Client { + return &http.Client{Transport: prefixRecorder{prefix: prefix}} +} + +type prefixRecorder struct{ prefix *string } + +func (p prefixRecorder) RoundTrip(r *http.Request) (*http.Response, error) { + if r.URL.Path == "/api/v0beta1/storage/downloads" { + *p.prefix = r.URL.Query().Get("prefix") + } + return http.DefaultTransport.RoundTrip(r) +}