Repository navigation
Expand file tree
/
Copy pathchains.py
More file actions
395 lines (336 loc) · 16.6 KB
/
Copy pathchains.py
File metadata and controls
395 lines (336 loc) · 16.6 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
"""Traversals that chain across edge types — the queries a lookup table cannot answer.
Every question here starts at one kind of node and ends at another, crossing a
different relationship on the way:
Maintainer -MAINTAINS-> Package -REQUIRED_BY*-> the whole blast radius
Advisory -AFFECTS-> Package -REQUIRED_BY*-> everyone the CVE reaches
Package -SIMILAR_TO-> impostor -REQUIRED_BY*-> who already installed it
That shape is the entire argument for a graph. "Which packages does qix
maintain" is a join. "How many packages transitively depend on anything qix
maintains" is a traversal across two relationship types, and it changes answer
the moment any edge anywhere in the graph moves.
HydraDB 0.1.0 constraint that shapes all of this: a variable-length MATCH needs
a *fixed source id*, so each stage of a chain is its own query anchored at a
known vertex, and the stages are fired concurrently. It cannot return a path at
all — no length(path), no path binding — so where a concrete chain is required
it is rebuilt from the sidecar edge table and then re-verified hop by hop
against the graph, rather than asserted.
"""
import time
from concurrent.futures import ThreadPoolExecutor
from hydra import Hydra, RESULT_LIMIT, adv_id, maint_id, pkg_id
from blast import REACH_COUNT, REACH_NAMES, _depth
MAINTAINS = "MATCH (m {id: $id})-[:MAINTAINS]->(p) RETURN p.name"
AFFECTS = "MATCH (a {id: $id})-[:AFFECTS]->(p) RETURN p.name"
SIMILAR = "MATCH (p {id: $id})-[:SIMILAR_TO]->(q) RETURN q.name"
ADVISORY_NODE = ("MATCH (a:Advisory {id: $id}) "
"RETURN a.osv_id, a.severity, a.is_malware, a.summary")
ONE_HOP = "MATCH (t {id: $id})-[:REQUIRED_BY*1..1]->(v) RETURN DISTINCT v.name"
def _reach(h, name, depth, ecosystem="npm"):
rows = h.query(REACH_NAMES % depth, {"id": pkg_id(name, ecosystem),
"limit": RESULT_LIMIT})
return {r["v.name"] for r in rows if r.get("v.name")}
def _count(h, name, depth, ecosystem="npm"):
rows = h.query(REACH_COUNT % depth, {"id": pkg_id(name, ecosystem)})
return rows[0]["count(*)"] if rows else 0
# --------------------------------------------------------------------------
# 0. node expansion — what the explorer walks
# --------------------------------------------------------------------------
# Every relationship is stored in both directions, so any node can be expanded
# from a fixed source id. Without the reverse edges there would be no way to ask
# "who maintains this package" starting from the package, and the explorer would
# have to fall back to SQL — which is exactly what putting these in the graph
# was meant to stop.
EXPANSIONS = {
"package": [
("REQUIRED_BY", "dependents", "package"),
("SIMILAR_TO", "similar name", "package"),
("MAINTAINED_BY", "maintained by", "maintainer"),
("HAS_ADVISORY", "advisory", "advisory"),
],
"maintainer": [("MAINTAINS", "maintains", "package")],
"advisory": [("AFFECTS", "affects", "package")],
}
# Package ids carry their ecosystem; maintainer and advisory ids deliberately
# do not, because one human and one advisory are single nodes across all of them.
NEIGHBOURS = "MATCH (t {id: $id})-[:%s]->(v) RETURN DISTINCT v.name LIMIT $limit"
IDENTIFY = "MATCH (n {id: $id}) RETURN n.name, n.osv_id, n.is_malware, n.severity"
def node_id(name: str, kind: str, ecosystem: str = "npm") -> int:
if kind == "maintainer":
return maint_id(name)
if kind == "advisory":
return adv_id(name)
return pkg_id(name, ecosystem)
def expand(h: Hydra, name: str, kind: str = "package", limit: int = 40,
degree_for: bool = True, ecosystem: str = "npm"):
"""One node and its neighbours across every edge type it participates in.
This is the whole explorer in one call: the front end holds no model of the
graph, it just asks HydraDB what is adjacent to whatever was clicked.
"""
t0 = time.perf_counter()
kind = kind if kind in EXPANSIONS else "package"
root_id = node_id(name, kind, ecosystem)
ident = h.query(IDENTIFY, {"id": root_id})
if not ident or not ident[0].get("n.name"):
return {"found": False, "name": name, "kind": kind,
"message": f"no {kind} named '{name}' in the graph",
"latency_ms": round((time.perf_counter() - t0) * 1000, 1)}
specs = EXPANSIONS[kind]
def pull(spec):
rel, label, target_kind = spec
rows = h.query(NEIGHBOURS % rel, {"id": root_id, "limit": limit})
return rel, label, target_kind, [r["v.name"] for r in rows if r.get("v.name")]
with ThreadPoolExecutor(max_workers=len(specs)) as pool:
pulled = list(pool.map(pull, specs))
neighbours = []
for rel, label, target_kind, names in pulled:
for n in names:
neighbours.append({"name": n, "kind": target_kind,
"edge": rel, "edge_label": label})
# Size packages by how much depends on them; that is the one number that
# makes a node worth looking at.
if degree_for:
pkgs = [n for n in neighbours if n["kind"] == "package"][:limit]
if pkgs:
with ThreadPoolExecutor(max_workers=min(10, len(pkgs))) as pool:
for n, c in zip(pkgs, pool.map(
lambda x: _count(h, x["name"], 1), pkgs)):
n["dependents"] = c
# An advisory neighbour has to carry its own is_malware, or the explorer
# cannot colour it — and "this one is malware" is the entire point of
# having advisories on the canvas at all.
advs = [n for n in neighbours if n["kind"] == "advisory"]
if advs:
def identify(n):
rows = h.query(IDENTIFY, {"id": node_id(n["name"], "advisory")})
return rows[0] if rows else {}
with ThreadPoolExecutor(max_workers=min(10, len(advs))) as pool:
for n, row in zip(advs, pool.map(identify, advs)):
n["is_malware"] = bool(row.get("n.is_malware"))
n["severity"] = row.get("n.severity") or ""
row = ident[0]
return {
"found": True,
"node": {
"name": row.get("n.name"), "kind": kind,
"osv_id": row.get("n.osv_id"),
"is_malware": bool(row.get("n.is_malware")),
"severity": row.get("n.severity") or "",
"dependents": _count(h, name, 1) if kind == "package" else None,
},
"neighbours": neighbours,
"counts": {label: len(names) for _, label, _, names in pulled},
"queries": 2 + len(specs),
"latency_ms": round((time.perf_counter() - t0) * 1000, 1),
}
# --------------------------------------------------------------------------
# 1. attack surface of a maintainer
# --------------------------------------------------------------------------
def attack_surface(h: Hydra, maintainer: str, depth: int = 4,
max_packages: int = 40):
"""Everything one account can reach. Two hops, two edge types.
This is the number that matters after a maintainer is phished: not "how
many packages did they publish", but how much of npm sits downstream of
those packages. The September 2025 attack began with exactly one account.
"""
d = _depth(depth)
t0 = time.perf_counter()
rows = h.query(MAINTAINS, {"id": maint_id(maintainer)})
packages = sorted({r["p.name"] for r in rows if r.get("p.name")})
if not packages:
return {"maintainer": maintainer, "controls": [], "package_count": 0,
"total_exposed": 0, "depth": d,
"message": f"no packages recorded for maintainer '{maintainer}'",
"latency_ms": round((time.perf_counter() - t0) * 1000, 1)}
# The union matters, not the sum: two packages by the same author usually
# share much of their downstream, and adding the counts double-counts it.
scope = packages[:max_packages]
with ThreadPoolExecutor(max_workers=min(10, len(scope))) as pool:
reaches = list(pool.map(lambda p: (p, _reach(h, p, d)), scope))
union: set[str] = set()
for _, r in reaches:
union |= r
union -= set(packages)
controls = sorted(({"package": p, "exposed": len(r)} for p, r in reaches),
key=lambda x: -x["exposed"])
return {
"maintainer": maintainer,
"package_count": len(packages),
"controls": controls,
"analysed_packages": len(scope),
"truncated": len(packages) > len(scope),
"total_exposed": len(union),
"sum_of_parts": sum(c["exposed"] for c in controls),
"depth": d,
"headline": (f"{maintainer} controls {len(packages):,} package"
f"{'' if len(packages) == 1 else 's'}. "
f"{len(union):,} packages depend on them."),
"queries": 1 + len(scope),
"latency_ms": round((time.perf_counter() - t0) * 1000, 1),
}
# --------------------------------------------------------------------------
# 2. why am I exposed — the actual chain
# --------------------------------------------------------------------------
def why_exposed(h: Hydra, db, source: str, target: str, depth: int = 6,
ecosystem: str = "npm"):
"""Not "you are exposed" — the specific chain, hop by hop.
The depth at which `target` first becomes reachable is computed in the
graph and is authoritative. The concrete chain cannot be: HydraDB 0.1.0
returns no path binding, so it is rebuilt from the sidecar edge table and
then every hop is re-confirmed with a one-hop graph query. If the graph
disagrees with the reconstruction, the response says so rather than
printing a chain nobody verified.
"""
d = _depth(depth)
t0 = time.perf_counter()
found_at = None
seen: set[str] = set()
for k in range(1, d + 1):
reach = _reach(h, source, k)
if target in reach:
found_at = k
seen = reach
break
seen = reach
if found_at is None:
return {"from": source, "to": target, "connected": False,
"searched_depth": d, "reachable_at_depth": None,
"message": (f"{target} is not reachable from {source} within "
f"{d} hops ({len(seen):,} packages searched)"),
"latency_ms": round((time.perf_counter() - t0) * 1000, 1)}
# Rebuild backwards through the sidecar: who depends on whom.
chain = _reconstruct(db, source, target, found_at)
verified = _verify_chain(h, chain, ecosystem) if chain else False
return {
"from": source, "to": target, "connected": True,
"reachable_at_depth": found_at,
"path": chain,
"hops": len(chain) - 1 if chain else None,
"graph_verified": verified,
"explanation": _explain(chain) if chain else
(f"{target} is reachable from {source} in {found_at} hops, "
f"but the concrete chain could not be rebuilt from the "
f"sidecar edge table."),
"queries": found_at + (len(chain) - 1 if chain and verified else 0),
"latency_ms": round((time.perf_counter() - t0) * 1000, 1),
}
def _reconstruct(db, source, target, want_len):
"""BFS over the sidecar edges from source to target, bounded by want_len."""
frontier = [source]
parent: dict[str, str] = {}
seen = {source}
for _ in range(want_len):
if not frontier:
break
marks = ",".join("?" * len(frontier))
rows = db.execute(
f"SELECT dst, src FROM deps WHERE dst IN ({marks}) AND kind = 'prod'",
frontier).fetchall()
nxt = []
for dep, dependent in rows:
if dependent in seen:
continue
seen.add(dependent)
parent[dependent] = dep
if dependent == target:
chain = [target]
while chain[-1] in parent:
chain.append(parent[chain[-1]])
return list(reversed(chain))
nxt.append(dependent)
frontier = nxt
return None
def _verify_chain(h: Hydra, chain, ecosystem="npm"):
"""Confirm every hop really exists in the graph, not just the sidecar."""
if not chain or len(chain) < 2:
return False
def hop(i):
rows = h.query(ONE_HOP, {"id": pkg_id(chain[i], ecosystem)})
return chain[i + 1] in {r["v.name"] for r in rows if r.get("v.name")}
with ThreadPoolExecutor(max_workers=min(8, len(chain) - 1)) as pool:
return all(pool.map(hop, range(len(chain) - 1)))
def _explain(chain):
if not chain or len(chain) < 2:
return ""
steps = [f"{chain[i + 1]} depends on {chain[i]}" for i in range(len(chain) - 1)]
return " → ".join(chain) + ". " + "; ".join(steps) + "."
# --------------------------------------------------------------------------
# 3. blast radius of an advisory
# --------------------------------------------------------------------------
def blast_advisory(h: Hydra, osv_id: str, depth: int = 4):
"""Start at the CVE, not the package. Advisory -AFFECTS-> Package -REQUIRED_BY*->
An advisory usually names several packages, and asking "how far does
GHSA-xxxx actually reach" is a different question from asking about any one
of them.
"""
d = _depth(depth)
t0 = time.perf_counter()
aid = adv_id(osv_id)
meta_rows = h.query(ADVISORY_NODE, {"id": aid})
rows = h.query(AFFECTS, {"id": aid})
packages = sorted({r["p.name"] for r in rows if r.get("p.name")})
if not packages:
return {"osv_id": osv_id, "affects": [], "total_exposed": 0, "depth": d,
"message": f"{osv_id} is not in the graph",
"latency_ms": round((time.perf_counter() - t0) * 1000, 1)}
with ThreadPoolExecutor(max_workers=min(10, len(packages))) as pool:
reaches = list(pool.map(lambda p: (p, _reach(h, p, d)), packages[:40]))
union: set[str] = set()
for _, r in reaches:
union |= r
union -= set(packages)
meta = meta_rows[0] if meta_rows else {}
return {
"osv_id": osv_id,
"severity": meta.get("a.severity") or "",
"is_malware": bool(meta.get("a.is_malware")),
"summary": meta.get("a.summary") or "",
"affects": [{"package": p, "exposed": len(r)} for p, r in
sorted(reaches, key=lambda x: -len(x[1]))],
"package_count": len(packages),
"total_exposed": len(union),
"depth": d,
"headline": (f"{osv_id} directly affects {len(packages):,} package"
f"{'' if len(packages) == 1 else 's'}; "
f"{len(union):,} more depend on them."),
"queries": 2 + len(reaches),
"latency_ms": round((time.perf_counter() - t0) * 1000, 1),
}
# --------------------------------------------------------------------------
# 4. typosquat risk — a squat with dependents is an incident
# --------------------------------------------------------------------------
def typosquat_risk(h: Hydra, name: str, depth: int = 3,
ecosystem: str = "npm"):
"""SIMILAR_TO neighbours, then how many packages already pull each one.
A near-miss name nobody uses is trivia. A near-miss name with real
dependents means somebody has already installed the wrong thing.
"""
d = _depth(depth)
t0 = time.perf_counter()
rows = h.query(SIMILAR, {"id": pkg_id(name, ecosystem)})
neighbours = sorted({r["q.name"] for r in rows if r.get("q.name")})
if not neighbours:
return {"name": name, "neighbours": [], "at_risk": 0, "depth": d,
"message": (f"no similarly-named package in the graph for "
f"'{name}'"),
"latency_ms": round((time.perf_counter() - t0) * 1000, 1)}
with ThreadPoolExecutor(max_workers=min(10, len(neighbours))) as pool:
counts = list(pool.map(lambda n: (n, _count(h, n, d)), neighbours[:40]))
ranked = sorted(({"name": n, "dependents": c} for n, c in counts),
key=lambda x: -x["dependents"])
active = [r for r in ranked if r["dependents"] > 0]
return {
"name": name,
"neighbours": ranked,
"neighbour_count": len(neighbours),
"at_risk": len(active),
"worst": ranked[0] if ranked else None,
"depth": d,
"headline": (f"{len(neighbours):,} package"
f"{'' if len(neighbours) == 1 else 's'} in the graph sit one "
f"edit from {name}"
+ (f"; {ranked[0]['name']} already has "
f"{ranked[0]['dependents']:,} dependents."
if active else "; none of them have dependents.")),
"queries": 1 + len(counts),
"latency_ms": round((time.perf_counter() - t0) * 1000, 1),
}