From 919b09278e5eee9d4a5d0657ae385f2d1312e712 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Ren=C3=A9=20Zander?= <96420681+renezander030@users.noreply.github.com> Date: Thu, 8 Oct 2026 02:29:11 +0000 Subject: [PATCH] Prepare v0.10.0: unattended operations - draftcat doctor: read-only preflight of config, credentials, operator access, state store, listener ports and schedules, with --json. - Cron expressions and @daily-style schedules in a per-pipeline time zone, with optional catch_up for slots missed during downtime. - Interval schedules continue from the recorded run history. - One claim admits every run (timer, /run, Run-now, webhook). - pause_after_failures pauses a failing timer pipeline and notifies once. - SIGINT/SIGTERM drain within timeouts.shutdown_grace. - budgets.alert_at notifies at fractions of the daily caps. - Configured credentials are redacted from logs, notifications, stored errors, spans and OTLP exports. - Current-state Prometheus gauges. --- CHANGELOG.md | 20 + README.md | 8 +- budget.go | 4 + budget_alerts.go | 103 +++++ budget_alerts_test.go | 130 +++++++ config.yaml | 6 +- docs/operations.md | 147 ++++++++ doctor_cmd.go | 346 +++++++++++++++++ doctor_cmd_test.go | 161 ++++++++ engine_lifecycle.go | 189 ++++++++++ internal/config/config.go | 28 +- internal/obs/gauges.go | 64 ++++ internal/obs/gauges_test.go | 50 +++ internal/obs/metrics.go | 2 + internal/obs/obs.go | 4 +- internal/obs/otlp.go | 3 + internal/obs/redact_test.go | 27 ++ internal/redact/redact.go | 101 +++++ internal/redact/redact_test.go | 74 ++++ internal/schedule/schedule.go | 260 +++++++++++++ internal/schedule/schedule_test.go | 124 ++++++ internal/validate/validate.go | 45 ++- internal/validate/validate_schedule_test.go | 73 ++++ main.go | 258 ++++--------- package-lock.json | 4 +- package.json | 2 +- redaction_test.go | 102 +++++ relay_channel.go | 5 +- scheduler.go | 398 ++++++++++++++++++++ scheduler_test.go | 364 ++++++++++++++++++ version.go | 2 +- 31 files changed, 2907 insertions(+), 197 deletions(-) create mode 100644 budget_alerts.go create mode 100644 budget_alerts_test.go create mode 100644 docs/operations.md create mode 100644 doctor_cmd.go create mode 100644 doctor_cmd_test.go create mode 100644 engine_lifecycle.go create mode 100644 internal/obs/gauges.go create mode 100644 internal/obs/gauges_test.go create mode 100644 internal/obs/redact_test.go create mode 100644 internal/redact/redact.go create mode 100644 internal/redact/redact_test.go create mode 100644 internal/schedule/schedule.go create mode 100644 internal/schedule/schedule_test.go create mode 100644 internal/validate/validate_schedule_test.go create mode 100644 redaction_test.go create mode 100644 scheduler.go create mode 100644 scheduler_test.go diff --git a/CHANGELOG.md b/CHANGELOG.md index a2dc211..44a912d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,25 @@ # Changelog +## 0.10.0 + +Operations you can leave running: calendar schedules, a preflight, a clean shutdown and credentials that stay out of every log. + +- Add `draftcat doctor`, a read-only preflight of config, credentials, operator access, the state store, listener ports and schedules, with a fix for each finding and `--json` output. Exits 1 when a check fails. +- Schedule pipelines with five-field cron expressions and `@hourly`/`@daily`/`@weekly`/`@monthly`, evaluated in a per-pipeline `timezone`. `catch_up: true` runs a slot missed during downtime once at start. +- Continue interval schedules from the recorded run history across restarts. +- Admit every pipeline run, whether from the timer, `/run`, the Run-now button or a webhook, through one claim, keeping one run per pipeline in progress. +- Pause a timer pipeline after `pause_after_failures` consecutive failed runs and notify the operator once; the streak is restored at start and `/cron resume` clears it. +- Drain on `SIGINT`/`SIGTERM`: refuse new runs, close the webhook listener and wait up to `timeouts.shutdown_grace` (default 30s) for running pipelines. +- Notify the operator at `budgets.alert_at` fractions of the daily token and cost caps, once per threshold and UTC day, also across restarts. +- Replace configured credentials with `[REDACTED:]` in logs, operator notifications, stored run and webhook errors, JSON spans and OTLP exports. +- Export current-state Prometheus gauges: running pipelines, paused pipelines, failure streaks, open approvals, and today's tokens and spend against the caps. + +`/cron set` accepts cron expressions. See [running Draftcat unattended](docs/operations.md). + +### Upgrade notes + +All new settings are optional. Configurations without them behave as before, except that interval pipelines continue from their last recorded run instead of starting a fresh interval at boot, and the engine waits up to 30 seconds for running pipelines on shutdown (`timeouts.shutdown_grace: 0s` restores an immediate exit). An operator `/run` of a pipeline that is already running is refused instead of starting a second run. + ## 0.9.0 - Isolate pipeline token and cost totals and serialize model admission against settled daily usage. diff --git a/README.md b/README.md index d996884..71f1693 100644 --- a/README.md +++ b/README.md @@ -164,7 +164,7 @@ cp secrets.yaml.example secrets.yaml docker compose up ``` -Pipelines live in `config.yaml`, prompts in `skills/`. A SQLite store opens at `./state.db` on first boot. To add the EU-resident **voice AI** plugin: `go build -tags voice -o draftcat .` — the lean binary is unchanged when the tag is off. +Pipelines live in `config.yaml`, prompts in `skills/`. Run `draftcat doctor` to check credentials, operators, the state store and schedules before the first start. A SQLite store opens at `./state.db` on first boot. To add the EU-resident **voice AI** plugin: `go build -tags voice -o draftcat .` — the lean binary is unchanged when the tag is off. ## Deploy — where it runs @@ -282,6 +282,7 @@ budgets: per_day_tokens: 100000 per_day_cost: 5.00 # money cap, same unit as your model rates (0 = off) per_pipeline_cost: 0.50 + alert_at: [0.5, 0.8] # tell the operator at 50% and 80% of a daily cap observability: {spans: false} # or DRAFTCAT_TRACE=1 state: {path: ./state.db} @@ -357,6 +358,7 @@ Skills are YAML prompt templates in `skills/` with an `output_schema` the engine ```bash draftcat # run the engine (validates config first; refuses to start on errors) draftcat validate [--strict] # lint config + skills +draftcat doctor [--json] # read-only preflight: credentials, operators, state, ports, schedules draftcat test # dry-run against fixtures// (never touches real APIs) draftcat runs [pipeline] # recent runs + the approval decisions in each (--json to archive) draftcat pending # approval gates waiting on a human right now (--json) @@ -387,7 +389,9 @@ Approval rows can be made tamper-evident with signed receipts, so a later audit can verify which operator approved which payload hash. See [`docs/action-receipts.md`](docs/action-receipts.md). -A pipeline's `schedule` decides when it runs — an interval (`1h`), `manual` (operator `/run` only), or `webhook`. The `webhook` server is opt-in and opens no port unless enabled: +A pipeline's `schedule` decides when it runs — an interval (`1h`), a cron expression in the pipeline's time zone (`"0 8 * * 1-5"` with `timezone: Europe/Berlin`), `@daily`-style shortcuts, `manual` (operator `/run` only), or `webhook`. Schedules continue from the recorded run history across restarts, `pause_after_failures` stops a pipeline that keeps failing, and `SIGTERM` lets running pipelines finish within `timeouts.shutdown_grace`. See [running Draftcat unattended](docs/operations.md). + +The `webhook` server is opt-in and opens no port unless enabled: ```yaml webhook: {enabled: true, addr: 127.0.0.1:8088, secret_env: DRAFTCAT_WEBHOOK_SECRET} diff --git a/budget.go b/budget.go index e036d5c..f2d4d6e 100644 --- a/budget.go +++ b/budget.go @@ -31,6 +31,7 @@ type BudgetTracker struct { dayCostLimit float64 pipelineCostLimit float64 pipelineTokenLimit int + alerts *budgetAlerts } func (b *BudgetTracker) root() *BudgetTracker { @@ -145,6 +146,9 @@ func (b *BudgetTracker) recordUsageLocked(tokens int, cost float64) { r.costToday += cost b.tokensUsedPipeline += tokens b.costPipeline += cost + if r.alerts != nil { + r.alerts.observeLocked(r) + } } func (b *BudgetTracker) record(tokens int) { diff --git a/budget_alerts.go b/budget_alerts.go new file mode 100644 index 0000000..6c7fb58 --- /dev/null +++ b/budget_alerts.go @@ -0,0 +1,103 @@ +package main + +import ( + "fmt" + "log" + "sort" + "time" + + "github.com/renezander030/draftcat/internal/config" +) + +// budgetAlerts notifies the operator when today's usage crosses a configured +// fraction of a daily cap (budgets.alert_at). Each threshold fires once per +// UTC day and cap; the mark is kept in the state store so a restart does not +// repeat it. +type budgetAlerts struct { + thresholds []float64 // ascending, each in (0, 1) + tokenCap int + costCap float64 + notify func(string) + day string + fired map[string]bool +} + +// configureAlerts enables threshold alerts on the shared daily ledger. +// notify is called outside the budget lock. +func (b *BudgetTracker) configureAlerts(cfg *config.Config, notify func(string)) { + var th []float64 + for _, f := range cfg.Budgets.AlertAt { + if f > 0 && f < 1 { + th = append(th, f) + } + } + if len(th) == 0 || notify == nil || (cfg.Budgets.PerDayTokens <= 0 && cfg.Budgets.PerDayCost <= 0) { + return + } + sort.Float64s(th) + r := b.root() + r.mu.Lock() + defer r.mu.Unlock() + r.alerts = &budgetAlerts{ + thresholds: th, + tokenCap: cfg.Budgets.PerDayTokens, + costCap: cfg.Budgets.PerDayCost, + notify: notify, + fired: map[string]bool{}, + } +} + +// observeLocked runs with the root budget lock held, after usage changed. +func (a *budgetAlerts) observeLocked(r *BudgetTracker) { + day := r.dayStart.UTC().Format("2006-01-02") + if day != a.day { + a.day = day + a.fired = map[string]bool{} + } + if a.tokenCap > 0 { + used := float64(r.tokensUsedToday) / float64(a.tokenCap) + if f, ok := a.crossLocked(r, day, "tokens", used); ok { + a.send(fmt.Sprintf("[budget] %d%% of today's token cap used: %d of %d (UTC day %s).", pct(f), r.tokensUsedToday, a.tokenCap, day)) + } + } + if a.costCap > 0 { + used := r.costToday / a.costCap + if f, ok := a.crossLocked(r, day, "cost", used); ok { + a.send(fmt.Sprintf("[budget] %d%% of today's cost cap used: %.4f of %.4f (UTC day %s).", pct(f), r.costToday, a.costCap, day)) + } + } +} + +// crossLocked marks every threshold at or below used that has not fired today +// and returns the highest one newly crossed. +func (a *budgetAlerts) crossLocked(r *BudgetTracker, day, kind string, used float64) (float64, bool) { + var top float64 + crossed := false + for _, f := range a.thresholds { + if used < f { + break + } + key := fmt.Sprintf("%s:%s:%g", day, kind, f) + if a.fired[key] { + continue + } + a.fired[key] = true + if r.store != nil { + first, err := r.store.TryMarkSeen("_budget", "alert", key, time.Now()) + if err != nil { + log.Printf("[budget] alert mark %s: %v", key, err) + } else if !first { + continue // already sent before a restart + } + } + top, crossed = f, true + } + return top, crossed +} + +func (a *budgetAlerts) send(msg string) { + notify := a.notify + go notify(msg) +} + +func pct(f float64) int { return int(f*100 + 0.5) } diff --git a/budget_alerts_test.go b/budget_alerts_test.go new file mode 100644 index 0000000..32bc068 --- /dev/null +++ b/budget_alerts_test.go @@ -0,0 +1,130 @@ +package main + +import ( + "path/filepath" + "strings" + "testing" + "time" + + "github.com/renezander030/draftcat/internal/config" + statestore "github.com/renezander030/draftcat/internal/state" +) + +func alertCollector() (func(string), func(n int, within time.Duration) []string) { + ch := make(chan string, 16) + notify := func(m string) { ch <- m } + collect := func(n int, within time.Duration) []string { + var out []string + deadline := time.After(within) + for len(out) < n { + select { + case m := <-ch: + out = append(out, m) + case <-deadline: + return out + } + } + // Anything extra arriving shortly after is a duplicate. + select { + case m := <-ch: + out = append(out, m) + case <-time.After(50 * time.Millisecond): + } + return out + } + return notify, collect +} + +func alertConfig() *config.Config { + cfg := &config.Config{} + cfg.Budgets.PerDayCost = 10 + cfg.Budgets.PerDayTokens = 1000 + cfg.Budgets.AlertAt = []float64{0.8, 0.5} + return cfg +} + +func TestBudgetAlertFiresOncePerThreshold(t *testing.T) { + b := &BudgetTracker{dayStart: time.Now()} + notify, collect := alertCollector() + b.configureAlerts(alertConfig(), notify) + + b.RecordCost(4) // 40%: nothing + if got := collect(1, 100*time.Millisecond); len(got) != 0 { + t.Fatalf("alert below threshold: %q", got) + } + b.RecordCost(1.5) // 55% + got := collect(1, time.Second) + if len(got) != 1 || !strings.Contains(got[0], "50% of today's cost cap") { + t.Fatalf("at 55%%: %q", got) + } + b.RecordCost(0.1) // 56%: already sent + if got := collect(1, 100*time.Millisecond); len(got) != 0 { + t.Fatalf("repeated alert: %q", got) + } + b.RecordCost(3) // 86% + got = collect(1, time.Second) + if len(got) != 1 || !strings.Contains(got[0], "80% of today's cost cap") { + t.Fatalf("at 86%%: %q", got) + } +} + +func TestBudgetAlertJumpReportsHighestThreshold(t *testing.T) { + b := &BudgetTracker{dayStart: time.Now()} + notify, collect := alertCollector() + b.configureAlerts(alertConfig(), notify) + b.record(900) // 0 -> 90% of the token cap in one step + got := collect(1, time.Second) + if len(got) != 1 || !strings.Contains(got[0], "80% of today's token cap used: 900 of 1000") { + t.Fatalf("got %q", got) + } +} + +func TestBudgetAlertMarksSurviveRestart(t *testing.T) { + path := filepath.Join(t.TempDir(), "state.db") + st, err := statestore.OpenStateStore(path) + if err != nil { + t.Fatal(err) + } + b := &BudgetTracker{dayStart: time.Now()} + if err := b.attachStore(st); err != nil { + t.Fatal(err) + } + notify, collect := alertCollector() + b.configureAlerts(alertConfig(), notify) + b.RecordCost(6) + if got := collect(1, time.Second); len(got) != 1 { + t.Fatalf("first process: %q", got) + } + st.Close() + + st2, err := statestore.OpenStateStore(path) + if err != nil { + t.Fatal(err) + } + defer st2.Close() + b2 := &BudgetTracker{dayStart: time.Now()} + if err := b2.attachStore(st2); err != nil { + t.Fatal(err) + } + notify2, collect2 := alertCollector() + b2.configureAlerts(alertConfig(), notify2) + b2.RecordCost(0.5) // 65%: 50% was already sent before the restart + if got := collect2(1, 150*time.Millisecond); len(got) != 0 { + t.Fatalf("alert repeated after restart: %q", got) + } +} + +func TestBudgetAlertsOffWithoutCapOrThresholds(t *testing.T) { + cfg := &config.Config{} + cfg.Budgets.AlertAt = []float64{0.5} + b := &BudgetTracker{dayStart: time.Now()} + notify, collect := alertCollector() + b.configureAlerts(cfg, notify) + b.RecordCost(100) + if got := collect(1, 100*time.Millisecond); len(got) != 0 { + t.Fatalf("alert without a cap: %q", got) + } + if b.alerts != nil { + t.Fatal("alerts configured without a daily cap") + } +} diff --git a/config.yaml b/config.yaml index 7bec70b..133723b 100644 --- a/config.yaml +++ b/config.yaml @@ -61,11 +61,13 @@ budgets: per_step_tokens: 2048 per_pipeline_tokens: 10000 per_day_tokens: 100000 + # alert_at: [0.5, 0.8, 0.95] # notify the operator at these fractions of a daily cap timeouts: ai_call: 30s operator_approval: 4h pipeline_total: 5m + # shutdown_grace: 30s # SIGTERM waits this long for running pipelines # Structured span emission for auditing/observability. Off by default; the # DRAFTCAT_TRACE env var also enables it. One JSON line per pipeline + per step. @@ -159,7 +161,9 @@ pipelines: # Stale opportunity recovery — drafts follow-ups for dormant deals - name: opportunity-recovery - schedule: 24h + schedule: 24h # or a cron expression, e.g. "0 9 * * 1-5" + # timezone: Europe/Berlin # zone for cron schedules (default: host zone) + # pause_after_failures: 3 # pause after 3 failed runs in a row steps: - name: fetch-stale type: deterministic diff --git a/docs/operations.md b/docs/operations.md new file mode 100644 index 0000000..c38b771 --- /dev/null +++ b/docs/operations.md @@ -0,0 +1,147 @@ +# Running Draftcat unattended + +This guide covers the settings that matter once Draftcat runs as a service: +preflight, schedules, failure handling, shutdown, spend warnings, metrics and +credential redaction. + +## Preflight: `draftcat doctor` + +```bash +draftcat doctor # ./config.yaml and ./skills +draftcat doctor --config /etc/draftcat/config.yaml --skills /etc/draftcat/skills --json +``` + +`doctor` checks everything the engine needs at start and says what to do about +each problem: + +| Check | Fails when | +|---|---| +| config, validate | the file is unreadable or `draftcat validate` reports errors | +| telegram token, provider key | the configured environment variables are empty | +| operators | no `allowed_users` or no chat to reach them | +| relay | a relay is configured but its shared secret is empty | +| webhook | the listener is enabled but its secret is empty | +| state store | the file cannot be opened, or its directory is not writable | +| schedule | a schedule does not parse | + +Warnings cover unsigned receipts (`DRAFTCAT_APPROVAL_SECRET` empty), a +`secrets.yaml` readable by other users, busy listener ports, unsettled provider +usage and approval gates left open. Each timer pipeline is listed with its next +run time in its own time zone. + +`doctor` is read-only: it never creates or migrates the state store and never +calls a provider. It exits 1 when any check fails, so it can gate a deploy: + +```bash +draftcat doctor && systemctl restart draftcat +``` + +## Schedules + +```yaml +pipelines: + - name: morning-digest + schedule: "0 8 * * 1-5" # weekdays 08:00 + timezone: Europe/Berlin # IANA zone; default is the host's zone + catch_up: true # run a slot missed during downtime once at start + pause_after_failures: 3 # pause after 3 failed runs in a row + steps: [...] + - name: inbox-sweep + schedule: 30m # interval +``` + +`schedule` accepts: + +- an interval: any Go duration (`30m`, `4h`, `24h`); +- a five-field cron expression: minute, hour, day of month, month, day of week, + with `*`, lists (`9,13`), ranges (`1-5`), steps (`*/15`) and names + (`mon-fri`, `jan`). When day of month and day of week are both restricted, + either one matching is enough, as in standard cron; +- `@hourly`, `@daily`, `@weekly`, `@monthly`; +- `manual` (operator `/run` only) or `webhook`. + +Calendar schedules run at wall-clock time in `timezone`, across daylight-saving +changes: a slot inside a skipped hour runs once right after the jump. + +Schedules continue from the run history in the state store. An interval +pipeline runs one interval after its last recorded run, so a `24h` pipeline +keeps its daily cadence across restarts and an overdue one runs at the next +tick. A calendar pipeline waits for its next slot; with `catch_up: true` a slot +that passed while the engine was down runs once at start. + +`/cron` shows each pipeline's schedule, last run and next run. +`/cron set ` accepts intervals and cron expressions and refuses +a schedule that does not parse. + +## One run at a time + +Every start of a pipeline — timer, `/run`, the Run-now button, a webhook — takes +the same claim. A pipeline has at most one run in progress: a second `/run` +answers `Not started: (pipeline is already running)` and a second webhook +gets `409`. An operator's `/run` works while the pipeline's timer is paused; +webhooks and timers do not. + +## Pausing a failing pipeline + +With `pause_after_failures: N`, a timer pipeline that fails N runs in a row is +paused and the operator receives one message with the last error. Any +successful run resets the count. The count is rebuilt from the run history at +start, so a pipeline that was failing before a restart starts paused and the +operator is told again. `/cron resume ` clears the pause and the count. + +## Shutdown + +On `SIGINT` or `SIGTERM` the engine stops admitting runs, closes the webhook +listener after its in-flight requests, and waits for running pipelines: + +```yaml +timeouts: + shutdown_grace: 30s # default; 0s exits at once +``` + +Approval taps keep arriving during the wait, so a run blocked on a decision can +still complete. Set your service manager's stop timeout above the grace period +(`TimeoutStopSec=` for systemd, `stop_grace_period:` for Compose). A gate still +open when the grace period ends is closed out and reported at the next start. + +## Spend warnings + +```yaml +budgets: + per_day_tokens: 100000 + per_day_cost: 20 + alert_at: [0.5, 0.8, 0.95] +``` + +When today's usage (UTC) crosses a fraction of a daily cap, the operator +channel receives one message, for example +`[budget] 80% of today's cost cap used: 16.0400 of 20.0000 (UTC day 2026-10-08).` +Each threshold is sent once per day and cap, also across restarts. A jump past +several thresholds sends the highest. + +## Metrics + +With `observability.prometheus.enabled`, `/metrics` adds current-state gauges +next to the existing counters and histograms: + +| Gauge | Meaning | +|---|---| +| `draftcat_pipelines_running` | runs in progress | +| `draftcat_pipeline_paused{pipeline}` | 1 when the timer is paused | +| `draftcat_pipeline_consecutive_failures{pipeline}` | current failure streak | +| `draftcat_approvals_open` | approval gates waiting on a decision | +| `draftcat_budget_day_tokens`, `draftcat_budget_day_tokens_limit` | today's tokens and the cap | +| `draftcat_budget_day_cost`, `draftcat_budget_day_cost_limit` | today's spend and the cap | +| `draftcat_budget_unsettled_calls` | provider calls awaiting reconciliation | + +## Credential redaction + +The values of every credential variable the engine reads — the Telegram token, +the provider key, the webhook, relay and approval secrets, the GoHighLevel key +and OTLP header values — are replaced by `[REDACTED:]` in the process +log, operator notifications, stored run and webhook errors, JSON spans and OTLP +exports. Values shorter than eight characters are not registered. + +Approval drafts are shown to the operator exactly as they will be sent, so the +payload hash covers what the operator saw. Use a `model_policy` rule to hold or +refuse drafts that contain credentials. diff --git a/doctor_cmd.go b/doctor_cmd.go new file mode 100644 index 0000000..f405fa1 --- /dev/null +++ b/doctor_cmd.go @@ -0,0 +1,346 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "io" + "log" + "net" + "os" + "path/filepath" + "strings" + "time" + + "gopkg.in/yaml.v3" + + "github.com/renezander030/draftcat/internal/config" + "github.com/renezander030/draftcat/internal/schedule" + skillsapi "github.com/renezander030/draftcat/internal/skills" + statestore "github.com/renezander030/draftcat/internal/state" + "github.com/renezander030/draftcat/internal/validate" +) + +// doctorCheck is one line of `draftcat doctor`. Status is "ok", "warn" or +// "fail"; Hint says what to do when it is not ok. +type doctorCheck struct { + Name string `json:"name"` + Status string `json:"status"` + Detail string `json:"detail"` + Hint string `json:"hint,omitempty"` +} + +type doctorReport struct { + OK bool `json:"ok"` + Checks []doctorCheck `json:"checks"` +} + +func (r *doctorReport) add(name, status, detail, hint string) { + r.Checks = append(r.Checks, doctorCheck{Name: name, Status: status, Detail: detail, Hint: hint}) +} + +// runDoctor implements `draftcat doctor`: a read-only preflight of everything +// the engine needs at start — config, credentials, operator channel, state +// store, listeners and schedules. It never writes, never creates the state +// store and never contacts a provider. Exit 1 when any check fails. +func runDoctor(args []string, stdout io.Writer) int { + configPath := "config.yaml" + skillsDir := "skills" + jsonOut := false + for i := 0; i < len(args); i++ { + switch a := args[i]; a { + case "--json", "-json": + jsonOut = true + case "--config", "-config": + if i+1 >= len(args) { + fmt.Fprintln(os.Stderr, "doctor: --config requires a path") + return 2 + } + configPath = args[i+1] + i++ + case "--skills", "-skills": + if i+1 >= len(args) { + fmt.Fprintln(os.Stderr, "doctor: --skills requires a path") + return 2 + } + skillsDir = args[i+1] + i++ + case "-h", "--help", "help": + _, _ = io.WriteString(stdout, "Usage: draftcat doctor [--config path] [--skills dir] [--json]\n\n"+ + "Read-only preflight: config, credentials, operator channel, state store, listeners, schedules.\n"+ + "Exits 1 when a check fails.\n") + return 0 + default: + fmt.Fprintf(os.Stderr, "doctor: unknown option %q\n", a) + return 2 + } + } + // The report is the output; loader chatter would interleave with it. + log.SetOutput(io.Discard) + defer log.SetOutput(os.Stderr) + rep := doctor(configPath, skillsDir, time.Now()) + if jsonOut { + enc := json.NewEncoder(stdout) + enc.SetIndent("", " ") + _ = enc.Encode(rep) + } else { + printDoctor(stdout, rep) + } + if !rep.OK { + return 1 + } + return 0 +} + +func doctor(configPath, skillsDir string, now time.Time) (rep doctorReport) { + defer func() { + rep.OK = true + for _, c := range rep.Checks { + if c.Status == "fail" { + rep.OK = false + } + } + }() + + // #nosec G304 -- the operator's own --config argument. + data, err := os.ReadFile(configPath) + if err != nil { + rep.add("config", "fail", fmt.Sprintf("cannot read %s: %v", configPath, err), "pass --config or run from the directory holding config.yaml") + return rep + } + var cfg config.Config + if err := yaml.Unmarshal(data, &cfg); err != nil { + rep.add("config", "fail", fmt.Sprintf("cannot parse %s: %v", configPath, err), "fix the YAML syntax, then run draftcat validate") + return rep + } + rep.add("config", "ok", fmt.Sprintf("%s: %d pipeline(s)", configPath, len(cfg.Pipelines)), "") + + errs, warns := 0, 0 + var firstErr string + for _, f := range validate.CheckAtStartup(&cfg, skillsDir) { + if f.Level == "error" { + if errs == 0 { + firstErr = f.Path + ": " + f.Message + } + errs++ + } else { + warns++ + } + } + switch { + case errs > 0: + rep.add("validate", "fail", fmt.Sprintf("%d error(s), %d warning(s); first: %s", errs, warns, firstErr), "run draftcat validate for the full list") + case warns > 0: + rep.add("validate", "warn", fmt.Sprintf("%d warning(s)", warns), "run draftcat validate for details") + default: + rep.add("validate", "ok", "no findings", "") + } + + doctorSkills(&rep, skillsDir) + doctorChannel(&rep, &cfg) + doctorProvider(&rep, &cfg) + doctorSigning(&rep) + doctorSecretsFile(&rep) + doctorState(&rep, &cfg, now) + doctorListeners(&rep, &cfg) + doctorSchedules(&rep, &cfg, now) + return rep +} + +func envSet(name string) bool { return strings.TrimSpace(os.Getenv(name)) != "" } + +func doctorSkills(rep *doctorReport, dir string) { + reg, err := skillsapi.LoadSkills(dir) + if err != nil { + rep.add("skills", "warn", fmt.Sprintf("%s: %v", dir, err), "pass the skills directory as --skills ") + return + } + rep.add("skills", "ok", fmt.Sprintf("%s: %d skill(s)", dir, len(reg.List())), "") +} + +func doctorChannel(rep *doctorReport, cfg *config.Config) { + tokenEnv := firstNonEmpty(cfg.Telegram.TokenEnv, "DRAFTCAT_TG_TOKEN") + if envSet(tokenEnv) || envSet("DRAFTCAT_TG_TOKEN") { + rep.add("telegram token", "ok", tokenEnv+" is set", "") + } else { + rep.add("telegram token", "fail", tokenEnv+" is empty", "export "+tokenEnv+"=; the engine refuses to start without it") + } + chat := cfg.Telegram.ChatID != 0 || envSet("DRAFTCAT_TG_CHAT_ID") || len(cfg.Telegram.Security.AllowedUsers) > 0 || envSet("DRAFTCAT_TG_ALLOWED_USERS") + users := len(cfg.Telegram.Security.AllowedUsers) > 0 || envSet("DRAFTCAT_TG_ALLOWED_USERS") + switch { + case !users: + rep.add("operators", "fail", "no telegram.security.allowed_users", "set allowed_users in config or DRAFTCAT_TG_ALLOWED_USERS=; without an operator no approval can be decided") + case !chat: + rep.add("operators", "fail", "no chat_id", "set telegram.chat_id or DRAFTCAT_TG_CHAT_ID") + default: + rep.add("operators", "ok", "operator chat and allowed users configured", "") + } + if cfg.Relay.Enabled() { + secretEnv := firstNonEmpty(cfg.Relay.SecretEnv, "DRAFTCAT_RELAY_SECRET") + if envSet(secretEnv) || envSet("DRAFTCAT_RELAY_SECRET") { + rep.add("relay", "ok", fmt.Sprintf("approvals go to %s; %s is set", cfg.Relay.URL, secretEnv), "") + } else { + rep.add("relay", "fail", secretEnv+" is empty", "export the shared HMAC secret the relay was configured with") + } + } +} + +func doctorProvider(rep *doctorReport, cfg *config.Config) { + keyEnv := firstNonEmpty(cfg.Provider.APIKeyEnv, "OPENROUTER_API_KEY") + if envSet(keyEnv) || envSet("OPENROUTER_API_KEY") { + rep.add("provider key", "ok", keyEnv+" is set", "") + } else { + rep.add("provider key", "fail", keyEnv+" is empty", "export "+keyEnv+"=") + } + if cfg.GHL.APIKeyEnv != "" && !envSet(cfg.GHL.APIKeyEnv) { + rep.add("gohighlevel key", "warn", cfg.GHL.APIKeyEnv+" is empty", "GoHighLevel steps fail until it is set") + } + if cfg.Observ.OTLP.Enabled && cfg.Observ.OTLP.HeaderEnv != "" && !envSet(cfg.Observ.OTLP.HeaderEnv) { + rep.add("otlp headers", "warn", cfg.Observ.OTLP.HeaderEnv+" is empty", "the collector may reject unauthenticated spans") + } +} + +func doctorSigning(rep *doctorReport) { + if envSet("DRAFTCAT_APPROVAL_SECRET") { + rep.add("receipt signing", "ok", "DRAFTCAT_APPROVAL_SECRET is set; approval receipts are signed", "") + return + } + rep.add("receipt signing", "warn", "DRAFTCAT_APPROVAL_SECRET is empty; approval receipts are recorded unsigned", "export a long random secret so audit-verify and receipts verify can prove each decision") +} + +func doctorSecretsFile(rep *doctorReport) { + fi, err := os.Stat("secrets.yaml") + if err != nil { + return + } + if fi.Mode().Perm()&0o077 != 0 { + rep.add("secrets.yaml", "warn", fmt.Sprintf("mode %04o: readable by other users", fi.Mode().Perm()), "chmod 600 secrets.yaml") + return + } + rep.add("secrets.yaml", "ok", fmt.Sprintf("mode %04o", fi.Mode().Perm()), "") +} + +// doctorState inspects the state store without creating or migrating it. +func doctorState(rep *doctorReport, cfg *config.Config, now time.Time) { + path := strings.TrimSpace(os.Getenv("DRAFTCAT_STATE_PATH")) + if path == "" { + path = cfg.State.Path + } + if path == "" { + path = "./state.db" + } + if _, err := os.Stat(path); err != nil { + dir := filepath.Dir(path) + if fi, derr := os.Stat(dir); derr != nil || !fi.IsDir() { + rep.add("state store", "fail", fmt.Sprintf("%s does not exist and %s is not a directory", path, dir), "create the directory or set state.path / DRAFTCAT_STATE_PATH") + return + } + probe, perr := os.CreateTemp(dir, ".draftcat-doctor-*") + if perr != nil { + rep.add("state store", "fail", fmt.Sprintf("%s will be created, but %s is not writable: %v", path, dir, perr), "fix the directory permissions or point state.path at a writable volume") + return + } + name := probe.Name() + _ = probe.Close() + _ = os.Remove(name) + rep.add("state store", "ok", fmt.Sprintf("%s will be created at first start", path), "") + return + } + st, err := statestore.OpenStateStoreReadOnly(path) + if err != nil { + rep.add("state store", "fail", fmt.Sprintf("cannot open %s: %v", path, err), "check the file permissions; the engine needs read-write access") + return + } + defer func() { _ = st.Close() }() + day, err := st.BudgetDay(context.Background(), now) + if err != nil { + rep.add("state store", "warn", fmt.Sprintf("%s opened, but the usage ledger is unreadable: %v", path, err), "start the engine once to upgrade an older store") + return + } + rep.add("state store", "ok", fmt.Sprintf("%s: today %d tokens, %.4f spent", path, day.Tokens, day.Cost), "") + if day.Unsettled > 0 { + rep.add("budget", "warn", fmt.Sprintf("%d provider call(s) with unsettled usage block new model calls", day.Unsettled), "run draftcat budget status, then draftcat budget reconcile") + } + if open, err := st.OpenApprovals(); err == nil && len(open) > 0 { + rep.add("approvals", "warn", fmt.Sprintf("%d approval gate(s) still open", len(open)), "they are closed out and reported at the next engine start; see draftcat pending") + } +} + +func doctorListeners(rep *doctorReport, cfg *config.Config) { + if cfg.Webhook.Enabled { + addr := firstNonEmpty(cfg.Webhook.Addr, "127.0.0.1:8088") + if cfg.Webhook.SecretEnv == "" || !envSet(cfg.Webhook.SecretEnv) { + rep.add("webhook", "fail", "webhook.enabled but its secret env is empty", "set webhook.secret_env and export that variable; the engine refuses an unauthenticated trigger") + } else { + checkListen(rep, "webhook", addr) + } + } + if cfg.Observ.Prometheus.Enabled { + checkListen(rep, "metrics", firstNonEmpty(cfg.Observ.Prometheus.Addr, "127.0.0.1:9090")) + } +} + +func checkListen(rep *doctorReport, name, addr string) { + ln, err := (&net.ListenConfig{}).Listen(context.Background(), "tcp", addr) + if err != nil { + rep.add(name, "warn", fmt.Sprintf("%s is not free: %v", addr, err), "another process (or a running engine) holds the port") + return + } + _ = ln.Close() + rep.add(name, "ok", addr+" is free", "") +} + +func doctorSchedules(rep *doctorReport, cfg *config.Config, now time.Time) { + for _, p := range cfg.Pipelines { + if !schedule.IsTimer(p.Schedule) { + continue + } + sc, err := schedule.Parse(p.Schedule, p.Timezone) + if err != nil { + rep.add("schedule "+p.Name, "fail", err.Error(), "fix pipelines[].schedule") + continue + } + next := sc.Next(now) + if next.IsZero() { + rep.add("schedule "+p.Name, "warn", p.Schedule+" never fires", "check the day and month fields") + continue + } + detail := fmt.Sprintf("%s: next run %s", p.Schedule, next.In(sc.Location()).Format("Mon 02 Jan 2006 15:04 MST")) + if sc.IsInterval() { + detail = fmt.Sprintf("every %s (continues from the last recorded run)", sc.Interval()) + } + rep.add("schedule "+p.Name, "ok", detail, "") + } +} + +func printDoctor(w io.Writer, rep doctorReport) { + width := 0 + for _, c := range rep.Checks { + if len(c.Name) > width { + width = len(c.Name) + } + } + var sb strings.Builder + fails, warns := 0, 0 + for _, c := range rep.Checks { + sb.WriteString(fmt.Sprintf("%-4s %-*s %s\n", c.Status, width, c.Name, c.Detail)) + if c.Hint != "" && c.Status != "ok" { + sb.WriteString(fmt.Sprintf(" %-*s → %s\n", width, "", c.Hint)) + } + switch c.Status { + case "fail": + fails++ + case "warn": + warns++ + } + } + switch { + case fails > 0: + sb.WriteString(fmt.Sprintf("\n%d check(s) failed, %d warning(s). The engine will not start cleanly until the failures are fixed.\n", fails, warns)) + case warns > 0: + sb.WriteString(fmt.Sprintf("\nReady to start, with %d warning(s).\n", warns)) + default: + sb.WriteString("\nReady to start.\n") + } + _, _ = io.WriteString(w, sb.String()) +} diff --git a/doctor_cmd_test.go b/doctor_cmd_test.go new file mode 100644 index 0000000..e14443a --- /dev/null +++ b/doctor_cmd_test.go @@ -0,0 +1,161 @@ +package main + +import ( + "bytes" + "encoding/json" + "errors" + "net" + "os" + "path/filepath" + "strings" + "testing" + "time" + + statestore "github.com/renezander030/draftcat/internal/state" +) + +const doctorConfig = `telegram: + token_env: DOC_TG_TOKEN + chat_id: 42 + security: + allowed_users: [42] + max_input_length: 500 + rate_limit: 10 +provider: + type: openrouter + api_key_env: DOC_PROVIDER_KEY +models: + m: + model: x/y +roles: + drafter: m +state: + path: STATE +pipelines: + - name: morning + schedule: "0 8 * * 1-5" + timezone: Europe/Berlin + steps: + - name: fetch + type: deterministic +` + +func writeDoctorConfig(t *testing.T, statePath string) string { + t.Helper() + dir := t.TempDir() + p := filepath.Join(dir, "config.yaml") + if err := os.WriteFile(p, []byte(strings.Replace(doctorConfig, "STATE", statePath, 1)), 0o600); err != nil { + t.Fatal(err) + } + return p +} + +func checkByName(rep doctorReport, name string) doctorCheck { + for _, c := range rep.Checks { + if c.Name == name { + return c + } + } + return doctorCheck{} +} + +func TestDoctorReadyConfig(t *testing.T) { + t.Setenv("DOC_TG_TOKEN", "123:abc") + t.Setenv("DOC_PROVIDER_KEY", "sk-test") + t.Setenv("DRAFTCAT_APPROVAL_SECRET", "approval-secret") + t.Setenv("DRAFTCAT_STATE_PATH", "") + statePath := filepath.Join(t.TempDir(), "state.db") + cfgPath := writeDoctorConfig(t, statePath) + now := time.Date(2026, 10, 9, 9, 0, 0, 0, time.UTC) // Friday + rep := doctor(cfgPath, t.TempDir(), now) + if !rep.OK { + t.Fatalf("ready config reported failures: %+v", rep.Checks) + } + if c := checkByName(rep, "state store"); c.Status != "ok" || !strings.Contains(c.Detail, "will be created") { + t.Errorf("state store: %+v", c) + } + if _, err := os.Stat(statePath); !errors.Is(err, os.ErrNotExist) { + t.Errorf("doctor created the state store: %v", err) + } + if c := checkByName(rep, "schedule morning"); c.Status != "ok" || !strings.Contains(c.Detail, "Mon 12 Oct 2026 08:00 CEST") { + t.Errorf("schedule: %+v", c) + } +} + +func TestDoctorReportsMissingCredentials(t *testing.T) { + t.Setenv("DOC_TG_TOKEN", "") + t.Setenv("DRAFTCAT_TG_TOKEN", "") + t.Setenv("DOC_PROVIDER_KEY", "") + t.Setenv("OPENROUTER_API_KEY", "") + t.Setenv("DRAFTCAT_APPROVAL_SECRET", "") + t.Setenv("DRAFTCAT_STATE_PATH", "") + cfgPath := writeDoctorConfig(t, filepath.Join(t.TempDir(), "state.db")) + rep := doctor(cfgPath, t.TempDir(), time.Now()) + if rep.OK { + t.Fatal("missing credentials reported OK") + } + for name, want := range map[string]string{"telegram token": "fail", "provider key": "fail", "receipt signing": "warn"} { + if c := checkByName(rep, name); c.Status != want || c.Hint == "" { + t.Errorf("%s: %+v, want %s with a hint", name, c, want) + } + } + var out bytes.Buffer + if code := runDoctor([]string{"--config", cfgPath, "--skills", t.TempDir()}, &out); code != 1 { + t.Errorf("exit = %d, want 1", code) + } + if !strings.Contains(out.String(), "check(s) failed") { + t.Errorf("summary missing:\n%s", out.String()) + } +} + +func TestDoctorInspectsExistingStoreReadOnly(t *testing.T) { + t.Setenv("DOC_TG_TOKEN", "123:abc") + t.Setenv("DOC_PROVIDER_KEY", "sk-test") + t.Setenv("DRAFTCAT_STATE_PATH", "") + statePath := filepath.Join(t.TempDir(), "state.db") + st, err := statestore.OpenStateStore(statePath) + if err != nil { + t.Fatal(err) + } + if err := st.AddBudgetUsage(t.Context(), time.Now(), 120, 0.25, 0, 0); err != nil { + t.Fatal(err) + } + st.Close() + before, _ := os.Stat(statePath) + rep := doctor(writeDoctorConfig(t, statePath), t.TempDir(), time.Now()) + if c := checkByName(rep, "state store"); c.Status != "ok" || !strings.Contains(c.Detail, "120 tokens") { + t.Errorf("state store: %+v", c) + } + after, _ := os.Stat(statePath) + if !before.ModTime().Equal(after.ModTime()) || before.Size() != after.Size() { + t.Error("doctor modified the state store") + } +} + +func TestDoctorJSONAndBadConfig(t *testing.T) { + var out bytes.Buffer + missing := filepath.Join(t.TempDir(), "nope.yaml") + if code := runDoctor([]string{"--config", missing, "--json"}, &out); code != 1 { + t.Fatalf("exit = %d", code) + } + var rep doctorReport + if err := json.Unmarshal(out.Bytes(), &rep); err != nil { + t.Fatalf("not JSON: %v\n%s", err, out.String()) + } + if rep.OK || len(rep.Checks) != 1 || rep.Checks[0].Name != "config" || rep.Checks[0].Status != "fail" { + t.Fatalf("report: %+v", rep) + } +} + +func TestDoctorFlagsBusyPort(t *testing.T) { + ln, err := (&net.ListenConfig{}).Listen(t.Context(), "tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer ln.Close() + rep := doctorReport{} + checkListen(&rep, "webhook", ln.Addr().String()) + if rep.Checks[0].Status != "warn" { + t.Errorf("busy port: %+v", rep.Checks[0]) + } +} diff --git a/engine_lifecycle.go b/engine_lifecycle.go new file mode 100644 index 0000000..b612840 --- /dev/null +++ b/engine_lifecycle.go @@ -0,0 +1,189 @@ +package main + +import ( + "context" + "log" + "net/http" + "os" + "strings" + "time" + + "github.com/renezander030/draftcat/internal/config" + "github.com/renezander030/draftcat/internal/obs" + "github.com/renezander030/draftcat/internal/redact" + statestore "github.com/renezander030/draftcat/internal/state" +) + +const defaultShutdownGrace = 30 * time.Second + +// secretEnvNames lists every environment variable the engine reads a +// credential from, labelled by the name the operator configured. +func secretEnvNames(cfg *config.Config) []string { + names := []string{ + firstNonEmpty(cfg.Telegram.TokenEnv, "DRAFTCAT_TG_TOKEN"), + firstNonEmpty(cfg.Relay.SecretEnv, "DRAFTCAT_RELAY_SECRET"), + firstNonEmpty(cfg.Provider.APIKeyEnv, "OPENROUTER_API_KEY"), + "DRAFTCAT_APPROVAL_SECRET", + } + for _, n := range []string{cfg.Webhook.SecretEnv, cfg.GHL.APIKeyEnv, cfg.Observ.OTLP.HeaderEnv} { + if n != "" { + names = append(names, n) + } + } + return names +} + +// registerSecrets registers the value of every credential variable for +// redaction. OTLP header variables hold "k=v,k2=v2"; each value is +// registered on its own as well as the whole string. +func registerSecrets(cfg *config.Config) { + for _, name := range secretEnvNames(cfg) { + v := os.Getenv(name) + redact.Register(name, v) + if name == cfg.Observ.OTLP.HeaderEnv { + for _, hv := range obs.ParseHeaderEnv(v) { + redact.Register(name, hv) + // "Bearer " — the token alone may also be echoed. + if i := strings.IndexByte(hv, ' '); i > 0 { + redact.Register(name, hv[i+1:]) + } + } + } + } + // Secrets already resolved into the config, in case they came from a + // fallback variable. + redact.Register(firstNonEmpty(cfg.Telegram.TokenEnv, "DRAFTCAT_TG_TOKEN"), cfg.Telegram.Token()) + redact.Register(firstNonEmpty(cfg.Relay.SecretEnv, "DRAFTCAT_RELAY_SECRET"), cfg.Relay.Secret()) + redact.Register(firstNonEmpty(cfg.Provider.APIKeyEnv, "OPENROUTER_API_KEY"), cfg.Provider.APIKey()) +} + +func firstNonEmpty(vals ...string) string { + for _, v := range vals { + if strings.TrimSpace(v) != "" { + return v + } + } + return "" +} + +// restoreSchedules continues every pipeline's schedule and failure streak +// from the run history, and tells the operator about pipelines that start +// paused because their recorded streak reached pause_after_failures. +func restoreSchedules(sched *Scheduler, cfg *config.Config, store *statestore.StateStore, ch interface{ Send(string) error }) { + if store == nil { + return + } + for _, p := range cfg.Pipelines { + n := p.PauseAfterFailures + if n < 1 { + n = 1 + } + runs, err := store.RecentRuns(p.Name, n) + if err != nil { + log.Printf("[scheduler] run history for %s unavailable: %v", p.Name, err) + continue + } + if sched.Restore(p.Name, runs) { + log.Printf("[scheduler] %s starts paused: last %d runs failed", p.Name, len(runs)) + if ch != nil { + msg := autoPauseNotice(p.Name, len(runs), nil) + if len(runs) > 0 && runs[0].Error != "" { + msg = autoPauseNotice(p.Name, len(runs), errText(runs[0].Error)) + } + _ = ch.Send(msg) + } + } + } +} + +type errText string + +func (e errText) Error() string { return string(e) } + +// shutdownGrace returns timeouts.shutdown_grace, default 30s. +func shutdownGrace(cfg *config.Config) time.Duration { + if v := strings.TrimSpace(cfg.Timeouts.ShutdownGrace); v != "" { + if d, err := time.ParseDuration(v); err == nil && d >= 0 { + return d + } + } + return defaultShutdownGrace +} + +// drainEngine stops admitting runs, closes the webhook listener after its +// in-flight requests, and waits up to grace for admitted pipeline runs. +// Approval taps keep arriving while it waits, so a run blocked on a decision +// can still finish. An approval gate still open at exit is reconciled at the +// next start. +func drainEngine(sched *Scheduler, webhook *http.Server, grace time.Duration) bool { + sched.BeginDrain() + if webhook != nil { + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + if err := webhook.Shutdown(ctx); err != nil { + log.Printf("[shutdown] webhook listener: %v", err) + } + cancel() + } + running := sched.RunningNames() + if len(running) == 0 { + return true + } + log.Printf("[shutdown] waiting up to %s for %d running pipeline(s): %s", grace, len(running), strings.Join(running, ", ")) + if sched.Wait(grace) { + log.Printf("[shutdown] all pipeline runs finished") + return true + } + log.Printf("[shutdown] grace period elapsed with %s still running; open approval gates are reconciled at next start", strings.Join(sched.RunningNames(), ", ")) + return false +} + +// redactedError keeps the wrapped error for errors.Is/As while its message +// carries no registered secret. +type redactedError struct{ err error } + +func (e redactedError) Error() string { return redact.String(e.err.Error()) } +func (e redactedError) Unwrap() error { return e.err } + +// redactErr returns err with registered secrets removed from its message. +func redactErr(err error) error { + if err == nil { + return nil + } + return redactedError{err} +} + +// engineGauges reports current engine state for /metrics: runs in progress, +// paused pipelines, open approval gates and today's usage against the caps. +func engineGauges(sched *Scheduler, budget *BudgetTracker, cfg *config.Config, store *statestore.StateStore) []obs.Gauge { + var out []obs.Gauge + running, _ := sched.Counts() + out = append(out, obs.Gauge{Name: "draftcat_pipelines_running", Help: "Pipeline runs in progress.", Value: float64(running)}) + for _, ps := range sched.GetAll() { + v := 0.0 + if ps.Paused { + v = 1 + } + out = append(out, obs.Gauge{Name: "draftcat_pipeline_paused", Help: "1 when the pipeline's timer is paused.", Labels: []string{"pipeline", ps.Name}, Value: v}) + out = append(out, obs.Gauge{Name: "draftcat_pipeline_consecutive_failures", Help: "Consecutive failed runs of the pipeline.", Labels: []string{"pipeline", ps.Name}, Value: float64(ps.Failures)}) + } + if store != nil { + if open, err := store.OpenApprovals(); err == nil { + out = append(out, obs.Gauge{Name: "draftcat_approvals_open", Help: "Approval gates waiting on a decision.", Value: float64(len(open))}) + } + } + if budget != nil { + snap := budget.snapshot() + out = append(out, + obs.Gauge{Name: "draftcat_budget_day_tokens", Help: "Tokens used today (UTC).", Value: float64(snap.tokensUsedToday)}, + obs.Gauge{Name: "draftcat_budget_day_cost", Help: "Cost spent today (UTC), in the unit of the model rates.", Value: snap.costToday}, + obs.Gauge{Name: "draftcat_budget_unsettled_calls", Help: "Provider calls whose usage is not settled yet.", Value: float64(snap.unsettled)}, + ) + if cfg.Budgets.PerDayTokens > 0 { + out = append(out, obs.Gauge{Name: "draftcat_budget_day_tokens_limit", Help: "Daily token cap.", Value: float64(cfg.Budgets.PerDayTokens)}) + } + if cfg.Budgets.PerDayCost > 0 { + out = append(out, obs.Gauge{Name: "draftcat_budget_day_cost_limit", Help: "Daily cost cap.", Value: cfg.Budgets.PerDayCost}) + } + } + return out +} diff --git a/internal/config/config.go b/internal/config/config.go index 9a63e6c..eb36212 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -271,12 +271,21 @@ type BudgetConfig struct { // step. Pair with per_step_tokens to bound that overshoot. PerDayCost float64 `yaml:"per_day_cost"` PerPipelineCost float64 `yaml:"per_pipeline_cost"` + // AlertAt lists fractions of the daily caps (per_day_tokens, per_day_cost) + // at which the operator channel is notified, e.g. [0.5, 0.8, 0.95]. Each + // threshold notifies once per UTC day and cap, also across restarts. + // Empty = no early warnings. + AlertAt []float64 `yaml:"alert_at"` } type TimeoutConfig struct { AICall string `yaml:"ai_call"` OperatorApproval string `yaml:"operator_approval"` PipelineTotal string `yaml:"pipeline_total"` + // ShutdownGrace is how long SIGINT/SIGTERM waits for running pipelines to + // finish before the engine exits. New runs are refused during the wait. + // Default 30s; "0s" exits at once. + ShutdownGrace string `yaml:"shutdown_grace"` } // WebhookConfig enables the opt-in HTTP trigger server. A pipeline with @@ -347,9 +356,22 @@ type OTLPConfig struct { } type PipelineConfig struct { - Name string `yaml:"name"` - Schedule string `yaml:"schedule"` - Steps []StepConfig `yaml:"steps"` + Name string `yaml:"name"` + // Schedule is "manual", "webhook", an interval ("30m", "24h") or a + // calendar schedule: a five-field cron expression ("0 8 * * 1-5") or + // @hourly/@daily/@weekly/@monthly. See internal/schedule. + Schedule string `yaml:"schedule"` + // Timezone is the IANA zone calendar schedules run in ("Europe/Berlin"). + // Empty = the process's local zone. Ignored for intervals. + Timezone string `yaml:"timezone"` + // CatchUp runs a calendar slot that passed while the engine was down once + // at startup. Intervals always continue from the last recorded run. + CatchUp bool `yaml:"catch_up"` + // PauseAfterFailures pauses a timer-scheduled pipeline after this many + // consecutive failed runs and tells the operator once. 0 = never pause. + // /cron resume clears the count. + PauseAfterFailures int `yaml:"pause_after_failures"` + Steps []StepConfig `yaml:"steps"` } type StepConfig struct { diff --git a/internal/obs/gauges.go b/internal/obs/gauges.go new file mode 100644 index 0000000..99096e5 --- /dev/null +++ b/internal/obs/gauges.go @@ -0,0 +1,64 @@ +package obs + +import ( + "sort" + "strings" + "sync" +) + +// Gauge is one current-value sample rendered on /metrics. Labels are +// alternating key, value pairs. +type Gauge struct { + Name string + Help string + Labels []string + Value float64 +} + +var ( + gaugeMu sync.Mutex + gaugeSrc func() []Gauge +) + +// SetGaugeSource registers the function that reports current state (running +// pipelines, open approvals, today's spend) at scrape time. It is called on +// every scrape, outside the metric registry lock; nil removes it. +func SetGaugeSource(f func() []Gauge) { + gaugeMu.Lock() + defer gaugeMu.Unlock() + gaugeSrc = f +} + +func renderGauges(sb *strings.Builder) { + gaugeMu.Lock() + src := gaugeSrc + gaugeMu.Unlock() + if src == nil { + return + } + samples := src() + byName := map[string][]Gauge{} + help := map[string]string{} + var names []string + for _, g := range samples { + if _, ok := byName[g.Name]; !ok { + names = append(names, g.Name) + } + byName[g.Name] = append(byName[g.Name], g) + if g.Help != "" { + help[g.Name] = g.Help + } + } + sort.Strings(names) + for _, name := range names { + sb.WriteString("# HELP " + name + " " + help[name] + "\n") + sb.WriteString("# TYPE " + name + " gauge\n") + for _, g := range byName[name] { + sb.WriteString(name) + if l := lbls(g.Labels...); l != "" { + sb.WriteString("{" + l + "}") + } + sb.WriteString(" " + formatNum(g.Value) + "\n") + } + } +} diff --git a/internal/obs/gauges_test.go b/internal/obs/gauges_test.go new file mode 100644 index 0000000..e5aaa74 --- /dev/null +++ b/internal/obs/gauges_test.go @@ -0,0 +1,50 @@ +package obs + +import ( + "strings" + "testing" +) + +func TestGaugesRenderedAtScrape(t *testing.T) { + resetMetrics() + defer resetMetrics() + calls := 0 + SetGaugeSource(func() []Gauge { + calls++ + return []Gauge{ + {Name: "draftcat_pipeline_paused", Help: "1 when paused.", Labels: []string{"pipeline", "b"}, Value: 1}, + {Name: "draftcat_approvals_open", Help: "Open gates.", Value: 2}, + {Name: "draftcat_pipeline_paused", Labels: []string{"pipeline", "a\"x"}, Value: 0}, + } + }) + var sb strings.Builder + if err := WriteMetrics(&sb); err != nil { + t.Fatal(err) + } + out := sb.String() + for _, want := range []string{ + "# TYPE draftcat_approvals_open gauge\ndraftcat_approvals_open 2\n", + "# HELP draftcat_pipeline_paused 1 when paused.\n# TYPE draftcat_pipeline_paused gauge\n", + `draftcat_pipeline_paused{pipeline="b"} 1`, + `draftcat_pipeline_paused{pipeline="a\"x"} 0`, + } { + if !strings.Contains(out, want) { + t.Errorf("missing %q in:\n%s", want, out) + } + } + if strings.Count(out, "# TYPE draftcat_pipeline_paused gauge") != 1 { + t.Errorf("gauge family rendered twice:\n%s", out) + } + if calls != 1 { + t.Errorf("source called %d times", calls) + } +} + +func TestNoGaugesWithoutSource(t *testing.T) { + resetMetrics() + var sb strings.Builder + _ = WriteMetrics(&sb) + if strings.Contains(sb.String(), " gauge\n") { + t.Errorf("gauge rendered without a source:\n%s", sb.String()) + } +} diff --git a/internal/obs/metrics.go b/internal/obs/metrics.go index 33069af..3382a8a 100644 --- a/internal/obs/metrics.go +++ b/internal/obs/metrics.go @@ -128,6 +128,7 @@ func resetMetrics() { promOn = false otlpOn = false mu.Unlock() + SetGaugeSource(nil) fresh := newRegistry() reg.mu.Lock() defer reg.mu.Unlock() @@ -192,6 +193,7 @@ func WriteMetrics(w io.Writer) error { reg.aiCost.render(&sb) reg.approvals.render(&sb) reg.mu.Unlock() + renderGauges(&sb) _, err := io.WriteString(w, sb.String()) return err } diff --git a/internal/obs/obs.go b/internal/obs/obs.go index 24ab066..8636639 100644 --- a/internal/obs/obs.go +++ b/internal/obs/obs.go @@ -15,6 +15,8 @@ import ( "os" "sync" "time" + + "github.com/renezander030/draftcat/internal/redact" ) var ( @@ -155,5 +157,5 @@ func write(rec map[string]interface{}) { } mu.Lock() defer mu.Unlock() - _, _ = out.Write(append(b, '\n')) + _, _ = out.Write(append(redact.Bytes(b), '\n')) } diff --git a/internal/obs/otlp.go b/internal/obs/otlp.go index 0b6fe5f..0d84f98 100644 --- a/internal/obs/otlp.go +++ b/internal/obs/otlp.go @@ -18,6 +18,8 @@ import ( "strconv" "strings" "time" + + "github.com/renezander030/draftcat/internal/redact" ) // otlpSpan is the engine-neutral span record handed to the exporter. @@ -99,6 +101,7 @@ func (e *otlpExporter) export(ctx context.Context, spans []otlpSpan) error { if err != nil { return err } + body = redact.Bytes(body) req, err := http.NewRequestWithContext(ctx, http.MethodPost, e.endpoint, bytes.NewReader(body)) if err != nil { return err diff --git a/internal/obs/redact_test.go b/internal/obs/redact_test.go new file mode 100644 index 0000000..9bfc4c9 --- /dev/null +++ b/internal/obs/redact_test.go @@ -0,0 +1,27 @@ +package obs + +import ( + "bytes" + "strings" + "testing" + "time" + + "github.com/renezander030/draftcat/internal/redact" +) + +func TestSpansAreRedacted(t *testing.T) { + redact.Reset() + defer redact.Reset() + redact.Register("OPENROUTER_API_KEY", "sk-or-v1-spanleak0001") + var buf bytes.Buffer + Enable(&buf) + defer reset() + Pipeline("p").End("error", map[string]interface{}{"error": "401 for key sk-or-v1-spanleak0001"}) + EmitStep("p", "s", "ai", time.Now(), "error", map[string]interface{}{"error": "sk-or-v1-spanleak0001"}) + if strings.Contains(buf.String(), "sk-or-v1-spanleak0001") { + t.Fatalf("span leaked a secret: %s", buf.String()) + } + if strings.Count(buf.String(), "[REDACTED:OPENROUTER_API_KEY]") != 2 { + t.Fatalf("spans not labelled: %s", buf.String()) + } +} diff --git a/internal/redact/redact.go b/internal/redact/redact.go new file mode 100644 index 0000000..e6de2be --- /dev/null +++ b/internal/redact/redact.go @@ -0,0 +1,101 @@ +// Package redact keeps the engine's own credentials out of everything it +// writes for people and machines to read: the process log, operator +// notifications, persisted run errors and exported spans. +// +// The engine knows its secrets exactly — each one is read from a named +// environment variable — so redaction is a literal replace of registered +// values, not a guess at what a credential looks like. A value is registered +// once at startup and every output path runs through String or Writer. +package redact + +import ( + "io" + "sort" + "strings" + "sync" +) + +// MinLen is the shortest value that is registered. Shorter strings are too +// likely to occur in ordinary text, and a replace would mangle it while +// protecting nothing worth protecting. +const MinLen = 8 + +type entry struct { + value string + label string +} + +var ( + mu sync.RWMutex + entries []entry +) + +// Register records a secret value under a label (usually the environment +// variable it came from). It is replaced by "[REDACTED: