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() }