Files
glchat/internal/store/dm_calls.go
T
grendervill 50709396cd feat(dm): хранилище и миграция звонков в беседах (Фаза 7)
Звонок в личной и групповой беседе живёт в своей таблице: в voice_states
guild_id обязателен, а у беседы сервера нет. Миграция 00026 добавляет
dm_calls (статус, media, длительность, причина завершения) и
dm_call_participants (состояние участника и флаги микрофона, камеры и
экрана); уникальный частичный индекс держит один идущий звонок на беседу.

В DMParticipantProfiles сравнение с владельцем обёрнуто в COALESCE: у 1:1
dm_owner_id пуст, и NULL нельзя было прочитать в int — ручка звонка в
личной беседе падала на 500.

Тесты: жизненный цикл звонка в хранилище (гудки, active на втором
участнике, история, повторный звонок) и выборка сторожем.
2026-09-26 17:02:51 +03:00

486 lines
17 KiB
Go

package store
import (
"context"
"database/sql"
"time"
)
// Звонки в личных и групповых беседах (AGENT.md 7.8, Фаза 7): своя таблица, а
// не voice_states — там guild_id обязателен, а у беседы сервера нет. Комната
// LiveKit у звонка одна на звонок, имя собирает voice.CallRoomName.
// Статусы звонка.
const (
DMCallRinging = "ringing"
DMCallActive = "active"
DMCallEnded = "ended"
)
// Состояния участника звонка.
const (
DMCallInvited = "invited"
DMCallJoined = "joined"
DMCallDeclined = "declined"
DMCallLeft = "left"
DMCallMissed = "missed"
)
// Причины завершения звонка.
const (
DMCallReasonCompleted = "completed"
DMCallReasonCanceled = "canceled"
DMCallReasonDeclined = "declined"
DMCallReasonMissed = "missed"
DMCallReasonFailed = "failed"
)
// История звонков беседы: по умолчанию 50 записей на страницу, максимум 100.
const (
defaultDMCallHistory = 50
maxDMCallHistory = 100
)
// DMCall — звонок беседы: кто позвонил, когда и чем закончилось.
type DMCall struct {
ID uint64
ChannelID uint64
InitiatorID uint64
Status string
Media string
CreatedAt time.Time
StartedAt *time.Time
EndedAt *time.Time
EndedBy *uint64
EndReason string
}
// DMCallParticipant — участник звонка и его состояние.
type DMCallParticipant struct {
CallID uint64
UserID uint64
State string
InvitedAt time.Time
JoinedAt *time.Time
LeftAt *time.Time
SelfMute bool
SelfDeaf bool
Camera bool
Screen bool
}
// DMCallFlags — флаги участника: микрофон, звук, камера, шаринг экрана.
type DMCallFlags struct {
SelfMute bool
SelfDeaf bool
Camera bool
Screen bool
}
// DMCallWithParticipants — звонок вместе с составом: так его отдают API и
// Gateway, поэтому отдельный тип, а не два запроса у вызывающего кода.
type DMCallWithParticipants struct {
DMCall
Participants []DMCallParticipant
}
// DMCallDurationSeconds считает длительность разговора: от момента, когда
// звонок стал active (кто-то ответил), до завершения или до текущего времени.
func (s *Store) DMCallDurationSeconds(call *DMCall) int {
if call == nil || call.StartedAt == nil {
return 0
}
end := s.now()
if call.EndedAt != nil {
end = *call.EndedAt
}
if end.Before(*call.StartedAt) {
return 0
}
return int(end.Sub(*call.StartedAt).Seconds())
}
const dmCallColumns = `id, channel_id, initiator_id, status, media, created_at,
started_at, ended_at, ended_by, end_reason`
const dmCallParticipantColumns = `call_id, user_id, state, invited_at, joined_at, left_at,
self_mute, self_deaf, camera, screen`
// CreateDMCall открывает звонок: инициатор сразу в комнате, остальные
// приглашены. Второй одновременный звонок в беседе невозможен — его
// отклоняет уникальный индекс dm_calls_live_channel_idx.
func (s *Store) CreateDMCall(ctx context.Context, channelID, initiatorID uint64, media string, invitees []uint64) (*DMCallWithParticipants, error) {
callID := s.NextID()
now := s.Now()
if err := s.InTx(ctx, func(tx *sql.Tx) error {
var live int
if err := tx.QueryRowContext(ctx,
`SELECT COUNT(*) FROM dm_calls WHERE channel_id = ? AND status <> ?`,
int64(channelID), DMCallEnded).Scan(&live); err != nil {
return err
}
if live > 0 {
return ErrConflict
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO dm_calls (id, channel_id, initiator_id, status, media, created_at)
VALUES (?, ?, ?, ?, ?, ?)`,
int64(callID), int64(channelID), int64(initiatorID), DMCallRinging, media, now); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO dm_call_participants (call_id, user_id, state, invited_at, joined_at)
VALUES (?, ?, ?, ?, ?)`,
int64(callID), int64(initiatorID), DMCallJoined, now, now); err != nil {
return err
}
for _, inviteeID := range invitees {
if inviteeID == initiatorID {
continue
}
if _, err := tx.ExecContext(ctx, `
INSERT INTO dm_call_participants (call_id, user_id, state, invited_at)
VALUES (?, ?, ?, ?)
ON CONFLICT (call_id, user_id) DO NOTHING`,
int64(callID), int64(inviteeID), DMCallInvited, now); err != nil {
return err
}
}
return nil
}); err != nil {
return nil, mapError(err)
}
return s.GetDMCall(ctx, callID)
}
// GetDMCall возвращает звонок с составом.
func (s *Store) GetDMCall(ctx context.Context, callID uint64) (*DMCallWithParticipants, error) {
row := s.reader.QueryRowContext(ctx,
`SELECT `+dmCallColumns+` FROM dm_calls WHERE id = ?`, int64(callID))
call, err := scanDMCall(row)
if err != nil {
return nil, err
}
participants, err := s.listDMCallParticipants(ctx, callID)
if err != nil {
return nil, err
}
return &DMCallWithParticipants{DMCall: *call, Participants: participants}, nil
}
// GetLiveDMCall возвращает идущий звонок беседы: гудки или разговор.
func (s *Store) GetLiveDMCall(ctx context.Context, channelID uint64) (*DMCallWithParticipants, error) {
row := s.reader.QueryRowContext(ctx,
`SELECT `+dmCallColumns+` FROM dm_calls WHERE channel_id = ? AND status <> ? ORDER BY id DESC LIMIT 1`,
int64(channelID), DMCallEnded)
call, err := scanDMCall(row)
if err != nil {
return nil, err
}
participants, err := s.listDMCallParticipants(ctx, call.ID)
if err != nil {
return nil, err
}
return &DMCallWithParticipants{DMCall: *call, Participants: participants}, nil
}
// IsDMCallParticipant сообщает, звали ли пользователя в этот звонок.
func (s *Store) IsDMCallParticipant(ctx context.Context, callID, userID uint64) (bool, error) {
var exists int
err := s.reader.QueryRowContext(ctx,
`SELECT 1 FROM dm_call_participants WHERE call_id = ? AND user_id = ?`,
int64(callID), int64(userID)).Scan(&exists)
if err == sql.ErrNoRows {
return false, nil
}
if err != nil {
return false, err
}
return true, nil
}
// JoinDMCallParticipant переводит участника в состояние «в звонке» и
// запускает отсчёт разговора, когда в комнате оказывается второй человек.
func (s *Store) JoinDMCallParticipant(ctx context.Context, callID, userID uint64) (*DMCall, error) {
now := s.Now()
if err := s.InTx(ctx, func(tx *sql.Tx) error {
result, err := tx.ExecContext(ctx, `
UPDATE dm_call_participants SET state = ?, joined_at = ?, left_at = NULL
WHERE call_id = ? AND user_id = ?`,
DMCallJoined, now, int64(callID), int64(userID))
if err != nil {
return err
}
if affected, err := result.RowsAffected(); err == nil && affected == 0 {
return ErrNotFound
}
var joined int
if err := tx.QueryRowContext(ctx, `
SELECT COUNT(*) FROM dm_call_participants WHERE call_id = ? AND state = ?`,
int64(callID), DMCallJoined).Scan(&joined); err != nil {
return err
}
if joined < 2 {
return nil
}
_, err = tx.ExecContext(ctx, `
UPDATE dm_calls SET status = ?, started_at = COALESCE(started_at, ?)
WHERE id = ? AND status = ?`,
DMCallActive, now, int64(callID), DMCallRinging)
return err
}); err != nil {
return nil, mapError(err)
}
return s.GetDMCallRow(ctx, callID)
}
// LeaveDMCallParticipant помечает участника вышедшим, отклонившим или
// пропустившим звонок.
func (s *Store) LeaveDMCallParticipant(ctx context.Context, callID, userID uint64, state string) error {
_, err := s.writer.ExecContext(ctx, `
UPDATE dm_call_participants SET state = ?, left_at = ?
WHERE call_id = ? AND user_id = ?`,
state, s.Now(), int64(callID), int64(userID))
return mapError(err)
}
// CountDMCallParticipants считает участников звонка в заданном состоянии.
func (s *Store) CountDMCallParticipants(ctx context.Context, callID uint64, state string) (int, error) {
var count int
err := s.reader.QueryRowContext(ctx,
`SELECT COUNT(*) FROM dm_call_participants WHERE call_id = ? AND state = ?`,
int64(callID), state).Scan(&count)
return count, err
}
// UpdateDMCallFlags сохраняет флаги участника (микрофон, звук, камера, экран).
func (s *Store) UpdateDMCallFlags(ctx context.Context, callID, userID uint64, flags DMCallFlags) error {
result, err := s.writer.ExecContext(ctx, `
UPDATE dm_call_participants
SET self_mute = ?, self_deaf = ?, camera = ?, screen = ?
WHERE call_id = ? AND user_id = ?`,
boolToInt(flags.SelfMute), boolToInt(flags.SelfDeaf),
boolToInt(flags.Camera), boolToInt(flags.Screen),
int64(callID), int64(userID))
if err != nil {
return mapError(err)
}
if affected, err := result.RowsAffected(); err == nil && affected == 0 {
return ErrNotFound
}
return nil
}
// EndDMCall завершает звонок: не ответившие помечаются пропустившими, те, кто
// был в комнате, — вышедшими. Повторный вызов не меняет уже завершённый звонок.
func (s *Store) EndDMCall(ctx context.Context, callID uint64, endedBy *uint64, reason string) error {
now := s.Now()
var actor sql.NullInt64
if endedBy != nil {
actor = sql.NullInt64{Int64: int64(*endedBy), Valid: true}
}
return s.InTx(ctx, func(tx *sql.Tx) error {
if _, err := tx.ExecContext(ctx, `
UPDATE dm_calls SET status = ?, ended_at = ?, ended_by = ?, end_reason = ?
WHERE id = ? AND status <> ?`,
DMCallEnded, now, actor, reason, int64(callID), DMCallEnded); err != nil {
return err
}
if _, err := tx.ExecContext(ctx, `
UPDATE dm_call_participants SET state = ?, left_at = ?
WHERE call_id = ? AND state = ?`,
DMCallMissed, now, int64(callID), DMCallInvited); err != nil {
return err
}
_, err := tx.ExecContext(ctx, `
UPDATE dm_call_participants SET state = ?, left_at = ?
WHERE call_id = ? AND state = ?`,
DMCallLeft, now, int64(callID), DMCallJoined)
return err
})
}
// ListDMCalls отдаёт историю звонков беседы от новых к старым. beforeID —
// курсор: отдаются звонки с меньшим идентификатором.
func (s *Store) ListDMCalls(ctx context.Context, channelID uint64, limit int, beforeID uint64) ([]DMCallWithParticipants, error) {
if limit <= 0 || limit > maxDMCallHistory {
limit = defaultDMCallHistory
}
rows, err := s.reader.QueryContext(ctx, `
SELECT `+dmCallColumns+` FROM dm_calls
WHERE channel_id = ? AND (? = 0 OR id < ?)
ORDER BY id DESC LIMIT ?`,
int64(channelID), int64(beforeID), int64(beforeID), limit)
if err != nil {
return nil, err
}
defer rows.Close()
calls := make([]DMCallWithParticipants, 0, limit)
for rows.Next() {
call, err := scanDMCall(rows)
if err != nil {
return nil, err
}
calls = append(calls, DMCallWithParticipants{DMCall: *call})
}
if err := rows.Err(); err != nil {
return nil, err
}
if len(calls) == 0 {
return calls, nil
}
// Участники тянутся одним запросом на страницу: диапазон идентификаторов
// страницы непрерывен, поэтому лишних звонков в выборку не попадает.
page := make(map[uint64][]DMCallParticipant, len(calls))
participants, err := s.listDMCallParticipantsRange(ctx, channelID, calls[len(calls)-1].ID, beforeID)
if err != nil {
return nil, err
}
for _, participant := range participants {
page[participant.CallID] = append(page[participant.CallID], participant)
}
for index := range calls {
calls[index].Participants = page[calls[index].ID]
if calls[index].Participants == nil {
calls[index].Participants = []DMCallParticipant{}
}
}
return calls, nil
}
// ListStaleDMCalls ищет звонки для сторожа: гудки без ответа (ringingBefore) и
// разговоры, в которых не осталось никого (emptyBefore — время последнего
// обновления звонка). Оба условия — «строго раньше» указанного момента.
func (s *Store) ListStaleDMCalls(ctx context.Context, ringingBefore, emptyBefore time.Time) ([]DMCall, error) {
rows, err := s.reader.QueryContext(ctx, `
SELECT `+dmCallColumns+` FROM dm_calls c
WHERE c.status <> ?
AND (
(c.status = ? AND c.created_at < ?)
OR (c.status = ? AND NOT EXISTS (
SELECT 1 FROM dm_call_participants p
WHERE p.call_id = c.id AND p.state = ?
) AND c.created_at < ?)
)
ORDER BY c.id`,
DMCallEnded, DMCallRinging, s.Timestamp(ringingBefore), DMCallActive, DMCallJoined,
s.Timestamp(emptyBefore))
if err != nil {
return nil, err
}
defer rows.Close()
calls := make([]DMCall, 0, 4)
for rows.Next() {
call, err := scanDMCall(rows)
if err != nil {
return nil, err
}
calls = append(calls, *call)
}
return calls, rows.Err()
}
// CountLiveDMCalls считает звонки, которые ещё не завершены (диагностика).
func (s *Store) CountLiveDMCalls(ctx context.Context) (int, error) {
var count int
err := s.reader.QueryRowContext(ctx,
`SELECT COUNT(*) FROM dm_calls WHERE status <> ?`, DMCallEnded).Scan(&count)
return count, err
}
// GetDMCallRow возвращает звонок без состава.
func (s *Store) GetDMCallRow(ctx context.Context, callID uint64) (*DMCall, error) {
row := s.reader.QueryRowContext(ctx,
`SELECT `+dmCallColumns+` FROM dm_calls WHERE id = ?`, int64(callID))
return scanDMCall(row)
}
// listDMCallParticipants перечисляет состав звонка.
func (s *Store) listDMCallParticipants(ctx context.Context, callID uint64) ([]DMCallParticipant, error) {
rows, err := s.reader.QueryContext(ctx,
`SELECT `+dmCallParticipantColumns+` FROM dm_call_participants WHERE call_id = ? ORDER BY invited_at, user_id`,
int64(callID))
if err != nil {
return nil, err
}
defer rows.Close()
return scanDMCallParticipants(rows)
}
// listDMCallParticipantsRange перечисляет состав звонков страницы истории.
func (s *Store) listDMCallParticipantsRange(ctx context.Context, channelID, minCallID, beforeID uint64) ([]DMCallParticipant, error) {
rows, err := s.reader.QueryContext(ctx, `
SELECT p.call_id, p.user_id, p.state, p.invited_at, p.joined_at, p.left_at,
p.self_mute, p.self_deaf, p.camera, p.screen
FROM dm_call_participants p JOIN dm_calls c ON c.id = p.call_id
WHERE c.channel_id = ? AND c.id >= ? AND (? = 0 OR c.id < ?)
ORDER BY p.call_id DESC, p.invited_at, p.user_id`,
int64(channelID), int64(minCallID), int64(beforeID), int64(beforeID))
if err != nil {
return nil, err
}
defer rows.Close()
return scanDMCallParticipants(rows)
}
func scanDMCall(scanner interface{ Scan(...any) error }) (*DMCall, error) {
var (
call DMCall
createdAt string
startedAt, endedAt sql.NullString
endedBy sql.NullInt64
endReason sql.NullString
)
err := scanner.Scan(&call.ID, &call.ChannelID, &call.InitiatorID, &call.Status, &call.Media,
&createdAt, &startedAt, &endedAt, &endedBy, &endReason)
if err != nil {
return nil, mapError(err)
}
call.CreatedAt = parseTimestamp(createdAt)
if startedAt.Valid {
value := parseTimestamp(startedAt.String)
call.StartedAt = &value
}
if endedAt.Valid {
value := parseTimestamp(endedAt.String)
call.EndedAt = &value
}
call.EndedBy = optionalID(endedBy)
call.EndReason = endReason.String
return &call, nil
}
func scanDMCallParticipants(rows *sql.Rows) ([]DMCallParticipant, error) {
participants := make([]DMCallParticipant, 0, 4)
for rows.Next() {
var (
participant DMCallParticipant
invitedAt string
joinedAt, leftAt sql.NullString
selfMute, selfDeaf int
camera, screen int
)
if err := rows.Scan(&participant.CallID, &participant.UserID, &participant.State, &invitedAt,
&joinedAt, &leftAt, &selfMute, &selfDeaf, &camera, &screen); err != nil {
return nil, err
}
participant.InvitedAt = parseTimestamp(invitedAt)
if joinedAt.Valid {
value := parseTimestamp(joinedAt.String)
participant.JoinedAt = &value
}
if leftAt.Valid {
value := parseTimestamp(leftAt.String)
participant.LeftAt = &value
}
participant.SelfMute = selfMute == 1
participant.SelfDeaf = selfDeaf == 1
participant.Camera = camera == 1
participant.Screen = screen == 1
participants = append(participants, participant)
}
return participants, rows.Err()
}