Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions .env.example
Original file line number Diff line number Diff line change
Expand Up @@ -1968,6 +1968,12 @@ APP_LOG_TO_FILE=true
# Default: false (the dashboard/database proxy-log records retain full details).
# PROXY_LOG_INCLUDE_IPS=false

# Record the time to the first useful body byte of each upstream send on its
# proxy-log row (first_chunk_ms). Opt-in: when on, every upstream body is piped
# through a lazy passthrough stream and the row is patched once the byte lands.
# Default: false (only the free headers timing, headers_ms, is recorded).
# PROXY_LOG_FIRST_CHUNK_TIMING=false

# ═══════════════════════════════════════════════════════════════════════════════
# 17. MEMORY OPTIMIZATION (Low-RAM / Docker)
# ═══════════════════════════════════════════════════════════════════════════════
Expand Down
1 change: 1 addition & 0 deletions changelog.d/features/14892-attempt-timing-columns.md
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
- **feat(proxy-logs):** record per-attempt upstream headers duration on proxy log rows, plus opt-in first-chunk duration (`PROXY_LOG_FIRST_CHUNK_TIMING`), and show them in the proxy log detail ([#14892](https://github.com/diegosouzapw/OmniRoute/pull/14892)) — thanks @maxmad64bis
5 changes: 0 additions & 5 deletions config/quality/eslint-suppressions.json
Original file line number Diff line number Diff line change
Expand Up @@ -1782,11 +1782,6 @@
"count": 3
}
},
"src/shared/components/ProxyLogDetail.tsx": {
"@typescript-eslint/no-unused-vars": {
"count": 1
}
},
"src/shared/components/RequestLoggerV2.tsx": {
"@typescript-eslint/no-unused-vars": {
"count": 3
Expand Down
1 change: 1 addition & 0 deletions docs/reference/ENVIRONMENT.md
Original file line number Diff line number Diff line change
Expand Up @@ -919,6 +919,7 @@ The logging system writes to both stdout and rotated log files. All configuratio
| `CALL_LOG_PIPELINE_MAX_SIZE_KB` | `512` | Max pipeline call log artifact size in KB when `call_log_pipeline_enabled=true`. |
| `PROXY_LOGS_TABLE_MAX_ROWS` | `100000` | Max rows in the `proxy_logs` SQLite table before pruning. |
| `PROXY_LOG_INCLUDE_IPS` | `false` | Include client/egress IPs and account prefixes in `[ProxyEgress]` console logs. The dashboard/database proxy-log records retain full details. |
| `PROXY_LOG_FIRST_CHUNK_TIMING` | `false` | Set to `"true"` or `"1"` to record `first_chunk_ms` (send start to first useful body byte) on proxy-log rows. Wraps each upstream body in a lazy passthrough stream and patches the row once the byte arrives; off by default so responses pass through untouched. |
| `APP_LOG_ROTATION_CHECK_INTERVAL_MS` | `60000` (1 min) | How often `src/lib/logRotation.ts` re-checks the active log file size. |
| `CHAT_LOG_TEXT_LIMIT` | `65536` | Max string length retained in chat log artifacts (default 64 KB). |
| `CHAT_LOG_ARRAY_TAIL_ITEMS` | `128` | Number of array items retained from the tail when truncating chat log payloads. |
Expand Down
177 changes: 154 additions & 23 deletions open-sse/utils/upstreamStatusCapture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -9,9 +9,23 @@ export type AttemptRecord = {
error?: string | null;
startedAt?: number;
durationMs?: number | null;
/** Send start -> response headers received; null when unknown (network throw). */
headersMs?: number | null;
/** Send start -> first useful body byte; null until the byte arrives. */
firstChunkMs?: number | null;
/**
* Fired once when the first useful body byte is read downstream, or with
* null when the body settles without one. Set by the capture wrapper when
* it envelopes the raw upstream body; read by the journal layer.
*/
onFirstChunk?: ((firstChunkMs: number | null) => void) | null;
/** True once the envelope replaced the raw body (single-wrap guard). */
bodyTracked?: boolean;
};

export type AttemptSink = {
proxy?: unknown;
rotationAccount?: string | null;
upstreamStatus?: number;
attempts?: AttemptRecord[];
};
Expand Down Expand Up @@ -40,6 +54,101 @@ function closeAttemptRecord(
record.durationMs = Date.now() - (record.startedAt ?? Date.now());
}

/**
* Non-negative integer durations only: clock skew or malformed values stay
* null so partially migrated databases never corrupt a row.
*/
export function sanitizeTimingMs(value: unknown): number | null {
return typeof value === "number" && Number.isInteger(value) && value >= 0 ? value : null;
}

/**
* Opt-in switch for the first-byte body envelope and its deferred row patch.
* Off by default: with it off the upstream Response passes through untouched
* (same object, no extra stream stage) and only the free headers timing is
* recorded. Read on every send so operators can flip it without a restart.
*/
export function isFirstChunkTimingEnabled(): boolean {
const raw = process.env.PROXY_LOG_FIRST_CHUNK_TIMING;
return raw === "true" || raw === "1";
}

function stampHeadersMs(record: AttemptRecord): void {
record.headersMs = sanitizeTimingMs(Date.now() - (record.startedAt ?? Date.now()));
}

/**
* Lazy first-byte envelope: wraps the raw upstream body in a passthrough
* stream that stamps the first useful byte when the downstream reader pulls
* it. Never reads by itself - the first downstream pull drives the first
* upstream read, so an unconsumed body stays untouched and timing stays null.
* Empty values pass through but never stamp; a body that settles without a
* useful byte notifies with null exactly once.
*/
function trackFirstChunk(record: AttemptRecord, response: Response): Response {
if (!response.body || record.bodyTracked) return response;
record.bodyTracked = true;
const startedAt = record.startedAt ?? Date.now();
let settled = false;
const notify = (firstChunkMs: number | null): void => {
if (settled) return;
settled = true;
record.firstChunkMs = sanitizeTimingMs(firstChunkMs);
try {
record.onFirstChunk?.(record.firstChunkMs);
} catch {
// Timing listeners are best-effort; never break the stream.
}
};
const stamp = (): void => notify(Date.now() - startedAt);
const upstream = response.body.getReader();
const wrapped = new ReadableStream<Uint8Array>({
async pull(controller) {
let result: ReadableStreamReadResult<Uint8Array>;
try {
result = await upstream.read();
} catch (error) {
notify(null);
controller.error(error);
return;
}
if (result.done) {
notify(null);
controller.close();
return;
}
if (result.value && result.value.length > 0) stamp();
if (result.value) controller.enqueue(result.value);
else controller.enqueue(new Uint8Array(0));
},
async cancel(reason) {
notify(null);
await upstream.cancel(reason).catch(() => {});
},
});
const enveloped = new Response(wrapped, {
status: response.status,
statusText: response.statusText,
headers: response.headers,
});
preserveResponseIdentity(enveloped, response);
return enveloped;
}

/**
* A reconstructed Response loses url/redirected/type (the constructor always
* yields "", false, "default"). Callers read `response.url` for redirect
* handling and logging, so carry the originals over the same way tlsClient.ts
* does for its adapted responses.
*/
function preserveResponseIdentity(target: Response, source: Response): void {
for (const key of ["url", "redirected", "type"] as const) {
const value = source[key];
if (value === undefined) continue;
Object.defineProperty(target, key, { value, configurable: true });
}
}

/**
* Wrap the process-wide fetch so the HTTP status the provider actually returned lands on
* the request's applied-proxy sink. Only calls that settle while a provider request is
Expand Down Expand Up @@ -70,23 +179,7 @@ export function withUpstreamStatusCapture<A extends unknown[]>(
return async (...args: A) => {
const sink = isDispatching() ? getSink() : undefined;
if (!sink) return inner(...args);
// One journal row per request actually sent: snapshot the outlet at entry,
// complete the row when the send settles. A stacked wrapper reuses the
// in-flight row instead of opening a second one. Sends outside a dispatch
// (side calls, background work) leave no trace, as before.
let record: AttemptRecord | undefined;
const inFlight = readFlight(sink);
if (inFlight) {
record = inFlight;
} else {
record = {
proxy: (sink as { proxy?: unknown }).proxy ?? null,
rotationAccount: (sink as { rotationAccount?: string | null }).rotationAccount ?? null,
startedAt: Date.now(),
};
(sink.attempts ??= []).push(record);
writeFlight(sink, record);
}
const record = openAttemptRecord(sink);
const active = record === readFlight(sink);
let response: Response;
try {
Expand All @@ -99,11 +192,49 @@ export function withUpstreamStatusCapture<A extends unknown[]>(
if (isDispatching()) sink.upstreamStatus = undefined;
throw error;
}
if (isDispatching() && active) {
closeAttemptRecord(record, { upstreamStatus: response.status, error: null });
writeFlight(sink, undefined);
}
if (isDispatching()) sink.upstreamStatus = response.status;
return response;
return settleAttemptResponse(sink, record, active, response, isDispatching);
};
}

/**
* Open or reuse the in-flight attempt record for one request actually sent:
* snapshot the outlet at entry, complete the row when the send settles. A
* stacked wrapper reuses the in-flight row instead of opening a second one.
* Sends outside a dispatch (side calls, background work) leave no trace.
*/
function openAttemptRecord(sink: AttemptSink): AttemptRecord {
const inFlight = readFlight(sink);
if (inFlight) return inFlight;
const record: AttemptRecord = {
proxy: (sink as { proxy?: unknown }).proxy ?? null,
rotationAccount: (sink as { rotationAccount?: string | null }).rotationAccount ?? null,
startedAt: Date.now(),
};
(sink.attempts ??= []).push(record);
writeFlight(sink, record);
return record;
}

/**
* Complete the attempt record once the send settles: stamp the headers timing
* on the upstream response, clear the flight marker, publish the status, and
* envelope the raw body for the lazy first-byte stamp.
*/
function settleAttemptResponse(
sink: AttemptSink,
record: AttemptRecord,
active: boolean,
response: Response,
isDispatching: () => boolean
): Response {
if (isDispatching() && active) {
closeAttemptRecord(record, { upstreamStatus: response.status, error: null });
stampHeadersMs(record);
writeFlight(sink, undefined);
}
if (isDispatching()) sink.upstreamStatus = response.status;
if (isDispatching() && active && isFirstChunkTimingEnabled()) {
return trackFirstChunk(record, response);
}
return response;
}
2 changes: 2 additions & 0 deletions src/i18n/messages/am.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "ከላይ የሚገኝ ሁኔታ",
"noResponse": "አንድ እንደማይመልስ",
"egressIp": "ወጣት አይፒ",
"headersAfter": "አርከቶች በኋላ",
"firstChunkAfter": "የመጀመሪያ ቁርጽ በኋ�",
"servedBy": "የተሰጠ በ",
"correlationId": "የተያያዘ መለያ"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/ar.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "حالة التيار العلوي",
"noResponse": "لا توجد استجابة",
"egressIp": "عنوان IP الصادر",
"headersAfter": "الرؤوس بعد",
"firstChunkAfter": "الكتلة الأولى بعد",
"servedBy": "مقدم بواسطة",
"correlationId": "معرف التوافق"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/az.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Yuxarı axın statusu",
"noResponse": "Cavab yoxdur",
"egressIp": "Çıxış IP",
"headersAfter": "Başlıqlar sonra",
"firstChunkAfter": "İlk hissə sonra",
"servedBy": "Təmin edən",
"correlationId": "Korrelyasiya ID"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/bg.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Статус на входящия",
"noResponse": "Няма отговор",
"egressIp": "Изходящ IP",
"headersAfter": "Заглавки след",
"firstChunkAfter": "Първи блок след",
"servedBy": "Обслужвано от",
"correlationId": "Идентификатор на корелацията"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/bn.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "আপস্ট্রিম স্থিতি",
"noResponse": "কোনো প্রতিক্রিয়া নেই",
"egressIp": "এগ্রেস আইপি",
"headersAfter": "হেডারস পরে",
"firstChunkAfter": "প্রথম চাঙ্ক পরে",
"servedBy": "দ্বারা পরিবেশন করা হয়েছে",
"correlationId": "সংশ্লেষণ আইডি"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/bs.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Status upravljačkog",
"noResponse": "Nema odgovora",
"egressIp": "Egress IP",
"headersAfter": "Zaglavlja nakon",
"firstChunkAfter": "Prvi dio nakon",
"servedBy": "Posluženo od",
"correlationId": "ID korelacije"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/cs.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Stav nahoře",
"noResponse": "Žádná odpověď",
"egressIp": "Egress IP",
"headersAfter": "Záhlaví po",
"firstChunkAfter": "První část po",
"servedBy": "Poskytováno",
"correlationId": "ID korelace"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/da.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Opstrøm status",
"noResponse": "Ingen respons",
"egressIp": "Udgående IP",
"headersAfter": "Headers efter",
"firstChunkAfter": "Første chunk efter",
"servedBy": "Betjent af",
"correlationId": "Korrelation ID"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/de.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Upstream-Status",
"noResponse": "Keine Antwort",
"egressIp": "Egress-IP",
"headersAfter": "Header nach",
"firstChunkAfter": "Erster Chunk danach",
"servedBy": "Bedient durch",
"correlationId": "Korrelation-ID"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/el.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Κατάσταση ανάντη",
"noResponse": "Καμία απάντηση",
"egressIp": "Διεύθυνση IP εξόδου",
"headersAfter": "Κεφαλίδες μετά",
"firstChunkAfter": "Πρώτο τμήμα μετά",
"servedBy": "Παρέχεται από",
"correlationId": "Αναγνωριστικό συσχέτισης"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/en.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Upstream status",
"noResponse": "No response",
"egressIp": "Egress IP",
"headersAfter": "Headers after",
"firstChunkAfter": "First chunk after",
"servedBy": "Served by",
"correlationId": "Correlation ID"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/es.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Estado de río arriba",
"noResponse": "Sin respuesta",
"egressIp": "IP de salida",
"headersAfter": "Encabezados después",
"firstChunkAfter": "Primer fragmento después",
"servedBy": "Servido por",
"correlationId": "ID de correlación"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/et.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Ülesvoolu olek",
"noResponse": "Ei vastata",
"egressIp": "Väljaminev IP",
"headersAfter": "Päised pärast",
"firstChunkAfter": "Esimene tükk pärast",
"servedBy": "Teenindab",
"correlationId": "Korreleerimise ID"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/fa.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "وضعیت بالا دستی",
"noResponse": "هیچ پاسخی",
"egressIp": "IP خروجی",
"headersAfter": "سرفصل‌ها پس از",
"firstChunkAfter": "اولین بخش پس از",
"servedBy": "توسط ارائه شده",
"correlationId": "شناسه همبستگی"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/fi.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Ylävirran tila",
"noResponse": "Ei vastausta",
"egressIp": "Lähtö-IP",
"headersAfter": "Otsakkeet jälkeen",
"firstChunkAfter": "Ensimmäinen lohko jälkeen",
"servedBy": "Palvelija",
"correlationId": "Korrelaatio-ID"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/fr.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "État amont",
"noResponse": "Aucune réponse",
"egressIp": "IP de sortie",
"headersAfter": "En-têtes après",
"firstChunkAfter": "Premier bloc après",
"servedBy": "Servi par",
"correlationId": "ID de corrélation"
},
Expand Down
2 changes: 2 additions & 0 deletions src/i18n/messages/ga.json
Original file line number Diff line number Diff line change
Expand Up @@ -14198,6 +14198,8 @@
"upstreamStatus": "Stádas uachtarach",
"noResponse": "Níl freagra",
"egressIp": "IP Egress",
"headersAfter": "Ceanntásca i ndiaidh",
"firstChunkAfter": "An chéad slabhra i ndiaidh",
"servedBy": "Freastalaithe ag",
"correlationId": "ID Comhoiriúnachta"
},
Expand Down
Loading
Loading