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