fix(page+bot): 리뷰 지적사항 반영 - 서버 멤버십·같은음성채널 인가(1,2), SSE 공유구독 fan-out(3), botRpc 즉시폴링(4), 큐삭제 encoded 검증(6), 진행바 anchor 타이머(7), 캐시 파싱 가드(8), key 안정화(9), 볼륨 롤백(10), SSE 자동재연결(11), body userId 정리(12)
This commit is contained in:
@@ -72,11 +72,14 @@ class RedisClientClass {
|
||||
if (!data.userId) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "userId를 찾을수 없습니다." }));
|
||||
const guild = await getGuildById(data.serverId);
|
||||
if (!guild) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "guild를 찾을수 없습니다." }));
|
||||
if (!(await this.isMember(guild, data.userId))) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "이 서버의 멤버가 아닙니다." }));
|
||||
let player = lavalinkManager.getPlayer(guild.id);
|
||||
const voiceChannel = await getVoiceChannelById(guild, data.userId);
|
||||
if (!player) {
|
||||
if (!voiceChannel) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "음성채널에 들어가서 이용해주세요." }));
|
||||
player = (await channelJoin(guild, voiceChannel.id)).player;
|
||||
} else if (!voiceChannel || voiceChannel.id !== player.voiceChannelId) {
|
||||
return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "봇과 같은 음성채널에 있어야 합니다." }));
|
||||
}
|
||||
if (!player) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "세션을 찾을수 없습니다." }));
|
||||
await lavalinkManager.search(guild.id, data.track.url, data.userId, player);
|
||||
@@ -88,11 +91,14 @@ class RedisClientClass {
|
||||
if (!data.userId) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "userId를 찾을수 없습니다." }));
|
||||
const guild = await getGuildById(data.serverId);
|
||||
if (!guild) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "guild를 찾을수 없습니다." }));
|
||||
if (!(await this.isMember(guild, data.userId))) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "이 서버의 멤버가 아닙니다." }));
|
||||
let player = lavalinkManager.getPlayer(guild.id);
|
||||
const voiceChannel = await getVoiceChannelById(guild, data.userId);
|
||||
if (!player) {
|
||||
if (!voiceChannel) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "음성채널에 들어가서 이용해주세요." }));
|
||||
player = (await channelJoin(guild, voiceChannel.id)).player;
|
||||
} else if (!voiceChannel || voiceChannel.id !== player.voiceChannelId) {
|
||||
return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "봇과 같은 음성채널에 있어야 합니다." }));
|
||||
}
|
||||
if (!player) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "세션을 찾을수 없습니다." }));
|
||||
await lavalinkManager.search(guild.id, data.playlistUrl, data.userId, player);
|
||||
@@ -102,8 +108,10 @@ class RedisClientClass {
|
||||
const resultKey = `player:now:${data.requestId}`;
|
||||
if (!data.serverId) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "serverId를 찾을수 없습니다." }));
|
||||
if (!data.userId) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "userId를 찾을수 없습니다." }));
|
||||
const nowGuild = await getGuildById(data.serverId);
|
||||
if (!nowGuild) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "guild를 찾을수 없습니다." }));
|
||||
if (!(await this.isMember(nowGuild, data.userId))) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "이 서버의 멤버가 아닙니다." }));
|
||||
const player = lavalinkManager.getPlayer(data.serverId);
|
||||
// if (!player) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "player를 찾을수 없습니다." }));
|
||||
await this.pub.setex(resultKey, 60, JSON.stringify({
|
||||
success: true,
|
||||
botPlayer: !!player,
|
||||
@@ -118,8 +126,10 @@ class RedisClientClass {
|
||||
const resultKey = `queue:list:${data.requestId}`;
|
||||
if (!data.serverId) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "serverId를 찾을수 없습니다." }));
|
||||
if (!data.userId) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "userId를 찾을수 없습니다." }));
|
||||
const qlGuild = await getGuildById(data.serverId);
|
||||
if (!qlGuild) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "guild를 찾을수 없습니다." }));
|
||||
if (!(await this.isMember(qlGuild, data.userId))) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "이 서버의 멤버가 아닙니다." }));
|
||||
const player = lavalinkManager.getPlayer(data.serverId);
|
||||
// if (!player) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "player를 찾을수 없습니다." }));
|
||||
await this.pub.setex(resultKey, 60, JSON.stringify({ success: true, queue: player?.queue?.slice(1) ?? [] }));
|
||||
}
|
||||
if (data.action === "queue_set") {
|
||||
@@ -152,6 +162,12 @@ class RedisClientClass {
|
||||
// queue[0]은 현재 재생중인 곡이므로 실제 대기열은 queue[1]부터 시작
|
||||
// numIndex는 대기열(queue[1]~) 기준이므로 실제 splice 위치<EC9C84><ECB998> numIndex+1
|
||||
if (numIndex >= context.player.queue.length - 1) return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "index가 대기열 범위를 초과합니다." }));
|
||||
// 인덱스 신뢰 대신, 클라이언트가 지우려던 곡(encoded)과 실제 대상이 같은지 확인.
|
||||
// SSE로 큐가 갱신되는 찰나 인덱스가 밀려 다른 곡이 삭제되는 것을 방지.
|
||||
const removeTarget = context.player.queue[numIndex + 1];
|
||||
if (data.encoded && removeTarget?.encoded && removeTarget.encoded !== data.encoded) {
|
||||
return await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "대기열이 변경되었습니다. 새로고침 후 다시 시도해주세요." }));
|
||||
}
|
||||
const [removedTrack] = context.player.queue.splice(numIndex + 1, 1);
|
||||
await this.pub.setex(resultKey, 60, JSON.stringify({ success: true, removedTrack }));
|
||||
context.player.setMsg();
|
||||
@@ -232,6 +248,16 @@ class RedisClientClass {
|
||||
Logger.log(`[Redis Pub] bot -> site 전송: ${event}`);
|
||||
}
|
||||
|
||||
/**
|
||||
* 요청한 userId가 해당 guild의 멤버인지 확인(캐시 우선, 없으면 단건 fetch).
|
||||
* 대시보드가 보낸 serverId를 그대로 신뢰하지 않기 위한 인가 검증.
|
||||
*/
|
||||
private async isMember(guild: Guild, userId: string): Promise<boolean> {
|
||||
if (guild.members.cache.has(userId)) return true;
|
||||
const fetched = await guild.members.fetch(userId).catch(() => null);
|
||||
return !!fetched;
|
||||
}
|
||||
|
||||
private async getContext(guildId: string, resultKey: string, userId: string): Promise<{
|
||||
ok: true;
|
||||
guild: Guild;
|
||||
@@ -243,6 +269,11 @@ class RedisClientClass {
|
||||
await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "guild를 찾을수 없습니다." }));
|
||||
return { ok: false };
|
||||
}
|
||||
// 인가: 요청자가 이 서버의 멤버여야 함 (남의 서버 제어 차단)
|
||||
if (!(await this.isMember(guild, userId))) {
|
||||
await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "이 서버의 멤버가 아닙니다." }));
|
||||
return { ok: false };
|
||||
}
|
||||
let player = lavalinkManager.getPlayer(guild.id);
|
||||
const voiceChannel = await getVoiceChannelById(guild, userId);
|
||||
if (!player) {
|
||||
@@ -251,6 +282,10 @@ class RedisClientClass {
|
||||
return { ok: false };
|
||||
}
|
||||
player = (await channelJoin(guild, voiceChannel.id)).player;
|
||||
} else if (!voiceChannel || voiceChannel.id !== player.voiceChannelId) {
|
||||
// 이미 재생 중이면 봇과 같은 음성채널에 있는 사람만 조작 가능
|
||||
await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "봇과 같은 음성채널에 있어야 조작할 수 있습니다." }));
|
||||
return { ok: false };
|
||||
}
|
||||
if (!player) {
|
||||
await this.pub.setex(resultKey, 60, JSON.stringify({ success: false, message: "player를 찾을수 없습니다." }));
|
||||
|
||||
@@ -12,6 +12,7 @@ import {
|
||||
interface QueueRemoveBody {
|
||||
serverId?: unknown;
|
||||
index?: unknown;
|
||||
encoded?: unknown;
|
||||
}
|
||||
|
||||
export async function POST(request: Request) {
|
||||
@@ -37,6 +38,8 @@ export async function POST(request: Request) {
|
||||
serverId: serverIdResult.value,
|
||||
userId,
|
||||
index: indexResult.value,
|
||||
// 인덱스-트랙 일치 검증용(봇이 대상 encoded 불일치 시 거절). 문자열일 때만 전달.
|
||||
encoded: typeof bodyResult.data.encoded === "string" ? bodyResult.data.encoded : undefined,
|
||||
},
|
||||
timeoutMs: 5000,
|
||||
});
|
||||
|
||||
@@ -55,7 +55,7 @@ export default function MainContent({
|
||||
}
|
||||
|
||||
let endpoint = "";
|
||||
const bodyData: Record<string, unknown> = { serverId: selectedServer.id, userId };
|
||||
const bodyData: Record<string, unknown> = { serverId: selectedServer.id };
|
||||
if (actionType === 'player_play') {
|
||||
endpoint = "/api/player/play";
|
||||
bodyData.track = track;
|
||||
@@ -108,8 +108,16 @@ export default function MainContent({
|
||||
setIsFetching(true);
|
||||
const cached = sessionStorage.getItem("filtered_servers");
|
||||
if (cached) {
|
||||
setServers(JSON.parse(cached));
|
||||
setIsFetching(false);
|
||||
try {
|
||||
const parsed = JSON.parse(cached);
|
||||
if (Array.isArray(parsed)) {
|
||||
setServers(parsed);
|
||||
setIsFetching(false);
|
||||
}
|
||||
} catch {
|
||||
// 캐시 손상 시 무시하고 아래 fetch로 새로 받는다.
|
||||
sessionStorage.removeItem("filtered_servers");
|
||||
}
|
||||
}
|
||||
|
||||
fetch("/api/servers")
|
||||
@@ -162,7 +170,7 @@ export default function MainContent({
|
||||
const hasAnyResults = searchResults.spotify.length > 0 || searchResults.youtubeMusic.length > 0 || searchResults.youtubeVideo.length > 0;
|
||||
|
||||
const renderTrackCard = (track: SearchTrack) => (
|
||||
<div key={track.videoId || track.id} className="bg-neutral-800/40 p-3 rounded-xl hover:bg-neutral-800 transition-all group border border-transparent hover:border-neutral-700 shadow-md w-full">
|
||||
<div key={track.videoId || track.id || track.url || track.title} className="bg-neutral-800/40 p-3 rounded-xl hover:bg-neutral-800 transition-all group border border-transparent hover:border-neutral-700 shadow-md w-full">
|
||||
<div className="aspect-square bg-neutral-700 rounded-md mb-2 relative overflow-hidden shadow-lg">
|
||||
{track.thumbnail && <img src={track.thumbnail} className="w-full h-full object-cover" alt={track.title} />}
|
||||
<button
|
||||
|
||||
@@ -22,6 +22,9 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
const [position, setPosition] = useState<number>(0);
|
||||
const [duration, setDuration] = useState<number>(0);
|
||||
const isDragging = useRef<boolean>(false); // 재생바를 드래그 중인지 여부
|
||||
// 진행바 계산 기준점: 이 시점(at)에 위치가 pos였다 → 현재위치 = pos + (now - at)
|
||||
const anchorRef = useRef<{ pos: number; at: number }>({ pos: 0, at: Date.now() });
|
||||
const prevVolumeRef = useRef<number>(50); // 드래그 시작 시점 볼륨(실패 시 롤백용)
|
||||
// 👇 [추가할 부분] 볼륨 바를 잡고 있는지 여부 추적
|
||||
const [isVolumeDragging, setIsVolumeDragging] = useState<boolean>(false);
|
||||
|
||||
@@ -38,7 +41,6 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
serverId: selectedServer.id,
|
||||
userId: userId,
|
||||
})
|
||||
});
|
||||
const data = await res.json();
|
||||
@@ -52,7 +54,9 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
setVolume(typeof data.volume === "number" ? data.volume : 50);
|
||||
// 드래그 중이 아닐 때만 서버 시간으로 동기화 (안 그러면 드래그할 때 튐)
|
||||
if (!isDragging.current) {
|
||||
setPosition(data.position || 0);
|
||||
const p = data.position || 0;
|
||||
setPosition(p);
|
||||
anchorRef.current = { pos: p, at: Date.now() }; // 기준점 재설정
|
||||
}
|
||||
} else {
|
||||
setTrack(null);
|
||||
@@ -84,15 +88,16 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
console.warn("SSE JSON 파싱 실패:", err);
|
||||
}
|
||||
};
|
||||
eventSource.onerror = (error) => {
|
||||
console.error("Player SSE 연결 오류:", error);
|
||||
eventSource.close();
|
||||
// 에러 시 close 하지 않는다 — EventSource가 자동 재연결하도록 둔다(일시적 네트워크 끊김 복구).
|
||||
eventSource.onerror = () => {
|
||||
console.warn("Player SSE 일시 오류 — 자동 재연결 대기");
|
||||
};
|
||||
return () => eventSource.close();
|
||||
}, [selectedServer, fetchNowPlaying]);
|
||||
|
||||
// 3. 🌟 로컬 1초 타이머 & 10초 서버 동기화 통합 (재생 중일 때만 작동!)
|
||||
// isPaused 상태를 ref 로 들고 있어서, interval 콜백이 항상 최신 값을 읽도록 처리.
|
||||
// 3. 진행바 시계: setInterval 누적(+1000ms) 방식은 탭 비활성/지터로 실제와 어긋나므로,
|
||||
// "기준점(anchor) + 실제 경과시간(Date.now())" 으로 계산해 드리프트를 없앤다.
|
||||
// anchor 는 서버 동기화/탐색/일시정지 전환 시점마다 갱신된다(아래 setAnchor).
|
||||
const isPausedRef = useRef(isPaused);
|
||||
useEffect(() => {
|
||||
isPausedRef.current = isPaused;
|
||||
@@ -100,25 +105,20 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
|
||||
useEffect(() => {
|
||||
if (!isPlaying) return;
|
||||
|
||||
// ① 1초마다 프론트엔드 단독으로 시계 굴리기 (부드러운 애니메이션용)
|
||||
const localInterval = setInterval(() => {
|
||||
// 250ms마다 anchor 기준으로 재계산(부드러운 진행)
|
||||
const tick = setInterval(() => {
|
||||
if (isPausedRef.current || isDragging.current) return;
|
||||
setPosition((prev) => {
|
||||
if (prev >= duration) return duration;
|
||||
return prev + 1000;
|
||||
});
|
||||
}, 1000);
|
||||
|
||||
// ② 10초마다 진짜 시간 서버에 물어보기 (오차 교정용)
|
||||
const a = anchorRef.current;
|
||||
const next = a.pos + (Date.now() - a.at);
|
||||
setPosition(duration > 0 ? Math.min(duration, next) : next);
|
||||
}, 250);
|
||||
// 15초마다 서버 시간으로 오차 교정(anchor 재설정은 fetchNowPlaying 내부에서)
|
||||
const syncInterval = setInterval(() => {
|
||||
if (isPausedRef.current || isDragging.current) return;
|
||||
fetchNowPlaying();
|
||||
}, 10000);
|
||||
|
||||
// 일시정지되거나 컴포넌트가 꺼지면 두 타이머 모두 깔끔하게 청소합니다.
|
||||
}, 15000);
|
||||
return () => {
|
||||
clearInterval(localInterval);
|
||||
clearInterval(tick);
|
||||
clearInterval(syncInterval);
|
||||
};
|
||||
}, [isPlaying, duration, fetchNowPlaying]);
|
||||
@@ -135,13 +135,13 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
const nextPaused = !isPaused;
|
||||
// UI 즉각 반영 (Optimistic UI)
|
||||
setIsPaused(nextPaused);
|
||||
anchorRef.current = { pos: position, at: Date.now() }; // 일시정지/재개 시점 기준 갱신
|
||||
try {
|
||||
const res = await fetch('/api/player/pause', {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
serverId: selectedServer.id,
|
||||
userId: userId,
|
||||
isPaused: nextPaused, // boolean 그대로 전송
|
||||
})
|
||||
});
|
||||
@@ -170,7 +170,6 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
serverId: selectedServer.id,
|
||||
userId: userId,
|
||||
})
|
||||
});
|
||||
// 성공하면 곧 SSE 이벤트가 와서 fetchNowPlaying을 트리거하겠지만, 즉각 반응을 위해 찔러줌
|
||||
@@ -189,6 +188,8 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
|
||||
// 🌟 [수정됨] e.target 대신 e.currentTarget을 사용해야 타입 에러가 나지 않습니다.
|
||||
const newPosition = Number(e.currentTarget.value);
|
||||
setPosition(newPosition);
|
||||
anchorRef.current = { pos: newPosition, at: Date.now() }; // 탐색 후 기준 갱신
|
||||
|
||||
try {
|
||||
await fetch('/api/player/seek', {
|
||||
@@ -197,7 +198,6 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
body: JSON.stringify({
|
||||
serverId: selectedServer.id,
|
||||
seek: newPosition,
|
||||
userId: userId,
|
||||
})
|
||||
});
|
||||
} catch (error) {
|
||||
@@ -221,17 +221,19 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
setVolume(finalVolume); // UI 즉시 반영
|
||||
|
||||
try {
|
||||
await fetch('/api/player/volume', {
|
||||
const res = await fetch('/api/player/volume', {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
serverId: selectedServer.id,
|
||||
userId: userId,
|
||||
volume: finalVolume,
|
||||
})
|
||||
});
|
||||
const data = await res.json().catch(() => ({}));
|
||||
if (!res.ok || !data.success) setVolume(prevVolumeRef.current); // 실패 시 롤백
|
||||
} catch (error) {
|
||||
console.error("볼륨 조절 에러:", error);
|
||||
setVolume(prevVolumeRef.current); // 실패 시 롤백
|
||||
}
|
||||
};
|
||||
|
||||
@@ -349,8 +351,8 @@ export default function PlayerBar({ selectedServer }: PlayerBarProps) {
|
||||
value={volume}
|
||||
disabled={!botPlayer || !isPlaying || !track}
|
||||
onChange={handleVolumeChange} // 눈에 보이는 볼륨만 즉시 변경
|
||||
onMouseDown={() => setIsVolumeDragging(true)}
|
||||
onTouchStart={() => setIsVolumeDragging(true)}
|
||||
onMouseDown={() => { prevVolumeRef.current = volume; setIsVolumeDragging(true); }}
|
||||
onTouchStart={() => { prevVolumeRef.current = volume; setIsVolumeDragging(true); }}
|
||||
onMouseUp={handleVolumeEnd} // 🌟 마우스를 뗐을 때 봇으로 전송
|
||||
onTouchEnd={handleVolumeEnd} // 🌟 스마트폰 터치를 뗐을 때 봇으로 전송
|
||||
className="absolute w-full h-1 opacity-0 cursor-pointer z-20"
|
||||
|
||||
@@ -70,9 +70,9 @@ export default function QueueSidebar({ selectedServer }: QueueSidebarProps) {
|
||||
}
|
||||
};
|
||||
|
||||
eventSource.onerror = (error) => {
|
||||
console.error("SSE 연결 오류:", error);
|
||||
eventSource.close();
|
||||
// 에러 시 close 하지 않는다 — EventSource 자동 재연결에 맡긴다(일시 끊김 복구).
|
||||
eventSource.onerror = () => {
|
||||
console.warn("Queue SSE 일시 오류 — 자동 재연결 대기");
|
||||
};
|
||||
|
||||
return () => {
|
||||
@@ -142,8 +142,9 @@ export default function QueueSidebar({ selectedServer }: QueueSidebarProps) {
|
||||
|
||||
const handleDelete = async (indexToRemove: number) => {
|
||||
if (!selectedServer) return;
|
||||
const userId = session?.user?.id;
|
||||
|
||||
// 지우려는 곡의 encoded 를 함께 보내 서버가 인덱스-트랙 일치를 검증(엉뚱한 곡 삭제 방지).
|
||||
const targetEncoded = queue[indexToRemove]?.encoded;
|
||||
const newQueue = queue.filter((_, index) => index !== indexToRemove);
|
||||
setQueue(newQueue);
|
||||
|
||||
@@ -153,8 +154,8 @@ export default function QueueSidebar({ selectedServer }: QueueSidebarProps) {
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify({
|
||||
serverId: selectedServer.id,
|
||||
userId: userId,
|
||||
index: indexToRemove,
|
||||
encoded: targetEncoded,
|
||||
})
|
||||
});
|
||||
const data = await res.json();
|
||||
@@ -185,7 +186,7 @@ export default function QueueSidebar({ selectedServer }: QueueSidebarProps) {
|
||||
|
||||
return (
|
||||
<div
|
||||
key={`${item.id || index}-${item.info?.title || ''}`}
|
||||
key={(item.encoded as string) || `${index}-${item.info?.title || ''}`}
|
||||
draggable
|
||||
onDragStart={() => {
|
||||
dragItem.current = index;
|
||||
|
||||
@@ -66,11 +66,18 @@ export async function botRpc(
|
||||
);
|
||||
|
||||
const deadline = Date.now() + timeoutMs;
|
||||
let interval = pollIntervalMs;
|
||||
// 봇은 보통 수십 ms 안에 응답하므로, 먼저 즉시 확인하고 촘촘히 폴링한다(초기 지연 제거).
|
||||
// (이상적으로는 봇이 reply 채널로 publish → 사이트가 구독하는 pub/sub 방식이나,
|
||||
// 봇의 30여 개 setex 지점을 모두 바꿔야 해 회귀 위험이 커서 폴링 최적화로 대체.)
|
||||
let interval = Math.min(pollIntervalMs, 30);
|
||||
let first = true;
|
||||
|
||||
while (Date.now() < deadline) {
|
||||
await sleep(interval);
|
||||
interval = Math.min(interval * 2, 400);
|
||||
if (!first) {
|
||||
await sleep(interval);
|
||||
interval = Math.min(Math.floor(interval * 1.6), 250);
|
||||
}
|
||||
first = false;
|
||||
|
||||
const reply = await Redis.get(resultKey);
|
||||
if (!reply) continue;
|
||||
|
||||
@@ -17,12 +17,47 @@ interface BotEventStreamOptions {
|
||||
clientEventType?: string;
|
||||
}
|
||||
|
||||
// ─────────────────────────────────────────────────────────────
|
||||
// 공유 구독자(fan-out): "bot-site" 채널은 프로세스당 구독자 1개만 두고,
|
||||
// 메모리에서 리스너들에게 분배한다. (SSE 연결마다 Redis 커넥션을 복제하지 않음)
|
||||
// ─────────────────────────────────────────────────────────────
|
||||
interface Listener {
|
||||
guildId: string;
|
||||
botEventName: string;
|
||||
send: () => void;
|
||||
}
|
||||
|
||||
const listeners = new Set<Listener>();
|
||||
const globalForSse = global as unknown as { botSiteSubscriber?: ReturnType<typeof Redis.duplicate> };
|
||||
|
||||
function ensureSubscriber() {
|
||||
if (globalForSse.botSiteSubscriber) return;
|
||||
const sub = Redis.duplicate();
|
||||
globalForSse.botSiteSubscriber = sub;
|
||||
sub.on("error", (err) => Logger.error(`[SSE] shared subscriber error: ${err.message}`));
|
||||
sub.subscribe("bot-site").catch((err) => Logger.error(`[SSE] subscribe 실패: ${String(err)}`));
|
||||
sub.on("message", (channel, message) => {
|
||||
if (channel !== "bot-site") return;
|
||||
let data: BotEvent;
|
||||
try {
|
||||
data = JSON.parse(message) as BotEvent;
|
||||
} catch (err) {
|
||||
Logger.warn(`[SSE] 잘못된 JSON: ${String(err)}`);
|
||||
return;
|
||||
}
|
||||
for (const l of listeners) {
|
||||
if (data.guildId === l.guildId && data.event === l.botEventName) {
|
||||
try { l.send(); } catch { /* 개별 클라이언트 전송 실패는 무시 */ }
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* 봇이 publish 하는 "bot-site" 채널을 구독해서 SSE 로 흘려보내는 공용 핸들러.
|
||||
* 봇이 publish 하는 "bot-site" 채널을 공유 구독자로 받아 SSE 로 흘려보내는 공용 핸들러.
|
||||
* - 인증 가드: 세션 없으면 401
|
||||
* - serverId 검증
|
||||
* - JSON.parse 안전 처리
|
||||
* - subscriber error / 클라이언트 abort 모두에서 깔끔히 정리
|
||||
* - 공유 구독자 + 메모리 fan-out (연결마다 Redis 커넥션 복제 안 함)
|
||||
* - keepalive ping (30초)
|
||||
*/
|
||||
export async function botEventStream(req: NextRequest, opts: BotEventStreamOptions): Promise<Response> {
|
||||
@@ -38,83 +73,41 @@ export async function botEventStream(req: NextRequest, opts: BotEventStreamOptio
|
||||
return new Response("Missing serverId", { status: 400 });
|
||||
}
|
||||
|
||||
const stream = new ReadableStream({
|
||||
async start(controller) {
|
||||
const subscriber = Redis.duplicate();
|
||||
let closed = false;
|
||||
const timers: NodeJS.Timeout[] = [];
|
||||
ensureSubscriber();
|
||||
|
||||
const cleanup = async () => {
|
||||
const stream = new ReadableStream({
|
||||
start(controller) {
|
||||
const encoder = new TextEncoder();
|
||||
let closed = false;
|
||||
|
||||
const cleanup = () => {
|
||||
if (closed) return;
|
||||
closed = true;
|
||||
for (const t of timers) clearInterval(t);
|
||||
try {
|
||||
await subscriber.unsubscribe("bot-site");
|
||||
} catch {
|
||||
/* noop */
|
||||
}
|
||||
try {
|
||||
await subscriber.quit();
|
||||
} catch {
|
||||
/* noop */
|
||||
}
|
||||
try {
|
||||
controller.close();
|
||||
} catch {
|
||||
/* 이미 닫혔을 수 있음 */
|
||||
}
|
||||
clearInterval(ping);
|
||||
listeners.delete(listener);
|
||||
try { controller.close(); } catch { /* 이미 닫힘 */ }
|
||||
};
|
||||
|
||||
subscriber.on("error", async (err) => {
|
||||
Logger.error(`[SSE:${opts.botEventName}] subscriber error: ${err.message}`);
|
||||
await cleanup();
|
||||
});
|
||||
|
||||
try {
|
||||
await subscriber.subscribe("bot-site");
|
||||
} catch (err) {
|
||||
Logger.error(`[SSE:${opts.botEventName}] subscribe 실패: ${String(err)}`);
|
||||
await cleanup();
|
||||
return;
|
||||
}
|
||||
|
||||
subscriber.on("message", (channel, message) => {
|
||||
if (channel !== "bot-site" || closed) return;
|
||||
let data: BotEvent;
|
||||
try {
|
||||
data = JSON.parse(message) as BotEvent;
|
||||
} catch (err) {
|
||||
Logger.warn(`[SSE:${opts.botEventName}] 잘못된 JSON: ${String(err)}`);
|
||||
return;
|
||||
}
|
||||
if (data.guildId !== serverId) return;
|
||||
if (data.event !== opts.botEventName) return;
|
||||
try {
|
||||
controller.enqueue(
|
||||
new TextEncoder().encode(`data: ${JSON.stringify({ type: clientEventType })}\n\n`),
|
||||
);
|
||||
} catch (err) {
|
||||
Logger.warn(`[SSE:${opts.botEventName}] enqueue 실패: ${String(err)}`);
|
||||
void cleanup();
|
||||
}
|
||||
});
|
||||
|
||||
// 30초마다 keep-alive 코멘트 전송 (프록시 timeout 방지)
|
||||
timers.push(
|
||||
setInterval(() => {
|
||||
const listener: Listener = {
|
||||
guildId: serverId,
|
||||
botEventName: opts.botEventName,
|
||||
send: () => {
|
||||
if (closed) return;
|
||||
try {
|
||||
controller.enqueue(new TextEncoder().encode(`: keep-alive\n\n`));
|
||||
} catch {
|
||||
void cleanup();
|
||||
}
|
||||
}, 30000),
|
||||
);
|
||||
controller.enqueue(encoder.encode(`data: ${JSON.stringify({ type: clientEventType })}\n\n`));
|
||||
},
|
||||
};
|
||||
listeners.add(listener);
|
||||
|
||||
// 클라이언트가 연결을 끊으면 정리
|
||||
req.signal.addEventListener("abort", () => {
|
||||
void cleanup();
|
||||
});
|
||||
const ping = setInterval(() => {
|
||||
if (closed) return;
|
||||
try {
|
||||
controller.enqueue(encoder.encode(`: keep-alive\n\n`));
|
||||
} catch {
|
||||
cleanup();
|
||||
}
|
||||
}, 30000);
|
||||
|
||||
req.signal.addEventListener("abort", cleanup);
|
||||
},
|
||||
});
|
||||
|
||||
|
||||
Reference in New Issue
Block a user