diff --git a/.github/workflows/testing.yml b/.github/workflows/testing.yml index 20c1cf5ea..d0bfa4f1a 100644 --- a/.github/workflows/testing.yml +++ b/.github/workflows/testing.yml @@ -30,6 +30,10 @@ jobs: run: | python -m pip install --upgrade pip python -m pip install .[ai-features,dev] + - name: Install sharkd (marker replay engine) + run: | + sudo apt-get update + sudo apt-get install -y wireshark-common - name: Lint with flake8 run: | diff --git a/docs/features/marker-tag-replay.md b/docs/features/marker-tag-replay.md index 4d148a5c8..51b35d3ac 100644 --- a/docs/features/marker-tag-replay.md +++ b/docs/features/marker-tag-replay.md @@ -12,14 +12,19 @@ See LICENSE file for licensing information. Replays traffic captured by [markers](marker-traffic-insight.md) **across links**, keyed by `tag`. Markers on different links that share a tag form one *distributed capture session*; once every marker under the tag is paused, their per-marker pcaps are merged into a single -timestamp-ordered timeline. The Web UI browses that timeline and fetches individual frames -on demand — each fetch decodes exactly one frame via `tshark` into a self-describing JSON -protocol tree. +timestamp-ordered timeline. The Web UI browses that timeline (with Wireshark-style packet +list columns) and fetches individual frames on demand — each fetch decodes exactly one +frame via the resident **sharkd** daemon into a self-describing JSON protocol tree. The unique observable: the delta between the same packet hitting two consecutive links measures the **intermediate node's forwarding latency** (host view) — something a single-link capture can never show. +**sharkd is a hard requirement** (part of the Wireshark package). Without it every replay +endpoint that needs the engine returns 501 — there is deliberately no degraded mode; one +engine, one rendering shape for the Web UI. (A tag whose sources captured nothing returns +an empty timeline without consulting the engine — an empty answer, not a degraded one.) + ## Architecture ```mermaid @@ -28,35 +33,71 @@ graph TB subgraph Controller["Controller (replay endpoints)"] GATE["Tag gate
(all markers under tag paused?)"] - SCAN["Timeline scan
(pcap record headers)"] - MAP["PDML → JSON
isomorphic mapper"] + SCAN["Timeline backbone
(pcap record-header scan,
merge ordering, hex reads)"] + SESS["sharkd sessions
(one per source pcap)"] end FS[("markers dir
{node}_{link}_{marker}.pcap")] - TS["tshark -T pdml
(one frame at a time)"] - TMP["/tmp scratch copy
(hardened-profile workaround)"] + TMP["/tmp scratch copies
(hardened-profile workaround)"] + SK["sharkd -
(resident JSON-RPC on stdio)"] - UI -->|"GET range / frames"| GATE + UI -->|"GET range / frames [?filter=&link=]"| GATE GATE --> SCAN - SCAN --> FS - UI -->|"GET frame detail (lazy)"| MAP - MAP --> TMP - TMP --> TS - MAP -->|"hex: raw bytes"| FS + SCAN -->|"columns / filter matches"| SESS + SESS --> TMP --> SK + SCAN -->|"hex: raw bytes"| FS + UI -->|"GET frame detail (lazy)"| SESS ``` -Two deliberately separated performance regimes: +Backbone vs engine, deliberately separated: -| Path | Work | tshark | -|------|------|--------| -| Timeline | 16-byte record-header scan per frame, cross-file merge sort | **never invoked** — browsing works without tshark | -| Frame detail | locate frame → `tshark -T pdml -Y "frame.number == N"` → map to JSON | forked once per frame the user opens | +- **Timeline backbone** (plain Python): the tag gate, 16-byte-per-frame pcap record-header + scan, cross-source merge ordering, canonical ts strings, and raw-bytes reads for the hex + view. Identity and ordering never depend on the engine. +- **Engine layer (sharkd)**: packet-list columns, display filters, and per-frame protocol + trees. One resident `sharkd -` process per source pcap (sharkd loads one file at a + time), spawned lazily on first use, addressed with one-line JSON-RPC on stdio. -tshark call count equals user clicks — no caching or rate limiting needed, and a tshark -failure (501/502) affects only that one frame, never timeline browsing. pcap files are -compute-side (`/project-files/markers/`); the initial scope is single-server -deployments (controller and compute in one process, direct file access) — remote computes -will reuse the existing capture-file proxy pattern. +Session lifecycle: each session is validated per request against the source pcap's +`(mtime, size)` — a mismatch (e.g. the capture node restarted and uBridge truncated the +pcap while paused) kills and respawns it. Sessions are refcounted while a request holds +them and the LRU bound (16) only ever evicts **idle** sessions, so neither a request's +own walk over a tag with more sources than the bound nor a concurrent request can have +its session killed mid-RPC. A single manager lock makes check-spawn atomic (concurrent +requests share one spawn instead of double-spawning a process nobody reaps). Each RPC +is serialized by a per-session lock (sharkd serves one request at a time), carries a +timeout and an **id check** — any transport failure (timeout, dead pipe, malformed or +stale reply, oversized line) kills the session for good, because a desynchronized +session would serve shifted results. sharkd reads a `/tmp` scratch **directory** +containing a copy of the pcap plus pinned Wireshark preferences (the column layout the +frames RPC is parsed against is a contract the server owns, immune to system or user +column customization) — hardened profiles (AppArmor &c.) deny it the project directory +and the user's home even though the server process can read both. The pcap copy runs +off the event loop, and server shutdown kills every session so no scratch copy leaks +into `/tmp` across restarts. + +### Resource bounds + +Three independent bounds, easily confused: + +| Bound | Value | Governs | +|-------|-------|---------| +| `_STREAM_LIMIT_BYTES` | 16 MB | Max **single JSON-RPC reply line** read from sharkd. A full 1000-row page measures ~190 KB against asyncio's 64 KB default — the trigger is the frame count in one pcap (≥ 1000), never the number of pcaps | +| `_FRAMES_PAGE` | 1000 | Rows fetched per `frames` RPC while draining one source's columns | +| `SESSION_MAX` | 16 | Resident **idle** sharkd sessions (LRU eviction) | + +The `range` frame list itself is uncapped by decision — the client owns the +rendering cost of a huge list. + +Process count follows usage, not the pcap inventory. sharkd loads one file per +process, so a session exists per source pcap **actually consulted** — a replay +request only touches its own tag's sources, and a session is either in use +(refcounted, held for a whole pagination drain rather than per RPC) or idle +(no holder — the only state the LRU cap ever evicts). At any moment the pool +holds at most `SESSION_MAX` idle sessions plus one per source being drained by +an in-flight request: a project with a thousand marker pcaps still runs at +most `16 + Σ(in-flight sources)` sharkd processes. Each process holds its pcap +in memory (~file size), which is the quantity the session cap primarily bounds. ## Business Process @@ -64,29 +105,22 @@ will reuse the existing capture-file proxy pattern. sequenceDiagram participant UI as Web UI participant C as Controller - participant PC as markers dir (pcaps) - participant TS as tshark + participant SK as sharkd session - Note over UI,PC: ① configure — same tag on every link's marker - UI->>C: POST markers (bpf, tag=666) on each link + Note over UI: ① configure — same tag on every link's marker + Note over C: ② capture — uBridge appends matches, replay forbidden (409) + Note over UI: ③ pause every marker under the tag - Note over PC: ② capture — uBridge appends matches (flushed per packet), replay forbidden - UI->>C: GET range - C--xC: 409 (a marker under tag 666 is still enabled) - - Note over UI: ③ pause — every marker under the tag - UI->>C: PUT markers {"enabled": false} × each - - Note over UI,TS: ④ replay - UI->>C: GET /markers/tags/666/replay/range - C->>PC: scan record headers, merge sort - C-->>UI: {start, end, sources, frames} - UI->>C: GET frames?ts=T&window_ms=W - C-->>UI: frames in [T, T+W] — or {"frames": []} - UI->>C: GET frame/detail?ts=…&node_id=…&link_id=…&marker=… - C->>PC: locate frame, read raw bytes (hex) - C->>TS: -T pdml (reads the /tmp copy) - TS-->>C: PDML + Note over UI,SK: ④ replay + UI->>C: GET /markers/tags/666/replay/range[?filter=…&link=…] + C->>C: gate → scan record headers → merge order + C->>SK: frames {filter, skip, limit} → columns + matches + C-->>UI: {start, end, sources, frames[] with src/dst/proto/info/bg/fg} + UI->>C: GET frames?ts=T&window_ms=W (paging — {"frames": []} on a miss) + UI->>C: GET frame/detail?ts=…&node_id=…&link_id=…&marker=…[&frame_number=…] + C->>C: hex straight from the pcap + C->>SK: frame {frame: N, proto: true} + SK-->>C: tree (keys renamed to the REST contract) C-->>UI: protocol tree + hex ``` @@ -110,96 +144,145 @@ markers at all → 404. - **Pause → resume → pause is fine.** The pcap accumulates the full history; replay covers everything up to the current pause point. - **The replay window ends when nodes restart.** A pcap's lifetime equals its uBridge's - lifetime: a fresh uBridge reinstalls every desired marker — paused ones too (install - first, then turn the filter off) — and uBridge opens the pcap with truncate semantics - (`pcap_dump_open`, not `_append`). Server restart + project reopen **without starting - nodes** is safe: nothing touches the files until a uBridge comes up (verified live). - Docker nodes effectively restart on server restart as well (stale-container cleanup), - so their window is shorter still. -- uBridge flushes every matched packet to the pcap immediately (`pcap_dump_flush` per - packet under a mutex — verified in the uBridge source), so a pause boundary never loses - tail frames. + lifetime: a fresh uBridge reinstalls every desired marker — paused ones too — and + uBridge opens the pcap with truncate semantics (`pcap_dump_open`, not `_append`). Server + restart + project reopen **without starting nodes** is safe: nothing touches the files + until a uBridge comes up. Docker nodes effectively restart on server restart as well + (stale-container cleanup), so their window is shorter still. ## API Endpoints -All read-only; JWT bearer token, privilege `Project.Audit`. +All read-only; JWT bearer token, privilege `Project.Audit`. All require sharkd — 501 +without it. | Method | Path | Description | |--------|------|-------------| -| GET | `/v3/projects/{pid}/markers/tags/{tag}/replay/range` | Timeline metadata + full frame list for the tag | -| GET | `/v3/projects/{pid}/markers/tags/{tag}/replay/frames?ts=&window_ms=&limit=` | Frames with ts in `[T, T+window]`, merged across sources | -| GET | `/v3/projects/{pid}/markers/tags/{tag}/replay/frame/detail?ts=&node_id=&link_id=&marker=` | Single frame: tshark protocol tree + raw hex | +| GET | `/v3/projects/{pid}/markers/tags/{tag}/replay/range[?filter=&link=]` | Timeline metadata + full merged frame list with packet-list columns | +| GET | `/v3/projects/{pid}/markers/tags/{tag}/replay/frames?ts=&window_ms=&limit=[&filter=&link=]` | Frames with ts in `[T, T+window]`, merged across sources | +| GET | `/v3/projects/{pid}/markers/tags/{tag}/replay/frame/detail?ts=&node_id=&link_id=&marker=[&frame_number=]` | Single frame: protocol tree + raw hex (lazy — one call per frame the user opens) | ### `range` — the timeline ```json { - "tag": 666, - "start": "1788196663.226372", - "end": "1788196713.706634", - "frame_count": 20, - "truncated": false, + "tag": 102, + "start": "1788369209.406812", + "end": "1788369219.249085", + "frame_count": 29, "sources": [ - { "node_id": "b764c434…", "link_id": "316ef8fd…", "marker": "global-def-…", - "data_link_type": "DLT_EN10MB", "count": 10 } + { "node_id": "47703cad…", "link_id": "2697a7c6…", "marker": "global-ospf", + "data_link_type": "DLT_EN10MB", "count": 4 } ], "frames": [ - { "ts": "1788196663.226372", "len": 98, "node_id": "b764c434…", - "link_id": "316ef8fd…", "marker": "global-def-…", "frame_number": 1 } + { "ts": "1788369209.406812", "len": 114, + "node_id": "47703cad…", "link_id": "2697a7c6…", + "marker": "global-ospf", "frame_number": 1, + "src": "10.0.12.1", "dst": "224.0.0.5", + "proto": "OSPF", "info": "Hello Packet", + "bg": "fff3d6", "fg": "12272e" } ] } ``` -- `frames` is the **full, merged, time-ordered list** (cap 5000) — the Web UI lays out the - whole timeline from one request. Over the cap, `frames` is omitted and per-second - `buckets` are returned instead with `truncated: true`. +- `frames` is the **full, merged, time-ordered list** — one request lays out the whole + timeline, **deliberately uncapped**: rendering a huge list is the client's concern, + and the window endpoint exists for incremental views. There is no `truncated` flag + and no bucket histogram; `frame_count` always equals `len(frames)`. +- Every frame entry carries the Wireshark packet-list columns — `src` / `dst` / `proto` + / `info` — plus the **coloring-rule hints `bg` / `fg`** (Wireshark's own palette + decisions, so the UI can color rows exactly like Wireshark without shipping the + colorization engine). Columns are `null` only for a frame the engine could not describe. - Each frame entry carries `(node_id, link_id, marker, frame_number)` — the locating - tuple for the detail request. + tuple for the detail request, and the link association for timeline/topology rendering + (`link_id` joins the Web UI's own link objects). -### `frames` — point / window query +### Display filter -A time with no frames is a normal, successful answer — an empty array, no sentinel strings: +`?filter=` on both `range` and `frames` is a Wireshark display filter, +applied **before** counting and slicing — `start` / `end` / `frame_count` / `frames` +are all computed on the matching frames only. Filtered frames keep +their original pcap frame numbers. The filter travels as one argv-style element (never +through a shell) and is capped at 2000 characters. An invalid expression is a **400** +whose message carries sharkd's original error text — suitable for inline display in the +filter bar, and distinct from the 409 gate / 404 unknown-tag semantics. Only sharkd's +filter-rejection error maps to 400; any other engine failure while a filter is set is a +502 (the filter was fine — the engine was not). + +### Capture-source selection + +`?link=` on both `range` and `frames` narrows the frame stream to one +capture source — a pure identity filter applied **before** any engine work (only the +selected link's pcap gets a sharkd pass), AND-composing with `filter`. Windows and the +histogram therefore always agree with the link-filtered view. Two boundaries by +design: + +- **`sources` is the stable inventory of the tag**: every capture source is always + listed, with engine-free **total** counts, unaffected by `link` / `filter` — a + source dropdown must not shrink when the view narrows. +- **An unknown `link_id` matches nothing**: `frame_count: 0`, `start: null`, empty + `frames` — the same shape as a zero-match display filter, deliberately not a 404. + +### `frames` — point / window query (paging) + +A time with no frames is a normal, successful answer — an empty array, no sentinel +strings: ```json GET …/replay/frames?ts=1788196700.000&window_ms=500 → { "frames": [] } ``` +Paging is deliberately **ts + window_ms only** (no offset/limit over the filtered set): +the merge spans multiple pcaps, so slicing happens server-side on the merged stream +either way, windows align with timeline semantics, and the gate freezes the data (the +window answer is deterministic). + ### `frame/detail` — lazy single-frame decode -Invoked only when the user opens a frame. The `ts` must be the **exact string received in -the timeline/frame list** (round-tripped verbatim — never re-serialized through a float); -`node_id + link_id + marker` identify the pcap. The server re-resolves the ts against the -file, guarding against a capture rebuilt between the timeline view and this click. +Invoked only when the user opens a frame. The `ts` must be the **exact string received +in the timeline/frame list** (round-tripped verbatim — never re-serialized through a +float); `node_id + link_id + marker` identify the pcap. The server re-resolves the ts +against the file, guarding against a capture rebuilt between the timeline view and this +click. Since ts is **not unique within one pcap** (same-microsecond frames are kept +deliberately), the optional `frame_number` — carried by every frame list entry — +disambiguates them: it must still land on the exact ts, and without it the first ts +match decodes. ```json { - "ts": "1788196663.226372", - "source": { "node_id": "b764c434…", "link_id": "316ef8fd…", - "marker": "global-def-…", "frame_number": 1 }, - "tshark_version": "TShark (Wireshark) 4.6.7 …", - "field_count": 85, - "hex": "45000062…", + "ts": "1788369209.406812", + "source": { "node_id": "47703cad…", "link_id": "2697a7c6…", + "marker": "global-ospf", "frame_number": 1 }, + "field_count": 89, + "hex": "01005e000005…", "tree": [ - { "element": "proto", "name": "ip", - "showname": "Internet Protocol Version 4, Src: 10.1.10.101, Dst: 203.0.113.1", - "children": [ - { "element": "field", "name": "ip.ttl", "show": "64", - "showname": "Time to Live: 64", "value": "40", "size": "1", - "pos": "22", "children": [] } + { "element": "proto", "label": "Internet Protocol Version 4, …", "children": [ + { "element": "field", "name": "ip.ttl", "label": "Time to Live: 1", + "filter_expr": "ip.ttl == 1", "pos": 22, "size": 1, "children": [] } ] } ] } ``` -- `tree` mirrors PDML **isomorphically**: every ``/`` becomes a node, every - XML attribute (`name`, `show`, `showname`, `value`, `size`, `pos`, `hide`, `mask`, - `unmaskedvalue`) becomes a JSON key, plus one structural key `element` (proto/field). - Nothing is selected out, nothing interpreted, and **all values stay strings** — numeric - conversion is the client's business. -- `hex` is the raw frame bytes read straight from the pcap (not via tshark); keeping - `pos`/`size` on every field enables Wireshark-style *click field → highlight bytes*. -- `field_count` is the mapped node count — a client-side sanity check. +The tree is sharkd's protocol tree with **keys renamed into the REST contract** — a +closed, protocol-independent key set (census-verified across ICMP / TCP / VLAN+OSPF +trees and pinned by a test), with values untouched: + +| sharkd | Contract key | Meaning | +|--------|--------------|---------| +| `t` | `element` | node type (`proto`, …) | +| `l` | `label` | display text | +| `fn` | `name` | field name (`ip.ttl`) | +| `f` | `filter_expr` | **ready-made display filter with the value baked in** — click-to-filter | +| `h` | `pos` + `size` | byte range — click field → highlight hex bytes | +| `s` | `expert` | expert severity name (`Chat` / `Warn` / …) | +| `g` | `generated` | generated-by-Wireshark flag | +| `n` | `children` | nested fields | +| `e` | *(dropped)* | Wireshark-internal hf id, unstable across versions | + +Unknown keys from a newer Wireshark pass through verbatim (never silently dropped); a +census test flags new keys for naming. `hex` is the raw frame bytes read straight from +the pcap; `field_count` is the mapped node count (client-side sanity check). ## Ordering and timestamps @@ -209,45 +292,54 @@ file, guarding against a capture rebuilt between the timeline view and this clic environment. - The sort key is `(ts, source file, frame_number)` — ts alone is **not** unique (two links can hit the same microsecond); the tiebreaker yields a stable, determined order - instead of a fictional one. Index structures must never use ts as a dict key, or - same-microsecond frames silently overwrite each other. + instead of a fictional one. - The cross-link delta is the intermediate node's end-to-end forwarding latency (veth/TAP → guest protocol stack → back to host), typically hundreds of microseconds to milliseconds. UI labels should read *node forwarding latency (host view)*, not link propagation delay. A live capture pair confirmed it end-to-end: same `ip.id`, TTL 64→63, 509 µs between two links. -## Fidelity guarantee (PDML → JSON) - -The conversion is an isomorphic structure map, not a semantic transform, with two rules: -**map every attribute** and **keep values as strings**. A round-trip test enforces both: -PDML element count equals JSON node count, and every XML attribute survives with an -identical JSON value (`tests/controller/test_marker_replay.py`). The raw frame bytes — -the one thing PDML genuinely does not contain — are covered by `hex` read directly from -the pcap. - ## Error Responses +All error bodies are `{"message": "…"}` (the app's unified format). + | Status | Description | |--------|-------------| +| 400 | Invalid display filter (sharkd's filter rejection only; message carries its original text) or filter longer than 2000 chars | | 401 | Not authenticated | -| 404 | Tag has no markers in the project; detail source unknown, or ts does not match the file (the capture may have been rebuilt) | +| 404 | Tag has no markers in the project; detail source unknown; `ts`/`frame_number` do not match the file (the capture may have been rebuilt); or the pcap vanished mid-request | | 409 | Tag gate: a marker under the tag is still `enabled: true` (the response lists them) | -| 501 | tshark not installed / unavailable — affects detail only; the timeline never needs tshark | -| 502 | tshark failed or timed out (10 s); truncated output never reaches the mapper | +| 501 | sharkd not installed / unavailable — replay is unavailable, no degraded mode | +| 502 | sharkd failed, timed out (10 s per RPC), or answered out of sync — the session is killed and re-spawned on the next request | ## Notes - **Heterogeneous link types coexist.** Frames are never merged into a single pcap (mergecap is deliberately not used) — each frame carries its source and is decoded individually, so Ethernet and serial (cHDLC/PPP) markers can share one timeline. - Malformed packets are tshark's problem: it emits `[Malformed Packet]` as regular PDML - and carries on. -- **Hardened tshark profiles.** openSUSE-style profiles (AppArmor &c.) can deny tshark - access to the project directory and the user's home even though the server process can - read both. The detail path therefore copies the pcap to a real file under `/tmp` - (not a symlink — the profile resolves real paths) and gives tshark a scratch `HOME`; - the copy is unlinked afterwards. The hex view still reads the original file. + Malformed packets are dissected like any other; sharkd marks them in the tree. +- **Live validation (2026-09, 9-link OSPF project, 29 frames over 9 sources).** Cold + `range` (spawning all sharkd sessions) answered in 0.84 s with full columns and + Wireshark coloring; filters verified in all four regimes (match / zero-match with + `start: null`, invalid expression → 400 with sharkd's text, oversized → 400); window + hit and miss behaved per contract; a frame detail returned 89 nodes with + `filter_expr: "ip.ttl == 1"` (OSPF multicast TTL) and byte ranges for hex + highlighting. +- **Columns are re-fetched per request** (~90 ms per source against a loaded session). + The data is frozen while the gate passes, so a cache keyed on `(mtime, size)` is a + natural follow-up if list latency ever matters at many-source scale. +- **Session invalidation is cheap and total.** Every request stats the source pcap; a + rewritten file (mtime/size change) respawns the session — a paused-but-restarted + capture can never serve stale dissect state. Transport failures are equally total: a + timeout, dead pipe, malformed or stale reply (id mismatch) kills the session instead + of leaving a poisoned one resident, and concurrent requests for the same pcap share a + single spawn. +- **The packet-list layout is pinned, not assumed.** sharkd runs with a scratch `HOME` + whose Wireshark preferences fix `gui.column.format` to exactly the four columns the + frame entries carry (personal config overrides any `/etc/wireshark` customization), so + the column indexes the server parses are a contract it owns. +- **Server shutdown kills every resident session** and drops its `/tmp` scratch + directory — nothing accumulates across restarts. - **Tag type.** REST and the `marker.match` WS event both carry `tag` as `int` (the listener normalizes); replay keys on that int value. - **Follow-ups.** Remote-compute support via the existing capture-file proxy pattern; diff --git a/gns3server/api/routes/controller/projects.py b/gns3server/api/routes/controller/projects.py index 0595bde7b..26fb7d042 100644 --- a/gns3server/api/routes/controller/projects.py +++ b/gns3server/api/routes/controller/projects.py @@ -45,7 +45,7 @@ from gns3server.controller.controller_error import ControllerError, ControllerBa from gns3server.controller.import_project import import_project as import_controller_project from gns3server.controller.export_project import export_project as export_controller_project from gns3server.controller import marker_replay -from gns3server.controller.marker_replay import TsharkError, TsharkMissingError +from gns3server.controller.marker_replay import SharkdError, SharkdMissingError from gns3server.utils.asyncio import aiozipstream from gns3server.utils.path import is_safe_path from gns3server.db.repositories.templates import TemplatesRepository @@ -221,36 +221,73 @@ def get_project_markers(project: Project = Depends(dep_project)) -> dict: return project.markers +async def _replay_response(awaitable): + """Shared engine-error mapping for the replay endpoints: 501 when sharkd + (the hard engine requirement) is unavailable, 502 when it fails. Data + state errors (409 gate / 404 unknown tag) and filter errors (400) map + through the global handlers before this.""" + + try: + return await awaitable + except SharkdMissingError as e: + raise HTTPException(status_code=status.HTTP_501_NOT_IMPLEMENTED, detail=str(e)) + except SharkdError as e: + raise HTTPException(status_code=status.HTTP_502_BAD_GATEWAY, detail=str(e)) + + @router.get( "/{project_id}/markers/tags/{tag}/replay/range", dependencies=[Depends(has_privilege("Project.Audit"))], ) -def replay_tag_range(tag: int, project: Project = Depends(dep_project)) -> dict: +async def replay_tag_range( + tag: int, + filter: Optional[str] = None, + link: Optional[str] = None, + project: Project = Depends(dep_project), +) -> dict: """ Aggregate replay timeline for a tag: merges the pcap of every marker - carrying ``tag`` into one timestamp-ordered view (design reference: - ``marker_replay`` module docstring). + carrying ``tag`` into one timestamp-ordered view. Every frame entry + carries Wireshark-style columns (``src`` / ``dst`` / ``proto`` / ``info`` + plus coloring hints ``bg`` / ``fg``). The tag gate applies: every marker under the tag must be paused (``enabled: false``) — 409 otherwise. The response carries the timeline - bounds, per-source stats, and the full merged frame list while under the - frame cap (5000); above it the list is replaced by per-second buckets. + bounds, per-source stats, and the full merged frame list — deliberately + uncapped (rendering a huge list is the client's concern; the window + endpoint exists for incremental views). + + ``filter`` is an optional Wireshark display filter applied **before** + counting and slicing — start / end / frame_count / frames are + all computed on the matching frames only. An invalid expression is a 400 + carrying sharkd's original error text (for inline display in the UI + filter bar). + + ``link`` narrows the frame stream to one capture source (link_id), + AND-composing with ``filter``; ``sources`` always lists the tag's full + inventory regardless. An unknown link_id yields an empty timeline (same + shape as a zero-match filter), not a 404. Requires sharkd — 501 without + it. Required privilege: Project.Audit """ - return marker_replay.build_timeline(project, tag) + return await _replay_response( + marker_replay.build_timeline(project, tag, filter_expr=filter, link_id=link) + ) @router.get( "/{project_id}/markers/tags/{tag}/replay/frames", dependencies=[Depends(has_privilege("Project.Audit"))], ) -def replay_tag_frames( +async def replay_tag_frames( tag: int, ts: str, window_ms: int = 100, limit: int = 1000, + filter: Optional[str] = None, + link: Optional[str] = None, project: Project = Depends(dep_project), ) -> dict: """ @@ -259,12 +296,20 @@ def replay_tag_frames( ``{"frames": []}``. The tag gate applies (409 while any marker captures). ``ts`` must be the exact string returned by the range response — never - re-serialize it through a float. + re-serialize it through a float. ``filter`` and ``link`` (optional + display filter / capture-source link_id) have the same semantics as on + the range endpoint — windows and the histogram always agree. Requires + sharkd — 501 without it. Required privilege: Project.Audit """ - return marker_replay.query_frames(project, tag, ts, window_ms=window_ms, limit=limit) + return await _replay_response( + marker_replay.query_frames( + project, tag, ts, window_ms=window_ms, limit=limit, + filter_expr=filter, link_id=link, + ) + ) @router.get( @@ -277,29 +322,31 @@ async def replay_tag_frame_detail( node_id: str, link_id: str, marker: str, + frame_number: Optional[int] = None, project: Project = Depends(dep_project), ) -> dict: """ Decode exactly one frame (lazy — invoked when the user opens a frame, never by the timeline itself): raw bytes for the hex view read straight - from the pcap, protocol tree from ``tshark -T pdml`` mapped isomorphically - to JSON (every PDML attribute survives, values stay strings). + from the pcap, protocol tree from the resident sharkd session with keys + renamed into the REST contract (``element`` / ``label`` / ``name`` / + ``filter_expr`` / ``pos`` + ``size`` / ``expert`` / ``generated`` / + ``children``) — values untouched. ``ts`` must be the exact string from the frame list; ``node_id`` + - ``link_id`` + ``marker`` identify the source pcap. The tag gate applies. + ``link_id`` + ``marker`` identify the source pcap. ``frame_number`` + (optional, from the frame list entry) disambiguates same-microsecond + frames on one link — without it the first ts match decodes. The tag gate + applies. Requires sharkd — 501 without it. Required privilege: Project.Audit """ - try: - return await marker_replay.decode_frame(project, tag, ts, node_id, link_id, marker) - except TsharkMissingError: - raise HTTPException( - status_code=status.HTTP_501_NOT_IMPLEMENTED, - detail="tshark is not installed on this server — frame detail is unavailable", + return await _replay_response( + marker_replay.decode_frame( + project, tag, ts, node_id, link_id, marker, frame_number=frame_number ) - except TsharkError as e: - raise HTTPException(status_code=status.HTTP_502_BAD_GATEWAY, detail=str(e)) + ) # --------------------------------------------------------------------------- diff --git a/gns3server/controller/marker_replay.py b/gns3server/controller/marker_replay.py index 993bb31f8..a3b9366e3 100644 --- a/gns3server/controller/marker_replay.py +++ b/gns3server/controller/marker_replay.py @@ -16,61 +16,109 @@ # along with this program. If not, see . """ -Tag-keyed aggregate replay over paused markers' pcap files. +Tag-keyed aggregate replay over paused markers' pcaps, powered by sharkd. Markers on different links sharing a ``tag`` form one distributed capture session. This module merges their per-marker pcaps (``/project-files/markers/{node_id}_{link_id}_{name}.pcap``) into a -single timestamp-ordered timeline and decodes individual frames on demand. +single timestamp-ordered timeline and decodes frames on demand. -Two deliberately separated performance regimes: +Engine: **sharkd is a hard requirement** (the Wireshark resident daemon, +driven over one-line JSON-RPC on stdin/stdout). Without it every replay +endpoint returns 501 — there is deliberately no degraded mode. Timeline +assembly (gate, pcap record-header scan, merge ordering, hex reads) is plain +Python, but it is an implementation detail, not an availability promise. -* the timeline path reads only the 16-byte pcap record headers — tshark is - never invoked, so browsing works even where tshark is not installed; -* the detail path runs one ``tshark -T pdml`` per frame the caller asks about - (call count = user clicks) and maps the XML to JSON isomorphically — every - PDML attribute survives as a JSON key, values stay strings, nothing is - selected out or interpreted. +Session model — sharkd loads one file at a time, so there is one resident +session **per source pcap**: spawned lazily on first use, fed a ``/tmp`` +scratch directory (hardened profiles deny sharkd the project directory; a +real copy, not a symlink, since the profile resolves real paths) with a +scratch ``HOME`` whose Wireshark preferences pin the packet-list column +layout, validated per request against the original's ``(mtime, size)`` — +a mismatch (e.g. the capture node restarted and uBridge truncated the pcap +while paused) kills and respawns the session. A per-session asyncio lock +serializes RPCs (sharkd serves one request at a time), each with a timeout +and an id check; any transport failure (timeout, dead pipe, malformed or +stale reply) kills the session for good — a desynced session must never +serve shifted results. A single manager lock makes check-spawn atomic (no +double spawn under concurrent requests) and sessions are refcounted while +in use, so the LRU cap only ever evicts idle ones. Timestamps are uBridge's userspace ``gettimeofday`` at match time (µs, a value measured after the packet has crossed the kernel twice — the last digit or two are scheduling noise). A timestamp is NOT a unique key: the merge sorts by ``(ts, source file, frame number)`` and index structures must never use ts alone as a dict key, or same-microsecond frames silently -overwrite each other. +overwrite each other. The canonical ts strings travel to clients verbatim +and must be round-tripped verbatim. """ import asyncio +import json import logging import os import shutil import struct import tempfile -import xml.etree.ElementTree as ET +import time +from contextlib import asynccontextmanager from .controller_error import ControllerError, ControllerNotFoundError, ControllerBadRequestError log = logging.getLogger(__name__) -# Full frame list is embedded in the range response while under this cap; -# above it the response degrades to start/end + per-second buckets so the -# client is never flooded by a high-traffic BPF. -FRAME_LIST_CAP = 5000 +# One JSON-RPC per sharkd session, serialized by a per-session lock. +RPC_TIMEOUT_SECONDS = 10.0 -# One tshark decode per user click: a generous ceiling, not a rate limiter. -TSHARK_TIMEOUT_SECONDS = 10.0 +# Resident sharkd sessions are bounded; only IDLE sessions are ever evicted +# (least-recently-used first), so a request's own walk over a tag with more +# sources than the cap never kills a session out from under it — the +# population is trimmed back as uses drop instead. +SESSION_MAX = 16 + +# Display filters travel as one argv element (never through a shell) and are +# capped to keep absurd expressions off the command line. +FILTER_MAX_LENGTH = 2000 + +# Batch size when draining sharkd's `frames` RPC (columns / filter matches). +_FRAMES_PAGE = 1000 + +# A full `frames` page (1000 rows of pinned columns) measures ~190 KB against +# the 64 KB default StreamReader limit — over it, readline() raises and clears +# the stream, failing the request AND desyncing the resident session. Raise +# the ceiling well above any page (or single-frame tree) we can produce. +_STREAM_LIMIT_BYTES = 16 * 1024 * 1024 + +# sharkd's JSON-RPC error code for a rejected display filter — the only error +# that belongs to the client (400); everything else is an engine fault (502). +_ERR_INVALID_FILTER = -13002 -class TsharkMissingError(ControllerError): - """tshark is not installed (or not on PATH) — frame detail unavailable.""" +class SharkdMissingError(ControllerError): + """sharkd is not installed (or not on PATH) — replay is unavailable (501).""" -class TsharkError(ControllerError): - """tshark exited non-zero / timed out / produced unusable output.""" +class SharkdError(ControllerError): + """sharkd failed, timed out, or produced an unusable response (502).""" + + +class FilterError(ControllerBadRequestError): + """sharkd rejected the display filter — its message is carried verbatim + so the UI can show it inline in the filter bar (400, distinct from the + 409 gate / 404 unknown-tag semantics).""" + + +class _SharkdRpcError(Exception): + """Internal: a JSON-RPC error object from sharkd (code + message).""" + + def __init__(self, code, message): + super().__init__(message) + self.code = code + self.message = message # --------------------------------------------------------------------------- -# pcap record-header scanning (timeline path — no tshark) +# pcap record-header scanning (timeline backbone — engine-free) # --------------------------------------------------------------------------- # magic → (byte order, timestamp unit). Both pcap families uBridge can write @@ -140,8 +188,8 @@ def scan_pcap_frames(path): def read_frame_bytes(path, frame_number): """ Read one frame's raw bytes (hex view) straight from the pcap — never via - tshark. ``frame_number`` is 1-based (the same number tshark's - ``frame.number`` filter uses). + the engine. ``frame_number`` is 1-based (the same number sharkd's frame + RPC uses). """ frames = scan_pcap_frames(path) @@ -156,6 +204,394 @@ def read_frame_bytes(path, frame_number): return f.read(incl_len).hex() +# --------------------------------------------------------------------------- +# sharkd: process environment and scratch copies +# --------------------------------------------------------------------------- + +# The packet-list layout is pinned through Wireshark preferences in the +# scratch HOME: personal config overrides any system-wide customization +# (/etc/wireshark &c.), so the column indexes in _columns_for are a contract +# we own rather than an environment default. Exactly the four columns the +# frame entries consume — nothing else rides along in every page. +_PINNED_COLUMNS = ( + 'gui.column.format: "Source", "%s", "Destination", "%d", ' + '"Protocol", "%p", "Info", "%i"\n' +) + + +async def _prepare_scratch(pcap): + """ + Build a scratch directory under the system temp dir holding the pcap copy + for sharkd to load plus the pinned column preferences. Hardened profiles + (AppArmor &c.) can deny sharkd the project directory / the user's home + while still allowing /tmp — a real copy, deliberately not a symlink, + since the profile resolves real paths. Caller must rmtree the directory. + + :returns: ``(scratch_dir, scratch_pcap_path)`` + """ + + scratch_dir = tempfile.mkdtemp(prefix="gns3-replay-") + scratch = os.path.join(scratch_dir, "capture.pcap") + # A pcap-sized copy has no business stalling the event loop — a 1 GB + # capture must not freeze every other request for the duration. + await asyncio.to_thread(shutil.copyfile, pcap, scratch) + prefs_dir = os.path.join(scratch_dir, ".config", "wireshark") + os.makedirs(prefs_dir, exist_ok=True) + with open(os.path.join(prefs_dir, "preferences"), "w") as f: + f.write(_PINNED_COLUMNS) + return scratch_dir, scratch + + +def _engine_env(scratch_dir): + """Scratch HOME (and XDG config) so sharkd never even tries to read the + user's home — and always reads OUR pinned preferences instead.""" + + env = dict(os.environ) + env["HOME"] = scratch_dir + env["XDG_CONFIG_HOME"] = os.path.join(scratch_dir, ".config") + return env + + +# --------------------------------------------------------------------------- +# sharkd: tree key renaming (the only transformation between sharkd and the +# REST contract — a closed, protocol-independent key set; values untouched) +# --------------------------------------------------------------------------- + +# Census-verified across ICMP / TCP / VLAN+OSPF trees: sharkd emits exactly +# these structural keys on every node regardless of protocol. Protocol +# semantics live in VALUES (field names, labels, filter expressions), which +# are never touched. +_KEY_RENAME = { + "t": "element", # node type ("proto", …) + "l": "label", # display text + "fn": "name", # field name (e.g. "ip.ttl") + "f": "filter_expr", # ready-made display filter with the value baked in + "s": "expert", # expert severity name ("Chat", "Warn", …) + "g": "generated", # generated-by-wireshark flag + "n": "children", # nested fields +} +# "h" → pos + size (byte range for hex highlighting) — handled specially. +# "e" is sharkd's internal header-field registry id — unstable across +# Wireshark versions and useless for rendering, so it is dropped. +_DROPPED_KEYS = {"e"} + + +def _rename_value(value): + """Recursive pass-through: rename known keys, drop none but 'e', + copy unknown keys verbatim (a future Wireshark adding a key never + silently loses data — the census test flags it for naming).""" + + if isinstance(value, dict): + out = {} + for key, item in value.items(): + if key in _DROPPED_KEYS: + continue + if key == "h" and isinstance(item, list) and len(item) == 2: + out["pos"], out["size"] = item + else: + out[_KEY_RENAME.get(key, key)] = _rename_value(item) + return out + if isinstance(value, list): + return [_rename_value(item) for item in value] + return value + + +def _count_tree_nodes(value): + if isinstance(value, dict): + return 1 + sum(_count_tree_nodes(item) for item in value.values()) + if isinstance(value, list): + return sum(_count_tree_nodes(item) for item in value) + return 0 + + +# --------------------------------------------------------------------------- +# sharkd sessions (one resident process per source pcap) +# --------------------------------------------------------------------------- + +class _SharkdSession: + """A resident `sharkd -` process with one pcap loaded, addressed through + line-oriented JSON-RPC. Serialized by an asyncio lock (sharkd serves one + request at a time). Any transport-level failure aborts the session for + good (see _abort) — a timed-out or desynced session must never answer + later requests with shifted results.""" + + def __init__(self, pcap, scratch_dir, scratch, proc, stat): + self.pcap = pcap + self.scratch_dir = scratch_dir + self.scratch = scratch + self.proc = proc + self.mtime_ns = stat.st_mtime_ns + self.size = stat.st_size + self.last_used = time.monotonic() + self.lock = asyncio.Lock() + self._next_id = 0 + # Refcount of in-use holders (the manager's `session` context); the + # LRU cap only ever evicts sessions with zero uses, and a detached + # session (replaced in the manager, e.g. its pcap was rebuilt) closes + # itself when its last holder releases it. + self._uses = 0 + self._detached = False + + def matches(self, stat): + """True while the source pcap is byte-identical to what was loaded — + a fresh (mtime, size) would serve a rebuilt/truncated capture.""" + + return stat.st_mtime_ns == self.mtime_ns and stat.st_size == self.size + + def alive(self): + return self.proc.returncode is None + + def touch(self): + self.last_used = time.monotonic() + + async def rpc(self, method, params): + async with self.lock: + self._next_id += 1 + request_id = self._next_id + request = {"jsonrpc": "2.0", "id": request_id, "method": method, "params": params} + try: + self.proc.stdin.write((json.dumps(request) + "\n").encode()) + await self.proc.stdin.drain() + raw = await asyncio.wait_for(self.proc.stdout.readline(), RPC_TIMEOUT_SECONDS) + except asyncio.TimeoutError: + await self._abort(f"sharkd timed out after {RPC_TIMEOUT_SECONDS:.0f}s on {method!r}") + except (OSError, ValueError) as e: + # ValueError: a reply line over the StreamReader limit — the + # buffer is cleared, so the session is desynced either way. + await self._abort(f"sharkd session died on {method!r}: {e}") + if not raw: + await self._abort(f"sharkd closed the session during {method!r}") + try: + response = json.loads(raw) + except ValueError as e: + await self._abort(f"Malformed sharkd response: {e}") + # The echoed id proves this reply belongs to THIS request: a stale + # reply left over from a timed-out predecessor (or any desync) + # must never be served as fresh data. + if response.get("id") != request_id: + await self._abort(f"sharkd reply id mismatch on {method!r} (session desynchronized)") + if "error" in response: + error = response["error"] + raise _SharkdRpcError(error.get("code"), str(error.get("message", ""))) + return response.get("result") + + async def _abort(self, message): + """Kill the session for good (process + scratch copy) and raise — the + manager replaces it on the next acquire. Idempotent with close().""" + + await self.close() + raise SharkdError(message) + + async def close(self): + try: + if self.proc.returncode is None: + self.proc.kill() + # Bounded wait: a kill that fails to reap (blocked signals, + # mocked os.kill in tests, a wedged process) must never hang + # the caller — leak the process with a log instead. + try: + await asyncio.wait_for(self.proc.wait(), timeout=2.0) + except asyncio.TimeoutError: + log.warning("sharkd session for %s did not exit after kill", self.pcap) + except ProcessLookupError: + pass + shutil.rmtree(self.scratch_dir, ignore_errors=True) + + +class _SharkdManager: + """Resident sharkd sessions keyed by source pcap path. + + One asyncio lock covers check-spawn-insert-reap: the check-then-spawn + window spans the pcap copy, the fork and the load RPC, so without it two + concurrent requests for the same pcap double-spawn and the loser leaks + (process + scratch copy) forever. Sessions are refcounted while in use + (the `session` context manager) and the LRU cap only evicts IDLE + sessions — a tag with more sources than the cap, or a concurrent request, + can never have its session killed mid-RPC; the population is trimmed + back to the cap as uses drop (temporary overshoot is allowed).""" + + def __init__(self): + self._sessions = {} + self._mu = asyncio.Lock() + + @asynccontextmanager + async def session(self, pcap): + session = await self._acquire(pcap) + try: + yield session + finally: + await self._release(session) + + async def _acquire(self, pcap): + if shutil.which("sharkd") is None: + raise SharkdMissingError( + "sharkd is not available on this server — marker replay requires sharkd " + "(part of the Wireshark package)" + ) + try: + stat = os.stat(pcap) + except FileNotFoundError: + # Deleted between the directory listing and here (marker removed, + # project cleaned up) — not a server fault. + raise ControllerNotFoundError(f"Capture file {os.path.basename(pcap)} no longer exists") + async with self._mu: + session = self._sessions.get(pcap) + if session is not None and session.matches(stat) and session.alive(): + session._uses += 1 + session.touch() + return session + if session is not None: + # Rebuilt/truncated source or a dead session: replace it. An + # in-use one closes itself when its last holder releases it. + del self._sessions[pcap] + if session._uses == 0: + await session.close() + else: + session._detached = True + spawned = await self._spawn(pcap, stat) + spawned._uses += 1 + self._sessions[pcap] = spawned + victims = self._evict_locked() + for victim in victims: + await victim.close() + return spawned + + def _evict_locked(self): + """Caller holds ``_mu``. Pop (not close — that can block on the reap) + least-recently-used IDLE sessions down to the cap; if everything is + in use, allow the overshoot rather than killing a live request.""" + + victims = [] + while len(self._sessions) > SESSION_MAX: + idle = [s for s in self._sessions.values() if s._uses == 0] + if not idle: + break + victim = min(idle, key=lambda s: s.last_used) + del self._sessions[victim.pcap] + victims.append(victim) + return victims + + async def _release(self, session): + async with self._mu: + session._uses -= 1 + victims = [] + if session._uses == 0: + if session._detached: + # Replaced in the manager while still held — now orphaned. + victims.append(session) + elif not session.alive(): + # Failed mid-use (its RPC aborted it) — drop the corpse. + if self._sessions.get(session.pcap) is session: + del self._sessions[session.pcap] + else: + victims = self._evict_locked() + for victim in victims: + await victim.close() + + async def _spawn(self, pcap, stat): + scratch_dir, scratch = await _prepare_scratch(pcap) + try: + proc = await asyncio.create_subprocess_exec( + "sharkd", "-", + stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.DEVNULL, env=_engine_env(scratch_dir), + limit=_STREAM_LIMIT_BYTES, + ) + except OSError as e: + shutil.rmtree(scratch_dir, ignore_errors=True) + raise SharkdError(f"Could not run sharkd: {e}") + session = _SharkdSession(pcap, scratch_dir, scratch, proc, stat) + try: + await session.rpc("load", {"file": scratch}) + except Exception as e: + await session.close() + raise SharkdError(f"sharkd failed to load {os.path.basename(pcap)}: {e}") + return session + + async def close_all(self): + async with self._mu: + sessions = list(self._sessions.values()) + self._sessions.clear() + for session in sessions: + await session.close() + + +_manager = None + + +def _get_manager(): + global _manager + if _manager is None: + _manager = _SharkdManager() + return _manager + + +async def close_sessions(): + """Server shutdown hook: kill every resident session and drop its scratch + copy, so nothing leaks into /tmp across restarts.""" + + global _manager + manager = _manager + if manager is not None: + await manager.close_all() + _manager = None + + +# --------------------------------------------------------------------------- +# Columns + display filter (sharkd `frames` RPC) +# --------------------------------------------------------------------------- + +async def _columns_for(pcap, filter_expr): + """ + One resident `frames` pass over the source pcap. Returns + ``{frame_number: {src, dst, proto, info, bg, fg}}``; with a + ``filter_expr`` the keys are exactly the matching frame numbers, so one + pass both filters and enriches the merge. + + The column layout is pinned by the scratch-HOME preferences (see + _PINNED_COLUMNS): ``c[0]``=src, ``c[1]``=dst, ``c[2]``=proto, ``c[3]``=info + is a layout we own, not an environment default. The whole pagination + loop holds the session so the LRU cap cannot evict it between pages. + """ + + manager = _get_manager() + async with manager.session(pcap) as session: + columns = {} + skip = 0 + while True: + # sharkd rejects skip=0 ("must be a positive integer") — only send + # it once there is actually something to skip. + params = {"limit": _FRAMES_PAGE} + if skip: + params["skip"] = skip + if filter_expr is not None: + params["filter"] = filter_expr + try: + rows = await session.rpc("frames", params) + except _SharkdRpcError as e: + if filter_expr is not None and e.code == _ERR_INVALID_FILTER: + raise FilterError(f"Invalid display filter: {e.message}") + # Anything else is an engine fault, not the client's filter. + raise SharkdError(f"sharkd frames failed on {os.path.basename(pcap)}: {e.message}") + for row in rows or []: + try: + frame_number = int(row.get("num")) + except (TypeError, ValueError): + continue + cols = row.get("c") or [] + columns[frame_number] = { + "src": cols[0] or None if len(cols) > 0 else None, + "dst": cols[1] or None if len(cols) > 1 else None, + "proto": cols[2] or None if len(cols) > 2 else None, + "info": cols[3] or None if len(cols) > 3 else None, + "bg": row.get("bg"), + "fg": row.get("fg"), + } + if len(rows or []) < _FRAMES_PAGE: + return columns + skip += len(rows or []) + + # --------------------------------------------------------------------------- # Tag gate + timeline assembly # --------------------------------------------------------------------------- @@ -199,12 +635,24 @@ def gate_tag(project, tag): return entries -def _merged_frames(project, entries): +async def _merged_frames(project, entries, filter_expr=None, link_id=None): """ - Scan every source pcap and merge into one list sorted by + Scan every source pcap's record headers, ask sharkd for columns (and, + with a filter, the matching set), and merge into one list sorted by ``(ts, source file, frame number)`` — ts alone is not unique (two links can hit the same microsecond); the tiebreaker yields a stable, determined - order instead of a fictional one. + order instead of a fictional one. With a filter, only frames sharkd + matched survive, keeping their original pcap frame numbers. + + ``link_id`` narrows the frame stream to one capture source **before** any + engine work (a pure identity filter — only the selected link's pcap gets + a sharkd pass) and AND-composes with ``filter_expr``. An unknown link + matches nothing: an empty stream, same shape as a zero-match display + filter — deliberately not a 404. + + ``sources`` is the stable inventory of the tag: EVERY capture source is + listed with engine-free total counts, unaffected by ``link_id`` / + ``filter_expr`` — the inventory must not shrink when the view narrows. """ markers_dir = project.markers_directory @@ -215,18 +663,35 @@ def _merged_frames(project, entries): markers_dir, f"{entry['node_id']}_{entry['link_id']}_{entry['marker']}.pcap" ) frames = scan_pcap_frames(pcap) if os.path.exists(pcap) else [] + # Inventory first: every source, engine-free totals. sources.append({**{k: entry[k] for k in ("node_id", "link_id", "marker", "data_link_type")}, "count": len(frames)}) + if link_id and entry["link_id"] != link_id: + continue # link narrows the stream before any engine work + if frames: + columns = await _columns_for(pcap, filter_expr) + else: + columns = {} + source_key = f"{entry['node_id']}_{entry['link_id']}_{entry['marker']}" for frame_number, (sec, usec, incl_len) in enumerate(frames, start=1): + if filter_expr is not None and frame_number not in columns: + continue + cols = columns.get(frame_number, {}) merged.append({ "ts": _format_ts(sec, usec), "ts_us": sec * 1_000_000 + usec, + "_source": source_key, "len": incl_len, "node_id": entry["node_id"], "link_id": entry["link_id"], "marker": entry["marker"], "frame_number": frame_number, - "_source": f"{entry['node_id']}_{entry['link_id']}_{entry['marker']}", + "src": cols.get("src"), + "dst": cols.get("dst"), + "proto": cols.get("proto"), + "info": cols.get("info"), + "bg": cols.get("bg"), + "fg": cols.get("fg"), }) merged.sort(key=lambda f: (f["ts_us"], f["_source"], f["frame_number"])) for frame in merged: @@ -235,46 +700,48 @@ def _merged_frames(project, entries): return merged, sources -def build_timeline(project, tag, frame_cap=FRAME_LIST_CAP): +def _validate_filter(filter_expr): + if filter_expr is not None and len(filter_expr) > FILTER_MAX_LENGTH: + raise ControllerBadRequestError( + f"Display filter too long (max {FILTER_MAX_LENGTH} characters)" + ) + + +async def build_timeline(project, tag, filter_expr=None, link_id=None): """ - The ``range`` response: timeline bounds, per-source stats, and (under - ``frame_cap``) the full merged frame list for one-request timeline - layout. Over the cap the list is replaced by per-second buckets. + The ``range`` response: timeline bounds, per-source stats, and the full + merged frame list — deliberately uncapped (a flat list is the whole + contract; rendering a huge list is the client's concern, and the window + endpoint exists for incremental views). With a ``filter_expr`` and/or a + ``link_id`` every figure is computed on the matching frames only + (``sources`` stays the full tag inventory). """ + _validate_filter(filter_expr) entries = gate_tag(project, tag) - frames, sources = _merged_frames(project, entries) + frames, sources = await _merged_frames(project, entries, filter_expr, link_id=link_id) - response = { + return { "tag": tag, "start": frames[0]["ts"] if frames else None, "end": frames[-1]["ts"] if frames else None, "frame_count": len(frames), - "truncated": len(frames) > frame_cap, "sources": sources, + "frames": frames, } - if len(frames) <= frame_cap: - response["frames"] = frames - else: - buckets = {} - for frame in frames: - second = _parse_ts(frame["ts"]) // 1_000_000 - buckets[second] = buckets.get(second, 0) + 1 - response["buckets"] = [ - {"ts": _format_ts(second, 0), "count": count} - for second, count in sorted(buckets.items()) - ] - return response -def query_frames(project, tag, ts, window_ms=100, limit=1000): +async def query_frames(project, tag, ts, window_ms=100, limit=1000, filter_expr=None, link_id=None): """ - Frames with ts in ``[T, T+window_ms]`` merged across sources. A time with - no frames is a normal, successful answer — ``{"frames": []}``. + Frames with ts in ``[T, T+window_ms]`` merged across sources — narrowed + by ``filter_expr`` / ``link_id`` with the same semantics as the range + endpoint, so windowed seconds and the histogram always agree. A time + with no frames is a normal, successful answer — ``{"frames": []}``. """ + _validate_filter(filter_expr) entries = gate_tag(project, tag) - frames, _sources = _merged_frames(project, entries) + frames, _sources = await _merged_frames(project, entries, filter_expr, link_id=link_id) start_us = _parse_ts(ts) end_us = start_us + max(window_ms, 0) * 1000 @@ -283,70 +750,22 @@ def query_frames(project, tag, ts, window_ms=100, limit=1000): # --------------------------------------------------------------------------- -# Frame detail (tshark path — lazy, one frame per call) +# Frame detail (lazy — one frame per call, via the resident session) # --------------------------------------------------------------------------- -async def _tshark_version(): - try: - proc = await asyncio.create_subprocess_exec( - "tshark", "--version", - stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.DEVNULL, - ) - stdout, _ = await asyncio.wait_for(proc.communicate(), timeout=TSHARK_TIMEOUT_SECONDS) - return stdout.decode(errors="replace").splitlines()[0].strip() - except (OSError, asyncio.TimeoutError, IndexError): - raise TsharkMissingError("tshark is not available on this server") - - -def _tshark_scratch_copy(pcap): - """ - Copy the pcap to a scratch file under the system temp dir for tshark to - read. Hardened tshark profiles (AppArmor &c.) can deny it access to the - project directory / the user's home while still allowing /tmp — a real - copy, deliberately not a symlink, since the profile resolves real paths. - Caller must unlink the returned path. - """ - - fd, scratch = tempfile.mkstemp(suffix=".pcap", prefix="gns3-replay-") - os.close(fd) - shutil.copyfile(pcap, scratch) - return scratch - - -def _tshark_env(): - """Scratch HOME so tshark never even tries to read the user's home.""" - - env = dict(os.environ) - env["HOME"] = tempfile.gettempdir() - return env - - -def _pdml_to_nodes(element): - """ - Isomorphic PDML → JSON mapping: every XML attribute becomes a JSON key - verbatim (values stay strings), children nest under ``children``. The - element tag ("proto"/"field") is carried as ``element`` — the one - structural key beyond the attributes, so a renderer can tell a protocol - group from a leaf field (geninfo's tagless names make names unreliable). - """ - - return { - "element": element.tag, - **element.attrib, - "children": [_pdml_to_nodes(child) for child in element], - } - - -def _count_nodes(nodes): - return 1 + sum(_count_nodes(child) for child in nodes.get("children", [])) - - -async def decode_frame(project, tag, ts, node_id, link_id, marker): +async def decode_frame(project, tag, ts, node_id, link_id, marker, frame_number=None): """ Decode exactly one frame: locate its pcap by source identity, verify the round-tripped ts still matches the file (guards a rebuild between the - timeline view and this click), read the raw bytes for the hex view, and - map tshark's PDML of that single frame to JSON. + timeline view and this click), read the raw bytes for the hex view + straight from the pcap, and rename the sharkd protocol tree into the + REST contract (closed key set, values untouched). + + ``ts`` is not unique within one pcap (same-microsecond frames are kept + deliberately); an explicit ``frame_number`` from the frame list entry + disambiguates them and must still land on the exact ts. Without one the + first ts match decodes — fine unless two frames share a microsecond on + the same link. """ entries = gate_tag(project, tag) @@ -365,64 +784,45 @@ async def decode_frame(project, tag, ts, node_id, link_id, marker): raise ControllerNotFoundError(f"No capture file for marker '{marker}' (nothing ever matched)") frames = scan_pcap_frames(pcap) - # The ts must be the exact string the timeline returned; find the frame - # it identifies rather than trusting any position hint from the client. - frame_number = next( - (i for i, (sec, usec, _len) in enumerate(frames, start=1) - if _format_ts(sec, usec) == ts), - None, + rebuilt_message = ( + f"No frame at ts {ts} in marker '{marker}' (the capture may have been rebuilt)" ) if frame_number is None: - raise ControllerNotFoundError( - f"No frame at ts {ts} in marker '{marker}' (the capture may have been rebuilt)" + # The ts must be the exact string the timeline returned; find the frame + # it identifies rather than trusting any position hint from the client. + frame_number = next( + (i for i, (sec, usec, _len) in enumerate(frames, start=1) + if _format_ts(sec, usec) == ts), + None, ) + if frame_number is None: + raise ControllerNotFoundError(rebuilt_message) + else: + if not 1 <= frame_number <= len(frames): + raise ControllerNotFoundError(rebuilt_message) + sec, usec, _len = frames[frame_number - 1] + if _format_ts(sec, usec) != ts: + raise ControllerNotFoundError(rebuilt_message) raw_hex = read_frame_bytes(pcap, frame_number) + if raw_hex is None: + # The header scan listed it but the bytes are gone — mid-write tail. + raise ControllerNotFoundError(rebuilt_message) - if shutil.which("tshark") is None: - raise TsharkMissingError("tshark is not installed — frame detail is unavailable") - version = await _tshark_version() - - # Hand tshark a scratch copy under the temp dir: hardened profiles may - # deny it the project directory even though this process can read it - # (the hex view above reads the original directly). - scratch = _tshark_scratch_copy(pcap) - try: + async with _get_manager().session(pcap) as session: try: - proc = await asyncio.create_subprocess_exec( - "tshark", "-r", scratch, "-T", "pdml", "-Y", f"frame.number == {frame_number}", - stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, - env=_tshark_env(), - ) - stdout, stderr = await asyncio.wait_for( - proc.communicate(), timeout=TSHARK_TIMEOUT_SECONDS - ) - except OSError as e: - raise TsharkError(f"Could not run tshark: {e}") - except asyncio.TimeoutError: - raise TsharkError(f"tshark timed out after {TSHARK_TIMEOUT_SECONDS:.0f}s") - if proc.returncode != 0 or not stdout.strip(): - # Never feed truncated/failed output to the mapper. - raise TsharkError(f"tshark failed: {stderr.decode(errors='replace').strip()[:500]}") - finally: - try: - os.unlink(scratch) - except OSError: - pass + result = await session.rpc("frame", {"frame": frame_number, "proto": True}) + except _SharkdRpcError as e: + # The frame range was validated against the file above, so an + # engine error here is a real fault (502), not a client 404. + raise SharkdError(f"sharkd frame failed: {e.message}") - try: - root = ET.fromstring(stdout) - except ET.ParseError as e: - raise TsharkError(f"Malformed PDML from tshark: {e}") - - packet = root.find("./packet") - tree = [_pdml_to_nodes(child) for child in packet] if packet is not None else [] + tree = _rename_value(result.get("tree", [])) return { "ts": ts, "source": {"node_id": node_id, "link_id": link_id, "marker": marker, "frame_number": frame_number}, - "tshark_version": version, - "field_count": sum(_count_nodes(node) for node in tree), + "field_count": _count_tree_nodes(tree), "hex": raw_hex, "tree": tree, } diff --git a/gns3server/core/tasks.py b/gns3server/core/tasks.py index f8e8dadc9..0b745a067 100644 --- a/gns3server/core/tasks.py +++ b/gns3server/core/tasks.py @@ -111,6 +111,11 @@ async def shutdown(app: FastAPI) -> None: auto_discover_images_task_handle.cancel() await HTTPClient.close_session() await MarkerManager.instance().stop() + # Kill resident sharkd sessions (marker replay) and drop their /tmp + # scratch copies before the process exits. + from gns3server.controller import marker_replay + + await marker_replay.close_sessions() await Controller.instance().stop() for module in MODULES: diff --git a/tests/api/routes/controller/test_controller.py b/tests/api/routes/controller/test_controller.py index ddd33aa8e..43d672302 100644 --- a/tests/api/routes/controller/test_controller.py +++ b/tests/api/routes/controller/test_controller.py @@ -29,13 +29,18 @@ pytestmark = pytest.mark.asyncio class TestControllerRoutes: - async def test_shutdown_local(self, app: FastAPI, client: AsyncClient, config: Config) -> None: - - os.kill = MagicMock() + async def test_shutdown_local(self, app: FastAPI, client: AsyncClient, config: Config, monkeypatch) -> None: + + # monkeypatch (not bare assignment): a global `os.kill = MagicMock()` + # is never restored and poisons every later test that kills a + # subprocess — resident sharkd sessions would wait() forever on an + # immortal process. + kill_mock = MagicMock() + monkeypatch.setattr(os, "kill", kill_mock) config.settings.Server.local = True response = await client.post(app.url_path_for("shutdown")) assert response.status_code == status.HTTP_204_NO_CONTENT - assert os.kill.called + assert kill_mock.called async def test_shutdown_non_local(self, app: FastAPI, client: AsyncClient, config: Config) -> None: diff --git a/tests/api/routes/controller/test_marker_replay.py b/tests/api/routes/controller/test_marker_replay.py index 9a19e291d..b5db7398d 100644 --- a/tests/api/routes/controller/test_marker_replay.py +++ b/tests/api/routes/controller/test_marker_replay.py @@ -16,27 +16,39 @@ # along with this program. If not, see . """ -HTTP-route tests for the tag replay endpoints: the tag gate (409 while any -marker captures, 404 unknown tag), the merged timeline, window queries -(empty window = success), and the lazy frame detail (ts guard, isomorphic -JSON, 501 when tshark is unavailable). +HTTP-route tests for the tag replay endpoints (sharkd edition): the tag gate +(409 while any marker captures, 404 unknown tag), the merged timeline with +columns, window queries (empty window = success), the display-filter +parameter (400 with sharkd's error text on a bad expression), the lazy frame +detail with the renamed sharkd tree, and 501 when sharkd is unavailable +(hard engine requirement, no degraded mode). """ import shutil from unittest.mock import patch import pytest +import pytest_asyncio from fastapi import FastAPI, status from httpx import AsyncClient from gns3server.controller.project import Project from gns3server.controller.udp_link import UDPLink +from gns3server.controller import marker_replay -from tests.controller.test_marker_replay import _write_pcap, _icmp_frame +from tests.controller.test_marker_replay import _write_pcap, _icmp_frame, _tcp_syn_frame, _cols pytestmark = pytest.mark.asyncio -tshark_present = pytest.mark.skipif(shutil.which("tshark") is None, reason="tshark not installed") +sharkd_present = pytest.mark.skipif(shutil.which("sharkd") is None, reason="sharkd not installed") + + +@pytest_asyncio.fixture +async def no_residual_sessions(): + yield + manager = marker_replay._manager + if manager is not None: + await manager.close_all() def _add_marker(project, tag, enabled, node_id, frames=None): @@ -77,18 +89,48 @@ class TestReplayRoutes: ) assert response.status_code == status.HTTP_404_NOT_FOUND - async def test_range_merges_sources_in_ts_order(self, app: FastAPI, client: AsyncClient, project: Project) -> None: + async def test_range_501_without_sharkd(self, app: FastAPI, client: AsyncClient, project: Project) -> None: + + # A non-empty pcap: the engine must be consulted, and without sharkd + # the whole feature is unavailable (hard requirement, no degraded mode). + _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ + (1693472000, 123456, _icmp_frame()), + ]) + + with patch("gns3server.controller.marker_replay.shutil.which", return_value=None): + response = await client.get( + app.url_path_for("replay_tag_range", project_id=project.id, tag=7) + ) + assert response.status_code == status.HTTP_501_NOT_IMPLEMENTED + # The app's HTTPException handler unifies the body as {"message": …}. + assert "sharkd" in response.json()["message"] + + async def test_range_merges_sources_in_ts_order( + self, app: FastAPI, client: AsyncClient, project: Project, monkeypatch + ) -> None: + + # Columns injected hermetically — the merge/order/tiebreak contract + # must not depend on the engine being installed. + r1, r2 = UDPLink(project), UDPLink(project) + project._links.update({r1.id: r1, r2.id: r2}) + + def _wire(link, node_id, frames): + link._markers["icmp"] = {"bpf": "icmp", "tag": 7, "enabled": False, "color": None, + "highlight_duration": None, "capture_node_id": node_id, + "direction": None, "data_link_type": "DLT_EN10MB"} + _write_pcap(f"{project.markers_directory}/{node_id}_{link.id}_icmp.pcap", frames) # r1→r2 captures at t1 and t3; r2→r3 captures at t2 and t3 (same µs # as source A's t3 — the tiebreak must keep both frames). - _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ - (1693472000, 500000, b"a" * 60), - (1693472002, 0, b"a" * 60), - ]) - _add_marker(project, tag=7, enabled=False, node_id="n2", frames=[ - (1693472001, 0, b"b" * 60), - (1693472002, 0, b"b" * 60), - ]) + _wire(r1, "n1", [(1693472000, 500000, b"a" * 60), (1693472002, 0, b"a" * 60)]) + _wire(r2, "n2", [(1693472001, 0, b"b" * 60), (1693472002, 0, b"b" * 60)]) + + async def fake_columns(pcap, filter_expr): + name = pcap.rsplit("/", 1)[-1] + src = "10.0.0.1" if name.startswith("n1") else "10.0.0.2" + return {1: _cols(src=src), 2: _cols(src=src)} + + monkeypatch.setattr(marker_replay, "_columns_for", fake_columns) response = await client.get( app.url_path_for("replay_tag_range", project_id=project.id, tag=7) @@ -99,22 +141,30 @@ class TestReplayRoutes: assert body["frame_count"] == 4 assert body["start"] == "1693472000.500000" assert body["end"] == "1693472002.000000" - assert body["truncated"] is False assert [f["node_id"] for f in body["frames"]] == ["n1", "n2", "n1", "n2"] assert [f["ts"] for f in body["frames"]] == [ "1693472000.500000", "1693472001.000000", "1693472002.000000", "1693472002.000000", ] + # Wireshark-style columns ride along on every frame entry. + assert body["frames"][0]["src"] == "10.0.0.1" + assert body["frames"][1]["src"] == "10.0.0.2" + assert body["frames"][0]["proto"] == "ICMP" assert len(body["sources"]) == 2 async def test_frames_window_miss_is_empty_success( - self, app: FastAPI, client: AsyncClient, project: Project + self, app: FastAPI, client: AsyncClient, project: Project, monkeypatch ) -> None: _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ (1693472000, 0, b"a" * 60), ]) + async def fake_columns(pcap, filter_expr): + return {1: _cols()} + + monkeypatch.setattr(marker_replay, "_columns_for", fake_columns) + response = await client.get( app.url_path_for("replay_tag_frames", project_id=project.id, tag=7), params={"ts": "1693472001.000000", "window_ms": 100}, @@ -122,13 +172,20 @@ class TestReplayRoutes: assert response.status_code == status.HTTP_200_OK assert response.json() == {"frames": []} - async def test_frames_window_hit(self, app: FastAPI, client: AsyncClient, project: Project) -> None: + async def test_frames_window_hit( + self, app: FastAPI, client: AsyncClient, project: Project, monkeypatch + ) -> None: _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ (1693472000, 0, b"a" * 60), (1693472000, 150000, b"a" * 60), ]) + async def fake_columns(pcap, filter_expr): + return {1: _cols(), 2: _cols()} + + monkeypatch.setattr(marker_replay, "_columns_for", fake_columns) + response = await client.get( app.url_path_for("replay_tag_frames", project_id=project.id, tag=7), params={"ts": "1693472000.000000", "window_ms": 150}, @@ -151,22 +208,105 @@ class TestReplayRoutes: assert response.status_code == status.HTTP_404_NOT_FOUND assert "rebuilt" in response.json()["message"] - async def test_detail_501_without_tshark(self, app: FastAPI, client: AsyncClient, project: Project) -> None: + async def test_range_link_param_narrows_and_keeps_sources( + self, app: FastAPI, client: AsyncClient, project: Project, monkeypatch + ) -> None: - link = _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ + r1, r2 = UDPLink(project), UDPLink(project) + project._links.update({r1.id: r1, r2.id: r2}) + + def _wire(link, node_id, frames): + link._markers["icmp"] = {"bpf": "icmp", "tag": 7, "enabled": False, "color": None, + "highlight_duration": None, "capture_node_id": node_id, + "direction": None, "data_link_type": "DLT_EN10MB"} + _write_pcap(f"{project.markers_directory}/{node_id}_{link.id}_icmp.pcap", frames) + + _wire(r1, "n1", [(1693472000, 0, b"a" * 60), (1693472002, 0, b"a" * 60)]) + _wire(r2, "n2", [(1693472001, 0, b"b" * 60)]) + + async def fake_columns(pcap, filter_expr): + return {1: _cols(), 2: _cols()} + + monkeypatch.setattr(marker_replay, "_columns_for", fake_columns) + + response = await client.get( + app.url_path_for("replay_tag_range", project_id=project.id, tag=7), + params={"link": r1.id}, + ) + body = response.json() + assert body["frame_count"] == 2 + assert {f["link_id"] for f in body["frames"]} == {r1.id} + # The source dropdown keeps the whole tag inventory. + assert sorted(s["count"] for s in body["sources"]) == [1, 2] + + # Unknown link: 200 with an empty timeline, same shape as a zero-match + # filter — not a 404. + response = await client.get( + app.url_path_for("replay_tag_range", project_id=project.id, tag=7), + params={"link": "00000000-0000-0000-0000-000000000000"}, + ) + body = response.json() + assert body["frame_count"] == 0 and body["frames"] == [] and body["start"] is None + assert len(body["sources"]) == 2 + + @sharkd_present + async def test_range_columns_and_filter_end_to_end( + self, app: FastAPI, client: AsyncClient, project: Project, no_residual_sessions + ) -> None: + + _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ + (1693472000, 123456, _icmp_frame()), + (1693472001, 0, _tcp_syn_frame()), + ]) + + # Unfiltered: both frames with real engine columns. + response = await client.get( + app.url_path_for("replay_tag_range", project_id=project.id, tag=7) + ) + assert response.status_code == status.HTTP_200_OK + frames = response.json()["frames"] + assert [f["proto"] for f in frames] == ["ICMP", "TCP"] + assert frames[0]["src"] == "10.0.0.1" and frames[0]["dst"] == "10.0.0.3" + assert "Echo" in frames[0]["info"] + assert frames[0]["frame_number"] == 1 + + # Display filter applied before counting: only the TCP frame survives. + response = await client.get( + app.url_path_for("replay_tag_range", project_id=project.id, tag=7), + params={"filter": "tcp"}, + ) + body = response.json() + assert body["frame_count"] == 1 + assert [f["frame_number"] for f in body["frames"]] == [2] # original pcap identity + assert body["start"] == "1693472001.000000" + + @sharkd_present + async def test_range_invalid_filter_is_400_with_sharkd_text( + self, app: FastAPI, client: AsyncClient, project: Project, no_residual_sessions + ) -> None: + + _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ (1693472000, 123456, _icmp_frame()), ]) - with patch("gns3server.controller.marker_replay.shutil.which", return_value=None): - response = await client.get( - app.url_path_for("replay_tag_frame_detail", project_id=project.id, tag=7), - params={"ts": "1693472000.123456", "node_id": "n1", - "link_id": link.id, "marker": "icmp"}, - ) - assert response.status_code == status.HTTP_501_NOT_IMPLEMENTED + response = await client.get( + app.url_path_for("replay_tag_range", project_id=project.id, tag=7), + params={"filter": "this is (not valid"}, + ) + assert response.status_code == status.HTTP_400_BAD_REQUEST + assert "filter" in response.json()["message"].lower() - @tshark_present - async def test_detail_decodes_single_frame(self, app: FastAPI, client: AsyncClient, project: Project) -> None: + # Oversized filters are rejected before reaching the engine. + response = await client.get( + app.url_path_for("replay_tag_range", project_id=project.id, tag=7), + params={"filter": "x" * (marker_replay.FILTER_MAX_LENGTH + 1)}, + ) + assert response.status_code == status.HTTP_400_BAD_REQUEST + + @sharkd_present + async def test_detail_decodes_single_frame( + self, app: FastAPI, client: AsyncClient, project: Project, no_residual_sessions + ) -> None: link = _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ (1693472000, 123456, _icmp_frame()), @@ -182,9 +322,65 @@ class TestReplayRoutes: assert body["source"]["frame_number"] == 1 assert body["hex"] == _icmp_frame().hex() assert body["field_count"] > 10 - assert "tshark" in body["tshark_version"].lower() - ip = next(p for p in body["tree"] if p.get("name") == "ip") - ttl = next(f for f in ip["children"] if f.get("name") == "ip.ttl") - # Values arrive as strings, exactly as tshark emitted them. - assert ttl["show"] == "64" and ttl["showname"] == "Time to Live: 64" + def find(node, name): + stack = node if isinstance(node, list) else [node] + for child in stack: + if child.get("name") == name: + return child + deep = find(child.get("children", []), name) + if deep is not None: + return deep + return None + + ttl = find(body["tree"], "ip.ttl") + assert ttl["label"] == "Time to Live: 64" + assert ttl["filter_expr"] == "ip.ttl == 64" + assert ttl["pos"] == 22 and ttl["size"] == 1 + + @sharkd_present + async def test_detail_frame_number_disambiguates_same_ts( + self, app: FastAPI, client: AsyncClient, project: Project, no_residual_sessions + ) -> None: + + # Two frames in the same microsecond: only the explicit frame number + # (from the frame list entry) tells them apart. + link = _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ + (1693472000, 123456, _icmp_frame()), + (1693472000, 123456, _tcp_syn_frame()), + ]) + common = {"ts": "1693472000.123456", "node_id": "n1", + "link_id": link.id, "marker": "icmp"} + + response = await client.get( + app.url_path_for("replay_tag_frame_detail", project_id=project.id, tag=7), + params=common, + ) + assert response.status_code == status.HTTP_200_OK + assert response.json()["source"]["frame_number"] == 1 + assert response.json()["hex"] == _icmp_frame().hex() + + response = await client.get( + app.url_path_for("replay_tag_frame_detail", project_id=project.id, tag=7), + params={**common, "frame_number": 2}, + ) + assert response.status_code == status.HTTP_200_OK + assert response.json()["source"]["frame_number"] == 2 + assert response.json()["hex"] == _tcp_syn_frame().hex() + + @sharkd_present + async def test_detail_501_when_sharkd_disappears( + self, app: FastAPI, client: AsyncClient, project: Project + ) -> None: + + link = _add_marker(project, tag=7, enabled=False, node_id="n1", frames=[ + (1693472000, 123456, _icmp_frame()), + ]) + + with patch("gns3server.controller.marker_replay.shutil.which", return_value=None): + response = await client.get( + app.url_path_for("replay_tag_frame_detail", project_id=project.id, tag=7), + params={"ts": "1693472000.123456", "node_id": "n1", + "link_id": link.id, "marker": "icmp"}, + ) + assert response.status_code == status.HTTP_501_NOT_IMPLEMENTED diff --git a/tests/controller/test_marker_replay.py b/tests/controller/test_marker_replay.py index e33fb9a8c..19888bbd0 100644 --- a/tests/controller/test_marker_replay.py +++ b/tests/controller/test_marker_replay.py @@ -16,36 +16,52 @@ # along with this program. If not, see . """ -Unit tests for the tag-keyed aggregate replay module (controller layer): +Unit tests for the tag-keyed aggregate replay module (sharkd edition): * pcap record-header scanning (ts extraction, truncated-tail tolerance, nanosecond-magic normalization) and raw-bytes reads for the hex view -* the tag gate (404 unknown tag, 409 while any marker captures) -* timeline assembly: cross-source merge ordered by (ts, source, frame number), - same-microsecond tiebreak, missing-pcap sources, frame-cap degradation to - per-second buckets -* window queries (inclusive bounds, empty-window success) -* the lazy frame detail: round-tripped-ts guard, and — where tshark is - installed — the isomorphic PDML → JSON mapping with a round-trip check - (element count and attribute coverage). +* the tag gate (404 unknown tag, 409 while any marker captures) — engine-free +* timeline assembly with injected columns: cross-source merge ordered by + (ts, source, frame number), same-microsecond tiebreak, the uncapped full + frame list, display-filter application before count/slice +* the tree key renaming (closed census key set, values untouched, unknown + keys pass through verbatim, internal hf ids dropped) +* the resident sharkd sessions — spawn/load/reuse, (mtime, size) + invalidation, LRU bound (idle-only eviction), single spawn under + concurrent acquires, timeout/stale-reply aborts — against the real sharkd + where installed, plus the frame detail end-to-end (hex + renamed tree + + filter expressions, same-microsecond disambiguation). """ +import asyncio import os import shutil import struct +import sys +from types import SimpleNamespace +from unittest.mock import patch import pytest -from types import SimpleNamespace -from gns3server.controller.controller_error import ControllerError, ControllerNotFoundError +from gns3server.controller.controller_error import ( + ControllerBadRequestError, + ControllerError, + ControllerNotFoundError, +) +from gns3server.controller import marker_replay from gns3server.controller.marker_replay import ( + FilterError, + SharkdError, + SharkdMissingError, build_timeline, decode_frame, query_frames, read_frame_bytes, scan_pcap_frames, + _count_tree_nodes, _format_ts, _parse_ts, + _rename_value, ) pytestmark = pytest.mark.asyncio @@ -53,6 +69,8 @@ pytestmark = pytest.mark.asyncio PCAP_MAGIC_US = 0xA1B2C3D4 PCAP_MAGIC_NS = 0xA1B23C4D +sharkd_present = pytest.mark.skipif(shutil.which("sharkd") is None, reason="sharkd not installed") + def _write_pcap(path, frames, magic=PCAP_MAGIC_US, snaplen=65535): """frames: list of (sec, frac, payload bytes); frac is µs (or ns for the ns magic).""" @@ -64,27 +82,39 @@ def _write_pcap(path, frames, magic=PCAP_MAGIC_US, snaplen=65535): f.write(payload) +def _cksum(data): + if len(data) % 2: + data = data + b"\x00" # RFC 1071 odd-length padding + s = 0 + for i in range(0, len(data), 2): + s += (data[i] << 8) + data[i + 1] + while s >> 16: + s = (s & 0xFFFF) + (s >> 16) + return (~s) & 0xFFFF + + def _icmp_frame(): """A minimal well-formed ICMP echo request (10.0.0.1 → 10.0.0.3).""" - def cksum(data): - if len(data) % 2: - data = data + b"\x00" # RFC 1071 odd-length padding - s = 0 - for i in range(0, len(data), 2): - s += (data[i] << 8) + data[i + 1] - while s >> 16: - s = (s & 0xFFFF) + (s >> 16) - return (~s) & 0xFFFF - icmp = bytes([8, 0, 0, 0]) + struct.pack(">HHH", 1, 1, 0) + b"payload12" - icmp = icmp[:2] + struct.pack(">H", cksum(icmp)) + icmp[4:] + icmp = icmp[:2] + struct.pack(">H", _cksum(icmp)) + icmp[4:] ip0 = struct.pack(">BBHHHBBH4s4s", 0x45, 0, 20 + len(icmp), 1, 0, 64, 1, 0, bytes([10, 0, 0, 1]), bytes([10, 0, 0, 3])) - ip = ip0[:10] + struct.pack(">H", cksum(ip0)) + ip0[12:] + ip = ip0[:10] + struct.pack(">H", _cksum(ip0)) + ip0[12:] return bytes.fromhex("0200000000020200000000010800") + ip + icmp +def _tcp_syn_frame(): + """A minimal TCP SYN (10.0.0.1:472 → 10.0.0.3:22).""" + + tcp = struct.pack(">HHIIBBHHH", 472, 22, 0, 0, 0x50, 0x02, 64240, 0, 0) + ip0 = struct.pack(">BBHHHBBH4s4s", 0x45, 0, 20 + len(tcp), 2, 0, 64, 6, 0, + bytes([10, 0, 0, 1]), bytes([10, 0, 0, 3])) + ip = ip0[:10] + struct.pack(">H", _cksum(ip0)) + ip0[12:] + tcp = tcp[:16] + struct.pack(">H", _cksum(ip0[12:] + tcp)) + tcp[18:] + return bytes.fromhex("0200000000020200000000010800") + ip + tcp + + def _fake_project(tmp_path, markers, markers_dir=None): """markers: the flat project.markers shape ({'link/name': {..., node_id}}).""" @@ -98,8 +128,56 @@ def _marker_entry(tag, enabled=True, node_id="node-1"): "node_id": node_id} +def _fake_columns(monkeypatch, mapping): + """Inject canned sharkd columns: mapping pcap-basename → {frame#: cols or None}. + + Frames absent from the dict get no columns; when the caller passes a + filter, the injected set IS the matching set (mirroring the real engine). + """ + + async def fake_columns_for(pcap, filter_expr): + return mapping.get(os.path.basename(pcap), {}) + + monkeypatch.setattr(marker_replay, "_columns_for", fake_columns_for) + + +def _cols(src="10.0.0.1", dst="10.0.0.3", proto="ICMP", info="Echo (ping) request"): + return {"src": src, "dst": dst, "proto": proto, "info": info, + "bg": "ffffff", "fg": "000000"} + + +class _FakeSession: + """Duck-typed _SharkdSession for manager-level tests (no engine).""" + + def __init__(self, pcap): + self.pcap = pcap + self.last_used = 0 + self.closed = False + self._uses = 0 + self._detached = False + self.mtime_ns, self.size = 0, 0 + + def matches(self, stat): + return True + + def alive(self): + return not self.closed + + def touch(self): + pass + + async def close(self): + self.closed = True + + +def _pretend_sharkd(monkeypatch): + """The manager-level tests run entirely against fakes — the sharkd binary + must not be a precondition (the suite stays green without the engine).""" + monkeypatch.setattr(marker_replay.shutil, "which", lambda name: "/usr/bin/sharkd") + + # --------------------------------------------------------------------------- -# pcap scanning +# pcap scanning (engine-free backbone) # --------------------------------------------------------------------------- class TestScanPcap: @@ -140,15 +218,15 @@ class TestScanPcap: # --------------------------------------------------------------------------- -# Tag gate + timeline +# Tag gate (engine-free — raised before any sharkd work) # --------------------------------------------------------------------------- -class TestGateAndTimeline: +class TestGate: async def test_unknown_tag_404(self, tmp_path): project = _fake_project(tmp_path, {"linkA/icmp": _marker_entry(tag=1)}) with pytest.raises(ControllerNotFoundError): - build_timeline(project, tag=7) + await build_timeline(project, tag=7) async def test_gate_409_while_capturing(self, tmp_path): project = _fake_project(tmp_path, { @@ -156,9 +234,22 @@ class TestGateAndTimeline: "linkB/icmp": _marker_entry(tag=7, enabled=True, node_id="n2"), }) with pytest.raises(ControllerError, match="linkB"): - build_timeline(project, tag=7) + await build_timeline(project, tag=7) - async def test_merge_orders_by_ts_with_stable_tiebreak(self, tmp_path): + async def test_filter_length_capped(self, tmp_path, monkeypatch): + _fake_columns(monkeypatch, {}) + project = _fake_project(tmp_path, {"linkA/icmp": _marker_entry(tag=7, enabled=False)}) + with pytest.raises(ControllerBadRequestError, match="too long"): + await build_timeline(project, tag=7, filter_expr="x" * (marker_replay.FILTER_MAX_LENGTH + 1)) + + +# --------------------------------------------------------------------------- +# Timeline assembly (columns injected — no engine needed) +# --------------------------------------------------------------------------- + +class TestTimeline: + + async def test_merge_orders_by_ts_with_stable_tiebreak(self, tmp_path, monkeypatch): # Two sources, deliberately interleaved in time, colliding on one µs. _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [ (1693472000, 500000, b"a" * 60), # t1 sourceA @@ -168,12 +259,16 @@ class TestGateAndTimeline: (1693472001, 000000, b"b" * 60), # t2 sourceB (1693472002, 000000, b"b" * 60), # t3 sourceB — same µs as t3 sourceA ]) + _fake_columns(monkeypatch, { + "n1_linkA_icmp.pcap": {1: _cols(), 2: _cols()}, + "n2_linkB_icmp.pcap": {1: _cols(src="10.0.0.2"), 2: _cols(src="10.0.0.2")}, + }) project = _fake_project(tmp_path, { "linkA/icmp": _marker_entry(tag=7, enabled=False, node_id="n1"), "linkB/icmp": _marker_entry(tag=7, enabled=False, node_id="n2"), }) - timeline = build_timeline(project, tag=7) + timeline = await build_timeline(project, tag=7) assert timeline["frame_count"] == 4 assert timeline["start"] == "1693472000.500000" assert timeline["end"] == "1693472002.000000" @@ -181,65 +276,410 @@ class TestGateAndTimeline: # Same-microsecond pair keeps both frames (a ts dict key would drop one). assert [f["ts"] for f in timeline["frames"]][2:] == ["1693472002.000000"] * 2 assert [f["frame_number"] for f in timeline["frames"]] == [1, 1, 2, 2] + # Columns ride along verbatim. + assert timeline["frames"][0]["src"] == "10.0.0.1" + assert timeline["frames"][1]["src"] == "10.0.0.2" + assert timeline["frames"][0]["proto"] == "ICMP" assert {s["count"] for s in timeline["sources"]} == {2} - async def test_missing_pcap_is_zero_count_source(self, tmp_path): + async def test_columns_missing_for_a_frame_still_lists_it(self, tmp_path, monkeypatch): + # A frame the (injected) engine did not describe keeps its place with + # null columns — the timeline backbone never depends on the engine. + _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [(1693472000, 0, b"a" * 60)]) + _fake_columns(monkeypatch, {}) # engine describes nothing + project = _fake_project(tmp_path, {"linkA/icmp": _marker_entry(tag=7, enabled=False, node_id="n1")}) + + timeline = await build_timeline(project, tag=7) + frame = timeline["frames"][0] + assert frame["ts"] == "1693472000.000000" + assert frame["src"] is None and frame["info"] is None + + async def test_missing_pcap_is_zero_count_source(self, tmp_path, monkeypatch): + _fake_columns(monkeypatch, {}) project = _fake_project(tmp_path, {"linkA/icmp": _marker_entry(tag=7, enabled=False)}) - timeline = build_timeline(project, tag=7) + timeline = await build_timeline(project, tag=7) assert timeline["frame_count"] == 0 assert timeline["start"] is None and timeline["end"] is None assert timeline["frames"] == [] assert timeline["sources"][0]["count"] == 0 - async def test_over_cap_degrades_to_buckets(self, tmp_path): + async def test_frame_list_is_uncapped(self, tmp_path, monkeypatch): + # The list is the whole contract — no truncation flag, no buckets, + # however many frames the tag holds. _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [ (1693472000, 0, b"a" * 60), (1693472000, 500000, b"a" * 60), (1693472001, 0, b"a" * 60), ]) + _fake_columns(monkeypatch, {"n1_linkA_icmp.pcap": {1: _cols(), 2: _cols(), 3: _cols()}}) project = _fake_project(tmp_path, {"linkA/icmp": _marker_entry(tag=7, enabled=False, node_id="n1")}) - timeline = build_timeline(project, tag=7, frame_cap=2) - assert timeline["truncated"] is True - assert "frames" not in timeline - assert timeline["buckets"] == [ - {"ts": "1693472000.000000", "count": 2}, - {"ts": "1693472001.000000", "count": 1}, + timeline = await build_timeline(project, tag=7) + assert timeline["frame_count"] == 3 + assert len(timeline["frames"]) == 3 + assert "truncated" not in timeline and "buckets" not in timeline + + async def test_filter_applies_before_count_and_slice(self, tmp_path, monkeypatch): + # Three frames; the injected "matching set" (what a real engine would + # return for the filter) contains only frames 1 and 3. + _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [ + (1693472000, 0, b"a" * 60), + (1693472001, 0, b"a" * 60), + (1693472002, 0, b"a" * 60), + ]) + _fake_columns(monkeypatch, {"n1_linkA_icmp.pcap": {1: _cols(), 3: _cols(proto="TCP")}}) + project = _fake_project(tmp_path, {"linkA/icmp": _marker_entry(tag=7, enabled=False, node_id="n1")}) + + timeline = await build_timeline(project, tag=7, filter_expr="tcp") + assert timeline["frame_count"] == 2 + assert timeline["start"] == "1693472000.000000" + assert timeline["end"] == "1693472002.000000" + # frame numbers keep their ORIGINAL pcap identity through the filter. + assert [f["frame_number"] for f in timeline["frames"]] == [1, 3] + # sources[] is the stable inventory: engine-free TOTAL counts, + # unaffected by link/filter (the WebUI source dropdown must not + # shrink when the view narrows). + assert timeline["sources"][0]["count"] == 3 + + +class TestLinkFilter: + """``link_id`` narrows the frame stream to one capture source before + counting/slicing/bucketing; sources stay the full inventory.""" + + def _project(self, tmp_path, monkeypatch, columns_override=None): + _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [ + (1693472000, 500000, b"a" * 60), + (1693472002, 000000, b"a" * 60), + ]) + _write_pcap(tmp_path / "n2_linkB_icmp.pcap", [ + (1693472001, 000000, b"b" * 60), + (1693472002, 000000, b"b" * 60), + ]) + + async def fake_columns(pcap, filter_expr): + if columns_override is not None: + return await columns_override(os.path.basename(pcap), filter_expr) + src = "10.0.0.1" if "n1" in pcap else "10.0.0.2" + return {1: _cols(src=src), 2: _cols(src=src)} + + monkeypatch.setattr(marker_replay, "_columns_for", fake_columns) + return _fake_project(tmp_path, { + "linkA/icmp": _marker_entry(tag=7, enabled=False, node_id="n1"), + "linkB/icmp": _marker_entry(tag=7, enabled=False, node_id="n2"), + }) + + async def test_link_narrows_before_count_and_slice(self, tmp_path, monkeypatch): + timeline = await build_timeline(self._project(tmp_path, monkeypatch), tag=7, link_id="linkA") + assert timeline["frame_count"] == 2 + assert timeline["start"] == "1693472000.500000" + assert [f["link_id"] for f in timeline["frames"]] == ["linkA", "linkA"] + # sources stay the FULL inventory with engine-free totals. + assert sorted((s["link_id"], s["count"]) for s in timeline["sources"]) == [ + ("linkA", 2), ("linkB", 2) ] + async def test_unknown_link_is_empty_success(self, tmp_path, monkeypatch): + timeline = await build_timeline(self._project(tmp_path, monkeypatch), tag=7, link_id="nope") + assert timeline["frame_count"] == 0 + assert timeline["start"] is None and timeline["end"] is None + assert timeline["frames"] == [] + assert len(timeline["sources"]) == 2 + + async def test_link_and_filter_compose_as_and(self, tmp_path, monkeypatch): + # The injected "matching set" contains only frame 2 per source — + # combined with link=linkA the view is exactly that one frame. + async def override(basename, filter_expr): + assert filter_expr == "tcp" + return {2: _cols(proto="TCP")} + + project = self._project(tmp_path, monkeypatch, columns_override=override) + timeline = await build_timeline(project, tag=7, filter_expr="tcp", link_id="linkA") + assert timeline["frame_count"] == 1 + assert timeline["frames"][0]["frame_number"] == 2 + assert timeline["frames"][0]["link_id"] == "linkA" + + async def test_empty_link_string_means_absent(self, tmp_path, monkeypatch): + timeline = await build_timeline(self._project(tmp_path, monkeypatch), tag=7, link_id="") + assert timeline["frame_count"] == 4 + + async def test_query_frames_window_over_link_stream(self, tmp_path, monkeypatch): + project = self._project(tmp_path, monkeypatch) + result = await query_frames(project, tag=7, ts="1693472002.000000", + window_ms=0, link_id="linkB") + assert [f["link_id"] for f in result["frames"]] == ["linkB"] -# --------------------------------------------------------------------------- -# Window query -# --------------------------------------------------------------------------- class TestQueryFrames: - def _project(self, tmp_path): + def _project(self, tmp_path, monkeypatch): _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [ (1693472000, 0, b"a" * 60), (1693472000, 150000, b"a" * 60), (1693472005, 0, b"a" * 60), ]) + _fake_columns(monkeypatch, {"n1_linkA_icmp.pcap": {i: _cols() for i in (1, 2, 3)}}) return _fake_project(tmp_path, {"linkA/icmp": _marker_entry(tag=7, enabled=False, node_id="n1")}) - async def test_window_inclusive_bounds(self, tmp_path): - result = query_frames(self._project(tmp_path), tag=7, ts="1693472000.000000", window_ms=150) + async def test_window_inclusive_bounds(self, tmp_path, monkeypatch): + result = await query_frames(self._project(tmp_path, monkeypatch), tag=7, + ts="1693472000.000000", window_ms=150) assert [f["ts"] for f in result["frames"]] == ["1693472000.000000", "1693472000.150000"] - async def test_window_miss_is_empty_success(self, tmp_path): - result = query_frames(self._project(tmp_path), tag=7, ts="1693472001.000000", window_ms=100) + async def test_window_miss_is_empty_success(self, tmp_path, monkeypatch): + result = await query_frames(self._project(tmp_path, monkeypatch), tag=7, + ts="1693472001.000000", window_ms=100) assert result == {"frames": []} - async def test_limit_applies(self, tmp_path): - result = query_frames(self._project(tmp_path), tag=7, ts="1693472000.000000", - window_ms=150, limit=1) + async def test_limit_applies(self, tmp_path, monkeypatch): + result = await query_frames(self._project(tmp_path, monkeypatch), tag=7, + ts="1693472000.000000", window_ms=150, limit=1) assert len(result["frames"]) == 1 # --------------------------------------------------------------------------- -# Frame detail (tshark path) +# Tree key renaming # --------------------------------------------------------------------------- -tshark_present = pytest.mark.skipif(shutil.which("tshark") is None, reason="tshark not installed") +class TestRename: + + async def test_renames_closed_key_set_and_drops_hf_id(self): + node = { + "t": "proto", "l": "Time to Live: 64", "fn": "ip.ttl", + "f": "ip.ttl == 64", "h": [22, 1], "s": None, "g": False, + "e": 8472, + "n": [{"l": "nested", "h": [23, 2], "n": []}], + } + renamed = _rename_value(node) + assert renamed == { + "element": "proto", "label": "Time to Live: 64", "name": "ip.ttl", + "filter_expr": "ip.ttl == 64", "pos": 22, "size": 1, + "expert": None, "generated": False, + "children": [{"label": "nested", "pos": 23, "size": 2, "children": []}], + } + + async def test_unknown_keys_pass_through_verbatim(self): + # A future Wireshark adding a key must never silently lose data. + node = {"l": "x", "future_key": {"deep": [1, 2]}, "n": []} + renamed = _rename_value(node) + assert renamed["future_key"] == {"deep": [1, 2]} + + async def test_count_tree_nodes_counts_dicts_only(self): + assert _count_tree_nodes({"a": [{"b": 1}, "str", 3]}) == 2 + + +# --------------------------------------------------------------------------- +# sharkd sessions + engine-backed behaviour +# --------------------------------------------------------------------------- + +class TestSessions: + + async def test_missing_sharkd_raises_501_error(self, tmp_path): + _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [(1693472000, 0, b"a" * 60)]) + project = _fake_project(tmp_path, {"linkA/icmp": _marker_entry(tag=7, enabled=False, node_id="n1")}) + with patch("gns3server.controller.marker_replay.shutil.which", return_value=None): + with pytest.raises(SharkdMissingError): + await build_timeline(project, tag=7) + + async def test_cap_evicts_idle_lru_only(self, tmp_path, monkeypatch): + _pretend_sharkd(monkeypatch) + manager = marker_replay._SharkdManager() + fakes = {} + + async def fake_spawn(pcap, stat): + session = _FakeSession(pcap) + fakes[pcap] = session + return session + + with patch.object(manager, "_spawn", side_effect=fake_spawn): + # cap + 2 distinct pcaps, each acquired and released (idle again). + for i in range(marker_replay.SESSION_MAX + 2): + pcap = tmp_path / f"pcap{i}" + _write_pcap(pcap, [(1693472000, 0, b"a" * 60)]) + async with manager.session(str(pcap)): + pass + # Bounded to SESSION_MAX; the earliest (least recently used) closed. + assert len(manager._sessions) == marker_replay.SESSION_MAX + assert fakes[str(tmp_path / "pcap0")].closed is True + assert fakes[str(tmp_path / "pcap1")].closed is True + assert fakes[str(tmp_path / "pcap2")].closed is False + + async def test_in_use_session_survives_cap_pressure(self, tmp_path, monkeypatch): + _pretend_sharkd(monkeypatch) + manager = marker_replay._SharkdManager() + + async def fake_spawn(pcap, stat): + return _FakeSession(pcap) + + with patch.object(manager, "_spawn", side_effect=fake_spawn): + pcap0 = tmp_path / "pcap0" + _write_pcap(pcap0, [(1693472000, 0, b"a" * 60)]) + async with manager.session(str(pcap0)) as held: + # More sources than the cap while pcap0 is held open. + for i in range(1, marker_replay.SESSION_MAX + 4): + pcap = tmp_path / f"pcap{i}" + _write_pcap(pcap, [(1693472000, 0, b"a" * 60)]) + async with manager.session(str(pcap)): + pass + # A held session is never the eviction victim (no mid-RPC kill). + assert held.closed is False + assert str(pcap0) in manager._sessions + # Released → the population is trimmed back under the cap. + assert len(manager._sessions) <= marker_replay.SESSION_MAX + await manager.close_all() + + async def test_concurrent_acquire_spawns_once(self, tmp_path, monkeypatch): + _pretend_sharkd(monkeypatch) + manager = marker_replay._SharkdManager() + spawns = [] + + async def slow_spawn(pcap, stat): + spawns.append(pcap) + await asyncio.sleep(0.05) # widen the check-then-spawn window + return _FakeSession(pcap) + + pcap = tmp_path / "n1_linkA_icmp.pcap" + _write_pcap(pcap, [(1693472000, 0, b"a" * 60)]) + with patch.object(manager, "_spawn", side_effect=slow_spawn): + first, second = await asyncio.gather( + manager._acquire(str(pcap)), manager._acquire(str(pcap)) + ) + assert first is second # concurrent requests share one spawn + assert len(spawns) == 1 + await manager.close_all() + + async def test_rpc_timeout_kills_session(self, tmp_path, monkeypatch): + monkeypatch.setattr(marker_replay, "RPC_TIMEOUT_SECONDS", 0.2) + pcap = tmp_path / "n1_linkA_icmp.pcap" + _write_pcap(pcap, [(1693472000, 0, b"a" * 60)]) + # A process that never answers: the session must die (never serve the + # late reply to a later request) instead of staying resident. + proc = await asyncio.create_subprocess_exec( + "sleep", "60", + stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, + ) + session = marker_replay._SharkdSession( + str(pcap), str(tmp_path / "scratch"), "unused", proc, os.stat(str(pcap)) + ) + with pytest.raises(SharkdError, match="timed out"): + await session.rpc("frames", {}) + assert session.alive() is False + assert proc.returncode is not None + + async def test_stale_reply_id_aborts_session(self, tmp_path): + # A "sharkd" that always answers with a foreign id — every reply is a + # stale one as far as the caller is concerned; serving it would be + # shifted data, so the session must abort instead. + fake = tmp_path / "fake_sharkd.py" + fake.write_text( + "import json, sys\n" + "for line in sys.stdin:\n" + " json.loads(line)\n" + " sys.stdout.write(json.dumps({'jsonrpc': '2.0', 'id': 424242, 'result': []}) + '\\n')\n" + " sys.stdout.flush()\n" + ) + pcap = tmp_path / "n1_linkA_icmp.pcap" + _write_pcap(pcap, [(1693472000, 0, b"a" * 60)]) + proc = await asyncio.create_subprocess_exec( + sys.executable, str(fake), + stdin=asyncio.subprocess.PIPE, stdout=asyncio.subprocess.PIPE, + ) + session = marker_replay._SharkdSession( + str(pcap), str(tmp_path / "scratch"), "unused", proc, os.stat(str(pcap)) + ) + with pytest.raises(SharkdError, match="id mismatch"): + await session.rpc("frames", {}) + assert session.alive() is False + assert proc.returncode is not None + + async def test_filter_error_only_for_the_filter_code(self, tmp_path, monkeypatch): + _pretend_sharkd(monkeypatch) + manager = marker_replay._SharkdManager() + monkeypatch.setattr(marker_replay, "_manager", manager) + pcap = tmp_path / "n1_linkA_icmp.pcap" + _write_pcap(pcap, [(1693472000, 0, b"a" * 60)]) + + class _RpcFail(_FakeSession): + def __init__(self, pcap, code): + super().__init__(pcap) + self.code = code + + async def rpc(self, method, params): + raise marker_replay._SharkdRpcError(self.code, "boom") + + manager._sessions[str(pcap)] = _RpcFail(str(pcap), marker_replay._ERR_INVALID_FILTER) + with pytest.raises(FilterError): + await marker_replay._columns_for(str(pcap), "icmp") + # Same failure WITHOUT a filter is an engine fault, not a 400. + manager._sessions[str(pcap)] = _RpcFail(str(pcap), marker_replay._ERR_INVALID_FILTER) + with pytest.raises(SharkdError): + await marker_replay._columns_for(str(pcap), None) + # A non-filter engine error with a filter set is NOT the client's fault. + manager._sessions[str(pcap)] = _RpcFail(str(pcap), -9999) + with pytest.raises(SharkdError): + await marker_replay._columns_for(str(pcap), "icmp") + + @sharkd_present + async def test_real_pagination_over_1000_frames(self, tmp_path): + # A full page measures well over the 64 KB default StreamReader limit + # (~190 KB with realistic columns) — pagination must not blow up the + # stream (nor desync the resident session). + pcap = tmp_path / "big.pcap" + _write_pcap(pcap, [(1693472000 + i, 0, _icmp_frame()) for i in range(1001)]) + manager = marker_replay._get_manager() + try: + columns = await marker_replay._columns_for(str(pcap), None) + assert len(columns) == 1001 + assert columns[1001]["src"] == "10.0.0.1" + assert columns[1001]["proto"] == "ICMP" + finally: + await manager.close_all() + + @sharkd_present + async def test_real_session_columns_and_filter(self, tmp_path): + _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [ + (1693472000, 123456, _icmp_frame()), + (1693472001, 0, _tcp_syn_frame()), + ]) + manager = marker_replay._get_manager() + try: + columns = await marker_replay._columns_for(str(tmp_path / "n1_linkA_icmp.pcap"), None) + assert columns[1]["src"] == "10.0.0.1" + assert columns[1]["proto"] == "ICMP" + assert "Echo" in columns[1]["info"] + assert columns[2]["proto"] == "TCP" + + only_tcp = await marker_replay._columns_for(str(tmp_path / "n1_linkA_icmp.pcap"), "tcp") + assert set(only_tcp) == {2} + finally: + await manager.close_all() + + @sharkd_present + async def test_real_session_invalid_filter_raises_filter_error(self, tmp_path): + _write_pcap(tmp_path / "n1_linkA_icmp.pcap", [(1693472000, 123456, _icmp_frame())]) + manager = marker_replay._get_manager() + try: + with pytest.raises(FilterError): + await marker_replay._columns_for(str(tmp_path / "n1_linkA_icmp.pcap"), "this is (not valid") + finally: + await manager.close_all() + + @sharkd_present + async def test_real_session_respawns_on_stat_change(self, tmp_path): + pcap = tmp_path / "n1_linkA_icmp.pcap" + _write_pcap(pcap, [(1693472000, 123456, _icmp_frame())]) + manager = marker_replay._get_manager() + try: + async with manager.session(str(pcap)) as first: + pass + async with manager.session(str(pcap)) as warm: + assert warm is first # warm reuse + # Rewrite the file (simulating a truncated/rebuilt capture). + _write_pcap(pcap, [(1693473000, 0, _icmp_frame())]) + async with manager.session(str(pcap)) as second: + assert second is not first + columns = await marker_replay._columns_for(str(pcap), None) + assert set(columns) == {1} + finally: + await manager.close_all() class TestDecodeFrame: @@ -264,83 +704,113 @@ class TestDecodeFrame: await decode_frame(project, tag=7, ts="1693472000.123456", node_id="nobody", link_id="linkA", marker="icmp") - async def test_decode_feeds_tshark_a_scratch_copy(self, tmp_path): - """Hardened tshark profiles deny the project dir — tshark must read a - /tmp copy (a real copy, not a symlink) that is unlinked afterwards.""" - - import tempfile - from unittest.mock import patch, AsyncMock - - observed = [] - PDML = (b'' - b'' - b'') - - class FakeProc: - returncode = 0 - - async def communicate(self): - return PDML, b"" - - async def fake_exec(*args, **kwargs): - r_index = args.index("-r") - observed.append((args[r_index + 1], kwargs.get("env"))) - return FakeProc() - + async def test_explicit_frame_number_out_of_range_404(self, tmp_path): project = self._project(tmp_path) - with patch("gns3server.controller.marker_replay.shutil.which", return_value="tshark"), \ - patch("gns3server.controller.marker_replay._tshark_version", AsyncMock(return_value="tshark 4.6.7")), \ - patch("gns3server.controller.marker_replay.asyncio.create_subprocess_exec", side_effect=fake_exec): + with pytest.raises(ControllerNotFoundError, match="rebuilt"): + await decode_frame(project, tag=7, ts="1693472000.123456", + node_id="n1", link_id="linkA", marker="icmp", frame_number=5) + + async def test_explicit_frame_number_ts_mismatch_404(self, tmp_path): + # The frame number must still land on the exact round-tripped ts — + # a rebuilt capture cannot be decoded by stale coordinates. + project = self._project(tmp_path) + with pytest.raises(ControllerNotFoundError, match="rebuilt"): + await decode_frame(project, tag=7, ts="1693472001.000000", + node_id="n1", link_id="linkA", marker="icmp", frame_number=1) + + async def test_hex_read_failure_is_404_not_null_hex(self, tmp_path, monkeypatch): + project = self._project(tmp_path) + monkeypatch.setattr(marker_replay, "read_frame_bytes", lambda path, n: None) + with pytest.raises(ControllerNotFoundError, match="rebuilt"): + await decode_frame(project, tag=7, ts="1693472000.123456", + node_id="n1", link_id="linkA", marker="icmp") + + @sharkd_present + async def test_same_microsecond_frames_disambiguated_by_frame_number(self, tmp_path): + # ts is not unique within one pcap: two frames share a microsecond, + # and only the explicit frame number tells them apart. + pcap = tmp_path / "n1_linkA_icmp.pcap" + icmp, tcp = _icmp_frame(), _tcp_syn_frame() + _write_pcap(pcap, [(1693472000, 123456, icmp), (1693472000, 123456, tcp)]) + project = _fake_project(tmp_path, { + "linkA/icmp": _marker_entry(tag=7, enabled=False, node_id="n1"), + }) + manager = marker_replay._get_manager() + try: + default = await decode_frame(project, tag=7, ts="1693472000.123456", + node_id="n1", link_id="linkA", marker="icmp") + second = await decode_frame(project, tag=7, ts="1693472000.123456", + node_id="n1", link_id="linkA", marker="icmp", + frame_number=2) + finally: + await manager.close_all() + # Without a frame number: first ts match. + assert default["source"]["frame_number"] == 1 + assert default["hex"] == icmp.hex() + # With one: exactly that frame's bytes and tree. + assert second["source"]["frame_number"] == 2 + assert second["hex"] == tcp.hex() + assert second["field_count"] > 10 + + @sharkd_present + async def test_decode_end_to_end(self, tmp_path): + manager = marker_replay._get_manager() + try: + project = self._project(tmp_path) detail = await decode_frame(project, tag=7, ts="1693472000.123456", node_id="n1", link_id="linkA", marker="icmp") - - assert detail["field_count"] == 2 # proto + field from the canned PDML - (scratch, env), = observed - original = str(tmp_path / "n1_linkA_icmp.pcap") - assert scratch != original - assert scratch.startswith(tempfile.gettempdir()) and scratch.endswith(".pcap") - assert env["HOME"] == tempfile.gettempdir() - assert not os.path.exists(scratch) # cleaned up after the decode - - @tshark_present - async def test_decode_isomorphic_mapping(self, tmp_path): - import asyncio - import xml.etree.ElementTree as ET - - project = self._project(tmp_path) - detail = await decode_frame(project, tag=7, ts="1693472000.123456", - node_id="n1", link_id="linkA", marker="icmp") + finally: + await manager.close_all() assert detail["source"]["frame_number"] == 1 assert detail["hex"] == _icmp_frame().hex() - assert detail["field_count"] > 0 - assert "tshark" in detail["tshark_version"].lower() + assert detail["field_count"] > 10 - # Round-trip fidelity: node count equals the PDML element count - # (protos + fields, excluding the container itself)… - proc = await asyncio.create_subprocess_exec( - "tshark", "-r", str(tmp_path / "n1_linkA_icmp.pcap"), "-T", "pdml", - stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.DEVNULL, - ) - stdout, _ = await proc.communicate() - packet = ET.fromstring(stdout).find("./packet") - xml_elements = [e for e in packet.iter() if e is not packet] - assert detail["field_count"] == len(xml_elements) + # Renamed tree: every label/name present, filter expressions with + # values baked in, byte ranges for hex highlighting. + def find(node, name): + for child in node if isinstance(node, list) else node.get("children", []): + if child.get("name") == name: + return child + deep = find(child, name) + if deep is not None: + return deep + return None - # …and every XML attribute survives verbatim as a JSON string key. - def walk(element, node): - for key, value in element.attrib.items(): - assert node.get(key) == value - assert all(isinstance(v, str) for k, v in node.items() if k != "children") - for child, child_node in zip(element, node["children"]): - walk(child, child_node) + ttl = find(detail["tree"], "ip.ttl") + assert ttl is not None + assert ttl["label"] == "Time to Live: 64" + assert ttl["filter_expr"] == "ip.ttl == 64" + assert ttl["pos"] == 22 and ttl["size"] == 1 - for element, node in zip(packet, detail["tree"]): - walk(element, node) + @sharkd_present + async def test_key_census_closed_set(self, tmp_path): + """The rename correctness guarantee: across protocol-diverse real + trees, sharkd's raw key set stays within the census-known keys.""" - # Values stay strings (no numeric re-typing). - ttl = next( - f for p in detail["tree"] if p.get("name") == "ip" - for f in p["children"] if f.get("name") == "ip.ttl" - ) - assert ttl["show"] == "64" and ttl["showname"] == "Time to Live: 64" + pcap = tmp_path / "proto_mix.pcap" + _write_pcap(pcap, [ + (1693472000, 100000, _icmp_frame()), + (1693472001, 0, _tcp_syn_frame()), + ]) + manager = marker_replay._get_manager() + try: + known = set(marker_replay._KEY_RENAME) | {"h", "e"} + seen = set() + + def walk(value): + if isinstance(value, dict): + seen.update(value.keys()) + for item in value.values(): + walk(item) + elif isinstance(value, list): + for item in value: + walk(item) + + async with manager.session(str(pcap)) as session: + for frame_number in (1, 2): + result = await session.rpc("frame", {"frame": frame_number, "proto": True}) + walk(result.get("tree", [])) + assert seen <= known, f"unknown sharkd keys appeared: {seen - known}" + finally: + await manager.close_all()