fix(distributed): make per-node backend upgrade actually upgrade (#10838)

* test(core/http): make the suite's HTTP port overridable

app_test.go and openresponses_test.go hardcoded 127.0.0.1:9090. When
another service already listens on 9090 the suite does not fail fast:
the server goroutine logs the bind error and the specs then poll
whatever is squatting the port until Eventually times out. On machines
where 9090 is permanently taken this makes the pre-commit coverage gate
impossible to pass.

Introduce testHTTPAddr, defaulting to 127.0.0.1:9090 (what CI has
always used) and overridable via LOCALAI_TEST_HTTP_PORT for local runs.

Assisted-by: Claude:claude-fable-5 golangci-lint
Signed-off-by: Ettore Di Giacinto <mudler@localai.io>

* fix(distributed): make per-node backend upgrade actually upgrade

The node detail page's Upgrade button reused the node-scoped install
path (POST /api/nodes/:id/backends/install). That fires NATS
backend.install with force=false, and the worker's install handler is
deliberately "ensure installed": when the backend binary already exists
on disk it short-circuits without touching the gallery. Since only an
installed backend can be upgraded, the whole chain was a guaranteed
successful no-op - the UI then toasted "backend upgraded" without even
waiting for the async job.

Route upgrades through the real force-reinstall path instead:

- BackendManager.UpgradeBackend now receives the ManagementOp (like
  InstallBackend already did) so implementations can honor
  op.TargetNodeID.
- DistributedBackendManager.UpgradeBackend scopes the backend.upgrade
  fan-out to op.TargetNodeID when set, and errors when the target node
  does not report the backend as installed.
- New POST /api/nodes/:id/backends/upgrade endpoint enqueues an
  Upgrade=true node-scoped op (async 202 + jobID, mirroring install).
- NodeDetail UI calls the new endpoint and reports the dispatch
  ("Upgrading ... on this node...") instead of claiming success; the
  Operations panel tracks the actual job.

Verified against a live local cluster (NATS + Postgres + two workers):
the target worker stops the running process, force-reinstalls from the
gallery and re-downloads the OCI image; the second worker receives no
backend.upgrade event; upgrading a backend missing from the target node
fails the job with a clear error.

Assisted-by: Claude:claude-fable-5 golangci-lint
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>
This commit is contained in:
LocalAI [bot]
2026-07-15 09:16:55 +02:00
committed by GitHub
parent 40763d1181
commit 3601174ce0
16 changed files with 316 additions and 88 deletions

View File

@@ -215,8 +215,11 @@ func (uc *UpgradeChecker) runCheck(ctx context.Context) {
var err error
if bm != nil {
// Background auto-upgrade: no live admin watching a progress bar,
// so opID is empty and the distributed path skips progress streaming.
err = bm.UpgradeBackend(ctx, "", name, nil)
// so op.ID is empty and the distributed path skips progress streaming.
err = bm.UpgradeBackend(ctx, &galleryop.ManagementOp[gallery.GalleryBackend, any]{
GalleryElementName: name,
Upgrade: true,
}, nil)
} else {
err = gallery.UpgradeBackend(ctx, uc.systemState, uc.modelLoader,
uc.galleries, name, nil, uc.appConfig.RequireBackendIntegrity)

View File

@@ -383,13 +383,13 @@ var _ = Describe("API test", func() {
Expect(err).ToNot(HaveOccurred())
go func() {
if err := app.Start("127.0.0.1:9090"); err != nil && err != http.ErrServerClosed {
if err := app.Start(testHTTPAddr); err != nil && err != http.ErrServerClosed {
xlog.Error("server error", "error", err)
}
}()
defaultConfig := openai.DefaultConfig(apiKey)
defaultConfig.BaseURL = "http://127.0.0.1:9090/v1"
defaultConfig.BaseURL = "http://" + testHTTPAddr + "/v1"
client2 = openaigo.NewClient("")
client2.BaseURL = defaultConfig.BaseURL
@@ -418,7 +418,7 @@ var _ = Describe("API test", func() {
Context("Auth Tests", func() {
It("Should fail if the api key is missing", func() {
err, sc := postInvalidRequest("http://127.0.0.1:9090/models/available")
err, sc := postInvalidRequest("http://" + testHTTPAddr + "/models/available")
Expect(err).ToNot(BeNil())
Expect(sc).To(Equal(401))
})
@@ -427,7 +427,7 @@ var _ = Describe("API test", func() {
Context("URL routing Tests", func() {
It("Should support reverse-proxy when unauthenticated", func() {
err, sc, body := getRequest("http://127.0.0.1:9090/myprefix/", http.Header{
err, sc, body := getRequest("http://"+testHTTPAddr+"/myprefix/", http.Header{
"X-Forwarded-Proto": {"https"},
"X-Forwarded-Host": {"example.org"},
"X-Forwarded-Prefix": {"/myprefix/"},
@@ -441,7 +441,7 @@ var _ = Describe("API test", func() {
It("Should support reverse-proxy when authenticated", func() {
err, sc, body := getRequest("http://127.0.0.1:9090/myprefix/", http.Header{
err, sc, body := getRequest("http://"+testHTTPAddr+"/myprefix/", http.Header{
"Authorization": {bearerKey},
"X-Forwarded-Proto": {"https"},
"X-Forwarded-Host": {"example.org"},
@@ -459,7 +459,7 @@ var _ = Describe("API test", func() {
// requests them through the proxy.
It("Should support reverse-proxy when prefix is stripped by the proxy", func() {
err, sc, body := getRequest("http://127.0.0.1:9090/app", http.Header{
err, sc, body := getRequest("http://"+testHTTPAddr+"/app", http.Header{
"X-Forwarded-Proto": {"https"},
"X-Forwarded-Host": {"example.org"},
"X-Forwarded-Prefix": {"/myprefix"},
@@ -477,7 +477,7 @@ var _ = Describe("API test", func() {
// from a foreign origin. BasePathPrefix must reject these via
// SafeForwardedPrefix and fall back to "/".
It("Should ignore an unsafe X-Forwarded-Prefix and not poison asset URLs", func() {
err, sc, body := getRequest("http://127.0.0.1:9090/app", http.Header{
err, sc, body := getRequest("http://"+testHTTPAddr+"/app", http.Header{
"X-Forwarded-Proto": {"https"},
"X-Forwarded-Host": {"example.org"},
"X-Forwarded-Prefix": {"//evil.com"},
@@ -492,13 +492,13 @@ var _ = Describe("API test", func() {
Context("Applying models", func() {
It("applies models from a gallery", func() {
models, err := getModels("http://127.0.0.1:9090/models/available")
models, err := getModels("http://" + testHTTPAddr + "/models/available")
Expect(err).To(BeNil())
Expect(len(models)).To(Equal(2), fmt.Sprint(models))
Expect(models[0].Installed).To(BeFalse(), fmt.Sprint(models))
Expect(models[1].Installed).To(BeFalse(), fmt.Sprint(models))
response := postModelApplyRequest("http://127.0.0.1:9090/models/apply", modelApplyRequest{
response := postModelApplyRequest("http://"+testHTTPAddr+"/models/apply", modelApplyRequest{
ID: "test@bert2",
})
@@ -507,7 +507,7 @@ var _ = Describe("API test", func() {
uuid := response["uuid"].(string)
resp := map[string]any{}
Eventually(func() bool {
response := getModelStatus("http://127.0.0.1:9090/models/jobs/" + uuid)
response := getModelStatus("http://" + testHTTPAddr + "/models/jobs/" + uuid)
fmt.Println(response)
resp = response
return response["processed"].(bool)
@@ -526,7 +526,7 @@ var _ = Describe("API test", func() {
Expect(content["usage"]).To(ContainSubstring("You can test this model with curl like this"))
Expect(content["foo"]).To(Equal("bar"))
models, err = getModels("http://127.0.0.1:9090/models/available")
models, err = getModels("http://" + testHTTPAddr + "/models/available")
Expect(err).To(BeNil())
Expect(len(models)).To(Equal(2), fmt.Sprint(models))
Expect(models[0].Name).To(Or(Equal("bert"), Equal("bert2")))
@@ -541,7 +541,7 @@ var _ = Describe("API test", func() {
})
It("overrides models", func() {
response := postModelApplyRequest("http://127.0.0.1:9090/models/apply", modelApplyRequest{
response := postModelApplyRequest("http://"+testHTTPAddr+"/models/apply", modelApplyRequest{
URL: bertEmbeddingsURL,
Name: "bert",
Overrides: map[string]any{
@@ -554,7 +554,7 @@ var _ = Describe("API test", func() {
uuid := response["uuid"].(string)
Eventually(func() bool {
response := getModelStatus("http://127.0.0.1:9090/models/jobs/" + uuid)
response := getModelStatus("http://" + testHTTPAddr + "/models/jobs/" + uuid)
return response["processed"].(bool)
}, "360s", "10s").Should(Equal(true))
@@ -567,7 +567,7 @@ var _ = Describe("API test", func() {
Expect(content["backend"]).To(Equal("llama"))
})
It("apply models without overrides", func() {
response := postModelApplyRequest("http://127.0.0.1:9090/models/apply", modelApplyRequest{
response := postModelApplyRequest("http://"+testHTTPAddr+"/models/apply", modelApplyRequest{
URL: bertEmbeddingsURL,
Name: "bert",
Overrides: map[string]any{},
@@ -578,7 +578,7 @@ var _ = Describe("API test", func() {
uuid := response["uuid"].(string)
Eventually(func() bool {
response := getModelStatus("http://127.0.0.1:9090/models/jobs/" + uuid)
response := getModelStatus("http://" + testHTTPAddr + "/models/jobs/" + uuid)
return response["processed"].(bool)
}, "360s", "10s").Should(Equal(true))
@@ -622,14 +622,14 @@ parameters:
}
var response schema.GalleryResponse
err := postRequestResponseJSON("http://127.0.0.1:9090/models/import-uri", &importReq, &response)
err := postRequestResponseJSON("http://"+testHTTPAddr+"/models/import-uri", &importReq, &response)
Expect(err).ToNot(HaveOccurred())
Expect(response.ID).ToNot(BeEmpty())
uuid := response.ID
resp := map[string]any{}
Eventually(func() bool {
response := getModelStatus("http://127.0.0.1:9090/models/jobs/" + uuid)
response := getModelStatus("http://" + testHTTPAddr + "/models/jobs/" + uuid)
resp = response
return response["processed"].(bool)
}, "360s", "10s").Should(Equal(true))
@@ -657,7 +657,7 @@ parameters:
}
var response schema.GalleryResponse
err := postRequestResponseJSON("http://127.0.0.1:9090/models/import-uri", &importReq, &response)
err := postRequestResponseJSON("http://"+testHTTPAddr+"/models/import-uri", &importReq, &response)
// The endpoint should return an error immediately
Expect(err).To(HaveOccurred())
Expect(err.Error()).To(ContainSubstring("failed to discover model config"))
@@ -693,14 +693,14 @@ parameters:
}
var response schema.GalleryResponse
err := postRequestResponseJSON("http://127.0.0.1:9090/models/import-uri", &importReq, &response)
err := postRequestResponseJSON("http://"+testHTTPAddr+"/models/import-uri", &importReq, &response)
Expect(err).ToNot(HaveOccurred())
Expect(response.ID).ToNot(BeEmpty())
uuid := response.ID
resp := map[string]any{}
Eventually(func() bool {
response := getModelStatus("http://127.0.0.1:9090/models/jobs/" + uuid)
response := getModelStatus("http://" + testHTTPAddr + "/models/jobs/" + uuid)
resp = response
return response["processed"].(bool)
}, "360s", "10s").Should(Equal(true))
@@ -763,13 +763,13 @@ chat_template_kwargs:
app, err = API(localAIApp)
Expect(err).ToNot(HaveOccurred())
go func() {
if err := app.Start("127.0.0.1:9090"); err != nil && err != http.ErrServerClosed {
if err := app.Start(testHTTPAddr); err != nil && err != http.ErrServerClosed {
xlog.Error("server error", "error", err)
}
}()
defaultConfig := openai.DefaultConfig("")
defaultConfig.BaseURL = "http://127.0.0.1:9090/v1"
defaultConfig.BaseURL = "http://" + testHTTPAddr + "/v1"
client2 = openaigo.NewClient("")
client2.BaseURL = defaultConfig.BaseURL
@@ -813,7 +813,7 @@ chat_template_kwargs:
// Mock-backend is registered via SetExternalBackend so it appears
// alongside any built-in entries; verifying that string proves the
// endpoint is wired up regardless of which real backends exist.
resp, err := http.Get("http://127.0.0.1:9090/system")
resp, err := http.Get("http://" + testHTTPAddr + "/system")
Expect(err).ToNot(HaveOccurred())
Expect(resp.StatusCode).To(Equal(200))
dat, err := io.ReadAll(resp.Body)
@@ -851,7 +851,7 @@ chat_template_kwargs:
} `json:"message"`
} `json:"choices"`
}
err := postRequestResponseJSON("http://127.0.0.1:9090/v1/chat/completions", &reqBody, &chatResp)
err := postRequestResponseJSON("http://"+testHTTPAddr+"/v1/chat/completions", &reqBody, &chatResp)
Expect(err).ToNot(HaveOccurred())
Expect(chatResp.Choices).ToNot(BeEmpty())
@@ -889,14 +889,14 @@ chat_template_kwargs:
}
var createResp map[string]any
err := postRequestResponseJSON("http://127.0.0.1:9090/api/agent/tasks", &taskBody, &createResp)
err := postRequestResponseJSON("http://"+testHTTPAddr+"/api/agent/tasks", &taskBody, &createResp)
Expect(err).ToNot(HaveOccurred())
Expect(createResp["id"]).ToNot(BeEmpty())
taskID := createResp["id"].(string)
// Get the task
var task schema.Task
resp, err := http.Get("http://127.0.0.1:9090/api/agent/tasks/" + taskID)
resp, err := http.Get("http://" + testHTTPAddr + "/api/agent/tasks/" + taskID)
Expect(err).ToNot(HaveOccurred())
Expect(resp.StatusCode).To(Equal(200))
body, _ := io.ReadAll(resp.Body)
@@ -904,7 +904,7 @@ chat_template_kwargs:
Expect(task.Name).To(Equal("Test Task"))
// List tasks
resp, err = http.Get("http://127.0.0.1:9090/api/agent/tasks")
resp, err = http.Get("http://" + testHTTPAddr + "/api/agent/tasks")
Expect(err).ToNot(HaveOccurred())
Expect(resp.StatusCode).To(Equal(200))
var tasks []schema.Task
@@ -914,18 +914,18 @@ chat_template_kwargs:
// Update task
taskBody["name"] = "Updated Task"
err = putRequestJSON("http://127.0.0.1:9090/api/agent/tasks/"+taskID, &taskBody)
err = putRequestJSON("http://"+testHTTPAddr+"/api/agent/tasks/"+taskID, &taskBody)
Expect(err).ToNot(HaveOccurred())
// Verify update
resp, err = http.Get("http://127.0.0.1:9090/api/agent/tasks/" + taskID)
resp, err = http.Get("http://" + testHTTPAddr + "/api/agent/tasks/" + taskID)
Expect(err).ToNot(HaveOccurred())
body, _ = io.ReadAll(resp.Body)
json.Unmarshal(body, &task)
Expect(task.Name).To(Equal("Updated Task"))
// Delete task
req, _ := http.NewRequest("DELETE", "http://127.0.0.1:9090/api/agent/tasks/"+taskID, nil)
req, _ := http.NewRequest("DELETE", "http://"+testHTTPAddr+"/api/agent/tasks/"+taskID, nil)
req.Header.Set("Authorization", bearerKey)
resp, err = http.DefaultClient.Do(req)
Expect(err).ToNot(HaveOccurred())
@@ -942,7 +942,7 @@ chat_template_kwargs:
}
var createResp map[string]any
err := postRequestResponseJSON("http://127.0.0.1:9090/api/agent/tasks", &taskBody, &createResp)
err := postRequestResponseJSON("http://"+testHTTPAddr+"/api/agent/tasks", &taskBody, &createResp)
Expect(err).ToNot(HaveOccurred())
taskID := createResp["id"].(string)
@@ -953,14 +953,14 @@ chat_template_kwargs:
}
var jobResp schema.JobExecutionResponse
err = postRequestResponseJSON("http://127.0.0.1:9090/api/agent/jobs/execute", &jobBody, &jobResp)
err = postRequestResponseJSON("http://"+testHTTPAddr+"/api/agent/jobs/execute", &jobBody, &jobResp)
Expect(err).ToNot(HaveOccurred())
Expect(jobResp.JobID).ToNot(BeEmpty())
jobID := jobResp.JobID
// Get job status
var job schema.Job
resp, err := http.Get("http://127.0.0.1:9090/api/agent/jobs/" + jobID)
resp, err := http.Get("http://" + testHTTPAddr + "/api/agent/jobs/" + jobID)
Expect(err).ToNot(HaveOccurred())
Expect(resp.StatusCode).To(Equal(200))
body, _ := io.ReadAll(resp.Body)
@@ -969,7 +969,7 @@ chat_template_kwargs:
Expect(job.TaskID).To(Equal(taskID))
// List jobs
resp, err = http.Get("http://127.0.0.1:9090/api/agent/jobs")
resp, err = http.Get("http://" + testHTTPAddr + "/api/agent/jobs")
Expect(err).ToNot(HaveOccurred())
Expect(resp.StatusCode).To(Equal(200))
var jobs []schema.Job
@@ -979,13 +979,13 @@ chat_template_kwargs:
// Cancel job (if still pending/running)
if job.Status == schema.JobStatusPending || job.Status == schema.JobStatusRunning {
req, _ := http.NewRequest("POST", "http://127.0.0.1:9090/api/agent/jobs/"+jobID+"/cancel", nil)
req, _ := http.NewRequest("POST", "http://"+testHTTPAddr+"/api/agent/jobs/"+jobID+"/cancel", nil)
req.Header.Set("Authorization", bearerKey)
resp, err = http.DefaultClient.Do(req)
Expect(err).ToNot(HaveOccurred())
if resp.StatusCode == http.StatusBadRequest {
// The worker can finish between the status read and cancellation request.
resp, err = http.Get("http://127.0.0.1:9090/api/agent/jobs/" + jobID)
resp, err = http.Get("http://" + testHTTPAddr + "/api/agent/jobs/" + jobID)
Expect(err).ToNot(HaveOccurred())
Expect(resp.StatusCode).To(Equal(http.StatusOK))
body, _ = io.ReadAll(resp.Body)
@@ -1012,13 +1012,13 @@ chat_template_kwargs:
}
var createResp map[string]any
err := postRequestResponseJSON("http://127.0.0.1:9090/api/agent/tasks", &taskBody, &createResp)
err := postRequestResponseJSON("http://"+testHTTPAddr+"/api/agent/tasks", &taskBody, &createResp)
Expect(err).ToNot(HaveOccurred())
// Execute by name
paramsBody := map[string]string{"param1": "value1"}
var jobResp schema.JobExecutionResponse
err = postRequestResponseJSON("http://127.0.0.1:9090/api/agent/tasks/Named Task/execute", &paramsBody, &jobResp)
err = postRequestResponseJSON("http://"+testHTTPAddr+"/api/agent/tasks/Named Task/execute", &paramsBody, &jobResp)
Expect(err).ToNot(HaveOccurred())
Expect(jobResp.JobID).ToNot(BeEmpty())
})
@@ -1078,13 +1078,13 @@ chat_template_kwargs:
Expect(err).ToNot(HaveOccurred())
go func() {
if err := app.Start("127.0.0.1:9090"); err != nil && err != http.ErrServerClosed {
if err := app.Start(testHTTPAddr); err != nil && err != http.ErrServerClosed {
xlog.Error("server error", "error", err)
}
}()
defaultConfig := openai.DefaultConfig("")
defaultConfig.BaseURL = "http://127.0.0.1:9090/v1"
defaultConfig.BaseURL = "http://" + testHTTPAddr + "/v1"
client2 = openaigo.NewClient("")
client2.BaseURL = defaultConfig.BaseURL
// Wait for API to be ready

View File

@@ -524,6 +524,71 @@ func InstallBackendOnNodeEndpoint(_ nodes.NodeCommandSender, galleryService *gal
}
}
// UpgradeBackendOnNodeEndpoint triggers a backend upgrade (force-reinstall)
// on a single worker node. Async like InstallBackendOnNodeEndpoint: enqueues
// a ManagementOp with Upgrade=true and TargetNodeID set, returns 202 + jobID.
// The gallery service routes Upgrade ops to
// DistributedBackendManager.UpgradeBackend, which fires the NATS
// backend.upgrade subject: the worker stops running processes for the
// backend and reinstalls from the gallery even when the artifact already
// exists on disk. Reusing the install path here would no-op: the worker's
// backend.install handler is "ensure installed" and short-circuits on an
// already-present binary (the "backend upgraded but nothing happens" bug).
//
// Only gallery-name upgrades are supported: the distributed upgrade path
// resolves galleries from server config, so unlike install there is no
// URI/name/alias or galleries-override surface.
func UpgradeBackendOnNodeEndpoint(galleryService *galleryop.GalleryService, opcache *galleryop.OpCache, appConfig *config.ApplicationConfig) echo.HandlerFunc {
return func(c echo.Context) error {
if galleryService == nil {
return c.JSON(http.StatusServiceUnavailable, nodeError(http.StatusServiceUnavailable, "gallery service not configured"))
}
nodeID := c.Param("id")
var req struct {
Backend string `json:"backend"`
}
if err := c.Bind(&req); err != nil {
return c.JSON(http.StatusBadRequest, nodeError(http.StatusBadRequest, "invalid request body"))
}
if req.Backend == "" {
return c.JSON(http.StatusBadRequest, nodeError(http.StatusBadRequest, "backend name required"))
}
jobUUID, err := uuid.NewUUID()
if err != nil {
return c.JSON(http.StatusInternalServerError, nodeError(http.StatusInternalServerError, "failed to generate job id"))
}
jobID := jobUUID.String()
// Node-scoped cache key so a concurrent upgrade of the same backend on
// another node doesn't stomp this job in opcache.
cacheKey := galleryop.NodeScopedKey(nodeID, req.Backend)
opcache.SetBackend(cacheKey, jobID)
ctx, cancelFunc := context.WithCancel(context.Background())
op := galleryop.ManagementOp[gallery.GalleryBackend, any]{
ID: jobID,
GalleryElementName: req.Backend,
Galleries: appConfig.BackendGalleries,
TargetNodeID: nodeID,
Upgrade: true,
Context: ctx,
CancelFunc: cancelFunc,
}
galleryService.StoreCancellation(jobID, cancelFunc)
go func() {
galleryService.BackendGalleryChannel <- op
}()
xlog.Info("Node-scoped backend upgrade dispatched", "node", nodeID, "backend", req.Backend, "jobID", jobID)
return c.JSON(http.StatusAccepted, map[string]string{
"jobID": jobID,
"statusUrl": "/api/backends/job/" + jobID,
"message": "backend upgrade started",
})
}
}
// DeleteBackendOnNodeEndpoint deletes a backend from a worker node via NATS.
func DeleteBackendOnNodeEndpoint(unloader nodes.NodeCommandSender) echo.HandlerFunc {
return func(c echo.Context) error {

View File

@@ -120,4 +120,54 @@ var _ = Describe("InstallBackendOnNodeEndpoint async behavior", func() {
Expect(opcache.Exists(galleryop.NodeScopedKey("node-xyz", "custom"))).To(BeTrue())
})
// The node detail page's per-node "Upgrade" button used to reuse the
// install path, which the worker treats as "ensure installed" and
// short-circuits when the backend already exists on disk - the original
// "backend upgraded but nothing happens" bug. Upgrades must dispatch an
// op with Upgrade=true so the gallery service routes it to the
// force-reinstall backend.upgrade path, scoped to the one node.
Describe("UpgradeBackendOnNodeEndpoint", func() {
It("returns 202 with a jobID and dispatches an Upgrade op scoped to the node", func() {
body := `{"backend": "llama-cpp"}`
req := httptest.NewRequest(http.MethodPost, "/api/nodes/node-xyz/backends/upgrade", bytes.NewBufferString(body))
req.Header.Set("Content-Type", "application/json")
rec := httptest.NewRecorder()
c := e.NewContext(req, rec)
c.SetParamNames("id")
c.SetParamValues("node-xyz")
handler := localai.UpgradeBackendOnNodeEndpoint(galleryService, opcache, appCfg)
Expect(handler(c)).To(Succeed())
Expect(rec.Code).To(Equal(http.StatusAccepted))
var resp map[string]any
Expect(json.Unmarshal(rec.Body.Bytes(), &resp)).To(Succeed())
Expect(resp["jobID"]).To(BeAssignableToTypeOf(""))
Expect(resp["jobID"].(string)).ToNot(BeEmpty())
Expect(resp["statusUrl"]).To(Equal("/api/backends/job/" + resp["jobID"].(string)))
var op galleryop.ManagementOp[gallery.GalleryBackend, any]
Eventually(dispatched, "2s").Should(Receive(&op))
Expect(op.Upgrade).To(BeTrue(), "op must take the force-reinstall upgrade path")
Expect(op.TargetNodeID).To(Equal("node-xyz"))
Expect(op.GalleryElementName).To(Equal("llama-cpp"))
Expect(opcache.Exists(galleryop.NodeScopedKey("node-xyz", "llama-cpp"))).To(BeTrue())
Expect(opcache.IsBackendOp(galleryop.NodeScopedKey("node-xyz", "llama-cpp"))).To(BeTrue())
})
It("returns 400 when no backend name is supplied", func() {
req := httptest.NewRequest(http.MethodPost, "/api/nodes/node-xyz/backends/upgrade", bytes.NewBufferString(`{}`))
req.Header.Set("Content-Type", "application/json")
rec := httptest.NewRecorder()
c := e.NewContext(req, rec)
c.SetParamNames("id")
c.SetParamValues("node-xyz")
handler := localai.UpgradeBackendOnNodeEndpoint(galleryService, opcache, appCfg)
Expect(handler(c)).To(Succeed())
Expect(rec.Code).To(Equal(http.StatusBadRequest))
})
})
})

View File

@@ -85,14 +85,14 @@ var _ = Describe("Open Responses API", func() {
Expect(err).ToNot(HaveOccurred())
go func() {
if err := app.Start("127.0.0.1:9090"); err != nil && err != http.ErrServerClosed {
if err := app.Start(testHTTPAddr); err != nil && err != http.ErrServerClosed {
xlog.Error("server error", "error", err)
}
}()
// Wait for API to be ready
Eventually(func() error {
resp, err := http.Get("http://127.0.0.1:9090/healthz")
resp, err := http.Get("http://" + testHTTPAddr + "/healthz")
if err != nil {
return err
}
@@ -131,7 +131,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -154,7 +154,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -179,7 +179,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -204,7 +204,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -232,7 +232,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -277,7 +277,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -305,7 +305,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -333,7 +333,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -364,7 +364,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -394,7 +394,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -422,7 +422,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -454,7 +454,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -490,7 +490,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -526,7 +526,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -575,7 +575,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -626,7 +626,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -660,7 +660,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -694,7 +694,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -716,7 +716,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -763,7 +763,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -792,7 +792,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -835,7 +835,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -871,7 +871,7 @@ var _ = Describe("Open Responses API", func() {
payload1, err := json.Marshal(reqBody1)
Expect(err).ToNot(HaveOccurred())
req1, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload1))
req1, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload1))
Expect(err).ToNot(HaveOccurred())
req1.Header.Set("Content-Type", "application/json")
req1.Header.Set("Authorization", bearerKey)
@@ -905,7 +905,7 @@ var _ = Describe("Open Responses API", func() {
payload2, err := json.Marshal(reqBody2)
Expect(err).ToNot(HaveOccurred())
req2, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload2))
req2, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload2))
Expect(err).ToNot(HaveOccurred())
req2.Header.Set("Content-Type", "application/json")
req2.Header.Set("Authorization", bearerKey)
@@ -933,7 +933,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)
@@ -969,7 +969,7 @@ var _ = Describe("Open Responses API", func() {
payload1, err := json.Marshal(reqBody1)
Expect(err).ToNot(HaveOccurred())
req1, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload1))
req1, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload1))
Expect(err).ToNot(HaveOccurred())
req1.Header.Set("Content-Type", "application/json")
req1.Header.Set("Authorization", bearerKey)
@@ -1022,7 +1022,7 @@ var _ = Describe("Open Responses API", func() {
payload2, err := json.Marshal(reqBody2)
Expect(err).ToNot(HaveOccurred())
req2, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload2))
req2, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload2))
Expect(err).ToNot(HaveOccurred())
req2.Header.Set("Content-Type", "application/json")
req2.Header.Set("Authorization", bearerKey)
@@ -1048,7 +1048,7 @@ var _ = Describe("Open Responses API", func() {
payload, err := json.Marshal(reqBody)
Expect(err).ToNot(HaveOccurred())
req, err := http.NewRequest("POST", "http://127.0.0.1:9090/v1/responses", bytes.NewBuffer(payload))
req, err := http.NewRequest("POST", "http://"+testHTTPAddr+"/v1/responses", bytes.NewBuffer(payload))
Expect(err).ToNot(HaveOccurred())
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", bearerKey)

View File

@@ -55,7 +55,10 @@ export default function NodeDetail() {
const resume = async () => { try { await nodesApi.resume(id); addToast('Node resumed', 'success'); refresh() } catch (e) { addToast(e.message, 'error') } }
const remove = async () => { try { await nodesApi.delete(id); addToast('Node removed', 'success'); navigate('/app/nodes') } catch (e) { addToast(e.message, 'error') } }
const unload = async (name) => { try { await nodesApi.unloadModel(id, name); addToast(`Model "${name}" unloaded`, 'success'); refresh() } catch (e) { addToast(e.message, 'error') } }
const upgradeBackend = async (name) => { try { await nodesApi.installBackend(id, name); addToast(`Backend "${name}" upgraded`, 'success'); refresh() } catch (e) { addToast(e.message, 'error') } }
// The upgrade runs async via the gallery job queue (202 + jobID); the
// global Operations panel tracks progress, so the toast only reports the
// dispatch, not completion.
const upgradeBackend = async (name) => { try { await nodesApi.upgradeBackend(id, name); addToast(`Upgrading "${name}" on this node...`, 'info'); setTimeout(refresh, 1200) } catch (e) { addToast(e.message, 'error') } }
const deleteBackend = async (name) => { try { await nodesApi.deleteBackend(id, name); addToast(`Backend "${name}" deleted`, 'success'); refresh() } catch (e) { addToast(e.message, 'error') } }
const addLabel = async (k, v) => { try { await nodesApi.mergeLabels(id, { [k]: v }); refresh() } catch (e) { addToast(e.message, 'error') } }
const delLabel = async (k) => { try { await nodesApi.deleteLabel(id, k); refresh() } catch (e) { addToast(e.message, 'error') } }

View File

@@ -567,6 +567,11 @@ export const nodesApi = {
...(opts.alias ? { alias: opts.alias } : {}),
...(opts.backend_galleries ? { backend_galleries: opts.backend_galleries } : {}),
}),
// upgradeBackend force-reinstalls a gallery backend on a single node. This
// is a distinct endpoint from installBackend: the worker treats install as
// "ensure installed" and no-ops when the backend already exists on disk,
// so an upgrade dispatched through install would silently do nothing.
upgradeBackend: (id, backend) => postJSON(API_CONFIG.endpoints.nodeBackendsUpgrade(id), { backend }),
deleteBackend: (id, backend) => postJSON(API_CONFIG.endpoints.nodeBackendsDelete(id), { backend }),
getBackendLogs: (id) => fetchJSON(API_CONFIG.endpoints.nodeBackendLogs(id)),
getBackendLogLines: (id, modelId) => fetchJSON(API_CONFIG.endpoints.nodeBackendLogsModel(id, modelId)),

View File

@@ -137,6 +137,7 @@ export const API_CONFIG = {
nodeHeartbeat: (id) => `/api/nodes/${id}/heartbeat`,
nodeBackends: (id) => `/api/nodes/${id}/backends`,
nodeBackendsInstall: (id) => `/api/nodes/${id}/backends/install`,
nodeBackendsUpgrade: (id) => `/api/nodes/${id}/backends/upgrade`,
nodeBackendsDelete: (id) => `/api/nodes/${id}/backends/delete`,
nodeBackendLogs: (id) => `/api/nodes/${id}/backend-logs`,
nodeBackendLogsModel: (id, modelId) => `/api/nodes/${id}/backend-logs/${encodeURIComponent(modelId)}`,

View File

@@ -90,6 +90,10 @@ func RegisterNodeAdminRoutes(e *echo.Echo, registry *nodes.NodeRegistry, unloade
// Backend management on workers
admin.GET("/:id/backends", localai.ListBackendsOnNodeEndpoint(unloader, registry))
admin.POST("/:id/backends/install", localai.InstallBackendOnNodeEndpoint(unloader, galleryService, opcache, appConfig))
// Upgrade is a distinct route (not install) because the worker's
// backend.install handler short-circuits when the backend already exists
// on disk; only the Upgrade op path force-reinstalls.
admin.POST("/:id/backends/upgrade", localai.UpgradeBackendOnNodeEndpoint(galleryService, opcache, appConfig))
admin.POST("/:id/backends/delete", localai.DeleteBackendOnNodeEndpoint(unloader))
// Model management on workers

View File

@@ -0,0 +1,17 @@
package http_test
import "os"
// testHTTPAddr is the loopback address the in-process API server binds in
// these suites. The 9090 default matches what CI has always used; set
// LOCALAI_TEST_HTTP_PORT when something else already listens on 9090 locally,
// otherwise the pre-commit coverage gate can never pass on that machine (the
// suite would poll whatever service squats the port and time out).
var testHTTPAddr = "127.0.0.1:" + testHTTPPort()
func testHTTPPort() string {
if p := os.Getenv("LOCALAI_TEST_HTTP_PORT"); p != "" {
return p
}
return "9090"
}

View File

@@ -71,7 +71,7 @@ func (g *GalleryService) backendHandler(op *ManagementOp[gallery.GalleryBackend,
var err error
if op.Upgrade {
err = g.backendManager.UpgradeBackend(ctx, op.ID, op.GalleryElementName, progressCallback)
err = g.backendManager.UpgradeBackend(ctx, op, progressCallback)
} else if op.Delete {
err = g.backendManager.DeleteBackend(op.GalleryElementName)
} else {

View File

@@ -20,7 +20,7 @@ type BackendManager interface {
InstallBackend(ctx context.Context, op *ManagementOp[gallery.GalleryBackend, any], progressCb ProgressCallback) error
DeleteBackend(name string) error
ListBackends() (gallery.SystemBackends, error)
UpgradeBackend(ctx context.Context, opID, name string, progressCb ProgressCallback) error
UpgradeBackend(ctx context.Context, op *ManagementOp[gallery.GalleryBackend, any], progressCb ProgressCallback) error
CheckUpgrades(ctx context.Context) (map[string]gallery.UpgradeInfo, error)
// IsDistributed reports whether installs fan out across worker nodes.
// The HTTP layer uses this to refuse hardware-specific (non-meta) installs

View File

@@ -101,11 +101,12 @@ func (b *LocalBackendManager) ListBackends() (gallery.SystemBackends, error) {
return gallery.ListSystemBackends(b.systemState)
}
// UpgradeBackend ignores opID: a single-node install reports progress through
// the local progressCb already; opID only matters for distributed per-node
// streaming (see DistributedBackendManager.UpgradeBackend).
func (b *LocalBackendManager) UpgradeBackend(ctx context.Context, _ string, name string, progressCb ProgressCallback) error {
return gallery.UpgradeBackend(ctx, b.systemState, b.modelLoader, b.backendGalleries, name, progressCb, b.requireBackendIntegrity)
// UpgradeBackend ignores op.ID and op.TargetNodeID: a single-node install
// reports progress through the local progressCb already, and there is only
// one node to target. Both fields only matter for distributed per-node
// streaming/scoping (see DistributedBackendManager.UpgradeBackend).
func (b *LocalBackendManager) UpgradeBackend(ctx context.Context, op *ManagementOp[gallery.GalleryBackend, any], progressCb ProgressCallback) error {
return gallery.UpgradeBackend(ctx, b.systemState, b.modelLoader, b.backendGalleries, op.GalleryElementName, progressCb, b.requireBackendIntegrity)
}
func (b *LocalBackendManager) CheckUpgrades(ctx context.Context) (map[string]gallery.UpgradeInfo, error) {

View File

@@ -533,7 +533,9 @@ func (d *DistributedBackendManager) InstallBackend(ctx context.Context, op *gall
// backend.upgrade, we try the legacy backend.install Force=true path so a
// new master + old worker still converges. Drop the fallback once every
// worker in the fleet is on 2026-05-08 or newer.
func (d *DistributedBackendManager) UpgradeBackend(ctx context.Context, opID, name string, progressCb galleryop.ProgressCallback) error {
func (d *DistributedBackendManager) UpgradeBackend(ctx context.Context, op *galleryop.ManagementOp[gallery.GalleryBackend, any], progressCb galleryop.ProgressCallback) error {
opID := op.ID
name := op.GalleryElementName
galleriesJSON, _ := json.Marshal(d.backendGalleries)
installed, err := d.ListBackends()
@@ -548,6 +550,16 @@ func (d *DistributedBackendManager) UpgradeBackend(ctx context.Context, opID, na
for _, n := range entry.Nodes {
targetNodeIDs[n.NodeID] = true
}
// Node-scoped upgrade (node detail page): restrict the fan-out to the one
// requested node, but only if that node actually reports the backend
// installed: upgrading a backend a node never had fails at the gallery
// and leaves a forever-retrying pending_backend_ops row.
if op.TargetNodeID != "" {
if !targetNodeIDs[op.TargetNodeID] {
return fmt.Errorf("backend %q is not installed on node %s", name, op.TargetNodeID)
}
targetNodeIDs = map[string]bool{op.TargetNodeID: true}
}
result, err := d.enqueueAndDrainBackendOp(ctx, opID, OpBackendUpgrade, name, galleriesJSON, targetNodeIDs, func(node BackendNode) error {
// Per-node progress sink: fan each worker download tick into the legacy

View File

@@ -317,7 +317,7 @@ func (stubLocalBackendManager) DeleteBackend(_ string) error { return gallery.Er
func (stubLocalBackendManager) ListBackends() (gallery.SystemBackends, error) {
return gallery.SystemBackends{}, nil
}
func (stubLocalBackendManager) UpgradeBackend(_ context.Context, _ string, _ string, _ galleryop.ProgressCallback) error {
func (stubLocalBackendManager) UpgradeBackend(_ context.Context, _ *galleryop.ManagementOp[gallery.GalleryBackend, any], _ galleryop.ProgressCallback) error {
return nil
}
func (stubLocalBackendManager) CheckUpgrades(_ context.Context) (map[string]gallery.UpgradeInfo, error) {
@@ -753,6 +753,11 @@ var _ = Describe("DistributedBackendManager", func() {
})
Describe("UpgradeBackend", func() {
// upgradeOp builds the minimal ManagementOp an upgrade caller enqueues:
// just the element name (cluster-wide) or name + TargetNodeID (node-scoped).
upgradeOp := func(name string) *galleryop.ManagementOp[gallery.GalleryBackend, any] {
return &galleryop.ManagementOp[gallery.GalleryBackend, any]{GalleryElementName: name, Upgrade: true}
}
// scriptInstalled tells the worker(s) named in `nodeIDs` to claim
// `backend` is installed when DistributedBackendManager.ListBackends()
// fans out backend.list. Anything not scripted defaults to an empty
@@ -782,7 +787,7 @@ var _ = Describe("DistributedBackendManager", func() {
mc.scriptReply(messaging.SubjectNodeBackendUpgrade(n2.ID),
messaging.BackendUpgradeReply{Success: false, Error: "registry unauthorized"})
err := mgr.UpgradeBackend(ctx, "", "vllm-development", nil)
err := mgr.UpgradeBackend(ctx, upgradeOp("vllm-development"), nil)
Expect(err).To(HaveOccurred())
Expect(err.Error()).To(ContainSubstring("worker-a"))
Expect(err.Error()).To(ContainSubstring("image manifest not found"))
@@ -797,7 +802,7 @@ var _ = Describe("DistributedBackendManager", func() {
scriptInstalled("vllm-development", n1.ID)
mc.scriptReply(messaging.SubjectNodeBackendUpgrade(n1.ID),
messaging.BackendUpgradeReply{Success: true})
Expect(mgr.UpgradeBackend(ctx, "", "vllm-development", nil)).To(Succeed())
Expect(mgr.UpgradeBackend(ctx, upgradeOp("vllm-development"), nil)).To(Succeed())
})
})
@@ -819,7 +824,7 @@ var _ = Describe("DistributedBackendManager", func() {
// if the manager attempts it, the scripted-client default returns
// fakeNoRespondersErr and the assertion below fails loudly.
Expect(mgr.UpgradeBackend(ctx, "", "cpu-insightface-development", nil)).To(Succeed())
Expect(mgr.UpgradeBackend(ctx, upgradeOp("cpu-insightface-development"), nil)).To(Succeed())
mc.mu.Lock()
defer mc.mu.Unlock()
@@ -830,12 +835,73 @@ var _ = Describe("DistributedBackendManager", func() {
})
})
// Node-scoped upgrade: the node detail page upgrades a backend on ONE
// node. op.TargetNodeID restricts the fan-out the same way it does for
// InstallBackend - without it the per-node button silently upgraded the
// whole cluster (or, through the install path it used to share, did
// nothing at all).
Context("when op.TargetNodeID scopes the upgrade to a single node", func() {
It("sends backend.upgrade only to the target node", func() {
n1 := registerHealthyBackend("worker-a", "10.0.0.1:50051")
n2 := registerHealthyBackend("worker-b", "10.0.0.2:50051")
scriptInstalled("vllm-development", n1.ID, n2.ID)
mc.scriptReply(messaging.SubjectNodeBackendUpgrade(n1.ID),
messaging.BackendUpgradeReply{Success: true})
mc.scriptReply(messaging.SubjectNodeBackendUpgrade(n2.ID),
messaging.BackendUpgradeReply{Success: true})
op := upgradeOp("vllm-development")
op.TargetNodeID = n2.ID
Expect(mgr.UpgradeBackend(ctx, op, nil)).To(Succeed())
mc.mu.Lock()
defer mc.mu.Unlock()
upgraded := map[string]bool{}
for _, call := range mc.calls {
if call.Subject == messaging.SubjectNodeBackendUpgrade(n1.ID) {
upgraded[n1.ID] = true
}
if call.Subject == messaging.SubjectNodeBackendUpgrade(n2.ID) {
upgraded[n2.ID] = true
}
}
Expect(upgraded).To(HaveKey(n2.ID), "target node never received backend.upgrade")
Expect(upgraded).ToNot(HaveKey(n1.ID), "upgrade leaked to a non-target node")
})
It("errors when the target node does not have the backend installed", func() {
has := registerHealthyBackend("worker-a", "10.0.0.1:50051")
lacks := registerHealthyBackend("worker-b", "10.0.0.2:50051")
scriptInstalled("vllm-development", has.ID)
scriptNoBackends(lacks.ID)
mc.scriptReply(messaging.SubjectNodeBackendUpgrade(has.ID),
messaging.BackendUpgradeReply{Success: true})
op := upgradeOp("vllm-development")
op.TargetNodeID = lacks.ID
err := mgr.UpgradeBackend(ctx, op, nil)
Expect(err).To(HaveOccurred())
Expect(err.Error()).To(ContainSubstring("not installed on node"))
mc.mu.Lock()
defer mc.mu.Unlock()
for _, call := range mc.calls {
Expect(call.Subject).ToNot(Equal(messaging.SubjectNodeBackendUpgrade(has.ID)),
"a node-scoped upgrade for %s must not touch other nodes", lacks.Name)
Expect(call.Subject).ToNot(Equal(messaging.SubjectNodeBackendUpgrade(lacks.ID)),
"the target node lacks the backend; nothing should be sent")
}
})
})
Context("when no node has the backend installed", func() {
It("returns a clear error and never attempts an install request", func() {
n1 := registerHealthyBackend("worker-a", "10.0.0.1:50051")
scriptNoBackends(n1.ID)
err := mgr.UpgradeBackend(ctx, "", "vllm-development", nil)
err := mgr.UpgradeBackend(ctx, upgradeOp("vllm-development"), nil)
Expect(err).To(HaveOccurred())
Expect(err.Error()).To(ContainSubstring("not installed on any node"))
@@ -865,7 +931,7 @@ var _ = Describe("DistributedBackendManager", func() {
func(req messaging.BackendInstallRequest) bool { return req.Force },
messaging.BackendInstallReply{Success: true, Address: "10.0.0.1:50100"})
Expect(mgr.UpgradeBackend(ctx, "", "vllm-development", nil)).To(Succeed())
Expect(mgr.UpgradeBackend(ctx, upgradeOp("vllm-development"), nil)).To(Succeed())
})
It("returns the upgrade error when it is not ErrNoResponders", func() {
@@ -875,7 +941,7 @@ var _ = Describe("DistributedBackendManager", func() {
mc.scriptReply(messaging.SubjectNodeBackendUpgrade(n.ID),
messaging.BackendUpgradeReply{Success: false, Error: "disk full"})
err := mgr.UpgradeBackend(ctx, "", "vllm-development", nil)
err := mgr.UpgradeBackend(ctx, upgradeOp("vllm-development"), nil)
Expect(err).To(HaveOccurred())
Expect(err.Error()).To(ContainSubstring("disk full"))
})

View File

@@ -333,6 +333,7 @@ Used by the WebUI and admin API consumers. Requires admin authentication.
| `POST` | `/api/nodes/:id/drain` | Admin-drain a worker |
| `POST` | `/api/nodes/:id/approve` | Approve a pending worker node |
| `POST` | `/api/nodes/:id/backends/install` | Install a backend on a worker |
| `POST` | `/api/nodes/:id/backends/upgrade` | Upgrade (force-reinstall) a backend on a worker |
| `POST` | `/api/nodes/:id/backends/delete` | Delete a backend from a worker |
| `POST` | `/api/nodes/:id/models/unload` | Unload a model from a worker |
| `POST` | `/api/nodes/:id/models/delete` | Delete model files from a worker |