Files
mudler's LocalAI [bot]andEttore Di Giacinto 82c191afad fix(distributed): keep model replicas config-consistent (#11664)
* docs: design configurable copy buffering

Document the context-aware copy buffer option and its validation plan.

Assisted-by: Codex:gpt-5

* docs: design durable distributed staging operations

Assisted-by: Codex:gpt-5

* docs: design distributed model config revisions

Assisted-by: Codex:GPT-5 [apply_patch] [exec_command]

* feat(config): add stable model revisions

Hash typed model configuration and effective protobuf options deterministically for distributed revision comparisons.

Assisted-by: Codex:GPT-5 [apply_patch] [exec_command]

* feat(worker): acknowledge exact model stops

Assisted-by: Codex:GPT-5 [apply_patch] [exec_command]

* feat(nodes): track model config revisions

Assisted-by: Codex:GPT-5 [apply_patch]

* fix(distributed): retry quarantined model cleanup

Stop quarantined replicas by exact process identity, retain failed cleanup as durable capped retries, and compare-and-delete only the claimed registry row. Process one sufficiently leased row at a time so multiple frontends cannot duplicate slow cleanup work.

Assisted-by: Codex:gpt-5

* fix(distributed): bind loads to config revisions

Assisted-by: Codex: GPT-5 [OpenAI Codex]

* fix(modeladmin): apply config revisions consistently

Route model edits, patches, state changes, deletion, and peer refreshes through the same revision lifecycle. Quarantine stale replicas before exact cleanup and report durable pending cleanup without failing successful config writes.

Assisted-by: Codex: GPT-5 [OpenAI Codex]

* feat(distributed): expose model config revision state

Document replica revision observability and durable cleanup behavior. Keep pending cleanup explicit in model mutation responses and verify endpoint contracts expose revision state without serialized load options.

Assisted-by: Codex:GPT-5 [OpenAI Codex]

* test(distributed): cover model revision convergence

Exercise cross-frontend quarantine, stale replay rejection, exact cleanup retry, worker re-registration, and current-generation replica convergence against the distributed PostgreSQL harness.

Assisted-by: Codex:gpt-5

* fix(distributed): pass config revision CI checks

Keep configured gallery sources out of authoritative runtime snapshots only after validating their real schema, and harden rollback snapshots against symlink races and non-regular files.

Assisted-by: Codex: GPT-5 [OpenAI Codex]

---------

Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
2026-08-22 22:44:03 +02:00

90 lines
3.4 KiB
Go

package modeladmin
import (
"context"
"errors"
"github.com/mudler/LocalAI/core/services/nodes"
. "github.com/onsi/ginkgo/v2"
. "github.com/onsi/gomega"
)
type lifecycleRegistry struct {
models map[string][]nodes.NodeModel
calls []string
err error
}
func (r *lifecycleRegistry) AdvanceModelConfigRevisions(_ context.Context, transitions []nodes.ModelConfigRevisionTransition) ([]nodes.NodeModel, error) {
for _, transition := range transitions {
r.calls = append(r.calls, transition.ModelName+":"+transition.ConfigRevision)
}
if r.err != nil {
return nil, r.err
}
var quarantined []nodes.NodeModel
for _, transition := range transitions {
quarantined = append(quarantined, r.models[transition.ModelName]...)
}
return quarantined, nil
}
type lifecycleCleanup struct {
registry *lifecycleRegistry
seen []nodes.NodeModel
pending int
}
func (c *lifecycleCleanup) Cleanup(_ context.Context, replicas []nodes.NodeModel, _ bool) int {
Expect(c.registry.calls).ToNot(BeEmpty(), "quarantine must precede worker cleanup")
c.seen = append(c.seen, replicas...)
return c.pending
}
var _ = Describe("DistributedModelRevisionLifecycle", func() {
It("advances the registry before cleanup and reports incomplete exact stops", func() {
registry := &lifecycleRegistry{models: map[string][]nodes.NodeModel{
"model": {{ID: "stale", ModelName: "model", State: "unloading"}},
}}
cleanup := &lifecycleCleanup{registry: registry, pending: 1}
lifecycle := NewDistributedModelRevisionLifecycle(registry, cleanup)
pending, err := lifecycle.ApplyConfigRevisions(context.Background(), []ModelRevisionTransition{{ModelName: "model", ConfigRevision: "rev-new"}})
Expect(err).ToNot(HaveOccurred())
Expect(pending).To(Equal(1))
Expect(registry.calls).To(Equal([]string{"model:rev-new"}))
Expect(cleanup.seen).To(HaveLen(1))
})
It("quarantines the old identity and establishes the renamed identity", func() {
registry := &lifecycleRegistry{models: map[string][]nodes.NodeModel{
"old": {{ID: "old-replica", ModelName: "old"}},
}}
cleanup := &lifecycleCleanup{registry: registry}
lifecycle := NewDistributedModelRevisionLifecycle(registry, cleanup)
_, err := lifecycle.ApplyConfigRevisions(context.Background(), []ModelRevisionTransition{
{ModelName: "old", ConfigRevision: DeletedModelConfigRevision("old"), Disabled: true},
{ModelName: "new", ConfigRevision: "rev-renamed"},
})
Expect(err).ToNot(HaveOccurred())
Expect(registry.calls).To(Equal([]string{"old:" + DeletedModelConfigRevision("old"), "new:rev-renamed"}))
Expect(cleanup.seen).To(ConsistOf(nodes.NodeModel{ID: "old-replica", ModelName: "old"}))
})
It("does not clean up a partially advanced rename when the atomic transition fails", func() {
registry := &lifecycleRegistry{err: errors.New("injected rename transition failure")}
cleanup := &lifecycleCleanup{registry: registry}
lifecycle := NewDistributedModelRevisionLifecycle(registry, cleanup)
pending, err := lifecycle.ApplyConfigRevisions(context.Background(), []ModelRevisionTransition{
{ModelName: "old", ConfigRevision: DeletedModelConfigRevision("old"), Disabled: true},
{ModelName: "new", ConfigRevision: "rev-renamed"},
})
Expect(err).To(MatchError(ContainSubstring("injected rename transition failure")))
Expect(pending).To(BeZero())
Expect(registry.calls).To(Equal([]string{"old:" + DeletedModelConfigRevision("old"), "new:rev-renamed"}))
Expect(cleanup.seen).To(BeEmpty())
})
})