package main
import (
"context"
"fmt"
"io"
"net/http"
"os"
"os/signal"
"strings"
"sync"
"syscall"
"time"
json "github.com/bytedance/sonic"
_ "github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/nbd-wtf/go-nostr"
"github.com/redis/go-redis/v9"
)
// --- 1. Modelos e DTOs ---
type NIP11Info struct {
Name string `json:"name"`
Description string `json:"description"`
Pubkey string `json:"pubkey"`
Contact string `json:"contact"`
SupportedNIPs []int `json:"supported_nips"`
Software string `json:"software"`
Version string `json:"version"`
}
type RelayRecord struct {
URL string `json:"url"`
Name string `json:"name"`
Contact string `json:"contact"`
NIPs []int `json:"nips"`
NIP11 json.NoCopyRawMessage `json:"nip11"`
LastAccessed time.Time `json:"last_accessed"`
IsActive bool `json:"is_active"`
ResponseTime int64 `json:"response_time_ms"`
}
// --- 2. Interfaces (SOLID - DIP) ---
type QueueService interface {
Push(ctx context.Context, url string) error
Pop(ctx context.Context) (string, error)
IsVisited(ctx context.Context, url string) (bool, error)
}
type LockService interface {
Acquire(ctx context.Context, key string, ttl time.Duration) (bool, error)
Release(ctx context.Context, key string) error
}
type Store interface {
Save(ctx context.Context, record RelayRecord) error
Migrate(ctx context.Context) error
}
// Nova interface para lidar com os servidores Blossom
type BlossomStore interface {
AddServer(ctx context.Context, url string) error
GetAllServers(ctx context.Context) ([]string, error)
}
// --- 3. Implementação Postgres (pgx) ---
type PostgresStore struct {
pool *pgxpool.Pool
}
func NewPostgresStore(connStr string) (*PostgresStore, error) {
config, err := pgxpool.ParseConfig(connStr)
if err != nil {
return nil, err
}
pool, err := pgxpool.NewWithConfig(context.Background(), config)
return &PostgresStore{pool: pool}, err
}
func (p *PostgresStore) Migrate(ctx context.Context) error {
query := `
CREATE TABLE IF NOT EXISTS relays (
url TEXT PRIMARY KEY,
name TEXT,
contact TEXT,
nip11 JSONB,
nips INT[],
last_accessed TIMESTAMP WITH TIME ZONE,
is_active BOOLEAN,
response_time_ms BIGINT
);
CREATE INDEX IF NOT EXISTS idx_relays_nips ON relays USING GIN (nips);
CREATE INDEX IF NOT EXISTS idx_relays_name ON relays (name);
CREATE INDEX IF NOT EXISTS idx_relays_contact ON relays (contact);
CREATE INDEX IF NOT EXISTS idx_relays_last_accessed ON relays (last_accessed);
CREATE INDEX IF NOT EXISTS idx_relays_active ON relays (is_active);
`
_, err := p.pool.Exec(ctx, query)
return err
}
func (p *PostgresStore) Save(ctx context.Context, r RelayRecord) error {
if len(r.NIP11) == 0 || string(r.NIP11) == "null" {
r.NIP11 = json.NoCopyRawMessage("{}")
}
query := `
INSERT INTO relays (url, name, contact, nip11, nips, last_accessed, is_active, response_time_ms)
VALUES ($1, $2, $3, $4, $5, $6, $7, $8)
ON CONFLICT (url) DO UPDATE SET
name = EXCLUDED.name,
contact = EXCLUDED.contact,
nip11 = EXCLUDED.nip11,
nips = EXCLUDED.nips,
last_accessed = EXCLUDED.last_accessed,
is_active = EXCLUDED.is_active,
response_time_ms = EXCLUDED.response_time_ms;
`
_, err := p.pool.Exec(ctx, query,
r.URL, r.Name, r.Contact, r.NIP11, r.NIPs, r.LastAccessed, r.IsActive, r.ResponseTime,
)
return err
}
// --- 4. Implementação Redis (Queue, Lock & BlossomStore) ---
type RedisService struct {
client *redis.Client
}
func (r *RedisService) Push(ctx context.Context, url string) error {
return r.client.LPush(ctx, "crawler:queue", url).Err()
}
func (r *RedisService) Pop(ctx context.Context) (string, error) {
res, err := r.client.BRPop(ctx, 0, "crawler:queue").Result()
if err != nil {
return "", err
}
return res[1], nil
}
func (r *RedisService) IsVisited(ctx context.Context, url string) (bool, error) {
added, err := r.client.SAdd(ctx, "crawler:visited", url).Result()
return added == 0, err
}
func (r *RedisService) Acquire(ctx context.Context, key string, ttl time.Duration) (bool, error) {
return r.client.SetNX(ctx, "lock:"+key, "1", ttl).Result()
}
func (r *RedisService) Release(ctx context.Context, key string) error {
return r.client.Del(ctx, "lock:"+key).Err()
}
// Métodos do BlossomStore
func (r *RedisService) AddServer(ctx context.Context, url string) error {
return r.client.SAdd(ctx, "blossom:servers", url).Err()
}
func (r *RedisService) GetAllServers(ctx context.Context) ([]string, error) {
return r.client.SMembers(ctx, "blossom:servers").Result()
}
// --- 5. Crawler Core ---
type Crawler struct {
queue QueueService
lock LockService
store Store
blossomStore BlossomStore // Novo storage para servidores Blossom
httpClient *http.Client
maxWorkers int
}
func (c *Crawler) fetchNIP11(ctx context.Context, url string) (json.NoCopyRawMessage, *NIP11Info) {
httpURL := strings.Replace(strings.Replace(url, "wss://", "https://", 1), "ws://", "http://", 1)
req, err := http.NewRequestWithContext(ctx, http.MethodGet, httpURL, nil)
if err != nil {
return nil, nil
}
req.Header.Add("Accept", "application/nostr+json")
resp, err := c.httpClient.Do(req)
if err != nil {
return nil, nil
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
return nil, nil
}
body, err := io.ReadAll(resp.Body)
if err != nil || len(body) == 0 {
return nil, nil
}
if !json.Valid(body) {
return nil, nil
}
var info NIP11Info
if err := json.Unmarshal(body, &info); err != nil {
return body, nil
}
return body, &info
}
func (c *Crawler) processRelay(ctx context.Context, url string) {
locked, _ := c.lock.Acquire(ctx, url, 1*time.Minute)
if !locked {
return
}
defer c.lock.Release(ctx, url)
fmt.Printf("🔍 Processando: %s\n", url)
record := RelayRecord{
URL: url,
LastAccessed: time.Now(),
IsActive: false,
}
start := time.Now()
relay, err := nostr.RelayConnect(ctx, url)
if err == nil {
record.IsActive = true
record.ResponseTime = time.Since(start).Milliseconds()
// Refatorado: Filtro para Relay List (10002) E Blossom Servers (10063)
sub, _ := relay.Subscribe(ctx, nostr.Filters{{Kinds: []int{10002, 10063}, Limit: 50}})
if sub != nil {
go func() {
for evt := range sub.Events {
switch evt.Kind {
case 10002: // Descoberta de novos relays Nostr
for _, tag := range evt.Tags {
if tag[0] == "r" && len(tag) > 1 {
c.Enqueue(ctx, tag[1])
}
}
case 10063: // Servidores Blossom
for _, tag := range evt.Tags {
// Servidores Blossom normalmente usam a tag "server"
if tag[0] == "server" && len(tag) > 1 {
fmt.Printf("🌸 Blossom Server encontrado: %s\n", tag[1])
c.blossomStore.AddServer(ctx, tag[1])
}
}
}
}
}()
time.Sleep(3 * time.Second)
}
relay.Close()
}
raw, info := c.fetchNIP11(ctx, url)
record.NIP11 = raw
if info != nil {
record.Name = info.Name
record.Contact = info.Contact
record.NIPs = info.SupportedNIPs
}
if err := c.store.Save(ctx, record); err != nil {
fmt.Printf("❌ Erro DB [%s]: %v\n", url, err)
}
}
func (c *Crawler) Enqueue(ctx context.Context, url string) {
if !strings.HasPrefix(url, "ws") {
return
}
visited, _ := c.queue.IsVisited(ctx, url)
if !visited {
c.queue.Push(ctx, url)
}
}
func (c *Crawler) Start(ctx context.Context) {
var wg sync.WaitGroup
for i := 0; i < c.maxWorkers; i++ {
wg.Add(1)
go func() {
defer wg.Done()
for {
select {
case <-ctx.Done():
return
default:
url, err := c.queue.Pop(ctx)
if err != nil {
continue
}
c.processRelay(ctx, url)
}
}
}()
}
wg.Wait()
}
// --- 6. Main e Setup ---
func main() {
// Refatorado: Context com cancelamento para suportar Graceful Shutdown
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
redisAddr := "localhost:6379"
pgConnStr := "postgres://postgres:Strong@P4ssword@localhost:5432/nostr_crawler"
rdb := redis.NewClient(&redis.Options{Addr: redisAddr})
redisSvc := &RedisService{client: rdb}
pgStore, err := NewPostgresStore(pgConnStr)
if err != nil {
panic(err)
}
// Descomente caso deseje rodar a migração
// if err := pgStore.Migrate(ctx); err != nil {
// panic(err)
// }
crawler := &Crawler{
queue: redisSvc,
lock: redisSvc,
store: pgStore,
blossomStore: redisSvc, // Injetando o Redis como armazenamento Blossom
maxWorkers: 20,
httpClient: &http.Client{Timeout: 10 * time.Second},
}
seeds := []string{
"wss://relay.damus.io",
"wss://nos.lol",
"wss://relay.snort.social",
}
for _, s := range seeds {
crawler.Enqueue(ctx, s)
}
// Configuração do Graceful Shutdown (Captura Ctrl+C)
sigs := make(chan os.Signal, 1)
signal.Notify(sigs, syscall.SIGINT, syscall.SIGTERM)
go func() {
<-sigs // Bloqueia até receber um sinal de interrupção
fmt.Println("\n🛑 Encerrando crawler. Iniciando exportação dos servidores Blossom...")
servers, err := redisSvc.GetAllServers(context.Background())
if err == nil && len(servers) > 0 {
file, err := os.Create("blossom_servers_export.txt")
if err == nil {
defer file.Close()
for _, server := range servers {
file.WriteString(server + "\n")
}
fmt.Printf("✅ %d servidores Blossom exportados com sucesso para 'blossom_servers_export.txt'!\n", len(servers))
} else {
fmt.Printf("❌ Erro ao criar arquivo de exportação: %v\n", err)
}
} else {
fmt.Println("⚠️ Nenhum servidor Blossom foi encontrado durante a execução.")
}
cancel() // Cancela o context, fazendo os workers pararem
}()
fmt.Println("🚀 Crawler Multi-Instância Iniciado. Pressione Ctrl+C para encerrar e exportar a lista de Blossom servers.")
crawler.Start(ctx)
}
https://nosbin.com/nevent1qqs9sv9nnv7uehgm2wz2ltvk7g00rvk35jt3ahzm8uk9qxd36a0eekqpzemhxue69uhkzarvv9ejumn0wd68ytnvv9hxgqg4waehxw309ajkgetw9ehx7um5wghxcctwvsq3wamnwvaz7tmwdaehgu3wvekhgtnhd9azucnf0gq3gamnwvaz7tmwdaehgu3wdau8gu3wv3jhvqgswaehxw309ahx7um5wgh8w6twv5q3jamnwvaz7tmwdaehgu3w0fjkyetyv4jjucmvda6kgqghwaehxw309aex2mrp0yhxxatjwfjkuapwveukjqg5waehxw309aex2mrp0yhxgctdw4eju6t0qyt8wumn8ghj7un9d3shjtnwdaeky6tw9e3k7mgprfmhxue69uhhyetvv9ujummjv9hxwetsd9kxctnyv4mqzxrhwden5te0wfjkccte9eekummjwsh8xmmrd9skcxnq09l
