Building out notification update channels
Some checks are pending
check / check (push) Waiting to run

This commit is contained in:
KS Jannette
2026-03-01 03:23:08 -05:00
parent d07dc971fe
commit c494bf5f53
18 changed files with 1282 additions and 187 deletions

View File

@@ -14,3 +14,7 @@ POLL_INTERVAL_MS=60000
# Firebase
FIREBASE_PROJECT_ID=koin-ping
# Email notifications (Resend — https://resend.com)
# RESEND_API_KEY=re_xxxxxxxxxxxx
# EMAIL_FROM=Koin Ping <alerts@yourdomain.com>

View File

@@ -14,6 +14,7 @@ import (
"github.com/kjannette/koin-ping/backend-go/internal/handlers"
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
"github.com/kjannette/koin-ping/backend-go/internal/models"
"github.com/kjannette/koin-ping/backend-go/internal/services"
)
const (
@@ -48,10 +49,15 @@ func main() {
alertEventModel := models.NewAlertEventModel(pool)
notifConfigModel := models.NewNotificationConfigModel(pool)
emailDigestSvc := services.NewEmailDigestService(
cfg.ResendAPIKey, cfg.EmailFrom, alertEventModel, notifConfigModel,
)
addressHandler := handlers.NewAddressHandler(addressModel)
alertRuleHandler := handlers.NewAlertRuleHandler(alertRuleModel, addressModel)
alertEventHandler := handlers.NewAlertEventHandler(alertEventModel)
notifConfigHandler := handlers.NewNotificationConfigHandler(notifConfigModel)
notifConfigHandler := handlers.NewNotificationConfigHandler(notifConfigModel, cfg)
emailDigestHandler := handlers.NewEmailDigestHandler(emailDigestSvc, notifConfigModel)
mux := http.NewServeMux()
b := cfg.APIBasePath // e.g. "/v1"
@@ -89,6 +95,14 @@ func main() {
middleware.Authenticate(http.HandlerFunc(notifConfigHandler.UpdateConfig)))
mux.Handle("DELETE "+b+"/notification-config",
middleware.Authenticate(http.HandlerFunc(notifConfigHandler.DeleteConfig)))
mux.Handle("POST "+b+"/notification-config/test",
middleware.Authenticate(http.HandlerFunc(notifConfigHandler.TestChannels)))
// Authenticated routes — email digest
mux.Handle("POST "+b+"/email/setup",
middleware.Authenticate(http.HandlerFunc(emailDigestHandler.SetupEmail)))
mux.Handle("POST "+b+"/email/digest",
middleware.Authenticate(http.HandlerFunc(emailDigestHandler.SendDigest)))
handler := corsMiddleware(mux)

View File

@@ -58,7 +58,10 @@ func main() {
notifConfigModel := models.NewNotificationConfigModel(pool)
observer := services.NewObserverService(eth, addressModel, checkpointModel)
evaluator := services.NewEvaluatorService(eth, alertRuleModel, alertEventModel, addressModel, notifConfigModel)
evaluator := services.NewEvaluatorService(
eth, alertRuleModel, alertEventModel, addressModel, notifConfigModel,
cfg.ResendAPIKey, cfg.EmailFrom,
)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()

View File

@@ -0,0 +1,2 @@
ALTER TABLE user_notification_configs
ADD COLUMN IF NOT EXISTS slack_webhook_url TEXT;

View File

@@ -54,8 +54,9 @@ CREATE TABLE user_notification_configs (
user_id VARCHAR(128) PRIMARY KEY,
discord_webhook_url TEXT, -- Discord webhook URL (nullable)
telegram_chat_id VARCHAR(128), -- Telegram chat ID (nullable)
telegram_bot_token VARCHAR(255), -- Telegram bot token (nullable, future use)
telegram_bot_token VARCHAR(255), -- Telegram bot token (nullable)
email VARCHAR(255), -- Email for notifications (nullable)
slack_webhook_url TEXT, -- Slack incoming webhook URL (nullable)
notification_enabled BOOLEAN DEFAULT TRUE, -- Master on/off switch
created_at TIMESTAMP DEFAULT NOW(),
updated_at TIMESTAMP DEFAULT NOW()

View File

@@ -28,6 +28,8 @@ type Config struct {
EthRPCURL string
PollIntervalMS int
NodeEnv string
ResendAPIKey string
EmailFrom string
}
// Load reads configuration from environment variables and returns a Config.
@@ -45,6 +47,8 @@ func Load() (*Config, error) {
EthRPCURL: os.Getenv("ETH_RPC_URL"),
PollIntervalMS: getEnvInt("POLL_INTERVAL_MS", defaultPollIntervalMS),
NodeEnv: getEnv("NODE_ENV", "development"),
ResendAPIKey: os.Getenv("RESEND_API_KEY"),
EmailFrom: getEnv("EMAIL_FROM", "Koin Ping <alerts@koinping.com>"),
}
if cfg.PollIntervalMS < minPollIntervalMS {

View File

@@ -103,8 +103,9 @@ type NotificationConfig struct {
UserID string `json:"user_id"` //nolint:tagliatelle
DiscordWebhookURL *string `json:"discord_webhook_url"` //nolint:tagliatelle
TelegramChatID *string `json:"telegram_chat_id"` //nolint:tagliatelle
TelegramBotToken *string `json:"telegram_bot_token,omitempty"` //nolint:tagliatelle
TelegramBotToken *string `json:"telegram_bot_token"` //nolint:tagliatelle
Email *string `json:"email"`
SlackWebhookURL *string `json:"slack_webhook_url"` //nolint:tagliatelle
NotificationEnabled bool `json:"notification_enabled"` //nolint:tagliatelle
CreatedAt *time.Time `json:"created_at,omitempty"` //nolint:tagliatelle
UpdatedAt *time.Time `json:"updated_at,omitempty"` //nolint:tagliatelle

View File

@@ -0,0 +1,101 @@
package handlers
import (
"log"
"net/http"
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
"github.com/kjannette/koin-ping/backend-go/internal/models"
"github.com/kjannette/koin-ping/backend-go/internal/services"
)
type EmailDigestHandler struct {
digestSvc *services.EmailDigestService
configs *models.NotificationConfigModel
}
func NewEmailDigestHandler(
digestSvc *services.EmailDigestService,
configs *models.NotificationConfigModel,
) *EmailDigestHandler {
return &EmailDigestHandler{digestSvc: digestSvc, configs: configs}
}
// SetupEmail reads the user's email from their notification config and sends
// a confirmation message via Resend to verify the integration works.
func (h *EmailDigestHandler) SetupEmail(w http.ResponseWriter, r *http.Request) {
userID := middleware.GetUserID(r.Context())
if !h.digestSvc.Configured() {
writeError(w, http.StatusServiceUnavailable, "EMAIL_NOT_CONFIGURED",
"Email service is not configured on the server")
return
}
cfg, err := h.configs.GetConfig(r.Context(), userID)
if err != nil {
log.Printf("Error getting notification config for email setup: %v", err)
writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR",
"Failed to load notification config")
return
}
if cfg == nil || cfg.Email == nil || *cfg.Email == "" {
writeError(w, http.StatusBadRequest, "VALIDATION_ERROR",
"Save an email address in notification settings first")
return
}
if err := h.digestSvc.SetupEmail(*cfg.Email); err != nil {
log.Printf("Email setup failed for user %s: %v", userID, err)
writeError(w, http.StatusBadGateway, "EMAIL_SEND_FAILED",
"Failed to send confirmation email — check server email config")
return
}
log.Printf("Email setup confirmation sent to user %s (%s)", userID, *cfg.Email)
writeJSON(w, http.StatusOK, map[string]any{
"success": true,
"email": *cfg.Email,
"message": "Confirmation email sent",
})
}
// SendDigest compiles and sends a digest of recent alerts to the user's email.
func (h *EmailDigestHandler) SendDigest(w http.ResponseWriter, r *http.Request) {
userID := middleware.GetUserID(r.Context())
if !h.digestSvc.Configured() {
writeError(w, http.StatusServiceUnavailable, "EMAIL_NOT_CONFIGURED",
"Email service is not configured on the server")
return
}
cfg, err := h.configs.GetConfig(r.Context(), userID)
if err != nil {
log.Printf("Error getting notification config for digest: %v", err)
writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR",
"Failed to load notification config")
return
}
if cfg == nil || cfg.Email == nil || *cfg.Email == "" {
writeError(w, http.StatusBadRequest, "VALIDATION_ERROR",
"No email address configured")
return
}
if err := h.digestSvc.SendDigest(r.Context(), userID, *cfg.Email); err != nil {
log.Printf("Digest send failed for user %s: %v", userID, err)
writeError(w, http.StatusBadGateway, "DIGEST_SEND_FAILED",
"Failed to send digest email")
return
}
log.Printf("Digest sent to user %s (%s)", userID, *cfg.Email)
writeJSON(w, http.StatusOK, map[string]any{
"success": true,
"email": *cfg.Email,
"message": "Digest email sent",
})
}

View File

@@ -7,19 +7,22 @@ import (
"regexp"
"strings"
"github.com/kjannette/koin-ping/backend-go/internal/config"
"github.com/kjannette/koin-ping/backend-go/internal/domain"
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
"github.com/kjannette/koin-ping/backend-go/internal/models"
"github.com/kjannette/koin-ping/backend-go/internal/notifications"
)
var emailRe = regexp.MustCompile(`^[^\s@]+@[^\s@]+\.[^\s@]+$`)
type NotificationConfigHandler struct {
configs *models.NotificationConfigModel
cfg *config.Config
}
func NewNotificationConfigHandler(configs *models.NotificationConfigModel) *NotificationConfigHandler {
return &NotificationConfigHandler{configs: configs}
func NewNotificationConfigHandler(configs *models.NotificationConfigModel, cfg *config.Config) *NotificationConfigHandler {
return &NotificationConfigHandler{configs: configs, cfg: cfg}
}
func (h *NotificationConfigHandler) GetConfig(w http.ResponseWriter, r *http.Request) {
@@ -51,7 +54,9 @@ func (h *NotificationConfigHandler) UpdateConfig(w http.ResponseWriter, r *http.
var body struct {
DiscordWebhookURL *string `json:"discord_webhook_url"`
TelegramChatID *string `json:"telegram_chat_id"`
TelegramBotToken *string `json:"telegram_bot_token"`
Email *string `json:"email"`
SlackWebhookURL *string `json:"slack_webhook_url"`
NotificationEnabled *bool `json:"notification_enabled"`
}
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
@@ -63,7 +68,8 @@ func (h *NotificationConfigHandler) UpdateConfig(w http.ResponseWriter, r *http.
log.Printf("User %s updating notification config", userID)
if body.DiscordWebhookURL == nil && body.TelegramChatID == nil &&
body.Email == nil && body.NotificationEnabled == nil {
body.TelegramBotToken == nil && body.Email == nil &&
body.SlackWebhookURL == nil && body.NotificationEnabled == nil {
writeError(w, http.StatusBadRequest, "VALIDATION_ERROR",
"At least one configuration field must be provided")
return
@@ -76,6 +82,13 @@ func (h *NotificationConfigHandler) UpdateConfig(w http.ResponseWriter, r *http.
return
}
if body.SlackWebhookURL != nil && *body.SlackWebhookURL != "" &&
!strings.HasPrefix(*body.SlackWebhookURL, "https://hooks.slack.com/") {
writeError(w, http.StatusBadRequest, "VALIDATION_ERROR",
"Invalid Slack webhook URL format")
return
}
if body.Email != nil && *body.Email != "" && !emailRe.MatchString(*body.Email) {
writeError(w, http.StatusBadRequest, "VALIDATION_ERROR",
"Invalid email address format")
@@ -90,7 +103,9 @@ func (h *NotificationConfigHandler) UpdateConfig(w http.ResponseWriter, r *http.
cfg := domain.NotificationConfig{
DiscordWebhookURL: body.DiscordWebhookURL,
TelegramChatID: body.TelegramChatID,
TelegramBotToken: body.TelegramBotToken,
Email: body.Email,
SlackWebhookURL: body.SlackWebhookURL,
NotificationEnabled: enabled,
}
@@ -126,3 +141,76 @@ func (h *NotificationConfigHandler) DeleteConfig(w http.ResponseWriter, r *http.
log.Println("Notification config deleted")
w.WriteHeader(http.StatusNoContent)
}
// TestChannels sends a test message to all configured notification channels.
func (h *NotificationConfigHandler) TestChannels(w http.ResponseWriter, r *http.Request) {
userID := middleware.GetUserID(r.Context())
log.Printf("User %s testing notification channels", userID)
cfg, err := h.configs.GetConfig(r.Context(), userID)
if err != nil {
log.Printf("Error getting notification config for test: %v", err)
writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR", "Failed to get notification config")
return
}
if cfg == nil {
writeError(w, http.StatusNotFound, "NOT_FOUND", "No notification configuration found")
return
}
type channelResult struct {
Channel string `json:"channel"`
Success bool `json:"success"`
Error string `json:"error,omitempty"`
}
var results []channelResult
if cfg.DiscordWebhookURL != nil && *cfg.DiscordWebhookURL != "" {
ok, testErr := notifications.TestDiscordWebhook(*cfg.DiscordWebhookURL)
res := channelResult{Channel: "discord", Success: ok}
if testErr != nil {
res.Error = testErr.Error()
}
results = append(results, res)
}
if cfg.TelegramBotToken != nil && *cfg.TelegramBotToken != "" &&
cfg.TelegramChatID != nil && *cfg.TelegramChatID != "" {
ok, testErr := notifications.TestTelegramWebhook(*cfg.TelegramBotToken, *cfg.TelegramChatID)
res := channelResult{Channel: "telegram", Success: ok}
if testErr != nil {
res.Error = testErr.Error()
}
results = append(results, res)
}
if cfg.SlackWebhookURL != nil && *cfg.SlackWebhookURL != "" {
ok, testErr := notifications.TestSlackWebhook(*cfg.SlackWebhookURL)
res := channelResult{Channel: "slack", Success: ok}
if testErr != nil {
res.Error = testErr.Error()
}
results = append(results, res)
}
if cfg.Email != nil && *cfg.Email != "" {
ok, testErr := notifications.TestEmailNotification(
h.cfg.ResendAPIKey, h.cfg.EmailFrom, *cfg.Email,
)
res := channelResult{Channel: "email", Success: ok}
if testErr != nil {
res.Error = testErr.Error()
}
results = append(results, res)
}
if len(results) == 0 {
writeError(w, http.StatusBadRequest, "NO_CHANNELS",
"No notification channels are configured")
return
}
writeJSON(w, http.StatusOK, map[string]any{"results": results})
}

View File

@@ -21,12 +21,12 @@ func (m *NotificationConfigModel) GetConfig(ctx context.Context, userID string)
var c domain.NotificationConfig
err := m.pool.QueryRow(ctx,
`SELECT user_id, discord_webhook_url, telegram_chat_id, telegram_bot_token,
email, notification_enabled, created_at, updated_at
email, slack_webhook_url, notification_enabled, created_at, updated_at
FROM user_notification_configs
WHERE user_id = $1`,
userID,
).Scan(&c.UserID, &c.DiscordWebhookURL, &c.TelegramChatID, &c.TelegramBotToken,
&c.Email, &c.NotificationEnabled, &c.CreatedAt, &c.UpdatedAt)
&c.Email, &c.SlackWebhookURL, &c.NotificationEnabled, &c.CreatedAt, &c.UpdatedAt)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
@@ -40,22 +40,24 @@ func (m *NotificationConfigModel) UpsertConfig(ctx context.Context, userID strin
var c domain.NotificationConfig
err := m.pool.QueryRow(ctx,
`INSERT INTO user_notification_configs
(user_id, discord_webhook_url, telegram_chat_id, telegram_bot_token, email, notification_enabled, updated_at)
VALUES ($1, $2, $3, $4, $5, $6, NOW())
(user_id, discord_webhook_url, telegram_chat_id, telegram_bot_token,
email, slack_webhook_url, notification_enabled, updated_at)
VALUES ($1, $2, $3, $4, $5, $6, $7, NOW())
ON CONFLICT (user_id)
DO UPDATE SET
discord_webhook_url = COALESCE($2, user_notification_configs.discord_webhook_url),
telegram_chat_id = COALESCE($3, user_notification_configs.telegram_chat_id),
telegram_bot_token = COALESCE($4, user_notification_configs.telegram_bot_token),
email = COALESCE($5, user_notification_configs.email),
notification_enabled = $6,
slack_webhook_url = COALESCE($6, user_notification_configs.slack_webhook_url),
notification_enabled = $7,
updated_at = NOW()
RETURNING user_id, discord_webhook_url, telegram_chat_id, telegram_bot_token,
email, notification_enabled, created_at, updated_at`,
email, slack_webhook_url, notification_enabled, created_at, updated_at`,
userID, cfg.DiscordWebhookURL, cfg.TelegramChatID, cfg.TelegramBotToken,
cfg.Email, cfg.NotificationEnabled,
cfg.Email, cfg.SlackWebhookURL, cfg.NotificationEnabled,
).Scan(&c.UserID, &c.DiscordWebhookURL, &c.TelegramChatID, &c.TelegramBotToken,
&c.Email, &c.NotificationEnabled, &c.CreatedAt, &c.UpdatedAt)
&c.Email, &c.SlackWebhookURL, &c.NotificationEnabled, &c.CreatedAt, &c.UpdatedAt)
if err != nil {
return nil, err
}
@@ -75,7 +77,8 @@ func (m *NotificationConfigModel) Remove(ctx context.Context, userID string) (bo
func (m *NotificationConfigModel) ListEnabled(ctx context.Context) ([]domain.NotificationConfig, error) {
rows, err := m.pool.Query(ctx,
`SELECT user_id, discord_webhook_url, telegram_chat_id, email
`SELECT user_id, discord_webhook_url, telegram_chat_id, telegram_bot_token,
email, slack_webhook_url
FROM user_notification_configs
WHERE notification_enabled = TRUE`,
)
@@ -87,7 +90,8 @@ func (m *NotificationConfigModel) ListEnabled(ctx context.Context) ([]domain.Not
var configs []domain.NotificationConfig
for rows.Next() {
var c domain.NotificationConfig
if err := rows.Scan(&c.UserID, &c.DiscordWebhookURL, &c.TelegramChatID, &c.Email); err != nil {
if err := rows.Scan(&c.UserID, &c.DiscordWebhookURL, &c.TelegramChatID,
&c.TelegramBotToken, &c.Email, &c.SlackWebhookURL); err != nil {
return nil, err
}
c.NotificationEnabled = true

View File

@@ -0,0 +1,142 @@
package notifications
import (
"bytes"
"encoding/json"
"fmt"
"log"
"net/http"
"time"
)
const emailHTTPTimeoutSeconds = 10
var emailHTTPClient = &http.Client{ //nolint:gochecknoglobals
Timeout: emailHTTPTimeoutSeconds * time.Second,
}
type resendPayload struct {
From string `json:"from"`
To string `json:"to"`
Subject string `json:"subject"`
HTML string `json:"html"`
}
func SendEmailNotification(apiKey, fromAddress, toAddress, message string, meta AlertMetadata) (bool, error) {
if apiKey == "" {
log.Printf("Skipping email notification: RESEND_API_KEY not configured")
return false, nil
}
subject := fmt.Sprintf("Koin Ping Alert: %s", alertTypeLabel(meta.AlertType))
txLink := ""
if meta.TxHash != "" {
txLink = fmt.Sprintf(
`<p><a href="https://etherscan.io/tx/%s">View on Etherscan</a></p>`,
meta.TxHash,
)
}
html := fmt.Sprintf(`
<div style="font-family: sans-serif; max-width: 600px; margin: 0 auto;">
<h2 style="color: #333;">Koin Ping Alert</h2>
<p style="font-size: 16px;">%s</p>
<table style="margin: 16px 0; border-collapse: collapse;">
<tr>
<td style="padding: 4px 12px 4px 0; color: #666;">Address</td>
<td style="padding: 4px 0;">%s</td>
</tr>
<tr>
<td style="padding: 4px 12px 4px 0; color: #666;">Blockchain</td>
<td style="padding: 4px 0; font-family: monospace; font-size: 13px;">%s</td>
</tr>
</table>
%s
<hr style="border: none; border-top: 1px solid #eee; margin: 24px 0;" />
<p style="font-size: 12px; color: #999;">Sent by Koin Ping</p>
</div>`,
message, meta.AddressLabel, meta.Address, txLink)
payload := resendPayload{
From: fromAddress,
To: toAddress,
Subject: subject,
HTML: html,
}
body, err := json.Marshal(payload)
if err != nil {
return false, fmt.Errorf("marshal email payload: %w", err)
}
req, err := http.NewRequest(http.MethodPost, "https://api.resend.com/emails", bytes.NewReader(body))
if err != nil {
return false, fmt.Errorf("create email request: %w", err)
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+apiKey)
resp, err := emailHTTPClient.Do(req)
if err != nil {
log.Printf("Failed to send email notification: %v", err)
return false, err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
log.Printf("Resend API failed: HTTP %d", resp.StatusCode)
return false, fmt.Errorf("resend API failed: HTTP %d", resp.StatusCode)
}
return true, nil
}
func TestEmailNotification(apiKey, fromAddress, toAddress string) (bool, error) {
if apiKey == "" {
return false, fmt.Errorf("email not configured: RESEND_API_KEY not set") //nolint:err113
}
payload := resendPayload{
From: fromAddress,
To: toAddress,
Subject: "Koin Ping — Test Notification",
HTML: `<p>Your email alerts are configured correctly!</p><p style="font-size:12px;color:#999;">Sent by Koin Ping</p>`,
}
body, err := json.Marshal(payload)
if err != nil {
return false, err
}
req, err := http.NewRequest(http.MethodPost, "https://api.resend.com/emails", bytes.NewReader(body))
if err != nil {
return false, err
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+apiKey)
resp, err := emailHTTPClient.Do(req)
if err != nil {
log.Printf("Email test failed: %v", err)
return false, err
}
defer resp.Body.Close()
return resp.StatusCode >= 200 && resp.StatusCode < 300, nil
}
func alertTypeLabel(alertType string) string {
switch alertType {
case "incoming_tx":
return "Incoming Transaction"
case "outgoing_tx":
return "Outgoing Transaction"
case "large_transfer":
return "Large Transfer"
case "balance_below":
return "Balance Below Threshold"
default:
return "Alert"
}
}

View File

@@ -0,0 +1,116 @@
package notifications
import (
"bytes"
"encoding/json"
"fmt"
"log"
"net/http"
"time"
)
const slackHTTPTimeoutSeconds = 10
var slackHTTPClient = &http.Client{ //nolint:gochecknoglobals
Timeout: slackHTTPTimeoutSeconds * time.Second,
}
type slackAttachment struct {
Color string `json:"color"`
Title string `json:"title"`
Text string `json:"text"`
Fields []slackField `json:"fields"`
Footer string `json:"footer"`
Ts int64 `json:"ts"`
}
type slackField struct {
Title string `json:"title"`
Value string `json:"value"`
Short bool `json:"short"`
}
type slackPayload struct {
Text string `json:"text,omitempty"`
Attachments []slackAttachment `json:"attachments,omitempty"`
}
func SendSlackNotification(webhookURL, message string, meta AlertMetadata) (bool, error) {
fields := []slackField{
{Title: "Address", Value: meta.AddressLabel, Short: true},
{Title: "Blockchain Address", Value: fmt.Sprintf("`%s`", meta.Address), Short: false},
}
if meta.TxHash != "" {
fields = append(fields, slackField{
Title: "Transaction",
Value: fmt.Sprintf("<https://etherscan.io/tx/%s|View on Etherscan>", meta.TxHash),
Short: false,
})
}
payload := slackPayload{
Attachments: []slackAttachment{
{
Color: slackColorForAlertType(meta.AlertType),
Title: "Koin Ping Alert",
Text: message,
Fields: fields,
Footer: "Koin Ping",
Ts: time.Now().Unix(),
},
},
}
body, err := json.Marshal(payload)
if err != nil {
return false, fmt.Errorf("marshal slack payload: %w", err)
}
resp, err := slackHTTPClient.Post(webhookURL, "application/json", bytes.NewReader(body))
if err != nil {
log.Printf("Failed to send Slack notification: %v", err)
return false, err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
log.Printf("Slack webhook failed: HTTP %d", resp.StatusCode)
return false, fmt.Errorf("slack webhook failed: HTTP %d", resp.StatusCode)
}
return true, nil
}
func TestSlackWebhook(webhookURL string) (bool, error) {
payload := slackPayload{
Text: "Koin Ping test notification — Your Slack alerts are configured correctly!",
}
body, err := json.Marshal(payload)
if err != nil {
return false, err
}
resp, err := slackHTTPClient.Post(webhookURL, "application/json", bytes.NewReader(body))
if err != nil {
log.Printf("Slack webhook test failed: %v", err)
return false, err
}
defer resp.Body.Close()
return resp.StatusCode >= 200 && resp.StatusCode < 300, nil
}
func slackColorForAlertType(alertType string) string {
switch alertType {
case "incoming_tx":
return "#00ff00"
case "outgoing_tx":
return "#ff9900"
case "large_transfer", "balance_below":
return "#ff0000"
default:
return "#0099ff"
}
}

View File

@@ -0,0 +1,109 @@
package notifications
import (
"bytes"
"encoding/json"
"fmt"
"log"
"net/http"
"time"
)
const telegramHTTPTimeoutSeconds = 10
var telegramHTTPClient = &http.Client{ //nolint:gochecknoglobals
Timeout: telegramHTTPTimeoutSeconds * time.Second,
}
type telegramPayload struct {
ChatID string `json:"chat_id"`
Text string `json:"text"`
ParseMode string `json:"parse_mode"`
}
func SendTelegramNotification(botToken, chatID, message string, meta AlertMetadata) (bool, error) {
text := fmt.Sprintf("*Koin Ping Alert*\n\n%s\n\n*Address:* %s\n`%s`",
escapeMarkdown(message), escapeMarkdown(meta.AddressLabel), meta.Address)
if meta.TxHash != "" {
text += fmt.Sprintf("\n\n[View on Etherscan](https://etherscan.io/tx/%s)", meta.TxHash)
}
payload := telegramPayload{
ChatID: chatID,
Text: text,
ParseMode: "Markdown",
}
body, err := json.Marshal(payload)
if err != nil {
return false, fmt.Errorf("marshal telegram payload: %w", err)
}
url := fmt.Sprintf("https://api.telegram.org/bot%s/sendMessage", botToken)
resp, err := telegramHTTPClient.Post(url, "application/json", bytes.NewReader(body))
if err != nil {
log.Printf("Failed to send Telegram notification: %v", err)
return false, err
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
log.Printf("Telegram API failed: HTTP %d", resp.StatusCode)
return false, fmt.Errorf("telegram API failed: HTTP %d", resp.StatusCode)
}
return true, nil
}
func TestTelegramWebhook(botToken, chatID string) (bool, error) {
payload := telegramPayload{
ChatID: chatID,
Text: "Koin Ping test notification — Your Telegram alerts are configured correctly!",
ParseMode: "Markdown",
}
body, err := json.Marshal(payload)
if err != nil {
return false, err
}
url := fmt.Sprintf("https://api.telegram.org/bot%s/sendMessage", botToken)
resp, err := telegramHTTPClient.Post(url, "application/json", bytes.NewReader(body))
if err != nil {
log.Printf("Telegram test failed: %v", err)
return false, err
}
defer resp.Body.Close()
return resp.StatusCode >= 200 && resp.StatusCode < 300, nil
}
func escapeMarkdown(s string) string {
replacer := []struct{ old, new string }{
{"_", "\\_"}, {"*", "\\*"}, {"[", "\\["}, {"]", "\\]"},
{"(", "\\("}, {")", "\\)"}, {"~", "\\~"}, {"`", "\\`"},
{">", "\\>"}, {"#", "\\#"}, {"+", "\\+"}, {"-", "\\-"},
{"=", "\\="}, {"|", "\\|"}, {"{", "\\{"}, {"}", "\\}"},
{".", "\\."}, {"!", "\\!"},
}
result := s
for _, r := range replacer {
result = replaceAll(result, r.old, r.new)
}
return result
}
func replaceAll(s, old, new string) string {
out := ""
for i := 0; i < len(s); i++ {
if string(s[i]) == old {
out += new
} else {
out += string(s[i])
}
}
return out
}

View File

@@ -0,0 +1,199 @@
package services
import (
"bytes"
"context"
"encoding/json"
"fmt"
"log"
"net/http"
"time"
"github.com/kjannette/koin-ping/backend-go/internal/models"
)
const (
resendAPIURL = "https://api.resend.com/emails"
emailHTTPTimeout = 10 * time.Second
defaultDigestMaxItems = 50
)
var digestHTTPClient = &http.Client{Timeout: emailHTTPTimeout} //nolint:gochecknoglobals
// EmailDigestService handles email setup and digest sending via Resend.
type EmailDigestService struct {
apiKey string
fromAddress string
alertEvents *models.AlertEventModel
notifCfgs *models.NotificationConfigModel
}
func NewEmailDigestService(
apiKey, fromAddress string,
alertEvents *models.AlertEventModel,
notifCfgs *models.NotificationConfigModel,
) *EmailDigestService {
return &EmailDigestService{
apiKey: apiKey,
fromAddress: fromAddress,
alertEvents: alertEvents,
notifCfgs: notifCfgs,
}
}
// Configured returns true when the Resend API key is present.
func (s *EmailDigestService) Configured() bool {
return s.apiKey != ""
}
// SetupEmail validates the email works by sending a welcome/confirmation
// message via Resend. Called when a user saves their email in notification settings.
func (s *EmailDigestService) SetupEmail(toAddress string) error {
if !s.Configured() {
return fmt.Errorf("email service not configured: RESEND_API_KEY not set") //nolint:err113
}
html := `
<div style="font-family: sans-serif; max-width: 600px; margin: 0 auto;">
<h2 style="color: #333;">Welcome to Koin Ping Email Alerts</h2>
<p>Your email has been successfully configured for alert notifications.</p>
<p>You will receive alert digests at this address when events are triggered
on your watched addresses.</p>
<hr style="border: none; border-top: 1px solid #eee; margin: 24px 0;" />
<p style="font-size: 12px; color: #999;">Sent by Koin Ping</p>
</div>`
return s.send(toAddress, "Koin Ping — Email Alerts Configured", html)
}
// SendDigest compiles recent alert events for a user and sends a digest email.
func (s *EmailDigestService) SendDigest(ctx context.Context, userID, toAddress string) error {
if !s.Configured() {
return fmt.Errorf("email service not configured: RESEND_API_KEY not set") //nolint:err113
}
events, err := s.alertEvents.ListRecentByUser(ctx, userID, defaultDigestMaxItems)
if err != nil {
return fmt.Errorf("fetch alert events: %w", err)
}
if len(events) == 0 {
log.Printf("No recent alerts for user %s — skipping digest", userID)
return nil
}
var rows string
for _, e := range events {
label := "—"
if e.AddressLabel != nil {
label = *e.AddressLabel
}
txLink := "—"
if e.TxHash != nil {
txLink = fmt.Sprintf(
`<a href="https://etherscan.io/tx/%s" style="color:#0066cc;">%s…</a>`,
*e.TxHash, (*e.TxHash)[:10],
)
}
rows += fmt.Sprintf(`
<tr>
<td style="padding:6px 8px; border-bottom:1px solid #eee;">%s</td>
<td style="padding:6px 8px; border-bottom:1px solid #eee;">%s</td>
<td style="padding:6px 8px; border-bottom:1px solid #eee;">%s</td>
<td style="padding:6px 8px; border-bottom:1px solid #eee; font-size:12px; color:#666;">%s</td>
</tr>`,
label, e.Message, txLink,
e.Timestamp.Format("Jan 2 15:04 UTC"),
)
}
html := fmt.Sprintf(`
<div style="font-family: sans-serif; max-width: 700px; margin: 0 auto;">
<h2 style="color: #333;">Koin Ping — Alert Digest</h2>
<p>Here are your recent alerts (%d total):</p>
<table style="width:100%%; border-collapse:collapse; font-size:14px;">
<thead>
<tr style="background:#f5f5f5;">
<th style="padding:8px; text-align:left;">Address</th>
<th style="padding:8px; text-align:left;">Alert</th>
<th style="padding:8px; text-align:left;">Tx</th>
<th style="padding:8px; text-align:left;">Time</th>
</tr>
</thead>
<tbody>%s</tbody>
</table>
<hr style="border:none; border-top:1px solid #eee; margin:24px 0;" />
<p style="font-size:12px; color:#999;">Sent by Koin Ping</p>
</div>`, len(events), rows)
subject := fmt.Sprintf("Koin Ping Digest — %d alerts", len(events))
return s.send(toAddress, subject, html)
}
// SendDigestsForAllUsers sends a digest email to every user that has
// notifications enabled and an email configured.
func (s *EmailDigestService) SendDigestsForAllUsers(ctx context.Context) (int, error) {
if !s.Configured() {
return 0, nil
}
configs, err := s.notifCfgs.ListEnabled(ctx)
if err != nil {
return 0, fmt.Errorf("list enabled configs: %w", err)
}
sent := 0
for _, cfg := range configs {
if cfg.Email == nil || *cfg.Email == "" {
continue
}
if err := s.SendDigest(ctx, cfg.UserID, *cfg.Email); err != nil {
log.Printf("Failed to send digest to user %s: %v", cfg.UserID, err)
continue
}
sent++
}
return sent, nil
}
type resendEmailPayload struct {
From string `json:"from"`
To string `json:"to"`
Subject string `json:"subject"`
HTML string `json:"html"`
}
func (s *EmailDigestService) send(to, subject, html string) error {
payload := resendEmailPayload{
From: s.fromAddress,
To: to,
Subject: subject,
HTML: html,
}
body, err := json.Marshal(payload)
if err != nil {
return fmt.Errorf("marshal email payload: %w", err)
}
req, err := http.NewRequest(http.MethodPost, resendAPIURL, bytes.NewReader(body))
if err != nil {
return fmt.Errorf("create request: %w", err)
}
req.Header.Set("Content-Type", "application/json")
req.Header.Set("Authorization", "Bearer "+s.apiKey)
resp, err := digestHTTPClient.Do(req)
if err != nil {
return fmt.Errorf("send email via Resend: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("Resend API returned HTTP %d", resp.StatusCode) //nolint:err113
}
return nil
}

View File

@@ -18,6 +18,8 @@ type EvaluatorService struct {
alertEvents *models.AlertEventModel
addresses *models.AddressModel
notifConfigs *models.NotificationConfigModel
resendAPIKey string
emailFrom string
}
func NewEvaluatorService(
@@ -26,6 +28,8 @@ func NewEvaluatorService(
alertEvents *models.AlertEventModel,
addresses *models.AddressModel,
notifConfigs *models.NotificationConfigModel,
resendAPIKey string,
emailFrom string,
) *EvaluatorService {
return &EvaluatorService{
eth: eth,
@@ -33,6 +37,8 @@ func NewEvaluatorService(
alertEvents: alertEvents,
addresses: addresses,
notifConfigs: notifConfigs,
resendAPIKey: resendAPIKey,
emailFrom: emailFrom,
}
}
@@ -187,25 +193,56 @@ func (s *EvaluatorService) sendNotification(ctx context.Context, userID, message
return
}
if notifConfig == nil || !notifConfig.NotificationEnabled || notifConfig.DiscordWebhookURL == nil {
if notifConfig == nil || !notifConfig.NotificationEnabled {
return
}
sent, err := notifications.SendDiscordNotification(
*notifConfig.DiscordWebhookURL,
message,
notifications.AlertMetadata{
TxHash: obs.Hash,
AddressLabel: addressLabel,
AlertType: string(rule.Type),
Address: address,
},
)
meta := notifications.AlertMetadata{
TxHash: obs.Hash,
AddressLabel: addressLabel,
AlertType: string(rule.Type),
Address: address,
}
if err != nil || !sent {
log.Printf("Discord notification failed for user %s: %v", userID, err)
} else {
log.Printf("Discord notification sent to user %s", userID)
if notifConfig.DiscordWebhookURL != nil && *notifConfig.DiscordWebhookURL != "" {
sent, sendErr := notifications.SendDiscordNotification(*notifConfig.DiscordWebhookURL, message, meta)
if sendErr != nil || !sent {
log.Printf("Discord notification failed for user %s: %v", userID, sendErr)
} else {
log.Printf("Discord notification sent to user %s", userID)
}
}
if notifConfig.TelegramBotToken != nil && *notifConfig.TelegramBotToken != "" &&
notifConfig.TelegramChatID != nil && *notifConfig.TelegramChatID != "" {
sent, sendErr := notifications.SendTelegramNotification(
*notifConfig.TelegramBotToken, *notifConfig.TelegramChatID, message, meta,
)
if sendErr != nil || !sent {
log.Printf("Telegram notification failed for user %s: %v", userID, sendErr)
} else {
log.Printf("Telegram notification sent to user %s", userID)
}
}
if notifConfig.SlackWebhookURL != nil && *notifConfig.SlackWebhookURL != "" {
sent, sendErr := notifications.SendSlackNotification(*notifConfig.SlackWebhookURL, message, meta)
if sendErr != nil || !sent {
log.Printf("Slack notification failed for user %s: %v", userID, sendErr)
} else {
log.Printf("Slack notification sent to user %s", userID)
}
}
if notifConfig.Email != nil && *notifConfig.Email != "" {
sent, sendErr := notifications.SendEmailNotification(
s.resendAPIKey, s.emailFrom, *notifConfig.Email, message, meta,
)
if sendErr != nil || !sent {
log.Printf("Email notification failed for user %s: %v", userID, sendErr)
} else {
log.Printf("Email notification sent to user %s", userID)
}
}
}