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
7 changes: 6 additions & 1 deletion docs/src/content/docs/user-guide/notifications.md
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,12 @@ pause uses none of their retries and is not reported as a delivery failure.
## Delivery retries and health

Scan changes, scan failures, cancellations, timeouts, stalled cycles, and
recovery events can all generate notifications. Definitive provider failures
recovery events can all generate notifications. A scan that stops because
EdgeWatch stopped, for example during an upgrade or restart, is recorded as
canceled with the reason "scan interrupted because EdgeWatch stopped" and
appears in Activity as **Scan interrupted**, but sends no notification. A
daemon that keeps stopping is still reported by the job's silence alert.
Definitive provider failures
are retried durably for up to 15 attempts over roughly 77 hours; the delay
doubles from two minutes and caps at 12 hours. A restart preserves each
delivery's retry schedule. Terminal failures are visible in the console
Expand Down
6 changes: 6 additions & 0 deletions docs/src/content/docs/user-guide/scanning.md
Original file line number Diff line number Diff line change
Expand Up @@ -73,6 +73,12 @@ next cycle. Accepting an incident, approving or resetting the baseline, or
changing the monitored scope discards paused progress, and the next run starts
a fresh cycle.

**Cancel scan** stops the scanner, and the job page and Overview show
**Cancellation requested** until the scan ends. Once EdgeWatch is saving a
scan's result, the scan can no longer be canceled. A scan that stops because
EdgeWatch itself stopped is recorded as interrupted rather than canceled; see
[notifications](/user-guide/notifications/#delivery-retries-and-health).

## Runtime capabilities

The default Compose configuration grants `NET_RAW`. Naabu SYN additionally
Expand Down
45 changes: 40 additions & 5 deletions internal/app/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -84,8 +84,30 @@ type activeRun struct {
cancel context.CancelFunc
// tenant is the ID of the run's tenant, set by registerRun.
tenant string
// finalizing is set once the result is being saved. A cancellation can
// no longer change the outcome then.
finalizing bool
}

// interruptedByShutdown reports whether scanCtx ended because the run
// context ctx ended, as it does when the daemon stops or loses its lease,
// rather than because someone asked to cancel this scan.
func (r *activeRun) interruptedByShutdown(ctx, scanCtx context.Context) bool {
if !errors.Is(scanCtx.Err(), context.Canceled) || ctx.Err() == nil {
return false
}
if r == nil {
return true
}
r.mu.RLock()
defer r.mu.RUnlock()
return !r.scan.CancelRequested
}

// ScanInterruptedMessage is the error of a scan that stopped because
// EdgeWatch stopped while it ran.
const ScanInterruptedMessage = "scan interrupted because EdgeWatch stopped"

type cronSlogLogger struct{ logger *slog.Logger }

func (l cronSlogLogger) Info(msg string, keysAndValues ...interface{}) {
Expand All @@ -108,6 +130,10 @@ func (l cronSlogLogger) Error(err error, msg string, keysAndValues ...interface{
// accepted because the daemon is stopping.
var ErrShuttingDown = errors.New("application is shutting down")

// ErrScanFinalizing is returned by CancelScan once a scan is saving its
// result, when cancelling it would no longer change the outcome.
var ErrScanFinalizing = errors.New("scan is saving its result and can no longer be canceled")

// ErrScanWorkBudget is returned before a lease is acquired when a job's
// estimated probe count exceeds the deployment guard and the job has not
// explicitly opted into high-cost work.
Expand Down Expand Up @@ -966,6 +992,10 @@ func (a *App) runJobWithQueueMarker(ctx context.Context, scope store.TenantScope
if errors.Is(scanCtx.Err(), context.Canceled) {
scan.Status = "canceled"
scan.Error = "scan canceled"
if run.interruptedByShutdown(ctx, scanCtx) {
scan.Error = ScanInterruptedMessage
scan.Interrupted = true
}
} else if errors.Is(scanCtx.Err(), context.DeadlineExceeded) || errors.Is(scanErr, context.DeadlineExceeded) {
scan.Status = "timed_out"
scan.Error = "scan timed out"
Expand All @@ -983,7 +1013,7 @@ func (a *App) runJobWithQueueMarker(ctx context.Context, scope store.TenantScope
// target scopes and will not advance a baseline from this scan.
engine.MarkIncompleteScan(&scan)
}
a.updateActivePhase(scan.ID, "finalizing")
a.beginActiveFinalization(scan.ID)
persistTimeout := scanPersistenceTimeout(len(scan.Snapshot.Hosts))
if a.persistenceBudget != nil {
persistTimeout = a.persistenceBudget(len(scan.Snapshot.Hosts))
Expand Down Expand Up @@ -1153,8 +1183,10 @@ func (a *App) CancelScan(scope store.TenantScope, id string) error {
if run.cancel == nil {
return store.ErrNotFound
}
run.cancel()
run.scan.Phase = "cancelling"
if run.finalizing {
return ErrScanFinalizing
}
run.requestCancelLocked()
return nil
}

Expand Down Expand Up @@ -1232,7 +1264,9 @@ func (a *App) updateActiveProgress(id string, progress scanner.Progress) {
}
}

func (a *App) updateActivePhase(id, phase string) {
// beginActiveFinalization reports the finalizing phase and refuses later
// cancellation requests, which could no longer change the result.
func (a *App) beginActiveFinalization(id string) {
value, ok := a.running.Load(id)
if !ok {
return
Expand All @@ -1242,8 +1276,9 @@ func (a *App) updateActivePhase(id, phase string) {
return
}
run.mu.Lock()
run.scan.Phase = phase
run.scan.Phase = "finalizing"
run.scan.ProcessAlive = false
run.finalizing = true
run.mu.Unlock()
}

Expand Down
6 changes: 3 additions & 3 deletions internal/app/coverage_more_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,10 +68,10 @@ func TestProgressPercentAndActiveRunUpdates(t *testing.T) {
t.Fatalf("active progress = %#v", got)
}
a.updateActiveProgress("scan", scanner.Progress{CompletedProbes: 30, CompletedInvocations: 2, ProcessAlive: false})
a.updateActivePhase("missing", "ignored")
a.updateActivePhase("scan", "finalizing")
a.beginActiveFinalization("missing")
a.beginActiveFinalization("scan")
got = run.snapshot()
if got.CompletedProbes != 30 || got.CompletedInvocations != 2 || got.ProcessAlive || got.Phase != "finalizing" {
if got.CompletedProbes != 30 || got.CompletedInvocations != 2 || got.ProcessAlive || got.Phase != "finalizing" || !run.finalizing {
t.Fatalf("active phase update = %#v", got)
}
if (&activeRun{}).snapshot().ElapsedSeconds != 0 {
Expand Down
8 changes: 8 additions & 0 deletions internal/app/resumable.go
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,10 @@ func (a *App) runResumableAttempt(ctx, scanCtx context.Context, ts *store.Tenant
if errors.Is(scanCtx.Err(), context.Canceled) {
scan.Status = "canceled"
scan.Error = "scan canceled while creating the scan plan"
if run.interruptedByShutdown(ctx, scanCtx) {
scan.Error = ScanInterruptedMessage + " while creating the scan plan"
scan.Interrupted = true
}
} else if errors.Is(scanCtx.Err(), context.DeadlineExceeded) || errors.Is(planErr, context.DeadlineExceeded) {
scan.Status = "timed_out"
scan.Error = "scan timed out while creating the scan plan"
Expand Down Expand Up @@ -441,6 +445,10 @@ func (a *App) runResumableAttempt(ctx, scanCtx context.Context, ts *store.Tenant
} else if canceled {
scan.Status = "canceled"
scan.Error = "scan canceled; progress was saved"
if run.interruptedByShutdown(ctx, scanCtx) {
scan.Error = ScanInterruptedMessage + "; progress was saved"
scan.Interrupted = true
}
} else {
scan.Status = "timed_out"
scan.Error = "scan timed out; progress was saved for the next trigger"
Expand Down
199 changes: 199 additions & 0 deletions internal/app/scan_interruption_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,199 @@
package app

import (
"context"
"encoding/json"
"errors"
"io"
"log/slog"
"testing"
"time"

"github.com/crypt0rr/edgewatch/internal/config"
"github.com/crypt0rr/edgewatch/internal/model"
"github.com/crypt0rr/edgewatch/internal/scanner"
"github.com/crypt0rr/edgewatch/internal/store"
"github.com/crypt0rr/edgewatch/internal/store/storetest"
)

func interruptionTestApp(t *testing.T) (*App, *store.Store) {
t.Helper()
s, err := store.Open(storetest.FreshPath(t))
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = s.Close() })
cfg := &config.Config{
Version: 1, Database: "test", Retention: config.Duration(24 * time.Hour),
Scheduler: config.Scheduler{MaxConcurrent: 1},
Web: config.Web{Listen: "127.0.0.1:8080"},
Notifications: config.Notifications{URLs: []string{"generic://localhost/edgewatch?disabletls=yes&template=json"}},
}
a, err := New(cfg, s, "missing", slog.New(slog.NewTextHandler(io.Discard, nil)))
if err != nil {
t.Fatal(err)
}
return a, s
}

func outboxCount(t *testing.T, s *store.Store) int {
t.Helper()
var count int
if err := s.DB.QueryRow(`SELECT COUNT(*) FROM outbox`).Scan(&count); err != nil {
t.Fatal(err)
}
return count
}

func TestShutdownInterruptedScanIsRecordedWithoutNotification(t *testing.T) {
t.Parallel()
a, s := interruptionTestApp(t)
record, err := defaultTenant(s).CreateJob(context.Background(), config.NormalizeJob(config.Job{
Name: "interrupted", Schedule: "0 * * * *", Timezone: "UTC", Targets: []string{"127.0.0.1"},
TCP: &config.Protocol{Ports: "1", Mode: "connect"}, Timing: "balanced", Timeout: config.Duration(time.Minute),
}))
if err != nil {
t.Fatal(err)
}
blocking := &blockingScanner{started: make(chan struct{}), release: make(chan struct{})}
a.Scanner = blocking
runCtx, stop := context.WithCancel(context.Background())
done := make(chan struct{})
var scan model.Scan
var events []model.Event
var runErr error
go func() {
scan, events, runErr = a.RunJobRecord(runCtx, record)
close(done)
}()
select {
case <-blocking.started:
case <-time.After(2 * time.Second):
t.Fatal("scan did not start")
}
// The daemon stopping cancels the run context, not this scan.
stop()
select {
case <-done:
case <-time.After(2 * time.Second):
t.Fatal("interrupted scan did not finish")
}
if !errors.Is(runErr, context.Canceled) || scan.Status != "canceled" || scan.Error != ScanInterruptedMessage {
t.Fatalf("interrupted scan = status %q error %q err %v", scan.Status, scan.Error, runErr)
}
if len(events) != 1 || events[0].Type != model.EventScanInterrupted || events[0].Message != "Scan interrupted because EdgeWatch stopped" {
t.Fatalf("interrupted scan events = %#v", events)
}
if count := outboxCount(t, s); count != 0 {
t.Fatalf("an interrupted scan queued %d notifications", count)
}
var raw []byte
if err := s.DB.QueryRow(`SELECT payload_json FROM events ORDER BY id DESC LIMIT 1`).Scan(&raw); err != nil {
t.Fatal(err)
}
var stored model.Event
if err := json.Unmarshal(raw, &stored); err != nil || stored.Type != model.EventScanInterrupted {
t.Fatalf("activity history = %s (%v), want the interruption", raw, err)
}
}

func TestShutdownInterruptedResumableScanKeepsProgressWithoutNotification(t *testing.T) {
t.Parallel()
a, s := interruptionTestApp(t)
probe := &resumableTestScanner{}
a.Scanner = probe
record, err := defaultTenant(s).CreateJob(context.Background(), config.NormalizeJob(config.Job{
Name: "broad", Schedule: "0 * * * *", Timezone: "UTC", Targets: []string{"192.0.2.1"},
TCP: &config.Protocol{Ports: "1-2", Mode: "syn"}, Timeout: config.Duration(time.Minute), ResumeWindow: config.Duration(time.Hour),
}))
if err != nil {
t.Fatal(err)
}
runCtx, stop := context.WithCancel(context.Background())
done := make(chan struct{})
var scan model.Scan
var events []model.Event
go func() {
scan, events, _ = a.RunJobRecord(runCtx, record)
close(done)
}()
// The second unit blocks until its context ends; stop once it runs.
deadline := time.Now().Add(5 * time.Second)
for {
probe.mu.Lock()
started := probe.calls[1] > 0
probe.mu.Unlock()
if started {
break
}
if time.Now().After(deadline) {
t.Fatal("second work unit did not start")
}
time.Sleep(5 * time.Millisecond)
}
stop()
select {
case <-done:
case <-time.After(5 * time.Second):
t.Fatal("interrupted resumable scan did not finish")
}
if scan.Status != "canceled" || scan.Error != ScanInterruptedMessage+"; progress was saved" || scan.CycleStatus != "paused" {
t.Fatalf("interrupted resumable scan = status %q error %q cycle %q", scan.Status, scan.Error, scan.CycleStatus)
}
if len(events) != 1 || events[0].Type != model.EventScanInterrupted {
t.Fatalf("interrupted resumable scan events = %#v", events)
}
if count := outboxCount(t, s); count != 0 {
t.Fatalf("an interrupted resumable scan queued %d notifications", count)
}
}

func TestCancelRequestOutlivesProgressAndEndsAtFinalization(t *testing.T) {
t.Parallel()
a := &App{}
scope := store.DefaultTenantScope()
canceled := false
run := &activeRun{scan: model.ActiveScan{ID: "scan", Phase: "scanning"}, cancel: func() { canceled = true }}
a.registerRun(scope.ID(), "scan", run)
if err := a.CancelScan(scope, "scan"); err != nil || !canceled {
t.Fatalf("CancelScan = %v, canceled %v", err, canceled)
}
// Nmap and Naabu report one last progress update when their process is
// stopped. The phase follows it; the cancellation request stays visible.
a.updateActiveProgress("scan", scanner.Progress{Phase: "scanning"})
if got := run.snapshot(); !got.CancelRequested || got.Phase != "scanning" {
t.Fatalf("after progress = cancel_requested %v phase %q", got.CancelRequested, got.Phase)
}
a.beginActiveFinalization("scan")
if got := run.snapshot(); !got.CancelRequested || got.Phase != "finalizing" {
t.Fatalf("after finalization began = cancel_requested %v phase %q", got.CancelRequested, got.Phase)
}
if err := a.CancelScan(scope, "scan"); !errors.Is(err, ErrScanFinalizing) {
t.Fatalf("CancelScan during finalization = %v, want ErrScanFinalizing", err)
}
}

func TestInterruptedByShutdownNeedsAStoppedRunAndNoCancelRequest(t *testing.T) {
t.Parallel()
live := context.Background()
stopped := canceledContext()
if (&activeRun{}).interruptedByShutdown(live, stopped) {
t.Fatal("a scan canceled while the run context is live was read as a shutdown")
}
if !(&activeRun{}).interruptedByShutdown(stopped, stopped) {
t.Fatal("a scan stopped with its run context was not read as a shutdown")
}
if (&activeRun{scan: model.ActiveScan{CancelRequested: true}}).interruptedByShutdown(stopped, stopped) {
t.Fatal("a requested cancellation during shutdown was read as a shutdown")
}
var missing *activeRun
if !missing.interruptedByShutdown(stopped, stopped) {
t.Fatal("a scan without an active run was not read as a shutdown")
}
timedOut, cancel := context.WithTimeout(live, 0)
defer cancel()
<-timedOut.Done()
if (&activeRun{}).interruptedByShutdown(stopped, timedOut) {
t.Fatal("a timed-out scan was read as a shutdown")
}
}
8 changes: 8 additions & 0 deletions internal/app/units.go
Original file line number Diff line number Diff line change
Expand Up @@ -166,9 +166,17 @@ func (a *App) registerRun(tenantID, id string, run *activeRun) {
func (run *activeRun) requestCancel() {
run.mu.Lock()
defer run.mu.Unlock()
run.requestCancelLocked()
}

// requestCancelLocked cancels the run and records that the cancellation was
// requested, so the scan is reported as canceled rather than interrupted.
// The caller holds run.mu.
func (run *activeRun) requestCancelLocked() {
if run.cancel != nil {
run.cancel()
run.scan.Phase = "cancelling"
run.scan.CancelRequested = true
}
}

Expand Down
Loading
Loading