refactor(realtime): 重构实时会议Socket会话与DTO
- 重构 OpenRealtimeSocketSessionCommand,将集合类型替换为明确字段 - 完善 RealtimeMeetingSocketSessionServiceImpl 逻辑及 WebSocket 配置 - 清理前端 MeetingCreateDrawer 中未使用的 Antd 组件导入 - 重构 HotWordServiceImplTest 测试用例并优化 logback 配置格式dev_na
parent
36c1461c81
commit
58d647d8c8
|
|
@ -591,9 +591,9 @@ public class AndroidMeetingController {
|
||||||
|
|
||||||
private Meeting requireOperableOfflineMeeting(Long meetingId, AndroidAuthContext authContext, LoginUser loginUser) {
|
private Meeting requireOperableOfflineMeeting(Long meetingId, AndroidAuthContext authContext, LoginUser loginUser) {
|
||||||
Meeting meeting = meetingAccessService.requireMeetingIgnoreTenant(meetingId);
|
Meeting meeting = meetingAccessService.requireMeetingIgnoreTenant(meetingId);
|
||||||
if (!MeetingConstants.TYPE_OFFLINE.equals(meeting.getMeetingType())) {
|
// if (!MeetingConstants.TYPE_OFFLINE.equals(meeting.getMeetingType())) {
|
||||||
throw new RuntimeException("当前会议不是离线会议");
|
// throw new RuntimeException("当前会议不是离线会议");
|
||||||
}
|
// }
|
||||||
if (authContext == null || authContext.getDeviceId() == null || authContext.getDeviceId().isBlank()) {
|
if (authContext == null || authContext.getDeviceId() == null || authContext.getDeviceId().isBlank()) {
|
||||||
throw new RuntimeException("设备ID不能为空");
|
throw new RuntimeException("设备ID不能为空");
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -98,7 +98,9 @@ public class AndroidMeetingRealtimeController {
|
||||||
authContext.getTenantId(),
|
authContext.getTenantId(),
|
||||||
authContext.getUserId(),
|
authContext.getUserId(),
|
||||||
resolveCreatorName(authContext),
|
resolveCreatorName(authContext),
|
||||||
MeetingTerminalEnum.CUSTOM_TERMINAL.getCode()
|
MeetingTerminalEnum.CUSTOM_TERMINAL.getCode(),
|
||||||
|
authContext.getDeviceId(),
|
||||||
|
resolveSourceDeviceMode(authContext)
|
||||||
);
|
);
|
||||||
|
|
||||||
RealtimeMeetingSessionStatusVO status = realtimeMeetingSessionStateService.getStatus(meeting.getId());
|
RealtimeMeetingSessionStatusVO status = realtimeMeetingSessionStateService.getStatus(meeting.getId());
|
||||||
|
|
@ -274,6 +276,12 @@ public class AndroidMeetingRealtimeController {
|
||||||
: "android:" + authContext.getDeviceId().trim();
|
: "android:" + authContext.getDeviceId().trim();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private String resolveSourceDeviceMode(AndroidAuthContext authContext) {
|
||||||
|
return authContext.isAnonymous()
|
||||||
|
? MeetingConstants.DEVICE_MODE_PUBLIC
|
||||||
|
: MeetingConstants.DEVICE_MODE_PRIVATE;
|
||||||
|
}
|
||||||
|
|
||||||
private String normalize(String value) {
|
private String normalize(String value) {
|
||||||
return normalize(value, null);
|
return normalize(value, null);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,8 @@ import lombok.Getter;
|
||||||
public enum MeetingPushTypeEnum {
|
public enum MeetingPushTypeEnum {
|
||||||
PUBLIC_MEETING_LOGIN_CONFIRM("PUBLIC_MEETING_LOGIN_CONFIRM", "公有设备扫码登录确认消息"),
|
PUBLIC_MEETING_LOGIN_CONFIRM("PUBLIC_MEETING_LOGIN_CONFIRM", "公有设备扫码登录确认消息"),
|
||||||
MEETING_PENDING("MEETING_PENDING", "待开始会议通知"),
|
MEETING_PENDING("MEETING_PENDING", "待开始会议通知"),
|
||||||
MEETING_STATUS_CHANGED("MEETING_STATUS_CHANGED", "会议状态变更通知");
|
MEETING_STATUS_CHANGED("MEETING_STATUS_CHANGED", "会议状态变更通知"),
|
||||||
|
MEETING_REALTIME_TRANSCRIPT("MEETING_REALTIME_TRANSCRIPT", "实时会议转写推送");
|
||||||
|
|
||||||
private final String code;
|
private final String code;
|
||||||
private final String desc;
|
private final String desc;
|
||||||
|
|
|
||||||
|
|
@ -17,24 +17,19 @@ import com.imeeting.service.realtime.RealtimeAsrChannel;
|
||||||
import com.imeeting.service.realtime.RealtimeAsrChannelContext;
|
import com.imeeting.service.realtime.RealtimeAsrChannelContext;
|
||||||
import com.imeeting.service.realtime.RealtimeAsrChannelFactory;
|
import com.imeeting.service.realtime.RealtimeAsrChannelFactory;
|
||||||
import com.imeeting.service.realtime.RealtimeMeetingAudioStorageService;
|
import com.imeeting.service.realtime.RealtimeMeetingAudioStorageService;
|
||||||
import com.imeeting.service.realtime.RealtimeMeetingTranscriptCacheService;
|
|
||||||
import com.imeeting.service.realtime.RealtimeMeetingEventHub;
|
|
||||||
import com.imeeting.service.realtime.RealtimeMeetingEventRelay;
|
import com.imeeting.service.realtime.RealtimeMeetingEventRelay;
|
||||||
import com.imeeting.service.realtime.RealtimeMeetingSubscriptionQuota;
|
|
||||||
import com.imeeting.service.realtime.RealtimeMeetingMetrics;
|
import com.imeeting.service.realtime.RealtimeMeetingMetrics;
|
||||||
|
import com.imeeting.service.realtime.RealtimeMeetingPushSubscriptionManager;
|
||||||
import com.imeeting.support.redis.RealtimeMeetingPublisherLease;
|
import com.imeeting.support.redis.RealtimeMeetingPublisherLease;
|
||||||
import io.grpc.Status;
|
import io.grpc.Status;
|
||||||
import io.grpc.stub.StreamObserver;
|
import io.grpc.stub.StreamObserver;
|
||||||
import lombok.RequiredArgsConstructor;
|
import lombok.RequiredArgsConstructor;
|
||||||
import lombok.extern.slf4j.Slf4j;
|
import lombok.extern.slf4j.Slf4j;
|
||||||
import org.springframework.stereotype.Service;
|
import org.springframework.stereotype.Service;
|
||||||
import org.springframework.scheduling.annotation.Scheduled;
|
|
||||||
import org.springframework.beans.factory.annotation.Value;
|
|
||||||
import org.springframework.web.socket.CloseStatus;
|
import org.springframework.web.socket.CloseStatus;
|
||||||
|
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.concurrent.ConcurrentHashMap;
|
import java.util.concurrent.ConcurrentHashMap;
|
||||||
import java.util.function.Consumer;
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Android 实时会议音频双向流。会议完成仍由 REST complete 接口负责。
|
* Android 实时会议音频双向流。会议完成仍由 REST complete 接口负责。
|
||||||
|
|
@ -55,25 +50,11 @@ public class AndroidRealtimeMeetingGrpcService
|
||||||
private final RealtimeMeetingAudioStorageService audioStorageService;
|
private final RealtimeMeetingAudioStorageService audioStorageService;
|
||||||
private final AiModelService aiModelService;
|
private final AiModelService aiModelService;
|
||||||
private final RealtimeMeetingPublisherLease publisherLease;
|
private final RealtimeMeetingPublisherLease publisherLease;
|
||||||
private final RealtimeMeetingEventHub eventHub;
|
|
||||||
private final RealtimeMeetingEventRelay eventRelay;
|
private final RealtimeMeetingEventRelay eventRelay;
|
||||||
private final RealtimeMeetingSubscriptionQuota subscriptionQuota;
|
private final RealtimeMeetingPushSubscriptionManager pushSubscriptionManager;
|
||||||
private final RealtimeMeetingTranscriptCacheService transcriptCacheService;
|
|
||||||
private final Map<Long, GrpcSession> sessions = new ConcurrentHashMap<>();
|
private final Map<Long, GrpcSession> sessions = new ConcurrentHashMap<>();
|
||||||
private final Map<String, GrpcSession> subscriberSessions = new ConcurrentHashMap<>();
|
private final Map<String, GrpcSession> subscriberSessions = new ConcurrentHashMap<>();
|
||||||
|
|
||||||
@Value("${imeeting.realtime.max-subscribers-per-meeting:200}")
|
|
||||||
private int maxSubscribersPerMeeting = 200;
|
|
||||||
|
|
||||||
@Scheduled(fixedDelay = 60_000)
|
|
||||||
public void renewSubscriberQuotaSlots() {
|
|
||||||
subscriberSessions.values().forEach(session -> {
|
|
||||||
if (!session.cleaned) {
|
|
||||||
subscriptionQuota.renew(session.meetingId, session.subscriptionMember);
|
|
||||||
}
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public StreamObserver<ClientMessage> stream(StreamObserver<ServerMessage> responseObserver) {
|
public StreamObserver<ClientMessage> stream(StreamObserver<ServerMessage> responseObserver) {
|
||||||
return new StreamObserver<>() {
|
return new StreamObserver<>() {
|
||||||
|
|
@ -285,14 +266,9 @@ public class AndroidRealtimeMeetingGrpcService
|
||||||
}
|
}
|
||||||
|
|
||||||
private void detachSubscriber(GrpcSession target) {
|
private void detachSubscriber(GrpcSession target) {
|
||||||
if (target.eventConsumer != null) {
|
if (target.pushSubscribed) {
|
||||||
eventHub.unsubscribe(target.meetingId, target.eventConsumer);
|
target.pushSubscribed = false;
|
||||||
target.eventConsumer = null;
|
pushSubscriptionManager.unsubscribe(target.connectionId);
|
||||||
}
|
|
||||||
if (target.subscriptionMember != null) {
|
|
||||||
subscriptionQuota.release(target.meetingId, target.subscriptionMember);
|
|
||||||
target.subscriptionMember = null;
|
|
||||||
RealtimeMeetingMetrics.subscriberClosed("grpc");
|
|
||||||
}
|
}
|
||||||
subscriberSessions.remove(target.connectionId, target);
|
subscriberSessions.remove(target.connectionId, target);
|
||||||
}
|
}
|
||||||
|
|
@ -349,28 +325,25 @@ public class AndroidRealtimeMeetingGrpcService
|
||||||
GrpcSession next = new GrpcSession(meetingId, connectionId, null, observer, null);
|
GrpcSession next = new GrpcSession(meetingId, connectionId, null, observer, null);
|
||||||
next.publisher = false;
|
next.publisher = false;
|
||||||
next.canPublish = canPublish;
|
next.canPublish = canPublish;
|
||||||
next.subscriptionMember = subscriptionQuota.member(connectionId);
|
String subscriberDeviceId = connect.getDeviceId().trim();
|
||||||
if (!subscriptionQuota.tryAcquire(meetingId, next.subscriptionMember)) {
|
String rejectReason = pushSubscriptionManager.subscribe(
|
||||||
RealtimeMeetingMetrics.subscriberRejected("grpc");
|
meetingId, connectionId, subscriberDeviceId, meeting.getTitle());
|
||||||
throw new IllegalStateException("会议实时订阅人数已达上限");
|
if (rejectReason != null) {
|
||||||
}
|
throw new IllegalStateException(rejectReason);
|
||||||
next.eventConsumer = next::sendTranscript;
|
|
||||||
if (!eventHub.trySubscribe(meetingId, next.eventConsumer, maxSubscribersPerMeeting)) {
|
|
||||||
subscriptionQuota.release(meetingId, next.subscriptionMember);
|
|
||||||
RealtimeMeetingMetrics.subscriberRejected("grpc");
|
|
||||||
throw new IllegalStateException("会议实时订阅人数已达上限");
|
|
||||||
}
|
}
|
||||||
|
next.pushSubscribed = true;
|
||||||
try {
|
try {
|
||||||
observer.onNext(ServerMessage.newBuilder()
|
observer.onNext(ServerMessage.newBuilder()
|
||||||
.setConnectAck(ConnectResponse.newBuilder().setSuccess(true).setMessage("实时转写订阅成功"))
|
.setConnectAck(ConnectResponse.newBuilder()
|
||||||
|
.setSuccess(true)
|
||||||
|
.setMessage("实时转写订阅成功,转写经推送通道下发")
|
||||||
|
.build())
|
||||||
.build());
|
.build());
|
||||||
replayCachedTranscripts(meetingId, next);
|
pushSubscriptionManager.replayHistory(meetingId, subscriberDeviceId, meeting.getTitle());
|
||||||
subscriberSessions.put(connectionId, next);
|
subscriberSessions.put(connectionId, next);
|
||||||
RealtimeMeetingMetrics.subscriberOpened("grpc");
|
|
||||||
return next;
|
return next;
|
||||||
} catch (RuntimeException ex) {
|
} catch (RuntimeException ex) {
|
||||||
eventHub.unsubscribe(meetingId, next.eventConsumer);
|
detachSubscriber(next);
|
||||||
subscriptionQuota.release(meetingId, next.subscriptionMember);
|
|
||||||
throw ex;
|
throw ex;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -428,7 +401,6 @@ public class AndroidRealtimeMeetingGrpcService
|
||||||
} catch (Exception ex) {
|
} catch (Exception ex) {
|
||||||
sessions.remove(meetingId, next);
|
sessions.remove(meetingId, next);
|
||||||
publisherLease.release(meetingId, publisherFenceToken);
|
publisherLease.release(meetingId, publisherFenceToken);
|
||||||
eventHub.unsubscribe(meetingId, next.eventConsumer);
|
|
||||||
try {
|
try {
|
||||||
channel.closeMeeting(next.context);
|
channel.closeMeeting(next.context);
|
||||||
} catch (Exception closeEx) {
|
} catch (Exception closeEx) {
|
||||||
|
|
@ -466,19 +438,6 @@ public class AndroidRealtimeMeetingGrpcService
|
||||||
return message == null || message.isBlank() ? fallback : message;
|
return message == null || message.isBlank() ? fallback : message;
|
||||||
}
|
}
|
||||||
|
|
||||||
private void replayCachedTranscripts(Long meetingId, GrpcSession session) {
|
|
||||||
try {
|
|
||||||
for (var item : transcriptCacheService.listOrderedItems(meetingId)) {
|
|
||||||
session.send(ServerMessage.newBuilder()
|
|
||||||
.setTranscript(TextMessage.newBuilder()
|
|
||||||
.setText(com.imeeting.service.realtime.impl.LocalRealtimeAsrChannel.buildFrontendTranscriptMessage(item)))
|
|
||||||
.build());
|
|
||||||
}
|
|
||||||
} catch (Exception ex) {
|
|
||||||
log.warn("Failed to replay realtime transcripts to subscriber, meetingId={}", meetingId, ex);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
private static final class GrpcSession {
|
private static final class GrpcSession {
|
||||||
private final Long meetingId;
|
private final Long meetingId;
|
||||||
private final String connectionId;
|
private final String connectionId;
|
||||||
|
|
@ -491,8 +450,7 @@ public class AndroidRealtimeMeetingGrpcService
|
||||||
private volatile boolean publisher;
|
private volatile boolean publisher;
|
||||||
private boolean canPublish;
|
private boolean canPublish;
|
||||||
private Runnable cleanupAction;
|
private Runnable cleanupAction;
|
||||||
private Consumer<String> eventConsumer;
|
private volatile boolean pushSubscribed;
|
||||||
private String subscriptionMember;
|
|
||||||
|
|
||||||
private GrpcSession(Long meetingId, String connectionId, RealtimeAsrChannel channel,
|
private GrpcSession(Long meetingId, String connectionId, RealtimeAsrChannel channel,
|
||||||
StreamObserver<ServerMessage> observer, String publisherFenceToken) {
|
StreamObserver<ServerMessage> observer, String publisherFenceToken) {
|
||||||
|
|
|
||||||
|
|
@ -29,6 +29,14 @@ public interface MeetingCommandService {
|
||||||
|
|
||||||
MeetingVO createRealtimeMeeting(CreateRealtimeMeetingCommand command, Long tenantId, Long creatorId, String creatorName, String meetingSource);
|
MeetingVO createRealtimeMeeting(CreateRealtimeMeetingCommand command, Long tenantId, Long creatorId, String creatorName, String meetingSource);
|
||||||
|
|
||||||
|
MeetingVO createRealtimeMeeting(CreateRealtimeMeetingCommand command,
|
||||||
|
Long tenantId,
|
||||||
|
Long creatorId,
|
||||||
|
String creatorName,
|
||||||
|
String meetingSource,
|
||||||
|
String sourceDeviceCode,
|
||||||
|
String sourceDeviceMode);
|
||||||
|
|
||||||
MeetingVO createPublicDeviceMeeting(PublicDeviceMeetingCreateCommand command,
|
MeetingVO createPublicDeviceMeeting(PublicDeviceMeetingCreateCommand command,
|
||||||
Long tenantId,
|
Long tenantId,
|
||||||
Long creatorId,
|
Long creatorId,
|
||||||
|
|
|
||||||
|
|
@ -246,6 +246,18 @@ public class MeetingCommandServiceImpl implements MeetingCommandService {
|
||||||
@Override
|
@Override
|
||||||
@Transactional(rollbackFor = Exception.class)
|
@Transactional(rollbackFor = Exception.class)
|
||||||
public MeetingVO createRealtimeMeeting(CreateRealtimeMeetingCommand command, Long tenantId, Long creatorId, String creatorName, String meetingSource) {
|
public MeetingVO createRealtimeMeeting(CreateRealtimeMeetingCommand command, Long tenantId, Long creatorId, String creatorName, String meetingSource) {
|
||||||
|
return createRealtimeMeeting(command, tenantId, creatorId, creatorName, meetingSource, null, null);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
@Transactional(rollbackFor = Exception.class)
|
||||||
|
public MeetingVO createRealtimeMeeting(CreateRealtimeMeetingCommand command,
|
||||||
|
Long tenantId,
|
||||||
|
Long creatorId,
|
||||||
|
String creatorName,
|
||||||
|
String meetingSource,
|
||||||
|
String sourceDeviceCode,
|
||||||
|
String sourceDeviceMode) {
|
||||||
RealtimeMeetingRuntimeProfile runtimeProfile = resolveCreateProfile(command, tenantId, creatorId);
|
RealtimeMeetingRuntimeProfile runtimeProfile = resolveCreateProfile(command, tenantId, creatorId);
|
||||||
Long hostUserId = resolveHostUserId(command.getHostUserId(), creatorId);
|
Long hostUserId = resolveHostUserId(command.getHostUserId(), creatorId);
|
||||||
String resolvedCreatorName = resolveMeetingUserName(creatorId, creatorName);
|
String resolvedCreatorName = resolveMeetingUserName(creatorId, creatorName);
|
||||||
|
|
@ -254,7 +266,7 @@ public class MeetingCommandServiceImpl implements MeetingCommandService {
|
||||||
Meeting meeting = meetingDomainSupport.initMeeting(command.getTitle(), command.getMeetingTime(), command.getParticipants(), command.getTags(),
|
Meeting meeting = meetingDomainSupport.initMeeting(command.getTitle(), command.getMeetingTime(), command.getParticipants(), command.getTags(),
|
||||||
null, MeetingConstants.TYPE_REALTIME, meetingSource, tenantId, creatorId, resolvedCreatorName,
|
null, MeetingConstants.TYPE_REALTIME, meetingSource, tenantId, creatorId, resolvedCreatorName,
|
||||||
hostUserId, hostName, runtimeProfile.getResolvedSummaryModelId(), runtimeProfile.getResolvedPromptId(),
|
hostUserId, hostName, runtimeProfile.getResolvedSummaryModelId(), runtimeProfile.getResolvedPromptId(),
|
||||||
runtimeProfile.getResolvedHotWordGroupId(), summaryDetailLevel, 0);
|
runtimeProfile.getResolvedHotWordGroupId(), summaryDetailLevel, 0, sourceDeviceCode, sourceDeviceMode);
|
||||||
meetingService.save(meeting);
|
meetingService.save(meeting);
|
||||||
Long chapterModelId = command.getChapterModelId() != null ? command.getChapterModelId() : runtimeProfile.getResolvedSummaryModelId();
|
Long chapterModelId = command.getChapterModelId() != null ? command.getChapterModelId() : runtimeProfile.getResolvedSummaryModelId();
|
||||||
createChapterTaskIfEnabled(
|
createChapterTaskIfEnabled(
|
||||||
|
|
|
||||||
|
|
@ -41,4 +41,8 @@ public final class RealtimeMeetingMetrics {
|
||||||
public static void subscriberDeliveryFailed() {
|
public static void subscriberDeliveryFailed() {
|
||||||
Metrics.counter("imeeting.realtime.subscriber.delivery.failed").increment();
|
Metrics.counter("imeeting.realtime.subscriber.delivery.failed").increment();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public static void pushDeliveryMissed() {
|
||||||
|
Metrics.counter("imeeting.realtime.subscriber.push.delivery.missed").increment();
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue