chore: add runtime observability ops
This commit is contained in:
+39
-8
@@ -143,6 +143,13 @@ function createRuntimeRouter({
|
||||
readNumber(values[4]) * 2
|
||||
);
|
||||
};
|
||||
const workerScoreFromValues = (values = []) => (
|
||||
readNumber(values[0]) * 3 +
|
||||
readNumber(values[1]) +
|
||||
readNumber(values[2]) * 0.01 +
|
||||
readNumber(values[3]) * 5 +
|
||||
readNumber(values[4]) * 2
|
||||
);
|
||||
return {
|
||||
async pickTarget(fallbackPick, orderedTargets = targets) {
|
||||
const client = await getClient();
|
||||
@@ -178,25 +185,32 @@ function createRuntimeRouter({
|
||||
if (!client || !target) return;
|
||||
const workerId = workerIdForTarget(target);
|
||||
const streamKey = sessionId ? key('stream', sessionId, 'status') : null;
|
||||
const now = String(Date.now());
|
||||
const multi = client
|
||||
.multi()
|
||||
.incr(key('worker', workerId, 'active_streams'))
|
||||
.set(key('worker', workerId, 'heartbeat'), String(Date.now()), { EX: 30 });
|
||||
.incr(key('worker', workerId, 'stream_open_count'))
|
||||
.set(key('worker', workerId, 'last_stream_started_at'), now)
|
||||
.set(key('worker', workerId, 'heartbeat'), now, { EX: 30 });
|
||||
if (streamKey) multi.set(streamKey, 'active', { EX: 60 * 60 });
|
||||
await multi.exec().catch(() => null);
|
||||
},
|
||||
async streamEnded(sessionId, target) {
|
||||
async streamEnded(sessionId, target, { status = 'closed' } = {}) {
|
||||
const client = await getClient();
|
||||
if (!client || !target) return;
|
||||
const workerId = workerIdForTarget(target);
|
||||
const streamKey = sessionId ? key('stream', sessionId, 'status') : null;
|
||||
const activeKey = key('worker', workerId, 'active_streams');
|
||||
const now = String(Date.now());
|
||||
const nextValue = await client.decr(activeKey).catch(() => null);
|
||||
const multi = client
|
||||
.multi()
|
||||
.set(key('worker', workerId, 'heartbeat'), String(Date.now()), { EX: 30 });
|
||||
.set(key('worker', workerId, 'last_stream_ended_at'), now)
|
||||
.set(key('worker', workerId, 'heartbeat'), now, { EX: 30 });
|
||||
if (Number(nextValue ?? 0) < 0) multi.set(activeKey, '0');
|
||||
if (streamKey) multi.set(streamKey, 'closed', { EX: 600 });
|
||||
if (status === 'aborted') multi.incr(key('worker', workerId, 'stream_abort_count'));
|
||||
if (status === 'error') multi.incr(key('worker', workerId, 'stream_error_count'));
|
||||
if (streamKey) multi.set(streamKey, status, { EX: 600 });
|
||||
await multi.exec().catch(() => null);
|
||||
},
|
||||
async getStatus() {
|
||||
@@ -214,6 +228,11 @@ function createRuntimeRouter({
|
||||
key('worker', workerId, 'memory_pressure'),
|
||||
key('worker', workerId, 'heartbeat'),
|
||||
key('worker', workerId, 'drain'),
|
||||
key('worker', workerId, 'stream_open_count'),
|
||||
key('worker', workerId, 'stream_abort_count'),
|
||||
key('worker', workerId, 'stream_error_count'),
|
||||
key('worker', workerId, 'last_stream_started_at'),
|
||||
key('worker', workerId, 'last_stream_ended_at'),
|
||||
])
|
||||
.catch(() => [])
|
||||
: [];
|
||||
@@ -227,6 +246,12 @@ function createRuntimeRouter({
|
||||
memoryPressure: readNumber(values?.[4]),
|
||||
heartbeat: values?.[5] ? Number(values[5]) : null,
|
||||
drain: /^(1|true|yes)$/i.test(String(values?.[6] ?? '')),
|
||||
streamOpenCount: readNumber(values?.[7]),
|
||||
streamAbortCount: readNumber(values?.[8]),
|
||||
streamErrorCount: readNumber(values?.[9]),
|
||||
lastStreamStartedAt: values?.[10] ? Number(values[10]) : null,
|
||||
lastStreamEndedAt: values?.[11] ? Number(values[11]) : null,
|
||||
score: workerScoreFromValues(values),
|
||||
});
|
||||
}
|
||||
return {
|
||||
@@ -861,9 +886,9 @@ export function createTkmindProxy({
|
||||
await runtimeRouter.streamStarted(sessionId, target).catch(() => null);
|
||||
}
|
||||
|
||||
async function markStreamEnded(sessionId, target) {
|
||||
async function markStreamEnded(sessionId, target, options = {}) {
|
||||
if (!runtimeRouter || !sessionId || !target) return;
|
||||
await runtimeRouter.streamEnded(sessionId, target).catch(() => null);
|
||||
await runtimeRouter.streamEnded(sessionId, target, options).catch(() => null);
|
||||
}
|
||||
|
||||
async function getRuntimeStatus() {
|
||||
@@ -1526,12 +1551,18 @@ export function createTkmindProxy({
|
||||
}, 20000);
|
||||
res.on('drain', () => source.resume());
|
||||
await markStreamStarted(sessionId, sessionTarget);
|
||||
let streamCloseStatus = 'closed';
|
||||
try {
|
||||
await pipeline(source, linkSanitizer, billingTransform, clientSink);
|
||||
try {
|
||||
await pipeline(source, linkSanitizer, billingTransform, clientSink);
|
||||
} catch (err) {
|
||||
streamCloseStatus = clientClosed || upstreamAbort.signal.aborted ? 'aborted' : 'error';
|
||||
throw err;
|
||||
}
|
||||
} finally {
|
||||
clearInterval(keepalive);
|
||||
req.off('close', abortUpstream);
|
||||
await markStreamEnded(sessionId, sessionTarget);
|
||||
await markStreamEnded(sessionId, sessionTarget, { status: streamCloseStatus });
|
||||
if (!res.writableEnded) res.end();
|
||||
}
|
||||
} catch (err) {
|
||||
|
||||
Reference in New Issue
Block a user