Compare commits
7 Commits
v1.1.1
...
b64f029892
| Author | SHA1 | Date | |
|---|---|---|---|
| b64f029892 | |||
| f9718e6077 | |||
| 8fdf4a1bdf | |||
| 750a06384b | |||
| 2110c7dab4 | |||
| 8ddd3456e1 | |||
| 74cb7f660f |
@@ -296,6 +296,9 @@ func (d *dispatcher) handle(ctx context.Context, env api.Envelope, tx wsclient.S
|
|||||||
}
|
}
|
||||||
go d.handleTreeList(ctx, env.ID, p, tx)
|
go d.handleTreeList(ctx, env.ID, p, tx)
|
||||||
|
|
||||||
|
case api.MsgSnapshotsRefresh:
|
||||||
|
go d.refreshSnapshots(ctx, tx)
|
||||||
|
|
||||||
case api.MsgScheduleSet:
|
case api.MsgScheduleSet:
|
||||||
var p api.ScheduleSetPayload
|
var p api.ScheduleSetPayload
|
||||||
if err := env.UnmarshalPayload(&p); err != nil {
|
if err := env.UnmarshalPayload(&p); err != nil {
|
||||||
@@ -405,6 +408,22 @@ func (d *dispatcher) handle(ctx context.Context, env api.Envelope, tx wsclient.S
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (d *dispatcher) refreshSnapshots(ctx context.Context, tx wsclient.Sender) {
|
||||||
|
creds, err := d.secrets.Load()
|
||||||
|
if err != nil || creds.Empty() {
|
||||||
|
slog.Warn("ws agent: snapshots.refresh unavailable", "err", err)
|
||||||
|
return
|
||||||
|
}
|
||||||
|
r := runner.New(runner.Config{
|
||||||
|
ResticBin: d.resticBin, ResticVersion: d.resticVer,
|
||||||
|
RepoURL: creds.URL, RepoUsername: creds.Username, RepoPassword: creds.Password,
|
||||||
|
SupportsRestoreNoOwnership: d.resticSupportsNoOwnership,
|
||||||
|
}, tx, time.Second)
|
||||||
|
if err := r.RefreshSnapshots(ctx); err != nil {
|
||||||
|
slog.Warn("ws agent: snapshots.refresh failed", "err", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// handleTreeList runs `restic ls --json <snapshot> <path>` and ships
|
// handleTreeList runs `restic ls --json <snapshot> <path>` and ships
|
||||||
// the matching tree.list.result envelope back, correlated by the
|
// the matching tree.list.result envelope back, correlated by the
|
||||||
// request envelope's ID. Errors (missing creds, restic failure)
|
// request envelope's ID. Errors (missing creds, restic failure)
|
||||||
|
|||||||
@@ -77,6 +77,12 @@ func (r *Runner) resticEnv() restic.Env {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// RefreshSnapshots reconciles the server's cached projection without running
|
||||||
|
// a mutating repository job.
|
||||||
|
func (r *Runner) RefreshSnapshots(ctx context.Context) error {
|
||||||
|
return r.reportSnapshots(ctx, r.resticEnv())
|
||||||
|
}
|
||||||
|
|
||||||
// sendStarted ships a job.started envelope.
|
// sendStarted ships a job.started envelope.
|
||||||
func (r *Runner) sendStarted(jobID string, kind api.JobKind, startedAt time.Time) {
|
func (r *Runner) sendStarted(jobID string, kind api.JobKind, startedAt time.Time) {
|
||||||
env, _ := api.Marshal(api.MsgJobStarted, jobID, api.JobStartedPayload{
|
env, _ := api.Marshal(api.MsgJobStarted, jobID, api.JobStartedPayload{
|
||||||
@@ -226,14 +232,9 @@ func (r *Runner) RunBackup(ctx context.Context, jobID string, paths, excludes, t
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
r.sendFinished(ctx, jobID, finishedAt, err, statsBlob)
|
|
||||||
|
|
||||||
// On a successful backup, refresh the server's snapshot projection.
|
// On a successful backup, refresh the server's snapshot projection.
|
||||||
// We do this *after* job.finished so the UI sees the job land first;
|
// Do this before job.finished so a failure in terminal reporting cannot
|
||||||
// the snapshot list is a follow-up that the host detail page polls
|
// prevent the independently useful projection refresh.
|
||||||
// or the dashboard sees on its next refresh. A failure here is
|
|
||||||
// logged but doesn't fail the job — the next successful backup will
|
|
||||||
// catch the projection up.
|
|
||||||
if err == nil {
|
if err == nil {
|
||||||
if rerr := r.reportSnapshots(ctx, env); rerr != nil {
|
if rerr := r.reportSnapshots(ctx, env); rerr != nil {
|
||||||
slog.Warn("runner: snapshots.report failed", "job_id", jobID, "err", rerr)
|
slog.Warn("runner: snapshots.report failed", "job_id", jobID, "err", rerr)
|
||||||
@@ -242,6 +243,7 @@ func (r *Runner) RunBackup(ctx context.Context, jobID string, paths, excludes, t
|
|||||||
slog.Warn("runner: stats.report after backup failed", "job_id", jobID, "err", rerr)
|
slog.Warn("runner: stats.report after backup failed", "job_id", jobID, "err", rerr)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
r.sendFinished(ctx, jobID, finishedAt, err, statsBlob)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("runner backup: %w", err)
|
return fmt.Errorf("runner backup: %w", err)
|
||||||
@@ -282,8 +284,6 @@ func (r *Runner) RunForget(ctx context.Context, jobID string, groups []restic.Fo
|
|||||||
var seq atomic.Int64
|
var seq atomic.Int64
|
||||||
err := env.RunForget(ctx, groups, dryRun, r.streamHandler(jobID, &seq))
|
err := env.RunForget(ctx, groups, dryRun, r.streamHandler(jobID, &seq))
|
||||||
finishedAt := time.Now().UTC()
|
finishedAt := time.Now().UTC()
|
||||||
r.sendFinished(ctx, jobID, finishedAt, err, nil)
|
|
||||||
|
|
||||||
// Refresh the server's snapshot projection — forget rewrites the
|
// Refresh the server's snapshot projection — forget rewrites the
|
||||||
// index so the host's snapshot list almost certainly shrunk.
|
// index so the host's snapshot list almost certainly shrunk.
|
||||||
if err == nil {
|
if err == nil {
|
||||||
@@ -292,6 +292,7 @@ func (r *Runner) RunForget(ctx context.Context, jobID string, groups []restic.Fo
|
|||||||
"job_id", jobID, "err", rerr)
|
"job_id", jobID, "err", rerr)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
r.sendFinished(ctx, jobID, finishedAt, err, nil)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("runner forget: %w", err)
|
return fmt.Errorf("runner forget: %w", err)
|
||||||
@@ -318,6 +319,9 @@ func (r *Runner) RunPrune(ctx context.Context, jobID string) error {
|
|||||||
if rerr := r.reportStats(ctx, env, api.RepoStatsPayload{LastPruneAt: &pruneAt}); rerr != nil {
|
if rerr := r.reportStats(ctx, env, api.RepoStatsPayload{LastPruneAt: &pruneAt}); rerr != nil {
|
||||||
slog.Warn("runner: stats.report after prune failed", "job_id", jobID, "err", rerr)
|
slog.Warn("runner: stats.report after prune failed", "job_id", jobID, "err", rerr)
|
||||||
}
|
}
|
||||||
|
if rerr := r.reportSnapshots(ctx, env); rerr != nil {
|
||||||
|
slog.Warn("runner: snapshots.report after prune failed", "job_id", jobID, "err", rerr)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
r.sendFinished(ctx, jobID, finishedAt, err, nil)
|
r.sendFinished(ctx, jobID, finishedAt, err, nil)
|
||||||
|
|||||||
@@ -116,7 +116,8 @@ func envelopeOrder(envs []api.Envelope) []api.MessageType {
|
|||||||
// TestRunPruneShipsExpectedEnvelopes drives RunPrune with a fake
|
// TestRunPruneShipsExpectedEnvelopes drives RunPrune with a fake
|
||||||
// binary that prints "prune" on stdout (for the log.stream envelope)
|
// binary that prints "prune" on stdout (for the log.stream envelope)
|
||||||
// and emits valid stats JSON so reportStats can populate size fields.
|
// and emits valid stats JSON so reportStats can populate size fields.
|
||||||
// Expected sequence: job.started → log.stream → repo.stats → job.finished.
|
// Expected sequence: job.started → log.stream → repo.stats → snapshots.report
|
||||||
|
// → job.finished.
|
||||||
func TestRunPruneShipsExpectedEnvelopes(t *testing.T) {
|
func TestRunPruneShipsExpectedEnvelopes(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
|
|
||||||
@@ -126,6 +127,7 @@ func TestRunPruneShipsExpectedEnvelopes(t *testing.T) {
|
|||||||
case "$1" in
|
case "$1" in
|
||||||
prune) echo "prune" ;;
|
prune) echo "prune" ;;
|
||||||
stats) echo '`+statsJSON+`' ;;
|
stats) echo '`+statsJSON+`' ;;
|
||||||
|
snapshots) echo "[]" ;;
|
||||||
*) echo "unknown: $*" ;;
|
*) echo "unknown: $*" ;;
|
||||||
esac
|
esac
|
||||||
`)
|
`)
|
||||||
@@ -138,7 +140,7 @@ esac
|
|||||||
|
|
||||||
order := envelopeOrder(tx.envs)
|
order := envelopeOrder(tx.envs)
|
||||||
// Confirm landmark envelope types appear in the required order.
|
// Confirm landmark envelope types appear in the required order.
|
||||||
wantTypes := []api.MessageType{api.MsgJobStarted, api.MsgLogStream, api.MsgRepoStats, api.MsgJobFinished}
|
wantTypes := []api.MessageType{api.MsgJobStarted, api.MsgLogStream, api.MsgRepoStats, api.MsgSnapshotsRpt, api.MsgJobFinished}
|
||||||
positions := map[api.MessageType]int{}
|
positions := map[api.MessageType]int{}
|
||||||
for i, mt := range order {
|
for i, mt := range order {
|
||||||
if _, seen := positions[mt]; !seen {
|
if _, seen := positions[mt]; !seen {
|
||||||
@@ -379,6 +381,15 @@ func TestRunInitShipsStartedAndFinished(t *testing.T) {
|
|||||||
_ = firstEnvOfType(t, tx.envs, api.MsgJobFinished)
|
_ = firstEnvOfType(t, tx.envs, api.MsgJobFinished)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func firstIndexOfType(envs []api.Envelope, typ api.MessageType) int {
|
||||||
|
for i, env := range envs {
|
||||||
|
if env.Type == typ {
|
||||||
|
return i
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return -1
|
||||||
|
}
|
||||||
|
|
||||||
// TestRunForgetShipsStartedAndFinished confirms the refactored
|
// TestRunForgetShipsStartedAndFinished confirms the refactored
|
||||||
// RunForget still produces job.started and job.finished envelopes.
|
// RunForget still produces job.started and job.finished envelopes.
|
||||||
func TestRunForgetShipsStartedAndFinished(t *testing.T) {
|
func TestRunForgetShipsStartedAndFinished(t *testing.T) {
|
||||||
@@ -402,5 +413,9 @@ esac
|
|||||||
t.Fatalf("RunForget: %v", err)
|
t.Fatalf("RunForget: %v", err)
|
||||||
}
|
}
|
||||||
_ = firstEnvOfType(t, tx.envs, api.MsgJobStarted)
|
_ = firstEnvOfType(t, tx.envs, api.MsgJobStarted)
|
||||||
_ = firstEnvOfType(t, tx.envs, api.MsgJobFinished)
|
finished := firstIndexOfType(tx.envs, api.MsgJobFinished)
|
||||||
|
refreshed := firstIndexOfType(tx.envs, api.MsgSnapshotsRpt)
|
||||||
|
if refreshed < 0 || finished < 0 || refreshed >= finished {
|
||||||
|
t.Fatalf("snapshot refresh must precede terminal report: %v", envelopeOrder(tx.envs))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -108,6 +108,7 @@ func connectOnce(ctx context.Context, cfg Config, handle Handler) error {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("dial: %w", err)
|
return fmt.Errorf("dial: %w", err)
|
||||||
}
|
}
|
||||||
|
conn.SetReadLimit(api.MaxWebSocketMessageBytes)
|
||||||
// On a successful upgrade coder/websocket transfers ownership of the
|
// On a successful upgrade coder/websocket transfers ownership of the
|
||||||
// response stream to conn and deliberately sets res.Body to nil. Closing
|
// response stream to conn and deliberately sets res.Body to nil. Closing
|
||||||
// the connection below releases that stream.
|
// the connection below releases that stream.
|
||||||
|
|||||||
@@ -2,12 +2,16 @@ package wsclient
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"encoding/json"
|
||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/coder/websocket"
|
"github.com/coder/websocket"
|
||||||
|
|
||||||
|
"gitea.dcglab.co.uk/steve/restic-manager/internal/api"
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestConnectOnceCleanDisconnectDoesNotPanic(t *testing.T) {
|
func TestConnectOnceCleanDisconnectDoesNotPanic(t *testing.T) {
|
||||||
@@ -44,3 +48,56 @@ func TestConnectOnceCleanDisconnectDoesNotPanic(t *testing.T) {
|
|||||||
t.Fatalf("server websocket: %v", err)
|
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.
|
// are evaluated — always-on hosts' stale_schedule stays a no-op.
|
||||||
const staleBackupThreshold = 7 * 24 * time.Hour
|
const staleBackupThreshold = 7 * 24 * time.Hour
|
||||||
|
|
||||||
|
const (
|
||||||
|
defaultStuckJobThreshold = 6 * time.Hour
|
||||||
|
longStuckJobThreshold = 24 * time.Hour
|
||||||
|
)
|
||||||
|
|
||||||
// JobFinishedEvent carries everything the engine needs to evaluate
|
// JobFinishedEvent carries everything the engine needs to evaluate
|
||||||
// the failed-X rules. Pushed via Engine.NotifyJobFinished from the
|
// the failed-X rules. Pushed via Engine.NotifyJobFinished from the
|
||||||
// MarkJobFinished site.
|
// MarkJobFinished site.
|
||||||
@@ -53,6 +58,7 @@ type Engine struct {
|
|||||||
// we raise. Configurable for tests; default 15m.
|
// we raise. Configurable for tests; default 15m.
|
||||||
agentOfflineFloor time.Duration
|
agentOfflineFloor time.Duration
|
||||||
tickPeriod time.Duration
|
tickPeriod time.Duration
|
||||||
|
stuckThresholds map[string]time.Duration
|
||||||
|
|
||||||
closeOnce sync.Once
|
closeOnce sync.Once
|
||||||
done chan struct{}
|
done chan struct{}
|
||||||
@@ -69,6 +75,11 @@ func NewEngine(st *store.Store, hub *notification.Hub) *Engine {
|
|||||||
hostUp: make(chan string, 32),
|
hostUp: make(chan string, 32),
|
||||||
agentOfflineFloor: 15 * time.Minute,
|
agentOfflineFloor: 15 * time.Minute,
|
||||||
tickPeriod: 60 * time.Second,
|
tickPeriod: 60 * time.Second,
|
||||||
|
stuckThresholds: map[string]time.Duration{
|
||||||
|
"backup": longStuckJobThreshold,
|
||||||
|
"restore": longStuckJobThreshold,
|
||||||
|
"check": longStuckJobThreshold,
|
||||||
|
},
|
||||||
done: make(chan struct{}),
|
done: make(chan struct{}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -125,6 +136,10 @@ func (e *Engine) NotifyHostOnline(hostID string) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func (e *Engine) handleJobFinished(ctx context.Context, ev JobFinishedEvent) {
|
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
|
// Determine which kind/severity pair this job maps to. Jobs not
|
||||||
// listed here (init, unlock, restore, diff) produce no alerts in v1.
|
// listed here (init, unlock, restore, diff) produce no alerts in v1.
|
||||||
var kind, severity string
|
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 {
|
if _, err := e.store.CleanupExpiredOIDCState(ctx, now.Add(-5*time.Minute)); err != nil {
|
||||||
slog.Warn("alert: cleanup expired oidc state", "err", err)
|
slog.Warn("alert: cleanup expired oidc state", "err", err)
|
||||||
}
|
}
|
||||||
|
e.evaluateStuckJobs(ctx, now)
|
||||||
|
|
||||||
hosts, err := e.store.ListHosts(ctx)
|
hosts, err := e.store.ListHosts(ctx)
|
||||||
if err != nil {
|
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
|
// roundDur returns a human-readable duration string, rounding to the
|
||||||
// nearest minute. Durations under a minute are reported as "less than
|
// nearest minute. Durations under a minute are reported as "less than
|
||||||
// a minute".
|
// a minute".
|
||||||
|
|||||||
@@ -36,6 +36,10 @@ const (
|
|||||||
// KindAgentOffline is raised when a host's last_seen_at is older
|
// KindAgentOffline is raised when a host's last_seen_at is older
|
||||||
// than the 15-minute floor and resolved when the host reconnects.
|
// than the 15-minute floor and resolved when the host reconnects.
|
||||||
KindAgentOffline = "agent_offline"
|
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
|
// raiseAndNotify is the standard raise pattern: store.RaiseOrTouch
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package alert
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"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) {
|
func TestEngineCheckFailedSeverityCritical(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
eng, st, hostID := setupEngine(t)
|
eng, st, hostID := setupEngine(t)
|
||||||
|
|||||||
@@ -10,6 +10,13 @@ import (
|
|||||||
// (not iota ints) makes traffic readable in logs and packet captures.
|
// (not iota ints) makes traffic readable in logs and packet captures.
|
||||||
type MessageType string
|
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.
|
// Agent → server message types.
|
||||||
const (
|
const (
|
||||||
MsgHello MessageType = "hello"
|
MsgHello MessageType = "hello"
|
||||||
@@ -35,6 +42,7 @@ const (
|
|||||||
MsgConfigUpdate MessageType = "config.update"
|
MsgConfigUpdate MessageType = "config.update"
|
||||||
MsgCommandUpdate MessageType = "command.update"
|
MsgCommandUpdate MessageType = "command.update"
|
||||||
MsgTreeList MessageType = "tree.list" // sync RPC: list a snapshot's children
|
MsgTreeList MessageType = "tree.list" // sync RPC: list a snapshot's children
|
||||||
|
MsgSnapshotsRefresh MessageType = "snapshots.refresh"
|
||||||
)
|
)
|
||||||
|
|
||||||
// Envelope is the framing for every WS message in either direction.
|
// Envelope is the framing for every WS message in either direction.
|
||||||
|
|||||||
@@ -8,7 +8,9 @@ import (
|
|||||||
"net/netip"
|
"net/netip"
|
||||||
"runtime"
|
"runtime"
|
||||||
"strings"
|
"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/config"
|
||||||
"gitea.dcglab.co.uk/steve/restic-manager/internal/server/metrics"
|
"gitea.dcglab.co.uk/steve/restic-manager/internal/server/metrics"
|
||||||
"gitea.dcglab.co.uk/steve/restic-manager/internal/store"
|
"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
|
return metrics.Snapshot{}, err
|
||||||
}
|
}
|
||||||
bySeverity := map[string]int{"info": 0, "warning": 0, "critical": 0}
|
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 {
|
for _, a := range open {
|
||||||
bySeverity[a.Severity]++
|
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
|
reg := s.deps.Metrics
|
||||||
if reg == nil {
|
if reg == nil {
|
||||||
reg = metrics.NewRegistry() // empty histogram block
|
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
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -262,6 +262,7 @@ func (s *Server) routes(r chi.Router) {
|
|||||||
r.Post("/api/hosts/{id}/repo/unlock", s.handleRunRepoUnlock)
|
r.Post("/api/hosts/{id}/repo/unlock", s.handleRunRepoUnlock)
|
||||||
r.Post("/api/jobs/{id}/cancel", s.handleCancelJob)
|
r.Post("/api/jobs/{id}/cancel", s.handleCancelJob)
|
||||||
r.Post("/api/hosts/{id}/snapshots/diff", s.handleSnapshotDiff)
|
r.Post("/api/hosts/{id}/snapshots/diff", s.handleSnapshotDiff)
|
||||||
|
r.Post("/api/hosts/{id}/snapshots/refresh", s.handleRefreshHostSnapshots)
|
||||||
|
|
||||||
// HTMX form variants outside /api.
|
// HTMX form variants outside /api.
|
||||||
r.Post("/hosts/{id}/snapshots/diff", s.handleSnapshotDiff)
|
r.Post("/hosts/{id}/snapshots/diff", s.handleSnapshotDiff)
|
||||||
|
|||||||
@@ -5,6 +5,10 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/go-chi/chi/v5"
|
"github.com/go-chi/chi/v5"
|
||||||
|
"github.com/oklog/ulid/v2"
|
||||||
|
|
||||||
|
"gitea.dcglab.co.uk/steve/restic-manager/internal/api"
|
||||||
|
"gitea.dcglab.co.uk/steve/restic-manager/internal/store"
|
||||||
)
|
)
|
||||||
|
|
||||||
// snapshotView is the public JSON shape for a snapshot. Matches the
|
// snapshotView is the public JSON shape for a snapshot. Matches the
|
||||||
@@ -26,6 +30,7 @@ type listSnapshotsResponse struct {
|
|||||||
HostID string `json:"host_id"`
|
HostID string `json:"host_id"`
|
||||||
Count int `json:"count"`
|
Count int `json:"count"`
|
||||||
RefreshedAt *time.Time `json:"refreshed_at,omitempty"`
|
RefreshedAt *time.Time `json:"refreshed_at,omitempty"`
|
||||||
|
Stale bool `json:"stale"`
|
||||||
Snapshots []snapshotView `json:"snapshots"`
|
Snapshots []snapshotView `json:"snapshots"`
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -45,7 +50,8 @@ func (s *Server) handleListHostSnapshots(w stdhttp.ResponseWriter, r *stdhttp.Re
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if _, err := s.deps.Store.GetHost(r.Context(), hostID); err != nil {
|
host, err := s.deps.Store.GetHost(r.Context(), hostID)
|
||||||
|
if err != nil {
|
||||||
writeJSONError(w, stdhttp.StatusNotFound, "host_not_found", "")
|
writeJSONError(w, stdhttp.StatusNotFound, "host_not_found", "")
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -61,10 +67,13 @@ func (s *Server) handleListHostSnapshots(w stdhttp.ResponseWriter, r *stdhttp.Re
|
|||||||
Count: len(snaps),
|
Count: len(snaps),
|
||||||
Snapshots: make([]snapshotView, len(snaps)),
|
Snapshots: make([]snapshotView, len(snaps)),
|
||||||
}
|
}
|
||||||
if len(snaps) > 0 {
|
out.RefreshedAt = host.SnapshotRefreshedAt
|
||||||
t := snaps[0].RefreshedAt
|
mutationAt, err := s.deps.Store.LatestSuccessfulRepoMutation(r.Context(), hostID)
|
||||||
out.RefreshedAt = &t
|
if err != nil {
|
||||||
|
writeJSONError(w, stdhttp.StatusInternalServerError, "internal", "")
|
||||||
|
return
|
||||||
}
|
}
|
||||||
|
out.Stale = mutationAt != nil && (out.RefreshedAt == nil || out.RefreshedAt.Before(*mutationAt))
|
||||||
for i, sn := range snaps {
|
for i, sn := range snaps {
|
||||||
out.Snapshots[i] = snapshotView{
|
out.Snapshots[i] = snapshotView{
|
||||||
ID: sn.ID,
|
ID: sn.ID,
|
||||||
@@ -80,3 +89,31 @@ func (s *Server) handleListHostSnapshots(w stdhttp.ResponseWriter, r *stdhttp.Re
|
|||||||
|
|
||||||
writeJSON(w, stdhttp.StatusOK, out)
|
writeJSON(w, stdhttp.StatusOK, out)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *Server) handleRefreshHostSnapshots(w stdhttp.ResponseWriter, r *stdhttp.Request) {
|
||||||
|
user, ok := s.requireUser(r)
|
||||||
|
if !ok {
|
||||||
|
writeJSONError(w, stdhttp.StatusUnauthorized, "unauthorised", "")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
hostID := chi.URLParam(r, "id")
|
||||||
|
if _, err := s.deps.Store.GetHost(r.Context(), hostID); err != nil {
|
||||||
|
writeJSONError(w, stdhttp.StatusNotFound, "host_not_found", "")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
if s.deps.Hub == nil || !s.deps.Hub.Connected(hostID) {
|
||||||
|
writeJSONError(w, stdhttp.StatusConflict, "host_offline", "agent is not currently connected")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
env, _ := api.Marshal(api.MsgSnapshotsRefresh, ulid.Make().String(), nil)
|
||||||
|
if err := s.deps.Hub.Send(r.Context(), hostID, env); err != nil {
|
||||||
|
writeJSONError(w, stdhttp.StatusConflict, "host_offline", err.Error())
|
||||||
|
return
|
||||||
|
}
|
||||||
|
now := time.Now().UTC()
|
||||||
|
_ = s.deps.Store.AppendAudit(r.Context(), store.AuditEntry{
|
||||||
|
ID: ulid.Make().String(), UserID: &user.ID, Actor: "user",
|
||||||
|
Action: "host.snapshots_refresh", TargetKind: ptr("host"), TargetID: &hostID, TS: now,
|
||||||
|
})
|
||||||
|
writeJSON(w, stdhttp.StatusAccepted, map[string]string{"status": "refresh_requested"})
|
||||||
|
}
|
||||||
|
|||||||
@@ -0,0 +1,83 @@
|
|||||||
|
package http
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"encoding/json"
|
||||||
|
stdhttp "net/http"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"gitea.dcglab.co.uk/steve/restic-manager/internal/api"
|
||||||
|
"gitea.dcglab.co.uk/steve/restic-manager/internal/store"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestSnapshotsFreshnessAndExplicitRefresh(t *testing.T) {
|
||||||
|
t.Parallel()
|
||||||
|
srv, ts, st := rawTestServerWithUI(t)
|
||||||
|
hostID, token := enrolHostForUI(t, srv, st, "snapshot-refresh-host")
|
||||||
|
c := agentDial(t, srv, ts, hostID, token)
|
||||||
|
sendHello(t, c, "snapshot-refresh-host")
|
||||||
|
_ = drainUntil(t, c, api.MsgScheduleSet)
|
||||||
|
cookie := loginAsAdmin(t, st)
|
||||||
|
|
||||||
|
mutationAt := time.Now().UTC().Add(-time.Minute).Truncate(time.Millisecond)
|
||||||
|
if err := st.CreateJob(context.Background(), store.Job{
|
||||||
|
ID: "mutation-job", HostID: hostID, Kind: "forget", ActorKind: "user", CreatedAt: mutationAt.Add(-time.Minute),
|
||||||
|
}); err != nil {
|
||||||
|
t.Fatalf("create mutation: %v", err)
|
||||||
|
}
|
||||||
|
if err := st.MarkJobFinished(context.Background(), "mutation-job", "succeeded", 0, nil, "", mutationAt); err != nil {
|
||||||
|
t.Fatalf("finish mutation: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
get := func() listSnapshotsResponse {
|
||||||
|
req, _ := stdhttp.NewRequest(stdhttp.MethodGet, ts.URL+"/api/hosts/"+hostID+"/snapshots", nil)
|
||||||
|
req.AddCookie(cookie)
|
||||||
|
res, err := stdhttp.DefaultClient.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("get snapshots: %v", err)
|
||||||
|
}
|
||||||
|
defer res.Body.Close()
|
||||||
|
var body listSnapshotsResponse
|
||||||
|
if err := json.NewDecoder(res.Body).Decode(&body); err != nil {
|
||||||
|
t.Fatalf("decode snapshots: %v", err)
|
||||||
|
}
|
||||||
|
return body
|
||||||
|
}
|
||||||
|
if body := get(); !body.Stale || body.RefreshedAt != nil {
|
||||||
|
t.Fatalf("unrefreshed projection should be stale: %+v", body)
|
||||||
|
}
|
||||||
|
|
||||||
|
refreshedAt := mutationAt.Add(time.Second)
|
||||||
|
if err := st.ReplaceHostSnapshots(context.Background(), hostID, nil, refreshedAt); err != nil {
|
||||||
|
t.Fatalf("replace empty: %v", err)
|
||||||
|
}
|
||||||
|
if body := get(); body.Stale || body.RefreshedAt == nil || !body.RefreshedAt.Equal(refreshedAt) {
|
||||||
|
t.Fatalf("fresh empty projection reported incorrectly: %+v", body)
|
||||||
|
}
|
||||||
|
|
||||||
|
req, _ := stdhttp.NewRequest(stdhttp.MethodPost, ts.URL+"/api/hosts/"+hostID+"/snapshots/refresh", nil)
|
||||||
|
req.AddCookie(cookie)
|
||||||
|
res, err := stdhttp.DefaultClient.Do(req)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("request refresh: %v", err)
|
||||||
|
}
|
||||||
|
defer res.Body.Close()
|
||||||
|
if res.StatusCode != stdhttp.StatusAccepted {
|
||||||
|
t.Fatalf("refresh status = %d, want 202", res.StatusCode)
|
||||||
|
}
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
|
||||||
|
defer cancel()
|
||||||
|
_, raw, err := c.Read(ctx)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("read refresh command: %v", err)
|
||||||
|
}
|
||||||
|
var env api.Envelope
|
||||||
|
if err := json.Unmarshal(raw, &env); err != nil {
|
||||||
|
t.Fatalf("decode envelope: %v", err)
|
||||||
|
}
|
||||||
|
if env.Type != api.MsgSnapshotsRefresh {
|
||||||
|
t.Fatalf("message type = %q, want %q", env.Type, api.MsgSnapshotsRefresh)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -155,6 +155,8 @@ type Snapshot struct {
|
|||||||
BuildCommit string
|
BuildCommit string
|
||||||
GoVersion string
|
GoVersion string
|
||||||
JobDurationRows []HistogramRow
|
JobDurationRows []HistogramRow
|
||||||
|
StuckJobs int
|
||||||
|
OldestStuckAge time.Duration
|
||||||
}
|
}
|
||||||
|
|
||||||
// SnapshotWith builds a Snapshot from raw inputs and the registry's
|
// 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",
|
fmt.Fprintf(&b, "rm_build_info{version=%q,commit=%q,go_version=%q} 1\n",
|
||||||
s.BuildVersion, s.BuildCommit, s.GoVersion)
|
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 -------------------------------------------------
|
// --- Per-host gauges -------------------------------------------------
|
||||||
// Stable order: by host id.
|
// Stable order: by host id.
|
||||||
hosts := append([]HostRow(nil), s.Hosts...)
|
hosts := append([]HostRow(nil), s.Hosts...)
|
||||||
|
|||||||
@@ -114,6 +114,8 @@ func TestRenderGolden(t *testing.T) {
|
|||||||
snap := r.SnapshotWith(hosts,
|
snap := r.SnapshotWith(hosts,
|
||||||
map[string]int{"info": 0, "warning": 1, "critical": 0},
|
map[string]int{"info": 0, "warning": 1, "critical": 0},
|
||||||
"v1.2.3", "deadbeef", "go1.25.0")
|
"v1.2.3", "deadbeef", "go1.25.0")
|
||||||
|
snap.StuckJobs = 2
|
||||||
|
snap.OldestStuckAge = 90 * time.Minute
|
||||||
|
|
||||||
var buf bytes.Buffer
|
var buf bytes.Buffer
|
||||||
if err := Render(&buf, snap); err != nil {
|
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="info"} 0`,
|
||||||
`rm_active_alerts{severity="critical"} 0`,
|
`rm_active_alerts{severity="critical"} 0`,
|
||||||
`rm_build_info{version="v1.2.3",commit="deadbeef",go_version="go1.25.0"} 1`,
|
`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="01H0001",host="alpha"} 1`,
|
||||||
`rm_host_agent_online{host_id="01H0002",host="bravo"} 0`,
|
`rm_host_agent_online{host_id="01H0002",host="bravo"} 0`,
|
||||||
`rm_host_last_backup_timestamp_seconds{host_id="01H0001",host="alpha"} 1700000000`,
|
`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)
|
slog.Warn("ws accept failed", "err", err, "host_id", host.ID)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
conn.SetReadLimit(api.MaxWebSocketMessageBytes)
|
||||||
|
|
||||||
c := NewConn(host.ID, conn)
|
c := NewConn(host.ID, conn)
|
||||||
// Keep agents alive across NAT boxes; coder/websocket
|
// 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")
|
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) {
|
func TestWSRejectsOldProtocol(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
url, token, _, _, _ := setupTestHub(t)
|
url, token, _, _, _ := setupTestHub(t)
|
||||||
|
|||||||
+12
-4
@@ -44,7 +44,7 @@ func (s *Store) LookupHostByAgentToken(ctx context.Context, tokenHash string) (*
|
|||||||
repo_size_bytes, snapshot_count, open_alert_count,
|
repo_size_bytes, snapshot_count, open_alert_count,
|
||||||
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
|
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
|
||||||
pre_hook_default, post_hook_default,
|
pre_hook_default, post_hook_default,
|
||||||
repo_status, repo_status_error, always_on
|
repo_status, repo_status_error, always_on, snapshot_refreshed_at
|
||||||
FROM hosts WHERE agent_token_hash = ?`,
|
FROM hosts WHERE agent_token_hash = ?`,
|
||||||
tokenHash)
|
tokenHash)
|
||||||
return scanHost(row)
|
return scanHost(row)
|
||||||
@@ -59,7 +59,7 @@ func (s *Store) GetHost(ctx context.Context, id string) (*Host, error) {
|
|||||||
repo_size_bytes, snapshot_count, open_alert_count,
|
repo_size_bytes, snapshot_count, open_alert_count,
|
||||||
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
|
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
|
||||||
pre_hook_default, post_hook_default,
|
pre_hook_default, post_hook_default,
|
||||||
repo_status, repo_status_error, always_on
|
repo_status, repo_status_error, always_on, snapshot_refreshed_at
|
||||||
FROM hosts WHERE id = ?`, id)
|
FROM hosts WHERE id = ?`, id)
|
||||||
return scanHost(row)
|
return scanHost(row)
|
||||||
}
|
}
|
||||||
@@ -227,7 +227,7 @@ func (s *Store) ListHosts(ctx context.Context) ([]Host, error) {
|
|||||||
repo_size_bytes, snapshot_count, open_alert_count,
|
repo_size_bytes, snapshot_count, open_alert_count,
|
||||||
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
|
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
|
||||||
pre_hook_default, post_hook_default,
|
pre_hook_default, post_hook_default,
|
||||||
repo_status, repo_status_error, always_on
|
repo_status, repo_status_error, always_on, snapshot_refreshed_at
|
||||||
FROM hosts ORDER BY name`)
|
FROM hosts ORDER BY name`)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("store: list hosts: %w", err)
|
return nil, fmt.Errorf("store: list hosts: %w", err)
|
||||||
@@ -268,6 +268,7 @@ func scanHostRow(s hostScanner) (*Host, error) {
|
|||||||
bwUp, bwDown sql.NullInt64
|
bwUp, bwDown sql.NullInt64
|
||||||
preHook, postHook sql.NullString
|
preHook, postHook sql.NullString
|
||||||
alwaysOn int
|
alwaysOn int
|
||||||
|
snapshotRefreshedAt sql.NullString
|
||||||
)
|
)
|
||||||
err := s.Scan(&h.ID, &h.Name, &h.OS, &h.Arch,
|
err := s.Scan(&h.ID, &h.Name, &h.OS, &h.Arch,
|
||||||
&h.AgentVersion, &h.ResticVersion, &h.ProtocolVersion,
|
&h.AgentVersion, &h.ResticVersion, &h.ProtocolVersion,
|
||||||
@@ -276,7 +277,7 @@ func scanHostRow(s hostScanner) (*Host, error) {
|
|||||||
&h.RepoSizeBytes, &h.SnapshotCount, &h.OpenAlertCount,
|
&h.RepoSizeBytes, &h.SnapshotCount, &h.OpenAlertCount,
|
||||||
&h.AppliedScheduleVersion, &bwUp, &bwDown,
|
&h.AppliedScheduleVersion, &bwUp, &bwDown,
|
||||||
&preHook, &postHook,
|
&preHook, &postHook,
|
||||||
&h.RepoStatus, &h.RepoStatusError, &alwaysOn)
|
&h.RepoStatus, &h.RepoStatusError, &alwaysOn, &snapshotRefreshedAt)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if errors.Is(err, sql.ErrNoRows) {
|
if errors.Is(err, sql.ErrNoRows) {
|
||||||
return nil, ErrNotFound
|
return nil, ErrNotFound
|
||||||
@@ -332,6 +333,13 @@ func scanHostRow(s hostScanner) (*Host, error) {
|
|||||||
h.PostHookDefault = postHook.String
|
h.PostHookDefault = postHook.String
|
||||||
}
|
}
|
||||||
h.AlwaysOn = alwaysOn != 0
|
h.AlwaysOn = alwaysOn != 0
|
||||||
|
if snapshotRefreshedAt.Valid {
|
||||||
|
t, err := time.Parse(time.RFC3339Nano, snapshotRefreshedAt.String)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("store: parse snapshot_refreshed_at: %w", err)
|
||||||
|
}
|
||||||
|
h.SnapshotRefreshedAt = &t
|
||||||
|
}
|
||||||
return &h, nil
|
return &h, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -27,6 +27,73 @@ type Job struct {
|
|||||||
CreatedAt time.Time
|
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
|
||||||
|
}
|
||||||
|
|
||||||
|
// LatestSuccessfulRepoMutation returns the newest completion time for a job
|
||||||
|
// that can change the repository's snapshot projection.
|
||||||
|
func (s *Store) LatestSuccessfulRepoMutation(ctx context.Context, hostID string) (*time.Time, error) {
|
||||||
|
var raw sql.NullString
|
||||||
|
err := s.db.QueryRowContext(ctx, `
|
||||||
|
SELECT MAX(finished_at) FROM jobs
|
||||||
|
WHERE host_id = ? AND status = 'succeeded'
|
||||||
|
AND kind IN ('backup', 'forget', 'prune')`, hostID).Scan(&raw)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("store: latest successful repo mutation: %w", err)
|
||||||
|
}
|
||||||
|
if !raw.Valid {
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
|
t, err := time.Parse(time.RFC3339Nano, raw.String)
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("store: parse latest repo mutation: %w", err)
|
||||||
|
}
|
||||||
|
return &t, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// 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
|
// CreateJob inserts a queued job. The agent will mark it running
|
||||||
// when it actually starts work. ScheduledID is set when the job
|
// when it actually starts work. ScheduledID is set when the job
|
||||||
// originates from a cron fire (actor_kind="schedule"); nil for
|
// originates from a cron fire (actor_kind="schedule"); nil for
|
||||||
|
|||||||
@@ -7,6 +7,48 @@ import (
|
|||||||
"time"
|
"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) {
|
func TestLatestJobByKind(t *testing.T) {
|
||||||
t.Parallel()
|
t.Parallel()
|
||||||
s := openTestStore(t)
|
s := openTestStore(t)
|
||||||
|
|||||||
@@ -0,0 +1 @@
|
|||||||
|
ALTER TABLE hosts ADD COLUMN snapshot_refreshed_at TEXT;
|
||||||
@@ -69,8 +69,8 @@ func (s *Store) ReplaceHostSnapshots(ctx context.Context, hostID string, snaps [
|
|||||||
}
|
}
|
||||||
|
|
||||||
if _, err := tx.ExecContext(ctx,
|
if _, err := tx.ExecContext(ctx,
|
||||||
`UPDATE hosts SET snapshot_count = ? WHERE id = ?`,
|
`UPDATE hosts SET snapshot_count = ?, snapshot_refreshed_at = ? WHERE id = ?`,
|
||||||
len(snaps), hostID); err != nil {
|
len(snaps), when.UTC().Format(time.RFC3339Nano), hostID); err != nil {
|
||||||
return fmt.Errorf("store: update host snapshot_count: %w", err)
|
return fmt.Errorf("store: update host snapshot_count: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -138,7 +138,8 @@ func TestReplaceHostSnapshotsEmpty(t *testing.T) {
|
|||||||
t.Fatalf("replace 1: %v", err)
|
t.Fatalf("replace 1: %v", err)
|
||||||
}
|
}
|
||||||
// Then empty — host has been wiped.
|
// Then empty — host has been wiped.
|
||||||
if err := s.ReplaceHostSnapshots(ctx, hostID, nil, time.Now().UTC()); err != nil {
|
refreshedAt := time.Now().UTC().Truncate(time.Millisecond)
|
||||||
|
if err := s.ReplaceHostSnapshots(ctx, hostID, nil, refreshedAt); err != nil {
|
||||||
t.Fatalf("replace empty: %v", err)
|
t.Fatalf("replace empty: %v", err)
|
||||||
}
|
}
|
||||||
out, err := s.ListSnapshotsByHost(ctx, hostID)
|
out, err := s.ListSnapshotsByHost(ctx, hostID)
|
||||||
@@ -152,4 +153,7 @@ func TestReplaceHostSnapshotsEmpty(t *testing.T) {
|
|||||||
if h.SnapshotCount != 0 {
|
if h.SnapshotCount != 0 {
|
||||||
t.Errorf("snapshot_count should reset to 0, got %d", h.SnapshotCount)
|
t.Errorf("snapshot_count should reset to 0, got %d", h.SnapshotCount)
|
||||||
}
|
}
|
||||||
|
if h.SnapshotRefreshedAt == nil || !h.SnapshotRefreshedAt.Equal(refreshedAt) {
|
||||||
|
t.Errorf("empty projection refresh time = %v, want %s", h.SnapshotRefreshedAt, refreshedAt)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -77,6 +77,7 @@ type Host struct {
|
|||||||
LastBackupStatus *string
|
LastBackupStatus *string
|
||||||
RepoSizeBytes int64
|
RepoSizeBytes int64
|
||||||
SnapshotCount int
|
SnapshotCount int
|
||||||
|
SnapshotRefreshedAt *time.Time
|
||||||
OpenAlertCount int
|
OpenAlertCount int
|
||||||
AppliedScheduleVersion int64
|
AppliedScheduleVersion int64
|
||||||
// Host-wide bandwidth caps applied to every restic invocation
|
// Host-wide bandwidth caps applied to every restic invocation
|
||||||
|
|||||||
Reference in New Issue
Block a user