diff --git a/core/application/upgrade_checker.go b/core/application/upgrade_checker.go index d537695d3..e18548a4e 100644 --- a/core/application/upgrade_checker.go +++ b/core/application/upgrade_checker.go @@ -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) diff --git a/core/http/app_test.go b/core/http/app_test.go index fa6423292..36072b96e 100644 --- a/core/http/app_test.go +++ b/core/http/app_test.go @@ -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", ¶msBody, &jobResp) + err = postRequestResponseJSON("http://"+testHTTPAddr+"/api/agent/tasks/Named Task/execute", ¶msBody, &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 diff --git a/core/http/endpoints/localai/nodes.go b/core/http/endpoints/localai/nodes.go index 71b4cbb11..2fb4ac530 100644 --- a/core/http/endpoints/localai/nodes.go +++ b/core/http/endpoints/localai/nodes.go @@ -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 { diff --git a/core/http/endpoints/localai/nodes_install_async_test.go b/core/http/endpoints/localai/nodes_install_async_test.go index c3ae9745a..ad3b58475 100644 --- a/core/http/endpoints/localai/nodes_install_async_test.go +++ b/core/http/endpoints/localai/nodes_install_async_test.go @@ -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)) + }) + }) }) diff --git a/core/http/openresponses_test.go b/core/http/openresponses_test.go index 4e6eca7b7..f30674362 100644 --- a/core/http/openresponses_test.go +++ b/core/http/openresponses_test.go @@ -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) diff --git a/core/http/react-ui/src/pages/NodeDetail.jsx b/core/http/react-ui/src/pages/NodeDetail.jsx index f61ba0576..6414cb4c3 100644 --- a/core/http/react-ui/src/pages/NodeDetail.jsx +++ b/core/http/react-ui/src/pages/NodeDetail.jsx @@ -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') } } diff --git a/core/http/react-ui/src/utils/api.js b/core/http/react-ui/src/utils/api.js index 9e9e92500..12bc9366b 100644 --- a/core/http/react-ui/src/utils/api.js +++ b/core/http/react-ui/src/utils/api.js @@ -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)), diff --git a/core/http/react-ui/src/utils/config.js b/core/http/react-ui/src/utils/config.js index 0fa0703b3..a1c3c9fa8 100644 --- a/core/http/react-ui/src/utils/config.js +++ b/core/http/react-ui/src/utils/config.js @@ -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)}`, diff --git a/core/http/routes/nodes.go b/core/http/routes/nodes.go index e35bea240..fbbcce081 100644 --- a/core/http/routes/nodes.go +++ b/core/http/routes/nodes.go @@ -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 diff --git a/core/http/testport_test.go b/core/http/testport_test.go new file mode 100644 index 000000000..6b282b5c2 --- /dev/null +++ b/core/http/testport_test.go @@ -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" +} diff --git a/core/services/galleryop/backends.go b/core/services/galleryop/backends.go index 0ad1e64c9..5400f7692 100644 --- a/core/services/galleryop/backends.go +++ b/core/services/galleryop/backends.go @@ -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 { diff --git a/core/services/galleryop/managers.go b/core/services/galleryop/managers.go index 87720696c..6435a9963 100644 --- a/core/services/galleryop/managers.go +++ b/core/services/galleryop/managers.go @@ -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 diff --git a/core/services/galleryop/managers_local.go b/core/services/galleryop/managers_local.go index 7d6b75879..10fb29491 100644 --- a/core/services/galleryop/managers_local.go +++ b/core/services/galleryop/managers_local.go @@ -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) { diff --git a/core/services/nodes/managers_distributed.go b/core/services/nodes/managers_distributed.go index 42fbab1aa..127425b1a 100644 --- a/core/services/nodes/managers_distributed.go +++ b/core/services/nodes/managers_distributed.go @@ -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 diff --git a/core/services/nodes/managers_distributed_test.go b/core/services/nodes/managers_distributed_test.go index d32fb1ecd..b83200eeb 100644 --- a/core/services/nodes/managers_distributed_test.go +++ b/core/services/nodes/managers_distributed_test.go @@ -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")) }) diff --git a/docs/content/features/distributed-mode.md b/docs/content/features/distributed-mode.md index fc3462412..adf0e3292 100644 --- a/docs/content/features/distributed-mode.md +++ b/docs/content/features/distributed-mode.md @@ -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 |