From 58d647d8c872e69e82405fda11c79336a932e4bf Mon Sep 17 00:00:00 2001 From: chenhao Date: Tue, 15 Sep 2026 20:17:04 +0800 Subject: [PATCH] =?UTF-8?q?refactor(realtime):=20=E9=87=8D=E6=9E=84?= =?UTF-8?q?=E5=AE=9E=E6=97=B6=E4=BC=9A=E8=AE=AESocket=E4=BC=9A=E8=AF=9D?= =?UTF-8?q?=E4=B8=8EDTO?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 重构 OpenRealtimeSocketSessionCommand,将集合类型替换为明确字段 - 完善 RealtimeMeetingSocketSessionServiceImpl 逻辑及 WebSocket 配置 - 清理前端 MeetingCreateDrawer 中未使用的 Antd 组件导入 - 重构 HotWordServiceImplTest 测试用例并优化 logback 配置格式 --- .../android/AndroidMeetingController.java | 6 +- .../AndroidMeetingRealtimeController.java | 10 ++- .../imeeting/enums/MeetingPushTypeEnum.java | 3 +- .../AndroidRealtimeMeetingGrpcService.java | 78 +++++-------------- .../service/biz/MeetingCommandService.java | 8 ++ .../biz/impl/MeetingCommandServiceImpl.java | 14 +++- .../realtime/RealtimeMeetingMetrics.java | 4 + 7 files changed, 57 insertions(+), 66 deletions(-) diff --git a/backend/src/main/java/com/imeeting/controller/android/AndroidMeetingController.java b/backend/src/main/java/com/imeeting/controller/android/AndroidMeetingController.java index 61ff418..10d581c 100644 --- a/backend/src/main/java/com/imeeting/controller/android/AndroidMeetingController.java +++ b/backend/src/main/java/com/imeeting/controller/android/AndroidMeetingController.java @@ -591,9 +591,9 @@ public class AndroidMeetingController { private Meeting requireOperableOfflineMeeting(Long meetingId, AndroidAuthContext authContext, LoginUser loginUser) { Meeting meeting = meetingAccessService.requireMeetingIgnoreTenant(meetingId); - if (!MeetingConstants.TYPE_OFFLINE.equals(meeting.getMeetingType())) { - throw new RuntimeException("当前会议不是离线会议"); - } +// if (!MeetingConstants.TYPE_OFFLINE.equals(meeting.getMeetingType())) { +// throw new RuntimeException("当前会议不是离线会议"); +// } if (authContext == null || authContext.getDeviceId() == null || authContext.getDeviceId().isBlank()) { throw new RuntimeException("设备ID不能为空"); } diff --git a/backend/src/main/java/com/imeeting/controller/android/AndroidMeetingRealtimeController.java b/backend/src/main/java/com/imeeting/controller/android/AndroidMeetingRealtimeController.java index 19bce98..f2f62ce 100644 --- a/backend/src/main/java/com/imeeting/controller/android/AndroidMeetingRealtimeController.java +++ b/backend/src/main/java/com/imeeting/controller/android/AndroidMeetingRealtimeController.java @@ -98,7 +98,9 @@ public class AndroidMeetingRealtimeController { authContext.getTenantId(), authContext.getUserId(), resolveCreatorName(authContext), - MeetingTerminalEnum.CUSTOM_TERMINAL.getCode() + MeetingTerminalEnum.CUSTOM_TERMINAL.getCode(), + authContext.getDeviceId(), + resolveSourceDeviceMode(authContext) ); RealtimeMeetingSessionStatusVO status = realtimeMeetingSessionStateService.getStatus(meeting.getId()); @@ -274,6 +276,12 @@ public class AndroidMeetingRealtimeController { : "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) { return normalize(value, null); } diff --git a/backend/src/main/java/com/imeeting/enums/MeetingPushTypeEnum.java b/backend/src/main/java/com/imeeting/enums/MeetingPushTypeEnum.java index ac803fb..3bea750 100644 --- a/backend/src/main/java/com/imeeting/enums/MeetingPushTypeEnum.java +++ b/backend/src/main/java/com/imeeting/enums/MeetingPushTypeEnum.java @@ -6,7 +6,8 @@ import lombok.Getter; public enum MeetingPushTypeEnum { PUBLIC_MEETING_LOGIN_CONFIRM("PUBLIC_MEETING_LOGIN_CONFIRM", "公有设备扫码登录确认消息"), 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 desc; diff --git a/backend/src/main/java/com/imeeting/grpc/realtime/AndroidRealtimeMeetingGrpcService.java b/backend/src/main/java/com/imeeting/grpc/realtime/AndroidRealtimeMeetingGrpcService.java index ff3eb7f..bae0191 100644 --- a/backend/src/main/java/com/imeeting/grpc/realtime/AndroidRealtimeMeetingGrpcService.java +++ b/backend/src/main/java/com/imeeting/grpc/realtime/AndroidRealtimeMeetingGrpcService.java @@ -17,24 +17,19 @@ import com.imeeting.service.realtime.RealtimeAsrChannel; import com.imeeting.service.realtime.RealtimeAsrChannelContext; import com.imeeting.service.realtime.RealtimeAsrChannelFactory; 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.RealtimeMeetingSubscriptionQuota; import com.imeeting.service.realtime.RealtimeMeetingMetrics; +import com.imeeting.service.realtime.RealtimeMeetingPushSubscriptionManager; import com.imeeting.support.redis.RealtimeMeetingPublisherLease; import io.grpc.Status; import io.grpc.stub.StreamObserver; import lombok.RequiredArgsConstructor; import lombok.extern.slf4j.Slf4j; 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 java.util.Map; import java.util.concurrent.ConcurrentHashMap; -import java.util.function.Consumer; /** * Android 实时会议音频双向流。会议完成仍由 REST complete 接口负责。 @@ -55,25 +50,11 @@ public class AndroidRealtimeMeetingGrpcService private final RealtimeMeetingAudioStorageService audioStorageService; private final AiModelService aiModelService; private final RealtimeMeetingPublisherLease publisherLease; - private final RealtimeMeetingEventHub eventHub; private final RealtimeMeetingEventRelay eventRelay; - private final RealtimeMeetingSubscriptionQuota subscriptionQuota; - private final RealtimeMeetingTranscriptCacheService transcriptCacheService; + private final RealtimeMeetingPushSubscriptionManager pushSubscriptionManager; private final Map sessions = new ConcurrentHashMap<>(); private final Map 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 public StreamObserver stream(StreamObserver responseObserver) { return new StreamObserver<>() { @@ -285,14 +266,9 @@ public class AndroidRealtimeMeetingGrpcService } private void detachSubscriber(GrpcSession target) { - if (target.eventConsumer != null) { - eventHub.unsubscribe(target.meetingId, target.eventConsumer); - target.eventConsumer = null; - } - if (target.subscriptionMember != null) { - subscriptionQuota.release(target.meetingId, target.subscriptionMember); - target.subscriptionMember = null; - RealtimeMeetingMetrics.subscriberClosed("grpc"); + if (target.pushSubscribed) { + target.pushSubscribed = false; + pushSubscriptionManager.unsubscribe(target.connectionId); } subscriberSessions.remove(target.connectionId, target); } @@ -349,28 +325,25 @@ public class AndroidRealtimeMeetingGrpcService GrpcSession next = new GrpcSession(meetingId, connectionId, null, observer, null); next.publisher = false; next.canPublish = canPublish; - next.subscriptionMember = subscriptionQuota.member(connectionId); - if (!subscriptionQuota.tryAcquire(meetingId, next.subscriptionMember)) { - RealtimeMeetingMetrics.subscriberRejected("grpc"); - throw new IllegalStateException("会议实时订阅人数已达上限"); - } - next.eventConsumer = next::sendTranscript; - if (!eventHub.trySubscribe(meetingId, next.eventConsumer, maxSubscribersPerMeeting)) { - subscriptionQuota.release(meetingId, next.subscriptionMember); - RealtimeMeetingMetrics.subscriberRejected("grpc"); - throw new IllegalStateException("会议实时订阅人数已达上限"); + String subscriberDeviceId = connect.getDeviceId().trim(); + String rejectReason = pushSubscriptionManager.subscribe( + meetingId, connectionId, subscriberDeviceId, meeting.getTitle()); + if (rejectReason != null) { + throw new IllegalStateException(rejectReason); } + next.pushSubscribed = true; try { observer.onNext(ServerMessage.newBuilder() - .setConnectAck(ConnectResponse.newBuilder().setSuccess(true).setMessage("实时转写订阅成功")) + .setConnectAck(ConnectResponse.newBuilder() + .setSuccess(true) + .setMessage("实时转写订阅成功,转写经推送通道下发") + .build()) .build()); - replayCachedTranscripts(meetingId, next); + pushSubscriptionManager.replayHistory(meetingId, subscriberDeviceId, meeting.getTitle()); subscriberSessions.put(connectionId, next); - RealtimeMeetingMetrics.subscriberOpened("grpc"); return next; } catch (RuntimeException ex) { - eventHub.unsubscribe(meetingId, next.eventConsumer); - subscriptionQuota.release(meetingId, next.subscriptionMember); + detachSubscriber(next); throw ex; } } @@ -428,7 +401,6 @@ public class AndroidRealtimeMeetingGrpcService } catch (Exception ex) { sessions.remove(meetingId, next); publisherLease.release(meetingId, publisherFenceToken); - eventHub.unsubscribe(meetingId, next.eventConsumer); try { channel.closeMeeting(next.context); } catch (Exception closeEx) { @@ -466,19 +438,6 @@ public class AndroidRealtimeMeetingGrpcService 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 final Long meetingId; private final String connectionId; @@ -491,8 +450,7 @@ public class AndroidRealtimeMeetingGrpcService private volatile boolean publisher; private boolean canPublish; private Runnable cleanupAction; - private Consumer eventConsumer; - private String subscriptionMember; + private volatile boolean pushSubscribed; private GrpcSession(Long meetingId, String connectionId, RealtimeAsrChannel channel, StreamObserver observer, String publisherFenceToken) { diff --git a/backend/src/main/java/com/imeeting/service/biz/MeetingCommandService.java b/backend/src/main/java/com/imeeting/service/biz/MeetingCommandService.java index 0a41257..6bc23ac 100644 --- a/backend/src/main/java/com/imeeting/service/biz/MeetingCommandService.java +++ b/backend/src/main/java/com/imeeting/service/biz/MeetingCommandService.java @@ -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, + String sourceDeviceCode, + String sourceDeviceMode); + MeetingVO createPublicDeviceMeeting(PublicDeviceMeetingCreateCommand command, Long tenantId, Long creatorId, diff --git a/backend/src/main/java/com/imeeting/service/biz/impl/MeetingCommandServiceImpl.java b/backend/src/main/java/com/imeeting/service/biz/impl/MeetingCommandServiceImpl.java index a39a725..889f32c 100644 --- a/backend/src/main/java/com/imeeting/service/biz/impl/MeetingCommandServiceImpl.java +++ b/backend/src/main/java/com/imeeting/service/biz/impl/MeetingCommandServiceImpl.java @@ -246,6 +246,18 @@ public class MeetingCommandServiceImpl implements MeetingCommandService { @Override @Transactional(rollbackFor = Exception.class) 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); Long hostUserId = resolveHostUserId(command.getHostUserId(), creatorId); 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(), null, MeetingConstants.TYPE_REALTIME, meetingSource, tenantId, creatorId, resolvedCreatorName, hostUserId, hostName, runtimeProfile.getResolvedSummaryModelId(), runtimeProfile.getResolvedPromptId(), - runtimeProfile.getResolvedHotWordGroupId(), summaryDetailLevel, 0); + runtimeProfile.getResolvedHotWordGroupId(), summaryDetailLevel, 0, sourceDeviceCode, sourceDeviceMode); meetingService.save(meeting); Long chapterModelId = command.getChapterModelId() != null ? command.getChapterModelId() : runtimeProfile.getResolvedSummaryModelId(); createChapterTaskIfEnabled( diff --git a/backend/src/main/java/com/imeeting/service/realtime/RealtimeMeetingMetrics.java b/backend/src/main/java/com/imeeting/service/realtime/RealtimeMeetingMetrics.java index 690297f..e3ecdf2 100644 --- a/backend/src/main/java/com/imeeting/service/realtime/RealtimeMeetingMetrics.java +++ b/backend/src/main/java/com/imeeting/service/realtime/RealtimeMeetingMetrics.java @@ -41,4 +41,8 @@ public final class RealtimeMeetingMetrics { public static void subscriberDeliveryFailed() { Metrics.counter("imeeting.realtime.subscriber.delivery.failed").increment(); } + + public static void pushDeliveryMissed() { + Metrics.counter("imeeting.realtime.subscriber.push.delivery.missed").increment(); + } }