From 63dadf62395245d4151e629f24ed5e21c0f0371d Mon Sep 17 00:00:00 2001 From: mrhid6 Date: Mon, 20 Jul 2026 15:18:35 +0100 Subject: [PATCH] feat(api): server-run log fetch and SSE stream endpoints --- server/internal/api/workflows.go | 112 +++++++++++++++++++++++++++++++ 1 file changed, 112 insertions(+) diff --git a/server/internal/api/workflows.go b/server/internal/api/workflows.go index ac761ac..de8780c 100644 --- a/server/internal/api/workflows.go +++ b/server/internal/api/workflows.go @@ -3,7 +3,11 @@ package api import ( "fmt" "net/http" + "os" + "regexp" "strconv" + "strings" + "time" "github.com/gin-gonic/gin" "github.com/mrhid6/vantage/server/internal/models" @@ -26,6 +30,114 @@ func registerWorkflowRoutes(g *gin.RouterGroup) { g.GET("/runs/:runId", getRun) g.POST("/runs/:runId/cancel", cancelRun) + g.GET("/runs/:runId/servers/:serverId/logs", getServerRunLog) + g.GET("/runs/:runId/servers/:serverId/logs/stream", streamServerRunLog) +} + +var uuidLike = regexp.MustCompile(`^[a-zA-Z0-9-]{1,64}$`) + +func getServerRunLog(c *gin.Context) { + runID, serverID := c.Param("runId"), c.Param("serverId") + if !uuidLike.MatchString(runID) || !uuidLike.MatchString(serverID) { + c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"}) + return + } + path := services.ServerRunLogPath(runID, serverID) + b, err := os.ReadFile(path) + if err != nil { + c.JSON(http.StatusNotFound, gin.H{"error": "no logs"}) + return + } + c.Data(http.StatusOK, "text/plain; charset=utf-8", b) +} + +func streamServerRunLog(c *gin.Context) { + runID, serverID := c.Param("runId"), c.Param("serverId") + if !uuidLike.MatchString(runID) || !uuidLike.MatchString(serverID) { + c.JSON(http.StatusBadRequest, gin.H{"error": "invalid id"}) + return + } + path := services.ServerRunLogPath(runID, serverID) + + c.Writer.Header().Set("Content-Type", "text/event-stream") + c.Writer.Header().Set("Cache-Control", "no-cache") + c.Writer.Header().Set("Connection", "keep-alive") + c.Writer.Header().Set("X-Accel-Buffering", "no") + + flusher, ok := c.Writer.(http.Flusher) + if !ok { + c.JSON(http.StatusInternalServerError, gin.H{"error": "stream unsupported"}) + return + } + + var offset int64 + sendNew := func() { + f, err := os.Open(path) + if err != nil { + return // file may not exist yet; keep waiting + } + defer f.Close() + if _, err := f.Seek(offset, 0); err != nil { + return + } + buf := make([]byte, 32*1024) + for { + n, _ := f.Read(buf) + if n <= 0 { + break + } + offset += int64(n) + // SSE data frame; split on newlines to keep frames well-formed. + for _, line := range splitSSE(buf[:n]) { + _, _ = c.Writer.WriteString("data: " + line + "\n") + } + _, _ = c.Writer.WriteString("\n") + flusher.Flush() + } + } + + ctx := c.Request.Context() + ticker := time.NewTicker(500 * time.Millisecond) + defer ticker.Stop() + for { + sendNew() + if serverRunTerminal(runID, serverID) { + sendNew() // final drain + _, _ = c.Writer.WriteString("event: done\ndata: end\n\n") + flusher.Flush() + return + } + select { + case <-ctx.Done(): + return + case <-ticker.C: + } + } +} + +// serverRunTerminal reports whether the given server-run has reached a terminal status. +func serverRunTerminal(runID, serverID string) bool { + r, err := services.GetRun(runID) + if err != nil { + return true + } + for _, sr := range r.ServerRuns { + if sr.ServerID == serverID { + switch sr.Status { + case "success", "failed", "skipped", "cancelled": + return true + } + return false + } + } + return true +} + +// splitSSE turns a raw byte slice into SSE-safe payload lines (newlines become +// separate data lines; carriage returns stripped). +func splitSSE(b []byte) []string { + s := strings.ReplaceAll(string(b), "\r", "") + return strings.Split(s, "\n") } func listSteps(c *gin.Context) {