diff --git a/README.md b/README.md index 9dd1631..148490c 100644 --- a/README.md +++ b/README.md @@ -101,6 +101,12 @@ The preflight is read-only and reports the export scope before any snapshot is w ## Automation +For public server-size observations without a bot token, use +[`metrics collect`](docs/commands/metrics.md) with a separate metrics database. +It records approximate members and online presence; online presence does not +measure active posters. `metrics import` preserves historical observations and +`metrics status` reports their freshness. + Discrawl exposes stable JSON for launchers, agents, and CI: ```bash diff --git a/docs/README.md b/docs/README.md index 847855a..3fef445 100644 --- a/docs/README.md +++ b/docs/README.md @@ -20,6 +20,7 @@ Mirror Discord guilds into local SQLite. Search server history without depending - **Need DM search?** [`wiretap`](commands/wiretap.html) imports local Discord Desktop cache. - **Need multilingual lexical search?** Configure [language analyzers](guides/search-modes.html) and run [`lexical rebuild`](commands/lexical.html). - **Want semantic search?** Configure [Embeddings](guides/embeddings.html), then run [`embed`](commands/embed.html). +- **Need public server-size history without a bot?** [`metrics`](commands/metrics.html) records approximate membership and online presence in a separate database. - **Wiring an agent or launcher?** `discrawl metadata --json`, `discrawl status --json`, `discrawl diagnostics --json`, `discrawl coverage --json`, `discrawl failures --json`, `discrawl remote status`, and `discrawl doctor --json` expose the read-only crawlkit control surface. ## At a glance diff --git a/docs/commands/metrics.md b/docs/commands/metrics.md new file mode 100644 index 0000000..a77119d --- /dev/null +++ b/docs/commands/metrics.md @@ -0,0 +1,187 @@ +# `metrics` + +Record public Discord server size and online-presence observations without a bot +token. Metrics use an explicitly configured, separate SQLite database; these +commands do not open the message/member archive, import Desktop data, or generate +embeddings. + +## Usage + +```bash +discrawl metrics collect --config /absolute/path/metrics.json +discrawl metrics import --config /absolute/path/metrics.json < history.ndjson +discrawl metrics status --config /absolute/path/metrics.json +discrawl help metrics +``` + +The metrics `--config` is a JSON file supplied after the subcommand. It is separate +from Discrawl's normal TOML archive configuration. All three commands output JSON. +`status` is read-only and reports observation/event counts, the latest observation +sequence, and the most recent observation time (`null` for an empty store). + +## Configuration + +```json +{ + "database": "/absolute/path/discrawl-metrics.sqlite", + "targets": [ + {"entity": "openclaw", "target": "clawd"} + ] +} +``` + +`entity` is your series label; `target` is a public invite code, not an invite URL +or guild ID. Codes must be unique in the configuration. Other servers can be +configured with their own invite codes. + +The database path must be absolute. `collect` and `import` create a new metrics +database if that path does not exist. Existing databases must identify their +owner as `discrawl` and metrics version as `1`; unrelated, empty, or newer-version +databases are refused before write access. `status` never creates a database. + +No token or cookie is needed. The shared configuration fields `cookieJar` and +`tokenEnv` are accepted but unused by this collector. + +## Standalone metrics runtime + +Metrics can use a dedicated, versioned copy of the normal Discrawl executable. +Keep its path distinct from the executable used by an existing message collector. +For example, a private per-user installation can use: + +```text +~/.local/libexec/discrawl-metrics//discrawl +``` + +Keep the runtime directories and executable private to the owning user (`0700` +on Unix). Copy the validated artifact into a new version directory, preserve its +bytes and signature, then verify the installed SHA-256 against the artifact +receipt. On macOS, also verify the copied executable against the approved signing +identity and designated requirement. Record the source commit, hash, and +verification result with the installation. A locally signed build remains a +local build; it is not an official notarized release. + +Use the full versioned executable path and explicit metrics JSON configuration: + +```bash +metrics_runtime="$HOME/.local/libexec/discrawl-metrics//discrawl" +"$metrics_runtime" metrics status --config "$HOME/.config/discrawl/metrics.json" +``` + +Replace `` with the installed build's commit. Read-only `status` +verifies access to the existing metrics store before cutover. Installing this +runtime does not register schedules or run collection/imports. Those operations +belong to the metrics deployment owner, which selects the exact executable path +for each job. Existing message-collector binaries, command symlinks, settings, +and processes remain independently managed. + +### Hourly scheduling on macOS + +A user LaunchAgent can invoke the versioned executable directly. Keep the plist +private (`0600`) under `~/Library/LaunchAgents`, with private (`0700`) log +directories and pre-created log files (`0600`). Use absolute executable, +configuration, and log paths; plist strings do not expand `$HOME` or `~`. + +For an hourly job at minute 4, configure these launchd keys: + +| Key | Value | +| --- | --- | +| `Label` | A dedicated label, such as `org.example.discrawl-metrics` | +| `ProgramArguments` | Versioned executable, `metrics`, `collect`, `--config`, absolute metrics config path | +| `StartCalendarInterval` | Dictionary with integer `Minute` set to `4` | +| `RunAtLoad` | `true` for one collection when the job is bootstrapped | +| `KeepAlive` | `false`; provider failures must not trigger a restart loop | +| `Umask` | Integer `63` (octal `077`) | +| `StandardOutPath`, `StandardErrorPath` | Absolute paths to the private log files | + +Set `PATH=/usr/bin:/bin:/usr/sbin:/sbin` in `EnvironmentVariables`. Include +`DISCRAWL_NO_AUTO_UPDATE=1` and `DISCRAWL_NO_UPDATE_CHECK=1`; metrics commands already +bypass archive update hooks, and the job should retain its verified executable. +Do not put tokens or cookies in the plist. Public invite metrics need neither. + +Validate the plist with `plutil -lint`, then bootstrap it in the owning user's +GUI domain. `RunAtLoad` supplies the first run, so a second kickstart is not +needed. Verify `launchctl print` retains `Minute = 4`, the expected arguments, +and a successful exit. Confirm a current `metric_runs` row with `status = 'ok'` +and non-NULL observations for every configured target's required metrics. + +Launchd runs one instance of a label at a time. Native SQLite transactions +serialize database writes across independent connections; a competing write +must wait or fail without inserting a partial batch. This is transaction-level +write exclusion, not a mutex covering the preceding HTTP requests. Avoid adding +another scheduler for the same metrics job. + +If macOS requests access to the configured volume, report the permission gate +before claiming successful collection. Do not change identities, broaden access, +or repeatedly restart a blocked job. Database ownership and private permissions +remain required for scheduled operation. + +After owner-granted consent, recheck the original invocation read-only. Acceptance +requires exit code `0`, a corresponding successful `metric_runs` row, and +non-NULL required measurements. The hourly calendar definition must remain +loaded. A completed batch job normally shows `state = not running` while it +waits for its next scheduled time. Leave the accepted job loaded; verifying this +state does not require another collection or kickstart. + +## What is measured + +Each collection makes one public `GET /api/v10/invites/{code}?with_counts=true` +request per target and records two observations: + +| Metric | Discord field | Meaning | +| --- | --- | --- | +| `members` | `approximate_member_count` | Approximate total server membership | +| `online` | `approximate_presence_count` | Approximate online presence | + +Both use `kind: "counter"` and `provenance: "discord_invite_approximate"`. +Here, counter means a timestamped snapshot: either value can decrease. Online +presence does **not** count active posters, messages, daily active users, or +engagement. See Discord's [Invite object and Get Invite documentation](https://docs.discord.com/developers/resources/invite#get-invite). + +Zero is a valid count. Missing or invalid values, expired invites, HTTP failures, +and rate limits produce SQL `NULL`, never a fabricated zero. Valid fields and +other successful targets are retained. Such a run records `partial`, outputs +`"ok": false`, and exits nonzero. Responses are bounded in size, requests have a +30-second timeout, and redirects are refused. There is no automatic retry loop +or scheduler; invoke `collect` at the cadence appropriate for your application. + +## Importing history + +`import` reads one JSON object per line from stdin. Every row needs a stable, +nonempty `id` and an exact configured `entity`/`target` pair. Times use RFC 3339 +(fractional seconds and offsets are accepted). Unknown fields are rejected. + +```json +{"type":"metric","id":"historical-members-1","entity":"openclaw","target":"clawd","metric":"members","kind":"counter","ts":"2026-09-14T12:00:00Z","value":1200,"observed_at":"2026-09-14T12:00:00Z","provenance":"historical-import"} +{"type":"metric","id":"historical-online-1","entity":"openclaw","target":"clawd","metric":"online","kind":"counter","ts":"2026-09-14T12:00:00Z","value":null,"observed_at":"2026-09-14T12:00:00Z","provenance":"historical-import"} +{"type":"event","id":"historical-event-1","entity":"openclaw","target":"clawd","kind":"milestone","ts":"2026-09-14T12:00:00Z","label":"Example milestone","url":"https://example.org/milestone","observed_at":"2026-09-14T12:00:00Z","provenance":"historical-import"} +``` + +Metric rows use `kind: "counter"` or `"daily"` and a nonnegative numeric value or +`null`. Daily rows describe a completed UTC day. Revisions require different IDs: +consumers select the latest sequence for a daily series/day instead of summing +revisions. Import preserves supplied timestamps, provenance, repeated values, +decreases, zeroes, NULLs, and events. Only an already stored ID is ignored. + +Imports commit in batches of 500. If a later line is invalid, the JSON result +reports the number already committed and the command exits nonzero. Correct the +input and replay it: stored IDs make retries idempotent. A line is limited to 4 MiB. + +## Storage and delivery + +The metrics database contains: + +- `metric_meta(key, value)`: collector owner and schema version. +- `metric_observations(sequence, id, entity, target, metric, kind, ts, value, observed_at, provenance)`. +- `metric_events(sequence, id, entity, target, kind, ts, label, url, observed_at, provenance)`. +- `metric_runs(sequence, ts, status, rows_written)`: collection attempts. + +Each table's sequence advances independently; gaps are allowed. Consumers can +open this database read-only and keep separate delivery cursors for observations +and events. Import IDs are preserved; native collection derives stable IDs from +the complete observation, including observation time. Repeated collections keep +their own observations even when counts do not change. + +## See also + +- [`analytics`](analytics.html) for activity calculated from archived messages. +- [Data storage](../guides/data-storage.html). diff --git a/internal/cli/cli.go b/internal/cli/cli.go index 147d65d..a2c02ee 100644 --- a/internal/cli/cli.go +++ b/internal/cli/cli.go @@ -15,6 +15,7 @@ import ( "github.com/openclaw/crawlkit/embed" "github.com/openclaw/discrawl/internal/config" "github.com/openclaw/discrawl/internal/discord" + "github.com/openclaw/discrawl/internal/headlinemetrics" "github.com/openclaw/discrawl/internal/share" "github.com/openclaw/discrawl/internal/store" "github.com/openclaw/discrawl/internal/syncer" @@ -73,6 +74,11 @@ func Run(ctx context.Context, args []string, stdout, stderr io.Writer) (runErr e _, _ = io.WriteString(stdout, currentVersion()+"\n") return nil } + if rest[0] == "metrics" { + // Metrics have their own explicit config and database; do not initialize + // archive services, resolve credentials, or run archive update checks. + return headlinemetrics.Run(ctx, rest[1:], "discrawl", headlinemetrics.CollectDiscord, os.Stdin, stdout, stderr) + } level := slog.LevelInfo if global.Quiet { level = slog.LevelError @@ -133,6 +139,7 @@ var discrawlCommandSpecs = []discrawlCommandSpec{ {name: "messages", description: "List archived messages."}, {name: "digest", description: "Summarize recent archive activity."}, {name: "analytics", description: "Analyze archive activity and trends."}, + {name: "metrics", description: "Collect public invite counts in a separate metrics database."}, {name: "dms", description: "List local Discord Desktop conversations."}, {name: "mentions", description: "List archived mentions."}, {name: "attachments", description: "List or fetch archived attachments."}, diff --git a/internal/cli/metrics_test.go b/internal/cli/metrics_test.go new file mode 100644 index 0000000..57cdd44 --- /dev/null +++ b/internal/cli/metrics_test.go @@ -0,0 +1,38 @@ +package cli + +import ( + "bytes" + "encoding/json" + "os" + "path/filepath" + "testing" + + "github.com/openclaw/discrawl/internal/headlinemetrics" + "github.com/stretchr/testify/require" +) + +func TestMetricsHelpAndIndependentStatus(t *testing.T) { + for _, args := range [][]string{{"help", "metrics"}, {"help", "metrics", "collect"}, {"metrics"}, {"metrics", "collect", "--help"}, {"--json", "metrics", "status", "--help"}} { + var out bytes.Buffer + require.NoError(t, Run(t.Context(), args, &out, &bytes.Buffer{})) + require.Contains(t, out.String(), "Usage: discrawl metrics") + } + var out bytes.Buffer + require.NoError(t, Run(t.Context(), []string{"--help"}, &out, &bytes.Buffer{})) + require.Contains(t, out.String(), "metrics") + root := t.TempDir() + database := filepath.Join(root, "metrics.sqlite") + s, err := headlinemetrics.Open(t.Context(), database, "discrawl") + require.NoError(t, err) + require.NoError(t, s.Close()) + config := filepath.Join(root, "metrics.json") + b, err := json.Marshal(headlinemetrics.Config{Database: database, Targets: []headlinemetrics.Target{{Entity: "openclaw", Target: "clawd"}}}) + require.NoError(t, err) + require.NoError(t, os.WriteFile(config, b, 0o600)) + out.Reset() + // The invalid archive config must never be opened or initialized by metrics. + archiveConfig := filepath.Join(root, "absent-archive-config.toml") + require.NoError(t, Run(t.Context(), []string{"--config", archiveConfig, "--json", "metrics", "status", "--config", config}, &out, &bytes.Buffer{})) + require.JSONEq(t, `{"source":"discrawl","observations":0,"events":0,"sequence":0,"last_observed":null}`, out.String()) + require.NoFileExists(t, archiveConfig) +} diff --git a/internal/cli/output.go b/internal/cli/output.go index 0402f03..f14955d 100644 --- a/internal/cli/output.go +++ b/internal/cli/output.go @@ -12,6 +12,7 @@ import ( "time" "github.com/openclaw/discrawl/internal/discorddesktop" + "github.com/openclaw/discrawl/internal/headlinemetrics" "github.com/openclaw/discrawl/internal/media" "github.com/openclaw/discrawl/internal/report" "github.com/openclaw/discrawl/internal/share" @@ -113,6 +114,10 @@ func printCommandUsage(w io.Writer, args []string) error { } var commandUsage = map[string]string{ + "metrics": headlinemetrics.Usage, + "metrics collect": headlinemetrics.Usage, + "metrics import": headlinemetrics.Usage, + "metrics status": headlinemetrics.Usage, "lexical": "Usage: discrawl lexical rebuild\n\nRebuild configured language indexes locally, without contacting Discord.\n", "lexical rebuild": "Usage: discrawl lexical rebuild\n\nRun after enabling languages or replacing helpers, dictionaries, or Kiwi models.\n", "metadata": `Usage: discrawl metadata [--json] diff --git a/internal/headlinemetrics/discord.go b/internal/headlinemetrics/discord.go new file mode 100644 index 0000000..a1c7776 --- /dev/null +++ b/internal/headlinemetrics/discord.go @@ -0,0 +1,83 @@ +package headlinemetrics + +import ( + "context" + "encoding/json" + "errors" + "io" + "net/http" + "net/url" + "regexp" + "time" +) + +var inviteCode = regexp.MustCompile(`^[A-Za-z0-9_-]{1,100}$`) + +func CollectDiscord(ctx context.Context, c Config, ts string) ([]Row, error) { + client := &http.Client{ + Timeout: 30 * time.Second, + CheckRedirect: func(*http.Request, []*http.Request) error { return http.ErrUseLastResponse }, + } + return collectDiscord(ctx, c, ts, client) +} + +func collectDiscord(ctx context.Context, c Config, ts string, client *http.Client) ([]Row, error) { + rows := make([]Row, 0, len(c.Targets)*2) + failed := false + for _, t := range c.Targets { + members, online := inviteCounts(ctx, client, t.Target) + if members == nil || online == nil { + failed = true + } + rows = append(rows, Counter(t, "members", members, ts, "discord_invite_approximate"), Counter(t, "online", online, ts, "discord_invite_approximate")) + } + if failed { + return rows, errors.New("discord invite counts unavailable") + } + return rows, nil +} + +func inviteCounts(ctx context.Context, client *http.Client, code string) (*float64, *float64) { + if !inviteCode.MatchString(code) { + return nil, nil + } + req, err := http.NewRequestWithContext(ctx, http.MethodGet, "https://discord.com/api/v10/invites/"+url.PathEscape(code)+"?with_counts=true", nil) + if err != nil { + return nil, nil + } + req.Header.Set("Accept", "application/json") + res, err := client.Do(req) + if err != nil { + return nil, nil + } + defer func() { _ = res.Body.Close() }() + if res.StatusCode != http.StatusOK { + // A scheduler may try again later, including after 429; never convert + // permission, expiry, or transient HTTP errors into zero members. + return nil, nil + } + const maxBody = 2 * 1024 * 1024 + body, err := io.ReadAll(io.LimitReader(res.Body, maxBody+1)) + if err != nil || len(body) > maxBody { + return nil, nil + } + var data struct { + Guild *struct { + ID string `json:"id"` + } `json:"guild"` + Members *int64 `json:"approximate_member_count"` + Online *int64 `json:"approximate_presence_count"` + } + if json.Unmarshal(body, &data) != nil || data.Guild == nil || data.Guild.ID == "" { + return nil, nil + } + return approximateCount(data.Members), approximateCount(data.Online) +} + +func approximateCount(n *int64) *float64 { + if n == nil || *n < 0 || *n > 1<<53 { + return nil + } + value := float64(*n) + return &value +} diff --git a/internal/headlinemetrics/discord_test.go b/internal/headlinemetrics/discord_test.go new file mode 100644 index 0000000..d4929a1 --- /dev/null +++ b/internal/headlinemetrics/discord_test.go @@ -0,0 +1,122 @@ +package headlinemetrics + +import ( + "context" + "errors" + "io" + "net/http" + "strings" + "testing" + + "github.com/stretchr/testify/require" +) + +type roundTripFunc func(*http.Request) (*http.Response, error) + +func (f roundTripFunc) RoundTrip(r *http.Request) (*http.Response, error) { return f(r) } + +const sampleTime = "2026-09-15T12:00:00Z" + +func TestCollectDiscordPublicCountsAndUnknowns(t *testing.T) { + cases := []struct { + name string + status int + body string + members *float64 + online *float64 + }{ + {"zero is valid", 200, `{"guild":{"id":"123"},"approximate_member_count":20,"approximate_presence_count":0}`, new(20.0), new(0.0)}, + {"missing presence", 200, `{"guild":{"id":"123"},"approximate_member_count":12}`, new(12.0), nil}, + {"null members", 200, `{"guild":{"id":"123"},"approximate_member_count":null,"approximate_presence_count":3}`, nil, new(3.0)}, + {"negative count", 200, `{"guild":{"id":"123"},"approximate_member_count":-1,"approximate_presence_count":3}`, nil, new(3.0)}, + {"fractional count", 200, `{"guild":{"id":"123"},"approximate_member_count":1.5}`, nil, nil}, + {"unrepresentable count", 200, `{"guild":{"id":"123"},"approximate_member_count":9007199254740993,"approximate_presence_count":1}`, nil, new(1.0)}, + {"non-guild invite", 200, `{"approximate_member_count":20,"approximate_presence_count":5}`, nil, nil}, + {"malformed", 200, `{"guild":`, nil, nil}, + {"trailing JSON", 200, `{"guild":{"id":"123"}} {}`, nil, nil}, + {"oversized", 200, strings.Repeat(" ", 2*1024*1024+1), nil, nil}, + {"expired invite", 404, `{}`, nil, nil}, + {"rate limited", 429, `{"retry_after":3600}`, nil, nil}, + {"unavailable", 503, `{}`, nil, nil}, + } + for _, tc := range cases { + t.Run(tc.name, func(t *testing.T) { + calls := 0 + body := &closedBody{Reader: strings.NewReader(tc.body)} + client := &http.Client{Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) { + calls++ + require.Equal(t, http.MethodGet, req.Method) + require.Equal(t, "https://discord.com/api/v10/invites/example-alpha?with_counts=true", req.URL.String()) + require.Empty(t, req.Header.Get("Authorization")) + require.Empty(t, req.Header.Get("Cookie")) + return &http.Response{StatusCode: tc.status, Body: body, Header: http.Header{"Retry-After": []string{"3600"}}}, nil + })} + rows, err := collectDiscord(t.Context(), Config{Targets: []Target{{Entity: "sample-alpha", Target: "example-alpha"}}}, sampleTime, client) + require.Len(t, rows, 2) + require.Equal(t, tc.members, rows[0].Value) + require.Equal(t, tc.online, rows[1].Value) + require.Equal(t, "members", rows[0].Metric) + require.Equal(t, "online", rows[1].Metric) + require.Equal(t, "discord_invite_approximate", rows[1].Provenance) + require.Equal(t, "counter", rows[1].Kind) + require.Equal(t, sampleTime, rows[1].ObservedAt) + require.Equal(t, 1, calls) + require.True(t, body.closed) + if tc.members == nil || tc.online == nil { + require.Error(t, err) + } else { + require.NoError(t, err) + } + }) + } +} + +type closedBody struct { + io.Reader + closed bool +} + +func (b *closedBody) Close() error { b.closed = true; return nil } + +func TestCollectDiscordIgnoresCredentialsAndRefusesRedirects(t *testing.T) { + t.Setenv("DISCORD_BOT_TOKEN", "fixture-not-to-be-read") + old := http.DefaultTransport + t.Cleanup(func() { http.DefaultTransport = old }) + requested := []string{} + http.DefaultTransport = roundTripFunc(func(req *http.Request) (*http.Response, error) { + requested = append(requested, req.URL.String()) + require.Equal(t, "discord.com", req.URL.Host) + require.Empty(t, req.Header.Get("Authorization")) + require.Empty(t, req.Header.Get("Cookie")) + if strings.HasSuffix(req.URL.Path, "/example-alpha") { + return &http.Response{StatusCode: http.StatusFound, Header: http.Header{"Location": []string{"https://elsewhere.invalid/private"}}, Body: io.NopCloser(strings.NewReader(""))}, nil + } + return &http.Response{StatusCode: http.StatusOK, Body: io.NopCloser(strings.NewReader(`{"guild":{"id":"456"},"approximate_member_count":50,"approximate_presence_count":7}`))}, nil + }) + rows, err := CollectDiscord(t.Context(), Config{CookieJar: "/never/read", TokenEnv: "DISCORD_BOT_TOKEN", Targets: []Target{{"sample-alpha", "example-alpha"}, {"sample-beta", "example-beta"}}}, sampleTime) + require.Error(t, err) + require.Len(t, requested, 2) + require.Len(t, rows, 4) + require.Nil(t, rows[0].Value) + require.Equal(t, new(50.0), rows[2].Value) + require.Equal(t, "sample-beta", rows[2].Entity) +} + +func TestCollectDiscordInvalidTargetAndTransportFailure(t *testing.T) { + calls := 0 + client := &http.Client{Transport: roundTripFunc(func(req *http.Request) (*http.Response, error) { + calls++ + return nil, errors.New("fixture transport failure") + })} + rows, err := collectDiscord(t.Context(), Config{Targets: []Target{{"bad", "../private"}, {"valid", "example-alpha"}}}, sampleTime, client) + require.Error(t, err) + require.Len(t, rows, 4) + require.Equal(t, 1, calls) + for _, row := range rows { + require.Nil(t, row.Value) + } + ctx, cancel := context.WithCancel(t.Context()) + cancel() + _, err = collectDiscord(ctx, Config{Targets: []Target{{"valid", "example-alpha"}}}, sampleTime, &http.Client{}) + require.Error(t, err) +} diff --git a/internal/headlinemetrics/metrics.go b/internal/headlinemetrics/metrics.go new file mode 100644 index 0000000..b7012fe --- /dev/null +++ b/internal/headlinemetrics/metrics.go @@ -0,0 +1,373 @@ +// Package headlinemetrics stores public aggregate observations separately from +// Discord's message/member archive. It does not load archive configuration. +package headlinemetrics + +import ( + "bufio" + "context" + "crypto/sha256" + "database/sql" + "encoding/hex" + "encoding/json" + "errors" + "flag" + "fmt" + "io" + "math" + "os" + "path/filepath" + "strings" + "time" + + "github.com/openclaw/crawlkit/store" +) + +type Target struct { + Entity string `json:"entity"` + Target string `json:"target"` +} + +type Config struct { + Database string `json:"database"` + Targets []Target `json:"targets"` + CookieJar string `json:"cookieJar,omitempty"` + TokenEnv string `json:"tokenEnv,omitempty"` +} + +type Row struct { + Type string `json:"type"` + ID string `json:"id,omitempty"` + Entity string `json:"entity"` + Target string `json:"target"` + Metric string `json:"metric,omitempty"` + Kind string `json:"kind"` + TS string `json:"ts"` + Value *float64 `json:"value"` + ObservedAt string `json:"observed_at"` + Provenance string `json:"provenance"` + Label string `json:"label,omitempty"` + URL string `json:"url,omitempty"` +} + +type Collector func(context.Context, Config, string) ([]Row, error) + +func Counter(t Target, metric string, value *float64, ts, basis string) Row { + return Row{Type: "metric", Entity: t.Entity, Target: t.Target, Metric: metric, Kind: "counter", TS: ts, Value: value, ObservedAt: ts, Provenance: basis} +} + +const schema = ` +CREATE TABLE IF NOT EXISTS metric_meta(key TEXT PRIMARY KEY,value TEXT NOT NULL); +CREATE TABLE IF NOT EXISTS metric_observations(sequence INTEGER PRIMARY KEY AUTOINCREMENT,id TEXT NOT NULL UNIQUE,entity TEXT NOT NULL,target TEXT NOT NULL,metric TEXT NOT NULL,kind TEXT NOT NULL CHECK(kind IN ('counter','daily')),ts TEXT NOT NULL,value REAL,observed_at TEXT NOT NULL,provenance TEXT NOT NULL); +CREATE INDEX IF NOT EXISTS metric_series ON metric_observations(target,metric,ts,sequence); +CREATE TABLE IF NOT EXISTS metric_events(sequence INTEGER PRIMARY KEY AUTOINCREMENT,id TEXT NOT NULL UNIQUE,entity TEXT NOT NULL,target TEXT NOT NULL,kind TEXT NOT NULL,ts TEXT NOT NULL,label TEXT NOT NULL,url TEXT NOT NULL,observed_at TEXT NOT NULL,provenance TEXT NOT NULL); +CREATE TABLE IF NOT EXISTS metric_runs(sequence INTEGER PRIMARY KEY AUTOINCREMENT,ts TEXT NOT NULL,status TEXT NOT NULL,rows_written INTEGER NOT NULL); +` + +func checkOwner(ctx context.Context, db *sql.DB, owner string) error { + var identity, version string + err := db.QueryRowContext(ctx, `SELECT + (SELECT value FROM metric_meta WHERE key='owner'), + (SELECT value FROM metric_meta WHERE key='version')`).Scan(&identity, &version) + if err != nil || identity != owner || version != "1" { + return errors.New("refusing database: expected this metrics collector's owner and version 1") + } + return nil +} + +func openReadOnly(ctx context.Context, path, owner string) (*store.Store, error) { + if !filepath.IsAbs(path) { + return nil, errors.New("metrics database path must be absolute") + } + s, err := store.OpenReadOnly(ctx, path) + if err != nil { + return nil, err + } + if err := checkOwner(ctx, s.DB(), owner); err != nil { + _ = s.Close() + return nil, err + } + return s, nil +} + +func Open(ctx context.Context, path, owner string) (*store.Store, error) { + if !filepath.IsAbs(path) || strings.TrimSpace(owner) == "" { + return nil, errors.New("metrics database path must be absolute and owner must be set") + } + _, err := os.Lstat(path) + if err == nil { + // Validate before write access, including journal or permission changes. + read, err := openReadOnly(ctx, path, owner) + if err != nil { + return nil, err + } + if err := read.Close(); err != nil { + return nil, err + } + return store.Open(ctx, store.Options{Path: path, MaxOpenConns: 1, MaxIdleConns: 1}) + } + if !errors.Is(err, os.ErrNotExist) { + return nil, err + } + return initialize(ctx, path, owner) +} + +func initialize(ctx context.Context, path, owner string) (*store.Store, error) { + if err := os.MkdirAll(filepath.Dir(path), 0o700); err != nil { + return nil, err + } + // Reserve only a previously nonexistent path. A concurrent creator must + // validate the existing database, never initialize over the winning file. + f, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_EXCL, 0o600) + if errors.Is(err, os.ErrExist) { + return Open(ctx, path, owner) + } + if err != nil { + return nil, err + } + if err := f.Close(); err != nil { + return nil, err + } + s, err := store.Open(ctx, store.Options{Path: path, MaxOpenConns: 1, MaxIdleConns: 1}) + if err != nil { + return nil, err + } + err = s.WithTx(ctx, func(tx *sql.Tx) error { + if _, err := tx.ExecContext(ctx, schema); err != nil { + return err + } + _, err := tx.ExecContext(ctx, "INSERT INTO metric_meta VALUES('owner',?),('version','1')", owner) + return err + }) + if err != nil { + _ = s.Close() + return nil, err + } + return s, nil +} + +func Validate(r Row) error { + if r.Type != "metric" && r.Type != "event" { + return errors.New("invalid observation type") + } + if strings.TrimSpace(r.Entity) == "" || strings.TrimSpace(r.Target) == "" || len(r.Entity) > 200 || len(r.Target) > 300 || strings.TrimSpace(r.Provenance) == "" || strings.TrimSpace(r.Kind) == "" { + return errors.New("invalid observation identity") + } + for _, v := range []string{r.TS, r.ObservedAt} { + if _, err := time.Parse(time.RFC3339Nano, v); err != nil { + return errors.New("invalid observation time") + } + } + if r.Type == "metric" && (strings.TrimSpace(r.Metric) == "" || (r.Kind != "counter" && r.Kind != "daily") || r.Value != nil && (math.IsNaN(*r.Value) || math.IsInf(*r.Value, 0) || *r.Value < 0)) { + return errors.New("invalid metric") + } + return nil +} + +func Write(ctx context.Context, s *store.Store, rows []Row) (int, error) { + written := 0 + err := s.WithTx(ctx, func(tx *sql.Tx) error { + for _, r := range rows { + if err := Validate(r); err != nil { + return err + } + if r.ID == "" { + b, err := json.Marshal(r) + if err != nil { + return err + } + h := sha256.Sum256(b) + r.ID = hex.EncodeToString(h[:]) + } + // Only IDs deduplicate. Retain repeated values, decreases, unknowns, + // and revised daily rows; consumers select the latest daily sequence. + var result sql.Result + var err error + if r.Type == "metric" { + result, err = tx.ExecContext(ctx, "INSERT INTO metric_observations(id,entity,target,metric,kind,ts,value,observed_at,provenance) VALUES(?,?,?,?,?,?,?,?,?) ON CONFLICT(id) DO NOTHING", r.ID, r.Entity, r.Target, r.Metric, r.Kind, r.TS, r.Value, r.ObservedAt, r.Provenance) + } else { + result, err = tx.ExecContext(ctx, "INSERT INTO metric_events(id,entity,target,kind,ts,label,url,observed_at,provenance) VALUES(?,?,?,?,?,?,?,?,?) ON CONFLICT(id) DO NOTHING", r.ID, r.Entity, r.Target, r.Kind, r.TS, r.Label, r.URL, r.ObservedAt, r.Provenance) + } + if err != nil { + return err + } + n, err := result.RowsAffected() + if err != nil { + return err + } + written += int(n) + } + return nil + }) + if err != nil { + return 0, err // The entire batch rolled back. + } + return written, nil +} + +const Usage = `Usage: discrawl metrics --config /absolute/metrics.json + +Collect public invite approximate members and online presence into a separate +metrics database. Online presence is not a count of active posters. +Import reads scoped NDJSON history from stdin. Status is read-only. +No Discord token, archive configuration, or message/member archive is used. +Successful commands write JSON; partial collection writes NULLs and exits nonzero. +` + +func loadConfig(path string) (Config, error) { + var c Config + raw, err := os.ReadFile(path) + if err != nil { + return c, errors.New("metrics config unavailable") + } + d := json.NewDecoder(strings.NewReader(string(raw))) + d.DisallowUnknownFields() + if err := d.Decode(&c); err != nil { + return c, errors.New("invalid metrics config") + } + if d.Decode(new(any)) != io.EOF || !filepath.IsAbs(c.Database) || len(c.Targets) == 0 { + return c, errors.New("metrics config requires an absolute database path and targets") + } + seen := map[string]bool{} + for _, t := range c.Targets { + if strings.TrimSpace(t.Entity) == "" || len(t.Entity) > 200 || !inviteCode.MatchString(t.Target) || seen[t.Target] { + return c, errors.New("metrics targets require an entity and unique Discord invite code") + } + seen[t.Target] = true + } + return c, nil +} + +func Run(ctx context.Context, args []string, owner string, collect Collector, in io.Reader, out, errout io.Writer) error { + if len(args) == 0 || args[0] == "--help" || args[0] == "-h" || args[0] == "help" { + _, err := io.WriteString(out, Usage) + return err + } + command := args[0] + if command != "collect" && command != "import" && command != "status" { + return errors.New("unknown metrics command") + } + flags := flag.NewFlagSet("metrics "+command, flag.ContinueOnError) + flags.SetOutput(errout) + flags.Usage = func() { _, _ = io.WriteString(out, Usage) } + configPath := flags.String("config", "", "metrics JSON configuration (separate from archive config)") + if err := flags.Parse(args[1:]); err != nil { + if errors.Is(err, flag.ErrHelp) { + return nil + } + return err + } + if *configPath == "" || flags.NArg() != 0 { + return errors.New("provide --config METRICS_CONFIG with no positional arguments") + } + c, err := loadConfig(*configPath) + if err != nil { + return err + } + if command == "status" { + return status(ctx, c, owner, out) + } + s, err := Open(ctx, c.Database, owner) + if err != nil { + return err + } + defer func() { _ = s.Close() }() + allowed := map[string]string{} + for _, t := range c.Targets { + allowed[t.Target] = t.Entity + } + if command == "import" { + written, err := importRows(ctx, s, allowed, in) + if outputErr := json.NewEncoder(out).Encode(map[string]any{"source": owner, "command": command, "rows_written": written, "ok": err == nil}); outputErr != nil { + return outputErr + } + return err + } + ts := time.Now().UTC().Format(time.RFC3339Nano) + rows, collectionErr := collect(ctx, c, ts) + for _, r := range rows { + entity, ok := allowed[r.Target] + if !ok || entity != r.Entity { + return errors.New("collector returned invalid target") + } + } + // Persist partial observations after request cancellation, with a bounded + // final write that cannot reach the message/member archive. + writeCtx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + written, err := Write(writeCtx, s, rows) + if err != nil { + return err + } + runStatus := "ok" + if collectionErr != nil { + runStatus = "partial" + } + if _, err := s.DB().ExecContext(writeCtx, "INSERT INTO metric_runs(ts,status,rows_written) VALUES(?,?,?)", ts, runStatus, written); err != nil { + return err + } + if err := json.NewEncoder(out).Encode(map[string]any{"source": owner, "command": command, "rows_written": written, "ok": collectionErr == nil}); err != nil { + return err + } + if collectionErr != nil { + return errors.New("one or more metric reads failed; successful values and NULL observations were retained") + } + return nil +} + +func importRows(ctx context.Context, s *store.Store, allowed map[string]string, in io.Reader) (int, error) { + scanner := bufio.NewScanner(in) + scanner.Buffer(make([]byte, 65536), 4*1024*1024) + batch := make([]Row, 0, 500) + written, line := 0, 0 + for scanner.Scan() { + line++ + var row Row + d := json.NewDecoder(strings.NewReader(scanner.Text())) + d.DisallowUnknownFields() + if err := d.Decode(&row); err != nil || d.Decode(new(any)) != io.EOF { + return written, fmt.Errorf("invalid import JSON at line %d", line) + } + entity, ok := allowed[row.Target] + if !ok || entity != row.Entity || strings.TrimSpace(row.ID) == "" { + return written, fmt.Errorf("invalid import identity or scope at line %d", line) + } + if err := Validate(row); err != nil { + return written, fmt.Errorf("import line %d: %w", line, err) + } + batch = append(batch, row) + if len(batch) == 500 { + n, err := Write(ctx, s, batch) + if err != nil { + return written, err + } + written += n + batch = batch[:0] + } + } + if err := scanner.Err(); err != nil { + return written, fmt.Errorf("read import after line %d: %w", line, err) + } + n, err := Write(ctx, s, batch) + return written + n, err +} + +func status(ctx context.Context, c Config, owner string, out io.Writer) error { + s, err := openReadOnly(ctx, c.Database, owner) + if err != nil { + return err + } + defer func() { _ = s.Close() }() + var observations, events, sequence int64 + if err := s.DB().QueryRowContext(ctx, "SELECT count(*),coalesce(max(sequence),0) FROM metric_observations").Scan(&observations, &sequence); err != nil { + return err + } + if err := s.DB().QueryRowContext(ctx, "SELECT count(*) FROM metric_events").Scan(&events); err != nil { + return err + } + var last *string + err = s.DB().QueryRowContext(ctx, "SELECT observed_at FROM metric_observations ORDER BY julianday(observed_at) DESC,sequence DESC LIMIT 1").Scan(&last) + if err != nil && !errors.Is(err, sql.ErrNoRows) { + return err + } + return json.NewEncoder(out).Encode(map[string]any{"source": owner, "observations": observations, "events": events, "sequence": sequence, "last_observed": last}) +} diff --git a/internal/headlinemetrics/metrics_test.go b/internal/headlinemetrics/metrics_test.go new file mode 100644 index 0000000..36365ec --- /dev/null +++ b/internal/headlinemetrics/metrics_test.go @@ -0,0 +1,298 @@ +package headlinemetrics + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "math" + "os" + "path/filepath" + "strings" + "sync" + "testing" + "time" + + "github.com/openclaw/crawlkit/store" + "github.com/stretchr/testify/require" +) + +func newMetrics(t *testing.T) *store.Store { + t.Helper() + s, err := Open(t.Context(), filepath.Join(t.TempDir(), "nested", "metrics.sqlite"), "discrawl") + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, s.Close()) }) + return s +} + +func testConfig(t *testing.T, database string) string { + t.Helper() + b, err := json.Marshal(Config{Database: database, Targets: []Target{{"sample-alpha", "example-alpha"}, {"sample-beta", "example-beta"}}}) + require.NoError(t, err) + p := filepath.Join(t.TempDir(), "metrics.json") + require.NoError(t, os.WriteFile(p, b, 0o600)) + return p +} + +func TestOpenPreservesUnownedAndNewerDatabases(t *testing.T) { + for _, kind := range []string{"archive", "other-owner", "newer-version", "empty"} { + t.Run(kind, func(t *testing.T) { + path := filepath.Join(t.TempDir(), "existing.sqlite") + switch kind { + case "empty": + require.NoError(t, os.WriteFile(path, nil, 0o600)) + case "archive": + s, err := store.Open(t.Context(), store.Options{Path: path, Schema: "CREATE TABLE messages(id TEXT PRIMARY KEY,content TEXT); INSERT INTO messages VALUES('saved','preserve this archive');"}) + require.NoError(t, err) + require.NoError(t, s.Close()) + default: + s, err := Open(t.Context(), path, "discrawl") + require.NoError(t, err) + query := "UPDATE metric_meta SET value='another-collector' WHERE key='owner'" + if kind == "newer-version" { + query = "UPDATE metric_meta SET value='2' WHERE key='version'" + } + _, err = s.DB().ExecContext(t.Context(), query) + require.NoError(t, err) + require.NoError(t, s.Close()) + } + before, err := os.ReadFile(path) + require.NoError(t, err) + _, err = Open(t.Context(), path, "discrawl") + require.Error(t, err) + _, err = openReadOnly(t.Context(), path, "discrawl") + require.Error(t, err) + after, err := os.ReadFile(path) + require.NoError(t, err) + require.Equal(t, before, after) + }) + } +} + +func TestConcurrentInitializationHasOneOwner(t *testing.T) { + path := filepath.Join(t.TempDir(), "metrics.sqlite") + var wg sync.WaitGroup + results := make(chan string, 2) + start := make(chan struct{}) + for _, owner := range []string{"discrawl", "another-collector"} { + wg.Go(func() { + <-start + s, err := Open(t.Context(), path, owner) + if err == nil { + results <- owner + _ = s.Close() + } + }) + } + close(start) + wg.Wait() + close(results) + owners := []string{} + for owner := range results { + owners = append(owners, owner) + } + require.Len(t, owners, 1) + s, err := openReadOnly(t.Context(), path, owners[0]) + require.NoError(t, err) + require.NoError(t, s.Close()) + leftovers, err := filepath.Glob(filepath.Join(filepath.Dir(path), ".metrics-init-*")) + require.NoError(t, err) + require.Empty(t, leftovers) +} + +func TestOpenRejectsDanglingDatabaseSymlink(t *testing.T) { + root := t.TempDir() + path := filepath.Join(root, "metrics.sqlite") + destination := filepath.Join(root, "missing-archive.sqlite") + if err := os.Symlink(destination, path); err != nil { + t.Skipf("symlinks unavailable: %v", err) + } + _, err := Open(t.Context(), path, "discrawl") + require.Error(t, err) + require.NoFileExists(t, destination) + leftovers, err := filepath.Glob(filepath.Join(root, ".metrics-init-*")) + require.NoError(t, err) + require.Empty(t, leftovers) +} + +func TestConcurrentWritesUseNativeSQLiteLock(t *testing.T) { + first := newMetrics(t) + second, err := Open(t.Context(), first.Path(), "discrawl") + require.NoError(t, err) + t.Cleanup(func() { require.NoError(t, second.Close()) }) + tx, err := first.DB().BeginTx(t.Context(), nil) + require.NoError(t, err) + t.Cleanup(func() { _ = tx.Rollback() }) + _, err = tx.ExecContext(t.Context(), "INSERT INTO metric_meta(key,value) VALUES('test-writer','held')") + require.NoError(t, err) + row := Counter(Target{"sample-alpha", "example-alpha"}, "members", new(12.0), sampleTime, "fixture") + ctx, cancel := context.WithTimeout(t.Context(), 50*time.Millisecond) + defer cancel() + written, err := Write(ctx, second, []Row{row}) + require.Error(t, err) + require.Zero(t, written) + var count int + require.NoError(t, second.DB().QueryRowContext(t.Context(), "SELECT count(*) FROM metric_observations").Scan(&count)) + require.Zero(t, count) + require.NoError(t, tx.Rollback()) + written, err = Write(t.Context(), second, []Row{row}) + require.NoError(t, err) + require.Equal(t, 1, written) +} + +func TestWritePreservesHistoryAndRollsBackBadBatches(t *testing.T) { + s := newMetrics(t) + rows := []Row{} + for i, value := range []*float64{new(20.0), new(0.0), new(19.0), nil} { + r := Counter(Target{"sample-alpha", "example-alpha"}, "members", value, sampleTime, "historical-source") + r.ID = fmt.Sprintf("counter-%d", i) + rows = append(rows, r) + } + for i, value := range []float64{10, 9, 9} { + r := Counter(Target{"sample-alpha", "example-alpha"}, "members", new(value), "2026-09-14T00:00:00Z", "native-history") + r.ID, r.Kind, r.ObservedAt = fmt.Sprintf("daily-%d", i), "daily", sampleTime + rows = append(rows, r) + } + rows = append(rows, Row{Type: "event", ID: "event-1", Entity: "sample-beta", Target: "example-beta", Kind: "release", TS: sampleTime, ObservedAt: sampleTime, Provenance: "historical-source", Label: "release", URL: "https://example.test/release"}) + written, err := Write(t.Context(), s, rows) + require.NoError(t, err) + require.Equal(t, len(rows), written) + written, err = Write(t.Context(), s, rows) + require.NoError(t, err) + require.Zero(t, written) + var count int + require.NoError(t, s.DB().QueryRowContext(t.Context(), "SELECT count(*) FROM metric_observations").Scan(&count)) + require.Equal(t, 7, count) + var nulls int + require.NoError(t, s.DB().QueryRowContext(t.Context(), "SELECT count(*) FROM metric_observations WHERE value IS NULL").Scan(&nulls)) + require.Equal(t, 1, nulls) + var latest float64 + require.NoError(t, s.DB().QueryRowContext(t.Context(), "SELECT value FROM metric_observations WHERE kind='daily' ORDER BY sequence DESC LIMIT 1").Scan(&latest)) + require.InDelta(t, 9, latest, 0) + good := Counter(Target{"sample-alpha", "example-alpha"}, "online", new(3.0), sampleTime, "native") + bad := good + bad.Value = new(-1.0) + written, err = Write(t.Context(), s, []Row{good, bad}) + require.Error(t, err) + require.Zero(t, written) + written, err = Write(t.Context(), s, []Row{good}) + require.NoError(t, err) + require.Equal(t, 1, written) + written, err = Write(t.Context(), s, []Row{good}) + require.NoError(t, err) + require.Zero(t, written) +} + +func TestValidationRejectsInvalidObservations(t *testing.T) { + good := Counter(Target{"sample-alpha", "example-alpha"}, "members", new(1.0), sampleTime, "native") + for _, change := range []func(*Row){ + func(r *Row) { r.Type = "other" }, func(r *Row) { r.Entity = " " }, + func(r *Row) { r.TS = "yesterday" }, func(r *Row) { r.ObservedAt = "unknown" }, + func(r *Row) { r.Metric = "" }, func(r *Row) { r.Kind = "gauge" }, + func(r *Row) { r.Value = new(math.NaN()) }, func(r *Row) { r.Value = new(math.Inf(1)) }, + func(r *Row) { r.Type = "event"; r.Kind = "" }, + } { + row := good + change(&row) + require.Error(t, Validate(row)) + } +} + +func TestRunPartialCollectionAndReadOnlyStatus(t *testing.T) { + database := filepath.Join(t.TempDir(), "metrics.sqlite") + config := testConfig(t, database) + collect := func(_ context.Context, c Config, ts string) ([]Row, error) { + return []Row{Counter(c.Targets[0], "members", new(42.0), ts, "discord_invite_approximate"), Counter(c.Targets[1], "members", nil, ts, "discord_invite_approximate")}, errors.New("fixture unavailable") + } + var out bytes.Buffer + err := Run(t.Context(), []string{"collect", "--config", config}, "discrawl", collect, nil, &out, io.Discard) + require.ErrorContains(t, err, "NULL observations were retained") + require.JSONEq(t, `{"source":"discrawl","command":"collect","rows_written":2,"ok":false}`, out.String()) + s, err := openReadOnly(t.Context(), database, "discrawl") + require.NoError(t, err) + var status string + require.NoError(t, s.DB().QueryRowContext(t.Context(), "SELECT status FROM metric_runs").Scan(&status)) + require.Equal(t, "partial", status) + require.NoError(t, s.Close()) + before, err := os.ReadFile(database) + require.NoError(t, err) + out.Reset() + require.NoError(t, Run(t.Context(), []string{"status", "--config", config}, "discrawl", nil, nil, &out, io.Discard)) + var state map[string]any + require.NoError(t, json.Unmarshal(out.Bytes(), &state)) + require.InDelta(t, 2, state["observations"], 0) + require.IsType(t, "", state["last_observed"]) + after, err := os.ReadFile(database) + require.NoError(t, err) + require.Equal(t, before, after) +} + +func TestImportCanResumeCommittedBatchesAndPreserveNulls(t *testing.T) { + database := filepath.Join(t.TempDir(), "metrics.sqlite") + config := testConfig(t, database) + var input bytes.Buffer + enc := json.NewEncoder(&input) + for i := range 501 { + r := Counter(Target{"sample-alpha", "example-alpha"}, "members", nil, sampleTime, "history") + r.ID = fmt.Sprintf("history-%d", i) + require.NoError(t, enc.Encode(r)) + } + var out bytes.Buffer + err := Run(t.Context(), []string{"import", "--config", config}, "discrawl", nil, strings.NewReader(input.String()+"broken\n"), &out, io.Discard) + require.ErrorContains(t, err, "line 502") + require.JSONEq(t, `{"source":"discrawl","command":"import","rows_written":500,"ok":false}`, out.String()) + out.Reset() + require.NoError(t, Run(t.Context(), []string{"import", "--config", config}, "discrawl", nil, &input, &out, io.Discard)) + require.JSONEq(t, `{"source":"discrawl","command":"import","rows_written":1,"ok":true}`, out.String()) + s, err := openReadOnly(t.Context(), database, "discrawl") + require.NoError(t, err) + defer func() { require.NoError(t, s.Close()) }() + var count int + require.NoError(t, s.DB().QueryRowContext(t.Context(), "SELECT count(*) FROM metric_observations WHERE value IS NULL").Scan(&count)) + require.Equal(t, 501, count) +} + +func TestImportRejectsScopeMissingIDAndUnknownFields(t *testing.T) { + for _, mutation := range []func(map[string]any){ + func(r map[string]any) { r["target"] = "another-invite" }, + func(r map[string]any) { r["entity"] = "another-entity" }, + func(r map[string]any) { delete(r, "id") }, + func(r map[string]any) { r["unexpected"] = true }, + func(r map[string]any) { r["ts"] = "invalid" }, + } { + r := map[string]any{"type": "metric", "id": "one", "entity": "sample-alpha", "target": "example-alpha", "metric": "members", "kind": "counter", "ts": sampleTime, "observed_at": sampleTime, "provenance": "history", "value": nil} + mutation(r) + b, err := json.Marshal(r) + require.NoError(t, err) + config := testConfig(t, filepath.Join(t.TempDir(), "metrics.sqlite")) + err = Run(t.Context(), []string{"import", "--config", config}, "discrawl", nil, bytes.NewReader(b), io.Discard, io.Discard) + require.Error(t, err) + } +} + +func TestConfigHelpAndMissingDatabase(t *testing.T) { + for _, args := range [][]string{nil, {"help"}, {"-h"}, {"--help"}, {"collect", "--help"}} { + var out bytes.Buffer + require.NoError(t, Run(t.Context(), args, "discrawl", nil, nil, &out, io.Discard)) + require.Contains(t, out.String(), "Online presence is not a count of active posters") + } + for _, args := range [][]string{{"bad"}, {"collect"}, {"collect", "--bad"}, {"collect", "--config", "missing"}, {"status", "--config", "x", "extra"}} { + require.Error(t, Run(t.Context(), args, "discrawl", nil, nil, io.Discard, io.Discard)) + } + path := filepath.Join(t.TempDir(), "missing.sqlite") + config := testConfig(t, path) + require.Error(t, Run(t.Context(), []string{"status", "--config", config}, "discrawl", nil, nil, io.Discard, io.Discard)) + require.NoFileExists(t, path) + for _, raw := range []string{ + `{}`, `{"database":"relative","targets":[{"entity":"x","target":"example-alpha"}]}`, + `{"database":"/tmp/x","targets":[{"entity":"x","target":"../bad"}]}`, + `{"database":"/tmp/x","targets":[{"entity":"x","target":"example-alpha"},{"entity":"y","target":"example-alpha"}]}`, + `{"unexpected":true}`, `{} {}`, `invalid`, + } { + require.NoError(t, os.WriteFile(config, []byte(raw), 0o600)) + _, err := loadConfig(config) + require.Error(t, err) + } +}