mirror of
https://github.com/Routstr/routstrd.git
synced 2026-08-09 19:56:58 +00:00
provider manage instnace is now persisted across calls.
This commit is contained in:
@@ -3,6 +3,7 @@ import { type IncomingMessage, type ServerResponse } from "http";
|
||||
import {
|
||||
routeRequestsToNodeResponse,
|
||||
InsufficientBalanceError,
|
||||
ProviderManager,
|
||||
} from "@routstr/sdk";
|
||||
import type { UsageTrackingDriver } from "@routstr/sdk";
|
||||
import { logger } from "../../utils/logger";
|
||||
@@ -37,6 +38,7 @@ type DaemonDeps = {
|
||||
ensureProvidersBootstrapped: () => Promise<void>;
|
||||
getRoutstr21Models: (forceRefresh?: boolean) => Promise<any[]>;
|
||||
mode?: ClientMode;
|
||||
providerManager: ProviderManager;
|
||||
};
|
||||
|
||||
/**
|
||||
@@ -277,6 +279,7 @@ export function createDaemonRequestHandler(deps: {
|
||||
getRoutstr21Models: (forceRefresh?: boolean) => Promise<any[]>;
|
||||
mode?: "xcashu" | "apikeys";
|
||||
usageTrackingDriver: UsageTrackingDriver;
|
||||
providerManager: ProviderManager;
|
||||
}) {
|
||||
return async function handler(req: IncomingMessage, res: ServerResponse) {
|
||||
const host = req.headers.host || "localhost";
|
||||
@@ -1193,6 +1196,7 @@ export function createDaemonRequestHandler(deps: {
|
||||
mode: deps.mode,
|
||||
usageTrackingDriver: deps.usageTrackingDriver,
|
||||
sdkStore: deps.store,
|
||||
providerManager: deps.providerManager,
|
||||
res,
|
||||
});
|
||||
return;
|
||||
|
||||
@@ -2,6 +2,7 @@ import { createServer } from "http";
|
||||
import { existsSync } from "fs";
|
||||
import {
|
||||
ModelManager,
|
||||
ProviderManager,
|
||||
createDiscoveryAdapterFromStore,
|
||||
createProviderRegistryFromStore,
|
||||
createStorageAdapterFromStore,
|
||||
@@ -46,6 +47,8 @@ async function main(): Promise<void> {
|
||||
const providerRegistry = createProviderRegistryFromStore(store);
|
||||
const storageAdapter = createStorageAdapterFromStore(store);
|
||||
const modelManager = new ModelManager(discoveryAdapter);
|
||||
// Create shared ProviderManager for consistent failure tracking across all requests
|
||||
const providerManager = new ProviderManager(providerRegistry, store);
|
||||
const { ensureProvidersBootstrapped, getRoutstr21Models } =
|
||||
createModelService(modelManager);
|
||||
|
||||
@@ -72,6 +75,7 @@ async function main(): Promise<void> {
|
||||
getRoutstr21Models,
|
||||
mode: config.mode || "apikeys",
|
||||
usageTrackingDriver,
|
||||
providerManager,
|
||||
}),
|
||||
);
|
||||
|
||||
|
||||
@@ -1,98 +0,0 @@
|
||||
import { Transform } from "stream";
|
||||
import type { UsageData } from "./types";
|
||||
|
||||
export function createSSEParserTransform(
|
||||
onUsage: (usage: UsageData) => void,
|
||||
onResponseId?: (responseId: string) => void,
|
||||
): Transform {
|
||||
let buffer = "";
|
||||
|
||||
const maybeCaptureUsageFromJson = (jsonText: string): void => {
|
||||
try {
|
||||
const data = JSON.parse(jsonText) as any;
|
||||
const responseId = data.id;
|
||||
if (typeof responseId === "string" && responseId.trim().length > 0) {
|
||||
onResponseId?.(responseId.trim());
|
||||
}
|
||||
|
||||
if (data.usage) {
|
||||
const usageCost = data.usage.cost;
|
||||
const cost =
|
||||
typeof usageCost === "number"
|
||||
? usageCost
|
||||
: usageCost?.total_usd ??
|
||||
data.metadata?.routstr?.cost?.total_usd ??
|
||||
0;
|
||||
const msats =
|
||||
data.metadata?.routstr?.cost?.total_msats ??
|
||||
(typeof data.usage.cost_sats === "number"
|
||||
? data.usage.cost_sats * 1000
|
||||
: 0);
|
||||
onUsage({
|
||||
promptTokens: data.usage.prompt_tokens ?? 0,
|
||||
completionTokens: data.usage.completion_tokens ?? 0,
|
||||
totalTokens: data.usage.total_tokens ?? 0,
|
||||
cost,
|
||||
satsCost: msats / 1000,
|
||||
});
|
||||
}
|
||||
} catch {
|
||||
// Ignore non-JSON lines/events.
|
||||
}
|
||||
};
|
||||
|
||||
const processLine = (self: Transform, line: string): void => {
|
||||
const trimmed = line.trim();
|
||||
if (!trimmed) {
|
||||
return;
|
||||
}
|
||||
|
||||
if (trimmed === "data: [DONE]" || trimmed === "[DONE]") {
|
||||
self.push("data: [DONE]\n\n");
|
||||
return;
|
||||
}
|
||||
|
||||
if (trimmed.startsWith("data:")) {
|
||||
const dataStr = trimmed.startsWith("data: ")
|
||||
? trimmed.slice(6)
|
||||
: trimmed.slice(5).trimStart();
|
||||
if (dataStr === "[DONE]") {
|
||||
self.push("data: [DONE]\n\n");
|
||||
return;
|
||||
}
|
||||
maybeCaptureUsageFromJson(dataStr);
|
||||
self.push(`data: ${dataStr}\n\n`);
|
||||
return;
|
||||
}
|
||||
|
||||
if (trimmed.startsWith("{")) {
|
||||
maybeCaptureUsageFromJson(trimmed);
|
||||
self.push(`data: ${trimmed}\n\n`);
|
||||
return;
|
||||
}
|
||||
|
||||
self.push(line + "\n");
|
||||
};
|
||||
|
||||
return new Transform({
|
||||
transform(chunk, encoding, callback) {
|
||||
buffer += chunk.toString();
|
||||
|
||||
const lines = buffer.split(/\r?\n/);
|
||||
buffer = lines.pop() || "";
|
||||
|
||||
for (const line of lines) {
|
||||
processLine(this, line);
|
||||
}
|
||||
|
||||
callback();
|
||||
},
|
||||
flush(callback) {
|
||||
if (buffer.trim()) {
|
||||
processLine(this, buffer);
|
||||
}
|
||||
buffer = "";
|
||||
callback();
|
||||
},
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user