Skip to content
Projects
Groups
Snippets
Help
This project
Loading...
Sign in / Register
Toggle navigation
Y
yd-communication
Overview
Overview
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
xingmin
yd-communication
Commits
76ebc4a4
Commit
76ebc4a4
authored
Jul 30, 2026
by
zhangxingmin
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
push
parent
ed1e4e0a
Show whitespace changes
Inline
Side-by-side
Showing
35 changed files
with
2101 additions
and
30 deletions
+2101
-30
.idea/vcs.xml
+4
-0
yd-communication-api/pom.xml
+0
-6
yd-communication-api/src/main/java/com/yd/communication/api/config/WebSocketConfig.java
+14
-0
yd-communication-api/src/main/java/com/yd/communication/api/controller/ApiCoSessionController.java
+78
-0
yd-communication-api/src/main/java/com/yd/communication/api/controller/CoSessionController.java
+0
-18
yd-communication-api/src/main/java/com/yd/communication/api/service/ApiCoSessionService.java
+32
-0
yd-communication-api/src/main/java/com/yd/communication/api/service/impl/ApiCoSessionServiceImpl.java
+394
-0
yd-communication-api/src/main/java/com/yd/communication/api/websocket/CoWebSocketServer.java
+679
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/client/ApiCoSessionFeignClient.java
+69
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/constant/RedisConstants.java
+13
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/dto/RoomRedisInfoDTO.java
+30
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/enums/CoSessionStatusEnum.java
+32
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/enums/ControlHolderTypeEnum.java
+30
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/enums/RedisEnum.java
+42
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/fallback/ApiCoSessionFeignFallbackFactory.java
+53
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/request/CreateRequest.java
+74
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/request/EndSessionRequest.java
+13
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/request/JoinRequest.java
+46
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/request/TransferControlRequest.java
+23
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/response/CommonResponse.java
+21
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/response/CreateResponse.java
+31
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/response/JoinResponse.java
+48
-0
yd-communication-feign/src/main/java/com/yd/communication/feign/response/SessionDetailResponse.java
+107
-0
yd-communication-service/pom.xml
+6
-0
yd-communication-service/src/main/java/com/yd/communication/service/dao/CoSessionMapper.java
+14
-0
yd-communication-service/src/main/java/com/yd/communication/service/model/CoSession.java
+3
-2
yd-communication-service/src/main/java/com/yd/communication/service/model/RecordingTask.java
+3
-3
yd-communication-service/src/main/java/com/yd/communication/service/service/ICoOperationLogService.java
+3
-0
yd-communication-service/src/main/java/com/yd/communication/service/service/ICoSessionService.java
+12
-0
yd-communication-service/src/main/java/com/yd/communication/service/service/IRecordingTaskService.java
+3
-0
yd-communication-service/src/main/java/com/yd/communication/service/service/impl/CoDesensitizationRuleServiceImpl.java
+84
-0
yd-communication-service/src/main/java/com/yd/communication/service/service/impl/CoOperationLogServiceImpl.java
+22
-0
yd-communication-service/src/main/java/com/yd/communication/service/service/impl/CoSessionServiceImpl.java
+31
-1
yd-communication-service/src/main/java/com/yd/communication/service/service/impl/RecordingTaskServiceImpl.java
+71
-0
yd-communication-service/src/main/java/com/yd/communication/service/utils/RandomUtil.java
+16
-0
No files found.
.idea/vcs.xml
View file @
76ebc4a4
...
@@ -5,4 +5,7 @@
...
@@ -5,4 +5,7 @@
<list
/>
<list
/>
</option>
</option>
</component>
</component>
<component
name=
"VcsDirectoryMappings"
>
<mapping
directory=
"$PROJECT_DIR$"
vcs=
"Git"
/>
</component>
</project>
</project>
\ No newline at end of file
yd-communication-api/pom.xml
View file @
76ebc4a4
...
@@ -29,12 +29,6 @@
...
@@ -29,12 +29,6 @@
<artifactId>
spring-boot-starter-web
</artifactId>
<artifactId>
spring-boot-starter-web
</artifactId>
</dependency>
</dependency>
<!-- Spring Boot Starter WebSocket -->
<dependency>
<groupId>
org.springframework.boot
</groupId>
<artifactId>
spring-boot-starter-websocket
</artifactId>
</dependency>
<dependency>
<dependency>
<groupId>
org.springframework.boot
</groupId>
<groupId>
org.springframework.boot
</groupId>
<artifactId>
spring-boot-starter
</artifactId>
<artifactId>
spring-boot-starter
</artifactId>
...
...
yd-communication-api/src/main/java/com/yd/communication/api/config/WebSocketConfig.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
api
.
config
;
import
org.springframework.context.annotation.Bean
;
import
org.springframework.context.annotation.Configuration
;
import
org.springframework.web.socket.server.standard.ServerEndpointExporter
;
@Configuration
public
class
WebSocketConfig
{
@Bean
public
ServerEndpointExporter
serverEndpointExporter
()
{
return
new
ServerEndpointExporter
();
}
}
\ No newline at end of file
yd-communication-api/src/main/java/com/yd/communication/api/controller/ApiCoSessionController.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
api
.
controller
;
import
com.yd.common.result.Result
;
import
com.yd.communication.api.service.ApiCoSessionService
;
import
com.yd.communication.feign.client.ApiCoSessionFeignClient
;
import
com.yd.communication.feign.request.CreateRequest
;
import
com.yd.communication.feign.request.EndSessionRequest
;
import
com.yd.communication.feign.request.JoinRequest
;
import
com.yd.communication.feign.request.TransferControlRequest
;
import
com.yd.communication.feign.response.CommonResponse
;
import
com.yd.communication.feign.response.CreateResponse
;
import
com.yd.communication.feign.response.JoinResponse
;
import
com.yd.communication.feign.response.SessionDetailResponse
;
import
org.springframework.validation.annotation.Validated
;
import
org.springframework.web.bind.annotation.*
;
import
javax.annotation.Resource
;
/**
* 协同会话信息
*
* @author zxm
* @since 2026-07-28
*/
@RestController
@RequestMapping
(
"/coSession"
)
@Validated
public
class
ApiCoSessionController
implements
ApiCoSessionFeignClient
{
@Resource
private
ApiCoSessionService
apiCoSessionService
;
/**
* 客户创建会话(生成共享码)
* @param request
* @return
*/
public
Result
<
CreateResponse
>
create
(
CreateRequest
request
)
{
return
apiCoSessionService
.
create
(
request
);
}
/**
* 顾问加入会话(输入共享码加入房间,可以多次加入)
* @param request
* @return
*/
@PostMapping
(
"/join"
)
public
Result
<
JoinResponse
>
join
(
JoinRequest
request
)
{
return
apiCoSessionService
.
join
(
request
);
}
/**
* 获取会话详情
* @param bizId
* @return
*/
public
Result
<
SessionDetailResponse
>
get
(
String
bizId
)
{
return
apiCoSessionService
.
get
(
bizId
);
}
/**
* 结束协同会话(关闭共享,仅客户可调用)
* @return
*/
public
Result
<
CommonResponse
>
end
(
EndSessionRequest
request
)
{
return
apiCoSessionService
.
end
(
request
);
}
/**
* 切换控制权(仅参与者(顾问)可调用)
* @param request
* @return
*/
public
Result
<
CommonResponse
>
transferControl
(
TransferControlRequest
request
)
{
return
apiCoSessionService
.
transferControl
(
request
);
}
}
yd-communication-api/src/main/java/com/yd/communication/api/controller/CoSessionController.java
deleted
100644 → 0
View file @
ed1e4e0a
package
com
.
yd
.
communication
.
api
.
controller
;
import
org.springframework.web.bind.annotation.RequestMapping
;
import
org.springframework.web.bind.annotation.RestController
;
/**
* <p>
* 协同-会话表(通用) 前端控制器
* </p>
*
* @author zxm
* @since 2026-07-28
*/
@RestController
@RequestMapping
(
"/coSession"
)
public
class
CoSessionController
{
}
yd-communication-api/src/main/java/com/yd/communication/api/service/ApiCoSessionService.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
api
.
service
;
import
com.yd.common.result.Result
;
import
com.yd.communication.feign.dto.RoomRedisInfoDTO
;
import
com.yd.communication.feign.request.CreateRequest
;
import
com.yd.communication.feign.request.EndSessionRequest
;
import
com.yd.communication.feign.request.JoinRequest
;
import
com.yd.communication.feign.request.TransferControlRequest
;
import
com.yd.communication.feign.response.CommonResponse
;
import
com.yd.communication.feign.response.CreateResponse
;
import
com.yd.communication.feign.response.JoinResponse
;
import
com.yd.communication.feign.response.SessionDetailResponse
;
import
com.yd.communication.service.model.CoSession
;
public
interface
ApiCoSessionService
{
Result
<
CreateResponse
>
create
(
CreateRequest
request
);
Result
<
JoinResponse
>
join
(
JoinRequest
request
);
Result
<
SessionDetailResponse
>
get
(
String
bizId
);
Result
<
CommonResponse
>
end
(
EndSessionRequest
request
);
Result
<
CommonResponse
>
transferControl
(
TransferControlRequest
request
);
RoomRedisInfoDTO
getCurrentController
(
String
roomId
);
void
updateCurrentPage
(
String
roomId
,
String
newPageJson
,
String
operatorId
);
CoSession
getByRoomId
(
String
roomId
);
}
yd-communication-api/src/main/java/com/yd/communication/api/service/impl/ApiCoSessionServiceImpl.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
api
.
service
.
impl
;
import
com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper
;
import
com.yd.common.enums.CommonEnum
;
import
com.yd.common.exception.BusinessException
;
import
com.yd.common.result.Result
;
import
com.yd.common.utils.RandomStringGenerator
;
import
com.yd.common.utils.RedisUtil
;
import
com.yd.communication.api.service.ApiCoSessionService
;
import
com.yd.communication.feign.dto.RoomRedisInfoDTO
;
import
com.yd.communication.feign.enums.CoSessionStatusEnum
;
import
com.yd.communication.feign.enums.ControlHolderTypeEnum
;
import
com.yd.communication.feign.enums.RedisEnum
;
import
com.yd.communication.feign.request.CreateRequest
;
import
com.yd.communication.feign.request.EndSessionRequest
;
import
com.yd.communication.feign.request.JoinRequest
;
import
com.yd.communication.feign.request.TransferControlRequest
;
import
com.yd.communication.feign.response.CommonResponse
;
import
com.yd.communication.feign.response.CreateResponse
;
import
com.yd.communication.feign.response.JoinResponse
;
import
com.yd.communication.feign.response.SessionDetailResponse
;
import
com.yd.communication.service.model.CoSession
;
import
com.yd.communication.service.service.ICoSessionService
;
import
com.yd.communication.service.service.IRecordingTaskService
;
import
com.yd.communication.service.utils.RandomUtil
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.commons.lang3.StringUtils
;
import
org.springframework.beans.BeanUtils
;
import
org.springframework.stereotype.Service
;
import
org.springframework.transaction.annotation.Transactional
;
import
javax.annotation.Resource
;
import
java.time.LocalDateTime
;
import
java.util.UUID
;
@Slf4j
@Service
public
class
ApiCoSessionServiceImpl
implements
ApiCoSessionService
{
@Resource
private
ICoSessionService
iCoSessionService
;
@Resource
private
IRecordingTaskService
recordingService
;
@Resource
private
RedisUtil
redisUtil
;
/**
* 客户创建会话(生成共享码)
* @param request
* @return
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
CreateResponse
>
create
(
CreateRequest
request
)
{
CoSession
session
=
createSession
(
request
.
getScope
(),
request
.
getResourceType
(),
request
.
getResourceId
(),
request
.
getResourceInit
(),
request
.
getOwnerId
(),
request
.
getOwnerType
(),
request
.
getUserId
(),
request
.
getToken
()
);
if
(
session
==
null
)
{
return
Result
.
success
();
}
CreateResponse
response
=
new
CreateResponse
();
response
.
setRoomId
(
session
.
getRoomId
());
response
.
setRoomPwd
(
session
.
getRoomPwd
());
response
.
setSessionBizId
(
session
.
getCoSessionBizId
());
response
.
setStatus
(
session
.
getStatus
());
return
Result
.
success
(
response
);
}
/**
* 客户创建会话(生成共享码)
* @param scope
* @param resourceType
* @param resourceId
* @param resourceInit
* @param ownerId
* @param ownerType
* @return
*/
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
CoSession
createSession
(
String
scope
,
String
resourceType
,
String
resourceId
,
String
resourceInit
,
String
ownerId
,
String
ownerType
,
String
userId
,
String
token
)
{
//创建会话
CoSession
session
=
new
CoSession
();
session
.
setCoSessionBizId
(
RandomStringGenerator
.
generateBizId16
(
CommonEnum
.
UID_TYPE_CO_SESSION
.
getCode
()));
session
.
setCoSessionNo
(
"S"
+
System
.
currentTimeMillis
());
session
.
setScope
(
scope
);
session
.
setResourceType
(
resourceType
);
session
.
setResourceId
(
resourceId
);
session
.
setResourceInit
(
resourceInit
);
session
.
setOwnerId
(
ownerId
);
session
.
setOwnerType
(
ownerType
);
//房间号
String
roomId
=
"room_"
+
UUID
.
randomUUID
().
toString
().
substring
(
0
,
8
);
//房间密码
String
roomPwd
=
RandomUtil
.
generateNumericCode
(
6
);
session
.
setRoomId
(
roomId
);
session
.
setRoomPwd
(
roomPwd
);
session
.
setChannelPrefix
(
"co"
);
//初始化会话信息默认控制权在资源所有者类型(owner)
session
.
setControlHolderType
(
ControlHolderTypeEnum
.
OWNER
.
getItemValue
());
session
.
setControlHolderId
(
ownerId
);
//待开始状态
session
.
setStatus
(
CoSessionStatusEnum
.
DKS
.
getItemValue
());
//当前操作页面,实时更新
session
.
setCurrentPage
(
resourceInit
);
//页面访问轨迹
session
.
setPageHistory
(
"["
+
resourceInit
+
"]"
);
session
.
setCreatorId
(
ownerId
);
session
.
setUpdaterId
(
ownerId
);
iCoSessionService
.
save
(
session
);
RoomRedisInfoDTO
roomRedisInfoDTO
=
new
RoomRedisInfoDTO
();
roomRedisInfoDTO
.
setToken
(
token
);
roomRedisInfoDTO
.
setUserId
(
userId
);
roomRedisInfoDTO
.
setControlHolderType
(
session
.
getControlHolderType
());
roomRedisInfoDTO
.
setControlHolderId
(
session
.
getControlHolderId
());
// 存储当前协同房间缓存信息redis
redisUtil
.
setCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
,
roomRedisInfoDTO
,
RedisEnum
.
ROOM
.
getTimeout
(),
RedisEnum
.
ROOM
.
getTimeUnit
());
// 自动开启录制
if
(
autoStartRecording
())
{
// recordingService.startRecording(session.getCoSessionBizId(), roomId);
}
log
.
info
(
"创建会话成功, roomId={}, roomPwd={}"
,
roomId
,
roomPwd
);
//添加操作日志,协同-操作日志表 TODO
return
session
;
}
/**
* 顾问加入会话((输入共享码加入房间,可以多次加入))
* @param request
* @return
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
JoinResponse
>
join
(
JoinRequest
request
)
{
CoSession
session
=
joinSession
(
request
.
getRoomPwd
(),
request
.
getParticipantId
(),
request
.
getParticipantType
()
);
if
(
session
==
null
)
{
return
Result
.
success
();
}
JoinResponse
joinResponse
=
new
JoinResponse
();
joinResponse
.
setControlHolderId
(
session
.
getControlHolderId
());
joinResponse
.
setControlHolderType
(
session
.
getControlHolderType
());
joinResponse
.
setResourceInit
(
session
.
getResourceInit
());
joinResponse
.
setRoomId
(
session
.
getRoomId
());
joinResponse
.
setSessionBizId
(
session
.
getResourceInit
());
//获取资源所有者缓存中的登录信息
RoomRedisInfoDTO
roomRedisInfoDTO
=
redisUtil
.
getCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
session
.
getRoomId
());
if
(
roomRedisInfoDTO
==
null
)
{
throw
new
RuntimeException
(
"会话发起者登录信息失效,建议联系会话发起者再次发起"
);
}
joinResponse
.
setUserId
(
roomRedisInfoDTO
.
getUserId
());
joinResponse
.
setToken
(
roomRedisInfoDTO
.
getToken
());
return
Result
.
success
(
joinResponse
);
}
/**
* 顾问加入会话
* @param roomPwd
* @param participantId
* @param participantType
* @return
*/
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
CoSession
joinSession
(
String
roomPwd
,
String
participantId
,
String
participantType
)
{
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomPwd
,
roomPwd
)
.
last
(
" limit 1 "
);
CoSession
session
=
iCoSessionService
.
getOne
(
wrapper
);
if
(
session
==
null
)
{
throw
new
RuntimeException
(
"房间密码(共享码)错误,或客户未发起会话"
);
}
if
(
CoSessionStatusEnum
.
YJS
.
getItemValue
().
equals
(
session
.
getStatus
()))
{
//已结束,不能再次加入房间
throw
new
RuntimeException
(
"会话已结束,不能再次加入房间"
);
}
if
(
StringUtils
.
isNotBlank
(
session
.
getParticipantId
())
&&
!
session
.
getParticipantId
().
equals
(
participantId
))
{
//库里参与者ID不为空并且库中参与者ID和加入房间的参与者ID不相等,不能加入房间
throw
new
RuntimeException
(
"当前房间被占用,不能加入到房间"
);
}
if
(
CoSessionStatusEnum
.
DKS
.
getItemValue
().
equals
(
session
.
getStatus
()))
{
if
(
session
.
getStartTime
()
==
null
)
{
//开始时间为空时,说明当前处于待开始到进行中过渡阶段,设置开始时间
session
.
setStartTime
(
LocalDateTime
.
now
());
}
//待开始状态下,控制权要自动移交给参与者(顾问)
session
.
setControlHolderType
(
ControlHolderTypeEnum
.
PARTICIPANT
.
getItemValue
());
session
.
setControlHolderId
(
participantId
);
}
//更新会话状态:进行中
session
.
setStatus
(
CoSessionStatusEnum
.
JXZ
.
getItemValue
());
session
.
setParticipantId
(
participantId
);
session
.
setParticipantType
(
participantType
);
session
.
setUpdaterId
(
participantId
);
iCoSessionService
.
updateById
(
session
);
//添加操作日志,协同-操作日志表 TODO
return
session
;
}
/**
* 获取会话详情
* @param bizId
* @return
*/
@Override
public
Result
<
SessionDetailResponse
>
get
(
String
bizId
)
{
CoSession
session
=
iCoSessionService
.
getByBizId
(
bizId
);
if
(
session
==
null
)
{
return
Result
.
success
();
}
SessionDetailResponse
response
=
new
SessionDetailResponse
();
BeanUtils
.
copyProperties
(
session
,
response
);
return
Result
.
success
(
response
);
}
/**
* 结束会话(客户调用)
* @param request
* @return
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
CommonResponse
>
end
(
EndSessionRequest
request
)
{
CoSession
coSession
=
iCoSessionService
.
lambdaQuery
()
.
eq
(
CoSession:
:
getRoomId
,
request
.
getRoomId
())
.
last
(
" limit 1"
)
.
one
();
if
(
coSession
==
null
)
{
throw
new
BusinessException
(
"会话不存在"
);
}
//结束会话关闭共享,更新信息
coSession
.
setStatus
(
CoSessionStatusEnum
.
YJS
.
getItemValue
());
//结束时间
coSession
.
setEndTime
(
LocalDateTime
.
now
());
iCoSessionService
.
saveOrUpdate
(
coSession
);
//销毁redis房间缓存信息
redisUtil
.
deleteObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
coSession
.
getRoomId
());
//结束录制视频,并且更新录制任务表信息存档 TODO
//添加操作日志,协同-操作日志表 TODO
CommonResponse
response
=
new
CommonResponse
();
response
.
setMessage
(
"会话已结束"
);
return
Result
.
success
(
response
);
}
/**
* 切换控制权(顾问调用)
* @param request
* @return
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
CommonResponse
>
transferControl
(
TransferControlRequest
request
)
{
transferControlUp
(
request
.
getOprType
(),
request
.
getRoomId
());
CommonResponse
response
=
new
CommonResponse
();
response
.
setMessage
(
"控制权已切换"
);
return
Result
.
success
(
response
);
}
/**
* 切换控制权
* @param oprType
* @param roomId
*/
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
void
transferControlUp
(
Integer
oprType
,
String
roomId
)
{
//根据房间号查询会话信息
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomId
,
roomId
).
last
(
" limit 1 "
);
CoSession
session
=
iCoSessionService
.
getOne
(
wrapper
);
if
(
session
==
null
)
{
throw
new
RuntimeException
(
"会话不存在"
);
}
//移交控制权
//控制权持有者类型:owner(资源所有者类型)/participant(参与者类型)
if
(
oprType
==
1
)
{
//1-开启客户操作,控制权移交给资源所有者(客户)
session
.
setControlHolderType
(
ControlHolderTypeEnum
.
OWNER
.
getItemValue
());
session
.
setControlHolderId
(
session
.
getOwnerId
());
}
else
if
(
oprType
==
2
)
{
//2-关闭客户操作,控制权移交给参与者(顾问)
session
.
setControlHolderType
(
ControlHolderTypeEnum
.
PARTICIPANT
.
getItemValue
());
session
.
setControlHolderId
(
session
.
getParticipantId
());
}
session
.
setUpdaterId
(
session
.
getParticipantId
());
iCoSessionService
.
updateById
(
session
);
//更新房间缓存redis信息-更新控制权字段,用于WebSocket获取
RoomRedisInfoDTO
roomRedisInfoDTO
=
redisUtil
.
getCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
);
if
(
roomRedisInfoDTO
!=
null
)
{
roomRedisInfoDTO
.
setControlHolderId
(
session
.
getControlHolderId
());
roomRedisInfoDTO
.
setControlHolderType
(
session
.
getControlHolderType
());
redisUtil
.
setCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
,
roomRedisInfoDTO
,
RedisEnum
.
ROOM
.
getTimeout
(),
RedisEnum
.
ROOM
.
getTimeUnit
());
}
//添加操作日志,协同-操作日志表 TODO
}
/**
* 获取当前控制者
* @param roomId
* @return
*/
@Override
public
RoomRedisInfoDTO
getCurrentController
(
String
roomId
)
{
RoomRedisInfoDTO
roomRedisInfoDTO
=
redisUtil
.
getCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
);
if
(
roomRedisInfoDTO
==
null
||
(
roomRedisInfoDTO
!=
null
&&
StringUtils
.
isBlank
(
roomRedisInfoDTO
.
getControlHolderId
())))
{
//缓存信息或者缓存内的控制者ID为空的时候,单独去查询库更新缓存信息
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomId
,
roomId
).
last
(
" limit 1 "
);
CoSession
session
=
iCoSessionService
.
getOne
(
wrapper
);
if
(
session
==
null
)
{
throw
new
RuntimeException
(
"会话不存在"
);
}
if
(
roomRedisInfoDTO
==
null
)
{
roomRedisInfoDTO
=
new
RoomRedisInfoDTO
();
}
roomRedisInfoDTO
.
setControlHolderType
(
session
.
getControlHolderType
());
roomRedisInfoDTO
.
setControlHolderId
(
session
.
getControlHolderId
());
//查询出来的信息更新回缓存里面
redisUtil
.
setCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
,
roomRedisInfoDTO
,
RedisEnum
.
ROOM
.
getTimeout
(),
RedisEnum
.
ROOM
.
getTimeUnit
());
}
return
roomRedisInfoDTO
;
}
/**
* 更新当前页面
* @param roomId
* @param newPageJson
* @param operatorId
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
void
updateCurrentPage
(
String
roomId
,
String
newPageJson
,
String
operatorId
)
{
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomId
,
roomId
).
last
(
" limit 1 "
);
CoSession
session
=
iCoSessionService
.
getOne
(
wrapper
);
if
(
session
==
null
)
{
throw
new
RuntimeException
(
"会话不存在"
);
}
String
history
=
session
.
getPageHistory
();
if
(
history
==
null
||
history
.
equals
(
"[]"
)
||
history
.
isEmpty
())
{
history
=
"["
+
newPageJson
+
"]"
;
}
else
{
// 简单追加(生产环境建议用JSONArray处理)
history
=
history
.
substring
(
0
,
history
.
length
()
-
1
)
+
","
+
newPageJson
+
"]"
;
}
iCoSessionService
.
updateCurrentPageAndHistory
(
roomId
,
newPageJson
,
history
,
operatorId
);
}
/**
* 根据房间号获取会话信息
* @param roomId
* @return
*/
@Override
public
CoSession
getByRoomId
(
String
roomId
)
{
return
iCoSessionService
.
getByRoomId
(
roomId
);
}
private
boolean
autoStartRecording
()
{
return
true
;
// 从配置读取
}
}
yd-communication-api/src/main/java/com/yd/communication/api/websocket/CoWebSocketServer.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
api
.
websocket
;
import
com.alibaba.fastjson2.JSON
;
import
com.fasterxml.jackson.databind.JsonNode
;
import
com.fasterxml.jackson.databind.ObjectMapper
;
import
com.yd.communication.api.service.ApiCoSessionService
;
import
com.yd.communication.feign.dto.RoomRedisInfoDTO
;
import
com.yd.communication.feign.enums.ControlHolderTypeEnum
;
import
com.yd.communication.feign.request.EndSessionRequest
;
import
com.yd.communication.feign.request.TransferControlRequest
;
import
com.yd.communication.service.model.CoSession
;
import
com.yd.communication.service.service.ICoOperationLogService
;
import
lombok.extern.slf4j.Slf4j
;
import
org.apache.commons.lang3.StringUtils
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.data.redis.core.RedisTemplate
;
import
org.springframework.stereotype.Component
;
import
javax.annotation.PostConstruct
;
import
javax.websocket.*
;
import
javax.websocket.server.PathParam
;
import
javax.websocket.server.ServerEndpoint
;
import
java.io.IOException
;
import
java.net.URLDecoder
;
import
java.nio.charset.StandardCharsets
;
import
java.util.*
;
import
java.util.concurrent.ConcurrentHashMap
;
import
java.util.concurrent.CopyOnWriteArraySet
;
import
java.util.concurrent.TimeUnit
;
/**
* 协同 WebSocket 服务端(支持多节点部署)
* <p>
* 设计思路:
* 1. 本地内存 Map 管理当前节点的 Session(无法跨节点共享)
* 2. Redis 存储房间成员信息、当前状态(跨节点共享)
* 3. Redis Pub/Sub 实现跨节点实时消息广播
* 4. 新节点接入时,从 Redis 拉取最新状态,实现"状态同步"而非"历史重放"
* </p>
*
* @author zxm
* @date 2026-07-28
*/
@Component
// 将当前类注入 Spring 容器,使其成为 Bean
@ServerEndpoint
(
"/ws/{roomId}"
)
// 声明 WebSocket 端点,路径为 /ws/{roomId},{roomId} 是路径参数
@Slf4j
public
class
CoWebSocketServer
{
// ==================== 本地内存(当前节点) ====================
/**
* 房间 -> 当前节点内的 WebSocket 会话集合
* key: 房间号 (roomId)
* value: 该房间内当前节点上的所有会话(Session)
* 使用 ConcurrentHashMap 保证多线程安全,CopyOnWriteArraySet 保证读写分离
*/
private
static
final
Map
<
String
,
Set
<
Session
>>
ROOMS
=
new
ConcurrentHashMap
<>();
/**
* 会话 -> 房间号 反向映射
* key: WebSocket Session 对象
* value: 该会话所在的房间号
* 用于在收到消息时快速查找房间号
*/
private
static
final
Map
<
Session
,
String
>
SESSION_ROOM
=
new
ConcurrentHashMap
<>();
/**
* 会话 -> 用户ID 映射
* key: WebSocket Session 对象
* value: 该会话对应的用户ID
*/
private
static
final
Map
<
Session
,
String
>
SESSION_USER_ID
=
new
ConcurrentHashMap
<>();
/**
* 会话 -> 用户类型 映射
* key: WebSocket Session 对象
* value: 用户类型 (owner: 客户, participant: 顾问)
*/
private
static
final
Map
<
Session
,
String
>
SESSION_USER_TYPE
=
new
ConcurrentHashMap
<>();
// ==================== Spring Bean 静态注入 ====================
/**
* API 层会话服务(静态化,供 WebSocket 生命周期方法使用)
* 因为 WebSocket 端点不由 Spring 直接管理,需要静态注入
*/
private
static
ApiCoSessionService
sessionService
;
/**
* 操作日志服务(静态化)
*/
private
static
ICoOperationLogService
operationLogService
;
/**
* Redis 模板(静态化),用于跨节点通信和状态存储
*/
private
static
RedisTemplate
<
String
,
String
>
redisTemplate
;
/**
* 注入会话服务
* @param service ApiCoSessionService 实例
*/
@Autowired
public
void
setSessionService
(
ApiCoSessionService
service
)
{
CoWebSocketServer
.
sessionService
=
service
;
}
/**
* 注入操作日志服务
* @param service ICoOperationLogService 实例
*/
@Autowired
public
void
setOperationLogService
(
ICoOperationLogService
service
)
{
CoWebSocketServer
.
operationLogService
=
service
;
}
/**
* 注入 Redis 模板
* @param redisTemplate RedisTemplate 实例
*/
@Autowired
public
void
setRedisTemplate
(
RedisTemplate
<
String
,
String
>
redisTemplate
)
{
CoWebSocketServer
.
redisTemplate
=
redisTemplate
;
}
// ==================== Redis 键名常量 ====================
/**
* 房间成员信息 Redis Hash 键模板
* 实际键名:room:members:{roomId}
* 存储结构:Hash,field 为 sessionId,value 为成员信息 JSON
*/
private
static
final
String
ROOM_MEMBERS_KEY
=
"room:members:%s"
;
/**
* 房间当前页面状态 Redis String 键模板
* 实际键名:room:state:{roomId}:currentPage
* 存储内容:当前页面的 JSON 字符串
*/
private
static
final
String
ROOM_STATE_PAGE_KEY
=
"room:state:%s:currentPage"
;
/**
* 房间广播频道 Redis Pub/Sub 频道模板
* 实际频道名:room:channel:{roomId}
* 用于跨节点消息广播
*/
private
static
final
String
ROOM_BROADCAST_CHANNEL
=
"room:channel:%s"
;
// ==================== 新增:节点标识 ====================
/**
* 当前节点的唯一标识(IP:端口 或 UUID)
* 用于区分消息是否为本节点发出,避免回环
*/
private
static
final
String
NODE_ID
=
UUID
.
randomUUID
().
toString
()
+
"@"
+
System
.
currentTimeMillis
();
// ==================== 跨节点订阅(每个节点启动时执行) ====================
/**
* Bean 初始化完成后执行,启动 Redis 订阅线程
* 每个节点启动时都会订阅通配符频道 "room:channel:*"
* 用于接收其他节点发来的跨节点广播消息
*/
@PostConstruct
public
void
init
()
{
// 启动一个独立的后台线程,避免阻塞主线程
new
Thread
(()
->
{
// 无限循环,支持断线自动重连
while
(
true
)
{
try
{
// 通过 Redis 连接执行订阅命令
redisTemplate
.
execute
((
connection
)
->
{
// 订阅所有以 "room:channel:" 开头的频道
// 第二个参数是监听器:收到消息时调用 handleCrossNodeMessage 处理
connection
.
subscribe
(
(
message
,
pattern
)
->
handleCrossNodeMessage
(
new
String
(
message
.
getBody
())),
"room:channel:*"
.
getBytes
()
);
// 返回 null,因为 subscribe 是阻塞方法,执行到这里说明订阅已结束(异常断开)
return
null
;
},
true
);
// true 表示使用事务(此处无实际影响)
}
catch
(
Exception
e
)
{
// 订阅断开(如 Redis 连接超时、网络抖动),记录错误日志
log
.
error
(
"Redis 订阅断开,5秒后重试..."
,
e
);
try
{
// 等待 5 秒后重连
Thread
.
sleep
(
5000
);
}
catch
(
InterruptedException
ex
)
{
// 线程被中断,退出循环
Thread
.
currentThread
().
interrupt
();
break
;
}
}
}
}).
start
();
// 启动线程
log
.
info
(
"Redis 跨节点广播订阅已启动,当前节点ID: {}"
,
NODE_ID
);
}
/**
* 处理跨节点广播消息(由 Redis 订阅触发)
* 当其他节点向 Redis 频道发布消息时,此方法会被调用
* 注意:此方法仅负责将消息转发给本节点的 Session,不处理状态同步
*
* @param body Redis 消息体(JSON 字符串)
*/
private
void
handleCrossNodeMessage
(
String
body
)
{
try
{
// 创建 Jackson 对象映射器,解析 JSON
ObjectMapper
mapper
=
new
ObjectMapper
();
// 将 JSON 字符串解析为树节点
JsonNode
json
=
mapper
.
readTree
(
body
);
// 提取消息中的节点ID
String
sourceNodeId
=
json
.
has
(
"sourceNodeId"
)
?
json
.
get
(
"sourceNodeId"
).
asText
()
:
null
;
// 如果消息是本节点发出的,直接忽略,避免回环
if
(
NODE_ID
.
equals
(
sourceNodeId
))
{
log
.
debug
(
"忽略本节点发出的消息,sourceNodeId={}"
,
sourceNodeId
);
return
;
}
// 提取房间号
String
roomId
=
json
.
get
(
"roomId"
).
asText
();
// 提取消息内容
String
msg
=
json
.
get
(
"message"
).
asText
();
// 提取需要排除的会话 ID(即原始发送者,避免消息回环)
String
excludeSessionId
=
json
.
has
(
"excludeSessionId"
)
?
json
.
get
(
"excludeSessionId"
).
asText
()
:
null
;
// 从本地内存中获取该房间在本节点的所有会话
Set
<
Session
>
sessions
=
ROOMS
.
get
(
roomId
);
// 如果本节点没有该房间的会话,直接忽略(其他节点会处理)
if
(
sessions
==
null
||
sessions
.
isEmpty
())
{
return
;
}
// 遍历该房间在本节点的所有会话
for
(
Session
s
:
sessions
)
{
// 跳过需要排除的会话(即消息发送者所在的节点已处理,本节点不再转发给同一个人)
if
(
excludeSessionId
!=
null
&&
excludeSessionId
.
equals
(
s
.
getId
()))
{
continue
;
}
// 检查会话是否还处于打开状态
if
(
s
.
isOpen
())
{
// 发送消息给客户端
s
.
getBasicRemote
().
sendText
(
msg
);
}
}
}
catch
(
Exception
e
)
{
// 处理消息异常,记录日志但不影响其他节点
log
.
error
(
"处理跨节点消息异常"
,
e
);
}
}
// ==================== WebSocket 生命周期 ====================
/**
* Jackson 对象映射器,用于解析 JSON 消息
*/
private
final
ObjectMapper
objectMapper
=
new
ObjectMapper
();
/**
* WebSocket 连接建立时触发
* 1. 解析 URL 查询参数(userId, userType)
* 2. 将当前会话加入本地内存(ROOMS, SESSION_ROOM 等)
* 3. 将成员信息存入 Redis(跨节点共享)
* 4. 从 Redis 拉取当前房间状态并同步给新连接的用户(状态同步)
*
* @param session 当前 WebSocket 会话对象
* @param roomId 路径参数:房间号
*/
@OnOpen
public
void
onOpen
(
Session
session
,
@PathParam
(
"roomId"
)
String
roomId
)
{
// 1. 获取 URL 查询字符串(如 ?userId=123&userType=owner)
String
queryString
=
session
.
getQueryString
();
// 解析查询参数为 Map
Map
<
String
,
String
>
params
=
parseQueryString
(
queryString
);
// 获取用户ID,默认为 "unknown"
String
userId
=
params
.
getOrDefault
(
"userId"
,
"unknown"
);
// 获取用户类型,默认为 "unknown"(owner: 客户, participant: 顾问)
String
userType
=
params
.
getOrDefault
(
"userType"
,
"unknown"
);
// 2. 将当前会话加入本地内存
// computeIfAbsent: 如果房间不存在则创建新的 Set 集合
ROOMS
.
computeIfAbsent
(
roomId
,
k
->
new
CopyOnWriteArraySet
<>()).
add
(
session
);
// 记录会话对应的房间号
SESSION_ROOM
.
put
(
session
,
roomId
);
// 记录会话对应的用户ID
SESSION_USER_ID
.
put
(
session
,
userId
);
// 记录会话对应的用户类型
SESSION_USER_TYPE
.
put
(
session
,
userType
);
// 3. 将成员信息存入 Redis(供跨节点查询使用)
// 构造 Redis Hash 的键名:room:members:{roomId}
String
memberKey
=
String
.
format
(
ROOM_MEMBERS_KEY
,
roomId
);
// 构造成员信息 Map
Map
<
String
,
String
>
memberInfo
=
new
HashMap
<>();
memberInfo
.
put
(
"userId"
,
userId
);
memberInfo
.
put
(
"userType"
,
userType
);
memberInfo
.
put
(
"sessionId"
,
session
.
getId
());
// sessionId 作为唯一标识
// 以 sessionId 为 field,成员信息 JSON 为 value,存入 Redis Hash
redisTemplate
.
opsForHash
().
put
(
memberKey
,
session
.
getId
(),
JSON
.
toJSONString
(
memberInfo
));
// 设置 Hash 过期时间为 60 分钟(与会话超时一致)
redisTemplate
.
expire
(
memberKey
,
60
,
TimeUnit
.
MINUTES
);
// 4. 状态同步(核心功能):向新用户推送当前房间的完整状态
try
{
// 4.1 获取当前控制权信息(从 Redis 缓存中读取)
RoomRedisInfoDTO
dto
=
sessionService
.
getCurrentController
(
roomId
);
// 控制权持有者类型:owner 或 participant
String
controlHolderType
=
dto
.
getControlHolderType
();
// 控制权持有者ID
String
controlHolderId
=
dto
.
getControlHolderId
();
// 4.2 获取当前页面状态(从 Redis 中读取,由翻页操作实时更新)
String
stateKey
=
String
.
format
(
ROOM_STATE_PAGE_KEY
,
roomId
);
String
currentPageJson
=
redisTemplate
.
opsForValue
().
get
(
stateKey
);
// 如果 Redis 中没有页面状态,则从数据库读取最新状态
if
(
StringUtils
.
isBlank
(
currentPageJson
))
{
// 根据房间号查询数据库中的会话记录
CoSession
coSession
=
sessionService
.
getByRoomId
(
roomId
);
// 如果会话存在,取其中的当前页面 JSON
currentPageJson
=
coSession
!=
null
?
coSession
.
getCurrentPage
()
:
"{}"
;
// 将数据库中的状态回填到 Redis,供后续新节点快速同步
if
(
StringUtils
.
isNotBlank
(
currentPageJson
))
{
redisTemplate
.
opsForValue
().
set
(
stateKey
,
currentPageJson
,
60
,
TimeUnit
.
MINUTES
);
}
}
// 4.3 组装初始化同步消息(包含控制权和当前页面)
// 消息格式:{"action":"init_sync","holderType":"xxx","holderId":"xxx","currentPage":{...}}
String
initMsg
=
String
.
format
(
"{\"action\":\"init_sync\",\"holderType\":\"%s\",\"holderId\":\"%s\",\"currentPage\":%s}"
,
controlHolderType
,
controlHolderId
,
currentPageJson
);
// 发送初始化同步消息给刚连接的客户端
session
.
getBasicRemote
().
sendText
(
initMsg
);
// 4.4 单独再推送一次控制权消息(兼容前端只监听 control_transfer 的情况)
String
controlMsg
=
String
.
format
(
"{\"action\":\"control_transfer\",\"holderType\":\"%s\",\"holderId\":\"%s\"}"
,
controlHolderType
,
controlHolderId
);
session
.
getBasicRemote
().
sendText
(
controlMsg
);
// 记录日志:用户已加入并完成状态同步
log
.
info
(
"用户 {} 加入房间 {},已同步状态"
,
userId
,
roomId
);
}
catch
(
Exception
e
)
{
// 状态同步失败,记录错误但不影响连接建立
log
.
error
(
"状态同步失败,用户 {} 可能无法恢复最新状态"
,
userId
,
e
);
}
// 记录连接建立日志,包含当前节点该房间的会话数
log
.
info
(
"用户 {} ({}) 加入房间 {}"
,
userId
,
userType
,
roomId
);
}
/**
* 接收客户端发送的 WebSocket 消息(核心业务分发器)
* 根据 action 类型分发到不同的业务处理逻辑
*
* @param session 当前会话
* @param message 客户端发送的 JSON 字符串
*/
@OnMessage
public
void
onMessage
(
Session
session
,
String
message
)
{
// 根据会话查询对应的房间号
String
roomId
=
SESSION_ROOM
.
get
(
session
);
// 如果会话未关联房间,忽略消息
if
(
roomId
==
null
)
{
return
;
}
try
{
// 1. 解析 JSON 消息
JsonNode
json
=
objectMapper
.
readTree
(
message
);
// 获取操作类型
String
action
=
json
.
get
(
"action"
).
asText
();
// 获取当前操作用户ID
String
userId
=
SESSION_USER_ID
.
get
(
session
);
// 获取当前操作用户类型
String
userType
=
SESSION_USER_TYPE
.
get
(
session
);
// 2. 获取会话业务ID(用于操作日志记录)
CoSession
coSession
=
sessionService
.
getByRoomId
(
roomId
);
String
bizId
=
coSession
!=
null
?
coSession
.
getCoSessionBizId
()
:
null
;
// ========== 业务逻辑分发 ==========
// --- 操作1:控制权切换(仅顾问可操作) ---
if
(
"control_transfer"
.
equals
(
action
))
{
// 校验权限:只有参与者(顾问)才能切换控制权
if
(!
"participant"
.
equals
(
userType
))
{
sendError
(
session
,
"只有顾问可以切换控制权"
);
return
;
}
// 解析顾问切换后的控制权信息
String
newHolderType
=
json
.
get
(
"holderType"
).
asText
();
String
newHolderId
=
json
.
get
(
"holderId"
).
asText
();
// 构造控制权切换请求
TransferControlRequest
req
=
new
TransferControlRequest
();
req
.
setRoomId
(
roomId
);
// 判断切换后的持有者类型:如果是 owner(客户)则 opType=1(开启客户操作),否则 opType=2(关闭客户操作)
req
.
setOprType
(
ControlHolderTypeEnum
.
OWNER
.
getItemValue
().
equals
(
newHolderType
)
?
1
:
2
);
// 调用服务层切换控制权(更新数据库 + Redis)
sessionService
.
transferControl
(
req
);
// 广播控制权变更消息给房间内所有用户(包括自己)
broadcast
(
roomId
,
message
,
session
);
// 记录操作日志(用于合规审计)
if
(
bizId
!=
null
)
{
operationLogService
.
log
(
bizId
,
userId
,
userType
,
userId
,
action
,
message
,
null
,
null
);
}
log
.
info
(
"控制权切换成功: roomId={}, holder={}:{}"
,
roomId
,
newHolderType
,
newHolderId
);
return
;
}
// --- 操作2:结束共享(仅客户可操作) ---
if
(
"end_sharing"
.
equals
(
action
))
{
// 校验权限:只有所有者(客户)才能结束共享
if
(!
"owner"
.
equals
(
userType
))
{
sendError
(
session
,
"只有客户可以结束共享"
);
return
;
}
// 构造结束会话请求
EndSessionRequest
req
=
new
EndSessionRequest
();
req
.
setRoomId
(
roomId
);
// 调用服务层结束会话(更新状态 + 停止录制 + 清理Redis)
sessionService
.
end
(
req
);
// 广播结束消息给所有用户
broadcast
(
roomId
,
"{\"action\":\"end_sharing\"}"
,
null
);
// 清理 Redis 中的页面状态
redisTemplate
.
delete
(
String
.
format
(
ROOM_STATE_PAGE_KEY
,
roomId
));
// 关闭房间所有连接并清理资源
closeRoom
(
roomId
);
log
.
info
(
"共享已结束: roomId={}"
,
roomId
);
return
;
}
// --- 操作3:脱敏开关(客户操作,全房间生效) ---
if
(
"desensitization_switch"
.
equals
(
action
))
{
// 广播脱敏状态给所有人(同步显示脱敏效果)
broadcast
(
roomId
,
message
,
session
);
// 记录操作日志
if
(
bizId
!=
null
)
{
operationLogService
.
log
(
bizId
,
userId
,
userType
,
userId
,
action
,
message
,
null
,
null
);
}
return
;
}
// --- 操作4:同步操作(翻页/滚动/缩放)——仅控制权持有者可操作 ---
if
(
"turn_the_page"
.
equals
(
action
)
||
"scroll"
.
equals
(
action
)
||
"scaling"
.
equals
(
action
))
{
// 4.1 获取当前控制权持有者信息
RoomRedisInfoDTO
dto
=
sessionService
.
getCurrentController
(
roomId
);
String
holderType
=
dto
.
getControlHolderType
();
String
holderId
=
dto
.
getControlHolderId
();
// 4.2 校验:当前用户是否持有控制权
if
(!(
holderType
.
equals
(
userType
)
&&
holderId
.
equals
(
userId
)))
{
sendError
(
session
,
"您没有控制权,无法操作"
);
return
;
}
// 4.3 如果是翻页或滚动操作,需要更新数据库和 Redis 状态
if
(
"turn_the_page"
.
equals
(
action
)
||
"scroll"
.
equals
(
action
))
{
// 获取当前页面 JSON(从消息中提取)
String
currentPage
=
json
.
has
(
"currentPage"
)
?
json
.
get
(
"currentPage"
).
toString
()
:
"{}"
;
// 更新数据库中的当前页面和页面历史轨迹
sessionService
.
updateCurrentPage
(
roomId
,
currentPage
,
userId
);
// 更新 Redis 状态(供新节点连接时同步)
String
stateKey
=
String
.
format
(
ROOM_STATE_PAGE_KEY
,
roomId
);
redisTemplate
.
opsForValue
().
set
(
stateKey
,
currentPage
,
60
,
TimeUnit
.
MINUTES
);
}
// 4.4 广播操作指令给房间内其他用户(排除自己,避免回环)
broadcast
(
roomId
,
message
,
session
);
// 4.5 记录操作日志(合规审计)
if
(
bizId
!=
null
)
{
operationLogService
.
log
(
bizId
,
userId
,
userType
,
userId
,
action
,
message
,
null
,
null
);
}
return
;
}
// --- 未知操作:记录警告日志 ---
log
.
warn
(
"未知 action: {}"
,
action
);
}
catch
(
Exception
e
)
{
// 处理消息异常,向客户端返回错误信息
log
.
error
(
"处理消息异常"
,
e
);
sendError
(
session
,
"服务器处理异常"
);
}
}
/**
* WebSocket 连接关闭时触发
* 1. 从本地内存中移除该会话
* 2. 从 Redis 中移除该成员信息
* 3. 如果房间为空,清理本地房间映射
*
* @param session 关闭的会话
*/
@OnClose
public
void
onClose
(
Session
session
)
{
// 从反向映射中移除会话,获取该会话对应的房间号
String
roomId
=
SESSION_ROOM
.
remove
(
session
);
if
(
roomId
!=
null
)
{
// 从本地内存中获取该房间的会话集合
Set
<
Session
>
sessions
=
ROOMS
.
get
(
roomId
);
if
(
sessions
!=
null
)
{
// 从集合中移除当前会话
sessions
.
remove
(
session
);
// 如果集合为空,移除该房间的映射,释放内存
if
(
sessions
.
isEmpty
())
{
ROOMS
.
remove
(
roomId
);
}
}
// 从 Redis 中移除该成员信息
String
memberKey
=
String
.
format
(
ROOM_MEMBERS_KEY
,
roomId
);
redisTemplate
.
opsForHash
().
delete
(
memberKey
,
session
.
getId
());
}
// 移除会话对应的用户ID映射
SESSION_USER_ID
.
remove
(
session
);
// 移除会话对应的用户类型映射
SESSION_USER_TYPE
.
remove
(
session
);
// 记录连接关闭日志
log
.
info
(
"WebSocket 关闭: sessionId={}"
,
session
.
getId
());
}
/**
* WebSocket 发生异常时触发
*
* @param session 发生异常的会话
* @param error 异常信息
*/
@OnError
public
void
onError
(
Session
session
,
Throwable
error
)
{
log
.
error
(
"WebSocket 错误"
,
error
);
}
// ==================== 广播方法 ====================
/**
* 向房间内所有成员广播消息(支持跨节点)
* 1. 本节点直接发送:遍历本地 ROOMS 中的会话,直接发送消息
* 2. 跨节点广播:将消息发布到 Redis,其他节点的订阅者收到后会转发给各自的客户端
*
* @param roomId 房间号
* @param message JSON 消息内容
* @param exclude 本节点需要排除的会话(通常为消息发送者,避免回环)
*/
private
void
broadcast
(
String
roomId
,
String
message
,
Session
exclude
)
{
// 1. 本节点直接发送消息给当前节点内的所有会话
Set
<
Session
>
localSessions
=
ROOMS
.
get
(
roomId
);
if
(
localSessions
!=
null
&&
!
localSessions
.
isEmpty
())
{
for
(
Session
s
:
localSessions
)
{
// 跳过需要排除的会话(发送者自己)
if
(
s
!=
exclude
&&
s
.
isOpen
())
{
try
{
// 发送消息
s
.
getBasicRemote
().
sendText
(
message
);
}
catch
(
IOException
e
)
{
// 发送失败,忽略(日志不打印,避免刷屏)
}
}
}
}
// 2. 发布到 Redis,通知其他节点广播
String
channel
=
String
.
format
(
ROOM_BROADCAST_CHANNEL
,
roomId
);
Map
<
String
,
String
>
data
=
new
HashMap
<>();
data
.
put
(
"roomId"
,
roomId
);
// 房间号
data
.
put
(
"message"
,
message
);
// 消息内容
data
.
put
(
"sourceNodeId"
,
NODE_ID
);
// 标记消息来源节点
if
(
exclude
!=
null
)
{
data
.
put
(
"excludeSessionId"
,
exclude
.
getId
());
// 需要排除的 sessionId
}
// 将数据转为 JSON 并发布到 Redis 频道
redisTemplate
.
convertAndSend
(
channel
,
JSON
.
toJSONString
(
data
));
}
/**
* 向指定会话发送错误消息
*
* @param session 目标会话
* @param errorMsg 错误描述信息
*/
private
void
sendError
(
Session
session
,
String
errorMsg
)
{
try
{
// 发送 JSON 格式的错误消息
session
.
getBasicRemote
().
sendText
(
"{\"error\":\""
+
errorMsg
+
"\"}"
);
}
catch
(
IOException
e
)
{
// 发送失败,忽略
}
}
/**
* 关闭房间所有连接并清理 Redis 状态
* 在结束共享时调用
*
* @param roomId 房间号
*/
private
void
closeRoom
(
String
roomId
)
{
// 1. 从本地内存移除该房间并获取所有会话
Set
<
Session
>
sessions
=
ROOMS
.
remove
(
roomId
);
if
(
sessions
!=
null
)
{
// 遍历所有会话,逐个关闭连接
for
(
Session
s
:
sessions
)
{
try
{
// 正常关闭连接,附带关闭原因
s
.
close
(
new
CloseReason
(
CloseReason
.
CloseCodes
.
NORMAL_CLOSURE
,
"会话结束"
));
}
catch
(
IOException
e
)
{
// 关闭失败,忽略
}
}
}
// 2. 清理 Redis 中的成员信息 Hash
redisTemplate
.
delete
(
String
.
format
(
ROOM_MEMBERS_KEY
,
roomId
));
// 3. 清理 Redis 中的页面状态
redisTemplate
.
delete
(
String
.
format
(
ROOM_STATE_PAGE_KEY
,
roomId
));
// 注意:控制权缓存由 sessionService.end() 负责清理
}
// ==================== 工具方法 ====================
/**
* 解析 URL 查询参数字符串
* 示例输入: "userId=123&userType=owner"
* 示例输出: {"userId": "123", "userType": "owner"}
*
* @param queryString 原始查询字符串(不含问号)
* @return 解析后的参数键值对 Map
*/
private
Map
<
String
,
String
>
parseQueryString
(
String
queryString
)
{
Map
<
String
,
String
>
result
=
new
HashMap
<>();
// 如果查询字符串为空,返回空 Map
if
(
queryString
==
null
||
queryString
.
isEmpty
())
{
return
result
;
}
// 按 & 符号分割多个参数
for
(
String
param
:
queryString
.
split
(
"&"
))
{
// 按 = 分割键和值
String
[]
pair
=
param
.
split
(
"="
,
2
);
// 解码键名
String
key
=
pair
.
length
>
0
?
decode
(
pair
[
0
])
:
""
;
// 解码值
String
value
=
pair
.
length
>
1
?
decode
(
pair
[
1
])
:
""
;
// 如果键不为空,存入 Map
if
(!
key
.
isEmpty
())
{
result
.
put
(
key
,
value
);
}
}
return
result
;
}
/**
* URL 解码(UTF-8)
*
* @param value 需要解码的字符串
* @return 解码后的字符串,如果解码失败则返回原值
*/
private
String
decode
(
String
value
)
{
try
{
// 使用 UTF-8 进行 URL 解码
return
URLDecoder
.
decode
(
value
,
StandardCharsets
.
UTF_8
.
name
());
}
catch
(
Exception
e
)
{
// 解码失败,返回原值
return
value
;
}
}
}
\ No newline at end of file
yd-communication-feign/src/main/java/com/yd/communication/feign/client/ApiCoSessionFeignClient.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
client
;
import
com.yd.common.result.Result
;
import
com.yd.communication.feign.fallback.ApiCoSessionFeignFallbackFactory
;
import
com.yd.communication.feign.request.CreateRequest
;
import
com.yd.communication.feign.request.EndSessionRequest
;
import
com.yd.communication.feign.request.JoinRequest
;
import
com.yd.communication.feign.request.TransferControlRequest
;
import
com.yd.communication.feign.response.CommonResponse
;
import
com.yd.communication.feign.response.CreateResponse
;
import
com.yd.communication.feign.response.JoinResponse
;
import
com.yd.communication.feign.response.SessionDetailResponse
;
import
org.springframework.cloud.openfeign.FeignClient
;
import
org.springframework.validation.annotation.Validated
;
import
org.springframework.web.bind.annotation.*
;
/**
* 通信服务-协同会话信息 Feign 客户端
*
* @author zxm
* @date 2026-07-28
*/
@FeignClient
(
name
=
"yd-communication-api"
,
path
=
"/communication/api/coSession"
,
fallbackFactory
=
ApiCoSessionFeignFallbackFactory
.
class
)
public
interface
ApiCoSessionFeignClient
{
/**
* 客户创建协同会话(生成共享码)
*
* @param request 创建会话请求
* @return 会话信息
*/
@PostMapping
(
"/create"
)
Result
<
CreateResponse
>
create
(
@Validated
@RequestBody
CreateRequest
request
);
/**
* 顾问加入协同会话(输入共享码)
*
* @param request 加入会话请求
* @return 会话详情
*/
@PostMapping
(
"/join"
)
Result
<
JoinResponse
>
join
(
@Validated
@RequestBody
JoinRequest
request
);
/**
* 获取协同会话详情
*
* @param bizId 会话业务ID
* @return 完整会话信息
*/
@GetMapping
(
"/{bizId}"
)
Result
<
SessionDetailResponse
>
get
(
@PathVariable
(
"bizId"
)
String
bizId
);
/**
* 结束协同会话(关闭共享,仅客户可调用)
*
*/
@PostMapping
(
"/end"
)
Result
<
CommonResponse
>
end
(
@Validated
@RequestBody
EndSessionRequest
request
);
/**
* 切换控制权(仅参与者(顾问)可调用)
* @param request
* @return
*/
@PostMapping
(
"/control/transfer"
)
Result
<
CommonResponse
>
transferControl
(
@Validated
@RequestBody
TransferControlRequest
request
);
}
\ No newline at end of file
yd-communication-feign/src/main/java/com/yd/communication/feign/constant/RedisConstants.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
constant
;
/**
* redis的key前缀常量
*/
public
class
RedisConstants
{
/**
* 协同房间缓存信息redis前缀
*/
public
static
final
String
ROOM_KEY_PREFIX
=
"room:"
;
}
yd-communication-feign/src/main/java/com/yd/communication/feign/dto/RoomRedisInfoDTO.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
dto
;
import
lombok.Data
;
/**
* 存储当前房间内的一些缓存字段信息
*/
@Data
public
class
RoomRedisInfoDTO
{
/**
* 控制权持有者类型:owner(资源所有者类型)/participant(参与者类型)
*/
private
String
controlHolderType
;
/**
* 控制权持有者ID(具体人的ID)
*/
private
String
controlHolderId
;
/**
* 资源所有者的登录用户ID
*/
private
String
userId
;
/**
* 资源所有者的登录token信息
*/
private
String
token
;
}
yd-communication-feign/src/main/java/com/yd/communication/feign/enums/CoSessionStatusEnum.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
enums
;
/**
* 协同会话状态枚举
*/
public
enum
CoSessionStatusEnum
{
DKS
(
"待开始"
,
"1"
),
JXZ
(
"进行中"
,
"2"
),
YJS
(
"已结束"
,
"3"
),
YCS
(
"已超时"
,
"4"
),
;
//字典项标签(名称)
private
String
itemLabel
;
//字典项值
private
String
itemValue
;
//构造函数
CoSessionStatusEnum
(
String
itemLabel
,
String
itemValue
)
{
this
.
itemLabel
=
itemLabel
;
this
.
itemValue
=
itemValue
;
}
public
String
getItemLabel
()
{
return
itemLabel
;
}
public
String
getItemValue
()
{
return
itemValue
;
}
}
yd-communication-feign/src/main/java/com/yd/communication/feign/enums/ControlHolderTypeEnum.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
enums
;
/**
* 控制权持有者类型枚举
*/
public
enum
ControlHolderTypeEnum
{
OWNER
(
"资源所有者类型"
,
"owner"
),
PARTICIPANT
(
"参与者类型"
,
"participant"
),
;
//字典项标签(名称)
private
String
itemLabel
;
//字典项值
private
String
itemValue
;
//构造函数
ControlHolderTypeEnum
(
String
itemLabel
,
String
itemValue
)
{
this
.
itemLabel
=
itemLabel
;
this
.
itemValue
=
itemValue
;
}
public
String
getItemLabel
()
{
return
itemLabel
;
}
public
String
getItemValue
()
{
return
itemValue
;
}
}
yd-communication-feign/src/main/java/com/yd/communication/feign/enums/RedisEnum.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
enums
;
import
com.yd.communication.feign.constant.RedisConstants
;
import
java.util.concurrent.TimeUnit
;
/**
* redis枚举
*/
public
enum
RedisEnum
{
//协同房间缓存信息redis参数
ROOM
(
RedisConstants
.
ROOM_KEY_PREFIX
,
600000000
,
TimeUnit
.
MINUTES
),
;
//缓存key前缀
private
String
prefix
;
//缓存过期时长
private
Integer
timeout
;
//缓存过期时长单位
private
TimeUnit
timeUnit
;
RedisEnum
(
String
prefix
,
Integer
timeout
,
TimeUnit
timeUnit
)
{
this
.
prefix
=
prefix
;
this
.
timeout
=
timeout
;
this
.
timeUnit
=
timeUnit
;
}
public
String
getPrefix
()
{
return
prefix
;
}
public
Integer
getTimeout
()
{
return
timeout
;
}
public
TimeUnit
getTimeUnit
()
{
return
timeUnit
;
}
}
yd-communication-feign/src/main/java/com/yd/communication/feign/fallback/ApiCoSessionFeignFallbackFactory.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
fallback
;
import
com.yd.common.result.Result
;
import
com.yd.communication.feign.client.ApiCoSessionFeignClient
;
import
com.yd.communication.feign.request.CreateRequest
;
import
com.yd.communication.feign.request.EndSessionRequest
;
import
com.yd.communication.feign.request.JoinRequest
;
import
com.yd.communication.feign.request.TransferControlRequest
;
import
com.yd.communication.feign.response.CommonResponse
;
import
com.yd.communication.feign.response.CreateResponse
;
import
com.yd.communication.feign.response.JoinResponse
;
import
com.yd.communication.feign.response.SessionDetailResponse
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.cloud.openfeign.FallbackFactory
;
import
org.springframework.stereotype.Component
;
/**
* 通信服务-协同会话信息Feign降级处理
*/
@Slf4j
@Component
public
class
ApiCoSessionFeignFallbackFactory
implements
FallbackFactory
<
ApiCoSessionFeignClient
>
{
@Override
public
ApiCoSessionFeignClient
create
(
Throwable
cause
)
{
return
new
ApiCoSessionFeignClient
()
{
@Override
public
Result
<
CreateResponse
>
create
(
CreateRequest
request
)
{
return
null
;
}
@Override
public
Result
<
JoinResponse
>
join
(
JoinRequest
request
)
{
return
null
;
}
@Override
public
Result
<
SessionDetailResponse
>
get
(
String
bizId
)
{
return
null
;
}
@Override
public
Result
<
CommonResponse
>
end
(
EndSessionRequest
request
)
{
return
null
;
}
@Override
public
Result
<
CommonResponse
>
transferControl
(
TransferControlRequest
request
)
{
return
null
;
}
};
}
}
yd-communication-feign/src/main/java/com/yd/communication/feign/request/CreateRequest.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
request
;
import
lombok.Data
;
import
javax.validation.constraints.NotBlank
;
/**
* 创建协同会话请求对象
* 客户调用此接口生成共享码,开启协同讲解
*/
@Data
public
class
CreateRequest
{
/**
* 协同作用域
* single: 单个资源(如单份报告、资讯)
* global: 全局协同(如小程序全局协同,切换页面自动跟随)
*/
@NotBlank
(
message
=
"协同作用域不能为空"
)
private
String
scope
;
/**
* 资源类型
* 用于区分不同的业务资源类型
* 可选值:report(报告)、news(资讯)、mini_program(小程序)等
*/
@NotBlank
(
message
=
"资源类型不能为空"
)
private
String
resourceType
;
/**
* 资源业务ID
* 具体资源的唯一标识
* 如:report-报告ID、news-资讯ID、mini_program-小程序应用标识
*/
@NotBlank
(
message
=
"资源业务ID不能为空"
)
private
String
resourceId
;
/**
* 资源初始化JSON串
* 创建会话时记录当前所在页面的完整信息,后续不做修改,用于历史追溯
* 示例:{"url":"https://mini.xxx.com/pages/index/index?userId=xxx"}
*/
@NotBlank
(
message
=
"资源初始化JSON串不能为空"
)
private
String
resourceInit
;
/**
* 资源所有者ID(即客户ID)
* 报告/资源的归属人,通常为发起协同的客户
*/
@NotBlank
(
message
=
"资源所有者ID不能为空"
)
private
String
ownerId
;
/**
* 资源所有者类型
* 默认:customer(客户)
* 可扩展:member(会员)、user(普通用户)等
*/
@NotBlank
(
message
=
"资源所有者类型不能为空"
)
private
String
ownerType
;
/**
* 资源所有者的登录用户ID
*/
@NotBlank
(
message
=
"资源所有者的登录用户ID不能为空"
)
private
String
userId
;
/**
* 资源所有者的登录token信息
*/
@NotBlank
(
message
=
"资源所有者的登录token信息不能为空"
)
private
String
token
;
}
\ No newline at end of file
yd-communication-feign/src/main/java/com/yd/communication/feign/request/EndSessionRequest.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
request
;
import
lombok.Data
;
@Data
public
class
EndSessionRequest
{
/**
* 房间号不能为空
*/
private
String
roomId
;
}
yd-communication-feign/src/main/java/com/yd/communication/feign/request/JoinRequest.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
request
;
import
lombok.Data
;
import
javax.validation.constraints.NotBlank
;
/**
* 加入协同会话请求对象
* 顾问输入房间号和共享码加入协同
*/
@Data
public
class
JoinRequest
{
// /**
// * 协同房间号(房间ID)
// * 客户创建会话时生成的唯一房间标识
// * 示例:room_abc12345
// */
// private String roomId;
/**
* 协同房间密码(共享码)
* 客户生成共享会话时返回的6位数字密码
* 顾问需要输入此密码才能加入协同
*/
@NotBlank
(
message
=
"房间密码不能为空"
)
private
String
roomPwd
;
/**
* 参与者ID(即顾问ID)
* 加入协同的顾问/专家的唯一标识
*/
@NotBlank
(
message
=
"参与者ID不能为空"
)
private
String
participantId
;
/**
* 参与者类型
* 默认:consultant(顾问)
* 可扩展:expert(专家)、trainer(培训师)等
*/
@NotBlank
(
message
=
"参与者类型不能为空"
)
private
String
participantType
;
}
\ No newline at end of file
yd-communication-feign/src/main/java/com/yd/communication/feign/request/TransferControlRequest.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
request
;
import
lombok.Data
;
import
javax.validation.constraints.NotBlank
;
import
javax.validation.constraints.NotNull
;
@Data
public
class
TransferControlRequest
{
/**
* 操作类型:1-开启客户操作 2-关闭客户操作
*/
@NotNull
(
message
=
"操作类型不能为空"
)
private
Integer
oprType
;
/**
* 房间号
*/
@NotBlank
(
message
=
"房间号不能为空"
)
private
String
roomId
;
}
yd-communication-feign/src/main/java/com/yd/communication/feign/response/CommonResponse.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
response
;
import
lombok.AllArgsConstructor
;
import
lombok.Data
;
import
lombok.NoArgsConstructor
;
/**
* 通用操作响应对象
* 用于结束会话、切换控制权等无需返回业务数据的操作
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public
class
CommonResponse
{
/**
* 操作结果消息
*/
private
String
message
;
}
\ No newline at end of file
yd-communication-feign/src/main/java/com/yd/communication/feign/response/CreateResponse.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
response
;
import
lombok.Data
;
/**
* 创建协同会话响应对象
*/
@Data
public
class
CreateResponse
{
/**
* 会话唯一业务ID
*/
private
String
sessionBizId
;
/**
* 房间号
*/
private
String
roomId
;
/**
* 共享码(6位数字)
*/
private
String
roomPwd
;
/**
* 会话状态:0-进行中,1-已结束,2-已超时
*/
private
String
status
;
}
\ No newline at end of file
yd-communication-feign/src/main/java/com/yd/communication/feign/response/JoinResponse.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
response
;
import
lombok.Data
;
import
javax.validation.constraints.NotBlank
;
/**
* 加入协同会话响应对象
*/
@Data
public
class
JoinResponse
{
/**
* 会话唯一业务ID
*/
private
String
sessionBizId
;
/**
* 房间号
*/
private
String
roomId
;
/**
* 资源初始化JSON(用于初始加载)
*/
private
String
resourceInit
;
/**
* 控制权持有者类型:owner-客户,participant-顾问
*/
private
String
controlHolderType
;
/**
* 控制权持有者ID
*/
private
String
controlHolderId
;
/**
* 资源所有者的登录用户ID
*/
private
String
userId
;
/**
* 资源所有者的登录token信息
*/
private
String
token
;
}
\ No newline at end of file
yd-communication-feign/src/main/java/com/yd/communication/feign/response/SessionDetailResponse.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
feign
.
response
;
import
lombok.Data
;
import
java.time.LocalDateTime
;
/**
* 协同会话详情响应对象
*/
@Data
public
class
SessionDetailResponse
{
/**
* 数据库主键
*/
private
Long
id
;
/**
* 协同-会话表唯一业务ID
*/
private
String
coSessionBizId
;
/**
* 会话编号
*/
private
String
coSessionNo
;
/**
* 协同作用域:single-单资源,global-全域
*/
private
String
scope
;
/**
* 资源类型:report/news/mini_program
*/
private
String
resourceType
;
/**
* 资源业务ID(如报告ID、小程序标识等)
*/
private
String
resourceId
;
/**
* 资源初始化JSON(创建时页面快照)
*/
private
String
resourceInit
;
/**
* 所有者类型:customer-客户
*/
private
String
ownerType
;
/**
* 所有者ID(客户ID)
*/
private
String
ownerId
;
/**
* 参与者类型:consultant-顾问
*/
private
String
participantType
;
/**
* 参与者ID(顾问ID)
*/
private
String
participantId
;
/**
* 房间号
*/
private
String
roomId
;
/**
* 控制权持有者类型:owner-客户,participant-顾问
*/
private
String
controlHolderType
;
/**
* 控制权持有者ID
*/
private
String
controlHolderId
;
/**
* 会话状态:0-进行中,1-已结束,2-已超时
*/
private
String
status
;
/**
* 会话开始时间
*/
private
LocalDateTime
startTime
;
/**
* 会话结束时间
*/
private
LocalDateTime
endTime
;
/**
* 当前页面JSON(实时更新)
*/
private
String
currentPage
;
/**
* 页面访问历史轨迹(JSON数组)
*/
private
String
pageHistory
;
}
\ No newline at end of file
yd-communication-service/pom.xml
View file @
76ebc4a4
...
@@ -50,6 +50,12 @@
...
@@ -50,6 +50,12 @@
<artifactId>
freemarker
</artifactId>
<artifactId>
freemarker
</artifactId>
</dependency>
</dependency>
<!-- Spring Boot Starter WebSocket -->
<dependency>
<groupId>
org.springframework.boot
</groupId>
<artifactId>
spring-boot-starter-websocket
</artifactId>
</dependency>
<dependency>
<dependency>
<groupId>
com.yd
</groupId>
<groupId>
com.yd
</groupId>
<artifactId>
yd-communication-feign
</artifactId>
<artifactId>
yd-communication-feign
</artifactId>
...
...
yd-communication-service/src/main/java/com/yd/communication/service/dao/CoSessionMapper.java
View file @
76ebc4a4
...
@@ -2,6 +2,10 @@ package com.yd.communication.service.dao;
...
@@ -2,6 +2,10 @@ package com.yd.communication.service.dao;
import
com.yd.communication.service.model.CoSession
;
import
com.yd.communication.service.model.CoSession
;
import
com.baomidou.mybatisplus.core.mapper.BaseMapper
;
import
com.baomidou.mybatisplus.core.mapper.BaseMapper
;
import
org.apache.ibatis.annotations.Param
;
import
org.apache.ibatis.annotations.Update
;
import
java.time.LocalDateTime
;
/**
/**
* <p>
* <p>
...
@@ -13,4 +17,14 @@ import com.baomidou.mybatisplus.core.mapper.BaseMapper;
...
@@ -13,4 +17,14 @@ import com.baomidou.mybatisplus.core.mapper.BaseMapper;
*/
*/
public
interface
CoSessionMapper
extends
BaseMapper
<
CoSession
>
{
public
interface
CoSessionMapper
extends
BaseMapper
<
CoSession
>
{
@Update
(
"UPDATE co_session SET status = #{status}, end_time = #{endTime} WHERE room_id = #{roomId}"
)
int
updateStatusAndEndTimeByRoomId
(
@Param
(
"roomId"
)
String
roomId
,
@Param
(
"status"
)
Integer
status
,
@Param
(
"endTime"
)
LocalDateTime
endTime
);
@Update
(
"UPDATE co_session SET current_page = #{currentPage}, page_history = #{pageHistory}, updater_id = #{updaterId} WHERE room_id = #{roomId}"
)
int
updateCurrentPageAndHistory
(
@Param
(
"roomId"
)
String
roomId
,
@Param
(
"currentPage"
)
String
currentPage
,
@Param
(
"pageHistory"
)
String
pageHistory
,
@Param
(
"updaterId"
)
String
updaterId
);
}
}
yd-communication-service/src/main/java/com/yd/communication/service/model/CoSession.java
View file @
76ebc4a4
...
@@ -121,10 +121,10 @@ public class CoSession implements Serializable {
...
@@ -121,10 +121,10 @@ public class CoSession implements Serializable {
private
String
controlHolderId
;
private
String
controlHolderId
;
/**
/**
*
0-进行中,1-已结束,2
-已超时
*
1-待开始,2-进行中,3-已结束,4
-已超时
*/
*/
@TableField
(
"status"
)
@TableField
(
"status"
)
private
Integer
status
;
private
String
status
;
/**
/**
* 开始时间
* 开始时间
...
@@ -185,4 +185,5 @@ public class CoSession implements Serializable {
...
@@ -185,4 +185,5 @@ public class CoSession implements Serializable {
*/
*/
@TableField
(
"update_time"
)
@TableField
(
"update_time"
)
private
LocalDateTime
updateTime
;
private
LocalDateTime
updateTime
;
}
}
yd-communication-service/src/main/java/com/yd/communication/service/model/RecordingTask.java
View file @
76ebc4a4
...
@@ -49,7 +49,7 @@ public class RecordingTask implements Serializable {
...
@@ -49,7 +49,7 @@ public class RecordingTask implements Serializable {
private
String
taskId
;
private
String
taskId
;
/**
/**
* 业务类型:co
-
session/meeting
* 业务类型:co
_
session/meeting
*/
*/
@TableField
(
"biz_type"
)
@TableField
(
"biz_type"
)
private
String
bizType
;
private
String
bizType
;
...
@@ -85,10 +85,10 @@ public class RecordingTask implements Serializable {
...
@@ -85,10 +85,10 @@ public class RecordingTask implements Serializable {
private
String
layout
;
private
String
layout
;
/**
/**
*
0-初始化,1-录制中,2-已停止,3-失败,4
-已归档
*
1-初始化,2-录制中,3-已停止,4-失败,5
-已归档
*/
*/
@TableField
(
"status"
)
@TableField
(
"status"
)
private
Integer
status
;
private
String
status
;
/**
/**
* 录制开始时间
* 录制开始时间
...
...
yd-communication-service/src/main/java/com/yd/communication/service/service/ICoOperationLogService.java
View file @
76ebc4a4
...
@@ -13,4 +13,7 @@ import com.baomidou.mybatisplus.extension.service.IService;
...
@@ -13,4 +13,7 @@ import com.baomidou.mybatisplus.extension.service.IService;
*/
*/
public
interface
ICoOperationLogService
extends
IService
<
CoOperationLog
>
{
public
interface
ICoOperationLogService
extends
IService
<
CoOperationLog
>
{
void
log
(
String
bizId
,
String
operatorId
,
String
operatorType
,
String
operatorName
,
String
action
,
String
content
,
String
deviceNumber
,
String
ip
);
}
}
yd-communication-service/src/main/java/com/yd/communication/service/service/ICoSessionService.java
View file @
76ebc4a4
...
@@ -2,6 +2,7 @@ package com.yd.communication.service.service;
...
@@ -2,6 +2,7 @@ package com.yd.communication.service.service;
import
com.yd.communication.service.model.CoSession
;
import
com.yd.communication.service.model.CoSession
;
import
com.baomidou.mybatisplus.extension.service.IService
;
import
com.baomidou.mybatisplus.extension.service.IService
;
import
org.springframework.transaction.annotation.Transactional
;
/**
/**
* <p>
* <p>
...
@@ -13,4 +14,15 @@ import com.baomidou.mybatisplus.extension.service.IService;
...
@@ -13,4 +14,15 @@ import com.baomidou.mybatisplus.extension.service.IService;
*/
*/
public
interface
ICoSessionService
extends
IService
<
CoSession
>
{
public
interface
ICoSessionService
extends
IService
<
CoSession
>
{
/**
* 根据房间ID获取会话(WebSocket 使用)
*/
CoSession
getByRoomId
(
String
roomId
);
/**
* 根据业务ID获取会话
*/
CoSession
getByBizId
(
String
bizId
);
int
updateCurrentPageAndHistory
(
String
roomId
,
String
currentPage
,
String
pageHistory
,
String
updaterId
);
}
}
yd-communication-service/src/main/java/com/yd/communication/service/service/IRecordingTaskService.java
View file @
76ebc4a4
...
@@ -13,4 +13,7 @@ import com.baomidou.mybatisplus.extension.service.IService;
...
@@ -13,4 +13,7 @@ import com.baomidou.mybatisplus.extension.service.IService;
*/
*/
public
interface
IRecordingTaskService
extends
IService
<
RecordingTask
>
{
public
interface
IRecordingTaskService
extends
IService
<
RecordingTask
>
{
String
startRecording
(
String
bizId
,
String
roomId
);
void
stopRecording
(
String
taskId
);
}
}
yd-communication-service/src/main/java/com/yd/communication/service/service/impl/CoDesensitizationRuleServiceImpl.java
View file @
76ebc4a4
package
com
.
yd
.
communication
.
service
.
service
.
impl
;
package
com
.
yd
.
communication
.
service
.
service
.
impl
;
import
com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper
;
import
com.fasterxml.jackson.databind.JsonNode
;
import
com.fasterxml.jackson.databind.ObjectMapper
;
import
com.yd.communication.service.model.CoDesensitizationRule
;
import
com.yd.communication.service.model.CoDesensitizationRule
;
import
com.yd.communication.service.dao.CoDesensitizationRuleMapper
;
import
com.yd.communication.service.dao.CoDesensitizationRuleMapper
;
import
com.yd.communication.service.service.ICoDesensitizationRuleService
;
import
com.yd.communication.service.service.ICoDesensitizationRuleService
;
import
com.baomidou.mybatisplus.extension.service.impl.ServiceImpl
;
import
com.baomidou.mybatisplus.extension.service.impl.ServiceImpl
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.beans.factory.annotation.Autowired
;
import
org.springframework.stereotype.Service
;
import
org.springframework.stereotype.Service
;
import
javax.annotation.Resource
;
import
java.util.List
;
/**
/**
* <p>
* <p>
* 协同-脱敏设置表(通用) 服务实现类
* 协同-脱敏设置表(通用) 服务实现类
...
@@ -14,7 +22,83 @@ import org.springframework.stereotype.Service;
...
@@ -14,7 +22,83 @@ import org.springframework.stereotype.Service;
* @author zxm
* @author zxm
* @since 2026-07-28
* @since 2026-07-28
*/
*/
@Slf4j
@Service
@Service
public
class
CoDesensitizationRuleServiceImpl
extends
ServiceImpl
<
CoDesensitizationRuleMapper
,
CoDesensitizationRule
>
implements
ICoDesensitizationRuleService
{
public
class
CoDesensitizationRuleServiceImpl
extends
ServiceImpl
<
CoDesensitizationRuleMapper
,
CoDesensitizationRule
>
implements
ICoDesensitizationRuleService
{
@Resource
private
CoDesensitizationRuleMapper
ruleMapper
;
private
final
ObjectMapper
objectMapper
=
new
ObjectMapper
();
public
String
applyDesensitization
(
String
rawData
,
String
resourceType
,
String
resourceId
)
{
// 获取规则(先特定资源,再默认)
LambdaQueryWrapper
<
CoDesensitizationRule
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoDesensitizationRule:
:
getResourceType
,
resourceType
)
.
eq
(
CoDesensitizationRule:
:
getEnabled
,
1
)
.
orderByAsc
(
CoDesensitizationRule:
:
getSortOrder
);
List
<
CoDesensitizationRule
>
rules
;
// 先查特定资源
wrapper
.
eq
(
CoDesensitizationRule:
:
getResourceId
,
resourceId
);
rules
=
ruleMapper
.
selectList
(
wrapper
);
if
(
rules
.
isEmpty
())
{
// 查默认规则
wrapper
.
clear
();
wrapper
.
eq
(
CoDesensitizationRule:
:
getResourceType
,
resourceType
)
.
eq
(
CoDesensitizationRule:
:
getIsDefault
,
1
)
.
eq
(
CoDesensitizationRule:
:
getEnabled
,
1
)
.
orderByAsc
(
CoDesensitizationRule:
:
getSortOrder
);
rules
=
ruleMapper
.
selectList
(
wrapper
);
}
if
(
rules
.
isEmpty
())
{
return
rawData
;
}
try
{
JsonNode
root
=
objectMapper
.
readTree
(
rawData
);
for
(
CoDesensitizationRule
rule
:
rules
)
{
String
fieldPath
=
rule
.
getFieldPath
();
// 简单处理:假设路径是$.xxx,直接取字段
String
fieldName
=
fieldPath
.
replace
(
"$."
,
""
);
if
(
root
.
has
(
fieldName
))
{
String
original
=
root
.
get
(
fieldName
).
asText
();
String
masked
=
maskValue
(
original
,
rule
.
getMaskType
(),
rule
.
getMaskConfig
());
// 由于JsonNode不可变,简单返回,生产需构建新节点
}
}
return
objectMapper
.
writeValueAsString
(
root
);
}
catch
(
Exception
e
)
{
log
.
error
(
"脱敏失败"
,
e
);
return
rawData
;
}
}
private
String
maskValue
(
String
original
,
String
maskType
,
String
maskConfig
)
{
try
{
JsonNode
config
=
objectMapper
.
readTree
(
maskConfig
!=
null
?
maskConfig
:
"{}"
);
int
prefix
=
config
.
has
(
"prefix"
)
?
config
.
get
(
"prefix"
).
asInt
()
:
0
;
int
suffix
=
config
.
has
(
"suffix"
)
?
config
.
get
(
"suffix"
).
asInt
()
:
0
;
String
replaceChar
=
config
.
has
(
"replace_char"
)
?
config
.
get
(
"replace_char"
).
asText
()
:
"*"
;
switch
(
maskType
)
{
case
"mask"
:
case
"partial"
:
if
(
original
.
length
()
<=
prefix
+
suffix
)
return
original
;
int
middle
=
original
.
length
()
-
prefix
-
suffix
;
return
original
.
substring
(
0
,
prefix
)
+
// replaceChar.repeat(middle) +
original
.
substring
(
original
.
length
()
-
suffix
);
case
"hide"
:
return
"****"
;
case
"replace"
:
return
config
.
has
(
"replace_char"
)
?
config
.
get
(
"replace_char"
).
asText
()
:
"***"
;
default
:
return
original
;
}
}
catch
(
Exception
e
)
{
return
original
;
}
}
}
}
yd-communication-service/src/main/java/com/yd/communication/service/service/impl/CoOperationLogServiceImpl.java
View file @
76ebc4a4
...
@@ -4,6 +4,7 @@ import com.yd.communication.service.model.CoOperationLog;
...
@@ -4,6 +4,7 @@ import com.yd.communication.service.model.CoOperationLog;
import
com.yd.communication.service.dao.CoOperationLogMapper
;
import
com.yd.communication.service.dao.CoOperationLogMapper
;
import
com.yd.communication.service.service.ICoOperationLogService
;
import
com.yd.communication.service.service.ICoOperationLogService
;
import
com.baomidou.mybatisplus.extension.service.impl.ServiceImpl
;
import
com.baomidou.mybatisplus.extension.service.impl.ServiceImpl
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.stereotype.Service
;
import
org.springframework.stereotype.Service
;
/**
/**
...
@@ -14,7 +15,28 @@ import org.springframework.stereotype.Service;
...
@@ -14,7 +15,28 @@ import org.springframework.stereotype.Service;
* @author zxm
* @author zxm
* @since 2026-07-28
* @since 2026-07-28
*/
*/
@Slf4j
@Service
@Service
public
class
CoOperationLogServiceImpl
extends
ServiceImpl
<
CoOperationLogMapper
,
CoOperationLog
>
implements
ICoOperationLogService
{
public
class
CoOperationLogServiceImpl
extends
ServiceImpl
<
CoOperationLogMapper
,
CoOperationLog
>
implements
ICoOperationLogService
{
@Override
public
void
log
(
String
bizId
,
String
operatorId
,
String
operatorType
,
String
operatorName
,
String
action
,
String
content
,
String
deviceNumber
,
String
ip
)
{
CoOperationLog
log
=
new
CoOperationLog
();
log
.
setBizType
(
"co_session"
);
log
.
setBizId
(
bizId
);
log
.
setOperatorId
(
operatorId
);
log
.
setOperatorType
(
operatorType
);
log
.
setOperatorName
(
operatorName
);
log
.
setAction
(
action
);
log
.
setActionCategory
(
"control"
);
log
.
setContent
(
content
);
log
.
setOperatorDeviceNumber
(
deviceNumber
);
log
.
setOperatorIp
(
ip
);
log
.
setCreatorId
(
operatorId
);
log
.
setUpdaterId
(
operatorId
);
this
.
save
(
log
);
}
}
}
yd-communication-service/src/main/java/com/yd/communication/service/service/impl/CoSessionServiceImpl.java
View file @
76ebc4a4
...
@@ -2,9 +2,11 @@ package com.yd.communication.service.service.impl;
...
@@ -2,9 +2,11 @@ package com.yd.communication.service.service.impl;
import
com.yd.communication.service.model.CoSession
;
import
com.yd.communication.service.model.CoSession
;
import
com.yd.communication.service.dao.CoSessionMapper
;
import
com.yd.communication.service.dao.CoSessionMapper
;
import
com.yd.communication.service.service.ICoSessionService
;
import
com.baomidou.mybatisplus.extension.service.impl.ServiceImpl
;
import
com.baomidou.mybatisplus.extension.service.impl.ServiceImpl
;
import
com.yd.communication.service.service.ICoSessionService
;
import
org.springframework.stereotype.Service
;
import
org.springframework.stereotype.Service
;
import
com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper
;
import
lombok.extern.slf4j.Slf4j
;
/**
/**
* <p>
* <p>
...
@@ -15,6 +17,34 @@ import org.springframework.stereotype.Service;
...
@@ -15,6 +17,34 @@ import org.springframework.stereotype.Service;
* @since 2026-07-28
* @since 2026-07-28
*/
*/
@Service
@Service
@Slf4j
public
class
CoSessionServiceImpl
extends
ServiceImpl
<
CoSessionMapper
,
CoSession
>
implements
ICoSessionService
{
public
class
CoSessionServiceImpl
extends
ServiceImpl
<
CoSessionMapper
,
CoSession
>
implements
ICoSessionService
{
/**
* 根据房间ID获取会话(WebSocket 使用)
*/
@Override
public
CoSession
getByRoomId
(
String
roomId
)
{
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomId
,
roomId
)
.
eq
(
CoSession:
:
getIsDeleted
,
0
);
return
this
.
getOne
(
wrapper
);
}
/**
* 根据业务ID获取会话
*/
@Override
public
CoSession
getByBizId
(
String
bizId
)
{
return
this
.
lambdaQuery
()
.
eq
(
CoSession:
:
getCoSessionBizId
,
bizId
)
.
last
(
" limit 1 "
)
.
one
();
}
@Override
public
int
updateCurrentPageAndHistory
(
String
roomId
,
String
currentPage
,
String
pageHistory
,
String
updaterId
){
return
baseMapper
.
updateCurrentPageAndHistory
(
roomId
,
currentPage
,
pageHistory
,
updaterId
);
}
}
}
yd-communication-service/src/main/java/com/yd/communication/service/service/impl/RecordingTaskServiceImpl.java
View file @
76ebc4a4
package
com
.
yd
.
communication
.
service
.
service
.
impl
;
package
com
.
yd
.
communication
.
service
.
service
.
impl
;
import
com.baomidou.mybatisplus.core.conditions.query.LambdaQueryWrapper
;
import
com.yd.common.enums.CommonEnum
;
import
com.yd.common.utils.RandomStringGenerator
;
import
com.yd.communication.service.model.RecordingTask
;
import
com.yd.communication.service.model.RecordingTask
;
import
com.yd.communication.service.dao.RecordingTaskMapper
;
import
com.yd.communication.service.dao.RecordingTaskMapper
;
import
com.yd.communication.service.service.IRecordingTaskService
;
import
com.yd.communication.service.service.IRecordingTaskService
;
import
com.baomidou.mybatisplus.extension.service.impl.ServiceImpl
;
import
com.baomidou.mybatisplus.extension.service.impl.ServiceImpl
;
import
lombok.extern.slf4j.Slf4j
;
import
org.springframework.stereotype.Service
;
import
org.springframework.stereotype.Service
;
import
java.time.LocalDateTime
;
import
java.util.UUID
;
/**
/**
* <p>
* <p>
...
@@ -14,7 +20,72 @@ import org.springframework.stereotype.Service;
...
@@ -14,7 +20,72 @@ import org.springframework.stereotype.Service;
* @author zxm
* @author zxm
* @since 2026-07-28
* @since 2026-07-28
*/
*/
@Slf4j
@Service
@Service
public
class
RecordingTaskServiceImpl
extends
ServiceImpl
<
RecordingTaskMapper
,
RecordingTask
>
implements
IRecordingTaskService
{
public
class
RecordingTaskServiceImpl
extends
ServiceImpl
<
RecordingTaskMapper
,
RecordingTask
>
implements
IRecordingTaskService
{
/**
* 初始化录制信息(协同生成共享码的时候就初始化信息)
* @param bizId
* @param roomId
* @return
*/
@Override
public
String
startRecording
(
String
bizId
,
String
roomId
)
{
//任务ID
String
taskId
=
"agora_"
+
System
.
currentTimeMillis
();
RecordingTask
task
=
new
RecordingTask
();
//录制任务表唯一业务ID
task
.
setRecordingTaskBizId
(
RandomStringGenerator
.
generateBizId16
(
CommonEnum
.
UID_TYPE_RECORDING_TASK
.
getCode
()));
//任务编号
task
.
setTaskNo
(
"R"
+
System
.
currentTimeMillis
());
task
.
setTaskId
(
taskId
);
//关联的任务类型: 协同会话
task
.
setBizType
(
"co_session"
);
//关联的任务类型表的ID: 协同会话表唯一业务ID
task
.
setBizId
(
bizId
);
//房间号
task
.
setRoomId
(
roomId
);
//RTC频道名
task
.
setChannel
(
"channel_"
+
roomId
);
//录制模式
task
.
setRecordingMode
(
"mix"
);
//布局
task
.
setLayout
(
"grid"
);
//1-初始化
task
.
setStatus
(
"1"
);
task
.
setStartTime
(
LocalDateTime
.
now
());
task
.
setCreatorId
(
"system"
);
this
.
save
(
task
);
log
.
info
(
"启动录制成功, taskId={}, roomId={}"
,
taskId
,
roomId
);
return
taskId
;
}
@Override
public
void
stopRecording
(
String
taskId
)
{
LambdaQueryWrapper
<
RecordingTask
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
RecordingTask:
:
getTaskId
,
taskId
);
RecordingTask
task
=
this
.
getOne
(
wrapper
);
if
(
task
==
null
)
{
throw
new
RuntimeException
(
"录制任务不存在"
);
}
// 模拟停止录制
task
.
setStatus
(
"3"
);
task
.
setStopTime
(
LocalDateTime
.
now
());
task
.
setFileUrl
(
"https://oss.example.com/recordings/"
+
taskId
+
".mp4"
);
task
.
setFileDuration
(
120
);
task
.
setFileSize
(
1024000L
);
task
.
setFileMd5
(
UUID
.
randomUUID
().
toString
().
substring
(
0
,
32
));
task
.
setFileFormat
(
"mp4"
);
task
.
setStorageType
(
"oss"
);
task
.
setStorageBucket
(
"coordination-recordings"
);
task
.
setStoragePath
(
"/recordings/"
+
taskId
+
".mp4"
);
task
.
setUpdaterId
(
"system"
);
this
.
updateById
(
task
);
log
.
info
(
"停止录制成功, taskId={}"
,
taskId
);
}
}
}
yd-communication-service/src/main/java/com/yd/communication/service/utils/RandomUtil.java
0 → 100644
View file @
76ebc4a4
package
com
.
yd
.
communication
.
service
.
utils
;
import
java.security.SecureRandom
;
public
class
RandomUtil
{
private
static
final
SecureRandom
random
=
new
SecureRandom
();
public
static
String
generateNumericCode
(
int
length
)
{
StringBuilder
sb
=
new
StringBuilder
();
for
(
int
i
=
0
;
i
<
length
;
i
++)
{
sb
.
append
(
random
.
nextInt
(
10
));
}
return
sb
.
toString
();
}
}
\ No newline at end of file
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment