88import java .util .List ;
99import java .util .Map ;
1010import java .util .concurrent .ConcurrentHashMap ;
11- import java .util .concurrent .CopyOnWriteArrayList ;
1211
1312import lombok .extern .slf4j .Slf4j ;
1413import org .springframework .stereotype .Component ;
@@ -27,40 +26,39 @@ public class SseConnectionManager {
2726 private static final Long DEFAULT_TIMEOUT = 60L * 1000 * 60 ; // 60๋ถ
2827 private static final String SSE_EVENT_NAME = "notification" ;
2928
30- // Key: memberId, Value: List of SseEmitters (๋ค์ค ๋๋ฐ์ด์ค ์ง์ )
29+ // Key: memberId, Value: SseEmitter (๋จ์ผ ๋๋ฐ์ด์ค)
3130 private final Map <Long , SseEmitter > emitters = new ConcurrentHashMap <>();
3231
3332 public SseEmitter createConnection (Long memberId ) {
34- // ๊ธฐ์กด ์ฐ๊ฒฐ์ด ์์ผ๋ฉด ์ข
๋ฃ (์ ์ฐ๊ฒฐ ์ฐ์ )
35- SseEmitter existingEmitter = emitters .get (memberId );
36- if (existingEmitter != null ) {
37- existingEmitter .complete ();
38- log .info ("Existing SSE connection closed for member: {}" , memberId );
39- }
40-
4133 SseEmitter emitter = new SseEmitter (DEFAULT_TIMEOUT );
42- emitters .put (memberId , emitter );
34+
35+ // ๊ธฐ์กด ์ฐ๊ฒฐ์ด ์์ด๋ ๊ฐ์ ์ข
๋ฃํ์ง ์์ - ๊ทธ๋ฅ ์ ์ฐ๊ฒฐ๋ก ๊ต์ฒด
36+ // ๊ธฐ์กด emitter๋ ์์ฐ์ค๋ฝ๊ฒ timeout/completion ์ฒ๋ฆฌ๋จ
37+ SseEmitter oldEmitter = emitters .put (memberId , emitter );
38+ if (oldEmitter != null ) {
39+ log .info ("Replacing existing SSE connection for member: {} (old connection will timeout naturally)" , memberId );
40+ }
4341
4442 log .info ("SSE connection created for member: {}" , memberId );
4543
46- // ์ฐ๊ฒฐ ์๋ฃ ์ ์ ๊ฑฐ
44+ // ์ฐ๊ฒฐ ์๋ฃ ์ ์ ๊ฑฐ (ํ์ฌ ํ์ฑ emitter์ธ ๊ฒฝ์ฐ๋ง)
4745 emitter .onCompletion (
4846 () -> {
49- removeEmitter (memberId );
47+ removeEmitterIfMatch (memberId , emitter );
5048 log .info ("SSE connection completed for member: {}" , memberId );
5149 });
5250
53- // ํ์์์ ์ ์ ๊ฑฐ
51+ // ํ์์์ ์ ์ ๊ฑฐ (ํ์ฌ ํ์ฑ emitter์ธ ๊ฒฝ์ฐ๋ง)
5452 emitter .onTimeout (
5553 () -> {
56- removeEmitter (memberId );
54+ removeEmitterIfMatch (memberId , emitter );
5755 log .warn ("SSE connection timeout for member: {}" , memberId );
5856 });
5957
60- // ์๋ฌ ์ ์ ๊ฑฐ
58+ // ์๋ฌ ์ ์ ๊ฑฐ (ํ์ฌ ํ์ฑ emitter์ธ ๊ฒฝ์ฐ๋ง)
6159 emitter .onError (
6260 e -> {
63- removeEmitter (memberId );
61+ removeEmitterIfMatch (memberId , emitter );
6462 log .error ("SSE connection error for member: {}, error: {}" , memberId , e .getMessage ());
6563 });
6664
@@ -69,7 +67,7 @@ public SseEmitter createConnection(Long memberId) {
6967 emitter .send (SseEmitter .event ().name ("connected" ).data ("SSE connection established" ));
7068 } catch (IOException e ) {
7169 log .error ("Failed to send initial connection event to member: {}" , memberId , e );
72- removeEmitter (memberId );
70+ removeEmitterIfMatch (memberId , emitter );
7371 throw new NotificationHandler (NotificationErrorStatus .SSE_CONNECTION_FAILED );
7472 }
7573
@@ -95,7 +93,7 @@ public boolean sendToMember(Long memberId, Object data) {
9593 emitter .send (SseEmitter .event ().name (SSE_EVENT_NAME ).data (data ));
9694 log .info ("Notification sent to member: {}" , memberId );
9795 return true ;
98- } catch (IOException e ) {
96+ } catch (IOException | IllegalStateException e ) {
9997 log .error ("Failed to send notification to member: {}, removing dead emitter" , memberId , e );
10098 removeEmitter (memberId );
10199 return false ;
@@ -202,6 +200,14 @@ private void removeEmitter(Long memberId) {
202200 emitters .remove (memberId );
203201 }
204202
203+ /**
204+ * ํน์ Emitter๊ฐ ํ์ฌ ํ์ฑ emitter์ธ ๊ฒฝ์ฐ์๋ง ์ ๊ฑฐ.
205+ * ์ด๋ฏธ ์ ์ฐ๊ฒฐ๋ก ๊ต์ฒด๋ ๊ฒฝ์ฐ ์ด์ emitter์ ์ฝ๋ฐฑ์ด ์ emitter๋ฅผ ์ ๊ฑฐํ์ง ์๋๋ก ํจ.
206+ */
207+ private void removeEmitterIfMatch (Long memberId , SseEmitter emitter ) {
208+ emitters .remove (memberId , emitter );
209+ }
210+
205211 /**
206212 * ๋ชจ๋ ์ฐ๊ฒฐ ํด์ (์๋ฒ ์ข
๋ฃ ์ ์ฌ์ฉ).
207213 */
0 commit comments