mirror of
https://github.com/mudler/LocalAI.git
synced 2026-07-30 09:57:57 -04:00
materializeLocked built a download task for every file in the resolved snapshot unconditionally. A completed file is promoted from .downloads/<hash> into snapshot/<path> and its blob deleted, so on any re-entry (a controller pod roll, a resubmit, a crash) the new pass built a task whose .downloads blob no longer existed, re-downloaded the whole file from Hugging Face, and its AfterDownload even removed the already-complete snapshot copy first. The only resume that worked was the downloader's per-file .partial resume for a file caught mid-transfer; completed files were never skipped. Production consequence: installing longcat-video-avatar-1.5 (~35 GB after allow_patterns) on a cluster whose controller rolls hourly (Flux image automation) never converged across ~14 hours. Each roll restarted from the first shard; the completed bytes on disk were repeatedly deleted and re-fetched, and the artifact never promoted. curl of the same files from inside the pod ran fine, proving the loss was the materializer re-fetching, not the network. Before building a task, check whether the file is already materialized and verified in this staging tree's snapshot/ and, if so, keep it and count it complete instead of downloading. "Materialized" means a regular file of the expected size that passes the same verifyDownloadedFile check the download path uses, so the kept manifest entry is byte-for-byte identical to a fresh one and integrity is re-checked. The manifest requires a SHA-256 for every file and non-LFS files carry none to borrow, so a hash is unavoidable for the manifest anyway; a full re-hash of local disk is still orders of magnitude cheaper than re-downloading, and the downloader re-verifies any file it does fetch. Manifest entries are now written at their snapshot index rather than appended in completion order, so a mix of skipped and downloaded files keeps the resolved order that committedResult and staging read. The unconditional root.Remove(destination) now runs only on the fresh-download path; a kept file survives. Skips are logged at INFO with count and bytes so an operator can see resume working. This is the resume-side counterpart to the sibling defects on this path: read/write error conflation and transient retry (#10985), hash-verify progress accounting and silent success on an expired deadline (#11026), and the response-header hang (#11053). The download machinery resumed a single in-flight file; the materializer above it still threw away every completed file on restart. It also makes orphan-partial adoption worth its cost: an adopted tree's completed files were re-downloaded anyway until now. Assisted-by: Claude Code:claude-opus-4-8[1m] [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
272 lines
11 KiB
Go
272 lines
11 KiB
Go
package modelartifacts_test
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"os"
|
|
"path/filepath"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"sync/atomic"
|
|
|
|
. "github.com/onsi/ginkgo/v2"
|
|
. "github.com/onsi/gomega"
|
|
|
|
hfapi "github.com/mudler/LocalAI/pkg/huggingface-api"
|
|
"github.com/mudler/LocalAI/pkg/modelartifacts"
|
|
)
|
|
|
|
type fakeSnapshotResolver struct {
|
|
mu sync.Mutex
|
|
snapshot hfapi.Snapshot
|
|
err error
|
|
calls int
|
|
}
|
|
|
|
func (f *fakeSnapshotResolver) ResolveSnapshot(context.Context, hfapi.SnapshotRequest) (hfapi.Snapshot, error) {
|
|
f.mu.Lock()
|
|
defer f.mu.Unlock()
|
|
f.calls++
|
|
return f.snapshot, f.err
|
|
}
|
|
|
|
func (f *fakeSnapshotResolver) callCount() int {
|
|
f.mu.Lock()
|
|
defer f.mu.Unlock()
|
|
return f.calls
|
|
}
|
|
|
|
var _ = Describe("controller artifact materializer", func() {
|
|
It("requires an injected token before resolving a private artifact", func() {
|
|
resolver := &fakeSnapshotResolver{}
|
|
_, err := modelartifacts.NewManager(resolver).Ensure(context.Background(), GinkgoT().TempDir(),
|
|
modelartifacts.Spec{Source: modelartifacts.Source{
|
|
Type: modelartifacts.SourceTypeHuggingFace, Repo: "owner/private", TokenEnv: modelartifacts.HuggingFaceTokenEnv,
|
|
}})
|
|
Expect(err).To(MatchError(ContainSubstring("non-empty HF_TOKEN")))
|
|
Expect(resolver.callCount()).To(BeZero())
|
|
})
|
|
|
|
It("downloads, commits, and reuses a pinned snapshot", func() {
|
|
weight := []byte("weight-bytes")
|
|
sum := sha256.Sum256(weight)
|
|
var requests atomic.Int32
|
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
requests.Add(1)
|
|
Expect(r.Header.Get("Authorization")).To(Equal("Bearer hf-secret"))
|
|
w.Header().Set("Content-Length", "12")
|
|
_, _ = w.Write(weight)
|
|
}))
|
|
DeferCleanup(server.Close)
|
|
|
|
resolver := &fakeSnapshotResolver{snapshot: hfapi.Snapshot{
|
|
Endpoint: "https://huggingface.co", Repo: "owner/repo",
|
|
RequestedRevision: "main", ResolvedRevision: "0123456789abcdef0123456789abcdef01234567",
|
|
Files: []hfapi.SnapshotFile{{Path: "nested/model.safetensors", Size: 12, LFSOID: hex.EncodeToString(sum[:]), URL: server.URL + "/model"}},
|
|
}}
|
|
manager := modelartifacts.NewManager(resolver, modelartifacts.WithHuggingFaceToken("hf-secret"))
|
|
modelsPath := GinkgoT().TempDir()
|
|
spec := modelartifacts.Spec{Source: modelartifacts.Source{Type: "huggingface", Repo: "owner/repo", TokenEnv: "HF_TOKEN"}}
|
|
|
|
result, err := manager.Ensure(context.Background(), modelsPath, spec)
|
|
Expect(err).NotTo(HaveOccurred())
|
|
Expect(result.CacheHit).To(BeFalse())
|
|
Expect(result.Spec.Resolved.Revision).To(Equal("0123456789abcdef0123456789abcdef01234567"))
|
|
// A snapshot with a single file records it as the primary file so the
|
|
// load target is the file, not the snapshot directory.
|
|
Expect(result.Spec.Resolved.PrimaryFile).To(Equal("nested/model.safetensors"))
|
|
Expect(result.RelativePath).To(HavePrefix(".artifacts/huggingface/"))
|
|
Expect(os.ReadFile(filepath.Join(modelsPath, filepath.FromSlash(result.RelativePath), "nested", "model.safetensors"))).To(Equal(weight))
|
|
|
|
manifestBytes, err := os.ReadFile(filepath.Join(modelsPath, filepath.Dir(filepath.FromSlash(result.RelativePath)), "manifest.json"))
|
|
Expect(err).NotTo(HaveOccurred())
|
|
Expect(string(manifestBytes)).NotTo(ContainSubstring("hf-secret"))
|
|
|
|
cached, err := manager.Ensure(context.Background(), modelsPath, result.Spec)
|
|
Expect(err).NotTo(HaveOccurred())
|
|
Expect(cached.CacheHit).To(BeTrue())
|
|
Expect(requests.Load()).To(Equal(int32(1)))
|
|
Expect(resolver.callCount()).To(Equal(1))
|
|
})
|
|
|
|
It("leaves the primary file unset for a multi-file snapshot", func() {
|
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { _, _ = w.Write([]byte("x")) }))
|
|
DeferCleanup(server.Close)
|
|
resolver := &fakeSnapshotResolver{snapshot: hfapi.Snapshot{
|
|
Endpoint: "https://huggingface.co", Repo: "owner/repo",
|
|
ResolvedRevision: "0123456789abcdef0123456789abcdef01234567",
|
|
Files: []hfapi.SnapshotFile{
|
|
{Path: "config.json", Size: 1, URL: server.URL + "/config"},
|
|
{Path: "model.safetensors", Size: 1, URL: server.URL + "/model"},
|
|
},
|
|
}}
|
|
result, err := modelartifacts.NewManager(resolver).Ensure(context.Background(), GinkgoT().TempDir(),
|
|
modelartifacts.Spec{Source: modelartifacts.Source{Type: "huggingface", Repo: "owner/repo"}})
|
|
Expect(err).NotTo(HaveOccurred())
|
|
// Multi-file snapshots are consumed as a directory (e.g. transformers), so
|
|
// no single file is promoted.
|
|
Expect(result.Spec.Resolved.PrimaryFile).To(BeEmpty())
|
|
})
|
|
|
|
It("rejects a path escape before opening a destination", func() {
|
|
resolver := &fakeSnapshotResolver{snapshot: hfapi.Snapshot{
|
|
Endpoint: "https://huggingface.co", Repo: "owner/repo",
|
|
ResolvedRevision: "0123456789abcdef0123456789abcdef01234567",
|
|
Files: []hfapi.SnapshotFile{{Path: "../escape", Size: 1, URL: "https://example.invalid/file"}},
|
|
}}
|
|
modelsPath := GinkgoT().TempDir()
|
|
_, err := modelartifacts.NewManager(resolver).Ensure(context.Background(), modelsPath,
|
|
modelartifacts.Spec{Source: modelartifacts.Source{Type: "huggingface", Repo: "owner/repo"}})
|
|
Expect(err).To(MatchError(ContainSubstring("unsafe Hub path")))
|
|
Expect(filepath.Join(modelsPath, "escape")).NotTo(BeAnExistingFile())
|
|
})
|
|
|
|
It("does not commit after cancellation", func() {
|
|
started := make(chan struct{})
|
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
close(started)
|
|
<-r.Context().Done()
|
|
}))
|
|
DeferCleanup(server.Close)
|
|
resolver := &fakeSnapshotResolver{snapshot: hfapi.Snapshot{
|
|
Endpoint: "https://huggingface.co", Repo: "owner/repo",
|
|
ResolvedRevision: "0123456789abcdef0123456789abcdef01234567",
|
|
Files: []hfapi.SnapshotFile{{Path: "model.bin", Size: 10, URL: server.URL}},
|
|
}}
|
|
ctx, cancel := context.WithCancel(context.Background())
|
|
done := make(chan error, 1)
|
|
modelsPath := GinkgoT().TempDir()
|
|
go func() {
|
|
_, err := modelartifacts.NewManager(resolver).Ensure(ctx, modelsPath,
|
|
modelartifacts.Spec{Source: modelartifacts.Source{Type: "huggingface", Repo: "owner/repo"}})
|
|
done <- err
|
|
}()
|
|
<-started
|
|
cancel()
|
|
Expect(<-done).To(MatchError(context.Canceled))
|
|
entries, err := filepath.Glob(filepath.Join(modelsPath, ".artifacts", "huggingface", "*"))
|
|
Expect(err).NotTo(HaveOccurred())
|
|
Expect(entries).To(BeEmpty())
|
|
})
|
|
|
|
It("resumes an interrupted materialization without re-downloading completed files", func() {
|
|
contents := map[string][]byte{
|
|
"a/first.bin": []byte("first-file-bytes"),
|
|
"b/second.bin": []byte("second-file-bytes-longer"),
|
|
"third.bin": []byte("third"),
|
|
}
|
|
// order fixes both the download sequence and the expected manifest order.
|
|
order := []string{"a/first.bin", "b/second.bin", "third.bin"}
|
|
oid := func(b []byte) string { sum := sha256.Sum256(b); return hex.EncodeToString(sum[:]) }
|
|
|
|
var mu sync.Mutex
|
|
fetched := map[string]int{}
|
|
// During run 1 this file 404s, interrupting the pass after the first file
|
|
// has already completed and been promoted into the staging snapshot.
|
|
failing := "b/second.bin"
|
|
armed := true
|
|
|
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
|
name := strings.TrimPrefix(r.URL.Path, "/file/")
|
|
body, ok := contents[name]
|
|
if !ok {
|
|
w.WriteHeader(http.StatusNotFound)
|
|
return
|
|
}
|
|
mu.Lock()
|
|
arm := armed
|
|
mu.Unlock()
|
|
if arm && name == failing {
|
|
w.WriteHeader(http.StatusNotFound)
|
|
return
|
|
}
|
|
mu.Lock()
|
|
fetched[name]++
|
|
mu.Unlock()
|
|
w.Header().Set("Content-Length", strconv.Itoa(len(body)))
|
|
_, _ = w.Write(body)
|
|
}))
|
|
DeferCleanup(server.Close)
|
|
|
|
files := make([]hfapi.SnapshotFile, 0, len(order))
|
|
for _, p := range order {
|
|
files = append(files, hfapi.SnapshotFile{
|
|
Path: p, Size: int64(len(contents[p])), LFSOID: oid(contents[p]),
|
|
URL: server.URL + "/file/" + p,
|
|
})
|
|
}
|
|
resolver := &fakeSnapshotResolver{snapshot: hfapi.Snapshot{
|
|
Endpoint: "https://huggingface.co", Repo: "owner/repo",
|
|
ResolvedRevision: "0123456789abcdef0123456789abcdef01234567",
|
|
Files: files,
|
|
}}
|
|
// A single manager (one writer identity) re-enters the same staging tree on
|
|
// the second call, exactly as a resubmit against a still-idle partial does.
|
|
manager := modelartifacts.NewManager(resolver)
|
|
modelsPath := GinkgoT().TempDir()
|
|
spec := modelartifacts.Spec{Source: modelartifacts.Source{Type: "huggingface", Repo: "owner/repo"}}
|
|
|
|
_, err := manager.Ensure(context.Background(), modelsPath, spec)
|
|
Expect(err).To(HaveOccurred())
|
|
mu.Lock()
|
|
Expect(fetched["a/first.bin"]).To(Equal(1), "the first file should have completed in run 1")
|
|
armed = false
|
|
fetched = map[string]int{}
|
|
mu.Unlock()
|
|
|
|
result, err := manager.Ensure(context.Background(), modelsPath, spec)
|
|
Expect(err).NotTo(HaveOccurred())
|
|
|
|
mu.Lock()
|
|
defer mu.Unlock()
|
|
// The decisive assertion: the already-completed file is NOT re-downloaded.
|
|
Expect(fetched).NotTo(HaveKey("a/first.bin"))
|
|
// Only the files missing after the interruption are fetched on resume.
|
|
Expect(fetched).To(HaveKey("b/second.bin"))
|
|
Expect(fetched).To(HaveKey("third.bin"))
|
|
|
|
// The manifest lists every file, in order, with the correct size and hash.
|
|
paths := make([]string, 0, len(result.Manifest.Files))
|
|
byPath := map[string]modelartifacts.ManifestFile{}
|
|
for _, f := range result.Manifest.Files {
|
|
paths = append(paths, f.Path)
|
|
byPath[f.Path] = f
|
|
}
|
|
Expect(paths).To(Equal(order))
|
|
for _, p := range order {
|
|
Expect(byPath[p].Size).To(Equal(int64(len(contents[p]))))
|
|
Expect(byPath[p].SHA256).To(Equal(oid(contents[p])))
|
|
}
|
|
// Every file, skipped or freshly downloaded, is present on disk after commit.
|
|
for _, p := range order {
|
|
Expect(os.ReadFile(filepath.Join(modelsPath, filepath.FromSlash(result.RelativePath), filepath.FromSlash(p)))).To(Equal(contents[p]))
|
|
}
|
|
})
|
|
|
|
It("emits resolving, downloading, verifying, and committing phases", func() {
|
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { _, _ = w.Write([]byte("x")) }))
|
|
DeferCleanup(server.Close)
|
|
resolver := &fakeSnapshotResolver{snapshot: hfapi.Snapshot{
|
|
Endpoint: "https://huggingface.co", Repo: "owner/repo",
|
|
ResolvedRevision: "0123456789abcdef0123456789abcdef01234567",
|
|
Files: []hfapi.SnapshotFile{{Path: "config.json", Size: 1, URL: server.URL}},
|
|
}}
|
|
var phases []modelartifacts.Phase
|
|
ctx := modelartifacts.WithProgressSink(context.Background(), func(event modelartifacts.ProgressEvent) {
|
|
if len(phases) == 0 || phases[len(phases)-1] != event.Phase {
|
|
phases = append(phases, event.Phase)
|
|
}
|
|
})
|
|
_, err := modelartifacts.NewManager(resolver).Ensure(ctx, GinkgoT().TempDir(),
|
|
modelartifacts.Spec{Source: modelartifacts.Source{Type: "huggingface", Repo: "owner/repo"}})
|
|
Expect(err).NotTo(HaveOccurred())
|
|
Expect(phases).To(Equal([]modelartifacts.Phase{
|
|
modelartifacts.PhaseResolving, modelartifacts.PhaseDownloading, modelartifacts.PhaseVerifying, modelartifacts.PhaseCommitting,
|
|
}))
|
|
})
|
|
})
|