feat: add AI limits MQTT service

This commit is contained in:
Gemma
2026-05-21 09:32:37 +02:00
commit 0788cffc88
7 changed files with 738 additions and 0 deletions

3
.gitignore vendored Normal file
View File

@@ -0,0 +1,3 @@
.env
*.log
ai-limits-mqtt

15
Dockerfile Normal file
View File

@@ -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"]

52
README.md Normal file
View File

@@ -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.

12
docker-compose.yml Normal file
View File

@@ -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

11
go.mod Normal file
View File

@@ -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
)

8
go.sum Normal file
View File

@@ -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=

637
main.go Normal file
View File

@@ -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])
}