From 87f091beb8778444b4783dea23dde858c1991a67 Mon Sep 17 00:00:00 2001 From: tiennm99 Date: Fri, 19 Jun 2026 14:39:19 +0700 Subject: [PATCH] feat: scaffold openai status bot --- .env.example | 8 + .gitignore | 9 ++ Dockerfile | 13 ++ README.md | 55 +++++++ cmd/openai-status-bot/main.go | 59 +++++++ docker-compose.yml | 21 +++ docs/setup-guide.md | 40 +++++ docs/system-architecture.md | 42 +++++ go.mod | 10 ++ go.sum | 10 ++ internal/bot/bot.go | 153 ++++++++++++++++++ internal/bot/format.go | 116 ++++++++++++++ internal/bot/format_test.go | 34 ++++ internal/config/config.go | 97 ++++++++++++ internal/openai/client.go | 63 ++++++++ internal/openai/client_test.go | 46 ++++++ internal/openai/models.go | 58 +++++++ internal/poller/format.go | 79 ++++++++++ internal/poller/format_test.go | 44 ++++++ internal/poller/poller.go | 172 +++++++++++++++++++++ internal/redisstore/store.go | 117 ++++++++++++++ internal/redisstore/store_test.go | 32 ++++ internal/telegram/client.go | 94 +++++++++++ internal/telegram/types.go | 19 +++ plans/260619-1401-initial-scaffold/plan.md | 27 ++++ 25 files changed, 1418 insertions(+) create mode 100644 .env.example create mode 100644 .gitignore create mode 100644 Dockerfile create mode 100644 README.md create mode 100644 cmd/openai-status-bot/main.go create mode 100644 docker-compose.yml create mode 100644 docs/setup-guide.md create mode 100644 docs/system-architecture.md create mode 100644 go.mod create mode 100644 go.sum create mode 100644 internal/bot/bot.go create mode 100644 internal/bot/format.go create mode 100644 internal/bot/format_test.go create mode 100644 internal/config/config.go create mode 100644 internal/openai/client.go create mode 100644 internal/openai/client_test.go create mode 100644 internal/openai/models.go create mode 100644 internal/poller/format.go create mode 100644 internal/poller/format_test.go create mode 100644 internal/poller/poller.go create mode 100644 internal/redisstore/store.go create mode 100644 internal/redisstore/store_test.go create mode 100644 internal/telegram/client.go create mode 100644 internal/telegram/types.go create mode 100644 plans/260619-1401-initial-scaffold/plan.md diff --git a/.env.example b/.env.example new file mode 100644 index 0000000..bec0de6 --- /dev/null +++ b/.env.example @@ -0,0 +1,8 @@ +TELEGRAM_BOT_TOKEN= +REDIS_ADDR=localhost:6379 +REDIS_PASSWORD= +REDIS_DB=0 +OPENAI_STATUS_BASE_URL=https://status.openai.com +POLL_INTERVAL=1m +HTTP_TIMEOUT=10s +LOG_LEVEL=info diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..ba28e32 --- /dev/null +++ b/.gitignore @@ -0,0 +1,9 @@ +.env +.env.* +!.env.example + +/openai-status-bot +/dist/ +/tmp/ + +*.log diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..4e09b9c --- /dev/null +++ b/Dockerfile @@ -0,0 +1,13 @@ +FROM golang:1.24-alpine AS build + +WORKDIR /src +COPY go.mod go.sum* ./ +RUN go mod download +COPY . . +RUN CGO_ENABLED=0 GOOS=linux go build -o /out/openai-status-bot ./cmd/openai-status-bot + +FROM alpine:3.21 +RUN adduser -D -H appuser +USER appuser +COPY --from=build /out/openai-status-bot /usr/local/bin/openai-status-bot +ENTRYPOINT ["openai-status-bot"] diff --git a/README.md b/README.md new file mode 100644 index 0000000..18cac83 --- /dev/null +++ b/README.md @@ -0,0 +1,55 @@ +# openai-status-bot + +Telegram bot that watches [OpenAI Status](https://status.openai.com/) every minute and sends updates to subscribed chats. State and subscriptions are stored in Redis. + +## Features + +- Checks OpenAI status on a configurable interval, default `1m` +- Notifies subscribers about new incident updates +- Notifies subscribers when component status changes +- Uses Redis for subscribers, component status checkpoints, and seen incident updates +- Supports Telegram supergroup topics via `message_thread_id` +- Includes Docker Compose for local Redis + bot runtime + +## Bot Commands + +| Command | Description | +|---------|-------------| +| `/start` | Subscribe current chat or topic | +| `/stop` | Unsubscribe current chat or topic | +| `/status` | Show current OpenAI status | +| `/components` | Show all OpenAI component statuses | +| `/history [count]` | Show recent incidents, default 5, max 10 | +| `/help` | Show command help | + +## Quick Start + +```bash +cp .env.example .env +# edit .env and set TELEGRAM_BOT_TOKEN +docker compose up --build +``` + +For local development without Docker: + +```bash +go mod tidy +go run ./cmd/openai-status-bot +``` + +## Configuration + +| Variable | Default | Description | +|----------|---------|-------------| +| `TELEGRAM_BOT_TOKEN` | required | Telegram bot token from BotFather | +| `REDIS_ADDR` | `localhost:6379` | Redis address | +| `REDIS_PASSWORD` | empty | Redis password | +| `REDIS_DB` | `0` | Redis database number | +| `OPENAI_STATUS_BASE_URL` | `https://status.openai.com` | OpenAI status page base URL | +| `POLL_INTERVAL` | `1m` | Status check interval | +| `HTTP_TIMEOUT` | `10s` | HTTP request timeout | +| `LOG_LEVEL` | `info` | `debug`, `info`, `warn`, or `error` | + +## Notes + +The first successful poll seeds Redis and does not send historical incidents. Notifications start from later changes. diff --git a/cmd/openai-status-bot/main.go b/cmd/openai-status-bot/main.go new file mode 100644 index 0000000..bfa06ec --- /dev/null +++ b/cmd/openai-status-bot/main.go @@ -0,0 +1,59 @@ +package main + +import ( + "context" + "log/slog" + "os" + "os/signal" + "syscall" + + "github.com/redis/go-redis/v9" + "github.com/tiennm99/openai-status-bot/internal/bot" + "github.com/tiennm99/openai-status-bot/internal/config" + openai "github.com/tiennm99/openai-status-bot/internal/openai" + "github.com/tiennm99/openai-status-bot/internal/poller" + "github.com/tiennm99/openai-status-bot/internal/redisstore" + "github.com/tiennm99/openai-status-bot/internal/telegram" +) + +func main() { + cfg, err := config.LoadFromEnv() + if err != nil { + slog.Error("load config", "error", err) + os.Exit(1) + } + + logger := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{ + Level: cfg.LogLevel, + })) + slog.SetDefault(logger) + + ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer stop() + + redisClient := redis.NewClient(&redis.Options{ + Addr: cfg.RedisAddr, + Password: cfg.RedisPassword, + DB: cfg.RedisDB, + }) + if err := redisClient.Ping(ctx).Err(); err != nil { + logger.Error("connect redis", "addr", cfg.RedisAddr, "error", err) + os.Exit(1) + } + defer redisClient.Close() + + store := redisstore.New(redisClient) + statusClient := openai.NewClient(cfg.OpenAIStatusBaseURL, cfg.HTTPTimeout) + telegramClient := telegram.NewClient(cfg.TelegramBotToken, cfg.HTTPTimeout) + + statusPoller := poller.NewRunner(statusClient, store, telegramClient, cfg.PollInterval, logger) + commandBot := bot.New(telegramClient, statusClient, store, logger) + + go statusPoller.Run(ctx) + + logger.Info("openai status bot started", "poll_interval", cfg.PollInterval.String()) + if err := commandBot.Run(ctx); err != nil && ctx.Err() == nil { + logger.Error("telegram bot stopped", "error", err) + os.Exit(1) + } +} diff --git a/docker-compose.yml b/docker-compose.yml new file mode 100644 index 0000000..97d4849 --- /dev/null +++ b/docker-compose.yml @@ -0,0 +1,21 @@ +services: + redis: + image: redis:7-alpine + ports: + - "6379:6379" + volumes: + - redis-data:/data + command: ["redis-server", "--appendonly", "yes"] + + bot: + build: . + env_file: + - .env + environment: + REDIS_ADDR: redis:6379 + depends_on: + - redis + restart: unless-stopped + +volumes: + redis-data: diff --git a/docs/setup-guide.md b/docs/setup-guide.md new file mode 100644 index 0000000..5480a42 --- /dev/null +++ b/docs/setup-guide.md @@ -0,0 +1,40 @@ +# Setup Guide + +## Prerequisites + +- Go 1.24+ +- Redis 7+ +- Telegram bot token from BotFather + +## Local Docker Run + +```bash +cp .env.example .env +# set TELEGRAM_BOT_TOKEN in .env +docker compose up --build +``` + +## Local Go Run + +Start Redis: + +```bash +docker run --rm -p 6379:6379 redis:7-alpine +``` + +Run bot: + +```bash +cp .env.example .env +# export TELEGRAM_BOT_TOKEN or source .env with your shell workflow +go run ./cmd/openai-status-bot +``` + +## Verification + +```bash +go test ./... +go build ./cmd/openai-status-bot +``` + +Then send `/start` to the Telegram bot and run `/status`. diff --git a/docs/system-architecture.md b/docs/system-architecture.md new file mode 100644 index 0000000..3b1694c --- /dev/null +++ b/docs/system-architecture.md @@ -0,0 +1,42 @@ +# System Architecture + +## Overview + +`openai-status-bot` is one Go process with two loops: + +- Telegram long polling loop for user commands. +- OpenAI status polling loop, default every minute. + +Redis stores subscribers and polling checkpoints. + +## Data Flow + +1. User sends `/start` in Telegram. +2. Bot stores chat ID and optional topic thread ID in Redis. +3. Poller fetches OpenAI status JSON: + - `GET /api/v2/summary.json` + - `GET /api/v2/incidents.json` +4. Poller compares fetched state with Redis checkpoints. +5. New incident updates or component status changes are sent to all subscribers. + +## Redis Keys + +| Key | Type | Purpose | +|-----|------|---------| +| `openai-status:subscribers` | set | Telegram chat or topic subscribers | +| `openai-status:component-statuses` | hash | Last seen component status by component ID | +| `openai-status:incident-updates` | set | Seen incident update IDs | +| `openai-status:initialized` | string | Baseline seed marker | + +Subscriber set members are `chatID` or `chatID:threadID`. + +## Runtime + +The service uses Telegram `getUpdates`, so it does not need a public webhook URL. Docker Compose starts Redis and the bot. + +## Failure Behavior + +- First successful poll seeds state without notification. +- OpenAI fetch failures are logged and retried on the next interval. +- Telegram send failures are logged; other subscribers still receive messages. +- Redis connection failure at startup exits the process. diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..b9ce74c --- /dev/null +++ b/go.mod @@ -0,0 +1,10 @@ +module github.com/tiennm99/openai-status-bot + +go 1.24 + +require github.com/redis/go-redis/v9 v9.17.0 + +require ( + github.com/cespare/xxhash/v2 v2.3.0 // indirect + github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f // indirect +) diff --git a/go.sum b/go.sum new file mode 100644 index 0000000..72d71ad --- /dev/null +++ b/go.sum @@ -0,0 +1,10 @@ +github.com/bsm/ginkgo/v2 v2.12.0 h1:Ny8MWAHyOepLGlLKYmXG4IEkioBysk6GpaRTLC8zwWs= +github.com/bsm/ginkgo/v2 v2.12.0/go.mod h1:SwYbGRRDovPVboqFv0tPTcG1sN61LM1Z4ARdbAV9g4c= +github.com/bsm/gomega v1.27.10 h1:yeMWxP2pV2fG3FgAODIY8EiRE3dy0aeFYt4l7wh6yKA= +github.com/bsm/gomega v1.27.10/go.mod h1:JyEr/xRbxbtgWNi8tIEVPUYZ5Dzef52k01W3YH0H+O0= +github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= +github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f h1:lO4WD4F/rVNCu3HqELle0jiPLLBs70cWOduZpkS1E78= +github.com/dgryski/go-rendezvous v0.0.0-20200823014737-9f7001d12a5f/go.mod h1:cuUVRXasLTGF7a8hSLbxyZXjz+1KgoB3wDUb6vlszIc= +github.com/redis/go-redis/v9 v9.17.0 h1:K6E+ZlYN95KSMmZeEQPbU/c++wfmEvfFB17yEAq/VhM= +github.com/redis/go-redis/v9 v9.17.0/go.mod h1:u410H11HMLoB+TP67dz8rL9s6QW2j76l0//kSOd3370= diff --git a/internal/bot/bot.go b/internal/bot/bot.go new file mode 100644 index 0000000..1799237 --- /dev/null +++ b/internal/bot/bot.go @@ -0,0 +1,153 @@ +package bot + +import ( + "context" + "log/slog" + "strings" + "time" + + openai "github.com/tiennm99/openai-status-bot/internal/openai" + "github.com/tiennm99/openai-status-bot/internal/redisstore" + "github.com/tiennm99/openai-status-bot/internal/telegram" +) + +type TelegramClient interface { + GetUpdates(ctx context.Context, offset int64, timeoutSeconds int) ([]telegram.Update, error) + SendText(ctx context.Context, chatID int64, threadID *int, text string) error +} + +type StatusClient interface { + FetchSummary(ctx context.Context) (openai.Summary, error) + FetchIncidents(ctx context.Context) (openai.IncidentsResponse, error) +} + +type Store interface { + AddSubscriber(ctx context.Context, sub redisstore.Subscriber) error + RemoveSubscriber(ctx context.Context, sub redisstore.Subscriber) error +} + +type Bot struct { + telegramClient TelegramClient + statusClient StatusClient + store Store + logger *slog.Logger +} + +func New(telegramClient TelegramClient, statusClient StatusClient, store Store, logger *slog.Logger) *Bot { + return &Bot{ + telegramClient: telegramClient, + statusClient: statusClient, + store: store, + logger: logger, + } +} + +func (b *Bot) Run(ctx context.Context) error { + var offset int64 + + for { + select { + case <-ctx.Done(): + return nil + default: + } + + updates, err := b.telegramClient.GetUpdates(ctx, offset, 50) + if err != nil { + if ctx.Err() != nil { + return nil + } + b.logger.Warn("get telegram updates", "error", err) + time.Sleep(3 * time.Second) + continue + } + + for _, update := range updates { + offset = update.UpdateID + 1 + if update.Message == nil { + continue + } + b.handleMessage(ctx, *update.Message) + } + } +} + +func (b *Bot) handleMessage(ctx context.Context, message telegram.Message) { + if !strings.HasPrefix(strings.TrimSpace(message.Text), "/") { + return + } + + command, fields := normalizeCommand(message.Text) + switch command { + case "/start": + b.subscribe(ctx, message) + case "/stop": + b.unsubscribe(ctx, message) + case "/status": + b.replyStatus(ctx, message) + case "/components": + b.replyComponents(ctx, message) + case "/history": + b.replyHistory(ctx, message, parseHistoryCount(fields)) + case "/help": + b.reply(ctx, message, helpText) + default: + b.reply(ctx, message, "Unknown command. Use /help.") + } +} + +func (b *Bot) subscribe(ctx context.Context, message telegram.Message) { + sub := redisstore.NewSubscriber(message.Chat.ID, message.MessageThreadID) + if err := b.store.AddSubscriber(ctx, sub); err != nil { + b.logger.Error("subscribe", "error", err) + b.reply(ctx, message, "Could not subscribe right now.") + return + } + b.reply(ctx, message, "Subscribed to OpenAI status updates.") +} + +func (b *Bot) unsubscribe(ctx context.Context, message telegram.Message) { + sub := redisstore.NewSubscriber(message.Chat.ID, message.MessageThreadID) + if err := b.store.RemoveSubscriber(ctx, sub); err != nil { + b.logger.Error("unsubscribe", "error", err) + b.reply(ctx, message, "Could not unsubscribe right now.") + return + } + b.reply(ctx, message, "Unsubscribed from OpenAI status updates.") +} + +func (b *Bot) replyStatus(ctx context.Context, message telegram.Message) { + summary, err := b.statusClient.FetchSummary(ctx) + if err != nil { + b.logger.Error("fetch status", "error", err) + b.reply(ctx, message, "Could not fetch OpenAI status right now.") + return + } + b.reply(ctx, message, formatStatus(summary)) +} + +func (b *Bot) replyComponents(ctx context.Context, message telegram.Message) { + summary, err := b.statusClient.FetchSummary(ctx) + if err != nil { + b.logger.Error("fetch components", "error", err) + b.reply(ctx, message, "Could not fetch OpenAI components right now.") + return + } + b.reply(ctx, message, formatComponents(summary)) +} + +func (b *Bot) replyHistory(ctx context.Context, message telegram.Message, count int) { + incidents, err := b.statusClient.FetchIncidents(ctx) + if err != nil { + b.logger.Error("fetch incidents", "error", err) + b.reply(ctx, message, "Could not fetch OpenAI incident history right now.") + return + } + b.reply(ctx, message, formatHistory(incidents.Incidents, count)) +} + +func (b *Bot) reply(ctx context.Context, message telegram.Message, text string) { + if err := b.telegramClient.SendText(ctx, message.Chat.ID, message.MessageThreadID, text); err != nil { + b.logger.Warn("send telegram reply", "chat_id", message.Chat.ID, "error", err) + } +} diff --git a/internal/bot/format.go b/internal/bot/format.go new file mode 100644 index 0000000..5aab1b1 --- /dev/null +++ b/internal/bot/format.go @@ -0,0 +1,116 @@ +package bot + +import ( + "fmt" + "strconv" + "strings" + + openai "github.com/tiennm99/openai-status-bot/internal/openai" + "github.com/tiennm99/openai-status-bot/internal/poller" +) + +const helpText = `OpenAI Status Bot + +/start - subscribe this chat or topic +/stop - unsubscribe this chat or topic +/status - show current OpenAI status +/components - show all component statuses +/history [count] - show recent incidents, default 5, max 10 +/help - show this help` + +func formatStatus(summary openai.Summary) string { + lines := []string{ + "OpenAI status", + "", + fmt.Sprintf("Overall: %s", summary.Status.Description), + } + + degraded := make([]string, 0) + for _, component := range summary.Components { + if component.Group || component.Status == "operational" { + continue + } + degraded = append(degraded, fmt.Sprintf("- %s: %s", component.Name, poller.StatusLabel(component.Status))) + } + + if len(degraded) == 0 { + lines = append(lines, "", "All listed components are operational.") + } else { + lines = append(lines, "", "Affected components:") + lines = append(lines, degraded...) + } + lines = append(lines, "", "https://status.openai.com/") + + return strings.Join(lines, "\n") +} + +func formatComponents(summary openai.Summary) string { + lines := []string{"OpenAI components", ""} + for _, component := range summary.Components { + if component.Group { + continue + } + lines = append(lines, fmt.Sprintf("- %s: %s", component.Name, poller.StatusLabel(component.Status))) + } + return truncateMessage(strings.Join(lines, "\n")) +} + +func formatHistory(incidents []openai.Incident, count int) string { + if len(incidents) == 0 { + return "No recent incidents found." + } + if count > len(incidents) { + count = len(incidents) + } + + lines := []string{"Recent OpenAI incidents", ""} + for i := 0; i < count; i++ { + incident := incidents[i] + lines = append(lines, fmt.Sprintf( + "%d. %s\n Status: %s | Impact: %s", + i+1, + incident.Name, + poller.StatusLabel(incident.Status), + poller.StatusLabel(incident.Impact), + )) + } + return truncateMessage(strings.Join(lines, "\n")) +} + +func parseHistoryCount(fields []string) int { + const ( + defaultCount = 5 + maxCount = 10 + ) + if len(fields) < 2 { + return defaultCount + } + count, err := strconv.Atoi(fields[1]) + if err != nil || count < 1 { + return defaultCount + } + if count > maxCount { + return maxCount + } + return count +} + +func normalizeCommand(text string) (string, []string) { + fields := strings.Fields(strings.TrimSpace(text)) + if len(fields) == 0 { + return "", nil + } + command := strings.ToLower(fields[0]) + if at := strings.Index(command, "@"); at >= 0 { + command = command[:at] + } + return command, fields +} + +func truncateMessage(value string) string { + const telegramLimit = 3900 + if len(value) <= telegramLimit { + return value + } + return strings.TrimSpace(value[:telegramLimit-3]) + "..." +} diff --git a/internal/bot/format_test.go b/internal/bot/format_test.go new file mode 100644 index 0000000..79b6e3f --- /dev/null +++ b/internal/bot/format_test.go @@ -0,0 +1,34 @@ +package bot + +import "testing" + +func TestNormalizeCommandStripsBotUsername(t *testing.T) { + command, fields := normalizeCommand("/history@OpenAIStatusBot 10") + if command != "/history" { + t.Fatalf("command = %q, want /history", command) + } + if len(fields) != 2 || fields[1] != "10" { + t.Fatalf("fields = %#v", fields) + } +} + +func TestParseHistoryCount(t *testing.T) { + tests := []struct { + name string + fields []string + want int + }{ + {name: "default", fields: []string{"/history"}, want: 5}, + {name: "valid", fields: []string{"/history", "3"}, want: 3}, + {name: "invalid", fields: []string{"/history", "abc"}, want: 5}, + {name: "max", fields: []string{"/history", "99"}, want: 10}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := parseHistoryCount(tt.fields); got != tt.want { + t.Fatalf("parseHistoryCount(%v) = %d, want %d", tt.fields, got, tt.want) + } + }) + } +} diff --git a/internal/config/config.go b/internal/config/config.go new file mode 100644 index 0000000..2a737a3 --- /dev/null +++ b/internal/config/config.go @@ -0,0 +1,97 @@ +package config + +import ( + "fmt" + "log/slog" + "os" + "strconv" + "strings" + "time" +) + +type Config struct { + TelegramBotToken string + RedisAddr string + RedisPassword string + RedisDB int + OpenAIStatusBaseURL string + PollInterval time.Duration + HTTPTimeout time.Duration + LogLevel slog.Level +} + +func LoadFromEnv() (Config, error) { + cfg := Config{ + TelegramBotToken: strings.TrimSpace(os.Getenv("TELEGRAM_BOT_TOKEN")), + RedisAddr: getEnv("REDIS_ADDR", "localhost:6379"), + RedisPassword: os.Getenv("REDIS_PASSWORD"), + OpenAIStatusBaseURL: strings.TrimRight(getEnv("OPENAI_STATUS_BASE_URL", "https://status.openai.com"), "/"), + } + if cfg.TelegramBotToken == "" { + return Config{}, fmt.Errorf("TELEGRAM_BOT_TOKEN is required") + } + + var err error + cfg.RedisDB, err = parseIntEnv("REDIS_DB", 0) + if err != nil { + return Config{}, err + } + cfg.PollInterval, err = parseDurationEnv("POLL_INTERVAL", time.Minute) + if err != nil { + return Config{}, err + } + cfg.HTTPTimeout, err = parseDurationEnv("HTTP_TIMEOUT", 10*time.Second) + if err != nil { + return Config{}, err + } + cfg.LogLevel = parseLogLevel(getEnv("LOG_LEVEL", "info")) + + return cfg, nil +} + +func getEnv(key, fallback string) string { + if value := strings.TrimSpace(os.Getenv(key)); value != "" { + return value + } + return fallback +} + +func parseIntEnv(key string, fallback int) (int, error) { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback, nil + } + parsed, err := strconv.Atoi(value) + if err != nil { + return 0, fmt.Errorf("%s must be an integer: %w", key, err) + } + return parsed, nil +} + +func parseDurationEnv(key string, fallback time.Duration) (time.Duration, error) { + value := strings.TrimSpace(os.Getenv(key)) + if value == "" { + return fallback, nil + } + parsed, err := time.ParseDuration(value) + if err != nil { + return 0, fmt.Errorf("%s must be a duration: %w", key, err) + } + if parsed <= 0 { + return 0, fmt.Errorf("%s must be positive", key) + } + return parsed, nil +} + +func parseLogLevel(value string) slog.Level { + switch strings.ToLower(strings.TrimSpace(value)) { + case "debug": + return slog.LevelDebug + case "warn", "warning": + return slog.LevelWarn + case "error": + return slog.LevelError + default: + return slog.LevelInfo + } +} diff --git a/internal/openai/client.go b/internal/openai/client.go new file mode 100644 index 0000000..4f7fbb4 --- /dev/null +++ b/internal/openai/client.go @@ -0,0 +1,63 @@ +package openai + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "strings" + "time" +) + +type Client struct { + baseURL string + httpClient *http.Client +} + +func NewClient(baseURL string, timeout time.Duration) *Client { + return &Client{ + baseURL: strings.TrimRight(baseURL, "/"), + httpClient: &http.Client{ + Timeout: timeout, + }, + } +} + +func (c *Client) FetchSummary(ctx context.Context) (Summary, error) { + var summary Summary + if err := c.getJSON(ctx, "/api/v2/summary.json", &summary); err != nil { + return Summary{}, err + } + return summary, nil +} + +func (c *Client) FetchIncidents(ctx context.Context) (IncidentsResponse, error) { + var incidents IncidentsResponse + if err := c.getJSON(ctx, "/api/v2/incidents.json", &incidents); err != nil { + return IncidentsResponse{}, err + } + return incidents, nil +} + +func (c *Client) getJSON(ctx context.Context, path string, target any) error { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.baseURL+path, nil) + if err != nil { + return err + } + req.Header.Set("Accept", "application/json") + req.Header.Set("User-Agent", "openai-status-bot/1.0") + + res, err := c.httpClient.Do(req) + if err != nil { + return err + } + defer res.Body.Close() + + if res.StatusCode < 200 || res.StatusCode >= 300 { + return fmt.Errorf("openai status API returned %s", res.Status) + } + if err := json.NewDecoder(res.Body).Decode(target); err != nil { + return fmt.Errorf("decode openai status response: %w", err) + } + return nil +} diff --git a/internal/openai/client_test.go b/internal/openai/client_test.go new file mode 100644 index 0000000..97ca842 --- /dev/null +++ b/internal/openai/client_test.go @@ -0,0 +1,46 @@ +package openai + +import ( + "context" + "net/http" + "net/http/httptest" + "testing" + "time" +) + +func TestFetchSummary(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/api/v2/summary.json" { + t.Fatalf("unexpected path %s", r.URL.Path) + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{ + "page": {"id": "page", "name": "OpenAI"}, + "status": {"description": "All Systems Operational", "indicator": "none"}, + "components": [{"id": "c1", "name": "Codex Web", "status": "operational"}], + "incidents": [] + }`)) + })) + defer server.Close() + + client := NewClient(server.URL, time.Second) + summary, err := client.FetchSummary(context.Background()) + if err != nil { + t.Fatalf("FetchSummary returned error: %v", err) + } + if summary.Status.Indicator != "none" || len(summary.Components) != 1 { + t.Fatalf("unexpected summary: %+v", summary) + } +} + +func TestFetchIncidentsReturnsHTTPError(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + http.Error(w, "nope", http.StatusBadGateway) + })) + defer server.Close() + + client := NewClient(server.URL, time.Second) + if _, err := client.FetchIncidents(context.Background()); err == nil { + t.Fatal("expected HTTP error") + } +} diff --git a/internal/openai/models.go b/internal/openai/models.go new file mode 100644 index 0000000..0125a89 --- /dev/null +++ b/internal/openai/models.go @@ -0,0 +1,58 @@ +package openai + +type Page struct { + ID string `json:"id"` + Name string `json:"name"` + URL string `json:"url"` + UpdatedAt string `json:"updated_at"` +} + +type OverallStatus struct { + Description string `json:"description"` + Indicator string `json:"indicator"` +} + +type Component struct { + ID string `json:"id"` + Name string `json:"name"` + Status string `json:"status"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` + Position int `json:"position"` + PageID string `json:"page_id"` + Group bool `json:"group"` +} + +type Summary struct { + Page Page `json:"page"` + Status OverallStatus `json:"status"` + Components []Component `json:"components"` + Incidents []Incident `json:"incidents"` +} + +type Incident struct { + ID string `json:"id"` + Name string `json:"name"` + Status string `json:"status"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` + ResolvedAt string `json:"resolved_at"` + Impact string `json:"impact"` + PageID string `json:"page_id"` + IncidentUpdates []IncidentUpdate `json:"incident_updates"` +} + +type IncidentUpdate struct { + ID string `json:"id"` + Body string `json:"body"` + CreatedAt string `json:"created_at"` + DisplayAt string `json:"display_at"` + IncidentID string `json:"incident_id"` + Status string `json:"status"` + UpdatedAt string `json:"updated_at"` +} + +type IncidentsResponse struct { + Page Page `json:"page"` + Incidents []Incident `json:"incidents"` +} diff --git a/internal/poller/format.go b/internal/poller/format.go new file mode 100644 index 0000000..62819bc --- /dev/null +++ b/internal/poller/format.go @@ -0,0 +1,79 @@ +package poller + +import ( + "fmt" + "strings" + + openai "github.com/tiennm99/openai-status-bot/internal/openai" +) + +func FormatComponentChange(component openai.Component, previousStatus string) string { + return fmt.Sprintf( + "OpenAI component status changed\n\nComponent: %s\nStatus: %s -> %s\n\n%s", + component.Name, + StatusLabel(previousStatus), + StatusLabel(component.Status), + "https://status.openai.com/", + ) +} + +func FormatIncidentUpdate(incident openai.Incident, update openai.IncidentUpdate) string { + body := strings.TrimSpace(update.Body) + if body == "" { + body = "No update message provided." + } + + return fmt.Sprintf( + "OpenAI incident update\n\n%s\nStatus: %s\nImpact: %s\n\n%s\n\n%s", + incident.Name, + StatusLabel(update.Status), + StatusLabel(incident.Impact), + truncate(body, 2600), + "https://status.openai.com/incidents/"+incident.ID, + ) +} + +func StatusLabel(value string) string { + switch strings.ToLower(strings.TrimSpace(value)) { + case "": + return "Unknown" + case "operational": + return "Operational" + case "degraded_performance": + return "Degraded performance" + case "partial_outage": + return "Partial outage" + case "major_outage": + return "Major outage" + case "under_maintenance": + return "Under maintenance" + case "none": + return "None" + case "minor": + return "Minor" + case "major": + return "Major" + case "critical": + return "Critical" + case "investigating": + return "Investigating" + case "identified": + return "Identified" + case "monitoring": + return "Monitoring" + case "resolved": + return "Resolved" + default: + return strings.ReplaceAll(value, "_", " ") + } +} + +func truncate(value string, limit int) string { + if len(value) <= limit { + return value + } + if limit <= 3 { + return value[:limit] + } + return strings.TrimSpace(value[:limit-3]) + "..." +} diff --git a/internal/poller/format_test.go b/internal/poller/format_test.go new file mode 100644 index 0000000..f0c5d04 --- /dev/null +++ b/internal/poller/format_test.go @@ -0,0 +1,44 @@ +package poller + +import ( + "strings" + "testing" + + openai "github.com/tiennm99/openai-status-bot/internal/openai" +) + +func TestStatusLabel(t *testing.T) { + tests := map[string]string{ + "operational": "Operational", + "degraded_performance": "Degraded performance", + "partial_outage": "Partial outage", + "minor": "Minor", + "": "Unknown", + } + + for input, expected := range tests { + if got := StatusLabel(input); got != expected { + t.Fatalf("StatusLabel(%q) = %q, want %q", input, got, expected) + } + } +} + +func TestFormatIncidentUpdateIncludesCoreFields(t *testing.T) { + message := FormatIncidentUpdate( + openai.Incident{ + ID: "inc_123", + Name: "Codex outage", + Impact: "major", + }, + openai.IncidentUpdate{ + Status: "identified", + Body: "We found the issue.", + }, + ) + + for _, want := range []string{"Codex outage", "Identified", "Major", "We found the issue.", "inc_123"} { + if !strings.Contains(message, want) { + t.Fatalf("message missing %q: %s", want, message) + } + } +} diff --git a/internal/poller/poller.go b/internal/poller/poller.go new file mode 100644 index 0000000..079f015 --- /dev/null +++ b/internal/poller/poller.go @@ -0,0 +1,172 @@ +package poller + +import ( + "context" + "log/slog" + "time" + + openai "github.com/tiennm99/openai-status-bot/internal/openai" + "github.com/tiennm99/openai-status-bot/internal/redisstore" +) + +type StatusClient interface { + FetchSummary(ctx context.Context) (openai.Summary, error) + FetchIncidents(ctx context.Context) (openai.IncidentsResponse, error) +} + +type Store interface { + ComponentStatuses(ctx context.Context) (map[string]string, error) + HasIncidentUpdate(ctx context.Context, updateID string) (bool, error) + IsInitialized(ctx context.Context) (bool, error) + ListSubscribers(ctx context.Context) ([]redisstore.Subscriber, error) + MarkIncidentUpdate(ctx context.Context, updateID string) error + SaveComponentStatus(ctx context.Context, componentID, status string) error + SetInitialized(ctx context.Context) error +} + +type Notifier interface { + SendMessage(ctx context.Context, sub redisstore.Subscriber, text string) error +} + +type Runner struct { + statusClient StatusClient + store Store + notifier Notifier + interval time.Duration + logger *slog.Logger +} + +func NewRunner(statusClient StatusClient, store Store, notifier Notifier, interval time.Duration, logger *slog.Logger) *Runner { + return &Runner{ + statusClient: statusClient, + store: store, + notifier: notifier, + interval: interval, + logger: logger, + } +} + +func (r *Runner) Run(ctx context.Context) { + r.checkAndLog(ctx) + + ticker := time.NewTicker(r.interval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + r.checkAndLog(ctx) + } + } +} + +func (r *Runner) CheckOnce(ctx context.Context) error { + initialized, err := r.store.IsInitialized(ctx) + if err != nil { + return err + } + + summary, err := r.statusClient.FetchSummary(ctx) + if err != nil { + return err + } + incidents, err := r.statusClient.FetchIncidents(ctx) + if err != nil { + return err + } + + messages, err := r.collectComponentMessages(ctx, summary, initialized) + if err != nil { + return err + } + incidentMessages, err := r.collectIncidentMessages(ctx, incidents, initialized) + if err != nil { + return err + } + messages = append(messages, incidentMessages...) + + if initialized && len(messages) > 0 { + r.notifyAll(ctx, messages) + } + if !initialized { + if err := r.store.SetInitialized(ctx); err != nil { + return err + } + r.logger.Info("seeded status baseline") + } + return nil +} + +func (r *Runner) collectComponentMessages(ctx context.Context, summary openai.Summary, initialized bool) ([]string, error) { + knownStatuses, err := r.store.ComponentStatuses(ctx) + if err != nil { + return nil, err + } + + messages := make([]string, 0) + for _, component := range summary.Components { + if component.ID == "" || component.Group { + continue + } + previousStatus, found := knownStatuses[component.ID] + if initialized && found && previousStatus != component.Status { + messages = append(messages, FormatComponentChange(component, previousStatus)) + } + if err := r.store.SaveComponentStatus(ctx, component.ID, component.Status); err != nil { + return nil, err + } + } + return messages, nil +} + +func (r *Runner) collectIncidentMessages(ctx context.Context, response openai.IncidentsResponse, initialized bool) ([]string, error) { + messages := make([]string, 0) + for _, incident := range response.Incidents { + for i := len(incident.IncidentUpdates) - 1; i >= 0; i-- { + update := incident.IncidentUpdates[i] + if update.ID == "" { + continue + } + seen, err := r.store.HasIncidentUpdate(ctx, update.ID) + if err != nil { + return nil, err + } + if initialized && !seen { + messages = append(messages, FormatIncidentUpdate(incident, update)) + } + if !seen { + if err := r.store.MarkIncidentUpdate(ctx, update.ID); err != nil { + return nil, err + } + } + } + } + return messages, nil +} + +func (r *Runner) notifyAll(ctx context.Context, messages []string) { + subscribers, err := r.store.ListSubscribers(ctx) + if err != nil { + r.logger.Error("list subscribers", "error", err) + return + } + if len(subscribers) == 0 { + return + } + + for _, message := range messages { + for _, subscriber := range subscribers { + if err := r.notifier.SendMessage(ctx, subscriber, message); err != nil { + r.logger.Warn("send telegram message", "subscriber", subscriber.Key(), "error", err) + } + } + } +} + +func (r *Runner) checkAndLog(ctx context.Context) { + if err := r.CheckOnce(ctx); err != nil && ctx.Err() == nil { + r.logger.Error("poll openai status", "error", err) + } +} diff --git a/internal/redisstore/store.go b/internal/redisstore/store.go new file mode 100644 index 0000000..97ee595 --- /dev/null +++ b/internal/redisstore/store.go @@ -0,0 +1,117 @@ +package redisstore + +import ( + "context" + "fmt" + "strconv" + "strings" + + "github.com/redis/go-redis/v9" +) + +const ( + subscribersKey = "openai-status:subscribers" + componentStatusesKey = "openai-status:component-statuses" + incidentUpdatesKey = "openai-status:incident-updates" + initializedKey = "openai-status:initialized" +) + +type Subscriber struct { + ChatID int64 + ThreadID *int +} + +func NewSubscriber(chatID int64, threadID *int) Subscriber { + var copiedThreadID *int + if threadID != nil { + value := *threadID + copiedThreadID = &value + } + return Subscriber{ChatID: chatID, ThreadID: copiedThreadID} +} + +func (s Subscriber) Key() string { + if s.ThreadID == nil { + return strconv.FormatInt(s.ChatID, 10) + } + return fmt.Sprintf("%d:%d", s.ChatID, *s.ThreadID) +} + +func ParseSubscriberKey(value string) (Subscriber, error) { + parts := strings.Split(value, ":") + if len(parts) != 1 && len(parts) != 2 { + return Subscriber{}, fmt.Errorf("invalid subscriber key %q", value) + } + + chatID, err := strconv.ParseInt(parts[0], 10, 64) + if err != nil { + return Subscriber{}, fmt.Errorf("invalid chat ID: %w", err) + } + if len(parts) == 1 { + return Subscriber{ChatID: chatID}, nil + } + + threadID, err := strconv.Atoi(parts[1]) + if err != nil { + return Subscriber{}, fmt.Errorf("invalid thread ID: %w", err) + } + return NewSubscriber(chatID, &threadID), nil +} + +type Store struct { + client redis.UniversalClient +} + +func New(client redis.UniversalClient) *Store { + return &Store{client: client} +} + +func (s *Store) AddSubscriber(ctx context.Context, sub Subscriber) error { + return s.client.SAdd(ctx, subscribersKey, sub.Key()).Err() +} + +func (s *Store) RemoveSubscriber(ctx context.Context, sub Subscriber) error { + return s.client.SRem(ctx, subscribersKey, sub.Key()).Err() +} + +func (s *Store) ListSubscribers(ctx context.Context) ([]Subscriber, error) { + keys, err := s.client.SMembers(ctx, subscribersKey).Result() + if err != nil { + return nil, err + } + + subscribers := make([]Subscriber, 0, len(keys)) + for _, key := range keys { + sub, err := ParseSubscriberKey(key) + if err != nil { + return nil, err + } + subscribers = append(subscribers, sub) + } + return subscribers, nil +} + +func (s *Store) IsInitialized(ctx context.Context) (bool, error) { + count, err := s.client.Exists(ctx, initializedKey).Result() + return count > 0, err +} + +func (s *Store) SetInitialized(ctx context.Context) error { + return s.client.Set(ctx, initializedKey, "1", 0).Err() +} + +func (s *Store) ComponentStatuses(ctx context.Context) (map[string]string, error) { + return s.client.HGetAll(ctx, componentStatusesKey).Result() +} + +func (s *Store) SaveComponentStatus(ctx context.Context, componentID, status string) error { + return s.client.HSet(ctx, componentStatusesKey, componentID, status).Err() +} + +func (s *Store) HasIncidentUpdate(ctx context.Context, updateID string) (bool, error) { + return s.client.SIsMember(ctx, incidentUpdatesKey, updateID).Result() +} + +func (s *Store) MarkIncidentUpdate(ctx context.Context, updateID string) error { + return s.client.SAdd(ctx, incidentUpdatesKey, updateID).Err() +} diff --git a/internal/redisstore/store_test.go b/internal/redisstore/store_test.go new file mode 100644 index 0000000..5c78347 --- /dev/null +++ b/internal/redisstore/store_test.go @@ -0,0 +1,32 @@ +package redisstore + +import "testing" + +func TestSubscriberKeyRoundTripWithoutThread(t *testing.T) { + sub := NewSubscriber(12345, nil) + parsed, err := ParseSubscriberKey(sub.Key()) + if err != nil { + t.Fatalf("ParseSubscriberKey returned error: %v", err) + } + if parsed.ChatID != sub.ChatID || parsed.ThreadID != nil { + t.Fatalf("parsed subscriber = %+v, want %+v", parsed, sub) + } +} + +func TestSubscriberKeyRoundTripWithThreadZero(t *testing.T) { + threadID := 0 + sub := NewSubscriber(-10012345, &threadID) + parsed, err := ParseSubscriberKey(sub.Key()) + if err != nil { + t.Fatalf("ParseSubscriberKey returned error: %v", err) + } + if parsed.ChatID != sub.ChatID || parsed.ThreadID == nil || *parsed.ThreadID != threadID { + t.Fatalf("parsed subscriber = %+v, want thread ID %d", parsed, threadID) + } +} + +func TestParseSubscriberKeyRejectsInvalidValue(t *testing.T) { + if _, err := ParseSubscriberKey("abc:def:ghi"); err == nil { + t.Fatal("expected invalid key error") + } +} diff --git a/internal/telegram/client.go b/internal/telegram/client.go new file mode 100644 index 0000000..625994b --- /dev/null +++ b/internal/telegram/client.go @@ -0,0 +1,94 @@ +package telegram + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "net/http" + "time" + + "github.com/tiennm99/openai-status-bot/internal/redisstore" +) + +type Client struct { + baseURL string + httpClient *http.Client +} + +func NewClient(token string, timeout time.Duration) *Client { + return &Client{ + baseURL: "https://api.telegram.org/bot" + token, + httpClient: &http.Client{ + Timeout: timeout + 60*time.Second, + }, + } +} + +func (c *Client) GetUpdates(ctx context.Context, offset int64, timeoutSeconds int) ([]Update, error) { + payload := map[string]any{ + "offset": offset, + "timeout": timeoutSeconds, + "allowed_updates": []string{"message"}, + } + + var updates []Update + if err := c.postJSON(ctx, "/getUpdates", payload, &updates); err != nil { + return nil, err + } + return updates, nil +} + +func (c *Client) SendMessage(ctx context.Context, sub redisstore.Subscriber, text string) error { + return c.SendText(ctx, sub.ChatID, sub.ThreadID, text) +} + +func (c *Client) SendText(ctx context.Context, chatID int64, threadID *int, text string) error { + payload := map[string]any{ + "chat_id": chatID, + "text": text, + } + if threadID != nil { + payload["message_thread_id"] = *threadID + } + + var result json.RawMessage + return c.postJSON(ctx, "/sendMessage", payload, &result) +} + +func (c *Client) postJSON(ctx context.Context, path string, payload any, target any) error { + body, err := json.Marshal(payload) + if err != nil { + return err + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, c.baseURL+path, bytes.NewReader(body)) + if err != nil { + return err + } + req.Header.Set("Content-Type", "application/json") + + res, err := c.httpClient.Do(req) + if err != nil { + return err + } + defer res.Body.Close() + + var envelope struct { + OK bool `json:"ok"` + Result json.RawMessage `json:"result"` + Description string `json:"description"` + } + if err := json.NewDecoder(res.Body).Decode(&envelope); err != nil { + return fmt.Errorf("decode telegram response: %w", err) + } + if !envelope.OK { + return fmt.Errorf("telegram API %s: %s", res.Status, envelope.Description) + } + if target != nil { + if err := json.Unmarshal(envelope.Result, target); err != nil { + return fmt.Errorf("decode telegram result: %w", err) + } + } + return nil +} diff --git a/internal/telegram/types.go b/internal/telegram/types.go new file mode 100644 index 0000000..64ea722 --- /dev/null +++ b/internal/telegram/types.go @@ -0,0 +1,19 @@ +package telegram + +type Update struct { + UpdateID int64 `json:"update_id"` + Message *Message `json:"message"` +} + +type Message struct { + MessageID int64 `json:"message_id"` + MessageThreadID *int `json:"message_thread_id,omitempty"` + Text string `json:"text"` + Chat Chat `json:"chat"` +} + +type Chat struct { + ID int64 `json:"id"` + Type string `json:"type"` + Title string `json:"title,omitempty"` +} diff --git a/plans/260619-1401-initial-scaffold/plan.md b/plans/260619-1401-initial-scaffold/plan.md new file mode 100644 index 0000000..842eec7 --- /dev/null +++ b/plans/260619-1401-initial-scaffold/plan.md @@ -0,0 +1,27 @@ +# Initial Scaffold Plan + +## Status + +Completed. + +## Phases + +| Phase | Status | Output | +|-------|--------|--------| +| Scaffold Go project | Completed | `go.mod`, source layout, Docker files | +| Implement bot runtime | Completed | Telegram long polling, commands | +| Implement status checks | Completed | OpenAI JSON client, one-minute poller | +| Implement Redis storage | Completed | Subscribers and checkpoint keys | +| Verify | Completed | Unit tests and build | + +## Acceptance Criteria + +- Go project exists at `/config/workspace/tiennm99/openai-status-bot`. +- Uses Redis as database for subscribers and poll state. +- Checks OpenAI status every minute by default. +- Telegram users can subscribe and query current status. +- Tests and build pass. + +## Unresolved Questions + +- None.