-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathWsWebFluxHandle.java
44 lines (38 loc) · 1.72 KB
/
WsWebFluxHandle.java
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
package com.example;
import lombok.extern.slf4j.Slf4j;
import org.springframework.stereotype.Component;
import org.springframework.web.reactive.socket.WebSocketHandler;
import org.springframework.web.reactive.socket.WebSocketMessage;
import org.springframework.web.reactive.socket.WebSocketSession;
import reactor.core.publisher.Mono;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
@Slf4j
@Component
public class WsWebFluxHandle implements WebSocketHandler {
private static final Map<String, WebSocketSession> SESSION_MAP = new ConcurrentHashMap<>();
@Override
public Mono<Void> handle(WebSocketSession session) {
SESSION_MAP.put(session.getId(), session);
log.info("link: {}", session.getId());
return session.receive()
.doOnComplete(() -> {
log.info("close: {}", session.getId());
SESSION_MAP.remove(session.getId());
})
.doOnError(throwable -> {
log.error("error[{}]: ", session.getId(), throwable);
SESSION_MAP.remove(session.getId());
})
.doOnNext(message -> {
log.info("receive[{}]: {}", session.getId(), message.getPayloadAsText());
WebSocketMessage msg = session.textMessage(session.getId() + ": " + message.getPayloadAsText());
SESSION_MAP.keySet()
.stream()
.filter(s -> !s.equals(session.getId()))
.map(SESSION_MAP::get)
.forEach(webSocketSession -> webSocketSession.send(Mono.just(msg)).subscribe());
})
.then();
}
}