Added retry mechanism to proxy for transient upstream 5xx erors
This commit is contained in:
+125
-103
@@ -257,121 +257,143 @@ router.post('/epg/:sourceId/channels', async (req, res) => {
|
|||||||
* This handles CORS for streams that don't allow cross-origin
|
* This handles CORS for streams that don't allow cross-origin
|
||||||
*/
|
*/
|
||||||
router.get('/stream', async (req, res) => {
|
router.get('/stream', async (req, res) => {
|
||||||
try {
|
const maxRetries = 2;
|
||||||
let { url } = req.query;
|
let lastError = null;
|
||||||
if (!url) {
|
|
||||||
return res.status(400).json({ error: 'URL required' });
|
|
||||||
}
|
|
||||||
|
|
||||||
// Forward some headers to be more "transparent" back to the origin
|
for (let attempt = 1; attempt <= maxRetries; attempt++) {
|
||||||
const isPluto = url.includes('pluto.tv');
|
try {
|
||||||
|
let { url } = req.query;
|
||||||
const headers = {
|
if (!url) {
|
||||||
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
return res.status(400).json({ error: 'URL required' });
|
||||||
'Accept': '*/*',
|
|
||||||
'Accept-Language': 'en-US,en;q=0.9',
|
|
||||||
// Using https and matching the origin of the request
|
|
||||||
'Origin': isPluto ? 'https://pluto.tv' : new URL(url).origin,
|
|
||||||
'Referer': isPluto ? 'https://pluto.tv/' : new URL(url).origin + '/'
|
|
||||||
};
|
|
||||||
|
|
||||||
const response = await fetch(url, { headers });
|
|
||||||
if (!response.ok) {
|
|
||||||
console.error(`Upstream error for ${url.substring(0, 80)}...: ${response.status} ${response.statusText}`);
|
|
||||||
if (response.status === 403) {
|
|
||||||
const errorBody = await response.text().catch(() => 'N/A');
|
|
||||||
console.error(`403 Response body: ${errorBody.substring(0, 200)}`);
|
|
||||||
}
|
|
||||||
return res.status(response.status).send(`Failed to fetch stream: ${response.statusText}`);
|
|
||||||
}
|
|
||||||
|
|
||||||
const contentType = response.headers.get('content-type') || '';
|
|
||||||
res.set('Access-Control-Allow-Origin', '*');
|
|
||||||
|
|
||||||
// Create an async iterator for the response body
|
|
||||||
const iterator = response.body[Symbol.asyncIterator]();
|
|
||||||
const first = await iterator.next();
|
|
||||||
|
|
||||||
if (first.done) {
|
|
||||||
res.set('Content-Type', contentType || 'application/octet-stream');
|
|
||||||
return res.end();
|
|
||||||
}
|
|
||||||
|
|
||||||
const firstChunk = Buffer.from(first.value);
|
|
||||||
|
|
||||||
// Peek at first bytes to check for HLS manifest ({ #EXTM3U })
|
|
||||||
const textPrefix = firstChunk.subarray(0, 7).toString('utf8');
|
|
||||||
const contentLooksLikeHls = textPrefix === '#EXTM3U';
|
|
||||||
|
|
||||||
if (contentLooksLikeHls) {
|
|
||||||
// HLS Manifest: We must read the WHOLE manifest to rewrite it
|
|
||||||
const chunks = [firstChunk];
|
|
||||||
|
|
||||||
// Consume the rest of the stream
|
|
||||||
let result = await iterator.next();
|
|
||||||
while (!result.done) {
|
|
||||||
chunks.push(Buffer.from(result.value));
|
|
||||||
result = await iterator.next();
|
|
||||||
}
|
}
|
||||||
|
|
||||||
const buffer = Buffer.concat(chunks);
|
// Forward some headers to be more "transparent" back to the origin
|
||||||
const finalUrl = response.url || url;
|
const isPluto = url.includes('pluto.tv');
|
||||||
console.log(`[Proxy] Processing HLS manifest from: ${finalUrl.substring(0, 80)}...`);
|
|
||||||
res.set('Content-Type', 'application/vnd.apple.mpegurl');
|
|
||||||
|
|
||||||
let manifest = buffer.toString('utf-8');
|
const headers = {
|
||||||
|
'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/120.0.0.0 Safari/537.36',
|
||||||
|
'Accept': '*/*',
|
||||||
|
'Accept-Language': 'en-US,en;q=0.9',
|
||||||
|
// Using https and matching the origin of the request
|
||||||
|
'Origin': isPluto ? 'https://pluto.tv' : new URL(url).origin,
|
||||||
|
'Referer': isPluto ? 'https://pluto.tv/' : new URL(url).origin + '/'
|
||||||
|
};
|
||||||
|
|
||||||
const finalUrlObj = new URL(finalUrl);
|
const response = await fetch(url, { headers });
|
||||||
const baseUrl = finalUrlObj.origin + finalUrlObj.pathname.substring(0, finalUrlObj.pathname.lastIndexOf('/') + 1);
|
|
||||||
|
|
||||||
manifest = manifest.split('\n').map(line => {
|
// Retry on 5xx errors (transient upstream issues)
|
||||||
const trimmed = line.trim();
|
if (response.status >= 500 && attempt < maxRetries) {
|
||||||
if (trimmed === '' || trimmed.startsWith('#')) {
|
console.log(`[Proxy] Upstream 5xx error (attempt ${attempt}/${maxRetries}), retrying in 500ms...`);
|
||||||
if (trimmed.includes('URI="')) {
|
await new Promise(r => setTimeout(r, 500));
|
||||||
return line.replace(/URI="([^"]+)"/g, (match, p1) => {
|
continue;
|
||||||
try {
|
}
|
||||||
const absoluteUrl = new URL(p1, baseUrl).href;
|
|
||||||
return `URI="${req.protocol}://${req.get('host')}${req.baseUrl}/stream?url=${encodeURIComponent(absoluteUrl)}"`;
|
if (!response.ok) {
|
||||||
} catch (e) { return match; }
|
console.error(`Upstream error for ${url.substring(0, 80)}...: ${response.status} ${response.statusText}`);
|
||||||
});
|
if (response.status === 403) {
|
||||||
}
|
const errorBody = await response.text().catch(() => 'N/A');
|
||||||
return line;
|
console.error(`403 Response body: ${errorBody.substring(0, 200)}`);
|
||||||
|
}
|
||||||
|
return res.status(response.status).send(`Failed to fetch stream: ${response.statusText}`);
|
||||||
|
}
|
||||||
|
|
||||||
|
const contentType = response.headers.get('content-type') || '';
|
||||||
|
res.set('Access-Control-Allow-Origin', '*');
|
||||||
|
|
||||||
|
// Create an async iterator for the response body
|
||||||
|
const iterator = response.body[Symbol.asyncIterator]();
|
||||||
|
const first = await iterator.next();
|
||||||
|
|
||||||
|
if (first.done) {
|
||||||
|
res.set('Content-Type', contentType || 'application/octet-stream');
|
||||||
|
return res.end();
|
||||||
|
}
|
||||||
|
|
||||||
|
const firstChunk = Buffer.from(first.value);
|
||||||
|
|
||||||
|
// Peek at first bytes to check for HLS manifest ({ #EXTM3U })
|
||||||
|
const textPrefix = firstChunk.subarray(0, 7).toString('utf8');
|
||||||
|
const contentLooksLikeHls = textPrefix === '#EXTM3U';
|
||||||
|
|
||||||
|
if (contentLooksLikeHls) {
|
||||||
|
// HLS Manifest: We must read the WHOLE manifest to rewrite it
|
||||||
|
const chunks = [firstChunk];
|
||||||
|
|
||||||
|
// Consume the rest of the stream
|
||||||
|
let result = await iterator.next();
|
||||||
|
while (!result.done) {
|
||||||
|
chunks.push(Buffer.from(result.value));
|
||||||
|
result = await iterator.next();
|
||||||
}
|
}
|
||||||
|
|
||||||
// Stream URL handling
|
const buffer = Buffer.concat(chunks);
|
||||||
try {
|
const finalUrl = response.url || url;
|
||||||
let absoluteUrl;
|
console.log(`[Proxy] Processing HLS manifest from: ${finalUrl.substring(0, 80)}...`);
|
||||||
if (trimmed.startsWith('http://') || trimmed.startsWith('https://')) {
|
res.set('Content-Type', 'application/vnd.apple.mpegurl');
|
||||||
absoluteUrl = trimmed;
|
|
||||||
} else {
|
let manifest = buffer.toString('utf-8');
|
||||||
absoluteUrl = new URL(trimmed, baseUrl).href;
|
|
||||||
|
const finalUrlObj = new URL(finalUrl);
|
||||||
|
const baseUrl = finalUrlObj.origin + finalUrlObj.pathname.substring(0, finalUrlObj.pathname.lastIndexOf('/') + 1);
|
||||||
|
|
||||||
|
manifest = manifest.split('\n').map(line => {
|
||||||
|
const trimmed = line.trim();
|
||||||
|
if (trimmed === '' || trimmed.startsWith('#')) {
|
||||||
|
if (trimmed.includes('URI="')) {
|
||||||
|
return line.replace(/URI="([^"]+)"/g, (match, p1) => {
|
||||||
|
try {
|
||||||
|
const absoluteUrl = new URL(p1, baseUrl).href;
|
||||||
|
return `URI="${req.protocol}://${req.get('host')}${req.baseUrl}/stream?url=${encodeURIComponent(absoluteUrl)}"`;
|
||||||
|
} catch (e) { return match; }
|
||||||
|
});
|
||||||
|
}
|
||||||
|
return line;
|
||||||
}
|
}
|
||||||
return `${req.protocol}://${req.get('host')}${req.baseUrl}/stream?url=${encodeURIComponent(absoluteUrl)}`;
|
|
||||||
} catch (e) { return line; }
|
|
||||||
}).join('\n');
|
|
||||||
|
|
||||||
return res.send(manifest);
|
// Stream URL handling
|
||||||
|
try {
|
||||||
|
let absoluteUrl;
|
||||||
|
if (trimmed.startsWith('http://') || trimmed.startsWith('https://')) {
|
||||||
|
absoluteUrl = trimmed;
|
||||||
|
} else {
|
||||||
|
absoluteUrl = new URL(trimmed, baseUrl).href;
|
||||||
|
}
|
||||||
|
return `${req.protocol}://${req.get('host')}${req.baseUrl}/stream?url=${encodeURIComponent(absoluteUrl)}`;
|
||||||
|
} catch (e) { return line; }
|
||||||
|
}).join('\n');
|
||||||
|
|
||||||
|
return res.send(manifest);
|
||||||
|
}
|
||||||
|
|
||||||
|
// Binary content (Video Segment): Efficient Pipe
|
||||||
|
console.log(`[Proxy] Piping binary stream (${contentType})`);
|
||||||
|
res.set('Content-Type', contentType || 'application/octet-stream');
|
||||||
|
|
||||||
|
// Write the chunk we peeked
|
||||||
|
res.write(firstChunk);
|
||||||
|
|
||||||
|
// Stream the rest
|
||||||
|
// Create a readable stream from the iterator
|
||||||
|
const restOfStream = Readable.from(iterator);
|
||||||
|
|
||||||
|
// Pipe to response
|
||||||
|
restOfStream.pipe(res);
|
||||||
|
return; // Success - exit the retry loop
|
||||||
|
|
||||||
|
} catch (err) {
|
||||||
|
lastError = err;
|
||||||
|
console.error(`Stream proxy error (attempt ${attempt}/${maxRetries}):`, err.message);
|
||||||
|
if (attempt < maxRetries) {
|
||||||
|
console.log('[Proxy] Retrying after error...');
|
||||||
|
await new Promise(r => setTimeout(r, 500));
|
||||||
|
continue;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// Binary content (Video Segment): Efficient Pipe
|
// All retries failed
|
||||||
console.log(`[Proxy] Piping binary stream (${contentType})`);
|
if (!res.headersSent) {
|
||||||
res.set('Content-Type', contentType || 'application/octet-stream');
|
res.status(500).json({ error: lastError?.message || 'Stream proxy failed after retries' });
|
||||||
|
|
||||||
// Write the chunk we peeked
|
|
||||||
res.write(firstChunk);
|
|
||||||
|
|
||||||
// Stream the rest
|
|
||||||
// Create a readable stream from the iterator
|
|
||||||
const restOfStream = Readable.from(iterator);
|
|
||||||
|
|
||||||
// Pipe to response
|
|
||||||
restOfStream.pipe(res);
|
|
||||||
|
|
||||||
} catch (err) {
|
|
||||||
console.error('Stream proxy error:', err);
|
|
||||||
if (!res.headersSent) {
|
|
||||||
res.status(500).json({ error: err.message });
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user