fix(worker): resolve temporary paths in tests (#11944)

Capacity guards reject symlink components. On macOS, temporary paths
start with /var, which links to /private/var, so the new staging tests
fail before exercising cleanup or capacity accounting.

Resolve the fixture directories before building guarded paths. Keep
explicit symlinks within the fixtures for containment tests.

Assisted-by: Codex:gpt-6

Co-authored-by: localai-org-maint-bot <306269227+localai-org-maint-bot@users.noreply.github.com>
This commit is contained in:
localai-org-maint-botandlocalai-org-maint-bot authored and GitHub committed 2026-09-09 08:51:09 +02:00
1 parent bf93008ef3
commit be5342ef05
4 files changed
+52 -43

No files matched your search

+23 -23
View File
@@ -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())
@@ -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)
@@ -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{}
@@ -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
}