Compare commits
10 Commits
resend-sdk
...
configureD
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
35107c0a9d | ||
|
|
9935817fa8 | ||
|
|
8d5716bdb7 | ||
|
|
aafe96d9a1 | ||
|
|
93168209b5 | ||
|
|
6840af1a47 | ||
|
|
7c36f3c214 | ||
|
|
c707b82f36 | ||
|
|
37ff644aa7 | ||
|
|
c80ed89c0a |
130
README.md
130
README.md
@@ -1,9 +1,10 @@
|
|||||||
# Koin Ping
|
# Koin Ping
|
||||||
|
A lightweight on-chain monitoring and alerting system designed to give users situational awareness over blockchain addresses they care about.
|
||||||
|
|
||||||
|
# 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.
|
||||||
|
|
||||||
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.
|
|
||||||
|
|
||||||
## Getting Started
|
## Getting Started
|
||||||
|
|
||||||
@@ -99,6 +100,127 @@ to the database, and dispatches Discord notifications.
|
|||||||
Firebase, communicates with the API via fetch, and renders the address/alert
|
Firebase, communicates with the API via fetch, and renders the address/alert
|
||||||
management UI.
|
management UI.
|
||||||
|
|
||||||
|
## Setting Up Your Alert Platforms
|
||||||
|
|
||||||
|
Koin Ping can send real-time alerts to **Telegram**, **Discord**, **Slack**, and **Email**. Each channel is configured per-user through the **Notification Settings** panel on the Alerts page.
|
||||||
|
|
||||||
|
Below are step-by-step guides for setting up each platform.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### Telegram
|
||||||
|
|
||||||
|
To receive alerts via Telegram, you need to create a bot and get your chat ID.
|
||||||
|
|
||||||
|
#### 1. Create a Telegram Bot
|
||||||
|
|
||||||
|
1. Open Telegram and search for **@BotFather** (look for the blue verified checkmark).
|
||||||
|
2. Open the conversation with BotFather and send: `/newbot`
|
||||||
|
3. BotFather will ask for a **display name** — enter something like `Koin Ping Alerts`.
|
||||||
|
4. BotFather will ask for a **username** — it must end in `bot`, e.g. `MyKoinPingBot`.
|
||||||
|
5. BotFather will reply with your **Bot Token** — a string that looks like `123456789:ABCdefGHIjklMNOpqrSTUvwxYZ`. Copy it.
|
||||||
|
|
||||||
|
#### 2. Get Your Chat ID
|
||||||
|
|
||||||
|
1. In Telegram, search for the bot username you just created and open the chat.
|
||||||
|
2. Tap **Start** or send any message (e.g. `hello`).
|
||||||
|
3. Open the following URL in your browser, replacing `YOUR_BOT_TOKEN` with the token from step 1:
|
||||||
|
|
||||||
|
```
|
||||||
|
https://api.telegram.org/botYOUR_BOT_TOKEN/getUpdates
|
||||||
|
```
|
||||||
|
|
||||||
|
4. In the JSON response, find the `"chat"` object — the `"id"` field is your **Chat ID** (a numeric value).
|
||||||
|
|
||||||
|
> **Tip:** If the `"result"` array is empty, make sure you sent a message to your bot first, then refresh the page.
|
||||||
|
|
||||||
|
#### 3. Save in Koin Ping
|
||||||
|
|
||||||
|
1. Go to the **Alerts** page in Koin Ping.
|
||||||
|
2. In the **Notification Settings** panel, find the **Telegram** section.
|
||||||
|
3. Paste your **Bot Token** and **Chat ID** into the corresponding fields.
|
||||||
|
4. Click **Save Settings**.
|
||||||
|
5. Click **Test All Channels** to verify — you should receive a test message from your bot in Telegram.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### Discord
|
||||||
|
|
||||||
|
To receive alerts in a Discord channel, you need to create a webhook.
|
||||||
|
|
||||||
|
#### 1. Create a Discord Webhook
|
||||||
|
|
||||||
|
1. Open **Discord** and navigate to the server where you want to receive alerts.
|
||||||
|
2. Go to the channel you want alerts posted in (or create a new one, e.g. `#koin-alerts`).
|
||||||
|
3. Click the **gear icon** next to the channel name to open **Channel Settings**.
|
||||||
|
4. Go to **Integrations** > **Webhooks**.
|
||||||
|
5. Click **New Webhook**.
|
||||||
|
6. Give it a name (e.g. `Koin Ping`) and click **Copy Webhook URL**.
|
||||||
|
|
||||||
|
The URL will look like: `https://discord.com/api/webhooks/123456789/AbCdEfGhIjKl...`
|
||||||
|
|
||||||
|
#### 2. Save in Koin Ping
|
||||||
|
|
||||||
|
1. Go to the **Alerts** page in Koin Ping.
|
||||||
|
2. In the **Notification Settings** panel, find the **Discord** section.
|
||||||
|
3. Paste your **Webhook URL** into the field.
|
||||||
|
4. Click **Save Settings**.
|
||||||
|
5. Click **Test All Channels** to verify — you should see a test message appear in your Discord channel.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### Slack
|
||||||
|
|
||||||
|
To receive alerts in a Slack channel, you need to create an Incoming Webhook.
|
||||||
|
|
||||||
|
#### 1. Create a Slack App with Incoming Webhooks
|
||||||
|
|
||||||
|
1. Go to [https://api.slack.com/apps](https://api.slack.com/apps) and click **Create New App**.
|
||||||
|
2. Choose **From scratch**, give it a name (e.g. `Koin Ping`), and select your workspace.
|
||||||
|
3. In the app settings, go to **Incoming Webhooks** in the left sidebar.
|
||||||
|
4. Toggle **Activate Incoming Webhooks** to **On**.
|
||||||
|
5. Click **Add New Webhook to Workspace**.
|
||||||
|
6. Select the channel where you want alerts posted (e.g. `#koin-alerts`) and click **Allow**.
|
||||||
|
7. Copy the **Webhook URL** — it will start with `https://hooks.slack.com/services/...`
|
||||||
|
|
||||||
|
#### 2. Save in Koin Ping
|
||||||
|
|
||||||
|
1. Go to the **Alerts** page in Koin Ping.
|
||||||
|
2. In the **Notification Settings** panel, find the **Slack** section.
|
||||||
|
3. Paste your **Webhook URL** into the field.
|
||||||
|
4. Click **Save Settings**.
|
||||||
|
5. Click **Test All Channels** to verify — you should see a test message appear in your Slack channel.
|
||||||
|
|
||||||
|
---
|
||||||
|
|
||||||
|
### Email
|
||||||
|
|
||||||
|
Koin Ping sends email alerts and digests via [Resend](https://resend.com). To enable email notifications, the application owner must configure a Resend account and verified sending domain.
|
||||||
|
|
||||||
|
#### For Application Owners (Deployment Setup)
|
||||||
|
|
||||||
|
1. Create an account at [resend.com](https://resend.com).
|
||||||
|
2. Go to [resend.com/domains](https://resend.com/domains) and add your sending domain.
|
||||||
|
3. Add the DNS records Resend provides (SPF, DKIM, etc.) to your domain's DNS settings.
|
||||||
|
4. Once verified, create an API key at [resend.com/api-keys](https://resend.com/api-keys).
|
||||||
|
5. Set the following in your `.env` file:
|
||||||
|
|
||||||
|
```
|
||||||
|
RESEND_API_KEY=re_your_api_key_here
|
||||||
|
EMAIL_FROM=Your App <alerts@yourdomain.com>
|
||||||
|
```
|
||||||
|
|
||||||
|
6. Restart the API server and poller to pick up the changes.
|
||||||
|
|
||||||
|
#### For Users
|
||||||
|
|
||||||
|
1. Go to the **Alerts** page in Koin Ping.
|
||||||
|
2. In the **Notification Settings** panel, find the **Email** section.
|
||||||
|
3. Enter the email address where you want to receive alerts.
|
||||||
|
4. Click **Save Settings**.
|
||||||
|
5. Click **Verify Email** to receive a confirmation message.
|
||||||
|
6. Click **Test All Channels** to verify — you should receive a test alert at your email address.
|
||||||
|
|
||||||
## License
|
## License
|
||||||
|
|
||||||
MIT. See [LICENSE](LICENSE).
|
MIT. See [LICENSE](LICENSE).
|
||||||
|
|||||||
BIN
backend-go/api
Executable file
BIN
backend-go/api
Executable file
Binary file not shown.
@@ -47,6 +47,7 @@ func main() {
|
|||||||
addressModel := models.NewAddressModel(pool)
|
addressModel := models.NewAddressModel(pool)
|
||||||
alertRuleModel := models.NewAlertRuleModel(pool)
|
alertRuleModel := models.NewAlertRuleModel(pool)
|
||||||
alertEventModel := models.NewAlertEventModel(pool)
|
alertEventModel := models.NewAlertEventModel(pool)
|
||||||
|
checkpointModel := models.NewCheckpointModel(pool)
|
||||||
notifConfigModel := models.NewNotificationConfigModel(pool)
|
notifConfigModel := models.NewNotificationConfigModel(pool)
|
||||||
|
|
||||||
emailDigestSvc := services.NewEmailDigestService(
|
emailDigestSvc := services.NewEmailDigestService(
|
||||||
@@ -58,13 +59,14 @@ func main() {
|
|||||||
alertEventHandler := handlers.NewAlertEventHandler(alertEventModel)
|
alertEventHandler := handlers.NewAlertEventHandler(alertEventModel)
|
||||||
notifConfigHandler := handlers.NewNotificationConfigHandler(notifConfigModel, cfg)
|
notifConfigHandler := handlers.NewNotificationConfigHandler(notifConfigModel, cfg)
|
||||||
emailDigestHandler := handlers.NewEmailDigestHandler(emailDigestSvc, notifConfigModel)
|
emailDigestHandler := handlers.NewEmailDigestHandler(emailDigestSvc, notifConfigModel)
|
||||||
|
statusHandler := handlers.NewStatusHandler(checkpointModel)
|
||||||
|
|
||||||
mux := http.NewServeMux()
|
mux := http.NewServeMux()
|
||||||
b := cfg.APIBasePath // e.g. "/v1"
|
b := cfg.APIBasePath // e.g. "/v1"
|
||||||
|
|
||||||
// Public routes
|
// Public routes
|
||||||
mux.HandleFunc("GET "+b+"/health", handlers.HealthCheck)
|
mux.HandleFunc("GET "+b+"/health", handlers.HealthCheck)
|
||||||
mux.HandleFunc("GET "+b+"/status", handlers.SystemStatus)
|
mux.HandleFunc("GET "+b+"/status", statusHandler.GetStatus)
|
||||||
|
|
||||||
// Authenticated routes — addresses
|
// Authenticated routes — addresses
|
||||||
mux.Handle("POST "+b+"/addresses",
|
mux.Handle("POST "+b+"/addresses",
|
||||||
@@ -73,6 +75,8 @@ func main() {
|
|||||||
middleware.Authenticate(http.HandlerFunc(addressHandler.List)))
|
middleware.Authenticate(http.HandlerFunc(addressHandler.List)))
|
||||||
mux.Handle("DELETE "+b+"/addresses/{addressId}",
|
mux.Handle("DELETE "+b+"/addresses/{addressId}",
|
||||||
middleware.Authenticate(http.HandlerFunc(addressHandler.Remove)))
|
middleware.Authenticate(http.HandlerFunc(addressHandler.Remove)))
|
||||||
|
mux.Handle("PATCH "+b+"/addresses/{addressId}",
|
||||||
|
middleware.Authenticate(http.HandlerFunc(addressHandler.UpdateLabel)))
|
||||||
|
|
||||||
// Authenticated routes for alert rules
|
// Authenticated routes for alert rules
|
||||||
mux.Handle("POST "+b+"/addresses/{addressId}/alerts",
|
mux.Handle("POST "+b+"/addresses/{addressId}/alerts",
|
||||||
|
|||||||
@@ -62,6 +62,7 @@ func main() {
|
|||||||
eth, alertRuleModel, alertEventModel, addressModel, notifConfigModel,
|
eth, alertRuleModel, alertEventModel, addressModel, notifConfigModel,
|
||||||
cfg.ResendAPIKey, cfg.EmailFrom,
|
cfg.ResendAPIKey, cfg.EmailFrom,
|
||||||
)
|
)
|
||||||
|
digestSvc := services.NewEmailDigestService(cfg.ResendAPIKey, cfg.EmailFrom, alertEventModel, notifConfigModel)
|
||||||
|
|
||||||
ctx, cancel := context.WithCancel(context.Background())
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
defer cancel()
|
defer cancel()
|
||||||
@@ -78,12 +79,14 @@ func main() {
|
|||||||
}()
|
}()
|
||||||
|
|
||||||
interval := time.Duration(cfg.PollIntervalMS) * time.Millisecond
|
interval := time.Duration(cfg.PollIntervalMS) * time.Millisecond
|
||||||
|
digestInterval := time.Duration(cfg.DigestIntervalHours) * time.Hour
|
||||||
|
|
||||||
log.Println(strings.Repeat("=", separatorWidth))
|
log.Println(strings.Repeat("=", separatorWidth))
|
||||||
log.Println("Koin Ping Observer Poller Starting")
|
log.Println("Koin Ping Observer Poller Starting")
|
||||||
log.Println(strings.Repeat("=", separatorWidth))
|
log.Println(strings.Repeat("=", separatorWidth))
|
||||||
log.Printf("RPC URL: %s", cfg.EthRPCURL)
|
log.Printf("RPC URL: %s", cfg.EthRPCURL)
|
||||||
log.Printf("Poll Interval: %dms (%ds)", cfg.PollIntervalMS, cfg.PollIntervalMS/msPerSecond)
|
log.Printf("Poll Interval: %dms (%ds)", cfg.PollIntervalMS, cfg.PollIntervalMS/msPerSecond)
|
||||||
|
log.Printf("Digest Interval: %dh", cfg.DigestIntervalHours)
|
||||||
log.Println(strings.Repeat("=", separatorWidth))
|
log.Println(strings.Repeat("=", separatorWidth))
|
||||||
|
|
||||||
runCycle(ctx, observer, evaluator)
|
runCycle(ctx, observer, evaluator)
|
||||||
@@ -91,6 +94,9 @@ func main() {
|
|||||||
ticker := time.NewTicker(interval)
|
ticker := time.NewTicker(interval)
|
||||||
defer ticker.Stop()
|
defer ticker.Stop()
|
||||||
|
|
||||||
|
digestTicker := time.NewTicker(digestInterval)
|
||||||
|
defer digestTicker.Stop()
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
@@ -99,6 +105,13 @@ func main() {
|
|||||||
return
|
return
|
||||||
case <-ticker.C:
|
case <-ticker.C:
|
||||||
runCycle(ctx, observer, evaluator)
|
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)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
3
backend-go/infra/migrations/004_alert_event_dedup.sql
Normal file
3
backend-go/infra/migrations/004_alert_event_dedup.sql
Normal 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;
|
||||||
@@ -13,6 +13,7 @@ const (
|
|||||||
defaultDBPort = 5432
|
defaultDBPort = 5432
|
||||||
defaultPollIntervalMS = 60000
|
defaultPollIntervalMS = 60000
|
||||||
minPollIntervalMS = 1000
|
minPollIntervalMS = 1000
|
||||||
|
defaultDigestIntervalHours = 24
|
||||||
)
|
)
|
||||||
|
|
||||||
type Config struct {
|
type Config struct {
|
||||||
@@ -30,6 +31,7 @@ type Config struct {
|
|||||||
NodeEnv string
|
NodeEnv string
|
||||||
ResendAPIKey string
|
ResendAPIKey string
|
||||||
EmailFrom string
|
EmailFrom string
|
||||||
|
DigestIntervalHours int
|
||||||
}
|
}
|
||||||
|
|
||||||
// Load reads configuration from environment variables and returns a Config.
|
// Load reads configuration from environment variables and returns a Config.
|
||||||
@@ -49,6 +51,7 @@ func Load() (*Config, error) {
|
|||||||
NodeEnv: getEnv("NODE_ENV", "development"),
|
NodeEnv: getEnv("NODE_ENV", "development"),
|
||||||
ResendAPIKey: os.Getenv("RESEND_API_KEY"),
|
ResendAPIKey: os.Getenv("RESEND_API_KEY"),
|
||||||
EmailFrom: getEnv("EMAIL_FROM", "Koin Ping <alerts@koinping.com>"),
|
EmailFrom: getEnv("EMAIL_FROM", "Koin Ping <alerts@koinping.com>"),
|
||||||
|
DigestIntervalHours: getEnvInt("DIGEST_INTERVAL_HOURS", defaultDigestIntervalHours),
|
||||||
}
|
}
|
||||||
|
|
||||||
if cfg.PollIntervalMS < minPollIntervalMS {
|
if cfg.PollIntervalMS < minPollIntervalMS {
|
||||||
|
|||||||
@@ -92,6 +92,44 @@ func (h *AddressHandler) List(w http.ResponseWriter, r *http.Request) {
|
|||||||
writeJSON(w, http.StatusOK, addresses)
|
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.
|
// Remove handles DELETE requests to remove a tracked address.
|
||||||
func (h *AddressHandler) Remove(w http.ResponseWriter, r *http.Request) {
|
func (h *AddressHandler) Remove(w http.ResponseWriter, r *http.Request) {
|
||||||
userID := middleware.GetUserID(r.Context())
|
userID := middleware.GetUserID(r.Context())
|
||||||
|
|||||||
@@ -4,7 +4,6 @@ import (
|
|||||||
"log"
|
"log"
|
||||||
"net/http"
|
"net/http"
|
||||||
"strconv"
|
"strconv"
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"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/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))
|
log.Printf("Found %d alert events for user", len(events))
|
||||||
|
|
||||||
// MVP scaffolding: return mock data if DB is empty
|
if events == nil {
|
||||||
if len(events) == 0 {
|
events = []domain.AlertEvent{}
|
||||||
events = mockEvents(limit)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
writeJSON(w, http.StatusOK, events)
|
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
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -1,10 +1,67 @@
|
|||||||
package handlers
|
package handlers
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"log"
|
||||||
"net/http"
|
"net/http"
|
||||||
"time"
|
"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) {
|
func HealthCheck(w http.ResponseWriter, r *http.Request) {
|
||||||
writeJSON(w, http.StatusOK, map[string]interface{}{
|
writeJSON(w, http.StatusOK, map[string]interface{}{
|
||||||
"status": "ok",
|
"status": "ok",
|
||||||
@@ -12,12 +69,3 @@ func HealthCheck(w http.ResponseWriter, r *http.Request) {
|
|||||||
"service": "koin-ping-backend",
|
"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),
|
|
||||||
})
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -107,6 +107,24 @@ func (m *AddressModel) FindByID(ctx context.Context, id int, userID *string) (*d
|
|||||||
return &a, nil
|
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) {
|
func (m *AddressModel) Remove(ctx context.Context, id int, userID string) (bool, error) {
|
||||||
tag, err := m.pool.Exec(ctx,
|
tag, err := m.pool.Exec(ctx,
|
||||||
`DELETE FROM addresses WHERE id = $1 AND user_id = $2`,
|
`DELETE FROM addresses WHERE id = $1 AND user_id = $2`,
|
||||||
|
|||||||
@@ -2,7 +2,9 @@ package models
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
"errors"
|
||||||
|
|
||||||
|
"github.com/jackc/pgx/v5"
|
||||||
"github.com/jackc/pgx/v5/pgxpool"
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"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,
|
err := m.pool.QueryRow(ctx,
|
||||||
`INSERT INTO alert_events (alert_rule_id, message, address_label, tx_hash)
|
`INSERT INTO alert_events (alert_rule_id, message, address_label, tx_hash)
|
||||||
VALUES ($1, $2, $3, $4)
|
VALUES ($1, $2, $3, $4)
|
||||||
|
ON CONFLICT DO NOTHING
|
||||||
RETURNING id, alert_rule_id, message, address_label, tx_hash, timestamp`,
|
RETURNING id, alert_rule_id, message, address_label, tx_hash, timestamp`,
|
||||||
alertRuleID, message, addressLabel, txHash,
|
alertRuleID, message, addressLabel, txHash,
|
||||||
).Scan(&e.ID, &e.AlertRuleID, &e.Message, &e.AddressLabel, &e.TxHash, &e.Timestamp)
|
).Scan(&e.ID, &e.AlertRuleID, &e.Message, &e.AddressLabel, &e.TxHash, &e.Timestamp)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
if errors.Is(err, pgx.ErrNoRows) {
|
||||||
|
// Duplicate silently skipped by ON CONFLICT DO NOTHING
|
||||||
|
return nil, nil
|
||||||
|
}
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return &e, nil
|
return &e, nil
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package models
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"github.com/jackc/pgx/v5"
|
||||||
"github.com/jackc/pgx/v5/pgxpool"
|
"github.com/jackc/pgx/v5/pgxpool"
|
||||||
@@ -17,6 +18,20 @@ func NewCheckpointModel(pool *pgxpool.Pool) *CheckpointModel {
|
|||||||
return &CheckpointModel{pool: pool}
|
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.
|
// 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) {
|
func (m *CheckpointModel) GetLastCheckedBlock(ctx context.Context, addressID int) (int, bool, error) {
|
||||||
var block int
|
var block int
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package notifications
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
@@ -25,11 +26,15 @@ var discordHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
|||||||
Timeout: discordHTTPTimeoutSeconds * time.Second,
|
Timeout: discordHTTPTimeoutSeconds * time.Second,
|
||||||
}
|
}
|
||||||
|
|
||||||
type AlertMetadata struct {
|
// DiscordNotifier sends alert notifications via a Discord webhook.
|
||||||
TxHash string
|
type DiscordNotifier struct {
|
||||||
AddressLabel string
|
WebhookURL string
|
||||||
AlertType string
|
}
|
||||||
Address 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 {
|
type discordEmbed struct {
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package notifications
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
@@ -9,6 +10,19 @@ import (
|
|||||||
"time"
|
"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
|
const emailHTTPTimeoutSeconds = 10
|
||||||
|
|
||||||
var emailHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
var emailHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
||||||
|
|||||||
16
backend-go/internal/notifications/notifier.go
Normal file
16
backend-go/internal/notifications/notifier.go
Normal 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
|
||||||
|
}
|
||||||
@@ -2,6 +2,7 @@ package notifications
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
@@ -9,6 +10,17 @@ import (
|
|||||||
"time"
|
"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
|
const slackHTTPTimeoutSeconds = 10
|
||||||
|
|
||||||
var slackHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
var slackHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
||||||
|
|||||||
@@ -2,6 +2,7 @@ package notifications
|
|||||||
|
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
@@ -9,6 +10,18 @@ import (
|
|||||||
"time"
|
"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
|
const telegramHTTPTimeoutSeconds = 10
|
||||||
|
|
||||||
var telegramHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
var telegramHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
||||||
|
|||||||
@@ -4,7 +4,9 @@ import (
|
|||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"log"
|
||||||
"math/big"
|
"math/big"
|
||||||
"net/http"
|
"net/http"
|
||||||
"strings"
|
"strings"
|
||||||
@@ -13,7 +15,11 @@ import (
|
|||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
const rpcTimeoutMS = 30000
|
const (
|
||||||
|
rpcTimeoutMS = 30000
|
||||||
|
rpcMaxRetries = 3
|
||||||
|
rpcRetryBaseMS = 1000
|
||||||
|
)
|
||||||
|
|
||||||
type JsonRpcEthereum struct {
|
type JsonRpcEthereum struct {
|
||||||
rpcURL string
|
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 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))
|
req, err := http.NewRequestWithContext(ctx, http.MethodPost, j.rpcURL, bytes.NewReader(body))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("create RPC request: %w", err)
|
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()
|
defer resp.Body.Close()
|
||||||
|
|
||||||
|
// 429 and 5xx are transient; other non-200 are permanent.
|
||||||
if resp.StatusCode != http.StatusOK {
|
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
|
var rpcResp rpcResponse
|
||||||
@@ -88,12 +137,25 @@ func (j *JsonRpcEthereum) callRPC(ctx context.Context, method string, params ...
|
|||||||
}
|
}
|
||||||
|
|
||||||
if rpcResp.Error != nil {
|
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
|
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) {
|
func (j *JsonRpcEthereum) GetLatestBlockNumber(ctx context.Context) (int, error) {
|
||||||
result, err := j.callRPC(ctx, "eth_blockNumber")
|
result, err := j.callRPC(ctx, "eth_blockNumber")
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
||||||
@@ -12,6 +13,12 @@ import (
|
|||||||
"github.com/kjannette/koin-ping/backend-go/internal/wei"
|
"github.com/kjannette/koin-ping/backend-go/internal/wei"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
notificationTimeout = 30 * time.Second
|
||||||
|
notificationMaxRetries = 3
|
||||||
|
notificationRetryBase = time.Second
|
||||||
|
)
|
||||||
|
|
||||||
type EvaluatorService struct {
|
type EvaluatorService struct {
|
||||||
eth ethereum.EthereumObserver
|
eth ethereum.EthereumObserver
|
||||||
alertRules *models.AlertRuleModel
|
alertRules *models.AlertRuleModel
|
||||||
@@ -167,25 +174,85 @@ func (s *EvaluatorService) fireAlert(ctx context.Context, rule domain.AlertRule,
|
|||||||
message := s.buildMessage(rule, obs)
|
message := s.buildMessage(rule, obs)
|
||||||
txHash := &obs.Hash
|
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 {
|
if err != nil {
|
||||||
return err
|
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)
|
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 {
|
if addr != nil {
|
||||||
|
userID := addr.UserID
|
||||||
|
address := addr.Address
|
||||||
go func() {
|
go func() {
|
||||||
s.sendNotification(
|
notifCtx, cancel := context.WithTimeout(context.Background(), notificationTimeout)
|
||||||
ctx, addr.UserID, message, obs, addressLabel, rule, addr.Address,
|
defer cancel()
|
||||||
)
|
s.sendNotification(notifCtx, userID, message, obs, addressLabel, rule, address)
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
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) {
|
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)
|
notifConfig, err := s.notifConfigs.GetConfig(ctx, userID)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -204,44 +271,11 @@ func (s *EvaluatorService) sendNotification(ctx context.Context, userID, message
|
|||||||
Address: address,
|
Address: address,
|
||||||
}
|
}
|
||||||
|
|
||||||
if notifConfig.DiscordWebhookURL != nil && *notifConfig.DiscordWebhookURL != "" {
|
for _, n := range s.buildNotifiers(notifConfig) {
|
||||||
sent, sendErr := notifications.SendDiscordNotification(*notifConfig.DiscordWebhookURL, message, meta)
|
if err := sendWithRetry(ctx, n, message, meta); err != nil {
|
||||||
if sendErr != nil || !sent {
|
log.Printf("Notification channel failed for user %s after retries: %v", userID, err)
|
||||||
log.Printf("Discord notification failed for user %s: %v", userID, sendErr)
|
|
||||||
} else {
|
} else {
|
||||||
log.Printf("Discord notification sent to user %s", userID)
|
log.Printf("Notification sent to user %s via %T", userID, n)
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
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)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -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
|
* Delete a tracked address
|
||||||
* @param {number} addressId - Address ID to delete
|
* @param {number} addressId - Address ID to delete
|
||||||
|
|||||||
@@ -1,11 +1,13 @@
|
|||||||
import { useState, useEffect } from "react";
|
import { useState, useEffect } from "react";
|
||||||
import AddressForm from "../components/AddressForm";
|
import AddressForm from "../components/AddressForm";
|
||||||
import { getAddresses, createAddress } from "../api/addresses";
|
import { getAddresses, createAddress, deleteAddress, updateAddress } from "../api/addresses";
|
||||||
|
|
||||||
export default function Addresses() {
|
export default function Addresses() {
|
||||||
const [addresses, setAddresses] = useState([]);
|
const [addresses, setAddresses] = useState([]);
|
||||||
const [loading, setLoading] = useState(true);
|
const [loading, setLoading] = useState(true);
|
||||||
const [error, setError] = useState(null);
|
const [error, setError] = useState(null);
|
||||||
|
const [editingId, setEditingId] = useState(null);
|
||||||
|
const [editLabel, setEditLabel] = useState("");
|
||||||
|
|
||||||
// Load addresses on mount
|
// Load addresses on mount
|
||||||
useEffect(() => {
|
useEffect(() => {
|
||||||
@@ -29,15 +31,52 @@ export default function Addresses() {
|
|||||||
async function handleAddressSubmit(data) {
|
async function handleAddressSubmit(data) {
|
||||||
try {
|
try {
|
||||||
const newAddress = await createAddress(data);
|
const newAddress = await createAddress(data);
|
||||||
// Append new address to state
|
|
||||||
setAddresses((prev) => [...prev, newAddress]);
|
setAddresses((prev) => [...prev, newAddress]);
|
||||||
setError(null); // Clear any previous errors
|
setError(null);
|
||||||
} catch (err) {
|
} catch (err) {
|
||||||
setError(err.message);
|
setError(err.message);
|
||||||
console.error("Failed to create address:", err);
|
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 (
|
return (
|
||||||
<div style={{ maxWidth: "800px", margin: "0 auto", padding: "2rem" }}>
|
<div style={{ maxWidth: "800px", margin: "0 auto", padding: "2rem" }}>
|
||||||
<h1>Tracked Addresses</h1>
|
<h1>Tracked Addresses</h1>
|
||||||
@@ -68,14 +107,62 @@ export default function Addresses() {
|
|||||||
backgroundColor: "#333",
|
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={{
|
style={{
|
||||||
fontWeight: "bold",
|
background: "#444",
|
||||||
marginBottom: "0.25rem",
|
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>
|
||||||
|
)}
|
||||||
<div
|
<div
|
||||||
style={{
|
style={{
|
||||||
fontFamily: "monospace",
|
fontFamily: "monospace",
|
||||||
@@ -85,6 +172,24 @@ export default function Addresses() {
|
|||||||
>
|
>
|
||||||
{addr.address}
|
{addr.address}
|
||||||
</div>
|
</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>
|
</li>
|
||||||
))}
|
))}
|
||||||
</ul>
|
</ul>
|
||||||
|
|||||||
Reference in New Issue
Block a user