From 4ef80138634aaba8adeb279545bca3b109071d9d Mon Sep 17 00:00:00 2001 From: cwiklik Date: Thu, 8 Oct 2026 11:42:18 -0400 Subject: [PATCH 1/2] feat: Populate the BYTES column for ordinary traffic MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit An operator hit a ContextWindowExceededError from LiteLLM and wanted to know how big the request Cortex had just forwarded actually was. They found agentop's BYTES column, switched it on, and it was blank on every row: the column only ever had a figure on an opaque tunnel's close row, which is what its own help text said. To get the number at all they yanked the event to a file and measured the file — which is the event's JSON, not the body that went on the wire. Every request row now carries the size of the body it forwarded and every response row the size of the body that came back, on all three listeners. BYTES goes default-on in the same change: #1200 made it the one opt-in column explicitly because only tunnel rows had a figure, and that reason is gone. Response bytes are summed in a new pipeline.Context.ResponseBytes rather than measured at the record site, because len(pctx.ResponseBody) is a buffer and not a tally — extproc replaces it per chunk on the SSE path and truncates it at its cap, and neither proxy's streaming relay fills it at all. Reading it would have reported a streamed inference response as the size of its final chunk, on exactly the traffic the column was asked for. The proxies count at one choke point each, through a new counting reader (core/listener/internal/bodycount), rather than per arm. Counting per arm put them 32 bytes below extproc on a four-event SSE response: the re-framing arm sees sseframe payloads with the `data: ` prefixes and blank-line separators already stripped. A parity fixture now pins all three listeners to the same figure for the same body. Zero means the listener counted nothing — not that the body was empty. Bodies are buffered only when a plugin asks, and both proxies' fully-unbuffered relays copy after the row is appended, so those rows report zero and agentop renders them blank. Content-Length is not a substitute: it is -1 on anything chunked and deleted outright on every streaming arm, so it would be a confident guess exactly where it was wrong. Fixes #1309 Assisted-By: Claude (Anthropic AI) Signed-off-by: cwiklik --- CLAUDE.md | 3 +- cmd/agentop/README.md | 52 +++--- cmd/agentop/tui/events_column_sizing_test.go | 23 ++- cmd/agentop/tui/events_columns.go | 19 ++- cmd/agentop/tui/events_columns_test.go | 7 +- cmd/agentop/tui/events_pane.go | 29 +++- cmd/agentop/tui/settings.go | 29 ++-- cmd/agentop/tui/settings_test.go | 42 +++-- cmd/agentop/tui/tunnel_close_test.go | 47 +++++- core/listener/extproc/bytescount_test.go | 144 +++++++++++++++++ core/listener/extproc/server.go | 13 ++ core/listener/forwardproxy/bytescount_test.go | 150 ++++++++++++++++++ core/listener/forwardproxy/server.go | 28 +++- core/listener/internal/bodycount/bodycount.go | 44 +++++ core/listener/parity/drivers_test.go | 42 +++++ core/listener/parity/parity_test.go | 131 +++++++++++++++ core/listener/reverseproxy/bytescount_test.go | 114 +++++++++++++ core/listener/reverseproxy/server.go | 19 +++ core/pipeline/context.go | 26 +++ core/pipeline/session.go | 44 ++++- 20 files changed, 923 insertions(+), 83 deletions(-) create mode 100644 core/listener/extproc/bytescount_test.go create mode 100644 core/listener/forwardproxy/bytescount_test.go create mode 100644 core/listener/internal/bodycount/bodycount.go create mode 100644 core/listener/reverseproxy/bytescount_test.go diff --git a/CLAUDE.md b/CLAUDE.md index 32f4d1fc5..23f72cd3b 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -577,7 +577,8 @@ Every event on `/v1/sessions/{id}` and `/v1/events` carries: reports the same situation one layer down as `tunnelReason: "dial-failed"` with `error.kind: "dial_failed"`. - `requestedHost` — present only when a plugin redirected the request (`pctx.Redirect`): the host the client asked for. `host` is where the request actually went, and usage, the cost ledger and pricing all follow `host`. agentop's detail pane shows both on a `redirected:` line. -- `tunnel`, `tunnelReason`, `bytesUp`, `bytesDown` — an opaque CONNECT (or transparent-redirect) tunnel records two rows sharing a `requestId`: the open (`phase: "request"`, `tunnelReason` saying why the bytes stayed opaque) and, when the tunnel ends, the close (`phase: "response"`). The close carries the CONNECT's own `statusCode` — 200, or 502 with the dial error in `error` when the destination could not be reached (`tunnelReason: "dial-failed"`) — plus `durationMs` for how long the tunnel stayed open and the bytes it carried each way (up = client to destination). It is not the destination's status: that travels inside the client's end-to-end TLS. A bridged tunnel's open is recorded with its first decrypted request — in that request's session, directly before it, stamped with its time — and records no close, because the request carries its own response; a bridged tunnel that recorded no request gets its open and a close when it ends. Tunnel rows are kept out of `/v1/usage`: a tunnel's lifetime is not a request latency. +- `bytesUp`, `bytesDown` — **body** bytes, one direction per row: a request row carries `bytesUp`, the body **as forwarded** (so a `tool-prune`-style rewrite is reflected, since the count is taken after the pipeline ran), and a response row carries `bytesDown`, the body **as it came back from the destination** — not as delivered, because extproc records the row before it emits a response mutator's replacement. Headers are not counted. A tunnel close is the one row carrying both. **Zero means the listener counted nothing, not that the body was empty**, and there is no third value for "unknown": every listener buffers bodies only when a plugin asks, so an unbuffered request reports zero, and so does a response relayed straight through — both proxies copy that one after the row is appended. agentop renders a zero side blank rather than as `0B`, and that is why its BYTES column can be trusted as a size when it shows one. All three listeners count the same thing for the same response, which took some doing: extproc sums the chunks Envoy hands it, while the proxies count off the upstream body at one choke point each, because the arm that re-frames SSE sees `sseframe` payloads with the `data: ` prefixes and blank-line separators already stripped — 32 bytes short on a four-event response. `core/listener/internal/bodycount` and `core/listener/parity`'s bytes fixtures hold that line. +- `tunnel`, `tunnelReason` — an opaque CONNECT (or transparent-redirect) tunnel records two rows sharing a `requestId`: the open (`phase: "request"`, `tunnelReason` saying why the bytes stayed opaque) and, when the tunnel ends, the close (`phase: "response"`). The close carries the CONNECT's own `statusCode` — 200, or 502 with the dial error in `error` when the destination could not be reached (`tunnelReason: "dial-failed"`) — plus `durationMs` for how long the tunnel stayed open and, in `bytesUp`/`bytesDown`, the bytes it carried each way (up = client to destination; opaque bytes have no body/header split, so these are everything that crossed). It is not the destination's status: that travels inside the client's end-to-end TLS. A bridged tunnel's open is recorded with its first decrypted request — in that request's session, directly before it, stamped with its time — and records no close, because the request carries its own response; a bridged tunnel that recorded no request gets its open and a close when it ends. Tunnel rows are kept out of `/v1/usage`: a tunnel's lifetime is not a request latency. - `httpMethod`, `httpPath` — the HTTP verb and path, so a request no parser recognized is still identifiable rather than showing only a host. Distinct from the `method` inside `a2a` / `mcp`, which is a protocol method name. On an opaque tunnel `httpMethod` is `CONNECT` and `httpPath` is absent — opaque bytes carry no request line. The path is query-stripped and percent-decoded, so query-borne credentials never reach the timeline, but a secret in a path *segment* (a bot token, a webhook path) does survive on this unauthenticated surface — worth knowing before exporting events off-box. ### Invocation action vocabulary diff --git a/cmd/agentop/README.md b/cmd/agentop/README.md index c1d44e7af..c8e9b808b 100644 --- a/cmd/agentop/README.md +++ b/cmd/agentop/README.md @@ -969,7 +969,7 @@ agentop is for, and the other three are surfaces you visit and leave. 100% of rows. Avoided spend now lives per session in the sessions table and per series in the `$` breakdown; volume readings live in the Usage pane. - **Events**: per-session event table. `c` opens a column picker — a popup with - a checkbox and a one-line description per column, since twelve abbreviated + a checkbox and a one-line description per column, since thirteen abbreviated headers are not self-describing. The picker is also where sorting lives: `s` orders the table by the column @@ -987,47 +987,59 @@ agentop is for, and the other three are surfaces you visit and leave. toward. Sorting never changes the `#` exchange pairing or the per-row token and cost figures; it reorders the finished rows only. - The twelve default columns together need ~168 terminal columns, so the table + The thirteen default columns together need ~185 terminal columns, so the table drops what does not fit and the footer says how many (`→ N more columns`). Columns carry a keep rank rather than being equally expendable: DIR, DURATION, TOKENS and COST give way first, while `#` and HOST survive longest. That is what makes HOST usable at 80 columns despite being last in display order — it is the column most people open this pane for. - Twelve of the thirteen are on by default: `#` (exchange number, shared by a - request and its response), TIME, DIR, PHASE, ACTION, PLUGIN, METHOD, STATUS, - DURATION, TOKENS, COST, HOST. The thirteenth, BYTES, is opt-in through the - picker (`c`), because only an opaque tunnel's close row has a figure for it. - On a narrow terminal the low-ranked ones are hidden rather than turned off, - so widening the window brings them back without touching the picker. + All thirteen are on by default: `#` (exchange number, shared by a request and + its response), TIME, DIR, PHASE, ACTION, PLUGIN, METHOD, STATUS, DURATION, + BYTES, TOKENS, COST, HOST. On a narrow terminal the low-ranked ones are hidden + rather than turned off, so widening the window brings them back without + touching the picker. + + BYTES was the one opt-in column until recently, on the grounds that only an + opaque tunnel's close row had a figure for it. Ordinary rows now carry one too + — a request row the body size it forwarded, a response row the body size that + came back — so the reason for hiding it is gone and the picker is where you go + to turn it *off*. A blank cell means the proxy counted nothing, which is not + the same as a body of zero bytes: a body no plugin asked to buffer is relayed + unmeasured and reads exactly as a body-less GET does. Live-updates while in view — if the cursor is on the last row, it auto-follows new events. - All twelve columns, on a terminal wide enough for them. Each `#` appears + All thirteen columns, on a terminal wide enough for them. Each `#` appears twice — once for the request, once for its response — which is how a row with no STATUS is read as "still in flight" rather than "failed": ``` agentop · ctx-abc-1234… - # TIME DIR PHASE ACTION PLUGIN METHOD STATUS DURATION TOKENS COST HOST - 1 14:23:07.41 in req allow jwt-validation weather-agent - 1 14:23:07.52 in resp — — 200 118ms weather-agent - 2 14:23:07.71 out req observe inference-parser claude-sonnet-5 681,300(−9.9k) $0.2546(−$0.0037) api.anthropic.com - 2 14:23:08.91 out resp — — claude-sonnet-5 200 1.20s 412 api.anthropic.com - 3 14:23:09.01 out req modify token-exchange tools/call github-tool-mcp - 3 14:23:09.10 out resp — — tools/call 503 96ms github-tool-mcp + # TIME DIR PHASE ACTION PLUGIN METHOD STATUS DURATION BYTES TOKENS COST HOST + 1 14:23:07.41 in req allow jwt-validation ↑1.4kB weather-agent + 1 14:23:07.52 in resp — — 200 118ms ↓842B weather-agent + 2 14:23:07.71 out req observe inference-parser claude-sonnet-5 ↑118.6kB 681,300(−9.9k) $0.2546(−$0.0037) api.anthropic.com + 2 14:23:08.91 out resp — — claude-sonnet-5 200 1.20s ↓12.4kB 412 api.anthropic.com + 3 14:23:09.01 out req modify token-exchange tools/call ↑318B github-tool-mcp + 3 14:23:09.10 out resp — — tools/call 503 96ms github-tool-mcp ● connected 2.1 events/sec [sort: DURATION▼] [filter: anthropic] - [↑↓] nav [b/f] page [↵] detail [c] columns [u] usage [s] hide passthru/skip [p] pause [/] filter [esc] back · → 4 more columns ([c] to choose) [?] keys [q] quit + [↑↓] nav [b/f] page [↵] detail [c] columns [u] usage [s] hide passthru/skip [p] pause [/] filter [esc] back · [?] keys [q] quit ``` + Exchange 3's response has no BYTES figure because there was nothing to count: + token-exchange could not reach the IdP, so the 503 is the proxy's own and no + upstream body ever arrived. + `—` in ACTION and PLUGIN means no plugin acted on that message; a `tunnel` there is an opaque CONNECT, where METHOD is blank because opaque bytes carry no request line. The tunnel's STATUS arrives on its `resp` row when it closes, with DURATION for how long it stayed open and, in BYTES, what it carried each - way. That STATUS is the + way — the one row that reports both directions at once, since opaque bytes + have no body-and-headers split to separate. That STATUS is the CONNECT's own (200, or 502 when the destination could not be reached), never the destination's, which travels inside the client's TLS. The TOKENS and COST figures on a request row carry what `tool-prune` saved in parentheses — `−` for a counted saving, @@ -1396,7 +1408,7 @@ filter: github-tool | `usage.window` | string | `10m0s` | usage-pane window: `10m0s`, `1h0m0s` or `6h0m0s` | | `usage.group` | string | `none` | usage-pane breakdown. `[b]` cycles `none`, `status`, `method`, `plugin`, `host`; a hand-edited file may also use `model`, `endpoint`, `session` or `agent` | -List a column only to change it — the twelve are all visible until you hide one, and +List a column only to change it — the thirteen are all visible until you hide one, and an unrecognised name or sort column is ignored. Inside an entry, always write `visible:` explicitly: it is optional to the parser but reads as `false`, so `- name: COST` on its own hides COST rather than showing it. @@ -1422,7 +1434,7 @@ Column ids, in display order — the same headers the picker shows: | `METHOD` | protocol operation: model name, MCP or A2A method | | `STATUS` | HTTP status of the response | | `DURATION` | how long the exchange took | -| `BYTES` | bytes an opaque tunnel carried: ↑ sent, ↓ received — off by default | +| `BYTES` | body bytes: ↑ request sent, ↓ response received (tunnels show both); blank when the proxy counted nothing | | `TOKENS` | tokens used, and what `tool-prune` saved | | `COST` | estimated cost, and what `tool-prune` saved | | `HOST` | host the message was sent to | diff --git a/cmd/agentop/tui/events_column_sizing_test.go b/cmd/agentop/tui/events_column_sizing_test.go index 2843e7920..15cfda041 100644 --- a/cmd/agentop/tui/events_column_sizing_test.go +++ b/cmd/agentop/tui/events_column_sizing_test.go @@ -239,11 +239,17 @@ func TestEventsTable_WiderFigureKeepsScrollPosition(t *testing.T) { } // The picker's "(no room)" marker is judged against the widths the table is using. At -// 150 columns the default set fits only because TOKENS and COST are sized to their +// this width the default set fits only because TOKENS and COST are sized to their // figures; judged at their declared widths, COST would be marked as having no room while // the table beside it showed it. +// +// 169 and not 150: #1309 turned BYTES on by default, which is 19 more columns of +// declared width (17 + cellPadding) and does not size to its content, so the window +// where sizing is what makes the defaults fit moved up by exactly that. The control +// assertion below is what keeps this honest — at declared widths this width must still +// be too narrow, or the test proves nothing. func TestColumnPicker_NoRoomAgreesWithTheSizedTable(t *testing.T) { - const width = 150 + const width = 169 m := sizingModel(t, width, sizingExchange(t, "r1", "api.anthropic.com", time.Now(), 1_048_576, nil)) if m.eventColsDropped != 0 { @@ -264,14 +270,19 @@ func TestColumnPicker_NoRoomAgreesWithTheSizedTable(t *testing.T) { } // A widening that crosses a fit boundary changes the column COUNT, not only widths, and -// must not re-anchor the pane either. At 150 columns the fixture's sized defaults fit; the -// first counted saving widens TOKENS and COST past that, a column is dropped, and the +// must not re-anchor the pane either. At the narrow width the fixture's sized defaults fit; +// the first counted saving widens TOKENS and COST past that, a column is dropped, and the // headings change. Then the reverse: a wider terminal lets the column back on, so the count // grows. Each direction takes a different SetRows/SetColumns order, and the two together are // the only ones a toggle or a resize can produce. +// +// Both widths moved up by BYTES' 19 declared columns when #1309 turned it on by default — +// the test needs a width where all thirteen fit and a wider one that still has room after +// the saving, not these particular numbers. The two Fatalf guards below catch it if a +// future column moves the boundary again. func TestEventsTable_WideningPastTheTerminalKeepsScrollPosition(t *testing.T) { m := cursorModel(t, 40) - m.width = 150 + m.width = 169 m.rebuildEventsTable() if m.eventColsDropped != 0 { t.Fatalf("%d columns dropped at %d before any saving; the fixture should start with all of them", @@ -292,7 +303,7 @@ func TestEventsTable_WideningPastTheTerminalKeepsScrollPosition(t *testing.T) { check func(n int) bool }{ {"a saving pushes a column off", func() {}, func(n int) bool { return n < before }}, - {"a wider terminal lets it back", func() { m.width = 200 }, func(n int) bool { return n == before }}, + {"a wider terminal lets it back", func() { m.width = 219 }, func(n int) bool { return n == before }}, } { step.apply() m.rebuildEventsTable() diff --git a/cmd/agentop/tui/events_columns.go b/cmd/agentop/tui/events_columns.go index edbc6ef42..a9e00fbeb 100644 --- a/cmd/agentop/tui/events_columns.go +++ b/cmd/agentop/tui/events_columns.go @@ -276,11 +276,20 @@ var eventColumns = []eventColumn{ desc: "how long the exchange took", cell: func(c cellContext) string { return durationCell(*c.row.event) }, sortKey: durationSortKey}, - // Off by default: only an opaque tunnel's close row has a figure, so a column on for - // everyone would spend columns of every terminal on a mostly blank one. The total is - // the sort key, so a descending sort finds the tunnel that carried the most. - {id: colBytes, width: 17, defaultOn: false, keep: keepLow, rightAlign: true, - desc: "bytes an opaque tunnel carried: ↑ sent, ↓ received", + // On by default since #1309. It was off because only an opaque tunnel's close row had + // a figure, so the column cost every terminal 17 columns of mostly blank — and an + // operator asking how big the request that overran a model's context window was found + // it blank on every row and resorted to measuring the yanked event's JSON, which is + // not the body. Ordinary request and response rows now carry a figure, so the reason + // for opt-in is gone. + // + // 17 is the tunnel close row's pair at its widest ("↑999.9GB ↓999.9GB"); an ordinary + // row reports one side and is strictly narrower. keepLow stays — it is still the + // first column a narrow terminal drops. The total is the sort key, which is also + // right for a one-sided row, so a descending sort finds the biggest payload either + // way. + {id: colBytes, width: 17, defaultOn: true, keep: keepLow, rightAlign: true, + desc: "body bytes: ↑ request sent, ↓ response received (tunnels show both)", cell: func(c cellContext) string { return bytesCell(*c.row.event) }, sortKey: func(c cellContext) sortValue { return numKey(c.row.event.BytesUp + c.row.event.BytesDown) }}, // 18: sized for a SEVEN-digit prompt with a counted saving, "1,048,576(−12,300)". diff --git a/cmd/agentop/tui/events_columns_test.go b/cmd/agentop/tui/events_columns_test.go index 4acd23597..af61ed2af 100644 --- a/cmd/agentop/tui/events_columns_test.go +++ b/cmd/agentop/tui/events_columns_test.go @@ -548,9 +548,10 @@ func TestColumnPicker_CursorStaysVisibleWhenClipped(t *testing.T) { // // bubbles pads each cell on both sides (Padding(0, 1) on Cell and Header), so a // column occupies width+2. Modelling it as +1 under-counted by one per column: a -// row of all twelve rendered at 168 against a computed 156, so the table wrapped -// at 80 and — worse, between 156 and 167 — reported dropped==0 while up to twelve -// columns sat off the edge, with no footer count and no "(no room)" marker. +// row of all the columns there were then — twelve, before #1309 turned BYTES on — +// rendered at 168 against a computed 156, so the table wrapped at 80 and — worse, +// between 156 and 167 — reported dropped==0 while up to twelve columns sat off the +// edge, with no footer count and no "(no room)" marker. func TestFitColumns_RenderedWidthNeverExceedsTerminal(t *testing.T) { // Several selections, not just all-on. The all-on case alone let a real bug // through: fitColumns credited back width+1 while columnsWidth charged diff --git a/cmd/agentop/tui/events_pane.go b/cmd/agentop/tui/events_pane.go index 8e5ac0321..535f17094 100644 --- a/cmd/agentop/tui/events_pane.go +++ b/cmd/agentop/tui/events_pane.go @@ -829,15 +829,34 @@ func generatedTokensCell(e *pipeline.SessionEvent) string { return formatCount(n) } -// bytesCell renders a tunnel close row's byte counts, up (client to destination) then -// down. Blank on every other row: only an opaque tunnel's close has counts. They are -// absent on the wire when zero, so a zero on one side is shown only once the other -// side proves this row carries counts at all. +// bytesCell renders a row's byte counts, up (client to destination) then down. +// +// A tunnel's close row is the one that reports BOTH, and it keeps the pair rendering +// even where one side is zero: the counts are absent on the wire when zero, so a +// genuine zero one way is shown only once the other side proves the row carries +// counts at all. Every other row reports a single side — a request knows what it +// forwarded, a response what came back — so printing the pair there would put a +// fabricated "↓0B" next to every request size. +// +// Blank when nothing was counted, which is NOT the same as a body of zero bytes: +// an unbuffered body is relayed unmeasured and reads here exactly as a body-less GET +// does. See pipeline.SessionEvent.BytesUp — formatBytes renders 0 as "0B", never +// blank, so the decision is this gate's alone. func bytesCell(e pipeline.SessionEvent) string { if e.BytesUp == 0 && e.BytesDown == 0 { return "" } - return "↑" + formatBytes(e.BytesUp) + " ↓" + formatBytes(e.BytesDown) + if e.Tunnel { + return "↑" + formatBytes(e.BytesUp) + " ↓" + formatBytes(e.BytesDown) + } + var parts []string + if e.BytesUp > 0 { + parts = append(parts, "↑"+formatBytes(e.BytesUp)) + } + if e.BytesDown > 0 { + parts = append(parts, "↓"+formatBytes(e.BytesDown)) + } + return strings.Join(parts, " ") } // formatBytes renders a byte count in decimal units, 1kB being 1000 bytes. Like diff --git a/cmd/agentop/tui/settings.go b/cmd/agentop/tui/settings.go index 4aafecac3..f321990a3 100644 --- a/cmd/agentop/tui/settings.go +++ b/cmd/agentop/tui/settings.go @@ -36,16 +36,18 @@ type UserSettings struct { // EventSettings is the events-table view state. type EventSettings struct { // Columns records only DEVIATIONS from the defaults — a column absent here takes - // its own default, which is visible for every column but the opt-in BYTES. Two + // its own default, which since #1309 is visible for every column there is. Two // consequences, both deliberate: // // A column added in a later build shows up for everyone who already has a // config file, instead of starting hidden because their file predates it. And // the file stays legible: turn two columns off and it has two entries, not - // twelve. + // thirteen. // // The alternative (a list of what is ON) reads more obviously but inverts both - // of those, since every column but one is defaultOn. + // of those, since every column is defaultOn. Note the scheme still has to + // round-trip a default-OFF column — a future one, or BYTES again — which is why + // the tests for that exercise a synthetic column rather than a real one. Columns []ColumnSetting `yaml:"columns,omitempty"` // SortColumn is the events-table sort column, by the same stable id the picker // shows as a header (#865). Empty — the zero value, and what every file written @@ -121,9 +123,10 @@ type ColumnSetting struct { // a phantom key would make the former report a selection the latter cannot render. // // Absent means "unchanged", so the column falls back to its own defaultOn. Every -// column but BYTES is defaultOn, which is what makes "absent means visible" true of -// them — but the fallback is the default, not a literal true, so BYTES stays off for -// a file that predates it, exactly as it does in defaultColumnSelection. +// column is defaultOn since #1309, which is what makes "absent means visible" true +// of them today — but the fallback is the column's default and not a literal true, +// so the next default-off column needs no change here, exactly as in +// defaultColumnSelection. func (u UserSettings) columnSelection() map[eventColumnID]bool { // Both directions, not just the off-list: a column whose defaultOn is false has // to be turnable ON by the file, or the picker could not persist enabling it. @@ -135,9 +138,10 @@ func (u UserSettings) columnSelection() map[eventColumnID]bool { for _, c := range eventColumns { // The fallback is c.defaultOn, not an unconditional true: "absent from the file" // means "I never changed this", so the answer is whatever the built-in default - // is. BYTES, the first column with defaultOn:false, is why hardcoding true here - // would be a trap: it would render ON for every user, fresh installs included, - // and disagree with defaultColumnSelection. + // is. BYTES was the first column with defaultOn:false and is why hardcoding true + // here would be a trap — it would render such a column ON for every user, fresh + // installs included, and disagree with defaultColumnSelection. BYTES is on by + // default now (#1309), so nothing exercises that today but the tests. if v, ok := set[string(c.id)]; ok { out[c.id] = v continue @@ -147,7 +151,7 @@ func (u UserSettings) columnSelection() map[eventColumnID]bool { // A file turning every column off would render a table with no columns. The // picker cannot produce that state (its toggle re-seeds the defaults), but a // hand-edited file can, and selectedColumns' fallback only rescues the render — - // the picker's checkboxes would still show twelve empty boxes. Re-seed here so + // the picker's checkboxes would still show thirteen empty boxes. Re-seed here so // the two agree, matching what the toggle handler already does. if !anyColumnSelected(out) { return defaultColumnSelection() @@ -182,7 +186,7 @@ func (u UserSettings) sortSelection() (eventColumnID, bool) { // Emits only the columns that differ from their own defaultOn, so an untouched // selection serialises to nothing rather than to a `visible: true` entry per column — // the file should record what the user changed: the default columns turned off, and -// BYTES if it was turned on. +// any default-off column turned on. // // Ordered by eventColumns, not by map iteration: Go randomises the latter, which // would rewrite the file with the same entries reshuffled on every save. That @@ -193,7 +197,8 @@ func columnSettingsFrom(sel map[eventColumnID]bool) []ColumnSetting { // Compared against the column's own default, in both directions: turning ON a // defaultOn:false column is as much a deviation as turning off a default one, // and recording only the off-list would make enabling such a column - // unpersistable. BYTES is that column. + // unpersistable. BYTES was that column until #1309 turned it on by default; + // the round-trip is still required, and is covered with a synthetic column. if sel[c.id] != c.defaultOn { out = append(out, ColumnSetting{Name: string(c.id), Visible: sel[c.id]}) } diff --git a/cmd/agentop/tui/settings_test.go b/cmd/agentop/tui/settings_test.go index c3af5a754..101284dc6 100644 --- a/cmd/agentop/tui/settings_test.go +++ b/cmd/agentop/tui/settings_test.go @@ -21,9 +21,15 @@ func resetSettingsForTest(t *testing.T) { // TestColumnSelection_AbsentColumnsDefaultOn is the headline rule of the file // format: the config records deviations, so a column the file does not mention takes -// its own default — VISIBLE for every default column. Inverting this would mean a -// column added in a later agentop starts hidden for everyone who already has a config -// file; ignoring the default would switch on the opt-in BYTES for all of them. +// its own default. Inverting this would mean a column added in a later agentop starts +// hidden for everyone who already has a config file. +// +// Asserted against each column's own defaultOn rather than against a literal true, +// which is the same distinction columnSelection itself draws: every column is +// defaultOn today (#1309 turned the last one on), so a test written as "absent means +// visible" would pass now and quietly stop testing the rule the next time a column +// ships opt-in. TestColumnSettings_RoundTripsADefaultOffColumn covers that case with +// a synthetic column. func TestColumnSelection_AbsentColumnsDefaultOn(t *testing.T) { s := UserSettings{Events: EventSettings{Columns: []ColumnSetting{ {Name: string(colCost), Visible: false}, @@ -42,24 +48,34 @@ func TestColumnSelection_AbsentColumnsDefaultOn(t *testing.T) { t.Errorf("column %q is absent from the config; want its default %v, got %v", c.id, c.defaultOn, sel[c.id]) } } - if sel[colBytes] { - t.Error("BYTES is opt-in, but a config that never mentions it turned it on") - } } // TestColumnSettings_EnablingAnOptInColumnPersists: turning ON a defaultOn:false column -// is a deviation too, so it must be written and read back. BYTES is the first such -// column; before it this path had no real case to test. +// is a deviation too, so it must be written and read back. Recording only the off-list +// would make it unpersistable — the user would enable the column, close the picker, and +// find it off again next launch. +// +// Exercised through a synthetic column because no real one is default-off any more: +// BYTES was, until #1309 populated it for ordinary traffic and turned it on. The scheme +// still has to carry such a column, so this keeps testing it rather than retiring with +// the last real example. func TestColumnSettings_EnablingAnOptInColumnPersists(t *testing.T) { + prev := eventColumns + t.Cleanup(func() { eventColumns = prev }) + eventColumns = []eventColumn{ + {id: colTime, width: 12, defaultOn: true, cell: func(cellContext) string { return "" }}, + {id: "OPTIONAL", width: 8, defaultOn: false, cell: func(cellContext) string { return "" }}, + } + sel := defaultColumnSelection() - sel[colBytes] = true + sel["OPTIONAL"] = true written := columnSettingsFrom(sel) - if len(written) != 1 || written[0].Name != string(colBytes) || !written[0].Visible { - t.Fatalf("enabling BYTES wrote %+v, want exactly [{BYTES true}]", written) + if len(written) != 1 || written[0] != (ColumnSetting{Name: "OPTIONAL", Visible: true}) { + t.Fatalf("enabling OPTIONAL wrote %+v, want exactly [{OPTIONAL true}]", written) } back := UserSettings{Events: EventSettings{Columns: written}}.columnSelection() - if !back[colBytes] { - t.Error("BYTES was enabled and saved, but reads back off") + if !back["OPTIONAL"] { + t.Error("OPTIONAL was enabled and saved, but reads back off") } } diff --git a/cmd/agentop/tui/tunnel_close_test.go b/cmd/agentop/tui/tunnel_close_test.go index a3d754790..850ca6ce2 100644 --- a/cmd/agentop/tui/tunnel_close_test.go +++ b/cmd/agentop/tui/tunnel_close_test.go @@ -158,13 +158,15 @@ func TestDurationCell_MinutesAndHours(t *testing.T) { } } -func TestBytesCell(t *testing.T) { +// A tunnel's close row reports both directions, and keeps doing so where one of them +// is a real zero. +func TestBytesCell_TunnelReportsBothWays(t *testing.T) { for _, tc := range []struct { name string up, down int64 want string }{ - {"not a tunnel close", 0, 0, ""}, + {"nothing counted", 0, 0, ""}, {"both ways", 4210, 18230, "↑4.2kB ↓18.2kB"}, {"small counts stay exact", 5, 7, "↑5B ↓7B"}, // Counts are absent on the wire when zero, so one side at zero is shown as 0 @@ -175,16 +177,45 @@ func TestBytesCell(t *testing.T) { {"carries at the tier boundary", 999_950, 999_949, "↑1.0MB ↓999.9kB"}, } { t.Run(tc.name, func(t *testing.T) { - if got := bytesCell(pipeline.SessionEvent{BytesUp: tc.up, BytesDown: tc.down}); got != tc.want { + ev := pipeline.SessionEvent{Tunnel: true, BytesUp: tc.up, BytesDown: tc.down} + if got := bytesCell(ev); got != tc.want { t.Errorf("bytesCell(up=%d, down=%d) = %q, want %q", tc.up, tc.down, got, tc.want) } }) } } -// Off by default: only opaque tunnels' close rows have a figure, so on by default it -// would spend columns of every terminal on a mostly blank column. -func TestBytesColumn_IsOptInAndSortsByTotal(t *testing.T) { +// An ordinary row knows one side — a request what it forwarded, a response what came +// back (#1309) — so it reports that side alone. The pair rendering would put a +// fabricated "↓0B" beside every request size, which is the opposite of the tunnel +// case above: there a zero is a measured zero, here it is an absence. +func TestBytesCell_OrdinaryRowReportsOneSide(t *testing.T) { + for _, tc := range []struct { + name string + up, down int64 + want string + }{ + {"request only", 121_406, 0, "↑121.4kB"}, + {"response only", 0, 2310, "↓2.3kB"}, + // Nothing counted on either side — an unbuffered body, or a body-less GET. + // Blank and not "↑0B ↓0B", and not "0B": the proxy does not know. + {"neither counted", 0, 0, ""}, + // Both sides is not a shape the listeners produce on one row, but the renderer + // must not drop half of one if they ever did. + {"both, defensively", 10, 20, "↑10B ↓20B"}, + } { + t.Run(tc.name, func(t *testing.T) { + ev := pipeline.SessionEvent{BytesUp: tc.up, BytesDown: tc.down} + if got := bytesCell(ev); got != tc.want { + t.Errorf("bytesCell(up=%d, down=%d) = %q, want %q", tc.up, tc.down, got, tc.want) + } + }) + } +} + +// On by default since #1309: ordinary request and response rows carry a figure now, so +// the reason it shipped opt-in — only tunnel close rows had one — is gone. +func TestBytesColumn_IsOnByDefaultAndSortsByTotal(t *testing.T) { var col *eventColumn for i := range eventColumns { if eventColumns[i].id == colBytes { @@ -194,8 +225,8 @@ func TestBytesColumn_IsOptInAndSortsByTotal(t *testing.T) { if col == nil { t.Fatal("no BYTES column") } - if col.defaultOn { - t.Error("BYTES is on by default; it should be opt-in") + if !col.defaultOn { + t.Error("BYTES is opt-in; it should be on by default") } if col.sortKey == nil { t.Fatal("BYTES has no sort key") diff --git a/core/listener/extproc/bytescount_test.go b/core/listener/extproc/bytescount_test.go new file mode 100644 index 000000000..8b0f8087a --- /dev/null +++ b/core/listener/extproc/bytescount_test.go @@ -0,0 +1,144 @@ +package extproc + +import ( + "context" + "testing" + + extprocv3 "github.com/envoyproxy/go-control-plane/envoy/service/ext_proc/v3" + + "github.com/rossoctl/cortex/core/pipeline" + "github.com/rossoctl/cortex/core/session" +) + +// TestBytes_SplitResponseBodySumsEveryMessage is this listener's reason for +// counting response bytes in an accumulator instead of measuring the buffer at the +// record site (#1309). +// +// Envoy may deliver one response body in several ResponseBody messages, and on the +// SSE arm handleResponseBody REPLACES pctx.ResponseBody each time rather than +// appending — it carries the SSE tail forward, not the whole response. So by the +// time the row is recorded the buffer holds the LAST message only, and +// len(pctx.ResponseBody) would report a streamed inference response as the size of +// its final chunk. On this two-message fixture that is a figure roughly half the +// truth, reported with no indication anything was missing — on exactly the traffic +// the BYTES column was asked for. +// +// Reuses splitStreamRequests from the split-body cost test rather than scripting a +// second two-message stream: it is already the shape that gets this wrong, and +// pinning both the charge and the byte count against one fixture keeps them from +// drifting into disagreement about what that response was. +// +// The parity suite cannot stand in for this. Swapping the accumulator back for +// len(pctx.ResponseBody) leaves listener/parity's bytes fixtures GREEN — its driver +// delivers each body in a single message, where the buffer and the sum are the same +// number. Only a multi-message body tells them apart, and only this listener has +// one. +func TestBytes_SplitResponseBodySumsEveryMessage(t *testing.T) { + srv, store := newStreamedServer(t) + reqs := splitStreamRequests() + + var ( + wantUp int64 + wantDown int64 + lastDown int64 + ) + for _, r := range reqs { + switch m := r.Request.(type) { + case *extprocv3.ProcessingRequest_RequestBody: + wantUp += int64(len(m.RequestBody.Body)) + case *extprocv3.ProcessingRequest_ResponseBody: + lastDown = int64(len(m.ResponseBody.Body)) + wantDown += lastDown + } + } + if lastDown == wantDown { + t.Fatalf("fixture delivers the response in one message; it cannot distinguish the sum from the trailing chunk") + } + + stream := &mockStream{ctx: context.Background(), requests: reqs} + _ = srv.Process(stream) + if stream.recvIdx != len(reqs) { + t.Fatalf("consumed %d of %d messages; the listener bailed", stream.recvIdx, len(reqs)) + } + + ev := responseEvent(t, store) + if ev == nil { + t.Fatal("no outbound response row recorded for a split streamed body") + } + switch ev.BytesDown { + case wantDown: + case lastDown: + t.Errorf("BytesDown = %d — the TRAILING MESSAGE only, which is all pctx.ResponseBody "+ + "holds once the SSE arm has replaced it. The whole body was %d bytes", + ev.BytesDown, wantDown) + default: + t.Errorf("BytesDown = %d, want %d (the sum of every ResponseBody message)", ev.BytesDown, wantDown) + } + + v := store.View(session.DefaultSessionID) + var reqRow *pipeline.SessionEvent + for i := range v.Events { + if v.Events[i].Phase == pipeline.SessionRequest && v.Events[i].Direction == pipeline.Outbound { + reqRow = &v.Events[i] + break + } + } + if reqRow == nil { + t.Fatalf("no outbound request row recorded; events = %+v", v.Events) + } + if reqRow.BytesUp != wantUp || reqRow.BytesDown != 0 { + t.Errorf("request row = ↑%d ↓%d, want ↑%d ↓0", reqRow.BytesUp, reqRow.BytesDown, wantUp) + } +} + +// TestBytes_HeaderOnlyResponseReportsZero pins the honest zero on the listener that +// has the least room to do better: Envoy sends no ResponseBody message for a +// response that ends on its headers, so there is nothing for the accumulator to add +// and nothing to infer from. +// +// Zero means the listener counted nothing — the event contract has no third value +// for "unknown", and agentop renders it blank rather than as 0B. A content-length +// header would be a plausible-looking substitute and is the wrong one: it is absent +// or -1 on everything streamed, which is most of what passes through here. +func TestBytes_HeaderOnlyResponseReportsZero(t *testing.T) { + srv, store := newStreamedServer(t) + reqs := []*extprocv3.ProcessingRequest{ + {Request: &extprocv3.ProcessingRequest_RequestHeaders{ + RequestHeaders: &extprocv3.HttpHeaders{ + Headers: makeHeaders( + "x-authbridge-direction", "outbound", + ":method", "GET", + ":path", "/v1/models", + ":authority", "litellm.local", + ), + EndOfStream: true, + }, + }}, + {Request: &extprocv3.ProcessingRequest_ResponseHeaders{ + ResponseHeaders: &extprocv3.HttpHeaders{ + Headers: makeHeaders( + ":status", "200", + // Load-bearing, not scenery: this listener records a row only + // when the pipeline left something on pctx, and on a response + // with no body the gateway's cost header is the only thing the + // inference-parser can claim. Without it there is no row to + // assert the zero on — which is a different fact from the zero. + "x-litellm-response-cost", "0.002", + ), + // The one fact that makes this a header-only response. + EndOfStream: true, + }, + }}, + } + + stream := &mockStream{ctx: context.Background(), requests: reqs} + _ = srv.Process(stream) + + ev := responseEvent(t, store) + if ev == nil { + t.Fatal("no outbound response row recorded for a header-only response") + } + if ev.BytesUp != 0 || ev.BytesDown != 0 { + t.Errorf("response row = ↑%d ↓%d, want ↑0 ↓0 — no body message arrived, so there is nothing to count", ev.BytesUp, ev.BytesDown) + } +} diff --git a/core/listener/extproc/server.go b/core/listener/extproc/server.go index 4f08f08c7..93ea94bf6 100644 --- a/core/listener/extproc/server.go +++ b/core/listener/extproc/server.go @@ -357,6 +357,7 @@ func (s *Server) recordInboundSession(pctx *pipeline.Context) { Direction: pipeline.Inbound, Phase: pipeline.SessionRequest, RequestID: pctx.RequestID(), + BytesUp: int64(len(pctx.Body)), A2A: pipeline.SnapshotA2A(pctx.Extensions.A2A), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseRequest), @@ -441,6 +442,7 @@ func (s *Server) recordInboundResponseSession(pctx *pipeline.Context) { Direction: pipeline.Inbound, Phase: pipeline.SessionResponse, RequestID: pctx.RequestID(), + BytesDown: pctx.ResponseBytes, A2A: pipeline.SnapshotA2A(pctx.Extensions.A2A), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseResponse), @@ -474,6 +476,7 @@ func (s *Server) recordOutboundResponseSession(pctx *pipeline.Context) { Direction: pipeline.Outbound, Phase: pipeline.SessionResponse, RequestID: pctx.RequestID(), + BytesDown: pctx.ResponseBytes, MCP: pipeline.SnapshotMCP(pctx.Extensions.MCP), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseResponse), @@ -525,6 +528,7 @@ func (s *Server) recordOutboundSession(pctx *pipeline.Context) { Direction: pipeline.Outbound, Phase: pipeline.SessionRequest, RequestID: pctx.RequestID(), + BytesUp: int64(len(pctx.Body)), MCP: pipeline.SnapshotMCP(pctx.Extensions.MCP), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseRequest), @@ -840,6 +844,15 @@ func (s *Server) handleResponseBody(ctx context.Context, body []byte, pctx *pipe } } + // Count before either arm touches the buffer. Neither arm's result is a tally: the SSE arm + // REPLACES ResponseBody with carry + this message, and the non-SSE arm stops growing at + // maxBodySize. So the length of that field is the trailing chunk on one path and a floor on + // the other, while this is the whole body as Envoy relayed it — which is what BytesDown + // reports. On the shipped Envoy config an oversized response never reaches here at all + // (Envoy's own buffer limit refuses it first, see #1325), leaving this zero, and zero is + // read as "not counted" rather than as an empty body. + pctx.ResponseBytes += int64(len(body)) + // A ResponseBody message is a chunk of a byte stream and NOT a unit of anything else: nothing // aligns Envoy's chunk boundaries with the body's own structure. Both arms exist because of // that, and differ only in how much has to be kept. diff --git a/core/listener/forwardproxy/bytescount_test.go b/core/listener/forwardproxy/bytescount_test.go new file mode 100644 index 000000000..86d7bece1 --- /dev/null +++ b/core/listener/forwardproxy/bytescount_test.go @@ -0,0 +1,150 @@ +package forwardproxy + +import ( + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/rossoctl/cortex/core/pipeline" + "github.com/rossoctl/cortex/core/session" +) + +// bytesOf drives one request through a forward proxy built on the given plugins +// and returns the byte counts the request and response rows reported (#1309). +// +// Returned as the two rows rather than four numbers because the split is the +// contract: a request row counts what was forwarded and nothing else, a response +// row counts what came back and nothing else, and a listener that put both on one +// row would read as a round trip that never happened. +func bytesOf(t *testing.T, plugins []pipeline.Plugin, req *http.Request, upstream http.HandlerFunc) (reqRow, respRow pipeline.SessionEvent) { + t.Helper() + p, err := pipeline.New(plugins) + if err != nil { + t.Fatalf("pipeline.New: %v", err) + } + backend := httptest.NewServer(upstream) + defer backend.Close() + + store := session.New(5*time.Minute, 100, 0) + defer store.Close() + + srv := &Server{OutboundPipeline: pipeline.NewHolder(p), Sessions: store, Client: http.DefaultClient} + proxy := httptest.NewServer(srv.Handler()) + defer proxy.Close() + + // The caller builds the request against a placeholder host so it can be + // written before the backend exists; retarget it here. + target := backend.URL + req.URL.Path + retargeted, err := http.NewRequest(req.Method, target, req.Body) + if err != nil { + t.Fatalf("http.NewRequest: %v", err) + } + retargeted.Header = req.Header + retargeted.ContentLength = req.ContentLength + + client := &http.Client{Transport: &http.Transport{Proxy: http.ProxyURL(mustParseURL(proxy.URL))}} + resp, err := client.Do(retargeted) + if err != nil { + t.Fatalf("request failed: %v", err) + } + // Drained, not just closed: the streamed arms count as the relay reads, so an + // undrained body would under-report by whatever was still in flight. + drainAndClose(t, resp) + + v := store.View(session.DefaultSessionID) + if v == nil || len(v.Events) != 2 { + t.Fatalf("want a request row and a response row, got %+v", v) + } + return v.Events[0], v.Events[1] +} + +func drainAndClose(t *testing.T, resp *http.Response) { + t.Helper() + buf := make([]byte, 4096) + for { + if _, err := resp.Body.Read(buf); err != nil { + break + } + } + resp.Body.Close() +} + +// TestBytes_BufferedPairCountsEachDirectionOnItsOwnRow is the forward proxy's own +// anchor for the figures the parity suite compares across listeners: with a body +// plugin in the chain, both directions are buffered and both counts are exact. +func TestBytes_BufferedPairCountsEachDirectionOnItsOwnRow(t *testing.T) { + const reqBody = `{"method":"tools/call","id":1}` + const respBody = `{"result":"ok"}` + + req, _ := http.NewRequest("POST", "http://placeholder/mcp", strings.NewReader(reqBody)) + req.Header.Set("Content-Type", "application/json") + + reqRow, respRow := bytesOf(t, []pipeline.Plugin{&bodyRecorderPlugin{}}, req, + func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, respBody) + }) + + if reqRow.BytesUp != int64(len(reqBody)) || reqRow.BytesDown != 0 { + t.Errorf("request row = ↑%d ↓%d, want ↑%d ↓0", reqRow.BytesUp, reqRow.BytesDown, len(reqBody)) + } + if respRow.BytesDown != int64(len(respBody)) || respRow.BytesUp != 0 { + t.Errorf("response row = ↑%d ↓%d, want ↑0 ↓%d", respRow.BytesUp, respRow.BytesDown, len(respBody)) + } +} + +// TestBytes_PassthroughSSECountsTheFramingToo covers streamPassthrough, the arm an +// event stream takes when no plugin is a StreamingResponder — relayed +// byte-for-byte, never buffered, and reached by no parity fixture, every one of +// which puts a streaming responder in the chain. +// +// The figure is the WIRE length, framing included, which is the whole reason the +// count is taken off the upstream body rather than inside each arm: this arm has +// no frames to add up, and the arm that does sees them with `data: ` and the +// blank-line separators already stripped. See listener/internal/bodycount. +func TestBytes_PassthroughSSECountsTheFramingToo(t *testing.T) { + const wire = "data: {\"type\":\"message_start\"}\n\ndata: [DONE]\n\n" + + req, _ := http.NewRequest("GET", "http://placeholder/events", nil) + + _, respRow := bytesOf(t, []pipeline.Plugin{}, req, + func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "text/event-stream") + fmt.Fprint(w, wire) + if f, ok := w.(http.Flusher); ok { + f.Flush() + } + }) + + if respRow.BytesDown != int64(len(wire)) { + t.Errorf("response row ↓%d, want ↓%d (the framing counts: %q)", respRow.BytesDown, len(wire), wire) + } +} + +// TestBytes_UnbufferedResponseReportsZeroRatherThanAGuess pins the one arm the +// counting wrapper cannot reach: a plain response that no plugin asked to buffer +// and that is not an event stream is relayed by an io.Copy at the bottom of +// handleHTTP, which runs AFTER the response row has been recorded. +// +// Zero is the honest answer and the event contract's own: zero means the listener +// counted nothing, not that the body was empty, and agentop renders it blank. The +// tempting alternative is Content-Length, which is -1 on anything chunked and +// deleted outright on every streaming arm — a guess that would be wrong exactly +// where it mattered. Pinned so that if the ordering here ever changes, it changes +// deliberately. +func TestBytes_UnbufferedResponseReportsZeroRatherThanAGuess(t *testing.T) { + req, _ := http.NewRequest("GET", "http://placeholder/plain", nil) + + _, respRow := bytesOf(t, []pipeline.Plugin{}, req, + func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"never":"buffered"}`) + }) + + if respRow.BytesDown != 0 { + t.Errorf("response row ↓%d, want ↓0 — the relay copies after the row is recorded, so any figure here is counted too late to be this row's", respRow.BytesDown) + } +} diff --git a/core/listener/forwardproxy/server.go b/core/listener/forwardproxy/server.go index abef1c728..c3f7bfbc2 100644 --- a/core/listener/forwardproxy/server.go +++ b/core/listener/forwardproxy/server.go @@ -24,6 +24,7 @@ import ( "errors" "github.com/rossoctl/cortex/core/listener/httpx" + "github.com/rossoctl/cortex/core/listener/internal/bodycount" "github.com/rossoctl/cortex/core/listener/internal/bodyread" "github.com/rossoctl/cortex/core/listener/internal/sessionevent" "github.com/rossoctl/cortex/core/listener/internal/sseframe" @@ -587,6 +588,23 @@ func (s *Server) serveOutbound(w http.ResponseWriter, r *http.Request, tl *tunne pctx.StatusCode = resp.StatusCode pctx.ResponseHeaders = resp.Header.Clone() + // BytesDown for the response row (#1309), counted HERE and nowhere else. Every + // arm below reads this one body — buffered whole, re-framed as SSE, relayed + // chunk by chunk, or copied straight through at the bottom of this function — + // so one wrapper covers all of them, and the arms that go on to replace + // resp.Body with a bytes.Reader over what they already read cannot + // double-count. Counted at the wire and not per arm because the SSE arm's + // frames have their `data: ` framing stripped by then, which is 32 bytes short + // of what extproc reports for the same response: see bodycount. + // + // What the wrapper does NOT reach is the final io.Copy at the bottom of this + // function — the arm a response takes when nothing buffered it and it was not + // streamed. The row is recorded above that copy, so those bytes are counted + // after the figure has been read and the row reports zero, which agentop + // renders blank rather than wrong. reverseproxy's unbuffered relay has the same + // shape for the same reason. + resp.Body = bodycount.Wrap(resp.Body, &pctx.ResponseBytes) + // SkipHosts: bypass response-phase pipeline + recording entirely. // Stream the upstream body straight through to the caller. Falls // out below to the unconditional header copy + io.Copy. @@ -949,6 +967,7 @@ func (s *Server) recordOutboundRequestEvent(tl *tunnelLog, pctx *pipeline.Contex Direction: pipeline.Outbound, Phase: pipeline.SessionRequest, RequestID: pctx.RequestID(), + BytesUp: int64(len(pctx.Body)), MCP: pipeline.SnapshotMCP(pctx.Extensions.MCP), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseRequest), @@ -1176,6 +1195,7 @@ func (s *Server) recordOutboundResponseEvent(pctx *pipeline.Context, statusCode Direction: pipeline.Outbound, Phase: pipeline.SessionResponse, RequestID: pctx.RequestID(), + BytesDown: pctx.ResponseBytes, MCP: pipeline.SnapshotMCP(pctx.Extensions.MCP), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseResponse), @@ -1341,7 +1361,6 @@ func (s *Server) handleStreamingResponse(w http.ResponseWriter, r *http.Request, flusher.Flush() reader := sseframe.NewReader(idleReader(resp.Body, streamReadIdleTimeout), maxBodySize) - bytesWritten := 0 for { frame, err := reader.ReadFrame() if err == io.EOF { @@ -1352,7 +1371,11 @@ func (s *Server) handleStreamingResponse(w http.ResponseWriter, r *http.Request, // already received some frames; the cleanest signal is to // close the connection and log. We can't promote this to // 502 — headers are sent. - slog.Warn("forward-proxy: streaming response read error", "host", r.Host, "error", err, "bytesWritten", bytesWritten) + // bytesRead, not the bytesWritten this used to keep in a local: the + // local was declared after the recording defer and never reached the + // session event, so the figure is now taken off the counting wrapper + // around resp.Body — which is bytes received, not frames relayed. + slog.Warn("forward-proxy: streaming response read error", "host", r.Host, "error", err, "bytesRead", pctx.ResponseBytes) break } @@ -1382,7 +1405,6 @@ func (s *Server) handleStreamingResponse(w http.ResponseWriter, r *http.Request, break } flusher.Flush() - bytesWritten += len(frame) } } diff --git a/core/listener/internal/bodycount/bodycount.go b/core/listener/internal/bodycount/bodycount.go new file mode 100644 index 000000000..cf4b17563 --- /dev/null +++ b/core/listener/internal/bodycount/bodycount.go @@ -0,0 +1,44 @@ +// Package bodycount counts the bytes a listener reads off an upstream response +// body, so every listener reports the same figure for the same response (#1309). +// +// It exists because the proxies read a response body through four different arms +// — buffered whole, re-framed as SSE, relayed chunk by chunk, or copied straight +// through — and three of those have a natural-looking number to count that is NOT +// the body: the SSE arm sees sseframe payloads with the `data: ` prefixes and +// blank-line separators already stripped, and a buffered read stops at the +// listener's cap. extproc has no such ambiguity (Envoy hands it the raw chunks), +// so counting at the one place where the raw bytes enter the proxy is what makes +// the three listeners agree. Wrap once, as early as the body is in hand, and let +// every arm downstream read through the wrapper. +package bodycount + +import "io" + +// Wrap returns rc with the length of every successful read added to *n. The +// pointer is a pipeline.Context field at both call sites: the arms that consume +// the body are several frames below the wrap, and the record site is above it +// again, so there is nowhere to thread a return value through. +// +// Not safe for concurrent reads, which is not a constraint in practice: a +// response body is read by the one goroutine serving its request, and that is the +// goroutine that records the row. +func Wrap(rc io.ReadCloser, n *int64) io.ReadCloser { + return &countingBody{rc: rc, n: n} +} + +type countingBody struct { + rc io.ReadCloser + n *int64 +} + +// Read counts what it hands back, including bytes returned alongside an error: +// io.Reader may do both, and those bytes did arrive. A body that breaks off +// mid-read therefore reports what was received before it broke rather than zero, +// which is the figure the 502 it produces is worth reading next to. +func (c *countingBody) Read(p []byte) (int, error) { + n, err := c.rc.Read(p) + *c.n += int64(n) + return n, err +} + +func (c *countingBody) Close() error { return c.rc.Close() } diff --git a/core/listener/parity/drivers_test.go b/core/listener/parity/drivers_test.go index 301738304..a3cc58ca3 100644 --- a/core/listener/parity/drivers_test.go +++ b/core/listener/parity/drivers_test.go @@ -141,6 +141,28 @@ type fixture struct { // where the fixture controls the request shape, so fixtures opt in // rather than inheriting a comparison they were not written for. expectedUpstream *upstreamSummary + + // expectedBytes, when non-nil, pins BytesUp and BytesDown EXACTLY on every + // listener rather than only comparing them (#1309). + // + // Absolute as well as pairwise because zero is the default and the pairwise + // diff cannot see a gap both legs share: three listeners that all count nothing + // agree perfectly, which is what BYTES looked like before this change and what + // a regression would look like after it. A pointer and not a plain pair of + // int64s for the same reason — zero is a value a fixture may want to pin (a + // body-less GET must report zero, not a length), so "assert nothing" needs to + // be a distinct state. + expectedBytes *bytesSummary +} + +// bytesSummary is a fixture's claim about the two byte counts on one row. +// +// Up and Down are populated per PHASE, not per row: a request row counts what was +// forwarded and a response row what came back, so a fixture asserting both runs +// the same shape twice with different wantPhase values. +type bytesSummary struct { + Up int64 + Down int64 } // contentType returns the fixture's response content-type or a sensible @@ -225,6 +247,24 @@ type observation struct { // did not ask (expectedUpstream) or when nothing reached the upstream at // all — a denial, where nil on every leg is the correct answer. Upstream *upstreamSummary + // BytesUp and BytesDown are the body sizes the event reported (#1309). A + // cross-listener obligation with nothing behind it until now: all three + // listeners count, each from its own machinery — extproc sums the chunks Envoy + // hands it, the proxies count off the upstream body — and the same body has to + // come out the same figure whichever shape the operator deployed. + // + // Compared as VALUES, unlike HasDuration above, and that is the point: a + // duration legitimately differs between two legs of one fixture because it is + // measured, while a body length is counted and must agree exactly. + // + // It earned its keep on the first run. The proxies originally tallied inside + // each response arm, where the SSE arm has sseframe PAYLOADS in hand — `data: ` + // prefixes and blank-line separators already stripped — so on the four-event + // reads-body-sse fixture extproc reported 137 and reverseproxy 105. Both + // numbers are defensible in isolation; only one of them can be what BYTES + // means. See listener/internal/bodycount for where the proxies count now. + BytesUp int64 + BytesDown int64 // Inference is the token report the event carried, nil when it carried none. // // IT IS NOT A RESTATEMENT OF THE COST RECORD, which travels separately in @@ -348,6 +388,8 @@ func observe(t *testing.T, store *session.Store, wantDir pipeline.Direction, wan Phase: ev.Phase.String(), StatusCode: ev.StatusCode, HasDuration: ev.Duration > 0, + BytesUp: ev.BytesUp, + BytesDown: ev.BytesDown, PluginEventJSON: map[string]string{}, } if ev.Identity != nil { diff --git a/core/listener/parity/parity_test.go b/core/listener/parity/parity_test.go index e9086a1ed..524627895 100644 --- a/core/listener/parity/parity_test.go +++ b/core/listener/parity/parity_test.go @@ -222,6 +222,118 @@ func TestParity_ReadsBodySSE(t *testing.T) { assertParity(t, f, pipeline.SessionResponse, inboundListeners) } +// TestBytesParity_CountsAreTheWholeBody pins what the BYTES column says, as an +// absolute figure on every listener (#1309). The pairwise diff is a weak check for +// this one field: zero is the default, so three listeners that all count nothing +// agree — which is exactly the state the column was in before this change. +// +// What each case pins: +// +// - the buffered pair, both directions: the request row reports the body as +// FORWARDED and the response row the body as RECEIVED, each on its own row and +// neither on the other's. A listener that set both on one row would read as a +// round trip that never happened. +// - SSE: the WIRE bytes, so 137 and not the 105 bytes of payload the pipeline +// sees after sseframe strips the framing — and not the trailing chunk either, +// which is all extproc's response buffer holds by the time the row is +// recorded. The two wrong answers available here differ from the right one by +// specific numbers, which is why this is stated absolutely. +// - the body-less GET: zero on both sides, which agentop renders blank. Pinned +// because "nothing was counted" and "the body was empty" have to stay the same +// state as far as the event is concerned — there is no third value for +// "unknown", and a listener inventing a length for a request that carried none +// would be worse than the blank. +func TestBytesParity_CountsAreTheWholeBody(t *testing.T) { + reqBody := []byte(`{"prompt":"hello"}`) + respBody := []byte(`{"reply":"ok"}`) + + // One shape, run four times: each direction × each phase. The counts are + // per-phase, so a fixture cannot assert both halves at once. + // + // The two Record flags are load-bearing and not boilerplate: extproc appends a + // session event only when the pipeline left something worth recording on pctx + // (A2A, an Invocation, or a plugin-public Custom entry — recordInboundSession's + // gate), while the proxies record unconditionally. A spy with ReadsBody alone + // publishes nothing, so the extproc leg produces no row at all and the fixture + // fails on presence rather than on the figures it is here to pin. Giving the + // spy something to publish is the cheapest way to get both legs recording; it + // does not touch the byte counts, which the listener takes off the wire. + buffered := func(name string, dir pipeline.Direction, want *bytesSummary) fixture { + return fixture{ + name: "bytes-buffered-" + name, + direction: dir, + entries: []config.PluginEntry{spyEntry(spyPluginAStreaming, spyConfig{ + ReadsBody: true, + RecordRequestBody: true, + RecordResponseFrames: true, + })}, + method: "POST", + path: "/parity/echo", + reqBody: reqBody, + upstreamStatus: 200, + upstreamBody: respBody, + expectedBytes: want, + } + } + for _, tc := range []struct { + name string + dir pipeline.Direction + listeners []listenerRun + }{ + {"inbound", pipeline.Inbound, inboundListeners}, + {"outbound", pipeline.Outbound, outboundListeners}, + } { + assertParity(t, buffered(tc.name+"-request", tc.dir, + &bytesSummary{Up: int64(len(reqBody))}), pipeline.SessionRequest, tc.listeners) + assertParity(t, buffered(tc.name+"-response", tc.dir, + &bytesSummary{Down: int64(len(respBody))}), pipeline.SessionResponse, tc.listeners) + } + + // Same four events as TestParity_ReadsBodySSE, re-framed here rather than + // shared: that fixture's anchor is the payload the PLUGIN sees, this one's is + // the wire, and the whole point of the case is that the two differ. + var sse strings.Builder + for _, e := range []string{ + `{"type":"message_start"}`, + `{"type":"message_delta","usage":{"output_tokens":7}}`, + `{"type":"message_stop"}`, + `[DONE]`, + } { + sse.WriteString("data: ") + sse.WriteString(e) + sse.WriteString("\n\n") + } + assertParity(t, fixture{ + name: "bytes-sse", + direction: pipeline.Inbound, + entries: []config.PluginEntry{spyEntry(spyPluginAStreaming, spyConfig{ + ReadsBody: true, + RecordResponseFrames: true, + })}, + method: "POST", + path: "/parity/sse", + upstreamStatus: 200, + upstreamBody: []byte(sse.String()), + upstreamContentType: "text/event-stream", + expectedBytes: &bytesSummary{Down: int64(sse.Len())}, + }, pipeline.SessionResponse, inboundListeners) + + // A GET with no body on either side. ReadsBody still set, so the zero is the + // absence of bytes and not the absence of a reader. + assertParity(t, fixture{ + name: "bytes-none", + direction: pipeline.Inbound, + entries: []config.PluginEntry{spyEntry(spyPluginAStreaming, spyConfig{ + ReadsBody: true, + RecordRequestBody: true, + })}, + method: "GET", + path: "/parity/empty", + upstreamStatus: 204, + expectedBytes: &bytesSummary{}, + }, pipeline.SessionRequest, inboundListeners) +} + // TestParity_HeaderOnlyResponse is the shape this suite could not express // until the synthetic-body gate in runExtproc was fixed, and the one it // existed to catch. @@ -786,6 +898,12 @@ func assertParity(t *testing.T, f fixture, wantPhase pipeline.SessionPhase, list if f.expectedInvocations != nil && !reflect.DeepEqual(g.observed.Invocations, f.expectedInvocations) { t.Errorf("fixture %q listener %s: Invocations (exact, ordered)\n got: %s\n want: %s", f.name, g.listener, jsonPretty(g.observed.Invocations), jsonPretty(f.expectedInvocations)) } + if f.expectedBytes != nil { + got := bytesSummary{Up: g.observed.BytesUp, Down: g.observed.BytesDown} + if got != *f.expectedBytes { + t.Errorf("fixture %q listener %s: bytes counted\n got: %+v\n want: %+v", f.name, g.listener, got, *f.expectedBytes) + } + } if f.expectedUpstream != nil { switch { case g.observed.Upstream == nil: @@ -846,6 +964,19 @@ func observationDiff(a, b *observation) string { if a.HasDuration != b.HasDuration { return fmt.Sprintf("HasDuration: %v vs %v", a.HasDuration, b.HasDuration) } + // One listener reporting a body size the next one disagrees with. Values, not + // presence — see the observation fields. The same caveat as Inference below + // applies and is why TestBytesParity_CountsAreTheWholeBody states the figures + // absolutely as well: two listeners that both count nothing agree here, so a gap + // shared by every leg of a fixture passes this check. What it catches is the + // split, which is the likelier failure given that each listener counts from + // different machinery. + if a.BytesUp != b.BytesUp { + return fmt.Sprintf("BytesUp: %d vs %d", a.BytesUp, b.BytesUp) + } + if a.BytesDown != b.BytesDown { + return fmt.Sprintf("BytesDown: %d vs %d", a.BytesDown, b.BytesDown) + } // One listener putting a plugin's rewrite on the wire while the other // forwards the original bytes — the split case of the gap // fixture.expectedUpstream covers absolutely, and the one no session-event diff --git a/core/listener/reverseproxy/bytescount_test.go b/core/listener/reverseproxy/bytescount_test.go new file mode 100644 index 000000000..0b7169d96 --- /dev/null +++ b/core/listener/reverseproxy/bytescount_test.go @@ -0,0 +1,114 @@ +package reverseproxy + +import ( + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" + + "github.com/rossoctl/cortex/core/pipeline" + "github.com/rossoctl/cortex/core/session" +) + +// bytesOf drives one request through a reverse proxy built on the given plugins +// and returns the byte counts the request and response rows reported (#1309). +// +// Returned as the two rows rather than four numbers because the split is the +// contract: a request row counts what was forwarded and nothing else, a response +// row counts what came back and nothing else. +func bytesOf(t *testing.T, plugins []pipeline.Plugin, method, path, reqBody string, upstream http.HandlerFunc) (reqRow, respRow pipeline.SessionEvent) { + t.Helper() + p, err := pipeline.New(plugins) + if err != nil { + t.Fatalf("pipeline.New: %v", err) + } + backend := httptest.NewServer(upstream) + defer backend.Close() + + store := session.New(5*time.Minute, 100, 0) + defer store.Close() + + srv, err := NewServer(pipeline.NewHolder(p), store, backend.URL, nil) + if err != nil { + t.Fatalf("NewServer: %v", err) + } + proxy := httptest.NewServer(srv.Handler()) + defer proxy.Close() + + // A nil reader, not an empty one: an empty strings.Reader is still a body as + // far as net/http is concerned, and the body-less case is one this is here to + // measure. + req, _ := http.NewRequest(method, proxy.URL+path, nil) + if reqBody != "" { + req, _ = http.NewRequest(method, proxy.URL+path, strings.NewReader(reqBody)) + req.Header.Set("Content-Type", "application/json") + } + + resp, err := http.DefaultClient.Do(req) + if err != nil { + t.Fatalf("request failed: %v", err) + } + // Drained, not just closed: the streamed arms count as the relay reads, so an + // undrained body would under-report by whatever was still in flight. + buf := make([]byte, 4096) + for { + if _, err := resp.Body.Read(buf); err != nil { + break + } + } + resp.Body.Close() + + v := store.View(session.DefaultSessionID) + if v == nil || len(v.Events) != 2 { + t.Fatalf("want a request row and a response row, got %+v", v) + } + return v.Events[0], v.Events[1] +} + +// TestBytes_BufferedPairCountsEachDirectionOnItsOwnRow is the reverse proxy's own +// anchor for the figures the parity suite compares across listeners: with a body +// plugin in the chain, both directions are buffered and both counts are exact. +func TestBytes_BufferedPairCountsEachDirectionOnItsOwnRow(t *testing.T) { + const reqBody = `{"jsonrpc":"2.0","method":"message/send"}` + const respBody = `{"jsonrpc":"2.0","result":{}}` + + reqRow, respRow := bytesOf(t, []pipeline.Plugin{&bodyRecorderPlugin{}}, "POST", "/a2a", reqBody, + func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, respBody) + }) + + if reqRow.BytesUp != int64(len(reqBody)) || reqRow.BytesDown != 0 { + t.Errorf("request row = ↑%d ↓%d, want ↑%d ↓0", reqRow.BytesUp, reqRow.BytesDown, len(reqBody)) + } + if respRow.BytesDown != int64(len(respBody)) || respRow.BytesUp != 0 { + t.Errorf("response row = ↑%d ↓%d, want ↑0 ↓%d", respRow.BytesUp, respRow.BytesDown, len(respBody)) + } +} + +// TestBytes_UnbufferedResponseReportsZeroRatherThanAGuess pins the one arm the +// counting wrapper cannot reach: with no plugin asking for the response body, +// httputil.ReverseProxy copies it straight to the client — and it does that AFTER +// modifyResponse returns, which is after the response row has been appended +// inside it. +// +// Zero is the honest answer and the event contract's own: zero means the listener +// counted nothing, not that the body was empty, and agentop renders it blank. The +// tempting alternative is Content-Length, which is -1 on anything chunked and +// deleted outright on the streaming arm — a guess that would be wrong exactly +// where it mattered. Pinned so that if the ordering here ever changes, it changes +// deliberately. forwardproxy's unbuffered relay has the same shape for the same +// reason. +func TestBytes_UnbufferedResponseReportsZeroRatherThanAGuess(t *testing.T) { + _, respRow := bytesOf(t, []pipeline.Plugin{}, "GET", "/plain", "", + func(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"never":"buffered"}`) + }) + + if respRow.BytesDown != 0 { + t.Errorf("response row ↓%d, want ↓0 — ReverseProxy copies after the row is recorded, so any figure here is counted too late to be this row's", respRow.BytesDown) + } +} diff --git a/core/listener/reverseproxy/server.go b/core/listener/reverseproxy/server.go index eb295aa2e..63b65e00e 100644 --- a/core/listener/reverseproxy/server.go +++ b/core/listener/reverseproxy/server.go @@ -19,6 +19,7 @@ import ( "time" "github.com/rossoctl/cortex/core/listener/httpx" + "github.com/rossoctl/cortex/core/listener/internal/bodycount" "github.com/rossoctl/cortex/core/listener/internal/bodyread" "github.com/rossoctl/cortex/core/listener/internal/sessionevent" "github.com/rossoctl/cortex/core/listener/internal/sseframe" @@ -459,6 +460,7 @@ func (s *Server) handleRequest(w http.ResponseWriter, r *http.Request) { Direction: pipeline.Inbound, Phase: pipeline.SessionRequest, RequestID: pctx.RequestID(), + BytesUp: int64(len(pctx.Body)), A2A: pipeline.SnapshotA2A(pctx.Extensions.A2A), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseRequest), @@ -485,6 +487,21 @@ func (s *Server) modifyResponse(resp *http.Response) error { pctx.StatusCode = resp.StatusCode pctx.ResponseHeaders = resp.Header.Clone() + // BytesDown for the response row (#1309), counted HERE and nowhere else — at + // the wire, before either arm below gets at the body. The streaming arm sees + // sseframe payloads with the `data: ` framing already stripped, which is short + // of what extproc reports for the same response, and the buffered arm never + // runs on a streamed one; one wrapper here is the only place both are the same + // bytes. The buffered arm's replacement of resp.Body with a reader over what it + // already read means nothing is counted twice. See bodycount. + // + // What this does NOT reach is the fully-unbuffered relay: ReverseProxy copies + // that body after modifyResponse returns, which is after the response row is + // appended below, so it stays at zero and renders blank rather than wrong. + if resp.Body != nil { + resp.Body = bodycount.Wrap(resp.Body, &pctx.ResponseBytes) + } + // Branch on Content-Type per response. Streaming-aware pipelines on // text/event-stream responses (A2A message/stream, MCP tools/call // result over Streamable HTTP) replace resp.Body with a streaming @@ -618,6 +635,7 @@ func (s *Server) modifyResponse(resp *http.Response) error { Direction: pipeline.Inbound, Phase: pipeline.SessionResponse, RequestID: pctx.RequestID(), + BytesDown: pctx.ResponseBytes, A2A: pipeline.SnapshotA2A(pctx.Extensions.A2A), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseResponse), @@ -808,6 +826,7 @@ func (s *Server) recordInboundResponseEvent(pctx *pipeline.Context, statusCode i Direction: pipeline.Inbound, Phase: pipeline.SessionResponse, RequestID: pctx.RequestID(), + BytesDown: pctx.ResponseBytes, A2A: pipeline.SnapshotA2A(pctx.Extensions.A2A), Inference: pipeline.SnapshotInference(pctx.Extensions.Inference), Invocations: pipeline.SnapshotInvocations(pctx.Extensions.Invocations, pipeline.InvocationPhaseResponse), diff --git a/core/pipeline/context.go b/core/pipeline/context.go index b1ef82b95..b7072d4a3 100644 --- a/core/pipeline/context.go +++ b/core/pipeline/context.go @@ -215,6 +215,32 @@ type Context struct { ResponseHeaders http.Header ResponseBody []byte + // ResponseBytes is how many response-body bytes the listener observed on the + // wire, summed across every chunk. SessionEvent.BytesDown is recorded from it, + // and it is authoritative over len(ResponseBody) — which is a BUFFER, not a + // tally, and is wrong on three separate paths: extproc REPLACES it per chunk + // on the SSE arm (so it holds only the trailing chunk), extproc truncates it + // at the listener's cap, and neither proxy's streaming relay fills it at all. + // Reading the buffer's length at a record site would report a confident wrong + // number on exactly the streamed inference traffic the count is wanted for. + // + // ON THE WIRE is the part that makes the three listeners agree, and it is not + // the obvious reading: the proxies' SSE arms handle sseframe PAYLOADS, with + // the `data: ` prefixes and blank-line separators already stripped, while + // extproc counts the raw chunks Envoy hands it. Counting what each arm has in + // hand put the two 32 bytes apart on a four-event fixture. So the proxies + // count at one choke point per listener (see listener/internal/bodycount) and + // extproc sums its chunks, and all three report the body the destination sent. + // + // Zero means the listener counted nothing — no plugin asked for the body, so + // it was relayed without passing through a counted path — and NOT that the + // body was empty. Consumers render that as unknown rather than as 0 B; see + // SessionEvent.BytesDown. + // + // Written only by the listener, on the one goroutine serving the request, + // before the response row is recorded. Plugins do not write to it. + ResponseBytes int64 + Extensions Extensions // responseDelivered says the response has already reached the client, so a refusal diff --git a/core/pipeline/session.go b/core/pipeline/session.go index 1be05a4a5..95a13cf35 100644 --- a/core/pipeline/session.go +++ b/core/pipeline/session.go @@ -191,11 +191,39 @@ type SessionEvent struct { // carried no request at all. Tunnel bool - // BytesUp and BytesDown are how many bytes an opaque tunnel carried each way: - // up is client to destination, down is destination to client. Set only on a - // tunnel's close row, which is when the counts are known. Zero — and absent on - // the wire — everywhere else, including a bridged tunnel's close, whose bytes - // were TLS the bridge terminated rather than counted. + // BytesUp and BytesDown are how many bytes went each way: up is client to + // destination, down is destination to client. A row carries whichever side it + // is in a position to know, so the two never collide on one row: + // + // request row BytesUp — the request body as FORWARDED, which is + // len(pctx.Body) at record time and so + // reflects any plugin that rewrote it + // response row BytesDown — the response body the listener observed + // FROM THE DESTINATION, read from + // pctx.ResponseBytes + // tunnel close row both — the opaque bytes the tunnel carried + // + // The two sides differ on where the pipeline sits, and not by preference: the + // request figure is read at record time, after any rewrite, because that is + // when pctx.Body is what will be forwarded; the response figure is summed as + // the bytes arrive, because on a streamed response there is no later moment + // at which the whole body exists. A response mutator's rewrite is therefore + // not reflected — extproc records the row before it even emits the + // replacement — so BytesDown is what the destination sent. + // + // ZERO MEANS NOT COUNTED, NOT ZERO BYTES, and the field is absent from the + // wire when zero. Body buffering is gated on a pipeline capability in all + // three listeners, so a request no plugin wanted the body of is relayed + // without ever being measured — and a body-less GET is indistinguishable from + // it, which is why neither may render as a real 0 B. A bridged tunnel's close + // reports zero for the same reason: its bytes were TLS the bridge terminated + // rather than counted, and its decrypted inner requests carry their own + // figures. Consumers render zero as blank (agentop's bytesCell). + // + // Headers are deliberately not counted. The question these answer is how big + // a payload was — a prompt that overran a model's context window is the case + // they exist for — and a kilobyte of constant header noise on every row + // obscures it. They are therefore NOT a wire-cost figure. BytesUp int64 BytesDown int64 @@ -431,8 +459,10 @@ type sessionEventWire struct { // key, and a new agentop against an old proxy sees "" and renders exactly what it // renders today. TunnelReason TunnelReason `json:"tunnelReason,omitempty"` - // omitempty for the same skew reason, and because only a tunnel's close row - // has a count to report. + // omitempty for the same skew reason, and because a row reports only the side + // it knows: a request row has no BytesDown and a response row no BytesUp. The + // absent key and a counted zero are therefore the same wire bytes, which is + // intended — both mean "no figure", per SessionEvent.BytesUp. BytesUp int64 `json:"bytesUp,omitempty"` BytesDown int64 `json:"bytesDown,omitempty"` // omitempty for the same skew reason as TunnelReason above: an old agentop From 6ab1629f71b794a660afe52aab18f5600c73ce77 Mon Sep 17 00:00:00 2001 From: cwiklik Date: Fri, 9 Oct 2026 09:44:51 -0400 Subject: [PATCH 2/2] fix: Regenerate the README demo asset for the extra default column MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit `TestCommittedAssetIsCurrent` caught this, and the whole diff is one token: the footer hint goes from `→ 3 more columns` to `→ 4 more columns`. That is BYTES becoming default-on. It keeps `keep: keepLow`, so at the demo's terminal width it is the first column dropped — it does not appear in the animation, it just moves the off-screen count up by one. The staleness check compares content, not size, which is why the two byte counts in its failure message were identical at 71121. Assisted-By: Claude (Anthropic AI) Signed-off-by: cwiklik --- docs/assets/cortex-demo.svg | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/assets/cortex-demo.svg b/docs/assets/cortex-demo.svg index bcc4226d3..8259cee16 100644 --- a/docs/assets/cortex-demo.svg +++ b/docs/assets/cortex-demo.svg @@ -384,7 +384,7 @@ clipPath rect{width:100000px} ● connected0.0 events/sec feedback: https://github.com/rossoctl/cortex/issues/new/choose -… [s] hide passthru/skip [p] pause [/] filter [esc] back · → 3 more columns ([c] to choose) [?] keys [q] quit +… [s] hide passthru/skip [p] pause [/] filter [esc] back · → 4 more columns ([c] to choose) [?] keys [q] quit agentop · fix the retry handler (api-7f3c) · event