adding data collection
This commit is contained in:
@@ -0,0 +1,74 @@
|
||||
import pg from 'pg';
|
||||
|
||||
const { Pool } = pg;
|
||||
|
||||
export function createPool(connectionString = process.env.DATABASE_URL) {
|
||||
if (!connectionString) throw new Error('DATABASE_URL is required');
|
||||
return new Pool({ connectionString });
|
||||
}
|
||||
|
||||
export async function ensureSchema(pool) {
|
||||
await pool.query(`
|
||||
CREATE TABLE IF NOT EXISTS requests (
|
||||
id BIGSERIAL PRIMARY KEY,
|
||||
created_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
|
||||
method TEXT NOT NULL,
|
||||
endpoint TEXT NOT NULL,
|
||||
request_headers JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||
request_body JSONB,
|
||||
response_headers JSONB NOT NULL DEFAULT '{}'::jsonb,
|
||||
response_body TEXT,
|
||||
response_status INTEGER NOT NULL,
|
||||
duration_ms DOUBLE PRECISION,
|
||||
model TEXT,
|
||||
client_ip TEXT,
|
||||
stream BOOLEAN NOT NULL DEFAULT FALSE,
|
||||
training_data JSONB
|
||||
)
|
||||
`);
|
||||
|
||||
// Keep databases created by an earlier version compatible.
|
||||
await pool.query(`
|
||||
ALTER TABLE requests ADD COLUMN IF NOT EXISTS method TEXT;
|
||||
ALTER TABLE requests ADD COLUMN IF NOT EXISTS request_headers JSONB NOT NULL DEFAULT '{}'::jsonb;
|
||||
ALTER TABLE requests ADD COLUMN IF NOT EXISTS response_headers JSONB NOT NULL DEFAULT '{}'::jsonb;
|
||||
ALTER TABLE requests ADD COLUMN IF NOT EXISTS duration_ms DOUBLE PRECISION;
|
||||
ALTER TABLE requests ADD COLUMN IF NOT EXISTS training_data JSONB;
|
||||
`);
|
||||
await pool.query('CREATE INDEX IF NOT EXISTS requests_created_at_idx ON requests (created_at DESC, id DESC)');
|
||||
}
|
||||
|
||||
export async function insertRequest(pool, entry) {
|
||||
await pool.query(
|
||||
`INSERT INTO requests (
|
||||
method, endpoint, request_headers, request_body, response_headers,
|
||||
response_body, response_status, duration_ms, model, client_ip, stream, training_data
|
||||
) VALUES ($1, $2, $3::jsonb, $4::jsonb, $5::jsonb, $6, $7, $8, $9, $10, $11, $12::jsonb)`,
|
||||
[
|
||||
entry.method,
|
||||
entry.endpoint,
|
||||
JSON.stringify(entry.requestHeaders ?? {}),
|
||||
entry.requestBody == null ? null : JSON.stringify(entry.requestBody),
|
||||
JSON.stringify(entry.responseHeaders ?? {}),
|
||||
entry.responseBody ?? null,
|
||||
entry.responseStatus,
|
||||
entry.durationMs ?? null,
|
||||
entry.model ?? null,
|
||||
entry.clientIp ?? null,
|
||||
Boolean(entry.stream),
|
||||
entry.trainingData == null ? null : JSON.stringify(entry.trainingData),
|
||||
],
|
||||
);
|
||||
}
|
||||
|
||||
export function createRequestLogger(pool, { onError = console.error } = {}) {
|
||||
return {
|
||||
async log(entry) {
|
||||
try {
|
||||
await insertRequest(pool, entry);
|
||||
} catch (error) {
|
||||
onError(`Failed to store request: ${error.message}`);
|
||||
}
|
||||
},
|
||||
};
|
||||
}
|
||||
+11
@@ -0,0 +1,11 @@
|
||||
import path from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import dotenv from 'dotenv';
|
||||
|
||||
export const projectRoot = path.dirname(path.dirname(fileURLToPath(import.meta.url)));
|
||||
|
||||
export function loadEnv(envPath = path.join(projectRoot, '.env')) {
|
||||
return dotenv.config({ path: envPath });
|
||||
}
|
||||
|
||||
loadEnv();
|
||||
@@ -0,0 +1,90 @@
|
||||
import './env.js';
|
||||
import { once } from 'node:events';
|
||||
import { createPool, ensureSchema } from './db.js';
|
||||
import { createTrainingExample } from './training-data.js';
|
||||
|
||||
function usage() {
|
||||
return 'Usage: npm run export -- [limit]\n npm run export -- --limit <limit>';
|
||||
}
|
||||
|
||||
function parseLimit(args) {
|
||||
if (args.length === 0) return null;
|
||||
|
||||
let value;
|
||||
if (args.length === 1 && !args[0].startsWith('-')) value = args[0];
|
||||
else if (args.length === 2 && ['--limit', '-n'].includes(args[0])) value = args[1];
|
||||
else if (args.length === 1 && args[0].startsWith('--limit=')) value = args[0].slice('--limit='.length);
|
||||
else throw new Error(usage());
|
||||
|
||||
const limit = Number(value);
|
||||
if (!Number.isSafeInteger(limit) || limit <= 0) {
|
||||
throw new Error(`Limit must be a positive integer, got "${value}"\n${usage()}`);
|
||||
}
|
||||
return limit;
|
||||
}
|
||||
|
||||
let pool;
|
||||
try {
|
||||
const limit = parseLimit(process.argv.slice(2));
|
||||
pool = createPool();
|
||||
await ensureSchema(pool);
|
||||
let exported = 0;
|
||||
let cursor = null;
|
||||
while (limit === null || exported < limit) {
|
||||
const remaining = limit === null ? 1_000 : limit - exported;
|
||||
const batchSize = Math.min(1_000, Math.max(100, remaining * 2));
|
||||
const params = cursor ? [cursor.cursor_created_at, cursor.id, batchSize] : [batchSize];
|
||||
const cursorClause = cursor ? 'AND (created_at, id) < ($1::timestamptz, $2::bigint)' : '';
|
||||
const batchParameter = cursor ? '$3' : '$1';
|
||||
const result = await pool.query(
|
||||
`SELECT id, created_at, created_at::text AS cursor_created_at, endpoint, request_body, response_body,
|
||||
response_status, duration_ms, stream, training_data
|
||||
FROM requests
|
||||
WHERE (method = 'POST' OR method IS NULL)
|
||||
AND response_status BETWEEN 200 AND 299
|
||||
AND (
|
||||
endpoint LIKE '/oai/v1/chat/completions%'
|
||||
OR endpoint LIKE '/oai/v1/responses%'
|
||||
OR endpoint LIKE '/ant/v1/messages%'
|
||||
OR endpoint LIKE '/ant/v1/v1/messages%'
|
||||
)
|
||||
${cursorClause}
|
||||
ORDER BY created_at DESC, id DESC
|
||||
LIMIT ${batchParameter}`,
|
||||
params,
|
||||
);
|
||||
|
||||
for (const row of result.rows) {
|
||||
const example = row.training_data || createTrainingExample({
|
||||
endpoint: row.endpoint,
|
||||
requestBody: row.request_body,
|
||||
responseBody: row.response_body,
|
||||
responseStatus: row.response_status,
|
||||
durationMs: row.duration_ms,
|
||||
stream: row.stream,
|
||||
});
|
||||
if (!example) continue;
|
||||
const exportedExample = {
|
||||
...example,
|
||||
metadata: {
|
||||
...(example.metadata || {}),
|
||||
request_id: String(row.id),
|
||||
created_at: row.created_at,
|
||||
response_status: row.response_status,
|
||||
...(row.duration_ms == null ? {} : { duration_ms: row.duration_ms }),
|
||||
},
|
||||
};
|
||||
if (!process.stdout.write(`${JSON.stringify(exportedExample)}\n`)) await once(process.stdout, 'drain');
|
||||
exported += 1;
|
||||
if (limit !== null && exported >= limit) break;
|
||||
}
|
||||
if (result.rows.length < batchSize || (limit !== null && exported >= limit)) break;
|
||||
cursor = result.rows.at(-1);
|
||||
}
|
||||
console.error(`Exported ${exported} training request(s) as JSONL, newest first`);
|
||||
} catch (error) {
|
||||
console.error(`Error: ${error.message}`);
|
||||
process.exitCode = 1;
|
||||
} finally {
|
||||
if (pool) await pool.end();
|
||||
}
|
||||
+64
-1
@@ -14,9 +14,71 @@ import {
|
||||
translateOpenAIChunk,
|
||||
} from './anthropic.js';
|
||||
import { translateResponsesRequest, translateResponsesResponse, translateResponsesChunk } from './responses.js';
|
||||
import { createTrainingExample } from './training-data.js';
|
||||
|
||||
export const DEFAULT_UPSTREAM_BASE_URL = 'https://opencode.ai/zen/go/v1';
|
||||
|
||||
export function collectRequestData(requestLogger) {
|
||||
return (request, response, next) => {
|
||||
const startedAt = process.hrtime.bigint();
|
||||
const chunks = [];
|
||||
const originalWrite = response.write.bind(response);
|
||||
const originalEnd = response.end.bind(response);
|
||||
let stored = false;
|
||||
|
||||
const capture = (chunk, encoding) => {
|
||||
if (chunk === undefined || chunk === null || typeof chunk === 'function') return;
|
||||
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk, typeof encoding === 'string' ? encoding : undefined));
|
||||
};
|
||||
|
||||
response.write = (...args) => {
|
||||
capture(args[0], args[1]);
|
||||
return originalWrite(...args);
|
||||
};
|
||||
response.end = (...args) => {
|
||||
capture(args[0], args[1]);
|
||||
return originalEnd(...args);
|
||||
};
|
||||
|
||||
const store = () => {
|
||||
if (stored) return;
|
||||
stored = true;
|
||||
const responseBody = Buffer.concat(chunks).toString('utf8');
|
||||
const entry = {
|
||||
method: request.method,
|
||||
endpoint: request.originalUrl,
|
||||
requestHeaders: request.headers,
|
||||
requestBody: request.body ?? null,
|
||||
responseHeaders: response.getHeaders(),
|
||||
responseBody,
|
||||
responseStatus: response.statusCode,
|
||||
durationMs: Number(process.hrtime.bigint() - startedAt) / 1_000_000,
|
||||
model: request.body?.model ?? null,
|
||||
clientIp: clientIp(request),
|
||||
stream: Boolean(request.body?.stream),
|
||||
};
|
||||
try {
|
||||
entry.trainingData = createTrainingExample(entry);
|
||||
} catch (error) {
|
||||
console.error(`Failed to normalize request for training: ${error.message}`);
|
||||
entry.trainingData = null;
|
||||
}
|
||||
|
||||
try {
|
||||
Promise.resolve(requestLogger.log(entry)).catch((error) => {
|
||||
console.error(`Failed to store request: ${error.message}`);
|
||||
});
|
||||
} catch (error) {
|
||||
console.error(`Failed to store request: ${error.message}`);
|
||||
}
|
||||
};
|
||||
|
||||
response.once('finish', store);
|
||||
response.once('close', store);
|
||||
next();
|
||||
};
|
||||
}
|
||||
|
||||
function copyResponseHeaders(upstream, response) {
|
||||
const contentType = upstream.headers.get('content-type');
|
||||
const cacheControl = upstream.headers.get('cache-control');
|
||||
@@ -60,9 +122,10 @@ function withContextLimit(payload, maxContextTokens) {
|
||||
};
|
||||
}
|
||||
|
||||
export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBaseUrl = DEFAULT_UPSTREAM_BASE_URL, maxContextTokens = maxContextFromEnv(), rateLimiter = createRateLimiter(), pagesDirectory = defaultPagesDirectory } = {}) {
|
||||
export function createProxyApp({ keys, fetchImpl = globalThis.fetch, upstreamBaseUrl = DEFAULT_UPSTREAM_BASE_URL, maxContextTokens = maxContextFromEnv(), rateLimiter = createRateLimiter(), pagesDirectory = defaultPagesDirectory, requestLogger } = {}) {
|
||||
const rotator = createKeyRotator(keys);
|
||||
const app = express();
|
||||
if (requestLogger) app.use(collectRequestData(requestLogger));
|
||||
app.use(express.json({ limit: '10mb' }));
|
||||
|
||||
app.get('/', (_request, response) => {
|
||||
|
||||
+16
-4
@@ -1,13 +1,25 @@
|
||||
import path from 'node:path';
|
||||
import { fileURLToPath } from 'node:url';
|
||||
import { createProxyApp } from './proxy.js';
|
||||
import { projectRoot } from './env.js';
|
||||
import { createPool, createRequestLogger, ensureSchema } from './db.js';
|
||||
import { loadConfiguredKeys } from './key-store.js';
|
||||
|
||||
const projectRoot = path.dirname(path.dirname(fileURLToPath(import.meta.url)));
|
||||
const keys = await loadConfiguredKeys({ path: path.join(projectRoot, 'keys.txt') });
|
||||
const app = createProxyApp({ keys });
|
||||
const pool = createPool();
|
||||
await ensureSchema(pool);
|
||||
const app = createProxyApp({ keys, requestLogger: createRequestLogger(pool) });
|
||||
const port = Number(process.env.PORT || 4005);
|
||||
|
||||
app.listen(port, '0.0.0.0', () => {
|
||||
const server = app.listen(port, '0.0.0.0', () => {
|
||||
console.log(`OpenCode Go proxy listening at http://localhost:${port}/oai/v1`);
|
||||
});
|
||||
|
||||
async function shutdown() {
|
||||
server.close(async () => {
|
||||
await pool.end();
|
||||
process.exit(0);
|
||||
});
|
||||
}
|
||||
|
||||
process.once('SIGINT', shutdown);
|
||||
process.once('SIGTERM', shutdown);
|
||||
|
||||
@@ -0,0 +1,359 @@
|
||||
function parseJson(value) {
|
||||
if (value && typeof value === 'object') return value;
|
||||
if (typeof value !== 'string') return null;
|
||||
try { return JSON.parse(value); } catch { return null; }
|
||||
}
|
||||
|
||||
function parseSse(body = '') {
|
||||
const events = [];
|
||||
for (const block of body.split(/\r?\n\r?\n/)) {
|
||||
const data = block.split(/\r?\n/)
|
||||
.filter((line) => line.startsWith('data:'))
|
||||
.map((line) => line.slice(5).trimStart())
|
||||
.join('\n');
|
||||
if (!data || data === '[DONE]') continue;
|
||||
const parsed = parseJson(data);
|
||||
if (parsed) events.push(parsed);
|
||||
}
|
||||
return events;
|
||||
}
|
||||
|
||||
function normalizedContent(content) {
|
||||
if (content === undefined || content === null || typeof content === 'string') return content ?? null;
|
||||
if (!Array.isArray(content)) return String(content);
|
||||
return content.map((part) => {
|
||||
if (!part || typeof part !== 'object') return { type: 'text', text: String(part ?? '') };
|
||||
if (['input_text', 'output_text'].includes(part.type)) return { type: 'text', text: part.text || '' };
|
||||
if (part.type === 'image' && part.source?.type === 'base64') {
|
||||
return { type: 'image_url', image_url: { url: `data:${part.source.media_type};base64,${part.source.data}` } };
|
||||
}
|
||||
return structuredClone(part);
|
||||
});
|
||||
}
|
||||
|
||||
function toolCall(call = {}) {
|
||||
const fn = call.function || call;
|
||||
const args = fn.arguments ?? call.arguments ?? '{}';
|
||||
return {
|
||||
id: call.call_id || call.id || 'tool_call',
|
||||
type: 'function',
|
||||
function: {
|
||||
name: fn.name || call.name || '',
|
||||
arguments: typeof args === 'string' ? args : JSON.stringify(args),
|
||||
},
|
||||
};
|
||||
}
|
||||
|
||||
function normalizedMessage(message = {}) {
|
||||
const role = message.role === 'developer' ? 'system' : message.role;
|
||||
if (!['system', 'user', 'assistant', 'tool'].includes(role)) return null;
|
||||
const result = { role, content: normalizedContent(message.content) };
|
||||
if (message.name !== undefined) result.name = message.name;
|
||||
const reasoning = message.reasoning_content ?? (typeof message.reasoning === 'string' ? message.reasoning : undefined);
|
||||
if (reasoning) result.reasoning_content = reasoning;
|
||||
if (message.tool_calls?.length) result.tool_calls = message.tool_calls.map(toolCall);
|
||||
if (role === 'tool') result.tool_call_id = message.tool_call_id || message.tool_use_id || 'tool_call';
|
||||
return result;
|
||||
}
|
||||
|
||||
function textFromParts(content = []) {
|
||||
if (typeof content === 'string') return content;
|
||||
if (!Array.isArray(content)) return '';
|
||||
return content
|
||||
.filter((part) => ['text', 'input_text', 'output_text'].includes(part?.type))
|
||||
.map((part) => part.text || '')
|
||||
.join('');
|
||||
}
|
||||
|
||||
function reasoningFromItem(item = {}) {
|
||||
if (typeof item.reasoning_content === 'string') return item.reasoning_content;
|
||||
if (typeof item.text === 'string') return item.text;
|
||||
if (typeof item.content === 'string') return item.content;
|
||||
return (item.summary || []).map((part) => part?.text || '').join('');
|
||||
}
|
||||
|
||||
function assistantFromResponsesOutput(output = []) {
|
||||
const assistant = { role: 'assistant', content: null };
|
||||
let text = '';
|
||||
let reasoning = '';
|
||||
const calls = [];
|
||||
|
||||
for (const item of output || []) {
|
||||
if (item?.type === 'message') text += textFromParts(item.content);
|
||||
else if (item?.type === 'reasoning') reasoning += reasoningFromItem(item);
|
||||
else if (item?.type === 'function_call') calls.push(toolCall(item));
|
||||
}
|
||||
|
||||
if (text) assistant.content = text;
|
||||
if (reasoning) assistant.reasoning_content = reasoning;
|
||||
if (calls.length) assistant.tool_calls = calls;
|
||||
return text || reasoning || calls.length ? assistant : null;
|
||||
}
|
||||
|
||||
function inputFromResponses(body = {}) {
|
||||
const messages = [];
|
||||
if (body.instructions !== undefined) messages.push({ role: 'system', content: normalizedContent(body.instructions) });
|
||||
const input = Array.isArray(body.input) ? body.input : [{ type: 'message', role: 'user', content: body.input ?? '' }];
|
||||
|
||||
const ensureAssistant = () => {
|
||||
if (messages.at(-1)?.role !== 'assistant') messages.push({ role: 'assistant', content: null });
|
||||
return messages.at(-1);
|
||||
};
|
||||
|
||||
for (const item of input) {
|
||||
if (item?.type === 'function_call_output') {
|
||||
messages.push({ role: 'tool', tool_call_id: item.call_id || 'tool_call', content: normalizedContent(item.output ?? '') });
|
||||
} else if (item?.type === 'function_call') {
|
||||
const assistant = ensureAssistant();
|
||||
assistant.tool_calls = assistant.tool_calls || [];
|
||||
assistant.tool_calls.push(toolCall(item));
|
||||
} else if (item?.type === 'reasoning') {
|
||||
const reasoning = reasoningFromItem(item);
|
||||
if (reasoning) ensureAssistant().reasoning_content = (ensureAssistant().reasoning_content || '') + reasoning;
|
||||
} else if (item?.type === 'message' || item?.role) {
|
||||
const message = normalizedMessage({ ...item, role: item.role || 'user' });
|
||||
if (!message) continue;
|
||||
if (message.role === 'assistant' && messages.at(-1)?.role === 'assistant') {
|
||||
const pending = messages.at(-1);
|
||||
pending.content = message.content;
|
||||
if (message.reasoning_content) pending.reasoning_content = (pending.reasoning_content || '') + message.reasoning_content;
|
||||
if (message.tool_calls) pending.tool_calls = [...(pending.tool_calls || []), ...message.tool_calls];
|
||||
} else messages.push(message);
|
||||
}
|
||||
}
|
||||
return messages;
|
||||
}
|
||||
|
||||
function assistantFromChatResponse(responseBody, stream) {
|
||||
if (!stream) {
|
||||
const payload = parseJson(responseBody);
|
||||
return normalizedMessage(payload?.choices?.[0]?.message);
|
||||
}
|
||||
|
||||
let content = '';
|
||||
let reasoning = '';
|
||||
const calls = new Map();
|
||||
for (const payload of parseSse(responseBody)) {
|
||||
const delta = payload.choices?.[0]?.delta;
|
||||
if (!delta) continue;
|
||||
if (typeof delta.content === 'string') content += delta.content;
|
||||
const reasoningDelta = delta.reasoning_content ?? delta.reasoning;
|
||||
if (typeof reasoningDelta === 'string') reasoning += reasoningDelta;
|
||||
for (const call of delta.tool_calls || []) {
|
||||
const index = call.index ?? calls.size;
|
||||
const current = calls.get(index) || { id: '', name: '', arguments: '' };
|
||||
if (call.id) current.id += call.id;
|
||||
if (call.function?.name) current.name += call.function.name;
|
||||
if (call.function?.arguments) current.arguments += call.function.arguments;
|
||||
calls.set(index, current);
|
||||
}
|
||||
}
|
||||
|
||||
if (!content && !reasoning && !calls.size) return null;
|
||||
const assistant = { role: 'assistant', content: content || null };
|
||||
if (reasoning) assistant.reasoning_content = reasoning;
|
||||
if (calls.size) assistant.tool_calls = [...calls.values()].map((call, index) => toolCall({
|
||||
id: call.id || `tool_call_${index}`, function: { name: call.name, arguments: call.arguments || '{}' },
|
||||
}));
|
||||
return assistant;
|
||||
}
|
||||
|
||||
function assistantFromResponsesResponse(responseBody, stream) {
|
||||
if (!stream) return assistantFromResponsesOutput(parseJson(responseBody)?.output);
|
||||
const events = parseSse(responseBody);
|
||||
const completed = [...events].reverse().find((event) => event.type === 'response.completed');
|
||||
if (completed?.response?.output) return assistantFromResponsesOutput(completed.response.output);
|
||||
|
||||
let content = '';
|
||||
let reasoning = '';
|
||||
const calls = new Map();
|
||||
for (const event of events) {
|
||||
if (event.type === 'response.output_text.delta') content += event.delta || '';
|
||||
else if (event.type === 'response.reasoning_summary_text.delta') reasoning += event.delta || '';
|
||||
else if (event.type === 'response.output_item.added' && event.item?.type === 'function_call') {
|
||||
calls.set(event.item.id, { id: event.item.call_id || event.item.id, name: event.item.name || '', arguments: event.item.arguments || '' });
|
||||
} else if (event.type === 'response.function_call_arguments.delta') {
|
||||
const call = calls.get(event.item_id) || { id: event.item_id, name: '', arguments: '' };
|
||||
call.arguments += event.delta || '';
|
||||
calls.set(event.item_id, call);
|
||||
}
|
||||
}
|
||||
const assistant = { role: 'assistant', content: content || null };
|
||||
if (reasoning) assistant.reasoning_content = reasoning;
|
||||
if (calls.size) assistant.tool_calls = [...calls.values()].map(toolCall);
|
||||
return content || reasoning || calls.size ? assistant : null;
|
||||
}
|
||||
|
||||
function inputFromAnthropic(body = {}) {
|
||||
const messages = [];
|
||||
if (body.system !== undefined) messages.push({ role: 'system', content: normalizedContent(body.system) });
|
||||
for (const source of body.messages || []) {
|
||||
if (!Array.isArray(source.content)) {
|
||||
const message = normalizedMessage(source);
|
||||
if (message) messages.push(message);
|
||||
continue;
|
||||
}
|
||||
|
||||
const ordinary = source.content.filter((part) => !['thinking', 'tool_use', 'tool_result'].includes(part?.type));
|
||||
const thinking = source.content.filter((part) => part?.type === 'thinking').map((part) => part.thinking || '').join('');
|
||||
const calls = source.content.filter((part) => part?.type === 'tool_use').map((part) => toolCall({ id: part.id, name: part.name, arguments: part.input }));
|
||||
if (ordinary.length || thinking || calls.length) {
|
||||
const message = { role: source.role === 'developer' ? 'system' : source.role, content: ordinary.length ? normalizedContent(ordinary) : null };
|
||||
if (thinking) message.reasoning_content = thinking;
|
||||
if (calls.length) message.tool_calls = calls;
|
||||
messages.push(message);
|
||||
}
|
||||
for (const result of source.content.filter((part) => part?.type === 'tool_result')) {
|
||||
messages.push({ role: 'tool', tool_call_id: result.tool_use_id || 'tool_call', content: normalizedContent(result.content ?? '') });
|
||||
}
|
||||
}
|
||||
return messages;
|
||||
}
|
||||
|
||||
function assistantFromAnthropicResponse(responseBody, stream) {
|
||||
if (!stream) {
|
||||
const payload = parseJson(responseBody);
|
||||
const assistant = { role: 'assistant', content: null };
|
||||
let text = '';
|
||||
let reasoning = '';
|
||||
const calls = [];
|
||||
for (const block of payload?.content || []) {
|
||||
if (block.type === 'text') text += block.text || '';
|
||||
else if (block.type === 'thinking') reasoning += block.thinking || '';
|
||||
else if (block.type === 'tool_use') calls.push(toolCall({ id: block.id, name: block.name, arguments: block.input }));
|
||||
}
|
||||
if (text) assistant.content = text;
|
||||
if (reasoning) assistant.reasoning_content = reasoning;
|
||||
if (calls.length) assistant.tool_calls = calls;
|
||||
return text || reasoning || calls.length ? assistant : null;
|
||||
}
|
||||
|
||||
const blocks = new Map();
|
||||
for (const event of parseSse(responseBody)) {
|
||||
if (event.type === 'content_block_start') {
|
||||
const block = event.content_block || {};
|
||||
blocks.set(event.index, { type: block.type, text: block.text || '', thinking: block.thinking || '', id: block.id, name: block.name, arguments: '' });
|
||||
} else if (event.type === 'content_block_delta') {
|
||||
const block = blocks.get(event.index) || {};
|
||||
if (event.delta?.type === 'text_delta') block.text = (block.text || '') + (event.delta.text || '');
|
||||
else if (event.delta?.type === 'thinking_delta') block.thinking = (block.thinking || '') + (event.delta.thinking || '');
|
||||
else if (event.delta?.type === 'input_json_delta') block.arguments = (block.arguments || '') + (event.delta.partial_json || '');
|
||||
blocks.set(event.index, block);
|
||||
}
|
||||
}
|
||||
return assistantFromAnthropicResponse(JSON.stringify({ content: [...blocks.values()].map((block) => {
|
||||
if (block.type === 'tool_use') return { type: block.type, id: block.id, name: block.name, input: parseJson(block.arguments) || block.arguments || {} };
|
||||
return block;
|
||||
}) }), false);
|
||||
}
|
||||
|
||||
function definedEntries(object) {
|
||||
return Object.fromEntries(Object.entries(object).filter(([, value]) => value !== undefined && value !== null));
|
||||
}
|
||||
|
||||
function responseMetadata(pathname, responseBody, stream) {
|
||||
if (pathname === '/oai/v1/chat/completions') {
|
||||
const payloads = stream ? parseSse(responseBody) : [parseJson(responseBody)].filter(Boolean);
|
||||
const idPayload = payloads.find((payload) => payload.id) || {};
|
||||
const modelPayload = payloads.find((payload) => payload.model) || {};
|
||||
const finishPayload = [...payloads].reverse().find((payload) => payload.choices?.[0]?.finish_reason) || {};
|
||||
const usagePayload = [...payloads].reverse().find((payload) => payload.usage) || {};
|
||||
return definedEntries({
|
||||
id: idPayload.id,
|
||||
model: modelPayload.model,
|
||||
finish_reason: finishPayload.choices?.[0]?.finish_reason,
|
||||
usage: usagePayload.usage,
|
||||
cost: usagePayload.cost ?? usagePayload.usage?.cost,
|
||||
});
|
||||
}
|
||||
|
||||
if (pathname === '/oai/v1/responses') {
|
||||
const payload = stream
|
||||
? [...parseSse(responseBody)].reverse().find((event) => event.type === 'response.completed')?.response
|
||||
: parseJson(responseBody);
|
||||
return definedEntries({
|
||||
id: payload?.id,
|
||||
model: payload?.model,
|
||||
status: payload?.status,
|
||||
incomplete_details: payload?.incomplete_details,
|
||||
usage: payload?.usage,
|
||||
});
|
||||
}
|
||||
|
||||
const payloads = stream ? parseSse(responseBody) : [parseJson(responseBody)].filter(Boolean);
|
||||
if (!stream) {
|
||||
const payload = payloads[0] || {};
|
||||
return definedEntries({
|
||||
id: payload.id,
|
||||
model: payload.model,
|
||||
finish_reason: payload.stop_reason,
|
||||
stop_sequence: payload.stop_sequence,
|
||||
usage: payload.usage,
|
||||
});
|
||||
}
|
||||
const start = payloads.find((payload) => payload.type === 'message_start')?.message || {};
|
||||
const delta = [...payloads].reverse().find((payload) => payload.type === 'message_delta') || {};
|
||||
const usage = definedEntries({ ...(start.usage || {}), ...(delta.usage || {}) });
|
||||
return definedEntries({
|
||||
id: start.id,
|
||||
model: start.model,
|
||||
finish_reason: delta.delta?.stop_reason,
|
||||
stop_sequence: delta.delta?.stop_sequence,
|
||||
usage: Object.keys(usage).length ? usage : undefined,
|
||||
});
|
||||
}
|
||||
|
||||
function trainingMetadata(pathname, requestBody, responseBody, stream, responseStatus, durationMs) {
|
||||
const api = pathname === '/oai/v1/chat/completions'
|
||||
? 'openai_chat_completions'
|
||||
: pathname === '/oai/v1/responses'
|
||||
? 'openai_responses'
|
||||
: 'anthropic_messages';
|
||||
const parameterNames = [
|
||||
'temperature', 'top_p', 'top_k', 'min_p', 'max_tokens', 'max_output_tokens',
|
||||
'stop', 'stop_sequences', 'frequency_penalty', 'presence_penalty', 'n', 'seed',
|
||||
'logprobs', 'top_logprobs', 'response_format', 'parallel_tool_calls', 'reasoning',
|
||||
'thinking', 'verbosity', 'text', 'include', 'truncation', 'modalities', 'audio',
|
||||
'service_tier',
|
||||
];
|
||||
const parameters = definedEntries(Object.fromEntries(parameterNames.map((name) => [name, requestBody[name]])));
|
||||
const metadata = {
|
||||
schema_version: 1,
|
||||
api,
|
||||
endpoint: pathname,
|
||||
model: requestBody.model ?? null,
|
||||
stream: Boolean(stream),
|
||||
tools: structuredClone(requestBody.tools || []),
|
||||
response_status: responseStatus,
|
||||
response: responseMetadata(pathname, responseBody, stream),
|
||||
};
|
||||
if (Number.isFinite(durationMs)) metadata.duration_ms = durationMs;
|
||||
if (requestBody.tool_choice !== undefined) metadata.tool_choice = structuredClone(requestBody.tool_choice);
|
||||
if (requestBody.metadata !== undefined) metadata.request_metadata = structuredClone(requestBody.metadata);
|
||||
if (Object.keys(parameters).length) metadata.parameters = parameters;
|
||||
return metadata;
|
||||
}
|
||||
|
||||
export function createTrainingExample({ endpoint = '', requestBody, responseBody, responseStatus, durationMs, stream = false } = {}) {
|
||||
if (responseStatus < 200 || responseStatus >= 300 || !requestBody) return null;
|
||||
const pathname = endpoint.split('?')[0];
|
||||
let messages;
|
||||
let assistant;
|
||||
|
||||
if (pathname === '/oai/v1/chat/completions') {
|
||||
messages = (requestBody.messages || []).map(normalizedMessage).filter(Boolean);
|
||||
assistant = assistantFromChatResponse(responseBody, stream);
|
||||
} else if (pathname === '/oai/v1/responses') {
|
||||
messages = inputFromResponses(requestBody);
|
||||
assistant = assistantFromResponsesResponse(responseBody, stream);
|
||||
} else if (['/ant/v1/messages', '/ant/v1/v1/messages'].includes(pathname)) {
|
||||
messages = inputFromAnthropic(requestBody);
|
||||
assistant = assistantFromAnthropicResponse(responseBody, stream);
|
||||
} else return null;
|
||||
|
||||
if (!assistant || !messages.length) return null;
|
||||
return {
|
||||
messages: [...messages, assistant],
|
||||
metadata: trainingMetadata(pathname, requestBody, responseBody, stream, responseStatus, durationMs),
|
||||
};
|
||||
}
|
||||
Reference in New Issue
Block a user