Skip to content

Commit 9400b4b

Browse files
committed
feat: persist helm install lifecycle
1 parent fb82883 commit 9400b4b

15 files changed

Lines changed: 1806 additions & 88 deletions

‎controllers/helm.go‎

Lines changed: 122 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,16 +1,24 @@
11
package controllers
22

33
import (
4+
"context"
5+
"crypto/sha256"
46
"encoding/json"
7+
"errors"
58
"fmt"
69
"io"
710
"net/http"
11+
"strconv"
12+
"strings"
813
"time"
914

15+
"github.com/beego/beego/logs"
1016
"github.com/casosorg/casos/object"
1117
"github.com/casosorg/casos/store"
1218
)
1319

20+
const helmOperationTaskNotFoundCode = "helm_task_not_found"
21+
1422
// ---------- ArtifactHub proxy ----------
1523

1624
type ahSearchResult struct {
@@ -231,27 +239,137 @@ func (c *ApiController) InstallHelmChartStream() {
231239
c.StopRun()
232240
return
233241
}
242+
owner := helmOperationOwner(c)
243+
if owner == "" {
244+
c.Ctx.ResponseWriter.ResponseWriter.Header().Set("Content-Type", "text/event-stream")
245+
fmt.Fprint(c.Ctx.ResponseWriter.ResponseWriter, "data: ERROR: unable to identify Helm operation owner\n\n")
246+
c.StopRun()
247+
return
248+
}
234249

235250
w := c.Ctx.ResponseWriter.ResponseWriter
236251
w.Header().Set("Content-Type", "text/event-stream")
237252
w.Header().Set("Cache-Control", "no-cache")
238253
w.Header().Set("X-Accel-Buffering", "no")
239254
w.WriteHeader(http.StatusOK)
240255

241-
flusher, canFlush := w.(http.Flusher)
242256
ctx := c.Ctx.Request.Context()
243-
logCh := store.InstallHelmChartStream(ctx, cfg, req.ReleaseName, req.Namespace, req.ChartName, req.RepoURL, req.Version, req.ValuesYAML)
257+
task, err := object.CreateHelmOperationTask(owner, object.HelmOperationInstall, req.ReleaseName, req.Namespace, req.ChartName, req.Version)
258+
if err != nil {
259+
message := "unable to start Helm installation"
260+
if errors.Is(err, object.ErrHelmOperationAlreadyActive) {
261+
message = err.Error()
262+
} else {
263+
logs.Error("create Helm operation task: %v", err)
264+
}
265+
fmt.Fprintf(w, "data: ERROR: %s\n\n", message)
266+
c.StopRun()
267+
return
268+
}
269+
finishUnstartedTask := func(cause error) {
270+
finishCtx, cancel := context.WithTimeout(context.Background(), object.HelmOperationPersistenceTimeout)
271+
defer cancel()
272+
if finishErr := object.FinishHelmOperationTaskContext(finishCtx, task.Id, false, cause.Error()); finishErr != nil {
273+
logs.Error("finish unstarted Helm operation task %d: %v", task.Id, finishErr)
274+
}
275+
}
276+
if _, err := fmt.Fprintf(w, "data: TASK_ID:%d\n\n", task.Id); err != nil {
277+
finishUnstartedTask(fmt.Errorf("failed to send Helm operation task id: %w", err))
278+
c.StopRun()
279+
return
280+
}
281+
responseController := http.NewResponseController(w)
282+
if err := responseController.Flush(); err != nil {
283+
finishUnstartedTask(fmt.Errorf("failed to flush Helm operation task id: %w", err))
284+
c.StopRun()
285+
return
286+
}
287+
recorder := object.NewHelmOperationRecorder(task.Id)
288+
logCh := store.InstallHelmChartStream(ctx, recorder, cfg, req.ReleaseName, req.Namespace, req.ChartName, req.RepoURL, req.Version, req.ValuesYAML)
244289
for line := range logCh {
245290
if _, err := fmt.Fprintf(w, "data: %s\n\n", line); err != nil {
246291
break
247292
}
248-
if canFlush {
249-
flusher.Flush()
293+
if err := responseController.Flush(); err != nil {
294+
break
250295
}
251296
}
252297
c.StopRun()
253298
}
254299

300+
// GetHelmOperationTask returns a persisted install task and its log history so
301+
// an administrator can reconnect after an SSE stream is interrupted.
302+
// @router /api/get-helm-operation-task [get]
303+
func (c *ApiController) GetHelmOperationTask() {
304+
if c.RequireAdmin() {
305+
return
306+
}
307+
id, err := strconv.ParseInt(c.GetString("id"), 10, 64)
308+
if err != nil || id <= 0 {
309+
c.ResponseError("invalid task id")
310+
return
311+
}
312+
owner := helmOperationOwner(c)
313+
if owner == "" {
314+
c.ResponseError("unable to identify Helm operation owner")
315+
return
316+
}
317+
task, err := object.GetHelmOperationTaskForOwner(id, owner)
318+
if err != nil {
319+
logs.Error("get Helm operation task %d: %v", id, err)
320+
c.ResponseError("failed to load Helm operation task")
321+
return
322+
}
323+
if task == nil {
324+
c.ResponseError("Helm operation task not found", helmOperationTaskNotFoundCode)
325+
return
326+
}
327+
if object.IsHelmOperationTaskStale(task, time.Now().UTC()) {
328+
if err := object.ExpireStaleHelmOperationTask(id); err != nil {
329+
logs.Error("expire Helm operation task %d: %v", id, err)
330+
c.ResponseError("failed to load Helm operation task")
331+
return
332+
}
333+
task, err = object.GetHelmOperationTaskForOwner(id, owner)
334+
if err != nil {
335+
logs.Error("reload Helm operation task %d: %v", id, err)
336+
c.ResponseError("failed to load Helm operation task")
337+
return
338+
}
339+
if task == nil {
340+
c.ResponseError("Helm operation task not found", helmOperationTaskNotFoundCode)
341+
return
342+
}
343+
}
344+
taskLogs, err := object.GetHelmOperationLogs(id, 1000)
345+
if err != nil {
346+
logs.Error("get Helm operation task %d logs: %v", id, err)
347+
c.ResponseError("failed to load Helm operation task")
348+
return
349+
}
350+
c.ResponseOk(task, taskLogs)
351+
}
352+
353+
func helmOperationOwner(c *ApiController) string {
354+
if user := c.GetSessionUser(); user != nil {
355+
return canonicalHelmOperationOwner(user.Id, user.Owner, user.Name)
356+
}
357+
return ""
358+
}
359+
360+
func canonicalHelmOperationOwner(id, owner, name string) string {
361+
if id = strings.TrimSpace(id); id != "" {
362+
return id
363+
}
364+
owner = strings.TrimSpace(owner)
365+
name = strings.TrimSpace(name)
366+
if owner == "" || name == "" {
367+
return ""
368+
}
369+
digest := sha256.Sum256([]byte(owner + "\x00" + name))
370+
return fmt.Sprintf("casdoor:%x", digest)
371+
}
372+
255373
// UpgradeHelmRelease upgrades an existing Helm release.
256374
// @router /api/upgrade-helm-release [post]
257375
func (c *ApiController) UpgradeHelmRelease() {

‎controllers/helm_test.go‎

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,22 @@
1+
package controllers
2+
3+
import "testing"
4+
5+
func TestCanonicalHelmOperationOwner(t *testing.T) {
6+
if got := canonicalHelmOperationOwner(" stable-id ", "built-in", "ci-user"); got != "stable-id" {
7+
t.Fatalf("stable ID was not preferred: %q", got)
8+
}
9+
fallback := canonicalHelmOperationOwner("", "built-in", "ci-user")
10+
if fallback == "" || len(fallback) > 100 {
11+
t.Fatalf("invalid owner/name fallback: %q", fallback)
12+
}
13+
if fallback != canonicalHelmOperationOwner("", "built-in", "ci-user") {
14+
t.Fatal("owner/name fallback is not deterministic")
15+
}
16+
if fallback == canonicalHelmOperationOwner("", "other", "ci-user") {
17+
t.Fatal("different Casdoor identities produced the same owner")
18+
}
19+
if got := canonicalHelmOperationOwner("", "built-in", ""); got != "" {
20+
t.Fatalf("incomplete identity produced an owner: %q", got)
21+
}
22+
}

0 commit comments

Comments
 (0)