package server import ( "testing" "time" "git.smesh.lol/morly/pkg/relay/ratelimit" ) // --- constructor wiring --- var stReadyCalls int32 func stReady() { stReadyCalls++ } func TestNewWiresPools(t *testing.T) { c := stCfg() c.RelayURL = "wss://relay.example.com" c.IngestWorkers = 1 c.MediaProxyWorkers = 1 c.BlossomWorkers = 1 c.FreeWriteLimit = 5 c.FreeWriteWindow = 60 c.MaxConnPerIP = 3 db := stFakeDB() s := New(db, c) if s.selfHost != "relay.example.com/" { t.Fatalf("selfHost = %s", s.selfHost) return } if s.writeLimiter == nil { t.Fatal("a positive free-write limit must create a limiter") return } if s.workers.Len() != 1 || s.proxyPool.Len() != 1 || s.blossomPool.Len() != 1 { t.Fatalf("pools: ingest=%d proxy=%d blossom=%d", s.workers.Len(), s.proxyPool.Len(), s.blossomPool.Len()) return } if s.bcast == nil { t.Fatal("the broadcast worker must be started") return } if len(s.proxyBusyTime) != 1 { t.Fatalf("proxyBusyTime = %d, want 1", len(s.proxyBusyTime)) return } s.Close() // Zero worker counts skip the pools entirely. c2 := stCfg() c2.IngestWorkers = 0 c2.MediaProxyWorkers = 0 c2.BlossomWorkers = 0 s2 := New(stFakeDB(), c2) if s2.workers.Len() != 0 || s2.proxyPool.Len() != 0 || s2.blossomPool.Len() != 0 { t.Fatal("zero worker counts must not start pools") return } s2.Close() } func TestCloseWithStartedPools(t *testing.T) { s := stServer(stCfg()) s.db = stFakeDB() s.startIngestWorkers(1) s.startProxyWorkers(1) s.startBlossomWorkers(1) s.startBroadcastWorker() s.Close() s.Close() if !s.closed { t.Fatal("Close must set the closed flag") return } } // --- connection-open path --- func TestOnWSConnected(t *testing.T) { c := stCfg() c.RelayURL = "wss://relay.example.com" s, b := stSubSrv(c) defer b.Close() s.OnWSConnected(7) cs := s.conns[7] if cs == nil { t.Fatal("OnWSConnected must create the connection state") return } if cs.subs == nil { t.Fatal("OnWSConnected must initialise the subscription map") return } if len(cs.challenge) == 0 { t.Fatal("OnWSConnected must issue a NIP-42 challenge when RelayURL is set") return } // No relay URL means no challenge, but the state still exists. c2 := stCfg() s2, b2 := stSubSrv(c2) defer b2.Close() s2.OnWSConnected(8) if s2.conns[8] == nil { t.Fatal("OnWSConnected must create state without a relay URL") return } if s2.conns[8].challenge != nil { t.Fatal("no challenge must be issued when RelayURL is empty") return } } func TestWireOnReadyAndOnWSMessage(t *testing.T) { s := stServer(stCfg()) stReadyCalls = 0 s.OnReady = stReady s.wireOnReady() s.t.OnReady() if stReadyCalls != 1 { t.Fatalf("OnReady calls = %d, want 1", stReadyCalls) return } // No connection for the fd: the message is dropped by dispatch. s.OnWSMessage(1, []byte("[\"NOTICE\",\"x\"]")) } // --- broadcast respawn --- func TestRespawnBroadcastWorker(t *testing.T) { s, old := stSubSrv(stCfg()) fd := int32(4) cs := &cstate{subs: map[string]*sub{}} cs.authedPubkey = []byte("pubkey") cs.subs["s1"] = &sub{id: "s1", rawReq: []byte("[\"REQ\",\"s1\",{}]")} s.conns[fd] = cs s.respawnBroadcastWorker() old.Close() if s.bcast == nil { t.Fatal("respawnBroadcastWorker must install a new broadcaster") return } s.bcast.Close() } // --- poll loop --- func TestOnPollRunsAllDrains(t *testing.T) { s, b := stSubSrv(stCfg()) defer b.Close() s.OnPoll() } func TestOnTickReaperAndLimiterCleanup(t *testing.T) { s := stServer(stCfg()) s.writeLimiter = ratelimit.New(1.0, 1) // tickCount is a process-global; enough calls guarantee one multiple of 60. for i := 0; i < 61; i++ { s.OnTick() } if _, ok := s.asyncPending[1]; ok { t.Fatal("OnTick must reap stale async entries") return } } func TestProxyReapStuck(t *testing.T) { s := stProxySrv(1) old := s.proxyIn[0] s.proxyPool.Busy[0] = true s.proxyBusyTime[0] = time.Now().UnixNano() - 20_000_000_000 s.proxyReapStuck() if s.proxyPool.Busy[0] { t.Fatal("a stuck worker must be respawned idle") return } if s.proxyBusyTime[0] != 0 { t.Fatal("a reaped worker must clear its busy timestamp") return } close(old) close(s.proxyIn[0]) // A freshly busy worker is not reaped. s2 := stProxySrv(1) s2.proxyPool.Busy[0] = true s2.proxyBusyTime[0] = time.Now().UnixNano() s2.proxyReapStuck() if !s2.proxyPool.Busy[0] { t.Fatal("a fresh worker must not be reaped") return } } // --- database-unavailable REQ --- func TestHandleReqDBUnavailable(t *testing.T) { s, b := stSubSrv(stCfg()) defer b.Close() s.db = nil fd := int32(4) s.conns[fd] = &cstate{subs: map[string]*sub{}} s.handleReq(fd, []byte("[\"REQ\",\"s1\",{}]")) if len(s.conns[fd].subs) != 1 { t.Fatal("the subscription must register even when the database is gone") return } }