diff --git a/appview/pages/pages.go b/appview/pages/pages.go
index 5a625d3e..5ad48c66 100644
--- a/appview/pages/pages.go
+++ b/appview/pages/pages.go
@@ -1543,6 +1543,15 @@ func (p *Pages) LogLine(w io.Writer, params LogLineParams) error {
return p.executePlain("repo/pipelines/fragments/logLine", w, params)
}
+type WorkflowSymbolOOBParams struct {
+ Name string
+ Statuses models.WorkflowStatus
+}
+
+func (p *Pages) WorkflowSymbolOOB(w io.Writer, params WorkflowSymbolOOBParams) error {
+ return p.executePlain("repo/pipelines/fragments/workflowSymbolOOB", w, params)
+}
+
type WorkflowParams struct {
LoggedInUser *oauth.MultiAccountUser
RepoInfo repoinfo.RepoInfo
diff --git a/appview/pages/templates/repo/pipelines/fragments/workflowSymbolOOB.html b/appview/pages/templates/repo/pipelines/fragments/workflowSymbolOOB.html
new file mode 100644
index 00000000..baef7dea
--- /dev/null
+++ b/appview/pages/templates/repo/pipelines/fragments/workflowSymbolOOB.html
@@ -0,0 +1,5 @@
+{{ define "repo/pipelines/fragments/workflowSymbolOOB" }}
+
+ {{ template "repo/pipelines/fragments/workflowSymbol" .Statuses }}
+
+{{ end }}
diff --git a/appview/pages/templates/repo/pipelines/workflow.html b/appview/pages/templates/repo/pipelines/workflow.html
index 9600818f..a864a50f 100644
--- a/appview/pages/templates/repo/pipelines/workflow.html
+++ b/appview/pages/templates/repo/pipelines/workflow.html
@@ -49,7 +49,7 @@
{{ $kind := $lastStatus.Status.String }}
-
+
{{ template "repo/pipelines/fragments/workflowSymbol" $all }}
diff --git a/appview/pipelines/pipelines.go b/appview/pipelines/pipelines.go
index e1da9a1f..041474f0 100644
--- a/appview/pipelines/pipelines.go
+++ b/appview/pipelines/pipelines.go
@@ -29,15 +29,16 @@ import (
)
type Pipelines struct {
- repoResolver *reporesolver.RepoResolver
- idResolver *idresolver.Resolver
- config *config.Config
- oauth *oauth.OAuth
- pages *pages.Pages
- spindlestream *eventconsumer.Consumer
- db *db.DB
- enforcer *rbac.Enforcer
- logger *slog.Logger
+ repoResolver *reporesolver.RepoResolver
+ idResolver *idresolver.Resolver
+ config *config.Config
+ oauth *oauth.OAuth
+ pages *pages.Pages
+ spindlestream *eventconsumer.Consumer
+ pipelineNotifier *StatusNotifier
+ db *db.DB
+ enforcer *rbac.Enforcer
+ logger *slog.Logger
}
func (p *Pipelines) Router(mw *middleware.Middleware) http.Handler {
@@ -57,6 +58,7 @@ func New(
repoResolver *reporesolver.RepoResolver,
pages *pages.Pages,
spindlestream *eventconsumer.Consumer,
+ pipelineNotifier *StatusNotifier,
idResolver *idresolver.Resolver,
db *db.DB,
config *config.Config,
@@ -64,15 +66,16 @@ func New(
logger *slog.Logger,
) *Pipelines {
return &Pipelines{
- oauth: oauth,
- repoResolver: repoResolver,
- pages: pages,
- idResolver: idResolver,
- config: config,
- spindlestream: spindlestream,
- db: db,
- enforcer: enforcer,
- logger: logger,
+ oauth: oauth,
+ repoResolver: repoResolver,
+ pages: pages,
+ idResolver: idResolver,
+ config: config,
+ spindlestream: spindlestream,
+ pipelineNotifier: pipelineNotifier,
+ db: db,
+ enforcer: enforcer,
+ logger: logger,
}
}
@@ -212,6 +215,9 @@ func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) {
knot := f.Knot
rkey := singlePipeline.Rkey
+ statusCh := p.pipelineNotifier.Subscribe(singlePipeline.AtUri())
+ defer p.pipelineNotifier.Unsubscribe(singlePipeline.AtUri(), statusCh)
+
if spindle == "" || knot == "" || rkey == "" {
http.Error(w, "invalid repo info", http.StatusBadRequest)
return
@@ -329,6 +335,34 @@ func (p *Pipelines) Logs(w http.ResponseWriter, r *http.Request) {
return
}
+ case _, ok := <-statusCh:
+ if !ok {
+ continue
+ }
+ fresh, err := db.GetPipelineStatuses(
+ p.db,
+ 1,
+ orm.FilterEq("p.repo_did", f.RepoDid),
+ orm.FilterEq("p.id", pipelineId),
+ )
+ if err != nil || len(fresh) == 0 {
+ continue
+ }
+ for name, ws := range fresh[0].Statuses {
+ fragment.Reset()
+ if err = p.pages.WorkflowSymbolOOB(&fragment, pages.WorkflowSymbolOOBParams{
+ Name: name,
+ Statuses: ws,
+ }); err != nil {
+ l.Error("failed to render workflow symbol OOB", "err", err)
+ continue
+ }
+ if err = clientConn.WriteMessage(websocket.TextMessage, fragment.Bytes()); err != nil {
+ l.Error("error writing workflow symbol to client", "err", err)
+ return
+ }
+ }
+
case <-time.After(30 * time.Second):
l.Debug("sent keepalive")
if err = clientConn.WriteControl(websocket.PingMessage, []byte{}, time.Now().Add(time.Second)); err != nil {
diff --git a/appview/state/router.go b/appview/state/router.go
index 3d820a21..2fc0d8cc 100644
--- a/appview/state/router.go
+++ b/appview/state/router.go
@@ -394,6 +394,7 @@ func (s *State) PipelinesRouter(mw *middleware.Middleware) http.Handler {
s.repoResolver,
s.pages,
s.spindlestream,
+ s.pipelineNotifier,
s.idResolver,
s.db,
s.config,