|
1 | 1 | package controllers |
2 | 2 |
|
3 | 3 | import ( |
| 4 | + "context" |
| 5 | + "crypto/sha256" |
4 | 6 | "encoding/json" |
| 7 | + "errors" |
5 | 8 | "fmt" |
6 | 9 | "io" |
7 | 10 | "net/http" |
| 11 | + "strconv" |
| 12 | + "strings" |
8 | 13 | "time" |
9 | 14 |
|
| 15 | + "github.com/beego/beego/logs" |
10 | 16 | "github.com/casosorg/casos/object" |
11 | 17 | "github.com/casosorg/casos/store" |
12 | 18 | ) |
13 | 19 |
|
| 20 | +const helmOperationTaskNotFoundCode = "helm_task_not_found" |
| 21 | + |
14 | 22 | // ---------- ArtifactHub proxy ---------- |
15 | 23 |
|
16 | 24 | type ahSearchResult struct { |
@@ -231,27 +239,137 @@ func (c *ApiController) InstallHelmChartStream() { |
231 | 239 | c.StopRun() |
232 | 240 | return |
233 | 241 | } |
| 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 | + } |
234 | 249 |
|
235 | 250 | w := c.Ctx.ResponseWriter.ResponseWriter |
236 | 251 | w.Header().Set("Content-Type", "text/event-stream") |
237 | 252 | w.Header().Set("Cache-Control", "no-cache") |
238 | 253 | w.Header().Set("X-Accel-Buffering", "no") |
239 | 254 | w.WriteHeader(http.StatusOK) |
240 | 255 |
|
241 | | - flusher, canFlush := w.(http.Flusher) |
242 | 256 | 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) |
244 | 289 | for line := range logCh { |
245 | 290 | if _, err := fmt.Fprintf(w, "data: %s\n\n", line); err != nil { |
246 | 291 | break |
247 | 292 | } |
248 | | - if canFlush { |
249 | | - flusher.Flush() |
| 293 | + if err := responseController.Flush(); err != nil { |
| 294 | + break |
250 | 295 | } |
251 | 296 | } |
252 | 297 | c.StopRun() |
253 | 298 | } |
254 | 299 |
|
| 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 | + |
255 | 373 | // UpgradeHelmRelease upgrades an existing Helm release. |
256 | 374 | // @router /api/upgrade-helm-release [post] |
257 | 375 | func (c *ApiController) UpgradeHelmRelease() { |
|
0 commit comments