From 37ff644aa783520b5b276006c7a7f5170b6df0c0 Mon Sep 17 00:00:00 2001 From: KS Jannette Date: Sun, 1 Mar 2026 07:49:08 -0500 Subject: [PATCH] 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 --- backend-go/cmd/api/main.go | 6 +- backend-go/cmd/poller/main.go | 13 ++ .../migrations/004_alert_event_dedup.sql | 3 + backend-go/internal/config/config.go | 11 +- backend-go/internal/handlers/address.go | 38 +++++ backend-go/internal/handlers/alert_event.go | 41 +---- backend-go/internal/handlers/status.go | 66 ++++++-- backend-go/internal/models/address.go | 18 +++ backend-go/internal/models/alert_event.go | 7 + backend-go/internal/models/checkpoint.go | 15 ++ backend-go/internal/notifications/discord.go | 15 +- backend-go/internal/notifications/email.go | 14 ++ backend-go/internal/notifications/notifier.go | 16 ++ backend-go/internal/notifications/slack.go | 12 ++ backend-go/internal/notifications/telegram.go | 13 ++ .../internal/protocols/ethereum/jsonrpc.go | 68 ++++++++- backend-go/internal/services/evaluator.go | 118 ++++++++++----- frontend/src/api/addresses.jsx | 35 +++++ frontend/src/pages/Addresses.jsx | 143 +++++++++++++++--- 19 files changed, 530 insertions(+), 122 deletions(-) create mode 100644 backend-go/infra/migrations/004_alert_event_dedup.sql create mode 100644 backend-go/internal/notifications/notifier.go diff --git a/backend-go/cmd/api/main.go b/backend-go/cmd/api/main.go index 3068a33..31d87a8 100644 --- a/backend-go/cmd/api/main.go +++ b/backend-go/cmd/api/main.go @@ -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", diff --git a/backend-go/cmd/poller/main.go b/backend-go/cmd/poller/main.go index abeee6f..8fb8714 100644 --- a/backend-go/cmd/poller/main.go +++ b/backend-go/cmd/poller/main.go @@ -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) + } } } } diff --git a/backend-go/infra/migrations/004_alert_event_dedup.sql b/backend-go/infra/migrations/004_alert_event_dedup.sql new file mode 100644 index 0000000..b0cbffc --- /dev/null +++ b/backend-go/infra/migrations/004_alert_event_dedup.sql @@ -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; diff --git a/backend-go/internal/config/config.go b/backend-go/internal/config/config.go index 49370c3..cfa0019 100644 --- a/backend-go/internal/config/config.go +++ b/backend-go/internal/config/config.go @@ -13,6 +13,7 @@ const ( defaultDBPort = 5432 defaultPollIntervalMS = 60000 minPollIntervalMS = 1000 + defaultDigestIntervalHours = 24 ) type Config struct { @@ -28,8 +29,9 @@ type Config struct { EthRPCURL string PollIntervalMS int NodeEnv string - ResendAPIKey string - EmailFrom string + ResendAPIKey string + EmailFrom string + DigestIntervalHours int } // Load reads configuration from environment variables and returns a Config. @@ -47,8 +49,9 @@ 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 "), + ResendAPIKey: os.Getenv("RESEND_API_KEY"), + EmailFrom: getEnv("EMAIL_FROM", "Koin Ping "), + DigestIntervalHours: getEnvInt("DIGEST_INTERVAL_HOURS", defaultDigestIntervalHours), } if cfg.PollIntervalMS < minPollIntervalMS { diff --git a/backend-go/internal/handlers/address.go b/backend-go/internal/handlers/address.go index f42946c..8d6cf4c 100644 --- a/backend-go/internal/handlers/address.go +++ b/backend-go/internal/handlers/address.go @@ -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()) diff --git a/backend-go/internal/handlers/alert_event.go b/backend-go/internal/handlers/alert_event.go index 27aef80..0e22de2 100644 --- a/backend-go/internal/handlers/alert_event.go +++ b/backend-go/internal/handlers/alert_event.go @@ -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 -} diff --git a/backend-go/internal/handlers/status.go b/backend-go/internal/handlers/status.go index 621adaa..aac320a 100644 --- a/backend-go/internal/handlers/status.go +++ b/backend-go/internal/handlers/status.go @@ -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), - }) -} diff --git a/backend-go/internal/models/address.go b/backend-go/internal/models/address.go index fd1c8f9..42d16ad 100644 --- a/backend-go/internal/models/address.go +++ b/backend-go/internal/models/address.go @@ -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`, diff --git a/backend-go/internal/models/alert_event.go b/backend-go/internal/models/alert_event.go index a1cb7db..3f68553 100644 --- a/backend-go/internal/models/alert_event.go +++ b/backend-go/internal/models/alert_event.go @@ -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 diff --git a/backend-go/internal/models/checkpoint.go b/backend-go/internal/models/checkpoint.go index c6c60d6..0be5831 100644 --- a/backend-go/internal/models/checkpoint.go +++ b/backend-go/internal/models/checkpoint.go @@ -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 diff --git a/backend-go/internal/notifications/discord.go b/backend-go/internal/notifications/discord.go index b1bb788..61c576f 100644 --- a/backend-go/internal/notifications/discord.go +++ b/backend-go/internal/notifications/discord.go @@ -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 { diff --git a/backend-go/internal/notifications/email.go b/backend-go/internal/notifications/email.go index c431d9a..9bd74fd 100644 --- a/backend-go/internal/notifications/email.go +++ b/backend-go/internal/notifications/email.go @@ -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 diff --git a/backend-go/internal/notifications/notifier.go b/backend-go/internal/notifications/notifier.go new file mode 100644 index 0000000..b317fcb --- /dev/null +++ b/backend-go/internal/notifications/notifier.go @@ -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 +} diff --git a/backend-go/internal/notifications/slack.go b/backend-go/internal/notifications/slack.go index 989ffd6..ff990a4 100644 --- a/backend-go/internal/notifications/slack.go +++ b/backend-go/internal/notifications/slack.go @@ -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 diff --git a/backend-go/internal/notifications/telegram.go b/backend-go/internal/notifications/telegram.go index 4489b17..ba17d3d 100644 --- a/backend-go/internal/notifications/telegram.go +++ b/backend-go/internal/notifications/telegram.go @@ -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 diff --git a/backend-go/internal/protocols/ethereum/jsonrpc.go b/backend-go/internal/protocols/ethereum/jsonrpc.go index 1e80cf5..ae77a55 100644 --- a/backend-go/internal/protocols/ethereum/jsonrpc.go +++ b/backend-go/internal/protocols/ethereum/jsonrpc.go @@ -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 { diff --git a/backend-go/internal/services/evaluator.go b/backend-go/internal/services/evaluator.go index 99569d2..02b5079 100644 --- a/backend-go/internal/services/evaluator.go +++ b/backend-go/internal/services/evaluator.go @@ -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, ¬ifications.DiscordNotifier{WebhookURL: *cfg.DiscordWebhookURL}) + } + + if cfg.TelegramBotToken != nil && *cfg.TelegramBotToken != "" && + cfg.TelegramChatID != nil && *cfg.TelegramChatID != "" { + notifiers = append(notifiers, ¬ifications.TelegramNotifier{ + BotToken: *cfg.TelegramBotToken, + ChatID: *cfg.TelegramChatID, + }) + } + + if cfg.SlackWebhookURL != nil && *cfg.SlackWebhookURL != "" { + notifiers = append(notifiers, ¬ifications.SlackNotifier{WebhookURL: *cfg.SlackWebhookURL}) + } + + if cfg.Email != nil && *cfg.Email != "" { + notifiers = append(notifiers, ¬ifications.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) } } } diff --git a/frontend/src/api/addresses.jsx b/frontend/src/api/addresses.jsx index 5cc8ccd..f31de5d 100644 --- a/frontend/src/api/addresses.jsx +++ b/frontend/src/api/addresses.jsx @@ -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} 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 diff --git a/frontend/src/pages/Addresses.jsx b/frontend/src/pages/Addresses.jsx index 4783342..13434f0 100644 --- a/frontend/src/pages/Addresses.jsx +++ b/frontend/src/pages/Addresses.jsx @@ -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 (

Tracked Addresses

@@ -68,22 +107,88 @@ export default function Addresses() { backgroundColor: "#333", }} > -
- {addr.label || "Unlabeled"} -
-
- {addr.address} +
+
+ {editingId === addr.id ? ( +
+ setEditLabel(e.target.value)} + placeholder="Label (optional)" + style={{ + 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 + /> + + +
+ ) : ( +
+ + {addr.label || "Unlabeled"} + + +
+ )} +
+ {addr.address} +
+
+
))}