Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion publications/comnet/experiments/analysis/aggregate.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,8 @@
import sys
from pathlib import Path

from scipy import stats


def load_results(results_dir: Path) -> dict[str, list]:
experiments: dict[str, list] = {}
Expand Down Expand Up @@ -49,7 +51,8 @@ def stdev(values: list[float]) -> float:
def ci95(values: list[float]) -> float:
if len(values) < 2:
return 0.0
return 1.96 * stdev(values) / math.sqrt(len(values))
n = len(values)
return float(stats.t.ppf(0.975, n - 1)) * stdev(values) / math.sqrt(n)


def extract_metric(data: dict) -> dict[str, float]:
Expand Down
72 changes: 72 additions & 0 deletions publications/comnet/experiments/analysis/buildcheck.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,72 @@
import csv
import json
import statistics
import sys
from pathlib import Path

CONFIGS = ["tcp", "tls", "quic-main"]
BUILDS = ["bcfleet", "bclater"]
FANOUT = 8


def collect(base):
cells = {}
excluded = []
for directory in sorted(base.glob("03h_buildcheck_g*")):
accepted = {row["label"]: row for row in csv.DictReader(open(directory / "acceptance.csv"))}
for row in csv.DictReader(open(directory / "manifest.csv")):
gate = accepted.get(row["label"])
if gate is None or gate["pass"] != "True":
excluded.append((row["label"], gate["failures"] if gate else "not evaluated"))
continue
sub = json.load(open(directory / f"{row['label']}.json"))
pub = json.load(open(directory / f"{row['label']}_pub.json"))
cells.setdefault((row["phase"], row["config"]), []).append({
"group": directory.name[-2:],
"unique": sub["results"]["throughput_avg"] / FANOUT,
"offered": pub["results"]["offered_rate"],
"sha": row["binary_sha"],
})
return cells, excluded


def group_means(samples, field):
per_group = {}
for sample in samples:
per_group.setdefault(sample["group"], []).append(sample[field])
return {group: statistics.mean(values) for group, values in per_group.items()}


def within_group_ratio(numerator, denominator):
left, right = group_means(numerator, "unique"), group_means(denominator, "unique")
return [left[g] / right[g] for g in sorted(set(left) & set(right))]


def main(results_dir):
cells, excluded = collect(Path(results_dir))
print(f"excluded runs: {len(excluded)}")
for label, reason in excluded:
print(f" {label}: {reason}")
print("\n=== 10% loss, router arm, unpaced. unique msg/s pooled mean (per-group means); offered msg/s")
for config in CONFIGS:
for build in BUILDS:
samples = cells.get((build, config), [])
if not samples:
continue
groups = " ".join(f"{g}:{v:.0f}" for g, v in sorted(group_means(samples, "unique").items()))
shas = {s["sha"][:8] for s in samples}
print(f" {config:10s} {build:8s} n={len(samples)} unique {statistics.mean(s['unique'] for s in samples):7.1f} ({groups}) "
f"offered {statistics.mean(s['offered'] for s in samples):9.0f} sha {','.join(sorted(shas))}")
ratios = within_group_ratio(cells.get(("bclater", config), []), cells.get(("bcfleet", config), []))
if ratios:
print(f" {config:10s} later/fleet {statistics.mean(ratios):.3f} (range {min(ratios):.3f}-{max(ratios):.3f}, {len(ratios)} groups)")
for build in BUILDS:
for other in ("tls", "tcp"):
ratios = within_group_ratio(cells.get((build, "quic-main"), []), cells.get((build, other), []))
if ratios:
print(f" QUIC/{other} {build:8s} {statistics.mean(ratios):.3f} (range {min(ratios):.3f}-{max(ratios):.3f}, {len(ratios)} groups)")


if __name__ == "__main__":
default = Path(__file__).resolve().parent.parent / "results-v5"
main(sys.argv[1] if len(sys.argv) > 1 else default)
158 changes: 158 additions & 0 deletions publications/comnet/experiments/analysis/capped_ppub.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,158 @@
import csv
import json
import statistics
import sys
from pathlib import Path

CONFIGS = ["quic-ppub", "quic-main-ppub", "quic-ctl"]
LABELS = {
"quic-ppub": "per-publish delivery (client per-publish)",
"quic-main-ppub": "per-topic delivery (client per-publish)",
"quic-ctl": "control-only delivery (client control-only)",
}
RATES = [1250, 2500, 5000, 10000, 20000, 0]
CAP_PROBE_RATES = {10000, 20000}
FANOUT = 8
PACING_TOLERANCE = 0.02
SUB_SATURATION_PREFIX = "P3: sub busy"
METRICS = ["offered", "delivered", "ratio", "cpu", "cpu_us_per_delivered", "cpu_us_per_offered", "pkts_per_delivered"]


def load_json(path):
try:
return json.load(open(path))
except (OSError, json.JSONDecodeError):
return None


def active_cpu(path):
try:
values = [float(row["cpu_percent"]) for row in csv.DictReader(open(path)) if row.get("cpu_percent")]
except (OSError, ValueError):
return None
if not values:
return None
return statistics.mean(value for value in values if value >= 0.5 * max(values))


def acceptance(directory):
path = directory / "acceptance.csv"
if not path.exists():
raise SystemExit(f"{path} missing: run single_segment_accept.py {directory} first")
return {row["label"]: row for row in csv.DictReader(open(path))}


def collect(base):
cells = {}
excluded = []
for directory in sorted(base.glob("03f_capped_g*")):
accepted = acceptance(directory)
rows = {row["label"]: row for row in csv.DictReader(open(directory / "manifest.csv"))}.values()
for row in rows:
label = row["label"]
gate = accepted.get(label)
if gate is None:
excluded.append((label, "not evaluated by acceptance"))
continue
reasons = [reason for reason in gate["failures"].split("; ") if reason]
if any(not reason.startswith(SUB_SATURATION_PREFIX) for reason in reasons):
excluded.append((label, gate["failures"]))
continue
saturated = bool(reasons)
pub = load_json(directory / f"{label}_pub.json")
sub = load_json(directory / f"{label}.json")
if not pub or not sub:
excluded.append((label, "missing results"))
continue
target = int(row["rate"])
configured = int(pub.get("config", {}).get("rate", 0) or 0)
offered = pub["results"].get("offered_rate") or 0.0
if configured != target:
excluded.append((label, f"bench rate {configured} != planned {target}"))
continue
if target and abs(offered / target - 1) > PACING_TOLERANCE:
excluded.append((label, f"paced at {offered:.0f}/s, target {target}/s"))
continue
delivered = sub["results"]["throughput_avg"]
received = sub["results"].get("received") or 0
cpu = None if saturated else active_cpu(directory / f"{label}_broker_resources.csv")
skbs = float(gate["skbs"]) if gate.get("skbs") else None
cells.setdefault((row["config"], target, int(row["loss"])), []).append({
"group": directory.name[-2:],
"saturated": saturated,
"offered": offered,
"delivered": delivered,
"ratio": delivered / (FANOUT * offered) if offered else None,
"cpu": cpu,
"cpu_us_per_delivered": cpu / 100 / delivered * 1e6 if cpu and delivered else None,
"cpu_us_per_offered": cpu / 100 / offered * 1e6 if cpu and offered else None,
"pkts_per_delivered": skbs / received if skbs and received else None,
})
return cells, excluded


def summary(samples, field):
values = [s[field] for s in samples if s[field] is not None]
if not values:
return "-"
per_group = {}
for sample in samples:
if sample[field] is not None:
per_group.setdefault(sample["group"], []).append(sample[field])
groups = " ".join(f"{g}:{statistics.mean(v):.3g}" for g, v in sorted(per_group.items()))
return f"{statistics.mean(values):.3g} ({groups})"


def paired_differences(cells, loss, field, left, right):
out = []
for rate in RATES:
a, b = cells.get((left, rate, loss)), cells.get((right, rate, loss))
if not a or not b:
continue
diffs = []
for group in sorted({s["group"] for s in a} & {s["group"] for s in b}):
va = [s[field] for s in a if s["group"] == group and s[field] is not None]
vb = [s[field] for s in b if s["group"] == group and s[field] is not None]
if va and vb:
diffs.append(statistics.mean(va) / statistics.mean(vb))
if diffs:
out.append((rate, statistics.mean(diffs), min(diffs), max(diffs), len(diffs)))
return out


def main(results_dir):
cells, excluded = collect(Path(results_dir))
print(f"excluded runs: {len(excluded)}")
for label, reason in excluded:
print(f" {label}: {reason}")
for loss in (0, 1):
print(f"\n=== loss {loss}% (router arm), 8 subscribers. Values: pooled mean (per-group means g1..g3)")
for config in CONFIGS:
print(f" {LABELS[config]}")
for rate in RATES:
samples = cells.get((config, rate, loss))
if not samples:
continue
name = f"rate {rate}" if rate else "uncapped"
tags = []
if rate in CAP_PROBE_RATES and config == "quic-ppub":
tags.append("cap probe")
saturated = sum(s["saturated"] for s in samples)
if saturated:
tags.append(f"{saturated} sub-saturated (cpu omitted)")
print(f" {name:11s} n={len(samples):2d} {'[' + ', '.join(tags) + ']' if tags else ''}")
for field in METRICS:
print(f" {field:22s} {summary(samples, field)}")
for field in ("cpu_us_per_delivered", "pkts_per_delivered"):
for other in ("quic-main-ppub", "quic-ctl"):
rows = paired_differences(cells, loss, field, "quic-ppub", other)
if rows:
print(f" within-group ratio {field}: per-publish / {other}")
for rate, mean, low, high, n in rows:
name = f"rate {rate}" if rate else "uncapped"
print(f" {name:11s} {mean:5.2f} (range {low:.2f}-{high:.2f}, {n} groups)")


if __name__ == "__main__":
default = Path(__file__).resolve().parent.parent / "results-v5"
main(sys.argv[1] if len(sys.argv) > 1 else default)
114 changes: 114 additions & 0 deletions publications/comnet/experiments/analysis/e1_fairness.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,114 @@
import csv
import glob
import json
import statistics
import sys
from pathlib import Path


def broker_tx_mbit(broker_csv: Path):
rows = [r for r in csv.DictReader(open(broker_csv)) if r["net_tx_bytes"].isdigit()]
if len(rows) < 4:
return None
rows = rows[1:-1]
dt = float(rows[-1]["timestamp"]) - float(rows[0]["timestamp"])
db = int(rows[-1]["net_tx_bytes"]) - int(rows[0]["net_tx_bytes"])
return db * 8 / 1e6 / dt if dt > 0 else None


def iperf_mbit(iperf_json: Path):
end = json.load(open(iperf_json)).get("end", {})
if "sum_received" not in end:
return None
return end["sum_received"]["bits_per_second"] / 1e6


def arm_fairness(results_dir: Path, arm: str, loss: int):
prefix = f"{arm}_rate*_loss{loss}pct"
shares, iperfs, jains = [], [], []
for bc in sorted(glob.glob(str(results_dir / f"{prefix}_run*_broker_resources.csv"))):
run = Path(bc).name.replace("_broker_resources.csv", "")
ic = results_dir / f"{run}_iperf.json"
mc = results_dir / f"{run}_messages.csv"
tx = broker_tx_mbit(Path(bc))
ip = iperf_mbit(ic) if ic.exists() else None
if tx and ip and tx > ip:
shares.append(100 * (tx - ip) / tx)
iperfs.append(ip)
if mc.exists():
gp = per_connection_goodput(mc)
if len(gp) > 1:
jains.append(jain_index(list(gp.values())))
med = lambda v: round(statistics.median(v), 1) if v else None
return {"n": len(shares), "mqtt_wire_share_pct": med(shares),
"iperf_mbit": med(iperfs), "jain": (round(statistics.median(jains), 3) if jains else None)}


def analyze_dir(results_dir: Path):
arms = ["tcp-1conn", "tcp-Nconn", "quic-control", "quic-pertopic"]
print(f"E1 fairness (MQTT share of a bottleneck shared with one greedy TCP flow)")
print(f"{'arm':>14} {'loss':>5} | {'MQTT share%':>11} {'iperf Mbit':>11} {'intra-Jain':>11} {'n':>3}")
for arm in arms:
for loss in [0, 1]:
r = arm_fairness(results_dir, arm, loss)
if r["n"]:
print(f"{arm:>14} {loss:>4}% | {str(r['mqtt_wire_share_pct']):>11} "
f"{str(r['iperf_mbit']):>11} {str(r['jain']):>11} {r['n']:>3}")


def per_connection_goodput(messages_csv: Path, trim: float = 0.1):
receive_ns = {}
for row in csv.DictReader(open(messages_csv)):
conn = int(row["conn_idx"])
receive_ns.setdefault(conn, []).append(int(row["receive_ns"]))

all_ns = [ns for series in receive_ns.values() for ns in series]
if not all_ns:
return {}
lo, hi = min(all_ns), max(all_ns)
span = hi - lo
start = lo + int(span * trim)
end = hi - int(span * trim)
window_s = max((end - start) / 1e9, 1e-9)

goodput = {}
for conn, series in receive_ns.items():
in_window = sum(1 for ns in series if start <= ns <= end)
goodput[conn] = in_window / window_s
return goodput


def jain_index(values):
if not values:
return 0.0
n = len(values)
total = sum(values)
total_sq = sum(v * v for v in values)
if total_sq == 0:
return 0.0
return (total * total) / (n * total_sq)


def main(messages_csv: Path):
goodput = per_connection_goodput(messages_csv)
if not goodput:
print(f"no data in {messages_csv}")
return
flows = [goodput[c] for c in sorted(goodput)]
print(f"file: {messages_csv}")
print(f"connections: {len(flows)}")
for conn in sorted(goodput):
print(f" conn {conn}: {goodput[conn]:.1f} msg/s")
print(f"aggregate: {sum(flows):.1f} msg/s")
print(f"Jain fairness index: {jain_index(flows):.4f} (1.0 = perfectly fair)")


if __name__ == "__main__":
if len(sys.argv) < 2:
print(f"usage: {sys.argv[0]} <messages.csv | E1_fairness_dir>")
sys.exit(1)
target = Path(sys.argv[1])
if target.is_dir():
analyze_dir(target)
else:
main(target)
Original file line number Diff line number Diff line change
Expand Up @@ -16,8 +16,8 @@
save_figure,
)

LOSS_RATES = [0, 1, 2, 5]
LOSS_LABELS = ["0%", "1%", "2%", "5%"]
LOSS_RATES = [1, 2, 5]
LOSS_LABELS = ["1%", "2%", "5%"]
RUNS = range(1, 16)


Expand Down Expand Up @@ -64,13 +64,6 @@ def main(results_dir: Path, output_dir: Path):

fig, ax = plt.subplots(figsize=(7, 4.5))

x_offsets = {
"tcp": -0.15,
"quic-control": -0.05,
"quic-pertopic": 0.05,
"quic-perpub": 0.15,
}

group_positions = np.arange(len(LOSS_RATES))

for transport in TRANSPORT_ORDER:
Expand All @@ -82,7 +75,7 @@ def main(results_dir: Path, output_dir: Path):
mean, ci_half = compute_ci(data[transport][loss])
means.append(mean)
ci_halves.append(ci_half)
positions.append(group_positions[loss_idx] + x_offsets[transport])
positions.append(group_positions[loss_idx])

ax.errorbar(
positions,
Expand Down
Loading
Loading