Files
OwnCord/Server/ws/export_test.go
T
J3vbandClaude Opus 5 0594a130cb fix(ws): give the LiveKit health check its own HTTP transport (#1356)
NewLiveKitProcess built its health-check http.Client without a Transport,
so it fell back to the process-wide http.DefaultTransport.

httptest.Server.Close calls CloseIdleConnections on http.DefaultTransport
by design ("assume most users of httptest.Server will be using the standard
transport, so help them out"), and ws is full of t.Parallel tests that each
defer srv.Close(). Any one of them finishing while a health check held a
pooled connection severed that request:

  livekit_test.go:978: HealthCheck: livekit health check failed:
    Get "http://127.0.0.1:41343": net/http: HTTP/1.x transport connection
    broken: http: CloseIdleConnections called

That surfaced as an unrelated-looking CI failure on a TypeScript lint bump
(#1341). It is not purely a test artifact: in production the health check
also shared one connection pool with every other DefaultTransport user in
the server process.

Cloning DefaultTransport keeps its tuned defaults (proxy, dial and TLS
timeouts, HTTP/2) while giving the client a private pool.

Locked by TestHealthCheckClientOwnsItsTransport, which fails on the
unfixed constructor.

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-11 19:56:54 +02:00

414 lines
15 KiB
Go

// export_test.go exposes unexported functions and methods for use in external
// test packages (package ws_test). This file is compiled only during "go test".
package ws
import (
"context"
"encoding/json"
"fmt"
"net/http"
"os/exec"
"time"
"github.com/livekit/protocol/livekit"
"github.com/owncord/server/db"
)
// ─── hub sweep helpers ─────────────────────────────────────────────────────
// SweepStaleClientsForTest exposes sweepStaleClients for external tests.
func (h *Hub) SweepStaleClientsForTest() {
h.sweepStaleClients()
}
// SweepStaleVoiceStatesForTest exposes sweepStaleVoiceStates for external tests.
func (h *Hub) SweepStaleVoiceStatesForTest() {
h.sweepStaleVoiceStates()
}
// SweepRevokedSessionsForTest exposes sweepRevokedSessions for external tests.
func (h *Hub) SweepRevokedSessionsForTest() {
h.sweepRevokedSessions()
}
// SetClientLastActivityForTest overwrites a client's lastActivity timestamp.
func SetClientLastActivityForTest(c *Client, t time.Time) {
c.mu.Lock()
defer c.mu.Unlock()
c.lastActivity = t
}
// ─── client getter/setter helpers ──────────────────────────────────────────
// GetLastActivityForTest exposes Client.getLastActivity for external tests.
func GetLastActivityForTest(c *Client) time.Time {
return c.getLastActivity()
}
// ClearVoiceChIDForTest exposes Client.clearVoiceChID for external tests.
func ClearVoiceChIDForTest(c *Client) int64 {
return c.clearVoiceChID()
}
// SetVoiceChIDForTest sets the voice channel ID atomically, clearing the join
// token when leaving (chID 0) — the same contract production keeps via
// setVoiceState. Test-only: production has no set-channel-without-token path.
func SetVoiceChIDForTest(c *Client, chID int64) {
c.voiceMu.Lock()
defer c.voiceMu.Unlock()
c.voiceChID = chID
if chID == 0 {
c.voiceJoinToken = ""
}
}
// SetClientVoiceChID is an alias kept for existing tests.
func SetClientVoiceChID(c *Client, channelID int64) {
SetVoiceChIDForTest(c, channelID)
}
// SubscribedToVoiceTopicForTest reports whether c itself (identity compare,
// not just its userID) holds the subscription to channelID's voice topic.
func (h *Hub) SubscribedToVoiceTopicForTest(c *Client, channelID int64) bool {
h.pubsub.mu.RLock()
defer h.pubsub.mu.RUnlock()
return h.pubsub.topics[VoiceTopic(channelID)][c.userID] == c
}
// SubscribeVoiceTopicForTest subscribes c to channelID's voice topic, as the
// production voice_join flow does.
func (h *Hub) SubscribeVoiceTopicForTest(c *Client, channelID int64) {
h.pubsub.Subscribe(c, VoiceTopic(channelID))
}
// SetClientVoiceStateForTest sets both the voice channel and join token.
func SetClientVoiceStateForTest(c *Client, channelID int64, joinToken string) {
c.voiceMu.Lock()
defer c.voiceMu.Unlock()
c.voiceChID = channelID
c.voiceJoinToken = joinToken
}
// SetClientE2EEPubKeyForTest sets the E2EE public key on a client (no signature).
func SetClientE2EEPubKeyForTest(c *Client, key string) {
c.setE2EEPubKey(key, "")
}
// GetClientE2EEPubKeyForTest returns the E2EE public key from a client.
func GetClientE2EEPubKeyForTest(c *Client) string {
key, _ := c.getE2EEPubKey()
return key
}
// NewTestClient creates a client with a caller-supplied send channel; conn is nil.
func NewTestClient(hub *Hub, userID int64, send chan []byte) *Client {
return &Client{
hub: hub,
ctx: context.Background(),
userID: userID,
send: send,
sendHigh: send, // unified for test observability
sendLow: send,
}
}
// NewTestClientWithChannel creates a test client subscribed to a specific channel.
func NewTestClientWithChannel(hub *Hub, userID, channelID int64, send chan []byte) *Client {
return &Client{
hub: hub,
ctx: context.Background(),
userID: userID,
channelID: channelID,
send: send,
sendHigh: send, // unified for test observability
sendLow: send,
}
}
// NewTestClientWithUser creates a test client with an authenticated user record
// set. Use this when tests need the client to pass permission checks.
func NewTestClientWithUser(hub *Hub, user *db.User, channelID int64, send chan []byte) *Client {
return &Client{
hub: hub,
ctx: context.Background(),
userID: user.ID,
user: user,
channelID: channelID,
send: send,
sendHigh: send, // unified for test observability
sendLow: send,
}
}
// NewTestClientWithTokenHash creates a test client that carries a session token
// hash. Use this when tests need to exercise the periodic session-expiry check.
func NewTestClientWithTokenHash(hub *Hub, user *db.User, tokenHash string, channelID int64, send chan []byte) *Client {
return &Client{
hub: hub,
ctx: context.Background(),
userID: user.ID,
user: user,
tokenHash: tokenHash,
channelID: channelID,
send: send,
sendHigh: send, // unified for test observability
sendLow: send,
}
}
// RunningForTest reports whether the hub's Run loop has started.
func (h *Hub) RunningForTest() bool {
return h.running.Load()
}
// ClientUserIDForTest returns the client's user ID for external tests.
func ClientUserIDForTest(c *Client) int64 {
return c.userID
}
// ClientChannelIDForTest returns the client's currently focused channel for
// external tests.
func ClientChannelIDForTest(c *Client) int64 {
return c.getChannelID()
}
// TouchForTest exposes Client.touch for external tests.
func TouchForTest(c *Client) {
c.touch()
}
// RollbackVoiceJoinForTest exposes Hub.rollbackVoiceJoin for external tests.
func (h *Hub) RollbackVoiceJoinForTest(c *Client, channelID int64) {
h.rollbackVoiceJoin(context.Background(), c, channelID, true)
}
// LeaveVoiceChannelWithRetryForTest exposes leaveVoiceChannelWithRetry for external tests.
func LeaveVoiceChannelWithRetryForTest(h *Hub, userID int64, channelID int64, joinToken string) error {
return leaveVoiceChannelWithRetry(context.Background(), h, userID, channelID, joinToken)
}
// ─── livekit process/webhook helpers ───────────────────────────────────────
// GenerateConfigForTest exposes LiveKitProcess.generateConfig for external tests.
func (p *LiveKitProcess) GenerateConfigForTest() (string, error) {
return p.generateConfig()
}
// HTTPTransportForTest exposes the health-check client's transport so tests can
// assert it does not share http.DefaultTransport's connection pool.
func (p *LiveKitProcess) HTTPTransportForTest() http.RoundTripper {
return p.httpClient.Transport
}
// SetProcessCmdForTest sets cmd to a non-nil value to simulate "already running".
func (p *LiveKitProcess) SetProcessCmdForTest() {
p.mu.Lock()
defer p.mu.Unlock()
p.cmd = &exec.Cmd{}
}
// SetProcessStoppedForTest sets stopped=true to simulate a stopped process.
func (p *LiveKitProcess) SetProcessStoppedForTest() {
p.mu.Lock()
defer p.mu.Unlock()
p.stopped = true
}
// NewHubForTest creates a minimal Hub with no DB or limiter for webhook testing.
func NewHubForTest() *Hub {
return &Hub{
clients: make(map[int64]*Client),
pubsub: NewPubSub(),
topicLimiter: NewTopicRateLimiter(topicRateLimitPerSecond, time.Second),
}
}
// PubSubForTest exposes the hub's PubSub for external tests.
func (h *Hub) PubSubForTest() *PubSub {
return h.pubsub
}
// BuildAuthOKForTest exposes Hub.buildAuthOK for external tests.
// Defaults to replay_source="none" since most callers test the fresh-connect
// path; tests that care about the resume tier can call buildAuthOK directly.
func (h *Hub) BuildAuthOKForTest(user *db.User, roleName string) []byte {
return h.buildAuthOK(context.Background(), user, roleName, "none")
}
// RunMentionCountsInlineForTest makes the hub's MessageService apply mention
// counts synchronously instead of on a background goroutine, so a test can read
// the counts deterministically right after driving a chat_send through the hub.
func (h *Hub) RunMentionCountsInlineForTest() {
h.messageSvc.RunBackgroundInlineForTest()
}
// BuildReadyForTest exposes Hub.buildReady for external tests.
// Passes nil role so no channels are visible (fail-closed, BUG-094).
func (h *Hub) BuildReadyForTest(database *db.DB, userID int64) ([]byte, error) {
return h.buildReady(context.Background(), database, userID, nil)
}
// BuildReadyWithRoleForTest exposes Hub.buildReady with a role for external tests.
func (h *Hub) BuildReadyWithRoleForTest(database *db.DB, userID int64, role *db.Role) ([]byte, error) {
return h.buildReady(context.Background(), database, userID, role)
}
// ComputeAllowedChannelsForTest exposes Hub.computeAllowedChannels for external
// tests (the REST/WS channel-visibility agreement test).
func (h *Hub) ComputeAllowedChannelsForTest(database *db.DB, user *db.User) (map[int64]bool, error) {
return h.computeAllowedChannels(context.Background(), database, user)
}
// GetCachedSettingsForTest exposes Hub.getCachedSettings for external tests.
func (h *Hub) GetCachedSettingsForTest() (string, string) {
return h.getCachedSettings(context.Background())
}
// GetClientVoiceChIDForTest exposes Client.getVoiceChID for external tests.
func GetClientVoiceChIDForTest(c *Client) int64 {
return c.getVoiceChID()
}
// GetClientVoiceJoinTokenForTest reads the join token under voiceMu.
func GetClientVoiceJoinTokenForTest(c *Client) string {
c.voiceMu.Lock()
defer c.voiceMu.Unlock()
return c.voiceJoinToken
}
// ExpireSettingsCacheForTest forces the settings cache to appear stale so that
// the next call to getCachedSettings triggers a DB refresh.
func (h *Hub) ExpireSettingsCacheForTest() {
h.settingsMu.Lock()
defer h.settingsMu.Unlock()
h.settingsLastUpdate = time.Time{} // zero time — always older than any TTL
}
// ParseChannelIDForTest exposes parseChannelID for external tests.
func ParseChannelIDForTest(payload json.RawMessage) (int64, error) {
return parseChannelID(payload)
}
// BuildJSONForTest exposes buildJSON for external tests.
func BuildJSONForTest(v any) []byte {
return buildJSON(v)
}
// ParseIdentityForTest parses a LiveKit participant identity and discards the
// join token, exercising the production parseParticipantIdentity.
func ParseIdentityForTest(identity string) (int64, error) {
userID, _, err := parseParticipantIdentity(identity)
return userID, err
}
// ParseParticipantIdentityForTest exposes parseParticipantIdentity for tests.
func ParseParticipantIdentityForTest(identity string) (int64, string, error) {
return parseParticipantIdentity(identity)
}
// ParseRoomChannelIDForTest exposes parseRoomChannelID for external tests.
func ParseRoomChannelIDForTest(roomName string) (int64, error) {
return parseRoomChannelID(roomName)
}
// WsToHTTPForTest exposes wsToHTTP for external tests.
func WsToHTTPForTest(wsURL string) string {
return wsToHTTP(wsURL)
}
// RegisterNowForTest exposes registerNow for external tests so clients are
// visible immediately (no channel round-trip through hub.Run). No channels are
// readable, matching the hub-loop registration path.
func (h *Hub) RegisterNowForTest(c *Client) {
h.registerNow(c, nil)
}
// RegisterNowWithReadableForTest exposes registerNow with an explicit
// READ_MESSAGES channel set, as the handshake paths in serve.go supply it.
func (h *Hub) RegisterNowWithReadableForTest(c *Client, readableChannelIDs map[int64]bool) {
h.registerNow(c, readableChannelIDs)
}
// ClearVoiceStateForTest exposes clearVoiceState for external tests.
func (c *Client) ClearVoiceStateForTest() {
c.clearVoiceState()
}
// QualityBitrateForTest exposes qualityBitrate for external tests.
func QualityBitrateForTest(quality string) int {
return qualityBitrate(quality)
}
// BuildDMChannelOpenForTest exposes buildDMChannelOpenFor for external tests.
func BuildDMChannelOpenForTest(channelID int64, recipient *db.User) []byte {
return buildDMChannelOpenFor(channelID, recipient, 0)
}
// BuildDMChannelOpenInfoForTest exposes the group-aware buildDMChannelOpen.
func BuildDMChannelOpenInfoForTest(info db.DMChannelInfo) []byte {
return buildDMChannelOpen(info)
}
// BuildCallSignalForTest exposes buildCallSignal for external tests.
func BuildCallSignalForTest(msgType string, channelID, fromUserID int64, username string) []byte {
return buildCallSignal(msgType, channelID, fromUserID, username)
}
// HandleWebhookParticipantLeftForTest exposes handleWebhookParticipantLeft for
// external tests so they can simulate LiveKit webhook events without HTTP.
func (h *Hub) HandleWebhookParticipantLeftForTest(userID int64, channelID int64, joinToken string) {
identity := fmt.Sprintf("user-%d:%s", userID, joinToken)
roomName := fmt.Sprintf("channel-%d", channelID)
event := &livekit.WebhookEvent{
Event: "participant_left",
Participant: &livekit.ParticipantInfo{
Identity: identity,
},
Room: &livekit.Room{
Name: roomName,
},
}
h.handleWebhookParticipantLeft(context.Background(), event)
}
// HandleWebhookParticipantJoinedForTest exposes handleWebhookParticipantJoined
// for external tests. identity and roomName are passed raw so a test can feed
// malformed values through the same parse path a hostile webhook would.
func (h *Hub) HandleWebhookParticipantJoinedForTest(identity, roomName string) {
event := &livekit.WebhookEvent{
Event: "participant_joined",
Participant: &livekit.ParticipantInfo{Identity: identity},
Room: &livekit.Room{Name: roomName},
}
h.handleWebhookParticipantJoined(context.Background(), event)
}
// HandleWebhookParticipantJoinedEventForTest exposes
// handleWebhookParticipantJoined with a caller-built event so tests can cover
// the nil-participant and nil-room guards.
func (h *Hub) HandleWebhookParticipantJoinedEventForTest(event *livekit.WebhookEvent) {
h.handleWebhookParticipantJoined(context.Background(), event)
}
// MustFullResyncForTest exposes mustFullResync for external tests.
func (h *Hub) MustFullResyncForTest(lastSeq uint64) bool {
return h.mustFullResync(lastSeq)
}
// HasChannelPermForTest exposes Hub.hasChannelPerm for external tests.
func (h *Hub) HasChannelPermForTest(c *Client, channelID, perm int64) bool {
return h.hasChannelPerm(context.Background(), c, channelID, perm)
}
// BroadcastVoiceEventForTest exposes Hub.broadcastVoiceEvent for external
// tests so a load/soak test can drive the channelReadAudience-resolved
// voice_state/voice_leave fan-out directly, without a full LiveKit join
// round-trip.
func (h *Hub) BroadcastVoiceEventForTest(channelID int64, msg []byte) {
h.broadcastVoiceEvent(context.Background(), channelID, msg)
}
// MaxColdReplayForTest exposes the cold-tier replay row cap so tests can seed
// exactly enough events to hit it.
const MaxColdReplayForTest = maxColdReplay