2026-07-21 18:57:09 +08:00

839 lines
32 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package service_test
import (
"context"
"errors"
"fmt"
"strings"
"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
err error
}
type recordingIMUserMuter struct {
mu sync.Mutex
calls []bool
err error
}
func (m *recordingIMUserMuter) SetRoomGroupMemberMuted(_ context.Context, _ string, _ int64, muted bool) error {
m.mu.Lock()
defer m.mu.Unlock()
m.calls = append(m.calls, muted)
return m.err
}
func (m *recordingIMUserMuter) snapshot() []bool {
m.mu.Lock()
defer m.mu.Unlock()
return append([]bool(nil), m.calls...)
}
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 b.err
}
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{}
imMuter := &recordingIMUserMuter{}
svc := roomservice.New(roomservice.Config{
NodeID: "node-muted-rtc-test",
LeaseTTL: 10 * time.Second,
RankLimit: 20,
SnapshotEveryN: 1,
Clock: now,
RTCUserAudioBlocker: rtcBlocker,
IMUserMuter: imMuter,
}, 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)
}
if _, err := svc.MuteUser(ctx, &roomv1.MuteUserRequest{
Meta: rocketMeta(roomID, ownerID, "mute-user-guard"),
TargetUserId: mutedUserID,
Muted: true,
}); err != nil {
t.Fatalf("old mute command replay failed: %v", err)
}
if got := rtcBlocker.snapshot(); len(got) != 0 {
// MuteUser 只提交 Room Cell 与 outbox旧命令重放和新命令都不能把腾讯 RTC 慢调用带入同步请求。
t.Fatalf("room mute command must not call Tencent RTC before durable outbox processing: %+v", got)
}
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) != 1 || !directCalls[0] {
// MuteUser 的 RTC 投影全部异步;只有被篡改客户端触发 audio_started 时,服务端回调会立即 re-block。
t.Fatalf("unexpected pre-outbox RTC audio block sequence: %+v", directCalls)
}
if calls := imMuter.snapshot(); len(calls) != 0 {
// IM 禁言不能进入 MuteUser 同步路径;否则腾讯慢请求会拉长房间命令,且实现重新依赖外部服务可用性。
t.Fatalf("room mute command must not call Tencent IM before durable outbox processing: %+v", calls)
}
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)
}
}
imCalls := imMuter.snapshot()
if len(imCalls) != 4 {
t.Fatalf("each durable RoomUserMuted event must reach Tencent IM projection exactly once: %+v", imCalls)
}
for index, muted := range imCalls {
// IM 投影只由 outbox 触发;处理四条历史事件时当前 Room Cell 已解除禁言,
// 因此每次都必须写 MuteTime=0不能按旧事件载荷恢复永久禁言。
if muted {
t.Fatalf("stale IM mute outbox call %d rolled current state backward", index)
}
}
}
type blockingRTCUserAudioBlocker struct {
mu sync.Mutex
calls []bool
blocked bool
firstStarted chan struct{}
releaseFirst chan struct{}
}
func newBlockingRTCUserAudioBlocker() *blockingRTCUserAudioBlocker {
return &blockingRTCUserAudioBlocker{firstStarted: make(chan struct{}), releaseFirst: make(chan struct{})}
}
func (b *blockingRTCUserAudioBlocker) SetUserAudioBlockedByStrRoomID(ctx context.Context, _ string, _ int64, blocked bool) error {
b.mu.Lock()
b.calls = append(b.calls, blocked)
first := len(b.calls) == 1
b.mu.Unlock()
if first {
close(b.firstStarted)
select {
case <-ctx.Done():
return ctx.Err()
case <-b.releaseFirst:
}
}
b.mu.Lock()
b.blocked = blocked
b.mu.Unlock()
return nil
}
func (b *blockingRTCUserAudioBlocker) snapshot() ([]bool, bool) {
b.mu.Lock()
defer b.mu.Unlock()
return append([]bool(nil), b.calls...), b.blocked
}
func TestRTCUserAudioReconcileConvergesSlowMuteAndNewerUnmute(t *testing.T) {
ctx := context.Background()
repository := mysqltest.NewRepository(t)
now := &fixedRoomRocketClock{now: time.Date(2026, 7, 21, 9, 0, 0, 0, time.UTC)}
rtcBlocker := newBlockingRTCUserAudioBlocker()
svc := roomservice.New(roomservice.Config{
NodeID: "node-muted-rtc-race-test",
LeaseTTL: 10 * time.Second,
SnapshotEveryN: 1,
Clock: now,
RTCUserAudioBlocker: rtcBlocker,
}, router.NewMemoryDirectory(), repository, &rocketTestWallet{}, integration.NewNoopRoomEventPublisher(), integration.NewNoopOutboxPublisher())
roomID := "room-muted-rtc-race"
ownerID := int64(8491)
targetUserID := int64(8492)
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
joinRocketRoom(t, ctx, svc, roomID, targetUserID)
if _, err := svc.MuteUser(ctx, &roomv1.MuteUserRequest{
Meta: rocketMeta(roomID, ownerID, "slow-mute"),
TargetUserId: targetUserID,
Muted: true,
}); err != nil {
t.Fatalf("mute command failed: %v", err)
}
projectionDone := make(chan error, 1)
go func() {
projectionDone <- svc.ProcessPendingOutbox(ctx, roomservice.OutboxWorkerOptions{PublishTimeout: 3 * time.Second, BatchSize: 100})
}()
select {
case <-rtcBlocker.firstStarted:
case <-time.After(2 * time.Second):
t.Fatal("slow mute outbox never reached RTC blocker")
}
if _, err := svc.MuteUser(ctx, &roomv1.MuteUserRequest{
Meta: rocketMeta(roomID, ownerID, "fast-unmute"),
TargetUserId: targetUserID,
Muted: false,
}); err != nil {
t.Fatalf("newer unmute command failed: %v", err)
}
deadline := time.Now().Add(2 * time.Second)
for {
permission, err := svc.CheckSpeakPermission(ctx, &roomv1.CheckSpeakPermissionRequest{RoomId: roomID, UserId: targetUserID, AppCode: appcode.Default})
if err == nil && permission.GetAllowed() {
break
}
if time.Now().After(deadline) {
t.Fatalf("newer unmute did not commit while slow RTC mute was in flight: permission=%+v err=%v", permission, err)
}
time.Sleep(5 * time.Millisecond)
}
close(rtcBlocker.releaseFirst)
if err := <-projectionDone; err != nil {
t.Fatalf("slow mute outbox processing failed: %v", err)
}
calls, blocked := rtcBlocker.snapshot()
if blocked || len(calls) < 2 || !calls[0] || calls[len(calls)-1] {
t.Fatalf("RTC projection must finish at newer unmuted state: calls=%+v blocked=%t", calls, blocked)
}
}
func TestMuteProjectionFailuresDoNotBlockFactPublisherAndRemainRetryable(t *testing.T) {
ctx := context.Background()
repository := mysqltest.NewRepository(t)
now := &fixedRoomRocketClock{now: time.Date(2026, 7, 21, 10, 0, 0, 0, time.UTC)}
rtcBlocker := &recordingRTCUserAudioBlocker{err: errors.New("rtc unavailable")}
imMuter := &recordingIMUserMuter{err: errors.New("im unavailable")}
outboxPublisher := newRecordingRoomDirectIMPublisher()
svc := roomservice.New(roomservice.Config{
NodeID: "node-muted-rtc-outbox-test",
LeaseTTL: 10 * time.Second,
SnapshotEveryN: 1,
Clock: now,
RTCUserAudioBlocker: rtcBlocker,
IMUserMuter: imMuter,
}, router.NewMemoryDirectory(), repository, &rocketTestWallet{}, integration.NewNoopRoomEventPublisher(), outboxPublisher)
roomID := "room-muted-rtc-outbox"
ownerID := int64(8493)
targetUserID := int64(8494)
createRocketRoom(t, ctx, svc, roomID, ownerID, 9101)
joinRocketRoom(t, ctx, svc, roomID, targetUserID)
if _, err := svc.MuteUser(ctx, &roomv1.MuteUserRequest{
Meta: rocketMeta(roomID, ownerID, "mute-with-rtc-failure"),
TargetUserId: targetUserID,
Muted: true,
}); err != nil {
t.Fatalf("durable room mute must not depend on RTC availability: %v", err)
}
if calls := rtcBlocker.snapshot(); len(calls) != 0 {
t.Fatalf("room mute request must return before RTC projection: %+v", calls)
}
if calls := imMuter.snapshot(); len(calls) != 0 {
t.Fatalf("room mute request must return before IM projection: %+v", calls)
}
pending, err := repository.ListPendingOutbox(ctx, 100)
if err != nil {
t.Fatalf("list pending room outbox: %v", err)
}
var muteEventID string
for _, record := range pending {
if record.EventType == "RoomUserMuted" {
muteEventID = record.EventID
break
}
}
if muteEventID == "" {
t.Fatal("RoomUserMuted outbox event not found")
}
if err := svc.ProcessPendingOutbox(ctx, roomservice.OutboxWorkerOptions{
PublishTimeout: time.Second,
BatchSize: 100,
MaxRetryCount: 3,
InitialBackoff: time.Second,
MaxBackoff: time.Second,
}); err != nil {
t.Fatalf("process room outbox with RTC failure: %v", err)
}
if outboxPublisher.counts()["RoomUserMuted"] != 1 {
t.Fatalf("IM/RTC failures must not prevent generic RoomUserMuted publication: counts=%+v", outboxPublisher.counts())
}
if calls := rtcBlocker.snapshot(); len(calls) != 1 || !calls[0] {
t.Fatalf("RTC projection must attempt the latest muted state once: %+v", calls)
}
if calls := imMuter.snapshot(); len(calls) != 1 || !calls[0] {
t.Fatalf("IM projection must attempt the latest muted state once: %+v", calls)
}
record, exists := repository.OutboxRecord(muteEventID)
if !exists || record.Status != "retryable" || record.RetryCount != 1 {
t.Fatalf("failed cloud projections must keep durable event retryable: %+v exists=%t", record, exists)
}
if !strings.Contains(record.LastError, "im unavailable") || !strings.Contains(record.LastError, "rtc unavailable") {
t.Fatalf("retryable outbox must preserve both projection failures: %+v", record)
}
}
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
}