diff --git a/core/services/worker/ephemeral_capacity_test.go b/core/services/worker/ephemeral_capacity_test.go index a5b4fd1f4..640d17799 100644 --- a/core/services/worker/ephemeral_capacity_test.go +++ b/core/services/worker/ephemeral_capacity_test.go @@ -43,7 +43,7 @@ func (capacityShortWriter) Write(p []byte) (int, error) { var _ = Describe("EphemeralCapacityGuard", func() { It("derives bounded defaults and preserves positive overrides", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() limit, headroom, err := effectiveEphemeralCapacity([]string{root}, 0, -1) Expect(err).NotTo(HaveOccurred()) Expect(limit).To(BeNumerically(">", 0)) @@ -57,8 +57,8 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("accounts existing regular files without following symlinks", func() { - root := GinkgoT().TempDir() - outside := filepath.Join(GinkgoT().TempDir(), "outside.bin") + root := canonicalWorkerTempDir() + outside := filepath.Join(canonicalWorkerTempDir(), "outside.bin") Expect(os.WriteFile(filepath.Join(root, "existing.bin"), make([]byte, 6), 0o600)).To(Succeed()) Expect(os.WriteFile(outside, make([]byte, 100), 0o600)).To(Succeed()) Expect(os.Symlink(outside, filepath.Join(root, "outside-link"))).To(Succeed()) @@ -77,7 +77,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("serializes competing reservations", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 1, 0) Expect(err).NotTo(HaveOccurred()) @@ -109,7 +109,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("makes only an equal active reservation idempotent", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "nested", "payload.bin") @@ -128,7 +128,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("retains committed bytes when the same path starts another reservation", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "payload.bin") @@ -147,7 +147,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("retains startup-accounted bytes when the path is reserved", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() path := filepath.Join(root, "payload.bin") Expect(os.WriteFile(path, make([]byte, 4), 0o600)).To(Succeed()) guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) @@ -163,7 +163,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("commits the regular file's actual size", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "payload.bin") @@ -176,7 +176,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("preserves configured filesystem headroom", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 1<<30, 1<<62) Expect(err).NotTo(HaveOccurred()) @@ -189,7 +189,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("reserves bounded chunks before forwarding unknown-length input", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, ephemeralCapacityWriteChunk+1, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "payload.bin") @@ -210,7 +210,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("waits for an open bounded writer before committing", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "payload.bin") @@ -254,7 +254,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("does not share pending capacity between concurrent writers", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "payload.bin") @@ -290,7 +290,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("rolls back bytes the destination writer does not accept", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 5, 0) Expect(err).NotTo(HaveOccurred()) writer, err := guard.NewWriter(filepath.Join(root, "payload.bin"), capacityShortWriter{}) @@ -304,8 +304,8 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("rejects paths outside roots and through symlinks", func() { - root := GinkgoT().TempDir() - outside := GinkgoT().TempDir() + root := canonicalWorkerTempDir() + outside := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 100, 0) Expect(err).NotTo(HaveOccurred()) @@ -319,7 +319,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("supports recovery tree accounting without dropping active reservations", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) active := filepath.Join(root, "active", "payload.bin") @@ -335,7 +335,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("waits for pre-release reservations before request cleanup scans", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "audio", "request-1", "input.wav") @@ -357,7 +357,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("rejects staging after request cleanup begins", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) Expect(guard.BeginRequestRelease(context.Background(), "request-1")).To(Succeed()) @@ -370,7 +370,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("leaves a late commit recoverable when release times out", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "audio", "request-1", "late.wav") @@ -386,7 +386,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("bounds release markers without reopening registered work", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) Expect(guard.BeginRequestOperation("request-pinned")).To(Succeed()) @@ -411,7 +411,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("applies backpressure at the release-pin cap and clears ownership", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "audio", "request-target", "input.wav") @@ -438,7 +438,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("makes committed files recoverable when pin backpressure expires", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() guard, err := NewEphemeralCapacityGuard([]string{root}, 10, 0) Expect(err).NotTo(HaveOccurred()) path := filepath.Join(root, "audio", "request-target", "input.wav") @@ -459,7 +459,7 @@ var _ = Describe("EphemeralCapacityGuard", func() { }) It("rejects a registered cache-hit claim after pin backpressure expires", func() { - root := GinkgoT().TempDir() + root := canonicalWorkerTempDir() path := filepath.Join(root, "audio", "request-target", "input.wav") Expect(os.MkdirAll(filepath.Dir(path), 0o750)).To(Succeed()) Expect(os.WriteFile(path, []byte("data"), 0o600)).To(Succeed()) diff --git a/core/services/worker/ephemeral_cleanup_test.go b/core/services/worker/ephemeral_cleanup_test.go index c8cb9e150..a8bec3616 100644 --- a/core/services/worker/ephemeral_cleanup_test.go +++ b/core/services/worker/ephemeral_cleanup_test.go @@ -24,7 +24,7 @@ var _ = Describe("Worker ephemeral staging cleanup", func() { return dir } - BeforeEach(func() { stagingDir = GinkgoT().TempDir() }) + BeforeEach(func() { stagingDir = canonicalWorkerTempDir() }) It("removes staged request directories older than the TTL", func() { old := mkEphemeral("aaaa1111", 48*time.Hour) @@ -57,7 +57,7 @@ var _ = Describe("Worker ephemeral staging cleanup", func() { }) It("sweeps both transport roots by newest descendant and skips active requests", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() httpRoot := filepath.Join(stagingDir, "ephemeral") s3Root := filepath.Join(cacheDir, "ephemeral") guard, err := NewEphemeralCapacityGuard([]string{httpRoot, s3Root}, 8, 0) diff --git a/core/services/worker/file_staging_release_test.go b/core/services/worker/file_staging_release_test.go index 62dfcd3cb..1019d475d 100644 --- a/core/services/worker/file_staging_release_test.go +++ b/core/services/worker/file_staging_release_test.go @@ -87,7 +87,7 @@ func (m *releaseMessagingClient) Close() {} var _ = Describe("Worker exact-key staging release", func() { It("protects a startup-accounted HTTP cache hit through authenticated repeated probes", func() { - stagingDir := GinkgoT().TempDir() + stagingDir := canonicalWorkerTempDir() root := filepath.Join(stagingDir, "ephemeral") key := "ephemeral/audio/request-id/input.wav" remotePath := filepath.Join(stagingDir, filepath.FromSlash(key)) @@ -107,11 +107,11 @@ var _ = Describe("Worker exact-key staging release", func() { Expect(err).NotTo(HaveOccurred()) addr := listener.Addr().String() Expect(listener.Close()).To(Succeed()) - server, err := nodes.StartFileTransferServerWithCapacity(addr, stagingDir, GinkgoT().TempDir(), GinkgoT().TempDir(), "secret", 0, nil, guard) + server, err := nodes.StartFileTransferServerWithCapacity(addr, stagingDir, canonicalWorkerTempDir(), canonicalWorkerTempDir(), "secret", 0, nil, guard) Expect(err).NotTo(HaveOccurred()) DeferCleanup(nodes.ShutdownFileTransferServer, server) - localPath := filepath.Join(GinkgoT().TempDir(), "input.wav") + localPath := filepath.Join(canonicalWorkerTempDir(), "input.wav") Expect(os.WriteFile(localPath, content, 0o600)).To(Succeed()) stager := nodes.NewHTTPFileStager(func(string) (string, error) { return addr, nil }, "secret") for range 2 { @@ -129,7 +129,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("claims a startup-scanned cache hit against stale recovery until release", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() root := filepath.Join(cacheDir, "ephemeral") key := "ephemeral/audio/request-id/input.wav" cachePath := filepath.Join(cacheDir, filepath.FromSlash(key)) @@ -158,7 +158,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("downloads again when a cache file disappears while being claimed", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() key := "ephemeral/audio/request-id/input.wav" cachePath := filepath.Join(cacheDir, filepath.FromSlash(key)) Expect(os.MkdirAll(filepath.Dir(cachePath), 0o750)).To(Succeed()) @@ -176,7 +176,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("makes repeated cache-hit claims idempotent", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() root := filepath.Join(cacheDir, "ephemeral") key := "ephemeral/audio/request-id/input.wav" cachePath := filepath.Join(cacheDir, filepath.FromSlash(key)) @@ -199,7 +199,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("capacity-checks growth of a startup-scanned cache file", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() root := filepath.Join(cacheDir, "ephemeral") key := "ephemeral/audio/request-id/input.wav" cachePath := filepath.Join(cacheDir, filepath.FromSlash(key)) @@ -220,7 +220,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("reserves S3 object size before download and releases it with the exact key", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() root := filepath.Join(cacheDir, "ephemeral") store := &stagingObjectStore{payload: []byte("data")} fm, err := storage.NewFileManager(store, cacheDir) @@ -240,7 +240,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("rejects an oversized S3 object before starting its download", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() store := &stagingObjectStore{payload: []byte("oversized")} fm, err := storage.NewFileManager(store, cacheDir) Expect(err).NotTo(HaveOccurred()) @@ -253,7 +253,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("rolls back an S3 reservation when the download fails", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() root := filepath.Join(cacheDir, "ephemeral") store := &stagingObjectStore{payload: []byte("data"), getErr: errors.New("download failed")} fm, err := storage.NewFileManager(store, cacheDir) @@ -267,7 +267,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("removes only the exact cache file and upload sidecars", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() categoryDir := filepath.Join(cacheDir, "ephemeral", "request-id", "audio") Expect(os.MkdirAll(categoryDir, 0750)).To(Succeed()) target := filepath.Join(categoryDir, "input.wav") @@ -285,7 +285,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("succeeds for a missing file and prunes empty category and request directories", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() categoryDir := filepath.Join(cacheDir, "ephemeral", "request-id", "audio") Expect(os.MkdirAll(categoryDir, 0750)).To(Succeed()) @@ -298,8 +298,8 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("rejects traversal and symlink escapes", func() { - cacheDir := GinkgoT().TempDir() - outsideDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() + outsideDir := canonicalWorkerTempDir() outsidePath := filepath.Join(outsideDir, "input.wav") Expect(os.WriteFile(outsidePath, []byte("keep"), 0640)).To(Succeed()) requestDir := filepath.Join(cacheDir, "ephemeral", "request-id") @@ -319,7 +319,7 @@ var _ = Describe("Worker exact-key staging release", func() { It("rejects symlinked files and sidecars without deleting their targets", func() { for _, linkedName := range []string{"input.wav", "input.wav.sha256", "input.wav.sha256.target"} { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() categoryDir := filepath.Join(cacheDir, "ephemeral", "request-id", "audio") Expect(os.MkdirAll(categoryDir, 0750)).To(Succeed()) target := filepath.Join(categoryDir, "input.wav") @@ -336,7 +336,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("registers an exact release handler", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() path := filepath.Join(cacheDir, "ephemeral", "request-id", "audio", "input.wav") Expect(os.MkdirAll(filepath.Dir(path), 0750)).To(Succeed()) Expect(os.WriteFile(path, []byte("data"), 0640)).To(Succeed()) @@ -358,7 +358,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("releases a request batch through one worker message", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() keys := []string{ "ephemeral/audio/request-id/input.wav", "ephemeral/images/request-id/frame.jpg", @@ -387,7 +387,7 @@ var _ = Describe("Worker exact-key staging release", func() { }) It("returns validation errors through the release handler", func() { - cacheDir := GinkgoT().TempDir() + cacheDir := canonicalWorkerTempDir() fm, err := storage.NewFileManager(nil, cacheDir) Expect(err).NotTo(HaveOccurred()) client := &releaseMessagingClient{} diff --git a/core/services/worker/worker_suite_test.go b/core/services/worker/worker_suite_test.go index 64186d88f..e55a32f5c 100644 --- a/core/services/worker/worker_suite_test.go +++ b/core/services/worker/worker_suite_test.go @@ -1,6 +1,7 @@ package worker import ( + "path/filepath" "testing" . "github.com/onsi/ginkgo/v2" @@ -11,3 +12,11 @@ func TestWorker(t *testing.T) { RegisterFailHandler(Fail) RunSpecs(t, "Worker Suite") } + +// Capacity guards reject symlink components, including macOS /var -> /private/var. +func canonicalWorkerTempDir() string { + GinkgoHelper() + dir, err := filepath.EvalSymlinks(GinkgoT().TempDir()) + Expect(err).NotTo(HaveOccurred()) + return dir +}