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
83 changes: 83 additions & 0 deletions docs/adr/0013-peak-pricing-tier-resolved-at-record-time.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,83 @@
# Peak/off-peak pricing tier is resolved once, at usage-record time

Status: accepted.

DeepSeek switched to time-of-day (peak/off-peak) API pricing on 2026-09-10 — peak hours cost
2x off-peak, currently 01:00–04:00 and 06:00–10:00 UTC, Monday–Friday. `store.Offering` held
exactly one flat price triple, so before this change roughly half of all DeepSeek spend was
mispriced no matter what number was entered.

## Decision

Each offering carries an optional peak-tier price triple (`price_in_per_1m_peak`/
`price_out_per_1m_peak`/`price_cached_in_per_1m_peak`, per-field nil = falls back to the base
rate — the same convention `PriceCachedInPer1M` already used); the schedule itself
(`internal/pricing.Windows`) lives on the **provider**, not the offering, since a peak window is
a fact about the provider's billing policy shared by every model it serves, not a per-model
property. `router/routing.go`'s `offeringChain` copies both triples plus the parsed schedule
onto `Backend`/`ResolvedBackend` at its single existing copy point.

**The tier itself is picked exactly once, in `computeCostNative` (`router/usage.go`), at the
same instant `recordExternalUsage` stamps the event's timestamp** — not earlier, at
route-selection/request-start time. Concretely: `now := time.Now()` is captured once in
`recordExternalUsage` and used for both `ev.TS` and the tier lookup; the tier is recorded on the
event (`usage_events.price_tier`).

## Why not resolve the tier at request start

Nothing upstream of cost computation uses price for a routing decision — `select.go` sorts
candidate offerings by `priority` only, so resolving the tier early buys nothing operationally.
Resolving it early would actively create a correctness problem: a request whose response
completes after a tier boundary would then have a stored cost that *disagrees with the tier
derivable from the event's own stored timestamp* — and `compressor_summary_handlers.go`'s
historical savings estimators (`estimateRemoteCacheDiscountSaved`/
`estimateRemoteCompressionSaved`) re-price past events from exactly that timestamp. Pinning the
tier at record time means the ledger is self-consistent by construction: one function
(`computeCostNative`), one time input, shared by both the live-billing path and the historical
estimator (which prefers the event's own recorded `price_tier`, falling back to evaluating the
schedule against `ev.TS` only for rows written before this migration).

**Consequence, disclosed rather than hidden:** a request that starts in one tier and whose
response completes after a boundary is billed entirely at the tier in force at completion, not
a split or the starting tier. A 2x rate delta on a long streaming completion is exactly the
boundary-crossing case where this matters most — a reproducible, auditable rule (visible via the
recorded `price_tier`) beats guessing DeepSeek's own internal accounting for the same case.

## Other decisions folded in here

- **Half-open `[start, end)` window boundaries** — a request landing exactly on the window's end
time is off-peak, pinned by `internal/pricing`'s own tests.
- **No "has peak pricing" boolean.** A flag can disagree with the data it describes; the
per-field-nil convention is one source of truth.
- **Peak price validated for sign only, never `peak >= base`.** Whether a provider's peak tier
costs more or less than off-peak is the provider's own policy, not this app's to enforce.
- **`peak_active_now` is computed server-side only**, both on the repeatedly-polled
`GET /api/v1/providers` list (via `internal/providers.Service`'s existing injectable clock) and
on the one-off create/update echo (`httpapi.peakActiveNow`, plain `time.Now()`) — the window
math is never ported to TypeScript.

## Correction, 2026-09-13: timezone support

The original decision above said "UTC only, no timezone field — a timezone field the code
doesn't honor is worse than no field." That was right for DeepSeek's own published schedule (it
really is UTC), but wrong as a general policy: most providers publish peak hours in one local
zone ("9am–5pm Pacific"), and requiring the operator to hand-convert to UTC is exactly the trap
the original reasoning was trying to avoid — a fixed UTC offset entered once goes silently wrong
by an hour across the next DST transition, twice a year, with nothing to catch it.

**Fixed properly, not worked around:** `Windows` gained an optional `TZ` field (IANA zone name,
`""` = UTC, backward-compatible with DeepSeek's already-stored no-`tz` schedule). `Active`
converts the evaluation instant into that zone via `time.Time.In` — Go's real IANA tzdata, not a
stored offset — before any day-of-week/time-of-day comparison, so the same schedule is correct
in January and July without ever being re-entered, and a schedule quoted in a zone far from UTC
correctly shifts which calendar day it evaluates against (e.g. "Monday 00:00–04:00 Asia/Tokyo" is
active during Sunday afternoon UTC, because it's already Monday in Tokyo). `time.LoadLocation` is
memoized (`internal/pricing`'s package-level `locCache`) since `Active` runs on the remote-request
hot path and the stdlib doesn't cache zoneinfo parsing itself. No new migration — `peak_windows`
was already a JSON `TEXT` column; `tz` is just a new key in that same blob.

Frontend: `ProviderKeys`' `PeakWindowsEditor` gained a labeled Timezone field (a curated
common-zones dropdown + a "Custom…" free-text IANA name escape hatch) separate from the windows
JSON textarea, composing/splitting the two into the one stored JSON string — the timezone is a
single flat value worth a real input, unlike the windows list itself (still JSON; still no
bespoke hour-range picker, per the original reasoning, which stands).
5 changes: 5 additions & 0 deletions go/cmd/forge-compress/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -113,6 +113,11 @@ func loadConfig() (config, error) {
} else {
c.Compress.ByteThreshold = v
}
if v, err := intEnv("FORGE_COMPRESS_BATCH_SIZE", c.Compress.BatchSize); err != nil {
return config{}, err
} else {
c.Compress.BatchSize = v
}
if v, err := intEnv("FORGE_COMPRESS_FAILOPEN_BUDGET_MS", c.FailOpenBudgetMS); err != nil {
return config{}, err
} else {
Expand Down
58 changes: 54 additions & 4 deletions go/cmd/forge-compress/messages.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,18 @@ import (
// succeed. Returns the real ModernBERT-tokenizer token counts summed
// across every message actually run through the engine (0 for either if
// nothing in the request was tokenizable/compressible), for the caller to
// fold into the tokens_saved metric.
func compressMessages(engine *compress.Engine, body map[string]any, budget time.Duration) (originalTokens, compressedTokens int64, failOpenTimeout, failOpenError int64) {
// fold into the tokens_saved metric, plus outcomeSize: a per-message count
// of "outcome:size_tier" composite labels (see messageOutcomeSize) for the
// caller to fold into messagesByOutcomeSize — added 2026-09-11 so the
// bimodal "compression barely matters on chatty messages, matters enormously
// on huge ones" shape (the operator's own early-testing observation) is
// directly visible in production telemetry instead of inferred from a mean.
func compressMessages(engine *compress.Engine, body map[string]any, budget time.Duration) (originalTokens, compressedTokens int64, failOpenTimeout, failOpenError int64, outcomeSize map[string]int64) {
messagesRaw, ok := body["messages"].([]any)
if !ok {
return 0, 0, 0, 0
return 0, 0, 0, 0, nil
}
outcomeSize = make(map[string]int64)
for _, mRaw := range messagesRaw {
msg, ok := mRaw.(map[string]any)
if !ok {
Expand All @@ -43,8 +49,52 @@ func compressMessages(engine *compress.Engine, body map[string]any, budget time.
case "error":
failOpenError++
}
outcomeSize[messageOutcomeSize(res, reason, len(content))]++
}
return originalTokens, compressedTokens, failOpenTimeout, failOpenError, outcomeSize
}

// messageOutcomeSize classifies one message into a composite
// "outcome:size_tier" label value.
//
// outcome distinguishes gated_passthrough (below Engine.Config's
// MinWords/ByteThreshold — the engine skips tokenization entirely, so
// compress.Result.OriginalTokens stays 0, per that field's own doc comment)
// from alldrop_passthrough (tokenization and scoring both ran, every word
// scored below threshold — OriginalTokens is real/nonzero) using that
// existing documented invariant, rather than adding a new field to
// compress.Result just for this.
func messageOutcomeSize(res compress.Result, reason string, contentBytes int) string {
var outcome string
switch {
case reason == "timeout":
outcome = "failopen_timeout"
case reason == "error":
outcome = "failopen_error"
case res.Passthrough && res.OriginalTokens == 0:
outcome = "gated_passthrough"
case res.Passthrough:
outcome = "alldrop_passthrough"
default:
outcome = "compressed"
}
return outcome + ":" + sizeTier(contentBytes)
}

// sizeTier buckets a message's raw content size — deliberately coarse (4
// tiers) so the resulting label cardinality stays small regardless of
// traffic volume.
func sizeTier(bytes int) string {
switch {
case bytes < 2*1024:
return "small"
case bytes < 20*1024:
return "medium"
case bytes < 200*1024:
return "large"
default:
return "huge"
}
return originalTokens, compressedTokens, failOpenTimeout, failOpenError
}

// compressOne runs engine.Compress with a wall-clock fail-open budget and
Expand Down
105 changes: 103 additions & 2 deletions go/cmd/forge-compress/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,8 @@ import (
"os"
"sort"
"sync"

"github.com/jsaigou/the-forge/internal/statutil"
)

// selfRSSBytes reads this process's resident set from /proc/self/statm
Expand Down Expand Up @@ -56,18 +58,41 @@ type metrics struct {
ttfb histogram
latency histogram
overhead histogram
// overheadRing backs compress_overhead_ms_p50/p90/p99 — see sampleRing's
// doc comment for why overhead specifically gets a percentile view and
// ttfb/latency don't (yet): overhead is the one number this session's
// investigation found was actively misleading as a mean (dominated by a
// small share of huge messages, hiding that most real traffic barely
// pays the tax).
overheadRing *sampleRing

failOpenTimeout counter
failOpenError counter

requestsByProvider labelCounter
requestsByModel labelCounter
// messagesByOutcomeSize is keyed by a composite "outcome:size_tier"
// label value (e.g. "compressed:huge") rather than two independent
// label dimensions — this repo's label-sample storage
// (internal/store's compressor_label_samples) is a flat
// (label_key, label_value, metric) shape with one dimension per row, so
// a composite value is how a second dimension rides along without a
// schema change. outcome ∈ {compressed, gated_passthrough,
// alldrop_passthrough, failopen_timeout, failopen_error}; size_tier ∈
// {small, medium, large, huge} — see messageOutcomeSize in messages.go.
// Added 2026-09-11 to answer whether compression's real value is
// concentrated in a few huge messages (the operator's own early-testing
// finding) or spread evenly — something the prior mean-only metrics
// couldn't show.
messagesByOutcomeSize labelCounter
}

func newMetrics() *metrics {
return &metrics{
requestsByProvider: newLabelCounter(),
requestsByModel: newLabelCounter(),
overheadRing: newSampleRing(overheadRingCapacity),
requestsByProvider: newLabelCounter(),
requestsByModel: newLabelCounter(),
messagesByOutcomeSize: newLabelCounter(),
}
}

Expand Down Expand Up @@ -106,9 +131,11 @@ func (m *metrics) WriteTo(w io.Writer) (int64, error) {
writeHistogram(write, "compress_ttfb_ms", &m.ttfb)
writeHistogram(write, "compress_latency_ms", &m.latency)
writeHistogram(write, "compress_overhead_ms", &m.overhead)
writePercentiles(write, "compress_overhead_ms", m.overheadRing)

writeLabelCounter(write, "compress_requests_by_provider", "provider", &m.requestsByProvider)
writeLabelCounter(write, "compress_requests_by_model", "model", &m.requestsByModel)
writeLabelCounter(write, "compress_messages_total", "outcome_size", &m.messagesByOutcomeSize)

return n, nil
}
Expand Down Expand Up @@ -179,6 +206,80 @@ func (h *histogram) snapshot() (count int64, sum, min, max float64) {
return h.count, h.sum, h.min, h.max
}

const (
// overheadRingCapacity mirrors the collector's existing 120-sample
// sparkline-ring pattern (internal/collector/run.go's rings field) —
// this process has no other precedent for bounding an otherwise
// unbounded-lifetime sample set.
overheadRingCapacity = 256
// percentileMinSamples is this repo's established floor for trusting a
// percentile computed from a raw sample set — see
// internal/httpapi/cost_handlers.go's activeSingleSlotWallW gate and
// compressor_summary_handlers.go's prefillObservedMinSamples, both
// named "10" for the same reason: a couple of noisy early observations
// shouldn't produce a misleadingly-precise-looking figure.
percentileMinSamples = 10
)

// sampleRing is a fixed-capacity, thread-safe ring buffer of recent
// float64 samples. This binary's histogram accumulators are lifetime-since-
// process-start (never reset — see histogram's doc comment), so a plain
// growing []float64 isn't safe for a long-running process; a bounded ring
// gives "percentile of recent traffic" instead, which is what actually
// answers "is this request typical" during an incident.
type sampleRing struct {
mu sync.Mutex
buf []float64
next int
full bool
}

func newSampleRing(capacity int) *sampleRing {
return &sampleRing{buf: make([]float64, capacity)}
}

func (r *sampleRing) add(v float64) {
r.mu.Lock()
defer r.mu.Unlock()
r.buf[r.next] = v
r.next++
if r.next == len(r.buf) {
r.next = 0
r.full = true
}
}

// snapshot returns a copy of the samples currently held. Order doesn't
// matter — statutil.Percentile sorts its own copy.
func (r *sampleRing) snapshot() []float64 {
r.mu.Lock()
defer r.mu.Unlock()
if r.full {
out := make([]float64, len(r.buf))
copy(out, r.buf)
return out
}
out := make([]float64, r.next)
copy(out, r.buf[:r.next])
return out
}

// writePercentiles emits p50/p90/p99 for r under the "name_pNN" series
// names, below percentileMinSamples samples emits nothing at all — a
// missing series is the honest signal, never a percentile computed from too
// few points to mean anything. Stored/read as a latest-snapshot gauge (like
// histogram's own min/max), not summed or averaged across a window — same
// invariant documented at internal/store/store.go's CompressorSavingsSampleRow.
func writePercentiles(write func(string, ...any), name string, r *sampleRing) {
vals := r.snapshot()
if len(vals) < percentileMinSamples {
return
}
write("%s_p50 %g\n", name, statutil.Percentile(vals, 50))
write("%s_p90 %g\n", name, statutil.Percentile(vals, 90))
write("%s_p99 %g\n", name, statutil.Percentile(vals, 99))
}

// labelCounter is a set of independent counters keyed by one label value
// (e.g. provider name, model name).
type labelCounter struct {
Expand Down
9 changes: 7 additions & 2 deletions go/cmd/forge-compress/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -168,8 +168,10 @@ func (s *server) handleChatCompletions(w http.ResponseWriter, r *http.Request) {
}
compressStart := time.Now()
budget := time.Duration(s.cfg.FailOpenBudgetMS) * time.Millisecond
originalTokens, compressedTokens, foTimeout, foError := compressMessages(s.engine, body, budget)
s.metrics.overhead.observe(msSince(compressStart))
originalTokens, compressedTokens, foTimeout, foError, outcomeSize := compressMessages(s.engine, body, budget)
overheadMs := msSince(compressStart)
s.metrics.overhead.observe(overheadMs)
s.metrics.overheadRing.add(overheadMs)
s.metrics.tokensInput.add(originalTokens)
if originalTokens > compressedTokens {
s.metrics.tokensSaved.add(originalTokens - compressedTokens)
Expand All @@ -180,6 +182,9 @@ func (s *server) handleChatCompletions(w http.ResponseWriter, r *http.Request) {
if foError > 0 {
s.metrics.failOpenError.add(foError)
}
for label, count := range outcomeSize {
s.metrics.messagesByOutcomeSize.add(label, count)
}
if reencoded, err := json.Marshal(body); err == nil {
mutatedBody = reencoded
}
Expand Down
12 changes: 12 additions & 0 deletions go/cmd/forge-compress/server_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,18 @@ func (fakeScorer) Score(inputIDs, _ []int64) ([]float32, error) {
return scores, nil
}

func (f fakeScorer) ScoreBatch(inputIDs, attentionMask [][]int64) ([][]float32, error) {
out := make([][]float32, len(inputIDs))
for i := range inputIDs {
s, err := f.Score(inputIDs[i], attentionMask[i])
if err != nil {
return nil, err
}
out[i] = s
}
return out, nil
}

func testEngine() *compress.Engine {
return &compress.Engine{
Tokenizer: fakeTokenizer{},
Expand Down
17 changes: 17 additions & 0 deletions go/internal/collector/compressor.go
Original file line number Diff line number Diff line change
Expand Up @@ -61,12 +61,29 @@ type CompressorSample struct {
OverheadSumMsDelta float64
OverheadMinMsSinceStart *float64
OverheadMaxMsSinceStart *float64
// OverheadP50/P90/P99MsRecent are percentiles of the proxy's own recent
// (bounded ring, not lifetime) overhead samples — nil below the 10-sample
// floor (cmd/forge-compress/metrics.go's percentileMinSamples), same
// null-not-zero convention as the Min/Max gauges above. Unlike those,
// "recent" here does NOT mean "since process start" — it's the last
// ~256 requests, since the mean alone was found (2026-09-11) to hide a
// bimodal shape: most messages barely pay the compression tax, a few
// huge ones pay a lot.
OverheadP50MsRecent *float64
OverheadP90MsRecent *float64
OverheadP99MsRecent *float64

// RequestsByProviderDelta / RequestsByModelDelta are request COUNTS per
// label value, not token counts — Compressor's compressor_requests_by_{
// provider,model} metrics carry no token dimension.
RequestsByProviderDelta map[string]int64
RequestsByModelDelta map[string]int64
// MessagesByOutcomeSizeDelta is per-MESSAGE (not per-request) counts
// keyed by a composite "outcome:size_tier" label value (e.g.
// "compressed:huge") — see cmd/forge-compress/messages.go's
// messageOutcomeSize. Added 2026-09-11 to make compression's real
// value visible by content-size tier instead of only as a blended mean.
MessagesByOutcomeSizeDelta map[string]int64

// Provider cache metrics (scraped from compress_cache_read_tokens_total,
// compress_uncached_input_tokens_total, etc. — labelled by provider).
Expand Down
Loading
Loading