Compare commits

..
44 Commits
Author SHA1 Message Date
mrhid6 b5f30bc7c8 Merge feat/step-picker-modal: autosave, step-picker modal, Steps page
Server Deploy / deploy (push) Successful in 1m22s
2026-07-21 11:51:03 +01:00
mrhid6 3a0116248e fix(web): guard autosave against in-flight lost-update race 2026-07-21 11:50:53 +01:00
mrhid6 6af0a88841 feat(web): add Steps to main nav 2026-07-21 11:45:43 +01:00
mrhid6 e4c3fc24d3 feat(web): standalone Steps management page 2026-07-21 11:43:42 +01:00
mrhid6 e46d0edbf2 feat(web): replace builder step sidebar with Add-step modal 2026-07-21 11:41:37 +01:00
mrhid6 a4c4a72dbc feat(web): StepPickerModal step picker component 2026-07-21 11:37:47 +01:00
mrhid6 67d729b360 feat(web): api.stepUsage binding 2026-07-21 11:36:03 +01:00
mrhid6 da6d825f45 fix(server): TestMain must not exit(0) before m.Run(); skip DB test individually 2026-07-21 11:34:56 +01:00
mrhid6 93423e32e6 feat(server): step usage counts endpoint 2026-07-21 11:31:05 +01:00
mrhid6 d0442291f5 feat(web): workflow builder autosave with last-saved status 2026-07-21 11:27:55 +01:00
mrhid6 6c5472760b Merge branch 'feat/adhoc-steps-import-export': ad-hoc steps, step import/export, default steps, auto-derived outputs
Server Deploy / deploy (push) Successful in 1m26s
2026-07-21 10:42:02 +01:00
mrhid6 7c4a676742 fix: inline step secrets + inline input/output display + import body limit 2026-07-21 10:40:45 +01:00
mrhid6 fbda26a188 feat(web): ad-hoc inline steps + import-to-inline in workflow editor 2026-07-21 10:34:51 +01:00
mrhid6 90ce7af769 feat(web): step export/import, sync defaults, auto outputs, default badge 2026-07-21 10:31:28 +01:00
mrhid6 b543cd1b3d feat(web): api client for step import/export/inline/defaults 2026-07-21 10:28:41 +01:00
mrhid6 15c9da1b01 feat(server): default steps seed-on-boot + admin re-sync 2026-07-21 10:25:41 +01:00
mrhid6 baa7bb239d feat(server): step import/export/parse endpoints 2026-07-21 10:22:35 +01:00
mrhid6 7342c46d99 feat(server): auto-derive outputs on save + resolve inline steps 2026-07-21 10:20:12 +01:00
mrhid6 813f9e6fef feat(server): derive declared_outputs from script + slugify 2026-07-21 10:17:50 +01:00
mrhid6 434f14ae3a feat(server): inline step ref + workflow validation 2026-07-21 10:15:52 +01:00
mrhid6 8398fd2279 docs: implementation plan for adhoc steps, import/export, defaults, auto-outputs 2026-07-21 10:11:16 +01:00
mrhid6 56f06b9eaf docs: auto-derive declared_outputs from script scan 2026-07-21 10:05:45 +01:00
mrhid6 d9d241f83b docs: design for adhoc steps, step import/export, default steps 2026-07-21 09:59:34 +01:00
mrhid6 aee910c1f8 fix: fixed variable inputs
Server Deploy / deploy (push) Successful in 1m22s
2026-07-20 18:00:47 +01:00
mrhid6 bea545e873 feat: More verbose logging on workflow logs
Agent Release / build (push) Successful in 45s
Agent Release / msi (push) Successful in 42s
Server Deploy / deploy (push) Successful in 1m52s
2026-07-20 17:35:58 +01:00
mrhid6 82d7dde5f8 fix: Fixed style on workflow run
Server Deploy / deploy (push) Successful in 42s
2026-07-20 16:44:53 +01:00
mrhid6 397016ad68 feat: Updated workflow runs page
Server Deploy / deploy (push) Successful in 1m20s
2026-07-20 16:00:41 +01:00
mrhid6 39348c9491 feat(web): live SSE log tail and log retention setting
Server Deploy / deploy (push) Successful in 1m35s
2026-07-20 15:18:36 +01:00
mrhid6 63dadf6239 feat(api): server-run log fetch and SSE stream endpoints 2026-07-20 15:18:35 +01:00
mrhid6 d905c99d32 feat(server): stream step logs to files, drop log bodies from run docs 2026-07-20 15:18:35 +01:00
mrhid6 85e1baf59a feat(server): workflow log-writer registry, retention setting, sweeper 2026-07-20 15:18:35 +01:00
mrhid6 351ad59dd8 fix(web): reseed edit-workflow modal state on open
Server Deploy / deploy (push) Successful in 1m23s
2026-07-20 14:56:14 +01:00
mrhid6 dcc901b0d2 feat(web): workflow runs list page and navigation links
Server Deploy / deploy (push) Successful in 1m29s
2026-07-20 14:52:13 +01:00
mrhid6 99bf093f00 feat(web): rebuild workflow builder — mockup styling, drag-and-drop, inputs inspector 2026-07-20 14:48:29 +01:00
mrhid6 619ccd28cb feat(web): edit-workflow modal (name/targets/delete) 2026-07-20 14:44:59 +01:00
mrhid6 f22f0a4729 feat(web): edit-base-step modal with inputs/outputs editor 2026-07-20 14:42:13 +01:00
mrhid6 78194daf5f feat(web): builder tokens, Modal primitive, input-param types 2026-07-20 14:40:07 +01:00
mrhid6 f141767fc2 feat(workflows): inject step inputs into env; update returns workflow 2026-07-20 14:37:36 +01:00
mrhid6 05cd8e154b feat(workflows): step input params model + cascade step delete 2026-07-20 14:34:50 +01:00
mrhid6 004cc03ba6 docs: add workflow builder v2 plan 2026-07-20 14:33:46 +01:00
mrhid6 236e89989f docs: add workflow builder v2 spec 2026-07-20 14:30:55 +01:00
mrhid6 b0a2de8ca1 fix: Fixed workflow style topbar
Server Deploy / deploy (push) Successful in 1m20s
2026-07-20 13:55:17 +01:00
mrhid6 e9ac7be8c3 feat: updates to workflow page
Server Deploy / deploy (push) Successful in 36s
2026-07-20 13:35:35 +01:00
mrhid6 47690c58d9 fix: Fixed workflow id
Server Deploy / deploy (push) Successful in 1m14s
2026-07-20 12:49:43 +01:00
50 changed files with 4643 additions and 5392 deletions
+25
View File
@@ -35,10 +35,21 @@ func (w *streamWriter) Write(p []byte) (int, error) {
return len(p), nil
}
// WorkspacePath returns the per-run working directory for a workspace id. The
// same id always maps to the same path so RunStep and the cleanup command agree.
func WorkspacePath(workspaceID string) string {
return filepath.Join(os.TempDir(), "vantage-run-"+workspaceID)
}
// RunStep writes the script to a temp file, provides a WORKFLOW_ENV file for
// the script to append KEY=value output to, executes it under the requested
// interpreter, and streams output via emit, returning the terminal result
// with empty stdout/stderr but populated exit_code/output_env.
//
// When the command carries a WorkspaceId the step runs with that per-run working
// directory as its cwd (created here if missing); the server removes it once the
// run finishes. The script and env files always live in a private temp dir so
// they never leak into the shared workspace.
func RunStep(cmd *pb.RunStepCmd, emit func(seq uint64, data []byte)) *pb.StepResult {
res := &pb.StepResult{CommandId: "", OutputEnv: map[string]string{}}
@@ -50,6 +61,16 @@ func RunStep(cmd *pb.RunStepCmd, emit func(seq uint64, data []byte)) *pb.StepRes
}
defer os.RemoveAll(dir)
workDir := ""
if cmd.WorkspaceId != "" {
workDir = WorkspacePath(cmd.WorkspaceId)
if err := os.MkdirAll(workDir, 0700); err != nil {
res.ExitCode = 1
res.Stderr = "create workspace: " + err.Error()
return res
}
}
envFile := filepath.Join(dir, "workflow_env")
if err := os.WriteFile(envFile, nil, 0600); err != nil {
res.ExitCode = 1
@@ -91,6 +112,10 @@ func RunStep(cmd *pb.RunStepCmd, emit func(seq uint64, data []byte)) *pb.StepRes
c = exec.CommandContext(ctx, "bash", scriptPath)
}
if workDir != "" {
c.Dir = workDir
}
c.Env = append(os.Environ(), "WORKFLOW_ENV="+envFile)
for k, v := range cmd.Env {
c.Env = append(c.Env, k+"="+v)
+16 -6
View File
@@ -63,12 +63,19 @@ type ReportUpdatesResponse struct{}
type ApplyUpdatesCmd struct{}
type ServerCommand struct {
CommandId string `json:"command_id"`
GenerateKey *GenerateKeyCmd `json:"generate_key,omitempty"`
DeleteKey *DeleteKeyCmd `json:"delete_key,omitempty"`
UpdateAgent *UpdateAgentCmd `json:"update_agent,omitempty"`
ApplyUpdates *ApplyUpdatesCmd `json:"apply_updates,omitempty"`
RunStep *RunStepCmd `json:"run_step,omitempty"`
CommandId string `json:"command_id"`
GenerateKey *GenerateKeyCmd `json:"generate_key,omitempty"`
DeleteKey *DeleteKeyCmd `json:"delete_key,omitempty"`
UpdateAgent *UpdateAgentCmd `json:"update_agent,omitempty"`
ApplyUpdates *ApplyUpdatesCmd `json:"apply_updates,omitempty"`
RunStep *RunStepCmd `json:"run_step,omitempty"`
CleanupWorkspace *CleanupWorkspaceCmd `json:"cleanup_workspace,omitempty"`
}
// CleanupWorkspaceCmd tells the agent to recursively remove the run's working
// directory once all steps on that server have finished.
type CleanupWorkspaceCmd struct {
WorkspaceId string `json:"workspace_id"`
}
type DeleteKeyCmd struct {
@@ -110,6 +117,9 @@ type RunStepCmd struct {
Script string `json:"script"`
Env map[string]string `json:"env,omitempty"`
TimeoutSeconds int `json:"timeout_seconds,omitempty"`
// WorkspaceId names the per-run working directory the agent creates and uses
// as the step's cwd. Empty means run in the agent's default directory.
WorkspaceId string `json:"workspace_id,omitempty"`
}
type StepResult struct {
+13
View File
@@ -198,6 +198,9 @@ func connectAndHandleStream(ctx context.Context, cfg *config.Config) error {
if cmd.ApplyUpdates != nil {
go handleApplyUpdates(cfg, cmd)
}
if cmd.CleanupWorkspace != nil {
go handleCleanupWorkspace(cmd)
}
if cmd.RunStep != nil {
go func(rc *pb.RunStepCmd, cid string) {
emit := func(seq uint64, data []byte) {
@@ -286,6 +289,16 @@ func handleApplyUpdates(cfg *config.Config, cmd *pb.ServerCommand) {
_ = client.ReportUpdates(cfg.ServerID, cfg.AgentToken, nil)
}
func handleCleanupWorkspace(cmd *pb.ServerCommand) {
id := cmd.CleanupWorkspace.WorkspaceId
dir := agentexec.WorkspacePath(id)
if err := os.RemoveAll(dir); err != nil {
log.Printf("cleanup workspace %s failed (cmd=%s): %v", dir, cmd.CommandId, err)
return
}
log.Printf("removed run workspace %s (cmd=%s)", dir, cmd.CommandId)
}
func handleDeleteKey(cmd *pb.ServerCommand) {
label := cmd.DeleteKey.Label
keyPath := fmt.Sprintf("/root/.ssh/vantage_%s", strings.ReplaceAll(label, " ", "_"))
File diff suppressed because it is too large Load Diff
File diff suppressed because it is too large Load Diff
@@ -1,887 +0,0 @@
# Workflow Log Streaming Implementation Plan
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
**Goal:** Stream workflow step output live from agents to per-server-run log files on the server, tail them live in the UI over SSE, and auto-expire them on a configurable retention period.
**Architecture:** Agent streams interleaved stdout/stderr chunks over the existing `CommandStream` (`AgentMessage.StepOutput`). Server appends secret-masked chunks to `<logdir>/<run_id>/<server_id>.log` via a per-command log-writer registry, records a per-step byte offset, and stops persisting log bodies in Mongo. UI tails via an SSE endpoint while running and fetches the whole file after. An hourly sweeper deletes run-log dirs older than the retention setting.
**Tech Stack:** Go (gin, mongo-driver v2), hand-written JSON-codec gRPC structs (no protoc), Next.js 16 app-router + react-query + EventSource, MongoDB, local filesystem for logs.
## Global Constraints
- **No tests this iteration** — do not write `*_test.go` or frontend tests. Verify each task with `go build ./...`, `go vet ./...`, and (frontend) `npm run build`.
- gRPC uses a **JSON codec** — proto messages are hand-written Go structs in **two** files that must stay identical: `server/internal/grpc/pb/vantage.pb.go` and `agent/internal/grpc/pb/vantage.pb.go`. No codegen. Also update `proto/vantage/v1/vantage.proto` as documentation.
- Mongo access pattern: `db.Col("collection_name")` with `context.WithTimeout`. Follow `server/internal/services/workflows.go`.
- Secret values must never be written into log files unmasked — mask by literal `***` replacement at write time, boundary-safe via a carry buffer.
- Interpreter values are the literals `"bash"` and `"powershell"`.
- Go module path: `github.com/mrhid6/vantage`.
- Log dir from env `VANTAGE_WORKFLOW_LOG_DIR`, default `<data>/workflow-logs`; files `0600`, dirs `0700`.
- Retention default **30** days, stored `settings.workflow_log_retention_days`; `0`/negative = keep forever.
- The agent's stream `Send` is only safe through the existing per-connection mutex-guarded `send()` closure in `connectAndHandleStream` — all `StepOutput`/`StepResult` sends MUST go through it.
---
## Task 1: Proto/pb — StepOutputChunk
**Files:**
- Modify: `proto/vantage/v1/vantage.proto`
- Modify: `server/internal/grpc/pb/vantage.pb.go`
- Modify: `agent/internal/grpc/pb/vantage.pb.go`
**Interfaces:**
- Produces: `pb.StepOutputChunk{CommandId string, Seq uint64, Data []byte, Eof bool}`; `pb.AgentMessage` gains `StepOutput *StepOutputChunk`.
- [ ] **Step 1: Document in the proto file**
In `proto/vantage/v1/vantage.proto`, add to the `AgentMessage` oneof: `StepOutputChunk step_output = 6;` and add the message:
```protobuf
message StepOutputChunk {
string command_id = 1;
uint64 seq = 2;
bytes data = 3;
bool eof = 4;
}
```
- [ ] **Step 2: Add struct + field to server pb file**
In `server/internal/grpc/pb/vantage.pb.go`, add to `type AgentMessage struct { ... }`:
```go
StepOutput *StepOutputChunk `json:"step_output,omitempty"`
```
and add the new struct:
```go
type StepOutputChunk struct {
CommandId string `json:"command_id"`
Seq uint64 `json:"seq"`
Data []byte `json:"data,omitempty"`
Eof bool `json:"eof,omitempty"`
}
```
- [ ] **Step 3: Mirror identical additions into the agent pb file**
Apply the identical `AgentMessage.StepOutput` field and `StepOutputChunk` struct to `agent/internal/grpc/pb/vantage.pb.go`.
- [ ] **Step 4: Verify build**
Run: `cd server && go build ./... && cd ../agent && go build ./...`
Expected: both succeed.
- [ ] **Step 5: Commit**
```bash
git add proto/vantage/v1/vantage.proto server/internal/grpc/pb/vantage.pb.go agent/internal/grpc/pb/vantage.pb.go
git commit -m "feat(proto): add StepOutputChunk streaming message"
```
---
## Task 2: Agent — stream step output
**Files:**
- Modify: `agent/internal/exec/exec.go`
- Modify: `agent/internal/sync/sync.go` (the `cmd.RunStep != nil` goroutine)
**Interfaces:**
- Consumes: `pb.RunStepCmd`, `pb.StepResult`, `pb.StepOutputChunk` (Task 1).
- Produces: `exec.RunStep(cmd *pb.RunStepCmd, emit func(seq uint64, data []byte)) *pb.StepResult` — streams output via `emit`, returns terminal result with empty stdout/stderr but populated exit_code/output_env.
- [ ] **Step 1: Rework `exec.RunStep` to stream**
In `agent/internal/exec/exec.go`, change the signature and replace the two `bytes.Buffer`s with a single mutex-guarded streaming writer. Full new body of the run/capture section (keep the existing temp-dir, env-file, interpreter-selection, timeout, and `parseEnvFile` logic exactly as-is):
Add this type at package scope:
```go
// streamWriter forwards every write to emit() as an ordered chunk. Used as both
// Stdout and Stderr so output interleaves in real execution order. The mutex
// ensures a single stdout/stderr write is not interleaved mid-slice with another.
type streamWriter struct {
mu sync.Mutex
seq uint64
emit func(seq uint64, data []byte)
}
func (w *streamWriter) Write(p []byte) (int, error) {
w.mu.Lock()
defer w.mu.Unlock()
if w.emit != nil {
buf := make([]byte, len(p))
copy(buf, p)
w.emit(w.seq, buf)
w.seq++
}
return len(p), nil
}
```
Add `"sync"` to the imports. Change the signature to:
```go
func RunStep(cmd *pb.RunStepCmd, emit func(seq uint64, data []byte)) *pb.StepResult {
```
Replace the block that currently declares `var stdout, stderr bytes.Buffer`, assigns `c.Stdout`/`c.Stderr`, and sets `res.Stdout`/`res.Stderr` from them, with:
```go
sw := &streamWriter{emit: emit}
c.Stdout = sw
c.Stderr = sw
runErr := c.Run()
// stdout/stderr are streamed via emit, not returned in the result.
if ctx.Err() == context.DeadlineExceeded {
res.ExitCode = 124
res.Stderr = "[vantage] step timed out"
} else if ee, ok := runErr.(*exec.ExitError); ok {
res.ExitCode = ee.ExitCode()
} else if runErr != nil {
res.ExitCode = 1
res.Stderr = "[vantage] " + runErr.Error()
}
res.OutputEnv = parseEnvFile(envFile)
return res
```
Remove the now-unused `"bytes"` and `"bufio"` imports **only if** they are no longer referenced (`parseEnvFile` uses `bufio` + `os` — keep `bufio`; `bytes` is likely now unused — remove it if so). Verify with `go build`.
- [ ] **Step 2: Wire streaming into the agent loop**
In `agent/internal/sync/sync.go`, the `cmd.RunStep != nil` goroutine currently calls `agentexec.RunStep(rc)` and sends one `StepResult` via `send()`. Change it to pass an `emit` closure that streams chunks, then send an eof chunk, then the terminal result — all through the existing mutex-guarded `send()`:
```go
if cmd.RunStep != nil {
go func(rc *pb.RunStepCmd, cid string) {
emit := func(seq uint64, data []byte) {
_ = send(&pb.AgentMessage{
ServerId: cfg.ServerID,
AgentToken: cfg.AgentToken,
StepOutput: &pb.StepOutputChunk{CommandId: cid, Seq: seq, Data: data},
})
}
res := agentexec.RunStep(rc, emit)
res.CommandId = cid
// Final eof marker so the server closes the log file.
_ = send(&pb.AgentMessage{
ServerId: cfg.ServerID,
AgentToken: cfg.AgentToken,
StepOutput: &pb.StepOutputChunk{CommandId: cid, Eof: true},
})
_ = send(&pb.AgentMessage{
ServerId: cfg.ServerID,
AgentToken: cfg.AgentToken,
StepResult: res,
})
}(cmd.RunStep, cmd.CommandId)
continue
}
```
(Match the exact field names already used by the existing `send()` calls in this function — `cfg.ServerID`, `cfg.AgentToken`, and the `send` closure. If the existing RunStep branch used different local names, keep those.)
- [ ] **Step 3: Verify build**
Run: `cd agent && go build ./... && go vet ./...`
Expected: success. Resolve any leftover unused-import error from Step 1.
- [ ] **Step 4: Commit**
```bash
git add agent/internal/exec/exec.go agent/internal/sync/sync.go
git commit -m "feat(agent): stream step output chunks over CommandStream"
```
---
## Task 3: Server log-writer registry + retention sweeper
**Files:**
- Create: `server/internal/services/steplogs.go`
**Interfaces:**
- Consumes: `settings` service (retention), `db.Col("workflow_runs")` (sweeper), env `VANTAGE_WORKFLOW_LOG_DIR`.
- Produces:
- `WorkflowLogDir() string` — resolved base dir (env or default), created on first call.
- `ServerRunLogPath(runID, serverID string) string``<logdir>/<runID>/<serverID>.log`.
- `AppendMarker(runID, serverID, line string) (int64, error)` — appends a marker line, returns the byte offset **before** the write (the step's `log_offset`).
- `var StepLogs *stepLogRegistry` with `Open(commandID, path string, secrets []string) error`, `Append(commandID string, data []byte)`, `Close(commandID string)`.
- `StartLogSweeper()` — launches the hourly retention goroutine; also sweeps once immediately.
- [ ] **Step 1: Write the registry, paths, and sweeper**
```go
package services
import (
"os"
"path/filepath"
"strings"
"sync"
"time"
"github.com/mrhid6/vantage/server/internal/db"
"go.mongodb.org/mongo-driver/v2/bson"
)
// WorkflowLogDir returns the base directory for workflow step logs, creating it.
func WorkflowLogDir() string {
dir := os.Getenv("VANTAGE_WORKFLOW_LOG_DIR")
if dir == "" {
dir = filepath.Join("data", "workflow-logs")
}
_ = os.MkdirAll(dir, 0700)
return dir
}
// ServerRunLogPath is the per-server-run log file path.
func ServerRunLogPath(runID, serverID string) string {
return filepath.Join(WorkflowLogDir(), runID, serverID+".log")
}
// AppendMarker appends a line to the server-run log and returns the byte offset
// at which the write began (used as a step's log_offset).
func AppendMarker(runID, serverID, line string) (int64, error) {
path := ServerRunLogPath(runID, serverID)
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return 0, err
}
f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
return 0, err
}
defer f.Close()
off, _ := f.Seek(0, 2) // current end = offset before write
if _, err := f.WriteString(line); err != nil {
return off, err
}
return off, nil
}
// ---- streamed chunk writer, boundary-safe secret masking ----
type stepLogWriter struct {
mu sync.Mutex
f *os.File
carry []byte
secrets []string
maxSecret int
}
type stepLogRegistry struct {
mu sync.Mutex
writers map[string]*stepLogWriter
}
var StepLogs = &stepLogRegistry{writers: make(map[string]*stepLogWriter)}
// Open opens (append) the server-run file for a step's streamed chunks.
func (r *stepLogRegistry) Open(commandID, path string, secrets []string) error {
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return err
}
f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
return err
}
max := 0
for _, s := range secrets {
if len(s) > max {
max = len(s)
}
}
w := &stepLogWriter{f: f, secrets: secrets, maxSecret: max}
r.mu.Lock()
r.writers[commandID] = w
r.mu.Unlock()
return nil
}
func (r *stepLogRegistry) get(commandID string) *stepLogWriter {
r.mu.Lock()
defer r.mu.Unlock()
return r.writers[commandID]
}
// Append masks and writes a chunk, holding back the last maxSecret-1 bytes so a
// secret split across a chunk boundary is still masked on the next append/close.
func (r *stepLogRegistry) Append(commandID string, data []byte) {
w := r.get(commandID)
if w == nil {
return
}
w.mu.Lock()
defer w.mu.Unlock()
if len(w.secrets) == 0 || w.maxSecret <= 1 {
_, _ = w.f.Write(data)
return
}
buf := append(w.carry, data...)
hold := w.maxSecret - 1
if len(buf) <= hold {
w.carry = buf
return
}
flush := buf[:len(buf)-hold]
w.carry = append([]byte{}, buf[len(buf)-hold:]...)
_, _ = w.f.Write(maskBytes(flush, w.secrets))
}
// Close flushes the carry (masked) and closes the file.
func (r *stepLogRegistry) Close(commandID string) {
r.mu.Lock()
w := r.writers[commandID]
delete(r.writers, commandID)
r.mu.Unlock()
if w == nil {
return
}
w.mu.Lock()
defer w.mu.Unlock()
if len(w.carry) > 0 {
_, _ = w.f.Write(maskBytes(w.carry, w.secrets))
w.carry = nil
}
_ = w.f.Close()
}
func maskBytes(b []byte, secrets []string) []byte {
s := string(b)
for _, v := range secrets {
if v == "" {
continue
}
s = strings.ReplaceAll(s, v, "***")
}
return []byte(s)
}
// ---- retention sweeper ----
// StartLogSweeper sweeps expired run-log dirs hourly (and once now).
func StartLogSweeper() {
go func() {
sweepLogs()
t := time.NewTicker(time.Hour)
defer t.Stop()
for range t.C {
sweepLogs()
}
}()
}
func sweepLogs() {
days := retentionDays()
if days <= 0 {
return
}
cutoff := time.Now().AddDate(0, 0, -days)
base := WorkflowLogDir()
entries, err := os.ReadDir(base)
if err != nil {
return
}
for _, e := range entries {
if !e.IsDir() {
continue
}
runID := e.Name()
dir := filepath.Join(base, runID)
if runExpired(runID, dir, cutoff) {
_ = os.RemoveAll(dir)
}
}
}
// runExpired is true when the run finished before cutoff (falling back to dir
// mtime when the run doc is gone).
func runExpired(runID, dir string, cutoff time.Time) bool {
ctx, cancel := wfCtx()
defer cancel()
var run struct {
FinishedAt *time.Time `bson:"finished_at"`
}
err := db.Col("workflow_runs").FindOne(ctx, bson.M{"run_id": runID}).Decode(&run)
if err == nil {
if run.FinishedAt == nil {
return false // still running / never finished — keep
}
return run.FinishedAt.Before(cutoff)
}
// run doc gone: use dir mtime
if fi, e := os.Stat(dir); e == nil {
return fi.ModTime().Before(cutoff)
}
return false
}
func retentionDays() int {
if v, err := GetWorkflowLogRetentionDays(); err == nil {
return v
}
return 30
}
```
Note: `wfCtx` is defined in `workflows.go` (same package) — reuse it. `GetWorkflowLogRetentionDays` is added in Task 4; this file references it (same package, compiles together).
- [ ] **Step 2: Verify build**
Run: `cd server && go build ./... && go vet ./...`
Expected: FAIL — `GetWorkflowLogRetentionDays` undefined until Task 4. This is expected; proceed to commit the file so Task 4 completes it. (If you prefer a green build, do Task 4's settings accessor first, then return — but committing here is fine since Task 4 immediately follows.)
Actually to keep every commit buildable: **temporarily** add a local stub at the bottom of this file and remove it in Task 4:
```go
// TEMP stub, replaced in Task 4.
func GetWorkflowLogRetentionDays() (int, error) { return 30, nil }
```
Then `cd server && go build ./... && go vet ./...` must succeed.
- [ ] **Step 3: Commit**
```bash
git add server/internal/services/steplogs.go
git commit -m "feat(server): workflow log-writer registry, paths, retention sweeper"
```
---
## Task 4: Settings — retention accessor + startup wiring
**Files:**
- Modify: `server/internal/services/settings.go` (or wherever settings get/set lives — search `settings` collection usage)
- Modify: `server/internal/services/steplogs.go` (remove the temp stub)
- Modify: `server/cmd/main.go` (start the sweeper)
**Interfaces:**
- Produces: `GetWorkflowLogRetentionDays() (int, error)` (default 30 when unset), `SetWorkflowLogRetentionDays(int) error`. If settings are exposed as a single document/struct, add the field there and derive these accessors.
- [ ] **Step 1: Inspect the settings service**
Read the existing settings service (search for the `settings` collection: `grep -rn "\"settings\"" server/internal/services`). Determine whether settings are a typed struct document or key/value. Match that pattern.
- [ ] **Step 2: Add the retention accessor**
If settings are a **typed document** (e.g. a `GetSettings()/UpdateSettings()`), add a field `WorkflowLogRetentionDays int `bson:"workflow_log_retention_days" json:"workflow_log_retention_days"`` to the settings struct and implement:
```go
func GetWorkflowLogRetentionDays() (int, error) {
s, err := GetSettings() // use the real accessor name
if err != nil {
return 30, err
}
if s.WorkflowLogRetentionDays == 0 && /* unset sentinel */ !s.WorkflowLogRetentionSet {
return 30, nil
}
return s.WorkflowLogRetentionDays, nil
}
```
Simplify to match reality: if the settings doc uses zero-value-means-unset and you cannot distinguish "0 = keep forever" from "unset", store the retention as a pointer `*int` or default at read: **treat a missing field as 30, an explicit 0 as keep-forever.** Prefer `*int` in the struct so the three states (unset→30, 0→forever, N→N) are representable. Implement `GetWorkflowLogRetentionDays` to return 30 when the pointer is nil, else its value. `SetWorkflowLogRetentionDays(n int)` sets the pointer.
If settings are **key/value**, implement both accessors against that store with the same nil→30 / 0→forever semantics (store empty/absent = 30).
- [ ] **Step 3: Remove the temp stub from `steplogs.go`**
Delete the `// TEMP stub` `GetWorkflowLogRetentionDays` added in Task 3 so the real one is used.
- [ ] **Step 4: Start the sweeper at boot**
In `server/cmd/main.go`, next to `EnsureWorkflowIndexes()`, add `services.StartLogSweeper()`.
- [ ] **Step 5: Verify build**
Run: `cd server && go build ./... && go vet ./...`
Expected: success (real accessor now resolves the reference from Task 3).
- [ ] **Step 6: Commit**
```bash
git add server/internal/services/settings.go server/internal/services/steplogs.go server/cmd/main.go
git commit -m "feat(server): workflow log retention setting + sweeper startup"
```
---
## Task 5: Runner + model — write to files, drop log bodies from Mongo
**Files:**
- Modify: `server/internal/models/workflow.go` (`StepRun`)
- Modify: `server/internal/services/workflow_runner.go`
- Modify: `server/internal/grpc/server.go` (stream delivery of `StepOutput`)
**Interfaces:**
- Consumes: `StepLogs`, `AppendMarker`, `ServerRunLogPath` (Task 3), `pb.StepOutputChunk` (Task 1).
- Produces: runner writes markers + streams chunks to files; `StepRun.LogOffset` persisted; `StepRun.Stdout/Stderr` removed.
- [ ] **Step 1: Update the `StepRun` model**
In `server/internal/models/workflow.go`, in `type StepRun struct`:
- Remove the `Stdout` and `Stderr` fields.
- Add: `LogOffset int64 `bson:"log_offset" json:"log_offset"``
- [ ] **Step 2: Deliver StepOutput chunks in the gRPC receive loop**
In `server/internal/grpc/server.go`, after the existing `if m.StepResult != nil { services.StepResults.Deliver(m.StepResult) }` block, add:
```go
if m.StepOutput != nil {
if m.StepOutput.Eof {
services.StepLogs.Close(m.StepOutput.CommandId)
} else {
services.StepLogs.Append(m.StepOutput.CommandId, m.StepOutput.Data)
}
}
```
- [ ] **Step 3: Rework `runServer` to open logs + write markers, drop persisted bodies**
In `server/internal/services/workflow_runner.go`, `runServer`:
Inside the per-step loop, **before** `dispatchAndWait`, add marker + open (compute `secretVals` first, which already exists in the loop):
```go
// Write the step marker and remember the offset for later slicing.
marker := fmt.Sprintf("\n===== step %d: %s =====\n", step.Order, step.Name)
offset, _ := AppendMarker(runID, serverID, marker)
logPath := ServerRunLogPath(runID, serverID)
_ = StepLogs.Open(commandID_placeholder, logPath, secretsSlice(secretVals))
```
There is a chicken-and-egg with `commandID`: today `dispatchAndWait` generates the `commandID` internally. Refactor so the runner owns the `commandID`:
1. Change `dispatchAndWait(serverID string, cmd *pb.RunStepCmd)` to `dispatchAndWait(serverID, commandID string, cmd *pb.RunStepCmd)` and remove its internal `commandID := uuid.New().String()` (use the passed one).
2. In `runServer`, generate `commandID := uuid.New().String()` at the top of each attempt-group (before the marker/open), open the log with it, then call `dispatchAndWait(serverID, commandID, cmd)`.
3. After the step completes (result received), call `StepLogs.Close(commandID)` defensively (idempotent — the agent's eof usually closed it already; Close on a missing key is a no-op).
Add a helper to convert the `secretVals map[string]string` to a `[]string` of values:
```go
func secretsSlice(m map[string]string) []string {
out := make([]string, 0, len(m))
for _, v := range m {
out = append(out, v)
}
return out
}
```
Update `finishStep(...)` call + signature: **remove** the `stdout, stderr string` params and the `output_env` masking stays. Persist `log_offset` instead. New `finishStep`:
```go
func finishStep(runID, serverID string, order int, status string, attempts, exit int, logOffset int64, outEnv map[string]string) {
now := time.Now()
updateStep(runID, serverID, order, bson.M{
"server_runs.$[s].steps.$[t].status": status,
"server_runs.$[s].steps.$[t].attempts": attempts,
"server_runs.$[s].steps.$[t].exit_code": exit,
"server_runs.$[s].steps.$[t].log_offset": logOffset,
"server_runs.$[s].steps.$[t].output_env": outEnv,
"server_runs.$[s].steps.$[t].finished_at": now,
})
}
```
In the loop, after receiving `res`, drop the `stdout, stderr := ...` masking of `res.Stdout/res.Stderr` (those are now streamed to file). Keep the `outEnv` build **with existing masking** (`maskSecrets(v, allSecrets)` per the merged secret-leak fix) — `output_env`/`run_env` masking is unchanged. Call:
```go
finishStep(runID, serverID, i, status, attempts, exit, offset, outEnv)
```
where `offset` is the marker offset captured before dispatch. If `res == nil`, still write a short note to the file so failures are visible:
```go
if res == nil {
_, _ = AppendMarker(runID, serverID, "[vantage] agent did not return a result\n")
}
```
Remove the initial `StepRun{... Status:"queued"}` `Stdout/Stderr` references if any (the model no longer has them — the queued StepRun in `TriggerWorkflow` set only `Order/Name/Status/OutputEnv`, so no change needed there; verify).
Ensure `fmt` is imported (it already is).
- [ ] **Step 4: Verify build**
Run: `cd server && go build ./... && go vet ./...`
Expected: success. Fix any remaining references to the removed `Stdout`/`Stderr` fields or the old `finishStep`/`dispatchAndWait` signatures.
- [ ] **Step 5: Commit**
```bash
git add server/internal/models/workflow.go server/internal/services/workflow_runner.go server/internal/grpc/server.go
git commit -m "feat(server): stream step logs to files, drop log bodies from run docs"
```
---
## Task 6: REST — log fetch + SSE stream endpoints
**Files:**
- Modify: `server/internal/api/workflows.go`
**Interfaces:**
- Consumes: `ServerRunLogPath`, `GetRun` (existing).
- Produces: `GET /api/runs/:runId/servers/:serverId/logs` and `GET /api/runs/:runId/servers/:serverId/logs/stream` (SSE).
- [ ] **Step 1: Add the two handlers + routes**
In `registerWorkflowRoutes`, add:
```go
g.GET("/runs/:runId/servers/:serverId/logs", getServerRunLog)
g.GET("/runs/:runId/servers/:serverId/logs/stream", streamServerRunLog)
```
Add a UUID-ish validator and the handlers:
```go
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() bool {
f, err := os.Open(path)
if err != nil {
return true // file may not exist yet; keep waiting
}
defer f.Close()
if _, err := f.Seek(offset, 0); err != nil {
return true
}
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()
}
return true
}
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")
}
```
Add imports: `"os"`, `"regexp"`, `"strings"`, `"time"`, `"net/http"` (already present). Confirm `services.GetRun` and `ServerRun.Status`/`ServerID` fields exist (they do from the Workflows feature).
- [ ] **Step 2: Verify build**
Run: `cd server && go build ./... && go vet ./...`
Expected: success.
- [ ] **Step 3: Commit**
```bash
git add server/internal/api/workflows.go
git commit -m "feat(api): server-run log fetch and SSE stream endpoints"
```
---
## Task 7: Frontend — live SSE tail + retention setting
**Files:**
- Modify: `web/lib/api.ts`
- Modify: `web/app/workflows/[id]/runs/[runId]/page.tsx`
- Modify: `web/app/settings/page.tsx`
**Interfaces:**
- Consumes: SSE endpoint, logs endpoint, settings mutation.
- [ ] **Step 1: Update API types + helpers**
In `web/lib/api.ts`:
- In `StepRun`, remove `stdout` and `stderr`; add `log_offset: number`.
- Add: `getServerRunLog: (runId: string, serverId: string) => request<string>(...)` — but the logs endpoint returns `text/plain`, so add a dedicated fetch that reads text. If `request<T>` assumes JSON, add a sibling:
```ts
async getServerRunLog(runId: string, serverId: string): Promise<string> {
const res = await fetch(`${API_BASE}/api/runs/${runId}/servers/${serverId}/logs`, { credentials: "include" });
if (!res.ok) throw new Error("no logs");
return res.text();
},
```
(Use the file's real base-URL constant / credentials pattern — inspect how `request` builds URLs and mirror it. If the app is same-origin with a rewrite, a relative `/api/...` fetch is fine.)
- Export a helper to build the SSE URL: `serverRunLogStreamUrl(runId, serverId)` returning the `/api/runs/:runId/servers/:serverId/logs/stream` URL against the same base.
- In the Settings type, add `workflow_log_retention_days?: number | null`.
- [ ] **Step 2: Live tail in the run detail page**
In `web/app/workflows/[id]/runs/[runId]/page.tsx`:
- Remove all use of `st.stdout` / `st.stderr` (fields gone). Step `<details>` now show status/exit/attempts pills only.
- Add a per-server live terminal. For each `server_run`, render a `<pre>` and, while `sr.status === "running"`, subscribe via `EventSource`:
```tsx
function ServerLog({ runId, serverId, status }: { runId: string; serverId: string; status: string }) {
const [text, setText] = useState("");
const preRef = useRef<HTMLPreElement>(null);
const running = status === "running";
useEffect(() => {
if (running) {
const es = new EventSource(api.serverRunLogStreamUrl(runId, serverId), { withCredentials: true });
es.onmessage = (e) => setText((t) => t + e.data + "\n");
es.addEventListener("done", () => es.close());
es.onerror = () => es.close();
return () => es.close();
}
// terminal: fetch the whole file once
api.getServerRunLog(runId, serverId).then(setText).catch(() => setText(""));
}, [running, runId, serverId]);
useEffect(() => { preRef.current?.scrollTo(0, preRef.current.scrollHeight); }, [text]);
return (
<pre ref={preRef} className="mt-2 max-h-80 overflow-auto rounded bg-black/40 p-2 font-mono text-xs text-text-secondary whitespace-pre-wrap">
{text || (running ? "Waiting for output…" : "No output.")}
</pre>
);
}
```
Render `<ServerLog runId={run.run_id} serverId={sr.server_id} status={sr.status} />` inside each server card, below the step pills. Keep the existing react-query `refetchInterval` on the run (drives status pills); the SSE handles live text.
- [ ] **Step 3: Retention field in Settings**
In `web/app/settings/page.tsx`, add a "Workflow log retention (days)" number input bound to `workflow_log_retention_days`, saved through the existing settings save mutation. Add helper text: "0 = keep forever." Match the page's existing input styling.
- [ ] **Step 4: Verify build**
Run: `cd web && npm run build`
Expected: type-checks and builds. Fix any lingering `st.stdout`/`st.stderr` references.
- [ ] **Step 5: Commit**
```bash
git add web/lib/api.ts web/app/workflows/[id]/runs/[runId]/page.tsx web/app/settings/page.tsx
git commit -m "feat(web): live SSE log tail and log retention setting"
```
---
## Task 8: End-to-end verification
**Files:** none (verification only).
- [ ] **Step 1: Build everything**
Run: `cd server && go build ./... && go vet ./... && cd ../agent && go build ./... && go vet ./... && cd ../web && npm run build`
Expected: all succeed.
- [ ] **Step 2: Manual smoke (documented, run if an environment is available)**
With server + MongoDB + a connected agent:
1. Run a workflow with a step that emits output slowly (e.g. `for i in $(seq 1 10); do echo "line $i"; sleep 1; done`). Open the run detail page while running; confirm lines appear live (SSE), not only at the end.
2. Confirm `<logdir>/<run_id>/<server_id>.log` exists on the server with step markers and the output.
3. Confirm `workflow_runs` doc no longer stores stdout/stderr bodies; `steps[].log_offset` is set.
4. Add a secret ref and echo it; confirm the file shows `***`, including when the secret would straddle a chunk boundary.
5. Set retention to 0 in Settings → confirm sweeper keeps files; set to a small value and backdate a run's `finished_at` → confirm the dir is removed within the hour (or call `sweepLogs` path manually).
- [ ] **Step 3: Commit any fixes found**
```bash
git add -A
git commit -m "fix: workflow log streaming e2e fixes"
```
---
## Self-Review Notes
- **Spec coverage:** §3 proto → T1; §4 agent streaming → T2; §5.1 registry + §7 sweeper → T3; §7.1 setting + startup → T4; §5.3/§5.4 runner+model → T5; §6 REST/SSE → T6; §8 frontend → T7. Tests omitted per Global Constraints.
- **Masking** boundary-safe carry buffer in `StepLogs.Append`, flushed in `Close` (T3); `output_env`/`run_env` masking unchanged (T5 keeps the merged fix).
- **commandID ownership** moved to the runner so the log file can be opened before dispatch (T5) — mirrors the `StepResults.Await`-before-dispatch ordering.
- **Buildable commits:** T3 adds a temp stub for `GetWorkflowLogRetentionDays`, removed in T4.
- **Removed fields** `StepRun.Stdout/Stderr` — every reader updated in T5 (runner) and T7 (frontend).
- **Open follow-ups (out of scope):** per-step SSE channels, log download/zip, compression, pre-existing runs have no files.
```
File diff suppressed because it is too large Load Diff
@@ -1,241 +0,0 @@
# Vantage Web Console (Guacamole Replacement) — Design
**Date:** 2026-07-17
**Status:** Approved design, pre-implementation
## Goal
Add a browser-based remote-access console to Vantage — SSH, RDP, and VNC into
managed servers — as a self-hosted Guacamole replacement. Users select an SSH
key to connect over SSH. RDP targets are reachable from a new Windows agent that
registers the host and reports status. Windows agent ships as an MSI installer
produced by CI.
## Non-Goals (YAGNI)
- Session recording / replay (may be added later).
- Native Go RDP implementation (guacd handles protocol translation).
- Per-user Linux/Windows account management from the agent.
- Tunneling console traffic through the agent (direct network path assumed).
---
## Architecture
```
Browser (guacamole-common-js, vendored — no CDN)
│ Guacamole protocol over WebSocket
Go server: /api/console/tunnel (github.com/wwt/guac)
│ Guacamole protocol over TCP :4822
guacd container (Apache Guacamole daemon)
│ SSH :22 / RDP :3389 / VNC :5900 — direct to target IP
Target host (LAN / VPN line-of-sight from server)
```
- **Browser:** loads vendored `guacamole-common-js`, renders RDP/VNC display and
SSH terminal. No external CDN (matches existing infra rules).
- **Go server:** exposes a WebSocket tunnel endpoint using `github.com/wwt/guac`
(Go Guacamole tunnel library). No Java `guacamole-client` required.
- **guacd:** new container in `deploy/docker-compose.yml`, bound to the internal
docker network only, reachable by the server on `:4822`.
- **Network path:** guacd connects **directly** to the target IP. Requires the
central server to have network line-of-sight to hosts (homelab LAN / VPN). The
agent's outbound-only guarantee is unchanged — the console path is
server→target, not agent-mediated.
---
## Data Model Changes
### `keys` — extend to hold private material
```json
{
"key_id": "uuid",
"label": "dom-macbook",
"public_key": "ssh-ed25519 AAAA...",
"private_key_enc": "<AES-256-GCM ciphertext | null>",
"has_private": true,
"passphrase_enc": "<AES-256-GCM ciphertext | null>",
"fingerprint": "SHA256:...",
"source": "uploaded|generated",
"created_at": "ISODate"
}
```
- A key may be created from an uploaded **private+public** pair, upload of a
public key only, or agent generation.
- Agent key generation now also uploads `private_key_enc` (reuses the existing
AES-256 key used for at-rest encryption). Private key no longer stays local
only — it is stored encrypted so the console can reuse it.
- Optional `passphrase_enc` for passphrase-protected private keys.
- Console lists only keys where `has_private = true`.
### `servers` — extend with console metadata
```json
{
"...": "...existing fields...",
"os_type": "linux|windows",
"console_protocols": ["ssh"],
"ssh_port": 22,
"rdp_port": 3389
}
```
- `os_type` set at registration from the agent.
- `console_protocols` lists enabled protocols per server (`ssh`, `rdp`, `vnc`).
- Port fields default to standard ports, overridable in the UI.
### `console_sessions` — new collection (audit)
```json
{
"session_id": "uuid",
"server_id": "uuid",
"protocol": "ssh|rdp|vnc",
"key_id": "uuid | null",
"user": "who opened it",
"started_at": "ISODate",
"ended_at": "ISODate | null",
"client_ip": "string"
}
```
---
## Session Broker + Connection Flow
New service: `server/internal/services/console.go`.
1. Browser `POST /api/console/connect`
`{ server_id, protocol, key_id?, rdp_username?, rdp_password? }`.
2. Broker validates request, loads the server (host IP, port for protocol),
loads the key and **decrypts `private_key_enc` in memory only**.
3. Builds the guacd connection parameter map:
- **SSH:** `hostname`, `port`, `username`, `private-key` (decrypted),
`passphrase` (if any).
- **RDP:** `hostname`, `port`, `username`, `password`, `security=any`,
`ignore-cert=true`.
- **VNC:** `hostname`, `port`, `password`.
4. Creates a `console_sessions` document, returns a short-lived signed session
token.
5. Browser opens WebSocket `/api/console/tunnel?token=…`. The `wwt/guac` handler
validates the token, dials guacd `:4822`, and pipes bytes in both directions.
6. On socket close, the broker sets `ended_at` on the session doc.
### Security
- Decrypted private keys and RDP passwords are **never persisted, never logged,
never sent to the browser** — passed only to guacd.
- Session token: short TTL (~60s to open the WebSocket), single-use,
HMAC-signed, bound to the authenticated user.
- guacd is bound to the internal docker network only; not exposed publicly.
- At-rest encryption (`private_key_enc`, `passphrase_enc`) reuses the existing
AES-256 key already used for agent-generated private keys.
---
## Windows Agent
Same Go codebase as the Linux agent, with a reduced role: **register +
heartbeat + status only**. No `authorized_keys` management (meaningless on
Windows).
- Build target: `GOOS=windows GOARCH=amd64``vantage-agent-windows-amd64.exe`.
- Agent detects OS at registration and sends `os_type=windows`.
- The key-sync loop is disabled on Windows via a runtime OS check (or build tag)
— no `authorized_keys` writes are ever attempted.
- Config file: `C:\ProgramData\vantage\config.yaml`, locked down via ACL to the
equivalent of `0600`.
- Runs as a Windows service via **nssm**.
---
## Windows Installer (MSI)
Agent ships as a WiX v4 MSI produced in CI.
- **WiX v4** chosen because it is a dotnet tool that builds MSIs
**cross-platform** — runs on the Linux Gitea act_runner. (Inno Setup is
Windows-only and does not fit the runner.)
- MSI bundles `vantage-agent.exe`, installs it to `C:\Program Files\Vantage\`,
and registers the nssm service (ships nssm or uses a CustomAction).
- Accepts install parameters as MSI properties for silent/headless install:
```
msiexec /i vantage-agent.msi /qn SERVERID=<id> TOKEN=<token> SERVERURL=vantage..:9090
```
- GUI install (double-click) prompts for server-id / token / server-url via a
dialog.
### Two install paths
1. **Installer direct** — user downloads `vantage-agent.msi`, double-clicks,
fills the dialog. No script required.
2. **PowerShell one-liner** — served dynamically (like the existing bash
`/install`). Script downloads the `.msi`, verifies SHA-256, then runs
`msiexec /qn` with injected `SERVERID` / `TOKEN` / `SERVERURL`. Used by the
copy-paste "Add Server" flow.
The PowerShell script (`/install.ps1`) steps:
1. Detect arch.
2. Download `vantage-agent.msi` from the latest Gitea `agent/v*` release.
3. Verify SHA-256 against `checksums.txt`.
4. Run `msiexec /i vantage-agent.msi /qn SERVERID=.. TOKEN=.. SERVERURL=..`.
---
## Frontend Routes
| Route | Change |
| ------------------------- | ------------------------------------------------------------- |
| `/servers` | Show `os_type` badge, enabled console protocols |
| `/servers/[id]` | Add **Connect** button(s) per enabled protocol |
| `/servers/[id]/console` | New — full-screen console (guacamole-common-js), key picker |
| `/servers/new` | Offer Windows (MSI) vs Linux (bash) install instructions |
Console page: select protocol + SSH key (SSH) or enter RDP credentials, call
`/api/console/connect`, open the tunnel WebSocket, mount the Guacamole client.
---
## CI/CD Changes
### `agent-release.yml`
- Add `windows/amd64` build: `vantage-agent-windows-amd64.exe`.
- Add WiX v4 MSI build job → `vantage-agent.msi`.
- Add both to `checksums.txt` and release assets.
Release assets become:
- `vantage-agent-linux-amd64`
- `vantage-agent-linux-arm64`
- `vantage-agent-windows-amd64.exe`
- `vantage-agent.msi`
- `checksums.txt`
### `server-deploy.yml`
- Add guacd service to `deploy/docker-compose.yml` (deployed alongside server).
---
## New Dependencies
- **Go:** `github.com/wwt/guac` (Guacamole tunnel/WebSocket in Go).
- **Container:** `guacamole/guacd` official image.
- **Frontend:** vendored `guacamole-common-js` (no CDN).
- **CI:** WiX v4 dotnet tool; nssm binary bundled for the MSI.
---
## Open Implementation Notes
- Confirm `wwt/guac` API surface for connection-parameter passing and token auth
binding during implementation.
- nssm packaging inside MSI: bundle the nssm binary as a payload + CustomAction,
or run `sc.exe`-based service install if nssm proves awkward in WiX.
- ACL hardening of `C:\ProgramData\vantage\config.yaml` in the MSI CustomAction.
@@ -1,238 +0,0 @@
# Server Workflows — Design
**Date:** 2026-07-20
**Status:** Approved (design) — ready for implementation planning
**Scope:** Server Workflows only. Fleet Inventory and SaaS/local-auth are separate sub-projects with their own specs.
Approved UI mockup: three-pane builder (Step Library · Canvas · Inspector), env vars shown riding the wire between nodes.
---
## 1. Summary
Let operators compose **reusable shell steps** (Bash or PowerShell) into **workflows** and run them across many managed servers in parallel. Steps pass data to later steps through a `$WORKFLOW_ENV` file (GitHub-Actions style). Every run is recorded with full per-step logs. Steps can reference org secrets, injected as environment variables at runtime.
Builds directly on the existing `CommandStream` gRPC infrastructure (`dispatch.go`, `ServerCommand` oneof, agent command loop).
---
## 2. Locked decisions
| Topic | Decision |
|-------|----------|
| Data passing | Implicit. Every step's `$WORKFLOW_ENV` outputs merge into the run's env and are exposed to **all** later steps as `$KEY`. No explicit port wiring. |
| Failure model | Per-step policy: `stop` (default), `continue`, `retry` (with max attempt count). |
| Targets | Fan-out. Same step sequence runs on N target servers **in parallel**. Steps within one server run **sequentially**. |
| History/logs | Every run persisted: status, timing, per-server per-step stdout/stderr/exit code, captured output env. |
| Secrets | Steps declare needed secret keys; resolved from existing `secrets` store and injected as env vars at exec time. Never persisted into run logs. |
| Testing | **Skipped** for this iteration per request. No test files written. |
---
## 3. Data model (MongoDB)
### `workflow_steps` — reusable step library
```json
{
"_id": "ObjectId",
"step_id": "uuid",
"name": "Restart service",
"description": "Restart-Service by name, wait ready",
"interpreter": "bash | powershell",
"script": "Restart-Service vantage-api\n...",
"declared_outputs": ["STARTED_AT"], // documentation/UI hints; not enforced
"secret_refs": ["DEPLOY_TOKEN"], // secret keys this step needs injected
"org_id": "uuid", // for future multi-tenant; single-org for now
"created_at": "ISODate",
"updated_at": "ISODate"
}
```
### `workflows` — ordered composition
```json
{
"_id": "ObjectId",
"workflow_id": "uuid",
"name": "Deploy & Restart API",
"target_server_ids": ["uuid", "uuid"],
"steps": [
{
"step_id": "uuid", // reference to library step
"order": 0,
"on_failure": "stop | continue | retry",
"max_retries": 0, // used when on_failure = retry
"overrides": { // optional local fork of the library step
"script": null,
"secret_refs": null
}
}
],
"created_at": "ISODate",
"updated_at": "ISODate"
}
```
Editing a library step from the Inspector writes an `overrides` block on that workflow step (a local fork) rather than mutating the shared step.
### `workflow_runs` — execution records
```json
{
"_id": "ObjectId",
"run_id": "uuid",
"workflow_id": "uuid",
"workflow_snapshot": { }, // frozen copy of workflow + resolved steps at trigger time
"status": "running | success | failed | cancelled",
"triggered_by": "user-id",
"started_at": "ISODate",
"finished_at": "ISODate | null",
"server_runs": [
{
"server_id": "uuid",
"status": "queued | running | success | failed | skipped",
"started_at": "ISODate | null",
"finished_at": "ISODate | null",
"run_env": { "VERSION": "a1b9f0" }, // accumulated non-secret output env
"steps": [
{
"order": 0,
"name": "Git pull & build",
"status": "success | failed | running | queued | skipped",
"attempts": 1,
"exit_code": 0,
"stdout": "…",
"stderr": "…",
"output_env": { "VERSION": "a1b9f0" },
"started_at": "ISODate",
"finished_at": "ISODate"
}
]
}
]
}
```
Secret values are never written to `stdout`/`stderr`/`run_env` by us; masking of known secret values in captured output is applied before persistence.
---
## 4. gRPC protocol changes (`proto/vantage/v1/vantage.proto`)
### New command in the `ServerCommand` oneof
```protobuf
message RunStepCmd {
string interpreter = 1; // "bash" | "powershell"
string script = 2;
map<string, string> env = 3; // inputs = accumulated run env + injected secrets
int32 timeout_seconds = 4;
}
```
Add `RunStepCmd run_step = 6;` to the `ServerCommand` oneof.
### Richer result — new `AgentMessage` payload
Current `CommandResult{command_id, success, message}` is too thin. Add a dedicated step result:
```protobuf
message StepResult {
string command_id = 1;
int32 exit_code = 2;
string stdout = 3;
string stderr = 4;
map<string, string> output_env = 5; // parsed $WORKFLOW_ENV KEY=value lines
}
```
Add `StepResult step_result = 5;` to the `AgentMessage` oneof (alongside existing `ready` / `result`).
---
## 5. Agent execution (`agent/internal/...`)
New handler for `RunStepCmd` in the agent command loop:
1. Create a temp dir; create empty `WORKFLOW_ENV` file inside it.
2. Write `script` to a temp script file.
3. Build the process environment: inherited env + `cmd.env` (run env + secrets) + `WORKFLOW_ENV=<path to env file>`.
4. Execute:
- `bash``bash <script>`
- `powershell``pwsh -NoProfile -File <script>` (fallback `powershell.exe` on Windows if `pwsh` absent).
5. Capture stdout, stderr, exit code. Enforce `timeout_seconds` (kill on exceed → non-zero exit, stderr note).
6. Parse the `WORKFLOW_ENV` file: each `KEY=value` line becomes an `output_env` entry (last write wins; supports multi-line via simple `KEY<<EOF` heredoc form, optional for v1 — start with single-line `KEY=value`).
7. Reply with `StepResult`. Delete temp dir.
Agent runs as root (existing), so no privilege change. Script content is trusted operator input.
---
## 6. Server orchestration (`server/internal/services/workflows.go`)
Runner responsibilities:
1. On trigger: snapshot the workflow (resolve each library step + overrides), create a `workflow_runs` doc with one `server_run` per target, all `queued`.
2. Spawn one goroutine **per target server** (parallel fan-out). Each goroutine:
- Verifies the agent is connected (`Dispatcher.IsConnected`); if not → `server_run.status = skipped`, reason recorded.
- Maintains a `run_env map[string]string`, seeded empty.
- For each step in order:
- Resolve `secret_refs` from the secrets service → merge into the command env (kept separate from persisted `run_env`).
- Dispatch `RunStepCmd{env: run_env + secrets}` via a **correlated** send — needs a way to await the matching `StepResult` by `command_id` (see §7).
- On result: persist step record (stdout/stderr/exit, masked); merge `output_env` into `run_env`.
- Apply `on_failure` on non-zero exit: `stop` (fail server_run, break), `continue` (mark failed, proceed), `retry` (re-dispatch up to `max_retries`).
3. Aggregate: run `status = success` if all server_runs succeeded, else `failed`. Set `finished_at`.
### Concurrency / queue
- One workflow run per workflow at a time (reject or queue concurrent triggers — v1: reject with clear error).
- Per-server step dispatch is serial; servers are parallel.
---
## 7. Correlated command results
The existing dispatcher is fire-and-forget; workflows need request/response by `command_id`. Add a small **pending-result registry** alongside `Dispatcher`:
- `AwaitResult(commandID) <-chan *pb.StepResult` — registers a channel before dispatch.
- The `CommandStream` receive loop, on a `StepResult`, looks up the pending channel by `command_id` and delivers it (falls back to existing `CommandResult` handling for other command types).
- Timeout guard on the server side (step `timeout_seconds` + grace) so a dead agent can't hang a run.
This is additive; existing `CommandResult` flow for key/update commands is unchanged.
---
## 8. REST API (`server/internal/api/workflows.go`)
| Method + path | Purpose |
|---------------|---------|
| `GET /api/steps` / `POST` / `PUT /:id` / `DELETE /:id` | Reusable step library CRUD |
| `GET /api/workflows` / `POST` / `PUT /:id` / `DELETE /:id` | Workflow CRUD (name, targets, ordered steps) |
| `POST /api/workflows/:id/run` | Trigger a run; returns `run_id` |
| `GET /api/workflows/:id/runs` | Run history (summary list) |
| `GET /api/runs/:run_id` | Full run detail incl. per-server per-step logs |
| `POST /api/runs/:run_id/cancel` | Best-effort cancel |
Secrets are referenced by key only through these APIs; values never returned.
---
## 9. Frontend (`web/app/workflows/`)
- `/workflows` — list workflows, last run status/time, Run button.
- `/workflows/[id]` — the three-pane builder from the approved mockup:
- **Library** (left): reusable steps, `bash`/`pwsh` badges, search, add.
- **Canvas** (center): ordered nodes, env chips on wires, live status pills.
- **Inspector** (right): name, command editor, declared inputs/outputs, `secret_refs` picker, `on_failure` + retry count.
- `/workflows/[id]/runs/[runId]` — run detail: per-server columns, expandable per-step stdout/stderr, exit codes, timing. Live-updating while `running` (poll, consistent with existing 30s-poll ethos — or reuse whatever the console screen uses).
Reuse existing web components/styling patterns (there is already `servers`, `secrets`, `audit`, console UI to match).
---
## 10. Security notes
- Scripts are trusted operator input executed as root — same trust level as the existing console feature. No new sandbox in v1.
- Secret values injected as env only; masked from all persisted logs (`stdout`/`stderr`/`run_env`) by literal replacement before write.
- Run triggering and step/workflow CRUD gated behind existing auth (`server/internal/auth`).
- Audit: emit audit-log entries (existing `audit` service) on workflow create/edit/delete and run trigger.
---
## 11. Out of scope (this iteration)
- Tests (explicitly skipped).
- Branching/conditional steps, matrix per-server conditionals (fan-out only).
- Scheduled/cron triggers (manual run only for v1).
- Multi-org isolation enforcement (schema carries `org_id` for later; single-org behavior now).
- Artifact upload/collection beyond env vars.
@@ -1,193 +0,0 @@
# Workflow Log Streaming — Design
**Date:** 2026-07-20
**Status:** Approved (design) — ready for implementation planning
**Scope:** Stream step stdout/stderr live from agent to server-side log files, tail them live in the UI, and auto-expire them on a retention period. Enhancement to the already-merged Server Workflows feature. No auth/orgs, no inventory.
---
## 1. Summary
Today a workflow step buffers all stdout/stderr in agent RAM, ships it in one terminal `StepResult`, and the server persists the whole body into the `workflow_runs` Mongo document. Long/chatty steps risk: agent memory blow-up, the gRPC 4MB message ceiling, and the Mongo 16MB document cap.
Change to **live streaming**:
1. Agent streams output chunks over the existing `CommandStream` as the process runs.
2. Server appends chunks (secret-masked) to a **per-server-run log file** on disk — not Mongo.
3. UI tails the file live via **SSE** while a server-run is running; slices per-step by byte offset after completion.
4. A **retention sweeper** deletes old run-log directories on a configurable period (default 30 days, set in Settings).
`workflow_runs` documents shrink: they no longer carry `stdout`/`stderr` bodies, only status/exit/attempts/output_env/timestamps plus a per-step `log_offset`.
---
## 2. Locked decisions
| Topic | Decision |
|-------|----------|
| Transport | Reuse bidirectional `CommandStream`. New `AgentMessage.StepOutput` chunk message. |
| Chunk shape | `{command_id, seq, data, eof}`. Interleaved stdout+stderr in execution order. |
| Terminal result | `StepResult` still sent at step end, now carries only `exit_code` + `output_env` (no stdout/stderr). |
| Log granularity | **One file per server-run**: `<logdir>/<run_id>/<server_id>.log`, with a marker line before each step. |
| Streams | **Interleaved** — single synchronized writer on the agent, terminal-order output. |
| Masking | **Server-side** (agent can't tell secret env from normal env). Per-stream carry buffer of `maxSecretLen-1` bytes so a secret split across a chunk boundary still masks; flushed on EOF. |
| Live tail | **SSE** at per-server-run granularity: `GET /api/runs/:runId/servers/:serverId/logs/stream`. Post-run whole-file fetch + per-step offset slice. |
| Retention | `settings.workflow_log_retention_days`, default **30**, editable in `/settings`. Hourly sweeper deletes `<logdir>/<run_id>/` dirs older than retention by run `finished_at`. |
| Log dir | Env `VANTAGE_WORKFLOW_LOG_DIR`, default `<data>/workflow-logs`. Created `0700`. |
| Mongo | No log bodies in `workflow_runs`. Disk is the source of truth for output. |
---
## 3. gRPC protocol (`proto/vantage/v1/vantage.proto` + both `pb.go` files)
Add to the `AgentMessage` oneof: `StepOutputChunk step_output = 6;`
```protobuf
message StepOutputChunk {
string command_id = 1;
uint64 seq = 2; // monotonic per command_id, 0-based
bytes data = 3; // raw interleaved stdout+stderr bytes
bool eof = 4; // true on the final (empty) chunk
}
```
`StepResult` is unchanged in shape but `stdout`/`stderr` are now left empty by the agent (kept in the message for backward-compat / error notes only — server ignores them for log content). The server still reads `exit_code` and `output_env` from `StepResult`.
Hand-written JSON-codec struct added to **both** `server/internal/grpc/pb/vantage.pb.go` and `agent/internal/grpc/pb/vantage.pb.go`, identical. `AgentMessage` gains `StepOutput *StepOutputChunk` in both.
`data` is `[]byte` in the Go structs (JSON-codec base64-encodes it, which is fine).
---
## 4. Agent (`agent/internal/exec/exec.go`)
`RunStep` signature gains a chunk sink:
```go
func RunStep(cmd *pb.RunStepCmd, emit func(seq uint64, data []byte)) *pb.StepResult
```
- Replace the two `bytes.Buffer`s with a single `streamWriter` set as **both** `c.Stdout` and `c.Stderr`. Its `Write` takes a mutex (so stdout+stderr interleave without interleaving *within* a write), assigns the next `seq`, and calls `emit(seq, copyOfBytes)`. Chunks are whatever the OS pipe delivers (typically ≤64KB); no extra buffering/line-assembly.
- `StepResult` returns with `Stdout`/`Stderr` empty; `ExitCode` and `OutputEnv` populated as today (env parsing unchanged).
- On timeout/exec error, put the short note in `StepResult.Stderr` (terminal, not streamed) so the runner can still surface a failure reason even if nothing streamed.
Agent loop (`agent/internal/sync/sync.go`, the `cmd.RunStep != nil` goroutine): pass an `emit` closure that sends `AgentMessage{ServerId, AgentToken, StepOutput: &pb.StepOutputChunk{CommandId, Seq, Data}}` through the existing mutex-guarded `send()`. After `RunStep` returns, send a final `StepOutput{eof:true, seq:last+1}` then the terminal `StepResult` (both via `send()`). Ordering: all chunks, then eof, then StepResult.
---
## 5. Server write path
### 5.1 Log writer registry (`server/internal/services/steplogs.go`)
Parallel to `StepResults`. Keyed by `command_id`:
```go
type stepLogWriter struct {
f *os.File
mu sync.Mutex
carry []byte // held-back tail for boundary-safe masking
secrets []string // secret literals to mask
maxSecret int
}
var StepLogs = &stepLogRegistry{ ... }
func (r *stepLogRegistry) Open(commandID, path string, secrets []string) (*stepLogWriter, error)
func (r *stepLogRegistry) Append(commandID string, data []byte) // masked write
func (r *stepLogRegistry) Close(commandID string) // flush carry, close file
```
- `Append` masking: concatenate `carry+data`, mask all secret literals (`ReplaceAll(v,"***")`), then write everything except the last `maxSecret-1` bytes; keep those as the new `carry`. `Close` masks+writes the remaining carry. If `secrets` empty, write straight through (no carry).
- The file handle is opened append-only (`O_APPEND|O_CREATE|O_WRONLY`, `0600`); dir `0700`.
### 5.2 Stream delivery (`server/internal/grpc/server.go`)
In the receive loop, after the `m.StepResult` block, add:
```go
if m.StepOutput != nil {
if m.StepOutput.Eof {
services.StepLogs.Close(m.StepOutput.CommandId)
} else {
services.StepLogs.Append(m.StepOutput.CommandId, m.StepOutput.Data)
}
}
```
### 5.3 Runner changes (`server/internal/services/workflow_runner.go`)
- Resolve the log dir + run/server file path once per server-run; ensure `<logdir>/<run_id>/` exists.
- Before dispatching each step: write the step marker line to the file (`\n===== step <order>: <name> =====\n`), record the current file byte offset as the step's `log_offset` (persisted on the `StepRun`), and `StepLogs.Open(commandID, path, secretVals)` **before** `DispatchRunStep` (same ordering rule as `StepResults.Await`).
- `dispatchAndWait` no longer expects stdout/stderr in the result. On terminal `StepResult`, `StepLogs.Close(commandID)` is driven by the agent's eof; the runner also calls `Close` defensively on timeout/dispatch-failure (idempotent).
- **Drop** `stdout`/`stderr` from `finishStep` persistence. Masking of the streamed body is done in `Append`; `run_env`/`output_env` masking (existing, from the merged fix) stays.
- The shared server-run file is written by two writers that never overlap in time (steps are serial, and the runner writes each step marker *before* `StepLogs.Open`): (a) the runner writes markers directly to the path, serially between steps; (b) `StepLogs` writes chunks during a step. **Resolved approach:** `Open(commandID, path, secrets)` opens the path fresh with `O_APPEND|O_CREATE|O_WRONLY` for that step and `Close(commandID)` closes it on eof. One handle live at a time per server-run (serial steps guarantee this), so there is no shared-handle race and no ref-counting. The runner's marker write is a separate short `O_APPEND` open/write/close on the same path.
### 5.4 Data model (`server/internal/models/workflow.go`)
`StepRun`:
- **Remove** `Stdout`, `Stderr` string fields.
- **Add** `LogOffset int64 `bson:"log_offset" json:"log_offset"`` — byte offset in the server-run file where this step's marker begins.
`ServerRun` gains nothing structural (its file path is derivable: `<logdir>/<run_id>/<server_id>.log`).
---
## 6. REST API (`server/internal/api/workflows.go`)
- `GET /api/runs/:runId/servers/:serverId/logs` — returns the whole server-run log file (`text/plain`). 404 if absent. Used post-run and as SSE fallback.
- `GET /api/runs/:runId/servers/:serverId/logs/stream`**SSE**. Opens the file, streams existing content as `data:` events, then polls for appends (~500ms) emitting new bytes, until the server-run status is terminal (success/failed/skipped/cancelled) AND no more bytes, then sends a final `event: done` and closes. Sets `Content-Type: text/event-stream`, disables gin's buffering. Guards against path traversal (runId/serverId are used as literal path segments — validate they are UUIDs / contain no separators).
Log content served by these endpoints is already masked (masking happens at write time), so no masking needed on read.
---
## 7. Retention
### 7.1 Setting
`settings` collection gains `workflow_log_retention_days int` (default 30 when unset). Read/write via the existing settings service + surfaced in `/settings` UI as a number input. `0` or negative disables sweeping (keep forever) — document this.
### 7.2 Sweeper (`server/internal/services/steplogs.go` or `logsweeper.go`)
- `StartLogSweeper()` launched at server startup (next to index setup): hourly `time.Ticker`.
- Each tick: read retention setting; if ≤0 skip. Compute cutoff = `now - retentionDays`. For each `<logdir>/<run_id>/` dir, look up the run's `finished_at` (query `workflow_runs` by run_id); if finished and older than cutoff, `os.RemoveAll` the dir. Fallback to dir mtime if the run doc is gone.
- Also run once at startup.
---
## 8. Frontend
### 8.1 API client (`web/lib/api.ts`)
- `StepRun`: remove `stdout`/`stderr`; add `log_offset: number`.
- Add `getServerRunLog(runId, serverId): Promise<string>` (GET .../logs).
- SSE consumed directly via `EventSource` in the component (not through the `request` helper), URL built from the same base.
- Settings type gains `workflow_log_retention_days`.
### 8.2 Run detail (`web/app/workflows/[id]/runs/[runId]/page.tsx`)
- Per-server card: while the server-run is `running`, open an `EventSource` to the stream endpoint and render a live `<pre>` terminal that appends incoming chunks (auto-scroll). Close the source on `event: done`, unmount, or terminal status.
- After completion: fetch the whole file once and render it; step `<details>` still list status/exit/attempts pills. (Per-step slicing by `log_offset` is optional polish — v1 may show the whole server log under the card and keep step pills as the status summary.)
- Remove reliance on `st.stdout`/`st.stderr` (fields gone).
### 8.3 Settings (`web/app/settings/page.tsx`)
- Add a "Workflow log retention (days)" number input bound to `workflow_log_retention_days`, saved via the existing settings mutation. Note that `0` = keep forever.
---
## 9. Security
- Secret masking moves to the streaming write path but remains server-side and boundary-safe (carry buffer). Same `***` replacement.
- Log files `0600`, dirs `0700`, under a dedicated log dir.
- SSE/read endpoints validate `runId`/`serverId` as UUID-shaped path segments to prevent traversal; they are session-authed (same `apiGroup`).
- Terminal `StepResult.Stderr` (error notes only) is still masked before any persistence (it is no longer persisted as log body; if surfaced, mask against secretVals).
---
## 10. Out of scope
- Per-step (rather than per-server) live SSE channels.
- Log compression / rotation within a run, remote log storage (S3), download-as-zip.
- Full-text search over logs.
- Backfilling/migrating already-existing `workflow_runs` stdout/stderr into files (pre-existing runs keep whatever they had; new field just won't be set — acceptable, feature is new).
- Tests (skipped, consistent with the Workflows iteration).
```
@@ -0,0 +1,267 @@
# Ad-hoc Steps, Step Import/Export, and Default Steps — Design
Date: 2026-07-21
## Summary
Four related additions to the workflow step system:
1. **Ad-hoc steps** — steps defined inline in a single workflow, not written to the
shared step library.
2. **Import/Export** — single-step portable JSON (`vantage.step/v1`). Import can
target the shared library or a workflow as an inline ad-hoc step.
3. **Default steps** — JSON files in a bind-mounted directory, seeded into the
library on boot and re-syncable on demand. Org-ready for a future SaaS plan.
4. **Auto-derived outputs**`declared_outputs` is scanned from the script
(writes to `$WORKFLOW_ENV`) instead of being entered by hand.
Existing model: shared steps live in the `workflow_steps` collection; a
`Workflow.Steps[]` is a list of `WorkflowStepRef` that reference a library step by
`step_id` and may carry `Overrides` + `Inputs`. `resolveSteps` freezes each ref
into a `ResolvedStep` snapshot at run time.
## 1. Data model
`server/internal/models/workflow.go`.
### WorkflowStep
Add a provenance field:
```go
Source string `bson:"source" json:"source"` // "user" | "default"
Slug string `bson:"slug" json:"slug"` // kebab of name; stable key for default seeding
```
`Slug` is set for `source="default"` steps (used as the upsert key by the seeder).
For `source="user"` steps it may be empty. Existing steps default to
`source="user"` (absent field decodes to "").
### WorkflowStepRef
Add an inline definition. A ref is **either** a library ref (`StepID` set) **or**
ad-hoc (`Inline` set). Never both.
```go
type WorkflowStepRef struct {
StepID string `bson:"step_id,omitempty" json:"step_id"`
Inline *WorkflowStep `bson:"inline,omitempty" json:"inline,omitempty"`
Order int `bson:"order" json:"order"`
OnFailure string `bson:"on_failure" json:"on_failure"`
MaxRetries int `bson:"max_retries" json:"max_retries"`
Overrides *StepOverride `bson:"overrides,omitempty" json:"overrides,omitempty"` // library-ref only
Inputs map[string]string `bson:"inputs,omitempty" json:"inputs,omitempty"`
}
```
`Inline` reuses `WorkflowStep` (name, description, interpreter, script,
declared_inputs, declared_outputs, secret_refs). Its `ID`, `StepID`, `Slug`,
`Source`, and timestamps stay empty and are never persisted to `workflow_steps`.
Validation on workflow create/update: for each ref, exactly one of `StepID` /
`Inline` must be set. `Overrides` is ignored when `Inline` is set.
## 2. Resolve at run time
`server/internal/services/workflow_runner.go`, `resolveSteps`.
For each ref:
- If `ref.Inline != nil`: build `ResolvedStep` from `ref.Inline` directly
(name/interpreter/script/secret_refs), apply `ref.Inputs` against
`Inline.DeclaredInputs` defaults. Skip `getStep`, skip `Overrides`.
- Else: current library path unchanged (load step, apply overrides).
`ResolvedStep` output shape and the run snapshot are unchanged, so the runner and
the run-history UI need no changes.
`DeleteStep` cascade is unaffected — ad-hoc refs carry no `step_id`, so they never
match the cascade query.
## 3. Import / Export
Portable single-step JSON, `kind: "vantage.step/v1"`:
```json
{
"kind": "vantage.step/v1",
"name": "...",
"description": "...",
"interpreter": "bash",
"script": "...",
"declared_inputs": [ { "name": "...", "default": "...", "description": "..." } ],
"declared_outputs": ["..."],
"secret_refs": ["NAME"]
}
```
Export strips `_id`, `step_id`, `slug`, `source`, and timestamps.
Secret refs are exported as names only. On import, dangling secret refs are kept
verbatim (not auto-created).
### Service functions (`services/workflows.go`)
- `ExportStep(stepID string) ([]byte, error)` — load library step, marshal to the
v1 shape.
- `ParseStepDoc(b []byte) (models.WorkflowStep, error)` — validate `kind`, decode
into a `WorkflowStep` (no id/source). Shared by both import targets.
- `ImportStepToLibrary(b []byte) (*models.WorkflowStep, error)``ParseStepDoc`
then `CreateStep` (fresh `step_id`, `source="user"`).
Import-to-inline needs no new service fn: the web editor calls `ParseStepDoc`'s
API equivalent (see routes) and drops the returned step object into a new
`WorkflowStepRef.Inline` in the workflow it's editing, then saves the workflow
normally.
### Routes (`api/workflows.go`)
- `GET /api/steps/:id/export` — returns JSON as a downloadable attachment
(`Content-Disposition`).
- `POST /api/steps/import` — body is the v1 JSON; imports to library; returns the
created step. (Used by the "import to library" flow.)
- `POST /api/steps/parse` — body is the v1 JSON; validates and returns the
normalized step object **without** persisting. Used by "import to inline" so the
editor can insert it as an ad-hoc ref. (Keeps parsing/validation server-side.)
Audit: `workflow.step_imported` logged on library import.
## 4. Default steps (seed + admin re-sync)
Mirrors the existing `WorkflowLogDir` pattern — a bind-mounted directory, no Go
`embed`.
### Directory
```go
// DefaultStepsDir returns the directory holding default step JSON files, creating it.
func DefaultStepsDir() string {
dir := os.Getenv("VANTAGE_DEFAULT_STEPS_DIR")
if dir == "" {
dir = filepath.Join("data", "default-steps")
}
_ = os.MkdirAll(dir, 0700)
return dir
}
```
Compose already bind-mounts `./data:/data`. Set
`VANTAGE_DEFAULT_STEPS_DIR=/data/default-steps` for explicitness (optional).
Operator drops `*.json` (`vantage.step/v1`) files into that folder.
### Seeder
`SeedDefaultSteps() (created, updated int, err error)`:
1. Glob `DefaultStepsDir()/*.json`.
2. For each file: `ParseStepDoc`, compute `slug = kebab(name)`.
3. Upsert into `workflow_steps` keyed on `{ slug, source: "default" }`:
- absent → insert with fresh `step_id`, `source="default"`, `slug`. (`created++`)
- present → `$set` name/description/interpreter/script/declared_*/secret_refs +
`updated_at`. (`updated++`)
**Override rule (confirmed):** re-sync is authoritative for `source="default"`
steps and overwrites their content, reverting any user edits to those steps.
`source="user"` steps are never touched by the seeder, even on a slug collision
(the seeder query is scoped to `source: "default"`).
Add a partial unique index on `slug` where `source == "default"` (or enforce
uniqueness in the seeder loop) to keep default slugs unambiguous.
### Boot
Call `SeedDefaultSteps()` from server startup after `EnsureWorkflowIndexes()`
(alongside index setup in `server/cmd/main.go`). Log the created/updated counts;
a seed error is logged but non-fatal (server still boots).
### Route
- `POST /api/steps/seed-defaults` (admin) — runs `SeedDefaultSteps()`, returns
`{ "created": n, "updated": m }`. Audit `workflow.defaults_synced`.
### Org readiness
Signature stays global today. When Orgs land, `SeedDefaultSteps(orgID)` seeds
per-org and the upsert key becomes `{ org_id, slug, source }`. No schema churn
blocks that later change.
## 5. Auto-derived outputs
Today `WorkflowStep.DeclaredOutputs` is entered by hand and consumed only by the
UI (no runtime reads it — outputs are captured at run time by `parseEnvFile` on
the agent). Replace manual entry with a server-side scan of the script.
At run time the agent exposes an env file path in `$WORKFLOW_ENV` (bash) /
`$env:WORKFLOW_ENV` (powershell); a step emits an output by appending a
`KEY=value` line to it, e.g. `echo "test=123" >> $WORKFLOW_ENV`.
### Scanner
`services.DeriveOutputs(script string) []string`:
- Scan line by line. For each line that references `WORKFLOW_ENV`, extract every
`KEY=` assignment target on that line, where `KEY` matches
`[A-Za-z_][A-Za-z0-9_]*`.
- Covers the common forms across both interpreters (line mentions `WORKFLOW_ENV`
and contains `KEY=...`):
- `echo "test=123" >> $WORKFLOW_ENV`
- `echo "test=123" >> "$WORKFLOW_ENV"`
- `printf 'k=v\n' >> $WORKFLOW_ENV`
- `"k=v" >> $env:WORKFLOW_ENV` / `Add-Content $env:WORKFLOW_ENV "k=v"`
- Deduplicate, preserve first-seen order. Best-effort heuristic — false positives
are acceptable (they only widen the documented output list); it never affects
what the agent actually captures.
### Wiring
- `CreateStep` and `UpdateStep` set `DeclaredOutputs = DeriveOutputs(s.Script)`,
ignoring any client-sent value.
- Inline ad-hoc steps: `DeriveOutputs` is applied when the workflow is saved (for
each `ref.Inline`), so inline outputs are derived too.
- `ParseStepDoc` (import) also derives outputs, so `declared_outputs` in an
imported/exported file is informational and always recomputed on import.
- `SeedDefaultSteps` derives outputs the same way when upserting.
`DeclaredOutputs` stays in the model and JSON (still shown in the UI and used to
wire step-to-step input references), it is just no longer user-authored.
### Web
The step editor's "declared outputs" input becomes a read-only, auto-populated
display (derived from the script, refreshed on save / on script edit). No manual
add/remove.
## 6. Web
`web/app/workflows/[id]/page.tsx` and the steps list page.
- **Steps list:** per-row **Export** (downloads JSON) and top-level **Import**
(file picker → `POST /api/steps/import` → library). **Sync defaults** admin
button → `POST /api/steps/seed-defaults`, toast the counts. `source="default"`
rows get a "default" badge.
- **Workflow editor — Add step:** existing "add from library" plus **Add ad-hoc
step** (inline mini-form: name, interpreter, script, optional inputs) stored as
a `WorkflowStepRef.Inline`. Also **Import ad-hoc from file** → `POST
/api/steps/parse` → inserts the returned step as a new inline ref.
- Ad-hoc rows in the editor are visually distinguished from library refs (badge)
and are editable in place; library refs keep the existing override UI.
## Testing
- `resolveSteps`: inline ref resolves without touching the library; inputs apply
from `Inline.DeclaredInputs` defaults and ref overrides; library path unchanged.
- Workflow validation: rejects a ref with both `StepID` and `Inline`, and one with
neither.
- Export → import round-trips to an equivalent library step with a new `step_id`.
- `ParseStepDoc` rejects a wrong/missing `kind`.
- `SeedDefaultSteps`: insert-then-update idempotency; user steps untouched;
user-edited default step reverted on re-sync; counts correct.
- `DeleteStep` cascade ignores ad-hoc refs.
- `DeriveOutputs`: extracts keys from each interpreter form above, dedupes,
preserves order, ignores lines not referencing `WORKFLOW_ENV`; create/update/
import/seed all populate `declared_outputs` from it and ignore client input.
## Out of scope
- Whole-workflow export/import.
- Auto-creating secret refs on import.
- Multi-tenant Org model (design is forward-compatible only).
+8
View File
@@ -92,9 +92,16 @@ message ServerCommand {
UpdateAgentCmd update_agent = 4;
ApplyUpdatesCmd apply_updates = 5;
RunStepCmd run_step = 6;
CleanupWorkspaceCmd cleanup_workspace = 7;
}
}
// CleanupWorkspaceCmd tells the agent to recursively remove the run's working
// directory once all steps on that server have finished.
message CleanupWorkspaceCmd {
string workspace_id = 1;
}
message DeleteKeyCmd {
string label = 1;
}
@@ -117,6 +124,7 @@ message RunStepCmd {
string script = 2;
map<string, string> env = 3;
int32 timeout_seconds = 4;
string workspace_id = 5; // per-run working dir the agent creates & uses as cwd
}
message StepResult {
+8
View File
@@ -31,6 +31,14 @@ func main() {
log.Printf("warning: failed to ensure workflow indexes: %v", err)
}
if created, updated, err := services.SeedDefaultSteps(); err != nil {
log.Printf("warning: failed to seed default steps: %v", err)
} else {
log.Printf("default steps seeded: %d created, %d updated", created, updated)
}
services.StartLogSweeper()
redisAddr := getEnv("REDIS_ADDR", "localhost:6379")
if err := auth.InitRedis(redisAddr); err != nil {
log.Fatalf("failed to connect to Redis: %v", err)
+4 -3
View File
@@ -450,14 +450,15 @@ func getSettings(c *gin.Context) {
func saveSettings(c *gin.Context) {
var body struct {
Alerts models.AlertSettings `json:"alerts"`
Email models.EmailSettings `json:"email"`
Alerts models.AlertSettings `json:"alerts"`
Email models.EmailSettings `json:"email"`
WorkflowLogRetentionDays *int `json:"workflow_log_retention_days"`
}
if err := c.ShouldBindJSON(&body); err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
if err := services.SaveSettings(body.Alerts, body.Email); err != nil {
if err := services.SaveSettings(body.Alerts, body.Email, body.WorkflowLogRetentionDays); err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
+188 -1
View File
@@ -2,8 +2,13 @@ package api
import (
"fmt"
"io"
"net/http"
"os"
"regexp"
"strconv"
"strings"
"time"
"github.com/gin-gonic/gin"
"github.com/mrhid6/vantage/server/internal/models"
@@ -15,6 +20,11 @@ func registerWorkflowRoutes(g *gin.RouterGroup) {
g.POST("/steps", createStep)
g.PUT("/steps/:id", updateStep)
g.DELETE("/steps/:id", deleteStep)
g.GET("/steps/:id/export", exportStep)
g.POST("/steps/import", importStep)
g.POST("/steps/seed-defaults", seedDefaults)
g.GET("/steps/usage", stepUsage)
g.POST("/steps/parse", parseStep)
g.GET("/workflows", listWorkflows)
g.POST("/workflows", createWorkflow)
@@ -26,6 +36,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) {
@@ -37,6 +155,15 @@ func listSteps(c *gin.Context) {
c.JSON(http.StatusOK, steps)
}
func stepUsage(c *gin.Context) {
counts, err := services.StepUsageCounts()
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, counts)
}
func createStep(c *gin.Context) {
var s models.WorkflowStep
if err := c.ShouldBindJSON(&s); err != nil {
@@ -75,6 +202,61 @@ func deleteStep(c *gin.Context) {
c.JSON(http.StatusOK, gin.H{"deleted": true})
}
func exportStep(c *gin.Context) {
b, err := services.ExportStep(c.Param("id"))
if err != nil {
c.JSON(http.StatusNotFound, gin.H{"error": err.Error()})
return
}
c.Header("Content-Disposition", fmt.Sprintf("attachment; filename=step-%s.json", c.Param("id")))
c.Data(http.StatusOK, "application/json", b)
}
func seedDefaults(c *gin.Context) {
created, updated, err := services.SeedDefaultSteps()
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
services.LogEvent("workflow.defaults_synced", actorFromCtx(c), "", "", fmt.Sprintf("default steps synced: %d created, %d updated", created, updated))
c.JSON(http.StatusOK, gin.H{"created": created, "updated": updated})
}
const maxStepBodyBytes = 1 << 20 // 1 MiB
func importStep(c *gin.Context) {
c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, maxStepBodyBytes)
body, err := io.ReadAll(c.Request.Body)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
out, err := services.ImportStepToLibrary(body)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
services.LogEvent("workflow.step_imported", actorFromCtx(c), "", out.StepID, fmt.Sprintf("step '%s' imported", out.Name))
c.JSON(http.StatusCreated, out)
}
// parseStep validates a step doc and returns the normalized step WITHOUT
// persisting — used by the editor to insert an imported ad-hoc (inline) step.
func parseStep(c *gin.Context) {
c.Request.Body = http.MaxBytesReader(c.Writer, c.Request.Body, maxStepBodyBytes)
body, err := io.ReadAll(c.Request.Body)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
s, err := services.ParseStepDoc(body)
if err != nil {
c.JSON(http.StatusBadRequest, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, s)
}
func listWorkflows(c *gin.Context) {
wfs, err := services.ListWorkflows()
if err != nil {
@@ -119,7 +301,12 @@ func updateWorkflow(c *gin.Context) {
return
}
services.LogEvent("workflow.updated", actorFromCtx(c), "", c.Param("id"), "workflow updated")
c.JSON(http.StatusOK, gin.H{"updated": true})
updated, err := services.GetWorkflow(c.Param("id"))
if err != nil {
c.JSON(http.StatusInternalServerError, gin.H{"error": err.Error()})
return
}
c.JSON(http.StatusOK, updated)
}
func deleteWorkflow(c *gin.Context) {
+16 -6
View File
@@ -66,12 +66,19 @@ type ReportUpdatesResponse struct{}
type ApplyUpdatesCmd struct{}
type ServerCommand struct {
CommandId string `json:"command_id"`
GenerateKey *GenerateKeyCmd `json:"generate_key,omitempty"`
DeleteKey *DeleteKeyCmd `json:"delete_key,omitempty"`
UpdateAgent *UpdateAgentCmd `json:"update_agent,omitempty"`
ApplyUpdates *ApplyUpdatesCmd `json:"apply_updates,omitempty"`
RunStep *RunStepCmd `json:"run_step,omitempty"`
CommandId string `json:"command_id"`
GenerateKey *GenerateKeyCmd `json:"generate_key,omitempty"`
DeleteKey *DeleteKeyCmd `json:"delete_key,omitempty"`
UpdateAgent *UpdateAgentCmd `json:"update_agent,omitempty"`
ApplyUpdates *ApplyUpdatesCmd `json:"apply_updates,omitempty"`
RunStep *RunStepCmd `json:"run_step,omitempty"`
CleanupWorkspace *CleanupWorkspaceCmd `json:"cleanup_workspace,omitempty"`
}
// CleanupWorkspaceCmd tells the agent to recursively remove the run's working
// directory once all steps on that server have finished.
type CleanupWorkspaceCmd struct {
WorkspaceId string `json:"workspace_id"`
}
type DeleteKeyCmd struct {
@@ -113,6 +120,9 @@ type RunStepCmd struct {
Script string `json:"script"`
Env map[string]string `json:"env,omitempty"`
TimeoutSeconds int `json:"timeout_seconds,omitempty"`
// WorkspaceId names the per-run working directory the agent creates and uses
// as the step's cwd. Empty means run in the agent's default directory.
WorkspaceId string `json:"workspace_id,omitempty"`
}
type StepResult struct {
+7
View File
@@ -132,6 +132,13 @@ func (s *vantageServer) CommandStream(stream pb.Vantage_CommandStreamServer) err
if m.StepResult != nil {
services.StepResults.Deliver(m.StepResult)
}
if m.StepOutput != nil {
if m.StepOutput.Eof {
services.StepLogs.Close(m.StepOutput.CommandId)
} else {
services.StepLogs.Append(m.StepOutput.CommandId, m.StepOutput.Data)
}
}
}
}()
+2
View File
@@ -36,4 +36,6 @@ type Settings struct {
Alerts AlertSettings `bson:"alerts" json:"alerts"`
Email EmailSettings `bson:"email" json:"email"`
Secrets SecretsSettings `bson:"secrets" json:"secrets"`
// WorkflowLogRetentionDays: nil = default 30, 0 = keep forever, N = N days.
WorkflowLogRetentionDays *int `bson:"workflow_log_retention_days,omitempty" json:"workflow_log_retention_days,omitempty"`
}
+38 -23
View File
@@ -1,26 +1,41 @@
package models
import "time"
import (
"time"
"go.mongodb.org/mongo-driver/v2/bson"
)
type InputParam struct {
Name string `bson:"name" json:"name"`
Default string `bson:"default" json:"default"`
Description string `bson:"description" json:"description"`
}
type WorkflowStep struct {
ID string `bson:"_id,omitempty" json:"-"`
StepID string `bson:"step_id" json:"step_id"`
Name string `bson:"name" json:"name"`
Description string `bson:"description" json:"description"`
Interpreter string `bson:"interpreter" json:"interpreter"` // "bash" | "powershell"
Script string `bson:"script" json:"script"`
DeclaredOutputs []string `bson:"declared_outputs" json:"declared_outputs"`
SecretRefs []string `bson:"secret_refs" json:"secret_refs"`
CreatedAt time.Time `bson:"created_at" json:"created_at"`
UpdatedAt time.Time `bson:"updated_at" json:"updated_at"`
ID bson.ObjectID `bson:"_id,omitempty" json:"-"`
StepID string `bson:"step_id" json:"step_id"`
Name string `bson:"name" json:"name"`
Description string `bson:"description" json:"description"`
Interpreter string `bson:"interpreter" json:"interpreter"` // "bash" | "powershell"
Script string `bson:"script" json:"script"`
DeclaredOutputs []string `bson:"declared_outputs" json:"declared_outputs"`
DeclaredInputs []InputParam `bson:"declared_inputs" json:"declared_inputs"`
SecretRefs []string `bson:"secret_refs" json:"secret_refs"`
Source string `bson:"source" json:"source"` // "user" | "default"
Slug string `bson:"slug,omitempty" json:"slug,omitempty"`
CreatedAt time.Time `bson:"created_at" json:"created_at"`
UpdatedAt time.Time `bson:"updated_at" json:"updated_at"`
}
type WorkflowStepRef struct {
StepID string `bson:"step_id" json:"step_id"`
Order int `bson:"order" json:"order"`
OnFailure string `bson:"on_failure" json:"on_failure"` // "stop" | "continue" | "retry"
MaxRetries int `bson:"max_retries" json:"max_retries"`
Overrides *StepOverride `bson:"overrides,omitempty" json:"overrides,omitempty"`
StepID string `bson:"step_id,omitempty" json:"step_id,omitempty"`
Inline *WorkflowStep `bson:"inline,omitempty" json:"inline,omitempty"`
Order int `bson:"order" json:"order"`
OnFailure string `bson:"on_failure" json:"on_failure"` // "stop" | "continue" | "retry"
MaxRetries int `bson:"max_retries" json:"max_retries"`
Overrides *StepOverride `bson:"overrides,omitempty" json:"overrides,omitempty"`
Inputs map[string]string `bson:"inputs,omitempty" json:"inputs,omitempty"`
}
type StepOverride struct {
@@ -29,7 +44,7 @@ type StepOverride struct {
}
type Workflow struct {
ID string `bson:"_id,omitempty" json:"-"`
ID bson.ObjectID `bson:"_id,omitempty" json:"-"`
WorkflowID string `bson:"workflow_id" json:"workflow_id"`
Name string `bson:"name" json:"name"`
TargetServerIDs []string `bson:"target_server_ids" json:"target_server_ids"`
@@ -44,9 +59,10 @@ type ResolvedStep struct {
Name string `bson:"name" json:"name"`
Interpreter string `bson:"interpreter" json:"interpreter"`
Script string `bson:"script" json:"script"`
SecretRefs []string `bson:"secret_refs" json:"secret_refs"`
OnFailure string `bson:"on_failure" json:"on_failure"`
MaxRetries int `bson:"max_retries" json:"max_retries"`
SecretRefs []string `bson:"secret_refs" json:"secret_refs"`
OnFailure string `bson:"on_failure" json:"on_failure"`
MaxRetries int `bson:"max_retries" json:"max_retries"`
Inputs map[string]string `bson:"inputs" json:"inputs"`
}
type StepRun struct {
@@ -55,8 +71,7 @@ type StepRun struct {
Status string `bson:"status" json:"status"` // queued|running|success|failed|skipped
Attempts int `bson:"attempts" json:"attempts"`
ExitCode int `bson:"exit_code" json:"exit_code"`
Stdout string `bson:"stdout" json:"stdout"`
Stderr string `bson:"stderr" json:"stderr"`
LogOffset int64 `bson:"log_offset" json:"log_offset"`
OutputEnv map[string]string `bson:"output_env" json:"output_env"`
StartedAt *time.Time `bson:"started_at,omitempty" json:"started_at,omitempty"`
FinishedAt *time.Time `bson:"finished_at,omitempty" json:"finished_at,omitempty"`
@@ -73,7 +88,7 @@ type ServerRun struct {
}
type WorkflowRun struct {
ID string `bson:"_id,omitempty" json:"-"`
ID bson.ObjectID `bson:"_id,omitempty" json:"-"`
RunID string `bson:"run_id" json:"run_id"`
WorkflowID string `bson:"workflow_id" json:"workflow_id"`
Name string `bson:"name" json:"name"`
+94
View File
@@ -0,0 +1,94 @@
package services
import (
"os"
"path/filepath"
"time"
"github.com/google/uuid"
"github.com/mrhid6/vantage/server/internal/db"
"github.com/mrhid6/vantage/server/internal/models"
"go.mongodb.org/mongo-driver/v2/bson"
"go.mongodb.org/mongo-driver/v2/mongo/options"
)
// DefaultStepsDir returns the directory holding default step JSON files.
func DefaultStepsDir() string {
dir := os.Getenv("VANTAGE_DEFAULT_STEPS_DIR")
if dir == "" {
dir = filepath.Join("data", "default-steps")
}
_ = os.MkdirAll(dir, 0700)
return dir
}
// readDefaultStepFiles parses every *.json in the defaults dir into
// source=default library steps (with slug set). Non-json and invalid files are
// skipped silently; a slug is derived from the step name.
func readDefaultStepFiles() ([]models.WorkflowStep, error) {
matches, err := filepath.Glob(filepath.Join(DefaultStepsDir(), "*.json"))
if err != nil {
return nil, err
}
out := []models.WorkflowStep{}
for _, path := range matches {
b, err := os.ReadFile(path)
if err != nil {
continue
}
s, err := ParseStepDoc(b)
if err != nil {
continue
}
s.Source = "default"
s.Slug = Slugify(s.Name)
if s.Slug == "" {
continue
}
out = append(out, s)
}
return out, nil
}
// SeedDefaultSteps upserts default steps from disk keyed on {slug, source}.
// Re-sync overwrites default-step content; user steps are never touched.
func SeedDefaultSteps() (created, updated int, err error) {
steps, err := readDefaultStepFiles()
if err != nil {
return 0, 0, err
}
ctx, cancel := wfCtx()
defer cancel()
col := db.Col("workflow_steps")
for _, s := range steps {
filter := bson.M{"slug": s.Slug, "source": "default"}
set := bson.M{
"name": s.Name,
"description": s.Description,
"interpreter": s.Interpreter,
"script": s.Script,
"declared_outputs": s.DeclaredOutputs,
"declared_inputs": s.DeclaredInputs,
"secret_refs": s.SecretRefs,
"updated_at": time.Now(),
}
res, uerr := col.UpdateOne(ctx, filter, bson.M{
"$set": set,
"$setOnInsert": bson.M{
"step_id": uuid.New().String(),
"slug": s.Slug,
"source": "default",
"created_at": time.Now(),
},
}, options.UpdateOne().SetUpsert(true))
if uerr != nil {
return created, updated, uerr
}
if res.UpsertedCount > 0 {
created++
} else if res.ModifiedCount > 0 {
updated++
}
}
return created, updated, nil
}
+38
View File
@@ -0,0 +1,38 @@
package services
import (
"os"
"path/filepath"
"testing"
)
func TestDefaultStepsDirEnv(t *testing.T) {
dir := filepath.Join(t.TempDir(), "ds")
t.Setenv("VANTAGE_DEFAULT_STEPS_DIR", dir)
got := DefaultStepsDir()
if got != dir {
t.Fatalf("got %q want %q", got, dir)
}
if _, err := os.Stat(dir); err != nil {
t.Fatalf("dir not created: %v", err)
}
}
func TestReadDefaultStepFiles(t *testing.T) {
dir := t.TempDir()
t.Setenv("VANTAGE_DEFAULT_STEPS_DIR", dir)
good := `{"kind":"vantage.step/v1","name":"Ping Host","interpreter":"bash","script":"ping -c1 x=1 >> $WORKFLOW_ENV"}`
os.WriteFile(filepath.Join(dir, "ping.json"), []byte(good), 0600)
os.WriteFile(filepath.Join(dir, "notes.txt"), []byte("ignore me"), 0600)
steps, err := readDefaultStepFiles()
if err != nil {
t.Fatal(err)
}
if len(steps) != 1 {
t.Fatalf("want 1 step, got %d", len(steps))
}
if steps[0].Slug != "ping-host" || steps[0].Source != "default" {
t.Fatalf("bad seed step: %+v", steps[0])
}
}
+13
View File
@@ -68,6 +68,19 @@ func DispatchRunStep(serverID, commandID string, cmd *pb.RunStepCmd) error {
return Dispatcher.dispatch(serverID, &pb.ServerCommand{CommandId: commandID, RunStep: cmd})
}
// DispatchCleanupWorkspace tells a server's agent to remove a run's working
// directory. Best-effort and fire-and-forget: if the agent is gone the temp dir
// is reclaimed by the OS on reboot anyway.
func DispatchCleanupWorkspace(serverID, workspaceID string) {
if !Dispatcher.IsConnected(serverID) {
return
}
_ = Dispatcher.dispatch(serverID, &pb.ServerCommand{
CommandId: uuid.New().String(),
CleanupWorkspace: &pb.CleanupWorkspaceCmd{WorkspaceId: workspaceID},
})
}
// KeyGenParams carries all options for a generate-key command.
type KeyGenParams struct {
Label string
+37
View File
@@ -0,0 +1,37 @@
package services
import (
"testing"
"github.com/mrhid6/vantage/server/internal/models"
)
func TestResolveInlineStep(t *testing.T) {
ref := models.WorkflowStepRef{
Order: 2,
OnFailure: "",
Inline: &models.WorkflowStep{
Name: "adhoc",
Interpreter: "bash",
Script: "echo hi",
SecretRefs: []string{"TOKEN"},
DeclaredInputs: []models.InputParam{
{Name: "REGION", Default: "eu"},
},
},
Inputs: map[string]string{"REGION": "us"},
}
rs := resolveInlineStep(ref)
if rs.Name != "adhoc" || rs.Script != "echo hi" || rs.Order != 2 {
t.Fatalf("bad resolve: %+v", rs)
}
if rs.OnFailure != "stop" {
t.Fatalf("want default on_failure=stop, got %q", rs.OnFailure)
}
if rs.Inputs["REGION"] != "us" {
t.Fatalf("want input override us, got %q", rs.Inputs["REGION"])
}
if len(rs.SecretRefs) != 1 || rs.SecretRefs[0] != "TOKEN" {
t.Fatalf("bad secret refs: %v", rs.SecretRefs)
}
}
+19 -2
View File
@@ -100,7 +100,7 @@ func VerifySecretsReadToken(token string) bool {
return subtle.ConstantTimeCompare(expected, got[:]) == 1
}
func SaveSettings(alerts models.AlertSettings, email models.EmailSettings) error {
func SaveSettings(alerts models.AlertSettings, email models.EmailSettings, retentionDays *int) error {
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
@@ -111,14 +111,31 @@ func SaveSettings(alerts models.AlertSettings, email models.EmailSettings) error
email.SMTPPort = 587
}
set := bson.M{"alerts": alerts, "email": email}
if retentionDays != nil {
set["workflow_log_retention_days"] = *retentionDays
}
_, err := db.Col("settings").UpdateOne(ctx,
bson.M{},
bson.M{"$set": bson.M{"alerts": alerts, "email": email}},
bson.M{"$set": set},
options.UpdateOne().SetUpsert(true),
)
return err
}
// GetWorkflowLogRetentionDays returns the log retention in days: 30 when unset,
// 0 for keep-forever, or the configured value.
func GetWorkflowLogRetentionDays() (int, error) {
s, err := GetSettings()
if err != nil {
return 30, err
}
if s.WorkflowLogRetentionDays == nil {
return 30, nil
}
return *s.WorkflowLogRetentionDays, nil
}
func SendOfflineWebhook(webhookURL, hostname, serverID, ipAddress string) {
payload := map[string]any{
"event": "server.offline",
+86
View File
@@ -0,0 +1,86 @@
package services
import (
"encoding/json"
"fmt"
"github.com/mrhid6/vantage/server/internal/models"
)
const StepDocKind = "vantage.step/v1"
// StepDoc is the portable, id-free representation of a step.
type StepDoc struct {
Kind string `json:"kind"`
Name string `json:"name"`
Description string `json:"description"`
Interpreter string `json:"interpreter"`
Script string `json:"script"`
DeclaredOutputs []string `json:"declared_outputs"`
DeclaredInputs []models.InputParam `json:"declared_inputs"`
SecretRefs []string `json:"secret_refs"`
}
// ExportStepDoc builds a portable doc from a library step (ids/source stripped).
func ExportStepDoc(s models.WorkflowStep) StepDoc {
return StepDoc{
Kind: StepDocKind,
Name: s.Name,
Description: s.Description,
Interpreter: s.Interpreter,
Script: s.Script,
DeclaredOutputs: s.DeclaredOutputs,
DeclaredInputs: s.DeclaredInputs,
SecretRefs: s.SecretRefs,
}
}
// ParseStepDoc validates a v1 doc and returns a normalized (id-free) step with
// declared_outputs recomputed from the script.
func ParseStepDoc(b []byte) (models.WorkflowStep, error) {
var d StepDoc
if err := json.Unmarshal(b, &d); err != nil {
return models.WorkflowStep{}, fmt.Errorf("invalid step JSON: %w", err)
}
if d.Kind != StepDocKind {
return models.WorkflowStep{}, fmt.Errorf("unsupported kind %q (want %q)", d.Kind, StepDocKind)
}
if d.Name == "" || d.Interpreter == "" {
return models.WorkflowStep{}, fmt.Errorf("step name and interpreter are required")
}
if d.SecretRefs == nil {
d.SecretRefs = []string{}
}
if d.DeclaredInputs == nil {
d.DeclaredInputs = []models.InputParam{}
}
return models.WorkflowStep{
Name: d.Name,
Description: d.Description,
Interpreter: d.Interpreter,
Script: d.Script,
DeclaredOutputs: DeriveOutputs(d.Script),
DeclaredInputs: d.DeclaredInputs,
SecretRefs: d.SecretRefs,
}, nil
}
// ImportStepToLibrary parses a doc and persists it as a new user library step.
func ImportStepToLibrary(b []byte) (*models.WorkflowStep, error) {
s, err := ParseStepDoc(b)
if err != nil {
return nil, err
}
return CreateStep(s)
}
// ExportStep loads a library step and marshals it to a portable doc.
func ExportStep(stepID string) ([]byte, error) {
ctx, cancel := wfCtx()
defer cancel()
s, err := getStep(ctx, stepID)
if err != nil {
return nil, err
}
return json.MarshalIndent(ExportStepDoc(*s), "", " ")
}
+60
View File
@@ -0,0 +1,60 @@
package services
import (
"encoding/json"
"testing"
"github.com/mrhid6/vantage/server/internal/models"
)
func mkStep() models.WorkflowStep {
return models.WorkflowStep{
StepID: "should-not-export", Source: "default", Name: "Restart",
Interpreter: "bash", Script: "echo x=1 >> $WORKFLOW_ENV",
SecretRefs: []string{"TOK"},
}
}
func TestParseStepDocValid(t *testing.T) {
raw := `{"kind":"vantage.step/v1","name":"Restart","interpreter":"bash",
"script":"echo x=1 >> $WORKFLOW_ENV","declared_outputs":["stale"],
"declared_inputs":[{"name":"A","default":"1"}],"secret_refs":["TOK"]}`
s, err := ParseStepDoc([]byte(raw))
if err != nil {
t.Fatal(err)
}
if s.Name != "Restart" || s.Interpreter != "bash" {
t.Fatalf("bad parse: %+v", s)
}
// declared_outputs recomputed from script, ignoring the file's ["stale"].
if len(s.DeclaredOutputs) != 1 || s.DeclaredOutputs[0] != "x" {
t.Fatalf("outputs should be derived, got %v", s.DeclaredOutputs)
}
if s.StepID != "" || s.Source != "" {
t.Fatalf("parse must not set id/source")
}
}
func TestParseStepDocBadKind(t *testing.T) {
if _, err := ParseStepDoc([]byte(`{"kind":"nope","name":"x"}`)); err == nil {
t.Fatal("want error for bad kind")
}
}
func TestParseStepDocBadJSON(t *testing.T) {
if _, err := ParseStepDoc([]byte(`{`)); err == nil {
t.Fatal("want error for bad json")
}
}
func TestExportStepDocRoundTrip(t *testing.T) {
doc := ExportStepDoc(mkStep())
b, _ := json.Marshal(doc)
s, err := ParseStepDoc(b)
if err != nil {
t.Fatal(err)
}
if s.Name != "Restart" || s.Interpreter != "bash" {
t.Fatalf("round trip lost data: %+v", s)
}
}
+217
View File
@@ -0,0 +1,217 @@
package services
import (
"bytes"
"os"
"path/filepath"
"strings"
"sync"
"time"
"github.com/mrhid6/vantage/server/internal/db"
"go.mongodb.org/mongo-driver/v2/bson"
)
// WorkflowLogDir returns the base directory for workflow step logs, creating it.
func WorkflowLogDir() string {
dir := os.Getenv("VANTAGE_WORKFLOW_LOG_DIR")
if dir == "" {
dir = filepath.Join("data", "workflow-logs")
}
_ = os.MkdirAll(dir, 0700)
return dir
}
// ServerRunLogPath is the per-server-run log file path.
func ServerRunLogPath(runID, serverID string) string {
return filepath.Join(WorkflowLogDir(), runID, serverID+".log")
}
// logTS is the UTC timestamp prefix stamped on every log line. Stored in UTC
// (RFC3339, millisecond precision); the UI renders it in the viewer's timezone.
func logTS() string {
return time.Now().UTC().Format("2006-01-02T15:04:05.000") + "Z"
}
// AppendMarker writes a timestamped event line to the server-run log and returns
// the byte offset at which the write began (used as a step's log_offset).
func AppendMarker(runID, serverID, text string) (int64, error) {
path := ServerRunLogPath(runID, serverID)
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return 0, err
}
f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
return 0, err
}
defer f.Close()
off, _ := f.Seek(0, 2) // current end = offset before write
if _, err := f.WriteString("[" + logTS() + "] " + text + "\n"); err != nil {
return off, err
}
return off, nil
}
// ---- streamed chunk writer, boundary-safe secret masking ----
type stepLogWriter struct {
mu sync.Mutex
f *os.File
carry []byte // bytes of an as-yet-unterminated line
secrets []string
}
type stepLogRegistry struct {
mu sync.Mutex
writers map[string]*stepLogWriter
}
var StepLogs = &stepLogRegistry{writers: make(map[string]*stepLogWriter)}
// Open opens (append) the server-run file for a step's streamed chunks.
func (r *stepLogRegistry) Open(commandID, path string, secrets []string) error {
if err := os.MkdirAll(filepath.Dir(path), 0700); err != nil {
return err
}
f, err := os.OpenFile(path, os.O_APPEND|os.O_CREATE|os.O_WRONLY, 0600)
if err != nil {
return err
}
w := &stepLogWriter{f: f, secrets: secrets}
r.mu.Lock()
r.writers[commandID] = w
r.mu.Unlock()
return nil
}
func (r *stepLogRegistry) get(commandID string) *stepLogWriter {
r.mu.Lock()
defer r.mu.Unlock()
return r.writers[commandID]
}
// Append buffers chunks into whole lines, then writes each complete line with a
// UTC timestamp prefix and secret masking applied. Buffering by line means a
// secret split across a chunk boundary is always masked (the whole line is
// assembled first) and every line carries its own timestamp.
func (r *stepLogRegistry) Append(commandID string, data []byte) {
w := r.get(commandID)
if w == nil {
return
}
w.mu.Lock()
defer w.mu.Unlock()
buf := append(w.carry, data...)
for {
i := bytes.IndexByte(buf, '\n')
if i < 0 {
break
}
w.writeLine(buf[:i])
buf = buf[i+1:]
}
w.carry = append([]byte{}, buf...)
}
// writeLine emits one masked, timestamped log line. Caller holds w.mu.
func (w *stepLogWriter) writeLine(line []byte) {
masked := maskBytes(line, w.secrets)
_, _ = w.f.WriteString("[" + logTS() + "] ")
_, _ = w.f.Write(masked)
_, _ = w.f.WriteString("\n")
}
// Close flushes any trailing partial line and closes the file.
func (r *stepLogRegistry) Close(commandID string) {
r.mu.Lock()
w := r.writers[commandID]
delete(r.writers, commandID)
r.mu.Unlock()
if w == nil {
return
}
w.mu.Lock()
defer w.mu.Unlock()
if len(w.carry) > 0 {
w.writeLine(w.carry)
w.carry = nil
}
_ = w.f.Close()
}
func maskBytes(b []byte, secrets []string) []byte {
s := string(b)
for _, v := range secrets {
if v == "" {
continue
}
s = strings.ReplaceAll(s, v, "***")
}
return []byte(s)
}
// ---- retention sweeper ----
// StartLogSweeper sweeps expired run-log dirs hourly (and once now).
func StartLogSweeper() {
go func() {
sweepLogs()
t := time.NewTicker(time.Hour)
defer t.Stop()
for range t.C {
sweepLogs()
}
}()
}
func sweepLogs() {
days := retentionDays()
if days <= 0 {
return
}
cutoff := time.Now().AddDate(0, 0, -days)
base := WorkflowLogDir()
entries, err := os.ReadDir(base)
if err != nil {
return
}
for _, e := range entries {
if !e.IsDir() {
continue
}
runID := e.Name()
dir := filepath.Join(base, runID)
if runExpired(runID, dir, cutoff) {
_ = os.RemoveAll(dir)
}
}
}
// runExpired is true when the run finished before cutoff (falling back to dir
// mtime when the run doc is gone).
func runExpired(runID, dir string, cutoff time.Time) bool {
ctx, cancel := wfCtx()
defer cancel()
var run struct {
FinishedAt *time.Time `bson:"finished_at"`
}
err := db.Col("workflow_runs").FindOne(ctx, bson.M{"run_id": runID}).Decode(&run)
if err == nil {
if run.FinishedAt == nil {
return false // still running / never finished — keep
}
return run.FinishedAt.Before(cutoff)
}
// run doc gone: use dir mtime
if fi, e := os.Stat(dir); e == nil {
return fi.ModTime().Before(cutoff)
}
return false
}
func retentionDays() int {
if v, err := GetWorkflowLogRetentionDays(); err == nil {
return v
}
return 30
}
+44
View File
@@ -0,0 +1,44 @@
package services
import (
"regexp"
"strings"
)
// keyAssign matches an env-var assignment target: KEY= (captures KEY).
var keyAssign = regexp.MustCompile(`([A-Za-z_][A-Za-z0-9_]*)=`)
// DeriveOutputs scans a step script and returns the output keys it writes to
// $WORKFLOW_ENV. Best-effort: only lines that reference WORKFLOW_ENV are
// considered. Deduplicated, first-seen order preserved.
func DeriveOutputs(script string) []string {
out := []string{}
seen := map[string]bool{}
for _, line := range strings.Split(script, "\n") {
if !strings.Contains(line, "WORKFLOW_ENV") {
continue
}
for _, m := range keyAssign.FindAllStringSubmatch(line, -1) {
key := m[1]
// Skip the sentinel itself (e.g. "WORKFLOW_ENV=..." assignments).
if key == "WORKFLOW_ENV" || key == "env" {
continue
}
if seen[key] {
continue
}
seen[key] = true
out = append(out, key)
}
}
return out
}
var slugStrip = regexp.MustCompile(`[^a-z0-9]+`)
// Slugify converts a step name into a stable kebab-case slug.
func Slugify(name string) string {
s := strings.ToLower(name)
s = slugStrip.ReplaceAllString(s, "-")
return strings.Trim(s, "-")
}
+45
View File
@@ -0,0 +1,45 @@
package services
import (
"reflect"
"testing"
)
func TestDeriveOutputs(t *testing.T) {
script := `#!/bin/bash
echo "test=123" >> $WORKFLOW_ENV
echo "other=hi" >> "$WORKFLOW_ENV"
printf 'third=1\n' >> $WORKFLOW_ENV
echo "test=456" >> $WORKFLOW_ENV
echo "ignored=nope"
NORMAL=assignment
`
got := DeriveOutputs(script)
want := []string{"test", "other", "third"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("got %v want %v", got, want)
}
}
func TestDeriveOutputsPowershell(t *testing.T) {
script := `"result=ok" >> $env:WORKFLOW_ENV
Add-Content $env:WORKFLOW_ENV "count=5"`
got := DeriveOutputs(script)
want := []string{"result", "count"}
if !reflect.DeepEqual(got, want) {
t.Fatalf("got %v want %v", got, want)
}
}
func TestDeriveOutputsNone(t *testing.T) {
got := DeriveOutputs("echo hello\nNOPE=1")
if len(got) != 0 {
t.Fatalf("got %v want empty", got)
}
}
func TestSlugify(t *testing.T) {
if got := Slugify("Restart NGINX Service!"); got != "restart-nginx-service" {
t.Fatalf("got %q", got)
}
}
+19
View File
@@ -0,0 +1,19 @@
package services
import (
"fmt"
"github.com/mrhid6/vantage/server/internal/models"
)
// ValidateWorkflow checks each step ref sets exactly one of step_id / inline.
func ValidateWorkflow(w models.Workflow) error {
for i, ref := range w.Steps {
hasLib := ref.StepID != ""
hasInline := ref.Inline != nil
if hasLib == hasInline {
return fmt.Errorf("step %d: exactly one of step_id or inline must be set", i)
}
}
return nil
}
+29
View File
@@ -0,0 +1,29 @@
package services
import (
"testing"
"github.com/mrhid6/vantage/server/internal/models"
)
func TestValidateWorkflow(t *testing.T) {
inline := &models.WorkflowStep{Name: "x", Interpreter: "bash", Script: "echo hi"}
cases := []struct {
name string
ref models.WorkflowStepRef
wantErr bool
}{
{"library only", models.WorkflowStepRef{StepID: "abc"}, false},
{"inline only", models.WorkflowStepRef{Inline: inline}, false},
{"both set", models.WorkflowStepRef{StepID: "abc", Inline: inline}, true},
{"neither set", models.WorkflowStepRef{}, true},
}
for _, tc := range cases {
t.Run(tc.name, func(t *testing.T) {
err := ValidateWorkflow(models.Workflow{Steps: []models.WorkflowStepRef{tc.ref}})
if (err != nil) != tc.wantErr {
t.Fatalf("got err=%v want wantErr=%v", err, tc.wantErr)
}
})
}
}
+118 -14
View File
@@ -2,6 +2,7 @@ package services
import (
"fmt"
"os"
"strings"
"time"
@@ -82,10 +83,24 @@ func resolveSteps(wf *models.Workflow) ([]models.ResolvedStep, error) {
defer cancel()
out := make([]models.ResolvedStep, 0, len(wf.Steps))
for _, ref := range wf.Steps {
if ref.Inline != nil {
out = append(out, resolveInlineStep(ref))
continue
}
lib, err := getStep(ctx, ref.StepID)
if err != nil {
return nil, err
}
inputs := map[string]string{}
for _, p := range lib.DeclaredInputs {
if ref.Inputs != nil {
if v, ok := ref.Inputs[p.Name]; ok {
inputs[p.Name] = v
continue
}
}
inputs[p.Name] = p.Default
}
rs := models.ResolvedStep{
Order: ref.Order,
Name: lib.Name,
@@ -94,6 +109,7 @@ func resolveSteps(wf *models.Workflow) ([]models.ResolvedStep, error) {
SecretRefs: lib.SecretRefs,
OnFailure: ref.OnFailure,
MaxRetries: ref.MaxRetries,
Inputs: inputs,
}
if ref.Overrides != nil {
if ref.Overrides.Script != nil {
@@ -111,6 +127,35 @@ func resolveSteps(wf *models.Workflow) ([]models.ResolvedStep, error) {
return out, nil
}
// resolveInlineStep freezes an ad-hoc (inline) step ref into a ResolvedStep.
func resolveInlineStep(ref models.WorkflowStepRef) models.ResolvedStep {
in := ref.Inline
inputs := map[string]string{}
for _, p := range in.DeclaredInputs {
if ref.Inputs != nil {
if v, ok := ref.Inputs[p.Name]; ok {
inputs[p.Name] = v
continue
}
}
inputs[p.Name] = p.Default
}
onFailure := ref.OnFailure
if onFailure == "" {
onFailure = "stop"
}
return models.ResolvedStep{
Order: ref.Order,
Name: in.Name,
Interpreter: in.Interpreter,
Script: in.Script,
SecretRefs: in.SecretRefs,
OnFailure: onFailure,
MaxRetries: ref.MaxRetries,
Inputs: inputs,
}
}
// executeRun fans out one goroutine per server run and waits for all to finish.
func executeRun(runID string) {
run, err := GetRun(runID)
@@ -151,16 +196,20 @@ func runServer(runID string, srvIdx int, steps []models.ResolvedStep, serverID s
if !Dispatcher.IsConnected(serverID) {
fin := time.Now()
_, _ = AppendMarker(runID, serverID, "agent not connected — server skipped")
setServerRun(runID, srvIdx, bson.M{"server_runs.$.status": "skipped", "server_runs.$.finished_at": fin})
return
}
_, _ = AppendMarker(runID, serverID, fmt.Sprintf("run started on %s — %d step(s), workspace vantage-run-%s", serverID, len(steps), runID))
runEnv := map[string]string{}
allSecrets := map[string]string{}
serverFailed := false
for i, step := range steps {
startStep(runID, serverID, i, "running")
stepStart := time.Now()
var res *pb.StepResult
attempts := 0
maxAttempts := 1
@@ -173,7 +222,20 @@ func runServer(runID string, srvIdx int, steps []models.ResolvedStep, serverID s
for k, v := range secretVals {
allSecrets[k] = v
}
// Input values may template earlier step outputs and secrets, e.g.
// URL="http://example.com/$VersionNumber". Expand against runEnv (outputs
// threaded from prior steps) and this step's secrets before dispatch.
subst := map[string]string{}
for k, v := range runEnv {
subst[k] = v
}
for k, v := range secretVals {
subst[k] = v
}
cmdEnv := map[string]string{}
for k, v := range step.Inputs {
cmdEnv[k] = expandVars(v, subst)
}
for k, v := range runEnv {
cmdEnv[k] = v
}
@@ -181,60 +243,83 @@ func runServer(runID string, srvIdx int, steps []models.ResolvedStep, serverID s
cmdEnv[k] = v
}
// Write the step marker to the server-run log and remember the offset so
// the UI can slice this step's output later.
marker := fmt.Sprintf("===== step %d/%d: %s (%s) =====", step.Order+1, len(steps), step.Name, step.Interpreter)
offset, _ := AppendMarker(runID, serverID, marker)
logPath := ServerRunLogPath(runID, serverID)
secretsSlice := secretValues(secretVals)
commandID := uuid.New().String()
for attempts < maxAttempts {
attempts++
res = dispatchAndWait(serverID, &pb.RunStepCmd{
if attempts > 1 {
_, _ = AppendMarker(runID, serverID, fmt.Sprintf("retry %d/%d after failure", attempts-1, maxAttempts-1))
}
// Open a fresh writer per attempt; the agent's eof closes it, and the
// defensive Close below covers a missing result.
_ = StepLogs.Open(commandID, logPath, secretsSlice)
res = dispatchAndWait(serverID, commandID, &pb.RunStepCmd{
Interpreter: step.Interpreter,
Script: step.Script,
Env: cmdEnv,
TimeoutSeconds: 0,
WorkspaceId: runID,
})
StepLogs.Close(commandID) // idempotent; no-op if eof already closed it
if res != nil && res.ExitCode == 0 {
break
}
}
// Mask secret values before persisting.
stdout, stderr := "", ""
exit := 1
outEnv := map[string]string{} // masked copy, safe to persist
outEnv := map[string]string{} // masked copy, safe to persist
if res != nil {
stdout = maskSecrets(res.Stdout, allSecrets)
stderr = maskSecrets(res.Stderr, allSecrets)
exit = res.ExitCode
for k, v := range res.OutputEnv {
runEnv[k] = v // real, unmasked value threads forward to later steps
outEnv[k] = maskSecrets(v, allSecrets)
}
} else {
stderr = "[vantage] agent did not return a result"
_, _ = AppendMarker(runID, serverID, "agent did not return a result")
}
status := "success"
if exit != 0 {
status = "failed"
}
finishStep(runID, serverID, i, status, attempts, exit, stdout, stderr, outEnv)
finishStep(runID, serverID, i, status, attempts, exit, offset, outEnv)
dur := time.Since(stepStart).Round(time.Millisecond)
_, _ = AppendMarker(runID, serverID, fmt.Sprintf("step %d/%d %s — exit %d, %d attempt(s), %s",
step.Order+1, len(steps), status, exit, attempts, dur))
if exit != 0 {
switch step.OnFailure {
case "continue":
// keep going
_, _ = AppendMarker(runID, serverID, "on_failure=continue — proceeding to next step")
default: // "stop" or exhausted "retry"
serverFailed = true
}
if serverFailed {
_, _ = AppendMarker(runID, serverID, "stopping run — remaining steps skipped")
markRemainingSkipped(runID, serverID, i+1)
break
}
}
}
// Tell the agent to remove the run's working directory now that its steps are
// done (success or failure). Best-effort; the OS reclaims temp dirs anyway.
DispatchCleanupWorkspace(serverID, runID)
fin := time.Now()
status := "success"
if serverFailed {
status = "failed"
}
_, _ = AppendMarker(runID, serverID, fmt.Sprintf("run %s in %s — workspace removed",
status, fin.Sub(now).Round(time.Millisecond)))
// Persist only a masked copy of runEnv; the real (unmasked) runEnv was already
// used above to build cmdEnv for each step and must never be written to the DB.
maskedRunEnv := make(map[string]string, len(runEnv))
@@ -250,8 +335,7 @@ func runServer(runID string, srvIdx int, steps []models.ResolvedStep, serverID s
// dispatchAndWait registers a waiter, dispatches the step, and blocks for the
// result or a timeout.
func dispatchAndWait(serverID string, cmd *pb.RunStepCmd) *pb.StepResult {
commandID := uuid.New().String()
func dispatchAndWait(serverID, commandID string, cmd *pb.RunStepCmd) *pb.StepResult {
ch := StepResults.Await(commandID)
if err := DispatchRunStep(serverID, commandID, cmd); err != nil {
StepResults.Cancel(commandID)
@@ -270,6 +354,18 @@ func dispatchAndWait(serverID string, cmd *pb.RunStepCmd) *pb.StepResult {
}
}
// expandVars substitutes $VAR and ${VAR} references in an input value from the
// given lookup (prior step outputs and secrets). Unknown references expand to
// empty, matching shell behaviour; a literal "$" is written as "$$".
func expandVars(v string, lookup map[string]string) string {
return os.Expand(v, func(name string) string {
if name == "$" {
return "$"
}
return lookup[name]
})
}
func resolveSecrets(refs []string) map[string]string {
out := map[string]string{}
for _, ref := range refs {
@@ -322,19 +418,27 @@ func startStep(runID, serverID string, order int, status string) {
})
}
func finishStep(runID, serverID string, order int, status string, attempts, exit int, stdout, stderr string, outEnv map[string]string) {
func finishStep(runID, serverID string, order int, status string, attempts, exit int, logOffset int64, outEnv map[string]string) {
now := time.Now()
updateStep(runID, serverID, order, bson.M{
"server_runs.$[s].steps.$[t].status": status,
"server_runs.$[s].steps.$[t].attempts": attempts,
"server_runs.$[s].steps.$[t].exit_code": exit,
"server_runs.$[s].steps.$[t].stdout": stdout,
"server_runs.$[s].steps.$[t].stderr": stderr,
"server_runs.$[s].steps.$[t].log_offset": logOffset,
"server_runs.$[s].steps.$[t].output_env": outEnv,
"server_runs.$[s].steps.$[t].finished_at": now,
})
}
// secretValues returns just the values of a secret map, for masking log output.
func secretValues(m map[string]string) []string {
out := make([]string, 0, len(m))
for _, v := range m {
out = append(out, v)
}
return out
}
func markRemainingSkipped(runID, serverID string, fromOrder int) {
ctx, cancel := wfCtx()
defer cancel()
+106 -5
View File
@@ -25,6 +25,13 @@ func EnsureWorkflowIndexes() error {
}); err != nil {
return err
}
if _, err := db.Col("workflow_steps").Indexes().CreateOne(ctx, mongo.IndexModel{
Keys: bson.D{{Key: "slug", Value: 1}},
Options: options.Index().SetUnique(true).
SetPartialFilterExpression(bson.M{"source": "default"}),
}); err != nil {
return err
}
if _, err := db.Col("workflows").Indexes().CreateOne(ctx, mongo.IndexModel{
Keys: bson.D{{Key: "workflow_id", Value: 1}}, Options: options.Index().SetUnique(true),
}); err != nil {
@@ -54,18 +61,50 @@ func ListSteps() ([]models.WorkflowStep, error) {
return steps, nil
}
// StepUsageCounts returns, per library step_id, the number of distinct
// workflows that reference it. Inline steps have no step_id and are ignored.
func StepUsageCounts() (map[string]int, error) {
ctx, cancel := wfCtx()
defer cancel()
cur, err := db.Col("workflows").Find(ctx, bson.M{})
if err != nil {
return nil, err
}
defer cur.Close(ctx)
var wfs []models.Workflow
if err := cur.All(ctx, &wfs); err != nil {
return nil, err
}
counts := map[string]int{}
for _, w := range wfs {
seen := map[string]bool{}
for _, ref := range w.Steps {
if ref.StepID == "" || seen[ref.StepID] {
continue
}
seen[ref.StepID] = true
counts[ref.StepID]++
}
}
return counts, nil
}
func CreateStep(s models.WorkflowStep) (*models.WorkflowStep, error) {
ctx, cancel := wfCtx()
defer cancel()
s.StepID = uuid.New().String()
s.CreatedAt = time.Now()
s.UpdatedAt = s.CreatedAt
if s.DeclaredOutputs == nil {
s.DeclaredOutputs = []string{}
s.DeclaredOutputs = DeriveOutputs(s.Script)
if s.Source == "" {
s.Source = "user"
}
if s.SecretRefs == nil {
s.SecretRefs = []string{}
}
if s.DeclaredInputs == nil {
s.DeclaredInputs = []models.InputParam{}
}
if _, err := db.Col("workflow_steps").InsertOne(ctx, s); err != nil {
return nil, err
}
@@ -80,7 +119,8 @@ func UpdateStep(stepID string, s models.WorkflowStep) error {
"description": s.Description,
"interpreter": s.Interpreter,
"script": s.Script,
"declared_outputs": s.DeclaredOutputs,
"declared_outputs": DeriveOutputs(s.Script),
"declared_inputs": s.DeclaredInputs,
"secret_refs": s.SecretRefs,
"updated_at": time.Now(),
}})
@@ -90,8 +130,38 @@ func UpdateStep(stepID string, s models.WorkflowStep) error {
func DeleteStep(stepID string) error {
ctx, cancel := wfCtx()
defer cancel()
_, err := db.Col("workflow_steps").DeleteOne(ctx, bson.M{"step_id": stepID})
return err
if _, err := db.Col("workflow_steps").DeleteOne(ctx, bson.M{"step_id": stepID}); err != nil {
return err
}
// Cascade: remove this step from every workflow that references it, re-sequencing orders.
cur, err := db.Col("workflows").Find(ctx, bson.M{"steps.step_id": stepID})
if err != nil {
return err
}
defer cur.Close(ctx)
var wfs []models.Workflow
if err := cur.All(ctx, &wfs); err != nil {
return err
}
for _, w := range wfs {
kept := make([]models.WorkflowStepRef, 0, len(w.Steps))
for _, ref := range w.Steps {
if ref.StepID == stepID {
continue
}
kept = append(kept, ref)
}
for i := range kept {
kept[i].Order = i
}
if _, err := db.Col("workflows").UpdateOne(ctx,
bson.M{"workflow_id": w.WorkflowID},
bson.M{"$set": bson.M{"steps": kept, "updated_at": time.Now()}},
); err != nil {
return err
}
}
return nil
}
func getStep(ctx context.Context, stepID string) (*models.WorkflowStep, error) {
@@ -144,6 +214,10 @@ func CreateWorkflow(w models.Workflow) (*models.Workflow, error) {
if w.Steps == nil {
w.Steps = []models.WorkflowStepRef{}
}
if err := ValidateWorkflow(w); err != nil {
return nil, err
}
normalizeInlineSteps(&w)
if _, err := db.Col("workflows").InsertOne(ctx, w); err != nil {
return nil, err
}
@@ -153,6 +227,10 @@ func CreateWorkflow(w models.Workflow) (*models.Workflow, error) {
func UpdateWorkflow(id string, w models.Workflow) error {
ctx, cancel := wfCtx()
defer cancel()
if err := ValidateWorkflow(w); err != nil {
return err
}
normalizeInlineSteps(&w)
_, err := db.Col("workflows").UpdateOne(ctx, bson.M{"workflow_id": id}, bson.M{"$set": bson.M{
"name": w.Name,
"target_server_ids": w.TargetServerIDs,
@@ -162,6 +240,29 @@ func UpdateWorkflow(id string, w models.Workflow) error {
return err
}
// normalizeInlineSteps derives outputs for inline steps and strips fields that
// only belong to library steps.
func normalizeInlineSteps(w *models.Workflow) {
for i := range w.Steps {
in := w.Steps[i].Inline
if in == nil {
continue
}
in.DeclaredOutputs = DeriveOutputs(in.Script)
in.StepID = ""
in.Slug = ""
in.Source = ""
in.CreatedAt = time.Time{}
in.UpdatedAt = time.Time{}
if in.SecretRefs == nil {
in.SecretRefs = []string{}
}
if in.DeclaredInputs == nil {
in.DeclaredInputs = []models.InputParam{}
}
}
}
func DeleteWorkflow(id string) error {
ctx, cancel := wfCtx()
defer cancel()
@@ -0,0 +1,66 @@
package services
import (
"os"
"testing"
"github.com/mrhid6/vantage/server/internal/db"
"github.com/mrhid6/vantage/server/internal/models"
)
var mongoAvailable bool
func TestMain(m *testing.M) {
uri := os.Getenv("VANTAGE_TEST_MONGO_URI")
if uri == "" {
uri = "mongodb://localhost:27117"
}
if err := db.Connect(uri, "vantage_test"); err != nil {
// No MongoDB available in this environment; DB-backed tests will be skipped
// individually, but the rest of the package's tests must still run.
mongoAvailable = false
} else {
mongoAvailable = true
}
os.Exit(m.Run())
}
func mkUsageStep(name string) models.WorkflowStep {
return models.WorkflowStep{Name: name, Interpreter: "bash", Script: "echo hi"}
}
func mkWorkflowWithStep(name, stepID string) models.Workflow {
return models.Workflow{Name: name, Steps: []models.WorkflowStepRef{{StepID: stepID, Order: 0, OnFailure: "stop"}}}
}
func TestStepUsageCounts(t *testing.T) {
if !mongoAvailable {
t.Skip("mongo unavailable: set VANTAGE_TEST_MONGO_URI")
}
// A step used by two workflows, a step used by none.
used, err := CreateStep(mkUsageStep("used-step"))
if err != nil {
t.Fatal(err)
}
unused, err := CreateStep(mkUsageStep("unused-step"))
if err != nil {
t.Fatal(err)
}
if _, err := CreateWorkflow(mkWorkflowWithStep("wf-a", used.StepID)); err != nil {
t.Fatal(err)
}
if _, err := CreateWorkflow(mkWorkflowWithStep("wf-b", used.StepID)); err != nil {
t.Fatal(err)
}
counts, err := StepUsageCounts()
if err != nil {
t.Fatal(err)
}
if counts[used.StepID] != 2 {
t.Fatalf("used step: want 2, got %d", counts[used.StepID])
}
if counts[unused.StepID] != 0 {
t.Fatalf("unused step: want 0, got %d", counts[unused.StepID])
}
}
+19
View File
@@ -36,3 +36,22 @@ body {
::-webkit-scrollbar-thumb:hover {
background: #3e4160;
}
@keyframes led-pulse {
0%, 100% { opacity: 1; }
50% { opacity: 0.5; }
}
.led-pulse { animation: led-pulse 1.4s ease-in-out infinite; }
@keyframes cell-ring {
0%, 100% { box-shadow: 0 0 0 0 rgba(99, 102, 241, 0.5); }
50% { box-shadow: 0 0 0 4px rgba(99, 102, 241, 0); }
}
.cell-ring { animation: cell-ring 1.4s ease-in-out infinite; }
@keyframes caret-blink { 50% { opacity: 0; } }
.caret-blink { animation: caret-blink 1s step-end infinite; }
@media (prefers-reduced-motion: reduce) {
.led-pulse, .cell-ring, .caret-blink { animation: none; }
}
+33 -2
View File
@@ -168,6 +168,9 @@ export default function SettingsPage() {
const [toAddrs, setToAddrs] = useState(""); // comma-separated in UI
const [useTLS, setUseTLS] = useState(false);
// Workflow log retention (days). 0 = keep forever.
const [logRetentionDays, setLogRetentionDays] = useState(30);
const [saved, setSaved] = useState(false);
useEffect(() => {
@@ -183,11 +186,15 @@ export default function SettingsPage() {
setFromAddr(settings.email?.from_addr ?? "");
setToAddrs((settings.email?.to_addrs ?? []).join(", "));
setUseTLS(settings.email?.use_tls ?? false);
setLogRetentionDays(settings.workflow_log_retention_days ?? 30);
}, [settings]);
const { mutate: save, isPending } = useMutation({
mutationFn: (payload: { alerts: AlertSettings; email: EmailSettings }) =>
api.saveSettings(payload),
mutationFn: (payload: {
alerts: AlertSettings;
email: EmailSettings;
workflow_log_retention_days?: number | null;
}) => api.saveSettings(payload),
onSuccess: () => {
queryClient.invalidateQueries({ queryKey: ["settings"] });
setSaved(true);
@@ -218,6 +225,7 @@ export default function SettingsPage() {
to_addrs: toList,
use_tls: useTLS,
},
workflow_log_retention_days: logRetentionDays,
});
}
@@ -374,6 +382,29 @@ export default function SettingsPage() {
</div>
</Card>
{/* Workflow logs */}
<Card>
<CardHeader>
<CardTitle>Workflow Logs</CardTitle>
</CardHeader>
<p className="mb-5 text-sm text-text-secondary">
How long to keep workflow run logs on the server before they are
automatically deleted.
</p>
<Field
label="Log retention (days)"
hint="0 = keep forever. Applies to per-run step output logs."
>
<input
type="number"
min={0}
value={logRetentionDays}
onChange={(e) => setLogRetentionDays(Number(e.target.value))}
className="w-32 rounded-lg border border-border bg-surface-2 px-3 py-2 text-sm text-text-primary focus:border-accent/50 focus:outline-none focus:ring-1 focus:ring-accent/30"
/>
</Field>
</Card>
<div className="flex items-center gap-3">
<Button type="submit" variant="primary" loading={isPending}>
{saved ? "Saved!" : "Save Settings"}
+214
View File
@@ -0,0 +1,214 @@
"use client";
import { useMemo, useRef, useState } from "react";
import { useQuery, useQueryClient } from "@tanstack/react-query";
import { api, WorkflowStep } from "@/lib/api";
import { Button } from "@/components/ui";
import { EditStepModal } from "@/components/workflows/EditStepModal";
type Tab = "all" | "bash" | "powershell" | "default" | "shared";
const inputClass =
"w-full rounded-lg border border-border bg-surface-2 px-3 py-2 text-sm text-text-primary placeholder-text-secondary/50 focus:border-signal focus:outline-none focus:ring-1 focus:ring-signal";
function ShellBadge({ interpreter }: { interpreter: "bash" | "powershell" }) {
const isBash = interpreter === "bash";
return (
<span className={`rounded px-1.5 py-0.5 font-mono text-[10px] uppercase ${isBash ? "bg-bash/15 text-bash" : "bg-pwsh/15 text-pwsh"}`}>
{isBash ? "bash" : "pwsh"}
</span>
);
}
export default function StepsPage() {
const qc = useQueryClient();
const { data: steps } = useQuery({ queryKey: ["steps"], queryFn: api.listSteps });
const { data: usage } = useQuery({ queryKey: ["step-usage"], queryFn: api.stepUsage });
const [search, setSearch] = useState("");
const [tab, setTab] = useState<Tab>("all");
const [editOpen, setEditOpen] = useState(false);
const [editing, setEditing] = useState<WorkflowStep | null>(null);
const [importing, setImporting] = useState(false);
const [syncing, setSyncing] = useState(false);
const [error, setError] = useState<string | null>(null);
const [notice, setNotice] = useState<string | null>(null);
const fileRef = useRef<HTMLInputElement>(null);
const rows = useMemo(() => {
const q = search.toLowerCase();
return (steps ?? []).filter((s) => {
const matchesText = s.name.toLowerCase().includes(q) || (s.description ?? "").toLowerCase().includes(q);
const matchesTab =
tab === "all" ||
(tab === "bash" && s.interpreter === "bash") ||
(tab === "powershell" && s.interpreter === "powershell") ||
(tab === "default" && s.source === "default") ||
(tab === "shared" && s.source !== "default");
return matchesText && matchesTab;
});
}, [steps, search, tab]);
const openNew = () => {
setEditing(null);
setEditOpen(true);
};
const openEdit = (s: WorkflowStep) => {
setEditing(s);
setEditOpen(true);
};
const onImport = async (e: React.ChangeEvent<HTMLInputElement>) => {
const file = e.target.files?.[0];
if (!file) return;
setImporting(true);
setError(null);
try {
const doc = JSON.parse(await file.text());
await api.importStep(doc);
qc.invalidateQueries({ queryKey: ["steps"] });
setNotice("Step imported.");
} catch (err) {
setError((err as Error).message);
} finally {
setImporting(false);
e.target.value = "";
}
};
const onSync = async () => {
setSyncing(true);
setError(null);
try {
const { created, updated } = await api.seedDefaults();
qc.invalidateQueries({ queryKey: ["steps"] });
setNotice(`${created} created, ${updated} updated`);
} catch (err) {
setError((err as Error).message);
} finally {
setSyncing(false);
}
};
return (
<div className="p-8">
<div className="mb-6 flex items-center gap-3">
<div>
<h1 className="text-xl font-semibold text-text-primary">Steps</h1>
<p className="text-sm text-text-secondary">Reusable steps shared across all workflows.</p>
</div>
<div className="ml-auto flex items-center gap-2">
<input ref={fileRef} type="file" accept="application/json" className="hidden" onChange={onImport} />
<Button variant="secondary" size="sm" loading={syncing} onClick={onSync}>
Sync defaults
</Button>
<Button variant="secondary" size="sm" loading={importing} onClick={() => fileRef.current?.click()}>
Import
</Button>
<Button size="sm" onClick={openNew}>
+ New step
</Button>
</div>
</div>
{error && <div className="mb-4 rounded border border-danger/30 bg-danger/10 px-3 py-2 text-sm text-danger">{error}</div>}
{notice && <div className="mb-4 rounded border border-signal/30 bg-signal/10 px-3 py-2 text-sm text-signal">{notice}</div>}
<div className="mb-4 flex items-center gap-2">
<input className={`${inputClass} max-w-sm`} placeholder="Search steps…" value={search} onChange={(e) => setSearch(e.target.value)} />
<div className="flex gap-1.5">
{(["all", "bash", "powershell", "default", "shared"] as Tab[]).map((t) => (
<button
key={t}
onClick={() => setTab(t)}
className={`rounded-full border px-3 py-1 text-xs capitalize ${
tab === t ? "border-signal/50 bg-signal/15 text-signal" : "border-border bg-surface-2 text-text-secondary hover:text-text-primary"
}`}
>
{t === "powershell" ? "PowerShell" : t}
</button>
))}
</div>
</div>
<div className="overflow-x-auto rounded-lg border border-border">
<table className="w-full text-sm">
<thead>
<tr className="border-b border-border text-left text-[11px] uppercase tracking-wide text-text-secondary">
<th className="px-4 py-2.5 font-bold">Name</th>
<th className="px-4 py-2.5 font-bold">Shell</th>
<th className="px-4 py-2.5 font-bold">Source</th>
<th className="px-4 py-2.5 font-bold">Outputs</th>
<th className="px-4 py-2.5 font-bold">Used by</th>
<th className="px-4 py-2.5 text-right font-bold">Actions</th>
</tr>
</thead>
<tbody>
{rows.map((s) => {
const count = usage?.[s.step_id] ?? 0;
return (
<tr key={s.step_id} className="border-b border-border last:border-0">
<td className="px-4 py-3">
<div className="font-medium text-text-primary">{s.name}</div>
{s.description && <div className="text-xs text-text-secondary">{s.description}</div>}
</td>
<td className="px-4 py-3">
<ShellBadge interpreter={s.interpreter} />
</td>
<td className="px-4 py-3">
<span className="rounded bg-surface-2 px-1.5 py-0.5 font-mono text-[10px] uppercase text-text-secondary">
{s.source === "default" ? "default" : "shared"}
</span>
</td>
<td className="px-4 py-3">
<div className="flex flex-wrap gap-1">
{(s.declared_outputs ?? []).map((o) => (
<span key={o} className="rounded border border-signal/35 px-1.5 py-0.5 font-mono text-[10px] text-signal">
{o}
</span>
))}
</div>
</td>
<td className="px-4 py-3 text-text-secondary">
{count === 0 ? "—" : `${count} workflow${count === 1 ? "" : "s"}`}
</td>
<td className="px-4 py-3 text-right">
<div className="flex items-center justify-end gap-3 text-text-secondary">
<button onClick={() => openEdit(s)} className="hover:text-text-primary">
Edit
</button>
<a href={api.exportStepUrl(s.step_id)} download className="hover:text-text-primary">
Export
</a>
<button onClick={() => openEdit(s)} className="hover:text-danger">
Delete
</button>
</div>
</td>
</tr>
);
})}
{rows.length === 0 && (
<tr>
<td colSpan={6} className="px-4 py-8 text-center text-sm text-text-secondary">
No steps found.
</td>
</tr>
)}
</tbody>
</table>
</div>
<EditStepModal
key={editing?.step_id ?? "new"}
open={editOpen}
step={editing}
onClose={() => {
setEditOpen(false);
qc.invalidateQueries({ queryKey: ["steps"] });
qc.invalidateQueries({ queryKey: ["step-usage"] });
}}
/>
</div>
);
}
File diff suppressed because it is too large Load Diff
+425 -85
View File
@@ -1,102 +1,442 @@
"use client";
import { useParams } from "next/navigation";
import { useEffect, useMemo, useRef, useState } from "react";
import { useQuery, useQueryClient } from "@tanstack/react-query";
import { api, ServerRun, StepRun } from "@/lib/api";
import { Button, Badge, Card } from "@/components/ui";
import { api, ServerRun, StepRun, WorkflowRun } from "@/lib/api";
import { Button } from "@/components/ui";
type BadgeVariant = "success" | "warning" | "danger" | "neutral" | "accent";
// ---- status vocabulary ----------------------------------------------------
const statusVariant: Record<string, BadgeVariant> = {
success: "success",
failed: "danger",
running: "accent",
queued: "neutral",
skipped: "neutral",
cancelled: "warning",
type CellKind = "done" | "fail" | "run" | "wait" | "skip" | "warn";
function cellKind(status: string): CellKind {
switch (status) {
case "success":
return "done";
case "failed":
return "fail";
case "running":
return "run";
case "skipped":
return "skip";
case "cancelled":
return "warn";
default:
return "wait"; // queued / pending / missing
}
}
const cellGlyph: Record<CellKind, string> = {
done: "✓",
fail: "✕",
run: "●",
wait: "○",
skip: "",
warn: "!",
};
function StatusBadge({ status }: { status: string }) {
return <Badge variant={statusVariant[status] ?? "neutral"}>{status}</Badge>;
const cellClass: Record<CellKind, string> = {
done: "bg-success/15 text-success",
fail: "bg-danger/15 text-danger",
run: "bg-accent/15 text-accent",
wait: "text-border",
skip: "text-text-secondary",
warn: "bg-warning/15 text-warning",
};
// ---- run-level status pill ------------------------------------------------
type PillKind = "running" | "success" | "failed" | "neutral";
function pillKind(status: string): PillKind {
if (status === "running") return "running";
if (status === "success") return "success";
if (status === "failed" || status === "cancelled") return "failed";
return "neutral";
}
const pillClass: Record<PillKind, string> = {
running: "text-accent border-accent/40 bg-accent/10",
success: "text-success border-success/35 bg-success/10",
failed: "text-danger border-danger/35 bg-danger/10",
neutral: "text-text-secondary border-border bg-surface-2",
};
const pillLed: Record<PillKind, string> = {
running: "bg-accent led-pulse",
success: "bg-success",
failed: "bg-danger",
neutral: "bg-text-secondary",
};
function StatusPill({ status, small }: { status: string; small?: boolean }) {
const kind = pillKind(status);
return (
<span
className={`inline-flex items-center gap-2 rounded-full border font-mono font-semibold uppercase tracking-wide ${
small ? "px-2 py-0.5 text-[10px]" : "px-2.5 py-1 text-xs"
} ${pillClass[kind]}`}
>
<span className={`h-1.5 w-1.5 rounded-full ${pillLed[kind]}`} />
{status}
</span>
);
}
// ---- time helpers ---------------------------------------------------------
function fmtDuration(ms: number): string {
if (ms < 0) ms = 0;
const s = Math.floor(ms / 1000);
if (s < 60) return `${s}s`;
const m = Math.floor(s / 60);
const rem = s % 60;
if (m < 60) return `${m}m ${rem}s`;
const h = Math.floor(m / 60);
return `${h}h ${m % 60}m`;
}
function stepDuration(st: StepRun, running: boolean, now: number): string {
if (!st.started_at) return st.status === "queued" ? "queued" : "";
const start = new Date(st.started_at).getTime();
const end = st.finished_at ? new Date(st.finished_at).getTime() : running ? now : start;
return fmtDuration(end - start);
}
// ---- live log terminal ----------------------------------------------------
function LogTerminal({ runId, server }: { runId: string; server: ServerRun }) {
const [text, setText] = useState("");
const preRef = useRef<HTMLDivElement>(null);
const running = server.status === "running";
const serverId = server.server_id;
useEffect(() => {
setText("");
if (running) {
const es = new EventSource(api.serverRunLogStreamUrl(runId, serverId), {
withCredentials: true,
});
es.onmessage = (e) => setText((t) => t + e.data + "\n");
es.addEventListener("done", () => es.close());
es.onerror = () => es.close();
return () => es.close();
}
api.getServerRunLog(runId, serverId)
.then(setText)
.catch(() => setText(""));
}, [running, runId, serverId]);
useEffect(() => {
preRef.current?.scrollTo(0, preRef.current.scrollHeight);
}, [text]);
const activeStep = server.steps.find((s) => s.status === "running") ?? [...server.steps].reverse().find((s) => s.started_at);
return (
<div className="overflow-hidden rounded-xl border border-border bg-[#0a0b10]">
<div className="flex items-center justify-between gap-2 border-b border-border bg-surface px-4 py-3">
<span className="truncate font-mono text-[13px] font-semibold text-text-primary">
{activeStep ? activeStep.name : "Output"} <span className="font-normal text-text-secondary">{server.hostname}</span>
</span>
{running && (
<span className="inline-flex items-center gap-1.5 font-mono text-[10.5px] uppercase tracking-wide text-accent">
<span className="h-1.5 w-1.5 rounded-full bg-accent led-pulse" />
Streaming
</span>
)}
</div>
<div ref={preRef} className="max-h-[340px] overflow-auto whitespace-pre-wrap px-4 py-3.5 font-mono text-[12.5px] leading-relaxed text-text-secondary">
{text ? <LogLines text={text} /> : running ? "Waiting for output…" : "No output."}
{running && text && <span className="ml-0.5 inline-block h-3.5 w-[7px] translate-y-[2px] bg-accent caret-blink align-baseline" />}
</div>
</div>
);
}
// LogLines renders the raw server-run log, parsing each line's leading UTC
// timestamp ([2026-07-20T12:04:02.000Z]) and rendering it in the viewer's local
// timezone. Event markers (===== …) are highlighted so the run's shape scans.
const TS_RE = /^\[(\d{4}-\d{2}-\d{2}T[\d:.]+Z)\]\s?(.*)$/;
function LogLines({ text }: { text: string }) {
const lines = text.replace(/\n$/, "").split("\n");
return (
<>
{lines.map((line, i) => {
const m = TS_RE.exec(line);
if (!m) {
return (
<span key={i} className="block">
{line || " "}
</span>
);
}
const local = new Date(m[1]).toLocaleTimeString([], { hour12: false });
const body = m[2];
const isMarker = body.startsWith("=====");
return (
<span key={i} className="block">
<span className="select-none text-[#565b74]" title={m[1]}>
{local}{" "}
</span>
<span className={isMarker ? "font-semibold text-accent" : ""}>{body || " "}</span>
</span>
);
})}
</>
);
}
// ---- step list ------------------------------------------------------------
function StepList({ server, now }: { server: ServerRun; now: number }) {
const running = server.status === "running";
return (
<div className="overflow-hidden rounded-xl border border-border bg-surface">
<div className="flex items-center justify-between gap-2 border-b border-border px-4 py-3">
<span className="font-mono text-[13px] font-semibold text-text-primary">
Steps <span className="font-normal text-text-secondary">{server.steps.length}</span>
</span>
<StatusPill status={server.status} small />
</div>
<div className="flex flex-col gap-0.5 p-1.5">
{server.steps.map((st) => {
const kind = cellKind(st.status);
return (
<div
key={st.order}
className={`grid grid-cols-[20px_1fr_auto] items-center gap-2.5 rounded-lg px-3 py-2.5 text-[13px] hover:bg-surface-2 ${st.status === "running" ? "bg-accent/[0.06]" : ""}`}
>
<span className="text-right font-mono text-[11px] text-text-secondary">{String(st.order + 1).padStart(2, "0")}</span>
<span className="flex items-center gap-2 font-medium text-text-primary">
<span className={`font-mono ${cellClass[kind].replace(/bg-\S+/, "")}`}>{cellGlyph[kind]}</span>
{st.name}
</span>
<span className="text-right font-mono text-[10.5px] text-text-secondary">
{st.status === "failed" && <span className="text-danger">exit {st.exit_code} · </span>}
{st.attempts > 1 ? `${st.attempts} tries` : "1 try"}
{stepDuration(st, running, now) ? ` · ${stepDuration(st, running, now)}` : ""}
</span>
</div>
);
})}
{server.steps.length === 0 && <p className="px-3 py-2 text-xs text-text-secondary">No steps yet.</p>}
</div>
</div>
);
}
// ---- execution matrix (signature) -----------------------------------------
interface Column {
order: number;
name: string;
}
function buildColumns(run: WorkflowRun): Column[] {
const byOrder = new Map<number, string>();
for (const sr of run.server_runs) {
for (const st of sr.steps) {
if (!byOrder.has(st.order)) byOrder.set(st.order, st.name);
}
}
return [...byOrder.entries()].map(([order, name]) => ({ order, name })).sort((a, b) => a.order - b.order);
}
function ExecutionMatrix({ run, columns, selected, onSelect }: { run: WorkflowRun; columns: Column[]; selected: string; onSelect: (serverId: string) => void }) {
return (
<div className="overflow-x-auto rounded-xl border border-border bg-surface">
<table className="w-full border-collapse font-mono text-[12.5px]">
<thead>
<tr>
<th className="border-b border-border px-4 py-3 text-left align-bottom text-xs font-semibold uppercase tracking-wider text-text-primary">Server</th>
{columns.map((c) => (
<th key={c.order} className="whitespace-nowrap border-b border-border px-3.5 py-3 align-bottom text-[11px] font-medium text-text-secondary">
<span className="block text-[10px] text-border">{String(c.order + 1).padStart(2, "0")}</span>
{c.name}
</th>
))}
</tr>
</thead>
<tbody>
{run.server_runs.map((sr) => {
const byOrder = new Map(sr.steps.map((s) => [s.order, s]));
const isSel = sr.server_id === selected;
return (
<tr key={sr.server_id} onClick={() => onSelect(sr.server_id)} className={`cursor-pointer ${isSel ? "bg-accent/5" : "hover:bg-white/[0.02]"}`}>
<th className="min-w-[240px] border-b border-r border-border px-4 py-3 text-left font-medium text-text-primary">
<div className="flex items-center gap-3">
<span className="flex-1 whitespace-nowrap">{sr.hostname}</span>
<StatusPill status={sr.status} small />
</div>
</th>
{columns.map((c) => {
const st = byOrder.get(c.order);
const kind = st ? cellKind(st.status) : "wait";
return (
<td key={c.order} className="relative border-b border-r border-border last:border-r-0">
<span className="flex h-[54px] items-center justify-center">
<span className={`relative flex h-[26px] w-[26px] items-center justify-center rounded-md ${cellClass[kind]}`}>
{kind === "run" && <span className="absolute inset-0 rounded-md border border-accent/50 cell-ring" />}
{cellGlyph[kind]}
</span>
</span>
</td>
);
})}
</tr>
);
})}
</tbody>
</table>
</div>
);
}
// ---- page -----------------------------------------------------------------
function SectionLabel({ children }: { children: React.ReactNode }) {
return (
<div className="mb-3 mt-8 flex items-center gap-2.5 font-mono text-[11px] uppercase tracking-widest text-text-secondary">
{children}
<span className="h-px flex-1 bg-border" />
</div>
);
}
export default function RunDetail() {
const { runId } = useParams<{ runId: string }>();
const queryClient = useQueryClient();
const { runId } = useParams<{ runId: string }>();
const queryClient = useQueryClient();
const [selected, setSelected] = useState<string | null>(null);
const [now, setNow] = useState(() => Date.now());
const { data: run, isLoading } = useQuery({
queryKey: ["run", runId],
queryFn: () => api.getRun(runId),
refetchInterval: (query) => (query.state.data?.status === "running" ? 2000 : false),
});
const { data: run, isLoading } = useQuery({
queryKey: ["run", runId],
queryFn: () => api.getRun(runId),
refetchInterval: (query) => (query.state.data?.status === "running" ? 2000 : false),
});
const cancel = async () => {
await api.cancelRun(runId);
queryClient.invalidateQueries({ queryKey: ["run", runId] });
};
const running = run?.status === "running";
if (isLoading || !run) {
return <div className="p-8 text-text-secondary">Loading</div>;
}
// tick the elapsed clock while running
useEffect(() => {
if (!running) return;
const t = setInterval(() => setNow(Date.now()), 1000);
return () => clearInterval(t);
}, [running]);
return (
<div className="p-8">
<div className="mb-6 flex items-center justify-between">
<div>
<h1 className="text-2xl font-bold text-text-primary">{run.name}</h1>
<p className="mt-1 flex items-center gap-2 text-sm text-text-secondary">
<span>Run {run.run_id.slice(0, 8)}</span>
<StatusBadge status={run.status} />
</p>
const columns = useMemo(() => (run ? buildColumns(run) : []), [run]);
// default selection: first running server, else first server
const selectedServer = useMemo(() => {
if (!run || run.server_runs.length === 0) return null;
if (selected) {
const match = run.server_runs.find((s) => s.server_id === selected);
if (match) return match;
}
return run.server_runs.find((s) => s.status === "running") ?? run.server_runs[0];
}, [run, selected]);
const cancel = async () => {
await api.cancelRun(runId);
queryClient.invalidateQueries({ queryKey: ["run", runId] });
};
if (isLoading || !run) {
return <div className="p-8 text-text-secondary">Loading</div>;
}
const totalSteps = run.server_runs.reduce((n, s) => n + s.steps.length, 0);
const doneSteps = run.server_runs.reduce((n, s) => n + s.steps.filter((st) => st.status === "success").length, 0);
const succeeded = run.server_runs.filter((s) => s.status === "success").length;
const failed = run.server_runs.filter((s) => s.status === "failed" || s.status === "cancelled").length;
const startMs = run.started_at ? new Date(run.started_at).getTime() : now;
const endMs = run.finished_at ? new Date(run.finished_at).getTime() : now;
const elapsed = fmtDuration(endMs - startMs);
const ago = fmtDuration(now - startMs);
return (
<div className="mx-auto max-w-[1180px] p-8 pb-16">
{/* identity bar */}
<div className="flex flex-wrap items-start justify-between gap-6">
<div>
<div className="mb-2 font-mono text-xs uppercase tracking-wide text-text-secondary">Workflows / {run.name} / Runs</div>
<h1 className="text-[28px] font-semibold tracking-tight text-text-primary">{run.name}</h1>
<div className="mt-2.5 flex flex-wrap items-center gap-x-4 gap-y-1 font-mono text-[12.5px] text-text-secondary">
<span>
run <b className="font-medium text-text-primary">{run.run_id.slice(0, 8)}</b>
</span>
<span className="h-[3px] w-[3px] rounded-full bg-border" />
<span>
triggered by <b className="font-medium text-text-primary">{run.triggered_by || "—"}</b>
</span>
<span className="h-[3px] w-[3px] rounded-full bg-border" />
<span>
started <b className="font-medium text-text-primary">{ago}</b> ago
</span>
<span className="h-[3px] w-[3px] rounded-full bg-border" />
<span>
elapsed <b className="font-medium text-text-primary">{elapsed}</b>
</span>
</div>
</div>
<div className="flex items-center gap-3.5">
<StatusPill status={run.status} />
{running && (
<Button variant="danger" onClick={cancel}>
Cancel run
</Button>
)}
</div>
</div>
{/* summary strip */}
<div className="mt-6 grid grid-cols-2 gap-px overflow-hidden rounded-xl border border-border bg-border sm:grid-cols-4">
<div className="bg-surface px-[18px] py-4">
<div className="font-mono text-[10.5px] uppercase tracking-wider text-text-secondary">Servers</div>
<div className="mt-1 font-mono text-[22px] font-semibold tabular-nums text-text-primary">{run.server_runs.length}</div>
</div>
<div className="bg-surface px-[18px] py-4">
<div className="font-mono text-[10.5px] uppercase tracking-wider text-text-secondary">Succeeded</div>
<div className="mt-1 font-mono text-[22px] font-semibold tabular-nums text-success">
{succeeded}
<small className="text-sm font-medium text-text-secondary"> / {run.server_runs.length}</small>
</div>
</div>
<div className="bg-surface px-[18px] py-4">
<div className="font-mono text-[10.5px] uppercase tracking-wider text-text-secondary">Failed</div>
<div className={`mt-1 font-mono text-[22px] font-semibold tabular-nums ${failed > 0 ? "text-danger" : "text-text-primary"}`}>{failed}</div>
</div>
<div className="bg-surface px-[18px] py-4">
<div className="font-mono text-[10.5px] uppercase tracking-wider text-text-secondary">Steps done</div>
<div className="mt-1 font-mono text-[22px] font-semibold tabular-nums text-text-primary">
{doneSteps}
<small className="text-sm font-medium text-text-secondary"> / {totalSteps}</small>
</div>
</div>
</div>
{run.server_runs.length === 0 ? (
<p className="mt-8 text-text-secondary">No servers targeted by this run.</p>
) : (
<>
<SectionLabel>Execution matrix</SectionLabel>
<ExecutionMatrix run={run} columns={columns} selected={selectedServer?.server_id ?? ""} onSelect={setSelected} />
{selectedServer && (
<>
<SectionLabel>{selectedServer.hostname}&nbsp;·&nbsp;steps &amp; live output</SectionLabel>
<div className="grid grid-cols-1 items-start gap-4 md:grid-cols-[320px_1fr]">
<StepList server={selectedServer} now={now} />
<LogTerminal runId={run.run_id} server={selectedServer} />
</div>
</>
)}
</>
)}
</div>
{run.status === "running" && (
<Button variant="danger" onClick={cancel}>
Cancel
</Button>
)}
</div>
<div className="grid gap-4 md:grid-cols-2 lg:grid-cols-3">
{run.server_runs.map((sr: ServerRun) => (
<Card key={sr.server_id}>
<div className="mb-3 flex items-center justify-between">
<span className="font-medium text-text-primary">{sr.hostname}</span>
<StatusBadge status={sr.status} />
</div>
<div className="space-y-2">
{sr.steps.map((st: StepRun) => (
<details
key={st.order}
className="rounded-lg border border-border bg-surface-2 p-2"
>
<summary className="flex cursor-pointer items-center justify-between gap-2">
<span className="text-sm text-text-primary">{st.name}</span>
<span className="flex items-center gap-2">
<span className="text-xs text-text-secondary">
attempts: {st.attempts}
{st.status === "failed" ? ` · exit ${st.exit_code}` : ""}
</span>
<StatusBadge status={st.status} />
</span>
</summary>
{(st.stdout || st.stderr) && (
<pre className="mt-2 max-h-64 overflow-auto rounded bg-black/40 p-2 font-mono text-xs text-text-secondary">
{st.stdout}
{st.stderr ? `\n${st.stderr}` : ""}
</pre>
)}
</details>
))}
{sr.steps.length === 0 && (
<p className="text-xs text-text-secondary">No steps yet.</p>
)}
</div>
</Card>
))}
{run.server_runs.length === 0 && (
<p className="text-text-secondary">No servers targeted by this run.</p>
)}
</div>
</div>
);
);
}
+78
View File
@@ -0,0 +1,78 @@
"use client";
import Link from "next/link";
import { useParams } from "next/navigation";
import { useQuery } from "@tanstack/react-query";
import { api, WorkflowRun } from "@/lib/api";
import { Card, Table, Thead, Tbody, Tr, Th, Td, Badge } from "@/components/ui";
type BadgeVariant = "success" | "warning" | "danger" | "neutral" | "accent";
const statusVariant: Record<string, BadgeVariant> = {
success: "success",
failed: "danger",
running: "warning",
cancelled: "neutral",
queued: "neutral",
};
export default function WorkflowRunsPage() {
const { id } = useParams<{ id: string }>();
const { data: wf } = useQuery({ queryKey: ["workflow", id], queryFn: () => api.getWorkflow(id) });
const { data: runs, isLoading, error } = useQuery({ queryKey: ["runs", id], queryFn: () => api.listRuns(id) });
return (
<div className="p-8">
<div className="mb-6">
<Link href={`/workflows/${id}`} className="text-sm text-text-secondary hover:text-text-primary">
Back to builder
</Link>
<h1 className="mt-2 text-2xl font-bold text-text-primary">Runs · {wf?.name ?? ""}</h1>
</div>
<Card padding={false}>
{isLoading ? (
<div className="flex items-center justify-center py-20">
<div className="h-8 w-8 animate-spin rounded-full border-2 border-border border-t-accent" />
</div>
) : error ? (
<div className="py-20 text-center text-danger">Failed to load runs. Is the backend running?</div>
) : runs && runs.length > 0 ? (
<Table>
<Thead>
<Tr>
<Th>Run</Th>
<Th>Status</Th>
<Th>Started</Th>
<Th>By</Th>
<Th>Servers</Th>
</Tr>
</Thead>
<Tbody>
{runs.map((r: WorkflowRun) => (
<Tr key={r.run_id}>
<Td>
<Link
href={`/workflows/${id}/runs/${r.run_id}`}
className="font-mono text-text-primary hover:text-signal"
>
{r.run_id.slice(0, 8)}
</Link>
</Td>
<Td>
<Badge variant={statusVariant[r.status] ?? "neutral"}>{r.status}</Badge>
</Td>
<Td className="text-text-secondary">{new Date(r.started_at).toLocaleString()}</Td>
<Td className="text-text-secondary">{r.triggered_by}</Td>
<Td className="text-text-secondary">{r.server_runs.length}</Td>
</Tr>
))}
</Tbody>
</Table>
) : (
<div className="py-16 text-center text-text-secondary">No runs yet.</div>
)}
</Card>
</div>
);
}
+8 -3
View File
@@ -82,9 +82,14 @@ export default function WorkflowsPage() {
<span className="text-text-secondary">{w.steps.length}</span>
</Td>
<Td>
<Link href={`/workflows/${w.workflow_id}`}>
<Button variant="ghost" size="sm">Open </Button>
</Link>
<div className="flex items-center justify-end gap-2">
<Link href={`/workflows/${w.workflow_id}/runs`}>
<Button variant="ghost" size="sm">Runs</Button>
</Link>
<Link href={`/workflows/${w.workflow_id}`}>
<Button variant="ghost" size="sm">Open </Button>
</Link>
</div>
</Td>
</Tr>
))}
+9
View File
@@ -60,11 +60,20 @@ function SettingsIcon() {
);
}
function StepsIcon() {
return (
<svg className="h-5 w-5" fill="none" viewBox="0 0 24 24" stroke="currentColor" strokeWidth={1.5}>
<path strokeLinecap="round" strokeLinejoin="round" d="M6 6.75A.75.75 0 016.75 6h10.5a.75.75 0 010 1.5H6.75A.75.75 0 016 6.75zm0 5.25a.75.75 0 01.75-.75h10.5a.75.75 0 010 1.5H6.75A.75.75 0 016 12zm0 5.25a.75.75 0 01.75-.75h10.5a.75.75 0 010 1.5H6.75A.75.75 0 016 17.25zM3 6.75a.75.75 0 11-1.5 0 .75.75 0 011.5 0zM3 12a.75.75 0 11-1.5 0 .75.75 0 011.5 0zm0 5.25a.75.75 0 11-1.5 0 .75.75 0 011.5 0z" />
</svg>
);
}
const navItems: NavItem[] = [
{ href: "/servers", label: "Servers", icon: <ServerIcon /> },
{ href: "/keys", label: "SSH Keys", icon: <KeyIcon /> },
{ href: "/secrets", label: "Secrets", icon: <SecretIcon /> },
{ href: "/workflows", label: "Workflows", icon: <WorkflowIcon /> },
{ href: "/steps", label: "Steps", icon: <StepsIcon /> },
{ href: "/audit", label: "Audit Log", icon: <AuditIcon /> },
{ href: "/settings", label: "Settings", icon: <SettingsIcon /> },
];
+45
View File
@@ -0,0 +1,45 @@
"use client";
import { useEffect } from "react";
export function Modal({
open,
title,
onClose,
children,
wide,
}: {
open: boolean;
title: string;
onClose: () => void;
children: React.ReactNode;
wide?: boolean;
}) {
useEffect(() => {
if (!open) return;
const onKey = (e: KeyboardEvent) => e.key === "Escape" && onClose();
window.addEventListener("keydown", onKey);
return () => window.removeEventListener("keydown", onKey);
}, [open, onClose]);
if (!open) return null;
return (
<div className="fixed inset-0 z-50 flex items-center justify-center p-4">
<div className="absolute inset-0 bg-black/60" onClick={onClose} />
<div
className={`relative z-10 w-full ${wide ? "max-w-2xl" : "max-w-md"} max-h-[90vh] overflow-auto rounded-xl border border-border bg-surface shadow-2xl`}
role="dialog"
aria-modal="true"
>
<div className="flex items-center justify-between border-b border-border px-5 py-3">
<h2 className="text-sm font-bold text-text-primary">{title}</h2>
<button onClick={onClose} className="text-text-secondary hover:text-text-primary" aria-label="Close">
</button>
</div>
<div className="p-5">{children}</div>
</div>
</div>
);
}
+1
View File
@@ -2,3 +2,4 @@ export { Button } from "./Button";
export { Badge } from "./Badge";
export { Card, CardHeader, CardTitle } from "./Card";
export { Table, Thead, Tbody, Tr, Th, Td } from "./Table";
export { Modal } from "./Modal";
+108
View File
@@ -0,0 +1,108 @@
"use client";
import { useState } from "react";
import { useQueryClient } from "@tanstack/react-query";
import { api, WorkflowStep, InputParam } from "@/lib/api";
import { Button, Modal } from "@/components/ui";
const inputClass =
"w-full rounded-lg border border-border bg-surface-2 px-3 py-2 text-sm text-text-primary placeholder-text-secondary/50 focus:border-signal focus:outline-none focus:ring-1 focus:ring-signal";
export function EditStepModal({ open, step, onClose }: { open: boolean; step: WorkflowStep | null; onClose: () => void }) {
const qc = useQueryClient();
const [name, setName] = useState(step?.name ?? "");
const [interpreter, setInterpreter] = useState<"bash" | "powershell">(step?.interpreter ?? "bash");
const [script, setScript] = useState(step?.script ?? "");
const [inputs, setInputs] = useState<InputParam[]>(step?.declared_inputs ?? []);
const [busy, setBusy] = useState(false);
const [error, setError] = useState<string | null>(null);
// NOTE: because state is seeded from props, render the modal conditionally
// (parent mounts it only when opening) OR key it by step_id so it re-seeds.
const save = async () => {
setBusy(true); setError(null);
try {
const payload: Partial<WorkflowStep> = {
name: name.trim(), description: step?.description ?? "", interpreter, script,
declared_inputs: inputs.filter((i) => i.name.trim() !== ""),
secret_refs: step?.secret_refs ?? [],
};
if (step) await api.updateStep(step.step_id, payload);
else await api.createStep(payload);
qc.invalidateQueries({ queryKey: ["steps"] });
onClose();
} catch (e) { setError((e as Error).message); } finally { setBusy(false); }
};
const del = async () => {
if (!step || !window.confirm("Delete this step? It will be removed from every workflow that uses it.")) return;
setBusy(true); setError(null);
try {
await api.deleteStep(step.step_id);
qc.invalidateQueries({ queryKey: ["steps"] });
qc.invalidateQueries({ queryKey: ["workflow"] });
onClose();
} catch (e) { setError((e as Error).message); } finally { setBusy(false); }
};
return (
<Modal open={open} onClose={onClose} title={step ? "Edit base step" : "New step"} wide>
<div className="space-y-4">
{error && <div className="rounded border border-danger/30 bg-danger/10 px-3 py-2 text-sm text-danger">{error}</div>}
<p className="text-xs text-text-secondary">Reusable steps are shared across all workflows. Editing here changes it everywhere.</p>
<div>
<label className="mb-1 block text-xs uppercase text-text-secondary">Name</label>
<input className={inputClass} value={name} onChange={(e) => setName(e.target.value)} />
</div>
<div>
<label className="mb-1 block text-xs uppercase text-text-secondary">Interpreter</label>
<select className={inputClass} value={interpreter} onChange={(e) => setInterpreter(e.target.value as "bash" | "powershell")}>
<option value="bash">bash</option>
<option value="powershell">powershell</option>
</select>
</div>
<div>
<label className="mb-1 block text-xs uppercase text-text-secondary">Script</label>
<textarea className={`${inputClass} h-40 font-mono text-xs`} value={script} onChange={(e) => setScript(e.target.value)} />
<p className="mt-1 text-xs text-text-secondary">Write <code className="text-signal">KEY=value</code> to <code className="text-signal">$WORKFLOW_ENV</code> to expose it to later steps.</p>
</div>
<div>
<label className="mb-1 block text-xs uppercase text-text-secondary">Outputs</label>
<div className="mb-1 flex flex-wrap gap-1">
{(step?.declared_outputs ?? []).length === 0 && (
<p className="text-xs text-text-secondary">No declared outputs.</p>
)}
{(step?.declared_outputs ?? []).map((o) => (
<span key={o} className="flex items-center gap-1 rounded bg-signal px-2 py-0.5 font-mono text-[11px] text-signal-ink">
{o}
</span>
))}
</div>
<p className="text-xs text-text-secondary">Outputs are detected automatically from lines writing to $WORKFLOW_ENV.</p>
</div>
<div>
<label className="mb-1 block text-xs uppercase text-text-secondary">Inputs</label>
<div className="space-y-2">
{inputs.map((inp, i) => (
<div key={i} className="flex gap-2">
<input className={inputClass} placeholder="name" value={inp.name} onChange={(e) => setInputs(inputs.map((x, j) => j === i ? { ...x, name: e.target.value } : x))} />
<input className={inputClass} placeholder="default" value={inp.default} onChange={(e) => setInputs(inputs.map((x, j) => j === i ? { ...x, default: e.target.value } : x))} />
<input className={inputClass} placeholder="description" value={inp.description} onChange={(e) => setInputs(inputs.map((x, j) => j === i ? { ...x, description: e.target.value } : x))} />
<Button variant="ghost" size="sm" onClick={() => setInputs(inputs.filter((_, j) => j !== i))}></Button>
</div>
))}
</div>
<Button variant="ghost" size="sm" className="mt-2" onClick={() => setInputs([...inputs, { name: "", default: "", description: "" }])}>Add input</Button>
</div>
<div className="flex items-center justify-between pt-2">
{step ? <Button variant="danger" onClick={del} loading={busy}>Delete step</Button> : <span />}
<div className="flex gap-2">
<Button variant="ghost" onClick={onClose}>Cancel</Button>
<Button variant="primary" onClick={save} loading={busy} disabled={!name.trim()}>Save</Button>
</div>
</div>
</div>
</Modal>
);
}
@@ -0,0 +1,77 @@
"use client";
import { useState, useEffect } from "react";
import { useRouter } from "next/navigation";
import { useQuery } from "@tanstack/react-query";
import { api, Workflow } from "@/lib/api";
import { Button, Modal } from "@/components/ui";
const inputClass =
"w-full rounded-lg border border-border bg-surface-2 px-3 py-2 text-sm text-text-primary focus:border-signal focus:outline-none focus:ring-1 focus:ring-signal";
export function EditWorkflowModal({ open, workflow, onSaved, onClose }: { open: boolean; workflow: Workflow; onSaved: (w: Workflow) => void; onClose: () => void }) {
const router = useRouter();
const [name, setName] = useState(workflow.name);
const [targets, setTargets] = useState<string[]>(workflow.target_server_ids);
const [busy, setBusy] = useState(false);
const [error, setError] = useState<string | null>(null);
const { data: servers } = useQuery({ queryKey: ["servers"], queryFn: api.listServers });
useEffect(() => {
if (open) {
setName(workflow.name);
setTargets(workflow.target_server_ids);
}
}, [open, workflow]);
const toggle = (id: string) => setTargets((t) => (t.includes(id) ? t.filter((x) => x !== id) : [...t, id]));
const save = async () => {
setBusy(true); setError(null);
try {
const updated = await api.updateWorkflow(workflow.workflow_id, { ...workflow, name, target_server_ids: targets });
onSaved(updated); onClose();
} catch (e) { setError((e as Error).message); } finally { setBusy(false); }
};
const del = async () => {
if (!window.confirm("Delete this workflow? This cannot be undone.")) return;
setBusy(true); setError(null);
try { await api.deleteWorkflow(workflow.workflow_id); router.push("/workflows"); }
catch (e) { setError((e as Error).message); setBusy(false); }
};
return (
<Modal open={open} onClose={onClose} title="Edit workflow">
<div className="space-y-4">
{error && <div className="rounded border border-danger/30 bg-danger/10 px-3 py-2 text-sm text-danger">{error}</div>}
<div>
<label className="mb-1 block text-xs uppercase text-text-secondary">Name</label>
<input className={inputClass} value={name} onChange={(e) => setName(e.target.value)} />
</div>
<div>
<label className="mb-1 block text-xs uppercase text-text-secondary">Target servers</label>
<div className="flex flex-wrap gap-2">
{servers?.map((s) => {
const on = targets.includes(s.server_id);
return (
<label key={s.server_id} className={`flex cursor-pointer items-center gap-2 rounded-lg border px-2 py-1 text-sm ${on ? "border-signal bg-signal/10 text-text-primary" : "border-border text-text-secondary"}`}>
<input type="checkbox" className="accent-signal" checked={on} onChange={() => toggle(s.server_id)} />
{s.hostname}
</label>
);
})}
{servers && servers.length === 0 && <p className="text-xs text-text-secondary">No servers registered.</p>}
</div>
</div>
<div className="flex items-center justify-between pt-2">
<Button variant="danger" onClick={del} loading={busy}>Delete workflow</Button>
<div className="flex gap-2">
<Button variant="ghost" onClick={onClose}>Cancel</Button>
<Button variant="primary" onClick={save} loading={busy} disabled={!name.trim()}>Save</Button>
</div>
</div>
</div>
</Modal>
);
}
@@ -0,0 +1,191 @@
"use client";
import { useMemo, useRef, useState } from "react";
import { useQuery } from "@tanstack/react-query";
import { api, WorkflowStep } from "@/lib/api";
import { Modal } from "@/components/ui";
type Tab = "all" | "bash" | "powershell" | "adhoc";
function ShellBadge({ interpreter }: { interpreter: "bash" | "powershell" }) {
const isBash = interpreter === "bash";
return (
<span
className={`rounded px-1.5 py-0.5 font-mono text-[10px] uppercase ${
isBash ? "bg-bash/15 text-bash" : "bg-pwsh/15 text-pwsh"
}`}
>
{isBash ? "bash" : "pwsh"}
</span>
);
}
function DefaultBadge() {
return (
<span className="rounded bg-surface-2 px-1.5 py-0.5 font-mono text-[10px] uppercase text-text-secondary">
default
</span>
);
}
function StepCard({ step, onAdd }: { step: WorkflowStep; onAdd: () => void }) {
return (
<button
onClick={onAdd}
className="group relative rounded-[10px] border border-border bg-surface-2 p-3 text-left transition-colors hover:border-signal/55"
>
<span className="absolute right-3 top-3 text-xs font-semibold text-signal opacity-0 group-hover:opacity-100">
+ Add
</span>
<div className="mb-1.5 flex items-center gap-2">
<ShellBadge interpreter={step.interpreter} />
{step.source === "default" && <DefaultBadge />}
<span className="text-sm font-medium text-text-primary">{step.name}</span>
</div>
{step.description && <p className="text-xs text-text-secondary">{step.description}</p>}
<div className="mt-2 flex flex-wrap gap-1.5">
{(step.declared_inputs ?? []).map((p) => (
<span key={p.name} className="rounded border border-border bg-background px-1.5 py-0.5 font-mono text-[10px] text-text-secondary">
in {p.name}
</span>
))}
{(step.declared_outputs ?? []).map((o) => (
<span key={o} className="rounded border border-signal/35 bg-background px-1.5 py-0.5 font-mono text-[10px] text-signal">
out {o}
</span>
))}
</div>
</button>
);
}
export function StepPickerModal({
open,
onClose,
onSelect,
onAddAdhoc,
onImportAdhoc,
}: {
open: boolean;
onClose: () => void;
onSelect: (stepId: string) => void;
onAddAdhoc: () => void;
onImportAdhoc: (file: File) => void;
}) {
const { data: library } = useQuery({ queryKey: ["steps"], queryFn: api.listSteps });
const [search, setSearch] = useState("");
const [tab, setTab] = useState<Tab>("all");
const fileRef = useRef<HTMLInputElement>(null);
const filtered = useMemo(() => {
const q = search.toLowerCase();
return (library ?? []).filter(
(s) =>
(s.name.toLowerCase().includes(q) || (s.description ?? "").toLowerCase().includes(q)) &&
(tab === "all" || tab === "adhoc" ? true : s.interpreter === tab),
);
}, [library, search, tab]);
const group = (source: "default" | "shared", interp: "bash" | "powershell") =>
filtered.filter(
(s) => s.interpreter === interp && (source === "default" ? s.source === "default" : s.source !== "default"),
);
const groups: { label: string; steps: WorkflowStep[] }[] = [
{ label: "Default · Bash", steps: group("default", "bash") },
{ label: "Default · PowerShell", steps: group("default", "powershell") },
{ label: "Shared · Bash", steps: group("shared", "bash") },
{ label: "Shared · PowerShell", steps: group("shared", "powershell") },
];
const showLibrary = tab !== "adhoc";
const showAdhocCards = tab === "all" || tab === "adhoc";
return (
<Modal open={open} onClose={onClose} title="Add a step" wide>
<div className="space-y-4">
<input
autoFocus
className="w-full rounded-lg border border-border bg-surface-2 px-3 py-2 text-sm text-text-primary placeholder-text-secondary/50 focus:border-signal focus:outline-none focus:ring-1 focus:ring-signal"
placeholder="Search steps by name or description…"
value={search}
onChange={(e) => setSearch(e.target.value)}
/>
<div className="flex gap-1.5">
{(["all", "bash", "powershell", "adhoc"] as Tab[]).map((t) => (
<button
key={t}
onClick={() => setTab(t)}
className={`rounded-full border px-3 py-1 text-xs capitalize ${
tab === t
? "border-signal/50 bg-signal/15 text-signal"
: "border-border bg-surface-2 text-text-secondary hover:text-text-primary"
}`}
>
{t === "all" ? "All" : t === "powershell" ? "PowerShell" : t === "adhoc" ? "Ad-hoc" : "Bash"}
</button>
))}
</div>
{showAdhocCards && (
<div className="grid grid-cols-2 gap-2.5">
<button
onClick={onAddAdhoc}
className="flex min-h-[74px] items-center justify-center gap-2 rounded-[10px] border border-dashed border-border text-sm text-text-secondary hover:border-signal/55 hover:text-signal"
>
+ New ad-hoc step
</button>
<button
onClick={() => fileRef.current?.click()}
className="flex min-h-[74px] items-center justify-center gap-2 rounded-[10px] border border-dashed border-border text-sm text-text-secondary hover:border-signal/55 hover:text-signal"
>
Import ad-hoc from file
</button>
<input
ref={fileRef}
type="file"
accept="application/json"
className="hidden"
onChange={(e) => {
const f = e.target.files?.[0];
if (f) onImportAdhoc(f);
e.target.value = "";
}}
/>
</div>
)}
{showLibrary &&
groups.map(
(g) =>
g.steps.length > 0 && (
<div key={g.label}>
<div className="mb-2.5 flex items-center gap-2 text-[11px] font-bold uppercase tracking-wide text-text-secondary">
{g.label}
<span className="h-px flex-1 bg-border" />
</div>
<div className="grid grid-cols-2 gap-2.5">
{g.steps.map((s) => (
<StepCard key={s.step_id} step={s} onAdd={() => onSelect(s.step_id)} />
))}
</div>
</div>
),
)}
{showLibrary && filtered.length === 0 && (
<p className="text-sm text-text-secondary">No steps match your search.</p>
)}
<p className="text-xs text-text-secondary">
Click a card to append it to the workflow · manage the library on the{" "}
<a href="/steps" className="text-signal hover:underline">
Steps
</a>{" "}
page.
</p>
</div>
</Modal>
);
}
+59 -4
View File
@@ -95,6 +95,7 @@ export interface Settings {
alerts: AlertSettings;
email: EmailSettings;
secrets: SecretsSettings;
workflow_log_retention_days?: number | null;
}
export interface SecretGroupSummary {
@@ -132,6 +133,12 @@ export interface ServerWithKeys extends Server {
keys: (Assignment & { key: Key })[];
}
export interface InputParam {
name: string;
default: string;
description: string;
}
export interface WorkflowStep {
step_id: string;
name: string;
@@ -139,15 +146,20 @@ export interface WorkflowStep {
interpreter: "bash" | "powershell";
script: string;
declared_outputs: string[];
declared_inputs: InputParam[];
secret_refs: string[];
source?: "user" | "default";
slug?: string;
}
export interface WorkflowStepRef {
step_id: string;
step_id?: string;
inline?: WorkflowStep;
order: number;
on_failure: "stop" | "continue" | "retry";
max_retries: number;
overrides?: { script?: string; secret_refs?: string[] };
inputs?: Record<string, string>;
}
export interface Workflow {
@@ -163,8 +175,7 @@ export interface StepRun {
status: string;
attempts: number;
exit_code: number;
stdout: string;
stderr: string;
log_offset: number;
output_env: Record<string, string>;
started_at?: string;
finished_at?: string;
@@ -282,7 +293,11 @@ export const api = {
return request<Settings>("/settings");
},
saveSettings(settings: { alerts: AlertSettings; email: EmailSettings }): Promise<{ saved: boolean }> {
saveSettings(settings: {
alerts: AlertSettings;
email: EmailSettings;
workflow_log_retention_days?: number | null;
}): Promise<{ saved: boolean }> {
return request<{ saved: boolean }>("/settings", {
method: "PUT",
body: JSON.stringify(settings),
@@ -407,6 +422,34 @@ export const api = {
return request<void>(`/steps/${stepId}`, { method: "DELETE" });
},
exportStepUrl(stepId: string): string {
return `/api/steps/${stepId}/export`;
},
importStep(doc: unknown): Promise<WorkflowStep> {
return request<WorkflowStep>("/steps/import", {
method: "POST",
body: JSON.stringify(doc),
});
},
parseStep(doc: unknown): Promise<WorkflowStep> {
return request<WorkflowStep>("/steps/parse", {
method: "POST",
body: JSON.stringify(doc),
});
},
seedDefaults(): Promise<{ created: number; updated: number }> {
return request<{ created: number; updated: number }>("/steps/seed-defaults", {
method: "POST",
});
},
stepUsage(): Promise<Record<string, number>> {
return request<Record<string, number>>("/steps/usage");
},
// Workflows
listWorkflows(): Promise<Workflow[]> {
return request<Workflow[]>("/workflows");
@@ -452,4 +495,16 @@ export const api = {
cancelRun(runId: string): Promise<void> {
return request<void>(`/runs/${runId}/cancel`, { method: "POST" });
},
async getServerRunLog(runId: string, serverId: string): Promise<string> {
const res = await fetch(`/api/runs/${runId}/servers/${serverId}/logs`, {
credentials: "include",
});
if (!res.ok) throw new Error("no logs");
return res.text();
},
serverRunLogStreamUrl(runId: string, serverId: string): string {
return `/api/runs/${runId}/servers/${serverId}/logs/stream`;
},
};
+4
View File
@@ -21,6 +21,10 @@ const config: Config = {
warning: "#f59e0b",
danger: "#ef4444",
"danger-hover": "#dc2626",
bash: "#3fb950",
pwsh: "#5b9bff",
signal: "#f5a524",
"signal-ink": "#241800",
},
},
},
File diff suppressed because one or more lines are too long