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
3bec6694
Commit
3bec6694
authored
Jul 30, 2026
by
zhangxingmin
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
push
parent
2442dfc8
Show whitespace changes
Inline
Side-by-side
Showing
2 changed files
with
220 additions
and
93 deletions
+220
-93
yd-communication-api/src/main/java/com/yd/communication/api/service/impl/ApiCoSessionServiceImpl.java
+177
-84
yd-communication-api/src/main/java/com/yd/communication/api/websocket/CoWebSocketServer.java
+43
-9
No files found.
yd-communication-api/src/main/java/com/yd/communication/api/service/impl/ApiCoSessionServiceImpl.java
View file @
3bec6694
...
...
@@ -28,11 +28,11 @@ 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
{
...
...
@@ -48,12 +48,12 @@ public class ApiCoSessionServiceImpl implements ApiCoSessionService {
/**
* 客户创建会话(生成共享码)
* @param request
* @return
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
CreateResponse
>
create
(
CreateRequest
request
)
{
log
.
info
(
"【创建会话】收到创建请求, request={}"
,
request
);
try
{
CoSession
session
=
createSession
(
request
.
getScope
(),
request
.
getResourceType
(),
...
...
@@ -65,6 +65,7 @@ public class ApiCoSessionServiceImpl implements ApiCoSessionService {
request
.
getToken
()
);
if
(
session
==
null
)
{
log
.
warn
(
"【创建会话】创建失败,返回空会话"
);
return
Result
.
success
();
}
CreateResponse
response
=
new
CreateResponse
();
...
...
@@ -72,24 +73,25 @@ public class ApiCoSessionServiceImpl implements ApiCoSessionService {
response
.
setRoomPwd
(
session
.
getRoomPwd
());
response
.
setSessionBizId
(
session
.
getCoSessionBizId
());
response
.
setStatus
(
session
.
getStatus
());
log
.
info
(
"【创建会话】成功,roomId={}, roomPwd={}, sessionBizId={}"
,
session
.
getRoomId
(),
session
.
getRoomPwd
(),
session
.
getCoSessionBizId
());
return
Result
.
success
(
response
);
}
catch
(
Exception
e
)
{
log
.
error
(
"【创建会话】异常"
,
e
);
throw
e
;
}
}
/**
* 客户创建会话(生成共享码)
* @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
)
{
//创建会话
String
userId
,
String
token
)
{
log
.
info
(
"【创建会话-内部】开始创建, scope={}, resourceType={}, resourceId={}, ownerId={}, ownerType={}, userId={}"
,
scope
,
resourceType
,
resourceId
,
ownerId
,
ownerType
,
userId
);
try
{
// 创建会话
CoSession
session
=
new
CoSession
();
session
.
setCoSessionBizId
(
RandomStringGenerator
.
generateBizId16
(
CommonEnum
.
UID_TYPE_CO_SESSION
.
getCode
()));
session
.
setCoSessionNo
(
"S"
+
System
.
currentTimeMillis
());
...
...
@@ -100,62 +102,66 @@ public class ApiCoSessionServiceImpl implements ApiCoSessionService {
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
);
log
.
info
(
"【创建会话-内部】数据库保存成功, id={}"
,
session
.
getId
());
// 存入 Redis
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
());
redisUtil
.
setCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
,
roomRedisInfoDTO
,
RedisEnum
.
ROOM
.
getTimeout
(),
RedisEnum
.
ROOM
.
getTimeUnit
());
log
.
info
(
"【创建会话-内部】Redis缓存已设置, key={}"
,
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
);
// 自动开启录制
// 自动开启录制(当前注释)
if
(
autoStartRecording
())
{
// recordingService.startRecording(session.getCoSessionBizId(), roomId);
// recordingService.startRecording(session.getCoSessionBizId(), roomId);
log
.
info
(
"【创建会话-内部】自动录制已触发(未实际启动)"
);
}
log
.
info
(
"创建会话成功, roomId={}, roomPwd={}"
,
roomId
,
roomPwd
);
//添加操作日志,协同-操作日志表 TODO
log
.
info
(
"【创建会话-内部】创建成功, roomId={}, roomPwd={}"
,
roomId
,
roomPwd
);
return
session
;
}
catch
(
Exception
e
)
{
log
.
error
(
"【创建会话-内部】异常"
,
e
);
throw
e
;
}
}
/**
* 顾问加入会话((输入共享码加入房间,可以多次加入))
* @param request
* @return
* 顾问加入会话(输入共享码加入房间,可以多次加入)
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
JoinResponse
>
join
(
JoinRequest
request
)
{
log
.
info
(
"【加入会话】收到加入请求, request={}"
,
request
);
try
{
CoSession
session
=
joinSession
(
request
.
getRoomPwd
(),
request
.
getParticipantId
(),
request
.
getParticipantType
()
);
if
(
session
==
null
)
{
log
.
warn
(
"【加入会话】加入失败,返回空会话"
);
return
Result
.
success
();
}
JoinResponse
joinResponse
=
new
JoinResponse
();
...
...
@@ -165,187 +171,252 @@ public class ApiCoSessionServiceImpl implements ApiCoSessionService {
joinResponse
.
setRoomId
(
session
.
getRoomId
());
joinResponse
.
setSessionBizId
(
session
.
getResourceInit
());
//
获取资源所有者缓存中的登录信息
//
获取资源所有者缓存中的登录信息
RoomRedisInfoDTO
roomRedisInfoDTO
=
redisUtil
.
getCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
session
.
getRoomId
());
if
(
roomRedisInfoDTO
==
null
)
{
log
.
error
(
"【加入会话】会话发起者缓存信息不存在,roomId={}"
,
session
.
getRoomId
());
throw
new
RuntimeException
(
"会话发起者登录信息失效,建议联系会话发起者再次发起"
);
}
joinResponse
.
setUserId
(
roomRedisInfoDTO
.
getUserId
());
joinResponse
.
setToken
(
roomRedisInfoDTO
.
getToken
());
log
.
info
(
"【加入会话】成功, roomId={}, participantId={}, controlHolderType={}"
,
session
.
getRoomId
(),
session
.
getParticipantId
(),
session
.
getControlHolderType
());
return
Result
.
success
(
joinResponse
);
}
catch
(
Exception
e
)
{
log
.
error
(
"【加入会话】异常"
,
e
);
throw
e
;
}
}
/**
* 顾问加入会话
* @param roomPwd
* @param participantId
* @param participantType
* @return
* 顾问加入会话(内部方法)
*/
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
CoSession
joinSession
(
String
roomPwd
,
String
participantId
,
String
participantType
)
{
log
.
info
(
"【加入会话-内部】开始, roomPwd={}, participantId={}, participantType={}"
,
roomPwd
,
participantId
,
participantType
);
try
{
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomPwd
,
roomPwd
)
.
last
(
" limit 1 "
);
CoSession
session
=
iCoSessionService
.
getOne
(
wrapper
);
if
(
session
==
null
)
{
log
.
error
(
"【加入会话-内部】未找到会话,roomPwd={}"
,
roomPwd
);
throw
new
RuntimeException
(
"房间密码(共享码)错误,或客户未发起会话"
);
}
log
.
info
(
"【加入会话-内部】找到会话, roomId={}, status={}, currentHolder={}:{}"
,
session
.
getRoomId
(),
session
.
getStatus
(),
session
.
getControlHolderType
(),
session
.
getControlHolderId
());
if
(
CoSessionStatusEnum
.
YJS
.
getItemValue
().
equals
(
session
.
getStatus
()))
{
//已结束,不能再次加入房间
log
.
warn
(
"【加入会话-内部】会话已结束,roomId={}"
,
session
.
getRoomId
());
throw
new
RuntimeException
(
"会话已结束,不能再次加入房间"
);
}
if
(
StringUtils
.
isNotBlank
(
session
.
getParticipantId
())
&&
!
session
.
getParticipantId
().
equals
(
participantId
))
{
//库里参与者ID不为空并且库中参与者ID和加入房间的参与者ID不相等,不能加入房间
log
.
warn
(
"【加入会话-内部】房间被占用,当前参与者={}, 新参与者={}"
,
session
.
getParticipantId
(),
participantId
);
throw
new
RuntimeException
(
"当前房间被占用,不能加入到房间"
);
}
// 若为待开始状态,设置开始时间并转移控制权
if
(
CoSessionStatusEnum
.
DKS
.
getItemValue
().
equals
(
session
.
getStatus
()))
{
if
(
session
.
getStartTime
()
==
null
)
{
//开始时间为空时,说明当前处于待开始到进行中过渡阶段,设置开始时间
session
.
setStartTime
(
LocalDateTime
.
now
());
log
.
info
(
"【加入会话-内部】设置开始时间={}"
,
session
.
getStartTime
());
}
//待开始状态下,控制权要自动移交给参与者(顾问)
// 待开始状态下,控制权自动移交给参与者(顾问)
session
.
setControlHolderType
(
ControlHolderTypeEnum
.
PARTICIPANT
.
getItemValue
());
session
.
setControlHolderId
(
participantId
);
log
.
info
(
"【加入会话-内部】待开始状态,控制权移交给参与者={}"
,
participantId
);
}
else
{
// 如果会话已经是进行中,但可能控制权不在顾问,此处强制转移给顾问(根据业务需求,可调整)
// 如果不希望强制转移,可注释掉以下行
session
.
setControlHolderType
(
ControlHolderTypeEnum
.
PARTICIPANT
.
getItemValue
());
session
.
setControlHolderId
(
participantId
);
log
.
info
(
"【加入会话-内部】强制控制权移交给参与者={}"
,
participantId
);
}
//更新会话状态:
进行中
// 更新会话状态为
进行中
session
.
setStatus
(
CoSessionStatusEnum
.
JXZ
.
getItemValue
());
session
.
setParticipantId
(
participantId
);
session
.
setParticipantType
(
participantType
);
session
.
setUpdaterId
(
participantId
);
iCoSessionService
.
updateById
(
session
);
log
.
info
(
"【加入会话-内部】数据库更新成功,新状态={}"
,
session
.
getStatus
());
//更新 Redis 缓存,更新控制权字段信息
// 更新 Redis 缓存
RoomRedisInfoDTO
roomRedisInfoDTO
=
redisUtil
.
getCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
session
.
getRoomId
());
if
(
roomRedisInfoDTO
!=
null
)
{
roomRedisInfoDTO
.
setControlHolderType
(
session
.
getControlHolderType
());
roomRedisInfoDTO
.
setControlHolderId
(
session
.
getControlHolderId
());
redisUtil
.
setCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
session
.
getRoomId
(),
roomRedisInfoDTO
,
RedisEnum
.
ROOM
.
getTimeout
(),
RedisEnum
.
ROOM
.
getTimeUnit
());
log
.
info
(
"【加入会话-内部】Redis缓存更新成功,新控制者={}:{}"
,
session
.
getControlHolderType
(),
session
.
getControlHolderId
());
}
else
{
log
.
warn
(
"【加入会话-内部】Redis缓存不存在,可能已过期,roomId={}"
,
session
.
getRoomId
());
}
//添加操作日志,协同-操作日志表 TODO
return
session
;
}
catch
(
Exception
e
)
{
log
.
error
(
"【加入会话-内部】异常"
,
e
);
throw
e
;
}
}
/**
* 获取会话详情
* @param bizId
* @return
*/
@Override
public
Result
<
SessionDetailResponse
>
get
(
String
bizId
)
{
log
.
info
(
"【获取会话详情】bizId={}"
,
bizId
);
try
{
CoSession
session
=
iCoSessionService
.
getByBizId
(
bizId
);
if
(
session
==
null
)
{
log
.
warn
(
"【获取会话详情】未找到会话,bizId={}"
,
bizId
);
return
Result
.
success
();
}
SessionDetailResponse
response
=
new
SessionDetailResponse
();
BeanUtils
.
copyProperties
(
session
,
response
);
BeanUtils
.
copyProperties
(
session
,
response
);
log
.
info
(
"【获取会话详情】成功,roomId={}"
,
session
.
getRoomId
());
return
Result
.
success
(
response
);
}
catch
(
Exception
e
)
{
log
.
error
(
"【获取会话详情】异常"
,
e
);
throw
e
;
}
}
/**
* 结束会话(客户调用)
* @param request
* @return
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
CommonResponse
>
end
(
EndSessionRequest
request
)
{
log
.
info
(
"【结束会话】收到请求, roomId={}"
,
request
.
getRoomId
());
try
{
CoSession
coSession
=
iCoSessionService
.
lambdaQuery
()
.
eq
(
CoSession:
:
getRoomId
,
request
.
getRoomId
())
.
eq
(
CoSession:
:
getRoomId
,
request
.
getRoomId
())
.
last
(
" limit 1"
)
.
one
();
if
(
coSession
==
null
)
{
log
.
error
(
"【结束会话】会话不存在,roomId={}"
,
request
.
getRoomId
());
throw
new
BusinessException
(
"会话不存在"
);
}
//结束会话关闭共享,更新信息
log
.
info
(
"【结束会话】找到会话,当前状态={}"
,
coSession
.
getStatus
());
// 结束会话关闭共享,更新信息
coSession
.
setStatus
(
CoSessionStatusEnum
.
YJS
.
getItemValue
());
//结束时间
coSession
.
setEndTime
(
LocalDateTime
.
now
());
iCoSessionService
.
saveOrUpdate
(
coSession
);
log
.
info
(
"【结束会话】数据库更新成功,状态已结束"
);
//
销毁redis房间缓存信息
//
销毁redis房间缓存信息
redisUtil
.
deleteObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
coSession
.
getRoomId
());
log
.
info
(
"【结束会话】Redis缓存已删除"
);
//
结束录制视频,并且更新录制任务表信息存档 TODO
//
结束录制视频,并且更新录制任务表信息存档 TODO
//
添加操作日志,协同-操作日志表 TODO
//
添加操作日志,协同-操作日志表 TODO
CommonResponse
response
=
new
CommonResponse
();
response
.
setMessage
(
"会话已结束"
);
log
.
info
(
"【结束会话】成功, roomId={}"
,
request
.
getRoomId
());
return
Result
.
success
(
response
);
}
catch
(
Exception
e
)
{
log
.
error
(
"【结束会话】异常"
,
e
);
throw
e
;
}
}
/**
* 切换控制权(顾问调用)
* @param request
* @return
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
Result
<
CommonResponse
>
transferControl
(
TransferControlRequest
request
)
{
transferControlUp
(
request
.
getOprType
(),
request
.
getRoomId
());
log
.
info
(
"【切换控制权】收到请求, oprType={}, roomId={}"
,
request
.
getOprType
(),
request
.
getRoomId
());
try
{
transferControlUp
(
request
.
getOprType
(),
request
.
getRoomId
());
CommonResponse
response
=
new
CommonResponse
();
response
.
setMessage
(
"控制权已切换"
);
log
.
info
(
"【切换控制权】成功, roomId={}"
,
request
.
getRoomId
());
return
Result
.
success
(
response
);
}
catch
(
Exception
e
)
{
log
.
error
(
"【切换控制权】异常"
,
e
);
throw
e
;
}
}
/**
* 切换控制权
* @param oprType
* @param roomId
* 切换控制权(内部方法)
*/
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
void
transferControlUp
(
Integer
oprType
,
String
roomId
)
{
//根据房间号查询会话信息
public
void
transferControlUp
(
Integer
oprType
,
String
roomId
)
{
log
.
info
(
"【切换控制权-内部】开始, oprType={}, roomId={}"
,
oprType
,
roomId
);
try
{
// 根据房间号查询会话信息
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomId
,
roomId
).
last
(
" limit 1 "
);
CoSession
session
=
iCoSessionService
.
getOne
(
wrapper
);
if
(
session
==
null
)
{
log
.
error
(
"【切换控制权-内部】会话不存在, roomId={}"
,
roomId
);
throw
new
RuntimeException
(
"会话不存在"
);
}
log
.
info
(
"【切换控制权-内部】找到会话, 当前控制者={}:{}"
,
session
.
getControlHolderType
(),
session
.
getControlHolderId
());
//移交控制权
//控制权持有者类型:owner(资源所有者类型)/participant(参与者类型)
// 移交控制权
if
(
oprType
==
1
)
{
//
1-开启客户操作,控制权移交给资源所有者(客户)
//
1-开启客户操作,控制权移交给资源所有者(客户)
session
.
setControlHolderType
(
ControlHolderTypeEnum
.
OWNER
.
getItemValue
());
session
.
setControlHolderId
(
session
.
getOwnerId
());
}
else
if
(
oprType
==
2
)
{
//2-关闭客户操作,控制权移交给参与者(顾问)
log
.
info
(
"【切换控制权-内部】开启客户操作,控制权移交给所有者={}"
,
session
.
getOwnerId
());
}
else
if
(
oprType
==
2
)
{
// 2-关闭客户操作,控制权移交给参与者(顾问)
session
.
setControlHolderType
(
ControlHolderTypeEnum
.
PARTICIPANT
.
getItemValue
());
session
.
setControlHolderId
(
session
.
getParticipantId
());
log
.
info
(
"【切换控制权-内部】关闭客户操作,控制权移交给参与者={}"
,
session
.
getParticipantId
());
}
else
{
log
.
warn
(
"【切换控制权-内部】未知oprType={}, 忽略"
,
oprType
);
return
;
}
session
.
setUpdaterId
(
session
.
getParticipantId
());
iCoSessionService
.
updateById
(
session
);
log
.
info
(
"【切换控制权-内部】数据库更新成功,新控制者={}:{}"
,
session
.
getControlHolderType
(),
session
.
getControlHolderId
());
//更新房间缓存redis信息-更新控制权字段,用于WebSocket获取
// 更新房间缓存redis信息
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
());
redisUtil
.
setCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
,
roomRedisInfoDTO
,
RedisEnum
.
ROOM
.
getTimeout
(),
RedisEnum
.
ROOM
.
getTimeUnit
());
log
.
info
(
"【切换控制权-内部】Redis缓存更新成功"
);
}
else
{
log
.
warn
(
"【切换控制权-内部】Redis缓存不存在,可能已过期"
);
}
// 添加操作日志,协同-操作日志表 TODO
}
catch
(
Exception
e
)
{
log
.
error
(
"【切换控制权-内部】异常"
,
e
);
throw
e
;
}
//添加操作日志,协同-操作日志表 TODO
}
/**
* 获取当前控制者
* @param roomId
* @return
*/
@Override
public
RoomRedisInfoDTO
getCurrentController
(
String
roomId
)
{
log
.
info
(
"【获取当前控制者】roomId={}"
,
roomId
);
try
{
RoomRedisInfoDTO
roomRedisInfoDTO
=
redisUtil
.
getCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
);
if
(
roomRedisInfoDTO
==
null
||
(
roomRedisInfoDTO
!=
null
&&
StringUtils
.
isBlank
(
roomRedisInfoDTO
.
getControlHolderId
())))
{
//缓存信息或者缓存内的控制者ID为空的时候,单独去查询库更新缓存信息
log
.
info
(
"【获取当前控制者】Redis缓存不存在或控制者ID为空,从数据库查询"
);
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomId
,
roomId
).
last
(
" limit 1 "
);
CoSession
session
=
iCoSessionService
.
getOne
(
wrapper
);
if
(
session
==
null
)
{
log
.
error
(
"【获取当前控制者】会话不存在, roomId={}"
,
roomId
);
throw
new
RuntimeException
(
"会话不存在"
);
}
if
(
roomRedisInfoDTO
==
null
)
{
...
...
@@ -353,28 +424,37 @@ public class ApiCoSessionServiceImpl implements ApiCoSessionService {
}
roomRedisInfoDTO
.
setControlHolderType
(
session
.
getControlHolderType
());
roomRedisInfoDTO
.
setControlHolderId
(
session
.
getControlHolderId
());
//查询出来的信息更新回缓存里面
redisUtil
.
setCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
,
roomRedisInfoDTO
,
RedisEnum
.
ROOM
.
getTimeout
(),
RedisEnum
.
ROOM
.
getTimeUnit
());
// 查询出来的信息更新回缓存里面
redisUtil
.
setCacheObject
(
RedisEnum
.
ROOM
.
getPrefix
()
+
roomId
,
roomRedisInfoDTO
,
RedisEnum
.
ROOM
.
getTimeout
(),
RedisEnum
.
ROOM
.
getTimeUnit
());
log
.
info
(
"【获取当前控制者】从数据库加载并更新缓存,控制者={}:{}"
,
session
.
getControlHolderType
(),
session
.
getControlHolderId
());
}
else
{
log
.
info
(
"【获取当前控制者】从Redis获取,控制者={}:{}"
,
roomRedisInfoDTO
.
getControlHolderType
(),
roomRedisInfoDTO
.
getControlHolderId
());
}
return
roomRedisInfoDTO
;
}
catch
(
Exception
e
)
{
log
.
error
(
"【获取当前控制者】异常"
,
e
);
throw
e
;
}
}
/**
* 更新当前页面
* @param roomId
* @param newPageJson
* @param operatorId
*/
@Override
@Transactional
(
rollbackFor
=
Exception
.
class
)
public
void
updateCurrentPage
(
String
roomId
,
String
newPageJson
,
String
operatorId
)
{
log
.
info
(
"【更新当前页面】roomId={}, operatorId={}, newPageJson={}"
,
roomId
,
operatorId
,
newPageJson
);
try
{
LambdaQueryWrapper
<
CoSession
>
wrapper
=
new
LambdaQueryWrapper
<>();
wrapper
.
eq
(
CoSession:
:
getRoomId
,
roomId
).
last
(
" limit 1 "
);
CoSession
session
=
iCoSessionService
.
getOne
(
wrapper
);
if
(
session
==
null
)
{
log
.
error
(
"【更新当前页面】会话不存在, roomId={}"
,
roomId
);
throw
new
RuntimeException
(
"会话不存在"
);
}
String
history
=
session
.
getPageHistory
();
if
(
history
==
null
||
history
.
equals
(
"[]"
)
||
history
.
isEmpty
())
{
history
=
"["
+
newPageJson
+
"]"
;
...
...
@@ -382,22 +462,36 @@ public class ApiCoSessionServiceImpl implements ApiCoSessionService {
// 简单追加(生产环境建议用JSONArray处理)
history
=
history
.
substring
(
0
,
history
.
length
()
-
1
)
+
","
+
newPageJson
+
"]"
;
}
iCoSessionService
.
updateCurrentPageAndHistory
(
roomId
,
newPageJson
,
history
,
operatorId
);
log
.
info
(
"【更新当前页面】成功, 新历史={}"
,
history
);
}
catch
(
Exception
e
)
{
log
.
error
(
"【更新当前页面】异常"
,
e
);
throw
e
;
}
}
/**
* 根据房间号获取会话信息
* @param roomId
* @return
*/
@Override
public
CoSession
getByRoomId
(
String
roomId
)
{
return
iCoSessionService
.
getByRoomId
(
roomId
);
log
.
info
(
"【根据房间号获取会话】roomId={}"
,
roomId
);
try
{
CoSession
session
=
iCoSessionService
.
getByRoomId
(
roomId
);
if
(
session
==
null
)
{
log
.
warn
(
"【根据房间号获取会话】未找到会话, roomId={}"
,
roomId
);
}
else
{
log
.
info
(
"【根据房间号获取会话】找到会话, controlHolder={}:{}"
,
session
.
getControlHolderType
(),
session
.
getControlHolderId
());
}
return
session
;
}
catch
(
Exception
e
)
{
log
.
error
(
"【根据房间号获取会话】异常"
,
e
);
throw
e
;
}
}
private
boolean
autoStartRecording
()
{
return
true
;
// 从配置读取
}
}
\ No newline at end of file
yd-communication-api/src/main/java/com/yd/communication/api/websocket/CoWebSocketServer.java
View file @
3bec6694
...
...
@@ -515,21 +515,23 @@ public class CoWebSocketServer {
public
void
onMessage
(
Session
session
,
String
message
)
{
// 根据会话查询对应的房间号
String
roomId
=
SESSION_ROOM
.
get
(
session
);
// 如果会话未关联房间,忽略消息
if
(
roomId
==
null
)
{
log
.
warn
(
"【WebSocket】会话未关联房间,忽略消息"
);
return
;
}
// 打印收到的原始消息
log
.
info
(
"【WebSocket】收到消息,房间号: {}, 消息内容: {}"
,
roomId
,
message
);
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
);
log
.
info
(
"【WebSocket】解析结果: action={}, userId={}, userType={}"
,
action
,
userId
,
userType
);
// 2. 获取会话业务ID(用于操作日志记录)
CoSession
coSession
=
sessionService
.
getByRoomId
(
roomId
);
String
bizId
=
coSession
!=
null
?
coSession
.
getCoSessionBizId
()
:
null
;
...
...
@@ -538,26 +540,41 @@ public class CoWebSocketServer {
// --- 操作1:控制权切换(仅顾问可操作) ---
if
(
"control_transfer"
.
equals
(
action
))
{
log
.
info
(
"【控制权切换】开始处理,当前用户类型: {}"
,
userType
);
// 校验权限:只有参与者(顾问)才能切换控制权
if
(!
"participant"
.
equals
(
userType
))
{
log
.
warn
(
"【控制权切换】权限不足,非顾问用户尝试切换,userType={}"
,
userType
);
sendError
(
session
,
"只有顾问可以切换控制权"
);
return
;
}
// 解析顾问切换后的控制权信息
String
newHolderType
=
json
.
get
(
"holderType"
).
asText
();
String
newHolderId
=
json
.
get
(
"holderId"
).
asText
();
log
.
info
(
"【控制权切换】目标控制者: holderType={}, holderId={}"
,
newHolderType
,
newHolderId
);
// 构造控制权切换请求
TransferControlRequest
req
=
new
TransferControlRequest
();
req
.
setRoomId
(
roomId
);
// 判断切换后的持有者类型:如果是 owner(客户)则 opType=1(开启客户操作),否则 opType=2(关闭客户操作)
req
.
setOprType
(
ControlHolderTypeEnum
.
OWNER
.
getItemValue
().
equals
(
newHolderType
)
?
1
:
2
);
Integer
oprType
=
ControlHolderTypeEnum
.
OWNER
.
getItemValue
().
equals
(
newHolderType
)
?
1
:
2
;
req
.
setOprType
(
oprType
);
log
.
info
(
"【控制权切换】操作类型: {}"
,
oprType
);
// 调用服务层切换控制权(更新数据库 + Redis)
try
{
sessionService
.
transferControl
(
req
);
log
.
info
(
"【控制权切换】服务层调用成功,数据库和Redis已更新"
);
}
catch
(
Exception
e
)
{
log
.
error
(
"【控制权切换】服务层调用失败"
,
e
);
sendError
(
session
,
"切换控制权失败:"
+
e
.
getMessage
());
return
;
}
// 广播控制权变更消息给房间内所有用户(包括自己)
broadcast
(
roomId
,
message
,
session
);
log
.
info
(
"【控制权切换】广播消息已发送"
);
// 记录操作日志(用于合规审计)
if
(
bizId
!=
null
)
{
...
...
@@ -570,8 +587,11 @@ public class CoWebSocketServer {
// --- 操作2:结束共享(仅客户可操作) ---
if
(
"end_sharing"
.
equals
(
action
))
{
log
.
info
(
"【结束共享】开始处理,当前用户类型: {}"
,
userType
);
// 校验权限:只有所有者(客户)才能结束共享
if
(!
"owner"
.
equals
(
userType
))
{
log
.
warn
(
"【结束共享】权限不足,非客户用户尝试结束共享,userType={}"
,
userType
);
sendError
(
session
,
"只有客户可以结束共享"
);
return
;
}
...
...
@@ -579,26 +599,29 @@ public class CoWebSocketServer {
// 构造结束会话请求
EndSessionRequest
req
=
new
EndSessionRequest
();
req
.
setRoomId
(
roomId
);
// 调用服务层结束会话(更新状态 + 停止录制 + 清理Redis)
log
.
info
(
"【结束共享】调用服务层结束会话"
);
sessionService
.
end
(
req
);
// 广播结束消息给所有用户
broadcast
(
roomId
,
"{\"action\":\"end_sharing\"}"
,
null
);
log
.
info
(
"【结束共享】结束消息已广播"
);
// 清理 Redis 中的页面状态
redisTemplate
.
delete
(
String
.
format
(
ROOM_STATE_PAGE_KEY
,
roomId
));
log
.
info
(
"【结束共享】Redis页面状态已清理"
);
// 关闭房间所有连接并清理资源
closeRoom
(
roomId
);
log
.
info
(
"共享已结束: roomId={}"
,
roomId
);
return
;
}
// --- 操作3:脱敏开关(客户操作,全房间生效) ---
if
(
"desensitization_switch"
.
equals
(
action
))
{
log
.
info
(
"【脱敏开关】收到请求,enabled={}"
,
json
.
get
(
"enabled"
).
asBoolean
());
// 广播脱敏状态给所有人(同步显示脱敏效果)
broadcast
(
roomId
,
message
,
session
);
log
.
info
(
"【脱敏开关】广播消息已发送"
);
// 记录操作日志
if
(
bizId
!=
null
)
{
...
...
@@ -609,13 +632,20 @@ public class CoWebSocketServer {
// --- 操作4:同步操作(翻页/滚动/缩放)——仅控制权持有者可操作 ---
if
(
"turn_the_page"
.
equals
(
action
)
||
"scroll"
.
equals
(
action
)
||
"scaling"
.
equals
(
action
))
{
log
.
info
(
"【同步操作】开始处理,action={}, 当前用户类型={}"
,
action
,
userType
);
// 4.1 获取当前控制权持有者信息
RoomRedisInfoDTO
dto
=
sessionService
.
getCurrentController
(
roomId
);
String
holderType
=
dto
.
getControlHolderType
();
String
holderId
=
dto
.
getControlHolderId
();
log
.
info
(
"【同步操作】当前控制权: holderType={}, holderId={}"
,
holderType
,
holderId
);
// 4.2 校验:当前用户是否持有控制权
if
(!(
holderType
.
equals
(
userType
)
&&
holderId
.
equals
(
userId
)))
{
boolean
isController
=
holderType
.
equals
(
userType
)
&&
holderId
.
equals
(
userId
);
log
.
info
(
"【同步操作】是否持有控制权: {}"
,
isController
);
if
(!
isController
)
{
log
.
warn
(
"【同步操作】用户无控制权,拒绝操作,userType={}, userId={}"
,
userType
,
userId
);
sendError
(
session
,
"您没有控制权,无法操作"
);
return
;
}
...
...
@@ -624,17 +654,21 @@ public class CoWebSocketServer {
if
(
"turn_the_page"
.
equals
(
action
)
||
"scroll"
.
equals
(
action
))
{
// 获取当前页面 JSON(从消息中提取)
String
currentPage
=
json
.
has
(
"currentPage"
)
?
json
.
get
(
"currentPage"
).
toString
()
:
"{}"
;
log
.
info
(
"【同步操作】更新页面状态: {}"
,
currentPage
);
// 更新数据库中的当前页面和页面历史轨迹
sessionService
.
updateCurrentPage
(
roomId
,
currentPage
,
userId
);
log
.
info
(
"【同步操作】数据库已更新"
);
// 更新 Redis 状态(供新节点连接时同步)
String
stateKey
=
String
.
format
(
ROOM_STATE_PAGE_KEY
,
roomId
);
redisTemplate
.
opsForValue
().
set
(
stateKey
,
currentPage
,
60
,
TimeUnit
.
MINUTES
);
log
.
info
(
"【同步操作】Redis状态已更新"
);
}
// 4.4 广播操作指令给房间内其他用户(排除自己,避免回环)
broadcast
(
roomId
,
message
,
session
);
log
.
info
(
"【同步操作】广播消息已发送"
);
// 4.5 记录操作日志(合规审计)
if
(
bizId
!=
null
)
{
...
...
@@ -648,7 +682,7 @@ public class CoWebSocketServer {
}
catch
(
Exception
e
)
{
// 处理消息异常,向客户端返回错误信息
log
.
error
(
"处理
消息异常"
,
e
);
log
.
error
(
"处理
WebSocket消息异常,消息内容: {}"
,
message
,
e
);
sendError
(
session
,
"服务器处理异常"
);
}
}
...
...
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