Compare commits

...

4 Commits

Author SHA1 Message Date
KS Jannette
2c86bba235 general cleanup - removed old comments, enforced naming conventions etc
Some checks are pending
check / check (push) Waiting to run
2026-03-04 23:29:58 -05:00
S Jannette
44dad43f1d Merge pull request #17 from kjannette/poller-tweaks
updated evaluator service -
2026-03-04 23:16:56 -05:00
KS Jannette
2fe0e5b8e9 updated evaluator service - added semaphore.Weighted(5) and sync.WaitGroup etc to cap conncurrent requests; also WaitForNotifications() so caller blocks until notifications finish. jsonrpc.go -- finally, bumped rpcRetryBaseMS from 1000 to 2000 - RPC retries at 2s/4s/8s backoff rate
Some checks are pending
check / check (push) Waiting to run
2026-03-04 23:14:55 -05:00
S Jannette
8f08105246 Merge pull request #16 from kjannette/stripe-2
Stripe 2
2026-03-04 22:43:49 -05:00
55 changed files with 90 additions and 102 deletions

2
.gitignore vendored
View File

@@ -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

View File

@@ -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/

View File

@@ -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

View File

@@ -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/

View File

@@ -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>

View File

@@ -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).

View File

@@ -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 (

View File

@@ -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),

View File

@@ -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

View File

@@ -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 (

View File

@@ -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}$`)

View File

@@ -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.

View File

@@ -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")

View File

@@ -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 {

View File

@@ -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@]+$`)

View File

@@ -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.

View File

@@ -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

View File

@@ -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

View File

@@ -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 {

View File

@@ -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 {

View File

@@ -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 {

View File

@@ -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 {

View File

@@ -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 {

View File

@@ -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 {

View File

@@ -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

View File

@@ -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

View File

@@ -10,7 +10,6 @@ type AlertMetadata struct {
Address string Address string
} }
// Notifier is the interface implemented by all notification channels.
type Notifier interface { type Notifier interface {
Send(ctx context.Context, message string, meta AlertMetadata) error Send(ctx context.Context, message string, meta AlertMetadata) error
} }

View File

@@ -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 {

View File

@@ -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.

View File

@@ -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

View File

@@ -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

View File

@@ -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