31 Commits

Author SHA1 Message Date
steve bc5c99898c Merge pull request 'Stop polling idle Agent updates controls' (#52) from fix-agent-update-refresh into main 2026-08-22 12:33:45 +01:00
steve 80dac33f0a Stop polling idle agent update controls
CI / Test (rest) (pull_request) Successful in 22s
CI / Lint (pull_request) Successful in 11s
CI / Build (windows/amd64) (pull_request) Successful in 8s
CI / Test (store) (pull_request) Successful in 40s
CI / Build (linux/amd64) (pull_request) Successful in 7s
CI / Build (linux/arm64) (pull_request) Successful in 8s
CI / Test (server-http) (pull_request) Successful in 1m45s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m26s
2026-08-22 12:31:08 +01:00
steve facef5b254 Merge pull request 'Fix Agent updates filters and refresh state' (#51) from fix-issue-50-agent-update-filters into main 2026-08-22 12:06:44 +01:00
steve 6374201f6c Fix fleet update filters across refreshes
CI / Test (store) (pull_request) Successful in 5s
CI / Test (rest) (pull_request) Successful in 8s
CI / Build (windows/amd64) (pull_request) Successful in 7s
CI / Lint (pull_request) Successful in 11s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 7s
CI / Test (server-http) (pull_request) Successful in 1m35s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m26s
2026-08-22 12:04:26 +01:00
steve 9080826b06 Merge pull request 'Release v1.2.0 — fleet updates and reliability' (#49) from release-v1.2.0 into main
Release / Build + push image (push) Successful in 2m10s
2026-08-22 11:34:51 +01:00
steve 39a0875d58 test(alert): drain notifications before teardown
CI / Test (rest) (pull_request) Successful in 39s
CI / Lint (pull_request) Successful in 11s
CI / Test (store) (pull_request) Successful in 43s
CI / Build (windows/amd64) (pull_request) Successful in 7s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 8s
CI / Test (server-http) (pull_request) Successful in 1m34s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m27s
2026-08-22 11:31:56 +01:00
steve 320be34970 docs(changelog): prepare v1.2.0
CI / Test (store) (pull_request) Successful in 15s
CI / Lint (pull_request) Successful in 11s
CI / Test (rest) (pull_request) Failing after 33s
CI / Build (windows/amd64) (pull_request) Successful in 7s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 7s
CI / Test (server-http) (pull_request) Successful in 1m44s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m27s
2026-08-22 11:29:59 +01:00
steve 25aefe90d6 Use anonymously pullable CI images 2026-08-22 11:29:15 +01:00
steve 56d97f13af Merge pull request 'Add selectable canary-first agent updates' (#48) from fix-issue-43-agent-updates-workflow into main 2026-08-22 11:11:41 +01:00
steve 383bdb7d36 Add selectable canary-first fleet updates
CI / Test (rest) (pull_request) Successful in 39s
CI / Test (store) (pull_request) Successful in 42s
CI / Lint (pull_request) Successful in 11s
CI / Build (windows/amd64) (pull_request) Successful in 7s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 8s
CI / Test (server-http) (pull_request) Successful in 1m45s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m17s
2026-08-22 11:08:53 +01:00
steve b64f029892 Merge pull request 'Make snapshot projections fresh and self-diagnosing' (#47) from fix-issue-40-snapshot-reconciliation into main 2026-08-22 11:03:47 +01:00
steve f9718e6077 Make snapshot projections self-diagnosing
CI / Test (rest) (pull_request) Successful in 40s
CI / Test (store) (pull_request) Successful in 42s
CI / Lint (pull_request) Successful in 11s
CI / Build (windows/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 7s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Test (server-http) (pull_request) Successful in 1m32s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m26s
2026-08-22 11:00:46 +01:00
steve 8fdf4a1bdf Merge pull request 'Detect and alert on stuck jobs' (#46) from fix-issue-41-stuck-job-alerts into main 2026-08-22 10:56:31 +01:00
steve 750a06384b Close websocket handshake response in test
CI / Test (server-http) (pull_request) Successful in 5s
CI / Test (rest) (pull_request) Successful in 23s
CI / Lint (pull_request) Successful in 10s
CI / Build (windows/amd64) (pull_request) Successful in 7s
CI / Test (store) (pull_request) Successful in 41s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 6s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m24s
2026-08-22 10:53:51 +01:00
steve 2110c7dab4 Detect and alert on stuck jobs
CI / Test (rest) (pull_request) Successful in 39s
CI / Test (store) (pull_request) Successful in 40s
CI / Lint (pull_request) Failing after 11s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 7s
CI / Build (windows/amd64) (pull_request) Successful in 25s
CI / Test (server-http) (pull_request) Successful in 1m39s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m16s
2026-08-22 10:52:13 +01:00
steve 8ddd3456e1 Merge pull request 'Fix WebSocket limit for large agent messages' (#45) 2026-08-22 10:48:31 +01:00
steve 74cb7f660f fix(ws): raise bounded agent message limit
CI / Test (store) (pull_request) Successful in 6s
CI / Lint (pull_request) Failing after 20s
CI / Build (windows/amd64) (pull_request) Successful in 8s
CI / Test (rest) (pull_request) Successful in 38s
CI / Build (linux/amd64) (pull_request) Successful in 7s
CI / Build (linux/arm64) (pull_request) Successful in 8s
CI / Test (server-http) (pull_request) Successful in 1m28s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m24s
2026-08-22 10:46:00 +01:00
steve 1131f4330c Merge pull request 'Release v1.1.1 — forget reliability fixes' (#42)
Release / Build + push image (push) Successful in 2m12s
2026-08-22 10:12:53 +01:00
steve 337472a819 docs(changelog): prepare v1.1.1
CI / Test (store) (pull_request) Successful in 5s
CI / Test (rest) (pull_request) Successful in 8s
CI / Build (windows/amd64) (pull_request) Successful in 8s
CI / Lint (pull_request) Successful in 11s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 25s
CI / Test (server-http) (pull_request) Successful in 1m30s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m37s
2026-08-22 10:10:12 +01:00
steve b315932cc8 chore(web): normalize generated stylesheet 2026-08-22 10:09:42 +01:00
steve f6fa84d7d8 docs: expand issue contribution guidance 2026-08-22 10:09:42 +01:00
steve 904c522a23 Merge pull request 'Fix manual forget payload and dry-run execution' (#39) 2026-08-22 09:56:12 +01:00
steve 31f53d65a7 fix(forget): populate retention groups for manual runs
CI / Test (store) (pull_request) Successful in 5s
CI / Lint (pull_request) Successful in 10s
CI / Build (windows/amd64) (pull_request) Successful in 7s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 7s
CI / Test (rest) (pull_request) Successful in 39s
CI / Test (server-http) (pull_request) Successful in 1m37s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m16s
2026-08-22 09:53:31 +01:00
steve 528bdef433 Merge pull request 'Fix agent panic on WebSocket disconnect' (#38) 2026-08-22 09:49:37 +01:00
steve dfe082629f chore(lint): document websocket response ownership
CI / Test (rest) (pull_request) Successful in 38s
CI / Lint (pull_request) Successful in 7s
CI / Test (store) (pull_request) Successful in 39s
CI / Build (windows/amd64) (pull_request) Successful in 22s
CI / Build (linux/amd64) (pull_request) Successful in 8s
CI / Build (linux/arm64) (pull_request) Successful in 7s
CI / Test (server-http) (pull_request) Successful in 1m38s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m22s
2026-08-22 09:46:50 +01:00
steve 39aff83837 fix(agent): avoid panic on websocket disconnect
CI / Test (server-http) (pull_request) Successful in 5s
CI / Test (rest) (pull_request) Successful in 7s
CI / Test (store) (pull_request) Successful in 5s
CI / Build (windows/amd64) (pull_request) Successful in 7s
CI / Lint (pull_request) Failing after 10s
CI / Build (linux/arm64) (pull_request) Successful in 8s
CI / Build (linux/amd64) (pull_request) Successful in 24s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m36s
2026-08-22 09:45:20 +01:00
steve 27be28ee9c Merge pull request 'docs: correct admin-credentials help text (fixes #34)' (#35) from docs/admin-creds-help-text into main 2026-08-21 22:39:04 +01:00
Steve Cliff a8a6fdfab5 docs: address review — make S3/B2/SFTP/local guidance actionable
CI / Test (store) (pull_request) Successful in 37s
CI / Test (rest) (pull_request) Successful in 48s
CI / Build (windows/amd64) (pull_request) Successful in 8s
CI / Lint (pull_request) Successful in 20s
CI / Build (linux/arm64) (pull_request) Successful in 24s
CI / Build (linux/amd64) (pull_request) Successful in 26s
CI / Test (server-http) (pull_request) Successful in 1m37s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m16s
The previous wording told operators to leave the slot blank and then
warned that doing so disables prune, which is contradictory. Prune is
gated on the admin slot for every backend, on both paths:
scheduled (maintenance_dispatch.go:52) skips silently, manual
(repo_ops.go:39) returns 400 admin_creds_required.

Also reworded the opening line, which still said 'only needed for
rest-server repos' while the new text asks other backends to fill it in.
2026-08-21 22:36:14 +01:00
Steve Cliff e9df802478 docs: correct admin-credentials help text
CI / Test (rest) (pull_request) Successful in 1m42s
CI / Test (store) (pull_request) Successful in 1m45s
CI / Lint (pull_request) Successful in 28s
CI / Build (windows/amd64) (pull_request) Successful in 29s
CI / Test (server-http) (pull_request) Successful in 2m27s
CI / Build (linux/amd64) (pull_request) Successful in 30s
CI / Build (linux/arm64) (pull_request) Successful in 26s
e2e / Playwright vs docker-compose (pull_request) Successful in 1m30s
Forget does not use admin credentials. Only JobPrune builds a runner
against the admin slot (cmd/agent/main.go:622-645); JobForget runs on
the everyday runner, as the comment at main.go:517-524 states.

Also corrects the claim that leaving this blank lets everyday creds
handle prune: maintenance_dispatch.go:52-56 skips prune entirely with
'prune skipped — no admin creds' on any backend, with no fallback.

Refs #34
2026-08-21 21:01:39 +01:00
steve 6c6b962e24 Merge pull request 'De-flake TestDrainPendingSerializesPerHost (CI stability)' (#33) from fix-flaky-server-http-tests into main
Reviewed-on: #33
2026-06-16 15:44:47 +01:00
steve e64075d5d7 test(pending-drain): de-flake TestDrainPendingSerializesPerHost
CI / Test (store) (pull_request) Successful in 8s
CI / Test (rest) (pull_request) Successful in 12s
CI / Build (windows/amd64) (pull_request) Successful in 15s
CI / Lint (pull_request) Successful in 19s
CI / Build (linux/amd64) (pull_request) Successful in 12s
CI / Build (linux/arm64) (pull_request) Successful in 44s
CI / Test (server-http) (pull_request) Successful in 2m55s
e2e / Playwright vs docker-compose (pull_request) Successful in 2m45s
Keep the test WS client actively reading (a real agent always is) so
the server-side conn stays registered under parallel load, and drain to
completion via condition polling instead of asserting one-shot
completeness. The conn could be dropped/unregistered under CI load,
making DrainPending correctly no-op (conn==nil) and the test observe a
partial/empty drain. -race confirms no production data race; the
exactly-5-jobs assertion (proving the per-host mutex blocks
double-dispatch) is unchanged. Verified: 0 failures over 25 loaded runs
+ 4 -race iterations.
2026-06-16 13:29:47 +01:00
47 changed files with 1366 additions and 165 deletions
+6 -15
View File
@@ -3,7 +3,7 @@
# Notes for anyone editing this file:
#
# Custom runner image
# Every job runs inside `gitea.dcglab.co.uk/steve/ci-runner-go`
# Every job runs inside `docker.dcglab.co.uk/public/ci-runner-go:latest`
# (recipe: https://gitea.dcglab.co.uk/steve/ci/src/branch/main/images/ci-runner-go).
# That image already ships:
# * Go on PATH at /usr/local/go/bin (so `actions/setup-go` is
@@ -11,8 +11,8 @@
# otherwise re-download Go on every job)
# * Node.js + npm (used by docs / e2e workflows)
# * Docker CLI, Buildx, Compose v2 (used by docker-build steps)
# When bumping the Go floor, push a new ci-runner-go image with
# the matching Go version and bump the date pin in IMAGE below.
# The runner image is intentionally tracked via `latest`: it is shared CI
# infrastructure, and updates should propagate to consuming repositories.
#
# Self-hosted runner expectations
# Each runner host bind-mounts persistent volumes for
@@ -71,10 +71,7 @@ jobs:
name: Test (${{ matrix.name }})
runs-on: ubuntu-latest
container:
image: docker.dcglab.co.uk/ci-runner-go:2026-05-15
credentials:
username: ${{ secrets.ZOT_USERNAME }}
password: ${{ secrets.ZOT_PASSWORD }}
image: docker.dcglab.co.uk/public/ci-runner-go:latest
strategy:
fail-fast: false
matrix:
@@ -110,10 +107,7 @@ jobs:
name: Lint
runs-on: ubuntu-latest
container:
image: docker.dcglab.co.uk/ci-runner-go:2026-05-15
credentials:
username: ${{ secrets.ZOT_USERNAME }}
password: ${{ secrets.ZOT_PASSWORD }}
image: docker.dcglab.co.uk/public/ci-runner-go:latest
steps:
- uses: actions/checkout@v4
- uses: golangci/golangci-lint-action@v7
@@ -130,10 +124,7 @@ jobs:
name: Build (${{ matrix.goos }}/${{ matrix.goarch }})
runs-on: ubuntu-latest
container:
image: docker.dcglab.co.uk/ci-runner-go:2026-05-15
credentials:
username: ${{ secrets.ZOT_USERNAME }}
password: ${{ secrets.ZOT_PASSWORD }}
image: docker.dcglab.co.uk/public/ci-runner-go:latest
strategy:
fail-fast: false
matrix:
+1 -1
View File
@@ -31,7 +31,7 @@ jobs:
e2e:
name: Playwright vs docker-compose
runs-on: ubuntu-latest
container: gitea.dcglab.co.uk/steve/ci-runner-go:2026-05-08
container: docker.dcglab.co.uk/public/ci-runner-go:latest
timeout-minutes: 15
steps:
- uses: actions/checkout@v4
+5 -6
View File
@@ -12,7 +12,9 @@
# plus install.sh / install.ps1 / the systemd unit baked in under
# /opt/restic-manager/dist (the read-only fallback path the server
# handlers use when <DataDir>/... is empty).
# * Pushes to zot OCI registry (docker.dcglab.co.uk).
# * Pushes to the public namespace in the zot OCI registry
# (docker.dcglab.co.uk/public/restic-manager). Pulls are anonymous;
# pushes remain authenticated.
#
# Tag fan-out
# * tag push: :vX.Y.Z, :X.Y, :X
@@ -29,7 +31,7 @@ on:
env:
REGISTRY: docker.dcglab.co.uk
IMAGE_NAME: restic-manager
IMAGE_NAME: public/restic-manager
# Force bash as the default shell — see ci.yml header.
defaults:
@@ -41,10 +43,7 @@ jobs:
name: Build + push image
runs-on: ubuntu-latest
container:
image: docker.dcglab.co.uk/ci-runner-go:2026-05-15
credentials:
username: ${{ secrets.ZOT_USERNAME }}
password: ${{ secrets.ZOT_PASSWORD }}
image: docker.dcglab.co.uk/public/ci-runner-go:latest
steps:
- uses: actions/checkout@v4
+63 -1
View File
@@ -6,6 +6,57 @@ and the project follows [Semantic Versioning](https://semver.org/).
## [Unreleased]
## [1.2.0] - 2026-08-22
### Added
- Fleet-wide agent updates are now discoverable from Settings and the
dashboard, with arbitrary host subsets, name/tag/version/eligibility
filtering, explicit exclusion reasons, server-side membership validation,
and a canary-first pause after the first verified reconnect. Halted rolls
expose retry and resume actions. ([#43])
- Running jobs with no recent persisted activity are detected server-side and
surfaced through deduplicated alerts. Kind-aware thresholds protect normal
long-running work, terminal jobs self-resolve, and Prometheus exports stuck
job count and oldest inactivity age. ([#41])
- Snapshot projections now retain host-level refresh timestamps—including
authoritative empty reports—expose a derived stale flag, refresh after
backup, forget, and prune, and support an explicit operator reconciliation
endpoint. ([#40])
### Fixed
- Raised the bounded WebSocket read limit on both peers so ordinary large
snapshot reports and restic events no longer disconnect otherwise healthy
agents. Regression coverage exercises payloads beyond the library's former
32 KiB default in both directions. ([#44])
- Fleet-update timeout verification now performs a final authoritative read of
the agent version delivered by the reconnect `hello`, avoiding a poll/deadline
race and reporting the last observed version when verification fails. ([#43])
### Changed
- CI runner images and release images use the anonymous-pull `public`
namespace at `docker.dcglab.co.uk`; registry authentication remains required
only for publishing.
## [1.1.1] - 2026-08-22
### Fixed
- Prevented agents from panicking when a WebSocket connection ends. The
WebSocket library transfers ownership of a successful upgrade stream to
the connection and leaves the HTTP response body nil; attempting to close
that body caused agents to restart and left forget jobs permanently stuck
in `running`. ([#37])
- Manual forget jobs now receive the same per-source-group retention policies
as scheduled forget jobs. The API also supports a validated `--dry-run`
option through the complete server-to-agent-to-restic path, and rejects
hosts without configured retention before creating a job. ([#36])
- Corrected the admin-credentials help text to reflect that forget uses normal
append-only credentials and blank admin credentials do not provide a
fallback for prune. ([#34])
## [1.1.0] - 2026-06-15
### Added
@@ -123,5 +174,16 @@ with a web UI, JSON API, and self-updating agent fleet.
go vet, golangci-lint).
- Threat model published (`docs/threat-model.md`).
[Unreleased]: https://gitea.dcglab.co.uk/steve/restic-manager/compare/v1.0.0...HEAD
[Unreleased]: https://gitea.dcglab.co.uk/steve/restic-manager/compare/v1.2.0...HEAD
[1.2.0]: https://gitea.dcglab.co.uk/steve/restic-manager/compare/v1.1.1...v1.2.0
[1.1.1]: https://gitea.dcglab.co.uk/steve/restic-manager/compare/v1.1.0...v1.1.1
[1.1.0]: https://gitea.dcglab.co.uk/steve/restic-manager/releases/tag/v1.1.0
[1.0.0]: https://gitea.dcglab.co.uk/steve/restic-manager/releases/tag/v1.0.0
[#37]: https://gitea.dcglab.co.uk/steve/restic-manager/issues/37
[#36]: https://gitea.dcglab.co.uk/steve/restic-manager/issues/36
[#34]: https://gitea.dcglab.co.uk/steve/restic-manager/issues/34
[#40]: https://gitea.dcglab.co.uk/steve/restic-manager/issues/40
[#41]: https://gitea.dcglab.co.uk/steve/restic-manager/issues/41
[#43]: https://gitea.dcglab.co.uk/steve/restic-manager/issues/43
[#44]: https://gitea.dcglab.co.uk/steve/restic-manager/issues/44
+37 -11
View File
@@ -2,7 +2,8 @@
Thanks for your interest in restic-manager. This document covers how
to set up a development environment, the conventions the project
follows, and how patches make it from your machine into `main`.
follows, and how to contribute through issues as well as patches that
make it from your machine into `main`.
## Project status and scope
@@ -108,6 +109,32 @@ admin user.
## Workflow
### Opening an issue
Issues are contributions too. Use them to report a bug, suggest a
feature, improve the documentation, or start a design discussion even
if you do not plan to submit a patch.
Before opening one, search the existing issues and check `tasks.md` to
see whether the topic is already tracked. Then choose the closest issue
template:
- [Bug report](./.gitea/issue_template/bug_report.md) for behaviour that
does not match the documentation or expected operation.
- [Feature request](./.gitea/issue_template/feature_request.md) for a new
capability or a change to existing behaviour.
Give the issue a specific title, keep it to one problem or proposal,
and complete the relevant template fields. If no template is an exact
fit, open a regular issue and explain the context, desired outcome, and
any alternatives you have considered. Maintainers may ask follow-up
questions or close requests that duplicate existing work or fall
outside the project's scope.
Security-sensitive reports are the exception: follow the
[SECURITY.md](./SECURITY.md) disclosure process and do not open a public
issue.
### Before opening a PR
1. **Open an issue first** for non-trivial changes. The design is
@@ -136,25 +163,24 @@ The PR template asks for:
### Reporting bugs
Open an issue with:
Use the bug report issue template and include:
- restic-manager version (`server --version`) and agent version.
- restic version on the affected host.
- Steps to reproduce.
- Server and agent logs (sanitise any tokens before pasting).
Security-sensitive bugs go through the [SECURITY.md](./SECURITY.md)
disclosure path instead — please don't open a public issue for
them.
For security-sensitive bugs, use the private disclosure process noted
above.
### Suggesting features
Open an issue describing the use case (not just the proposed
solution). The roadmap in `tasks.md` shows where the project is
heading; if the suggestion fits a future phase we'll wire it in
there. If it falls outside the project's scope (multi-tenancy, SaaS,
non-restic backends — see `spec.md` §2 non-goals) we'll say so
early to save your time.
Use the feature request issue template and describe the use case (not
just the proposed solution). The roadmap in `tasks.md` shows where the
project is heading; if the suggestion fits a future phase we'll wire it
in there. If it falls outside the project's scope (multi-tenancy, SaaS,
non-restic backends — see `spec.md` §2 non-goals) we'll say so early to
save your time.
## Code of conduct
+26 -2
View File
@@ -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)
case api.MsgSnapshotsRefresh:
go d.refreshSnapshots(ctx, tx)
case api.MsgScheduleSet:
var p api.ScheduleSetPayload
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
}
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
// the matching tree.list.result envelope back, correlated by the
// request envelope's ID. Errors (missing creds, restic failure)
@@ -601,6 +620,11 @@ func (d *dispatcher) runJob(ctx context.Context, p api.CommandRunPayload, tx wsc
failJob(p, tx, "forget: command.run carried no forget_groups (server didn't populate them)")
return fmt.Errorf("forget: command.run carried no forget_groups (server didn't populate them)")
}
if len(p.Args) > 1 || (len(p.Args) == 1 && p.Args[0] != "--dry-run") {
failJob(p, tx, "forget: command.run carried unsupported arguments")
return fmt.Errorf("forget: command.run carried unsupported arguments")
}
dryRun := len(p.Args) == 1
groups := make([]restic.ForgetGroup, 0, len(p.ForgetGroups))
for _, g := range p.ForgetGroups {
groups = append(groups, restic.ForgetGroup{
@@ -615,9 +639,9 @@ func (d *dispatcher) runJob(ctx context.Context, p api.CommandRunPayload, tx wsc
},
})
}
slog.Info("agent: accepting forget job", "job_id", p.JobID, "groups", len(groups))
slog.Info("agent: accepting forget job", "job_id", p.JobID, "groups", len(groups), "dry_run", dryRun)
spawn("forget", func(jobCtx context.Context) error {
return r.RunForget(jobCtx, p.JobID, groups)
return r.RunForget(jobCtx, p.JobID, groups, dryRun)
})
case api.JobPrune:
// Prune may require admin creds (delete authority on rest-server).
+15 -11
View File
@@ -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.
func (r *Runner) sendStarted(jobID string, kind api.JobKind, startedAt time.Time) {
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.
// We do this *after* job.finished so the UI sees the job land first;
// the snapshot list is a follow-up that the host detail page polls
// 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.
// Do this before job.finished so a failure in terminal reporting cannot
// prevent the independently useful projection refresh.
if err == nil {
if rerr := r.reportSnapshots(ctx, env); rerr != nil {
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)
}
}
r.sendFinished(ctx, jobID, finishedAt, err, statsBlob)
if err != nil {
return fmt.Errorf("runner backup: %w", err)
@@ -274,16 +276,14 @@ func (r *Runner) RunInit(ctx context.Context, jobID string) error {
// snapshot projection (forget rewrites the snapshot index — the
// host's snapshot list shrinks). Snapshot refresh runs once after
// every group completes, not per-group.
func (r *Runner) RunForget(ctx context.Context, jobID string, groups []restic.ForgetGroup) error {
func (r *Runner) RunForget(ctx context.Context, jobID string, groups []restic.ForgetGroup, dryRun bool) error {
startedAt := time.Now().UTC()
r.sendStarted(jobID, api.JobForget, startedAt)
env := r.resticEnv()
var seq atomic.Int64
err := env.RunForget(ctx, groups, r.streamHandler(jobID, &seq))
err := env.RunForget(ctx, groups, dryRun, r.streamHandler(jobID, &seq))
finishedAt := time.Now().UTC()
r.sendFinished(ctx, jobID, finishedAt, err, nil)
// Refresh the server's snapshot projection — forget rewrites the
// index so the host's snapshot list almost certainly shrunk.
if err == nil {
@@ -292,6 +292,7 @@ func (r *Runner) RunForget(ctx context.Context, jobID string, groups []restic.Fo
"job_id", jobID, "err", rerr)
}
}
r.sendFinished(ctx, jobID, finishedAt, err, nil)
if err != nil {
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 {
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)
+19 -4
View File
@@ -116,7 +116,8 @@ func envelopeOrder(envs []api.Envelope) []api.MessageType {
// TestRunPruneShipsExpectedEnvelopes drives RunPrune with a fake
// binary that prints "prune" on stdout (for the log.stream envelope)
// 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) {
t.Parallel()
@@ -126,6 +127,7 @@ func TestRunPruneShipsExpectedEnvelopes(t *testing.T) {
case "$1" in
prune) echo "prune" ;;
stats) echo '`+statsJSON+`' ;;
snapshots) echo "[]" ;;
*) echo "unknown: $*" ;;
esac
`)
@@ -138,7 +140,7 @@ esac
order := envelopeOrder(tx.envs)
// 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{}
for i, mt := range order {
if _, seen := positions[mt]; !seen {
@@ -379,6 +381,15 @@ func TestRunInitShipsStartedAndFinished(t *testing.T) {
_ = 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
// RunForget still produces job.started and job.finished envelopes.
func TestRunForgetShipsStartedAndFinished(t *testing.T) {
@@ -398,9 +409,13 @@ esac
Tag: "documents",
Policy: restic.ForgetPolicy{KeepLast: &keepLast},
}}
if err := r.RunForget(context.Background(), "job-forget", groups); err != nil {
if err := r.RunForget(context.Background(), "job-forget", groups, false); err != nil {
t.Fatalf("RunForget: %v", err)
}
_ = 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))
}
}
+5 -5
View File
@@ -103,15 +103,15 @@ func connectOnce(ctx context.Context, cfg Config, handle Handler) error {
}
dialCtx, cancel := context.WithTimeout(ctx, 30*time.Second)
conn, res, err := websocket.Dial(dialCtx, wsURL, dialOpts)
conn, _, err := websocket.Dial(dialCtx, wsURL, dialOpts) //nolint:bodyclose // successful upgrades have a nil response body owned by conn
cancel()
if err != nil {
return fmt.Errorf("dial: %w", err)
}
// websocket.Dial returns the upgrade response separately from the
// conn. Body is empty on a successful upgrade but Go's net/http
// still expects it closed to release the connection.
defer func() { _ = res.Body.Close() }()
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.
defer conn.CloseNow() //nolint:errcheck
// Send hello.
+103
View File
@@ -0,0 +1,103 @@
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) {
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
// Wait for the agent hello so Dial and the first client write have both
// completed before ending the connection normally.
if _, _, err := conn.Read(r.Context()); err != nil {
serverErr <- err
return
}
serverErr <- conn.Close(websocket.StatusNormalClosure, "test complete")
}))
defer srv.Close()
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
defer cancel()
err := connectOnce(ctx, Config{
ServerURL: srv.URL,
AgentToken: "test-token",
HeartbeatPeriod: time.Hour,
}, nil)
if err == nil {
t.Fatal("connectOnce returned nil after server disconnected")
}
if err := <-serverErr; err != nil {
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")
}
}
+68 -1
View File
@@ -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,11 +58,23 @@ type Engine struct {
// we raise. Configurable for tests; default 15m.
agentOfflineFloor time.Duration
tickPeriod time.Duration
stuckThresholds map[string]time.Duration
closeOnce sync.Once
notifyWG sync.WaitGroup
done chan struct{}
}
func (e *Engine) dispatchNotification(ctx context.Context, payload notification.Payload) {
e.notifyWG.Add(1)
go func() {
defer e.notifyWG.Done()
e.hub.Dispatch(ctx, payload)
}()
}
func (e *Engine) waitNotifications() { e.notifyWG.Wait() }
// NewEngine builds the engine. agentOfflineFloor + tickPeriod default
// to 15min and 60s respectively when zero.
func NewEngine(st *store.Store, hub *notification.Hub) *Engine {
@@ -69,7 +86,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 +147,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 +236,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 +284,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".
+8 -4
View File
@@ -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
@@ -56,7 +60,7 @@ func (e *Engine) raiseAndNotify(ctx context.Context, hostID, kind, dedupKey, sev
if err == nil {
hostName = host.Name
}
go e.hub.Dispatch(ctx, notification.Payload{
e.dispatchNotification(ctx, notification.Payload{
Event: notification.EventRaised,
AlertID: id,
Severity: severity,
@@ -81,7 +85,7 @@ func (e *Engine) Acknowledge(ctx context.Context, alertID, userID string, when t
return nil //nolint:nilerr
}
p := alertPayload(ctx, e.store, notification.EventAcknowledged, a)
go e.hub.Dispatch(context.WithoutCancel(ctx), p)
e.dispatchNotification(context.WithoutCancel(ctx), p)
return nil
}
@@ -95,7 +99,7 @@ func (e *Engine) Resolve(ctx context.Context, alertID string, when time.Time) er
return nil
}
p := alertPayload(ctx, e.store, notification.EventResolved, a)
go e.hub.Dispatch(context.WithoutCancel(ctx), p)
e.dispatchNotification(context.WithoutCancel(ctx), p)
return nil
}
@@ -160,7 +164,7 @@ func (e *Engine) resolveAndNotify(ctx context.Context, hostID, kind, dedupKey st
if a.Kind != kind || a.DedupKey != dedupKey {
continue
}
go e.hub.Dispatch(ctx, notification.Payload{
e.dispatchNotification(ctx, notification.Payload{
Event: notification.EventResolved,
AlertID: a.ID,
Severity: a.Severity,
+61
View File
@@ -3,6 +3,7 @@ package alert
import (
"context"
"path/filepath"
"strings"
"testing"
"time"
@@ -24,6 +25,7 @@ func setupEngine(t *testing.T) (*Engine, *store.Store, string) {
aead, _ := crypto.NewAEAD(key)
hub := notification.NewHub(st, aead, "https://rm.example")
eng := NewEngine(st, hub)
t.Cleanup(eng.waitNotifications)
hostID := ulid.Make().String()
if err := st.CreateHost(context.Background(), store.Host{
ID: hostID, Name: "alfa-01", OS: "linux", Arch: "amd64",
@@ -69,6 +71,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)
+14 -6
View File
@@ -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"
@@ -29,12 +36,13 @@ const (
// Server → agent message types.
const (
MsgCommandRun MessageType = "command.run"
MsgCommandCancel MessageType = "command.cancel"
MsgScheduleSet MessageType = "schedule.set"
MsgConfigUpdate MessageType = "config.update"
MsgCommandUpdate MessageType = "command.update"
MsgTreeList MessageType = "tree.list" // sync RPC: list a snapshot's children
MsgCommandRun MessageType = "command.run"
MsgCommandCancel MessageType = "command.cancel"
MsgScheduleSet MessageType = "schedule.set"
MsgConfigUpdate MessageType = "config.update"
MsgCommandUpdate MessageType = "command.update"
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.
+4 -1
View File
@@ -322,7 +322,7 @@ type ForgetGroup struct {
// any keep-* would delete every snapshot in the tagged set).
// Returns the first error encountered, or nil when every group runs
// to a clean exit.
func (e Env) RunForget(ctx context.Context, groups []ForgetGroup, handle LineHandler) error {
func (e Env) RunForget(ctx context.Context, groups []ForgetGroup, dryRun bool, handle LineHandler) error {
if len(groups) == 0 {
return fmt.Errorf("restic forget: refusing to run with no groups (would be a no-op)")
}
@@ -332,6 +332,9 @@ func (e Env) RunForget(ctx context.Context, groups []ForgetGroup, handle LineHan
}
args := []string{"forget", "--json", "--tag", g.Tag}
args = append(args, g.Policy.args()...)
if dryRun {
args = append(args, "--dry-run")
}
cmd := e.resticCmd(ctx, args...)
if err := runWithPump(cmd, handle); err != nil {
return err
+20
View File
@@ -60,6 +60,26 @@ func TestRunPruneInvokesPrune(t *testing.T) {
t.Fatalf("expected 'prune' in captured output; got: %v", *lines)
}
func TestRunForgetDryRunArgument(t *testing.T) {
bin := setupScriptBin(t, `echo "$@"`)
env := Env{Bin: bin}
lines, h := captureLines()
keepLast := 1
groups := []ForgetGroup{{
Tag: "documents",
Policy: ForgetPolicy{KeepLast: &keepLast},
}}
if err := env.RunForget(context.Background(), groups, true, h); err != nil {
t.Fatalf("RunForget: %v", err)
}
for _, line := range *lines {
if strings.Contains(line, "forget --json --tag documents --keep-last 1 --dry-run") {
return
}
}
t.Fatalf("expected forget invocation with --dry-run; got: %v", *lines)
}
// --- B2: RunCheck ---
func TestRunCheckLockSniff(t *testing.T) {
+31 -1
View File
@@ -71,6 +71,16 @@ func NewWorker(st *store.Store, hub Hub, disp Dispatcher, alerts AlertRaiser) *W
// worker goroutine. Returns the new fleet_update_id on success.
// store.ErrFleetUpdateRunning bubbles up unchanged.
func (w *Worker) Start(ctx context.Context, userID, targetVersion string, hostIDs []string) (string, error) {
return w.start(ctx, userID, targetVersion, hostIDs, false)
}
// StartCanary starts the same sequential rollout but pauses after the first
// verified host so an operator can inspect it before starting the remainder.
func (w *Worker) StartCanary(ctx context.Context, userID, targetVersion string, hostIDs []string) (string, error) {
return w.start(ctx, userID, targetVersion, hostIDs, true)
}
func (w *Worker) start(ctx context.Context, userID, targetVersion string, hostIDs []string, pauseAfterFirst bool) (string, error) {
if userID == "" || targetVersion == "" {
return "", errors.New("fleetupdate: userID and targetVersion required")
}
@@ -85,6 +95,7 @@ func (w *Worker) Start(ctx context.Context, userID, targetVersion string, hostID
StartedByUserID: userID,
TargetVersion: targetVersion,
Status: "running",
PauseAfterFirst: pauseAfterFirst,
}, hostIDs); err != nil {
return "", err
}
@@ -137,6 +148,13 @@ func (w *Worker) run(ctx context.Context, fuID, userID, targetVersion string) {
next := pending[0]
w.processHost(ctx, fuID, userID, next)
if fu.PauseAfterFirst && next.Position == 0 {
updated, _, gerr := w.store.GetFleetUpdate(ctx, fuID)
if gerr == nil && updated.Status == "running" {
_ = w.store.HaltFleetUpdate(ctx, fuID, "canary succeeded; review the host, then resume remaining agents", time.Now().UTC())
}
return
}
}
}
@@ -195,7 +213,19 @@ func (w *Worker) processHost(ctx context.Context, fuID, userID string, slot stor
return
}
}
reason := fmt.Sprintf("timeout waiting for %s to reach %s", hostID, w.targetVersion)
// One authoritative final read closes the poll/deadline race. AgentVersion
// is written directly from the reconnect hello handshake.
lastVersion := "unknown"
if h, err := w.store.GetHost(ctx, hostID); err == nil && h != nil {
if h.AgentVersion == w.targetVersion {
_ = w.store.SetFleetUpdateHostStatus(ctx, fuID, hostID, "succeeded", "", jobID)
return
}
if h.AgentVersion != "" {
lastVersion = h.AgentVersion
}
}
reason := fmt.Sprintf("timeout waiting for %s reconnect hello at %s (last reported %s)", hostID, w.targetVersion, lastVersion)
_ = w.store.SetFleetUpdateHostStatus(ctx, fuID, hostID, "failed", reason, jobID)
w.halt(ctx, fuID, reason)
}
@@ -160,6 +160,35 @@ func TestWorkerTwoHostsBothSucceed(t *testing.T) {
}
}
func TestWorkerCanaryPausesAfterFirstVerifiedHost(t *testing.T) {
st := openStore(t)
uid := mustCreateAdmin(t, st)
h1 := mustCreateHost(t, st, "canary", "v0")
h2 := mustCreateHost(t, st, "remainder", "v0")
hub := &fakeHub{online: map[string]bool{h1: true, h2: true}}
disp := &fakeDispatcher{st: st, target: "v2", delayMS: 20}
alerts := &recAlert{}
w := NewWorker(st, hub, disp, alerts)
w.pollPeriod = 10 * time.Millisecond
w.hostTimeout = time.Second
fuID, err := w.StartCanary(context.Background(), uid, "v2", []string{h1, h2})
if err != nil {
t.Fatalf("start canary: %v", err)
}
fu := waitForStatus(t, st, fuID, "halted", 2*time.Second)
if fu.HaltedReason != "canary succeeded; review the host, then resume remaining agents" {
t.Fatalf("halt reason: %q", fu.HaltedReason)
}
_, hosts, _ := st.GetFleetUpdate(context.Background(), fuID)
if hosts[0].Status != "succeeded" || hosts[1].Status != "pending" {
t.Fatalf("canary statuses: %+v", hosts)
}
if len(alerts.reasons) != 0 {
t.Fatalf("successful canary pause must not alert: %v", alerts.reasons)
}
}
func TestWorkerSecondHostTimesOutHalts(t *testing.T) {
st := openStore(t)
uid := mustCreateAdmin(t, st)
+109 -12
View File
@@ -32,6 +32,13 @@ import (
type fleetUpdateStartReq struct {
TargetVersion string `json:"target_version,omitempty"`
HostIDs []string `json:"host_ids,omitempty"`
CanaryFirst bool `json:"canary_first,omitempty"`
}
type fleetUpdateCandidate struct {
Host store.Host
Eligible bool
Reason string
}
// fleetUpdateHostView is one row in the JSON response for GET
@@ -56,6 +63,7 @@ type fleetUpdateView struct {
CurrentHostID string `json:"current_host_id,omitempty"`
HaltedReason string `json:"halted_reason,omitempty"`
CompletedAt *string `json:"completed_at,omitempty"`
CanaryFirst bool `json:"canary_first"`
Hosts []fleetUpdateHostView `json:"hosts"`
}
@@ -65,6 +73,7 @@ type fleetUpdateView struct {
type fleetUpdatePage struct {
// Idle-state fields.
OutOfDateHosts []store.Host // online hosts whose version != target
Candidates []fleetUpdateCandidate
TargetVersion string
// Active-state fields. Nil when no fleet update has ever run.
@@ -100,6 +109,11 @@ func (s *Server) handleAPIFleetUpdateStart(w stdhttp.ResponseWriter, r *stdhttp.
if target == "" {
target = version.Version
}
if target != version.Version {
writeJSONError(w, stdhttp.StatusUnprocessableEntity, "unsupported_target_version",
"fleet updates can only target the running server version")
return
}
hostIDs := body.HostIDs
if len(hostIDs) == 0 {
derived, err := s.deriveOutOfDateOnlineHostIDs(r.Context(), target)
@@ -108,6 +122,19 @@ func (s *Server) handleAPIFleetUpdateStart(w stdhttp.ResponseWriter, r *stdhttp.
return
}
hostIDs = derived
} else {
validated, reasons, err := s.validateFleetUpdateHostIDs(r.Context(), target, hostIDs)
if err != nil {
writeJSONError(w, stdhttp.StatusInternalServerError, "internal", err.Error())
return
}
if len(reasons) > 0 {
writeJSON(w, stdhttp.StatusUnprocessableEntity, map[string]any{
"code": "ineligible_hosts", "message": "one or more selected hosts are ineligible", "reasons": reasons,
})
return
}
hostIDs = validated
}
if len(hostIDs) == 0 {
writeJSONError(w, stdhttp.StatusConflict, "no_hosts_eligible",
@@ -115,7 +142,20 @@ func (s *Server) handleAPIFleetUpdateStart(w stdhttp.ResponseWriter, r *stdhttp.
return
}
fuID, err := s.deps.FleetWorker.Start(r.Context(), user.ID, target, hostIDs)
var fuID string
var err error
if body.CanaryFirst && len(hostIDs) > 1 {
worker, supported := s.deps.FleetWorker.(interface {
StartCanary(context.Context, string, string, []string) (string, error)
})
if !supported {
writeJSONError(w, stdhttp.StatusServiceUnavailable, "canary_unavailable", "")
return
}
fuID, err = worker.StartCanary(r.Context(), user.ID, target, hostIDs)
} else {
fuID, err = s.deps.FleetWorker.Start(r.Context(), user.ID, target, hostIDs)
}
if err != nil {
if errors.Is(err, store.ErrFleetUpdateRunning) {
writeJSONError(w, stdhttp.StatusConflict, "fleet_update_in_progress", err.Error())
@@ -129,6 +169,7 @@ func (s *Server) handleAPIFleetUpdateStart(w stdhttp.ResponseWriter, r *stdhttp.
"fleet_update_id": fuID,
"target_version": target,
"host_count": len(hostIDs),
"canary_first": body.CanaryFirst,
})
_ = s.deps.Store.AppendAudit(r.Context(), store.AuditEntry{
ID: ulid.Make().String(), UserID: &user.ID, Actor: "user",
@@ -209,6 +250,7 @@ func (s *Server) handleAPIFleetUpdateGet(w stdhttp.ResponseWriter, r *stdhttp.Re
Status: fu.Status,
CurrentHostID: fu.CurrentHostID,
HaltedReason: fu.HaltedReason,
CanaryFirst: fu.PauseAfterFirst,
Hosts: make([]fleetUpdateHostView, 0, len(hosts)),
}
if fu.CompletedAt != nil {
@@ -286,6 +328,27 @@ func (s *Server) buildFleetUpdatePage(r *stdhttp.Request) (fleetUpdatePage, erro
}
for _, h := range hosts {
page.HostNames[h.ID] = h.Name
candidate := fleetUpdateCandidate{Host: h}
switch {
case h.AgentVersion == "":
candidate.Reason = "version unknown"
case h.AgentVersion == page.TargetVersion:
candidate.Reason = "already current"
case s.deps.Hub == nil || !s.deps.Hub.Connected(h.ID):
candidate.Reason = "offline"
default:
updating, uerr := s.deps.Store.RunningUpdateJobForHost(r.Context(), h.ID)
if uerr != nil {
return page, uerr
}
if updating != "" {
candidate.Reason = "update already running"
} else {
candidate.Eligible = true
page.OutOfDateHosts = append(page.OutOfDateHosts, h)
}
}
page.Candidates = append(page.Candidates, candidate)
}
active, err := s.deps.Store.ActiveFleetUpdate(r.Context())
@@ -328,20 +391,54 @@ func (s *Server) buildFleetUpdatePage(r *stdhttp.Request) (fleetUpdatePage, erro
}
}
// Idle list (or "still out of date" reference even when an active
// roll is running — cheap to compute, harmless to attach).
for _, h := range hosts {
if h.Status != "online" {
continue
}
if h.AgentVersion == "" || h.AgentVersion == page.TargetVersion {
continue
}
page.OutOfDateHosts = append(page.OutOfDateHosts, h)
}
return page, nil
}
// validateFleetUpdateHostIDs deduplicates an explicit selection while
// preserving review order and rejects the entire request if membership has
// changed or any selected host is not dispatchable.
func (s *Server) validateFleetUpdateHostIDs(ctx context.Context, target string, requested []string) ([]string, map[string]string, error) {
hosts, err := s.deps.Store.ListHosts(ctx)
if err != nil {
return nil, nil, err
}
byID := make(map[string]store.Host, len(hosts))
for _, h := range hosts {
byID[h.ID] = h
}
seen := map[string]bool{}
validated := make([]string, 0, len(requested))
reasons := map[string]string{}
for _, id := range requested {
if seen[id] {
continue
}
seen[id] = true
h, ok := byID[id]
switch {
case !ok:
reasons[id] = "host not found"
case h.AgentVersion == "":
reasons[id] = "agent version unknown"
case h.AgentVersion == target:
reasons[id] = "already at target version"
case s.deps.Hub == nil || !s.deps.Hub.Connected(id):
reasons[id] = "host offline"
default:
jobID, jerr := s.deps.Store.RunningUpdateJobForHost(ctx, id)
if jerr != nil {
return nil, nil, jerr
}
if jobID != "" {
reasons[id] = "update already in progress"
} else {
validated = append(validated, id)
}
}
}
return validated, reasons, nil
}
// deriveOutOfDateOnlineHostIDs returns the list of host IDs that
// (a) are online (Hub.Connected) and (b) have an agent_version that's
// non-empty AND != target. Used by the start endpoint when the caller
+128
View File
@@ -6,7 +6,9 @@ import (
"bytes"
"context"
"encoding/json"
"io"
stdhttp "net/http"
"strings"
"sync"
"testing"
"time"
@@ -19,6 +21,67 @@ import (
"gitea.dcglab.co.uk/steve/restic-manager/internal/version"
)
func TestFleetUpdateIdlePageDoesNotPollOrRequireTypedConfirmation(t *testing.T) {
t.Parallel()
_, baseURL, st := newTestServerWithUI(t)
cookie := loginAsAdmin(t, st)
_ = makeHost(t, st, "fleet-filter-host")
req, _ := stdhttp.NewRequest("GET", baseURL+"/settings/fleet-update", nil)
req.AddCookie(cookie)
res, err := stdhttp.DefaultClient.Do(req)
if err != nil {
t.Fatalf("get fleet update page: %v", err)
}
defer res.Body.Close()
if res.StatusCode != stdhttp.StatusOK {
t.Fatalf("status: got %d, want 200", res.StatusCode)
}
raw, _ := io.ReadAll(res.Body)
body := string(raw)
for _, want := range []string{
`id="fleet-filter-name"`,
`id="fleet-filter-version"`,
`id="fleet-filter-state"`,
"r.style.display=",
} {
if !strings.Contains(body, want) {
t.Errorf("page missing %q", want)
}
}
if strings.Contains(body, "Type selected count to confirm") || strings.Contains(body, `id="fleet-update-confirm"`) {
t.Error("page still renders the typed-count confirmation")
}
if strings.Contains(body, `hx-trigger="every 3s`) {
t.Error("idle selection page must not poll and replace operator input")
}
}
func TestFleetUpdateRunningPagePollsForProgress(t *testing.T) {
t.Parallel()
_, baseURL, st := newTestServerWithUI(t)
cookie, userID := loginAsAdminWithID(t, st)
hostID := makeHost(t, st, "fleet-running-host")
if err := st.CreateFleetUpdate(context.Background(), store.FleetUpdate{
ID: ulid.Make().String(), StartedByUserID: userID, TargetVersion: version.Version,
}, []string{hostID}); err != nil {
t.Fatalf("create fleet update: %v", err)
}
req, _ := stdhttp.NewRequest("GET", baseURL+"/settings/fleet-update", nil)
req.AddCookie(cookie)
res, err := stdhttp.DefaultClient.Do(req)
if err != nil {
t.Fatalf("get fleet update page: %v", err)
}
defer res.Body.Close()
raw, _ := io.ReadAll(res.Body)
if body := string(raw); !strings.Contains(body, `hx-trigger="every 3s`) {
t.Error("running rollout page must poll for progress")
}
}
// fakeFleetWorker stands in for *fleetupdate.Worker in HTTP tests.
// It records what was passed to Start/Cancel and lets tests inject
// canned errors. Satisfies the FleetWorker interface in
@@ -50,6 +113,10 @@ func (f *fakeFleetWorker) Start(_ context.Context, userID, target string, hostID
return f.startID, nil
}
func (f *fakeFleetWorker) StartCanary(ctx context.Context, userID, target string, hostIDs []string) (string, error) {
return f.Start(ctx, userID, target, hostIDs)
}
func (f *fakeFleetWorker) Cancel(_ context.Context, id string) error {
f.mu.Lock()
defer f.mu.Unlock()
@@ -190,6 +257,67 @@ func TestFleetUpdateStartDerivesHostIDsWhenEmpty(t *testing.T) {
}
}
func TestFleetUpdateStartDeduplicatesExplicitSelection(t *testing.T) {
t.Parallel()
srv, ts, st := rawTestServer(t)
worker := &fakeFleetWorker{startID: ulid.Make().String()}
srv.deps.FleetWorker = worker
cookie := loginAsAdmin(t, st)
hostID := helloOnlineHost(t, srv, st, "duplicate-host", "v0")
raw, _ := json.Marshal(map[string]any{"host_ids": []string{hostID, hostID}})
req, _ := stdhttp.NewRequest("POST", ts.URL+"/api/fleet/update", bytes.NewReader(raw))
req.AddCookie(cookie)
req.Header.Set("Content-Type", "application/json")
res, err := stdhttp.DefaultClient.Do(req)
if err != nil {
t.Fatalf("do: %v", err)
}
defer res.Body.Close()
if res.StatusCode != stdhttp.StatusAccepted {
t.Fatalf("status: %d", res.StatusCode)
}
if got := worker.startCalls[0].HostIDs; len(got) != 1 || got[0] != hostID {
t.Fatalf("deduplicated ids: %v", got)
}
}
func TestFleetUpdateStartRejectsUnknownAndIneligibleHosts(t *testing.T) {
t.Parallel()
srv, ts, st := rawTestServer(t)
worker := &fakeFleetWorker{startID: ulid.Make().String()}
srv.deps.FleetWorker = worker
cookie := loginAsAdmin(t, st)
offline := makeHost(t, st, "offline-host")
if err := st.MarkHostHello(context.Background(), offline, "v0", "0.17", api.CurrentProtocolVersion, time.Now().UTC()); err != nil {
t.Fatalf("mark offline host: %v", err)
}
raw, _ := json.Marshal(map[string]any{"host_ids": []string{offline, "does-not-exist"}})
req, _ := stdhttp.NewRequest("POST", ts.URL+"/api/fleet/update", bytes.NewReader(raw))
req.AddCookie(cookie)
req.Header.Set("Content-Type", "application/json")
res, err := stdhttp.DefaultClient.Do(req)
if err != nil {
t.Fatalf("do: %v", err)
}
defer res.Body.Close()
if res.StatusCode != stdhttp.StatusUnprocessableEntity {
t.Fatalf("status: got %d, want 422", res.StatusCode)
}
var body struct {
Code string `json:"code"`
Reasons map[string]string `json:"reasons"`
}
if err := json.NewDecoder(res.Body).Decode(&body); err != nil {
t.Fatalf("decode: %v", err)
}
if body.Code != "ineligible_hosts" || len(body.Reasons) != 2 {
t.Fatalf("structured reasons: %+v", body)
}
if len(worker.startCalls) != 0 {
t.Fatal("worker must not start with invalid membership")
}
}
func TestFleetUpdateCancelHappyPath(t *testing.T) {
t.Parallel()
srv, ts, st := rawTestServer(t)
+24 -2
View File
@@ -65,10 +65,32 @@ func (s *Server) handleRunNow(w stdhttp.ResponseWriter, r *stdhttp.Request) {
func (s *Server) dispatchJob(ctx context.Context, user *store.User,
hostID string, kind api.JobKind, args []string,
) (res runNowResponse, status int, code, msg string) {
return s.dispatchJobWithPayload(ctx, user, hostID, kind, nil, api.CommandRunPayload{
payload := api.CommandRunPayload{
Kind: kind,
Args: args,
})
}
if kind == api.JobForget {
if !validForgetArgs(args) {
return res, stdhttp.StatusBadRequest, "invalid_args",
"forget accepts no arguments other than --dry-run"
}
var ok bool
var err error
payload, ok, err = s.buildForgetPayloadForHost(ctx, hostID)
if err != nil {
return res, stdhttp.StatusInternalServerError, "internal", ""
}
if !ok {
return res, stdhttp.StatusUnprocessableEntity, "no_retention_policy",
"host has no source groups with a retention policy"
}
payload.Args = args
}
return s.dispatchJobWithPayload(ctx, user, hostID, kind, nil, payload)
}
func validForgetArgs(args []string) bool {
return len(args) == 0 || (len(args) == 1 && args[0] == "--dry-run")
}
// dispatchJobWithPayload is dispatchJob's variant that lets callers
+106
View File
@@ -0,0 +1,106 @@
package http
import (
"bytes"
"context"
"encoding/json"
stdhttp "net/http"
"testing"
"time"
"github.com/oklog/ulid/v2"
"gitea.dcglab.co.uk/steve/restic-manager/internal/api"
"gitea.dcglab.co.uk/steve/restic-manager/internal/store"
)
func TestRunNowForgetShipsRetentionGroupsAndDryRun(t *testing.T) {
t.Parallel()
srv, ts, st := rawTestServer(t)
hostID, token := enrolHostForWS(t, srv, st, "manual-forget-host")
seedInitJob(t, st, hostID)
keepDaily := 7
if err := st.CreateSourceGroup(context.Background(), &store.SourceGroup{
ID: ulid.Make().String(),
HostID: hostID,
Name: "documents",
Includes: []string{"/home/documents"},
RetentionPolicy: store.RetentionPolicy{KeepDaily: &keepDaily},
}); err != nil {
t.Fatalf("create source group: %v", err)
}
c := agentDial(t, srv, ts, hostID, token)
sendHello(t, c, "manual-forget-host")
_ = drainUntil(t, c, api.MsgScheduleSet)
body, err := json.Marshal(runNowRequest{Kind: api.JobForget, Args: []string{"--dry-run"}})
if err != nil {
t.Fatalf("marshal request: %v", err)
}
req, err := stdhttp.NewRequest(stdhttp.MethodPost, ts.URL+"/api/hosts/"+hostID+"/jobs", bytes.NewReader(body))
if err != nil {
t.Fatalf("new request: %v", err)
}
req.Header.Set("Content-Type", "application/json")
req.AddCookie(loginAsAdmin(t, st))
res, err := stdhttp.DefaultClient.Do(req)
if err != nil {
t.Fatalf("post run-now: %v", err)
}
defer res.Body.Close()
if res.StatusCode != stdhttp.StatusAccepted {
t.Fatalf("status: got %d, want %d", res.StatusCode, stdhttp.StatusAccepted)
}
got := readNextCommandRun(t, c, time.Now().Add(2*time.Second))
if got == nil {
t.Fatal("no command.run received")
}
if len(got.Args) != 1 || got.Args[0] != "--dry-run" {
t.Fatalf("Args: got %q, want [--dry-run]", got.Args)
}
if len(got.ForgetGroups) != 1 {
t.Fatalf("ForgetGroups: got %d, want 1", len(got.ForgetGroups))
}
group := got.ForgetGroups[0]
if group.Tag != "documents" || group.Policy.KeepDaily == nil || *group.Policy.KeepDaily != 7 {
t.Fatalf("ForgetGroups[0]: got %+v", group)
}
}
func TestRunNowForgetRejectsHostWithoutRetention(t *testing.T) {
t.Parallel()
srv, ts, st := rawTestServer(t)
hostID, token := enrolHostForWS(t, srv, st, "no-manual-retention-host")
seedInitJob(t, st, hostID)
c := agentDial(t, srv, ts, hostID, token)
sendHello(t, c, "no-manual-retention-host")
_ = drainUntil(t, c, api.MsgScheduleSet)
body := []byte(`{"kind":"forget","args":["--dry-run"]}`)
req, err := stdhttp.NewRequest(stdhttp.MethodPost, ts.URL+"/api/hosts/"+hostID+"/jobs", bytes.NewReader(body))
if err != nil {
t.Fatalf("new request: %v", err)
}
req.Header.Set("Content-Type", "application/json")
req.AddCookie(loginAsAdmin(t, st))
res, err := stdhttp.DefaultClient.Do(req)
if err != nil {
t.Fatalf("post run-now: %v", err)
}
defer res.Body.Close()
if res.StatusCode != stdhttp.StatusUnprocessableEntity {
t.Fatalf("status: got %d, want %d", res.StatusCode, stdhttp.StatusUnprocessableEntity)
}
var jobs int
if err := st.DB().QueryRow(`SELECT COUNT(*) FROM jobs WHERE host_id = ? AND kind = 'forget'`, hostID).Scan(&jobs); err != nil {
t.Fatalf("count forget jobs: %v", err)
}
if jobs != 0 {
t.Fatalf("forget jobs: got %d, want 0", jobs)
}
}
+10 -6
View File
@@ -37,7 +37,12 @@ func (s *Server) DispatchMaintenance(ctx context.Context, decisions []maintenanc
}
switch d.Kind {
case "forget":
payload, ok := s.buildForgetPayloadForHost(ctx, d.HostID)
payload, ok, err := s.buildForgetPayloadForHost(ctx, d.HostID)
if err != nil {
slog.Warn("maintenance: list source groups failed",
"host_id", d.HostID, "err", err)
continue
}
if !ok {
slog.Info("maintenance: forget skipped — no source groups with retention",
"host_id", d.HostID)
@@ -88,11 +93,10 @@ func (s *Server) DispatchMaintenance(ctx context.Context, decisions []maintenanc
// that has a non-empty retention policy and builds a CommandRunPayload
// with ForgetGroups populated. Returns ok=false if the host has no
// such groups (the dispatcher then skips this kind).
func (s *Server) buildForgetPayloadForHost(ctx context.Context, hostID string) (api.CommandRunPayload, bool) {
func (s *Server) buildForgetPayloadForHost(ctx context.Context, hostID string) (api.CommandRunPayload, bool, error) {
groups, err := s.deps.Store.ListSourceGroupsByHost(ctx, hostID)
if err != nil {
slog.Warn("maintenance: list source groups failed", "host_id", hostID, "err", err)
return api.CommandRunPayload{}, false
return api.CommandRunPayload{}, false, err
}
fg := make([]api.ForgetGroup, 0, len(groups))
for _, g := range groups {
@@ -105,9 +109,9 @@ func (s *Server) buildForgetPayloadForHost(ctx context.Context, hostID string) (
})
}
if len(fg) == 0 {
return api.CommandRunPayload{}, false
return api.CommandRunPayload{}, false, nil
}
return api.CommandRunPayload{ForgetGroups: fg}, true
return api.CommandRunPayload{ForgetGroups: fg}, true, nil
}
func isEmptyRetention(p store.RetentionPolicy) bool {
+27 -1
View File
@@ -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
}
+36 -17
View File
@@ -512,11 +512,27 @@ func TestDrainPendingSerializesPerHost(t *testing.T) {
// Connect the agent so DrainPending can dispatch.
c := agentDial(t, srv, ts, hostID, token)
sendHello(t, c, "serialise-host")
// Drain the on-hello goroutine's pass first (no pending rows yet),
// then wait for the schedule.set so the connection is fully settled.
// Wait for the on-hello push to settle.
_ = drainUntil(t, c, api.MsgScheduleSet)
// Insert 5 pending rows now that the on-hello drain has already run.
// A real agent is always in a read loop. Keep this test client
// reading in the background for the rest of the test: without an
// active reader the server-side conn can be dropped under parallel
// load, which unregisters it from the hub and makes DrainPending
// no-op (conn == nil) — the historical source of this test's
// flakiness (it would observe 0 or a partial drain). The reader also
// consumes the command.run envelopes our drains emit.
readerCtx, stopReader := context.WithCancel(context.Background())
defer stopReader()
go func() {
for {
if _, _, err := c.Read(readerCtx); err != nil {
return
}
}
}()
// Insert 5 due pending rows.
now := time.Now().UTC()
for i := range 5 {
pid := ulid.Make().String()
@@ -533,7 +549,8 @@ func TestDrainPendingSerializesPerHost(t *testing.T) {
}
}
// Spawn 10 goroutines all calling DrainPending concurrently.
// Fire 10 concurrent DrainPending calls. The per-host mutex must
// ensure each row is dispatched at most once (no double-dispatch).
var wg sync.WaitGroup
for range 10 {
wg.Add(1)
@@ -544,24 +561,26 @@ func TestDrainPendingSerializesPerHost(t *testing.T) {
}
wg.Wait()
// Drain any envelopes the agent received so we don't block below.
// We read with short timeouts and stop when the connection goes quiet.
drainDeadline := time.Now().Add(500 * time.Millisecond)
for time.Now().Before(drainDeadline) {
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond)
_, _, err := c.Read(ctx)
cancel()
if err != nil {
break
}
// Drain to completion. The fire-and-forget on-hello DrainPending
// shares the same per-host mutex and can hold it during the burst,
// leaving rows for a later pass — exactly how production drains
// (repeatedly, via the 30s tick / on reconnect). Re-drain until the
// queue is empty; because every drain is still serialised, each row
// is dispatched at most once, so the exactly-5 job count below proves
// there was no double-dispatch.
deadline := time.Now().Add(5 * time.Second)
for countPendingForHost(t, st, hostID) > 0 && time.Now().Before(deadline) {
srv.DrainPending(context.Background(), hostID)
time.Sleep(10 * time.Millisecond)
}
// All 5 pending rows must be gone.
// All 5 pending rows must be drained.
if n := countPendingForHost(t, st, hostID); n != 0 {
t.Errorf("pending rows after concurrent drain: got %d, want 0", n)
t.Errorf("pending rows after drain-to-completion: got %d, want 0", n)
}
// Exactly 5 backup job rows (one per pending row), not 10+ from a race.
// Exactly 5 backup job rows (one per pending row) — never more, which
// would mean the per-host mutex failed to prevent double-dispatch.
var n int
_ = st.DB().QueryRow(
`SELECT COUNT(*) FROM jobs WHERE host_id = ? AND kind = 'backup' AND actor_kind = 'schedule'`,
+1
View File
@@ -262,6 +262,7 @@ func (s *Server) routes(r chi.Router) {
r.Post("/api/hosts/{id}/repo/unlock", s.handleRunRepoUnlock)
r.Post("/api/jobs/{id}/cancel", s.handleCancelJob)
r.Post("/api/hosts/{id}/snapshots/diff", s.handleSnapshotDiff)
r.Post("/api/hosts/{id}/snapshots/refresh", s.handleRefreshHostSnapshots)
// HTMX form variants outside /api.
r.Post("/hosts/{id}/snapshots/diff", s.handleSnapshotDiff)
+41 -4
View File
@@ -5,6 +5,10 @@ import (
"time"
"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
@@ -26,6 +30,7 @@ type listSnapshotsResponse struct {
HostID string `json:"host_id"`
Count int `json:"count"`
RefreshedAt *time.Time `json:"refreshed_at,omitempty"`
Stale bool `json:"stale"`
Snapshots []snapshotView `json:"snapshots"`
}
@@ -45,7 +50,8 @@ func (s *Server) handleListHostSnapshots(w stdhttp.ResponseWriter, r *stdhttp.Re
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", "")
return
}
@@ -61,10 +67,13 @@ func (s *Server) handleListHostSnapshots(w stdhttp.ResponseWriter, r *stdhttp.Re
Count: len(snaps),
Snapshots: make([]snapshotView, len(snaps)),
}
if len(snaps) > 0 {
t := snaps[0].RefreshedAt
out.RefreshedAt = &t
out.RefreshedAt = host.SnapshotRefreshedAt
mutationAt, err := s.deps.Store.LatestSuccessfulRepoMutation(r.Context(), hostID)
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 {
out.Snapshots[i] = snapshotView{
ID: sn.ID,
@@ -80,3 +89,31 @@ func (s *Server) handleListHostSnapshots(w stdhttp.ResponseWriter, r *stdhttp.Re
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"})
}
+83
View File
@@ -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)
}
}
+9
View File
@@ -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...)
+4
View File
@@ -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`,
+1
View File
@@ -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
+65
View File
@@ -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)
+11 -7
View File
@@ -42,9 +42,9 @@ func (st *Store) CreateFleetUpdate(ctx context.Context, fu FleetUpdate, hostIDs
}
if _, err := tx.ExecContext(ctx,
`INSERT INTO fleet_updates (id, started_at, started_by_user_id, target_version, status)
VALUES (?, ?, ?, ?, ?)`,
fu.ID, fu.StartedAt.UTC().Format(time.RFC3339Nano), fu.StartedByUserID, fu.TargetVersion, fu.Status,
`INSERT INTO fleet_updates (id, started_at, started_by_user_id, target_version, status, pause_after_first)
VALUES (?, ?, ?, ?, ?, ?)`,
fu.ID, fu.StartedAt.UTC().Format(time.RFC3339Nano), fu.StartedByUserID, fu.TargetVersion, fu.Status, fu.PauseAfterFirst,
); err != nil {
return fmt.Errorf("store: insert fleet_updates: %w", err)
}
@@ -67,12 +67,13 @@ func (st *Store) ActiveFleetUpdate(ctx context.Context) (*FleetUpdate, error) {
var current sql.NullString
var halted sql.NullString
var completedAt sql.NullString
var pauseAfterFirst int
err := st.db.QueryRowContext(ctx,
`SELECT id, started_at, started_by_user_id, target_version, status,
current_host_id, halted_reason, completed_at
current_host_id, halted_reason, completed_at, pause_after_first
FROM fleet_updates WHERE status = 'running' LIMIT 1`).
Scan(&fu.ID, &startedAt, &fu.StartedByUserID, &fu.TargetVersion, &fu.Status,
&current, &halted, &completedAt)
&current, &halted, &completedAt, &pauseAfterFirst)
if errors.Is(err, sql.ErrNoRows) {
return nil, nil
}
@@ -82,6 +83,7 @@ func (st *Store) ActiveFleetUpdate(ctx context.Context) (*FleetUpdate, error) {
fu.StartedAt, _ = time.Parse(time.RFC3339Nano, startedAt)
fu.CurrentHostID = current.String
fu.HaltedReason = halted.String
fu.PauseAfterFirst = pauseAfterFirst != 0
if completedAt.Valid {
t, _ := time.Parse(time.RFC3339Nano, completedAt.String)
fu.CompletedAt = &t
@@ -97,12 +99,13 @@ func (st *Store) GetFleetUpdate(ctx context.Context, id string) (*FleetUpdate, [
var current sql.NullString
var halted sql.NullString
var completedAt sql.NullString
var pauseAfterFirst int
err := st.db.QueryRowContext(ctx,
`SELECT id, started_at, started_by_user_id, target_version, status,
current_host_id, halted_reason, completed_at
current_host_id, halted_reason, completed_at, pause_after_first
FROM fleet_updates WHERE id = ?`, id).
Scan(&fu.ID, &startedAt, &fu.StartedByUserID, &fu.TargetVersion, &fu.Status,
&current, &halted, &completedAt)
&current, &halted, &completedAt, &pauseAfterFirst)
if errors.Is(err, sql.ErrNoRows) {
return nil, nil, ErrNotFound
}
@@ -112,6 +115,7 @@ func (st *Store) GetFleetUpdate(ctx context.Context, id string) (*FleetUpdate, [
fu.StartedAt, _ = time.Parse(time.RFC3339Nano, startedAt)
fu.CurrentHostID = current.String
fu.HaltedReason = halted.String
fu.PauseAfterFirst = pauseAfterFirst != 0
if completedAt.Valid {
t, _ := time.Parse(time.RFC3339Nano, completedAt.String)
fu.CompletedAt = &t
+12 -4
View File
@@ -44,7 +44,7 @@ func (s *Store) LookupHostByAgentToken(ctx context.Context, tokenHash string) (*
repo_size_bytes, snapshot_count, open_alert_count,
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
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 = ?`,
tokenHash)
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,
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
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)
return scanHost(row)
}
@@ -227,7 +227,7 @@ func (s *Store) ListHosts(ctx context.Context) ([]Host, error) {
repo_size_bytes, snapshot_count, open_alert_count,
applied_schedule_version, bandwidth_up_kbps, bandwidth_down_kbps,
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`)
if err != nil {
return nil, fmt.Errorf("store: list hosts: %w", err)
@@ -268,6 +268,7 @@ func scanHostRow(s hostScanner) (*Host, error) {
bwUp, bwDown sql.NullInt64
preHook, postHook sql.NullString
alwaysOn int
snapshotRefreshedAt sql.NullString
)
err := s.Scan(&h.ID, &h.Name, &h.OS, &h.Arch,
&h.AgentVersion, &h.ResticVersion, &h.ProtocolVersion,
@@ -276,7 +277,7 @@ func scanHostRow(s hostScanner) (*Host, error) {
&h.RepoSizeBytes, &h.SnapshotCount, &h.OpenAlertCount,
&h.AppliedScheduleVersion, &bwUp, &bwDown,
&preHook, &postHook,
&h.RepoStatus, &h.RepoStatusError, &alwaysOn)
&h.RepoStatus, &h.RepoStatusError, &alwaysOn, &snapshotRefreshedAt)
if err != nil {
if errors.Is(err, sql.ErrNoRows) {
return nil, ErrNotFound
@@ -332,6 +333,13 @@ func scanHostRow(s hostScanner) (*Host, error) {
h.PostHookDefault = postHook.String
}
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
}
+67
View File
@@ -27,6 +27,73 @@ 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
}
// 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
// when it actually starts work. ScheduledID is set when the job
// originates from a cron fire (actor_kind="schedule"); nil for
+42
View File
@@ -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)
@@ -0,0 +1 @@
ALTER TABLE hosts ADD COLUMN snapshot_refreshed_at TEXT;
@@ -0,0 +1 @@
ALTER TABLE fleet_updates ADD COLUMN pause_after_first INTEGER NOT NULL DEFAULT 0;
+2 -2
View File
@@ -69,8 +69,8 @@ func (s *Store) ReplaceHostSnapshots(ctx context.Context, hostID string, snaps [
}
if _, err := tx.ExecContext(ctx,
`UPDATE hosts SET snapshot_count = ? WHERE id = ?`,
len(snaps), hostID); err != nil {
`UPDATE hosts SET snapshot_count = ?, snapshot_refreshed_at = ? WHERE id = ?`,
len(snaps), when.UTC().Format(time.RFC3339Nano), hostID); err != nil {
return fmt.Errorf("store: update host snapshot_count: %w", err)
}
+5 -1
View File
@@ -138,7 +138,8 @@ func TestReplaceHostSnapshotsEmpty(t *testing.T) {
t.Fatalf("replace 1: %v", err)
}
// 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)
}
out, err := s.ListSnapshotsByHost(ctx, hostID)
@@ -152,4 +153,7 @@ func TestReplaceHostSnapshotsEmpty(t *testing.T) {
if h.SnapshotCount != 0 {
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)
}
}
+2
View File
@@ -77,6 +77,7 @@ type Host struct {
LastBackupStatus *string
RepoSizeBytes int64
SnapshotCount int
SnapshotRefreshedAt *time.Time
OpenAlertCount int
AppliedScheduleVersion int64
// Host-wide bandwidth caps applied to every restic invocation
@@ -230,6 +231,7 @@ type FleetUpdate struct {
CurrentHostID string
HaltedReason string
CompletedAt *time.Time
PauseAfterFirst bool
}
// FleetUpdateHost is one host's slot in a fleet update. Position is
+1 -1
View File
@@ -69,7 +69,7 @@
{{/* ---------- Hosts-behind hero tile (P6-18) ---------- */}}
{{if gt $page.UpdatesBehind 0}}
<div class="pt-4">
<a href="?updates=behind" class="hero-tile hero-tile--amber" style="display:inline-flex;">
<a href="/settings/fleet-update" class="hero-tile hero-tile--amber" style="display:inline-flex;">
<span class="hero-num">{{$page.UpdatesBehind}}</span>
<span class="hero-label">{{if eq $page.UpdatesBehind 1}}host behind{{else}}hosts behind{{end}} · review →</span>
</a>
+13 -3
View File
@@ -8,18 +8,18 @@
<div class="crumbs pt-6">
<a href="/">Dashboard</a><span class="sep">/</span>
<a href="/settings">Settings</a><span class="sep">/</span>
<span class="text-ink-mid">fleet update</span>
<span class="text-ink-mid">agent updates</span>
</div>
{{/* page header */}}
<div class="flex items-baseline justify-between mt-3.5">
<div>
<h1 class="text-[22px] font-medium tracking-[-0.005em]">
Fleet update
Agent updates
<span class="text-ink-fade font-normal text-[14px] ml-2 mono">target {{$page.TargetVersion}}</span>
</h1>
<p class="text-ink-mute text-[12px] mt-1 max-w-[760px] leading-[1.55]">
Rolling, sequential agent self-update. One host at a time, halts on first failure,
Rolling, sequential agent self-update. One host at a time, halts on first failure,
cancellable mid-roll. Only online hosts whose <span class="mono">agent_version</span>
differs from the server are eligible.
</p>
@@ -28,5 +28,15 @@
{{template "fleet_update_inner" .}}
<script>
function fleetBoxes(){return Array.from(document.querySelectorAll('.fleet-host:not(:disabled)'))}
function fleetReview(){const n=fleetBoxes().filter(x=>x.checked).length;document.getElementById('fleet-selected-count').textContent=n;document.getElementById('fleet-timeout').textContent=n+' × 95s';document.getElementById('fleet-update-start-btn').disabled=n===0}
function fleetFilter(){const q=document.getElementById('fleet-filter-name').value.toLowerCase(),v=document.getElementById('fleet-filter-version').value.toLowerCase(),s=document.getElementById('fleet-filter-state').value;document.querySelectorAll('.fleet-candidate').forEach(r=>r.style.display=r.dataset.name.toLowerCase().includes(q)&&r.dataset.version.toLowerCase().includes(v)&&(s==='all'||r.dataset.state===s)?'':'none')}
function fleetSelectVisible(on){document.querySelectorAll('.fleet-candidate').forEach(r=>{const x=r.querySelector('.fleet-host:not(:disabled)');if(x&&r.style.display!=='none')x.checked=on});fleetReview()}
function fleetSelectOnly(id){fleetBoxes().forEach(x=>x.checked=x.value===id);fleetReview()}
async function fleetStart(e){e.preventDefault();const ids=fleetBoxes().filter(x=>x.checked).map(x=>x.value),out=document.getElementById('fleet-start-error');out.textContent='';const res=await fetch('/api/fleet/update',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({host_ids:ids,canary_first:document.getElementById('fleet-canary').checked})});if(res.ok){location.reload();return false}const body=await res.json();out.textContent=body.message||body.code||'Unable to start';return false}
if(document.getElementById('fleet-selected-count'))fleetReview()
</script>
</div>
{{end}}
+8 -4
View File
@@ -82,10 +82,14 @@
<div class="text-[12px] text-ok mb-3 mono">✓ saved</div>
{{end}}
<p class="text-[12.5px] text-ink-mid leading-[1.6] mb-4 max-w-[640px]">
Only needed for rest-server repos that distinguish an append-only
user (everyday backups) from a delete-capable user (prune /
forget). For S3 / B2 / SFTP / local, leave this blank — the
everyday repo credentials handle prune too.
Required for prune. On rest-server repos this is the
delete-capable user, as distinct from the append-only user used
for everyday backups. Note that <strong>forget</strong> always
runs with the everyday repo credentials, so those must have
delete authority. For S3 / B2 / SFTP / local, enter the same
delete-capable repository credentials here if you want prune
enabled. <strong>Prune is skipped when admin credentials are
unset</strong>, on any backend.
</p>
<div class="grid grid-cols-2 gap-4">
<div>
+1
View File
@@ -40,6 +40,7 @@
{{if not $page.Form}}<span class="mono text-ink-fade text-[11px] ml-1">{{len $page.Channels}}</span>{{end}}
</a>
<a href="/settings/users" class="sub-tab {{if eq $page.ActiveTab "users"}}active{{end}}">Users</a>
<a href="/settings/fleet-update" class="sub-tab">Agent updates</a>
<span class="sub-tab text-ink-fade cursor-default" title="lands later">Authentication</span>
</div>
+40 -31
View File
@@ -7,11 +7,11 @@
*/}}
{{define "fleet_update_inner"}}
{{$page := .Page}}
<div id="fleet-update-panel" class="mt-5"
<div id="fleet-update-panel" class="mt-5"{{if and $page.Active (eq $page.Active.Status "running")}}
hx-get="{{$page.PollURL}}"
hx-trigger="every 3s [document.visibilityState==='visible']"
hx-select="#fleet-update-panel"
hx-swap="outerHTML">
hx-swap="outerHTML"{{end}}>
{{if and $page.Active (eq $page.Active.Status "running")}}
@@ -32,6 +32,7 @@
</div>
<div class="text-[11.5px] text-ink-mute mt-1">
target <span class="mono text-ink-mid">{{$page.Active.TargetVersion}}</span>
{{if $page.Active.PauseAfterFirst}} · canary-first{{end}}
· started <span class="mono text-ink-mid">{{relTime $page.Active.StartedAt}}</span>
{{if $page.Active.CurrentHostID}}
· waiting on <span class="mono text-ink-mid">{{index $page.HostNames $page.Active.CurrentHostID}}</span>
@@ -68,6 +69,12 @@
{{if $page.Active.HaltedReason}}
<div class="text-[12px] text-bad mt-2">{{$page.Active.HaltedReason}}</div>
{{end}}
{{if eq $page.Active.Status "halted"}}
<div class="mt-3 flex gap-2 text-[12px]">
<a class="btn" href="#fleet-update-start-form">Resume remaining as new rollout</a>
{{range $page.ActiveRows}}{{if eq .Status "failed"}}<button class="btn" type="button" onclick="fleetSelectOnly('{{.HostID}}');document.getElementById('fleet-update-start-form').scrollIntoView()">Retry {{.HostName}}</button>{{end}}{{end}}
</div>
{{end}}
</div>
{{template "fleet_update_rows" $page}}
@@ -122,11 +129,11 @@
{{define "fleet_update_idle_panel"}}
{{$page := .}}
<div class="panel rounded-[7px] px-5 py-4">
{{if eq (len $page.OutOfDateHosts) 0}}
{{if eq (len $page.Candidates) 0}}
<div class="flex items-center gap-3">
<span class="dot dot-online"></span>
<div>
<div class="text-ink text-[14px] font-medium">All hosts are up to date.</div>
<div class="text-ink text-[14px] font-medium">No hosts enrolled.</div>
<div class="text-ink-mute text-[12px] mt-0.5">
Every online agent matches server version <span class="mono">{{$page.TargetVersion}}</span>.
</div>
@@ -137,34 +144,36 @@
<h2 class="text-[14px] font-medium">{{len $page.OutOfDateHosts}} host{{if ne (len $page.OutOfDateHosts) 1}}s{{end}} out of date</h2>
<span class="mono text-[11px] text-ink-fade">target {{$page.TargetVersion}}</span>
</div>
<ul class="mt-3 space-y-1 text-[12px]">
{{range $page.OutOfDateHosts}}
<li class="flex items-center gap-3">
<span class="dot dot-online"></span>
<span class="mono text-ink">{{.Name}}</span>
<span class="mono text-ink-mute">{{if .AgentVersion}}{{.AgentVersion}}{{else}}—{{end}} → {{$page.TargetVersion}}</span>
</li>
{{end}}
</ul>
<form id="fleet-update-start-form" class="mt-4 flex items-center gap-3"
hx-post="/api/fleet/update"
hx-headers='{"Content-Type":"application/json"}'
hx-vals='{}'
hx-swap="none"
hx-on::after-request="if(event.detail.successful) location.reload()">
<label class="text-[11.5px] text-ink-mute">
Type the count
<span class="mono text-ink-mid">({{len $page.OutOfDateHosts}})</span>
to enable Start:
<div class="mt-3 flex flex-wrap gap-2">
<input id="fleet-filter-name" class="field text-[12px]" placeholder="Filter name or tag" oninput="fleetFilter()">
<input id="fleet-filter-version" class="field text-[12px] mono" placeholder="Version" oninput="fleetFilter()">
<select id="fleet-filter-state" class="field text-[12px]" onchange="fleetFilter()">
<option value="all">All states</option><option value="eligible">Eligible</option><option value="excluded">Excluded</option>
</select>
<button type="button" class="btn" onclick="fleetSelectVisible(true)">Select visible eligible</button>
<button type="button" class="btn" onclick="fleetSelectVisible(false)">Clear selection</button>
</div>
<div class="panel mt-3 rounded-[7px] overflow-hidden">
{{range $page.Candidates}}
<label class="fleet-candidate grid items-center px-3 py-2 hairline text-[12px]"
data-name="{{.Host.Name}} {{range .Host.Tags}}{{.}} {{end}}" data-version="{{.Host.AgentVersion}}" data-state="{{if .Eligible}}eligible{{else}}excluded{{end}}"
style="grid-template-columns: 28px 1.4fr .8fr .8fr 1.2fr;gap:12px">
<input class="fleet-host" type="checkbox" value="{{.Host.ID}}" {{if .Eligible}}checked onchange="fleetReview()"{{else}}disabled{{end}}>
<span class="mono">{{.Host.Name}}</span>
<span class="mono text-ink-mute">{{if .Host.AgentVersion}}{{.Host.AgentVersion}}{{else}}unknown{{end}}</span>
<span class="mono text-ink-mute">{{$page.TargetVersion}}</span>
<span class="text-ink-mute">{{if .Eligible}}eligible{{else}}{{.Reason}}{{end}}</span>
</label>
<input type="text" id="fleet-update-confirm" class="field mono text-[12.5px]"
style="width: 80px; padding: 5px 8px;"
oninput="document.getElementById('fleet-update-start-btn').disabled = (this.value !== '{{len $page.OutOfDateHosts}}');"
autocomplete="off" />
<button type="submit" id="fleet-update-start-btn" class="btn btn-amber" disabled>
Start fleet update
</button>
{{end}}
</div>
<div class="mt-4 text-[12px] text-ink-mute">
<span id="fleet-selected-count">{{len $page.OutOfDateHosts}}</span> selected · sequential, halts on first failure · worst-case
<span id="fleet-timeout" class="mono">{{len $page.OutOfDateHosts}} × 95s</span>
</div>
<form class="mt-3 flex items-center gap-3" onsubmit="return fleetStart(event)">
<label class="text-[11.5px] text-ink-mute"><input id="fleet-canary" type="checkbox" checked> Pause after first host for canary review</label>
<button id="fleet-update-start-btn" class="btn btn-amber">Start agent update</button>
<span id="fleet-start-error" class="text-bad text-[12px]"></span>
</form>
{{end}}
</div>