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
60 changes: 52 additions & 8 deletions asap-precompute-go/precompute.go
Original file line number Diff line number Diff line change
Expand Up @@ -143,6 +143,20 @@ var (
ErrSketchTypeMismatch = errors.New("precompute: envelope sketch_type does not match config")
)

// LatencyObserver is the host-neutral hook that the runtime invokes
// for every Observe call with the wall-clock duration spent inside
// Observe (matchers, sketch factory, observer dispatch, window
// admission). Adapters wire this to a Prometheus / OTel histogram so
// the deployed shim publishes per-observation latency continuously,
// not just under `testing.B` (Phase 2.11B gap #3 — closes the
// deployment-level confirmation of ADR-0002 §"Performance contract").
//
// The hook is invoked exactly once per Observe call regardless of
// outcome (success, ErrSeriesCapExceeded, ErrLateData, matcher miss).
// Implementations must be cheap — the hook runs on the hot path. Nil
// observers (the default) cost a single nil-check.
type LatencyObserver func(d time.Duration)

// Precompute is the host-neutral runtime described in
// design-doc §6.2. One Precompute instance owns one sketch type
// (see config.SketchType); a deployment with multiple sketch types
Expand Down Expand Up @@ -184,6 +198,11 @@ type Precompute interface {
UpdateConfig(cs *PrecomputeConfigSet)
// Stats returns the live counters; safe to call concurrently.
Stats() *PrecomputeStats
// SetLatencyObserver installs (or replaces) the per-Observe
// latency hook. Pass nil to disable. Safe to call concurrently
// with Observe; the runtime stores the function pointer atomically.
// See LatencyObserver godoc for semantics.
SetLatencyObserver(fn LatencyObserver)
// Shutdown flushes any in-progress state; intended for the
// shim's Shutdown path to run a final Tick before returning.
Shutdown(ctx context.Context) error
Expand All @@ -199,14 +218,15 @@ type Precompute interface {
//
// No global mutex around the Precompute itself.
type precompute struct {
cfg atomic.Pointer[PrecomputeConfig]
sketchFactory SketchFactory
observer SketchObserver
window *windowState
snapshotCache *SnapshotCache
stats *PrecomputeStats
sketchType SketchType
closed atomic.Bool
cfg atomic.Pointer[PrecomputeConfig]
sketchFactory SketchFactory
observer SketchObserver
window *windowState
snapshotCache *SnapshotCache
stats *PrecomputeStats
sketchType SketchType
closed atomic.Bool
latencyObserver atomic.Pointer[LatencyObserver]
}

// New constructs a Precompute given an initial config, a sketch
Expand Down Expand Up @@ -239,6 +259,14 @@ func (p *precompute) activeConfig() *PrecomputeConfig {

// Observe implements Precompute.Observe.
func (p *precompute) Observe(obs *Observation) error {
// Time the entire Observe path including early-exit branches —
// the deployed-stack consumer (a Prom histogram on each shim)
// wants the same envelope `testing.B` measures, not just the
// "happy path admit" subset. Closes Phase 2.11B gap #3.
if fn := p.latencyObserver.Load(); fn != nil && *fn != nil {
start := time.Now()
defer func() { (*fn)(time.Since(start)) }()
}
if p.closed.Load() {
return errors.New("precompute: instance is closed")
}
Expand Down Expand Up @@ -472,6 +500,22 @@ func (p *precompute) Stats() *PrecomputeStats {
return p.stats
}

// SetLatencyObserver implements Precompute.SetLatencyObserver. The
// hook is stored in an atomic pointer so concurrent Observe calls
// see a coherent snapshot without locking; replacing the hook never
// races with the Observe deferred-call.
func (p *precompute) SetLatencyObserver(fn LatencyObserver) {
if fn == nil {
// Storing a nil-valued LatencyObserver pointer is harmless
// (Observe nil-checks the dereferenced function), but storing
// a nil pointer makes the hot-path Load return nil so the
// caller skips the deferred-call entirely. Cheaper.
p.latencyObserver.Store(nil)
return
}
p.latencyObserver.Store(&fn)
}

// Shutdown implements Precompute.Shutdown.
func (p *precompute) Shutdown(ctx context.Context) error {
if !p.closed.CompareAndSwap(false, true) {
Expand Down
71 changes: 71 additions & 0 deletions asap-precompute-go/precompute_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -450,3 +450,74 @@ func TestSketchSubtraitAssertion(t *testing.T) {
}
_ = cs.EstimateCardinality()
}

// TestPrecompute_LatencyObserver verifies that the per-Observe
// latency hook fires exactly once per Observe call (regardless of
// outcome) and never on Observe-bypass paths like Tick. Closes the
// runtime side of Phase 2.11B gap #3.
func TestPrecompute_LatencyObserver(t *testing.T) {
t.Parallel()
cfg := &PrecomputeConfig{
AggID: 7,
SketchType: SketchTypeDDSketch,
Mode: Tumbling,
Window: WindowSpec{Size: 10 * time.Second},
}
p := New(cfg, newFakeFactory(), &fakeObserver{})

var (
samples int
total time.Duration
)
p.SetLatencyObserver(func(d time.Duration) {
samples++
total += d
})

// 5 normal observations.
for i := 0; i < 5; i++ {
if err := p.Observe(&Observation{
TimestampMs: uint64(1_000 + i*100),
Metric: "m",
Labels: []KeyValue{{Key: "k", Value: "v"}},
Value: FloatValue(float64(i + 1)),
}); err != nil {
t.Fatalf("observe %d: %v", i, err)
}
}
// Tick is intentionally NOT counted by the per-observation gate
// (ADR-0002 §"Performance contract" pins per-Observe latency,
// not per-flush). The hook should NOT fire here.
_ = p.Tick(10_000)
// One more Observe after the window rotates — should still fire
// the hook even though the call ends up touching a fresh sketch.
if err := p.Observe(&Observation{
TimestampMs: 11_000,
Metric: "m",
Labels: []KeyValue{{Key: "k", Value: "v"}},
Value: FloatValue(10),
}); err != nil {
t.Fatalf("observe post-tick: %v", err)
}

if samples != 6 {
t.Errorf("LatencyObserver samples: want 6, got %d", samples)
}
if total <= 0 {
t.Errorf("LatencyObserver total duration: want >0, got %v", total)
}

// Disable the hook; subsequent Observe must not call back.
p.SetLatencyObserver(nil)
if err := p.Observe(&Observation{
TimestampMs: 12_000,
Metric: "m",
Labels: []KeyValue{{Key: "k", Value: "v"}},
Value: FloatValue(11),
}); err != nil {
t.Fatalf("observe after disable: %v", err)
}
if samples != 6 {
t.Errorf("LatencyObserver samples after disable: want 6, got %d", samples)
}
}
22 changes: 22 additions & 0 deletions deploy/docker-compose/baseline-b3-delta.yml
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,30 @@
# AGENT_CONFIG=sketchcol-agent-b3-delta.yaml \
# docker compose -f base.yml -f agents-N1.yml \
# -f baseline-b3-delta.yml up -d
#
# Workload profile (see docs/phase-2-perf-deployment.md "Gaps" #6):
# the legacy 1k cardinality × 10 Hz default ran the agent at ~0.25%
# of one core, three orders of magnitude below saturation. At that
# load CPU diffs sit inside measurement noise so a 10% shim
# regression is invisible. The override below pins B3 to a
# saturating profile — 1e5 cardinality × 100 Hz — chosen to keep the
# agent in the 50–70% one-core band where CPU deltas discriminate.
# Each value can be overridden from the host shell (e.g.
# `EXPORTER_FREQ_HZ=50 docker compose ...`); the override here only
# changes the default fallback. Phase 2.11B gap #4.

services:
controller:
labels:
asap.baseline: "b3-delta"

fake-exporter:
environment:
# 1e5 cardinality × 100 Hz ≈ 1e7 obs/s into the agent's input
# ring. Empirically lands the agent at ~50–70% of one core on
# the Threadripper hardware described in docs/phase-2-perf-
# deployment.md "Hardware". Tune downward if the host caps out
# before the agent does — agent CPU is the gate metric for the
# shim regression check, not host CPU.
EXPORTER_CARDINALITY: "${EXPORTER_CARDINALITY:-100000}"
EXPORTER_FREQ_HZ: "${EXPORTER_FREQ_HZ:-100}"
42 changes: 33 additions & 9 deletions deploy/scripts/measure-baseline.py
Original file line number Diff line number Diff line change
Expand Up @@ -178,8 +178,19 @@ def docker_stats_with_bytes_rate(


# Prometheus query templates. `w` is the rate window (e.g. "1m").
# Agent (v0.141) uses `_total` suffix on counters; gateway
# (v0.108) doesn't — hence the duplicated-looking queries.
#
# Both agent and gateway are on otelcol v0.141 today (the gateway was
# bumped along with the agent during the Phase 2 shim PRs); both use
# the `_total` suffix on counters and the `_bytes` suffix on the
# process-RSS gauge. The original measure-baseline.py was written
# against a v0.108 gateway and v0.141 agent — that asymmetry no
# longer holds, so the gateway queries now mirror the agent shape.
# Closes Phase 2.11B gap #1.
#
# Each gateway query is wrapped in `or` against the v0.108 (no-suffix)
# variant so this script keeps producing rows when run against a
# legacy gateway image (e.g. someone replaying an old worktree).
# When neither variant exists the result is NaN, same as before.
QUERIES: dict[str, str] = {
# ── Agent tier (source) ─────────────────────────────────────
"agent_cpu_cores": (
Expand Down Expand Up @@ -221,20 +232,33 @@ def docker_stats_with_bytes_rate(
")"
),
# ── Gateway tier (destination-1) ────────────────────────────
# Gateway is v0.108, no `_total` suffix on process counters.
"gateway_cpu_cores": "rate(otelcol_process_cpu_seconds{{job=\"gateway\"}}[{w}])",
"gateway_rss_mib": "otelcol_process_memory_rss{{job=\"gateway\"}} / 1024 / 1024",
# Gateway v0.108 metric names drop the `_total` suffix that
# v0.141 adds — use the no-suffix variant here.
# v0.141 names with `_total` suffix; legacy v0.108 names appended
# via PromQL `or` so old worktrees keep producing data.
"gateway_cpu_cores": (
"rate(otelcol_process_cpu_seconds_total{{job=\"gateway\"}}[{w}])"
" or rate(otelcol_process_cpu_seconds{{job=\"gateway\"}}[{w}])"
),
"gateway_rss_mib": (
"otelcol_process_memory_rss_bytes{{job=\"gateway\"}} / 1024 / 1024"
" or otelcol_process_memory_rss{{job=\"gateway\"}} / 1024 / 1024"
),
"gateway_points_per_s": (
"rate(otelcol_receiver_accepted_metric_points"
"rate(otelcol_receiver_accepted_metric_points_total"
" {{job=\"gateway\"}}[{w}])"
" or rate(otelcol_receiver_accepted_metric_points"
" {{job=\"gateway\"}}[{w}])"
),
"gateway_out_series_per_s": (
"rate(otelcol_exporter_sent_metric_points"
"rate(otelcol_exporter_sent_metric_points_total"
" {{job=\"gateway\"}}[{w}])"
" or rate(otelcol_exporter_sent_metric_points"
" {{job=\"gateway\"}}[{w}])"
),
# ── Backend tier (destination-2) ────────────────────────────
# NOTE: these only populate under an ingest+query soak. An ingest-
# only soak (like the Phase 2.11B audit) leaves them at NaN. This
# is a design-level gap, not a query bug — see
# docs/phase-2-perf-deployment.md "Gaps" section #2.
"backend_samples_per_s": "rate(asap_ingest_samples_total[{w}])",
"backend_query_p99_ms": (
"1000 * histogram_quantile(0.99, "
Expand Down
Loading