publish.go raw

   1  package main
   2  
   3  import (
   4  	"context"
   5  	"fmt"
   6  	"log"
   7  	"strings"
   8  	"sync"
   9  	"sync/atomic"
  10  	"time"
  11  
  12  	"git.mleku.dev/mleku/dendrite/pkg/grammar"
  13  	"git.mleku.dev/mleku/dendrite/pkg/nostr"
  14  )
  15  
  16  // Broadcaster publishes verdict events to every relay in the pool.
  17  type Broadcaster struct {
  18  	pool *RelayPool
  19  	blob *BlobInfo // may be nil if blossom upload failed
  20  }
  21  
  22  // NewBroadcaster creates a broadcaster backed by the relay pool.
  23  func NewBroadcaster(pool *RelayPool, blob *BlobInfo) *Broadcaster {
  24  	return &Broadcaster{pool: pool, blob: blob}
  25  }
  26  
  27  // PublishVerdict composes a verdict reply, signs it with a fresh throwaway
  28  // keypair, and blasts it to every known relay. The private key is discarded
  29  // after signing. Each verdict has a unique, unmutable identity.
  30  func (b *Broadcaster) PublishVerdict(ctx context.Context, original *nostr.Event, verdict grammar.Verdict) {
  31  	id, err := nostr.NewIdentity()
  32  	if err != nil {
  33  		log.Printf("keygen: %v", err)
  34  		return
  35  	}
  36  
  37  	// Tag the original event as a reply. Include the watch relay as hint.
  38  	tags := [][]string{
  39  		{"e", original.ID, "wss://relay.orly.dev", "reply"},
  40  		{"p", original.Pubkey},
  41  		{"t", "dendrite"},
  42  		{"t", "ai-detection"},
  43  	}
  44  	if verdict.TrollLabel != "" {
  45  		tags = append(tags, []string{"t", "troll-detection"})
  46  	}
  47  	ev := &nostr.Event{
  48  		CreatedAt: time.Now().Unix(),
  49  		Kind:      1,
  50  		Tags:      tags,
  51  		Content:   b.verdictContent(verdict),
  52  	}
  53  
  54  	if err := ev.Sign(id.PrivKeyHex()); err != nil {
  55  		log.Printf("sign verdict: %v", err)
  56  		return
  57  	}
  58  
  59  	urls := b.pool.URLs()
  60  
  61  	// Log the verdict as a nevent URI so operators can look it up directly.
  62  	nevent, err := nostr.Nevent(ev.ID, []string{"wss://relay.orly.dev"}, ev.Pubkey, ev.Kind)
  63  	if err != nil {
  64  		log.Printf("broadcasting verdict %s to %d relays for event %s",
  65  			ev.ID[:12], len(urls), original.ID[:12])
  66  	} else {
  67  		log.Printf("verdict nostr:%s", nevent)
  68  		log.Printf("broadcasting to %d relays for event %s", len(urls), original.ID[:12])
  69  	}
  70  
  71  	// Fan out to all relays concurrently.
  72  	var accepted, rejected, failed atomic.Int32
  73  	var wg sync.WaitGroup
  74  
  75  	for _, url := range urls {
  76  		wg.Add(1)
  77  		go func(url string) {
  78  			defer wg.Done()
  79  			publishOne(ctx, url, ev, &accepted, &rejected, &failed)
  80  		}(url)
  81  	}
  82  
  83  	// Wait for all publishes to complete (or timeout).
  84  	done := make(chan struct{})
  85  	go func() {
  86  		wg.Wait()
  87  		close(done)
  88  	}()
  89  
  90  	select {
  91  	case <-done:
  92  	case <-time.After(30 * time.Second):
  93  	}
  94  
  95  	origNevent, _ := nostr.Nevent(original.ID, []string{"wss://relay.orly.dev"}, original.Pubkey, original.Kind)
  96  	log.Printf("result: accepted=%d rejected=%d failed=%d (of %d relays) re: nostr:%s",
  97  		accepted.Load(), rejected.Load(), failed.Load(), len(urls), origNevent)
  98  }
  99  
 100  // verdictContent composes the full verdict text including source and blob links.
 101  func (b *Broadcaster) verdictContent(v grammar.Verdict) string {
 102  	var sb strings.Builder
 103  	sb.WriteString(v.String())
 104  	if b.blob != nil && b.blob.SHA256 != "" {
 105  		fmt.Fprintf(&sb, "binary sha256: %s\n", b.blob.SHA256)
 106  		if len(b.blob.URLs) > 0 {
 107  			fmt.Fprintf(&sb, "download: %s\n", b.blob.URLs[0])
 108  		}
 109  	}
 110  	return sb.String()
 111  }
 112  
 113  func publishOne(ctx context.Context, url string, ev *nostr.Event, accepted, rejected, failed *atomic.Int32) {
 114  	pubCtx, cancel := context.WithTimeout(ctx, 10*time.Second)
 115  	defer cancel()
 116  
 117  	c, err := nostr.Connect(pubCtx, url)
 118  	if err != nil {
 119  		failed.Add(1)
 120  		return
 121  	}
 122  	defer c.Disconnect()
 123  
 124  	go c.Listen(pubCtx)
 125  
 126  	if err := c.Publish(pubCtx, ev); err != nil {
 127  		failed.Add(1)
 128  		return
 129  	}
 130  
 131  	select {
 132  	case ok := <-c.OKs:
 133  		if ok.Accepted {
 134  			accepted.Add(1)
 135  		} else {
 136  			rejected.Add(1)
 137  		}
 138  	case <-pubCtx.Done():
 139  		failed.Add(1)
 140  	}
 141  }
 142