server_lifecycle_test.mx raw

   1  package server
   2  
   3  import (
   4  	"testing"
   5  	"time"
   6  
   7  	"git.smesh.lol/morly/pkg/relay/ratelimit"
   8  )
   9  
  10  // --- constructor wiring ---
  11  
  12  var stReadyCalls int32
  13  
  14  func stReady() { stReadyCalls++ }
  15  
  16  func TestNewWiresPools(t *testing.T) {
  17  	c := stCfg()
  18  	c.RelayURL = "wss://relay.example.com"
  19  	c.IngestWorkers = 1
  20  	c.MediaProxyWorkers = 1
  21  	c.BlossomWorkers = 1
  22  	c.FreeWriteLimit = 5
  23  	c.FreeWriteWindow = 60
  24  	c.MaxConnPerIP = 3
  25  	db := stFakeDB()
  26  
  27  	s := New(db, c)
  28  	if s.selfHost != "relay.example.com/" {
  29  		t.Fatalf("selfHost = %s", s.selfHost)
  30  		return
  31  	}
  32  	if s.writeLimiter == nil {
  33  		t.Fatal("a positive free-write limit must create a limiter")
  34  		return
  35  	}
  36  	if s.workers.Len() != 1 || s.proxyPool.Len() != 1 || s.blossomPool.Len() != 1 {
  37  		t.Fatalf("pools: ingest=%d proxy=%d blossom=%d", s.workers.Len(), s.proxyPool.Len(), s.blossomPool.Len())
  38  		return
  39  	}
  40  	if s.bcast == nil {
  41  		t.Fatal("the broadcast worker must be started")
  42  		return
  43  	}
  44  	if len(s.proxyBusyTime) != 1 {
  45  		t.Fatalf("proxyBusyTime = %d, want 1", len(s.proxyBusyTime))
  46  		return
  47  	}
  48  	s.Close()
  49  
  50  	// Zero worker counts skip the pools entirely.
  51  	c2 := stCfg()
  52  	c2.IngestWorkers = 0
  53  	c2.MediaProxyWorkers = 0
  54  	c2.BlossomWorkers = 0
  55  	s2 := New(stFakeDB(), c2)
  56  	if s2.workers.Len() != 0 || s2.proxyPool.Len() != 0 || s2.blossomPool.Len() != 0 {
  57  		t.Fatal("zero worker counts must not start pools")
  58  		return
  59  	}
  60  	s2.Close()
  61  }
  62  
  63  func TestCloseWithStartedPools(t *testing.T) {
  64  	s := stServer(stCfg())
  65  	s.db = stFakeDB()
  66  	s.startIngestWorkers(1)
  67  	s.startProxyWorkers(1)
  68  	s.startBlossomWorkers(1)
  69  	s.startBroadcastWorker()
  70  
  71  	s.Close()
  72  	s.Close()
  73  	if !s.closed {
  74  		t.Fatal("Close must set the closed flag")
  75  		return
  76  	}
  77  }
  78  
  79  // --- connection-open path ---
  80  
  81  func TestOnWSConnected(t *testing.T) {
  82  	c := stCfg()
  83  	c.RelayURL = "wss://relay.example.com"
  84  	s, b := stSubSrv(c)
  85  	defer b.Close()
  86  	s.OnWSConnected(7)
  87  	cs := s.conns[7]
  88  	if cs == nil {
  89  		t.Fatal("OnWSConnected must create the connection state")
  90  		return
  91  	}
  92  	if cs.subs == nil {
  93  		t.Fatal("OnWSConnected must initialise the subscription map")
  94  		return
  95  	}
  96  	if len(cs.challenge) == 0 {
  97  		t.Fatal("OnWSConnected must issue a NIP-42 challenge when RelayURL is set")
  98  		return
  99  	}
 100  
 101  	// No relay URL means no challenge, but the state still exists.
 102  	c2 := stCfg()
 103  	s2, b2 := stSubSrv(c2)
 104  	defer b2.Close()
 105  	s2.OnWSConnected(8)
 106  	if s2.conns[8] == nil {
 107  		t.Fatal("OnWSConnected must create state without a relay URL")
 108  		return
 109  	}
 110  	if s2.conns[8].challenge != nil {
 111  		t.Fatal("no challenge must be issued when RelayURL is empty")
 112  		return
 113  	}
 114  }
 115  
 116  func TestWireOnReadyAndOnWSMessage(t *testing.T) {
 117  	s := stServer(stCfg())
 118  	stReadyCalls = 0
 119  	s.OnReady = stReady
 120  	s.wireOnReady()
 121  	s.t.OnReady()
 122  	if stReadyCalls != 1 {
 123  		t.Fatalf("OnReady calls = %d, want 1", stReadyCalls)
 124  		return
 125  	}
 126  	// No connection for the fd: the message is dropped by dispatch.
 127  	s.OnWSMessage(1, []byte("[\"NOTICE\",\"x\"]"))
 128  }
 129  
 130  // --- broadcast respawn ---
 131  
 132  func TestRespawnBroadcastWorker(t *testing.T) {
 133  	s, old := stSubSrv(stCfg())
 134  	fd := int32(4)
 135  	cs := &cstate{subs: map[string]*sub{}}
 136  	cs.authedPubkey = []byte("pubkey")
 137  	cs.subs["s1"] = &sub{id: "s1", rawReq: []byte("[\"REQ\",\"s1\",{}]")}
 138  	s.conns[fd] = cs
 139  
 140  	s.respawnBroadcastWorker()
 141  	old.Close()
 142  	if s.bcast == nil {
 143  		t.Fatal("respawnBroadcastWorker must install a new broadcaster")
 144  		return
 145  	}
 146  	s.bcast.Close()
 147  }
 148  
 149  // --- poll loop ---
 150  
 151  func TestOnPollRunsAllDrains(t *testing.T) {
 152  	s, b := stSubSrv(stCfg())
 153  	defer b.Close()
 154  	s.OnPoll()
 155  }
 156  
 157  func TestOnTickReaperAndLimiterCleanup(t *testing.T) {
 158  	s := stServer(stCfg())
 159  	s.writeLimiter = ratelimit.New(1.0, 1)
 160  	// tickCount is a process-global; enough calls guarantee one multiple of 60.
 161  	for i := 0; i < 61; i++ {
 162  		s.OnTick()
 163  	}
 164  	if _, ok := s.asyncPending[1]; ok {
 165  		t.Fatal("OnTick must reap stale async entries")
 166  		return
 167  	}
 168  }
 169  
 170  func TestProxyReapStuck(t *testing.T) {
 171  	s := stProxySrv(1)
 172  	old := s.proxyIn[0]
 173  	s.proxyPool.Busy[0] = true
 174  	s.proxyBusyTime[0] = time.Now().UnixNano() - 20_000_000_000
 175  	s.proxyReapStuck()
 176  	if s.proxyPool.Busy[0] {
 177  		t.Fatal("a stuck worker must be respawned idle")
 178  		return
 179  	}
 180  	if s.proxyBusyTime[0] != 0 {
 181  		t.Fatal("a reaped worker must clear its busy timestamp")
 182  		return
 183  	}
 184  	close(old)
 185  	close(s.proxyIn[0])
 186  
 187  	// A freshly busy worker is not reaped.
 188  	s2 := stProxySrv(1)
 189  	s2.proxyPool.Busy[0] = true
 190  	s2.proxyBusyTime[0] = time.Now().UnixNano()
 191  	s2.proxyReapStuck()
 192  	if !s2.proxyPool.Busy[0] {
 193  		t.Fatal("a fresh worker must not be reaped")
 194  		return
 195  	}
 196  }
 197  
 198  // --- database-unavailable REQ ---
 199  
 200  func TestHandleReqDBUnavailable(t *testing.T) {
 201  	s, b := stSubSrv(stCfg())
 202  	defer b.Close()
 203  	s.db = nil
 204  	fd := int32(4)
 205  	s.conns[fd] = &cstate{subs: map[string]*sub{}}
 206  	s.handleReq(fd, []byte("[\"REQ\",\"s1\",{}]"))
 207  	if len(s.conns[fd].subs) != 1 {
 208  		t.Fatal("the subscription must register even when the database is gone")
 209  		return
 210  	}
 211  }
 212