mirror of
https://github.com/mudler/LocalAI.git
synced 2026-07-30 18:09:05 -04:00
flock(2) on CIFS/SMB returns EACCES when another client holds the lock: the kernel maps STATUS_LOCK_NOT_GRANTED and STATUS_FILE_LOCK_CONFLICT to -EACCES and never produces EWOULDBLOCK on that path. gofrs/flock only recognises EWOULDBLOCK as contention, so TryLockContext returned a bare "permission denied" and Ensure aborted. Both replicas then fell back to legacy loading, which makes the worker download the whole repo in-band inside LoadModel and blow the remote-load deadline. Replace TryLockContext with an explicit wait loop over a new Locker interface, classifying EWOULDBLOCK/EAGAIN/EACCES/EBUSY as contention. EACCES is ambiguous at the syscall boundary but not here: the lock file is already open O_CREATE|O_RDWR, so a real permission problem would have failed the open with an *fs.PathError, and flock(2) documents no EACCES on Linux at all. The wait is bounded (DefaultLockWait, overridable via WithLockWait), so even a misclassification degrades to a delay. On timeout the committed result is re-checked before reporting the new ErrLockContended, so a peer that finished the work still wins. Locker also exists so the contention path is testable without a network filesystem: nothing in CI can make flock(2) return EACCES on demand. Raise the fallback to error for a managedArtifactBackends backend, via a shared config.LogArtifactFallback used by both call sites. For those backends the legacy path is not graceful degradation, and the operator otherwise sees only a timeout with no causal link. The fallback stays non-fatal. Drop the os.Chmod(layout.Lock, 0o600) after acquisition: flock.New already creates the file 0600, and the chmod was gratuitous risk on a nounix mount that ignores modes. Fixes #10981 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>
170 lines
5.8 KiB
Go
170 lines
5.8 KiB
Go
package modelartifacts_test
|
|
|
|
import (
|
|
"context"
|
|
"crypto/sha256"
|
|
"encoding/hex"
|
|
"net/http"
|
|
"net/http/httptest"
|
|
"sync"
|
|
"syscall"
|
|
"time"
|
|
|
|
. "github.com/onsi/ginkgo/v2"
|
|
. "github.com/onsi/gomega"
|
|
|
|
hfapi "github.com/mudler/LocalAI/pkg/huggingface-api"
|
|
"github.com/mudler/LocalAI/pkg/modelartifacts"
|
|
)
|
|
|
|
// lockAttempt is one scripted outcome of Locker.TryLock. Attempts beyond the
|
|
// script acquire the lock.
|
|
type lockAttempt struct {
|
|
locked bool
|
|
err error
|
|
}
|
|
|
|
type scriptedLocker struct {
|
|
mu sync.Mutex
|
|
script []lockAttempt
|
|
attempts int
|
|
unlocks int
|
|
onAttempt func(attempt int)
|
|
}
|
|
|
|
func (s *scriptedLocker) TryLock() (bool, error) {
|
|
s.mu.Lock()
|
|
s.attempts++
|
|
attempt := s.attempts
|
|
var outcome lockAttempt
|
|
if attempt <= len(s.script) {
|
|
outcome = s.script[attempt-1]
|
|
} else {
|
|
outcome = lockAttempt{locked: true}
|
|
}
|
|
callback := s.onAttempt
|
|
s.mu.Unlock()
|
|
if callback != nil {
|
|
callback(attempt)
|
|
}
|
|
return outcome.locked, outcome.err
|
|
}
|
|
|
|
func (s *scriptedLocker) Unlock() error {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
s.unlocks++
|
|
return nil
|
|
}
|
|
|
|
func (s *scriptedLocker) attemptCount() int {
|
|
s.mu.Lock()
|
|
defer s.mu.Unlock()
|
|
return s.attempts
|
|
}
|
|
|
|
// singleFileResolver serves one small file so a full Ensure can run end to end
|
|
// without the test caring about download mechanics.
|
|
func singleFileResolver() *fakeSnapshotResolver {
|
|
weight := []byte("weight-bytes")
|
|
sum := sha256.Sum256(weight)
|
|
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) {
|
|
_, _ = w.Write(weight)
|
|
}))
|
|
DeferCleanup(server.Close)
|
|
return &fakeSnapshotResolver{snapshot: hfapi.Snapshot{
|
|
Endpoint: "https://huggingface.co", Repo: "owner/repo",
|
|
RequestedRevision: "main", ResolvedRevision: "0123456789abcdef0123456789abcdef01234567",
|
|
Files: []hfapi.SnapshotFile{{
|
|
Path: "model.safetensors", Size: int64(len(weight)),
|
|
LFSOID: hex.EncodeToString(sum[:]), URL: server.URL + "/model",
|
|
}},
|
|
}}
|
|
}
|
|
|
|
var specUnderTest = modelartifacts.Spec{Source: modelartifacts.Source{Type: "huggingface", Repo: "owner/repo"}}
|
|
|
|
// CIFS/SMB maps STATUS_LOCK_NOT_GRANTED and STATUS_FILE_LOCK_CONFLICT to
|
|
// EACCES, so a contended flock(2) on a shared /models never produces the
|
|
// EWOULDBLOCK that gofrs/flock recognises. Treating that as a hard failure made
|
|
// both replicas abandon materialization and degrade to an in-band download of
|
|
// the whole repo inside LoadModel (#10981).
|
|
var _ = Describe("artifact lock contention on network filesystems", func() {
|
|
It("waits out CIFS-style EACCES contention and then materializes", func() {
|
|
locker := &scriptedLocker{script: []lockAttempt{
|
|
{err: syscall.EACCES},
|
|
{err: syscall.EACCES},
|
|
}}
|
|
manager := modelartifacts.NewManager(singleFileResolver(),
|
|
modelartifacts.WithLocker(func(string) modelartifacts.Locker { return locker }),
|
|
modelartifacts.WithLockWait(10*time.Second))
|
|
|
|
result, err := manager.Ensure(context.Background(), GinkgoT().TempDir(), specUnderTest)
|
|
Expect(err).NotTo(HaveOccurred())
|
|
Expect(result.RelativePath).To(HavePrefix(".artifacts/huggingface/"))
|
|
Expect(locker.attemptCount()).To(Equal(3))
|
|
})
|
|
|
|
It("waits out an EWOULDBLOCK contention reported without an error", func() {
|
|
locker := &scriptedLocker{script: []lockAttempt{{}, {}}}
|
|
manager := modelartifacts.NewManager(singleFileResolver(),
|
|
modelartifacts.WithLocker(func(string) modelartifacts.Locker { return locker }),
|
|
modelartifacts.WithLockWait(10*time.Second))
|
|
|
|
_, err := manager.Ensure(context.Background(), GinkgoT().TempDir(), specUnderTest)
|
|
Expect(err).NotTo(HaveOccurred())
|
|
Expect(locker.attemptCount()).To(Equal(3))
|
|
})
|
|
|
|
It("does not retry a lock error that is not contention", func() {
|
|
locker := &scriptedLocker{script: []lockAttempt{{err: syscall.ENOLCK}}}
|
|
manager := modelartifacts.NewManager(singleFileResolver(),
|
|
modelartifacts.WithLocker(func(string) modelartifacts.Locker { return locker }),
|
|
modelartifacts.WithLockWait(10*time.Second))
|
|
|
|
_, err := manager.Ensure(context.Background(), GinkgoT().TempDir(), specUnderTest)
|
|
Expect(err).To(MatchError(syscall.ENOLCK))
|
|
Expect(locker.attemptCount()).To(Equal(1))
|
|
})
|
|
|
|
It("adopts the peer's snapshot when contention outlasts the wait window", func() {
|
|
modelsPath := GinkgoT().TempDir()
|
|
resolver := singleFileResolver()
|
|
// The "peer" replica commits the snapshot while we are still waiting,
|
|
// which is the whole point of waiting: the result exists even though we
|
|
// never held the lock ourselves.
|
|
peer := modelartifacts.NewManager(resolver,
|
|
modelartifacts.WithLocker(func(string) modelartifacts.Locker { return &scriptedLocker{} }))
|
|
locker := &scriptedLocker{
|
|
script: []lockAttempt{{err: syscall.EACCES}, {err: syscall.EACCES}, {err: syscall.EACCES}, {err: syscall.EACCES}},
|
|
onAttempt: func(attempt int) {
|
|
if attempt != 2 {
|
|
return
|
|
}
|
|
_, err := peer.Ensure(context.Background(), modelsPath, specUnderTest)
|
|
Expect(err).NotTo(HaveOccurred())
|
|
},
|
|
}
|
|
manager := modelartifacts.NewManager(resolver,
|
|
modelartifacts.WithLocker(func(string) modelartifacts.Locker { return locker }),
|
|
modelartifacts.WithLockWait(300*time.Millisecond))
|
|
|
|
result, err := manager.Ensure(context.Background(), modelsPath, specUnderTest)
|
|
Expect(err).NotTo(HaveOccurred())
|
|
Expect(result.CacheHit).To(BeTrue())
|
|
})
|
|
|
|
It("reports contention distinctly when the peer never commits", func() {
|
|
locker := &scriptedLocker{script: []lockAttempt{
|
|
{err: syscall.EACCES}, {err: syscall.EACCES}, {err: syscall.EACCES},
|
|
{err: syscall.EACCES}, {err: syscall.EACCES}, {err: syscall.EACCES},
|
|
}}
|
|
manager := modelartifacts.NewManager(singleFileResolver(),
|
|
modelartifacts.WithLocker(func(string) modelartifacts.Locker { return locker }),
|
|
modelartifacts.WithLockWait(200*time.Millisecond))
|
|
|
|
_, err := manager.Ensure(context.Background(), GinkgoT().TempDir(), specUnderTest)
|
|
Expect(err).To(MatchError(modelartifacts.ErrLockContended))
|
|
})
|
|
})
|