mirror of
https://github.com/mudler/LocalAI.git
synced 2026-07-30 18:09:05 -04:00
* feat(ui): record finished gallery operations in a bounded history ring The operations panel drops an operation the moment it succeeds, so a user who steps away cannot tell whether an install finished, failed or was never started. OpCache now keeps the last 50 terminal operations, recorded from the point where an op leaves the cache. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * test(ui): pin the history ring's dedupe, outcome order and start stamp Review of the history ring found four gaps. The dedupe guard and the bounded seen set were unreachable through the exported API and so had no coverage; an in-package spec file now drives opHistory directly. The outcome switch claimed an ordering was load bearing that nothing pinned, so an errored op that never reached Processed now has a spec. Two behaviour fixes come with it. StartedAt was the zero time for ops recovered from the store or replicated from a peer, since neither path stamps a start time, which would have rendered as a two-millennia duration; it now falls back to the finish time. Reusing a cache key with a fresh job ID orphaned the previous stamp, so Set and SetBackend now drop it. The comment on the outcome switch described a state the code cannot be in: CancelOperation sets Cancelled and Processed synchronously before the handler removes the entry, so status.Cancelled already covers the cancel endpoint. The !Processed clause stays for the dismiss endpoint firing on an in-flight op, and the comments now say so. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(ui): record operations that end on a peer replica The NATS end event is the only signal a replica gets for an install another replica ran. Record from applyEnd too, deduped by job ID so the originating replica does not record its own broadcast twice. Three start-stamp defects in the same path go with it. applyEnd now drops the stamp unconditionally, since recordTerminal only cleans up on the path where it found a cache key and an end event can overtake the local Set. applyStart drops the stamp of the job whose cache key it replaces, which a peer-driven retry previously stranded. And recordTerminal reads the stamp once instead of testing Exists and then reading, so a concurrent record for the same job can no longer delete the stamp between the two and let the zero time overwrite the finish-time fallback, which the Activity page would render as a two-millennia run. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): do not guess the outcome of a peer operation with no local status A replica that restarts mid-operation hydrates its OpCache keys from PostgreSQL, but gallery statuses are in-memory only and come back empty. The end broadcast then landed on recordTerminal's nil-status branch, which reads a missing status as queued-and-removed and filed a successful install as cancelled. That reading is right locally and wrong on the peer path, where a missing status means the outcome was never held here. recordTerminal now takes the source of the terminal event and records nothing when the peer path finds no status, restoring what the replica did before the end event started recording. The local path is unchanged. Also move the ApplyEndForTest seam to the conventional export_test.go, and stop the dedupe spec from claiming to guard the ring's seen set: the local delete removes the status keys, so the broadcast that follows returns before reaching it. An in-package spec that calls recordTerminal twice does the pinning. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(api): add GET and DELETE /api/operations/history Admin gated like the rest of the operations API. The live /api/operations payload is unchanged so the one second poll stays small. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(ui): expose operation history through OperationsContext Fetched on demand and when the live list shrinks, never on the one second poll interval. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): detect operation departure by identity and ignore committing ops in the ETA gate Refetching history on a shrinking live count missed a completion that coincided with a start, which is the common case during a batch install. Track the live job IDs instead, so any departure triggers the refetch regardless of how the count moved. An operation that has finished downloading stays live at currentBytes == totalBytes for the whole commit and install phase and can never produce an estimate, so counting it in the all-or-nothing gate blanked every other operation's time remaining for as long as it lasted. Only operations still moving bytes get a vote. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): let only downloading operations gate the time remaining estimate Verifying pins an operation flat below its total for the whole sha256 pass: the AfterDownload hook reports completedBytes plus the finished file against a total summed over every file, then hashes synchronously without emitting progress. Files download sequentially, so a 15 shard model enters that window 14 times, and a byte comparison cannot see it because the counter is genuinely below the total throughout. Gating on phase closes resolving, verifying, committing and persisting in one predicate, so a quiet neighbour no longer blanks every other operation's estimate for minutes at a time. The byte clauses stay: a producer can report downloading with bytes already at the total. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(ui): collapse the operations bar to a single line Four concurrent installs used to take four rows above every page. The strip now shows one operation, failure first, with a counter linking to Activity. The close button hides the strip and no longer cancels an install. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): keep the operations strip from widening the page and from muting a failure A long install error made the strip report a 1600px minimum width, which sized main-content to fit and gave every page under it a horizontal scrollbar. Inline-size containment plus shrinkable detail and bytes cells keep it inside the viewport. Hiding is no longer able to swallow the hidden job's own failure, a completed removal or staging says so instead of claiming an install, a cancelling operation renders as cancelling, and the live region no longer covers the per-second percentage. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): shrink main-content instead of containing the strip, and expose progress min-width on .main-content is what actually lets a long install error shrink, and unlike inline-size containment it has no browser support floor and no latent collapse if the strip ever lands in a shrink-to-fit context. It matches what .app-layout-chat .main-content already does, and it clears pre-existing horizontal overflow on narrow viewports as a side effect. The progress track is now a labelled progressbar, so assistive tech can read the value on demand rather than losing it to the aria-hidden that stopped the live region re-announcing every poll. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(ui): add the live operation card for the Activity page Carries the detail the one-line strip has to drop: phase, bytes, the per-node breakdown for cluster installs, and a labelled Cancel button. Cancelling is destructive, so it gets a labelled button rather than a glyph. A cancelling operation drops its progress bar and its time estimate, the same call the strip makes: a percentage still climbing under "Cancelling" reads as the cancel not having taken. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): give the operation card a verb, live node disclosure and its per-node detail The card carried no verb, so an install, a removal and a staging op rendered as spinner plus name plus kind tag and were indistinguishable. It now runs the same verb and icon chain as the one-line strip, which is what stops the page that is meant to carry more detail from carrying less. The auto-expand default was evaluated once at mount. An operation is listed as soon as it is admitted but its nodes are filled in only when the fan-out starts reporting, so a card mounted at creation latched on the empty list and stayed collapsed. The default is a live expression now, and state holds only an explicit choice. Also: an optional onRetry gates a Retry button, so the page can own the install reconstruction without the card ever showing a control with nothing behind it; the disclosure moved above the region it controls and gained aria-controls; the toggle is gated at more than one node so the count is never "1 nodes"; an unmapped node status is passed through instead of being relabelled "Queued"; error text is clamped with the full string in the title; and file_name plus the per-node progress bar are rendered again, reviving three CSS rules that had gone dead along with the detail they styled. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(ui): add the Activity page Live operations, unacknowledged failures and the record of what finished, at /app/activity in the Operate console. Cancelling an install now lives here behind a labelled button rather than on the strip, and a failed install can be retried: the retry dismisses the failure first so it still reaches the record, then reissues the model, backend or node-scoped backend install. The sidebar Operate entry carries the operation count. The console rail is only rendered on an Operate route and can be collapsed, so a badge there could vanish while operations were still running. Two follow-ups from review fold in here: a failed removal or staging job no longer reports a failed install on either the card or the strip, and the card's error text can shrink so one unbroken token cannot widen the card. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): dismiss operations by job, and stop the Activity page contradicting itself Dismissing resolved the job by display id, but /api/operations strips the "node:<nodeID>:" prefix before emitting, so a local install and a node-scoped install of one backend arrive as two jobs sharing one id. Dismissing by id retired whichever came first. That defeated the guarantee retry was built around: with the wrong job dismissed, the reinstall overwrote the acted-on failure's opcache entry in place, bypassing recordTerminal, while an unrelated failure vanished from Needs attention. dismissFailedOp, the card's dismiss control and the strip now all pass the jobID, which is what the endpoint takes. A filter matching nothing rendered the "nothing has ever run" empty state while the header counted the records the filter had hidden. The empty state is now gated on the All chip and a narrowed view gets its own message plus a way back; the header counts the instance rather than the chip, so selecting Backends no longer reports "Nothing running" over running model installs. Also: the summary drops a zero clause instead of rendering "0 needs attention" on the happy path and pluralises both counts; a record duration is floored at "< 1s" and rejected above a day, so a zero-value start stamp cannot render a span of millennia and a zero span cannot render "installed in" with nothing after it; a deletion cancelled mid-flight reports the cancellation rather than claiming it was removed; and the retry variant comment names the fix instead of calling the gap closed. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: document the Activity page and the operations history endpoints Adds an Activity page under Operations covering the one-line operations strip, the /app/activity sections and filters, per-operation cancel, retry and dismiss, the in-memory 50-entry record, and the sidebar count. Documents GET and DELETE /api/operations/history, and fills the gap in the admin-only endpoint list, which also omitted the pre-existing POST /api/operations/:jobID/dismiss. Corrects the distributed-mode install-watching section: the per-node breakdown now lives on the Activity page rather than on the strip, which rolls a fan-out up into a single phrase. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: correct nine details in the Activity page documentation The operations strip never renders a file name: its detail line is the error, the node roll-up, the target node, the phase or the queued note. Drops the stale clause in the distributed-mode section, where the per-node bullet is now the only place a file name is described. Scopes the phase vocabulary to artifact-backed gallery models, since a plain GGUF install emits no phase. Corrects the per-node list: the toggle exists for any fan-out of two or more workers and the four-node threshold only governs whether it starts open, while the N nodes tag needs more than one node. Notes that a cancelled operation can sit in the live section reading Cancelling, that cluster staging never reaches the record, and that Clear history appears only when the record has something in it. Names the operations response envelope, with a JSON example, so callers do not index a bare array, and stops describing the icon-only dismiss control as a labelled button. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: drop the unreachable Cancelling state and scope the byte claims An operation can only report isCancelled while it is unprocessed, but every writer of Cancelled sets Processed in the same breath, on the peer path as much as the local one, and the cache evicts cancelled entries before the handler sees them. The state cannot reach the page, so the live section is described again as running or queued operations. Byte counts come from the artifact bridge alone, the same producer as the phase, so a plain GGUF install, a removal and a backend install report none. Scopes both to artifact-backed gallery models and leaves the verb, the name and the percentage as what every operation shows. A worker backend install reports its bytes through fields the operations payload does not carry, so the distributed section now describes the percentage and the node roll-up, with per-file counts pointed at the per-node detail. Also: staging jobs carry no error, so they never reach Needs attention and Retry never had a staging case to exclude; an install that involves workers is no longer called node-scoped, which this page uses for node-targeted installs; and the record timestamps carry nanoseconds. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * docs: state only the verb and the name as unconditional on the strip The percentage is as conditional as the bytes were: it renders only for a running operation that has reported progress, so a queued operation, a failed one and a removal never carry it. A removal in particular sits at progress zero for its whole visible life, since the delete path reports none and its completion is filtered out. Both the strip and the card paragraphs now lead with what always shows and list the rest as conditions. The Cluster chip matches on a node list that finished operations do not carry, so a fan-out install leaves the chip once it reaches the record. Scoped that claim to the live sections. Two more of the same shape, found by re-reading each clause alone: the strip also appears for a failure, which is not running, and the four-second hold only applies when nothing replaces the operation that just finished. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): stop reporting a cancelled install as installed, and make queued real Three defects that all trace to one root cause: `isCancelled: true` is unreachable from /api/operations. Every writer of Cancelled=true also sets Processed=true, the handler skips Processed && Cancelled, and OpCache.GetStatus evicts a cancelled op before the handler iterates it. Cancelling the last running operation put a green "Installed model X" on the strip for four seconds: the completion hold was guarded by `!previous.isCancelled`, which is dead. A cancellation deletes the operation server side, so the strip sees exactly what it sees on a completion, and nothing in the payload separates the two. The signal now comes from the side that issued the cancel: the operations context remembers the job IDs it cancelled (pruned after a minute) and the strip asks before it holds anything. A cancelled operation goes as soon as it stops; the record already reports it as cancelled. isQueued was set only when the gallery status was missing, but markQueued publishes a "queued" status at admission, so a queued op has a status for its whole queued life and the state was unreachable outside a microsecond window. Every operation waiting behind a running install rendered as "Installing model X" with a spinner. The queued phase is now the signal, via an exported PhaseQueued and a nil-safe OpStatus.IsQueued() next to the writer. With those two fixed, the Cancelling state has no way to be entered: cancelling is instantaneous from the API's point of view. Its branches, CSS, locale key and the isCancelled field itself are removed rather than left for a future reader to assume they work. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(ui): keep a removal a removal, and say what an install is doing OpStatus.Deletion was set once, at admission, and lost on the next status write: UpdateStatus replaces the whole status and only carried Nodes forward. Every later writer (the worker's first write, the progress ticks, the failure path) leaves the field at its zero value, so the flag survived only the queued window, and both surfaces test isQueued first. The reachable consequence is that a failed removal reported itself as a failed install, which is exactly the shape the Activity page offers Retry for, and Retry installs: pressing it on a removal that failed re-downloaded the model. A running delete also rendered as "Installing model X" with a spinner, and a successful one as "Installed model X". Carry Deletion forward the way Nodes already is. A job is a delete or an install for its whole life; an unset flag means "no new information", not "this is an install". Pinned by Go specs on both the service and /api/operations: the existing Playwright specs were green only because they stubbed a payload the server could not emit. Also restore the operation's own status message on the Activity card. Phases and byte counters exist only on the managed-artifact path, so a legacy files: gallery model and every backend install rendered a sub-row with nothing in it but the verb. The strip stays terse on purpose. And give the strip's name a min-width floor: overflow: hidden zeroes its automatic minimum, so a long error squeezed the name down to "mod…" and the identity of the thing that broke was the first thing lost. primaryOperation is made module-private: its comment claimed the Activity page selected the same operation, but that page shows all of them, partitioned into failed and running, and never imported it. Assisted-by: Claude Code:Opus 5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * feat(activity): read the operations record from PostgreSQL The Activity page's record of finished installs and removals was a 50-entry in-memory ring per frontend replica. In distributed mode that is the wrong place for it: each replica keeps its own copy, a replica added by a scale-out or a rolling deploy starts empty and never backfills, and "Clear history" clears only the replica that served the request, so the record reappears on the next poll routed elsewhere. The data is already in gallery_operations. Read it from there. GalleryStore gains ListTerminal and ClearTerminal, sharing a lifted terminalStatuses set with CleanOld so there is one definition of "finished". ListTerminal orders by updated_at, when the operation reached its terminal status, because the record reports what finished and when. OpCache.History and ClearHistory dispatch on whether a store is wired, so the HTTP handlers and the OpRecord JSON shape are unchanged and the page needed no change. A failed store read falls back to the local ring rather than blanking the page, and ClearHistory empties the ring as well so a database blip cannot resurrect a record the admin just cleared. The name derivation in recordTerminal is lifted into operationDisplayName and used by both paths, so the ring and the store cannot name the same operation differently. Also fixes a pre-existing bug the store path made visible: the backend channel hardcoded op_type "backend_install" even for a removal, while the model channel derives model_install/model_delete from op.Delete. Both channels carry the same ManagementOp, whose Delete field the backend handler already branches on, so the backend channel now derives backend_delete the same way. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(activity): keep a cancelled operation cancelled, and report a failed clear Review follow-up on the store-backed Activity record. A cancelled install was recorded as a failure. The cancel handler persists "cancelled" synchronously, then the handler goroutine unwinds with the context error and Start hands that to updateError unconditionally, which overwrote the row with "failed: context canceled". The page rendered a cancelled install as a red failure card offering Retry, with a raw context error as the reason. Fixed in GalleryStore rather than in Start, because an operation finishes once and the paths that retire one are not mutually exclusive: UpdateStatus now refuses to rewrite a row that already reached a terminal status. That also pins updated_at to when the operation really finished, which is the key the record is ordered by, and Create's upsert now freezes the same columns so a worker dequeuing an operation the admin cancelled while it was queued cannot reopen it as pending. ClearHistory returned nothing, so a failed delete logged a warning while the handler still answered 200. The admin watched the record clear and come back on the next fetch with nothing said about why. It now returns the error, the DELETE handler answers 500, and the store is cleared before the local ring so a failure leaves the fallback record intact rather than faking an empty one. Hydrate is the only reader that decides from op_type whether an operation is a removal, and it tested for "model_delete" exactly, so the backend_delete added in the previous commit hydrated as an install: a replica restarting during a backend removal rendered "Installing backend X". Both discriminations now go through IsDeleteOpType/IsBackendOpType so a fifth op_type cannot silently read as an install in whichever consumer was missed. Also: the backend channel now persists Cancellable as !op.Delete, matching the model channel; IsBackend falls back to the op_type prefix, since is_backend_op is only written by UpsertCacheKey and the rows needing the name fallback were reporting backend operations as models; and an unrecognized terminal status is logged rather than quietly filed as a success, which is what the comment already claimed. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(activity): keep a reaped operation correctable by its real outcome The terminal-status freeze added in the previous commit was too wide. It froze "failed" alongside "completed" and "cancelled", and the stale reaper writes "failed" onto operations that are still going to run. The gallery worker is a single goroutine consuming both channels serially, so an operation queued behind a large download sits in "pending" with nothing bumping updated_at, and ReapStaleOperations gives up on it after 30 minutes. That used to be self-healing: the worker dequeued it, Create reset the row to "pending", and the operation reported its real outcome. With the freeze the row stayed "failed" forever while the install ran and succeeded underneath it: a red failure card offering Retry for a model that is installed, omitted from ListActive so no replica hydrates it, and no longer deduped cluster-wide by FindDuplicate. Freeze on ("completed", "cancelled") instead. That is all the cancelled-install fix ever needed, and it leaves a failure correctable by what actually happened. The set is separate from terminalStatuses, which ListTerminal, ClearTerminal and CleanOld all still want in full, because the two mean different things: a failure can be superseded by a real outcome, a completion or a cancellation is the real outcome. UpdateStatus now writes the error column unconditionally, so a corrected outcome drops the previous attempt's reason rather than being recorded as completed while still carrying "stale operation reaped" as its error. Also adds the route-level spec for the 500 branch of DELETE /api/operations/history, and trims a comment that credited the persisted cancellable column with more than it survives long enough to do. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(activity): offer Cancel in the phase that can honour it The cancellable flag was set at both ends of an operation's life and was wrong at both, in opposite directions. A queued operation is cancellable whatever it is. EnqueueModelOp and EnqueueBackendOp select on the operation context, so cancelling one that is still waiting releases the delivery goroutine and abandonQueued retires it: the worker never sees it, nothing is downloaded, nothing is deleted. markQueued nevertheless wrote Cancellable: !deletion, so a queued removal reported cancellable: false and the UI hid the Cancel button in the one window where pressing it both works and leaves no trace. A removal queued behind a large install was stuck there until the install finished. A running removal is not cancellable at all. DeleteModel and DeleteBackend take no context, and modelHandler only checks the operation context after the call returns, so a "cancelled" verdict would land after the model was already gone. Both handlers nevertheless wrote Cancellable: true unconditionally at entry, ahead of the op.Delete branch, offering a Cancel button the server cannot honour. So the queued phase is more cancellable than the running phase, which is the reverse of the usual shape. markQueued now reports true unconditionally, and the handler-entry writes report !op.Delete. Both sites carry a comment saying why, because reading either one alone suggests the other is a bug. GalleryStore.Create keeps !op.Delete: it runs at dequeue, so its value already describes the running phase. Its comment now says so. Specs cover queued removal, queued install, running removal and running install through the handlers, plus the queued-removal case through /api/operations where the flag is consumed, plus the behaviour the whole asymmetry rests on: a removal cancelled while queued never reaches the worker and deletes nothing. No existing spec asserted the old values. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Write] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> * fix(activity): clamp the installer message, and add a real-binary e2e spec Running the page against a real local-ai showed the legacy installer message wrapping to three lines and dominating the card: it embeds an absolute file path, so it is both long and a single unbreakable token. One line, ellipsised, full text in the title, matching what the error string already does. The spec that found it runs with no route stubbing at all. Every other spec here stubs /api/operations, which is how a payload the server cannot emit (isDeletion true on a live operation) stayed green through a full review while the UI rendered a removal as an install. It is skipped unless LOCALAI_REAL_BINARY is set, so CI is unaffected. Assisted-by: Claude Code:claude-opus-5 [Read] [Edit] [Bash] Signed-off-by: Ettore Di Giacinto <mudler@localai.io> --------- Signed-off-by: Ettore Di Giacinto <mudler@localai.io> Co-authored-by: Ettore Di Giacinto <mudler@localai.io>
745 lines
26 KiB
Go
745 lines
26 KiB
Go
package galleryop
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"errors"
|
|
"fmt"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/mudler/LocalAI/core/config"
|
|
"github.com/mudler/LocalAI/core/services/distributed"
|
|
"github.com/mudler/LocalAI/core/services/messaging"
|
|
"github.com/mudler/LocalAI/pkg/xsync"
|
|
"github.com/mudler/xlog"
|
|
)
|
|
|
|
type ManagementOp[T any, E any] struct {
|
|
ID string
|
|
GalleryElementName string
|
|
Delete bool
|
|
|
|
Req T
|
|
|
|
// If specified, we install directly the gallery element
|
|
GalleryElement *E
|
|
|
|
Galleries []config.Gallery
|
|
BackendGalleries []config.Gallery
|
|
|
|
// Context for cancellation support
|
|
Context context.Context
|
|
CancelFunc context.CancelFunc
|
|
|
|
// External backend installation parameters (for OCI/URL/path)
|
|
// These are used when installing backends from external sources rather than galleries
|
|
ExternalURI string // The OCI image, URL, or path
|
|
ExternalName string // Custom name for the backend
|
|
ExternalAlias string // Custom alias for the backend
|
|
|
|
// TargetNodeID scopes a backend install/upgrade to a single worker node.
|
|
// Empty means fan out to every healthy backend node (the previous behavior).
|
|
// Set by InstallBackendOnNodeEndpoint so an admin can install a hardware-specific
|
|
// build on one node without touching the rest of the cluster.
|
|
TargetNodeID string
|
|
|
|
// Variant pins a model install to one of the gallery entry's declared
|
|
// variants, by that variant's model name. Empty means auto-select: LocalAI
|
|
// picks the largest variant this host's backend support and memory can
|
|
// actually run, and falls back to the entry's own build.
|
|
//
|
|
// A name that is not among the entry's variants fails the install rather
|
|
// than quietly auto-selecting, so a typo cannot masquerade as a choice.
|
|
Variant string
|
|
|
|
// Upgrade is true if this is an upgrade operation (not a fresh install)
|
|
Upgrade bool
|
|
|
|
// Force reinstalls a backend even when it is already installed and
|
|
// runnable. Without it a backend install op is idempotent — API clients
|
|
// that ensure a backend exists on every boot must not trigger a full
|
|
// artifact re-download each time. The UI's explicit "Reinstall backend"
|
|
// action sets it.
|
|
Force bool
|
|
}
|
|
|
|
type OpStatus struct {
|
|
Deletion bool `json:"deletion"` // Deletion is true if the operation is a deletion
|
|
FileName string `json:"file_name"`
|
|
Error error `json:"-"` // see MarshalJSON: serialized to "error" as a string
|
|
Processed bool `json:"processed"`
|
|
Message string `json:"message"`
|
|
Progress float64 `json:"progress"`
|
|
Phase string `json:"phase,omitempty"`
|
|
CurrentBytes int64 `json:"current_bytes,omitempty"`
|
|
TotalBytes int64 `json:"total_bytes,omitempty"`
|
|
TotalFileSize string `json:"file_size"`
|
|
DownloadedFileSize string `json:"downloaded_size"`
|
|
GalleryElementName string `json:"gallery_element_name"`
|
|
Cancelled bool `json:"cancelled"` // Cancelled is true if the operation was cancelled
|
|
Cancellable bool `json:"cancellable"` // Cancellable is true if the operation can be cancelled
|
|
|
|
// Nodes is the per-node breakdown for a fanned-out backend install.
|
|
// Populated by DistributedBackendManager (per-node terminal status)
|
|
// and by the Phase 2 progress bridge (per-byte ticks). The
|
|
// /api/operations handler surfaces this so the UI can render an
|
|
// expandable per-node view of an in-flight install.
|
|
Nodes []NodeProgress `json:"nodes,omitempty"`
|
|
}
|
|
|
|
// opStatusWire is the JSON shape used when an OpStatus crosses a process
|
|
// boundary (NATS broadcast). The Error field on OpStatus is an `error`
|
|
// interface, which json.Marshal flattens to `{}` because the concrete error
|
|
// type usually has no exported fields — so a failed install replicated to a
|
|
// peer frontend would arrive with a nil error and the UI would never surface
|
|
// the failure. opStatusWire serializes the error as its Error() string and
|
|
// reconstructs it on read.
|
|
type opStatusWire struct {
|
|
Deletion bool `json:"deletion"`
|
|
FileName string `json:"file_name"`
|
|
ErrorMessage string `json:"error,omitempty"`
|
|
Processed bool `json:"processed"`
|
|
Message string `json:"message"`
|
|
Progress float64 `json:"progress"`
|
|
Phase string `json:"phase,omitempty"`
|
|
CurrentBytes int64 `json:"current_bytes,omitempty"`
|
|
TotalBytes int64 `json:"total_bytes,omitempty"`
|
|
TotalFileSize string `json:"file_size"`
|
|
DownloadedFileSize string `json:"downloaded_size"`
|
|
GalleryElementName string `json:"gallery_element_name"`
|
|
Cancelled bool `json:"cancelled"`
|
|
Cancellable bool `json:"cancellable"`
|
|
Nodes []NodeProgress `json:"nodes,omitempty"`
|
|
}
|
|
|
|
func (o OpStatus) MarshalJSON() ([]byte, error) {
|
|
w := opStatusWire{
|
|
Deletion: o.Deletion,
|
|
FileName: o.FileName,
|
|
Processed: o.Processed,
|
|
Message: o.Message,
|
|
Progress: o.Progress,
|
|
Phase: o.Phase,
|
|
CurrentBytes: o.CurrentBytes,
|
|
TotalBytes: o.TotalBytes,
|
|
TotalFileSize: o.TotalFileSize,
|
|
DownloadedFileSize: o.DownloadedFileSize,
|
|
GalleryElementName: o.GalleryElementName,
|
|
Cancelled: o.Cancelled,
|
|
Cancellable: o.Cancellable,
|
|
Nodes: o.Nodes,
|
|
}
|
|
if o.Error != nil {
|
|
w.ErrorMessage = o.Error.Error()
|
|
}
|
|
return json.Marshal(w)
|
|
}
|
|
|
|
func (o *OpStatus) UnmarshalJSON(data []byte) error {
|
|
var w opStatusWire
|
|
if err := json.Unmarshal(data, &w); err != nil {
|
|
return err
|
|
}
|
|
o.Deletion = w.Deletion
|
|
o.FileName = w.FileName
|
|
o.Processed = w.Processed
|
|
o.Message = w.Message
|
|
o.Progress = w.Progress
|
|
o.Phase = w.Phase
|
|
o.CurrentBytes = w.CurrentBytes
|
|
o.TotalBytes = w.TotalBytes
|
|
o.TotalFileSize = w.TotalFileSize
|
|
o.DownloadedFileSize = w.DownloadedFileSize
|
|
o.GalleryElementName = w.GalleryElementName
|
|
o.Cancelled = w.Cancelled
|
|
o.Cancellable = w.Cancellable
|
|
o.Nodes = w.Nodes
|
|
if w.ErrorMessage != "" {
|
|
o.Error = errors.New(w.ErrorMessage)
|
|
} else {
|
|
o.Error = nil
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// OpCacheEvent is the NATS payload broadcast by frontend replicas when an
|
|
// admin operation is admitted (SubjectGalleryOpStart) or dismissed
|
|
// (SubjectGalleryOpEnd). Peers merge these into their local OpCache so a
|
|
// load-balanced /api/operations poll never returns an empty list while a
|
|
// peer is mid-install.
|
|
type OpCacheEvent struct {
|
|
JobID string `json:"job_id"`
|
|
CacheKey string `json:"cache_key"`
|
|
IsBackend bool `json:"is_backend"`
|
|
}
|
|
|
|
// GalleryProgressEvent is the NATS payload for an OpStatus broadcast. It
|
|
// wraps OpStatus with the opID/JobID so subscribers reading the wildcard
|
|
// subject don't need to parse it back out of the NATS subject string.
|
|
type GalleryProgressEvent struct {
|
|
JobID string `json:"job_id"`
|
|
Status *OpStatus `json:"status"`
|
|
}
|
|
|
|
// GalleryCancelEvent is the NATS payload for a gallery cancellation. The
|
|
// local cancellation func may live on a different frontend replica than the
|
|
// one that received the UI cancel button click; the broadcast subscriber
|
|
// runs the cancel func on whichever replica registered it.
|
|
type GalleryCancelEvent struct {
|
|
JobID string `json:"id"`
|
|
}
|
|
|
|
// NodeStatus values shared between NodeProgress (per-node tick) and the
|
|
// NodeOpStatus surfaced by DistributedBackendManager's fan-out. Defined
|
|
// as exported constants so producers (the manager, the progress bridge)
|
|
// and consumers (the /api/operations handler, the React OperationsBar
|
|
// through its JSON contract) stay in sync via a single source of truth.
|
|
const (
|
|
NodeStatusQueued = "queued" // node accepted the intent but install has not started
|
|
NodeStatusDownloading = "downloading" // worker is actively pulling the OCI image
|
|
NodeStatusRunningOnWorker = "running_on_worker" // NATS round-trip timed out but worker is still installing
|
|
NodeStatusSuccess = "success" // install completed on this node
|
|
NodeStatusError = "error" // install failed on this node
|
|
)
|
|
|
|
// NodeProgress is a single node's contribution to a backend install
|
|
// operation. Populated by DistributedBackendManager (per-node terminal
|
|
// status) and by the Phase 2 progress bridge (per-byte ticks). Read by
|
|
// the /api/operations handler so the UI can render an expandable
|
|
// per-node breakdown.
|
|
//
|
|
// Status holds one of the NodeStatus* constants above.
|
|
type NodeProgress struct {
|
|
NodeID string `json:"node_id"`
|
|
NodeName string `json:"node_name"`
|
|
Status string `json:"status"`
|
|
FileName string `json:"file_name,omitempty"`
|
|
Current string `json:"current,omitempty"`
|
|
Total string `json:"total,omitempty"`
|
|
Percentage float64 `json:"percentage"`
|
|
Phase string `json:"phase,omitempty"`
|
|
Error string `json:"error,omitempty"`
|
|
}
|
|
|
|
type OpCache struct {
|
|
status *xsync.SyncedMap[string, string]
|
|
backendOps *xsync.SyncedMap[string, bool] // Tracks which operations are backend operations
|
|
galleryService *GalleryService
|
|
|
|
// Finished operations, for GET /api/operations/history. started stamps the
|
|
// start time when an op enters the cache so the record can report duration.
|
|
history *opHistory
|
|
started *xsync.SyncedMap[string, time.Time]
|
|
|
|
// Distributed sync (nil when standalone).
|
|
mu sync.RWMutex
|
|
nats messaging.MessagingClient
|
|
store *distributed.GalleryStore
|
|
subs []messaging.Subscription
|
|
}
|
|
|
|
func NewOpCache(galleryService *GalleryService) *OpCache {
|
|
return &OpCache{
|
|
status: xsync.NewSyncedMap[string, string](),
|
|
backendOps: xsync.NewSyncedMap[string, bool](),
|
|
galleryService: galleryService,
|
|
history: newOpHistory(DefaultHistorySize),
|
|
started: xsync.NewSyncedMap[string, time.Time](),
|
|
}
|
|
}
|
|
|
|
// SetMessagingClient enables cross-replica OpCache sync. Once set, Set/
|
|
// SetBackend/DeleteUUID publish OpCacheEvent messages that peer OpCaches
|
|
// merge into their local maps. Call Start after this to subscribe.
|
|
func (m *OpCache) SetMessagingClient(nc messaging.MessagingClient) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
m.nats = nc
|
|
}
|
|
|
|
// SetGalleryStore enables PostgreSQL-backed OpCache persistence.
|
|
// Set/SetBackend upsert the cache_key + is_backend_op columns; Start
|
|
// hydrates the in-memory maps from active rows so a freshly-started
|
|
// replica does not return an empty /api/operations payload while a peer
|
|
// is mid-install.
|
|
func (m *OpCache) SetGalleryStore(s *distributed.GalleryStore) {
|
|
m.mu.Lock()
|
|
defer m.mu.Unlock()
|
|
m.store = s
|
|
}
|
|
|
|
// Start hydrates the in-memory maps from PostgreSQL (if a store was wired)
|
|
// and subscribes to the broadcast subjects (if NATS was wired). It returns
|
|
// the first subscribe error; hydration errors are logged but non-fatal so
|
|
// the frontend still comes up.
|
|
//
|
|
// Safe to call exactly once after SetMessagingClient / SetGalleryStore. The
|
|
// ctx parameter is reserved for future cancellation — current subscriptions
|
|
// live for the lifetime of the OpCache and are released by Close.
|
|
func (m *OpCache) Start(_ context.Context) error {
|
|
m.mu.RLock()
|
|
store := m.store
|
|
nc := m.nats
|
|
m.mu.RUnlock()
|
|
|
|
if store != nil {
|
|
if err := m.hydrateFromStore(store); err != nil {
|
|
xlog.Warn("OpCache hydrate failed; starting empty", "error", err)
|
|
}
|
|
}
|
|
|
|
if nc == nil {
|
|
return nil
|
|
}
|
|
|
|
startSub, err := messaging.SubscribeJSON(nc, messaging.SubjectGalleryOpStart, func(evt OpCacheEvent) {
|
|
m.applyStart(evt)
|
|
})
|
|
if err != nil {
|
|
return err
|
|
}
|
|
endSub, err := messaging.SubscribeJSON(nc, messaging.SubjectGalleryOpEnd, func(evt OpCacheEvent) {
|
|
m.applyEnd(evt)
|
|
})
|
|
if err != nil {
|
|
if uerr := startSub.Unsubscribe(); uerr != nil {
|
|
xlog.Warn("failed to unsubscribe partial OpCache subscription", "error", uerr)
|
|
}
|
|
return err
|
|
}
|
|
|
|
m.mu.Lock()
|
|
m.subs = append(m.subs, startSub, endSub)
|
|
m.mu.Unlock()
|
|
return nil
|
|
}
|
|
|
|
// Close drops all NATS subscriptions. Safe to call multiple times.
|
|
func (m *OpCache) Close() {
|
|
m.mu.Lock()
|
|
subs := m.subs
|
|
m.subs = nil
|
|
m.mu.Unlock()
|
|
for _, s := range subs {
|
|
if err := s.Unsubscribe(); err != nil {
|
|
xlog.Warn("OpCache unsubscribe failed", "error", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (m *OpCache) hydrateFromStore(store *distributed.GalleryStore) error {
|
|
ops, err := store.ListActive()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
for _, op := range ops {
|
|
if op.CacheKey == "" {
|
|
continue
|
|
}
|
|
m.status.Set(op.CacheKey, op.ID)
|
|
if op.IsBackendOp {
|
|
m.backendOps.Set(op.CacheKey, true)
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// applyStart merges an inbound OpStart event into the local maps. Idempotent:
|
|
// receiving our own broadcast is a harmless re-assignment of the same value.
|
|
func (m *OpCache) applyStart(evt OpCacheEvent) {
|
|
if evt.CacheKey == "" || evt.JobID == "" {
|
|
return
|
|
}
|
|
// A retry admitted on a peer arrives as a start event reusing the key, which
|
|
// strands the stamp of whichever job the key pointed at here. Only the drop
|
|
// half of stampStart runs: a replicated op deliberately carries no stamp and
|
|
// reports as zero-length through recordTerminal's finish-time fallback.
|
|
m.dropReplacedStamp(evt.CacheKey, evt.JobID)
|
|
m.status.Set(evt.CacheKey, evt.JobID)
|
|
if evt.IsBackend {
|
|
m.backendOps.Set(evt.CacheKey, true)
|
|
}
|
|
}
|
|
|
|
// applyEnd removes any entries whose jobID matches the event. Idempotent.
|
|
func (m *OpCache) applyEnd(evt OpCacheEvent) {
|
|
if evt.JobID == "" {
|
|
return
|
|
}
|
|
// Record before the keys go. The history ring dedupes by job ID, so the
|
|
// originating replica recording locally and then receiving its own
|
|
// broadcast still produces one entry.
|
|
m.recordTerminal(evt.JobID, terminalPeer)
|
|
for _, k := range m.status.Keys() {
|
|
if m.status.Get(k) == evt.JobID {
|
|
m.status.Delete(k)
|
|
m.backendOps.Delete(k)
|
|
}
|
|
}
|
|
// recordTerminal only drops the stamp on the path where it found a cache
|
|
// key; an end event that overtakes the local Set finds none. The operation
|
|
// is over cluster-wide either way, so the stamp goes unconditionally.
|
|
m.started.Delete(evt.JobID)
|
|
}
|
|
|
|
func (m *OpCache) Set(key string, value string) {
|
|
m.stampStart(key, value)
|
|
m.status.Set(key, value)
|
|
m.persistAndBroadcastStart(key, value, false)
|
|
}
|
|
|
|
// SetBackend sets a key-value pair and marks it as a backend operation
|
|
func (m *OpCache) SetBackend(key string, value string) {
|
|
m.stampStart(key, value)
|
|
m.status.Set(key, value)
|
|
m.backendOps.Set(key, true)
|
|
m.persistAndBroadcastStart(key, value, true)
|
|
}
|
|
|
|
// stampStart records when jobID started, so the history record can report a
|
|
// duration.
|
|
func (m *OpCache) stampStart(key, jobID string) {
|
|
m.dropReplacedStamp(key, jobID)
|
|
m.started.Set(jobID, time.Now())
|
|
}
|
|
|
|
// dropReplacedStamp forgets the start time of the job a cache key is about to
|
|
// stop pointing at. Retrying an operation reuses the key with a fresh job ID,
|
|
// so without this the map would grow for the lifetime of the process: the
|
|
// callers in ui_api.go are not uniformly guarded by Exists. Must run before
|
|
// status.Set, which is what the previous job ID is read from.
|
|
func (m *OpCache) dropReplacedStamp(key, jobID string) {
|
|
if prev := m.status.Get(key); prev != "" && prev != jobID {
|
|
m.started.Delete(prev)
|
|
}
|
|
}
|
|
|
|
func (m *OpCache) persistAndBroadcastStart(key, value string, isBackend bool) {
|
|
m.mu.RLock()
|
|
store := m.store
|
|
nc := m.nats
|
|
m.mu.RUnlock()
|
|
|
|
if store != nil {
|
|
if err := store.UpsertCacheKey(value, key, isBackend); err != nil {
|
|
xlog.Warn("OpCache failed to persist cache key", "job_id", value, "error", err)
|
|
}
|
|
}
|
|
if nc != nil {
|
|
if err := nc.Publish(messaging.SubjectGalleryOpStart, OpCacheEvent{
|
|
JobID: value,
|
|
CacheKey: key,
|
|
IsBackend: isBackend,
|
|
}); err != nil {
|
|
xlog.Warn("OpCache failed to broadcast start", "job_id", value, "error", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
// IsBackendOp returns true if the given key is a backend operation
|
|
func (m *OpCache) IsBackendOp(key string) bool {
|
|
return m.backendOps.Get(key)
|
|
}
|
|
|
|
func (m *OpCache) Get(key string) string {
|
|
return m.status.Get(key)
|
|
}
|
|
|
|
// terminalSource is the path that retired an operation. It exists for one
|
|
// decision: what a missing gallery status means. Locally it means the op was
|
|
// queued and removed before anything ran; on the peer path it means this
|
|
// replica never held the outcome, which is a different thing wearing the same
|
|
// signal.
|
|
type terminalSource int
|
|
|
|
const (
|
|
terminalLocal terminalSource = iota
|
|
terminalPeer
|
|
)
|
|
|
|
// recordTerminal appends a finished operation to the history ring. It must run
|
|
// BEFORE the cache entry is deleted: the gallery key is the only place the
|
|
// display name, the backend flag and the node scoping live.
|
|
//
|
|
// Safe to call for an unknown job ID (no key, no record) and safe to call
|
|
// twice for the same job (the ring dedupes), which is what makes it usable
|
|
// from both the local delete path and the NATS end event.
|
|
func (m *OpCache) recordTerminal(jobID string, src terminalSource) {
|
|
if jobID == "" {
|
|
return
|
|
}
|
|
|
|
key := ""
|
|
for _, k := range m.status.Keys() {
|
|
if m.status.Get(k) == jobID {
|
|
key = k
|
|
break
|
|
}
|
|
}
|
|
if key == "" {
|
|
return
|
|
}
|
|
|
|
// A replica that restarted mid-operation hydrates its cache keys from
|
|
// PostgreSQL, but gallery statuses are in-memory only and come back empty.
|
|
// The end broadcast then arrives with nothing to read the outcome from, and
|
|
// guessing would file a successful install as cancelled. Record nothing:
|
|
// absence beats wrong data in a "what just happened" view, and it is exactly
|
|
// what this replica produced before the end event started recording.
|
|
status := m.galleryService.GetStatus(jobID)
|
|
if status == nil && src == terminalPeer {
|
|
return
|
|
}
|
|
|
|
rec := OpRecord{
|
|
ID: key,
|
|
JobID: jobID,
|
|
IsBackend: m.backendOps.Get(key),
|
|
TaskType: "installation",
|
|
Outcome: OutcomeCompleted,
|
|
FinishedAt: time.Now(),
|
|
}
|
|
// hydrateFromStore and applyStart populate status without a stamp, so an op
|
|
// recovered from PostgreSQL or replicated from a peer has no start time.
|
|
// Reporting the finish time makes such a record a zero-length operation
|
|
// rather than one that appears to have run since year one.
|
|
//
|
|
// Read the stamp once and reject a zero value rather than testing Exists and
|
|
// then reading: a concurrent recordTerminal for the same job can delete the
|
|
// stamp between the two, and the read that follows returns the zero time,
|
|
// which would overwrite the fallback with year one.
|
|
rec.StartedAt = rec.FinishedAt
|
|
if started := m.started.Get(jobID); !started.IsZero() {
|
|
rec.StartedAt = started
|
|
}
|
|
|
|
rec.Name, rec.NodeID = operationDisplayName(key)
|
|
|
|
// Outcome order matters: an error outweighs everything else, because an op
|
|
// can carry both an error and an unfinished status and the failure is what
|
|
// the user needs to see.
|
|
//
|
|
// An operation that left the cache while it was still unprocessed never got
|
|
// to finish, so it is not a success. That is what the dismiss endpoint
|
|
// produces when it fires on an op that is still in flight, and what a status
|
|
// that exists but never started looks like. The cancel endpoint is already
|
|
// covered by status.Cancelled, which CancelOperation sets synchronously
|
|
// before the handler removes the entry.
|
|
if status != nil {
|
|
if status.Deletion {
|
|
rec.TaskType = "deletion"
|
|
}
|
|
switch {
|
|
case status.Error != nil:
|
|
rec.Outcome = OutcomeFailed
|
|
rec.Error = status.Error.Error()
|
|
case status.Cancelled || !status.Processed:
|
|
rec.Outcome = OutcomeCancelled
|
|
}
|
|
} else {
|
|
// Queued but never started, then removed.
|
|
rec.Outcome = OutcomeCancelled
|
|
}
|
|
|
|
m.history.add(rec)
|
|
m.started.Delete(jobID)
|
|
}
|
|
|
|
// History returns finished operations, newest first.
|
|
//
|
|
// With a gallery store wired the record comes from PostgreSQL, so every replica
|
|
// answers with the same history: the in-memory ring only holds what this
|
|
// replica happened to serve, which makes the page's contents depend on which
|
|
// replica the poll was routed to and leaves a replica added by a scale-out or a
|
|
// rolling deploy blank forever.
|
|
func (m *OpCache) History() []OpRecord {
|
|
m.mu.RLock()
|
|
store := m.store
|
|
m.mu.RUnlock()
|
|
|
|
if store == nil {
|
|
return m.history.list()
|
|
}
|
|
|
|
ops, err := store.ListTerminal(DefaultHistorySize)
|
|
if err != nil {
|
|
// A transient database failure must not blank the page. The ring holds
|
|
// a subset of the same record, so it is a strictly better answer than
|
|
// nothing.
|
|
xlog.Warn("OpCache failed to read the operation record; falling back to the local ring", "error", err)
|
|
return m.history.list()
|
|
}
|
|
|
|
records := make([]OpRecord, 0, len(ops))
|
|
for _, op := range ops {
|
|
records = append(records, recordFromStore(op))
|
|
}
|
|
return records
|
|
}
|
|
|
|
// ClearHistory empties the record. Live operations are untouched.
|
|
//
|
|
// With a gallery store wired this clears the record for the whole cluster,
|
|
// which is the only way "Clear history" can mean anything: clearing one
|
|
// replica's ring leaves the record to reappear on the next poll routed
|
|
// elsewhere.
|
|
//
|
|
// The error is returned rather than logged because the caller is an HTTP
|
|
// handler and the rows are what the next read returns: reporting success on a
|
|
// failed delete makes the record vanish from the page and come straight back on
|
|
// the next fetch, with nothing said about why.
|
|
func (m *OpCache) ClearHistory() error {
|
|
m.mu.RLock()
|
|
store := m.store
|
|
m.mu.RUnlock()
|
|
|
|
if store == nil {
|
|
m.history.clear()
|
|
return nil
|
|
}
|
|
|
|
// Store first: the ring is only ever a fallback for a failed read, so
|
|
// emptying it before the rows are confirmed gone would report an empty
|
|
// record while the real one is still there.
|
|
if err := store.ClearTerminal(); err != nil {
|
|
return fmt.Errorf("clearing the persisted operation record: %w", err)
|
|
}
|
|
m.history.clear()
|
|
return nil
|
|
}
|
|
|
|
func (m *OpCache) DeleteUUID(uuid string) {
|
|
// Before the keys go: they carry the name and the backend flag.
|
|
m.recordTerminal(uuid, terminalLocal)
|
|
|
|
deleted := false
|
|
for _, k := range m.status.Keys() {
|
|
if m.status.Get(k) == uuid {
|
|
m.status.Delete(k)
|
|
m.backendOps.Delete(k) // Also clean up the backend flag
|
|
deleted = true
|
|
}
|
|
}
|
|
if !deleted {
|
|
return
|
|
}
|
|
m.mu.RLock()
|
|
nc := m.nats
|
|
m.mu.RUnlock()
|
|
if nc != nil {
|
|
if err := nc.Publish(messaging.SubjectGalleryOpEnd, OpCacheEvent{JobID: uuid}); err != nil {
|
|
xlog.Warn("OpCache failed to broadcast end", "job_id", uuid, "error", err)
|
|
}
|
|
}
|
|
}
|
|
|
|
func (m *OpCache) Map() map[string]string {
|
|
return m.status.Map()
|
|
}
|
|
|
|
func (m *OpCache) Exists(key string) bool {
|
|
return m.status.Exists(key)
|
|
}
|
|
|
|
func (m *OpCache) GetStatus() (map[string]string, map[string]string) {
|
|
taskTypes := map[string]string{}
|
|
processingModelsData := map[string]string{}
|
|
|
|
// Iterate a snapshot (Keys() copies) and build a fresh result map. We must
|
|
// NOT delete from m.Map() during the range: Map() returns the live internal
|
|
// map by reference, so a bare delete here would be an unsynchronized write
|
|
// to a map four HTTP handlers read every ~1s — a concurrent-map-write crash.
|
|
// Collect evictions and apply them via the locked DeleteUUID after the loop.
|
|
var evict []string
|
|
for _, k := range m.status.Keys() {
|
|
v := m.status.Get(k)
|
|
if v == "" {
|
|
continue // raced with a concurrent Delete
|
|
}
|
|
status := m.galleryService.GetStatus(v)
|
|
// Terminal ops must not keep showing as "processing". Cleanup was
|
|
// previously only triggered by a client polling /api/backends/job/:uid,
|
|
// but the Manage-page Reinstall/Upgrade buttons never poll, so completed
|
|
// ops leaked into processingBackends forever and the card spun
|
|
// "reinstalling" indefinitely. Evict here on the list read (the UI always
|
|
// calls this). DeleteUUID broadcasts the eviction so peer replicas converge.
|
|
//
|
|
// We evict ONLY a clean success (progress 100 + "completed", matching the
|
|
// job-poll's historical delete condition) or a cancellation. Deliberately
|
|
// NOT evicted:
|
|
// - failed ops (Error != nil): kept so /api/operations can surface the
|
|
// error and offer Dismiss.
|
|
// - the ErrWorkerStillInstalling soft-path (Processed=true, Error=nil,
|
|
// progress != 100): the worker is still installing in the background
|
|
// and the reconciler confirms the real outcome later — evicting it
|
|
// would hide an install that may still fail.
|
|
if status != nil && status.Processed &&
|
|
((status.Progress == 100 && status.Message == "completed") || status.Cancelled) {
|
|
evict = append(evict, v)
|
|
continue
|
|
}
|
|
processingModelsData[k] = v
|
|
taskTypes[k] = "Installation"
|
|
if status != nil && status.Deletion {
|
|
taskTypes[k] = "Deletion"
|
|
} else if status == nil {
|
|
taskTypes[k] = "Waiting"
|
|
}
|
|
}
|
|
|
|
for _, v := range evict {
|
|
m.DeleteUUID(v)
|
|
}
|
|
|
|
return processingModelsData, taskTypes
|
|
}
|
|
|
|
// operationDisplayName reduces an operation key to what the record shows: the
|
|
// node prefix ("node:<nodeID>:") detached into a node ID the page reports
|
|
// separately, and the gallery prefix ("<gallery>@") dropped.
|
|
//
|
|
// The in-memory ring and the store-backed record both go through this, so the
|
|
// same operation cannot end up named two different ways depending on which
|
|
// source answered.
|
|
func operationDisplayName(key string) (name, nodeID string) {
|
|
name = key
|
|
if id, backend, ok := ParseNodeScopedKey(key); ok {
|
|
nodeID = id
|
|
name = backend
|
|
}
|
|
if _, after, found := strings.Cut(name, "@"); found {
|
|
name = after
|
|
}
|
|
return name, nodeID
|
|
}
|
|
|
|
// NodeScopedKeyPrefix is the opcache key prefix used by InstallBackendOnNodeEndpoint
|
|
// so per-node installs do not collide on the bare backend name. Format:
|
|
// "node:<nodeID>:<backend>". Read by /api/operations to extract nodeID for the UI.
|
|
const NodeScopedKeyPrefix = "node:"
|
|
|
|
// NodeScopedKey returns the opcache key for a node-scoped backend operation.
|
|
// The prefix lets ParseNodeScopedKey detach the nodeID back out so the
|
|
// operations endpoint can surface it without storing nodeID separately.
|
|
func NodeScopedKey(nodeID, backend string) string {
|
|
return NodeScopedKeyPrefix + nodeID + ":" + backend
|
|
}
|
|
|
|
// ParseNodeScopedKey extracts (nodeID, backend) from a key built by NodeScopedKey.
|
|
// Returns ok=false for keys that lack the prefix or are missing the nodeID or
|
|
// backend segment. Backend names containing colons are preserved because we
|
|
// split on the first colon after the prefix only.
|
|
func ParseNodeScopedKey(key string) (nodeID, backend string, ok bool) {
|
|
rest, hasPrefix := strings.CutPrefix(key, NodeScopedKeyPrefix)
|
|
if !hasPrefix {
|
|
return "", "", false
|
|
}
|
|
nodeID, backend, ok = strings.Cut(rest, ":")
|
|
if !ok || nodeID == "" || backend == "" {
|
|
return "", "", false
|
|
}
|
|
return nodeID, backend, true
|
|
}
|