Compare commits
3
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
10d7ecd405 | ||
|
|
ff554fcff2 | ||
|
|
f8b253ba5e |
@@ -279,6 +279,25 @@ export async function initializeDiscordGateway() {
|
|||||||
});
|
});
|
||||||
|
|
||||||
process.on("uncaughtException", (err) => {
|
process.on("uncaughtException", (err) => {
|
||||||
|
const code =
|
||||||
|
typeof (err as NodeJS.ErrnoException).code === "string"
|
||||||
|
? (err as NodeJS.ErrnoException).code
|
||||||
|
: "";
|
||||||
|
// Transient stream-teardown errors (voice stop/disconnect races, child
|
||||||
|
// process stdin closed while we still write) are NOT fatal — crashing the
|
||||||
|
// gateway on EPIPE takes the whole bot offline mid-music. Log + continue.
|
||||||
|
if (
|
||||||
|
code === "EPIPE" ||
|
||||||
|
code === "ERR_STREAM_DESTROYED" ||
|
||||||
|
code === "ERR_STREAM_WRITE_AFTER_END" ||
|
||||||
|
code === "ECONNRESET"
|
||||||
|
) {
|
||||||
|
logger.warn(
|
||||||
|
{ error: err },
|
||||||
|
"Uncaught transient stream error — continuing",
|
||||||
|
);
|
||||||
|
return;
|
||||||
|
}
|
||||||
logger.error({ error: err }, "Uncaught exception");
|
logger.error({ error: err }, "Uncaught exception");
|
||||||
gracefulShutdown("uncaughtException");
|
gracefulShutdown("uncaughtException");
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -262,6 +262,9 @@ a=ice-lite
|
|||||||
[audioSection, videoSection, videoRtpMap].join("\n"),
|
[audioSection, videoSection, videoRtpMap].join("\n"),
|
||||||
"answer",
|
"answer",
|
||||||
);
|
);
|
||||||
|
console.log(
|
||||||
|
`[goLive:${this.constructor.name}] SELECT_PROTOCOL_ACK processed — remote answer set (${[audioSection, videoSection].join("\n").length}B)`,
|
||||||
|
);
|
||||||
this.emit("select_protocol_ack");
|
this.emit("select_protocol_ack");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -519,10 +522,14 @@ a=ice-lite
|
|||||||
const reconnect = () => {
|
const reconnect = () => {
|
||||||
const webRtcConn = this._webRtcWrapper.initWebRtc();
|
const webRtcConn = this._webRtcWrapper.initWebRtc();
|
||||||
webRtcConn.onStateChange((state) => {
|
webRtcConn.onStateChange((state) => {
|
||||||
|
console.log(`[goLive:${this.constructor.name}] pc state => ${state}`);
|
||||||
if (state === "closed" && !this._closed) reconnect();
|
if (state === "closed" && !this._closed) reconnect();
|
||||||
});
|
});
|
||||||
this._webRtcWrapper.onLocalDescription = (sdp) => {
|
this._webRtcWrapper.onLocalDescription = (sdp) => {
|
||||||
const rtc_connection_id = randomUUID();
|
const rtc_connection_id = randomUUID();
|
||||||
|
console.log(
|
||||||
|
`[goLive:${this.constructor.name}] sending SELECT_PROTOCOL (offer ${sdp.length}B, rtc_connection_id=${rtc_connection_id.slice(0, 8)})`,
|
||||||
|
);
|
||||||
this.sendOpcode(VoiceOpCodes.SELECT_PROTOCOL, {
|
this.sendOpcode(VoiceOpCodes.SELECT_PROTOCOL, {
|
||||||
protocol: "webrtc",
|
protocol: "webrtc",
|
||||||
codecs: Object.values(CodecPayloadType),
|
codecs: Object.values(CodecPayloadType),
|
||||||
|
|||||||
@@ -48,7 +48,19 @@ export class Streamer {
|
|||||||
this._client = client;
|
this._client = client;
|
||||||
// listen for gateway dispatch events
|
// listen for gateway dispatch events
|
||||||
this.client.on("raw", (packet) => {
|
this.client.on("raw", (packet) => {
|
||||||
this._gatewayEmitter.emit(packet.t, packet.d);
|
const t = packet.t as string;
|
||||||
|
if (
|
||||||
|
t === "STREAM_CREATE" ||
|
||||||
|
t === "STREAM_SERVER_UPDATE" ||
|
||||||
|
t === "VOICE_STATE_UPDATE" ||
|
||||||
|
t === "VOICE_SERVER_UPDATE"
|
||||||
|
) {
|
||||||
|
console.log(
|
||||||
|
`[goLive:Streamer] raw dispatch ${t}`,
|
||||||
|
JSON.stringify(packet.d).slice(0, 220),
|
||||||
|
);
|
||||||
|
}
|
||||||
|
this._gatewayEmitter.emit(t, packet.d);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -142,6 +154,13 @@ export class Streamer {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
this.signalStream();
|
this.signalStream();
|
||||||
|
const streamTimeout = setTimeout(() => {
|
||||||
|
reject(
|
||||||
|
new Error(
|
||||||
|
"Timed out waiting for STREAM_CREATE/STREAM_SERVER_UPDATE from Discord (stream handshake) — voice media session may not be active",
|
||||||
|
),
|
||||||
|
);
|
||||||
|
}, 12_000);
|
||||||
const {
|
const {
|
||||||
guildId: clientGuildId,
|
guildId: clientGuildId,
|
||||||
channelId: clientChannelId,
|
channelId: clientChannelId,
|
||||||
@@ -155,6 +174,7 @@ export class Streamer {
|
|||||||
clientUserId,
|
clientUserId,
|
||||||
clientChannelId,
|
clientChannelId,
|
||||||
(conn) => {
|
(conn) => {
|
||||||
|
clearTimeout(streamTimeout);
|
||||||
resolve(conn);
|
resolve(conn);
|
||||||
},
|
},
|
||||||
);
|
);
|
||||||
|
|||||||
@@ -79,6 +79,22 @@ export function transcodeToHighQualityOgg(
|
|||||||
);
|
);
|
||||||
|
|
||||||
input.pipe(proc.stdin);
|
input.pipe(proc.stdin);
|
||||||
|
// ffmpeg teardown closes stdin while the upstream source may still write —
|
||||||
|
// swallow EPIPE / destroyed-stream errors so they don't crash the gateway.
|
||||||
|
proc.stdin.on("error", (err: NodeJS.ErrnoException) => {
|
||||||
|
if (
|
||||||
|
err.code === "EPIPE" ||
|
||||||
|
err.code === "ERR_STREAM_DESTROYED" ||
|
||||||
|
err.code === "ERR_STREAM_WRITE_AFTER_END"
|
||||||
|
) {
|
||||||
|
logger.debug(
|
||||||
|
{ code: err.code },
|
||||||
|
"Transcode stdin closed during teardown",
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
logger.error({ error: err.message }, "Transcode stdin error");
|
||||||
|
}
|
||||||
|
});
|
||||||
activeProcesses.add(proc);
|
activeProcesses.add(proc);
|
||||||
|
|
||||||
const cleanup = () => {
|
const cleanup = () => {
|
||||||
@@ -248,6 +264,20 @@ export function resolveMediaUrl(
|
|||||||
// `--print` headers to stderr — pipe stdout immediately so the child
|
// `--print` headers to stderr — pipe stdout immediately so the child
|
||||||
// never blocks on a full pipe while we wait for the headers on stderr.
|
// never blocks on a full pipe while we wait for the headers on stderr.
|
||||||
const mediaStream = new PassThrough();
|
const mediaStream = new PassThrough();
|
||||||
|
// Teardown (player stop / ffmpeg exit) destroys this stream while
|
||||||
|
// yt-dlp may still push bytes — without a listener an EPIPE /
|
||||||
|
// ERR_STREAM_DESTROYED surfaces as an uncaughtException.
|
||||||
|
mediaStream.on("error", (err: NodeJS.ErrnoException) => {
|
||||||
|
if (
|
||||||
|
err.code === "EPIPE" ||
|
||||||
|
err.code === "ERR_STREAM_DESTROYED" ||
|
||||||
|
err.code === "ERR_STREAM_WRITE_AFTER_END"
|
||||||
|
) {
|
||||||
|
logger.debug({ code: err.code }, "Media stream closed during teardown");
|
||||||
|
} else {
|
||||||
|
logger.error({ error: err.message }, "Media stream error");
|
||||||
|
}
|
||||||
|
});
|
||||||
proc.stdout.pipe(mediaStream);
|
proc.stdout.pipe(mediaStream);
|
||||||
|
|
||||||
let stderrBuf = "";
|
let stderrBuf = "";
|
||||||
|
|||||||
@@ -56,6 +56,24 @@ export class VoiceTransmitter {
|
|||||||
// Create PCM input stream
|
// Create PCM input stream
|
||||||
this.pcmStream = new PassThrough();
|
this.pcmStream = new PassThrough();
|
||||||
this.pcmStream.setMaxListeners(32); // drain listeners accumulate during backpressure
|
this.pcmStream.setMaxListeners(32); // drain listeners accumulate during backpressure
|
||||||
|
// Voice teardown (stop / disconnect / ffmpeg exit) destroys this stream
|
||||||
|
// while Redis PCM messages may still be in flight. Without a listener,
|
||||||
|
// EPIPE / ERR_STREAM_DESTROYED / ERR_STREAM_WRITE_AFTER_END surface as
|
||||||
|
// an uncaughtException and crash the whole gateway.
|
||||||
|
this.pcmStream.on("error", (err: NodeJS.ErrnoException) => {
|
||||||
|
if (
|
||||||
|
err.code === "EPIPE" ||
|
||||||
|
err.code === "ERR_STREAM_DESTROYED" ||
|
||||||
|
err.code === "ERR_STREAM_WRITE_AFTER_END"
|
||||||
|
) {
|
||||||
|
logger.debug(
|
||||||
|
{ code: err.code },
|
||||||
|
"PCM stream closed during voice teardown — ignoring",
|
||||||
|
);
|
||||||
|
} else {
|
||||||
|
logger.error({ error: err.message }, "PCM stream error");
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
// Spawn FFmpeg to encode 24kHz mono PCM → OggOpus
|
// Spawn FFmpeg to encode 24kHz mono PCM → OggOpus
|
||||||
// Input: 24kHz mono s16le (raw PCM)
|
// Input: 24kHz mono s16le (raw PCM)
|
||||||
@@ -146,7 +164,12 @@ export class VoiceTransmitter {
|
|||||||
);
|
);
|
||||||
|
|
||||||
this.redisSub.on("message", (channel, message) => {
|
this.redisSub.on("message", (channel, message) => {
|
||||||
if (channel !== this.TRANSMIT_CHANNEL || !this.pcmStream) return;
|
if (
|
||||||
|
!this.isActive ||
|
||||||
|
channel !== this.TRANSMIT_CHANNEL ||
|
||||||
|
!this.pcmStream
|
||||||
|
)
|
||||||
|
return;
|
||||||
|
|
||||||
try {
|
try {
|
||||||
const data = JSON.parse(message);
|
const data = JSON.parse(message);
|
||||||
@@ -161,11 +184,21 @@ export class VoiceTransmitter {
|
|||||||
this.draining = false;
|
this.draining = false;
|
||||||
// Re-acquire stream reference (could have been replaced by restart)
|
// Re-acquire stream reference (could have been replaced by restart)
|
||||||
const currentStream = this.pcmStream;
|
const currentStream = this.pcmStream;
|
||||||
if (!currentStream) return;
|
if (!currentStream || !this.isActive) return;
|
||||||
// Flush queued chunks
|
// Flush queued chunks
|
||||||
while (this.backpressureQueue.length > 0) {
|
while (this.backpressureQueue.length > 0) {
|
||||||
const queued = this.backpressureQueue.shift()!;
|
const queued = this.backpressureQueue.shift()!;
|
||||||
if (!currentStream.write(queued)) break;
|
try {
|
||||||
|
if (!currentStream.write(queued)) break;
|
||||||
|
} catch (err) {
|
||||||
|
logger.debug(
|
||||||
|
{
|
||||||
|
error: err instanceof Error ? err.message : String(err),
|
||||||
|
},
|
||||||
|
"PCM flush write failed during teardown — ignoring",
|
||||||
|
);
|
||||||
|
break;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user