614 lines
24 KiB
Go
614 lines
24 KiB
Go
package service_test
|
||
|
||
import (
|
||
"context"
|
||
"fmt"
|
||
"sync"
|
||
"testing"
|
||
"time"
|
||
|
||
roomv1 "hyapp.local/api/proto/room/v1"
|
||
"hyapp/pkg/appcode"
|
||
"hyapp/pkg/xerr"
|
||
"hyapp/services/room-service/internal/integration"
|
||
"hyapp/services/room-service/internal/room/command"
|
||
roomservice "hyapp/services/room-service/internal/room/service"
|
||
"hyapp/services/room-service/internal/router"
|
||
"hyapp/services/room-service/internal/testutil/mysqltest"
|
||
)
|
||
|
||
type recordingRTCUserAudioBlocker struct {
|
||
mu sync.Mutex
|
||
calls []bool
|
||
}
|
||
|
||
func (b *recordingRTCUserAudioBlocker) SetUserAudioBlockedByStrRoomID(_ context.Context, _ string, _ int64, blocked bool) error {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
b.calls = append(b.calls, blocked)
|
||
return nil
|
||
}
|
||
|
||
func (b *recordingRTCUserAudioBlocker) snapshot() []bool {
|
||
b.mu.Lock()
|
||
defer b.mu.Unlock()
|
||
return append([]bool(nil), b.calls...)
|
||
}
|
||
|
||
func TestMutedUserCannotSpeakOrMicUpUntilUnmuted(t *testing.T) {
|
||
ctx := context.Background()
|
||
repository := mysqltest.NewRepository(t)
|
||
now := &fixedRoomRocketClock{now: time.Date(2026, 7, 21, 8, 0, 0, 0, time.UTC)}
|
||
rtcBlocker := &recordingRTCUserAudioBlocker{}
|
||
svc := roomservice.New(roomservice.Config{
|
||
NodeID: "node-muted-rtc-test",
|
||
LeaseTTL: 10 * time.Second,
|
||
RankLimit: 20,
|
||
SnapshotEveryN: 1,
|
||
Clock: now,
|
||
RTCUserAudioBlocker: rtcBlocker,
|
||
}, router.NewMemoryDirectory(), repository, &rocketTestWallet{}, integration.NewNoopRoomEventPublisher(), integration.NewNoopOutboxPublisher())
|
||
|
||
roomID := "room-muted-user-guard"
|
||
ownerID := int64(8481)
|
||
mutedUserID := int64(8482)
|
||
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
|
||
joinRocketRoom(t, ctx, svc, roomID, mutedUserID)
|
||
|
||
if _, err := svc.MuteUser(ctx, &roomv1.MuteUserRequest{
|
||
Meta: rocketMeta(roomID, ownerID, "mute-user-guard"),
|
||
TargetUserId: mutedUserID,
|
||
Muted: true,
|
||
}); err != nil {
|
||
t.Fatalf("mute user failed: %v", err)
|
||
}
|
||
|
||
mutedSpeak, err := svc.CheckSpeakPermission(ctx, &roomv1.CheckSpeakPermissionRequest{
|
||
RoomId: roomID,
|
||
UserId: mutedUserID,
|
||
AppCode: appcode.Default,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("check muted speak permission failed: %v", err)
|
||
}
|
||
if mutedSpeak.GetAllowed() || mutedSpeak.GetReason() != "user_muted" {
|
||
t.Fatalf("muted user must be rejected by public chat guard: %+v", mutedSpeak)
|
||
}
|
||
if _, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, mutedUserID, "muted-mic-up"),
|
||
SeatNo: 2,
|
||
}); !xerr.IsCode(err, xerr.PermissionDenied) {
|
||
t.Fatalf("muted user mic up must return permission denied: %v", err)
|
||
}
|
||
|
||
if _, err := svc.MuteUser(ctx, &roomv1.MuteUserRequest{
|
||
Meta: rocketMeta(roomID, ownerID, "unmute-user-guard"),
|
||
TargetUserId: mutedUserID,
|
||
Muted: false,
|
||
}); err != nil {
|
||
t.Fatalf("unmute user failed: %v", err)
|
||
}
|
||
|
||
unmutedSpeak, err := svc.CheckSpeakPermission(ctx, &roomv1.CheckSpeakPermissionRequest{
|
||
RoomId: roomID,
|
||
UserId: mutedUserID,
|
||
AppCode: appcode.Default,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("check unmuted speak permission failed: %v", err)
|
||
}
|
||
if !unmutedSpeak.GetAllowed() || unmutedSpeak.GetReason() != "" {
|
||
t.Fatalf("unmuted user must recover public chat permission: %+v", unmutedSpeak)
|
||
}
|
||
upResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, mutedUserID, "unmuted-mic-up"),
|
||
SeatNo: 2,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("unmuted user must recover mic up permission: %v", err)
|
||
}
|
||
|
||
mutedOnSeat, err := svc.MuteUser(ctx, &roomv1.MuteUserRequest{
|
||
Meta: rocketMeta(roomID, ownerID, "mute-current-rtc-seat"),
|
||
TargetUserId: mutedUserID,
|
||
Muted: true,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("mute current RTC seat failed: %v", err)
|
||
}
|
||
var currentSeat *roomv1.SeatState
|
||
for _, seat := range mutedOnSeat.GetRoom().GetMicSeats() {
|
||
if seat.GetUserId() == mutedUserID {
|
||
currentSeat = seat
|
||
break
|
||
}
|
||
}
|
||
if currentSeat == nil || !currentSeat.GetMicMuted() || currentSeat.GetSeatNo() != 2 {
|
||
t.Fatalf("room mute must project to the current RTC seat without removing it: %+v", currentSeat)
|
||
}
|
||
if _, err := svc.SetMicMute(ctx, &roomv1.SetMicMuteRequest{
|
||
Meta: rocketMeta(roomID, mutedUserID, "muted-user-open-mic"),
|
||
TargetUserId: mutedUserID,
|
||
Muted: false,
|
||
}); !xerr.IsCode(err, xerr.PermissionDenied) {
|
||
t.Fatalf("room-muted user must not reopen RTC mic: %v", err)
|
||
}
|
||
if _, err := svc.ConfirmMicPublishing(ctx, &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: rocketMeta(roomID, mutedUserID, "muted-publish-confirm"),
|
||
TargetUserId: mutedUserID,
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
RoomVersion: upResp.GetResult().GetRoomVersion(),
|
||
EventTimeMs: now.Now().UnixMilli(),
|
||
Source: "client",
|
||
}); !xerr.IsCode(err, xerr.PermissionDenied) {
|
||
t.Fatalf("room mute between MicUp and RTC confirmation must reject publishing: %v", err)
|
||
}
|
||
rtcEventResp, err := svc.ApplyRTCEvent(ctx, &roomv1.ApplyRTCEventRequest{
|
||
Meta: rocketMeta(roomID, mutedUserID, "muted-rtc-audio-started"),
|
||
TargetUserId: mutedUserID,
|
||
EventType: "audio_started",
|
||
EventTimeMs: now.Now().UnixMilli(),
|
||
Source: "tencent_rtc_callback",
|
||
ExternalEventId: "rtc-event-muted-audio-started",
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("muted RTC audio callback must be acknowledged after cloud re-block: %v", err)
|
||
}
|
||
if rtcEventResp.GetResult().GetApplied() {
|
||
t.Fatalf("muted RTC audio callback must not advance the seat to publishing: %+v", rtcEventResp)
|
||
}
|
||
|
||
if _, err := svc.MuteUser(ctx, &roomv1.MuteUserRequest{
|
||
Meta: rocketMeta(roomID, ownerID, "unmute-current-rtc-seat"),
|
||
TargetUserId: mutedUserID,
|
||
Muted: false,
|
||
}); err != nil {
|
||
t.Fatalf("unmute current RTC seat failed: %v", err)
|
||
}
|
||
if _, err := svc.SetMicMute(ctx, &roomv1.SetMicMuteRequest{
|
||
Meta: rocketMeta(roomID, mutedUserID, "unmuted-user-open-mic"),
|
||
TargetUserId: mutedUserID,
|
||
Muted: false,
|
||
}); err != nil {
|
||
t.Fatalf("room unmute must restore explicit RTC mic control: %v", err)
|
||
}
|
||
|
||
directCalls := rtcBlocker.snapshot()
|
||
if len(directCalls) != 5 || !directCalls[0] || directCalls[1] || !directCalls[2] || !directCalls[3] || directCalls[4] {
|
||
// 四次 MuteUser 分别投影 true/false/true/false;中间被篡改客户端触发的
|
||
// audio_started 必须额外 re-block=true,因此按调用顺序断言服务端防线完整。
|
||
t.Fatalf("unexpected direct RTC audio block sequence: %+v", directCalls)
|
||
}
|
||
|
||
if err := svc.ProcessPendingOutbox(ctx, roomservice.OutboxWorkerOptions{PublishTimeout: time.Second, BatchSize: 100}); err != nil {
|
||
t.Fatalf("process RTC mute outbox failed: %v", err)
|
||
}
|
||
allCalls := rtcBlocker.snapshot()
|
||
if len(allCalls)-len(directCalls) != 4 {
|
||
t.Fatalf("each durable RoomUserMuted event must reach RTC compensation: direct=%+v all=%+v", directCalls, allCalls)
|
||
}
|
||
for index, blocked := range allCalls[len(directCalls):] {
|
||
// outbox 中既有 mute 也有 unmute 历史事件;补偿必须读取当前已解除的 Room Cell
|
||
// 状态并全部收敛到 false,不能按旧事件载荷重新把用户禁言。
|
||
if blocked {
|
||
t.Fatalf("stale RTC mute outbox call %d rolled current state backward", index)
|
||
}
|
||
}
|
||
}
|
||
|
||
func TestMicUpReturnsExistingSeatForDuplicateUserCommand(t *testing.T) {
|
||
ctx := context.Background()
|
||
repository := mysqltest.NewRepository(t)
|
||
now := &fixedRoomRocketClock{now: time.Date(2026, 6, 4, 8, 0, 0, 0, time.UTC)}
|
||
svc := newRocketTestService(t, repository, &rocketTestWallet{}, now)
|
||
|
||
roomID := "room-mic-up-duplicate-user"
|
||
ownerID := int64(8501)
|
||
speakerID := int64(8502)
|
||
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
|
||
joinRocketRoom(t, ctx, svc, roomID, speakerID)
|
||
|
||
upResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-up-duplicate-user-first"),
|
||
SeatNo: 2,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("first mic up failed: %v", err)
|
||
}
|
||
|
||
duplicateResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-up-duplicate-user-second"),
|
||
SeatNo: 3,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("duplicate mic up must return current session instead of conflict: %v", err)
|
||
}
|
||
if duplicateResp.GetResult().GetApplied() {
|
||
t.Fatalf("duplicate mic up must be a no-op mutation: %+v", duplicateResp.GetResult())
|
||
}
|
||
if duplicateResp.GetRoom().GetVersion() != upResp.GetRoom().GetVersion() {
|
||
t.Fatalf("duplicate mic up must not advance room version: got=%d want=%d", duplicateResp.GetRoom().GetVersion(), upResp.GetRoom().GetVersion())
|
||
}
|
||
if duplicateResp.GetSeatNo() != 2 || duplicateResp.GetMicSessionId() != upResp.GetMicSessionId() || duplicateResp.GetPublishDeadlineMs() != upResp.GetPublishDeadlineMs() {
|
||
t.Fatalf("duplicate mic up must return original mic session: duplicate=%+v first=%+v", duplicateResp, upResp)
|
||
}
|
||
if originalSeat := seatByNo(duplicateResp.GetRoom(), 2); originalSeat == nil || originalSeat.GetUserId() != speakerID || originalSeat.GetMicSessionId() != upResp.GetMicSessionId() {
|
||
t.Fatalf("original seat must remain occupied by speaker: %+v", originalSeat)
|
||
}
|
||
if requestedSeat := seatByNo(duplicateResp.GetRoom(), 3); requestedSeat == nil || requestedSeat.GetUserId() != 0 {
|
||
// 重复 MicUp 只能表达“我已经在麦上”的幂等结果;换到 seat 3 必须走 ChangeMicSeat。
|
||
t.Fatalf("duplicate mic up must not move user to requested seat: %+v", requestedSeat)
|
||
}
|
||
}
|
||
|
||
func TestConfirmMicPublishingAcceptsFastClientClockBeforeServerDeadline(t *testing.T) {
|
||
ctx := context.Background()
|
||
repository := mysqltest.NewRepository(t)
|
||
now := &fixedRoomRocketClock{now: time.Date(2026, 6, 4, 8, 0, 0, 0, time.UTC)}
|
||
svc := newRocketTestService(t, repository, &rocketTestWallet{}, now)
|
||
|
||
roomID := "room-mic-confirm-fast-client-clock"
|
||
ownerID := int64(8551)
|
||
speakerID := int64(8552)
|
||
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
|
||
joinRocketRoom(t, ctx, svc, roomID, speakerID)
|
||
|
||
upResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-confirm-fast-clock-up"),
|
||
SeatNo: 2,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("mic up failed: %v", err)
|
||
}
|
||
|
||
arrivedAtMS := upResp.GetPublishDeadlineMs() - 1
|
||
fastClientEventMS := upResp.GetPublishDeadlineMs() + int64((5 * time.Second).Milliseconds())
|
||
now.now = time.UnixMilli(arrivedAtMS)
|
||
confirmMeta := rocketMeta(roomID, speakerID, "mic-confirm-fast-clock-confirm")
|
||
confirmResp, err := svc.ConfirmMicPublishing(ctx, &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: confirmMeta,
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
RoomVersion: upResp.GetRoom().GetVersion(),
|
||
EventTimeMs: fastClientEventMS,
|
||
Source: "client",
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("confirm mic publishing failed: %v", err)
|
||
}
|
||
if !confirmResp.GetResult().GetApplied() {
|
||
t.Fatalf("fast client event_time_ms must not be treated as timeout before server deadline: %+v", confirmResp.GetResult())
|
||
}
|
||
seat := seatByNo(confirmResp.GetRoom(), 2)
|
||
if seat == nil || seat.GetPublishState() != "publishing" {
|
||
t.Fatalf("seat must be confirmed as publishing: %+v", seat)
|
||
}
|
||
if seat.GetLastPublishEventTimeMs() != arrivedAtMS || seat.GetMicHeartbeatAtMs() != arrivedAtMS {
|
||
t.Fatalf("confirmation must store server accepted time, not fast client event time: seat=%+v accepted=%d client_event=%d", seat, arrivedAtMS, fastClientEventMS)
|
||
}
|
||
|
||
record, exists, err := repository.GetCommand(appcode.WithContext(ctx, appcode.Default), roomID, confirmMeta.GetCommandId())
|
||
if err != nil {
|
||
t.Fatalf("get confirm command log failed: %v", err)
|
||
}
|
||
if !exists {
|
||
t.Fatalf("confirm command log must be persisted")
|
||
}
|
||
decoded, err := command.Deserialize(command.ConfirmMicPublishing{}.Type(), record.Payload)
|
||
if err != nil {
|
||
t.Fatalf("decode confirm command log failed: %v", err)
|
||
}
|
||
confirmCmd, ok := decoded.(*command.ConfirmMicPublishing)
|
||
if !ok {
|
||
t.Fatalf("confirm command log has unexpected type: %T", decoded)
|
||
}
|
||
if confirmCmd.EventTimeMS != fastClientEventMS || confirmCmd.AcceptedEventTimeMS != arrivedAtMS {
|
||
t.Fatalf("confirm command log must keep raw client time and accepted server time: %+v raw=%d accepted=%d", confirmCmd, fastClientEventMS, arrivedAtMS)
|
||
}
|
||
}
|
||
|
||
func TestConfirmMicPublishingIgnoresStaleSessionVersionAndTimeout(t *testing.T) {
|
||
tests := []struct {
|
||
name string
|
||
adjust func(clock *fixedRoomRocketClock, upResp *roomv1.MicUpResponse)
|
||
request func(roomID string, speakerID int64, clock *fixedRoomRocketClock, upResp *roomv1.MicUpResponse) *roomv1.ConfirmMicPublishingRequest
|
||
}{
|
||
{
|
||
name: "old_session",
|
||
adjust: func(clock *fixedRoomRocketClock, _ *roomv1.MicUpResponse) {
|
||
clock.now = clock.now.Add(500 * time.Millisecond)
|
||
},
|
||
request: func(roomID string, speakerID int64, clock *fixedRoomRocketClock, upResp *roomv1.MicUpResponse) *roomv1.ConfirmMicPublishingRequest {
|
||
return &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-confirm-old-session"),
|
||
MicSessionId: upResp.GetMicSessionId() + "-old",
|
||
RoomVersion: upResp.GetRoom().GetVersion(),
|
||
EventTimeMs: clock.Now().UnixMilli(),
|
||
Source: "client",
|
||
}
|
||
},
|
||
},
|
||
{
|
||
name: "old_room_version",
|
||
adjust: func(clock *fixedRoomRocketClock, _ *roomv1.MicUpResponse) {
|
||
clock.now = clock.now.Add(500 * time.Millisecond)
|
||
},
|
||
request: func(roomID string, speakerID int64, clock *fixedRoomRocketClock, upResp *roomv1.MicUpResponse) *roomv1.ConfirmMicPublishingRequest {
|
||
return &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-confirm-old-room-version"),
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
RoomVersion: upResp.GetRoom().GetVersion() - 1,
|
||
EventTimeMs: clock.Now().UnixMilli(),
|
||
Source: "client",
|
||
}
|
||
},
|
||
},
|
||
{
|
||
name: "deadline_reached",
|
||
adjust: func(clock *fixedRoomRocketClock, upResp *roomv1.MicUpResponse) {
|
||
clock.now = time.UnixMilli(upResp.GetPublishDeadlineMs())
|
||
},
|
||
request: func(roomID string, speakerID int64, _ *fixedRoomRocketClock, upResp *roomv1.MicUpResponse) *roomv1.ConfirmMicPublishingRequest {
|
||
return &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-confirm-deadline-reached"),
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
RoomVersion: upResp.GetRoom().GetVersion(),
|
||
EventTimeMs: upResp.GetPublishDeadlineMs() - 1,
|
||
Source: "client",
|
||
}
|
||
},
|
||
},
|
||
}
|
||
|
||
for index, tt := range tests {
|
||
t.Run(tt.name, func(t *testing.T) {
|
||
ctx := context.Background()
|
||
repository := mysqltest.NewRepository(t)
|
||
now := &fixedRoomRocketClock{now: time.Date(2026, 6, 4, 8, 0, 0, 0, time.UTC)}
|
||
svc := newRocketTestService(t, repository, &rocketTestWallet{}, now)
|
||
|
||
roomID := fmt.Sprintf("room-mic-confirm-guard-%d", index)
|
||
ownerID := int64(8561 + index*10)
|
||
speakerID := int64(8562 + index*10)
|
||
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
|
||
joinRocketRoom(t, ctx, svc, roomID, speakerID)
|
||
|
||
upResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, speakerID, fmt.Sprintf("mic-confirm-guard-up-%d", index)),
|
||
SeatNo: 2,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("mic up failed: %v", err)
|
||
}
|
||
tt.adjust(now, upResp)
|
||
|
||
confirmResp, err := svc.ConfirmMicPublishing(ctx, tt.request(roomID, speakerID, now, upResp))
|
||
if err != nil {
|
||
t.Fatalf("stale confirm must be ignored without error: %v", err)
|
||
}
|
||
if confirmResp.GetResult().GetApplied() || confirmResp.GetRoom().GetVersion() != upResp.GetRoom().GetVersion() {
|
||
t.Fatalf("stale confirm must not advance room state: confirm=%+v up_version=%d", confirmResp.GetResult(), upResp.GetRoom().GetVersion())
|
||
}
|
||
seat := seatByNo(confirmResp.GetRoom(), 2)
|
||
if seat == nil || seat.GetPublishState() != "pending_publish" || seat.GetLastPublishEventTimeMs() != 0 || seat.GetMicHeartbeatAtMs() != 0 {
|
||
t.Fatalf("stale confirm must leave pending session untouched: %+v", seat)
|
||
}
|
||
})
|
||
}
|
||
}
|
||
|
||
func TestMicHeartbeatRefreshesPublishingSessionAndPresence(t *testing.T) {
|
||
ctx := context.Background()
|
||
repository := mysqltest.NewRepository(t)
|
||
now := &fixedRoomRocketClock{now: time.Date(2026, 6, 4, 8, 0, 0, 0, time.UTC)}
|
||
svc := newRocketTestService(t, repository, &rocketTestWallet{}, now)
|
||
|
||
roomID := "room-mic-heartbeat"
|
||
ownerID := int64(8601)
|
||
speakerID := int64(8602)
|
||
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
|
||
joinRocketRoom(t, ctx, svc, roomID, speakerID)
|
||
|
||
upResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-heartbeat-up"),
|
||
SeatNo: 2,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("mic up failed: %v", err)
|
||
}
|
||
now.now = now.now.Add(500 * time.Millisecond)
|
||
confirmResp, err := svc.ConfirmMicPublishing(ctx, &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-heartbeat-confirm"),
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
RoomVersion: upResp.GetRoom().GetVersion(),
|
||
EventTimeMs: now.Now().UnixMilli(),
|
||
Source: "client",
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("confirm mic publishing failed: %v", err)
|
||
}
|
||
confirmedSeat := seatByNo(confirmResp.GetRoom(), 2)
|
||
if confirmedSeat == nil || confirmedSeat.GetPublishState() != "publishing" || confirmedSeat.GetMicHeartbeatAtMs() == 0 {
|
||
t.Fatalf("publish confirmation must initialize heartbeat: %+v", confirmedSeat)
|
||
}
|
||
|
||
now.now = now.now.Add(3 * time.Second)
|
||
heartbeatAtMS := now.Now().UnixMilli()
|
||
heartbeatResp, err := svc.MicHeartbeat(ctx, &roomv1.MicHeartbeatRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-heartbeat-refresh"),
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("mic heartbeat failed: %v", err)
|
||
}
|
||
if heartbeatResp.GetSeatNo() != 2 || heartbeatResp.GetMicHeartbeatAtMs() != heartbeatAtMS {
|
||
t.Fatalf("heartbeat response mismatch: %+v want_at=%d", heartbeatResp, heartbeatAtMS)
|
||
}
|
||
seat := seatByNo(heartbeatResp.GetRoom(), 2)
|
||
if seat == nil || seat.GetMicHeartbeatAtMs() != heartbeatAtMS {
|
||
t.Fatalf("seat heartbeat must be visible in snapshot: %+v", seat)
|
||
}
|
||
user := onlineUserByID(heartbeatResp.GetRoom(), speakerID)
|
||
if user == nil || user.GetLastSeenAtMs() != heartbeatAtMS {
|
||
t.Fatalf("mic heartbeat must refresh room presence: %+v", user)
|
||
}
|
||
}
|
||
|
||
func TestConfirmMicPublishingIgnoresBackwardEventTime(t *testing.T) {
|
||
ctx := context.Background()
|
||
repository := mysqltest.NewRepository(t)
|
||
now := &fixedRoomRocketClock{now: time.Date(2026, 6, 4, 8, 0, 0, 0, time.UTC)}
|
||
svc := newRocketTestService(t, repository, &rocketTestWallet{}, now)
|
||
|
||
roomID := "room-mic-confirm-backward-event"
|
||
ownerID := int64(8651)
|
||
speakerID := int64(8652)
|
||
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
|
||
joinRocketRoom(t, ctx, svc, roomID, speakerID)
|
||
|
||
upResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-confirm-backward-up"),
|
||
SeatNo: 2,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("mic up failed: %v", err)
|
||
}
|
||
now.now = now.now.Add(time.Second)
|
||
acceptedEventMS := now.Now().UnixMilli()
|
||
confirmResp, err := svc.ConfirmMicPublishing(ctx, &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-confirm-backward-first"),
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
RoomVersion: upResp.GetRoom().GetVersion(),
|
||
EventTimeMs: acceptedEventMS,
|
||
Source: "client",
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("first confirm failed: %v", err)
|
||
}
|
||
if !confirmResp.GetResult().GetApplied() {
|
||
t.Fatalf("first confirm must apply: %+v", confirmResp.GetResult())
|
||
}
|
||
|
||
now.now = now.now.Add(time.Second)
|
||
backwardResp, err := svc.ConfirmMicPublishing(ctx, &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-confirm-backward-retry"),
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
RoomVersion: confirmResp.GetRoom().GetVersion(),
|
||
EventTimeMs: acceptedEventMS - 1,
|
||
Source: "client",
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("backward confirm must be ignored without error: %v", err)
|
||
}
|
||
if backwardResp.GetResult().GetApplied() || backwardResp.GetRoom().GetVersion() != confirmResp.GetRoom().GetVersion() {
|
||
t.Fatalf("backward event must not advance room state: backward=%+v confirmed_version=%d", backwardResp.GetResult(), confirmResp.GetRoom().GetVersion())
|
||
}
|
||
seat := seatByNo(backwardResp.GetRoom(), 2)
|
||
if seat == nil || seat.GetPublishState() != "publishing" || seat.GetLastPublishEventTimeMs() != acceptedEventMS {
|
||
t.Fatalf("backward event must not lower accepted event time: %+v want_event=%d", seat, acceptedEventMS)
|
||
}
|
||
}
|
||
|
||
func TestRTCAudioStoppedReleasesAfterFastClientConfirmWatermark(t *testing.T) {
|
||
ctx := context.Background()
|
||
repository := mysqltest.NewRepository(t)
|
||
now := &fixedRoomRocketClock{now: time.Date(2026, 6, 4, 8, 0, 0, 0, time.UTC)}
|
||
svc := newRocketTestService(t, repository, &rocketTestWallet{}, now)
|
||
|
||
roomID := "room-rtc-stop-after-fast-confirm"
|
||
ownerID := int64(8661)
|
||
speakerID := int64(8662)
|
||
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
|
||
joinRocketRoom(t, ctx, svc, roomID, speakerID)
|
||
|
||
upResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "rtc-stop-after-fast-confirm-up"),
|
||
SeatNo: 2,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("mic up failed: %v", err)
|
||
}
|
||
createdAtMS := upResp.GetPublishDeadlineMs() - int64((15 * time.Second).Milliseconds())
|
||
stopEventMS := createdAtMS + int64((2 * time.Second).Milliseconds())
|
||
acceptedAtMS := upResp.GetPublishDeadlineMs() - 1
|
||
now.now = time.UnixMilli(acceptedAtMS)
|
||
confirmResp, err := svc.ConfirmMicPublishing(ctx, &roomv1.ConfirmMicPublishingRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "rtc-stop-after-fast-confirm-confirm"),
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
RoomVersion: upResp.GetRoom().GetVersion(),
|
||
EventTimeMs: upResp.GetPublishDeadlineMs() + 5000,
|
||
Source: "client",
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("confirm mic publishing failed: %v", err)
|
||
}
|
||
confirmedSeat := seatByNo(confirmResp.GetRoom(), 2)
|
||
if confirmedSeat == nil || confirmedSeat.GetLastPublishEventTimeMs() != acceptedAtMS {
|
||
t.Fatalf("confirm must store accepted watermark before RTC stop test: %+v accepted=%d", confirmedSeat, acceptedAtMS)
|
||
}
|
||
|
||
now.now = now.now.Add(time.Second)
|
||
stopResp, err := svc.ApplyRTCEvent(ctx, &roomv1.ApplyRTCEventRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "rtc-stop-after-fast-confirm-stop"),
|
||
TargetUserId: speakerID,
|
||
EventType: "audio_stopped",
|
||
EventTimeMs: stopEventMS,
|
||
Reason: "rtc_audio_stopped:0",
|
||
Source: "tencent_rtc_callback",
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("rtc audio stopped failed: %v", err)
|
||
}
|
||
if !stopResp.GetResult().GetApplied() {
|
||
t.Fatalf("rtc stop must release current mic session even when provider event time is below client confirm accepted time: %+v", stopResp.GetResult())
|
||
}
|
||
seat := seatByNo(stopResp.GetRoom(), 2)
|
||
if seat == nil || seat.GetUserId() != 0 || seat.GetMicSessionId() != "" {
|
||
t.Fatalf("rtc stop must clear current mic seat: %+v", seat)
|
||
}
|
||
}
|
||
|
||
func TestMicHeartbeatRejectsPendingSession(t *testing.T) {
|
||
ctx := context.Background()
|
||
repository := mysqltest.NewRepository(t)
|
||
now := &fixedRoomRocketClock{now: time.Date(2026, 6, 4, 8, 0, 0, 0, time.UTC)}
|
||
svc := newRocketTestService(t, repository, &rocketTestWallet{}, now)
|
||
|
||
roomID := "room-mic-heartbeat-pending"
|
||
ownerID := int64(8701)
|
||
speakerID := int64(8702)
|
||
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
|
||
joinRocketRoom(t, ctx, svc, roomID, speakerID)
|
||
|
||
upResp, err := svc.MicUp(ctx, &roomv1.MicUpRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-heartbeat-pending-up"),
|
||
SeatNo: 1,
|
||
})
|
||
if err != nil {
|
||
t.Fatalf("mic up failed: %v", err)
|
||
}
|
||
if _, err := svc.MicHeartbeat(ctx, &roomv1.MicHeartbeatRequest{
|
||
Meta: rocketMeta(roomID, speakerID, "mic-heartbeat-pending-refresh"),
|
||
MicSessionId: upResp.GetMicSessionId(),
|
||
}); err == nil {
|
||
t.Fatalf("pending_publish session must not accept mic heartbeat before publish confirmation")
|
||
}
|
||
}
|
||
|
||
func seatByNo(snapshot *roomv1.RoomSnapshot, seatNo int32) *roomv1.SeatState {
|
||
for _, seat := range snapshot.GetMicSeats() {
|
||
if seat.GetSeatNo() == seatNo {
|
||
return seat
|
||
}
|
||
}
|
||
return nil
|
||
}
|
||
|
||
func onlineUserByID(snapshot *roomv1.RoomSnapshot, userID int64) *roomv1.RoomUser {
|
||
for _, user := range snapshot.GetOnlineUsers() {
|
||
if user.GetUserId() == userID {
|
||
return user
|
||
}
|
||
}
|
||
return nil
|
||
}
|