ship(P01): iter.Seq streaming merged into v0.3 milestone
---ci--- project: orca phase: 1 milestone: v0.3 status: complete requirements: covered: [REQ-022, REQ-030] partial: [] ---/ci--- P01: iter.Seq streaming for --watch flags. - JobRepo.Watch / NodeRepo.Watch: pull-based iter.Seq[[]*T] snapshot-per-tick (G-001) - Immediate first yield before ticker (G-002) - Table mode: clear-screen + re-render on change - JSON mode: init/update/delete events, one line per change - signal.NotifyContext on SIGINT/SIGTERM (D-023) - 16 tests (8 store + 8 CLI), all pass under -race - 4-layer verification passed
This commit is contained in:
+76
-1
@@ -5,6 +5,9 @@ import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
"os/signal"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
@@ -36,6 +39,7 @@ var (
|
||||
stopID string
|
||||
runTarget string
|
||||
runIDKey string
|
||||
jobWatch bool
|
||||
)
|
||||
|
||||
var jobRunCmd = &cobra.Command{
|
||||
@@ -112,8 +116,11 @@ var jobRunCmd = &cobra.Command{
|
||||
var jobListCmd = &cobra.Command{
|
||||
Use: "list",
|
||||
Short: "List all jobs",
|
||||
Long: "Display all jobs and their status.",
|
||||
Long: "Display all jobs and their status. Use --watch to stream updates until Ctrl-C.",
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
if jobWatch {
|
||||
return watchJobs(cmd)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
@@ -142,6 +149,73 @@ var jobListCmd = &cobra.Command{
|
||||
},
|
||||
}
|
||||
|
||||
func watchJobs(cmd *cobra.Command) error {
|
||||
ctx, cancel := signal.NotifyContext(cmd.Context(), os.Interrupt, syscall.SIGTERM)
|
||||
defer cancel()
|
||||
return watchJobsCtx(cmd, ctx)
|
||||
}
|
||||
|
||||
func watchJobsCtx(cmd *cobra.Command, ctx context.Context) error {
|
||||
db, closer, err := openDB()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer closer()
|
||||
|
||||
out := cmd.OutOrStdout()
|
||||
|
||||
if jsonOutput {
|
||||
seen := make(map[string]string)
|
||||
for snapshot := range store.NewJobRepo(db).Watch(ctx) {
|
||||
current := make(map[string]bool, len(snapshot))
|
||||
for _, j := range snapshot {
|
||||
current[j.ID] = true
|
||||
compact, _ := json.Marshal(j)
|
||||
key := string(compact)
|
||||
if prev, ok := seen[j.ID]; !ok || prev != key {
|
||||
event := "init"
|
||||
if ok {
|
||||
event = "update"
|
||||
}
|
||||
line, _ := json.Marshal(map[string]any{"event": event, "job": j})
|
||||
fmt.Fprintln(out, string(line))
|
||||
seen[j.ID] = key
|
||||
}
|
||||
}
|
||||
for id := range seen {
|
||||
if !current[id] {
|
||||
line, _ := json.Marshal(map[string]any{"event": "delete", "id": id})
|
||||
fmt.Fprintln(out, string(line))
|
||||
delete(seen, id)
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
prevTable := ""
|
||||
for snapshot := range store.NewJobRepo(db).Watch(ctx) {
|
||||
table := renderJobTable(snapshot)
|
||||
if table != prevTable {
|
||||
fmt.Fprint(out, "\033[2J\033[H")
|
||||
fmt.Fprint(out, table)
|
||||
prevTable = table
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func renderJobTable(jobs []*model.Job) string {
|
||||
if len(jobs) == 0 {
|
||||
return "No jobs.\n"
|
||||
}
|
||||
out := fmt.Sprintf("%-36s %-20s %-12s %-8s\n", "ID", "NAME", "STATUS", "EXIT")
|
||||
for _, j := range jobs {
|
||||
out += fmt.Sprintf("%-36s %-20s %-12s %-8d\n", j.ID, j.Name, j.Status, j.ExitCode)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
var jobStopCmd = &cobra.Command{
|
||||
Use: "stop [job-id]",
|
||||
Short: "Stop a running job",
|
||||
@@ -235,6 +309,7 @@ func init() {
|
||||
jobLogsCmd.Flags().StringVar(&stopID, "id", "", "job id")
|
||||
jobRunCmd.Flags().StringVar(&runTarget, "target", "", "pin job to a specific node id (overrides bin-packing)")
|
||||
jobRunCmd.Flags().StringVar(&runIDKey, "idempotency-key", "", "X-Orca-Idempotency-Key for cross-node dispatch dedupe")
|
||||
jobListCmd.Flags().BoolVar(&jobWatch, "watch", false, "stream jobs until Ctrl-C (table refresh or --json per-event)")
|
||||
|
||||
jobCmd.AddCommand(jobRunCmd)
|
||||
jobCmd.AddCommand(jobListCmd)
|
||||
|
||||
+76
-1
@@ -3,10 +3,13 @@ package cli
|
||||
import (
|
||||
"context"
|
||||
"database/sql"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"os"
|
||||
"os/signal"
|
||||
"path/filepath"
|
||||
"syscall"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
@@ -54,6 +57,7 @@ var (
|
||||
joinAddr string
|
||||
joinCAFinger string
|
||||
leaveID string
|
||||
nodeWatch bool
|
||||
)
|
||||
|
||||
var nodeCmd = &cobra.Command{
|
||||
@@ -155,8 +159,11 @@ var nodeLeaveCmd = &cobra.Command{
|
||||
var nodeListCmd = &cobra.Command{
|
||||
Use: "list",
|
||||
Short: "List all nodes in the orca registry",
|
||||
Long: "Display all registered nodes and their state.",
|
||||
Long: "Display all registered nodes and their state. Use --watch to stream updates until Ctrl-C.",
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
if nodeWatch {
|
||||
return watchNodes(cmd)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(cmd.Context(), 5*time.Second)
|
||||
defer cancel()
|
||||
|
||||
@@ -185,11 +192,79 @@ var nodeListCmd = &cobra.Command{
|
||||
},
|
||||
}
|
||||
|
||||
func watchNodes(cmd *cobra.Command) error {
|
||||
ctx, cancel := signal.NotifyContext(cmd.Context(), os.Interrupt, syscall.SIGTERM)
|
||||
defer cancel()
|
||||
return watchNodesCtx(cmd, ctx)
|
||||
}
|
||||
|
||||
func watchNodesCtx(cmd *cobra.Command, ctx context.Context) error {
|
||||
db, closer, err := openDB()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer closer()
|
||||
|
||||
out := cmd.OutOrStdout()
|
||||
|
||||
if jsonOutput {
|
||||
seen := make(map[string]string)
|
||||
for snapshot := range store.NewNodeRepo(db).Watch(ctx) {
|
||||
current := make(map[string]bool, len(snapshot))
|
||||
for _, n := range snapshot {
|
||||
current[n.ID] = true
|
||||
compact, _ := json.Marshal(n)
|
||||
key := string(compact)
|
||||
if prev, ok := seen[n.ID]; !ok || prev != key {
|
||||
event := "init"
|
||||
if ok {
|
||||
event = "update"
|
||||
}
|
||||
line, _ := json.Marshal(map[string]any{"event": event, "node": n})
|
||||
fmt.Fprintln(out, string(line))
|
||||
seen[n.ID] = key
|
||||
}
|
||||
}
|
||||
for id := range seen {
|
||||
if !current[id] {
|
||||
line, _ := json.Marshal(map[string]any{"event": "delete", "id": id})
|
||||
fmt.Fprintln(out, string(line))
|
||||
delete(seen, id)
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
prevTable := ""
|
||||
for snapshot := range store.NewNodeRepo(db).Watch(ctx) {
|
||||
table := renderNodeTable(snapshot)
|
||||
if table != prevTable {
|
||||
fmt.Fprint(out, "\033[2J\033[H")
|
||||
fmt.Fprint(out, table)
|
||||
prevTable = table
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func renderNodeTable(nodes []*model.Node) string {
|
||||
if len(nodes) == 0 {
|
||||
return "No nodes registered.\n"
|
||||
}
|
||||
out := fmt.Sprintf("%-36s %-20s %-22s %-10s\n", "ID", "NAME", "ADDRESS", "STATE")
|
||||
for _, n := range nodes {
|
||||
out += fmt.Sprintf("%-36s %-20s %-22s %-10s\n", n.ID, n.Name, n.Address, n.State)
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
func init() {
|
||||
nodeJoinCmd.Flags().StringVar(&joinName, "name", "", "node name (required)")
|
||||
nodeJoinCmd.Flags().StringVar(&joinAddr, "addr", "", "node address (default localhost:8443)")
|
||||
nodeJoinCmd.Flags().StringVar(&joinCAFinger, "ca-fingerprint", "", "pin CA cert SHA-256 (REQ-026); fails if on-disk CA doesn't match")
|
||||
nodeLeaveCmd.Flags().StringVar(&leaveID, "id", "", "node id")
|
||||
nodeListCmd.Flags().BoolVar(&nodeWatch, "watch", false, "stream nodes until Ctrl-C (table refresh or --json per-event)")
|
||||
|
||||
nodeCmd.AddCommand(nodeJoinCmd)
|
||||
nodeCmd.AddCommand(nodeLeaveCmd)
|
||||
|
||||
@@ -0,0 +1,287 @@
|
||||
package cli
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"git.cloudinit.dev/coreci/orca/internal/model"
|
||||
"git.cloudinit.dev/coreci/orca/internal/store"
|
||||
)
|
||||
|
||||
func TestWatchJobs_JSONStreaming(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
dbPath := filepath.Join(dir, "orca.db")
|
||||
db, err := store.Open(dbPath)
|
||||
if err != nil {
|
||||
t.Fatalf("open db: %v", err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
repo := store.NewJobRepo(db)
|
||||
bgCtx := context.Background()
|
||||
_ = repo.Insert(bgCtx, &model.Job{ID: "seed-job", Name: "seed", Spec: "t", Status: model.JobStatusPending})
|
||||
|
||||
t.Setenv("ORCA_DB", dbPath)
|
||||
|
||||
jsonOutput = true
|
||||
t.Cleanup(func() { jsonOutput = false })
|
||||
|
||||
var buf bytes.Buffer
|
||||
rootCmd.SetOut(&buf)
|
||||
rootCmd.SetErr(&buf)
|
||||
t.Cleanup(func() { rootCmd.SetOut(os.Stdout); rootCmd.SetErr(os.Stderr) })
|
||||
|
||||
ctx, cancel := context.WithCancel(bgCtx)
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- watchJobsCtx(rootCmd, ctx) }()
|
||||
|
||||
// First yield is immediate (G-002); wait for it.
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
|
||||
_ = repo.Insert(bgCtx, &model.Job{ID: "watch-job", Name: "watch", Spec: "t", Status: model.JobStatusPending})
|
||||
|
||||
// Wait for at least one ticker interval (default 1s) to capture the change.
|
||||
time.Sleep(1100 * time.Millisecond)
|
||||
cancel()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("watchJobsCtx did not return within 2s after cancel")
|
||||
}
|
||||
|
||||
output := buf.String()
|
||||
if !strings.Contains(output, `"event":"init"`) {
|
||||
t.Errorf("expected init event, got: %s", output)
|
||||
}
|
||||
if !strings.Contains(output, "watch-job") {
|
||||
t.Errorf("expected watch-job in output, got: %s", output)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWatchJobs_TableRefresh(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
dbPath := filepath.Join(dir, "orca.db")
|
||||
db, err := store.Open(dbPath)
|
||||
if err != nil {
|
||||
t.Fatalf("open db: %v", err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
repo := store.NewJobRepo(db)
|
||||
bgCtx := context.Background()
|
||||
_ = repo.Insert(bgCtx, &model.Job{ID: "seed-job", Name: "seed", Spec: "t", Status: model.JobStatusPending})
|
||||
|
||||
t.Setenv("ORCA_DB", dbPath)
|
||||
|
||||
jsonOutput = false
|
||||
t.Cleanup(func() { jsonOutput = false })
|
||||
|
||||
var buf bytes.Buffer
|
||||
rootCmd.SetOut(&buf)
|
||||
rootCmd.SetErr(&buf)
|
||||
t.Cleanup(func() { rootCmd.SetOut(os.Stdout); rootCmd.SetErr(os.Stderr) })
|
||||
|
||||
ctx, cancel := context.WithCancel(bgCtx)
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- watchJobsCtx(rootCmd, ctx) }()
|
||||
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
|
||||
_ = repo.Insert(bgCtx, &model.Job{ID: "table-job", Name: "table", Spec: "t", Status: model.JobStatusPending})
|
||||
|
||||
time.Sleep(1100 * time.Millisecond)
|
||||
cancel()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("watchJobsCtx did not return within 2s after cancel")
|
||||
}
|
||||
|
||||
output := buf.String()
|
||||
if !strings.Contains(output, "\033[2J\033[H") {
|
||||
t.Errorf("expected clear-screen escape in table watch output, got: %s", output)
|
||||
}
|
||||
if !strings.Contains(output, "table-job") {
|
||||
t.Errorf("expected table-job in output, got: %s", output)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWatchNodes_JSONStreaming(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
dbPath := filepath.Join(dir, "orca.db")
|
||||
db, err := store.Open(dbPath)
|
||||
if err != nil {
|
||||
t.Fatalf("open db: %v", err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
repo := store.NewNodeRepo(db)
|
||||
bgCtx := context.Background()
|
||||
_ = repo.Insert(bgCtx, &model.Node{
|
||||
ID: "seed-node", Name: "seed", Address: "addr",
|
||||
State: model.NodeStateReady, JoinedAt: time.Now().UTC(), LastSeen: time.Now().UTC(),
|
||||
})
|
||||
|
||||
t.Setenv("ORCA_DB", dbPath)
|
||||
|
||||
jsonOutput = true
|
||||
t.Cleanup(func() { jsonOutput = false })
|
||||
|
||||
var buf bytes.Buffer
|
||||
rootCmd.SetOut(&buf)
|
||||
rootCmd.SetErr(&buf)
|
||||
t.Cleanup(func() { rootCmd.SetOut(os.Stdout); rootCmd.SetErr(os.Stderr) })
|
||||
|
||||
ctx, cancel := context.WithCancel(bgCtx)
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- watchNodesCtx(rootCmd, ctx) }()
|
||||
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
|
||||
_ = repo.Insert(bgCtx, &model.Node{
|
||||
ID: "watch-node", Name: "watch", Address: "addr2",
|
||||
State: model.NodeStateReady, JoinedAt: time.Now().UTC(), LastSeen: time.Now().UTC(),
|
||||
})
|
||||
|
||||
time.Sleep(1100 * time.Millisecond)
|
||||
cancel()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("watchNodesCtx did not return within 2s after cancel")
|
||||
}
|
||||
|
||||
output := buf.String()
|
||||
initFound := false
|
||||
watchNodeFound := false
|
||||
for _, line := range strings.Split(output, "\n") {
|
||||
line = strings.TrimSpace(line)
|
||||
if line == "" {
|
||||
continue
|
||||
}
|
||||
var event map[string]any
|
||||
if err := json.Unmarshal([]byte(line), &event); err != nil {
|
||||
continue
|
||||
}
|
||||
if event["event"] == "init" {
|
||||
initFound = true
|
||||
if node, ok := event["node"].(map[string]any); ok {
|
||||
if node["id"] == "watch-node" {
|
||||
watchNodeFound = true
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
if !initFound {
|
||||
t.Errorf("expected init event in JSON stream, got: %s", output)
|
||||
}
|
||||
if !watchNodeFound {
|
||||
t.Errorf("expected watch-node in JSON stream, got: %s", output)
|
||||
}
|
||||
}
|
||||
|
||||
func TestWatchNodes_TableRefresh(t *testing.T) {
|
||||
dir := t.TempDir()
|
||||
dbPath := filepath.Join(dir, "orca.db")
|
||||
db, err := store.Open(dbPath)
|
||||
if err != nil {
|
||||
t.Fatalf("open db: %v", err)
|
||||
}
|
||||
defer db.Close()
|
||||
|
||||
repo := store.NewNodeRepo(db)
|
||||
bgCtx := context.Background()
|
||||
_ = repo.Insert(bgCtx, &model.Node{
|
||||
ID: "seed-node", Name: "seed", Address: "addr",
|
||||
State: model.NodeStateReady, JoinedAt: time.Now().UTC(), LastSeen: time.Now().UTC(),
|
||||
})
|
||||
|
||||
t.Setenv("ORCA_DB", dbPath)
|
||||
|
||||
jsonOutput = false
|
||||
t.Cleanup(func() { jsonOutput = false })
|
||||
|
||||
var buf bytes.Buffer
|
||||
rootCmd.SetOut(&buf)
|
||||
rootCmd.SetErr(&buf)
|
||||
t.Cleanup(func() { rootCmd.SetOut(os.Stdout); rootCmd.SetErr(os.Stderr) })
|
||||
|
||||
ctx, cancel := context.WithCancel(bgCtx)
|
||||
|
||||
done := make(chan error, 1)
|
||||
go func() { done <- watchNodesCtx(rootCmd, ctx) }()
|
||||
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
|
||||
_ = repo.Insert(bgCtx, &model.Node{
|
||||
ID: "table-node", Name: "table", Address: "addr2",
|
||||
State: model.NodeStateReady, JoinedAt: time.Now().UTC(), LastSeen: time.Now().UTC(),
|
||||
})
|
||||
|
||||
time.Sleep(1100 * time.Millisecond)
|
||||
cancel()
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
case <-time.After(2 * time.Second):
|
||||
t.Fatal("watchNodesCtx did not return within 2s after cancel")
|
||||
}
|
||||
|
||||
output := buf.String()
|
||||
if !strings.Contains(output, "\033[2J\033[H") {
|
||||
t.Errorf("expected clear-screen escape in table watch output, got: %s", output)
|
||||
}
|
||||
if !strings.Contains(output, "table-node") {
|
||||
t.Errorf("expected table-node in output, got: %s", output)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderJobTable(t *testing.T) {
|
||||
jobs := []*model.Job{
|
||||
{ID: "j1", Name: "alpha", Status: "running", ExitCode: 0},
|
||||
{ID: "j2", Name: "beta", Status: "done", ExitCode: 0},
|
||||
}
|
||||
out := renderJobTable(jobs)
|
||||
if !strings.Contains(out, "j1") || !strings.Contains(out, "alpha") {
|
||||
t.Errorf("renderJobTable missing job 1: %s", out)
|
||||
}
|
||||
if !strings.Contains(out, "j2") || !strings.Contains(out, "beta") {
|
||||
t.Errorf("renderJobTable missing job 2: %s", out)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderJobTableEmpty(t *testing.T) {
|
||||
out := renderJobTable(nil)
|
||||
if !strings.Contains(out, "No jobs") {
|
||||
t.Errorf("expected empty message, got: %s", out)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderNodeTable(t *testing.T) {
|
||||
nodes := []*model.Node{
|
||||
{ID: "n1", Name: "alpha", Address: "localhost:8443", State: "ready"},
|
||||
}
|
||||
out := renderNodeTable(nodes)
|
||||
if !strings.Contains(out, "n1") || !strings.Contains(out, "alpha") {
|
||||
t.Errorf("renderNodeTable missing node: %s", out)
|
||||
}
|
||||
}
|
||||
|
||||
func TestRenderNodeTableEmpty(t *testing.T) {
|
||||
out := renderNodeTable(nil)
|
||||
if !strings.Contains(out, "No nodes") {
|
||||
t.Errorf("expected empty message, got: %s", out)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user