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/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..c104f3cb 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:11:31.329Z", "trajectories": { "traj_4pvrlmqfnzng": { "title": "Review PR #278 in AgentWorkforce/relayfile", @@ -407,6 +407,20 @@ "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" + }, + "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/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..2b96a0c7 --- /dev/null +++ b/internal/mountsync/object_cache.go @@ -0,0 +1,204 @@ +package mountsync + +import ( + "crypto/sha256" + "encoding/base64" + "encoding/hex" + "os" + "path/filepath" + "regexp" + "sort" + "strings" + "sync" + "time" + "unicode/utf8" +) + +var objectHashPattern = regexp.MustCompile(`^[a-f0-9]{64}$`) + +// 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() + 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) 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 RemoteFile{}, false + } + 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 == "" { + return + } + data, ok := remoteFileBytes(file) + if !ok || int64(len(data)) > c.limit() { + return + } + 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-*") + 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, 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 remoteFileBytes(file RemoteFile) ([]byte, bool) { + if strings.EqualFold(file.Encoding, "base64") { + decoded, err := base64.StdEncoding.DecodeString(file.Content) + if err != nil { + return nil, false + } + return decoded, true + } + return []byte(file.Content), true +} diff --git a/internal/mountsync/object_cache_test.go b/internal/mountsync/object_cache_test.go new file mode 100644 index 00000000..da80505a --- /dev/null +++ b/internal/mountsync/object_cache_test.go @@ -0,0 +1,171 @@ +package mountsync + +import ( + "context" + "encoding/base64" + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +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("tampered"), 0o600); err != nil { + t.Fatal(err) + } + 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 d474d61d..ebce9742 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 { @@ -8016,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 @@ -8023,6 +8053,22 @@ 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.cachedBootstrapFile(job) + if !ok { + remaining = append(remaining, job) + continue + } + 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 +8148,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 +8318,14 @@ func (s *Syncer) readBootstrapFilesIndividuallyBatchEach(ctx context.Context, jo go func() { defer wg.Done() for job := range jobCh { + if cached, ok := s.cachedBootstrapFile(job); ok { + 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 +8363,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 +9373,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 +12543,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 +12551,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..61faac61 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/`. 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 - 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..df191aa2 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 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 ```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..fcf4d1c5 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,110 @@ 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: + """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 + # 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, server ETag) + self.paths: dict[str, tuple[dict[str, Any], str, str]] = {} + self.lock = threading.RLock() + + 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, 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]}, etag + + 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 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 None + content = str(value.get("content", "")) + encoding = str(value.get("encoding") or "utf-8") + if encoding == "base64": + import base64 + try: + size = len(base64.b64decode(content, validate=True)) + except ValueError: + return None + else: + size = len(content.encode()) + meta = {k: v for k, v in value.items() if k != "content"} + return f"{digest}:{encoding}", content, size, meta + + 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]: + 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] + + # --------------------------------------------------------------------------- # Helpers # --------------------------------------------------------------------------- @@ -137,12 +243,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 +386,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 +402,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 +478,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 +523,24 @@ def read_file( correlation_id: str | None = None, ) -> dict[str, Any]: query = _build_query({"path": path}) - return self._request( + key = f"{workspace_id}:{path}" + 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}", + headers=headers, 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, response.headers.get("ETag")) + return value def query_files( self, @@ -982,6 +1100,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 +1113,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 +1193,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 +1238,24 @@ 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}" + 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}", + headers=headers, 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, response.headers.get("ETag")) + 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..f8ea0b3f 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,78 @@ 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_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}, headers={"ETag": etag}), + httpx.Response(304, headers={"ETag": etag}), + 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") + third = client.read_file("ws_acme", "/f") + + 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} + 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)] + ) + 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") + + 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 + 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"}, 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") + + 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"} @@ -1051,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" diff --git a/packages/sdk/typescript/CHANGELOG.md b/packages/sdk/typescript/CHANGELOG.md index 27912b33..6403b9b4 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 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. - 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..9dea1230 100644 --- a/packages/sdk/typescript/src/client.test.ts +++ b/packages/sdk/typescript/src/client.test.ts @@ -1466,6 +1466,105 @@ describe("RelayFileClient — existing methods", () => { expect(res.content).toBe('{"id":48291}'); expect(res.revision).toBe("rev_3"); }); + + 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", + revision: "rev_1", + contentType: "text/plain", + 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, 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); + 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 () => { + 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, 200, { ETag: '"etag-a"' })) + .mockResolvedValueOnce(jsonResponse(b, 200, { ETag: '"etag-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); + 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, 200, { ETag: '"rev_1"' })) + .mockResolvedValueOnce(jsonResponse(second, 200, { ETag: '"rev_2"' })); + 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", + content: name.repeat(4), contentHash: name.repeat(64), + } satisfies FileReadResponse)); + 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"); + 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 ---- @@ -1545,7 +1644,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", @@ -1554,11 +1658,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, 200, { ETag: '"stale-etag"' })); } return Promise.resolve(jsonResponse(conflict, 409)); }); @@ -1566,7 +1670,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, @@ -1577,7 +1681,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 () => { @@ -1656,11 +1760,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, 200, { ETag: '"stale-etag"' })); } return Promise.resolve(jsonResponse(conflict, 409)); }); @@ -1668,12 +1772,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 151e5aed..08793e03 100644 --- a/packages/sdk/typescript/src/client.ts +++ b/packages/sdk/typescript/src/client.ts @@ -375,48 +375,109 @@ 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; +// 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; + /** 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 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); + get(key: string): { value: FileReadResponse; etag: string } | undefined { + const entry = this.paths.get(key); if (!entry) return undefined; - if (Date.now() > entry.expiresAt) { - this.entries.delete(key); + const object = this.objects.get(entry.objectKey); + if (!object) { + this.paths.delete(key); return undefined; } - 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); - } - } - this.entries.set(key, { value, expiresAt: Date.now() + this.ttlMs }); + this.objects.delete(entry.objectKey); + this.objects.set(entry.objectKey, object); + return { value: { ...entry.meta, content: object.content }, etag: entry.etag }; + } + + set(key: string, value: FileReadResponse, etag: string | undefined): void { + this.deletePath(key); + // 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 || !etag || 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, etag }); + while (this.totalBytes > this.maxBytes && this.objects.size > 0) { + const oldest = this.objects.keys().next().value; + if (oldest === undefined) break; + 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 { - 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.deletePath(key); + } + for (const key of this.inFlight.keys()) { + if (key.startsWith(`${workspaceId}:`) && key.endsWith(`:${path}`)) this.inFlight.delete(key); + } } getInFlight(key: string): Promise | undefined { @@ -429,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)); } }, () => { @@ -441,6 +502,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,22 +1824,32 @@ 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; + const fetch = this.performRequest({ method: "GET", path: `/v1/workspaces/${encodeURIComponent(input.workspaceId)}/fs/file${query}`, + headers: cached ? { "If-None-Match": cached.etag } : 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) => { + 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) { @@ -2845,6 +2921,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 +2970,7 @@ export class RelayFileClient { throw error; } - if (response.ok) { + if (response.ok || (params.allowNotModified && response.status === 304)) { return response; } @@ -2948,10 +3025,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 +3034,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/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() } 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;