Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
146 changes: 146 additions & 0 deletions coordinator/internal/passctrl/parkedretry_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,146 @@
package passctrl

import (
"testing"
"time"

"drsync/coordinator/internal/model"
)

// TestAdvanceAutoRetriesParkedShardsOnce: a pass that would otherwise block
// indefinitely on parked work now gets one automatic retry round first —
// many park causes (a mount blip, a brief NFS hiccup) are transient at the
// moment a shard exhausts MaxShardAttempts, and by the time the rest of a
// large pass has drained, conditions may well have changed.
// RetryParkedByJob resets the attempt counter, so this call must requeue the
// shard (not silently no-op), and advance() must not treat that requeue as a
// completed phase this same tick — the shard needs a later tick to actually
// be granted and drain.
func TestAdvanceAutoRetriesParkedShardsOnce(t *testing.T) {
c := newController(t)
job := makeJob(t, c, []byte(baseSpec))
pass, err := c.st.CreatePass(job.ID, 1, model.PassDirfix)
if err != nil {
t.Fatal(err)
}
parkOneShard(t, c, pass, "a/parked-dir")

if err := c.advance(job); err != nil {
t.Fatal(err)
}

counts, err := c.st.ShardStateCounts(pass.ID)
if err != nil {
t.Fatal(err)
}
if counts[model.ShardParked] != 0 {
t.Fatalf("parked count after first advance() = %d, want 0 (auto-retry should have requeued it)",
counts[model.ShardParked])
}
if counts[model.ShardQueued] != 1 {
t.Fatalf("queued count after first advance() = %d, want 1 (the retried shard)",
counts[model.ShardQueued])
}
// Pass must still be DIRFIX: the requeued shard has not been granted or
// drained yet, so the phase cannot have transitioned.
pass, err = c.st.PassByNo(job.ID, 1)
if err != nil {
t.Fatal(err)
}
if pass.State != model.PassDirfix {
t.Fatalf("pass state = %s, want still DIRFIX (retried shard not yet drained)", pass.State)
}
}

// TestAdvanceBlocksOnSecondParkAfterAutoRetry: if the auto-retried shard
// parks again, the pass must fall back to the original block-and-alert
// behavior rather than retrying forever — an unbounded auto-retry loop on a
// genuinely stuck shard (e.g. a real permissions problem) would just be an
// infinite loop wearing a different name, and would silently mask a job that
// needs operator attention.
func TestAdvanceBlocksOnSecondParkAfterAutoRetry(t *testing.T) {
c := newController(t)
job := makeJob(t, c, []byte(baseSpec))
pass, err := c.st.CreatePass(job.ID, 1, model.PassDirfix)
if err != nil {
t.Fatal(err)
}
parkOneShard(t, c, pass, "a/parked-dir")

if err := c.advance(job); err != nil {
t.Fatal(err)
}
// Simulate the retried shard failing again: lease it and park it once more.
leased, err := c.st.LeaseShards("agent-1", 10, time.Minute)
if err != nil {
t.Fatal(err)
}
if len(leased) != 1 {
t.Fatalf("leased %d shards, want 1 (the retried shard)", len(leased))
}
if err := c.st.ParkShard(leased[0].ID, leased[0].LeaseID, "EIO again"); err != nil {
t.Fatal(err)
}

if err := c.advance(job); err != nil {
t.Fatal(err)
}
counts, err := c.st.ShardStateCounts(pass.ID)
if err != nil {
t.Fatal(err)
}
if counts[model.ShardParked] != 1 {
t.Fatalf("parked count after second advance() = %d, want 1 (must stay parked, "+
"not retried a second time)", counts[model.ShardParked])
}
if counts[model.ShardQueued] != 0 {
t.Fatalf("queued count after second advance() = %d, want 0 (no further auto-retry)",
counts[model.ShardQueued])
}
pass, err = c.st.PassByNo(job.ID, 1)
if err != nil {
t.Fatal(err)
}
if pass.State != model.PassDirfix {
t.Fatalf("pass state = %s, want still DIRFIX (blocked on the re-parked shard)", pass.State)
}
}

// TestAdvanceProceedsAfterAutoRetrySucceeds: the common case the feature
// exists for — the retried shard succeeds this time, and the phase advances
// normally once it drains, with no operator intervention needed at all.
func TestAdvanceProceedsAfterAutoRetrySucceeds(t *testing.T) {
c := newController(t)
job := makeJob(t, c, []byte(baseSpec))
pass, err := c.st.CreatePass(job.ID, 1, model.PassDirfix)
if err != nil {
t.Fatal(err)
}
parkOneShard(t, c, pass, "a/parked-dir")

if err := c.advance(job); err != nil {
t.Fatal(err)
}
leased, err := c.st.LeaseShards("agent-1", 10, time.Minute)
if err != nil {
t.Fatal(err)
}
if len(leased) != 1 {
t.Fatalf("leased %d shards, want 1 (the retried shard)", len(leased))
}
if err := c.st.CompleteShard(leased[0].ID, leased[0].LeaseID, 0, nil, nil); err != nil {
t.Fatal(err)
}

if err := c.advance(job); err != nil {
t.Fatal(err)
}
drainReaps(t, c)
pass, err = c.st.PassByNo(job.ID, 1)
if err != nil {
t.Fatal(err)
}
if pass.State != model.PassLinkfix {
t.Fatalf("pass state = %s, want LINKFIX (retried shard completed, phase should advance)", pass.State)
}
}
45 changes: 44 additions & 1 deletion coordinator/internal/passctrl/passctrl.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,20 @@ type Controller struct {
// firing their own email, then flushed together once the window elapses.
parkedRollup map[string]*parkedRollupState

// parkedAutoRetried tracks pass IDs whose parked backlog has already been
// given one automatic retry round (see advance()'s parked-shard gate) —
// so a persistently-failing shard that re-parks after the retry blocks
// the pass for operator attention as before, instead of looping forever.
// Guarded by parkedAlertMu (same lock as the other parked-tracking maps;
// a dedicated mutex would only ever be taken alongside it in practice).
// In-memory only, same "safe direction to be wrong in" reasoning as
// parkedAlerted: a coordinator restart re-attempts one retry round for a
// pass that was already retried before the restart — a redundant but
// harmless extra attempt, not a correctness issue. Forgotten once a pass
// leaves parked-blocked state (advances or the job ends), so a job that
// runs again later gets its own fresh retry round.
parkedAutoRetried map[int64]bool

// reapCh feeds runReapWorker — see its doc comment for why reaping runs on
// its own goroutine instead of inline in advance().
reapCh chan reapRequest
Expand Down Expand Up @@ -129,7 +143,8 @@ func (c *Controller) jobTerminal(jobID int64) {
func New(st *store.Store, journalRoot string) *Controller {
return &Controller{st: st, journalRoot: journalRoot,
parkedAlerted: map[int64]bool{}, parkedRollup: map[string]*parkedRollupState{},
reapCh: make(chan reapRequest, reapChanBuffer), reapInflight: map[reapKey]bool{}}
parkedAutoRetried: map[int64]bool{},
reapCh: make(chan reapRequest, reapChanBuffer), reapInflight: map[reapKey]bool{}}
}

// SetNotifier wires an email sender for pass/job completion notifications. A
Expand Down Expand Up @@ -442,6 +457,34 @@ func (c *Controller) advance(job *store.Job) error {
return nil // phase still draining
}
if parked := counts[model.ShardParked]; parked > 0 {
// Give the backlog one automatic retry round before blocking for an
// operator: a parked shard exhausted MaxShardAttempts (5) grants, but
// many park causes are transient at the point they're hit (a mount
// blip, a brief NFS hiccup) rather than deterministic — by the time
// the rest of the pass has drained (potentially hours later on a
// large tree), conditions may well have changed, and RetryParkedByJob
// resets the attempt counter, so a retried shard gets the same fresh
// 5-attempt budget as any new shard rather than one bonus try. Only
// one round: parkedAutoRetried (keyed by pass, not shard — shard IDs
// change across the retry) stops this from firing every tick forever
// on a genuinely stuck shard, which would just be an infinite retry
// loop wearing a different name. If it re-parks after the one round,
// fall through to the existing block-and-alert behavior below.
c.parkedAlertMu.Lock()
alreadyRetried := c.parkedAutoRetried[pass.ID]
if !alreadyRetried {
c.parkedAutoRetried[pass.ID] = true
}
c.parkedAlertMu.Unlock()
if !alreadyRetried {
n, err := c.st.RetryParkedByJob(job.Name)
if err != nil {
return fmt.Errorf("auto-retry parked shards: %w", err)
}
slog.Info("auto-retried parked shards at end of pass", "job", job.Name,
"pass", pass.PassNo, "retried", n)
return nil // re-queued work needs a later tick to grant/drain
}
// Do not advance past parked work silently; operator resolves via API.
slog.Warn("pass blocked on parked shards", "job", job.Name,
"pass", pass.PassNo, "parked", parked)
Expand Down
31 changes: 26 additions & 5 deletions docs/DESIGN-coordinator.md
Original file line number Diff line number Diff line change
Expand Up @@ -91,9 +91,10 @@ PENDING ──▶ PROBING ──all probes ok──▶ SCANNING ──all shards
synced into the underlying rootfs. A missing/misordered mount or a stub on any host is
thus caught before bulk work runs — not just on whichever agent grabbed the root shard.
A failed probe parks (like any shard),
and the parked-shard guard holds the pass until the operator fixes the mount and
retries. Probes pinned to an agent that departs after seeding are pruned so the phase
is not stalled. An empty fleet skips probing (nobody to probe or grant work to).
and the parked-shard guard holds the pass — after one automatic retry round, see
§2.3 — until the operator fixes the mount and retries. Probes pinned to an agent
that departs after seeding are pruned so the phase is not stalled. An empty fleet
skips probing (nobody to probe or grant work to).
- `SCANNING` is the long phase: walk, diff, and copy are interleaved *per shard*, so
data starts moving seconds after pass start; there is no global "scan first" barrier.
- `DIRFIX` applies directory metadata deepest-first from the journal's dir records
Expand Down Expand Up @@ -124,8 +125,8 @@ PENDING ──▶ PROBING ──all probes ok──▶ SCANNING ──all shards
```
QUEUED ──grant──▶ LEASED ──ShardResult ok──▶ DONE
│ │
lease expiry┘ └─ShardResult(err)──▶ PARKED ──(operator retry / auto after
transient-window)──▶ QUEUED
lease expiry┘ └─ShardResult(err)──▶ PARKED ──(operator retry, or one
automatic round at pass end)──▶ QUEUED
QUEUED (attempt++)
```

Expand All @@ -134,6 +135,26 @@ QUEUED ──grant──▶ LEASED ──ShardResult ok──▶ DONE
diagnosis breadcrumbs instead of poisoning the fleet forever.
- Shards created by `ShardSplit` enter `QUEUED` in the same transaction that records
the split against the parent (ordering invariant, protocol doc §4.2).
- **One automatic retry round at the end of a phase, then block for the operator.**
`advance()`'s parked-shard guard (§2.2) used to hold every phase transition open
indefinitely the instant any shard parked — correct, but on a long-running pass a
job could sit for hours waiting on a park caused by something transient at the
moment it happened (a mount blip, a brief NFS hiccup), not a shard that will fail
the same way forever. `advance()` now calls `store.RetryParkedByJob` once per pass
(`Controller.parkedAutoRetried`, keyed by pass id, in-memory — a coordinator restart
re-attempts one round for a pass already retried before the restart, a redundant
but harmless extra attempt, same "safe direction to be wrong in" reasoning as
`parkedAlerted` below) before falling back to the original block-and-alert
behavior. `RetryParkedByJob` resets `attempt` to 0, so a retried shard gets the
same fresh 5-attempt budget as any new shard, not one bonus try. Deliberately
bounded to exactly one round: an unbounded auto-retry on a genuinely stuck shard
(a real permissions problem, a permanently unreachable mount) would just be an
infinite retry loop wearing a different name, and would silently mask a job that
actually needs operator attention. `checkParkedShards`' alerting (`passctrl.go`,
the tick-driven parked-shard email digest) composes with this unchanged — a shard
that leaves and re-enters PARKED state naturally clears and re-sets its
`parkedAlerted` entry, so the operator is still alerted, just once the automatic
round has already been given a chance to make the alert unnecessary.

## 3. Schema (SQLite)

Expand Down
46 changes: 39 additions & 7 deletions webui/console.html
Original file line number Diff line number Diff line change
Expand Up @@ -619,13 +619,17 @@ <h2>Jobs</h2><span class="eyebrow">migration lifecycle</span>
</section>

<section class="panel">
<div class="head"><h2>Aggregate throughput</h2><span class="eyebrow">bytes copied · fleet-wide · 90s</span></div>
<div class="head"><h2 id="tl-title">Aggregate throughput</h2><span class="eyebrow" id="tl-eyebrow">bytes copied · fleet-wide · 90s</span>
<span class="chips" id="tl-mode" style="margin-left:auto" role="tablist" aria-label="Timeline metric">
<button class="chip on" data-tl="bw" role="tab" aria-selected="true">throughput</button>
<button class="chip" data-tl="scan" role="tab" aria-selected="false">scan rate</button>
</span></div>
<div class="tl-wrap">
<div class="tl-top">
<div class="tl-big"><span id="tl-now">–</span> <span class="u"></span></div>
<div class="tl-leg"><span><i></i> <span id="tl-leg-u">throughput</span></span><span id="tl-peak">peak –</span></div>
</div>
<canvas id="tl" height="150" aria-label="Aggregate throughput timeline"></canvas>
<canvas id="tl" height="150" aria-label="Aggregate timeline"></canvas>
</div>
</section>

Expand Down Expand Up @@ -857,6 +861,15 @@ <h3 id="jmodal-title">New job</h3>
while (v >= 1000 && i < u.length - 1) { v /= 1000; i++; }
return [i === 0 ? String(Math.round(v)) : v.toFixed(v < 10 ? 2 : 1), u[i]];
}
// scanRate: entries/s → [value, unit] — the same [value, unit] shape as
// rate() so drawTL can treat both metrics identically, just with SI counts
// (1000-based, matching countFmt) instead of a byte/bit rate.
function scanRate(eps) {
eps = +eps || 0;
const u = ["/s","K/s","M/s","B/s"]; let i = 0;
while (eps >= 1000 && i < u.length - 1) { eps /= 1000; i++; }
return [i === 0 ? String(Math.round(eps)) : eps.toFixed(eps < 10 ? 2 : 1), u[i]];
}
function countFmt(n) {
n = +n || 0;
if (n >= 1e9) return (n / 1e9).toFixed(2) + "B";
Expand Down Expand Up @@ -1718,8 +1731,25 @@ <h3 id="jmodal-title">New job</h3>
: `<div class="jmeta" style="padding:6px 4px">no parked shards</div>`;
}

// ---------- throughput timeline (auto-scaling canvas) ----------
// ---------- aggregate timeline (auto-scaling canvas) ----------
// "bw" (default) plots bwHist through rate() (throughput, unit-mode aware);
// "scan" plots scanHist through scanRate() (entries/s). Same canvas, same
// draw routine — just a different history array and formatter, persisted
// like unitsMode so a reload keeps the operator's last choice.
let tlMode = localStorage.getItem("drsync.tlmode") === "scan" ? "scan" : "bw";
function applyTLMode(m) {
tlMode = m;
document.querySelectorAll("#tl-mode button").forEach(b => { const on = b.dataset.tl === m; b.classList.toggle("on", on); b.setAttribute("aria-selected", on); });
$("#tl-title").textContent = m === "scan" ? "Aggregate scan rate" : "Aggregate throughput";
$("#tl-eyebrow").textContent = m === "scan" ? "entries scanned · fleet-wide · 90s" : "bytes copied · fleet-wide · 90s";
cv.setAttribute("aria-label", m === "scan" ? "Aggregate scan rate timeline" : "Aggregate throughput timeline");
}
$("#tl-mode").addEventListener("click", e => {
const b = e.target.closest("button[data-tl]"); if (!b) return;
applyTLMode(b.dataset.tl); localStorage.setItem("drsync.tlmode", tlMode); drawTL();
});
const cv = $("#tl"), ctx = cv.getContext("2d");
applyTLMode(tlMode);
function sizeCanvas() {
const dpr = Math.min(devicePixelRatio || 1, 2), w = cv.clientWidth, h = 150;
cv.width = w * dpr; cv.height = h * dpr; ctx.setTransform(dpr, 0, 0, dpr, 0, 0);
Expand All @@ -1728,14 +1758,16 @@ <h3 id="jmodal-title">New job</h3>
const niceTop = v => { if (v <= 0) return 1; const p = Math.pow(10, Math.floor(Math.log10(v))), n = v / p;
return (n <= 1 ? 1 : n <= 2 ? 2 : n <= 5 ? 5 : 10) * p * 1.05; };
function drawTL() {
const fmt = tlMode === "scan" ? scanRate : rate;
const src = tlMode === "scan" ? scanHist : bwHist;
const { w, h } = sizeCanvas(); ctx.clearRect(0, 0, w, h);
const pad = { t:12, r:6, b:16, l:6 }, gw = w - pad.l - pad.r, gh = h - pad.t - pad.b;
const data = bwHist.length ? bwHist : [0], n = data.length, hi = niceTop(Math.max(...data, 1));
const data = src.length ? src : [0], n = data.length, hi = niceTop(Math.max(...data, 1));
const x = i => n < 2 ? pad.l + gw : pad.l + i / (n - 1) * gw;
const y = v => pad.t + (1 - v / hi) * gh;
ctx.strokeStyle = css("--grid"); ctx.lineWidth = 1; ctx.font = "10px ui-monospace, monospace"; ctx.fillStyle = css("--muted");
[0, hi / 2, hi].forEach(g => { ctx.beginPath(); ctx.moveTo(pad.l, y(g)); ctx.lineTo(w - pad.r, y(g)); ctx.stroke();
const [rv, ru] = rate(g); ctx.fillText(rv + " " + ru, pad.l + 2, y(g) - 3); });
const [rv, ru] = fmt(g); ctx.fillText(rv + " " + ru, pad.l + 2, y(g) - 3); });
const acc = css("--accent"), acc2 = css("--accent-2");
const grad = ctx.createLinearGradient(0, pad.t, 0, h - pad.b);
grad.addColorStop(0, hexA(acc, .34)); grad.addColorStop(1, hexA(acc, .02));
Expand All @@ -1746,10 +1778,10 @@ <h3 id="jmodal-title">New job</h3>
const ex = x(n - 1), ey = y(data[n - 1]);
ctx.beginPath(); ctx.arc(ex, ey, 6, 0, 7); ctx.fillStyle = hexA(acc2, .22); ctx.fill();
ctx.beginPath(); ctx.arc(ex, ey, 3.2, 0, 7); ctx.fillStyle = acc2; ctx.fill();
const [cvv, cu] = rate(data[n - 1]); $("#tl-now").textContent = cvv;
const [cvv, cu] = fmt(data[n - 1]); $("#tl-now").textContent = cvv;
const uEl = document.querySelector(".tl-big .u"); if (uEl) uEl.textContent = cu;
$("#tl-leg-u").textContent = cu;
const [pv, pu] = rate(Math.max(...data)); $("#tl-peak").textContent = "peak " + pv + " " + pu;
const [pv, pu] = fmt(Math.max(...data)); $("#tl-peak").textContent = "peak " + pv + " " + pu;
}
function hexA(hex, a) { hex = hex.replace("#", ""); if (hex.length === 3) hex = hex.split("").map(c => c + c).join("");
const n = parseInt(hex, 16); return `rgba(${(n >> 16) & 255},${(n >> 8) & 255},${n & 255},${a})`; }
Expand Down
27 changes: 27 additions & 0 deletions webui/test/console.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,33 @@ test("a closely-spaced poll does not blank the fleet's rates", async () => {
assert.ok(/Mbps|Gbps|Kbps|MiB|GiB|KiB/.test(fleet), "fleet throughput column blanked");
});

test("aggregate timeline toggles between throughput and scan rate", async () => {
// Default view is throughput: the scan-rate tab is present but inactive.
const bwTab = c.$('#tl-mode button[data-tl="bw"]');
const scanTab = c.$('#tl-mode button[data-tl="scan"]');
assert.ok(bwTab.classList.contains("on"), "throughput tab not selected by default");
assert.ok(!scanTab.classList.contains("on"), "scan-rate tab selected before being clicked");
assert.match(c.text("#tl-title"), /throughput/i);

const bwNow = c.text("#tl-now");
assert.notEqual(bwNow, "–", "timeline never populated in throughput mode");

scanTab.click();
await c.tick(50);

assert.ok(scanTab.classList.contains("on"), "scan-rate tab did not activate on click");
assert.ok(!bwTab.classList.contains("on"), "throughput tab still marked active after switching");
assert.match(c.text("#tl-title"), /scan rate/i);
const scanNow = c.text("#tl-now");
assert.notEqual(scanNow, "–", "timeline blanked after switching to scan rate");

// Switching back restores the throughput reading, proving the two
// histories are genuinely independent series, not one relabeled.
bwTab.click();
await c.tick(50);
assert.equal(c.text("#tl-now"), bwNow, "throughput reading changed after switching back");
});

// --------------------------------------------------------------------------
// Per-agent in-flight work
// --------------------------------------------------------------------------
Expand Down