From 0788cffc88eca8b32c859aa2311c4b55a0a0a07b Mon Sep 17 00:00:00 2001 From: Gemma Date: Thu, 21 May 2026 09:32:37 +0200 Subject: [PATCH] feat: add AI limits MQTT service --- .gitignore | 3 + Dockerfile | 15 ++ README.md | 52 ++++ docker-compose.yml | 12 + go.mod | 11 + go.sum | 8 + main.go | 637 +++++++++++++++++++++++++++++++++++++++++++++ 7 files changed, 738 insertions(+) create mode 100644 .gitignore create mode 100644 Dockerfile create mode 100644 README.md create mode 100644 docker-compose.yml create mode 100644 go.mod create mode 100644 go.sum create mode 100644 main.go diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..3f51cfa --- /dev/null +++ b/.gitignore @@ -0,0 +1,3 @@ +.env +*.log +ai-limits-mqtt diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..24cd810 --- /dev/null +++ b/Dockerfile @@ -0,0 +1,15 @@ +FROM golang:1.25-bookworm AS builder +WORKDIR /src +COPY go.mod go.sum* ./ +RUN go mod download +COPY . . +RUN /usr/local/go/bin/gofmt -w main.go && CGO_ENABLED=0 GOOS=linux /usr/local/go/bin/go build -trimpath -ldflags='-s -w' -o /out/ai-limits-mqtt . + +FROM node:24-bookworm-slim +RUN apt-get update \ + && apt-get install -y --no-install-recommends ca-certificates git ripgrep sqlite3 \ + && rm -rf /var/lib/apt/lists/* \ + && npm install -g @anthropic-ai/claude-code @openai/codex +COPY --from=builder /out/ai-limits-mqtt /usr/local/bin/ai-limits-mqtt +ENV NODE_ENV=production +CMD ["/usr/local/bin/ai-limits-mqtt"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..a5b220a --- /dev/null +++ b/README.md @@ -0,0 +1,52 @@ +# AI Limits MQTT + +Small Go service that publishes Claude Code and Codex CLI status to Home Assistant over MQTT. + +## What it reports + +The CLIs do not expose a stable standalone "remaining quota" API. This service therefore performs a small real CLI probe and classifies the result: + +- `ok`: CLI authenticated and a tiny probe completed. +- `limited`: output matched rate/usage/quota-limit errors. +- `auth_required`: login/auth is missing or invalid. +- `timeout` / `error`: probe failed for another reason. + +It also publishes any usage object returned by the CLI probe: + +- Claude: JSON result `usage` from `claude -p ... --output-format json`. +- Codex: JSONL `usage` from `codex exec --json` turn completion. + +It now also reports rolling local token usage for each CLI: + +- `session_5h.used_tokens` and `session_5h.percent_used` +- `week.used_tokens` and `week.percent_used` + +The percent fields are computed against env-configured token limits: + +- `CLAUDE_SESSION_5H_TOKEN_LIMIT` +- `CLAUDE_WEEK_TOKEN_LIMIT` +- `CODEX_SESSION_5H_TOKEN_LIMIT` +- `CODEX_WEEK_TOKEN_LIMIT` + +Set a limit to `0` to publish raw used tokens without a percentage. + +## MQTT topics + +- `ai_limits/availability` +- `ai_limits/claude/state` +- `ai_limits/codex/state` +- `ai_limits/summary/state` + +Home Assistant discovery is retained under: + +- `homeassistant/sensor/ai_limits_claude/config` +- `homeassistant/sensor/ai_limits_codex/config` +- `homeassistant/sensor/ai_limits_summary/config` + +## Run + +```bash +docker compose up -d --build +``` + +Secrets live in `.env` and are intentionally gitignored. diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..002fc44 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,12 @@ +services: + ai-limits-mqtt: + build: . + image: ai-limits-mqtt:latest + container_name: ai-limits-mqtt + restart: unless-stopped + env_file: + - .env + volumes: + - /root/.claude:/root/.claude + - /root/.claude.json:/root/.claude.json + - /root/.codex:/root/.codex diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..5c46bb3 --- /dev/null +++ b/go.mod @@ -0,0 +1,11 @@ +module ai-limits-mqtt + +go 1.23 + +require github.com/eclipse/paho.mqtt.golang v1.5.0 + +require ( + github.com/gorilla/websocket v1.5.3 // indirect + golang.org/x/net v0.27.0 // indirect + golang.org/x/sync v0.7.0 // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..b555dd0 --- /dev/null +++ b/go.sum @@ -0,0 +1,8 @@ +github.com/eclipse/paho.mqtt.golang v1.5.0 h1:EH+bUVJNgttidWFkLLVKaQPGmkTUfQQqjOsyvMGvD6o= +github.com/eclipse/paho.mqtt.golang v1.5.0/go.mod h1:du/2qNQVqJf/Sqs4MEL77kR8QTqANF7XU7Fk0aOTAgk= +github.com/gorilla/websocket v1.5.3 h1:saDtZ6Pbx/0u+bgYQ3q96pZgCzfhKXGPqt7kZ72aNNg= +github.com/gorilla/websocket v1.5.3/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE= +golang.org/x/net v0.27.0 h1:5K3Njcw06/l2y9vpGCSdcxWOYHOUk3dVNGDXN+FvAys= +golang.org/x/net v0.27.0/go.mod h1:dDi0PyhWNoiUOrAS8uXv/vnScO4wnHQO4mj9fn/RytE= +golang.org/x/sync v0.7.0 h1:YsImfSBoP9QPYL0xyKJPq0gcaJdG3rInoqxTWbfQu9M= +golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= diff --git a/main.go b/main.go new file mode 100644 index 0000000..8e2b5ce --- /dev/null +++ b/main.go @@ -0,0 +1,637 @@ +package main + +import ( + "bufio" + "bytes" + "context" + "crypto/sha1" + "encoding/hex" + "encoding/json" + "errors" + "fmt" + "log" + "math" + "os" + "os/exec" + "path/filepath" + "strconv" + "strings" + "sync" + "time" + + mqtt "github.com/eclipse/paho.mqtt.golang" +) + +type Config struct { + MQTTBroker string + MQTTUsername string + MQTTPassword string + MQTTClientID string + MQTTBaseTopic string + HADiscoveryPrefix string + ProbeInterval time.Duration + ProbeTimeout time.Duration + ClaudeCommand string + CodexCommand string + PublishDiscovery bool + RunProbeCommand bool + ClaudeProbePrompt string + CodexProbePrompt string + ClaudeExtraArgs []string + CodexExtraArgs []string + ClaudeHistoryDir string + CodexStateDB string + ClaudeSession5hLimit int64 + ClaudeWeekLimit int64 + CodexSession5hLimit int64 + CodexWeekLimit int64 +} + +type LimitWindow struct { + Window string `json:"window"` + UsedTokens int64 `json:"used_tokens"` + LimitTokens int64 `json:"limit_tokens,omitempty"` + PercentUsed *float64 `json:"percent_used,omitempty"` +} + +type ProviderStatus struct { + Provider string `json:"provider"` + Available bool `json:"available"` + Status string `json:"status"` + Message string `json:"message,omitempty"` + Limited bool `json:"limited"` + LimitResetText string `json:"limit_reset_text,omitempty"` + CheckedAt string `json:"checked_at"` + LatencyMS int64 `json:"latency_ms"` + ExitCode int `json:"exit_code"` + Version string `json:"version,omitempty"` + Auth string `json:"auth,omitempty"` + Usage map[string]interface{} `json:"usage,omitempty"` + Session5h LimitWindow `json:"session_5h"` + Week LimitWindow `json:"week"` +} + +type Summary struct { + Status string `json:"status"` + CheckedAt string `json:"checked_at"` + Claude string `json:"claude"` + Codex string `json:"codex"` +} + +func main() { + cfg := loadConfig() + client := connectMQTT(cfg) + defer client.Disconnect(500) + + publish(client, cfg.topic("availability"), "online", true) + if cfg.PublishDiscovery { + publishDiscovery(client, cfg) + } + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + statuses := probeAndPublish(ctx, client, cfg) + log.Printf("initial probe complete: claude=%s codex=%s", statuses["claude"].Status, statuses["codex"].Status) + + ticker := time.NewTicker(cfg.ProbeInterval) + defer ticker.Stop() + for range ticker.C { + probeAndPublish(ctx, client, cfg) + } +} + +func loadConfig() Config { + return Config{ + MQTTBroker: env("MQTT_BROKER", "tcp://192.168.33.33:1883"), + MQTTUsername: env("MQTT_USERNAME", ""), + MQTTPassword: env("MQTT_PASSWORD", ""), + MQTTClientID: env("MQTT_CLIENT_ID", "ai-limits-mqtt"), + MQTTBaseTopic: trimSlashes(env("MQTT_BASE_TOPIC", "ai_limits")), + HADiscoveryPrefix: trimSlashes(env("HA_DISCOVERY_PREFIX", "homeassistant")), + ProbeInterval: envDuration("PROBE_INTERVAL", 15*time.Minute), + ProbeTimeout: envDuration("PROBE_TIMEOUT", 180*time.Second), + ClaudeCommand: env("CLAUDE_CMD", "claude"), + CodexCommand: env("CODEX_CMD", "codex"), + PublishDiscovery: envBool("PUBLISH_DISCOVERY", true), + RunProbeCommand: envBool("RUN_PROBE_COMMAND", true), + ClaudeProbePrompt: env("CLAUDE_PROBE_PROMPT", "Reply with OK only."), + CodexProbePrompt: env("CODEX_PROBE_PROMPT", "Reply with OK only."), + ClaudeExtraArgs: splitArgs(env("CLAUDE_EXTRA_ARGS", "")), + CodexExtraArgs: splitArgs(env("CODEX_EXTRA_ARGS", "")), + ClaudeHistoryDir: env("CLAUDE_HISTORY_DIR", "/root/.claude/projects"), + CodexStateDB: env("CODEX_STATE_DB", "/root/.codex/state_5.sqlite"), + ClaudeSession5hLimit: envInt64("CLAUDE_SESSION_5H_TOKEN_LIMIT", 0), + ClaudeWeekLimit: envInt64("CLAUDE_WEEK_TOKEN_LIMIT", 0), + CodexSession5hLimit: envInt64("CODEX_SESSION_5H_TOKEN_LIMIT", 0), + CodexWeekLimit: envInt64("CODEX_WEEK_TOKEN_LIMIT", 0), + } +} + +func connectMQTT(cfg Config) mqtt.Client { + opts := mqtt.NewClientOptions().AddBroker(cfg.MQTTBroker).SetClientID(cfg.MQTTClientID) + opts.SetCleanSession(true) + opts.SetAutoReconnect(true) + opts.SetConnectRetry(true) + opts.SetOrderMatters(false) + opts.SetWill(cfg.topic("availability"), "offline", 1, true) + if cfg.MQTTUsername != "" { + opts.SetUsername(cfg.MQTTUsername) + opts.SetPassword(cfg.MQTTPassword) + } + client := mqtt.NewClient(opts) + if token := client.Connect(); token.Wait() && token.Error() != nil { + log.Fatalf("mqtt connect failed: %v", token.Error()) + } + return client +} + +func probeAndPublish(ctx context.Context, client mqtt.Client, cfg Config) map[string]ProviderStatus { + var wg sync.WaitGroup + results := make(map[string]ProviderStatus, 2) + var mu sync.Mutex + + for _, job := range []struct { + name string + fn func(context.Context, Config) ProviderStatus + }{ + {"claude", probeClaude}, + {"codex", probeCodex}, + } { + wg.Add(1) + go func(name string, fn func(context.Context, Config) ProviderStatus) { + defer wg.Done() + st := fn(ctx, cfg) + attachLimitUsage(ctx, cfg, &st) + mu.Lock() + results[name] = st + mu.Unlock() + publishJSON(client, cfg.topic(name, "state"), st, true) + }(job.name, job.fn) + } + wg.Wait() + + summary := Summary{CheckedAt: time.Now().UTC().Format(time.RFC3339)} + summary.Claude = results["claude"].Status + summary.Codex = results["codex"].Status + if summary.Claude == "ok" && summary.Codex == "ok" { + summary.Status = "ok" + } else if summary.Claude == "limited" || summary.Codex == "limited" { + summary.Status = "limited" + } else { + summary.Status = "degraded" + } + publishJSON(client, cfg.topic("summary", "state"), summary, true) + publish(client, cfg.topic("availability"), "online", true) + return results +} + +func probeClaude(parent context.Context, cfg Config) ProviderStatus { + st := baseStatus("claude") + st.Version = firstLine(runShort(parent, 10*time.Second, cfg.ClaudeCommand, "--version")) + auth := runShort(parent, 15*time.Second, cfg.ClaudeCommand, "auth", "status", "--text") + st.Auth = compact(auth, 240) + if strings.Contains(strings.ToLower(auth), "not logged") || strings.Contains(strings.ToLower(auth), "login") && strings.Contains(strings.ToLower(auth), "required") { + st.Status = "auth_required" + st.Message = compact(auth, 500) + return st + } + if !cfg.RunProbeCommand { + st.Available = true + st.Status = "auth_ok" + return st + } + + args := []string{"-p", cfg.ClaudeProbePrompt, "--output-format", "json", "--max-turns", "1", "--no-session-persistence"} + args = append(args, cfg.ClaudeExtraArgs...) + started := time.Now() + out, exitCode, err := runCommand(parent, cfg.ProbeTimeout, cfg.ClaudeCommand, args...) + st.LatencyMS = time.Since(started).Milliseconds() + st.ExitCode = exitCode + if err != nil && strings.TrimSpace(out) == "" { + st.Status = classifyFailure(out, err.Error()) + st.Message = compact(err.Error(), 500) + st.Limited = st.Status == "limited" + return st + } + var doc map[string]interface{} + if json.Unmarshal([]byte(lastJSON(out)), &doc) == nil { + if usage, ok := doc["usage"].(map[string]interface{}); ok { + st.Usage = usage + } + if isErr, _ := doc["is_error"].(bool); !isErr && stringField(doc, "subtype") == "success" { + st.Available = true + st.Status = "ok" + st.Message = stringField(doc, "result") + return st + } + } + st.Status = classifyFailure(out, "") + st.Limited = st.Status == "limited" + st.Message = compact(out, 800) + st.LimitResetText = extractReset(out) + return st +} + +func probeCodex(parent context.Context, cfg Config) ProviderStatus { + st := baseStatus("codex") + st.Version = firstLine(runShort(parent, 10*time.Second, cfg.CodexCommand, "--version")) + auth := runShort(parent, 15*time.Second, cfg.CodexCommand, "login", "status") + st.Auth = compact(auth, 240) + if strings.Contains(strings.ToLower(auth), "not logged") || strings.Contains(strings.ToLower(auth), "login") && strings.Contains(strings.ToLower(auth), "required") { + st.Status = "auth_required" + st.Message = compact(auth, 500) + return st + } + if !cfg.RunProbeCommand { + st.Available = true + st.Status = "auth_ok" + return st + } + + args := []string{"exec", "--skip-git-repo-check", "--ephemeral", "--json", cfg.CodexProbePrompt} + args = append(args, cfg.CodexExtraArgs...) + started := time.Now() + out, exitCode, err := runCommand(parent, cfg.ProbeTimeout, cfg.CodexCommand, args...) + st.LatencyMS = time.Since(started).Milliseconds() + st.ExitCode = exitCode + if err != nil && strings.TrimSpace(out) == "" { + st.Status = classifyFailure(out, err.Error()) + st.Message = compact(err.Error(), 500) + st.Limited = st.Status == "limited" + return st + } + usage := parseCodexUsage(out) + if usage != nil { + st.Usage = usage + } + if exitCode == 0 && strings.Contains(out, "turn.completed") { + st.Available = true + st.Status = "ok" + st.Message = "OK" + return st + } + st.Status = classifyFailure(out, "") + st.Limited = st.Status == "limited" + st.Message = compact(out, 800) + st.LimitResetText = extractReset(out) + return st +} + +func baseStatus(provider string) ProviderStatus { + return ProviderStatus{Provider: provider, Status: "unknown", CheckedAt: time.Now().UTC().Format(time.RFC3339), ExitCode: -1} +} + +func runCommand(parent context.Context, timeout time.Duration, name string, args ...string) (string, int, error) { + ctx, cancel := context.WithTimeout(parent, timeout) + defer cancel() + cmd := exec.CommandContext(ctx, name, args...) + cmd.Env = os.Environ() + var buf bytes.Buffer + cmd.Stdout = &buf + cmd.Stderr = &buf + err := cmd.Run() + out := buf.String() + if ctx.Err() != nil { + return out, -1, ctx.Err() + } + if err == nil { + return out, 0, nil + } + var ee *exec.ExitError + if errors.As(err, &ee) { + return out, ee.ExitCode(), err + } + return out, -1, err +} + +func runShort(parent context.Context, timeout time.Duration, name string, args ...string) string { + out, _, err := runCommand(parent, timeout, name, args...) + if err != nil && strings.TrimSpace(out) == "" { + return err.Error() + } + return out +} + +func classifyFailure(parts ...string) string { + s := strings.ToLower(strings.Join(parts, "\n")) + switch { + case strings.Contains(s, "rate_limit") || strings.Contains(s, "rate limit") || strings.Contains(s, "usage limit") || strings.Contains(s, "quota") || strings.Contains(s, "too many requests") || strings.Contains(s, "429"): + return "limited" + case strings.Contains(s, "not logged") || strings.Contains(s, "unauthorized") || strings.Contains(s, "authentication") || strings.Contains(s, "auth") && strings.Contains(s, "fail"): + return "auth_required" + case strings.Contains(s, "deadline exceeded") || strings.Contains(s, "timeout"): + return "timeout" + default: + return "error" + } +} + +func extractReset(s string) string { + for _, line := range strings.Split(s, "\n") { + l := strings.ToLower(line) + if strings.Contains(l, "reset") || strings.Contains(l, "try again") || strings.Contains(l, "wait") { + return compact(strings.TrimSpace(line), 240) + } + } + return "" +} + +func parseCodexUsage(out string) map[string]interface{} { + scanner := bufio.NewScanner(strings.NewReader(out)) + var usage map[string]interface{} + for scanner.Scan() { + var event map[string]interface{} + if json.Unmarshal([]byte(scanner.Text()), &event) != nil { + continue + } + if u, ok := event["usage"].(map[string]interface{}); ok { + usage = u + } + } + return usage +} + +func attachLimitUsage(ctx context.Context, cfg Config, st *ProviderStatus) { + now := time.Now() + sessionCutoff := now.Add(-5 * time.Hour) + weekCutoff := now.Add(-7 * 24 * time.Hour) + switch st.Provider { + case "claude": + st.Session5h = makeWindow("5h", readClaudeTokens(cfg.ClaudeHistoryDir, sessionCutoff), cfg.ClaudeSession5hLimit) + st.Week = makeWindow("7d", readClaudeTokens(cfg.ClaudeHistoryDir, weekCutoff), cfg.ClaudeWeekLimit) + case "codex": + st.Session5h = makeWindow("5h", readCodexTokens(ctx, cfg.CodexStateDB, sessionCutoff), cfg.CodexSession5hLimit) + st.Week = makeWindow("7d", readCodexTokens(ctx, cfg.CodexStateDB, weekCutoff), cfg.CodexWeekLimit) + default: + st.Session5h = makeWindow("5h", 0, 0) + st.Week = makeWindow("7d", 0, 0) + } +} + +func makeWindow(name string, used, limit int64) LimitWindow { + w := LimitWindow{Window: name, UsedTokens: used, LimitTokens: limit} + if limit > 0 { + pct := math.Round((float64(used)/float64(limit))*1000) / 10 + w.PercentUsed = &pct + } + return w +} + +func readClaudeTokens(root string, cutoff time.Time) int64 { + var total int64 + _ = filepath.WalkDir(root, func(path string, d os.DirEntry, err error) error { + if err != nil || d.IsDir() || !strings.HasSuffix(path, ".jsonl") { + return nil + } + f, err := os.Open(path) + if err != nil { + return nil + } + defer f.Close() + scanner := bufio.NewScanner(f) + scanner.Buffer(make([]byte, 0, 64*1024), 8*1024*1024) + for scanner.Scan() { + var event map[string]interface{} + if json.Unmarshal(scanner.Bytes(), &event) != nil { + continue + } + ts, ok := parseEventTime(event["timestamp"]) + if !ok || ts.Before(cutoff) { + continue + } + msg, _ := event["message"].(map[string]interface{}) + usage, _ := msg["usage"].(map[string]interface{}) + if usage == nil { + continue + } + total += usageTokenSum(usage) + } + return nil + }) + return total +} + +func usageTokenSum(usage map[string]interface{}) int64 { + keys := []string{"input_tokens", "output_tokens", "cache_creation_input_tokens", "cache_read_input_tokens", "reasoning_output_tokens"} + var total int64 + for _, key := range keys { + switch v := usage[key].(type) { + case float64: + total += int64(v) + case int64: + total += v + case int: + total += int64(v) + } + } + return total +} + +func parseEventTime(v interface{}) (time.Time, bool) { + s, ok := v.(string) + if !ok || s == "" { + return time.Time{}, false + } + for _, layout := range []string{time.RFC3339Nano, time.RFC3339} { + if ts, err := time.Parse(layout, s); err == nil { + return ts, true + } + } + return time.Time{}, false +} + +func readCodexTokens(ctx context.Context, db string, cutoff time.Time) int64 { + if db == "" { + return 0 + } + if _, err := os.Stat(db); err != nil { + return 0 + } + cutoffMS := cutoff.UnixMilli() + query := fmt.Sprintf("select coalesce(sum(tokens_used),0) from threads where updated_at_ms >= %d;", cutoffMS) + out, _, err := runCommand(ctx, 5*time.Second, "sqlite3", db, query) + if err != nil { + log.Printf("codex token usage query failed db=%s: %v", db, err) + return 0 + } + n, _ := strconv.ParseInt(strings.TrimSpace(out), 10, 64) + return n +} + +func lastJSON(out string) string { + lines := strings.Split(strings.TrimSpace(out), "\n") + for i := len(lines) - 1; i >= 0; i-- { + line := strings.TrimSpace(lines[i]) + if strings.HasPrefix(line, "{") && strings.HasSuffix(line, "}") { + return line + } + } + return strings.TrimSpace(out) +} + +func stringField(m map[string]interface{}, key string) string { + if v, ok := m[key].(string); ok { + return v + } + return "" +} + +func publishDiscovery(client mqtt.Client, cfg Config) { + for _, p := range []string{"claude", "codex"} { + name := "AI Limits " + titleWord(p) + unique := "ai_limits_" + p + configTopic := fmt.Sprintf("%s/sensor/%s/config", cfg.HADiscoveryPrefix, unique) + payload := map[string]interface{}{ + "name": name, + "unique_id": unique, + "state_topic": cfg.topic(p, "state"), + "value_template": "{{ value_json.status }}", + "json_attributes_topic": cfg.topic(p, "state"), + "availability_topic": cfg.topic("availability"), + "icon": "mdi:robot", + "device": discoveryDevice(), + } + publishJSON(client, configTopic, payload, true) + publishPercentDiscovery(client, cfg, p, "session_5h", "5h Used") + publishPercentDiscovery(client, cfg, p, "week", "Week Used") + } + summaryTopic := fmt.Sprintf("%s/sensor/ai_limits_summary/config", cfg.HADiscoveryPrefix) + publishJSON(client, summaryTopic, map[string]interface{}{ + "name": "AI Limits Summary", + "unique_id": "ai_limits_summary", + "state_topic": cfg.topic("summary", "state"), + "value_template": "{{ value_json.status }}", + "json_attributes_topic": cfg.topic("summary", "state"), + "availability_topic": cfg.topic("availability"), + "icon": "mdi:gauge", + "device": discoveryDevice(), + }, true) +} + +func publishPercentDiscovery(client mqtt.Client, cfg Config, provider, field, label string) { + unique := fmt.Sprintf("ai_limits_%s_%s_percent", provider, field) + configTopic := fmt.Sprintf("%s/sensor/%s/config", cfg.HADiscoveryPrefix, unique) + template := fmt.Sprintf("{{ value_json.%s.percent_used | default('unknown') }}", field) + publishJSON(client, configTopic, map[string]interface{}{ + "name": fmt.Sprintf("AI Limits %s %s", titleWord(provider), label), + "unique_id": unique, + "state_topic": cfg.topic(provider, "state"), + "value_template": template, + "availability_topic": cfg.topic("availability"), + "unit_of_measurement": "%", + "state_class": "measurement", + "icon": "mdi:gauge", + "device": discoveryDevice(), + }, true) +} + +func discoveryDevice() map[string]interface{} { + return map[string]interface{}{ + "identifiers": []string{"ai_limits_mqtt"}, + "name": "AI Limits Monitor", + "manufacturer": "Hermes", + "model": "Go MQTT probe", + } +} + +func publishJSON(client mqtt.Client, topic string, v interface{}, retained bool) { + b, err := json.Marshal(v) + if err != nil { + log.Printf("json marshal failed for %s: %v", topic, err) + return + } + publish(client, topic, string(b), retained) +} + +func publish(client mqtt.Client, topic, payload string, retained bool) { + token := client.Publish(topic, 1, retained, payload) + if !token.WaitTimeout(10 * time.Second) { + log.Printf("mqtt publish timed out topic=%s", topic) + return + } + if token.Error() != nil { + log.Printf("mqtt publish failed topic=%s: %v", topic, token.Error()) + } +} + +func (c Config) topic(parts ...string) string { + all := append([]string{c.MQTTBaseTopic}, parts...) + return strings.Join(all, "/") +} + +func env(key, fallback string) string { + if v := os.Getenv(key); v != "" { + return v + } + return fallback +} + +func envBool(key string, fallback bool) bool { + v := os.Getenv(key) + if v == "" { + return fallback + } + b, err := strconv.ParseBool(v) + if err != nil { + return fallback + } + return b +} + +func envInt64(key string, fallback int64) int64 { + v := os.Getenv(key) + if v == "" { + return fallback + } + n, err := strconv.ParseInt(v, 10, 64) + if err != nil { + return fallback + } + return n +} + +func envDuration(key string, fallback time.Duration) time.Duration { + v := os.Getenv(key) + if v == "" { + return fallback + } + if d, err := time.ParseDuration(v); err == nil { + return d + } + if n, err := strconv.Atoi(v); err == nil { + return time.Duration(n) * time.Second + } + return fallback +} + +func splitArgs(s string) []string { + fields := strings.Fields(s) + if len(fields) == 0 { + return nil + } + return fields +} + +func trimSlashes(s string) string { return strings.Trim(s, "/") } +func titleWord(s string) string { + if s == "" { + return s + } + return strings.ToUpper(s[:1]) + s[1:] +} +func firstLine(s string) string { + s = strings.TrimSpace(s) + if i := strings.IndexByte(s, '\n'); i >= 0 { + return s[:i] + } + return s +} +func compact(s string, max int) string { + s = strings.TrimSpace(strings.Join(strings.Fields(s), " ")) + if len(s) <= max { + return s + } + h := sha1.Sum([]byte(s)) + return s[:max] + "… sha1=" + hex.EncodeToString(h[:4]) +}