Files
OwnCord/Server/ws/serve_pumps.go
T
J3vbandClaude Fable 5 7be9ccd2f9 fix: batch of 22 correctness fixes across server and client (#1371)
* fix(voice): 4 defect(s) (OC-0008, OC-0009, OC-0042, OC-0080)

Guard LiveKit session state against supersession: bump the camera/screen
generation in leaveVoice and teardownForReconnect so an in-flight enable
discards its track, bail out of restoreLocalVoiceState when a newer room
claimed _room mid-await, and recheck isStateConnected in the auto-reconnect
tail.

* fix(ws): 1 defect(s) (OC-0019)

* fix(db): 1 defect(s) (OC-0023)

* fix(ws): 1 defect(s) (OC-0029)

* fix(ws): 1 defect(s) (OC-0032)

* fix(voice): 1 defect(s) (OC-0034)

* fix(admin): 1 defect(s) (OC-0035)

* fix(service): 2 defect(s) (OC-0036, OC-0128)

* fix(voice): 2 defect(s) (OC-0038, OC-0065)

OC-0038: the LiveKit participant_left webhook cleared the leaver's own
client voice state before broadcasting voice_leave, so the broadcast
audience (READ_MESSAGES holders union still-in-the-room participants)
could no longer see them. Voice membership is gated on CONNECT_VOICE
alone, so a participant without READ_MESSAGES never learned the server
had torn down their call. Extracted finishVoiceLeave's audience logic
into broadcastVoiceEventWithLeaver and used it on the webhook path.

OC-0065: handleWebhookParticipantJoined OR'd a GetVoiceState read error
into the same branch as "no matching row", so a transient DB failure
ejected a legitimate participant from the SFU mid-call. Now the read
error is logged and the check skipped, matching sweepStaleVoiceStates.

* fix(client): 1 defect(s) (OC-0041)

* fix(client): 1 defect(s) (OC-0043)

* fix(client): 1 defect(s) (OC-0046)

* fix(client): 1 defect(s) (OC-0047)

* fix(client): 1 defect(s) (OC-0049)

* fix(client): 1 defect(s) (OC-0108)

* fix(client): 2 defect(s) (OC-0111, OC-0143)

OC-0111: retry a presence_update dropped by the 1-per-10s limiter once the
window reopens, so auto-idle's return-to-online does not leave the server
and every other client stuck on idle.

OC-0143: pass apiConfig.host to the DM profile sidebar so per-user notes
are scoped per server, matching channel mutes, the NSFW gate and volume.

* test(ws): align aborted-switch test with OC-0034 no-resurrect behavior

The fix agent rewrote this pre-existing test (it locked the buggy restore
path) but the prove agent left it out of c67d25ed; committed state alone
failed go test ./ws/ without it.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>

---------

Co-authored-by: Claude Fable 5 <noreply@anthropic.com>
2026-08-14 16:15:32 +02:00

224 lines
6.6 KiB
Go

package ws
import (
"context"
"log/slog"
"time"
"github.com/coder/websocket"
"github.com/owncord/server/db"
)
// writePump drains the client's send channels and writes to the WebSocket.
// Priority ordering: high > normal > low. High-priority messages (DMs, mentions)
// are drained first. Normal messages (chat, reactions) come next. Low-priority
// messages (typing, presence) are only sent when no higher-priority work is pending.
func writePump(ctx context.Context, conn *websocket.Conn, c *Client) {
writeMsg := func(msg []byte) bool {
wCtx, cancel := context.WithTimeout(ctx, writeTimeout)
err := conn.Write(wCtx, websocket.MessageText, msg)
cancel()
if err != nil {
slog.Warn("ws writePump error", "user_id", c.userID, "err", err)
return false
}
return true
}
// drainChannel writes every message still buffered on ch without blocking.
// Returns false only when a write failed; empty or closed is true.
drainChannel := func(ch chan []byte) bool {
for {
select {
case msg, ok := <-ch:
if !ok {
return true
}
if !writeMsg(msg) {
return false
}
default:
return true
}
}
}
// drainAndClose flushes whatever the kick paths queued before closing the
// send channels (e.g. the BANNED error frame that makes the client clear
// its credentials) — serve.go and hub_broadcast.go both document that
// writePump drains remaining messages after closeSend. Returning on the
// first closed channel would drop those frames.
drainAndClose := func() {
if drainChannel(c.sendHigh) && drainChannel(c.send) {
drainChannel(c.sendLow)
}
_ = conn.Close(websocket.StatusNormalClosure, "")
}
for {
// Priority 1: drain all pending high-priority messages first.
select {
case msg, ok := <-c.sendHigh:
if !ok {
drainAndClose()
return
}
if !writeMsg(msg) {
return
}
continue
default:
}
// Priority 2: try high or normal, non-blocking. Go's select among
// ready cases is uniformly random, so sendLow cannot be a peer here —
// a case that fires the moment any low-priority frame is queued would
// let it win the coin flip against a pending normal frame roughly
// half the time, contradicting "low-priority messages are only sent
// when no higher-priority work is pending". Only fall through to
// sendLow (via the blocking select below) once this default proves
// neither high nor normal has anything ready right now.
select {
case msg, ok := <-c.sendHigh:
if !ok {
drainAndClose()
return
}
if !writeMsg(msg) {
return
}
continue
case msg, ok := <-c.send:
if !ok {
drainAndClose()
return
}
if !writeMsg(msg) {
return
}
continue
default:
}
// Priority 3: nothing high or normal is ready — block on all three
// (plus shutdown) so an idle connection still gets its typing/presence
// frames instead of busy-looping.
select {
case msg, ok := <-c.sendHigh:
if !ok {
drainAndClose()
return
}
if !writeMsg(msg) {
return
}
case msg, ok := <-c.send:
if !ok {
drainAndClose()
return
}
if !writeMsg(msg) {
return
}
case msg, ok := <-c.sendLow:
if !ok {
drainAndClose()
return
}
if !writeMsg(msg) {
return
}
case <-ctx.Done():
return
}
}
}
// readPump reads from the WebSocket and dispatches messages. Blocks until disconnect.
func readPump(ctx context.Context, conn *websocket.Conn, hub *Hub, c *Client) {
var lastReadErr error
defer func() {
// The connection is gone, so ctx is (or is about to be) cancelled.
// Teardown DB writes must still complete — a dead connection must not
// cancel its own cleanup — so detach cancellation but keep values.
cleanupCtx := context.WithoutCancel(ctx)
// Snapshot voice state BEFORE unregister to avoid TOCTOU with replacement connections.
voiceChID := c.getVoiceChID()
replaced := hub.unregisterNow(c)
if c.user != nil {
// Clean up voice state only when this was the user's final
// connection. A replacement connection owns the (transferred)
// voice session, and the join_token guard cannot tell the
// difference — the transfer keeps the same joined_at — so
// cleaning here would delete the replacement's DB row whenever
// teardown snapshots voiceChID before the transfer zeroes it.
if voiceChID != 0 && !replaced {
hub.handleVoiceLeave(cleanupCtx, c)
}
c.mu.Lock()
received := c.msgsReceived
sent := c.msgsSent
dropped := c.msgsDropped
c.mu.Unlock()
duration := time.Since(c.connectedAt)
attrs := []any{
"username", c.user.Username,
"user_id", c.userID,
"remote", c.remoteAddr,
"duration_s", int64(duration.Seconds()),
"msgs_received", received,
"msgs_sent", sent,
"msgs_dropped", dropped,
}
if voiceChID > 0 {
attrs = append(attrs, "voice_channel_id", voiceChID)
}
if replaced {
attrs = append(attrs, "replaced", true)
}
if lastReadErr != nil {
attrs = append(attrs, "last_error", lastReadErr.Error())
}
slog.Info("websocket disconnected", attrs...)
// shouldMarkOffline re-checks h.clients instead of trusting
// `replaced` alone: that flag was sampled before handleVoiceLeave,
// which can block for seconds, so a reconnect landing during that
// window would otherwise be invisible here and this dead
// connection's teardown would mark the live session's user
// offline (OC-0019).
if hub.shouldMarkOffline(c, replaced) {
// A real disconnect is offline for everyone, the user
// included, so this path needs no invisible mapping. The row,
// however, keeps a *chosen* status (idle/dnd/invisible)
// standing — that is what the next connect reads to avoid
// stamping the user back online. MarkUserDisconnected clears
// only the non-choice "online" and refreshes last_seen; the
// stale-choice problem it would otherwise create is handled at
// read time, where a member with no live connection renders
// offline no matter what the column says.
_ = hub.db.MarkUserDisconnected(cleanupCtx, c.userID)
// custom_status is nil, not c.user.CustomStatus: that field is a
// snapshot taken once at auth (client.go) and never updated, so
// broadcasting it here would resurrect a status the user changed
// or cleared mid-session. presentableMembers applies the same
// rule for a fresh ready payload (serve_ready.go) — a member with
// no live connection shows no custom status.
hub.BroadcastToAll(buildPresenceMsg(c.userID, db.StatusOffline, nil))
}
}
}()
for {
_, msg, err := conn.Read(ctx)
if err != nil {
lastReadErr = err
return
}
c.touch()
hub.handleMessage(c, msg)
}
}