Compare commits
5 Commits
v1.1.1
...
8fdf4a1bdf
| Author | SHA1 | Date | |
|---|---|---|---|
| 8fdf4a1bdf | |||
| 750a06384b | |||
| 2110c7dab4 | |||
| 8ddd3456e1 | |||
| 74cb7f660f |
@@ -108,6 +108,7 @@ func connectOnce(ctx context.Context, cfg Config, handle Handler) error {
|
||||
if err != nil {
|
||||
return fmt.Errorf("dial: %w", err)
|
||||
}
|
||||
conn.SetReadLimit(api.MaxWebSocketMessageBytes)
|
||||
// On a successful upgrade coder/websocket transfers ownership of the
|
||||
// response stream to conn and deliberately sets res.Body to nil. Closing
|
||||
// the connection below releases that stream.
|
||||
|
||||
@@ -2,12 +2,16 @@ package wsclient
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/coder/websocket"
|
||||
|
||||
"gitea.dcglab.co.uk/steve/restic-manager/internal/api"
|
||||
)
|
||||
|
||||
func TestConnectOnceCleanDisconnectDoesNotPanic(t *testing.T) {
|
||||
@@ -44,3 +48,56 @@ func TestConnectOnceCleanDisconnectDoesNotPanic(t *testing.T) {
|
||||
t.Fatalf("server websocket: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestConnectOnceAcceptsMessageLargerThanDefaultReadLimit(t *testing.T) {
|
||||
received := make(chan struct{}, 1)
|
||||
serverErr := make(chan error, 1)
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
conn, err := websocket.Accept(w, r, nil)
|
||||
if err != nil {
|
||||
serverErr <- err
|
||||
return
|
||||
}
|
||||
defer conn.CloseNow() //nolint:errcheck
|
||||
if _, _, err := conn.Read(r.Context()); err != nil {
|
||||
serverErr <- err
|
||||
return
|
||||
}
|
||||
env := api.Envelope{
|
||||
Type: api.MsgConfigUpdate,
|
||||
Payload: json.RawMessage(`{"padding":"` + strings.Repeat("x", 40*1024) + `"}`),
|
||||
}
|
||||
raw, _ := json.Marshal(env)
|
||||
serverErr <- conn.Write(r.Context(), websocket.MessageText, raw)
|
||||
<-r.Context().Done()
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
done := make(chan error, 1)
|
||||
go func() {
|
||||
done <- connectOnce(ctx, Config{
|
||||
ServerURL: srv.URL,
|
||||
AgentToken: "test-token",
|
||||
HeartbeatPeriod: time.Hour,
|
||||
}, func(_ context.Context, env api.Envelope, _ Sender) error {
|
||||
if env.Type == api.MsgConfigUpdate {
|
||||
received <- struct{}{}
|
||||
cancel()
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-received:
|
||||
case <-time.After(5 * time.Second):
|
||||
t.Fatal("agent did not receive oversized server message")
|
||||
}
|
||||
if err := <-serverErr; err != nil {
|
||||
t.Fatalf("server websocket: %v", err)
|
||||
}
|
||||
if err := <-done; err == nil {
|
||||
t.Fatal("connectOnce returned nil")
|
||||
}
|
||||
}
|
||||
|
||||
@@ -28,6 +28,11 @@ import (
|
||||
// are evaluated — always-on hosts' stale_schedule stays a no-op.
|
||||
const staleBackupThreshold = 7 * 24 * time.Hour
|
||||
|
||||
const (
|
||||
defaultStuckJobThreshold = 6 * time.Hour
|
||||
longStuckJobThreshold = 24 * time.Hour
|
||||
)
|
||||
|
||||
// JobFinishedEvent carries everything the engine needs to evaluate
|
||||
// the failed-X rules. Pushed via Engine.NotifyJobFinished from the
|
||||
// MarkJobFinished site.
|
||||
@@ -53,6 +58,7 @@ type Engine struct {
|
||||
// we raise. Configurable for tests; default 15m.
|
||||
agentOfflineFloor time.Duration
|
||||
tickPeriod time.Duration
|
||||
stuckThresholds map[string]time.Duration
|
||||
|
||||
closeOnce sync.Once
|
||||
done chan struct{}
|
||||
@@ -69,7 +75,12 @@ func NewEngine(st *store.Store, hub *notification.Hub) *Engine {
|
||||
hostUp: make(chan string, 32),
|
||||
agentOfflineFloor: 15 * time.Minute,
|
||||
tickPeriod: 60 * time.Second,
|
||||
done: make(chan struct{}),
|
||||
stuckThresholds: map[string]time.Duration{
|
||||
"backup": longStuckJobThreshold,
|
||||
"restore": longStuckJobThreshold,
|
||||
"check": longStuckJobThreshold,
|
||||
},
|
||||
done: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -125,6 +136,10 @@ func (e *Engine) NotifyHostOnline(hostID string) {
|
||||
}
|
||||
|
||||
func (e *Engine) handleJobFinished(ctx context.Context, ev JobFinishedEvent) {
|
||||
// A late terminal message is authoritative and clears any stuck alert for
|
||||
// this exact job regardless of its outcome or kind.
|
||||
e.resolveAndNotify(ctx, ev.HostID, KindJobStuck, ev.JobID, ev.When)
|
||||
|
||||
// Determine which kind/severity pair this job maps to. Jobs not
|
||||
// listed here (init, unlock, restore, diff) produce no alerts in v1.
|
||||
var kind, severity string
|
||||
@@ -210,6 +225,7 @@ func (e *Engine) tick(ctx context.Context, now time.Time) {
|
||||
if _, err := e.store.CleanupExpiredOIDCState(ctx, now.Add(-5*time.Minute)); err != nil {
|
||||
slog.Warn("alert: cleanup expired oidc state", "err", err)
|
||||
}
|
||||
e.evaluateStuckJobs(ctx, now)
|
||||
|
||||
hosts, err := e.store.ListHosts(ctx)
|
||||
if err != nil {
|
||||
@@ -257,6 +273,46 @@ func (e *Engine) tick(ctx context.Context, now time.Time) {
|
||||
}
|
||||
}
|
||||
|
||||
func (e *Engine) evaluateStuckJobs(ctx context.Context, now time.Time) {
|
||||
running, err := e.store.ListRunningJobActivity(ctx)
|
||||
if err != nil {
|
||||
slog.Warn("alert: tick list running jobs", "err", err)
|
||||
return
|
||||
}
|
||||
active := make(map[string]store.RunningJobActivity, len(running))
|
||||
for _, job := range running {
|
||||
active[job.JobID] = job
|
||||
threshold := defaultStuckJobThreshold
|
||||
if configured, ok := e.stuckThresholds[job.Kind]; ok {
|
||||
threshold = configured
|
||||
}
|
||||
age := now.Sub(job.LastActivity)
|
||||
if age < threshold {
|
||||
continue
|
||||
}
|
||||
e.raiseAndNotify(ctx, job.HostID, KindJobStuck, job.JobID, "warning",
|
||||
fmt.Sprintf("%s job %s is stuck: started %s, last activity %s (%s ago; threshold %s)",
|
||||
job.Kind, job.JobID, job.StartedAt.Format(time.RFC3339),
|
||||
job.LastActivity.Format(time.RFC3339), roundDur(age), threshold), now)
|
||||
}
|
||||
|
||||
// Self-heal alerts when another server path made a job terminal without
|
||||
// emitting JobFinishedEvent, or after a restart missed the event.
|
||||
alerts, err := e.store.ListAlerts(ctx, store.AlertFilter{Status: "open"})
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
acked, _ := e.store.ListAlerts(ctx, store.AlertFilter{Status: "acknowledged"})
|
||||
for _, item := range append(alerts, acked...) {
|
||||
if item.Kind != KindJobStuck || item.HostID == nil {
|
||||
continue
|
||||
}
|
||||
if _, ok := active[item.DedupKey]; !ok {
|
||||
e.resolveAndNotify(ctx, *item.HostID, KindJobStuck, item.DedupKey, now)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// roundDur returns a human-readable duration string, rounding to the
|
||||
// nearest minute. Durations under a minute are reported as "less than
|
||||
// a minute".
|
||||
|
||||
@@ -36,6 +36,10 @@ const (
|
||||
// KindAgentOffline is raised when a host's last_seen_at is older
|
||||
// than the 15-minute floor and resolved when the host reconnects.
|
||||
KindAgentOffline = "agent_offline"
|
||||
|
||||
// KindJobStuck is raised per job when a running job has no persisted
|
||||
// activity beyond its kind-specific threshold. The job ID is the dedup key.
|
||||
KindJobStuck = "job_stuck"
|
||||
)
|
||||
|
||||
// raiseAndNotify is the standard raise pattern: store.RaiseOrTouch
|
||||
|
||||
@@ -3,6 +3,7 @@ package alert
|
||||
import (
|
||||
"context"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
@@ -69,6 +70,65 @@ func TestEngineBackupFailedRaisesThenResolves(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
func TestEngineStuckJobRaisesDeduplicatesAndResolves(t *testing.T) {
|
||||
t.Parallel()
|
||||
eng, st, hostID := setupEngine(t)
|
||||
ctx := context.Background()
|
||||
now := time.Now().UTC().Truncate(time.Second)
|
||||
started := now.Add(-7 * time.Hour)
|
||||
if err := st.CreateJob(ctx, store.Job{ID: "stuck-job", HostID: hostID, Kind: "prune", ActorKind: "user", CreatedAt: started}); err != nil {
|
||||
t.Fatalf("create job: %v", err)
|
||||
}
|
||||
if err := st.MarkJobStarted(ctx, "stuck-job", started); err != nil {
|
||||
t.Fatalf("start job: %v", err)
|
||||
}
|
||||
|
||||
eng.evaluateStuckJobs(ctx, now)
|
||||
eng.evaluateStuckJobs(ctx, now.Add(time.Minute))
|
||||
open, err := st.ListAlerts(ctx, store.AlertFilter{Status: "open", HostID: hostID})
|
||||
if err != nil {
|
||||
t.Fatalf("list alerts: %v", err)
|
||||
}
|
||||
if len(open) != 1 || open[0].Kind != KindJobStuck || open[0].DedupKey != "stuck-job" {
|
||||
t.Fatalf("expected one deduplicated stuck alert, got %+v", open)
|
||||
}
|
||||
if !strings.Contains(open[0].Message, started.Format(time.RFC3339)) {
|
||||
t.Errorf("alert does not identify start/activity time: %q", open[0].Message)
|
||||
}
|
||||
|
||||
if err := st.MarkJobFinished(ctx, "stuck-job", "succeeded", 0, nil, "", now); err != nil {
|
||||
t.Fatalf("finish job: %v", err)
|
||||
}
|
||||
eng.evaluateStuckJobs(ctx, now.Add(2*time.Minute))
|
||||
open, _ = st.ListAlerts(ctx, store.AlertFilter{Status: "open", HostID: hostID})
|
||||
if len(open) != 0 {
|
||||
t.Fatalf("expected terminal job alert resolved, got %+v", open)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEngineStuckJobUsesLatestLogActivity(t *testing.T) {
|
||||
t.Parallel()
|
||||
eng, st, hostID := setupEngine(t)
|
||||
ctx := context.Background()
|
||||
now := time.Now().UTC().Truncate(time.Second)
|
||||
started := now.Add(-7 * time.Hour)
|
||||
if err := st.CreateJob(ctx, store.Job{ID: "active-job", HostID: hostID, Kind: "forget", ActorKind: "user", CreatedAt: started}); err != nil {
|
||||
t.Fatalf("create job: %v", err)
|
||||
}
|
||||
if err := st.MarkJobStarted(ctx, "active-job", started); err != nil {
|
||||
t.Fatalf("start job: %v", err)
|
||||
}
|
||||
if err := st.AppendJobLog(ctx, "active-job", 1, now.Add(-time.Hour), "stdout", "active"); err != nil {
|
||||
t.Fatalf("append log: %v", err)
|
||||
}
|
||||
|
||||
eng.evaluateStuckJobs(ctx, now)
|
||||
open, _ := st.ListAlerts(ctx, store.AlertFilter{Status: "open", HostID: hostID})
|
||||
if len(open) != 0 {
|
||||
t.Fatalf("recent activity should suppress stuck alert, got %+v", open)
|
||||
}
|
||||
}
|
||||
|
||||
func TestEngineCheckFailedSeverityCritical(t *testing.T) {
|
||||
t.Parallel()
|
||||
eng, st, hostID := setupEngine(t)
|
||||
|
||||
@@ -10,6 +10,13 @@ import (
|
||||
// (not iota ints) makes traffic readable in logs and packet captures.
|
||||
type MessageType string
|
||||
|
||||
// MaxWebSocketMessageBytes is the protocol-wide upper bound for one agent ↔
|
||||
// server envelope. Snapshot projections and restic JSON log events can exceed
|
||||
// coder/websocket's 32 KiB default on ordinary repositories, so both peers set
|
||||
// this limit explicitly. It remains bounded to protect either process from an
|
||||
// untrusted or malfunctioning peer allocating without limit.
|
||||
const MaxWebSocketMessageBytes int64 = 8 << 20
|
||||
|
||||
// Agent → server message types.
|
||||
const (
|
||||
MsgHello MessageType = "hello"
|
||||
|
||||
@@ -8,7 +8,9 @@ import (
|
||||
"net/netip"
|
||||
"runtime"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"gitea.dcglab.co.uk/steve/restic-manager/internal/alert"
|
||||
"gitea.dcglab.co.uk/steve/restic-manager/internal/server/config"
|
||||
"gitea.dcglab.co.uk/steve/restic-manager/internal/server/metrics"
|
||||
"gitea.dcglab.co.uk/steve/restic-manager/internal/store"
|
||||
@@ -173,13 +175,37 @@ func (s *Server) gatherMetricsSnapshot(ctx context.Context) (metrics.Snapshot, e
|
||||
return metrics.Snapshot{}, err
|
||||
}
|
||||
bySeverity := map[string]int{"info": 0, "warning": 0, "critical": 0}
|
||||
stuckJobs := 0
|
||||
var oldestStuckAge time.Duration
|
||||
now := time.Now().UTC()
|
||||
running, err := s.deps.Store.ListRunningJobActivity(ctx)
|
||||
if err != nil {
|
||||
return metrics.Snapshot{}, err
|
||||
}
|
||||
activityByJob := make(map[string]time.Time, len(running))
|
||||
for _, job := range running {
|
||||
activityByJob[job.JobID] = job.LastActivity
|
||||
}
|
||||
for _, a := range open {
|
||||
bySeverity[a.Severity]++
|
||||
if a.Kind == alert.KindJobStuck {
|
||||
stuckJobs++
|
||||
lastActivity, ok := activityByJob[a.DedupKey]
|
||||
if !ok {
|
||||
continue
|
||||
}
|
||||
if age := now.Sub(lastActivity); age > oldestStuckAge {
|
||||
oldestStuckAge = age
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
reg := s.deps.Metrics
|
||||
if reg == nil {
|
||||
reg = metrics.NewRegistry() // empty histogram block
|
||||
}
|
||||
return reg.SnapshotWith(hostRows, bySeverity, version.Version, version.Commit, runtime.Version()), nil
|
||||
snap := reg.SnapshotWith(hostRows, bySeverity, version.Version, version.Commit, runtime.Version())
|
||||
snap.StuckJobs = stuckJobs
|
||||
snap.OldestStuckAge = oldestStuckAge
|
||||
return snap, nil
|
||||
}
|
||||
|
||||
@@ -155,6 +155,8 @@ type Snapshot struct {
|
||||
BuildCommit string
|
||||
GoVersion string
|
||||
JobDurationRows []HistogramRow
|
||||
StuckJobs int
|
||||
OldestStuckAge time.Duration
|
||||
}
|
||||
|
||||
// SnapshotWith builds a Snapshot from raw inputs and the registry's
|
||||
@@ -205,6 +207,13 @@ func Render(w io.Writer, s Snapshot) error {
|
||||
fmt.Fprintf(&b, "rm_build_info{version=%q,commit=%q,go_version=%q} 1\n",
|
||||
s.BuildVersion, s.BuildCommit, s.GoVersion)
|
||||
|
||||
b.WriteString("# HELP rm_stuck_jobs Number of open job_stuck alerts.\n")
|
||||
b.WriteString("# TYPE rm_stuck_jobs gauge\n")
|
||||
fmt.Fprintf(&b, "rm_stuck_jobs %d\n", s.StuckJobs)
|
||||
b.WriteString("# HELP rm_oldest_stuck_job_age_seconds Time since last activity for the oldest open stuck job.\n")
|
||||
b.WriteString("# TYPE rm_oldest_stuck_job_age_seconds gauge\n")
|
||||
fmt.Fprintf(&b, "rm_oldest_stuck_job_age_seconds %.0f\n", s.OldestStuckAge.Seconds())
|
||||
|
||||
// --- Per-host gauges -------------------------------------------------
|
||||
// Stable order: by host id.
|
||||
hosts := append([]HostRow(nil), s.Hosts...)
|
||||
|
||||
@@ -114,6 +114,8 @@ func TestRenderGolden(t *testing.T) {
|
||||
snap := r.SnapshotWith(hosts,
|
||||
map[string]int{"info": 0, "warning": 1, "critical": 0},
|
||||
"v1.2.3", "deadbeef", "go1.25.0")
|
||||
snap.StuckJobs = 2
|
||||
snap.OldestStuckAge = 90 * time.Minute
|
||||
|
||||
var buf bytes.Buffer
|
||||
if err := Render(&buf, snap); err != nil {
|
||||
@@ -129,6 +131,8 @@ func TestRenderGolden(t *testing.T) {
|
||||
`rm_active_alerts{severity="info"} 0`,
|
||||
`rm_active_alerts{severity="critical"} 0`,
|
||||
`rm_build_info{version="v1.2.3",commit="deadbeef",go_version="go1.25.0"} 1`,
|
||||
`rm_stuck_jobs 2`,
|
||||
`rm_oldest_stuck_job_age_seconds 5400`,
|
||||
`rm_host_agent_online{host_id="01H0001",host="alpha"} 1`,
|
||||
`rm_host_agent_online{host_id="01H0002",host="bravo"} 0`,
|
||||
`rm_host_last_backup_timestamp_seconds{host_id="01H0001",host="alpha"} 1700000000`,
|
||||
|
||||
@@ -77,6 +77,7 @@ func AgentHandler(deps HandlerDeps) stdhttp.Handler {
|
||||
slog.Warn("ws accept failed", "err", err, "host_id", host.ID)
|
||||
return
|
||||
}
|
||||
conn.SetReadLimit(api.MaxWebSocketMessageBytes)
|
||||
|
||||
c := NewConn(host.ID, conn)
|
||||
// Keep agents alive across NAT boxes; coder/websocket
|
||||
|
||||
@@ -123,6 +123,71 @@ func TestWSHelloAndHeartbeat(t *testing.T) {
|
||||
t.Error("heartbeat did not update last_seen_at")
|
||||
}
|
||||
|
||||
func TestWSAcceptsSnapshotReportLargerThanDefaultReadLimit(t *testing.T) {
|
||||
t.Parallel()
|
||||
url, token, hostID, st, hub := setupTestHub(t)
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
|
||||
defer cancel()
|
||||
c, resp, err := websocket.Dial(ctx, url, &websocket.DialOptions{
|
||||
HTTPHeader: stdhttp.Header{"Authorization": []string{"Bearer " + token}},
|
||||
})
|
||||
if resp != nil && resp.Body != nil {
|
||||
defer resp.Body.Close() //nolint:errcheck
|
||||
}
|
||||
if err != nil {
|
||||
t.Fatalf("dial: %v", err)
|
||||
}
|
||||
defer c.CloseNow() //nolint:errcheck
|
||||
|
||||
hello, _ := api.Marshal(api.MsgHello, "", api.HelloPayload{
|
||||
ProtocolVersion: api.CurrentProtocolVersion,
|
||||
AgentVersion: "0.1.0",
|
||||
ResticVersion: "0.17.1",
|
||||
Hostname: "h1",
|
||||
OS: api.OSLinux,
|
||||
Arch: api.ArchAmd64,
|
||||
})
|
||||
helloRaw, _ := json.Marshal(hello)
|
||||
if err := c.Write(ctx, websocket.MessageText, helloRaw); err != nil {
|
||||
t.Fatalf("write hello: %v", err)
|
||||
}
|
||||
deadline := time.Now().Add(time.Second)
|
||||
for !hub.Connected(hostID) && time.Now().Before(deadline) {
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
|
||||
report, _ := api.Marshal(api.MsgSnapshotsRpt, "", api.SnapshotsReportPayload{
|
||||
Snapshots: []api.Snapshot{{
|
||||
ID: strings.Repeat("a", 64),
|
||||
ShortID: "aaaaaaaa",
|
||||
Time: time.Now().UTC(),
|
||||
Hostname: "h1",
|
||||
Paths: []string{"/" + strings.Repeat("long-path/", 5000)},
|
||||
}},
|
||||
})
|
||||
reportRaw, _ := json.Marshal(report)
|
||||
if len(reportRaw) <= 32*1024 {
|
||||
t.Fatalf("test payload is only %d bytes; must exceed old limit", len(reportRaw))
|
||||
}
|
||||
if int64(len(reportRaw)) >= api.MaxWebSocketMessageBytes {
|
||||
t.Fatalf("test payload %d exceeds protocol limit", len(reportRaw))
|
||||
}
|
||||
if err := c.Write(ctx, websocket.MessageText, reportRaw); err != nil {
|
||||
t.Fatalf("write snapshots.report: %v", err)
|
||||
}
|
||||
|
||||
deadline = time.Now().Add(2 * time.Second)
|
||||
for time.Now().Before(deadline) {
|
||||
host, err := st.GetHost(context.Background(), hostID)
|
||||
if err == nil && host.SnapshotCount == 1 {
|
||||
return
|
||||
}
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
}
|
||||
t.Fatal("oversized snapshots.report was not projected")
|
||||
}
|
||||
|
||||
func TestWSRejectsOldProtocol(t *testing.T) {
|
||||
t.Parallel()
|
||||
url, token, _, _, _ := setupTestHub(t)
|
||||
|
||||
@@ -27,6 +27,52 @@ type Job struct {
|
||||
CreatedAt time.Time
|
||||
}
|
||||
|
||||
// RunningJobActivity is the persisted activity summary used to detect jobs
|
||||
// whose terminal message was lost. LastActivity is the newer of started_at and
|
||||
// the most recent persisted log line; ephemeral progress events deliberately
|
||||
// do not extend it.
|
||||
type RunningJobActivity struct {
|
||||
JobID string
|
||||
HostID string
|
||||
Kind string
|
||||
StartedAt time.Time
|
||||
LastActivity time.Time
|
||||
}
|
||||
|
||||
// ListRunningJobActivity returns every running job with its latest persisted
|
||||
// activity. Jobs without started_at are excluded because they have not actually
|
||||
// entered the running state coherently and cannot be aged safely here.
|
||||
func (s *Store) ListRunningJobActivity(ctx context.Context) ([]RunningJobActivity, error) {
|
||||
rows, err := s.db.QueryContext(ctx, `
|
||||
SELECT j.id, j.host_id, j.kind, j.started_at,
|
||||
COALESCE(MAX(l.ts), j.started_at) AS last_activity
|
||||
FROM jobs j
|
||||
LEFT JOIN job_logs l ON l.job_id = j.id
|
||||
WHERE j.status = 'running' AND j.started_at IS NOT NULL
|
||||
GROUP BY j.id, j.host_id, j.kind, j.started_at
|
||||
ORDER BY j.started_at`)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("store: list running job activity: %w", err)
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
|
||||
var out []RunningJobActivity
|
||||
for rows.Next() {
|
||||
var item RunningJobActivity
|
||||
var started, activity string
|
||||
if err := rows.Scan(&item.JobID, &item.HostID, &item.Kind, &started, &activity); err != nil {
|
||||
return nil, fmt.Errorf("store: scan running job activity: %w", err)
|
||||
}
|
||||
item.StartedAt, _ = time.Parse(time.RFC3339Nano, started)
|
||||
item.LastActivity, _ = time.Parse(time.RFC3339Nano, activity)
|
||||
out = append(out, item)
|
||||
}
|
||||
if err := rows.Err(); err != nil {
|
||||
return nil, fmt.Errorf("store: iterate running job activity: %w", err)
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
// CreateJob inserts a queued job. The agent will mark it running
|
||||
// when it actually starts work. ScheduledID is set when the job
|
||||
// originates from a cron fire (actor_kind="schedule"); nil for
|
||||
|
||||
@@ -7,6 +7,48 @@ import (
|
||||
"time"
|
||||
)
|
||||
|
||||
func TestListRunningJobActivity(t *testing.T) {
|
||||
t.Parallel()
|
||||
s := openTestStore(t)
|
||||
ctx := context.Background()
|
||||
hostID := makeSchedHost(t, s)
|
||||
started := time.Now().UTC().Add(-8 * time.Hour).Truncate(time.Second)
|
||||
|
||||
for _, id := range []string{"running-idle", "running-active", "finished"} {
|
||||
if err := s.CreateJob(ctx, Job{ID: id, HostID: hostID, Kind: "backup", ActorKind: "user", CreatedAt: started}); err != nil {
|
||||
t.Fatalf("create %s: %v", id, err)
|
||||
}
|
||||
if err := s.MarkJobStarted(ctx, id, started); err != nil {
|
||||
t.Fatalf("start %s: %v", id, err)
|
||||
}
|
||||
}
|
||||
activity := started.Add(7 * time.Hour)
|
||||
if err := s.AppendJobLog(ctx, "running-active", 1, activity, "stdout", "still working"); err != nil {
|
||||
t.Fatalf("append log: %v", err)
|
||||
}
|
||||
if err := s.MarkJobFinished(ctx, "finished", "succeeded", 0, nil, "", activity); err != nil {
|
||||
t.Fatalf("finish: %v", err)
|
||||
}
|
||||
|
||||
got, err := s.ListRunningJobActivity(ctx)
|
||||
if err != nil {
|
||||
t.Fatalf("list activity: %v", err)
|
||||
}
|
||||
if len(got) != 2 {
|
||||
t.Fatalf("got %d rows, want 2: %+v", len(got), got)
|
||||
}
|
||||
byID := make(map[string]RunningJobActivity, len(got))
|
||||
for _, row := range got {
|
||||
byID[row.JobID] = row
|
||||
}
|
||||
if !byID["running-idle"].LastActivity.Equal(started) {
|
||||
t.Errorf("idle last activity = %s, want %s", byID["running-idle"].LastActivity, started)
|
||||
}
|
||||
if !byID["running-active"].LastActivity.Equal(activity) {
|
||||
t.Errorf("active last activity = %s, want %s", byID["running-active"].LastActivity, activity)
|
||||
}
|
||||
}
|
||||
|
||||
func TestLatestJobByKind(t *testing.T) {
|
||||
t.Parallel()
|
||||
s := openTestStore(t)
|
||||
|
||||
Reference in New Issue
Block a user