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
20 changes: 20 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -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:<VARIABLE>]` 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.
Expand Down
8 changes: 6 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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}
Expand Down Expand Up @@ -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 <pipeline> # dry-run against fixtures/<pipeline>/ (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)
Expand Down Expand Up @@ -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}
Expand Down
4 changes: 4 additions & 0 deletions budget.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ type BudgetTracker struct {
dayCostLimit float64
pipelineCostLimit float64
pipelineTokenLimit int
alerts *budgetAlerts
}

func (b *BudgetTracker) root() *BudgetTracker {
Expand Down Expand Up @@ -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) {
Expand Down
103 changes: 103 additions & 0 deletions budget_alerts.go
Original file line number Diff line number Diff line change
@@ -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) }
130 changes: 130 additions & 0 deletions budget_alerts_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
}
6 changes: 5 additions & 1 deletion config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading