Compare commits

..

4 Commits

Author SHA1 Message Date
KS Jannette
7c36f3c214 readme change
Some checks are pending
check / check (push) Waiting to run
2026-03-01 07:55:00 -05:00
KS Jannette
37ff644aa7 audit: fix issues #1-#9 (mock data, notifier interface, retry/timeout, digest scheduling, real status, RPC retry, dedup, address edit UI)
- #1: remove mock alert events that poisoned production responses
- #2: add Notifier interface + Send implementations for all channels
- #3/#4: notification goroutines now carry a 30s context timeout and retry up to 3x with exponential backoff
- #5: schedule email digest in poller via DIGEST_INTERVAL_HOURS (default 24h)
- #6: SystemStatus queries real checkpoint data; returns starting/active/idle with latestBlock and lag
- #7: Ethereum RPC client retries on network errors, 429, and 5xx with 1s/2s backoff
- #8: alert_events dedup index + ON CONFLICT DO NOTHING to prevent duplicate events on poller restart
- #9: PATCH /v1/addresses/{id} for label editing; frontend address list gains inline edit and remove buttons
2026-03-01 07:49:08 -05:00
S Jannette
c80ed89c0a Merge pull request #7 from kjannette/resend-sdk
add sdk
2026-03-01 04:37:58 -05:00
KS Jannette
6590cbf4be add sdk
Some checks are pending
check / check (push) Waiting to run
2026-03-01 04:31:32 -05:00
22 changed files with 536 additions and 126 deletions

View File

@@ -1,9 +1,8 @@
# Koin Ping
A lightweight on-chain monitoring and alerting system designed to give users situational awareness over blockchain addresses they care about.
Koin Ping is an MIT-licensed blockchain monitoring system by Steven Jannette
that polls Ethereum addresses for on-chain activity and delivers real-time
alerts to users via a Go REST API backend and a React single-page application
frontend.
# Overview
Koin Ping observes on-chain activity and notifies users when predefined conditions are met. It does not execute transactions, manage wallets, or speculate on prices.
## Getting Started

View File

@@ -47,6 +47,7 @@ func main() {
addressModel := models.NewAddressModel(pool)
alertRuleModel := models.NewAlertRuleModel(pool)
alertEventModel := models.NewAlertEventModel(pool)
checkpointModel := models.NewCheckpointModel(pool)
notifConfigModel := models.NewNotificationConfigModel(pool)
emailDigestSvc := services.NewEmailDigestService(
@@ -58,13 +59,14 @@ func main() {
alertEventHandler := handlers.NewAlertEventHandler(alertEventModel)
notifConfigHandler := handlers.NewNotificationConfigHandler(notifConfigModel, cfg)
emailDigestHandler := handlers.NewEmailDigestHandler(emailDigestSvc, notifConfigModel)
statusHandler := handlers.NewStatusHandler(checkpointModel)
mux := http.NewServeMux()
b := cfg.APIBasePath // e.g. "/v1"
// Public routes
mux.HandleFunc("GET "+b+"/health", handlers.HealthCheck)
mux.HandleFunc("GET "+b+"/status", handlers.SystemStatus)
mux.HandleFunc("GET "+b+"/status", statusHandler.GetStatus)
// Authenticated routes — addresses
mux.Handle("POST "+b+"/addresses",
@@ -73,6 +75,8 @@ func main() {
middleware.Authenticate(http.HandlerFunc(addressHandler.List)))
mux.Handle("DELETE "+b+"/addresses/{addressId}",
middleware.Authenticate(http.HandlerFunc(addressHandler.Remove)))
mux.Handle("PATCH "+b+"/addresses/{addressId}",
middleware.Authenticate(http.HandlerFunc(addressHandler.UpdateLabel)))
// Authenticated routes for alert rules
mux.Handle("POST "+b+"/addresses/{addressId}/alerts",

View File

@@ -62,6 +62,7 @@ func main() {
eth, alertRuleModel, alertEventModel, addressModel, notifConfigModel,
cfg.ResendAPIKey, cfg.EmailFrom,
)
digestSvc := services.NewEmailDigestService(cfg.ResendAPIKey, cfg.EmailFrom, alertEventModel, notifConfigModel)
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
@@ -78,12 +79,14 @@ func main() {
}()
interval := time.Duration(cfg.PollIntervalMS) * time.Millisecond
digestInterval := time.Duration(cfg.DigestIntervalHours) * time.Hour
log.Println(strings.Repeat("=", separatorWidth))
log.Println("Koin Ping Observer Poller Starting")
log.Println(strings.Repeat("=", separatorWidth))
log.Printf("RPC URL: %s", cfg.EthRPCURL)
log.Printf("Poll Interval: %dms (%ds)", cfg.PollIntervalMS, cfg.PollIntervalMS/msPerSecond)
log.Printf("Digest Interval: %dh", cfg.DigestIntervalHours)
log.Println(strings.Repeat("=", separatorWidth))
runCycle(ctx, observer, evaluator)
@@ -91,6 +94,9 @@ func main() {
ticker := time.NewTicker(interval)
defer ticker.Stop()
digestTicker := time.NewTicker(digestInterval)
defer digestTicker.Stop()
for {
select {
case <-ctx.Done():
@@ -99,6 +105,13 @@ func main() {
return
case <-ticker.C:
runCycle(ctx, observer, evaluator)
case <-digestTicker.C:
sent, digestErr := digestSvc.SendDigestsForAllUsers(ctx)
if digestErr != nil {
log.Printf("Email digest failed: %v", digestErr)
} else {
log.Printf("Sent %d email digests", sent)
}
}
}
}

View File

@@ -42,6 +42,7 @@ require (
github.com/jackc/puddle/v2 v2.2.2 // indirect
github.com/joho/godotenv v1.5.1 // indirect
github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10 // indirect
github.com/resend/resend-go/v3 v3.1.1 // indirect
github.com/spiffe/go-spiffe/v2 v2.6.0 // indirect
go.opentelemetry.io/auto/sdk v1.2.1 // indirect
go.opentelemetry.io/contrib/detectors/gcp v1.39.0 // indirect

View File

@@ -92,6 +92,8 @@ github.com/planetscale/vtprotobuf v0.6.1-0.20240319094008-0393e58bdf10/go.mod h1
github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2 h1:Jamvg5psRIccs7FGNTlIRMkT8wgtp5eCXdBlqhYGL6U=
github.com/pmezard/go-difflib v1.0.1-0.20181226105442-5d4384ee4fb2/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4=
github.com/resend/resend-go/v3 v3.1.1 h1:Uwpf/tZU+O/r/3nMWE6zUAMIG9dX/vTBS3wlQzYJKSw=
github.com/resend/resend-go/v3 v3.1.1/go.mod h1:iI7VA0NoGjWvsNii5iNC5Dy0llsI3HncXPejhniYzwE=
github.com/spiffe/go-spiffe/v2 v2.6.0 h1:l+DolpxNWYgruGQVV0xsfeya3CsC7m8iBzDnMpsbLuo=
github.com/spiffe/go-spiffe/v2 v2.6.0/go.mod h1:gm2SeUoMZEtpnzPNs2Csc0D/gX33k1xIx7lEzqblHEs=
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=

View File

@@ -0,0 +1,3 @@
CREATE UNIQUE INDEX IF NOT EXISTS idx_alert_events_dedup
ON alert_events (alert_rule_id, tx_hash)
WHERE tx_hash IS NOT NULL;

View File

@@ -13,6 +13,7 @@ const (
defaultDBPort = 5432
defaultPollIntervalMS = 60000
minPollIntervalMS = 1000
defaultDigestIntervalHours = 24
)
type Config struct {
@@ -30,6 +31,7 @@ type Config struct {
NodeEnv string
ResendAPIKey string
EmailFrom string
DigestIntervalHours int
}
// Load reads configuration from environment variables and returns a Config.
@@ -49,6 +51,7 @@ func Load() (*Config, error) {
NodeEnv: getEnv("NODE_ENV", "development"),
ResendAPIKey: os.Getenv("RESEND_API_KEY"),
EmailFrom: getEnv("EMAIL_FROM", "Koin Ping <alerts@koinping.com>"),
DigestIntervalHours: getEnvInt("DIGEST_INTERVAL_HOURS", defaultDigestIntervalHours),
}
if cfg.PollIntervalMS < minPollIntervalMS {

View File

@@ -92,6 +92,44 @@ func (h *AddressHandler) List(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, addresses)
}
// UpdateLabel handles PATCH requests to update an address label.
func (h *AddressHandler) UpdateLabel(w http.ResponseWriter, r *http.Request) {
userID := middleware.GetUserID(r.Context())
addressID, ok := parseIntParam(r.PathValue("addressId"))
if !ok {
writeError(w, http.StatusBadRequest, "VALIDATION_ERROR", "Invalid address ID")
return
}
var body struct {
Label *string `json:"label"`
}
if err := json.NewDecoder(r.Body).Decode(&body); err != nil {
writeError(w, http.StatusBadRequest, "VALIDATION_ERROR", "Invalid request body")
return
}
log.Printf("User %s updating label for address ID: %d", userID, addressID)
addr, err := h.addresses.UpdateLabel(r.Context(), addressID, userID, body.Label)
if err != nil {
log.Printf("Error updating address label: %v", err)
writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR", "Failed to update address")
return
}
if addr == nil {
writeError(w, http.StatusNotFound, "NOT_FOUND", "Address not found")
return
}
writeJSON(w, http.StatusOK, addr)
}
// Remove handles DELETE requests to remove a tracked address.
func (h *AddressHandler) Remove(w http.ResponseWriter, r *http.Request) {
userID := middleware.GetUserID(r.Context())

View File

@@ -4,7 +4,6 @@ import (
"log"
"net/http"
"strconv"
"time"
"github.com/kjannette/koin-ping/backend-go/internal/domain"
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
@@ -51,45 +50,9 @@ func (h *AlertEventHandler) List(w http.ResponseWriter, r *http.Request) {
log.Printf("Found %d alert events for user", len(events))
// MVP scaffolding: return mock data if DB is empty
if len(events) == 0 {
events = mockEvents(limit)
if events == nil {
events = []domain.AlertEvent{}
}
writeJSON(w, http.StatusOK, events)
}
func mockEvents(limit int) []domain.AlertEvent {
label1 := "Treasury Wallet"
label2 := "Cold Storage"
mocks := []domain.AlertEvent{
{
ID: 1,
AlertRuleID: 1,
Message: "Incoming transaction detected: 5.5 ETH received",
AddressLabel: &label1,
Timestamp: time.Now().Add(-2 * time.Hour),
},
{
ID: 2, //nolint:mnd
AlertRuleID: 2, //nolint:mnd
Message: "Balance dropped below threshold: Current balance 8.2 ETH",
AddressLabel: &label1,
Timestamp: time.Now().Add(-5 * time.Hour),
},
{
ID: 3, //nolint:mnd
AlertRuleID: 3, //nolint:mnd
Message: "Outgoing transaction detected: 2.0 ETH sent",
AddressLabel: &label2,
Timestamp: time.Now().Add(-24 * time.Hour),
},
}
if limit < len(mocks) {
return mocks[:limit]
}
return mocks
}

View File

@@ -1,10 +1,67 @@
package handlers
import (
"log"
"net/http"
"time"
"github.com/kjannette/koin-ping/backend-go/internal/models"
)
// StatusHandler handles the system status endpoint.
type StatusHandler struct {
checkpoints *models.CheckpointModel
}
// NewStatusHandler creates a new StatusHandler.
func NewStatusHandler(checkpoints *models.CheckpointModel) *StatusHandler {
return &StatusHandler{checkpoints: checkpoints}
}
// GetStatus returns real-time system status derived from checkpoint data.
func (h *StatusHandler) GetStatus(w http.ResponseWriter, r *http.Request) {
block, checkedAt, err := h.checkpoints.GetLatestBlock(r.Context())
if err != nil {
log.Printf("Error querying latest block: %v", err)
writeError(w, http.StatusInternalServerError, "INTERNAL_ERROR", "Failed to get system status")
return
}
latestBlock := 0
lag := 0
status := "starting"
if checkedAt != nil {
lag = int(time.Since(*checkedAt).Seconds())
if lag > 600 { //nolint:mnd
status = "idle"
} else {
status = "active"
}
}
if block != nil {
latestBlock = *block
}
writeJSON(w, http.StatusOK, map[string]interface{}{
"status": status,
"latestBlock": latestBlock,
"lag": lag,
"lastCheckedAt": checkedAtStr(checkedAt),
"timestamp": time.Now().UTC().Format(time.RFC3339),
})
}
func checkedAtStr(t *time.Time) string {
if t == nil {
return ""
}
return t.UTC().Format(time.RFC3339)
}
func HealthCheck(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]interface{}{
"status": "ok",
@@ -12,12 +69,3 @@ func HealthCheck(w http.ResponseWriter, r *http.Request) {
"service": "koin-ping-backend",
})
}
func SystemStatus(w http.ResponseWriter, r *http.Request) {
writeJSON(w, http.StatusOK, map[string]interface{}{
"latestBlock": 0,
"lag": 0,
"status": "healthy",
"timestamp": time.Now().UTC().Format(time.RFC3339),
})
}

View File

@@ -107,6 +107,24 @@ func (m *AddressModel) FindByID(ctx context.Context, id int, userID *string) (*d
return &a, nil
}
// UpdateLabel updates the label for an address owned by userID.
// Returns nil, nil if no row matched (address not found or not owned by user).
func (m *AddressModel) UpdateLabel(ctx context.Context, id int, userID string, label *string) (*domain.Address, error) {
var a domain.Address
err := m.pool.QueryRow(ctx,
`UPDATE addresses SET label = $3 WHERE id = $1 AND user_id = $2
RETURNING id, user_id, address, label, created_at`,
id, userID, label,
).Scan(&a.ID, &a.UserID, &a.Address, &a.Label, &a.CreatedAt)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
return nil, nil
}
return nil, err
}
return &a, nil
}
func (m *AddressModel) Remove(ctx context.Context, id int, userID string) (bool, error) {
tag, err := m.pool.Exec(ctx,
`DELETE FROM addresses WHERE id = $1 AND user_id = $2`,

View File

@@ -2,7 +2,9 @@ package models
import (
"context"
"errors"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/kjannette/koin-ping/backend-go/internal/domain"
)
@@ -71,10 +73,15 @@ func (m *AlertEventModel) Create(ctx context.Context, alertRuleID int, message s
err := m.pool.QueryRow(ctx,
`INSERT INTO alert_events (alert_rule_id, message, address_label, tx_hash)
VALUES ($1, $2, $3, $4)
ON CONFLICT DO NOTHING
RETURNING id, alert_rule_id, message, address_label, tx_hash, timestamp`,
alertRuleID, message, addressLabel, txHash,
).Scan(&e.ID, &e.AlertRuleID, &e.Message, &e.AddressLabel, &e.TxHash, &e.Timestamp)
if err != nil {
if errors.Is(err, pgx.ErrNoRows) {
// Duplicate silently skipped by ON CONFLICT DO NOTHING
return nil, nil
}
return nil, err
}
return &e, nil

View File

@@ -3,6 +3,7 @@ package models
import (
"context"
"errors"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
@@ -17,6 +18,20 @@ func NewCheckpointModel(pool *pgxpool.Pool) *CheckpointModel {
return &CheckpointModel{pool: pool}
}
// GetLatestBlock returns the highest last_checked_block and its timestamp across all addresses.
// Returns nil, nil, nil when no checkpoints exist yet.
func (m *CheckpointModel) GetLatestBlock(ctx context.Context) (*int, *time.Time, error) {
var block *int
var checkedAt *time.Time
err := m.pool.QueryRow(ctx,
`SELECT MAX(last_checked_block), MAX(last_checked_at) FROM address_checkpoints`,
).Scan(&block, &checkedAt)
if err != nil {
return nil, nil, err
}
return block, checkedAt, nil
}
// GetLastCheckedBlock returns the last checked block for an address, or -1 if never checked.
func (m *CheckpointModel) GetLastCheckedBlock(ctx context.Context, addressID int) (int, bool, error) {
var block int

View File

@@ -2,6 +2,7 @@ package notifications
import (
"bytes"
"context"
"encoding/json"
"fmt"
"log"
@@ -25,11 +26,15 @@ var discordHTTPClient = &http.Client{ //nolint:gochecknoglobals
Timeout: discordHTTPTimeoutSeconds * time.Second,
}
type AlertMetadata struct {
TxHash string
AddressLabel string
AlertType string
Address string
// DiscordNotifier sends alert notifications via a Discord webhook.
type DiscordNotifier struct {
WebhookURL string
}
// Send implements Notifier for Discord.
func (d *DiscordNotifier) Send(_ context.Context, message string, meta AlertMetadata) error {
_, err := SendDiscordNotification(d.WebhookURL, message, meta)
return err
}
type discordEmbed struct {

View File

@@ -2,6 +2,7 @@ package notifications
import (
"bytes"
"context"
"encoding/json"
"fmt"
"log"
@@ -9,6 +10,19 @@ import (
"time"
)
// EmailNotifier sends alert notifications via email (Resend).
type EmailNotifier struct {
APIKey string
From string
To string
}
// Send implements Notifier for email.
func (e *EmailNotifier) Send(_ context.Context, message string, meta AlertMetadata) error {
_, err := SendEmailNotification(e.APIKey, e.From, e.To, message, meta)
return err
}
const emailHTTPTimeoutSeconds = 10
var emailHTTPClient = &http.Client{ //nolint:gochecknoglobals

View File

@@ -0,0 +1,16 @@
package notifications
import "context"
// AlertMetadata holds context about the alert being sent.
type AlertMetadata struct {
TxHash string
AddressLabel string
AlertType string
Address string
}
// Notifier is the interface implemented by all notification channels.
type Notifier interface {
Send(ctx context.Context, message string, meta AlertMetadata) error
}

View File

@@ -2,6 +2,7 @@ package notifications
import (
"bytes"
"context"
"encoding/json"
"fmt"
"log"
@@ -9,6 +10,17 @@ import (
"time"
)
// SlackNotifier sends alert notifications via a Slack webhook.
type SlackNotifier struct {
WebhookURL string
}
// Send implements Notifier for Slack.
func (s *SlackNotifier) Send(_ context.Context, message string, meta AlertMetadata) error {
_, err := SendSlackNotification(s.WebhookURL, message, meta)
return err
}
const slackHTTPTimeoutSeconds = 10
var slackHTTPClient = &http.Client{ //nolint:gochecknoglobals

View File

@@ -2,6 +2,7 @@ package notifications
import (
"bytes"
"context"
"encoding/json"
"fmt"
"log"
@@ -9,6 +10,18 @@ import (
"time"
)
// TelegramNotifier sends alert notifications via Telegram.
type TelegramNotifier struct {
BotToken string
ChatID string
}
// Send implements Notifier for Telegram.
func (t *TelegramNotifier) Send(_ context.Context, message string, meta AlertMetadata) error {
_, err := SendTelegramNotification(t.BotToken, t.ChatID, message, meta)
return err
}
const telegramHTTPTimeoutSeconds = 10
var telegramHTTPClient = &http.Client{ //nolint:gochecknoglobals

View File

@@ -4,7 +4,9 @@ import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"log"
"math/big"
"net/http"
"strings"
@@ -13,7 +15,11 @@ import (
"github.com/kjannette/koin-ping/backend-go/internal/domain"
)
const rpcTimeoutMS = 30000
const (
rpcTimeoutMS = 30000
rpcMaxRetries = 3
rpcRetryBaseMS = 1000
)
type JsonRpcEthereum struct {
rpcURL string
@@ -66,6 +72,44 @@ func (j *JsonRpcEthereum) callRPC(ctx context.Context, method string, params ...
return nil, fmt.Errorf("marshal RPC request: %w", err)
}
return j.callWithRetry(ctx, method, body)
}
// callWithRetry executes a JSON-RPC POST with exponential backoff on transient errors.
// It retries on network errors, HTTP 429, and HTTP 5xx. It does NOT retry on RPC-level
// errors or other 4xx responses (those are permanent failures).
func (j *JsonRpcEthereum) callWithRetry(ctx context.Context, method string, body []byte) (json.RawMessage, error) {
var lastErr error
for attempt := range rpcMaxRetries {
if attempt > 0 {
wait := time.Duration(rpcRetryBaseMS*(1<<(attempt-1))) * time.Millisecond
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(wait):
}
log.Printf("Retrying RPC call [%s] (attempt %d/%d)", method, attempt+1, rpcMaxRetries)
}
result, err := j.doRPCCall(ctx, method, body)
if err == nil {
return result, nil
}
lastErr = err
// Permanent errors: do not retry
if isPermanentRPCError(err) {
return nil, err
}
log.Printf("Transient RPC error [%s] (attempt %d/%d): %v", method, attempt+1, rpcMaxRetries, err)
}
return nil, lastErr
}
func (j *JsonRpcEthereum) doRPCCall(ctx context.Context, method string, body []byte) (json.RawMessage, error) {
req, err := http.NewRequestWithContext(ctx, http.MethodPost, j.rpcURL, bytes.NewReader(body))
if err != nil {
return nil, fmt.Errorf("create RPC request: %w", err)
@@ -78,8 +122,13 @@ func (j *JsonRpcEthereum) callRPC(ctx context.Context, method string, params ...
}
defer resp.Body.Close()
// 429 and 5xx are transient; other non-200 are permanent.
if resp.StatusCode != http.StatusOK {
return nil, fmt.Errorf("HTTP %d: %s for %s", resp.StatusCode, resp.Status, method)
err := fmt.Errorf("HTTP %d: %s for %s", resp.StatusCode, resp.Status, method)
if resp.StatusCode == http.StatusTooManyRequests || resp.StatusCode >= 500 { //nolint:mnd
return nil, err // transient — will be retried
}
return nil, &permanentRPCError{err}
}
var rpcResp rpcResponse
@@ -88,12 +137,25 @@ func (j *JsonRpcEthereum) callRPC(ctx context.Context, method string, params ...
}
if rpcResp.Error != nil {
return nil, fmt.Errorf("RPC Error [%s]: %s (code: %d)", method, rpcResp.Error.Message, rpcResp.Error.Code)
// RPC-level errors are permanent (bad params, unsupported method, etc.)
return nil, &permanentRPCError{
fmt.Errorf("RPC Error [%s]: %s (code: %d)", method, rpcResp.Error.Message, rpcResp.Error.Code),
}
}
return rpcResp.Result, nil
}
type permanentRPCError struct{ cause error }
func (e *permanentRPCError) Error() string { return e.cause.Error() }
func (e *permanentRPCError) Unwrap() error { return e.cause }
func isPermanentRPCError(err error) bool {
var p *permanentRPCError
return errors.As(err, &p)
}
func (j *JsonRpcEthereum) GetLatestBlockNumber(ctx context.Context) (int, error) {
result, err := j.callRPC(ctx, "eth_blockNumber")
if err != nil {

View File

@@ -4,6 +4,7 @@ import (
"context"
"fmt"
"log"
"time"
"github.com/kjannette/koin-ping/backend-go/internal/domain"
"github.com/kjannette/koin-ping/backend-go/internal/models"
@@ -12,6 +13,12 @@ import (
"github.com/kjannette/koin-ping/backend-go/internal/wei"
)
const (
notificationTimeout = 30 * time.Second
notificationMaxRetries = 3
notificationRetryBase = time.Second
)
type EvaluatorService struct {
eth ethereum.EthereumObserver
alertRules *models.AlertRuleModel
@@ -167,25 +174,85 @@ func (s *EvaluatorService) fireAlert(ctx context.Context, rule domain.AlertRule,
message := s.buildMessage(rule, obs)
txHash := &obs.Hash
_, err = s.alertEvents.Create(ctx, rule.ID, message, &addressLabel, txHash)
event, err := s.alertEvents.Create(ctx, rule.ID, message, &addressLabel, txHash)
if err != nil {
return err
}
if event == nil {
log.Printf("[ALERT DEDUP] Rule %d (%s) - duplicate event skipped for TX: %s", rule.ID, rule.Type, obs.Hash)
return nil
}
log.Printf("[ALERT FIRED] Rule %d (%s) - %s - TX: %s", rule.ID, rule.Type, message, obs.Hash)
// Send Discord notification (non-fatal on failure)
if addr != nil {
userID := addr.UserID
address := addr.Address
go func() {
s.sendNotification(
ctx, addr.UserID, message, obs, addressLabel, rule, addr.Address,
)
notifCtx, cancel := context.WithTimeout(context.Background(), notificationTimeout)
defer cancel()
s.sendNotification(notifCtx, userID, message, obs, addressLabel, rule, address)
}()
}
return nil
}
func (s *EvaluatorService) buildNotifiers(cfg *domain.NotificationConfig) []notifications.Notifier {
var notifiers []notifications.Notifier
if cfg.DiscordWebhookURL != nil && *cfg.DiscordWebhookURL != "" {
notifiers = append(notifiers, &notifications.DiscordNotifier{WebhookURL: *cfg.DiscordWebhookURL})
}
if cfg.TelegramBotToken != nil && *cfg.TelegramBotToken != "" &&
cfg.TelegramChatID != nil && *cfg.TelegramChatID != "" {
notifiers = append(notifiers, &notifications.TelegramNotifier{
BotToken: *cfg.TelegramBotToken,
ChatID: *cfg.TelegramChatID,
})
}
if cfg.SlackWebhookURL != nil && *cfg.SlackWebhookURL != "" {
notifiers = append(notifiers, &notifications.SlackNotifier{WebhookURL: *cfg.SlackWebhookURL})
}
if cfg.Email != nil && *cfg.Email != "" {
notifiers = append(notifiers, &notifications.EmailNotifier{
APIKey: s.resendAPIKey,
From: s.emailFrom,
To: *cfg.Email,
})
}
return notifiers
}
func sendWithRetry(ctx context.Context, n notifications.Notifier, message string, meta notifications.AlertMetadata) error {
var lastErr error
for attempt := range notificationMaxRetries {
if attempt > 0 {
wait := notificationRetryBase * time.Duration(1<<(attempt-1))
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(wait):
}
}
if err := n.Send(ctx, message, meta); err != nil {
log.Printf("Notification attempt %d/%d failed: %v", attempt+1, notificationMaxRetries, err)
lastErr = err
continue
}
return nil
}
return lastErr
}
func (s *EvaluatorService) sendNotification(ctx context.Context, userID, message string, obs domain.ObservedTx, addressLabel string, rule domain.AlertRule, address string) {
notifConfig, err := s.notifConfigs.GetConfig(ctx, userID)
if err != nil {
@@ -204,44 +271,11 @@ func (s *EvaluatorService) sendNotification(ctx context.Context, userID, message
Address: address,
}
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)
for _, n := range s.buildNotifiers(notifConfig) {
if err := sendWithRetry(ctx, n, message, meta); err != nil {
log.Printf("Notification channel failed for user %s after retries: %v", userID, err)
} 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)
log.Printf("Notification sent to user %s via %T", userID, n)
}
}
}

View File

@@ -80,6 +80,41 @@ export async function getAddresses() {
}
}
/**
* Update an address (e.g. change its label)
* @param {number} addressId - Address ID to update
* @param {Object} data - Fields to update (e.g. { label: "New Label" })
* @returns {Promise<Object>} Updated address
*/
export async function updateAddress(addressId, data) {
try {
const headers = await getAuthHeaders();
const response = await fetch(`${API_BASE}/addresses/${addressId}`, {
method: "PATCH",
headers: headers,
body: JSON.stringify(data),
});
if (!response.ok) {
let errorMessage = "Failed to update address";
try {
const error = await response.json();
errorMessage = error.message || errorMessage;
} catch {
errorMessage = `Server error: ${response.status} ${response.statusText}`;
}
throw new Error(errorMessage);
}
return response.json();
} catch (error) {
if (error.message.includes("fetch")) {
throw new Error("Cannot connect to server. Is the backend running?");
}
throw error;
}
}
/**
* Delete a tracked address
* @param {number} addressId - Address ID to delete

View File

@@ -1,11 +1,13 @@
import { useState, useEffect } from "react";
import AddressForm from "../components/AddressForm";
import { getAddresses, createAddress } from "../api/addresses";
import { getAddresses, createAddress, deleteAddress, updateAddress } from "../api/addresses";
export default function Addresses() {
const [addresses, setAddresses] = useState([]);
const [loading, setLoading] = useState(true);
const [error, setError] = useState(null);
const [editingId, setEditingId] = useState(null);
const [editLabel, setEditLabel] = useState("");
// Load addresses on mount
useEffect(() => {
@@ -29,15 +31,52 @@ export default function Addresses() {
async function handleAddressSubmit(data) {
try {
const newAddress = await createAddress(data);
// Append new address to state
setAddresses((prev) => [...prev, newAddress]);
setError(null); // Clear any previous errors
setError(null);
} catch (err) {
setError(err.message);
console.error("Failed to create address:", err);
}
}
async function handleDelete(id, label) {
const displayName = label || "this address";
if (!window.confirm(`Remove "${displayName}"? This will also delete all associated alert rules.`)) {
return;
}
try {
await deleteAddress(id);
setAddresses((prev) => prev.filter((a) => a.id !== id));
setError(null);
} catch (err) {
setError(err.message);
console.error("Failed to delete address:", err);
}
}
function handleEditStart(addr) {
setEditingId(addr.id);
setEditLabel(addr.label ?? "");
}
async function handleEditSave(id) {
try {
const updated = await updateAddress(id, { label: editLabel || null });
setAddresses((prev) => prev.map((a) => (a.id === id ? updated : a)));
setEditingId(null);
setEditLabel("");
setError(null);
} catch (err) {
setError(err.message);
console.error("Failed to update address:", err);
}
}
function handleEditCancel() {
setEditingId(null);
setEditLabel("");
}
return (
<div style={{ maxWidth: "800px", margin: "0 auto", padding: "2rem" }}>
<h1>Tracked Addresses</h1>
@@ -68,14 +107,62 @@ export default function Addresses() {
backgroundColor: "#333",
}}
>
<div
<div style={{ display: "flex", justifyContent: "space-between", alignItems: "center" }}>
<div style={{ flex: 1 }}>
{editingId === addr.id ? (
<div style={{ display: "flex", gap: "0.5rem", alignItems: "center", marginBottom: "0.25rem" }}>
<input
value={editLabel}
onChange={(e) => setEditLabel(e.target.value)}
placeholder="Label (optional)"
style={{
fontWeight: "bold",
marginBottom: "0.25rem",
background: "#444",
border: "1px solid #666",
borderRadius: "3px",
color: "#fff",
padding: "0.25rem 0.5rem",
fontSize: "0.9rem",
}}
onKeyDown={(e) => {
if (e.key === "Enter") handleEditSave(addr.id);
if (e.key === "Escape") handleEditCancel();
}}
autoFocus
/>
<button
onClick={() => handleEditSave(addr.id)}
style={{ cursor: "pointer", padding: "0.25rem 0.6rem", fontSize: "0.85rem" }}
>
Save
</button>
<button
onClick={handleEditCancel}
style={{ cursor: "pointer", padding: "0.25rem 0.6rem", fontSize: "0.85rem", background: "transparent", color: "#aaa", border: "1px solid #555" }}
>
Cancel
</button>
</div>
) : (
<div style={{ display: "flex", alignItems: "center", gap: "0.5rem", marginBottom: "0.25rem" }}>
<span style={{ fontWeight: "bold" }}>
{addr.label || "Unlabeled"}
</span>
<button
onClick={() => handleEditStart(addr)}
style={{
cursor: "pointer",
background: "transparent",
border: "none",
color: "#6699cc",
fontSize: "0.8rem",
padding: "0",
textDecoration: "underline",
}}
>
{addr.label || "Unlabeled"}
Edit
</button>
</div>
)}
<div
style={{
fontFamily: "monospace",
@@ -85,6 +172,24 @@ export default function Addresses() {
>
{addr.address}
</div>
</div>
<button
onClick={() => handleDelete(addr.id, addr.label)}
style={{
cursor: "pointer",
background: "transparent",
border: "1px solid #884444",
color: "#cc6666",
borderRadius: "3px",
padding: "0.3rem 0.7rem",
fontSize: "0.85rem",
marginLeft: "1rem",
flexShrink: 0,
}}
>
Remove
</button>
</div>
</li>
))}
</ul>