Skip to content

Commit 869c07d

Browse files
committed
feat: persist helm install lifecycle
1 parent e295ede commit 869c07d

10 files changed

Lines changed: 603 additions & 18 deletions

File tree

‎controllers/helm.go‎

Lines changed: 65 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@ import (
55
"fmt"
66
"io"
77
"net/http"
8+
"strconv"
89
"time"
910

1011
"github.com/casosorg/casos/object"
@@ -240,7 +241,19 @@ func (c *ApiController) InstallHelmChartStream() {
240241

241242
flusher, canFlush := w.(http.Flusher)
242243
ctx := c.Ctx.Request.Context()
243-
logCh := store.InstallHelmChartStream(ctx, cfg, req.ReleaseName, req.Namespace, req.ChartName, req.RepoURL, req.Version, req.ValuesYAML)
244+
task, logCh, err := store.InstallHelmChartStream(helmOperationOwner(c), ctx, cfg, req.ReleaseName, req.Namespace, req.ChartName, req.RepoURL, req.Version, req.ValuesYAML)
245+
if err != nil {
246+
fmt.Fprintf(w, "data: ERROR: %s\n\n", err.Error())
247+
c.StopRun()
248+
return
249+
}
250+
if _, err := fmt.Fprintf(w, "data: TASK_ID:%d\n\n", task.Id); err != nil {
251+
c.StopRun()
252+
return
253+
}
254+
if canFlush {
255+
flusher.Flush()
256+
}
244257
for line := range logCh {
245258
if _, err := fmt.Fprintf(w, "data: %s\n\n", line); err != nil {
246259
break
@@ -252,6 +265,57 @@ func (c *ApiController) InstallHelmChartStream() {
252265
c.StopRun()
253266
}
254267

268+
// GetHelmOperationTask returns a persisted install task and its log history so
269+
// a browser can reconnect after an SSE stream is interrupted.
270+
// @router /api/get-helm-operation-task [get]
271+
func (c *ApiController) GetHelmOperationTask() {
272+
if c.RequireSignedIn() {
273+
return
274+
}
275+
id, err := strconv.ParseInt(c.GetString("id"), 10, 64)
276+
if err != nil || id <= 0 {
277+
c.ResponseError("invalid task id")
278+
return
279+
}
280+
task, err := object.GetHelmOperationTask(id)
281+
if err != nil {
282+
c.ResponseError(err.Error())
283+
return
284+
}
285+
if task == nil || task.Owner != helmOperationOwner(c) {
286+
c.ResponseError("Helm operation task not found")
287+
return
288+
}
289+
logs, err := object.GetHelmOperationLogs(id, 1000)
290+
if err != nil {
291+
c.ResponseError(err.Error())
292+
return
293+
}
294+
c.ResponseOk(task, logs)
295+
}
296+
297+
// GetHelmOperationTasks returns recent persisted install tasks for the user.
298+
// @router /api/get-helm-operation-tasks [get]
299+
func (c *ApiController) GetHelmOperationTasks() {
300+
if c.RequireSignedIn() {
301+
return
302+
}
303+
limit, _ := c.GetInt("limit", 50)
304+
tasks, err := object.GetHelmOperationTasks(helmOperationOwner(c), limit)
305+
if err != nil {
306+
c.ResponseError(err.Error())
307+
return
308+
}
309+
c.ResponseOk(tasks)
310+
}
311+
312+
func helmOperationOwner(c *ApiController) string {
313+
if user := c.GetSessionUser(); user != nil && user.Name != "" {
314+
return user.Name
315+
}
316+
return "admin"
317+
}
318+
255319
// UpgradeHelmRelease upgrades an existing Helm release.
256320
// @router /api/upgrade-helm-release [post]
257321
func (c *ApiController) UpgradeHelmRelease() {

‎main.go‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -27,6 +27,9 @@ func main() {
2727
object.InitFlag()
2828
object.InitAdapter()
2929
object.CreateTables()
30+
if err := object.FailActiveHelmOperationTasks("Helm operation was interrupted by service restart"); err != nil {
31+
logs.Warning("Helm operation task cleanup: %v", err)
32+
}
3033
object.InitSite()
3134
if err := object.SeedDefaultPolicies(); err != nil {
3235
logs.Warning("casbin seed: %v", err)

‎object/helm_operation.go‎

Lines changed: 263 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,263 @@
1+
package object
2+
3+
import (
4+
"fmt"
5+
"strings"
6+
"time"
7+
8+
"xorm.io/xorm"
9+
)
10+
11+
const (
12+
HelmOperationInstall = "install"
13+
14+
HelmOperationStatusPending = "pending"
15+
HelmOperationStatusRunning = "running"
16+
HelmOperationStatusSucceeded = "succeeded"
17+
HelmOperationStatusFailed = "failed"
18+
19+
HelmOperationPhaseQueued = "queued"
20+
HelmOperationPhaseLoading = "loading"
21+
HelmOperationPhaseInstalling = "installing"
22+
HelmOperationPhaseReady = "ready"
23+
HelmOperationPhaseFailed = "failed"
24+
25+
HelmOperationLogLevelInfo = "info"
26+
HelmOperationLogLevelError = "error"
27+
)
28+
29+
type HelmOperationTask struct {
30+
Id int64 `xorm:"pk autoincr" json:"id"`
31+
Owner string `xorm:"varchar(100) notnull index" json:"owner"`
32+
Operation string `xorm:"varchar(30) notnull" json:"operation"`
33+
ReleaseName string `xorm:"varchar(253) notnull index" json:"releaseName"`
34+
Namespace string `xorm:"varchar(253) notnull index" json:"namespace"`
35+
ChartName string `xorm:"varchar(253) notnull" json:"chartName"`
36+
Version string `xorm:"varchar(100)" json:"version"`
37+
Status string `xorm:"varchar(30) notnull index" json:"status"`
38+
Phase string `xorm:"varchar(30) notnull" json:"phase"`
39+
ErrorMsg string `xorm:"text" json:"errorMsg"`
40+
CreatedAt time.Time `json:"createdAt"`
41+
StartedAt time.Time `json:"startedAt"`
42+
FinishedAt time.Time `json:"finishedAt"`
43+
UpdatedAt time.Time `json:"updatedAt"`
44+
}
45+
46+
type HelmOperationLog struct {
47+
Id int64 `xorm:"pk autoincr" json:"id"`
48+
TaskId int64 `xorm:"notnull index" json:"taskId"`
49+
Level string `xorm:"varchar(20) notnull" json:"level"`
50+
Message string `xorm:"text" json:"message"`
51+
CreatedAt time.Time `json:"createdAt"`
52+
}
53+
54+
func CreateHelmOperationTask(owner, operation, releaseName, namespace, chartName, version string) (*HelmOperationTask, error) {
55+
owner = strings.TrimSpace(owner)
56+
operation = strings.TrimSpace(operation)
57+
releaseName = strings.TrimSpace(releaseName)
58+
namespace = strings.TrimSpace(namespace)
59+
chartName = strings.TrimSpace(chartName)
60+
if owner == "" || operation == "" || releaseName == "" || namespace == "" || chartName == "" {
61+
return nil, fmt.Errorf("owner, operation, releaseName, namespace, and chartName are required")
62+
}
63+
64+
result, err := withHelmOperationTransaction(func(session *xorm.Session) (interface{}, error) {
65+
active := &HelmOperationTask{}
66+
found, err := session.
67+
Where("namespace = ? AND release_name = ? AND status IN (?, ?)", namespace, releaseName, HelmOperationStatusPending, HelmOperationStatusRunning).
68+
Desc("id").
69+
ForUpdate().
70+
Get(active)
71+
if err != nil {
72+
return nil, err
73+
}
74+
if found {
75+
return nil, fmt.Errorf("Helm operation task %d is already active for %s/%s", active.Id, namespace, releaseName)
76+
}
77+
78+
now := time.Now().UTC()
79+
task := &HelmOperationTask{
80+
Owner: owner,
81+
Operation: operation,
82+
ReleaseName: releaseName,
83+
Namespace: namespace,
84+
ChartName: chartName,
85+
Version: strings.TrimSpace(version),
86+
Status: HelmOperationStatusPending,
87+
Phase: HelmOperationPhaseQueued,
88+
CreatedAt: now,
89+
UpdatedAt: now,
90+
}
91+
if _, err := session.Insert(task); err != nil {
92+
return nil, err
93+
}
94+
return task, nil
95+
})
96+
if err != nil {
97+
return nil, err
98+
}
99+
task, ok := result.(*HelmOperationTask)
100+
if !ok || task == nil {
101+
return nil, fmt.Errorf("create Helm operation task returned invalid result")
102+
}
103+
return task, nil
104+
}
105+
106+
func GetHelmOperationTask(id int64) (*HelmOperationTask, error) {
107+
task := &HelmOperationTask{Id: id}
108+
found, err := ormer.Engine.Get(task)
109+
if err != nil {
110+
return nil, err
111+
}
112+
if !found {
113+
return nil, nil
114+
}
115+
return task, nil
116+
}
117+
118+
func GetHelmOperationTasks(owner string, limit int) ([]*HelmOperationTask, error) {
119+
owner = strings.TrimSpace(owner)
120+
if owner == "" {
121+
return nil, fmt.Errorf("owner is required")
122+
}
123+
if limit <= 0 {
124+
limit = 50
125+
}
126+
if limit > 100 {
127+
return nil, fmt.Errorf("limit must not exceed 100")
128+
}
129+
tasks := []*HelmOperationTask{}
130+
err := ormer.Engine.Where("owner = ?", owner).Desc("id").Limit(limit).Find(&tasks)
131+
return tasks, err
132+
}
133+
134+
func StartHelmOperationTask(id int64, phase string) error {
135+
if !isValidHelmOperationPhase(phase) || phase == HelmOperationPhaseQueued || phase == HelmOperationPhaseFailed {
136+
return fmt.Errorf("invalid Helm operation start phase: %s", phase)
137+
}
138+
now := time.Now().UTC()
139+
affected, err := ormer.Engine.ID(id).
140+
Where("status = ?", HelmOperationStatusPending).
141+
Cols("status", "phase", "started_at", "updated_at").
142+
Update(&HelmOperationTask{Status: HelmOperationStatusRunning, Phase: phase, StartedAt: now, UpdatedAt: now})
143+
if err != nil {
144+
return err
145+
}
146+
if affected == 0 {
147+
return fmt.Errorf("Helm operation task %d is not pending", id)
148+
}
149+
return nil
150+
}
151+
152+
func UpdateHelmOperationTaskPhase(id int64, phase string) error {
153+
if !isValidHelmOperationPhase(phase) || phase == HelmOperationPhaseQueued || phase == HelmOperationPhaseFailed {
154+
return fmt.Errorf("invalid Helm operation phase: %s", phase)
155+
}
156+
_, err := ormer.Engine.ID(id).
157+
Where("status = ?", HelmOperationStatusRunning).
158+
Cols("phase", "updated_at").
159+
Update(&HelmOperationTask{Phase: phase, UpdatedAt: time.Now().UTC()})
160+
return err
161+
}
162+
163+
func FinishHelmOperationTask(id int64, success bool, errorMsg string) error {
164+
status := HelmOperationStatusSucceeded
165+
phase := HelmOperationPhaseReady
166+
if !success {
167+
status = HelmOperationStatusFailed
168+
phase = HelmOperationPhaseFailed
169+
}
170+
now := time.Now().UTC()
171+
affected, err := ormer.Engine.ID(id).
172+
Where("status IN (?, ?)", HelmOperationStatusPending, HelmOperationStatusRunning).
173+
Cols("status", "phase", "error_msg", "finished_at", "updated_at").
174+
Update(&HelmOperationTask{Status: status, Phase: phase, ErrorMsg: errorMsg, FinishedAt: now, UpdatedAt: now})
175+
if err != nil {
176+
return err
177+
}
178+
if affected == 0 {
179+
return fmt.Errorf("Helm operation task %d is already finished", id)
180+
}
181+
return nil
182+
}
183+
184+
func AddHelmOperationLog(taskID int64, level, message string) error {
185+
if level != HelmOperationLogLevelInfo && level != HelmOperationLogLevelError {
186+
return fmt.Errorf("invalid Helm operation log level: %s", level)
187+
}
188+
_, err := ormer.Engine.Insert(&HelmOperationLog{
189+
TaskId: taskID,
190+
Level: level,
191+
Message: message,
192+
CreatedAt: time.Now().UTC(),
193+
})
194+
return err
195+
}
196+
197+
func GetHelmOperationLogs(taskID int64, limit int) ([]*HelmOperationLog, error) {
198+
if taskID <= 0 {
199+
return nil, fmt.Errorf("invalid taskId")
200+
}
201+
if limit <= 0 {
202+
limit = 500
203+
}
204+
if limit > 1000 {
205+
return nil, fmt.Errorf("limit must not exceed 1000")
206+
}
207+
logs := []*HelmOperationLog{}
208+
err := ormer.Engine.Where("task_id = ?", taskID).Asc("id").Limit(limit).Find(&logs)
209+
return logs, err
210+
}
211+
212+
func FailActiveHelmOperationTasks(reason string) error {
213+
reason = strings.TrimSpace(reason)
214+
if reason == "" {
215+
reason = "Helm operation was interrupted by service restart"
216+
}
217+
now := time.Now().UTC()
218+
_, err := ormer.Engine.
219+
Where("status IN (?, ?)", HelmOperationStatusPending, HelmOperationStatusRunning).
220+
Cols("status", "phase", "error_msg", "finished_at", "updated_at").
221+
Update(&HelmOperationTask{
222+
Status: HelmOperationStatusFailed,
223+
Phase: HelmOperationPhaseFailed,
224+
ErrorMsg: reason,
225+
FinishedAt: now,
226+
UpdatedAt: now,
227+
})
228+
return err
229+
}
230+
231+
func isValidHelmOperationPhase(phase string) bool {
232+
switch phase {
233+
case HelmOperationPhaseQueued, HelmOperationPhaseLoading, HelmOperationPhaseInstalling, HelmOperationPhaseReady, HelmOperationPhaseFailed:
234+
return true
235+
default:
236+
return false
237+
}
238+
}
239+
240+
func withHelmOperationTransaction(fn func(*xorm.Session) (interface{}, error)) (interface{}, error) {
241+
session := ormer.Engine.NewSession()
242+
defer func() {
243+
if v := recover(); v != nil {
244+
_ = session.Rollback()
245+
session.Close()
246+
panic(v)
247+
}
248+
session.Close()
249+
}()
250+
if err := session.Begin(); err != nil {
251+
return nil, err
252+
}
253+
result, err := fn(session)
254+
if err != nil {
255+
_ = session.Rollback()
256+
return nil, err
257+
}
258+
if err := session.Commit(); err != nil {
259+
_ = session.Rollback()
260+
return nil, err
261+
}
262+
return result, nil
263+
}

‎object/ormer.go‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -179,6 +179,8 @@ func (a *Ormer) createTable() {
179179
_ = a.Engine.Sync2(new(MachineNodeDeployTask))
180180
_ = a.Engine.Sync2(new(MachineNodeDeployLog))
181181
_ = a.Engine.Sync2(new(MachineNodeDeployCredential))
182+
_ = a.Engine.Sync2(new(HelmOperationTask))
183+
_ = a.Engine.Sync2(new(HelmOperationLog))
182184
_ = a.Engine.Sync2(new(CasbinRule))
183185
_ = a.Engine.Sync2(new(TrivyScanResult))
184186
_ = a.Engine.Sync2(new(HelmRepo))

‎routers/router.go‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -150,6 +150,8 @@ func InitAPI() {
150150
beego.Router("/api/get-helm-releases", &controllers.ApiController{}, "GET:GetHelmReleases")
151151
beego.Router("/api/install-helm-chart", &controllers.ApiController{}, "POST:InstallHelmChart")
152152
beego.Router("/api/install-helm-chart-stream", &controllers.ApiController{}, "POST:InstallHelmChartStream")
153+
beego.Router("/api/get-helm-operation-task", &controllers.ApiController{}, "GET:GetHelmOperationTask")
154+
beego.Router("/api/get-helm-operation-tasks", &controllers.ApiController{}, "GET:GetHelmOperationTasks")
153155
beego.Router("/api/upgrade-helm-release", &controllers.ApiController{}, "POST:UpgradeHelmRelease")
154156
beego.Router("/api/rollback-helm-release", &controllers.ApiController{}, "POST:RollbackHelmRelease")
155157
beego.Router("/api/uninstall-helm-release", &controllers.ApiController{}, "POST:UninstallHelmRelease")

0 commit comments

Comments
 (0)