Compare commits
7 Commits
stripe-2
...
login-styl
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
068ca2f834 | ||
|
|
7f10bcb7de | ||
|
|
f91ca33752 | ||
|
|
2c86bba235 | ||
|
|
44dad43f1d | ||
|
|
2fe0e5b8e9 | ||
|
|
8f08105246 |
2
.gitignore
vendored
2
.gitignore
vendored
@@ -21,7 +21,7 @@ node_modules/
|
|||||||
*.key
|
*.key
|
||||||
|
|
||||||
# Go build artifacts
|
# Go build artifacts
|
||||||
backend-go/bin/
|
backend/bin/
|
||||||
*.exe
|
*.exe
|
||||||
*.exe~
|
*.exe~
|
||||||
*.dll
|
*.dll
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
node_modules/
|
node_modules/
|
||||||
frontend/dist/
|
frontend/dist/
|
||||||
frontend/build/
|
frontend/build/
|
||||||
backend-go/bin/
|
backend/bin/
|
||||||
*.lock
|
*.lock
|
||||||
Prompts/
|
Prompts/
|
||||||
.claude/
|
.claude/
|
||||||
|
|||||||
2
Makefile
2
Makefile
@@ -2,7 +2,7 @@
|
|||||||
test-go test-js lint-go lint-js fmt-go fmt-js fmt-check-go fmt-check-js \
|
test-go test-js lint-go lint-js fmt-go fmt-js fmt-check-go fmt-check-js \
|
||||||
build-go build-js
|
build-go build-js
|
||||||
|
|
||||||
GODIR := backend-go
|
GODIR := backend
|
||||||
JSDIR := frontend
|
JSDIR := frontend
|
||||||
PRETTIER := $(JSDIR)/node_modules/.bin/prettier
|
PRETTIER := $(JSDIR)/node_modules/.bin/prettier
|
||||||
|
|
||||||
|
|||||||
@@ -28,8 +28,8 @@ make hooks
|
|||||||
cd frontend && npm install && cd ..
|
cd frontend && npm install && cd ..
|
||||||
|
|
||||||
# Copy and fill in environment variables
|
# Copy and fill in environment variables
|
||||||
cp backend-go/.env.example backend-go/.env
|
cp backend/.env.example backend/.env
|
||||||
# edit backend-go/.env with your DATABASE_URL, FIREBASE_PROJECT_ID, ETH_RPC_URL
|
# edit backend/.env with your DATABASE_URL, FIREBASE_PROJECT_ID, ETH_RPC_URL
|
||||||
|
|
||||||
# Run checks (requires golangci-lint)
|
# Run checks (requires golangci-lint)
|
||||||
make check
|
make check
|
||||||
@@ -38,7 +38,7 @@ make check
|
|||||||
make run
|
make run
|
||||||
|
|
||||||
# Start the poller (separate terminal)
|
# Start the poller (separate terminal)
|
||||||
cd backend-go && go run ./cmd/poller
|
cd backend && go run ./cmd/poller
|
||||||
|
|
||||||
# Start the frontend dev server (separate terminal)
|
# Start the frontend dev server (separate terminal)
|
||||||
cd frontend && npm run dev
|
cd frontend && npm run dev
|
||||||
@@ -63,7 +63,7 @@ frontend:
|
|||||||
|
|
||||||
```
|
```
|
||||||
koin_ping_0.2.0/
|
koin_ping_0.2.0/
|
||||||
├── backend-go/ # Go monorepo root
|
├── backend/ # Go monorepo root
|
||||||
│ ├── cmd/api/ # HTTP REST API server
|
│ ├── cmd/api/ # HTTP REST API server
|
||||||
│ ├── cmd/poller/ # Blockchain polling daemon
|
│ ├── cmd/poller/ # Blockchain polling daemon
|
||||||
│ └── internal/
|
│ └── internal/
|
||||||
|
|||||||
@@ -1,20 +0,0 @@
|
|||||||
# Server
|
|
||||||
PORT=3001
|
|
||||||
API_BASE_PATH=/v1
|
|
||||||
NODE_ENV=development
|
|
||||||
|
|
||||||
# Database
|
|
||||||
DATABASE_URL=postgresql://user:password@localhost:5432/koin_ping
|
|
||||||
|
|
||||||
# Ethereum JSON-RPC
|
|
||||||
ETH_RPC_URL=https://mainnet.infura.io/v3/YOUR-PROJECT-ID
|
|
||||||
|
|
||||||
# Polling interval (ms, minimum 1000)
|
|
||||||
POLL_INTERVAL_MS=60000
|
|
||||||
|
|
||||||
# Firebase
|
|
||||||
FIREBASE_PROJECT_ID=koin-ping
|
|
||||||
|
|
||||||
# Email notifications (Resend — https://resend.com)
|
|
||||||
# RESEND_API_KEY=re_xxxxxxxxxxxx
|
|
||||||
# EMAIL_FROM=Koin Ping <alerts@yourdomain.com>
|
|
||||||
@@ -1,16 +0,0 @@
|
|||||||
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,16 +2,16 @@ Start DB:
|
|||||||
|
|
||||||
brew services start postgresql@15
|
brew services start postgresql@15
|
||||||
|
|
||||||
From the backend-go directory, you have a few options:
|
From the backend directory, you have a few options:
|
||||||
|
|
||||||
Option 1: Single command (both API + poller)
|
Option 1: Single command (both API + poller)
|
||||||
cd /Users/kjannette/workspace/koin_ping_0.2.0/backend-gomake dev-all
|
cd /Users/kjannette/workspace/koin_ping_0.2.0/backendmake dev-all
|
||||||
|
|
||||||
Option 2: Two separate terminals
|
Option 2: Two separate terminals
|
||||||
Terminal 1 (API server):
|
Terminal 1 (API server):
|
||||||
cd /Users/kjannette/workspace/koin_ping_0.2.0/backend-go go run ./cmd/api
|
cd /Users/kjannette/workspace/koin_ping_0.2.0/backend go run ./cmd/api
|
||||||
Terminal 2 (Poller):
|
Terminal 2 (Poller):
|
||||||
cd /Users/kjannette/workspace/koin_ping_0.2.0/backend-go go run ./cmd/poller
|
cd /Users/kjannette/workspace/koin_ping_0.2.0/backend go run ./cmd/poller
|
||||||
|
|
||||||
make run — Builds and runs the API server.
|
make run — Builds and runs the API server.
|
||||||
make dev — Runs the API server with auto-reload via air (falls back to go run if air isn't installed).
|
make dev — Runs the API server with auto-reload via air (falls back to go run if air isn't installed).
|
||||||
@@ -8,13 +8,13 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/joho/godotenv"
|
"github.com/joho/godotenv"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/config"
|
"github.com/kjannette/koin-ping/backend/internal/config"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/database"
|
"github.com/kjannette/koin-ping/backend/internal/database"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/firebase"
|
"github.com/kjannette/koin-ping/backend/internal/firebase"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/handlers"
|
"github.com/kjannette/koin-ping/backend/internal/handlers"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
|
"github.com/kjannette/koin-ping/backend/internal/middleware"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/services"
|
"github.com/kjannette/koin-ping/backend/internal/services"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -12,11 +12,11 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/joho/godotenv"
|
"github.com/joho/godotenv"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/config"
|
"github.com/kjannette/koin-ping/backend/internal/config"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/database"
|
"github.com/kjannette/koin-ping/backend/internal/database"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/protocols/ethereum"
|
"github.com/kjannette/koin-ping/backend/internal/protocols/ethereum"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/services"
|
"github.com/kjannette/koin-ping/backend/internal/services"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -138,6 +138,8 @@ func runCycle(
|
|||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
evaluator.WaitForNotifications()
|
||||||
|
|
||||||
duration := time.Since(startTime)
|
duration := time.Since(startTime)
|
||||||
log.Printf("[%s] Cycle complete: %d observations, %d alerts fired in %s",
|
log.Printf("[%s] Cycle complete: %d observations, %d alerts fired in %s",
|
||||||
time.Now().UTC().Format(time.RFC3339),
|
time.Now().UTC().Format(time.RFC3339),
|
||||||
@@ -1,4 +1,4 @@
|
|||||||
module github.com/kjannette/koin-ping/backend-go
|
module github.com/kjannette/koin-ping/backend
|
||||||
|
|
||||||
go 1.25.0
|
go 1.25.0
|
||||||
|
|
||||||
@@ -47,7 +47,7 @@ var ThresholdRequiredTypes = []AlertType{ //nolint:gochecknoglobals
|
|||||||
AlertBalanceBelow,
|
AlertBalanceBelow,
|
||||||
}
|
}
|
||||||
|
|
||||||
// IsValidAlertType returns true if the given string matches a known AlertType.
|
// returns true if the given string matches a known AlertType.
|
||||||
func IsValidAlertType(t string) bool {
|
func IsValidAlertType(t string) bool {
|
||||||
for _, v := range ValidAlertTypes {
|
for _, v := range ValidAlertTypes {
|
||||||
if string(v) == t {
|
if string(v) == t {
|
||||||
@@ -58,7 +58,6 @@ func IsValidAlertType(t string) bool {
|
|||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
|
|
||||||
// IsThresholdRequired returns true if the given AlertType requires a threshold.
|
|
||||||
func IsThresholdRequired(t AlertType) bool {
|
func IsThresholdRequired(t AlertType) bool {
|
||||||
for _, v := range ThresholdRequiredTypes {
|
for _, v := range ThresholdRequiredTypes {
|
||||||
if v == t {
|
if v == t {
|
||||||
@@ -93,7 +92,6 @@ type AddressCheckpoint struct {
|
|||||||
LastCheckedAt time.Time `json:"last_checked_at"` //nolint:tagliatelle
|
LastCheckedAt time.Time `json:"last_checked_at"` //nolint:tagliatelle
|
||||||
}
|
}
|
||||||
|
|
||||||
// CheckpointDetail combines checkpoint and address info for reporting.
|
|
||||||
type CheckpointDetail struct {
|
type CheckpointDetail struct {
|
||||||
AddressID int `json:"address_id"` //nolint:tagliatelle
|
AddressID int `json:"address_id"` //nolint:tagliatelle
|
||||||
Address string `json:"address"`
|
Address string `json:"address"`
|
||||||
@@ -102,7 +100,7 @@ type CheckpointDetail struct {
|
|||||||
LastCheckedAt time.Time `json:"last_checked_at"` //nolint:tagliatelle
|
LastCheckedAt time.Time `json:"last_checked_at"` //nolint:tagliatelle
|
||||||
}
|
}
|
||||||
|
|
||||||
// NotificationConfig holds a user's notification preferences.
|
// holds a user's notification preferences.
|
||||||
type NotificationConfig struct {
|
type NotificationConfig struct {
|
||||||
UserID string `json:"user_id"` //nolint:tagliatelle
|
UserID string `json:"user_id"` //nolint:tagliatelle
|
||||||
DiscordWebhookURL *string `json:"discord_webhook_url"` //nolint:tagliatelle
|
DiscordWebhookURL *string `json:"discord_webhook_url"` //nolint:tagliatelle
|
||||||
@@ -131,14 +129,12 @@ type NormalizedTx struct {
|
|||||||
TokenValue *string `json:"token_value,omitempty"` //nolint:tagliatelle
|
TokenValue *string `json:"token_value,omitempty"` //nolint:tagliatelle
|
||||||
}
|
}
|
||||||
|
|
||||||
// IsTokenTransfer returns true if this transaction represents an ERC-20 token transfer.
|
|
||||||
func (tx NormalizedTx) IsTokenTransfer() bool {
|
func (tx NormalizedTx) IsTokenTransfer() bool {
|
||||||
return tx.TokenContract != nil
|
return tx.TokenContract != nil
|
||||||
}
|
}
|
||||||
|
|
||||||
type Direction string
|
type Direction string
|
||||||
|
|
||||||
// String implements fmt.Stringer.
|
|
||||||
func (d Direction) String() string { return string(d) }
|
func (d Direction) String() string { return string(d) }
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -8,9 +8,9 @@ import (
|
|||||||
"regexp"
|
"regexp"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend/internal/domain"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
|
"github.com/kjannette/koin-ping/backend/internal/middleware"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
)
|
)
|
||||||
|
|
||||||
var ethAddressRe = regexp.MustCompile(`^0x[a-fA-F0-9]{40}$`)
|
var ethAddressRe = regexp.MustCompile(`^0x[a-fA-F0-9]{40}$`)
|
||||||
@@ -5,9 +5,9 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"strconv"
|
"strconv"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend/internal/domain"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
|
"github.com/kjannette/koin-ping/backend/internal/middleware"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
)
|
)
|
||||||
|
|
||||||
// AlertEventHandler handles HTTP requests for alert event history.
|
// AlertEventHandler handles HTTP requests for alert event history.
|
||||||
@@ -9,9 +9,9 @@ import (
|
|||||||
"strconv"
|
"strconv"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend/internal/domain"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
|
"github.com/kjannette/koin-ping/backend/internal/middleware"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
)
|
)
|
||||||
|
|
||||||
var errThresholdFormat = errors.New("unsupported threshold format")
|
var errThresholdFormat = errors.New("unsupported threshold format")
|
||||||
@@ -4,9 +4,9 @@ import (
|
|||||||
"log"
|
"log"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
|
"github.com/kjannette/koin-ping/backend/internal/middleware"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/services"
|
"github.com/kjannette/koin-ping/backend/internal/services"
|
||||||
)
|
)
|
||||||
|
|
||||||
type EmailDigestHandler struct {
|
type EmailDigestHandler struct {
|
||||||
@@ -7,11 +7,11 @@ import (
|
|||||||
"regexp"
|
"regexp"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/config"
|
"github.com/kjannette/koin-ping/backend/internal/config"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend/internal/domain"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
|
"github.com/kjannette/koin-ping/backend/internal/middleware"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/notifications"
|
"github.com/kjannette/koin-ping/backend/internal/notifications"
|
||||||
)
|
)
|
||||||
|
|
||||||
var emailRe = regexp.MustCompile(`^[^\s@]+@[^\s@]+\.[^\s@]+$`)
|
var emailRe = regexp.MustCompile(`^[^\s@]+@[^\s@]+\.[^\s@]+$`)
|
||||||
@@ -5,7 +5,7 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
)
|
)
|
||||||
|
|
||||||
// StatusHandler handles the system status endpoint.
|
// StatusHandler handles the system status endpoint.
|
||||||
@@ -10,9 +10,9 @@ import (
|
|||||||
checkoutsession "github.com/stripe/stripe-go/v82/checkout/session"
|
checkoutsession "github.com/stripe/stripe-go/v82/checkout/session"
|
||||||
"github.com/stripe/stripe-go/v82/webhook"
|
"github.com/stripe/stripe-go/v82/webhook"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/config"
|
"github.com/kjannette/koin-ping/backend/internal/config"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/middleware"
|
"github.com/kjannette/koin-ping/backend/internal/middleware"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
)
|
)
|
||||||
|
|
||||||
const webhookMaxBodyBytes = 65536
|
const webhookMaxBodyBytes = 65536
|
||||||
@@ -8,8 +8,8 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
fbauth "github.com/kjannette/koin-ping/backend-go/internal/firebase"
|
fbauth "github.com/kjannette/koin-ping/backend/internal/firebase"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
)
|
)
|
||||||
|
|
||||||
type contextKey string
|
type contextKey string
|
||||||
@@ -6,7 +6,7 @@ import (
|
|||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"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/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
type AddressModel struct {
|
type AddressModel struct {
|
||||||
@@ -6,7 +6,7 @@ import (
|
|||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"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/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
type AlertEventModel struct {
|
type AlertEventModel struct {
|
||||||
@@ -6,7 +6,7 @@ import (
|
|||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"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/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
type AlertRuleModel struct {
|
type AlertRuleModel struct {
|
||||||
@@ -7,7 +7,7 @@ import (
|
|||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"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/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
type CheckpointModel struct {
|
type CheckpointModel struct {
|
||||||
@@ -6,7 +6,7 @@ import (
|
|||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"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/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
type NotificationConfigModel struct {
|
type NotificationConfigModel struct {
|
||||||
@@ -6,7 +6,7 @@ import (
|
|||||||
|
|
||||||
"github.com/jackc/pgx/v5"
|
"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/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
type UserModel struct {
|
type UserModel struct {
|
||||||
@@ -21,17 +21,14 @@ const (
|
|||||||
colorBlue = 0x0099ff
|
colorBlue = 0x0099ff
|
||||||
)
|
)
|
||||||
|
|
||||||
// discordHTTPClient is a shared HTTP client with a timeout for Discord requests.
|
|
||||||
var discordHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
var discordHTTPClient = &http.Client{ //nolint:gochecknoglobals
|
||||||
Timeout: discordHTTPTimeoutSeconds * time.Second,
|
Timeout: discordHTTPTimeoutSeconds * time.Second,
|
||||||
}
|
}
|
||||||
|
|
||||||
// DiscordNotifier sends alert notifications via a Discord webhook.
|
// sends alert notifications via a Discord webhook.
|
||||||
type DiscordNotifier struct {
|
type DiscordNotifier struct {
|
||||||
WebhookURL string
|
WebhookURL string
|
||||||
}
|
}
|
||||||
|
|
||||||
// Send implements Notifier for Discord.
|
|
||||||
func (d *DiscordNotifier) Send(_ context.Context, message string, meta AlertMetadata) error {
|
func (d *DiscordNotifier) Send(_ context.Context, message string, meta AlertMetadata) error {
|
||||||
_, err := SendDiscordNotification(d.WebhookURL, message, meta)
|
_, err := SendDiscordNotification(d.WebhookURL, message, meta)
|
||||||
return err
|
return err
|
||||||
@@ -104,8 +101,12 @@ func SendDiscordNotification(webhookURL, message string, meta AlertMetadata) (bo
|
|||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
|
|
||||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||||
|
err := fmt.Errorf("discord webhook failed: HTTP %d", resp.StatusCode)
|
||||||
log.Printf("Discord webhook failed: HTTP %d", resp.StatusCode)
|
log.Printf("Discord webhook failed: HTTP %d", resp.StatusCode)
|
||||||
return false, fmt.Errorf("discord webhook failed: HTTP %d", resp.StatusCode)
|
if isPermanentStatusCode(resp.StatusCode) {
|
||||||
|
return false, &PermanentError{Err: err}
|
||||||
|
}
|
||||||
|
return false, err
|
||||||
}
|
}
|
||||||
|
|
||||||
return true, nil
|
return true, nil
|
||||||
@@ -10,14 +10,12 @@ import (
|
|||||||
"time"
|
"time"
|
||||||
)
|
)
|
||||||
|
|
||||||
// EmailNotifier sends alert notifications via email (Resend).
|
|
||||||
type EmailNotifier struct {
|
type EmailNotifier struct {
|
||||||
APIKey string
|
APIKey string
|
||||||
From string
|
From string
|
||||||
To string
|
To string
|
||||||
}
|
}
|
||||||
|
|
||||||
// Send implements Notifier for email.
|
|
||||||
func (e *EmailNotifier) Send(_ context.Context, message string, meta AlertMetadata) error {
|
func (e *EmailNotifier) Send(_ context.Context, message string, meta AlertMetadata) error {
|
||||||
_, err := SendEmailNotification(e.APIKey, e.From, e.To, message, meta)
|
_, err := SendEmailNotification(e.APIKey, e.From, e.To, message, meta)
|
||||||
return err
|
return err
|
||||||
@@ -99,8 +97,12 @@ func SendEmailNotification(apiKey, fromAddress, toAddress, message string, meta
|
|||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
|
|
||||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||||
|
err := fmt.Errorf("resend API failed: HTTP %d", resp.StatusCode)
|
||||||
log.Printf("Resend API failed: HTTP %d", resp.StatusCode)
|
log.Printf("Resend API failed: HTTP %d", resp.StatusCode)
|
||||||
return false, fmt.Errorf("resend API failed: HTTP %d", resp.StatusCode)
|
if isPermanentStatusCode(resp.StatusCode) {
|
||||||
|
return false, &PermanentError{Err: err}
|
||||||
|
}
|
||||||
|
return false, err
|
||||||
}
|
}
|
||||||
|
|
||||||
return true, nil
|
return true, nil
|
||||||
38
backend/internal/notifications/notifier.go
Normal file
38
backend/internal/notifications/notifier.go
Normal file
@@ -0,0 +1,38 @@
|
|||||||
|
package notifications
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
"net/http"
|
||||||
|
)
|
||||||
|
|
||||||
|
// AlertMetadata holds context about the alert being sent.
|
||||||
|
type AlertMetadata struct {
|
||||||
|
TxHash string
|
||||||
|
AddressLabel string
|
||||||
|
AlertType string
|
||||||
|
Address string
|
||||||
|
}
|
||||||
|
|
||||||
|
type Notifier interface {
|
||||||
|
Send(ctx context.Context, message string, meta AlertMetadata) error
|
||||||
|
}
|
||||||
|
|
||||||
|
// PermanentError wraps errors that should not be retried (e.g. 401, 403, 404).
|
||||||
|
type PermanentError struct{ Err error }
|
||||||
|
|
||||||
|
func (e *PermanentError) Error() string { return e.Err.Error() }
|
||||||
|
func (e *PermanentError) Unwrap() error { return e.Err }
|
||||||
|
|
||||||
|
func IsPermanent(err error) bool {
|
||||||
|
var p *PermanentError
|
||||||
|
return errors.As(err, &p)
|
||||||
|
}
|
||||||
|
|
||||||
|
func isPermanentStatusCode(code int) bool {
|
||||||
|
return code == http.StatusUnauthorized ||
|
||||||
|
code == http.StatusForbidden ||
|
||||||
|
code == http.StatusNotFound ||
|
||||||
|
code == http.StatusMethodNotAllowed ||
|
||||||
|
code == http.StatusGone
|
||||||
|
}
|
||||||
@@ -87,8 +87,12 @@ func SendSlackNotification(webhookURL, message string, meta AlertMetadata) (bool
|
|||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
|
|
||||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||||
|
err := fmt.Errorf("slack webhook failed: HTTP %d", resp.StatusCode)
|
||||||
log.Printf("Slack webhook failed: HTTP %d", resp.StatusCode)
|
log.Printf("Slack webhook failed: HTTP %d", resp.StatusCode)
|
||||||
return false, fmt.Errorf("slack webhook failed: HTTP %d", resp.StatusCode)
|
if isPermanentStatusCode(resp.StatusCode) {
|
||||||
|
return false, &PermanentError{Err: err}
|
||||||
|
}
|
||||||
|
return false, err
|
||||||
}
|
}
|
||||||
|
|
||||||
return true, nil
|
return true, nil
|
||||||
@@ -63,8 +63,12 @@ func SendTelegramNotification(botToken, chatID, message string, meta AlertMetada
|
|||||||
defer resp.Body.Close()
|
defer resp.Body.Close()
|
||||||
|
|
||||||
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
|
||||||
|
err := fmt.Errorf("telegram API failed: HTTP %d", resp.StatusCode)
|
||||||
log.Printf("Telegram API failed: HTTP %d", resp.StatusCode)
|
log.Printf("Telegram API failed: HTTP %d", resp.StatusCode)
|
||||||
return false, fmt.Errorf("telegram API failed: HTTP %d", resp.StatusCode)
|
if isPermanentStatusCode(resp.StatusCode) {
|
||||||
|
return false, &PermanentError{Err: err}
|
||||||
|
}
|
||||||
|
return false, err
|
||||||
}
|
}
|
||||||
|
|
||||||
return true, nil
|
return true, nil
|
||||||
@@ -12,13 +12,13 @@ import (
|
|||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
rpcTimeoutMS = 30000
|
rpcTimeoutMS = 30000
|
||||||
rpcMaxRetries = 3
|
rpcMaxRetries = 3
|
||||||
rpcRetryBaseMS = 1000
|
rpcRetryBaseMS = 2000
|
||||||
)
|
)
|
||||||
|
|
||||||
type JsonRpcEthereum struct {
|
type JsonRpcEthereum struct {
|
||||||
@@ -3,7 +3,7 @@ package ethereum
|
|||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend/internal/domain"
|
||||||
)
|
)
|
||||||
|
|
||||||
// EthereumObserver defines the interface for blockchain interaction.
|
// EthereumObserver defines the interface for blockchain interaction.
|
||||||
@@ -9,7 +9,7 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -20,7 +20,7 @@ const (
|
|||||||
|
|
||||||
var digestHTTPClient = &http.Client{Timeout: emailHTTPTimeout} //nolint:gochecknoglobals
|
var digestHTTPClient = &http.Client{Timeout: emailHTTPTimeout} //nolint:gochecknoglobals
|
||||||
|
|
||||||
// EmailDigestService handles email setup and digest sending via Resend.
|
// handles email setup and digest sending via Resend.
|
||||||
type EmailDigestService struct {
|
type EmailDigestService struct {
|
||||||
apiKey string
|
apiKey string
|
||||||
fromAddress string
|
fromAddress string
|
||||||
@@ -66,7 +66,6 @@ func (s *EmailDigestService) SetupEmail(toAddress string) error {
|
|||||||
return s.send(toAddress, "Koin Ping — Email Alerts Configured", html)
|
return s.send(toAddress, "Koin Ping — Email Alerts Configured", html)
|
||||||
}
|
}
|
||||||
|
|
||||||
// SendDigest compiles recent alert events for a user and sends a digest email.
|
|
||||||
func (s *EmailDigestService) SendDigest(ctx context.Context, userID, toAddress string) error {
|
func (s *EmailDigestService) SendDigest(ctx context.Context, userID, toAddress string) error {
|
||||||
if !s.Configured() {
|
if !s.Configured() {
|
||||||
return fmt.Errorf("email service not configured: RESEND_API_KEY not set") //nolint:err113
|
return fmt.Errorf("email service not configured: RESEND_API_KEY not set") //nolint:err113
|
||||||
@@ -131,8 +130,6 @@ func (s *EmailDigestService) SendDigest(ctx context.Context, userID, toAddress s
|
|||||||
return s.send(toAddress, subject, html)
|
return s.send(toAddress, subject, html)
|
||||||
}
|
}
|
||||||
|
|
||||||
// SendDigestsForAllUsers sends a digest email to every user that has
|
|
||||||
// notifications enabled and an email configured.
|
|
||||||
func (s *EmailDigestService) SendDigestsForAllUsers(ctx context.Context) (int, error) {
|
func (s *EmailDigestService) SendDigestsForAllUsers(ctx context.Context) (int, error) {
|
||||||
if !s.Configured() {
|
if !s.Configured() {
|
||||||
return 0, nil
|
return 0, nil
|
||||||
@@ -4,19 +4,23 @@ import (
|
|||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
|
"sync"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"golang.org/x/sync/semaphore"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/notifications"
|
"github.com/kjannette/koin-ping/backend/internal/domain"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/protocols/ethereum"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/wei"
|
"github.com/kjannette/koin-ping/backend/internal/notifications"
|
||||||
|
"github.com/kjannette/koin-ping/backend/internal/protocols/ethereum"
|
||||||
|
"github.com/kjannette/koin-ping/backend/internal/wei"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
notificationTimeout = 30 * time.Second
|
notificationTimeout = 30 * time.Second
|
||||||
notificationMaxRetries = 3
|
notificationMaxRetries = 3
|
||||||
notificationRetryBase = time.Second
|
notificationRetryBase = time.Second
|
||||||
|
maxConcurrentNotifications = 5
|
||||||
)
|
)
|
||||||
|
|
||||||
type EvaluatorService struct {
|
type EvaluatorService struct {
|
||||||
@@ -27,6 +31,8 @@ type EvaluatorService struct {
|
|||||||
notifConfigs *models.NotificationConfigModel
|
notifConfigs *models.NotificationConfigModel
|
||||||
resendAPIKey string
|
resendAPIKey string
|
||||||
emailFrom string
|
emailFrom string
|
||||||
|
notifSem *semaphore.Weighted
|
||||||
|
notifWg sync.WaitGroup
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewEvaluatorService(
|
func NewEvaluatorService(
|
||||||
@@ -46,6 +52,7 @@ func NewEvaluatorService(
|
|||||||
notifConfigs: notifConfigs,
|
notifConfigs: notifConfigs,
|
||||||
resendAPIKey: resendAPIKey,
|
resendAPIKey: resendAPIKey,
|
||||||
emailFrom: emailFrom,
|
emailFrom: emailFrom,
|
||||||
|
notifSem: semaphore.NewWeighted(maxConcurrentNotifications),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -189,7 +196,14 @@ func (s *EvaluatorService) fireAlert(ctx context.Context, rule domain.AlertRule,
|
|||||||
if addr != nil {
|
if addr != nil {
|
||||||
userID := addr.UserID
|
userID := addr.UserID
|
||||||
address := addr.Address
|
address := addr.Address
|
||||||
|
if err := s.notifSem.Acquire(ctx, 1); err != nil {
|
||||||
|
log.Printf("Failed to acquire notification semaphore for rule %d: %v", rule.ID, err)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
s.notifWg.Add(1)
|
||||||
go func() {
|
go func() {
|
||||||
|
defer s.notifSem.Release(1)
|
||||||
|
defer s.notifWg.Done()
|
||||||
notifCtx, cancel := context.WithTimeout(context.Background(), notificationTimeout)
|
notifCtx, cancel := context.WithTimeout(context.Background(), notificationTimeout)
|
||||||
defer cancel()
|
defer cancel()
|
||||||
s.sendNotification(notifCtx, userID, message, obs, addressLabel, rule, address)
|
s.sendNotification(notifCtx, userID, message, obs, addressLabel, rule, address)
|
||||||
@@ -199,6 +213,11 @@ func (s *EvaluatorService) fireAlert(ctx context.Context, rule domain.AlertRule,
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// WaitForNotifications blocks until all in-flight notification goroutines finish.
|
||||||
|
func (s *EvaluatorService) WaitForNotifications() {
|
||||||
|
s.notifWg.Wait()
|
||||||
|
}
|
||||||
|
|
||||||
func (s *EvaluatorService) buildNotifiers(cfg *domain.NotificationConfig) []notifications.Notifier {
|
func (s *EvaluatorService) buildNotifiers(cfg *domain.NotificationConfig) []notifications.Notifier {
|
||||||
var notifiers []notifications.Notifier
|
var notifiers []notifications.Notifier
|
||||||
|
|
||||||
@@ -242,8 +261,14 @@ func sendWithRetry(ctx context.Context, n notifications.Notifier, message string
|
|||||||
}
|
}
|
||||||
|
|
||||||
if err := n.Send(ctx, message, meta); err != nil {
|
if err := n.Send(ctx, message, meta); err != nil {
|
||||||
log.Printf("Notification attempt %d/%d failed: %v", attempt+1, notificationMaxRetries, err)
|
|
||||||
lastErr = err
|
lastErr = err
|
||||||
|
|
||||||
|
if notifications.IsPermanent(err) {
|
||||||
|
log.Printf("Permanent notification failure, skipping retries: %v", err)
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
|
||||||
|
log.Printf("Notification attempt %d/%d failed: %v", attempt+1, notificationMaxRetries, err)
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -5,9 +5,9 @@ import (
|
|||||||
"log"
|
"log"
|
||||||
"strings"
|
"strings"
|
||||||
|
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/domain"
|
"github.com/kjannette/koin-ping/backend/internal/domain"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/models"
|
"github.com/kjannette/koin-ping/backend/internal/models"
|
||||||
"github.com/kjannette/koin-ping/backend-go/internal/protocols/ethereum"
|
"github.com/kjannette/koin-ping/backend/internal/protocols/ethereum"
|
||||||
)
|
)
|
||||||
|
|
||||||
const maxBlocksPerRun = 100
|
const maxBlocksPerRun = 100
|
||||||
Reference in New Issue
Block a user