refactor: 重构sse消息推送,封装sseManger

This commit is contained in:
czhqwer
2025-03-30 14:03:42 +08:00
parent 1982a450a0
commit a38bfe0a69
9 changed files with 52 additions and 128 deletions
@@ -154,7 +154,7 @@ public class FileController {
} }
enableShare = enable; enableShare = enable;
Map<String, Object> data = MapUtil.of("enable", enable); Map<String, Object> data = MapUtil.of("enable", enable);
sseService.notifySystemEvent(new NotifyMessage("enableShare", data)); sseService.notification(new NotifyMessage("enableShare", data));
return Result.success(); return Result.success();
} }
@@ -18,16 +18,13 @@ public class SseController {
@Resource @Resource
private ISseService sseService; private ISseService sseService;
@GetMapping("/subscribeSharedFiles") /**
public SseEmitter subscribeSharedFiles(HttpServletRequest request) { * 通用订阅
*/
@GetMapping("/subscribe")
public SseEmitter subscribe(HttpServletRequest request) {
String remoteAddr = request.getRemoteAddr(); String remoteAddr = request.getRemoteAddr();
return sseService.subscribeSharedFiles(remoteAddr); return sseService.subscribe(remoteAddr);
}
@GetMapping("/subscribeSystemEvent")
public SseEmitter subscribeSystemEvent(HttpServletRequest request) {
String remoteAddr = request.getRemoteAddr();
return sseService.subscribeSystemEvent(remoteAddr);
} }
@@ -6,23 +6,13 @@ import org.springframework.web.servlet.mvc.method.annotation.SseEmitter;
public interface ISseService { public interface ISseService {
/** /**
* 订阅共享文件更新 * 通用订阅
*/ */
SseEmitter subscribeSharedFiles(String clientId); SseEmitter subscribe(String remoteAddr);
/** /**
* 订阅系统事件 * 通用通知
*/ */
SseEmitter subscribeSystemEvent(String clientId); void notification(NotifyMessage message);
/**
* 通知共享文件更新
*/
void notifySharedFileUpdate(String message);
/**
* 通知系统事件
*/
void notifySystemEvent(NotifyMessage message);
} }
@@ -58,7 +58,7 @@ public class AuthServiceImpl implements IAuthService {
data.put("password", password); data.put("password", password);
NotifyMessage message = new NotifyMessage("setPassword", data); NotifyMessage message = new NotifyMessage("setPassword", data);
sseService.notifySystemEvent(message); sseService.notification(message);
} }
@Override @Override
@@ -1,6 +1,7 @@
package cn.czh.service.impl; package cn.czh.service.impl;
import cn.czh.base.BusinessException; import cn.czh.base.BusinessException;
import cn.czh.dto.NotifyMessage;
import cn.czh.entity.SharedFile; import cn.czh.entity.SharedFile;
import cn.czh.entity.UploadFile; import cn.czh.entity.UploadFile;
import cn.czh.mapper.ShareFileMapper; import cn.czh.mapper.ShareFileMapper;
@@ -76,14 +77,22 @@ public class FileServiceImpl implements IFileService {
sharedFile.setFileIdentifier(identifier); sharedFile.setFileIdentifier(identifier);
sharedFile.setCreateTime(LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss"))); sharedFile.setCreateTime(LocalDateTime.now().format(DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")));
shareFileMapper.insert(sharedFile); shareFileMapper.insert(sharedFile);
sseService.notifySharedFileUpdate("add"); HashMap<String, Object> map = new HashMap<>();
map.put("fileIdentifier", identifier);
map.put("action", "add");
NotifyMessage message = new NotifyMessage("sharedFileUpdate", map);
sseService.notification(message);
} }
} }
@Override @Override
public void removeSharedFile(String fileIdentifier) { public void removeSharedFile(String fileIdentifier) {
shareFileMapper.delete(Wrappers.lambdaQuery(SharedFile.class).eq(SharedFile::getFileIdentifier, fileIdentifier)); shareFileMapper.delete(Wrappers.lambdaQuery(SharedFile.class).eq(SharedFile::getFileIdentifier, fileIdentifier));
sseService.notifySharedFileUpdate("remove"); HashMap<String, Object> map = new HashMap<>();
map.put("fileIdentifier", fileIdentifier);
map.put("action", "remote");
NotifyMessage message = new NotifyMessage("sharedFileUpdate", map);
sseService.notification(message);
} }
@Override @Override
@@ -16,41 +16,23 @@ import java.util.concurrent.ConcurrentHashMap;
@Service @Service
public class SseServiceImpl implements ISseService { public class SseServiceImpl implements ISseService {
private final Map<String, SseEmitter> sharedFilesEmitters = new ConcurrentHashMap<>(); private final Map<String, SseEmitter> clientEmitters = new ConcurrentHashMap<>();
private final Map<String, SseEmitter> systemEventEmitters = new ConcurrentHashMap<>();
@Override @Override
public SseEmitter subscribeSharedFiles(String clientId) { public SseEmitter subscribe(String remoteAddr) {
SseEmitter emitter = new SseEmitter(0L); // 设置超时时间 SseEmitter emitter = new SseEmitter(0L);
sharedFilesEmitters.put(clientId, emitter); clientEmitters.put(remoteAddr, emitter);
emitter.onCompletion(() -> sharedFilesEmitters.remove(clientId)); emitter.onCompletion(() -> clientEmitters.remove(remoteAddr));
emitter.onTimeout(() -> sharedFilesEmitters.remove(clientId)); emitter.onTimeout(() -> clientEmitters.remove(remoteAddr));
return emitter;
}
@Override
public SseEmitter subscribeSystemEvent(String clientId) {
SseEmitter emitter = new SseEmitter(0L); // 设置超时时间
systemEventEmitters.put(clientId, emitter);
emitter.onCompletion(() -> systemEventEmitters.remove(clientId));
emitter.onTimeout(() -> systemEventEmitters.remove(clientId));
return emitter; return emitter;
} }
@Async @Async
@Override @Override
public void notifySharedFileUpdate(String message) { public void notification(NotifyMessage message) {
sharedFilesEmitters.forEach((clientId, emitter) -> clientEmitters.forEach((clientId, emitter) ->
sendNotification(clientId, emitter, "sharedFileUpdate", message, sharedFilesEmitters)); sendNotification(clientId, emitter, message.getType(), JSONUtil.toJsonStr(message), clientEmitters));
}
@Async
@Override
public void notifySystemEvent(NotifyMessage message) {
systemEventEmitters.forEach((clientId, emitter) ->
sendNotification(clientId, emitter, "systemEvent", JSONUtil.toJsonStr(message), systemEventEmitters));
} }
private void sendNotification(String clientId, SseEmitter emitter, String event, String message, Map<String, SseEmitter> emittersMap) { private void sendNotification(String clientId, SseEmitter emitter, String event, String message, Map<String, SseEmitter> emittersMap) {
@@ -51,6 +51,7 @@
</template> </template>
<script> <script>
import sseManager from '@/utils/sse';
import { getPassword, setPassword } from '@/utils/api'; import { getPassword, setPassword } from '@/utils/api';
export default { export default {
@@ -64,10 +65,16 @@ export default {
tempPassword: '', tempPassword: '',
inputPassword: '', inputPassword: '',
showPasswordDialog: false, showPasswordDialog: false,
mainUser: false,
}; };
}, },
mounted() { mounted() {
this.getPassword(); this.getPassword();
sseManager.subscribe('setPassword', () => {
if (!this.mainUser) {
this.getPassword();
}
});
}, },
methods: { methods: {
async getPassword() { async getPassword() {
@@ -76,7 +83,7 @@ export default {
const dataStr = res.data.toString(); const dataStr = res.data.toString();
const password = dataStr.slice(0, -1); const password = dataStr.slice(0, -1);
const mainUser = dataStr.slice(-1) === '1'; const mainUser = dataStr.slice(-1) === '1';
this.mainUser = mainUser;
this.tempPassword = password; this.tempPassword = password;
if (password) { if (password) {
this.tempEnablePassword = true; this.tempEnablePassword = true;
@@ -68,7 +68,8 @@
</template> </template>
<script> <script>
import { getSharedFiles, enableShare, getShareStatus, shareAddress, unShareFile, subscribeSharedFiles, subscribeSystemEvent } from '@/utils/api'; import sseManager from '@/utils/sse'
import { getSharedFiles, enableShare, getShareStatus, shareAddress, unShareFile } from '@/utils/api';
export default { export default {
name: 'ShareFile', name: 'ShareFile',
@@ -80,8 +81,6 @@ export default {
fileList: [], fileList: [],
enableShare: false, enableShare: false,
isAdmin: false, isAdmin: false,
sharedFilesEventSource: null,
systemEventSource: null,
}; };
}, },
mounted() { mounted() {
@@ -89,26 +88,22 @@ export default {
this.getShareStatus(); this.getShareStatus();
this.getShareAddress(); this.getShareAddress();
this.sharedFilesEventSource = subscribeSharedFiles(() => { sseManager.subscribe('sharedFileUpdate', () => {
this.fetchFiles(); this.fetchFiles();
}); });
this.systemEventSource = subscribeSystemEvent((event) => { sseManager.subscribe('enableShare', (event) => {
if (event.type === 'enableShare') {
this.fileList = []; this.fileList = [];
this.enableShare = event.data.enable; this.enableShare = event.enable;
this.fetchFiles(); this.fetchFiles();
}
}); });
}, },
beforeDestroy() { beforeDestroy() {
if (this.sharedFilesEventSource) { sseManager.unsubscribe('sharedFileUpdate');
this.sharedFilesEventSource.close(); sseManager.unsubscribe('enableShare');
}
if (this.systemEventSource) {
this.systemEventSource.close();
}
}, },
methods: { methods: {
async getShareAddress() { async getShareAddress() {
try { try {
+1 -57
View File
@@ -1,4 +1,4 @@
import { http, baseUrl, httpExtra } from './http'; import { http, httpExtra } from './http';
/** /**
* 获取上传进度 * 获取上传进度
@@ -271,60 +271,6 @@ const setPassword = (password) => {
return http.post("/config/setPassword", formData); return http.post("/config/setPassword", formData);
} }
/**
* 订阅共享文件更新
* @param {Function} callback 文件更新的回调函数
* @returns {EventSource} 返回 EventSource 实例
*/
const subscribeSharedFiles = (callback) => {
const eventSource = new EventSource(`${baseUrl}/sse/subscribeSharedFiles`);
eventSource.addEventListener('sharedFileUpdate', (event) => {
callback(event.data);
});
eventSource.onerror = (error) => {
console.error("SSE连接异常:", error);
if (eventSource.readyState === EventSource.CLOSED) {
console.log("连接已关闭");
}
};
return eventSource;
};
/**
* 订阅系统事件更新
* @param {Function} callback 文件更新的回调函数
* @returns {EventSource} 返回 EventSource 实例
*/
const subscribeSystemEvent = (callback) => {
const eventSource = new EventSource(`${baseUrl}/sse/subscribeSystemEvent`);
eventSource.addEventListener('systemEvent', (event) => {
try {
const data = JSON.parse(event.data);
callback(data);
} catch (e) {
console.error('JSON解析失败:', e, '原始数据:', event.data);
}
});
eventSource.onmessage = (event) => {
console.log('[SSE] Received default message:', event.data);
};
eventSource.onerror = (error) => {
console.error("SSE连接异常:", error);
if (eventSource.readyState === EventSource.CLOSED) {
console.log("连接已关闭");
}
};
return eventSource;
}
export { export {
getUploadProgress, getUploadProgress,
createMultipartUpload, createMultipartUpload,
@@ -345,7 +291,5 @@ export {
setPassword, setPassword,
addSharedFile, addSharedFile,
deleteFile, deleteFile,
subscribeSharedFiles,
subscribeSystemEvent,
httpExtra httpExtra
}; };