From cafa926c6f07be6dfcf4a04fadd026db93902b84 Mon Sep 17 00:00:00 2001 From: Relayflow Date: Sun, 27 Sep 2026 20:02:05 +0000 Subject: [PATCH 1/6] feat: add polite content-hash client caching --- .../completed/2026-09/traj_6qsrb702mx1s.json | 65 +++++++++ .../completed/2026-09/traj_6qsrb702mx1s.md | 36 +++++ .trajectories/index.json | 9 +- cmd/relayfile-cli/main.go | 57 ++++++-- internal/mountfuse/wsinvalidate.go | 43 +++++- internal/mountsync/object_cache.go | 88 ++++++++++++ internal/mountsync/object_cache_test.go | 32 +++++ internal/mountsync/syncer.go | 134 +++++++++++------- internal/mountsync/syncer_test.go | 20 +-- packages/local-mount/CHANGELOG.md | 2 + packages/sdk/parity.json | 6 +- packages/sdk/python/README.md | 2 + packages/sdk/python/src/relayfile/__init__.py | 3 +- packages/sdk/python/src/relayfile/client.py | 88 ++++++++++-- packages/sdk/python/tests/test_client.py | 18 +++ packages/sdk/typescript/CHANGELOG.md | 4 + packages/sdk/typescript/src/client.test.ts | 38 +++++ packages/sdk/typescript/src/client.ts | 95 ++++++++----- packages/sdk/typescript/src/types.ts | 8 +- summary.md | 28 ++++ 20 files changed, 653 insertions(+), 123 deletions(-) create mode 100644 .trajectories/completed/2026-09/traj_6qsrb702mx1s.json create mode 100644 .trajectories/completed/2026-09/traj_6qsrb702mx1s.md create mode 100644 internal/mountsync/object_cache.go create mode 100644 internal/mountsync/object_cache_test.go create mode 100644 summary.md diff --git a/.trajectories/completed/2026-09/traj_6qsrb702mx1s.json b/.trajectories/completed/2026-09/traj_6qsrb702mx1s.json new file mode 100644 index 00000000..b5d9744b --- /dev/null +++ b/.trajectories/completed/2026-09/traj_6qsrb702mx1s.json @@ -0,0 +1,65 @@ +{ + "id": "traj_6qsrb702mx1s", + "version": 1, + "task": { + "title": "Implement content-hash client cache and jittered Retry-After handling" + }, + "status": "completed", + "startedAt": "2026-09-27T19:55:27.193Z", + "completedAt": "2026-09-27T20:01:45.765Z", + "agents": [ + { + "name": "default", + "role": "lead", + "joinedAt": "2026-09-27T19:55:27.421Z" + } + ], + "chapters": [ + { + "id": "chap_yvjv91bjz5w9", + "title": "Work", + "agentName": "default", + "startedAt": "2026-09-27T19:55:27.421Z", + "endedAt": "2026-09-27T20:01:45.765Z", + "events": [ + { + "ts": 1790538927421, + "type": "decision", + "content": "Use content hashes rather than the existing revision ETag for conditional cache identity: Use content hashes rather than the existing revision ETag for conditional cache identity", + "raw": { + "question": "Use content hashes rather than the existing revision ETag for conditional cache identity", + "chosen": "Use content hashes rather than the existing revision ETag for conditional cache identity", + "alternatives": [], + "reasoning": "Current ETag semantics are revision-based while contentHash is already in response bodies and tree entries." + }, + "significance": "high" + }, + { + "ts": 1790538927631, + "type": "decision", + "content": "Include the FUSE websocket invalidator but defer unrelated write and control-plane retry loops: Include the FUSE websocket invalidator but defer unrelated write and control-plane retry loops", + "raw": { + "question": "Include the FUSE websocket invalidator but defer unrelated write and control-plane retry loops", + "chosen": "Include the FUSE websocket invalidator but defer unrelated write and control-plane retry loops", + "alternatives": [], + "reasoning": "FUSE reproduces the shipped reconnect-storm path; setup retries, agents reconnect, outbox writes, and polite polling are outside the read fan-in failure mode." + }, + "significance": "high" + } + ] + } + ], + "retrospective": { + "summary": "Added SDK content-hash caches, persistent mount object reuse, polite full-jitter retries, bounded bootstrap fan-out, and cursor-safe listen reconnects", + "approach": "Standard approach", + "confidence": 0.85 + }, + "commits": [], + "filesChanged": [], + "projectId": "AgentWorkforce/relayfile", + "tags": [], + "_trace": { + "startRef": "2d32cfddbfc06e9e371462a50bebd0373e7343a6", + "endRef": "2d32cfddbfc06e9e371462a50bebd0373e7343a6" + } +} \ No newline at end of file diff --git a/.trajectories/completed/2026-09/traj_6qsrb702mx1s.md b/.trajectories/completed/2026-09/traj_6qsrb702mx1s.md new file mode 100644 index 00000000..551d3b59 --- /dev/null +++ b/.trajectories/completed/2026-09/traj_6qsrb702mx1s.md @@ -0,0 +1,36 @@ +# Trajectory: Implement content-hash client cache and jittered Retry-After handling + +> **Status:** ✅ Completed +> **Confidence:** 85% +> **Started:** September 27, 2026 at 07:55 PM +> **Completed:** September 27, 2026 at 08:01 PM + +--- + +## Summary + +Added SDK content-hash caches, persistent mount object reuse, polite full-jitter retries, bounded bootstrap fan-out, and cursor-safe listen reconnects + +**Approach:** Standard approach + +--- + +## Key Decisions + +### Use content hashes rather than the existing revision ETag for conditional cache identity +- **Chose:** Use content hashes rather than the existing revision ETag for conditional cache identity +- **Reasoning:** Current ETag semantics are revision-based while contentHash is already in response bodies and tree entries. + +### Include the FUSE websocket invalidator but defer unrelated write and control-plane retry loops +- **Chose:** Include the FUSE websocket invalidator but defer unrelated write and control-plane retry loops +- **Reasoning:** FUSE reproduces the shipped reconnect-storm path; setup retries, agents reconnect, outbox writes, and polite polling are outside the read fan-in failure mode. + +--- + +## Chapters + +### 1. Work +*Agent: default* + +- Use content hashes rather than the existing revision ETag for conditional cache identity: Use content hashes rather than the existing revision ETag for conditional cache identity +- Include the FUSE websocket invalidator but defer unrelated write and control-plane retry loops: Include the FUSE websocket invalidator but defer unrelated write and control-plane retry loops diff --git a/.trajectories/index.json b/.trajectories/index.json index 6182fc70..88975238 100644 --- a/.trajectories/index.json +++ b/.trajectories/index.json @@ -1,6 +1,6 @@ { "version": 1, - "lastUpdated": "2026-09-18T14:07:19.400Z", + "lastUpdated": "2026-09-27T20:01:45.824Z", "trajectories": { "traj_4pvrlmqfnzng": { "title": "Review PR #278 in AgentWorkforce/relayfile", @@ -407,6 +407,13 @@ "startedAt": "2026-09-18T13:58:18.225Z", "completedAt": "2026-09-18T14:07:19.346Z", "path": "/home/khaliqgant/Projects/AgentWorkforce/relayfile/.trajectories/completed/2026-09/traj_p2o6b2q8epj1.json" + }, + "traj_6qsrb702mx1s": { + "title": "Implement content-hash client cache and jittered Retry-After handling", + "status": "completed", + "startedAt": "2026-09-27T19:55:27.193Z", + "completedAt": "2026-09-27T20:01:45.765Z", + "path": "/home/daytona/.relayflow-v2-supervisor/durable/repository/.trajectories/completed/2026-09/traj_6qsrb702mx1s.json" } } } \ No newline at end of file diff --git a/cmd/relayfile-cli/main.go b/cmd/relayfile-cli/main.go index 6ce5dd91..6c148648 100644 --- a/cmd/relayfile-cli/main.go +++ b/cmd/relayfile-cli/main.go @@ -8785,25 +8785,34 @@ func runListen(args []string, stdout io.Writer) error { ) backoff := baseBackoff everConnected := false + lastCursor := "" for { - connected, err := runListenSession(rootCtx, cfg) - if connected { + cfg.dialURL = listenDialURLWithCursor(cfg.dialURL, lastCursor) + result, err := runListenSession(rootCtx, cfg) + if result.cursor != "" { + lastCursor = result.cursor + } + if result.connected { everConnected = true backoff = baseBackoff } if err == nil { return nil } - if !everConnected { + if !everConnected && result.status != http.StatusTooManyRequests && result.status != http.StatusServiceUnavailable { return fmt.Errorf("connect to event stream: %w", err) } + delay := time.Duration(mathrand.Float64() * float64(backoff)) + if result.retryAfter > delay { + delay = result.retryAfter + } if !*daemonized { - fmt.Fprintf(os.Stderr, "listen: stream error (%v); reconnecting in %s\n", err, backoff) + fmt.Fprintf(os.Stderr, "listen: stream error (%v); reconnecting in %s\n", err, delay) } select { case <-rootCtx.Done(): return nil - case <-time.After(backoff): + case <-time.After(delay): } if backoff < maxBackoff { backoff *= 2 @@ -8814,6 +8823,18 @@ func runListen(args []string, stdout io.Writer) error { } } +func listenDialURLWithCursor(rawURL, cursor string) string { + u, err := url.Parse(rawURL) + if err != nil || strings.TrimSpace(cursor) == "" { + return rawURL + } + q := u.Query() + q.Del("from") + q.Set("cursor", strings.TrimSpace(cursor)) + u.RawQuery = q.Encode() + return u.String() +} + // listenSessionConfig carries the per-session inputs for runListenSession so a // dropped connection can be re-established without recomputing filters. type listenSessionConfig struct { @@ -8833,15 +8854,27 @@ type listenSessionConfig struct { // bool reports whether the websocket was successfully established, so the // caller can tell a first-attempt dial failure (fatal) from a mid-flight // disconnect (retryable). -func runListenSession(rootCtx context.Context, cfg listenSessionConfig) (bool, error) { - conn, _, err := websocket.Dial(rootCtx, cfg.dialURL, &websocket.DialOptions{ +type listenSessionResult struct { + connected bool + cursor string + status int + retryAfter time.Duration +} + +func runListenSession(rootCtx context.Context, cfg listenSessionConfig) (listenSessionResult, error) { + conn, response, err := websocket.Dial(rootCtx, cfg.dialURL, &websocket.DialOptions{ HTTPClient: cfg.httpClient, HTTPHeader: http.Header{ "Authorization": []string{"Bearer " + cfg.token}, }, }) if err != nil { - return false, err + result := listenSessionResult{} + if response != nil { + result.status = response.StatusCode + result.retryAfter = parseRetryAfter(response.Header.Get("Retry-After")) + } + return result, err } defer conn.Close(websocket.StatusNormalClosure, "") @@ -8852,6 +8885,7 @@ func runListenSession(rootCtx context.Context, cfg listenSessionConfig) (bool, e conn.SetReadLimit(8 << 20) // 8 MiB stdout := cfg.stdout + lastCursor := "" for { var raw json.RawMessage if err := wsjson.Read(rootCtx, conn, &raw); err != nil { @@ -8863,15 +8897,18 @@ func runListenSession(rootCtx context.Context, cfg listenSessionConfig) (bool, e // would exit 0 and a `Restart=on-failure` supervisor unit would // not bring the listener back. if errors.Is(err, context.Canceled) || rootCtx.Err() != nil { - return true, nil + return listenSessionResult{connected: true, cursor: lastCursor}, nil } - return true, fmt.Errorf("event stream error: %w", err) + return listenSessionResult{connected: true, cursor: lastCursor}, fmt.Errorf("event stream error: %w", err) } var evt listenEvent if err := json.Unmarshal(raw, &evt); err != nil || evt.Type == "" || evt.Type == "pong" { continue } + if cursor := strings.TrimSpace(evt.EventID); cursor != "" { + lastCursor = cursor + } if cfg.typeFilter != "" && evt.Type != cfg.typeFilter { continue } diff --git a/internal/mountfuse/wsinvalidate.go b/internal/mountfuse/wsinvalidate.go index dc881fe9..3d5cf895 100644 --- a/internal/mountfuse/wsinvalidate.go +++ b/internal/mountfuse/wsinvalidate.go @@ -3,10 +3,13 @@ package mountfuse import ( "context" "encoding/json" + "errors" "fmt" "log" + "math/rand/v2" "net/http" "net/url" + "strconv" "strings" "time" @@ -30,6 +33,15 @@ type wsEvent struct { Path string `json:"path"` } +type wsDialError struct { + err error + status int + retryAfter time.Duration +} + +func (e *wsDialError) Error() string { return e.err.Error() } +func (e *wsDialError) Unwrap() error { return e.err } + // WSInvalidator connects to the relayfile WebSocket event stream and // invalidates cached entries in fsState when files change remotely. type WSInvalidator struct { @@ -98,12 +110,17 @@ func (w *WSInvalidator) Run(ctx context.Context) { return } - w.logger.Printf("mountfuse: ws disconnected: %v; reconnecting in %v", err, backoff) + delay := time.Duration(rand.Float64() * float64(backoff)) + var dialErr *wsDialError + if errors.As(err, &dialErr) && dialErr.retryAfter > delay { + delay = dialErr.retryAfter + } + w.logger.Printf("mountfuse: ws disconnected: %v; reconnecting in %v", err, delay) select { case <-ctx.Done(): return - case <-time.After(backoff): + case <-time.After(delay): } backoff *= 2 @@ -119,6 +136,10 @@ func isAuthError(err error) bool { if err == nil { return false } + var dialErr *wsDialError + if errors.As(err, &dialErr) { + return dialErr.status == http.StatusUnauthorized || dialErr.status == http.StatusForbidden + } msg := err.Error() return strings.Contains(msg, "status = 401") || strings.Contains(msg, "status = 403") || @@ -132,12 +153,15 @@ func (w *WSInvalidator) listenOnce(ctx context.Context) error { return fmt.Errorf("building ws url: %w", err) } - conn, _, err := websocket.Dial(ctx, wsURL, &websocket.DialOptions{ + conn, response, err := websocket.Dial(ctx, wsURL, &websocket.DialOptions{ HTTPHeader: http.Header{ "Authorization": []string{"Bearer " + w.currentToken()}, }, }) if err != nil { + if response != nil { + return &wsDialError{err: fmt.Errorf("ws dial: %w", err), status: response.StatusCode, retryAfter: parseWSRetryAfter(response.Header.Get("Retry-After"))} + } return fmt.Errorf("ws dial: %w", err) } defer conn.Close(websocket.StatusNormalClosure, "") @@ -161,6 +185,19 @@ func (w *WSInvalidator) listenOnce(ctx context.Context) error { } } +func parseWSRetryAfter(value string) time.Duration { + value = strings.TrimSpace(value) + if seconds, err := strconv.Atoi(value); err == nil && seconds >= 0 { + return time.Duration(seconds) * time.Second + } + if timestamp, err := http.ParseTime(value); err == nil { + if delay := time.Until(timestamp); delay > 0 { + return delay + } + } + return 0 +} + func (w *WSInvalidator) handleEvent(event wsEvent) { eventType := strings.ToLower(strings.TrimSpace(event.Type)) if event.Path == "" { diff --git a/internal/mountsync/object_cache.go b/internal/mountsync/object_cache.go new file mode 100644 index 00000000..2643951f --- /dev/null +++ b/internal/mountsync/object_cache.go @@ -0,0 +1,88 @@ +package mountsync + +import ( + "crypto/sha256" + "encoding/base64" + "encoding/hex" + "encoding/json" + "os" + "path/filepath" + "regexp" + "strings" +) + +var objectHashPattern = regexp.MustCompile(`^[a-f0-9]{64}$`) + +type objectCache struct{ root string } + +func defaultObjectCache() *objectCache { + home, err := os.UserHomeDir() + if err != nil || home == "" { + return nil + } + return &objectCache{root: filepath.Join(home, ".relayfile", "cache", "objects")} +} + +func normalizeObjectHash(value string) string { + value = strings.TrimPrefix(strings.ToLower(strings.TrimSpace(value)), "sha256:") + if !objectHashPattern.MatchString(value) { + return "" + } + return value +} + +func (c *objectCache) get(hash string) (RemoteFile, bool) { + var file RemoteFile + hash = normalizeObjectHash(hash) + if c == nil || hash == "" { + return file, false + } + data, err := os.ReadFile(filepath.Join(c.root, hash)) + if err != nil || json.Unmarshal(data, &file) != nil || normalizeObjectHash(file.ContentHash) != hash || remoteFileHash(file) != hash { + return RemoteFile{}, false + } + return file, true +} + +func (c *objectCache) put(file RemoteFile) { + hash := normalizeObjectHash(file.ContentHash) + if c == nil || hash == "" || remoteFileHash(file) != hash { + return + } + if os.MkdirAll(c.root, 0o700) != nil { + return + } + data, err := json.Marshal(file) + if err != nil { + return + } + tmp, err := os.CreateTemp(c.root, ".object-*") + if err != nil { + return + } + name := tmp.Name() + defer os.Remove(name) + _ = tmp.Chmod(0o600) + if _, err = tmp.Write(data); err == nil { + err = tmp.Sync() + } + if closeErr := tmp.Close(); err == nil { + err = closeErr + } + if err == nil { + _ = os.Rename(name, filepath.Join(c.root, hash)) + } +} + +func remoteFileHash(file RemoteFile) string { + content := []byte(file.Content) + if strings.EqualFold(file.Encoding, "base64") { + decoded, err := base64.StdEncoding.DecodeString(file.Content) + if err != nil { + return "" + } + content = decoded + } + sum := sha256.Sum256(content) + return hex.EncodeToString(sum[:]) +} diff --git a/internal/mountsync/object_cache_test.go b/internal/mountsync/object_cache_test.go new file mode 100644 index 00000000..e42d020e --- /dev/null +++ b/internal/mountsync/object_cache_test.go @@ -0,0 +1,32 @@ +package mountsync + +import ( + "os" + "path/filepath" + "testing" +) + +func TestObjectCacheRoundTripAndRejectsCorruption(t *testing.T) { + cache := &objectCache{root: filepath.Join(t.TempDir(), "objects")} + file := RemoteFile{Content: "hello", ContentHash: "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"} + cache.put(file) + + got, ok := cache.get(file.ContentHash) + if !ok || got.Content != "hello" { + t.Fatalf("cache get = %#v, %v", got, ok) + } + path := filepath.Join(cache.root, file.ContentHash) + info, err := os.Stat(path) + if err != nil { + t.Fatal(err) + } + if info.Mode().Perm() != 0o600 { + t.Fatalf("cache mode = %v", info.Mode().Perm()) + } + if err := os.WriteFile(path, []byte(`{"content":"tampered","contentHash":"`+file.ContentHash+`"}`), 0o600); err != nil { + t.Fatal(err) + } + if _, ok := cache.get(file.ContentHash); ok { + t.Fatal("corrupt object was accepted") + } +} diff --git a/internal/mountsync/syncer.go b/internal/mountsync/syncer.go index d474d61d..4296869d 100644 --- a/internal/mountsync/syncer.go +++ b/internal/mountsync/syncer.go @@ -14,6 +14,7 @@ import ( "fmt" "io" "math" + "math/rand" "mime" "net" "net/http" @@ -191,7 +192,7 @@ const ( defaultIncrementalReadNotReadyTTL = 5 * time.Minute defaultCursorResolutionAttempts = 3 defaultCursorRetryBaseDelay = 250 * time.Millisecond - defaultBootstrapReadWorkers = 16 + defaultBootstrapReadWorkers = 4 // Bulk bootstrap reads deliberately match the tree checkpoint size. One // request therefore replaces at most 32 point reads without creating a new // unbounded response surface. The decoded aggregate stays below 32 MiB and @@ -205,7 +206,7 @@ const ( defaultBulkReadMaxPathBytes = 4096 defaultBulkReadMaxPathsBytes = 32 << 10 defaultBulkReadMaxRequestBytes int64 = 64 << 10 - defaultIncrementalReadWorkers = 16 + defaultIncrementalReadWorkers = 4 defaultReceiptSettlementWorkers = 16 // fullTreeTraversalDepth bounds each tree request so the client can see and // prune a reserved .relay directory before the server reaches deep @@ -326,13 +327,18 @@ type HTTPError struct { Code string Message string Action string + Reason string } func (e *HTTPError) Error() string { + reason := "" + if e.Reason != "" { + reason = " (reason: " + e.Reason + ")" + } if e.Code != "" { - return fmt.Sprintf("http %d %s: %s", e.StatusCode, e.Code, e.Message) + return fmt.Sprintf("http %d %s: %s%s", e.StatusCode, e.Code, e.Message, reason) } - return fmt.Sprintf("http %d: %s", e.StatusCode, e.Message) + return fmt.Sprintf("http %d: %s%s", e.StatusCode, e.Message, reason) } // MalformedPaginationError reports a server response that cannot make a @@ -1289,12 +1295,16 @@ func (c *HTTPClient) ExportGithubWorkingTreeTar(ctx context.Context, workspaceID Code string `json:"code"` Message string `json:"message"` Action string `json:"action"` + Details struct { + Reason string `json:"reason"` + } `json:"details"` } _ = json.Unmarshal(payloadBytes, &errPayload) return GithubWorkingTreeTar{}, &HTTPError{ StatusCode: resp.StatusCode, Code: errPayload.Code, Message: errPayload.Message, + Reason: errPayload.Details.Reason, } } } @@ -1449,6 +1459,9 @@ func (c *HTTPClient) doBytesWithLimit(ctx context.Context, method, requestPath s Code string `json:"code"` Message string `json:"message"` Action string `json:"action"` + Details struct { + Reason string `json:"reason"` + } `json:"details"` } _ = json.Unmarshal(payloadBytes, &errPayload) @@ -1494,6 +1507,7 @@ func (c *HTTPClient) doBytesWithLimit(ctx context.Context, method, requestPath s Code: errPayload.Code, Message: errPayload.Message, Action: errPayload.Action, + Reason: errPayload.Details.Reason, } } } @@ -1535,14 +1549,17 @@ func (c *HTTPClient) refreshTokenAfterUnauthorized(attemptedToken string) bool { } type SyncerOptions struct { - WorkspaceID string - RemoteRoot string - LocalRoot string - StateFile string - StateDir string - MountKind string - ValidateState bool - EventProvider string + // ObjectCacheRoot overrides the content-addressed cache directory. The + // default HTTP mount cache is ~/.relayfile/cache/objects. + ObjectCacheRoot string + WorkspaceID string + RemoteRoot string + LocalRoot string + StateFile string + StateDir string + MountKind string + ValidateState bool + EventProvider string // ScopedChild identifies a Syncer whose local root is one child beneath a // catalog root. Provider paths that resemble catalog-only artifacts are // ordinary content there; exact mounts keep those artifacts reserved. @@ -1747,6 +1764,7 @@ func formatBytes(value uint64) string { type Syncer struct { client RemoteClient + objectCache *objectCache workspace string remoteRoot string localRoot string @@ -2609,8 +2627,15 @@ func NewSyncer(client RemoteClient, opts SyncerOptions) (*Syncer, error) { } } githubWorkingTree := detectGithubWorkingTreeMount(remoteRoot) + var cache *objectCache + if root := strings.TrimSpace(opts.ObjectCacheRoot); root != "" { + cache = &objectCache{root: root} + } else if _, ok := client.(*HTTPClient); ok { + cache = defaultObjectCache() + } return &Syncer{ client: client, + objectCache: cache, workspace: workspace, remoteRoot: remoteRoot, localRoot: localRoot, @@ -5976,9 +6001,6 @@ func (s *Syncer) scheduleWebSocketReconnectLocked(retryAfter time.Duration) { if delay <= 0 { delay = websocketReconnectDelay(s.wsReconnectFailures) } - if delay > defaultWebSocketReconnectMax { - delay = defaultWebSocketReconnectMax - } s.wsNextAttempt = time.Now().Add(delay) } @@ -6007,15 +6029,7 @@ func websocketReconnectDelay(failures int) time.Duration { break } } - jitter := time.Duration(time.Now().UnixNano()%int64(defaultWebSocketReconnectJitter*2)) - defaultWebSocketReconnectJitter - delay += jitter - if delay < defaultWebSocketReconnectBase { - return defaultWebSocketReconnectBase - } - if delay > defaultWebSocketReconnectMax { - return defaultWebSocketReconnectMax - } - return delay + return time.Duration(rand.Float64() * float64(delay)) } // listenerHealthLocked returns listener state independently of LastEventAt. @@ -7721,6 +7735,7 @@ func (s *Syncer) pullRemoteFullTree(ctx context.Context, conflicted map[string]s Index: len(readJobs), RemotePath: remotePath, Size: entry.Size, + Entry: entry, }) } var readErr error @@ -7938,6 +7953,7 @@ type bootstrapReadJob struct { RemotePath string Size int64 ForcePointRead bool + Entry TreeEntry } type bootstrapReadResult struct { @@ -8023,6 +8039,24 @@ func (s *Syncer) readBootstrapFiles(ctx context.Context, jobs []bootstrapReadJob // complete; callers that need incremental oversized reads stay on the bulk // path above. func (s *Syncer) readBootstrapFilesEach(ctx context.Context, jobs []bootstrapReadJob, prog bootstrapProgress, handle func(bootstrapReadResult) error) error { + if len(jobs) == 0 { + return nil + } + remaining := make([]bootstrapReadJob, 0, len(jobs)) + for _, job := range jobs { + cached, ok := s.objectCache.get(job.Entry.ContentHash) + if !ok { + remaining = append(remaining, job) + continue + } + cached.Path, cached.Revision = job.RemotePath, job.Entry.Revision + cached.Type, cached.Target, cached.Mode = job.Entry.Type, job.Entry.Target, job.Entry.Mode + prog.touch() + if err := handle(bootstrapReadResult{Index: job.Index, RemotePath: job.RemotePath, File: cached}); err != nil { + return err + } + } + jobs = remaining if len(jobs) == 0 { return nil } @@ -8102,20 +8136,16 @@ func (s *Syncer) readBootstrapFilesBulkEach(ctx context.Context, client bulkRead continue } prog.touch() + file := RemoteFile{ + Path: result.Path, Type: result.Type, Target: result.Target, Mode: result.Mode, + Revision: result.Revision, ContentType: result.ContentType, Content: result.Content, + Encoding: result.Encoding, ContentHash: result.ContentHash, + } + s.objectCache.put(file) if callbackErr := handle(bootstrapReadResult{ Index: job.Index, RemotePath: job.RemotePath, - File: RemoteFile{ - Path: result.Path, - Type: result.Type, - Target: result.Target, - Mode: result.Mode, - Revision: result.Revision, - ContentType: result.ContentType, - Content: result.Content, - Encoding: result.Encoding, - ContentHash: result.ContentHash, - }, + File: file, }); callbackErr != nil { return callbackErr } @@ -8276,8 +8306,16 @@ func (s *Syncer) readBootstrapFilesIndividuallyBatchEach(ctx context.Context, jo go func() { defer wg.Done() for job := range jobCh { + if cached, ok := s.objectCache.get(job.Entry.ContentHash); ok { + cached.Path, cached.Revision = job.RemotePath, job.Entry.Revision + cached.Type, cached.Target, cached.Mode = job.Entry.Type, job.Entry.Target, job.Entry.Mode + prog.touch() + resultCh <- bootstrapReadResult{Index: job.Index, RemotePath: job.RemotePath, File: cached} + continue + } file, err := s.client.ReadFile(readCtx, s.workspace, job.RemotePath) if err == nil { + s.objectCache.put(file) prog.touch() } resultCh <- bootstrapReadResult{ @@ -8315,8 +8353,8 @@ func bootstrapReadWorkers() int { if err != nil || v <= 0 { return defaultBootstrapReadWorkers } - if v > 64 { - return 64 + if v > 4 { + return 4 } return v } @@ -9325,8 +9363,8 @@ func incrementalReadWorkers() int { if err != nil || workers <= 0 { return defaultIncrementalReadWorkers } - if workers > 64 { - return 64 + if workers > 4 { + return 4 } return workers } @@ -12495,12 +12533,7 @@ func (c *HTTPClient) retryDelay(attempt int, retryAfterHeader string) time.Durat if maxDelay <= 0 { maxDelay = defaultRetryAfterMaxDelay } - if retryAfter := parseRetryAfter(retryAfterHeader); retryAfter > 0 { - if retryAfter > maxDelay { - return maxDelay - } - return retryAfter - } + retryAfter := parseRetryAfter(retryAfterHeader) delay := c.baseDelay if delay <= 0 { delay = 100 * time.Millisecond @@ -12508,13 +12541,18 @@ func (c *HTTPClient) retryDelay(attempt int, retryAfterHeader string) time.Durat for i := 1; i < attempt; i++ { delay *= 2 if delay >= maxDelay { - return maxDelay + delay = maxDelay + break } } if delay > maxDelay { - return maxDelay + delay = maxDelay } - return delay + jittered := time.Duration(rand.Float64() * float64(delay)) + if retryAfter > jittered { + return retryAfter + } + return jittered } func parseRetryAfter(header string) time.Duration { diff --git a/internal/mountsync/syncer_test.go b/internal/mountsync/syncer_test.go index 4c5b5799..1e5f8dac 100644 --- a/internal/mountsync/syncer_test.go +++ b/internal/mountsync/syncer_test.go @@ -65,8 +65,8 @@ func TestHTTPClientRetryDelayHonorsRetryAfter(t *testing.T) { if got := client.retryDelay(1, "30"); got != 30*time.Second { t.Fatalf("expected Retry-After 30s, got %s", got) } - if got := client.retryDelay(1, "999"); got != defaultRetryAfterMaxDelay { - t.Fatalf("expected Retry-After cap %s, got %s", defaultRetryAfterMaxDelay, got) + if got := client.retryDelay(1, "999"); got != 999*time.Second { + t.Fatalf("expected Retry-After floor 999s, got %s", got) } } @@ -155,10 +155,10 @@ func TestRedactSensitiveLogQueryValues(t *testing.T) { } func TestWebSocketReconnectDelayBounds(t *testing.T) { - if got := websocketReconnectDelay(1); got < defaultWebSocketReconnectBase || got > defaultWebSocketReconnectBase+defaultWebSocketReconnectJitter { + if got := websocketReconnectDelay(1); got < 0 || got > defaultWebSocketReconnectBase { t.Fatalf("first reconnect delay out of bounds: %s", got) } - if got := websocketReconnectDelay(20); got < defaultWebSocketReconnectMax-defaultWebSocketReconnectJitter || got > defaultWebSocketReconnectMax { + if got := websocketReconnectDelay(20); got < 0 || got > defaultWebSocketReconnectMax { t.Fatalf("capped reconnect delay out of bounds: %s", got) } } @@ -14083,11 +14083,13 @@ func TestTreeBootstrapPersistsWithinPageBeforeDeadlineAndResumes(t *testing.T) { localDir := t.TempDir() first, err := NewSyncer(client, SyncerOptions{ - WorkspaceID: "ws_neon_checkpoint", - RemoteRoot: "/neon/advisors/by-project", - LocalRoot: localDir, - StateFile: stateFile, - BootstrapTimeout: 210 * time.Millisecond, + WorkspaceID: "ws_neon_checkpoint", + RemoteRoot: "/neon/advisors/by-project", + LocalRoot: localDir, + StateFile: stateFile, + // Four polite bootstrap workers need eight 60ms waves to cross the + // 32-file durable checkpoint without completing the 140-file page. + BootstrapTimeout: 550 * time.Millisecond, BootstrapMaxFilesPerCycle: -1, FullPullEvery: -1, }) diff --git a/packages/local-mount/CHANGELOG.md b/packages/local-mount/CHANGELOG.md index 3c17309a..f63571c8 100644 --- a/packages/local-mount/CHANGELOG.md +++ b/packages/local-mount/CHANGELOG.md @@ -6,6 +6,8 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +- Bootstrap point reads are capped at four concurrent requests and reuse verified objects from `~/.relayfile/cache/objects/`. HTTP and WebSocket overload retries use full jitter without shortening `Retry-After`. + ### Fixed - Mount event pagination now fails closed with an actionable malformed-pagination error when the server repeats or cycles a non-empty `nextCursor`, instead of spinning until the reconcile timeout. diff --git a/packages/sdk/parity.json b/packages/sdk/parity.json index a98251b2..bd41c5a1 100644 --- a/packages/sdk/parity.json +++ b/packages/sdk/parity.json @@ -416,12 +416,14 @@ }, { "id": "client-read-cache", - "status": "planned", + "status": "both", "summary": "Client-side read cache with write-through invalidation.", "tsExports": [ "RelayFileReadCacheOptions" ], - "pyExports": [] + "pyExports": [ + "RelayFileReadCacheOptions" + ] }, { "id": "sync-socket-manager", diff --git a/packages/sdk/python/README.md b/packages/sdk/python/README.md index e9e8c8a6..3836496d 100644 --- a/packages/sdk/python/README.md +++ b/packages/sdk/python/README.md @@ -2,6 +2,8 @@ Python SDK for the RelayFile virtual filesystem API. +File reads use a 32 MiB content-addressed cache by default and revalidate cached hashes with `If-None-Match`. Configure it with `RelayFileReadCacheOptions(max_bytes=...)`, or pass `read_cache=False` to either client to disable it. Retry delays use full jitter and never retry before a valid server `Retry-After` value. + ## Install ```bash diff --git a/packages/sdk/python/src/relayfile/__init__.py b/packages/sdk/python/src/relayfile/__init__.py index 916ba43c..4337fefc 100644 --- a/packages/sdk/python/src/relayfile/__init__.py +++ b/packages/sdk/python/src/relayfile/__init__.py @@ -1,4 +1,4 @@ -from .client import RelayFileClient, AsyncRelayFileClient, RetryOptions +from .client import AsyncRelayFileClient, RelayFileClient, RelayFileReadCacheOptions, RetryOptions from .errors import ( CloudApiError, IntegrationConnectionTimeoutError, @@ -94,6 +94,7 @@ "RelayFileClient", "AsyncRelayFileClient", "RetryOptions", + "RelayFileReadCacheOptions", "RelayfileSetup", "WorkspaceHandle", "WORKSPACE_INTEGRATION_PROVIDERS", diff --git a/packages/sdk/python/src/relayfile/client.py b/packages/sdk/python/src/relayfile/client.py index e55bdc74..7465916f 100644 --- a/packages/sdk/python/src/relayfile/client.py +++ b/packages/sdk/python/src/relayfile/client.py @@ -2,7 +2,9 @@ import json import random +import threading import time +from collections import OrderedDict from dataclasses import dataclass from typing import Any, Callable, Union from urllib.parse import quote, urlencode @@ -39,6 +41,55 @@ class RetryOptions: jitter_ratio: float = 0.2 +@dataclass +class RelayFileReadCacheOptions: + """Byte-capped content-addressed read cache configuration.""" + + max_bytes: int = 32 * 1024 * 1024 + + +class _FileReadCache: + def __init__(self, options: RelayFileReadCacheOptions | None) -> None: + self.max_bytes = max(0, (options or RelayFileReadCacheOptions()).max_bytes) + self.total_bytes = 0 + self.objects: OrderedDict[str, tuple[dict[str, Any], int]] = OrderedDict() + self.paths: dict[str, str] = {} + self.lock = threading.RLock() + + def get(self, key: str) -> dict[str, Any] | None: + with self.lock: + digest = self.paths.get(key) + item = self.objects.get(digest or "") + if item is None: + return None + self.objects.move_to_end(digest) + return dict(item[0]) + + def put(self, key: str, value: dict[str, Any]) -> None: + digest = str(value.get("contentHash", "")).removeprefix("sha256:").lower() + if len(digest) != 64 or any(c not in "0123456789abcdef" for c in digest): + return + content = str(value.get("content", "")) + if value.get("encoding") == "base64": + import base64 + try: + size = len(base64.b64decode(content, validate=True)) + except ValueError: + return + else: + size = len(content.encode()) + with self.lock: + old = self.objects.pop(digest, None) + if old: + self.total_bytes -= old[1] + self.objects[digest] = (dict(value), size) + self.paths[key] = digest + self.total_bytes += size + while self.total_bytes > self.max_bytes and self.objects: + _, (_, removed_size) = self.objects.popitem(last=False) + self.total_bytes -= removed_size + + # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- @@ -137,12 +188,9 @@ def _parse_retry_after_ms(header: str | None) -> float | None: def _compute_delay(retry: RetryOptions, attempt: int, retry_after: str | None) -> float: parsed = _parse_retry_after_ms(retry_after) - if parsed is not None: - return min(retry.max_delay_ms, parsed) backoff = retry.base_delay_ms * (2 ** max(0, attempt - 1)) capped = min(retry.max_delay_ms, backoff) - factor = 1 + (random.random() * 2 - 1) * retry.jitter_ratio - return max(0, round(capped * factor)) + return max(parsed or 0, round(random.random() * capped)) def _should_retry(status: int, retries: int, max_retries: int) -> bool: @@ -283,6 +331,7 @@ def __init__( retry: RetryOptions | None = None, http_client: httpx.Client | None = None, cloud_base_url: str | None = None, + read_cache: RelayFileReadCacheOptions | bool = RelayFileReadCacheOptions(), ) -> None: self._base_url = base_url.rstrip("/") # Integration setup verbs (list_accessible_resources + @@ -298,6 +347,7 @@ def __init__( self._retry = _normalize_retry(retry) self._client = http_client or httpx.Client(timeout=timeout) self._owns_client = http_client is None + self._read_cache = None if read_cache is False else _FileReadCache(read_cache if isinstance(read_cache, RelayFileReadCacheOptions) else None) def close(self) -> None: if self._owns_client: @@ -373,7 +423,7 @@ def _request_response( continue raise - if resp.is_success: + if resp.is_success or resp.status_code == 304: return resp payload = _read_payload(resp) @@ -418,11 +468,21 @@ def read_file( correlation_id: str | None = None, ) -> dict[str, Any]: query = _build_query({"path": path}) - return self._request( + key = f"{workspace_id}:{path}" + cached = self._read_cache.get(key) if self._read_cache else None + headers = {"If-None-Match": f'"{str(cached["contentHash"]).removeprefix("sha256:")}"'} if cached and cached.get("contentHash") else None + response = self._request_response( "GET", f"/v1/workspaces/{_enc(workspace_id)}/fs/file{query}", + headers=headers, correlation_id=correlation_id, ) + if response.status_code == 304 and cached is not None: + return cached + value = _read_payload(response) + if self._read_cache and isinstance(value, dict): + self._read_cache.put(key, value) + return value def query_files( self, @@ -982,6 +1042,7 @@ def __init__( retry: RetryOptions | None = None, http_client: httpx.AsyncClient | None = None, cloud_base_url: str | None = None, + read_cache: RelayFileReadCacheOptions | bool = RelayFileReadCacheOptions(), ) -> None: self._base_url = base_url.rstrip("/") # See the sync client for why integration setup verbs target a @@ -994,6 +1055,7 @@ def __init__( self._retry = _normalize_retry(retry) self._client = http_client or httpx.AsyncClient(timeout=timeout) self._owns_client = http_client is None + self._read_cache = None if read_cache is False else _FileReadCache(read_cache if isinstance(read_cache, RelayFileReadCacheOptions) else None) async def aclose(self) -> None: if self._owns_client: @@ -1073,7 +1135,7 @@ async def _request_response( continue raise - if resp.is_success: + if resp.is_success or resp.status_code == 304: return resp payload = _read_payload(resp) @@ -1118,11 +1180,21 @@ async def read_file( correlation_id: str | None = None, ) -> dict[str, Any]: query = _build_query({"path": path}) - return await self._request( + key = f"{workspace_id}:{path}" + cached = self._read_cache.get(key) if self._read_cache else None + headers = {"If-None-Match": f'"{str(cached["contentHash"]).removeprefix("sha256:")}"'} if cached and cached.get("contentHash") else None + response = await self._request_response( "GET", f"/v1/workspaces/{_enc(workspace_id)}/fs/file{query}", + headers=headers, correlation_id=correlation_id, ) + if response.status_code == 304 and cached is not None: + return cached + value = _read_payload(response) + if self._read_cache and isinstance(value, dict): + self._read_cache.put(key, value) + return value async def query_files( self, diff --git a/packages/sdk/python/tests/test_client.py b/packages/sdk/python/tests/test_client.py index 4562415c..134abd23 100644 --- a/packages/sdk/python/tests/test_client.py +++ b/packages/sdk/python/tests/test_client.py @@ -17,6 +17,7 @@ QueueFullError, RelayFileApiError, RelayFileClient, + RelayFileReadCacheOptions, RetryOptions, RevisionConflictError, WritebackItem, @@ -85,6 +86,23 @@ def test_read_file(self) -> None: res = client.read_file("ws_acme", "/f.json") assert res["content"] == '{"id":1}' + @respx.mock + def test_read_file_revalidates_with_content_hash_and_handles_304(self) -> None: + digest = "a" * 64 + route = respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file").mock( + side_effect=[ + httpx.Response(200, json={"path": "/f", "revision": "rev_1", "contentType": "text/plain", "content": "hello", "contentHash": digest}), + httpx.Response(304), + ] + ) + client = RelayFileClient(BASE, "tok", read_cache=RelayFileReadCacheOptions(max_bytes=5)) + + first = client.read_file("ws_acme", "/f") + second = client.read_file("ws_acme", "/f") + + assert second == first + assert route.calls[1].request.headers["If-None-Match"] == f'"{digest}"' + @respx.mock def test_write_file(self) -> None: payload = {"opId": "op_1", "status": "queued", "targetRevision": "rev_4"} diff --git a/packages/sdk/typescript/CHANGELOG.md b/packages/sdk/typescript/CHANGELOG.md index 27912b33..f0e3af86 100644 --- a/packages/sdk/typescript/CHANGELOG.md +++ b/packages/sdk/typescript/CHANGELOG.md @@ -8,12 +8,16 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added +- File reads now use a configurable, byte-capped content-addressed cache (`readCache.maxBytes`, 32 MiB by default; `readCache: false` disables it), revalidate with `If-None-Match`, and serve `304` responses from cached content. + - New `@relayfile/sdk/relay-cli` subpath export: `createRelayCliSurface()` returns relayfile's mountable CLI surface (`id: 'relayfile'`, contract v1), which the `agent-relay` CLI mounts as `agent-relay file`. `commands` is a checked-in snapshot of the Go CLI's own command table (`relayfile __command-spec --json`, regenerated by `npm run gen:command-spec`), and `run(argv, io)` spawns the same Go binary `relayfile` does, returning its real exit code. No relayfile command is reimplemented. - Binary resolution (`resolveRelayfileBinary`, `platformBinaryName`, `genericBinaryName`, `findSourceCheckoutRoot`) and the Cloud sign-in preflight (`prepareCloudSession`, `announceSetupIntent`) now live here, moved out of the `relayfile` package's scripts. Both the `relayfile` bin shim and `agent-relay file` use these, so each exists once in the repo. - The `relayfile-cli` binary now ships as per-platform optional dependencies — `@relayfile/cli-darwin-arm64`, `@relayfile/cli-darwin-x64`, `@relayfile/cli-linux-arm64`, `@relayfile/cli-linux-x64`, `@relayfile/cli-win32-arm64`, `@relayfile/cli-win32-x64` — matching the existing `@relayfile/mount-*` packages. npm installs only the one matching the host, so there is no install-time download: installs work offline and in CI, and integrity comes from the registry. `resolveRelayfileBinary` looks there first, then falls back to the previous chain (`RELAYFILE_CLI_BIN`, the `relayfile` package's `bin/`, a source-checkout build, `go run`, `PATH`). ### Fixed +- Retries use full jitter and treat `Retry-After` as a minimum delay instead of clamping it to the client backoff cap. + - A consumer that depends on `@relayfile/sdk` without also depending on `relayfile` could not resolve the CLI binary at all, because only the `relayfile` package's `postinstall` fetched it. `agent-relay file ` failed on a clean `npm i -g agent-relay` for exactly this reason. The per-platform packages above close that gap. - When no binary can be found, the error now names the `@relayfile/cli-*` package to install for the current platform and how the optional dependency goes missing, instead of surfacing a bare `ENOENT`. Mounted as a CLI surface this returns exit code 127 through the host's `io` rather than throwing. - `run(argv, io)` now passes the binary's output through as raw bytes instead of UTF-8-decoded strings, so `relayfile export --format tar --output -` survives being mounted as a CLI surface. diff --git a/packages/sdk/typescript/src/client.test.ts b/packages/sdk/typescript/src/client.test.ts index 572db2a4..c0ce2787 100644 --- a/packages/sdk/typescript/src/client.test.ts +++ b/packages/sdk/typescript/src/client.test.ts @@ -1466,6 +1466,44 @@ describe("RelayFileClient — existing methods", () => { expect(res.content).toBe('{"id":48291}'); expect(res.revision).toBe("rev_3"); }); + + it("revalidates cached content by hash and serves a 304 from the byte cache", async () => { + const contentHash = "a".repeat(64); + const payload: FileReadResponse = { + path: "/cached.txt", + revision: "rev_1", + contentType: "text/plain", + content: "cached", + contentHash, + }; + const f = vi.fn() + .mockResolvedValueOnce(jsonResponse(payload)) + .mockResolvedValueOnce(jsonResponse(undefined, 304)); + const client = makeClient(f); + + await client.readFile("ws_acme", payload.path); + const cached = await client.readFile("ws_acme", payload.path); + + expect(cached).toEqual(payload); + expect((f.mock.calls[1]![1] as RequestInit).headers).toMatchObject({ + "If-None-Match": `"${contentHash}"`, + }); + }); + + it("evicts least-recently-used content when the decoded byte cap is exceeded", async () => { + const files = ["a", "b", "a"].map((name) => ({ + path: `/${name}.txt`, revision: "rev_1", contentType: "text/plain", + content: name.repeat(4), contentHash: name.repeat(64), + } satisfies FileReadResponse)); + const f = vi.fn().mockImplementation(() => Promise.resolve(jsonResponse(files.shift()))); + const client = new RelayFileClient({ baseUrl: "https://relay.test", token: "tok", fetchImpl: f, readCache: { maxBytes: 4 } }); + + await client.readFile("ws", "/a.txt"); + await client.readFile("ws", "/b.txt"); + await client.readFile("ws", "/a.txt"); + + expect((f.mock.calls[2]![1] as RequestInit).headers).not.toHaveProperty("If-None-Match"); + }); }); // ---- writeFile ---- diff --git a/packages/sdk/typescript/src/client.ts b/packages/sdk/typescript/src/client.ts index 151e5aed..d425b033 100644 --- a/packages/sdk/typescript/src/client.ts +++ b/packages/sdk/typescript/src/client.ts @@ -375,48 +375,64 @@ const changeLogSettings = new WeakMap>>>(); const fileReadCaches = new WeakMap(); -const DEFAULT_READ_CACHE_TTL_MS = 5_000; -const DEFAULT_READ_CACHE_MAX_ENTRIES = 500; +const DEFAULT_READ_CACHE_MAX_BYTES = 32 * 1024 * 1024; interface ReadCacheEntry { value: FileReadResponse; - expiresAt: number; + bytes: number; } class FileReadCache { - private readonly ttlMs: number; - private readonly maxEntries: number; - private readonly entries = new Map(); + private readonly maxBytes: number; + private totalBytes = 0; + private readonly objects = new Map(); + private readonly paths = new Map(); private readonly inFlight = new Map>(); constructor(options?: RelayFileReadCacheOptions) { - this.ttlMs = options?.ttlMs ?? DEFAULT_READ_CACHE_TTL_MS; - this.maxEntries = options?.maxEntries ?? DEFAULT_READ_CACHE_MAX_ENTRIES; + this.maxBytes = Math.max(0, Math.floor(options?.maxBytes ?? DEFAULT_READ_CACHE_MAX_BYTES)); } get(key: string): FileReadResponse | undefined { - const entry = this.entries.get(key); + const hash = this.paths.get(key); + if (!hash) return undefined; + const entry = this.objects.get(hash); if (!entry) return undefined; - if (Date.now() > entry.expiresAt) { - this.entries.delete(key); - return undefined; - } + this.objects.delete(hash); + this.objects.set(hash, entry); return entry.value; } set(key: string, value: FileReadResponse): void { - if (this.entries.size >= this.maxEntries && !this.entries.has(key)) { - const oldest = this.entries.keys().next().value; - if (oldest !== undefined) { - this.entries.delete(oldest); - } + const hash = normalizeContentHash(value.contentHash) ?? `legacy:${key}`; + if (this.maxBytes === 0) return; + const bytes = value.encoding === "base64" + ? Math.floor(value.content.length * 3 / 4) - (value.content.endsWith("==") ? 2 : value.content.endsWith("=") ? 1 : 0) + : new TextEncoder().encode(value.content).byteLength; + const previous = this.objects.get(hash); + if (previous) { + this.totalBytes -= previous.bytes; + this.objects.delete(hash); + } + this.objects.set(hash, { value, bytes }); + this.totalBytes += bytes; + this.paths.set(key, hash); + while (this.totalBytes > this.maxBytes && this.objects.size > 0) { + const oldest = this.objects.keys().next().value; + if (oldest === undefined) break; + const evicted = this.objects.get(oldest)!; + this.objects.delete(oldest); + this.totalBytes -= evicted.bytes; } - this.entries.set(key, { value, expiresAt: Date.now() + this.ttlMs }); } evict(workspaceId: string, path: string): void { - this.entries.delete(`${workspaceId}:${path}`); - this.inFlight.delete(`${workspaceId}:${path}`); + for (const key of this.paths.keys()) { + if (key.startsWith(`${workspaceId}:`) && key.endsWith(`:${path}`)) this.paths.delete(key); + } + for (const key of this.inFlight.keys()) { + if (key.startsWith(`${workspaceId}:`) && key.endsWith(`:${path}`)) this.inFlight.delete(key); + } } getInFlight(key: string): Promise | undefined { @@ -441,6 +457,11 @@ class FileReadCache { } } +function normalizeContentHash(hash: string | undefined): string | undefined { + const value = hash?.trim().replace(/^sha256:/i, "").toLowerCase(); + return value && /^[a-f0-9]{64}$/.test(value) ? value : undefined; +} + function getFileReadCache(client: RelayFileClient): FileReadCache | false { const cached = fileReadCaches.get(client); if (cached !== undefined) return cached; @@ -1758,23 +1779,25 @@ export class RelayFileClient { const cacheRaw = getFileReadCache(this); // Skip cache for fork-scoped reads (isolated state) and when cache is disabled. const cache: FileReadCache | undefined = cacheRaw !== false ? cacheRaw : undefined; - const cacheKey = (cache && !input.forkId) ? `${input.workspaceId}:${input.path}` : undefined; + const cacheKey = cache ? `${input.workspaceId}:${input.forkId ?? ""}:${input.path}` : undefined; if (cache && cacheKey) { - const hit = cache.get(cacheKey); - if (hit) return hit; const pending = cache.getInFlight(cacheKey); if (pending) return pending; } const query = buildQuery({ path: input.path, forkId: input.forkId }); - const fetch = this.request({ + const cached = cache && cacheKey ? cache.get(cacheKey) : undefined; + if (cached && !normalizeContentHash(cached.contentHash)) return cached; + const fetch = this.performRequest({ method: "GET", path: `/v1/workspaces/${encodeURIComponent(input.workspaceId)}/fs/file${query}`, + headers: cached?.contentHash ? { "If-None-Match": `"${normalizeContentHash(cached.contentHash) ?? cached.contentHash}"` } : undefined, correlationId: input.correlationId, signal: input.signal, - tokenOverride: (input as ReadFileInput & { token?: string }).token - }); + tokenOverride: (input as ReadFileInput & { token?: string }).token, + allowNotModified: Boolean(cached) + }).then(async (response) => response.status === 304 && cached ? cached : await this.readPayload(response) as FileReadResponse); if (cache && cacheKey) { cache.setInFlight(cacheKey, fetch); @@ -2845,6 +2868,7 @@ export class RelayFileClient { signal?: AbortSignal; accept?: string; tokenOverride?: string; + allowNotModified?: boolean; }): Promise { const existingCorrelationId = getHeaderValue(params.headers, "X-Correlation-Id"); const correlationId = existingCorrelationId ?? params.correlationId ?? generateCorrelationId(); @@ -2893,7 +2917,7 @@ export class RelayFileClient { throw error; } - if (response.ok) { + if (response.ok || (params.allowNotModified && response.status === 304)) { return response; } @@ -2948,10 +2972,8 @@ export class RelayFileClient { // unchanged (bounded by our own `maxDelayMs`). Only when the header is // absent or unparseable do we consult the body below, so a body hint never // silently overrides a shorter, explicit header the server already sent. - const retryAfterMs = this.parseRetryAfterMs(retryAfterHeader); - if (retryAfterMs !== null) { - return Math.min(this.retryOptions.maxDelayMs, retryAfterMs); - } + const retryAfterMs = this.parseRetryAfterMs(retryAfterHeader) + ?? this.parseRetryAfterSecondsFromBody(payload); // With no usable header, a 429 body can still advertise an explicit // backpressure delay as `details.retryAfterSeconds` (e.g. `workspace_busy` // when the workspace durable object is overloaded, or `queue_full`). Honor @@ -2959,15 +2981,10 @@ export class RelayFileClient { // to `maxDelayMs`, which governs our own exponential backoff. Truncating it // (maxDelayMs defaults to 2s vs. a typical 5s advertised delay) retries // into the still-busy resource and exhausts the retry budget. - const advertisedMs = this.parseRetryAfterSecondsFromBody(payload); - if (advertisedMs !== null) { - return Math.max(0, Math.min(RETRY_AFTER_MAX_MS, advertisedMs)); - } const backoff = this.retryOptions.baseDelayMs * Math.pow(2, Math.max(0, retryAttempt - 1)); const capped = Math.min(this.retryOptions.maxDelayMs, backoff); - const jitter = this.retryOptions.jitterRatio; - const factor = 1 + (Math.random() * 2 - 1) * jitter; - return Math.max(0, Math.round(capped * factor)); + const jittered = Math.round(Math.random() * capped); + return Math.max(retryAfterMs ?? 0, jittered); } /** diff --git a/packages/sdk/typescript/src/types.ts b/packages/sdk/typescript/src/types.ts index 07204890..883932e7 100644 --- a/packages/sdk/typescript/src/types.ts +++ b/packages/sdk/typescript/src/types.ts @@ -65,9 +65,11 @@ export interface ContentIdentity { } export interface RelayFileReadCacheOptions { - /** Cache TTL in ms. Default: 5000. */ + /** Maximum decoded content bytes retained by the content-addressed LRU. Default: 32 MiB. */ + maxBytes?: number; + /** @deprecated Content-hash entries are revalidated on every read. */ ttlMs?: number; - /** Max cached entries before LRU eviction. Default: 500. */ + /** @deprecated Eviction is byte-capped; use maxBytes. */ maxEntries?: number; } @@ -76,6 +78,8 @@ export interface FileReadResponse { revision: string; contentType: string; content: string; + /** SHA-256 content identity returned by Relayfile. */ + contentHash?: string; encoding?: "utf-8" | "base64"; provider?: string; providerObjectId?: string; diff --git a/summary.md b/summary.md new file mode 100644 index 00000000..12580415 --- /dev/null +++ b/summary.md @@ -0,0 +1,28 @@ +# PR summary + +## What changed + +- Replaced the TypeScript path/TTL cache with a configurable 32 MiB byte-capped content-addressed LRU, and added equivalent sync/async Python caching. Cached reads send `If-None-Match` using `contentHash` and serve `304` responses locally; caching can be disabled. +- Added a verified, atomic, permission-restricted mount object store at `~/.relayfile/cache/objects/`. Authorized bootstrap tree hashes are checked before body reads, so later mounts can materialize identical content without transferring it again. +- Reduced bootstrap and incremental point-read concurrency from 16 (previous environment ceiling 64) to a hard ceiling of 4. +- Changed SDK and Go transport retry delay calculation to full jitter while treating `Retry-After` as a minimum, including HTTP-date parsing already supported by each client. +- Made `relayfile listen` retry initial 429/503 handshakes, retain the last processed event cursor, and reconnect from that cursor with full jitter. Applied the same handshake `Retry-After`/jitter behavior to mount and FUSE WebSocket reconnects. +- Added overload `details.reason` propagation to Go `HTTPError`, updated SDK parity metadata, tests, and configuration/changelog documentation. + +## Scope decisions + +- Included `internal/mountfuse/wsinvalidate.go` because it is a shipped `/fs/ws` reconnect loop with the same fleet lockstep risk. +- Deferred SDK setup/control-plane retries, `packages/agents` reconnects, mount outbox writes, and CLI polite polling: they are one-shot, write-path, or polling flows rather than the file-read fan-in and event reconnect paths addressed here. +- No server handler or provider mutation path changed. The digest runtime contract therefore does not apply; filesystem-event emission and digest regeneration behavior are unchanged. + +## Validation + +- TypeScript SDK build and typecheck pass. All 113 relevant `client.test.ts` assertions pass; one pre-existing environment assertion expects Node to lack global `ErrorEvent`, which is false on the installed Node 25 runtime. +- Python SDK: 97 tests pass. +- Targeted Go cache/retry/WebSocket/listen tests pass for `internal/mountsync`, `internal/mountfuse`, and `cmd/relayfile-cli`. +- `scripts/check-contract-surface.sh` passes, including SDK parity. +- The complete TypeScript package suite additionally requires the downloaded Go toolchain to be on the subprocess `PATH`; unrelated launcher timing tests remain environment-sensitive. + +## Contract note + +The clients deliberately use the response body's `contentHash` for object identity. The currently documented server `ETag` is a revision identifier; the separate server change must redefine it to the quoted content hash and add `If-None-Match`/`304` handling before conditional requests can save network bodies against older servers. From 4f6ec87eecb8955559c40ba24c34ea55f6059ed8 Mon Sep 17 00:00:00 2001 From: Relayflow Date: Sun, 27 Sep 2026 20:11:37 +0000 Subject: [PATCH 2/6] fix(sdk): tolerate foreground mount state publication race --- .../completed/2026-09/traj_0d21k8ac1mdg.json | 53 +++++++++++++++++++ .../completed/2026-09/traj_0d21k8ac1mdg.md | 31 +++++++++++ .trajectories/index.json | 9 +++- packages/sdk/typescript/src/mount-launcher.ts | 14 +++++ 4 files changed, 106 insertions(+), 1 deletion(-) create mode 100644 .trajectories/completed/2026-09/traj_0d21k8ac1mdg.json create mode 100644 .trajectories/completed/2026-09/traj_0d21k8ac1mdg.md diff --git a/.trajectories/completed/2026-09/traj_0d21k8ac1mdg.json b/.trajectories/completed/2026-09/traj_0d21k8ac1mdg.json new file mode 100644 index 00000000..e6f70634 --- /dev/null +++ b/.trajectories/completed/2026-09/traj_0d21k8ac1mdg.json @@ -0,0 +1,53 @@ +{ + "id": "traj_0d21k8ac1mdg", + "version": 1, + "task": { + "title": "Repair SDK/mount cache and retry changes from branch check failures" + }, + "status": "completed", + "startedAt": "2026-09-27T20:02:30.999Z", + "completedAt": "2026-09-27T20:11:31.271Z", + "agents": [ + { + "name": "default", + "role": "lead", + "joinedAt": "2026-09-27T20:08:43.980Z" + } + ], + "chapters": [ + { + "id": "chap_tntqc5ip89bb", + "title": "Work", + "agentName": "default", + "startedAt": "2026-09-27T20:08:43.980Z", + "endedAt": "2026-09-27T20:11:31.271Z", + "events": [ + { + "ts": 1790539723981, + "type": "decision", + "content": "Poll for foreground mount state after exit 0 until readiness timeout: Poll for foreground mount state after exit 0 until readiness timeout", + "raw": { + "question": "Poll for foreground mount state after exit 0 until readiness timeout", + "chosen": "Poll for foreground mount state after exit 0 until readiness timeout", + "alternatives": [], + "reasoning": "The child can exit immediately after publishing state; one overlapping read can observe an absent/partial file and previously caused a false early-exit failure." + }, + "significance": "high" + } + ] + } + ], + "retrospective": { + "summary": "Provisioned CI toolchains locally, fixed foreground mount readiness publication race, and passed the full repository check.", + "approach": "Standard approach", + "confidence": 0.95 + }, + "commits": [], + "filesChanged": [], + "projectId": "AgentWorkforce/relayfile", + "tags": [], + "_trace": { + "startRef": "cafa926c6f07be6dfcf4a04fadd026db93902b84", + "endRef": "cafa926c6f07be6dfcf4a04fadd026db93902b84" + } +} \ No newline at end of file diff --git a/.trajectories/completed/2026-09/traj_0d21k8ac1mdg.md b/.trajectories/completed/2026-09/traj_0d21k8ac1mdg.md new file mode 100644 index 00000000..f397e2db --- /dev/null +++ b/.trajectories/completed/2026-09/traj_0d21k8ac1mdg.md @@ -0,0 +1,31 @@ +# Trajectory: Repair SDK/mount cache and retry changes from branch check failures + +> **Status:** ✅ Completed +> **Confidence:** 95% +> **Started:** September 27, 2026 at 08:02 PM +> **Completed:** September 27, 2026 at 08:11 PM + +--- + +## Summary + +Provisioned CI toolchains locally, fixed foreground mount readiness publication race, and passed the full repository check. + +**Approach:** Standard approach + +--- + +## Key Decisions + +### Poll for foreground mount state after exit 0 until readiness timeout +- **Chose:** Poll for foreground mount state after exit 0 until readiness timeout +- **Reasoning:** The child can exit immediately after publishing state; one overlapping read can observe an absent/partial file and previously caused a false early-exit failure. + +--- + +## Chapters + +### 1. Work +*Agent: default* + +- Poll for foreground mount state after exit 0 until readiness timeout: Poll for foreground mount state after exit 0 until readiness timeout diff --git a/.trajectories/index.json b/.trajectories/index.json index 88975238..c104f3cb 100644 --- a/.trajectories/index.json +++ b/.trajectories/index.json @@ -1,6 +1,6 @@ { "version": 1, - "lastUpdated": "2026-09-27T20:01:45.824Z", + "lastUpdated": "2026-09-27T20:11:31.329Z", "trajectories": { "traj_4pvrlmqfnzng": { "title": "Review PR #278 in AgentWorkforce/relayfile", @@ -414,6 +414,13 @@ "startedAt": "2026-09-27T19:55:27.193Z", "completedAt": "2026-09-27T20:01:45.765Z", "path": "/home/daytona/.relayflow-v2-supervisor/durable/repository/.trajectories/completed/2026-09/traj_6qsrb702mx1s.json" + }, + "traj_0d21k8ac1mdg": { + "title": "Repair SDK/mount cache and retry changes from branch check failures", + "status": "completed", + "startedAt": "2026-09-27T20:02:30.999Z", + "completedAt": "2026-09-27T20:11:31.271Z", + "path": "/home/daytona/.relayflow-v2-supervisor/durable/repository/.trajectories/completed/2026-09/traj_0d21k8ac1mdg.json" } } } \ No newline at end of file diff --git a/packages/sdk/typescript/src/mount-launcher.ts b/packages/sdk/typescript/src/mount-launcher.ts index e3bf2ae9..71fa94a4 100644 --- a/packages/sdk/typescript/src/mount-launcher.ts +++ b/packages/sdk/typescript/src/mount-launcher.ts @@ -379,6 +379,20 @@ class RelayfileMountProcessInstance await this.restartOnceMount() continue } + // A successful foreground --once process can exit immediately after + // publishing its final state. If the readiness poll overlaps that + // publication, give the state file the remainder of the existing + // readiness budget to become observable instead of reporting a false + // early-exit failure. + if ( + this.input.background === false && + this.exitCode === 0 && + !this.stopping && + this.now() < timeoutAt + ) { + await delay(this.readyPollIntervalMs) + continue + } throw this.buildEarlyExitError() } From 38888ddb3478853e773f5de2a0d67f4d3a52059a Mon Sep 17 00:00:00 2001 From: Relayflow Date: Sun, 27 Sep 2026 20:14:48 +0000 Subject: [PATCH 3/6] Keep relayflow working files out of the change --- summary.md | 28 ---------------------------- 1 file changed, 28 deletions(-) delete mode 100644 summary.md diff --git a/summary.md b/summary.md deleted file mode 100644 index 12580415..00000000 --- a/summary.md +++ /dev/null @@ -1,28 +0,0 @@ -# PR summary - -## What changed - -- Replaced the TypeScript path/TTL cache with a configurable 32 MiB byte-capped content-addressed LRU, and added equivalent sync/async Python caching. Cached reads send `If-None-Match` using `contentHash` and serve `304` responses locally; caching can be disabled. -- Added a verified, atomic, permission-restricted mount object store at `~/.relayfile/cache/objects/`. Authorized bootstrap tree hashes are checked before body reads, so later mounts can materialize identical content without transferring it again. -- Reduced bootstrap and incremental point-read concurrency from 16 (previous environment ceiling 64) to a hard ceiling of 4. -- Changed SDK and Go transport retry delay calculation to full jitter while treating `Retry-After` as a minimum, including HTTP-date parsing already supported by each client. -- Made `relayfile listen` retry initial 429/503 handshakes, retain the last processed event cursor, and reconnect from that cursor with full jitter. Applied the same handshake `Retry-After`/jitter behavior to mount and FUSE WebSocket reconnects. -- Added overload `details.reason` propagation to Go `HTTPError`, updated SDK parity metadata, tests, and configuration/changelog documentation. - -## Scope decisions - -- Included `internal/mountfuse/wsinvalidate.go` because it is a shipped `/fs/ws` reconnect loop with the same fleet lockstep risk. -- Deferred SDK setup/control-plane retries, `packages/agents` reconnects, mount outbox writes, and CLI polite polling: they are one-shot, write-path, or polling flows rather than the file-read fan-in and event reconnect paths addressed here. -- No server handler or provider mutation path changed. The digest runtime contract therefore does not apply; filesystem-event emission and digest regeneration behavior are unchanged. - -## Validation - -- TypeScript SDK build and typecheck pass. All 113 relevant `client.test.ts` assertions pass; one pre-existing environment assertion expects Node to lack global `ErrorEvent`, which is false on the installed Node 25 runtime. -- Python SDK: 97 tests pass. -- Targeted Go cache/retry/WebSocket/listen tests pass for `internal/mountsync`, `internal/mountfuse`, and `cmd/relayfile-cli`. -- `scripts/check-contract-surface.sh` passes, including SDK parity. -- The complete TypeScript package suite additionally requires the downloaded Go toolchain to be on the subprocess `PATH`; unrelated launcher timing tests remain environment-sensitive. - -## Contract note - -The clients deliberately use the response body's `contentHash` for object identity. The currently documented server `ETag` is a revision identifier; the separate server change must redefine it to the quoted content hash and add `If-None-Match`/`304` handling before conditional requests can save network bodies against older servers. From 35205e7c8a7a125b3d367b9e415f67318e414e50 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 28 Sep 2026 07:54:41 +0200 Subject: [PATCH 4/6] fix(cache): keep per-path metadata out of content-addressed stores Address review feedback on the client content-hash caches: - SDK (TS + Python): dedupe only content bytes by (hash, encoding); keep revision/path/semantics per cache key so two paths with identical bytes never swap metadata on a 304. - SDK (TS): never store hashless responses; they cannot be revalidated with If-None-Match, so they always go to the server instead of being served from cache indefinitely. - mount: the persistent object store now holds raw bytes only (no path, revision or content type from the workspace that first fetched them) and is byte-capped (1 GiB default) with LRU eviction by mtime. Co-Authored-By: Claude Opus 5.5 (1M context) --- internal/mountsync/object_cache.go | 154 +++++++++++++++++--- internal/mountsync/object_cache_test.go | 145 +++++++++++++++++- internal/mountsync/syncer.go | 22 ++- packages/local-mount/CHANGELOG.md | 2 +- packages/sdk/python/src/relayfile/client.py | 67 +++++++-- packages/sdk/python/tests/test_client.py | 35 +++++ packages/sdk/typescript/CHANGELOG.md | 2 +- packages/sdk/typescript/src/client.test.ts | 61 ++++++-- packages/sdk/typescript/src/client.ts | 95 ++++++++---- 9 files changed, 498 insertions(+), 85 deletions(-) diff --git a/internal/mountsync/object_cache.go b/internal/mountsync/object_cache.go index 2643951f..2b96a0c7 100644 --- a/internal/mountsync/object_cache.go +++ b/internal/mountsync/object_cache.go @@ -4,16 +4,36 @@ import ( "crypto/sha256" "encoding/base64" "encoding/hex" - "encoding/json" "os" "path/filepath" "regexp" + "sort" "strings" + "sync" + "time" + "unicode/utf8" ) var objectHashPattern = regexp.MustCompile(`^[a-f0-9]{64}$`) -type objectCache struct{ root string } +// defaultObjectCacheMaxBytes bounds the persistent content-addressed store. +// Objects are evicted least-recently-used (by mtime, refreshed on hit) once +// the directory grows past the cap. +const defaultObjectCacheMaxBytes int64 = 1 << 30 + +// objectCache is a content-addressed store of raw file bytes keyed by their +// SHA-256. It deliberately stores bytes only: path, revision, content type and +// every other piece of per-workspace metadata comes from the server response +// (tree entry or read) that authorized the caller to see the object, so the +// store never retains or replays metadata from another workspace. +type objectCache struct { + root string + maxBytes int64 + + mu sync.Mutex + sized bool + totalBytes int64 +} func defaultObjectCache() *objectCache { home, err := os.UserHomeDir() @@ -31,29 +51,63 @@ func normalizeObjectHash(value string) string { return value } -func (c *objectCache) get(hash string) (RemoteFile, bool) { - var file RemoteFile +func (c *objectCache) limit() int64 { + if c.maxBytes > 0 { + return c.maxBytes + } + return defaultObjectCacheMaxBytes +} + +// get returns the cached bytes for hash as a RemoteFile carrying only +// content, encoding and hash. Callers fill in the path metadata from the +// authorizing server response. encoding is the encoding the server reported +// for the path ("" when unknown); binary content is always base64-encoded. +func (c *objectCache) get(hash, encoding string) (RemoteFile, bool) { hash = normalizeObjectHash(hash) if c == nil || hash == "" { - return file, false + return RemoteFile{}, false } - data, err := os.ReadFile(filepath.Join(c.root, hash)) - if err != nil || json.Unmarshal(data, &file) != nil || normalizeObjectHash(file.ContentHash) != hash || remoteFileHash(file) != hash { + path := filepath.Join(c.root, hash) + data, err := os.ReadFile(path) + if err != nil { return RemoteFile{}, false } + sum := sha256.Sum256(data) + if hex.EncodeToString(sum[:]) != hash { + return RemoteFile{}, false + } + now := time.Now() + _ = os.Chtimes(path, now, now) + file := RemoteFile{ContentHash: hash} + if normalizeEncoding(encoding) == "base64" || !utf8.Valid(data) { + file.Content = base64.StdEncoding.EncodeToString(data) + file.Encoding = "base64" + } else { + file.Content = string(data) + } return file, true } func (c *objectCache) put(file RemoteFile) { hash := normalizeObjectHash(file.ContentHash) - if c == nil || hash == "" || remoteFileHash(file) != hash { + if c == nil || hash == "" { return } - if os.MkdirAll(c.root, 0o700) != nil { + data, ok := remoteFileBytes(file) + if !ok || int64(len(data)) > c.limit() { return } - data, err := json.Marshal(file) - if err != nil { + sum := sha256.Sum256(data) + if hex.EncodeToString(sum[:]) != hash { + return + } + target := filepath.Join(c.root, hash) + if _, err := os.Stat(target); err == nil { + now := time.Now() + _ = os.Chtimes(target, now, now) + return + } + if os.MkdirAll(c.root, 0o700) != nil { return } tmp, err := os.CreateTemp(c.root, ".object-*") @@ -69,20 +123,82 @@ func (c *objectCache) put(file RemoteFile) { if closeErr := tmp.Close(); err == nil { err = closeErr } - if err == nil { - _ = os.Rename(name, filepath.Join(c.root, hash)) + if err != nil || os.Rename(name, target) != nil { + return + } + c.added(int64(len(data))) +} + +// added accounts for a newly stored object and evicts least-recently-used +// objects once the store exceeds its byte cap. The running total is seeded +// from one directory scan and re-measured on every eviction, so objects +// written by concurrent mounts sharing the directory are accounted for. +func (c *objectCache) added(size int64) { + c.mu.Lock() + defer c.mu.Unlock() + if !c.sized { + c.totalBytes = 0 + for _, object := range c.listObjects() { + c.totalBytes += object.size + } + c.sized = true + } else { + c.totalBytes += size + } + if c.totalBytes <= c.limit() { + return + } + objects := c.listObjects() + total := int64(0) + for _, object := range objects { + total += object.size + } + sort.Slice(objects, func(i, j int) bool { return objects[i].modTime.Before(objects[j].modTime) }) + for _, object := range objects { + if total <= c.limit() { + break + } + // Removing an object another process is reading is safe: the reader + // either already holds the bytes (and re-verifies the hash) or misses. + if err := os.Remove(filepath.Join(c.root, object.name)); err == nil || os.IsNotExist(err) { + total -= object.size + } + } + c.totalBytes = total +} + +type objectCacheEntry struct { + name string + size int64 + modTime time.Time +} + +func (c *objectCache) listObjects() []objectCacheEntry { + entries, err := os.ReadDir(c.root) + if err != nil { + return nil + } + objects := make([]objectCacheEntry, 0, len(entries)) + for _, entry := range entries { + if !entry.Type().IsRegular() || !objectHashPattern.MatchString(entry.Name()) { + continue + } + info, err := entry.Info() + if err != nil { + continue + } + objects = append(objects, objectCacheEntry{name: entry.Name(), size: info.Size(), modTime: info.ModTime()}) } + return objects } -func remoteFileHash(file RemoteFile) string { - content := []byte(file.Content) +func remoteFileBytes(file RemoteFile) ([]byte, bool) { if strings.EqualFold(file.Encoding, "base64") { decoded, err := base64.StdEncoding.DecodeString(file.Content) if err != nil { - return "" + return nil, false } - content = decoded + return decoded, true } - sum := sha256.Sum256(content) - return hex.EncodeToString(sum[:]) + return []byte(file.Content), true } diff --git a/internal/mountsync/object_cache_test.go b/internal/mountsync/object_cache_test.go index e42d020e..da80505a 100644 --- a/internal/mountsync/object_cache_test.go +++ b/internal/mountsync/object_cache_test.go @@ -1,9 +1,13 @@ package mountsync import ( + "context" + "encoding/base64" "os" "path/filepath" + "strings" "testing" + "time" ) func TestObjectCacheRoundTripAndRejectsCorruption(t *testing.T) { @@ -11,7 +15,7 @@ func TestObjectCacheRoundTripAndRejectsCorruption(t *testing.T) { file := RemoteFile{Content: "hello", ContentHash: "2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824"} cache.put(file) - got, ok := cache.get(file.ContentHash) + got, ok := cache.get(file.ContentHash, "") if !ok || got.Content != "hello" { t.Fatalf("cache get = %#v, %v", got, ok) } @@ -23,10 +27,145 @@ func TestObjectCacheRoundTripAndRejectsCorruption(t *testing.T) { if info.Mode().Perm() != 0o600 { t.Fatalf("cache mode = %v", info.Mode().Perm()) } - if err := os.WriteFile(path, []byte(`{"content":"tampered","contentHash":"`+file.ContentHash+`"}`), 0o600); err != nil { + if err := os.WriteFile(path, []byte("tampered"), 0o600); err != nil { t.Fatal(err) } - if _, ok := cache.get(file.ContentHash); ok { + if _, ok := cache.get(file.ContentHash, ""); ok { t.Fatal("corrupt object was accepted") } } + +func TestObjectCacheStoresBytesOnlyNotWorkspaceMetadata(t *testing.T) { + cache := &objectCache{root: filepath.Join(t.TempDir(), "objects")} + file := RemoteFile{ + Path: "/workspace-a/secret-name.txt", + Revision: "rev_workspace_a", + ContentType: "application/x-workspace-a", + Content: "hello", + ContentHash: "sha256:2cf24dba5fb0a30e26e83b2ac5b9e29e1b161e5c1fa7425e73043362938b9824", + } + cache.put(file) + + data, err := os.ReadFile(filepath.Join(cache.root, normalizeObjectHash(file.ContentHash))) + if err != nil { + t.Fatal(err) + } + if string(data) != "hello" { + t.Fatalf("object store must hold raw bytes only, got %q", data) + } + got, ok := cache.get(file.ContentHash, "") + if !ok { + t.Fatal("expected cache hit") + } + if got.Path != "" || got.Revision != "" || got.ContentType != "" { + t.Fatalf("cache replayed another response's metadata: %#v", got) + } +} + +func TestObjectCacheBase64RoundTrip(t *testing.T) { + cache := &objectCache{root: filepath.Join(t.TempDir(), "objects")} + raw := []byte{0xff, 0x00, 0xfe} + hash := hashBytes(raw) + cache.put(RemoteFile{Content: base64.StdEncoding.EncodeToString(raw), Encoding: "base64", ContentHash: hash}) + + got, ok := cache.get(hash, "") + if !ok || got.Encoding != "base64" { + t.Fatalf("binary object should round-trip as base64: %#v, %v", got, ok) + } + decoded, err := base64.StdEncoding.DecodeString(got.Content) + if err != nil || string(decoded) != string(raw) { + t.Fatalf("decoded = %v, %v", decoded, err) + } +} + +func TestObjectCacheEvictsLeastRecentlyUsedOverByteCap(t *testing.T) { + cache := &objectCache{root: filepath.Join(t.TempDir(), "objects"), maxBytes: 10} + put := func(content string) string { + hash := hashBytes([]byte(content)) + cache.put(RemoteFile{Content: content, ContentHash: hash}) + return hash + } + a := put("aaaa") + b := put("bbbb") + old := time.Now().Add(-time.Hour) + if err := os.Chtimes(filepath.Join(cache.root, a), old, old); err != nil { + t.Fatal(err) + } + if err := os.Chtimes(filepath.Join(cache.root, b), old.Add(time.Minute), old.Add(time.Minute)); err != nil { + t.Fatal(err) + } + // Touch a so b becomes the least recently used object. + if _, ok := cache.get(a, ""); !ok { + t.Fatal("expected a to be cached") + } + c := put("cccc") + + if _, ok := cache.get(b, ""); ok { + t.Fatal("least recently used object should have been evicted") + } + for _, hash := range []string{a, c} { + if _, ok := cache.get(hash, ""); !ok { + t.Fatalf("object %s should remain cached", hash) + } + } + var total int64 + for _, object := range cache.listObjects() { + total += object.size + } + if total > cache.maxBytes { + t.Fatalf("store holds %d bytes, cap %d", total, cache.maxBytes) + } + + // An object larger than the whole cap is never stored. + big := strings.Repeat("x", 11) + if _, ok := cache.get(put(big), ""); ok { + t.Fatal("object larger than the cap must not be stored") + } +} + +func TestSecondMountBootstrapsCachedObjectsWithoutBodyReads(t *testing.T) { + cacheRoot := filepath.Join(t.TempDir(), "objects") + shared := "shared bytes" + sharedHash := hashBytes([]byte(shared)) + + first := &fakeClient{files: map[string]RemoteFile{ + "/docs/a.txt": {Path: "/docs/a.txt", Revision: "rev_first", ContentType: "text/plain", Content: shared, ContentHash: sharedHash}, + }} + firstSyncer, err := NewSyncer(first, SyncerOptions{WorkspaceID: "ws_first", RemoteRoot: "/", LocalRoot: t.TempDir(), ObjectCacheRoot: cacheRoot}) + if err != nil { + t.Fatal(err) + } + if err := firstSyncer.SyncOnce(context.Background()); err != nil { + t.Fatalf("first sync: %v", err) + } + if first.requestedReadCalls() == 0 { + t.Fatal("first mount should fetch bodies from the server") + } + + other := "other bytes" + second := &fakeClient{files: map[string]RemoteFile{ + "/notes/b.txt": {Path: "/notes/b.txt", Revision: "rev_second", ContentType: "text/plain", Content: shared, ContentHash: sharedHash}, + "/notes/c.txt": {Path: "/notes/c.txt", Revision: "rev_c", ContentType: "text/plain", Content: other, ContentHash: hashBytes([]byte(other))}, + }} + localDir := t.TempDir() + secondSyncer, err := NewSyncer(second, SyncerOptions{WorkspaceID: "ws_second", RemoteRoot: "/", LocalRoot: localDir, ObjectCacheRoot: cacheRoot}) + if err != nil { + t.Fatal(err) + } + if err := secondSyncer.SyncOnce(context.Background()); err != nil { + t.Fatalf("second sync: %v", err) + } + if got := second.readFileCallsByPath["/notes/b.txt"]; got != 0 { + t.Fatalf("cached object should not be re-read, got %d reads", got) + } + if got := second.readFileCallsByPath["/notes/c.txt"]; got != 1 { + t.Fatalf("uncached object should be read once, got %d", got) + } + data, err := os.ReadFile(filepath.Join(localDir, "notes", "b.txt")) + if err != nil || string(data) != shared { + t.Fatalf("materialized content = %q, %v", data, err) + } + if got := secondSyncer.state.Files["/notes/b.txt"].Revision; got != "rev_second" { + t.Fatalf("revision must come from this mount's tree entry, got %q", got) + } +} diff --git a/internal/mountsync/syncer.go b/internal/mountsync/syncer.go index 4296869d..ebce9742 100644 --- a/internal/mountsync/syncer.go +++ b/internal/mountsync/syncer.go @@ -8032,6 +8032,20 @@ func (s *Syncer) readBootstrapFiles(ctx context.Context, jobs []bootstrapReadJob return results } +// cachedBootstrapFile materializes a bootstrap job from the local +// content-addressed store when the authorizing tree entry names a hash the +// store already holds. The store holds bytes only; every piece of path +// metadata comes from this mount's own tree entry. +func (s *Syncer) cachedBootstrapFile(job bootstrapReadJob) (RemoteFile, bool) { + cached, ok := s.objectCache.get(job.Entry.ContentHash, job.Entry.Encoding) + if !ok { + return RemoteFile{}, false + } + cached.Path, cached.Revision = job.RemotePath, job.Entry.Revision + cached.Type, cached.Target, cached.Mode = job.Entry.Type, job.Entry.Target, job.Entry.Mode + return cached, true +} + // readBootstrapFilesEach dispatches bulk reads and oversized point reads in // index order. Oversized point reads are each passed as a singleton, so their // response body is released before the next oversized read begins. The @@ -8044,13 +8058,11 @@ func (s *Syncer) readBootstrapFilesEach(ctx context.Context, jobs []bootstrapRea } remaining := make([]bootstrapReadJob, 0, len(jobs)) for _, job := range jobs { - cached, ok := s.objectCache.get(job.Entry.ContentHash) + cached, ok := s.cachedBootstrapFile(job) if !ok { remaining = append(remaining, job) continue } - cached.Path, cached.Revision = job.RemotePath, job.Entry.Revision - cached.Type, cached.Target, cached.Mode = job.Entry.Type, job.Entry.Target, job.Entry.Mode prog.touch() if err := handle(bootstrapReadResult{Index: job.Index, RemotePath: job.RemotePath, File: cached}); err != nil { return err @@ -8306,9 +8318,7 @@ func (s *Syncer) readBootstrapFilesIndividuallyBatchEach(ctx context.Context, jo go func() { defer wg.Done() for job := range jobCh { - if cached, ok := s.objectCache.get(job.Entry.ContentHash); ok { - cached.Path, cached.Revision = job.RemotePath, job.Entry.Revision - cached.Type, cached.Target, cached.Mode = job.Entry.Type, job.Entry.Target, job.Entry.Mode + if cached, ok := s.cachedBootstrapFile(job); ok { prog.touch() resultCh <- bootstrapReadResult{Index: job.Index, RemotePath: job.RemotePath, File: cached} continue diff --git a/packages/local-mount/CHANGELOG.md b/packages/local-mount/CHANGELOG.md index f63571c8..61faac61 100644 --- a/packages/local-mount/CHANGELOG.md +++ b/packages/local-mount/CHANGELOG.md @@ -6,7 +6,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] -- Bootstrap point reads are capped at four concurrent requests and reuse verified objects from `~/.relayfile/cache/objects/`. HTTP and WebSocket overload retries use full jitter without shortening `Retry-After`. +- Bootstrap point reads are capped at four concurrent requests and reuse verified objects from `~/.relayfile/cache/objects/`. The store holds raw bytes only (no paths or other workspace metadata) and is capped at 1 GiB with least-recently-used eviction. HTTP and WebSocket overload retries use full jitter without shortening `Retry-After`. ### Fixed diff --git a/packages/sdk/python/src/relayfile/client.py b/packages/sdk/python/src/relayfile/client.py index 7465916f..d79eafb8 100644 --- a/packages/sdk/python/src/relayfile/client.py +++ b/packages/sdk/python/src/relayfile/client.py @@ -49,28 +49,43 @@ class RelayFileReadCacheOptions: class _FileReadCache: + """Content bytes are deduplicated by (hash, encoding); all other response + fields (revision, path, ...) are per-path metadata kept per cache key, so + two paths with identical bytes never swap metadata.""" + def __init__(self, options: RelayFileReadCacheOptions | None) -> None: self.max_bytes = max(0, (options or RelayFileReadCacheOptions()).max_bytes) self.total_bytes = 0 - self.objects: OrderedDict[str, tuple[dict[str, Any], int]] = OrderedDict() - self.paths: dict[str, str] = {} + # object key -> (content, decoded size, referencing path keys) + self.objects: OrderedDict[str, tuple[str, int, set[str]]] = OrderedDict() + # path key -> (metadata without content, object key) + self.paths: dict[str, tuple[dict[str, Any], str]] = {} self.lock = threading.RLock() def get(self, key: str) -> dict[str, Any] | None: with self.lock: - digest = self.paths.get(key) - item = self.objects.get(digest or "") + entry = self.paths.get(key) + if entry is None: + return None + meta, object_key = entry + item = self.objects.get(object_key) if item is None: + self.paths.pop(key, None) return None - self.objects.move_to_end(digest) - return dict(item[0]) + self.objects.move_to_end(object_key) + return {**meta, "content": item[0]} def put(self, key: str, value: dict[str, Any]) -> None: + with self.lock: + self._delete_path(key) digest = str(value.get("contentHash", "")).removeprefix("sha256:").lower() if len(digest) != 64 or any(c not in "0123456789abcdef" for c in digest): return + if self.max_bytes == 0: + return content = str(value.get("content", "")) - if value.get("encoding") == "base64": + encoding = str(value.get("encoding") or "utf-8") + if encoding == "base64": import base64 try: size = len(base64.b64decode(content, validate=True)) @@ -78,16 +93,38 @@ def put(self, key: str, value: dict[str, Any]) -> None: return else: size = len(content.encode()) + object_key = f"{digest}:{encoding}" + meta = {k: v for k, v in value.items() if k != "content"} with self.lock: - old = self.objects.pop(digest, None) - if old: - self.total_bytes -= old[1] - self.objects[digest] = (dict(value), size) - self.paths[key] = digest - self.total_bytes += size + item = self.objects.pop(object_key, None) + if item is None: + item = (content, size, set()) + self.total_bytes += size + item[2].add(key) + self.objects[object_key] = item + self.paths[key] = (meta, object_key) while self.total_bytes > self.max_bytes and self.objects: - _, (_, removed_size) = self.objects.popitem(last=False) - self.total_bytes -= removed_size + oldest = next(iter(self.objects)) + self._delete_object(oldest) + + def _delete_path(self, key: str) -> None: + entry = self.paths.pop(key, None) + if entry is None: + return + item = self.objects.get(entry[1]) + if item is None: + return + item[2].discard(key) + if not item[2]: + self._delete_object(entry[1]) + + def _delete_object(self, object_key: str) -> None: + item = self.objects.pop(object_key, None) + if item is None: + return + self.total_bytes -= item[1] + for key in item[2]: + self.paths.pop(key, None) # --------------------------------------------------------------------------- diff --git a/packages/sdk/python/tests/test_client.py b/packages/sdk/python/tests/test_client.py index 134abd23..4e997057 100644 --- a/packages/sdk/python/tests/test_client.py +++ b/packages/sdk/python/tests/test_client.py @@ -103,6 +103,41 @@ def test_read_file_revalidates_with_content_hash_and_handles_304(self) -> None: assert second == first assert route.calls[1].request.headers["If-None-Match"] == f'"{digest}"' + @respx.mock + def test_read_file_cache_keeps_per_path_metadata_for_identical_bytes(self) -> None: + digest = "b" * 64 + a = {"path": "/a", "revision": "rev_a", "contentType": "text/plain", "content": "hello", "contentHash": digest} + b = {"path": "/b", "revision": "rev_b", "contentType": "text/markdown", "content": "hello", "contentHash": digest} + respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file", params={"path": "/a"}).mock( + side_effect=[httpx.Response(200, json=a), httpx.Response(304)] + ) + respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file", params={"path": "/b"}).mock( + side_effect=[httpx.Response(200, json=b), httpx.Response(304)] + ) + client = RelayFileClient(BASE, "tok") + + client.read_file("ws_acme", "/a") + client.read_file("ws_acme", "/b") + + assert client.read_file("ws_acme", "/a") == a + assert client.read_file("ws_acme", "/b") == b + + @respx.mock + def test_read_file_never_serves_hashless_responses_from_cache(self) -> None: + route = respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file").mock( + side_effect=[ + httpx.Response(200, json={"path": "/f", "revision": "rev_1", "contentType": "text/plain", "content": "old"}), + httpx.Response(200, json={"path": "/f", "revision": "rev_2", "contentType": "text/plain", "content": "new"}), + ] + ) + client = RelayFileClient(BASE, "tok") + + client.read_file("ws_acme", "/f") + second = client.read_file("ws_acme", "/f") + + assert second["content"] == "new" + assert "If-None-Match" not in route.calls[1].request.headers + @respx.mock def test_write_file(self) -> None: payload = {"opId": "op_1", "status": "queued", "targetRevision": "rev_4"} diff --git a/packages/sdk/typescript/CHANGELOG.md b/packages/sdk/typescript/CHANGELOG.md index f0e3af86..aacc8819 100644 --- a/packages/sdk/typescript/CHANGELOG.md +++ b/packages/sdk/typescript/CHANGELOG.md @@ -8,7 +8,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added -- File reads now use a configurable, byte-capped content-addressed cache (`readCache.maxBytes`, 32 MiB by default; `readCache: false` disables it), revalidate with `If-None-Match`, and serve `304` responses from cached content. +- File reads now use a configurable, byte-capped content-addressed cache (`readCache.maxBytes`, 32 MiB by default; `readCache: false` disables it), revalidate with `If-None-Match`, and serve `304` responses from cached content. Only responses carrying a `contentHash` are cached (hashless responses always go to the server), and per-path metadata such as `revision` is kept separate from the deduplicated content bytes. - New `@relayfile/sdk/relay-cli` subpath export: `createRelayCliSurface()` returns relayfile's mountable CLI surface (`id: 'relayfile'`, contract v1), which the `agent-relay` CLI mounts as `agent-relay file`. `commands` is a checked-in snapshot of the Go CLI's own command table (`relayfile __command-spec --json`, regenerated by `npm run gen:command-spec`), and `run(argv, io)` spawns the same Go binary `relayfile` does, returning its real exit code. No relayfile command is reimplemented. - Binary resolution (`resolveRelayfileBinary`, `platformBinaryName`, `genericBinaryName`, `findSourceCheckoutRoot`) and the Cloud sign-in preflight (`prepareCloudSession`, `announceSetupIntent`) now live here, moved out of the `relayfile` package's scripts. Both the `relayfile` bin shim and `agent-relay file` use these, so each exists once in the repo. diff --git a/packages/sdk/typescript/src/client.test.ts b/packages/sdk/typescript/src/client.test.ts index c0ce2787..61ad2722 100644 --- a/packages/sdk/typescript/src/client.test.ts +++ b/packages/sdk/typescript/src/client.test.ts @@ -1490,6 +1490,42 @@ describe("RelayFileClient — existing methods", () => { }); }); + it("keeps per-path metadata separate when two paths share identical bytes", async () => { + const contentHash = "b".repeat(64); + const a: FileReadResponse = { path: "/a.txt", revision: "rev_a", contentType: "text/plain", content: "hello", contentHash }; + const b: FileReadResponse = { path: "/b.txt", revision: "rev_b", contentType: "text/markdown", content: "hello", contentHash }; + const f = vi.fn() + .mockResolvedValueOnce(jsonResponse(a)) + .mockResolvedValueOnce(jsonResponse(b)) + .mockResolvedValueOnce(jsonResponse(undefined, 304)) + .mockResolvedValueOnce(jsonResponse(undefined, 304)); + const client = makeClient(f); + + await client.readFile("ws_acme", "/a.txt"); + await client.readFile("ws_acme", "/b.txt"); + const againA = await client.readFile("ws_acme", "/a.txt"); + const againB = await client.readFile("ws_acme", "/b.txt"); + + expect(againA).toEqual(a); + expect(againB).toEqual(b); + }); + + it("never serves hashless responses from the cache", async () => { + const first: FileReadResponse = { path: "/legacy.txt", revision: "rev_1", contentType: "text/plain", content: "old" }; + const second: FileReadResponse = { path: "/legacy.txt", revision: "rev_2", contentType: "text/plain", content: "new" }; + const f = vi.fn() + .mockResolvedValueOnce(jsonResponse(first)) + .mockResolvedValueOnce(jsonResponse(second)); + const client = makeClient(f); + + await client.readFile("ws_acme", "/legacy.txt"); + const again = await client.readFile("ws_acme", "/legacy.txt"); + + expect(f).toHaveBeenCalledTimes(2); + expect((f.mock.calls[1]![1] as RequestInit).headers).not.toHaveProperty("If-None-Match"); + expect(again).toEqual(second); + }); + it("evicts least-recently-used content when the decoded byte cap is exceeded", async () => { const files = ["a", "b", "a"].map((name) => ({ path: `/${name}.txt`, revision: "rev_1", contentType: "text/plain", @@ -1583,7 +1619,12 @@ describe("RelayFileClient — existing methods", () => { revision: "rev_loser", contentType: "application/json", content: '{"pulls":["loser"]}', + contentHash: "c".repeat(64), }; + // Cached reads are revalidated with If-None-Match; an unconditional GET + // after a failed write proves the cache entry was invalidated. + const isConditional = (init: RequestInit) => + Boolean((init.headers as Record | undefined)?.["If-None-Match"]); const conflict = { code: "revision_conflict", message: "Conflict", @@ -1592,11 +1633,11 @@ describe("RelayFileClient — existing methods", () => { }; it("makes a real read request after a conflicting write instead of serving the populated cache", async () => { - let readRequests = 0; + const unconditionalReads: boolean[] = []; const f = vi.fn().mockImplementation((_url: string, init: RequestInit) => { if (init.method === "GET") { - readRequests += 1; - return Promise.resolve(jsonResponse(staleFile)); + unconditionalReads.push(!isConditional(init)); + return Promise.resolve(isConditional(init) ? jsonResponse(undefined, 304) : jsonResponse(staleFile)); } return Promise.resolve(jsonResponse(conflict, 409)); }); @@ -1604,7 +1645,7 @@ describe("RelayFileClient — existing methods", () => { await client.readFile(workspaceId, path); await client.readFile(workspaceId, path); - expect(readRequests).toBe(1); + expect(unconditionalReads).toEqual([true, false]); await expect(client.writeFile({ workspaceId, @@ -1615,7 +1656,7 @@ describe("RelayFileClient — existing methods", () => { })).rejects.toThrow(RevisionConflictError); await client.readFile(workspaceId, path); - expect(readRequests).toBe(2); + expect(unconditionalReads).toEqual([true, false, true]); }); it("observes the winning revision on an immediate read after a conflicting write", async () => { @@ -1694,11 +1735,11 @@ describe("RelayFileClient — existing methods", () => { }), }, ])("$name invalidates a populated cache when the request fails", async ({ attempt }) => { - let readRequests = 0; + const unconditionalReads: boolean[] = []; const f = vi.fn().mockImplementation((_url: string, init: RequestInit) => { if (init.method === "GET") { - readRequests += 1; - return Promise.resolve(jsonResponse(staleFile)); + unconditionalReads.push(!isConditional(init)); + return Promise.resolve(isConditional(init) ? jsonResponse(undefined, 304) : jsonResponse(staleFile)); } return Promise.resolve(jsonResponse(conflict, 409)); }); @@ -1706,12 +1747,12 @@ describe("RelayFileClient — existing methods", () => { await client.readFile(workspaceId, path); await client.readFile(workspaceId, path); - expect(readRequests).toBe(1); + expect(unconditionalReads).toEqual([true, false]); await expect(attempt(client)).rejects.toThrow(); await client.readFile(workspaceId, path); - expect(readRequests).toBe(2); + expect(unconditionalReads).toEqual([true, false, true]); }); }); diff --git a/packages/sdk/typescript/src/client.ts b/packages/sdk/typescript/src/client.ts index d425b033..10ee54df 100644 --- a/packages/sdk/typescript/src/client.ts +++ b/packages/sdk/typescript/src/client.ts @@ -377,16 +377,25 @@ const fileReadCaches = new WeakMap(); const DEFAULT_READ_CACHE_MAX_BYTES = 32 * 1024 * 1024; -interface ReadCacheEntry { - value: FileReadResponse; +// Content bytes are deduplicated by (hash, encoding); everything else in a +// read response (revision, path, semantics, ...) is per-path metadata and is +// kept per cache key so two paths with identical bytes never swap metadata. +interface ReadCacheObject { + content: string; bytes: number; + refs: Set; +} + +interface ReadCachePathEntry { + meta: Omit; + objectKey: string; } class FileReadCache { private readonly maxBytes: number; private totalBytes = 0; - private readonly objects = new Map(); - private readonly paths = new Map(); + private readonly objects = new Map(); + private readonly paths = new Map(); private readonly inFlight = new Map>(); constructor(options?: RelayFileReadCacheOptions) { @@ -394,41 +403,68 @@ class FileReadCache { } get(key: string): FileReadResponse | undefined { - const hash = this.paths.get(key); - if (!hash) return undefined; - const entry = this.objects.get(hash); + const entry = this.paths.get(key); if (!entry) return undefined; - this.objects.delete(hash); - this.objects.set(hash, entry); - return entry.value; + const object = this.objects.get(entry.objectKey); + if (!object) { + this.paths.delete(key); + return undefined; + } + this.objects.delete(entry.objectKey); + this.objects.set(entry.objectKey, object); + return { ...entry.meta, content: object.content }; } set(key: string, value: FileReadResponse): void { - const hash = normalizeContentHash(value.contentHash) ?? `legacy:${key}`; - if (this.maxBytes === 0) return; - const bytes = value.encoding === "base64" - ? Math.floor(value.content.length * 3 / 4) - (value.content.endsWith("==") ? 2 : value.content.endsWith("=") ? 1 : 0) - : new TextEncoder().encode(value.content).byteLength; - const previous = this.objects.get(hash); - if (previous) { - this.totalBytes -= previous.bytes; - this.objects.delete(hash); - } - this.objects.set(hash, { value, bytes }); - this.totalBytes += bytes; - this.paths.set(key, hash); + this.deletePath(key); + // Only responses carrying a verifiable content hash are retained: they are + // revalidated with If-None-Match on every read. Hashless (legacy) responses + // cannot be revalidated, so they are never served from the cache. + const hash = normalizeContentHash(value.contentHash); + if (!hash || this.maxBytes === 0) return; + const objectKey = `${hash}:${value.encoding ?? "utf-8"}`; + let object = this.objects.get(objectKey); + if (object) { + this.objects.delete(objectKey); + } else { + const bytes = value.encoding === "base64" + ? Math.floor(value.content.length * 3 / 4) - (value.content.endsWith("==") ? 2 : value.content.endsWith("=") ? 1 : 0) + : new TextEncoder().encode(value.content).byteLength; + object = { content: value.content, bytes, refs: new Set() }; + this.totalBytes += bytes; + } + object.refs.add(key); + this.objects.set(objectKey, object); + const { content: _content, ...meta } = value; + this.paths.set(key, { meta, objectKey }); while (this.totalBytes > this.maxBytes && this.objects.size > 0) { const oldest = this.objects.keys().next().value; if (oldest === undefined) break; - const evicted = this.objects.get(oldest)!; - this.objects.delete(oldest); - this.totalBytes -= evicted.bytes; + this.deleteObject(oldest); } } + private deletePath(key: string): void { + const entry = this.paths.get(key); + if (!entry) return; + this.paths.delete(key); + const object = this.objects.get(entry.objectKey); + if (!object) return; + object.refs.delete(key); + if (object.refs.size === 0) this.deleteObject(entry.objectKey); + } + + private deleteObject(objectKey: string): void { + const object = this.objects.get(objectKey); + if (!object) return; + this.objects.delete(objectKey); + this.totalBytes -= object.bytes; + for (const key of object.refs) this.paths.delete(key); + } + evict(workspaceId: string, path: string): void { - for (const key of this.paths.keys()) { - if (key.startsWith(`${workspaceId}:`) && key.endsWith(`:${path}`)) this.paths.delete(key); + for (const key of [...this.paths.keys()]) { + if (key.startsWith(`${workspaceId}:`) && key.endsWith(`:${path}`)) this.deletePath(key); } for (const key of this.inFlight.keys()) { if (key.startsWith(`${workspaceId}:`) && key.endsWith(`:${path}`)) this.inFlight.delete(key); @@ -1788,11 +1824,10 @@ export class RelayFileClient { const query = buildQuery({ path: input.path, forkId: input.forkId }); const cached = cache && cacheKey ? cache.get(cacheKey) : undefined; - if (cached && !normalizeContentHash(cached.contentHash)) return cached; const fetch = this.performRequest({ method: "GET", path: `/v1/workspaces/${encodeURIComponent(input.workspaceId)}/fs/file${query}`, - headers: cached?.contentHash ? { "If-None-Match": `"${normalizeContentHash(cached.contentHash) ?? cached.contentHash}"` } : undefined, + headers: cached?.contentHash ? { "If-None-Match": `"${normalizeContentHash(cached.contentHash)}"` } : undefined, correlationId: input.correlationId, signal: input.signal, tokenOverride: (input as ReadFileInput & { token?: string }).token, From be39be3036b99991e0b38fc3f2aec64aba835e14 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 28 Sep 2026 07:58:32 +0200 Subject: [PATCH 5/6] fix(sdk): treat the file-read ETag as opaque The server ETag for /fs/file is opaque and may encode more than the content hash, so building If-None-Match from contentHash never revalidates. Both SDKs now store the exact ETag the server returned for each path (refreshed from 304 responses) and echo it verbatim in If-None-Match. The local byte store stays keyed by the response body's contentHash. Responses without an ETag are not cached, since they cannot be revalidated. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/sdk/python/README.md | 2 +- packages/sdk/python/src/relayfile/client.py | 38 +++++++++------ packages/sdk/python/tests/test_client.py | 40 ++++++++++++---- packages/sdk/typescript/CHANGELOG.md | 2 +- packages/sdk/typescript/src/client.test.ts | 51 +++++++++++++++------ packages/sdk/typescript/src/client.ts | 40 +++++++++++----- 6 files changed, 124 insertions(+), 49 deletions(-) diff --git a/packages/sdk/python/README.md b/packages/sdk/python/README.md index 3836496d..df191aa2 100644 --- a/packages/sdk/python/README.md +++ b/packages/sdk/python/README.md @@ -2,7 +2,7 @@ Python SDK for the RelayFile virtual filesystem API. -File reads use a 32 MiB content-addressed cache by default and revalidate cached hashes with `If-None-Match`. Configure it with `RelayFileReadCacheOptions(max_bytes=...)`, or pass `read_cache=False` to either client to disable it. Retry delays use full jitter and never retry before a valid server `Retry-After` value. +File reads use a 32 MiB content-addressed cache by default and revalidate cached entries by echoing the server's opaque `ETag` verbatim in `If-None-Match`. Configure it with `RelayFileReadCacheOptions(max_bytes=...)`, or pass `read_cache=False` to either client to disable it. Retry delays use full jitter and never retry before a valid server `Retry-After` value. ## Install diff --git a/packages/sdk/python/src/relayfile/client.py b/packages/sdk/python/src/relayfile/client.py index d79eafb8..0adde87b 100644 --- a/packages/sdk/python/src/relayfile/client.py +++ b/packages/sdk/python/src/relayfile/client.py @@ -58,26 +58,32 @@ def __init__(self, options: RelayFileReadCacheOptions | None) -> None: self.total_bytes = 0 # object key -> (content, decoded size, referencing path keys) self.objects: OrderedDict[str, tuple[str, int, set[str]]] = OrderedDict() - # path key -> (metadata without content, object key) - self.paths: dict[str, tuple[dict[str, Any], str]] = {} + # path key -> (metadata without content, object key, server ETag) + self.paths: dict[str, tuple[dict[str, Any], str, str]] = {} self.lock = threading.RLock() - def get(self, key: str) -> dict[str, Any] | None: + def get(self, key: str) -> tuple[dict[str, Any], str] | None: + """Return (cached response, exact server ETag) for a path key.""" with self.lock: entry = self.paths.get(key) if entry is None: return None - meta, object_key = entry + meta, object_key, etag = entry item = self.objects.get(object_key) if item is None: self.paths.pop(key, None) return None self.objects.move_to_end(object_key) - return {**meta, "content": item[0]} + return {**meta, "content": item[0]}, etag - def put(self, key: str, value: dict[str, Any]) -> None: + def put(self, key: str, value: dict[str, Any], etag: str | None) -> None: + # The server ETag is opaque: it is stored exactly as returned and echoed + # verbatim in If-None-Match, never derived from contentHash. Responses + # without an ETag cannot be revalidated and are not cached. with self.lock: self._delete_path(key) + if not etag: + return digest = str(value.get("contentHash", "")).removeprefix("sha256:").lower() if len(digest) != 64 or any(c not in "0123456789abcdef" for c in digest): return @@ -102,7 +108,7 @@ def put(self, key: str, value: dict[str, Any]) -> None: self.total_bytes += size item[2].add(key) self.objects[object_key] = item - self.paths[key] = (meta, object_key) + self.paths[key] = (meta, object_key, etag) while self.total_bytes > self.max_bytes and self.objects: oldest = next(iter(self.objects)) self._delete_object(oldest) @@ -506,8 +512,9 @@ def read_file( ) -> dict[str, Any]: query = _build_query({"path": path}) key = f"{workspace_id}:{path}" - cached = self._read_cache.get(key) if self._read_cache else None - headers = {"If-None-Match": f'"{str(cached["contentHash"]).removeprefix("sha256:")}"'} if cached and cached.get("contentHash") else None + hit = self._read_cache.get(key) if self._read_cache else None + cached, etag = hit if hit else (None, None) + headers = {"If-None-Match": etag} if etag else None response = self._request_response( "GET", f"/v1/workspaces/{_enc(workspace_id)}/fs/file{query}", @@ -515,10 +522,12 @@ def read_file( correlation_id=correlation_id, ) if response.status_code == 304 and cached is not None: + if self._read_cache: + self._read_cache.put(key, cached, response.headers.get("ETag") or etag) return cached value = _read_payload(response) if self._read_cache and isinstance(value, dict): - self._read_cache.put(key, value) + self._read_cache.put(key, value, response.headers.get("ETag")) return value def query_files( @@ -1218,8 +1227,9 @@ async def read_file( ) -> dict[str, Any]: query = _build_query({"path": path}) key = f"{workspace_id}:{path}" - cached = self._read_cache.get(key) if self._read_cache else None - headers = {"If-None-Match": f'"{str(cached["contentHash"]).removeprefix("sha256:")}"'} if cached and cached.get("contentHash") else None + hit = self._read_cache.get(key) if self._read_cache else None + cached, etag = hit if hit else (None, None) + headers = {"If-None-Match": etag} if etag else None response = await self._request_response( "GET", f"/v1/workspaces/{_enc(workspace_id)}/fs/file{query}", @@ -1227,10 +1237,12 @@ async def read_file( correlation_id=correlation_id, ) if response.status_code == 304 and cached is not None: + if self._read_cache: + self._read_cache.put(key, cached, response.headers.get("ETag") or etag) return cached value = _read_payload(response) if self._read_cache and isinstance(value, dict): - self._read_cache.put(key, value) + self._read_cache.put(key, value, response.headers.get("ETag")) return value async def query_files( diff --git a/packages/sdk/python/tests/test_client.py b/packages/sdk/python/tests/test_client.py index 4e997057..c1c9dac3 100644 --- a/packages/sdk/python/tests/test_client.py +++ b/packages/sdk/python/tests/test_client.py @@ -87,11 +87,14 @@ def test_read_file(self) -> None: assert res["content"] == '{"id":1}' @respx.mock - def test_read_file_revalidates_with_content_hash_and_handles_304(self) -> None: + def test_read_file_revalidates_with_opaque_etag_and_handles_304(self) -> None: digest = "a" * 64 + # The server ETag is opaque (not the bare content hash); echo it verbatim. + etag = f'"{digest}:rev_1"' route = respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file").mock( side_effect=[ - httpx.Response(200, json={"path": "/f", "revision": "rev_1", "contentType": "text/plain", "content": "hello", "contentHash": digest}), + httpx.Response(200, json={"path": "/f", "revision": "rev_1", "contentType": "text/plain", "content": "hello", "contentHash": digest}, headers={"ETag": etag}), + httpx.Response(304, headers={"ETag": etag}), httpx.Response(304), ] ) @@ -99,20 +102,35 @@ def test_read_file_revalidates_with_content_hash_and_handles_304(self) -> None: first = client.read_file("ws_acme", "/f") second = client.read_file("ws_acme", "/f") + third = client.read_file("ws_acme", "/f") - assert second == first - assert route.calls[1].request.headers["If-None-Match"] == f'"{digest}"' + assert second == first == third + assert route.calls[1].request.headers["If-None-Match"] == etag + assert route.calls[2].request.headers["If-None-Match"] == etag + + @respx.mock + def test_read_file_without_etag_is_not_revalidated_from_cache(self) -> None: + payload = {"path": "/f", "revision": "rev_1", "contentType": "text/plain", "content": "hi", "contentHash": "d" * 64} + route = respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file").mock( + side_effect=[httpx.Response(200, json=payload), httpx.Response(200, json=payload)] + ) + client = RelayFileClient(BASE, "tok") + + client.read_file("ws_acme", "/f") + client.read_file("ws_acme", "/f") + + assert "If-None-Match" not in route.calls[1].request.headers @respx.mock def test_read_file_cache_keeps_per_path_metadata_for_identical_bytes(self) -> None: digest = "b" * 64 a = {"path": "/a", "revision": "rev_a", "contentType": "text/plain", "content": "hello", "contentHash": digest} b = {"path": "/b", "revision": "rev_b", "contentType": "text/markdown", "content": "hello", "contentHash": digest} - respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file", params={"path": "/a"}).mock( - side_effect=[httpx.Response(200, json=a), httpx.Response(304)] + a_route = respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file", params={"path": "/a"}).mock( + side_effect=[httpx.Response(200, json=a, headers={"ETag": '"etag-a"'}), httpx.Response(304)] ) - respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file", params={"path": "/b"}).mock( - side_effect=[httpx.Response(200, json=b), httpx.Response(304)] + b_route = respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file", params={"path": "/b"}).mock( + side_effect=[httpx.Response(200, json=b, headers={"ETag": '"etag-b"'}), httpx.Response(304)] ) client = RelayFileClient(BASE, "tok") @@ -121,13 +139,15 @@ def test_read_file_cache_keeps_per_path_metadata_for_identical_bytes(self) -> No assert client.read_file("ws_acme", "/a") == a assert client.read_file("ws_acme", "/b") == b + assert a_route.calls[1].request.headers["If-None-Match"] == '"etag-a"' + assert b_route.calls[1].request.headers["If-None-Match"] == '"etag-b"' @respx.mock def test_read_file_never_serves_hashless_responses_from_cache(self) -> None: route = respx.get(f"{BASE}/v1/workspaces/ws_acme/fs/file").mock( side_effect=[ - httpx.Response(200, json={"path": "/f", "revision": "rev_1", "contentType": "text/plain", "content": "old"}), - httpx.Response(200, json={"path": "/f", "revision": "rev_2", "contentType": "text/plain", "content": "new"}), + httpx.Response(200, json={"path": "/f", "revision": "rev_1", "contentType": "text/plain", "content": "old"}, headers={"ETag": '"rev_1"'}), + httpx.Response(200, json={"path": "/f", "revision": "rev_2", "contentType": "text/plain", "content": "new"}, headers={"ETag": '"rev_2"'}), ] ) client = RelayFileClient(BASE, "tok") diff --git a/packages/sdk/typescript/CHANGELOG.md b/packages/sdk/typescript/CHANGELOG.md index aacc8819..6403b9b4 100644 --- a/packages/sdk/typescript/CHANGELOG.md +++ b/packages/sdk/typescript/CHANGELOG.md @@ -8,7 +8,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Added -- File reads now use a configurable, byte-capped content-addressed cache (`readCache.maxBytes`, 32 MiB by default; `readCache: false` disables it), revalidate with `If-None-Match`, and serve `304` responses from cached content. Only responses carrying a `contentHash` are cached (hashless responses always go to the server), and per-path metadata such as `revision` is kept separate from the deduplicated content bytes. +- File reads now use a configurable, byte-capped content-addressed cache (`readCache.maxBytes`, 32 MiB by default; `readCache: false` disables it), revalidate by echoing the server's opaque `ETag` verbatim in `If-None-Match`, and serve `304` responses from cached content. Only responses carrying both an `ETag` and a `contentHash` are cached (others always go to the server), and per-path metadata such as `revision` is kept separate from the deduplicated content bytes. - New `@relayfile/sdk/relay-cli` subpath export: `createRelayCliSurface()` returns relayfile's mountable CLI surface (`id: 'relayfile'`, contract v1), which the `agent-relay` CLI mounts as `agent-relay file`. `commands` is a checked-in snapshot of the Go CLI's own command table (`relayfile __command-spec --json`, regenerated by `npm run gen:command-spec`), and `run(argv, io)` spawns the same Go binary `relayfile` does, returning its real exit code. No relayfile command is reimplemented. - Binary resolution (`resolveRelayfileBinary`, `platformBinaryName`, `genericBinaryName`, `findSourceCheckoutRoot`) and the Cloud sign-in preflight (`prepareCloudSession`, `announceSetupIntent`) now live here, moved out of the `relayfile` package's scripts. Both the `relayfile` bin shim and `agent-relay file` use these, so each exists once in the repo. diff --git a/packages/sdk/typescript/src/client.test.ts b/packages/sdk/typescript/src/client.test.ts index 61ad2722..9dea1230 100644 --- a/packages/sdk/typescript/src/client.test.ts +++ b/packages/sdk/typescript/src/client.test.ts @@ -1467,7 +1467,7 @@ describe("RelayFileClient — existing methods", () => { expect(res.revision).toBe("rev_3"); }); - it("revalidates cached content by hash and serves a 304 from the byte cache", async () => { + it("revalidates with the server's opaque ETag verbatim and serves a 304 from the byte cache", async () => { const contentHash = "a".repeat(64); const payload: FileReadResponse = { path: "/cached.txt", @@ -1476,18 +1476,38 @@ describe("RelayFileClient — existing methods", () => { content: "cached", contentHash, }; + // The ETag is opaque: it is not the bare content hash, and the client + // must echo it exactly rather than derive one from `contentHash`. + const etag = `"${contentHash}:rev_1"`; const f = vi.fn() - .mockResolvedValueOnce(jsonResponse(payload)) - .mockResolvedValueOnce(jsonResponse(undefined, 304)); + .mockResolvedValueOnce(jsonResponse(payload, 200, { ETag: etag })) + .mockResolvedValueOnce(jsonResponse(undefined, 304, { ETag: etag })) + .mockResolvedValueOnce(jsonResponse(undefined, 304, { ETag: etag })); const client = makeClient(f); await client.readFile("ws_acme", payload.path); const cached = await client.readFile("ws_acme", payload.path); + await client.readFile("ws_acme", payload.path); expect(cached).toEqual(payload); - expect((f.mock.calls[1]![1] as RequestInit).headers).toMatchObject({ - "If-None-Match": `"${contentHash}"`, - }); + for (const call of [1, 2]) { + expect((f.mock.calls[call]![1] as RequestInit).headers).toMatchObject({ "If-None-Match": etag }); + } + }); + + it("sends no If-None-Match when the server returned no ETag", async () => { + const payload: FileReadResponse = { + path: "/no-etag.txt", revision: "rev_1", contentType: "text/plain", content: "x", contentHash: "d".repeat(64), + }; + const f = vi.fn() + .mockResolvedValueOnce(jsonResponse(payload)) + .mockResolvedValueOnce(jsonResponse(payload)); + const client = makeClient(f); + + await client.readFile("ws_acme", payload.path); + await client.readFile("ws_acme", payload.path); + + expect((f.mock.calls[1]![1] as RequestInit).headers).not.toHaveProperty("If-None-Match"); }); it("keeps per-path metadata separate when two paths share identical bytes", async () => { @@ -1495,8 +1515,8 @@ describe("RelayFileClient — existing methods", () => { const a: FileReadResponse = { path: "/a.txt", revision: "rev_a", contentType: "text/plain", content: "hello", contentHash }; const b: FileReadResponse = { path: "/b.txt", revision: "rev_b", contentType: "text/markdown", content: "hello", contentHash }; const f = vi.fn() - .mockResolvedValueOnce(jsonResponse(a)) - .mockResolvedValueOnce(jsonResponse(b)) + .mockResolvedValueOnce(jsonResponse(a, 200, { ETag: '"etag-a"' })) + .mockResolvedValueOnce(jsonResponse(b, 200, { ETag: '"etag-b"' })) .mockResolvedValueOnce(jsonResponse(undefined, 304)) .mockResolvedValueOnce(jsonResponse(undefined, 304)); const client = makeClient(f); @@ -1508,14 +1528,16 @@ describe("RelayFileClient — existing methods", () => { expect(againA).toEqual(a); expect(againB).toEqual(b); + expect((f.mock.calls[2]![1] as RequestInit).headers).toMatchObject({ "If-None-Match": '"etag-a"' }); + expect((f.mock.calls[3]![1] as RequestInit).headers).toMatchObject({ "If-None-Match": '"etag-b"' }); }); it("never serves hashless responses from the cache", async () => { const first: FileReadResponse = { path: "/legacy.txt", revision: "rev_1", contentType: "text/plain", content: "old" }; const second: FileReadResponse = { path: "/legacy.txt", revision: "rev_2", contentType: "text/plain", content: "new" }; const f = vi.fn() - .mockResolvedValueOnce(jsonResponse(first)) - .mockResolvedValueOnce(jsonResponse(second)); + .mockResolvedValueOnce(jsonResponse(first, 200, { ETag: '"rev_1"' })) + .mockResolvedValueOnce(jsonResponse(second, 200, { ETag: '"rev_2"' })); const client = makeClient(f); await client.readFile("ws_acme", "/legacy.txt"); @@ -1531,7 +1553,10 @@ describe("RelayFileClient — existing methods", () => { path: `/${name}.txt`, revision: "rev_1", contentType: "text/plain", content: name.repeat(4), contentHash: name.repeat(64), } satisfies FileReadResponse)); - const f = vi.fn().mockImplementation(() => Promise.resolve(jsonResponse(files.shift()))); + const f = vi.fn().mockImplementation(() => { + const file = files.shift()!; + return Promise.resolve(jsonResponse(file, 200, { ETag: `"${file.contentHash}:${file.revision}"` })); + }); const client = new RelayFileClient({ baseUrl: "https://relay.test", token: "tok", fetchImpl: f, readCache: { maxBytes: 4 } }); await client.readFile("ws", "/a.txt"); @@ -1637,7 +1662,7 @@ describe("RelayFileClient — existing methods", () => { const f = vi.fn().mockImplementation((_url: string, init: RequestInit) => { if (init.method === "GET") { unconditionalReads.push(!isConditional(init)); - return Promise.resolve(isConditional(init) ? jsonResponse(undefined, 304) : jsonResponse(staleFile)); + return Promise.resolve(isConditional(init) ? jsonResponse(undefined, 304) : jsonResponse(staleFile, 200, { ETag: '"stale-etag"' })); } return Promise.resolve(jsonResponse(conflict, 409)); }); @@ -1739,7 +1764,7 @@ describe("RelayFileClient — existing methods", () => { const f = vi.fn().mockImplementation((_url: string, init: RequestInit) => { if (init.method === "GET") { unconditionalReads.push(!isConditional(init)); - return Promise.resolve(isConditional(init) ? jsonResponse(undefined, 304) : jsonResponse(staleFile)); + return Promise.resolve(isConditional(init) ? jsonResponse(undefined, 304) : jsonResponse(staleFile, 200, { ETag: '"stale-etag"' })); } return Promise.resolve(jsonResponse(conflict, 409)); }); diff --git a/packages/sdk/typescript/src/client.ts b/packages/sdk/typescript/src/client.ts index 10ee54df..08793e03 100644 --- a/packages/sdk/typescript/src/client.ts +++ b/packages/sdk/typescript/src/client.ts @@ -389,8 +389,16 @@ interface ReadCacheObject { interface ReadCachePathEntry { meta: Omit; objectKey: string; + /** Exact ETag the server returned for this path; opaque, echoed verbatim. */ + etag: string; } +// The ETag a read response arrived with. The server's ETag is opaque (it may +// encode more than the content hash, e.g. the revision), so the client never +// derives it from `contentHash`; it stores the header it was given and sends +// it back unchanged in If-None-Match. +const readResponseEtags = new WeakMap(); + class FileReadCache { private readonly maxBytes: number; private totalBytes = 0; @@ -402,7 +410,7 @@ class FileReadCache { this.maxBytes = Math.max(0, Math.floor(options?.maxBytes ?? DEFAULT_READ_CACHE_MAX_BYTES)); } - get(key: string): FileReadResponse | undefined { + get(key: string): { value: FileReadResponse; etag: string } | undefined { const entry = this.paths.get(key); if (!entry) return undefined; const object = this.objects.get(entry.objectKey); @@ -412,16 +420,17 @@ class FileReadCache { } this.objects.delete(entry.objectKey); this.objects.set(entry.objectKey, object); - return { ...entry.meta, content: object.content }; + return { value: { ...entry.meta, content: object.content }, etag: entry.etag }; } - set(key: string, value: FileReadResponse): void { + set(key: string, value: FileReadResponse, etag: string | undefined): void { this.deletePath(key); - // Only responses carrying a verifiable content hash are retained: they are - // revalidated with If-None-Match on every read. Hashless (legacy) responses - // cannot be revalidated, so they are never served from the cache. + // Only responses carrying both a server ETag (to revalidate with + // If-None-Match on every read) and a content hash (to key the byte store) + // are retained. Anything else cannot be revalidated, so it is never served + // from the cache. const hash = normalizeContentHash(value.contentHash); - if (!hash || this.maxBytes === 0) return; + if (!hash || !etag || this.maxBytes === 0) return; const objectKey = `${hash}:${value.encoding ?? "utf-8"}`; let object = this.objects.get(objectKey); if (object) { @@ -436,7 +445,7 @@ class FileReadCache { object.refs.add(key); this.objects.set(objectKey, object); const { content: _content, ...meta } = value; - this.paths.set(key, { meta, objectKey }); + this.paths.set(key, { meta, objectKey, etag }); while (this.totalBytes > this.maxBytes && this.objects.size > 0) { const oldest = this.objects.keys().next().value; if (oldest === undefined) break; @@ -481,7 +490,7 @@ class FileReadCache { (result) => { if (this.inFlight.get(key) === promise) { this.inFlight.delete(key); - this.set(key, result); + this.set(key, result, readResponseEtags.get(result)); } }, () => { @@ -1827,12 +1836,21 @@ export class RelayFileClient { const fetch = this.performRequest({ method: "GET", path: `/v1/workspaces/${encodeURIComponent(input.workspaceId)}/fs/file${query}`, - headers: cached?.contentHash ? { "If-None-Match": `"${normalizeContentHash(cached.contentHash)}"` } : undefined, + headers: cached ? { "If-None-Match": cached.etag } : undefined, correlationId: input.correlationId, signal: input.signal, tokenOverride: (input as ReadFileInput & { token?: string }).token, allowNotModified: Boolean(cached) - }).then(async (response) => response.status === 304 && cached ? cached : await this.readPayload(response) as FileReadResponse); + }).then(async (response) => { + const etag = response.headers.get("ETag") ?? undefined; + if (response.status === 304 && cached) { + readResponseEtags.set(cached.value, etag ?? cached.etag); + return cached.value; + } + const value = await this.readPayload(response) as FileReadResponse; + if (etag && value && typeof value === "object") readResponseEtags.set(value, etag); + return value; + }); if (cache && cacheKey) { cache.setInFlight(cacheKey, fetch); From 6364f3f01b18318a975316467643c7fae33f9e77 Mon Sep 17 00:00:00 2001 From: Khaliq Date: Mon, 28 Sep 2026 08:23:18 +0200 Subject: [PATCH 6/6] fix(sdk-python): make read-cache put atomic put() invalidated the path under the lock, released it to validate the body, then re-inserted. A concurrent put for the same path could leave a stale object reference whose later LRU eviction dropped the live path. Validate outside the lock, then invalidate + insert under one lock hold; eviction now only drops paths that still point at the evicted object. Co-Authored-By: Claude Opus 5.5 (1M context) --- packages/sdk/python/src/relayfile/client.py | 50 ++++++++++------- packages/sdk/python/tests/test_client.py | 62 +++++++++++++++++++++ 2 files changed, 93 insertions(+), 19 deletions(-) diff --git a/packages/sdk/python/src/relayfile/client.py b/packages/sdk/python/src/relayfile/client.py index 0adde87b..fcf4d1c5 100644 --- a/packages/sdk/python/src/relayfile/client.py +++ b/packages/sdk/python/src/relayfile/client.py @@ -80,15 +80,35 @@ def put(self, key: str, value: dict[str, Any], etag: str | None) -> None: # The server ETag is opaque: it is stored exactly as returned and echoed # verbatim in If-None-Match, never derived from contentHash. Responses # without an ETag cannot be revalidated and are not cached. + entry = self._prepare(value, etag) + # Invalidate and insert under one lock hold: releasing it in between + # lets a concurrent put for the same path leave a stale object + # reference whose later eviction would drop the live path entry. with self.lock: self._delete_path(key) - if not etag: - return + if entry is None: + return + object_key, content, size, meta = entry + item = self.objects.pop(object_key, None) + if item is None: + item = (content, size, set()) + self.total_bytes += size + item[2].add(key) + self.objects[object_key] = item + self.paths[key] = (meta, object_key, etag or "") + while self.total_bytes > self.max_bytes and self.objects: + oldest = next(iter(self.objects)) + self._delete_object(oldest) + + def _prepare( + self, value: dict[str, Any], etag: str | None + ) -> tuple[str, str, int, dict[str, Any]] | None: + """Validate and size a response outside the lock; None = don't cache.""" + if not etag or self.max_bytes == 0: + return None digest = str(value.get("contentHash", "")).removeprefix("sha256:").lower() if len(digest) != 64 or any(c not in "0123456789abcdef" for c in digest): - return - if self.max_bytes == 0: - return + return None content = str(value.get("content", "")) encoding = str(value.get("encoding") or "utf-8") if encoding == "base64": @@ -96,22 +116,11 @@ def put(self, key: str, value: dict[str, Any], etag: str | None) -> None: try: size = len(base64.b64decode(content, validate=True)) except ValueError: - return + return None else: size = len(content.encode()) - object_key = f"{digest}:{encoding}" meta = {k: v for k, v in value.items() if k != "content"} - with self.lock: - item = self.objects.pop(object_key, None) - if item is None: - item = (content, size, set()) - self.total_bytes += size - item[2].add(key) - self.objects[object_key] = item - self.paths[key] = (meta, object_key, etag) - while self.total_bytes > self.max_bytes and self.objects: - oldest = next(iter(self.objects)) - self._delete_object(oldest) + return f"{digest}:{encoding}", content, size, meta def _delete_path(self, key: str) -> None: entry = self.paths.pop(key, None) @@ -130,7 +139,10 @@ def _delete_object(self, object_key: str) -> None: return self.total_bytes -= item[1] for key in item[2]: - self.paths.pop(key, None) + entry = self.paths.get(key) + # Only drop paths that still point at this object. + if entry is not None and entry[1] == object_key: + del self.paths[key] # --------------------------------------------------------------------------- diff --git a/packages/sdk/python/tests/test_client.py b/packages/sdk/python/tests/test_client.py index c1c9dac3..f8ea0b3f 100644 --- a/packages/sdk/python/tests/test_client.py +++ b/packages/sdk/python/tests/test_client.py @@ -1124,3 +1124,65 @@ async def test_set_integration_metadata_rejects_malformed_cloud_response( await client.set_integration_metadata( "ws_acme", "jira", {"cloudId": "cloud-1"} ) + + +def _cache_consistent(cache) -> None: + for key, (_meta, object_key, _etag) in cache.paths.items(): + assert object_key in cache.objects + assert key in cache.objects[object_key][2] + for object_key, (_content, _size, refs) in cache.objects.items(): + for key in refs: + assert cache.paths[key][1] == object_key + assert cache.total_bytes == sum(item[1] for item in cache.objects.values()) + + +def test_read_cache_put_is_atomic_against_interleaved_put_for_same_path() -> None: + from relayfile.client import _FileReadCache + + cache = _FileReadCache(RelayFileReadCacheOptions(max_bytes=1024)) + hash_a, hash_b = "a" * 64, "b" * 64 + value_a = {"path": "/f", "revision": "r1", "contentHash": hash_a, "content": "AAAA"} + value_b = {"path": "/f", "revision": "r2", "contentHash": hash_b, "content": "BBBB"} + + prepare = cache._prepare + interleaved = False + + def racing_prepare(value, etag): + nonlocal interleaved + if not interleaved: + interleaved = True + # A concurrent reader stores a different body for the same path + # while this put is validating outside the lock. + cache.put("ws:/f", value_b, '"b"') + return prepare(value, etag) + + cache._prepare = racing_prepare + cache.put("ws:/f", value_a, '"a"') + cache._prepare = prepare + + _cache_consistent(cache) + hit = cache.get("ws:/f") + assert hit is not None and hit[0]["revision"] == "r1" + assert list(cache.objects) == [f"{hash_a}:utf-8"] + + # Evicting any other object must never drop the live path entry. + cache.put("ws:/other", {**value_b, "path": "/other"}, '"o"') + cache._delete_object(f"{hash_b}:utf-8") + _cache_consistent(cache) + assert cache.get("ws:/f") is not None + + +def test_read_cache_eviction_ignores_stale_path_references() -> None: + from relayfile.client import _FileReadCache + + cache = _FileReadCache(RelayFileReadCacheOptions(max_bytes=1024)) + live = {"path": "/f", "revision": "r2", "contentHash": "b" * 64, "content": "BBBB"} + cache.put("ws:/f", live, '"b"') + # Plant a stale object that still (wrongly) references the live path. + cache.objects[f"{'a' * 64}:utf-8"] = ("AAAA", 4, {"ws:/f"}) + cache.total_bytes += 4 + + cache._delete_object(f"{'a' * 64}:utf-8") + + hit = cache.get("ws:/f") + assert hit is not None and hit[0]["revision"] == "r2"