bunker: fix concurrency bugs, config.Clients/newSecret/cancelPreviousBunkerInfoPrint were accessed from multiple goroutines without a mutex, and nostrconnect subscriptions forwarded events into the pool-owned channel which panics once the pool closes it.

This commit is contained in:
Yasuhiro Matsumoto
2026-07-15 10:26:55 +09:00
parent 34281a9c8e
commit 44db7af6c0
+122 -99
View File
@@ -13,6 +13,7 @@ import (
"path/filepath" "path/filepath"
"slices" "slices"
"strings" "strings"
"sync"
"time" "time"
"fiatjaf.com/nostr" "fiatjaf.com/nostr"
@@ -235,12 +236,19 @@ var bunker = &cli.Command{
// it will be stored // it will be stored
newSecret := randString(12) newSecret := randString(12)
// guards config.Clients, newSecret and cancelPreviousBunkerInfoPrint, which are
// accessed from the socket goroutine, the per-request handler goroutines and here
var mu sync.Mutex
// static information // static information
pubkey := sec.Public() pubkey := sec.Public()
npub := nip19.EncodeNpub(pubkey) npub := nip19.EncodeNpub(pubkey)
// this function will be called every now and then // this function will be called every now and then
printBunkerInfo := func() { printBunkerInfo := func() {
mu.Lock()
defer mu.Unlock()
iqs := make(url.Values) iqs := make(url.Values)
maps.Copy(iqs, qs) maps.Copy(iqs, qs)
iqs.Set("secret", newSecret) iqs.Set("secret", newSecret)
@@ -342,70 +350,15 @@ var bunker = &cli.Command{
signer := nip46.NewStaticKeySigner(sec) signer := nip46.NewStaticKeySigner(sec)
signer.DefaultRelays = config.Relays signer.DefaultRelays = config.Relays
// unix socket nostrconnect:// handling
go func() {
for uri := range onSocketConnect(ctx, c) {
clientPublicKey, err := nostr.PubKeyFromHex(uri.Host)
if err != nil {
continue
}
log("- got nostrconnect:// request from '%s': %s\n", color.New(color.Bold, color.FgBlue).Sprint(clientPublicKey.Hex()), uri.String())
relays := uri.Query()["relay"]
// pre-authorize this client since the user has explicitly added it
if !slices.ContainsFunc(config.Clients, func(c BunkerConfigClient) bool {
return c.PubKey == clientPublicKey
}) {
config.Clients = append(config.Clients, BunkerConfigClient{
PubKey: clientPublicKey,
Name: uri.Query().Get("name"),
URL: uri.Query().Get("url"),
Icon: uri.Query().Get("icon"),
CustomRelays: relays,
})
}
if persist != nil {
persist()
}
resp, eventResponse, err := signer.HandleNostrConnectURI(ctx, uri)
if err != nil {
log("* failed to handle: %s\n", err)
continue
}
go func() {
for event := range sys.Pool.SubscribeMany(ctx, relays, nostr.Filter{
Kinds: []nostr.Kind{nostr.KindNostrConnect},
Tags: nostr.TagMap{"p": []string{pubkey.Hex()}},
Since: nostr.Now(),
LimitZero: true,
}, nostr.SubscriptionOptions{Label: "nak-bunker"}) {
events <- event
}
}()
time.Sleep(time.Millisecond * 25)
jresp, _ := json.MarshalIndent(resp, "", " ")
log("~ responding with %s\n", string(jresp))
for res := range sys.Pool.PublishMany(ctx, relays, eventResponse) {
if res.Error == nil {
log("* sent through %s\n", res.Relay.URL)
} else {
log("* failed to send through %s: %s\n", res.RelayURL, res.Error)
}
}
}
}()
// just a gimmick // just a gimmick
var cancelPreviousBunkerInfoPrint context.CancelFunc var cancelPreviousBunkerInfoPrint context.CancelFunc
_, cancel := context.WithCancel(ctx) _, cancel := context.WithCancel(ctx)
cancelPreviousBunkerInfoPrint = cancel cancelPreviousBunkerInfoPrint = cancel
signer.AuthorizeRequest = func(harmless bool, from nostr.PubKey, secret string) bool { signer.AuthorizeRequest = func(harmless bool, from nostr.PubKey, secret string) bool {
mu.Lock()
defer mu.Unlock()
if slices.ContainsFunc(config.Clients, func(b BunkerConfigClient) bool { return b.PubKey == from }) { if slices.ContainsFunc(config.Clients, func(b BunkerConfigClient) bool { return b.PubKey == from }) {
return true return true
} }
@@ -441,60 +394,130 @@ var bunker = &cli.Command{
return false return false
} }
for ie := range events { handleBunkerRequest := func(ie nostr.RelayEvent) {
mu.Lock()
cancelPreviousBunkerInfoPrint() // this prevents us from printing a million bunker info blocks cancelPreviousBunkerInfoPrint() // this prevents us from printing a million bunker info blocks
mu.Unlock()
go func() { // handle the NIP-46 request event
// handle the NIP-46 request event from := ie.Event.PubKey
from := ie.Event.PubKey req, resp, eventResponse, err := signer.HandleRequest(ctx, ie.Event)
req, resp, eventResponse, err := signer.HandleRequest(ctx, ie.Event) if err != nil {
if err != nil { if errors.Is(err, nip46.AlreadyHandled) {
if errors.Is(err, nip46.AlreadyHandled) {
return
}
log("< failed to handle request from %s: %s\n", from.Hex(), err.Error())
return return
} }
jreq, _ := json.MarshalIndent(req, "", " ") log("< failed to handle request from %s: %s\n", from.Hex(), err.Error())
log("- got request from '%s': %s\n", color.New(color.Bold, color.FgBlue).Sprint(from.Hex()), string(jreq)) return
jresp, _ := json.MarshalIndent(resp, "", " ") }
log("~ responding with %s\n", string(jresp))
// use custom relays if they are defined for this client jreq, _ := json.MarshalIndent(req, "", " ")
// (normally if the initial connection came from a nostrconnect:// URL) log("- got request from '%s': %s\n", color.New(color.Bold, color.FgBlue).Sprint(from.Hex()), string(jreq))
relays := config.Relays jresp, _ := json.MarshalIndent(resp, "", " ")
for _, c := range config.Clients { log("~ responding with %s\n", string(jresp))
if c.PubKey == from && len(c.CustomRelays) > 0 {
relays = c.CustomRelays // use custom relays if they are defined for this client
break // (normally if the initial connection came from a nostrconnect:// URL)
} relays := config.Relays
mu.Lock()
for _, c := range config.Clients {
if c.PubKey == from && len(c.CustomRelays) > 0 {
relays = c.CustomRelays
break
}
}
mu.Unlock()
for res := range sys.Pool.PublishMany(ctx, relays, eventResponse) {
if res.Error == nil {
log("* sent response through %s\n", res.Relay.URL)
} else {
log("* failed to send response through %s: %s\n", res.RelayURL, res.Error)
}
}
// just after handling one request we trigger this
go func() {
ctx, cancel := context.WithCancel(ctx)
defer cancel()
mu.Lock()
cancelPreviousBunkerInfoPrint = cancel
mu.Unlock()
// the idea is that we will print the bunker URL again so it is easier to copy-paste by users
// but we will only do if the bunker is inactive for more than 5 minutes
select {
case <-ctx.Done():
case <-time.After(time.Minute * 5):
log("\n")
printBunkerInfo()
}
}()
}
// unix socket nostrconnect:// handling
go func() {
for uri := range onSocketConnect(ctx, c) {
clientPublicKey, err := nostr.PubKeyFromHex(uri.Host)
if err != nil {
continue
}
log("- got nostrconnect:// request from '%s': %s\n", color.New(color.Bold, color.FgBlue).Sprint(clientPublicKey.Hex()), uri.String())
relays := uri.Query()["relay"]
// pre-authorize this client since the user has explicitly added it
mu.Lock()
if !slices.ContainsFunc(config.Clients, func(c BunkerConfigClient) bool {
return c.PubKey == clientPublicKey
}) {
config.Clients = append(config.Clients, BunkerConfigClient{
PubKey: clientPublicKey,
Name: uri.Query().Get("name"),
URL: uri.Query().Get("url"),
Icon: uri.Query().Get("icon"),
CustomRelays: relays,
})
} }
for res := range sys.Pool.PublishMany(ctx, relays, eventResponse) { if persist != nil {
if res.Error == nil { persist()
log("* sent response through %s\n", res.Relay.URL) }
} else { mu.Unlock()
log("* failed to send response through %s: %s\n", res.RelayURL, res.Error)
} resp, eventResponse, err := signer.HandleNostrConnectURI(ctx, uri)
if err != nil {
log("* failed to handle: %s\n", err)
continue
} }
// just after handling one request we trigger this
go func() { go func() {
ctx, cancel := context.WithCancel(ctx) for event := range sys.Pool.SubscribeMany(ctx, relays, nostr.Filter{
defer cancel() Kinds: []nostr.Kind{nostr.KindNostrConnect},
cancelPreviousBunkerInfoPrint = cancel Tags: nostr.TagMap{"p": []string{pubkey.Hex()}},
// the idea is that we will print the bunker URL again so it is easier to copy-paste by users Since: nostr.Now(),
// but we will only do if the bunker is inactive for more than 5 minutes LimitZero: true,
select { }, nostr.SubscriptionOptions{Label: "nak-bunker"}) {
case <-ctx.Done(): // handle directly instead of forwarding into the main events
case <-time.After(time.Minute * 5): // channel, which is owned (and eventually closed) by the pool
log("\n") go handleBunkerRequest(event)
printBunkerInfo()
} }
}() }()
}()
time.Sleep(time.Millisecond * 25)
jresp, _ := json.MarshalIndent(resp, "", " ")
log("~ responding with %s\n", string(jresp))
for res := range sys.Pool.PublishMany(ctx, relays, eventResponse) {
if res.Error == nil {
log("* sent through %s\n", res.Relay.URL)
} else {
log("* failed to send through %s: %s\n", res.RelayURL, res.Error)
}
}
}
}()
for ie := range events {
go handleBunkerRequest(ie)
} }
return nil return nil