Skip to content

Commit bbdcf66

Browse files
P2p call metrics (#174)
* feat: Enhance call metrics and documentation * refactor: Remove legacy cascade storage metrics keys for cleaner event data * feat: Pass send function to downloadArtifacts and restoreFileFromLayout for event streaming * feat: Add success rate percentage and noop metrics for store and retrieve operations * Increase TCP buffer size * Add handler request tracking per IP
1 parent 1b23e28 commit bbdcf66

31 files changed

Lines changed: 2190 additions & 812 deletions

File tree

‎docs/p2p-metrics-capture.md‎

Lines changed: 186 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,186 @@
1+
# P2P Metrics Capture — What Each Field Means and Where It’s Collected
2+
3+
This guide explains every field we emit in Cascade events, how it is measured, and exactly where it is captured in the code.
4+
5+
The design is minimal by intent:
6+
- Metrics are collected only for the first pass of Register (store) and for the active Download operation.
7+
- P2P APIs return errors only; per‑RPC details are captured via a small metrics package (`pkg/p2pmetrics`).
8+
- No aggregation; we only group raw RPC attempts by IP.
9+
10+
---
11+
12+
## Store (Register) Event
13+
14+
Event payload shape
15+
16+
```json
17+
{
18+
"store": {
19+
"duration_ms": 9876,
20+
"symbols_first_pass": 220,
21+
"symbols_total": 1200,
22+
"id_files_count": 14,
23+
"success_rate_pct": 82.5,
24+
"calls_by_ip": {
25+
"10.0.0.5": [
26+
{"ip": "10.0.0.5", "address": "A:4445", "keys": 100, "success": true, "duration_ms": 120},
27+
{"ip": "10.0.0.5", "address": "A:4445", "keys": 120, "success": false, "error": "timeout", "duration_ms": 300}
28+
]
29+
}
30+
}
31+
}
32+
```
33+
34+
### Fields
35+
36+
- `store.duration_ms`
37+
- Meaning: End‑to‑end elapsed time of the first‑pass store phase (Register’s storage section only).
38+
- Where captured: `supernode/services/cascade/adaptors/p2p.go`
39+
- A `time.Now()` timestamp is taken just before the first‑pass store function and measured on return.
40+
41+
- `store.symbols_first_pass`
42+
- Meaning: Number of symbols sent during the Register first pass (across the combined first batch and any immediate first‑pass symbol batches).
43+
- Where captured: `supernode/services/cascade/adaptors/p2p.go` via `p2pmetrics.SetStoreSummary(...)` using the value returned by `storeCascadeSymbolsAndData`.
44+
45+
- `store.symbols_total`
46+
- Meaning: Total symbols available in the symbol directory (before sampling). Used to contextualize the first‑pass coverage.
47+
- Where captured: Computed in `storeCascadeSymbolsAndData` and included in `SetStoreSummary`.
48+
49+
- `store.id_files_count`
50+
- Meaning: Number of redundant metadata files (ID files) sent in the first combined batch.
51+
- Where captured: `len(req.IDFiles)` in `StoreArtefacts`, passed to `SetStoreSummary`.
52+
53+
- `store.calls_by_ip`
54+
- Meaning: All raw network store RPC attempts grouped by the node IP.
55+
- Each array entry is a single RPC attempt with:
56+
- `ip` — Node IP (fallback to `address` if missing).
57+
- `address` — Node string `IP:port`.
58+
- `keys` — Number of items in that RPC attempt (metadata + first symbols for the first combined batch, symbols for subsequent batches within the first pass).
59+
- `success` — True if there was no transport error and no error message returned by the node response. Note: this flag does not explicitly check the `ResultOk` status; in rare cases, a non‑OK response with an empty error message may appear as `success` in metrics. (Internal success‑rate enforcement still uses explicit response status.)
60+
- `error` — Any error string captured; omitted when success.
61+
- `duration_ms` — RPC duration in milliseconds.
62+
- `noop` — Present and `true` when no store payload was sent to the node (empty batch for that node). Such entries are recorded as `success=true`, `keys=0`, with no `error`.
63+
- Where captured:
64+
- Emission point (P2P): `p2p/kademlia/dht.go::IterateBatchStore(...)`
65+
- After each node RPC returns, we call `p2pmetrics.RecordStore(taskID, Call{...})`. For nodes with no payload, a `noop: true` entry is emitted without sending a wire RPC.
66+
- `taskID` is read from the context via `p2pmetrics.TaskIDFromContext(ctx)`.
67+
- Grouping: `pkg/p2pmetrics/metrics.go`
68+
- `StartStoreCapture(taskID)` enables capture; `StopStoreCapture(taskID)` disables it.
69+
- Calls are grouped by `ip` (fallback to `address`) without further aggregation.
70+
71+
- `store.success_rate_pct`
72+
- Meaning: First‑pass store success rate computed from captured per‑RPC outcomes: successful responses divided by total recorded store RPC attempts, expressed as a percentage.
73+
- Where captured: Computed in `pkg/p2pmetrics/metrics.go::BuildStoreEventPayloadFromCollector` from `calls_by_ip` data.
74+
75+
### First‑Pass Success Threshold
76+
77+
- Internal enforcement only: if DHT first‑pass success rate is below 75%, `IterateBatchStore` returns an error.
78+
- We also emit `store.success_rate_pct` for analytics; the threshold only affects control flow (errors), not the emitted metric.
79+
- Code: `p2p/kademlia/dht.go::IterateBatchStore`.
80+
81+
### Scope Limits
82+
83+
- Background worker (which continues storing remaining symbols) is NOT captured — we don’t set a metrics task ID on those paths.
84+
85+
---
86+
87+
## Download Event
88+
89+
Event payload shape
90+
91+
```json
92+
{
93+
"retrieve": {
94+
"found_local": 42,
95+
"retrieve_ms": 2000,
96+
"decode_ms": 8000,
97+
"calls_by_ip": {
98+
"10.0.0.7": [
99+
{"ip": "10.0.0.7", "address": "B:4445", "keys": 13, "success": true, "duration_ms": 90}
100+
]
101+
}
102+
}
103+
}
104+
```
105+
106+
### Fields
107+
108+
- `retrieve.found_local`
109+
- Meaning: Number of items retrieved from local storage before any network calls.
110+
- Where captured: `p2p/kademlia/dht.go::BatchRetrieve(...)`
111+
- After `fetchAndAddLocalKeys`, we call `p2pmetrics.ReportFoundLocal(taskID, int(foundLocalCount))`.
112+
- `taskID` is read from context with `p2pmetrics.TaskIDFromContext(ctx)`.
113+
114+
- `retrieve.retrieve_ms`
115+
- Meaning: Time spent in network batch‑retrieve.
116+
- Where captured: `supernode/services/cascade/download.go`
117+
- Timestamp before `BatchRetrieve`, measured after it returns.
118+
119+
- `retrieve.decode_ms`
120+
- Meaning: Time spent decoding symbols and reconstructing the file.
121+
- Where captured: `supernode/services/cascade/download.go`
122+
- Timestamp before decode, measured after it returns.
123+
124+
- `retrieve.calls_by_ip`
125+
- Meaning: All raw per‑RPC retrieve attempts grouped by node IP.
126+
- Each array entry is a single RPC attempt with:
127+
- `ip`, `address` — Identifiers as available.
128+
- `keys` — Number of symbols returned by that node in that call.
129+
- `success` — True if the RPC completed without error (even if `keys == 0`). Transport/status errors remain `success=false` with an `error` message.
130+
- `error` — Error string when the RPC failed; omitted otherwise.
131+
- `duration_ms` — RPC duration in milliseconds.
132+
- `noop` — Present and `true` when no network request was actually sent to the node (e.g., all requested keys were already satisfied or deduped before issuing the call). Such entries are recorded as `success=true`, `keys=0`, with no `error`.
133+
- Where captured:
134+
- Emission point (P2P): `p2p/kademlia/dht.go::iterateBatchGetValues(...)`
135+
- Each node attempt records a `p2pmetrics.RecordRetrieve(taskID, Call{...})`. For attempts where no network RPC is sent, a `noop: true` entry is emitted.
136+
- `taskID` is extracted from context using `p2pmetrics.TaskIDFromContext(ctx)`.
137+
- Grouping: `pkg/p2pmetrics/metrics.go` (same grouping/fallback as store).
138+
139+
### Scope Limits
140+
141+
- Metrics are captured only for the active Download call (context is tagged in `download.go`).
142+
143+
---
144+
145+
## Context Tagging (Task ID)
146+
147+
- We use an explicit, metrics‑only context key defined in `pkg/p2pmetrics` to tag P2P calls with a task ID.
148+
- Setters: `p2pmetrics.WithTaskID(ctx, id)`.
149+
- Getters: `p2pmetrics.TaskIDFromContext(ctx)`.
150+
- Where it is set:
151+
- Store (first pass): `supernode/services/cascade/adaptors/p2p.go` wraps `StoreBatch` calls.
152+
- Download: `supernode/services/cascade/download.go` wraps `BatchRetrieve` call.
153+
154+
---
155+
156+
## Building and Emitting Events
157+
158+
- Store
159+
- `supernode/services/cascade/helper.go::emitArtefactsStored(...)`
160+
- Builds `store` payload via `p2pmetrics.BuildStoreEventPayloadFromCollector(taskID)`.
161+
- Includes `success_rate_pct` (first‑pass store success rate computed from captured per‑RPC outcomes) in addition to the minimal fields.
162+
- Emits the event.
163+
164+
- Download
165+
- `supernode/services/cascade/download.go`
166+
- Builds `retrieve` payload via `p2pmetrics.BuildDownloadEventPayloadFromCollector(actionID)`.
167+
- Emits the event.
168+
169+
---
170+
171+
## Quick File Map
172+
173+
- Capture + grouping: `supernode/pkg/p2pmetrics/metrics.go`
174+
- Store adaptor: `supernode/supernode/services/cascade/adaptors/p2p.go`
175+
- Store event: `supernode/supernode/services/cascade/helper.go`
176+
- Download flow: `supernode/supernode/services/cascade/download.go`
177+
- DHT store calls: `supernode/p2p/kademlia/dht.go::IterateBatchStore`
178+
- DHT retrieve calls: `supernode/p2p/kademlia/dht.go::BatchRetrieve` and `iterateBatchGetValues`
179+
180+
---
181+
182+
## Notes
183+
184+
- No P2P stats/snapshots are used to build events.
185+
- No aggregation is performed; we only group raw RPC attempts by IP.
186+
- First‑pass success rate is enforced internally (75% threshold) but not emitted as a metric.

0 commit comments

Comments
 (0)