P3-03: restic restore + diff execution path
Wires JobRestore and JobDiff end-to-end at the agent layer (the wizard
backend that drives this lands in the next slice).
- internal/api: JobRestore + JobDiff JobKind constants. CommandRunPayload
grows nullable Restore + Diff sub-payloads. RestorePayload carries
snapshot_id, paths, in_place, target_dir; DiffPayload carries
snapshot_a + snapshot_b.
- internal/restic.RunRestore wraps 'restic restore <sid> --target ...
[--no-ownership] [--include p]...' with --json. New pumpRestoreStdout
parses the per-line status / summary objects (drops raw status from
log.stream — the throttled job.progress envelope covers it). New
RestoreStatus + RestoreSummary types mirror restic's wire shape.
- internal/restic.RunDiff wraps 'restic diff --json <a> <b>'.
- internal/agent/runner: RunRestore translates RestoreStatus into
job.progress (mapping FilesRestored → FilesDone etc) with a small
estimateETA helper since restic doesn't provide ETA for restore.
RunDiff is a thin streamHandler wrapper.
- cmd/agent dispatcher gains JobRestore + JobDiff cases. Both reuse
the spawn() helper from P3-X1 so cancel just works.
- Drive-by fix: lastProgress was initialised to time.Now() so the
very first status event was suppressed by the 1s throttle if the
agent reported quickly. Initialise to time.Time{} (zero) so the
first event always emits. Affects backup + restore.
Tests:
- restore_test covers restore happy path (started → progress →
finished, kind=restore on the started envelope), in-place argv
asserts no --no-ownership, new-dir argv asserts --no-ownership +
--target + --include, diff produces the expected log.stream lines.
Restage block (CLAUDE.md) is deferred to the end of the restore
sub-phase so we restage once with all changes.
This commit is contained in:
@@ -156,7 +156,7 @@ func (r *Runner) RunBackup(ctx context.Context, jobID string, paths, excludes, t
|
||||
}
|
||||
|
||||
env := r.resticEnv()
|
||||
lastProgress := time.Now()
|
||||
lastProgress := time.Time{} // zero time → first status event always emits
|
||||
|
||||
handle := func(stream string, line string, ev any) {
|
||||
// Throttled progress events come from restic's `status` JSON.
|
||||
@@ -359,6 +359,102 @@ func (r *Runner) RunCheck(ctx context.Context, jobID string, subsetPct int) erro
|
||||
return nil
|
||||
}
|
||||
|
||||
// RunRestore executes a restic restore job and reports back via the
|
||||
// sender. paths is the operator-selected file/dir list to restore.
|
||||
// inPlace=true preserves uid/gid/mode and writes at "/"; inPlace=false
|
||||
// writes at targetDir with --no-ownership.
|
||||
//
|
||||
// Status events from restic are throttled into job.progress in the
|
||||
// same shape as backup; raw status lines are dropped from log.stream
|
||||
// (they would drown the log on a fast restore — the progress widget
|
||||
// already covers them).
|
||||
func (r *Runner) RunRestore(ctx context.Context, jobID, snapshotID string, paths []string, inPlace bool, targetDir string) error {
|
||||
startedAt := time.Now().UTC()
|
||||
r.sendStarted(jobID, api.JobRestore, startedAt)
|
||||
|
||||
env := r.resticEnv()
|
||||
var seq atomic.Int64
|
||||
lastProgress := time.Time{} // zero time → first status event always emits
|
||||
|
||||
handle := func(stream string, line string, ev any) {
|
||||
status, isStatus := ev.(restic.RestoreStatus)
|
||||
if !isStatus {
|
||||
now := time.Now().UTC()
|
||||
logEnv, _ := api.Marshal(api.MsgLogStream, "", api.LogStreamLine{
|
||||
JobID: jobID,
|
||||
Seq: seq.Add(1),
|
||||
TS: now,
|
||||
Stream: api.LogStream(stream),
|
||||
Payload: line,
|
||||
})
|
||||
_ = r.tx.Send(logEnv)
|
||||
}
|
||||
if isStatus {
|
||||
if time.Since(lastProgress) < r.progressMinPeriod {
|
||||
return
|
||||
}
|
||||
lastProgress = time.Now()
|
||||
progEnv, _ := api.Marshal(api.MsgJobProgress, jobID, api.JobProgressPayload{
|
||||
JobID: jobID,
|
||||
PercentDone: status.PercentDone,
|
||||
FilesDone: status.FilesRestored,
|
||||
TotalFiles: status.TotalFiles,
|
||||
BytesDone: status.BytesRestored,
|
||||
TotalBytes: status.TotalBytes,
|
||||
ETASeconds: estimateETA(status.BytesRestored, status.TotalBytes, status.SecondsElapsed),
|
||||
ThroughputBps: throughput(status.BytesRestored, status.SecondsElapsed),
|
||||
})
|
||||
_ = r.tx.Send(progEnv)
|
||||
}
|
||||
}
|
||||
|
||||
summary, err := env.RunRestore(ctx, snapshotID, paths, inPlace, targetDir, handle)
|
||||
finishedAt := time.Now().UTC()
|
||||
|
||||
var statsBlob json.RawMessage
|
||||
if summary != nil {
|
||||
statsBlob, _ = json.Marshal(summary)
|
||||
}
|
||||
r.sendFinished(ctx, jobID, finishedAt, err, statsBlob)
|
||||
if err != nil {
|
||||
return fmt.Errorf("runner restore: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// estimateETA computes an ETA in seconds based on current bytes
|
||||
// progress + elapsed seconds. Restic restore's --json doesn't emit an
|
||||
// ETA field of its own (unlike backup), so we approximate by linear
|
||||
// extrapolation. Returns 0 when we don't have enough data.
|
||||
func estimateETA(bytesDone, totalBytes, secondsElapsed int64) int64 {
|
||||
if bytesDone <= 0 || totalBytes <= 0 || secondsElapsed <= 0 || bytesDone >= totalBytes {
|
||||
return 0
|
||||
}
|
||||
rate := float64(bytesDone) / float64(secondsElapsed)
|
||||
if rate <= 0 {
|
||||
return 0
|
||||
}
|
||||
return int64(float64(totalBytes-bytesDone) / rate)
|
||||
}
|
||||
|
||||
// RunDiff executes `restic diff --json <a> <b>` and forwards output
|
||||
// as log.stream lines. No snapshot-list refresh, no stats update —
|
||||
// diff is purely informational.
|
||||
func (r *Runner) RunDiff(ctx context.Context, jobID, snapshotA, snapshotB string) error {
|
||||
startedAt := time.Now().UTC()
|
||||
r.sendStarted(jobID, api.JobDiff, startedAt)
|
||||
|
||||
env := r.resticEnv()
|
||||
var seq atomic.Int64
|
||||
err := env.RunDiff(ctx, snapshotA, snapshotB, r.streamHandler(jobID, &seq))
|
||||
finishedAt := time.Now().UTC()
|
||||
r.sendFinished(ctx, jobID, finishedAt, err, nil)
|
||||
if err != nil {
|
||||
return fmt.Errorf("runner diff: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// RunUnlock executes a `restic unlock` job. On success it ships a
|
||||
// repo.stats envelope with LockPresent=false so the UI banner clears.
|
||||
func (r *Runner) RunUnlock(ctx context.Context, jobID string) error {
|
||||
|
||||
Reference in New Issue
Block a user