14 Commits

Author SHA1 Message Date
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
15 changed files with 330 additions and 53 deletions
+24 -1
View File
@@ -6,6 +6,23 @@ and the project follows [Semantic Versioning](https://semver.org/).
## [Unreleased] ## [Unreleased]
## [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 ## [1.1.0] - 2026-06-15
### Added ### Added
@@ -123,5 +140,11 @@ with a web UI, JSON API, and self-updating agent fleet.
go vet, golangci-lint). go vet, golangci-lint).
- Threat model published (`docs/threat-model.md`). - 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.1.1...HEAD
[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 [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
+37 -11
View File
@@ -2,7 +2,8 @@
Thanks for your interest in restic-manager. This document covers how Thanks for your interest in restic-manager. This document covers how
to set up a development environment, the conventions the project 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 ## Project status and scope
@@ -108,6 +109,32 @@ admin user.
## Workflow ## 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 ### Before opening a PR
1. **Open an issue first** for non-trivial changes. The design is 1. **Open an issue first** for non-trivial changes. The design is
@@ -136,25 +163,24 @@ The PR template asks for:
### Reporting bugs ### Reporting bugs
Open an issue with: Use the bug report issue template and include:
- restic-manager version (`server --version`) and agent version. - restic-manager version (`server --version`) and agent version.
- restic version on the affected host. - restic version on the affected host.
- Steps to reproduce. - Steps to reproduce.
- Server and agent logs (sanitise any tokens before pasting). - Server and agent logs (sanitise any tokens before pasting).
Security-sensitive bugs go through the [SECURITY.md](./SECURITY.md) For security-sensitive bugs, use the private disclosure process noted
disclosure path instead — please don't open a public issue for above.
them.
### Suggesting features ### Suggesting features
Open an issue describing the use case (not just the proposed Use the feature request issue template and describe the use case (not
solution). The roadmap in `tasks.md` shows where the project is just the proposed solution). The roadmap in `tasks.md` shows where the
heading; if the suggestion fits a future phase we'll wire it in project is heading; if the suggestion fits a future phase we'll wire it
there. If it falls outside the project's scope (multi-tenancy, SaaS, 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 non-restic backends — see `spec.md` §2 non-goals) we'll say so early to
early to save your time. save your time.
## Code of conduct ## Code of conduct
+7 -2
View File
@@ -601,6 +601,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)") 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)") 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)) groups := make([]restic.ForgetGroup, 0, len(p.ForgetGroups))
for _, g := range p.ForgetGroups { for _, g := range p.ForgetGroups {
groups = append(groups, restic.ForgetGroup{ groups = append(groups, restic.ForgetGroup{
@@ -615,9 +620,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 { 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: case api.JobPrune:
// Prune may require admin creds (delete authority on rest-server). // Prune may require admin creds (delete authority on rest-server).
+2 -2
View File
@@ -274,13 +274,13 @@ func (r *Runner) RunInit(ctx context.Context, jobID string) error {
// snapshot projection (forget rewrites the snapshot index — the // snapshot projection (forget rewrites the snapshot index — the
// host's snapshot list shrinks). Snapshot refresh runs once after // host's snapshot list shrinks). Snapshot refresh runs once after
// every group completes, not per-group. // 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() startedAt := time.Now().UTC()
r.sendStarted(jobID, api.JobForget, startedAt) r.sendStarted(jobID, api.JobForget, startedAt)
env := r.resticEnv() env := r.resticEnv()
var seq atomic.Int64 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() finishedAt := time.Now().UTC()
r.sendFinished(ctx, jobID, finishedAt, err, nil) r.sendFinished(ctx, jobID, finishedAt, err, nil)
+1 -1
View File
@@ -398,7 +398,7 @@ esac
Tag: "documents", Tag: "documents",
Policy: restic.ForgetPolicy{KeepLast: &keepLast}, 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) t.Fatalf("RunForget: %v", err)
} }
_ = firstEnvOfType(t, tx.envs, api.MsgJobStarted) _ = firstEnvOfType(t, tx.envs, api.MsgJobStarted)
+4 -5
View File
@@ -103,15 +103,14 @@ func connectOnce(ctx context.Context, cfg Config, handle Handler) error {
} }
dialCtx, cancel := context.WithTimeout(ctx, 30*time.Second) 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() cancel()
if err != nil { if err != nil {
return fmt.Errorf("dial: %w", err) return fmt.Errorf("dial: %w", err)
} }
// websocket.Dial returns the upgrade response separately from the // On a successful upgrade coder/websocket transfers ownership of the
// conn. Body is empty on a successful upgrade but Go's net/http // response stream to conn and deliberately sets res.Body to nil. Closing
// still expects it closed to release the connection. // the connection below releases that stream.
defer func() { _ = res.Body.Close() }()
defer conn.CloseNow() //nolint:errcheck defer conn.CloseNow() //nolint:errcheck
// Send hello. // Send hello.
+46
View File
@@ -0,0 +1,46 @@
package wsclient
import (
"context"
"net/http"
"net/http/httptest"
"testing"
"time"
"github.com/coder/websocket"
)
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)
}
}
+4 -1
View File
@@ -322,7 +322,7 @@ type ForgetGroup struct {
// any keep-* would delete every snapshot in the tagged set). // any keep-* would delete every snapshot in the tagged set).
// Returns the first error encountered, or nil when every group runs // Returns the first error encountered, or nil when every group runs
// to a clean exit. // 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 { if len(groups) == 0 {
return fmt.Errorf("restic forget: refusing to run with no groups (would be a no-op)") 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 := []string{"forget", "--json", "--tag", g.Tag}
args = append(args, g.Policy.args()...) args = append(args, g.Policy.args()...)
if dryRun {
args = append(args, "--dry-run")
}
cmd := e.resticCmd(ctx, args...) cmd := e.resticCmd(ctx, args...)
if err := runWithPump(cmd, handle); err != nil { if err := runWithPump(cmd, handle); err != nil {
return err return err
+20
View File
@@ -60,6 +60,26 @@ func TestRunPruneInvokesPrune(t *testing.T) {
t.Fatalf("expected 'prune' in captured output; got: %v", *lines) 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 --- // --- B2: RunCheck ---
func TestRunCheckLockSniff(t *testing.T) { func TestRunCheckLockSniff(t *testing.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, func (s *Server) dispatchJob(ctx context.Context, user *store.User,
hostID string, kind api.JobKind, args []string, hostID string, kind api.JobKind, args []string,
) (res runNowResponse, status int, code, msg string) { ) (res runNowResponse, status int, code, msg string) {
return s.dispatchJobWithPayload(ctx, user, hostID, kind, nil, api.CommandRunPayload{ payload := api.CommandRunPayload{
Kind: kind, Kind: kind,
Args: args, 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 // 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 { switch d.Kind {
case "forget": 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 { if !ok {
slog.Info("maintenance: forget skipped — no source groups with retention", slog.Info("maintenance: forget skipped — no source groups with retention",
"host_id", d.HostID) "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 // that has a non-empty retention policy and builds a CommandRunPayload
// with ForgetGroups populated. Returns ok=false if the host has no // with ForgetGroups populated. Returns ok=false if the host has no
// such groups (the dispatcher then skips this kind). // 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) groups, err := s.deps.Store.ListSourceGroupsByHost(ctx, hostID)
if err != nil { if err != nil {
slog.Warn("maintenance: list source groups failed", "host_id", hostID, "err", err) return api.CommandRunPayload{}, false, err
return api.CommandRunPayload{}, false
} }
fg := make([]api.ForgetGroup, 0, len(groups)) fg := make([]api.ForgetGroup, 0, len(groups))
for _, g := range groups { for _, g := range groups {
@@ -105,9 +109,9 @@ func (s *Server) buildForgetPayloadForHost(ctx context.Context, hostID string) (
}) })
} }
if len(fg) == 0 { 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 { func isEmptyRetention(p store.RetentionPolicy) bool {
+36 -17
View File
@@ -512,11 +512,27 @@ func TestDrainPendingSerializesPerHost(t *testing.T) {
// Connect the agent so DrainPending can dispatch. // Connect the agent so DrainPending can dispatch.
c := agentDial(t, srv, ts, hostID, token) c := agentDial(t, srv, ts, hostID, token)
sendHello(t, c, "serialise-host") sendHello(t, c, "serialise-host")
// Drain the on-hello goroutine's pass first (no pending rows yet), // Wait for the on-hello push to settle.
// then wait for the schedule.set so the connection is fully settled.
_ = drainUntil(t, c, api.MsgScheduleSet) _ = 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() now := time.Now().UTC()
for i := range 5 { for i := range 5 {
pid := ulid.Make().String() 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 var wg sync.WaitGroup
for range 10 { for range 10 {
wg.Add(1) wg.Add(1)
@@ -544,24 +561,26 @@ func TestDrainPendingSerializesPerHost(t *testing.T) {
} }
wg.Wait() wg.Wait()
// Drain any envelopes the agent received so we don't block below. // Drain to completion. The fire-and-forget on-hello DrainPending
// We read with short timeouts and stop when the connection goes quiet. // shares the same per-host mutex and can hold it during the burst,
drainDeadline := time.Now().Add(500 * time.Millisecond) // leaving rows for a later pass — exactly how production drains
for time.Now().Before(drainDeadline) { // (repeatedly, via the 30s tick / on reconnect). Re-drain until the
ctx, cancel := context.WithTimeout(context.Background(), 100*time.Millisecond) // queue is empty; because every drain is still serialised, each row
_, _, err := c.Read(ctx) // is dispatched at most once, so the exactly-5 job count below proves
cancel() // there was no double-dispatch.
if err != nil { deadline := time.Now().Add(5 * time.Second)
break 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 { 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 var n int
_ = st.DB().QueryRow( _ = st.DB().QueryRow(
`SELECT COUNT(*) FROM jobs WHERE host_id = ? AND kind = 'backup' AND actor_kind = 'schedule'`, `SELECT COUNT(*) FROM jobs WHERE host_id = ? AND kind = 'backup' AND actor_kind = 'schedule'`,
File diff suppressed because one or more lines are too long
+8 -4
View File
@@ -82,10 +82,14 @@
<div class="text-[12px] text-ok mb-3 mono">✓ saved</div> <div class="text-[12px] text-ok mb-3 mono">✓ saved</div>
{{end}} {{end}}
<p class="text-[12.5px] text-ink-mid leading-[1.6] mb-4 max-w-[640px]"> <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 Required for prune. On rest-server repos this is the
user (everyday backups) from a delete-capable user (prune / delete-capable user, as distinct from the append-only user used
forget). For S3 / B2 / SFTP / local, leave this blank — the for everyday backups. Note that <strong>forget</strong> always
everyday repo credentials handle prune too. 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> </p>
<div class="grid grid-cols-2 gap-4"> <div class="grid grid-cols-2 gap-4">
<div> <div>