Files
OwnCord/Server/api/dm_handler.go
T
J3vbandClaude Opus 5 ea0430c5b0 fix: batch of correctness fixes across server and client (#1375)
* fix(client): 1 defect(s) (OC-0201)

* fix(service): 1 defect(s) (OC-0202)

HandleTyping built the per-user-per-channel rate-limit key before resolving the channel or checking read permission, so forged channel ids could pin unbounded dead entries in the shared process-wide RateLimiter.

* fix(client): 2 defect(s) (OC-0203, OC-0224)

* fix(server): 1 defect(s) (OC-0204)

* fix(ws): 2 defect(s) (OC-0205, OC-0211)

* fix(admin): 2 defect(s) (OC-0209, OC-0212)

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

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

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

Route handler-driven PresenceEvent through BroadcastToAll instead of BroadcastToAllLow so every source of a user's presence shares one ordered per-client FIFO.

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

PATCH /users/{id} combining banned + role_id committed and broadcast the ban before authorizing the role change, so a refused role change returned an error while leaving the target banned. Authorize the role change up front via the new ModerationService.AuthorizeRoleChange.

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

LinkAttachmentsToMessage no longer claims an attachment that is a user's live avatar (users.avatar points at it). Once message_id is set, handleServeFile's avatar branch (gated on ChannelID == nil) is unreachable and the file falls under the message's channel ACL / soft-delete state, permanently disagreeing with users.avatar about who may read it.

* fix(emoji): 1 defect(s) (OC-0217)

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

The data-copy phase of an HTTP proxy tunnel was unbounded. Steps 1-2 of
handle_connection (header read, TCP connect, TLS handshake) each run under
a 10s guard, but step 3 called io::copy_bidirectional with no deadline. A
remote that completes the TLS handshake and then neither responds nor
closes parks the spawned connection task, the loopback socket and the
remote TLS session indefinitely: copy_bidirectional only resolves once
BOTH directions finish, so closing the local side alone does not free it.

Wrap the copy in copy_with_deadline, a generic helper bounded by
DATA_PHASE_TIMEOUT (600s). The bound is deliberately far looser than the
10s setup guards because this phase carries the REST body, including
attachment and avatar uploads, so it must reclaim only genuinely stuck
connections rather than merely slow ones. The helper is generic over the
stream types so it can be exercised without a live TLS connection.

Regression test drives two in-memory duplex pairs whose far ends stay
alive, so neither half ever observes EOF and raw copy_bidirectional would
block forever; the test asserts the call resolves on its own deadline with
ErrorKind::TimedOut.

Claude-Session: https://claude.ai/code/session_01ENMDTh8gDLiHCaRFdMYRiL

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

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

UpdateNotifier scheduled its deferred update check with a setTimeout whose
handle was never retained, so destroy() could not cancel it. A component torn
down inside the 3s window (page swap / logout) still fired performCheck() and
issued a network update check against the old server URL. Retain the timer
handle and clear it in destroy().

* fix(dm): 1 defect(s) (OC-0222)

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

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

The Grant-Microphone retry's .finally hardcoded grantMicBtn.disabled = false, undoing updateFrozen()'s socket-down freeze when the WS socket dropped while the mic permission request was in flight. Delegate the state back to render().

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

handleApplyUpdate broadcasts a 'restarting in 5s' notice before the on-disk
swap. Every failure path in the swap returned silently, leaving clients
counting down to a restart that never happened. Extract the swap into
applyStagedUpdate and send a corrective 'update_aborted' broadcast from a
deferred guard on every path that does not reach the respawn.

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

PATCH /channels/{id} accepted a blank or whitespace-only name, leaving the
channel unidentifiable in clients. updateChannelRequest.validate() now
rejects it the way handleCreateChannel already did.

* fix(identity): 1 defect(s) (OC-0228)

* fix(admin): run deferred cleanup before the update restart exits

The fix batch left three golangci-lint findings and two prettier findings
that CI gates on.

applyStagedUpdate called os.Exit(0) in the same function that defers both
staged.Close() and the corrective "update_aborted" broadcast, so neither
ran (gocritic exitAfterDefer). Return a bool instead and let the caller
exit once those defers have run — on Windows, releasing the staged binary's
file handle is the reason the restart exists at all, so this is a real fix
rather than a lint appeasement. The exported test hook calls the function as
a statement, so the added result does not affect it.

Also modernize a bulk-insert loop to range-over-int, compare backup bytes
with bytes.Equal, and reflow two test files to prettier's output.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ENMDTh8gDLiHCaRFdMYRiL

* test(ws): pin the live presence path against the invisible custom-status leak

OC-0207 and OC-0211 are the same defect at two emitters: hub_broadcast.go's
BroadcastPresence (connect/reconnect) and event.go's presenceEvents (live
presence_update). The fix for OC-0211 closed both sites in one change, but
only the hub_broadcast side got a regression test.

This pins the event.go sibling: an invisible user's real custom status must
be blanked on the PresenceOthersEvent frame while the owner's own
PresenceSelfEvent still carries it. Without it, a later change could reopen
the live path while the committed test kept passing.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ENMDTh8gDLiHCaRFdMYRiL

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

* test(ws): silence a contextcheck false positive in the reconnect race test

RefreshChannelVisibility takes no context by design — it is reached through
the admin HubBroadcaster interface, which carries none, so it builds its own
internally. contextcheck flags the call only because the test closure around
it holds a ctx for its override write, so there is nothing to propagate.
Suppress at the call site rather than widen a production interface (and its
mocks) to satisfy a lint in a test.

golangci-lint v2.11.3 (the version ci.yml pins) now reports 0 issues.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01ENMDTh8gDLiHCaRFdMYRiL

---------

Co-authored-by: Claude <noreply@anthropic.com>
2026-08-15 16:30:05 +02:00

444 lines
16 KiB
Go

package api
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"net/http"
"github.com/go-chi/chi/v5"
"github.com/owncord/server/db"
"github.com/owncord/server/service"
"github.com/owncord/server/ws"
)
// DMBroadcaster is the interface needed to send WebSocket events from REST
// handlers. Satisfied by *ws.Hub.
type DMBroadcaster interface {
SendToUser(userID int64, msg []byte) bool
}
// dmVisibilityMarker is an optional DMBroadcaster capability: bump the hub's
// visibility watermark after an unsequenced, targeted event so a client that
// warm-reconnects across the gap takes the full-ready path instead of a
// sequenced-only replay that can never redeliver it. Reached by type
// assertion rather than being added to DMBroadcaster directly so the
// SendToUser-only test doubles in this package keep working. Satisfied in
// production by ws.Hub.MarkVisibilityChanged, which forwards to the same
// bumpVisibilityWatermark the WS-side dm_channel_open emitter uses
// (ws/emit.go).
type dmVisibilityMarker interface {
MarkVisibilityChanged()
}
// The production broadcaster must keep satisfying it: a type assertion that
// silently stops matching would turn the watermark bump back into the no-op
// this fixed, with nothing failing to say so — mirrors the dmVoiceEvictor
// assertion below for its sibling capability.
var _ dmVisibilityMarker = (*ws.Hub)(nil)
// markDMVisibilityChanged bumps the visibility watermark if broadcaster
// supports it. dm_channel_open/close are unsequenced and targeted, so a
// client that misses one via a dropped connection and then warm-reconnects
// never gets it redelivered by the ordinary seq-replay path — mirroring why
// the WS emitter of the same event (ws/emit.go DMChannelOpenEvent) forces
// this bump unconditionally, regardless of whether the send itself
// succeeded.
func markDMVisibilityChanged(broadcaster DMBroadcaster) {
if vm, ok := broadcaster.(dmVisibilityMarker); ok {
vm.MarkVisibilityChanged()
}
}
// dmVoiceEvictor is the DMBroadcaster capability used to evict a user's
// voice-call connection for one specific channel, leaving an unrelated call
// they may currently be in untouched (which the unconditional
// DisconnectFromVoice would not). It is kept out of DMBroadcaster itself and
// reached by type assertion so the handler stays usable with the
// SendToUser-only test doubles the package already has.
type dmVoiceEvictor interface {
DisconnectFromVoiceInChannel(ctx context.Context, userID, channelID int64) bool
}
// The production broadcaster must keep satisfying it: a type assertion that
// silently stops matching would turn the eviction back into the no-op this
// fixed, with nothing failing to say so.
var _ dmVoiceEvictor = (*ws.Hub)(nil)
// MountDMRoutes registers DM-related routes onto r.
// All routes require authentication.
// hub is used to send real-time WebSocket events on DM close.
func MountDMRoutes(r chi.Router, database *db.DB, svc *service.Services, broadcaster DMBroadcaster) {
r.Route("/api/v1/dms", func(r chi.Router) {
r.Use(AuthMiddleware(database))
r.Post("/", handleCreateDM(svc))
r.Post("/group", handleCreateGroupDM(svc, broadcaster))
r.Get("/", handleListDMs(svc))
r.Patch("/{channelId}", handleRenameGroupDM(svc, broadcaster))
r.Delete("/{channelId}", handleCloseDM(svc, broadcaster))
})
// User blocking routes — prevent DM creation and messaging.
r.Route("/api/v1/blocks", func(r chi.Router) {
r.Use(AuthMiddleware(database))
r.Get("/", handleListBlocks(svc))
r.Put("/{userId}", handleBlockUser(svc))
r.Delete("/{userId}", handleUnblockUser(svc))
})
}
// createDMRequest is the JSON body for POST /api/v1/dms.
type createDMRequest struct {
RecipientID int64 `json:"recipient_id"`
}
// createDMResponse is the JSON response for POST /api/v1/dms.
type createDMResponse struct {
ChannelID int64 `json:"channel_id"`
Recipient db.DMUser `json:"recipient"`
Created bool `json:"created"`
}
// createGroupDMRequest is the JSON body for POST /api/v1/dms/group.
type createGroupDMRequest struct {
RecipientIDs []int64 `json:"recipient_ids"`
Name string `json:"name"`
}
// renameDMRequest is the JSON body for PATCH /api/v1/dms/{channelId}.
type renameDMRequest struct {
Name string `json:"name"`
}
// listDMsResponse is the JSON response for GET /api/v1/dms.
type listDMsResponse struct {
DMChannels []db.DMChannelInfo `json:"dm_channels"`
}
// handleCreateDM creates or retrieves a DM channel with a recipient.
func handleCreateDM(svc *service.Services) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
user, ok := r.Context().Value(UserKey).(*db.User)
if !ok || user == nil {
writeJSON(w, http.StatusUnauthorized, errorResponse{
Error: "UNAUTHORIZED", Message: "authentication required",
})
return
}
var req createDMRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, errorResponse{
Error: "BAD_REQUEST", Message: "invalid request body",
})
return
}
result, err := svc.DMs.CreateDM(r.Context(), user.ID, req.RecipientID)
if err != nil {
writeServiceError(r.Context(), w, err)
return
}
avatarStr := ""
if result.Recipient.Avatar != nil {
avatarStr = *result.Recipient.Avatar
}
displayName := ""
if result.Recipient.DisplayName != nil {
displayName = *result.Recipient.DisplayName
}
dmUser := db.DMUser{
ID: result.Recipient.ID,
Username: result.Recipient.Username,
Avatar: avatarStr,
Status: db.StatusForViewer(result.Recipient.Status, result.Recipient.ID, user.ID),
DisplayName: displayName,
}
status := http.StatusOK
if result.Created {
status = http.StatusCreated
}
writeJSON(w, status, createDMResponse{
ChannelID: result.Channel.ID,
Recipient: dmUser,
Created: result.Created,
})
}
}
// handleListDMs returns all open DM channels for the authenticated user.
func handleListDMs(svc *service.Services) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
user, ok := r.Context().Value(UserKey).(*db.User)
if !ok || user == nil {
writeJSON(w, http.StatusUnauthorized, errorResponse{
Error: "UNAUTHORIZED", Message: "authentication required",
})
return
}
channels, err := svc.DMs.ListDMs(r.Context(), user.ID)
if err != nil {
writeServiceError(r.Context(), w, err)
return
}
writeJSON(w, http.StatusOK, listDMsResponse{DMChannels: channels})
}
}
// handleCloseDM removes a DM channel from the authenticated user's open list.
func handleCloseDM(svc *service.Services, broadcaster DMBroadcaster) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
user, ok := r.Context().Value(UserKey).(*db.User)
if !ok || user == nil {
writeJSON(w, http.StatusUnauthorized, errorResponse{
Error: "UNAUTHORIZED", Message: "authentication required",
})
return
}
channelID, ok := parseIDParam(w, r, "channelId")
if !ok {
return
}
result, err := svc.DMs.CloseDM(r.Context(), user.ID, channelID)
if err != nil {
writeServiceError(r.Context(), w, err)
return
}
// Notify via WebSocket so sidebar updates immediately.
if broadcaster != nil {
closeMsg := fmt.Appendf(nil, `{"type":%q,"payload":{"channel_id":%d}}`, ws.MsgTypeDMChannelClose, channelID)
if ok := broadcaster.SendToUser(user.ID, closeMsg); !ok {
slog.Debug("handleCloseDM: user not connected", "user_id", user.ID, "channel_id", channelID)
}
// dm_channel_close is unsequenced and targeted like
// dm_channel_open — see markDMVisibilityChanged.
markDMVisibilityChanged(broadcaster)
// A group leave changes the membership everyone else renders, so
// the survivors get a refreshed dm_channel_open rather than being
// left showing a member who has gone.
if result.Left && !result.ChannelDeleted {
broadcastDMOpen(r.Context(), svc, broadcaster, channelID, result.RemainingParticipantIDs)
}
// Leaving a group DM removes the caller from its membership but,
// without this, leaves them connected to its live voice call —
// they keep hearing and speaking to a room they are no longer a
// member of. Scoped to this channel so a leaver currently on an
// unrelated voice call is untouched. This also covers the
// last-participant case (ChannelDeleted): the row is already gone
// by now, so CleanupVoiceForChannel would read an FK-cascaded
// empty voice_states and do nothing, while the leaver — the only
// participant left — is evicted here.
if result.Left {
if ve, ok := broadcaster.(dmVoiceEvictor); ok {
ve.DisconnectFromVoiceInChannel(context.WithoutCancel(r.Context()), user.ID, channelID)
}
}
}
w.WriteHeader(http.StatusNoContent)
}
}
// broadcastDMOpen sends a per-viewer dm_channel_open for channelID to each of
// targetIDs. The payload differs per addressee (`recipient`/`recipients` are
// relative to who is reading), so it is rebuilt inside the loop.
//
// Failures are logged and skipped, never surfaced: the mutation that prompted
// this has already committed, and a client that misses the event re-derives
// the same state from its next `ready`.
func broadcastDMOpen(ctx context.Context, svc *service.Services, broadcaster DMBroadcaster, channelID int64, targetIDs []int64) {
if broadcaster == nil || len(targetIDs) == 0 {
return
}
// The mutation that led here has already committed, so this fan-out must
// survive the caller's request context being cancelled after that point
// (client disconnect mid-handler) — otherwise every DMSummaryFor lookup
// below fails with context.Canceled and no participant, including ones
// otherwise unaffected by the cancellation, ever receives the open.
ctx = context.WithoutCancel(ctx)
// dm_channel_open is unsequenced and targeted — a recipient who is
// offline or drops the connection right now can never have it replayed
// to them by the ordinary seq-based resume path, so a warm reconnect must
// be forced onto the full-ready path instead. Bumped once per call,
// unconditionally (not per-recipient SendToUser result): the ws emitter
// of this same event does the same (ws/emit.go), and this covers every
// caller — group create, rename refresh, and the group-leave refresh.
markDMVisibilityChanged(broadcaster)
for _, pid := range targetIDs {
summary, pErr := svc.DMs.DMSummaryFor(ctx, pid, channelID)
if pErr != nil {
slog.Debug("broadcastDMOpen: summary unavailable", "user_id", pid, "channel_id", channelID, "err", pErr)
continue
}
msg, mErr := json.Marshal(map[string]any{
"type": "dm_channel_open",
"payload": summary,
})
if mErr != nil {
slog.Warn("broadcastDMOpen: marshal failed", "err", mErr, "channel_id", channelID)
continue
}
if ok := broadcaster.SendToUser(pid, msg); !ok {
slog.Debug("broadcastDMOpen: user not connected", "user_id", pid, "channel_id", channelID)
}
}
}
// handleCreateGroupDM creates a group DM between the caller and 2..8 others.
func handleCreateGroupDM(svc *service.Services, broadcaster DMBroadcaster) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
user, ok := r.Context().Value(UserKey).(*db.User)
if !ok || user == nil {
writeJSON(w, http.StatusUnauthorized, errorResponse{
Error: "UNAUTHORIZED", Message: "authentication required",
})
return
}
var req createGroupDMRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, errorResponse{
Error: "BAD_REQUEST", Message: "invalid request body",
})
return
}
result, err := svc.DMs.CreateGroupDM(r.Context(), user.ID, req.RecipientIDs, req.Name)
if err != nil {
writeServiceError(r.Context(), w, err)
return
}
// Everyone gets the DM in their sidebar immediately, the creator
// included — the REST response is only the creator's copy, and a
// second window of theirs needs the event just as much as the others.
broadcastDMOpen(r.Context(), svc, broadcaster, result.Channel.ID, result.ParticipantIDs)
writeJSON(w, http.StatusCreated,
db.NewDMChannelInfo(result.Channel.ID, result.Channel.Name, true, result.Participants, user.ID))
}
}
// handleRenameGroupDM sets or clears a group DM's name. Participants only —
// there is no owner, so every member holds the same authority over it.
func handleRenameGroupDM(svc *service.Services, broadcaster DMBroadcaster) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
user, ok := r.Context().Value(UserKey).(*db.User)
if !ok || user == nil {
writeJSON(w, http.StatusUnauthorized, errorResponse{
Error: "UNAUTHORIZED", Message: "authentication required",
})
return
}
channelID, ok := parseIDParam(w, r, "channelId")
if !ok {
return
}
var req renameDMRequest
if err := json.NewDecoder(r.Body).Decode(&req); err != nil {
writeJSON(w, http.StatusBadRequest, errorResponse{
Error: "BAD_REQUEST", Message: "invalid request body",
})
return
}
if _, err := svc.DMs.RenameGroupDM(r.Context(), user.ID, channelID, req.Name); err != nil {
writeServiceError(r.Context(), w, err)
return
}
// The rename has already committed at this point, so this lookup must
// survive the caller's request context being cancelled right after
// that commit (client disconnect mid-handler) — same reasoning as
// broadcastDMOpen's own context.WithoutCancel, and the failure must be
// logged rather than silently dropping the fan-out (participants would
// keep rendering the stale name with no compensating resync, since
// dm_channel_open is unsequenced/targeted and can't be replayed).
bgCtx := context.WithoutCancel(r.Context())
participantIDs, pErr := svc.Channels.GetDMParticipantIDs(bgCtx, channelID)
if pErr != nil {
slog.Error("handleRenameGroupDM: participant lookup failed", "err", pErr, "channel_id", channelID)
} else {
broadcastDMOpen(bgCtx, svc, broadcaster, channelID, participantIDs)
}
summary, sErr := svc.DMs.DMSummaryFor(r.Context(), user.ID, channelID)
if sErr != nil {
writeServiceError(r.Context(), w, sErr)
return
}
writeJSON(w, http.StatusOK, summary)
}
}
// handleBlockUser blocks a user.
func handleBlockUser(svc *service.Services) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
user, _ := r.Context().Value(UserKey).(*db.User)
if user == nil {
writeJSON(w, http.StatusUnauthorized, errorResponse{Error: "UNAUTHORIZED", Message: "authentication required"})
return
}
targetID, ok := parseIDParam(w, r, "userId")
if !ok {
return
}
if err := svc.Blocks.BlockUser(r.Context(), user.ID, targetID); err != nil {
writeServiceError(r.Context(), w, err)
return
}
writeJSON(w, http.StatusOK, map[string]string{"message": "user blocked"})
}
}
// handleUnblockUser unblocks a user.
func handleUnblockUser(svc *service.Services) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
user, _ := r.Context().Value(UserKey).(*db.User)
if user == nil {
writeJSON(w, http.StatusUnauthorized, errorResponse{Error: "UNAUTHORIZED", Message: "authentication required"})
return
}
targetID, ok := parseIDParam(w, r, "userId")
if !ok {
return
}
if err := svc.Blocks.UnblockUser(r.Context(), user.ID, targetID); err != nil {
writeServiceError(r.Context(), w, err)
return
}
writeJSON(w, http.StatusOK, map[string]string{"message": "user unblocked"})
}
}
// handleListBlocks returns all blocked user IDs.
func handleListBlocks(svc *service.Services) http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
user, _ := r.Context().Value(UserKey).(*db.User)
if user == nil {
writeJSON(w, http.StatusUnauthorized, errorResponse{Error: "UNAUTHORIZED", Message: "authentication required"})
return
}
ids, err := svc.Blocks.ListBlocked(r.Context(), user.ID)
if err != nil {
writeServiceError(r.Context(), w, err)
return
}
writeJSON(w, http.StatusOK, map[string]any{"blocked_user_ids": ids})
}
}