Files
LocalAI/core/services/galleryop/service.go
Ettore Di Giacinto 0eab3eb7f1 fix(distributed): read a departed agent tunnel as the routing fact it is
Task 4 gave agent workers tunnels and deliberately left the NodeType skip in
HealthMonitor.tunnelDeparted, with a spec asserting that an agent node whose
presence reader answers PresenceGone is NOT marked unhealthy. That spec was
scaffolding. It was true while an agent worker took its jobs and its verbs over
the message bus: a departure row for one said nothing about whether it could
work, and an early bug in the new tunnel client could otherwise have demoted a
fleet of healthy agent workers.

There is no bus. An agent worker is reachable through its tunnel and through
nothing else, so a departed agent tunnel means exactly what a departed backend
tunnel means: no live replica holds it, the departure has outlived the reconnect
grace, and that is a routing fact the scheduler and a reaper may act on. The
skip would now hide the only symptom an unreachable agent worker has. This is
the deliberate removal Task 4's M6 predicted, and task-4-report.md is where that
mutation already stands recorded red against the spec this commit deletes.

The skip existed at ONE site. router_liveness.go has none: its candidates come
from queries that already filter node_type = 'backend'. The two skips in
managers_distributed.go stay, because an agent worker still runs no backend
processes, so it has no backend to list and no backend op to apply.

Two node types can depart now, which is why the second half exists. Before this,
one type could depart and every per-node cache a departure left stale was
dropped from wherever its owner happened to notice, so a reader could not tell
which caches a demotion invalidated by reading the demotion path. Departure gets
ONE notification point. DepartureNotifier is edge triggered, because the monitor
runs on a ticker and a departed node stays departed; its subscribers are NAMED,
because what has to be caught is a forgotten cache and a count can say only that
one of four is missing; and NewHealthMonitor takes it as a required positional
argument, so a caller that does not pass one fails to compile.

Four caches subscribe: prefix-cache affinity in every model, probe freshness at
every address, in-flight staging operations, and the per-node breakdown of every
open gallery operation. The prefix-cache one is registered only when
prefix-cache routing is enabled, so --distributed-prefix-cache=false stays a
true no-op. The notification carries the node's name as well as its id, because
the staging tracker keys on the name and the other two key on the id, and a
subscriber should not have to read the registry from inside an eviction hook.

A departure notification is an act on absence, so it fires only on the routing
fact. A tunnel lost inside the grace, a worker that never dialled, a presence
query that failed and a stale heartbeat all announce nothing, asserted per node
type. The stale-heartbeat branch is excluded on purpose: it already marks the
node offline, which deletes its rows and runs the registry's replica-removed
hooks, so firing there too would double-evict and make the notification mean two
different things at its subscribers.

Assisted-by: Claude Opus 5 [claude-code]
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>
2026-09-20 03:05:35 +00:00

940 lines
34 KiB
Go

package galleryop
import (
"context"
"errors"
"fmt"
"runtime/debug"
"sync"
"time"
"github.com/mudler/LocalAI/core/config"
"github.com/mudler/LocalAI/core/gallery"
"github.com/mudler/LocalAI/core/services/distributed"
"github.com/mudler/LocalAI/core/services/messaging"
"github.com/mudler/LocalAI/pkg/downloader"
"github.com/mudler/LocalAI/pkg/model"
"github.com/mudler/LocalAI/pkg/system"
"github.com/mudler/xlog"
)
type GalleryService struct {
appConfig *config.ApplicationConfig
sync.Mutex
ModelGalleryChannel chan ManagementOp[gallery.GalleryModel, gallery.ModelConfig]
BackendGalleryChannel chan ManagementOp[gallery.GalleryBackend, any]
modelLoader *model.ModelLoader
modelManager ModelManager
backendManager BackendManager
statuses map[string]*OpStatus
cancellations map[string]cancellationActions
// Distributed mode (nil when not in distributed mode).
//
// broadcaster is the deployment's fan-out carrier, and it is ONE field on
// purpose: the progress and cancel publishes and the wildcard
// subscriptions SubscribeBroadcasts opens all read it, so this service
// cannot end up publishing on one carrier and listening on another. A
// replica in that state shows every gallery operation it started and none
// of its peers', with no error anywhere.
broadcaster messaging.Broadcaster
galleryStore *distributed.GalleryStore
broadcastSubs []messaging.Subscription
// OnBackendOpCompleted is fired after every successful install/upgrade/delete
// on the backend channel. The Application wires this to UpgradeChecker.TriggerCheck
// so `/api/backends/upgrades` stops surfacing a backend as upgradeable the moment
// the worker finishes — previously the cache only refreshed on the 6-hour tick,
// making manual upgrades look like they failed even when they hadn't.
//
// In distributed mode the same hook also fires on peer replicas via the
// SubjectCacheInvalidateBackends subscriber, so every replica's
// UpgradeChecker stays in sync.
OnBackendOpCompleted func()
// OnModelsChanged is fired on peer replicas when SOMEONE else publishes
// SubjectCacheInvalidateModels. The Application wires this to
// ModelConfigLoader.LoadModelConfigsFromPath so a chat completion that
// load-balances onto this replica can find the just-installed model.
// The originating replica reloads inline (models.go) so it does not need
// the hook.
OnModelsChanged func(messaging.CacheInvalidateEvent)
modelRevisionLifecycle interface {
ApplyConfigRevisions(context.Context, []config.ModelConfigRevisionTransition) (int, error)
}
}
// SetModelRevisionLifecycle wires the distributed config-generation boundary
// into gallery deletion without coupling gallery operations to node internals.
func (g *GalleryService) SetModelRevisionLifecycle(lifecycle interface {
ApplyConfigRevisions(context.Context, []config.ModelConfigRevisionTransition) (int, error)
}) {
g.Lock()
defer g.Unlock()
g.modelRevisionLifecycle = lifecycle
}
func NewGalleryService(appConfig *config.ApplicationConfig, ml *model.ModelLoader) *GalleryService {
return &GalleryService{
appConfig: appConfig,
ModelGalleryChannel: make(chan ManagementOp[gallery.GalleryModel, gallery.ModelConfig]),
BackendGalleryChannel: make(chan ManagementOp[gallery.GalleryBackend, any]),
modelLoader: ml,
modelManager: NewLocalModelManager(appConfig, ml),
backendManager: NewLocalBackendManager(appConfig, ml),
statuses: make(map[string]*OpStatus),
cancellations: make(map[string]cancellationActions),
}
}
// SetModelManager replaces the model manager (e.g. with a distributed implementation).
func (g *GalleryService) SetModelManager(m ModelManager) {
g.Lock()
defer g.Unlock()
g.modelManager = m
}
// SetBackendManager replaces the backend manager (e.g. with a distributed implementation).
func (g *GalleryService) SetBackendManager(b BackendManager) {
g.Lock()
defer g.Unlock()
g.backendManager = b
}
// BackendManager returns the current backend manager. Callers like the
// periodic upgrade checker need this so they run CheckUpgrades through the
// distributed implementation (which asks workers) instead of the frontend's
// local filesystem — the latter is always empty in distributed deployments.
func (g *GalleryService) BackendManager() BackendManager {
g.Lock()
defer g.Unlock()
return g.backendManager
}
// ModelArtifactMaterializer returns the controller-only acquisition capability
// used by startup paths that install gallery entries outside the operation loop.
func (g *GalleryService) ModelArtifactMaterializer() config.ArtifactMaterializer {
if g == nil || g.appConfig == nil {
return nil
}
return g.appConfig.ModelArtifactMaterializer
}
// SetBroadcaster wires the deployment's fan-out carrier, which is PostgreSQL
// LISTEN/NOTIFY. Both halves of this service's cross-replica sync ride it: the
// progress and cancel publishes, and the wildcard subscriptions
// SubscribeBroadcasts opens.
//
// It takes messaging.Broadcaster, which is Publish plus Subscribe and nothing
// else. Neither half needs request/reply or a queue group, and a wider
// parameter here would let this service acquire one without anyone deciding to.
func (g *GalleryService) SetBroadcaster(b messaging.Broadcaster) {
g.Lock()
defer g.Unlock()
g.broadcaster = b
}
// SetGalleryStore sets the PostgreSQL gallery store for distributed persistence.
func (g *GalleryService) SetGalleryStore(s *distributed.GalleryStore) {
g.Lock()
defer g.Unlock()
g.galleryStore = s
}
// ListBackends returns installed backends via the backend manager.
// In standalone mode this checks the local filesystem; in distributed mode
// it aggregates from all healthy worker nodes.
func (g *GalleryService) ListBackends() (gallery.SystemBackends, error) {
g.Lock()
mgr := g.backendManager
g.Unlock()
return mgr.ListBackends()
}
// DeleteBackend delegates backend deletion to the backend manager, which in distributed
// mode fans out the deletion to worker nodes via NATS.
func (g *GalleryService) DeleteBackend(name string) error {
g.Lock()
mgr := g.backendManager
g.Unlock()
return mgr.DeleteBackend(name)
}
func (g *GalleryService) UpdateStatus(s string, op *OpStatus) {
g.Lock()
// Preserve any per-node entries already accumulated by UpdateNodeProgress:
// the legacy progressCb path (used by the Phase 2 install bridge) calls
// UpdateStatus with a fresh *OpStatus on every tick, which would otherwise
// wipe the Nodes slice and leave the UI flickering between one node and
// another. If the caller explicitly populates Nodes on the incoming op,
// that wins; an empty Nodes slice on the incoming op is treated as "no
// new per-node data" and the previous Nodes are carried forward.
if op != nil {
if prev := g.statuses[s]; prev != nil {
if len(op.Nodes) == 0 && len(prev.Nodes) > 0 {
op.Nodes = prev.Nodes
}
// A job is a delete or an install for its whole life. markQueued is
// the only writer that knows which; every later status omits the
// flag, so an unset value means "no new information", not "this is
// an install".
if !op.Deletion {
op.Deletion = prev.Deletion
}
}
}
g.statuses[s] = op
store := g.galleryStore
nc := g.broadcaster
g.Unlock()
// I/O happens after Unlock. The broadcast loops back into our own
// wildcard subscriber (mergeStatus), which would deadlock on this mutex
// if we still held it. Holding the lock across a PostgreSQL round-trip
// would also stall every concurrent reader on each progress tick.
if store != nil && op != nil {
if op.Processed {
status, errMsg := "completed", ""
if op.Error != nil {
status = "failed"
errMsg = op.Error.Error()
}
if op.Cancelled {
status = "cancelled"
}
if err := store.UpdateStatus(s, status, errMsg); err != nil {
xlog.Warn("Failed to persist gallery operation status", "op_id", s, "error", err)
}
} else {
if err := store.UpdateProgress(s, op.Progress, op.Message, op.DownloadedFileSize, op.Cancellable,
distributed.OperationProgressDetails{
Phase: op.Phase, CurrentBytes: op.CurrentBytes, TotalBytes: op.TotalBytes,
}); err != nil {
xlog.Warn("Failed to persist gallery operation progress", "op_id", s, "error", err)
}
}
}
// Broadcast progress in distributed mode. The payload wraps the OpStatus
// with the opID so peer replicas reading the wildcard subject don't need
// to parse it back out of the subject string.
//
// A progress event carries one entry per node, so on a fleet of a few tens
// of workers it outgrows the notification cap and travels as a spilled row
// instead. That is the carrier's ordinary path and not an error: the peer
// receives the same bytes either way.
if nc != nil {
if err := nc.Publish(messaging.SubjectGalleryProgress(s), GalleryProgressEvent{
JobID: s,
Status: op,
}); err != nil {
xlog.Warn("Failed to broadcast gallery progress", "op_id", s, "error", err)
}
}
}
// publishCacheInvalidate broadcasts a cache invalidation event so peer
// replicas refresh whatever in-memory state mirrors disk. No-op when no
// carrier is wired (standalone mode).
//
// An invalidation is not a hint and is never traded away for a cheaper
// publish: a peer that misses one keeps serving from a cache it believes is
// valid, which is the one reading of a missed message this programme forbids.
// It goes out through Publish, which spills a message too large for a
// notification rather than refusing it.
func (g *GalleryService) publishCacheInvalidate(subject string, evt messaging.CacheInvalidateEvent) {
g.Lock()
nc := g.broadcaster
g.Unlock()
if nc == nil {
return
}
if err := nc.Publish(subject, evt); err != nil {
xlog.Warn("Failed to broadcast cache invalidation", "subject", subject, "error", err)
}
}
// BroadcastModelsChanged notifies peer replicas that a model config was
// created, edited, or removed out-of-band of the gallery install/delete
// channel (e.g. the admin /models/edit, /models/import and
// /models/toggle-state endpoints, which write the YAML and reload only the
// local in-memory loader). Peers receive it via OnModelsChanged and refresh
// their own ModelConfigLoader so a request load-balanced to any replica sees
// the same config. No-op in standalone mode (no carrier).
//
// op is "install" for a create/edit (the element must be (re)loaded from
// disk) or "delete" for a removal (the element must be pruned from memory,
// which a reload-from-path cannot do because the loader is additive).
func (g *GalleryService) BroadcastModelsChanged(element, op string) {
g.BroadcastModelsChangedRevision(element, op, "")
}
// BroadcastModelsChangedRevision includes the accepted semantic generation so
// peers can apply the same registry transition idempotently.
func (g *GalleryService) BroadcastModelsChangedRevision(element, op, configRevision string) {
g.publishCacheInvalidate(messaging.SubjectCacheInvalidateModels, messaging.CacheInvalidateEvent{
Element: element,
Op: op,
ConfigRevision: configRevision,
})
}
// mergeStatus is the broadcast-side merge: it updates the in-memory map from
// a peer's GalleryProgressEvent without re-publishing to the carrier or re-writing
// to PostgreSQL. UpdateStatus is the local-write entry point and does both;
// mergeStatus is what the wildcard subscriber calls. Splitting them avoids
// an echo loop (replica publishes → its own subscriber receives → mergeStatus
// silently re-applies → no second publish).
func (g *GalleryService) mergeStatus(opID string, op *OpStatus) {
if op == nil {
return
}
g.Lock()
defer g.Unlock()
prev := g.statuses[opID]
// A cancellation is terminal and a progress tick is not, and the carrier
// puts no order on the two. The owning replica's last tick is published
// before the admin's cancel and can be DELIVERED after it, on the owner's
// own echo as readily as on a peer; merging it wholesale would clear
// Cancelled and leave the operation reading as still running on that
// replica while every other one shows it stopped. A missed or late message
// must never read as an operation that was not cancelled, so a stale tick
// is dropped rather than merged. A terminal status that carries the
// cancellation still merges, which is how the final "cancelled" message
// arrives.
if prev != nil && prev.Cancelled && !op.Cancelled {
return
}
if len(op.Nodes) == 0 {
if prev != nil && len(prev.Nodes) > 0 {
op.Nodes = prev.Nodes
}
}
g.statuses[opID] = op
}
// UpdateNodeProgress merges a per-node progress tick into OpStatus.Nodes,
// keyed by nodeID, and mirrors the latest values into the aggregate
// Progress / FileName / DownloadedFileSize / TotalFileSize / Message
// fields so the legacy single-bar OperationsBar view keeps working
// unchanged alongside the new per-node breakdown.
//
// We deliberately do NOT delegate the aggregate mirror to UpdateStatus
// here: UpdateStatus overwrites the entire OpStatus, which would clobber
// the Nodes slice we just merged into. Doing the merge + mirror under a
// single lock keeps both views consistent and concurrent-safe.
func (g *GalleryService) UpdateNodeProgress(opID, nodeID string, np NodeProgress) {
g.Lock()
defer g.Unlock()
status := g.statuses[opID]
if status == nil {
status = &OpStatus{}
g.statuses[opID] = status
}
merged := false
for i := range status.Nodes {
if status.Nodes[i].NodeID == nodeID {
status.Nodes[i] = np
merged = true
break
}
}
if !merged {
status.Nodes = append(status.Nodes, np)
}
// Mirror the latest tick into the legacy aggregate fields so the
// existing single-bar UI keeps rendering meaningful progress.
status.FileName = np.FileName
status.Progress = np.Percentage
status.DownloadedFileSize = np.Current
status.TotalFileSize = np.Total
if np.Phase != "" {
status.Message = np.Phase
}
}
// DropNodeProgress removes nodeID's per-node rows from every operation that is
// still open.
//
// A per-node row is a claim that a node is doing something. When the node
// departs, nothing will ever move that row off "downloading": the operation
// stays open in /api/operations with a bar that never advances, and the node it
// names is not in the fleet any more.
//
// Only OPEN operations. A processed operation's breakdown is the record of what
// each node did, and rewriting it because a node later left would be reporting
// history rather than state. The aggregate mirror fields are left alone for the
// same reason: they are the last tick that happened, not a claim about now.
func (g *GalleryService) DropNodeProgress(nodeID string) {
if g == nil || nodeID == "" {
return
}
g.Lock()
defer g.Unlock()
for _, status := range g.statuses {
if status == nil || status.Processed {
continue
}
kept := make([]NodeProgress, 0, len(status.Nodes))
for _, np := range status.Nodes {
if np.NodeID != nodeID {
kept = append(kept, np)
}
}
if len(kept) == len(status.Nodes) {
continue
}
// Replaced rather than written through, which is the rule GetStatus
// documents: a reader holding the old slice header sees a consistent
// older breakdown instead of a torn one.
status.Nodes = kept
}
}
// GetStatus returns a COPY of the operation's status, not the stored pointer.
//
// The copy is what makes the lock mean anything. Every caller of this and of
// GetAllStatus only reads what it gets back, but the broadcast subscribers
// mutate the stored OpStatus IN PLACE - applyCancel sets Cancelled on the
// struct a peer's event names, mergeStatus rewrites the fields a peer sent - so
// handing out the pointer let an /api/operations response be marshalled while a
// peer's cancel was being written into it. Returning the pointer under a mutex
// serialises the map lookup and nothing else.
//
// Nodes is shared with the stored status and not deep-copied: UpdateNodeProgress
// replaces that slice rather than writing through it, so a reader holding the
// old header sees a consistent older breakdown rather than a torn one.
func (g *GalleryService) GetStatus(s string) *OpStatus {
g.Lock()
defer g.Unlock()
status, ok := g.statuses[s]
if !ok || status == nil {
return nil
}
copied := *status
return &copied
}
// GetAllStatus returns a snapshot of every operation's status. Same rule as
// GetStatus, and for the same reason: the map and the statuses in it are both
// copied, because this one is handed straight to a JSON encoder while peers'
// broadcasts are still arriving.
func (g *GalleryService) GetAllStatus() map[string]*OpStatus {
g.Lock()
defer g.Unlock()
snapshot := make(map[string]*OpStatus, len(g.statuses))
for id, status := range g.statuses {
if status == nil {
continue
}
copied := *status
snapshot[id] = &copied
}
return snapshot
}
// ReapStaleOperations marks abandoned in-progress operations (pending/
// downloading/processing) older than `age` as failed, so an op orphaned by a
// replica that died mid-flight does not linger as "processing" forever. The
// store's CleanStale runs once on startup; this exposes it for periodic
// invocation (a post-startup orphan is otherwise not reaped until the next
// restart). No-op when no gallery store is wired. Returns rows reaped.
func (g *GalleryService) ReapStaleOperations(age time.Duration) (int64, error) {
g.Lock()
store := g.galleryStore
g.Unlock()
if store == nil {
return 0, nil
}
// Collect the IDs before the update: once CleanStale flips them to
// "failed" they no longer match the stale predicate.
staleIDs, err := store.ListStale(age)
if err != nil {
xlog.Warn("Failed to list stale gallery operations", "error", err)
}
n, err := store.CleanStale(age)
if err != nil {
return 0, err
}
if n > 0 {
xlog.Info("Reaped stale gallery operations", "count", n)
}
// The database row is only half the picture. GET /models/jobs/<id> and
// /api/operations read the in-memory statuses map, which is populated
// locally and via the progress broadcast and never expires. An op
// orphaned by a replica that died mid-download therefore kept serving its
// last frozen tick (phase=downloading, processed=false, error=none) on
// every replica indefinitely, long after the reaper had already given up
// on the row. Reconcile the in-memory copy with the reap.
for _, id := range staleIDs {
g.failStaleStatus(id)
}
return n, nil
}
// failStaleStatus flips a locally-cached in-progress status to a terminal
// failure after its store row was reaped. Statuses that already reached a
// terminal state are left alone so a genuine completion or cancellation that
// raced the reaper is not rewritten as a failure.
func (g *GalleryService) failStaleStatus(id string) {
g.Lock()
st, ok := g.statuses[id]
if !ok || st == nil || st.Processed {
g.Unlock()
return
}
elementName := st.GalleryElementName
g.Unlock()
xlog.Warn("Marking orphaned gallery operation as failed", "op_id", id, "element", elementName)
g.UpdateStatus(id, &OpStatus{
Processed: true,
Error: errors.New("stale operation reaped (abandoned by a crashed or restarted instance)"),
Message: "error: stale operation reaped (abandoned by a crashed or restarted instance)",
GalleryElementName: elementName,
Cancellable: false,
})
}
// CancelOperation cancels an in-progress operation by its ID.
//
// In distributed mode the UI's cancel click may land on a different replica
// than the one running the operation. We still publish the cancel event in
// that case — the peer holding the cancellation func picks it up via the
// SubjectGalleryCancelWildcard subscriber and runs it locally. The caller
// gets a non-error reply so the UI shows the cancel as accepted.
func (g *GalleryService) CancelOperation(id string) error {
return g.stopOperation(id, false)
}
// PauseOperation stops an in-progress download while preserving its partial
// file. Re-submitting the same install resumes it through the downloader's
// existing HTTP Range support.
func (g *GalleryService) PauseOperation(id string) error {
return g.stopOperation(id, true)
}
func (g *GalleryService) stopOperation(id string, pause bool) error {
g.Lock()
if status, ok := g.statuses[id]; ok && status.Cancelled {
g.Unlock()
return fmt.Errorf("operation %q is already cancelled", id)
}
actions, localExists := g.cancellations[id]
if localExists {
delete(g.cancellations, id)
}
nc := g.broadcaster
store := g.galleryStore
if !localExists && nc == nil {
g.Unlock()
return fmt.Errorf("operation %q not found or already completed", id)
}
if status, ok := g.statuses[id]; ok {
status.Cancelled = true
status.Processed = true
status.Message = map[bool]string{true: "paused", false: "cancelled"}[pause]
} else {
g.statuses[id] = &OpStatus{
Cancelled: true,
Processed: true,
Message: map[bool]string{true: "paused", false: "cancelled"}[pause],
Cancellable: false,
}
}
g.Unlock()
// Persist the terminal status so the cancel survives a restart. Without
// this the row stays in its active state and re-hydrates straight back into
// processingBackends on the next replica boot — the UI spins again on an op
// the admin already cancelled. The peer that broadcasts wins the write; a
// no-op when standalone (store nil).
if store != nil {
if err := store.Cancel(id); err != nil {
xlog.Warn("Failed to persist gallery operation cancellation", "op_id", id, "error", err)
}
}
// I/O and user-provided callback after Unlock — the cancel-wildcard
// subscriber loops back into applyCancel on this same replica, which
// would otherwise deadlock on g.Mutex.
stopFunc := actions.cancel
if pause {
stopFunc = actions.pause
}
if stopFunc != nil {
stopFunc()
}
if nc != nil {
if err := nc.Publish(messaging.SubjectGalleryCancel(id), GalleryCancelEvent{JobID: id, Pause: pause}); err != nil {
xlog.Warn("Failed to broadcast gallery cancel", "op_id", id, "error", err)
}
}
return nil
}
// applyCancel is the broadcast-side counterpart to CancelOperation. The
// wildcard subscriber calls it when a peer publishes a cancel event:
// run the local cancel func if we have one (no echo via the carrier), and reflect
// the cancellation in the local statuses map. Idempotent: a replica that
// already cancelled this op locally treats the inbound event as a no-op.
func (g *GalleryService) applyCancel(id string, pause bool) {
g.Lock()
actions, hasCancel := g.cancellations[id]
if hasCancel {
delete(g.cancellations, id)
}
if status, ok := g.statuses[id]; ok {
if status.Cancelled {
g.Unlock()
return
}
status.Cancelled = true
status.Processed = true
status.Message = map[bool]string{true: "paused", false: "cancelled"}[pause]
} else {
g.statuses[id] = &OpStatus{
Cancelled: true,
Processed: true,
Message: map[bool]string{true: "paused", false: "cancelled"}[pause],
Cancellable: false,
}
}
g.Unlock()
// Invoke the cancel func after Unlock so a callback that touches
// GalleryService doesn't re-enter the mutex.
stopFunc := actions.cancel
if pause {
stopFunc = actions.pause
}
if hasCancel && stopFunc != nil {
stopFunc()
}
}
// NewUserCancellableContext returns a child context whose CancelFunc cancels
// with the downloader.ErrUserCancelled cause. This lets the download layer
// distinguish a deliberate user cancel (discard the half-downloaded .partial)
// from an incidental cancellation such as process shutdown (keep the .partial
// so the next run resumes via Range instead of restarting from zero).
// NewUserCancellableContext creates distinct callbacks for destructive cancel
// and resume-safe pause while sharing one operation context.
func NewUserCancellableContext(parent context.Context) (context.Context, context.CancelFunc, context.CancelFunc) {
ctx, cancelCause := context.WithCancelCause(parent)
return ctx,
func() { cancelCause(downloader.ErrUserCancelled) },
func() { cancelCause(context.Canceled) }
}
type cancellationActions struct {
cancel context.CancelFunc
pause context.CancelFunc
}
func (g *GalleryService) storeCancellation(id string, cancelFunc, pauseFunc context.CancelFunc) {
g.Lock()
defer g.Unlock()
if pauseFunc == nil {
pauseFunc = cancelFunc
}
g.cancellations[id] = cancellationActions{cancel: cancelFunc, pause: pauseFunc}
}
// StoreCancellation is a public method to store a cancellation function for an operation
// This allows cancellation functions to be stored immediately when operations are created,
// enabling cancellation of queued operations that haven't started processing yet.
func (g *GalleryService) StoreCancellation(id string, cancelFunc context.CancelFunc) {
g.storeCancellation(id, cancelFunc, cancelFunc)
}
// StoreCancellationActions registers distinct destructive-cancel and
// resume-safe pause callbacks for an operation.
func (g *GalleryService) StoreCancellationActions(id string, cancelFunc, pauseFunc context.CancelFunc) {
g.storeCancellation(id, cancelFunc, pauseFunc)
}
// removeCancellation removes a cancellation function when operation completes
func (g *GalleryService) removeCancellation(id string) {
g.Lock()
defer g.Unlock()
delete(g.cancellations, id)
}
// runOpHandler runs one operation handler and converts a panic into an error.
//
// The gallery worker is a single goroutine consuming both channels serially. A
// panic anywhere in an install handler (a malformed gallery entry, a nil
// dereference in a backend-specific path) took down the entire process with it,
// and every queued operation went with it. Containing the panic to the
// operation that caused it keeps the consumer alive so subsequent operations
// are still picked up, and surfaces the failure on the op itself instead of as
// an unexplained restart.
func runOpHandler(fn func() error) (err error) {
defer func() {
if r := recover(); r != nil {
xlog.Error("Gallery operation handler panicked", "panic", r, "stack", string(debug.Stack()))
err = fmt.Errorf("gallery operation handler panicked: %v", r)
}
}()
return fn()
}
func (g *GalleryService) Start(c context.Context, cl *config.ModelConfigLoader, systemState *system.SystemState) error {
// updates the status with an error
var updateError func(id string, e error)
if !g.appConfig.OpaqueErrors {
updateError = func(id string, e error) {
g.UpdateStatus(id, &OpStatus{Error: e, Processed: true, Message: "error: " + e.Error()})
}
} else {
updateError = func(id string, _ error) {
g.UpdateStatus(id, &OpStatus{Error: fmt.Errorf("an error occurred"), Processed: true})
}
}
go func() {
for {
select {
case <-c.Done():
return
case op := <-g.BackendGalleryChannel:
// Create context if not provided
if op.Context == nil {
op.Context, op.CancelFunc, op.PauseFunc = NewUserCancellableContext(c)
g.storeCancellation(op.ID, op.CancelFunc, op.PauseFunc)
} else if op.CancelFunc != nil {
g.storeCancellation(op.ID, op.CancelFunc, op.PauseFunc)
}
// Create DB record for distributed tracking
if g.galleryStore != nil {
opType := "backend_install"
if op.Delete {
opType = "backend_delete"
}
if err := g.galleryStore.Create(&distributed.GalleryOperationRecord{
ID: op.ID,
GalleryElementName: op.GalleryElementName,
OpType: opType,
Status: "pending",
// Create runs at dequeue, so this is the running-phase
// value: a running delete is not cancellable, an install
// is, matching the model channel. The queued phase before
// this point is cancellable either way (see markQueued).
Cancellable: !op.Delete,
}); err != nil {
// Not fatal: the install still runs and the in-memory
// status still updates. Logged because without the row
// the cross-replica dedup guard and hydration cannot
// see this operation at all.
xlog.Warn("Failed to create gallery operation record", "op_id", op.ID, "error", err)
}
}
err := runOpHandler(func() error { return g.backendHandler(&op, systemState) })
if err != nil {
updateError(op.ID, err)
} else if g.OnBackendOpCompleted != nil {
// Let listeners (e.g. UpgradeChecker) refresh their view of
// installed state. Run off the worker goroutine so a slow
// callback doesn't stall the next queued operation.
go g.OnBackendOpCompleted()
}
g.removeCancellation(op.ID)
case op := <-g.ModelGalleryChannel:
// Create context if not provided
if op.Context == nil {
op.Context, op.CancelFunc, op.PauseFunc = NewUserCancellableContext(c)
g.storeCancellation(op.ID, op.CancelFunc, op.PauseFunc)
} else if op.CancelFunc != nil {
g.storeCancellation(op.ID, op.CancelFunc, op.PauseFunc)
}
// Create DB record for distributed tracking
if g.galleryStore != nil {
opType := "model_install"
if op.Delete {
opType = "model_delete"
}
if err := g.galleryStore.Create(&distributed.GalleryOperationRecord{
ID: op.ID,
GalleryElementName: op.GalleryElementName,
OpType: opType,
Status: "pending",
// Create runs at dequeue, so this is the running-phase
// value: a running delete is not cancellable, an install
// is. The queued phase before this point is cancellable
// either way (see markQueued).
Cancellable: !op.Delete,
}); err != nil {
xlog.Warn("Failed to create gallery operation record", "op_id", op.ID, "error", err)
}
}
err := runOpHandler(func() error { return g.modelHandler(&op, cl, systemState) })
if err != nil {
updateError(op.ID, err)
}
g.removeCancellation(op.ID)
}
}
}()
return nil
}
// SubscribeBroadcasts opens the wildcard subscriptions that keep this
// replica's in-memory statuses + cancellation state in sync with peers.
// Returns an error if the progress subscription fails; cancel-sub failures
// are not fatal but are logged.
//
// Hydrate should be called before this so the freshly-started replica has
// the pre-existing operations before live updates start flowing.
func (g *GalleryService) SubscribeBroadcasts() error {
g.Lock()
nc := g.broadcaster
g.Unlock()
if nc == nil {
return nil
}
progressSub, err := messaging.SubscribeJSON(nc, messaging.SubjectGalleryProgressWildcard, func(evt GalleryProgressEvent) {
if evt.JobID == "" || evt.Status == nil {
return
}
g.mergeStatus(evt.JobID, evt.Status)
})
if err != nil {
return fmt.Errorf("subscribing to gallery progress wildcard: %w", err)
}
cancelSub, err := messaging.SubscribeJSON(nc, messaging.SubjectGalleryCancelWildcard, func(evt GalleryCancelEvent) {
if evt.JobID == "" {
return
}
g.applyCancel(evt.JobID, evt.Pause)
})
if err != nil {
if uerr := progressSub.Unsubscribe(); uerr != nil {
xlog.Warn("failed to unsubscribe partial gallery progress sub", "error", uerr)
}
return fmt.Errorf("subscribing to gallery cancel wildcard: %w", err)
}
modelsSub, err := messaging.SubscribeJSON(nc, messaging.SubjectCacheInvalidateModels, func(evt messaging.CacheInvalidateEvent) {
g.Lock()
cb := g.OnModelsChanged
g.Unlock()
if cb != nil {
cb(evt)
}
})
if err != nil {
if uerr := progressSub.Unsubscribe(); uerr != nil {
xlog.Warn("failed to unsubscribe partial gallery progress sub", "error", uerr)
}
if uerr := cancelSub.Unsubscribe(); uerr != nil {
xlog.Warn("failed to unsubscribe partial gallery cancel sub", "error", uerr)
}
return fmt.Errorf("subscribing to models invalidation: %w", err)
}
backendsSub, err := messaging.SubscribeJSON(nc, messaging.SubjectCacheInvalidateBackends, func(_ messaging.CacheInvalidateEvent) {
g.Lock()
cb := g.OnBackendOpCompleted
g.Unlock()
if cb != nil {
// Run off-goroutine so a slow UpgradeChecker doesn't stall the
// carrier's delivery loop. Matches the local fire-after-install path.
go cb()
}
})
if err != nil {
if uerr := progressSub.Unsubscribe(); uerr != nil {
xlog.Warn("failed to unsubscribe partial gallery progress sub", "error", uerr)
}
if uerr := cancelSub.Unsubscribe(); uerr != nil {
xlog.Warn("failed to unsubscribe partial gallery cancel sub", "error", uerr)
}
if uerr := modelsSub.Unsubscribe(); uerr != nil {
xlog.Warn("failed to unsubscribe partial models sub", "error", uerr)
}
return fmt.Errorf("subscribing to backends invalidation: %w", err)
}
g.Lock()
g.broadcastSubs = append(g.broadcastSubs, progressSub, cancelSub, modelsSub, backendsSub)
g.Unlock()
return nil
}
// CloseBroadcasts drops the wildcard subscriptions. Safe to call multiple times.
func (g *GalleryService) CloseBroadcasts() {
g.Lock()
subs := g.broadcastSubs
g.broadcastSubs = nil
g.Unlock()
for _, s := range subs {
if err := s.Unsubscribe(); err != nil {
xlog.Warn("GalleryService unsubscribe failed", "error", err)
}
}
}
// Hydrate loads still-active operations from the GalleryStore into the
// in-memory statuses map so a freshly-started replica does not return an
// empty /api/operations payload while a peer is mid-install. Idempotent.
// No-op when no store is wired.
//
// The reconstructed OpStatus carries no Error type — the DB stores the
// message as a string and Hydrate surfaces it via errors.New so the UI's
// "operation failed" banner survives a frontend restart.
func (g *GalleryService) Hydrate() error {
g.Lock()
store := g.galleryStore
g.Unlock()
if store == nil {
return nil
}
ops, err := store.ListActive()
if err != nil {
return fmt.Errorf("listing active gallery ops: %w", err)
}
g.Lock()
defer g.Unlock()
for _, op := range ops {
// Skip rows that already have an in-memory status — the live
// broadcast subscriber will fill any gaps with fresher data.
if _, ok := g.statuses[op.ID]; ok {
continue
}
st := &OpStatus{
Message: op.Message,
Progress: op.Progress,
Phase: op.Phase,
CurrentBytes: op.CurrentBytes,
TotalBytes: op.TotalBytes,
FileName: op.FileName,
TotalFileSize: op.TotalFileSize,
DownloadedFileSize: op.DownloadedFileSize,
GalleryElementName: op.GalleryElementName,
Cancellable: op.Cancellable,
Deletion: IsDeleteOpType(op.OpType),
}
if op.Error != "" {
st.Error = errors.New(op.Error)
}
g.statuses[op.ID] = st
}
xlog.Info("Hydrated gallery service statuses from store", "count", len(ops))
return nil
}