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
5 changes: 4 additions & 1 deletion docs/src/content/docs/user-guide/jobs-baselines-incidents.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,10 @@ changes confirm an incident. When a security-impacting job setting changes,
EdgeWatch shows the affected scope and asks for explicit rebaselining. Schedule
and execution-tuning changes do not reset the baseline. A run that waits for a
free scan slot uses the job's settings when it starts; if the job is paused or
archived while a scheduled run waits, that run is skipped.
archived while a scheduled run waits, that run is skipped. **Cancel queued
scan** on the job page or the Overview withdraws a run that is still waiting
for a slot, so it never starts; a scan that has started is stopped with
**Cancel scan** instead.

**Pause schedule** on the job page stops a job's scheduled runs without
editing the job, and **Resume schedule** starts them again; **Scan now** keeps
Expand Down
81 changes: 74 additions & 7 deletions internal/app/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,37 @@ var ErrScanCycleStalled = errors.New("scan cycle is stalled; manual retry requir
// slot.
var ErrQueuedRunSkipped = errors.New("queued run skipped")

// ErrQueuedRunCanceled is returned by a run whose wait for a scan slot was
// canceled with CancelQueuedRun.
var ErrQueuedRunCanceled = errors.New("queued run was canceled")

// ErrRunNotQueued is returned by CancelQueuedRun when the job has no run
// waiting for a scan slot, for example because it has already started.
var ErrRunNotQueued = errors.New("job has no queued run")

// queuedRun is a run waiting for a scan slot. cancel ends its wait. started
// and canceled are guarded by mu, so a cancellation either ends the wait or
// is refused because the run has already taken its slot.
type queuedRun struct {
mu sync.Mutex
run model.QueuedRun
cancel context.CancelCauseFunc
started bool
canceled bool
}

// start marks the run as having taken its slot. It reports false when the
// run was canceled first.
func (q *queuedRun) start() bool {
q.mu.Lock()
defer q.mu.Unlock()
if q.canceled {
return false
}
q.started = true
return true
}

const (
scanPersistenceTimeoutFloor = 10 * time.Second
scanPersistenceTimeoutPerHost = 25 * time.Millisecond
Expand Down Expand Up @@ -573,6 +604,8 @@ func scanSkipReason(err error) string {
return "paused"
}
return "could_not_start"
case errors.Is(err, ErrQueuedRunCanceled):
return "canceled"
case errors.Is(err, ErrScanWorkBudget):
return "budget"
case errors.Is(err, store.ErrTenantNotActive):
Expand Down Expand Up @@ -794,22 +827,33 @@ func (a *App) runJobWithQueueMarker(ctx context.Context, scope store.TenantScope
if manual {
trigger = "manual"
}
queued := &model.QueuedRun{JobID: jobID, Job: job.Name, QueuedAt: time.Now().UTC(), Trigger: trigger, TenantID: scope.ID()}
// The wait has its own context, so CancelQueuedRun can end it without
// touching the run context the scan itself will use.
waitCtx, cancelWait := context.WithCancelCause(ctx)
defer cancelWait(nil)
queued := &queuedRun{run: model.QueuedRun{JobID: jobID, Job: job.Name, QueuedAt: time.Now().UTC(), Trigger: trigger, TenantID: scope.ID()}, cancel: cancelWait}
queuedInPool := false
releaseSlot, slotErr := a.slots.AcquireWithQueued(ctx, scope.ID(), func() {
releaseSlot, slotErr := a.slots.AcquireWithQueued(waitCtx, scope.ID(), func() {
a.queuedRuns.Store(jobID, queued)
queuedInPool = true
})
if slotErr != nil {
if queuedInPool {
a.queuedRuns.Delete(jobID)
a.queuedRuns.CompareAndDelete(jobID, queued)
}
if errors.Is(context.Cause(waitCtx), ErrQueuedRunCanceled) {
return model.Scan{}, nil, ErrQueuedRunCanceled
}
return model.Scan{}, nil, slotErr
}
if queuedInPool {
defer a.queuedRuns.Delete(jobID)
defer a.queuedRuns.CompareAndDelete(jobID, queued)
}
defer releaseSlot()
if !queued.start() {
// The cancellation arrived as the slot was granted.
return model.Scan{}, nil, ErrQueuedRunCanceled
}
ts, system := a.Store.Tenant(scope), a.Store.System()
var queuedErr error
if job, revision, queuedErr = a.queuedManagedJob(ctx, ts, job, jobID, revision, manual); queuedErr != nil {
Expand Down Expand Up @@ -1151,8 +1195,8 @@ func (a *App) ActiveScans(scope store.TenantScope) []model.ActiveScan {
func (a *App) QueuedRuns(scope store.TenantScope) []model.QueuedRun {
var queued []model.QueuedRun
a.queuedRuns.Range(func(_, value any) bool {
if run, ok := value.(*model.QueuedRun); ok && run.TenantID == scope.ID() {
queued = append(queued, *run)
if entry, ok := value.(*queuedRun); ok && entry.run.TenantID == scope.ID() {
queued = append(queued, entry.run)
}
return true
})
Expand All @@ -1165,6 +1209,29 @@ func (a *App) QueuedRuns(scope store.TenantScope) []model.QueuedRun {
return queued
}

// CancelQueuedRun ends the wait of the job's run that is queued for a scan
// slot in the tenant of scope. The run then never starts and is reported as
// skipped. A job of another tenant, or one without a queued run, is
// ErrRunNotQueued, exactly as a run that has already taken its slot.
func (a *App) CancelQueuedRun(scope store.TenantScope, jobID string) error {
value, ok := a.queuedRuns.Load(jobID)
if !ok {
return ErrRunNotQueued
}
entry, ok := value.(*queuedRun)
if !ok || !scope.Valid() || entry.run.TenantID != scope.ID() {
return ErrRunNotQueued
}
entry.mu.Lock()
defer entry.mu.Unlock()
if entry.started {
return ErrRunNotQueued
}
entry.canceled = true
entry.cancel(ErrQueuedRunCanceled)
return nil
}

// CancelScan requests cancellation of an active scan of the tenant of scope.
// The scanner owns the process context and will persist a canceled terminal
// record without mutating baseline or incident state. Another tenant's scan
Expand Down Expand Up @@ -1692,7 +1759,7 @@ func (a *App) startManagedScheduled(ctx context.Context, scope store.TenantScope
a.Logger.Warn("scheduled run skipped because resumable cycle is stalled; manual retry required", "job", record.Job.Name)
return
}
if errors.Is(runErr, ErrQueuedRunSkipped) {
if errors.Is(runErr, ErrQueuedRunSkipped) || errors.Is(runErr, ErrQueuedRunCanceled) {
a.Logger.Info("scheduled run skipped", "job", record.Job.Name, "reason", runErr)
return
}
Expand Down
115 changes: 115 additions & 0 deletions internal/app/queued_cancel_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
package app

import (
"context"
"errors"
"reflect"
"testing"
"time"

"github.com/crypt0rr/edgewatch/internal/model"
"github.com/crypt0rr/edgewatch/internal/store"
)

func TestCancelQueuedRunEndsTheWaitAndReportsTheSkip(t *testing.T) {
t.Parallel()
a, db := newLifecycleTestApp(t, schedulerFake{}, nil)
ctx, _ := a.BeginRun(context.Background())
record, err := defaultTenant(db).CreateJob(ctx, lifecycleJob("queued-cancel"))
if err != nil {
t.Fatal(err)
}
events := make(chan model.Event, 8)
a.SetEventHandler(func(event model.Event) { events <- event })
releaseSlot := holdScanSlot(t, a)
defer releaseSlot()
done := make(chan error, 1)
if err := a.StartManagedRun(defaultTenant(db), record.ID, func(_ model.Scan, _ []model.Event, err error) { done <- err }); err != nil {
t.Fatal(err)
}
waitForSlotQueue(t, a, 1)
scope := store.DefaultTenantScope()
if queued := a.QueuedRuns(scope); len(queued) != 1 || queued[0].JobID != record.ID {
t.Fatalf("queued runs = %#v", queued)
}

// A scope without a unit cannot cancel the run.
if err := a.CancelQueuedRun(store.TenantScope{}, record.ID); !errors.Is(err, ErrRunNotQueued) {
t.Fatalf("cancel without a unit = %v, want ErrRunNotQueued", err)
}
if err := a.CancelQueuedRun(scope, record.ID); err != nil {
t.Fatal(err)
}
select {
case err := <-done:
if !errors.Is(err, ErrQueuedRunCanceled) {
t.Fatalf("canceled queued run result = %v, want ErrQueuedRunCanceled", err)
}
case <-time.After(5 * time.Second):
t.Fatal("canceled queued run did not return")
}
if queued := a.QueuedRuns(scope); len(queued) != 0 {
t.Fatalf("queued runs after cancel = %#v", queued)
}
if err := a.CancelQueuedRun(scope, record.ID); !errors.Is(err, ErrRunNotQueued) {
t.Fatalf("second cancel = %v, want ErrRunNotQueued", err)
}
// The run never took a slot; the slot this test holds is still the only
// one in use.
want := slotSnapshot{Capacity: 1, InUse: 1, Keys: map[string]slotUsage{defaultSlotKey: {InUse: 1, Limit: 1}}}
if got := a.slots.CapacitySnapshot(); !reflect.DeepEqual(got, want) {
t.Fatalf("slots after cancel = %#v, want %#v", got, want)
}
var skipped []model.Event
for len(events) > 0 {
event := <-events
if event.Type == "scan.started" {
t.Fatalf("canceled queued run started: %+v", event)
}
if event.Type == "scan.skipped" {
skipped = append(skipped, event)
}
}
if len(skipped) != 1 || skipped[0].Reason != "canceled" || skipped[0].JobID != record.ID {
t.Fatalf("skip events = %+v, want one canceled skip", skipped)
}
scans, err := defaultTenant(db).ListJobScans(context.Background(), record.ID, 10)
if err != nil || len(scans) != 0 {
t.Fatalf("canceled queued run persisted scans: %#v, %v", scans, err)
}
}

func TestCancelQueuedRunRefusesARunThatTookItsSlot(t *testing.T) {
t.Parallel()
a := &App{}
scope := store.DefaultTenantScope()
canceled := false
entry := &queuedRun{run: model.QueuedRun{JobID: "job", TenantID: scope.ID()}, cancel: func(error) { canceled = true }}
a.queuedRuns.Store("job", entry)
if !entry.start() {
t.Fatal("an uncanceled run could not start")
}
if err := a.CancelQueuedRun(scope, "job"); !errors.Is(err, ErrRunNotQueued) || canceled {
t.Fatalf("cancel after start = %v (canceled %v), want ErrRunNotQueued", err, canceled)
}
if err := a.CancelQueuedRun(scope, "missing"); !errors.Is(err, ErrRunNotQueued) {
t.Fatalf("cancel of a job without a queued run = %v", err)
}
foreign := &queuedRun{run: model.QueuedRun{JobID: "foreign", TenantID: "another-unit"}, cancel: func(error) { canceled = true }}
a.queuedRuns.Store("foreign", foreign)
if err := a.CancelQueuedRun(scope, "foreign"); !errors.Is(err, ErrRunNotQueued) || canceled {
t.Fatalf("cancel of another unit's run = %v (canceled %v), want ErrRunNotQueued", err, canceled)
}

late := &queuedRun{run: model.QueuedRun{JobID: "late", TenantID: scope.ID()}, cancel: func(error) {}}
a.queuedRuns.Store("late", late)
if err := a.CancelQueuedRun(scope, "late"); err != nil {
t.Fatal(err)
}
if late.start() {
t.Fatal("a run canceled as its slot was granted still started")
}
if got := scanSkipReason(ErrQueuedRunCanceled); got != "canceled" {
t.Fatalf("skip reason = %q, want canceled", got)
}
}
1 change: 1 addition & 0 deletions internal/store/audit_category.go
Original file line number Diff line number Diff line change
Expand Up @@ -127,6 +127,7 @@ var auditActionCategories = map[string]string{
"notifications.updated": auditCategoryData,
"public_dashboard.updated": auditCategoryData,
"scan.cancel_requested": auditCategoryData,
"scan.queued_run_canceled": auditCategoryData,
"scan.cycle_discarded": auditCategoryData,
"scan.run_requested": auditCategoryData,
"scanner_profile.archived": auditCategoryData,
Expand Down
2 changes: 1 addition & 1 deletion internal/store/audit_category_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ var knownAuditActions = map[string][]string{
"notifications.test", "notifications.test_failed", "notifications.update_routing",
"notifications.updated", "notifications.incident_reminders_changed",
"public_dashboard.updated",
"scan.cancel_requested", "scan.cycle_discarded", "scan.run_requested",
"scan.cancel_requested", "scan.cycle_discarded", "scan.queued_run_canceled", "scan.run_requested",
"scanner_profile.archived", "scanner_profile.created", "scanner_profile.restored",
"scanner_profile.updated",
},
Expand Down
6 changes: 6 additions & 0 deletions internal/web/job_handlers.go
Original file line number Diff line number Diff line change
Expand Up @@ -596,6 +596,12 @@ func (s *Server) jobRoute(w http.ResponseWriter, r *http.Request, session store.
}
return
}
if len(parts) == 2 && parts[1] == "run" && r.Method == http.MethodDelete {
if job, ok := s.resolveJob(w, r, ts, id, jobMissingOnAnyError); ok {
s.cancelQueuedRun(w, r, session, ts, job)
}
return
}
if len(parts) == 2 && parts[1] == "scan-cycle" && r.Method == http.MethodGet {
if job, ok := s.resolveJob(w, r, ts, id, jobStoreErrorInternal); ok {
s.scanCycle(w, r, ts, job)
Expand Down
3 changes: 2 additions & 1 deletion internal/web/permissions.go
Original file line number Diff line number Diff line change
Expand Up @@ -310,7 +310,7 @@ func requiredJobPermission(path, method string) string {
return auth.PermissionJobsWrite
}
case "run":
if method == http.MethodPost {
if method == http.MethodPost || method == http.MethodDelete {
return auth.PermissionJobsRun
}
case "scan-cycle":
Expand Down Expand Up @@ -564,6 +564,7 @@ var apiRoutes = []apiRoute{
{Method: http.MethodPost, Template: "/jobs/{id}/pause", Permission: auth.PermissionJobsWrite, Mutates: true, Example: "/jobs/job-1/pause"},
{Method: http.MethodPost, Template: "/jobs/{id}/resume", Permission: auth.PermissionJobsWrite, Mutates: true, Example: "/jobs/job-1/resume"},
{Method: http.MethodPost, Template: "/jobs/{id}/run", Permission: auth.PermissionJobsRun, Mutates: true, Example: "/jobs/job-1/run"},
{Method: http.MethodDelete, Template: "/jobs/{id}/run", Permission: auth.PermissionJobsRun, Mutates: true, Example: "/jobs/job-1/run"},
{Method: http.MethodGet, Template: "/jobs/{id}/scan-cycle", Permission: auth.PermissionScansRead, Example: "/jobs/job-1/scan-cycle"},
{Method: http.MethodDelete, Template: "/jobs/{id}/scan-cycle/{cycle}", Permission: auth.PermissionJobsRun, Mutates: true, Example: "/jobs/job-1/scan-cycle/cycle-1"},
{Method: http.MethodGet, Template: "/jobs/{id}/scans", Permission: auth.PermissionScansRead, Example: "/jobs/job-1/scans"},
Expand Down
75 changes: 75 additions & 0 deletions internal/web/scan_control_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,10 +3,12 @@ package web
import (
"context"
"encoding/json"
"errors"
"io"
"log/slog"
"net/http"
"net/http/httptest"
"strings"
"sync"
"testing"
"time"
Expand Down Expand Up @@ -300,3 +302,76 @@ func TestActiveScanEndpointIncludesTenantScopedQueuedRuns(t *testing.T) {
t.Fatal("queued scan unexpectedly succeeded after cancellation")
}
}

func TestCancelQueuedRunRouteWithdrawsTheWaitingRun(t *testing.T) {
t.Parallel()
ctx := context.Background()
db, err := store.Open(storetest.FreshPath(t))
if err != nil {
t.Fatal(err)
}
defer db.Close()
cfg := &config.Config{Version: 1, Database: db.Path, Retention: config.Duration(24 * time.Hour), Scheduler: config.Scheduler{MaxConcurrent: 1}, Web: config.Web{Listen: "127.0.0.1:8080"}}
a, err := app.New(cfg, db, "missing-nmap", slog.New(slog.NewTextHandler(io.Discard, nil)))
if err != nil {
t.Fatal(err)
}
scanner := &blockingWebScanner{started: make(chan struct{})}
a.Scanner = scanner
job := config.NormalizeJob(config.Job{Name: "slot-holder", Schedule: "0 * * * *", Timezone: "UTC", Targets: []string{"127.0.0.1"}, TCP: &config.Protocol{Ports: "1", Mode: "connect"}, Timeout: config.Duration(time.Hour), Timing: "balanced"})
first, err := defaultTenant(db).CreateJob(ctx, job)
if err != nil {
t.Fatal(err)
}
job.Name = "queued-job"
second, err := defaultTenant(db).CreateJob(ctx, job)
if err != nil {
t.Fatal(err)
}
a.BeginRun(ctx)
defer a.StopRun()
if err := a.StartManagedRun(defaultTenant(db), first.ID, nil); err != nil {
t.Fatal(err)
}
select {
case <-scanner.started:
case <-time.After(2 * time.Second):
t.Fatal("slot-holding scan did not start")
}
secondDone := make(chan error, 1)
if err := a.StartManagedRun(defaultTenant(db), second.ID, func(_ model.Scan, _ []model.Event, err error) { secondDone <- err }); err != nil {
t.Fatal(err)
}
deadline := time.Now().Add(2 * time.Second)
for len(a.QueuedRuns(store.DefaultTenantScope())) == 0 && time.Now().Before(deadline) {
time.Sleep(5 * time.Millisecond)
}
server := NewServer(a, db, slog.New(slog.NewTextHandler(io.Discard, nil)))
session := store.Session{UserID: "admin", Username: "admin", Role: store.RoleAdministrator}
cancelRun := func(id string) *httptest.ResponseRecorder {
recorder := httptest.NewRecorder()
server.jobRoute(recorder, httptest.NewRequest(http.MethodDelete, "/api/v1/jobs/"+id+"/run", nil), session, defaultTenantStore(server), id+"/run")
return recorder
}
if response := cancelRun(second.ID); response.Code != http.StatusAccepted || !strings.Contains(response.Body.String(), `"status":"canceled"`) {
t.Fatalf("cancel queued run = %d %s", response.Code, response.Body.String())
}
select {
case err := <-secondDone:
if !errors.Is(err, app.ErrQueuedRunCanceled) {
t.Fatalf("queued run result = %v, want ErrQueuedRunCanceled", err)
}
case <-time.After(5 * time.Second):
t.Fatal("canceled queued run did not return")
}
// Nothing is queued now, and the running scan is not a queued run.
for _, id := range []string{second.ID, first.ID} {
want := wantErrorBody(t, "run_not_queued", "The job has no scan waiting for a slot.", nil)
if response := cancelRun(id); response.Code != http.StatusConflict || response.Body.String() != want {
t.Fatalf("cancel %s without a queued run = %d %q, want 409 %q", id, response.Code, response.Body.String(), want)
}
}
if active := a.ActiveScans(store.DefaultTenantScope()); len(active) != 1 || active[0].JobID != first.ID {
t.Fatalf("active scans after the queued cancel = %+v", active)
}
}
Loading
Loading