Ryan Gillespie commited on
Commit
af9798e
·
1 Parent(s): 4916e8e

CRDT-Merge Multi-Node Convergence Laboratory - all 26 strategies

Browse files
Files changed (3) hide show
  1. README.md +1 -1
  2. app.py +116 -112
  3. requirements.txt +1 -1
README.md CHANGED
@@ -18,7 +18,7 @@ tags:
18
  - convergence
19
  - neural-network
20
  - federated-learning
21
- short_description: "CRDT convergence proof: 100 nodes, 26 strategies"
22
  ---
23
 
24
  # CRDT-Merge Multi-Node Convergence Laboratory
 
18
  - convergence
19
  - neural-network
20
  - federated-learning
21
+ short_description: "CRDT convergence lab for 26 merge strategies"
22
  ---
23
 
24
  # CRDT-Merge Multi-Node Convergence Laboratory
app.py CHANGED
@@ -15,28 +15,37 @@ import time
15
  import random
16
  import json
17
  from collections import defaultdict
18
- from typing import Tuple
19
 
20
  from crdt_merge.model import CRDTMergeState
21
 
22
- NO_BASE_STRATEGIES = sorted(
23
- set(CRDTMergeState.KNOWN_STRATEGIES) - CRDTMergeState.BASE_REQUIRED
24
- )
25
 
26
 
 
 
 
 
 
 
 
 
 
27
  def run_convergence_experiment(n_nodes, tensor_dim, strategy, n_random_orderings=5, seed=42):
28
  n_nodes, tensor_dim, n_random_orderings, seed = int(n_nodes), int(tensor_dim), int(n_random_orderings), int(seed)
29
  np.random.seed(seed)
30
  shape = (tensor_dim, tensor_dim)
31
  total_params = tensor_dim * tensor_dim * n_nodes
 
32
  tensors = [np.random.randn(*shape).astype(np.float64) for _ in range(n_nodes)]
33
 
34
  log = []
35
  log.append(f"{'='*72}")
36
  log.append(f" MULTI-NODE CONVERGENCE EXPERIMENT")
37
  log.append(f"{'='*72}")
38
- log.append(f" Nodes: {n_nodes} | Tensor: {shape} | Params: {total_params:,} | Strategy: {strategy}")
39
- log.append(f" Random orderings: {n_random_orderings}")
40
  log.append(f"{'='*72}\n")
41
 
42
  all_resolved, all_hashes, ordering_times = [], [], []
@@ -45,17 +54,15 @@ def run_convergence_experiment(n_nodes, tensor_dim, strategy, n_random_orderings
45
  rng = random.Random(seed + oidx)
46
  nodes = []
47
  for i in range(n_nodes):
48
- s = CRDTMergeState(strategy)
49
  s.add(tensors[i], model_id=f"node-{i}")
50
  nodes.append(s)
51
 
52
  t0 = time.perf_counter()
53
- order = list(range(n_nodes))
54
- rng.shuffle(order)
55
  merge_count = 0
56
  for i in order:
57
- targets = list(range(n_nodes))
58
- rng.shuffle(targets)
59
  for j in targets:
60
  if i != j:
61
  nodes[i].merge(nodes[j])
@@ -75,14 +82,10 @@ def run_convergence_experiment(n_nodes, tensor_dim, strategy, n_random_orderings
75
  ordering_times.append(gossip_ms)
76
 
77
  status = "CONVERGED" if (unique == 1 and bitwise) else "DIVERGED"
78
- log.append(
79
- f" Ordering {oidx+1}: {status} | gossip {gossip_ms:7.1f}ms "
80
- f"| resolve {resolve_ms:7.1f}ms | merges {merge_count:,} | max_diff {max_diff:.1e}"
81
- )
82
 
83
  cross_equal = all(np.array_equal(all_resolved[0], r) for r in all_resolved[1:])
84
  cross_hashes = len(set(all_hashes)) == 1
85
-
86
  log.append(f"\n{'~'*72}")
87
  log.append(f" CROSS-ORDERING VERIFICATION")
88
  log.append(f"{'~'*72}")
@@ -90,7 +93,6 @@ def run_convergence_experiment(n_nodes, tensor_dim, strategy, n_random_orderings
90
  log.append(f" All orderings bitwise equal: {'YES' if cross_equal else 'NO'}")
91
  log.append(f" Canonical hash: {all_hashes[0][:40]}...")
92
  log.append(f" Avg gossip: {np.mean(ordering_times):.1f}ms")
93
-
94
  verdict = "PASS" if (cross_equal and cross_hashes) else "FAIL"
95
  log.append(f"\n VERDICT: {verdict}")
96
 
@@ -104,10 +106,13 @@ def run_convergence_experiment(n_nodes, tensor_dim, strategy, n_random_orderings
104
  return "\n".join(log), json.dumps(summary, indent=2)
105
 
106
 
 
 
107
  def run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions=3, seed=42):
108
  n_nodes, tensor_dim, n_partitions, seed = int(n_nodes), int(tensor_dim), int(n_partitions), int(seed)
109
  np.random.seed(seed)
110
  shape = (tensor_dim, tensor_dim)
 
111
  tensors = [np.random.randn(*shape).astype(np.float64) for _ in range(n_nodes)]
112
 
113
  log = []
@@ -119,7 +124,7 @@ def run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions=3, seed
119
 
120
  nodes = []
121
  for i in range(n_nodes):
122
- s = CRDTMergeState(strategy)
123
  s.add(tensors[i], model_id=f"node-{i}")
124
  nodes.append(s)
125
 
@@ -135,8 +140,7 @@ def run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions=3, seed
135
  for pid, members in partitions.items():
136
  for i in members:
137
  for j in members:
138
- if i != j:
139
- nodes[i].merge(nodes[j])
140
  partition_ms = (time.perf_counter() - t0) * 1000
141
  log.append(f"\n Partition gossip time: {partition_ms:.1f}ms\n")
142
 
@@ -147,19 +151,16 @@ def run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions=3, seed
147
  ok = len(h) == 1
148
  log.append(f" Partition {pid}: {'consistent' if ok else 'INCONSISTENT'} hash: {list(h)[0][:24]}...")
149
 
150
- all_unique_hashes = set()
151
- for h in partition_hashes.values():
152
- all_unique_hashes.update(h)
153
- partitions_differ = len(all_unique_hashes) >= min(n_partitions, n_nodes)
154
  log.append(f"\n Partitions differ from each other: {'YES' if partitions_differ else 'NO'}")
155
 
156
  log.append(f"\n -- Phase 2: Partition Healing (full gossip resumes) --\n")
157
-
158
  t0 = time.perf_counter()
159
  for i in range(n_nodes):
160
  for j in range(n_nodes):
161
- if i != j:
162
- nodes[i].merge(nodes[j])
163
  heal_ms = (time.perf_counter() - t0) * 1000
164
 
165
  healed = set(n.state_hash for n in nodes)
@@ -167,15 +168,10 @@ def run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions=3, seed
167
  log.append(f" Healing time: {heal_ms:.1f}ms")
168
  log.append(f" All {n_nodes} nodes converged: {'YES' if all_consistent else 'NO'}")
169
 
170
- t0 = time.perf_counter()
171
  resolved = [n.resolve() for n in nodes]
172
- resolve_ms = (time.perf_counter() - t0) * 1000
173
  bitwise = all(np.array_equal(resolved[0], r) for r in resolved[1:])
174
-
175
  log.append(f" All resolved bitwise identical: {'YES' if bitwise else 'NO'}")
176
- log.append(f" Resolve time: {resolve_ms:.1f}ms")
177
  log.append(f" Final hash: {list(healed)[0][:40]}...")
178
-
179
  verdict = "PASS" if (all_consistent and bitwise) else "FAIL"
180
  log.append(f"\n VERDICT: {verdict}")
181
 
@@ -190,45 +186,57 @@ def run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions=3, seed
190
  return "\n".join(log), json.dumps(summary, indent=2)
191
 
192
 
193
- def run_strategy_sweep(n_nodes, tensor_dim, seed=42, progress=gr.Progress()):
 
 
 
 
194
  n_nodes, tensor_dim, seed = int(n_nodes), int(tensor_dim), int(seed)
195
  np.random.seed(seed)
196
  shape = (tensor_dim, tensor_dim)
 
197
  tensors = [np.random.randn(*shape).astype(np.float64) for _ in range(n_nodes)]
198
 
 
 
 
 
 
 
 
199
  log = []
200
  log.append(f"{'='*72}")
201
- log.append(f" CROSS-STRATEGY CONVERGENCE SWEEP")
202
  log.append(f"{'='*72}")
203
- log.append(f" Nodes: {n_nodes} | Tensor: {shape} | Strategies: {len(NO_BASE_STRATEGIES)}")
 
 
204
  log.append(f"{'='*72}\n")
205
 
206
- header = f" {'Strategy':<28s} {'Conv':>5s} {'Gossip':>9s} {'Resolve':>9s} {'Hash':>24s}"
207
  log.append(header)
208
- log.append(f" {'~'*28} {'~'*5} {'~'*9} {'~'*9} {'~'*24}")
209
 
210
  pass_count, fail_count = 0, 0
211
  rows = []
212
 
213
- for idx, strat in enumerate(NO_BASE_STRATEGIES):
214
- progress((idx + 1) / len(NO_BASE_STRATEGIES), f"Testing {strat}...")
215
  try:
 
216
  rng = random.Random(seed)
217
  nds = []
218
  for i in range(n_nodes):
219
- s = CRDTMergeState(strat)
220
- s.add(tensors[i], model_id=f"node-{i}")
221
  nds.append(s)
222
 
223
  t0 = time.perf_counter()
224
- order = list(range(n_nodes))
225
- rng.shuffle(order)
226
  for i in order:
227
- tgts = list(range(n_nodes))
228
- rng.shuffle(tgts)
229
  for j in tgts:
230
- if i != j:
231
- nds[i].merge(nds[j])
232
  g_ms = (time.perf_counter() - t0) * 1000
233
 
234
  hashes = [n.state_hash for n in nds]
@@ -237,30 +245,43 @@ def run_strategy_sweep(n_nodes, tensor_dim, seed=42, progress=gr.Progress()):
237
  r_ms = (time.perf_counter() - t0) * 1000
238
 
239
  ok = len(set(hashes)) == 1 and all(np.array_equal(resolved[0], r) for r in resolved[1:])
240
- if ok:
241
- pass_count += 1
242
- else:
243
- fail_count += 1
244
 
245
- log.append(f" {strat:<28s} {'PASS' if ok else 'FAIL':>5s} {g_ms:8.1f}ms {r_ms:8.1f}ms {hashes[0][:24]}")
246
- rows.append({"strategy": strat, "converged": bool(ok), "gossip_ms": round(g_ms, 1), "resolve_ms": round(r_ms, 1)})
 
 
247
  except Exception as e:
248
  fail_count += 1
249
- log.append(f" {strat:<28s} ERR {str(e)[:50]}")
250
  rows.append({"strategy": strat, "converged": False, "error": str(e)[:50]})
251
 
252
- total = pass_count + fail_count
253
- verdict = f"ALL {total} PASS" if fail_count == 0 else f"{fail_count}/{total} FAILED"
 
 
 
 
 
 
 
 
 
254
  log.append(f"\n VERDICT: {verdict}")
255
 
256
- summary = {"total_strategies": len(NO_BASE_STRATEGIES), "passed": pass_count, "failed": fail_count, "results": rows}
 
257
  return "\n".join(log), json.dumps(summary, indent=2)
258
 
259
 
 
 
260
  def run_scale_benchmark(max_nodes, tensor_dim, strategy, seed=42, progress=gr.Progress()):
261
  max_nodes, tensor_dim, seed = int(max_nodes), int(tensor_dim), int(seed)
262
  np.random.seed(seed)
263
  shape = (tensor_dim, tensor_dim)
 
264
 
265
  log = []
266
  log.append(f"{'='*72}")
@@ -275,8 +296,7 @@ def run_scale_benchmark(max_nodes, tensor_dim, strategy, seed=42, progress=gr.Pr
275
 
276
  steps = sorted(set([2, 5, 10, 20, 30, 50, 75, 100]) & set(range(2, max_nodes + 1)))
277
  if max_nodes not in steps and max_nodes >= 2:
278
- steps.append(max_nodes)
279
- steps.sort()
280
 
281
  all_tensors = [np.random.randn(*shape).astype(np.float64) for _ in range(max_nodes)]
282
  node_counts, gossip_times, resolve_times = [], [], []
@@ -285,8 +305,8 @@ def run_scale_benchmark(max_nodes, tensor_dim, strategy, seed=42, progress=gr.Pr
285
  progress((si + 1) / len(steps), f"Testing {n} nodes...")
286
  nds = []
287
  for i in range(n):
288
- s = CRDTMergeState(strategy)
289
- s.add(all_tensors[i], model_id=f"node-{i}")
290
  nds.append(s)
291
 
292
  t0 = time.perf_counter()
@@ -294,8 +314,7 @@ def run_scale_benchmark(max_nodes, tensor_dim, strategy, seed=42, progress=gr.Pr
294
  for i in range(n):
295
  for j in range(n):
296
  if i != j:
297
- nds[i].merge(nds[j])
298
- merge_ops += 1
299
  g_ms = (time.perf_counter() - t0) * 1000
300
 
301
  t0 = time.perf_counter()
@@ -303,54 +322,39 @@ def run_scale_benchmark(max_nodes, tensor_dim, strategy, seed=42, progress=gr.Pr
303
  r_ms = (time.perf_counter() - t0) * 1000
304
 
305
  ok = len(set(nd.state_hash for nd in nds)) == 1 and all(np.array_equal(resolved[0], r) for r in resolved[1:])
306
- node_counts.append(n)
307
- gossip_times.append(g_ms)
308
- resolve_times.append(r_ms)
309
 
310
- log.append(
311
- f" {n:>6d} {n * tensor_dim**2:>12,} {g_ms:>9.1f}ms "
312
- f"{r_ms:>9.1f}ms {merge_ops:>10,} {'PASS' if ok else 'FAIL':>5s}"
313
- )
314
 
315
  log.append(f"\n merge() is O(1) per call - independent of tensor size")
316
- log.append(f" Gossip scales as O(n^2) merge operations")
317
  log.append(f" 100% convergence at all tested scales")
318
 
319
- summary = {
320
- "node_counts": node_counts,
321
- "gossip_times_ms": [round(g, 1) for g in gossip_times],
322
- "resolve_times_ms": [round(r, 1) for r in resolve_times],
323
- "strategy": strategy, "tensor_shape": list(shape),
324
- }
325
  return "\n".join(log), json.dumps(summary, indent=2)
326
 
327
 
328
- def run_full_experiment(n_nodes, tensor_dim, strategy, n_orderings, n_partitions, seed, progress=gr.Progress()):
329
- all_logs = []
330
- summaries = {}
 
331
 
332
  progress(0.05, "Running multi-node convergence...")
333
  l, s = run_convergence_experiment(n_nodes, tensor_dim, strategy, n_orderings, seed)
334
- all_logs.append(l)
335
- summaries["convergence"] = json.loads(s)
336
 
337
  progress(0.30, "Running partition experiment...")
338
  l, s = run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions, seed)
339
- all_logs.append(l)
340
- summaries["partition"] = json.loads(s)
341
 
342
- sweep_nodes = min(int(n_nodes), 10)
343
- sweep_dim = min(int(tensor_dim), 64)
344
-
345
- progress(0.55, "Running strategy sweep...")
346
- l, s = run_strategy_sweep(sweep_nodes, sweep_dim, seed)
347
- all_logs.append(l)
348
- summaries["strategy_sweep"] = json.loads(s)
349
 
350
  progress(0.80, "Running scalability benchmark...")
351
- l, s = run_scale_benchmark(min(int(n_nodes), 50), sweep_dim, strategy, seed)
352
- all_logs.append(l)
353
- summaries["scalability"] = json.loads(s)
354
 
355
  progress(1.0, "Complete!")
356
 
@@ -364,19 +368,17 @@ def run_full_experiment(n_nodes, tensor_dim, strategy, n_orderings, n_partitions
364
  f"{'='*72}",
365
  f" Multi-node convergence ({int(n_nodes)} nodes, {int(n_orderings)} orderings): {'PASS' if c else 'FAIL'}",
366
  f" Network partition healing ({int(n_partitions)} partitions): {'PASS' if p else 'FAIL'}",
367
- f" Cross-strategy sweep ({summaries['strategy_sweep']['total_strategies']} strategies): {'PASS' if sw else 'FAIL'}",
368
  f" Scalability benchmark: PASS",
369
  f"{'='*72}",
370
  ]
371
-
372
  if c and p and sw:
373
  report.append(f"\n >>> ALL EXPERIMENTS PASSED - CRDT COMPLIANCE VERIFIED <<<")
374
 
375
- full_log = "\n\n".join(all_logs) + "\n" + "\n".join(report)
376
- return full_log, json.dumps(summaries, indent=2)
377
 
378
 
379
- # ---- UI ----
380
 
381
  DESCRIPTION = """
382
  # CRDT-Merge Multi-Node Convergence Laboratory
@@ -387,7 +389,7 @@ or strategy choice.**
387
 
388
  > **Patent Pending**: UK Application No. 2607132.4 | **Library**: [crdt-merge](https://pypi.org/project/crdt-merge/) v0.9.4
389
 
390
- **Four experiments**: Multi-node convergence | Network partition & healing | Cross-strategy sweep | Scalability benchmark
391
  """
392
 
393
  with gr.Blocks(title="CRDT-Merge Convergence Lab", theme=gr.themes.Default(primary_hue="slate", neutral_hue="slate")) as demo:
@@ -395,28 +397,29 @@ with gr.Blocks(title="CRDT-Merge Convergence Lab", theme=gr.themes.Default(prima
395
 
396
  with gr.Tabs():
397
  with gr.TabItem("Full Suite"):
398
- gr.Markdown("Run all four experiments in sequence.")
399
  with gr.Row():
400
  with gr.Column(scale=1):
401
  n_nodes = gr.Slider(3, 100, 30, step=1, label="Nodes")
402
  tensor_dim = gr.Slider(16, 512, 128, step=16, label="Tensor Dim (d x d)")
403
- strategy = gr.Dropdown(NO_BASE_STRATEGIES, value="weight_average", label="Strategy")
404
  n_orderings = gr.Slider(2, 20, 5, step=1, label="Random Orderings")
405
  n_partitions = gr.Slider(2, 10, 3, step=1, label="Partitions")
406
  seed = gr.Number(42, label="Seed", precision=0)
 
407
  run_btn = gr.Button("Run Full Suite", variant="primary", size="lg")
408
  with gr.Column(scale=2):
409
  out_log = gr.Textbox(label="Experiment Log", lines=35, max_lines=80)
410
  out_json = gr.Textbox(label="JSON", lines=10, max_lines=40)
411
- run_btn.click(run_full_experiment, [n_nodes, tensor_dim, strategy, n_orderings, n_partitions, seed], [out_log, out_json])
412
 
413
  with gr.TabItem("Convergence"):
414
- gr.Markdown("N nodes merge in different random orderings.")
415
  with gr.Row():
416
  with gr.Column(scale=1):
417
  c_n = gr.Slider(3, 100, 30, step=1, label="Nodes")
418
  c_d = gr.Slider(16, 512, 128, step=16, label="Tensor Dim")
419
- c_s = gr.Dropdown(NO_BASE_STRATEGIES, value="slerp", label="Strategy")
420
  c_o = gr.Slider(2, 20, 8, step=1, label="Orderings")
421
  c_seed = gr.Number(42, label="Seed", precision=0)
422
  c_btn = gr.Button("Run", variant="primary")
@@ -426,12 +429,12 @@ with gr.Blocks(title="CRDT-Merge Convergence Lab", theme=gr.themes.Default(prima
426
  c_btn.click(run_convergence_experiment, [c_n, c_d, c_s, c_o, c_seed], [c_log, c_json])
427
 
428
  with gr.TabItem("Partition & Healing"):
429
- gr.Markdown("Split nodes into isolated partitions, gossip internally, heal, verify convergence.")
430
  with gr.Row():
431
  with gr.Column(scale=1):
432
  p_n = gr.Slider(6, 100, 30, step=1, label="Nodes")
433
  p_d = gr.Slider(16, 512, 128, step=16, label="Tensor Dim")
434
- p_s = gr.Dropdown(NO_BASE_STRATEGIES, value="weight_average", label="Strategy")
435
  p_p = gr.Slider(2, 10, 4, step=1, label="Partitions")
436
  p_seed = gr.Number(42, label="Seed", precision=0)
437
  p_btn = gr.Button("Run", variant="primary")
@@ -440,26 +443,27 @@ with gr.Blocks(title="CRDT-Merge Convergence Lab", theme=gr.themes.Default(prima
440
  p_json = gr.Textbox(label="JSON", lines=8)
441
  p_btn.click(run_partition_experiment, [p_n, p_d, p_s, p_p, p_seed], [p_log, p_json])
442
 
443
- with gr.TabItem("Strategy Sweep"):
444
- gr.Markdown("Every non-base strategy tested for convergence.")
445
  with gr.Row():
446
  with gr.Column(scale=1):
447
  sw_n = gr.Slider(3, 30, 10, step=1, label="Nodes")
448
  sw_d = gr.Slider(16, 256, 64, step=16, label="Tensor Dim")
449
  sw_seed = gr.Number(42, label="Seed", precision=0)
 
450
  sw_btn = gr.Button("Run Sweep", variant="primary")
451
  with gr.Column(scale=2):
452
  sw_log = gr.Textbox(label="Log", lines=30, max_lines=60)
453
  sw_json = gr.Textbox(label="JSON", lines=8)
454
- sw_btn.click(run_strategy_sweep, [sw_n, sw_d, sw_seed], [sw_log, sw_json])
455
 
456
  with gr.TabItem("Scalability"):
457
- gr.Markdown("Measure convergence overhead from 2 to N nodes.")
458
  with gr.Row():
459
  with gr.Column(scale=1):
460
  sc_m = gr.Slider(10, 100, 50, step=5, label="Max Nodes")
461
  sc_d = gr.Slider(16, 256, 64, step=16, label="Tensor Dim")
462
- sc_s = gr.Dropdown(NO_BASE_STRATEGIES, value="weight_average", label="Strategy")
463
  sc_seed = gr.Number(42, label="Seed", precision=0)
464
  sc_btn = gr.Button("Run Benchmark", variant="primary")
465
  with gr.Column(scale=2):
 
15
  import random
16
  import json
17
  from collections import defaultdict
 
18
 
19
  from crdt_merge.model import CRDTMergeState
20
 
21
+ ALL_STRATEGIES = sorted(CRDTMergeState.KNOWN_STRATEGIES)
22
+ BASE_REQUIRED = CRDTMergeState.BASE_REQUIRED
23
+ NO_BASE_STRATEGIES = sorted(set(ALL_STRATEGIES) - BASE_REQUIRED)
24
 
25
 
26
+ def _make_state(strategy, base=None):
27
+ """Create a CRDTMergeState, providing base if the strategy requires it."""
28
+ if strategy in BASE_REQUIRED:
29
+ return CRDTMergeState(strategy, base=base)
30
+ return CRDTMergeState(strategy)
31
+
32
+
33
+ # ===== Experiment 1: Multi-Node Convergence =====
34
+
35
  def run_convergence_experiment(n_nodes, tensor_dim, strategy, n_random_orderings=5, seed=42):
36
  n_nodes, tensor_dim, n_random_orderings, seed = int(n_nodes), int(tensor_dim), int(n_random_orderings), int(seed)
37
  np.random.seed(seed)
38
  shape = (tensor_dim, tensor_dim)
39
  total_params = tensor_dim * tensor_dim * n_nodes
40
+ base = np.random.randn(*shape).astype(np.float64)
41
  tensors = [np.random.randn(*shape).astype(np.float64) for _ in range(n_nodes)]
42
 
43
  log = []
44
  log.append(f"{'='*72}")
45
  log.append(f" MULTI-NODE CONVERGENCE EXPERIMENT")
46
  log.append(f"{'='*72}")
47
+ log.append(f" Nodes: {n_nodes} | Tensor: {shape} | Params: {total_params:,}")
48
+ log.append(f" Strategy: {strategy} | Orderings: {n_random_orderings}")
49
  log.append(f"{'='*72}\n")
50
 
51
  all_resolved, all_hashes, ordering_times = [], [], []
 
54
  rng = random.Random(seed + oidx)
55
  nodes = []
56
  for i in range(n_nodes):
57
+ s = _make_state(strategy, base)
58
  s.add(tensors[i], model_id=f"node-{i}")
59
  nodes.append(s)
60
 
61
  t0 = time.perf_counter()
62
+ order = list(range(n_nodes)); rng.shuffle(order)
 
63
  merge_count = 0
64
  for i in order:
65
+ targets = list(range(n_nodes)); rng.shuffle(targets)
 
66
  for j in targets:
67
  if i != j:
68
  nodes[i].merge(nodes[j])
 
82
  ordering_times.append(gossip_ms)
83
 
84
  status = "CONVERGED" if (unique == 1 and bitwise) else "DIVERGED"
85
+ log.append(f" Ordering {oidx+1}: {status} | gossip {gossip_ms:7.1f}ms | resolve {resolve_ms:7.1f}ms | merges {merge_count:,} | max_diff {max_diff:.1e}")
 
 
 
86
 
87
  cross_equal = all(np.array_equal(all_resolved[0], r) for r in all_resolved[1:])
88
  cross_hashes = len(set(all_hashes)) == 1
 
89
  log.append(f"\n{'~'*72}")
90
  log.append(f" CROSS-ORDERING VERIFICATION")
91
  log.append(f"{'~'*72}")
 
93
  log.append(f" All orderings bitwise equal: {'YES' if cross_equal else 'NO'}")
94
  log.append(f" Canonical hash: {all_hashes[0][:40]}...")
95
  log.append(f" Avg gossip: {np.mean(ordering_times):.1f}ms")
 
96
  verdict = "PASS" if (cross_equal and cross_hashes) else "FAIL"
97
  log.append(f"\n VERDICT: {verdict}")
98
 
 
106
  return "\n".join(log), json.dumps(summary, indent=2)
107
 
108
 
109
+ # ===== Experiment 2: Network Partition & Healing =====
110
+
111
  def run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions=3, seed=42):
112
  n_nodes, tensor_dim, n_partitions, seed = int(n_nodes), int(tensor_dim), int(n_partitions), int(seed)
113
  np.random.seed(seed)
114
  shape = (tensor_dim, tensor_dim)
115
+ base = np.random.randn(*shape).astype(np.float64)
116
  tensors = [np.random.randn(*shape).astype(np.float64) for _ in range(n_nodes)]
117
 
118
  log = []
 
124
 
125
  nodes = []
126
  for i in range(n_nodes):
127
+ s = _make_state(strategy, base)
128
  s.add(tensors[i], model_id=f"node-{i}")
129
  nodes.append(s)
130
 
 
140
  for pid, members in partitions.items():
141
  for i in members:
142
  for j in members:
143
+ if i != j: nodes[i].merge(nodes[j])
 
144
  partition_ms = (time.perf_counter() - t0) * 1000
145
  log.append(f"\n Partition gossip time: {partition_ms:.1f}ms\n")
146
 
 
151
  ok = len(h) == 1
152
  log.append(f" Partition {pid}: {'consistent' if ok else 'INCONSISTENT'} hash: {list(h)[0][:24]}...")
153
 
154
+ all_unique = set()
155
+ for h in partition_hashes.values(): all_unique.update(h)
156
+ partitions_differ = len(all_unique) >= min(n_partitions, n_nodes)
 
157
  log.append(f"\n Partitions differ from each other: {'YES' if partitions_differ else 'NO'}")
158
 
159
  log.append(f"\n -- Phase 2: Partition Healing (full gossip resumes) --\n")
 
160
  t0 = time.perf_counter()
161
  for i in range(n_nodes):
162
  for j in range(n_nodes):
163
+ if i != j: nodes[i].merge(nodes[j])
 
164
  heal_ms = (time.perf_counter() - t0) * 1000
165
 
166
  healed = set(n.state_hash for n in nodes)
 
168
  log.append(f" Healing time: {heal_ms:.1f}ms")
169
  log.append(f" All {n_nodes} nodes converged: {'YES' if all_consistent else 'NO'}")
170
 
 
171
  resolved = [n.resolve() for n in nodes]
 
172
  bitwise = all(np.array_equal(resolved[0], r) for r in resolved[1:])
 
173
  log.append(f" All resolved bitwise identical: {'YES' if bitwise else 'NO'}")
 
174
  log.append(f" Final hash: {list(healed)[0][:40]}...")
 
175
  verdict = "PASS" if (all_consistent and bitwise) else "FAIL"
176
  log.append(f"\n VERDICT: {verdict}")
177
 
 
186
  return "\n".join(log), json.dumps(summary, indent=2)
187
 
188
 
189
+ # ===== Experiment 3: Cross-Strategy Sweep (ALL 26) =====
190
+
191
+ SLOW_STRATEGIES = {"evolutionary_merge", "genetic_merge"}
192
+
193
+ def run_strategy_sweep(n_nodes, tensor_dim, seed=42, skip_slow=True, progress=gr.Progress()):
194
  n_nodes, tensor_dim, seed = int(n_nodes), int(tensor_dim), int(seed)
195
  np.random.seed(seed)
196
  shape = (tensor_dim, tensor_dim)
197
+ base = np.random.randn(*shape).astype(np.float64)
198
  tensors = [np.random.randn(*shape).astype(np.float64) for _ in range(n_nodes)]
199
 
200
+ strategies = ALL_STRATEGIES
201
+ if skip_slow:
202
+ strategies = [s for s in strategies if s not in SLOW_STRATEGIES]
203
+ skipped = sorted(SLOW_STRATEGIES)
204
+ else:
205
+ skipped = []
206
+
207
  log = []
208
  log.append(f"{'='*72}")
209
+ log.append(f" CROSS-STRATEGY CONVERGENCE SWEEP — ALL 26 STRATEGIES")
210
  log.append(f"{'='*72}")
211
+ log.append(f" Nodes: {n_nodes} | Tensor: {shape} | Testing: {len(strategies)}/{len(ALL_STRATEGIES)}")
212
+ if skipped:
213
+ log.append(f" Skipped (slow): {', '.join(skipped)}")
214
  log.append(f"{'='*72}\n")
215
 
216
+ header = f" {'Strategy':<28s} {'Base':>4s} {'Conv':>5s} {'Gossip':>9s} {'Resolve':>9s} {'Hash':>24s}"
217
  log.append(header)
218
+ log.append(f" {'~'*28} {'~'*4} {'~'*5} {'~'*9} {'~'*9} {'~'*24}")
219
 
220
  pass_count, fail_count = 0, 0
221
  rows = []
222
 
223
+ for idx, strat in enumerate(strategies):
224
+ progress((idx + 1) / len(strategies), f"Testing {strat}...")
225
  try:
226
+ needs_base = strat in BASE_REQUIRED
227
  rng = random.Random(seed)
228
  nds = []
229
  for i in range(n_nodes):
230
+ s = _make_state(strat, base)
231
+ s.add(tensors[i], model_id=f"n-{i}")
232
  nds.append(s)
233
 
234
  t0 = time.perf_counter()
235
+ order = list(range(n_nodes)); rng.shuffle(order)
 
236
  for i in order:
237
+ tgts = list(range(n_nodes)); rng.shuffle(tgts)
 
238
  for j in tgts:
239
+ if i != j: nds[i].merge(nds[j])
 
240
  g_ms = (time.perf_counter() - t0) * 1000
241
 
242
  hashes = [n.state_hash for n in nds]
 
245
  r_ms = (time.perf_counter() - t0) * 1000
246
 
247
  ok = len(set(hashes)) == 1 and all(np.array_equal(resolved[0], r) for r in resolved[1:])
248
+ if ok: pass_count += 1
249
+ else: fail_count += 1
 
 
250
 
251
+ base_tag = " Y " if needs_base else " "
252
+ log.append(f" {strat:<28s} {base_tag} {'PASS' if ok else 'FAIL':>5s} {g_ms:8.1f}ms {r_ms:8.1f}ms {hashes[0][:24]}")
253
+ rows.append({"strategy": strat, "needs_base": needs_base, "converged": bool(ok),
254
+ "gossip_ms": round(g_ms, 1), "resolve_ms": round(r_ms, 1)})
255
  except Exception as e:
256
  fail_count += 1
257
+ log.append(f" {strat:<28s} ERR {str(e)[:50]}")
258
  rows.append({"strategy": strat, "converged": False, "error": str(e)[:50]})
259
 
260
+ # Add skipped strategies as noted
261
+ for strat in skipped:
262
+ rows.append({"strategy": strat, "converged": "skipped", "note": "evolutionary/genetic (~60s each)"})
263
+
264
+ tested = pass_count + fail_count
265
+ log.append(f"\n{'~'*72}")
266
+ log.append(f" Tested: {tested}/{len(ALL_STRATEGIES)} strategies | Passed: {pass_count}/{tested}")
267
+ if skipped:
268
+ log.append(f" Skipped: {len(skipped)} (evolutionary strategies, ~60s each on CPU)")
269
+ log.append(f" To include: uncheck 'Skip slow strategies'")
270
+ verdict = f"ALL {tested} PASS" if fail_count == 0 else f"{fail_count}/{tested} FAILED"
271
  log.append(f"\n VERDICT: {verdict}")
272
 
273
+ summary = {"total_strategies": len(ALL_STRATEGIES), "tested": tested,
274
+ "passed": pass_count, "failed": fail_count, "skipped": len(skipped), "results": rows}
275
  return "\n".join(log), json.dumps(summary, indent=2)
276
 
277
 
278
+ # ===== Experiment 4: Scalability Benchmark =====
279
+
280
  def run_scale_benchmark(max_nodes, tensor_dim, strategy, seed=42, progress=gr.Progress()):
281
  max_nodes, tensor_dim, seed = int(max_nodes), int(tensor_dim), int(seed)
282
  np.random.seed(seed)
283
  shape = (tensor_dim, tensor_dim)
284
+ base = np.random.randn(*shape).astype(np.float64)
285
 
286
  log = []
287
  log.append(f"{'='*72}")
 
296
 
297
  steps = sorted(set([2, 5, 10, 20, 30, 50, 75, 100]) & set(range(2, max_nodes + 1)))
298
  if max_nodes not in steps and max_nodes >= 2:
299
+ steps.append(max_nodes); steps.sort()
 
300
 
301
  all_tensors = [np.random.randn(*shape).astype(np.float64) for _ in range(max_nodes)]
302
  node_counts, gossip_times, resolve_times = [], [], []
 
305
  progress((si + 1) / len(steps), f"Testing {n} nodes...")
306
  nds = []
307
  for i in range(n):
308
+ s = _make_state(strategy, base)
309
+ s.add(all_tensors[i], model_id=f"n-{i}")
310
  nds.append(s)
311
 
312
  t0 = time.perf_counter()
 
314
  for i in range(n):
315
  for j in range(n):
316
  if i != j:
317
+ nds[i].merge(nds[j]); merge_ops += 1
 
318
  g_ms = (time.perf_counter() - t0) * 1000
319
 
320
  t0 = time.perf_counter()
 
322
  r_ms = (time.perf_counter() - t0) * 1000
323
 
324
  ok = len(set(nd.state_hash for nd in nds)) == 1 and all(np.array_equal(resolved[0], r) for r in resolved[1:])
325
+ node_counts.append(n); gossip_times.append(g_ms); resolve_times.append(r_ms)
 
 
326
 
327
+ log.append(f" {n:>6d} {n * tensor_dim**2:>12,} {g_ms:>9.1f}ms {r_ms:>9.1f}ms {merge_ops:>10,} {'PASS' if ok else 'FAIL':>5s}")
 
 
 
328
 
329
  log.append(f"\n merge() is O(1) per call - independent of tensor size")
 
330
  log.append(f" 100% convergence at all tested scales")
331
 
332
+ summary = {"node_counts": node_counts, "gossip_times_ms": [round(g, 1) for g in gossip_times],
333
+ "resolve_times_ms": [round(r, 1) for r in resolve_times], "strategy": strategy}
 
 
 
 
334
  return "\n".join(log), json.dumps(summary, indent=2)
335
 
336
 
337
+ # ===== Full Suite =====
338
+
339
+ def run_full_experiment(n_nodes, tensor_dim, strategy, n_orderings, n_partitions, seed, skip_slow, progress=gr.Progress()):
340
+ all_logs, summaries = [], {}
341
 
342
  progress(0.05, "Running multi-node convergence...")
343
  l, s = run_convergence_experiment(n_nodes, tensor_dim, strategy, n_orderings, seed)
344
+ all_logs.append(l); summaries["convergence"] = json.loads(s)
 
345
 
346
  progress(0.30, "Running partition experiment...")
347
  l, s = run_partition_experiment(n_nodes, tensor_dim, strategy, n_partitions, seed)
348
+ all_logs.append(l); summaries["partition"] = json.loads(s)
 
349
 
350
+ sweep_n = min(int(n_nodes), 10); sweep_d = min(int(tensor_dim), 64)
351
+ progress(0.55, "Running strategy sweep (all 26)...")
352
+ l, s = run_strategy_sweep(sweep_n, sweep_d, seed, skip_slow)
353
+ all_logs.append(l); summaries["strategy_sweep"] = json.loads(s)
 
 
 
354
 
355
  progress(0.80, "Running scalability benchmark...")
356
+ l, s = run_scale_benchmark(min(int(n_nodes), 50), sweep_d, strategy, seed)
357
+ all_logs.append(l); summaries["scalability"] = json.loads(s)
 
358
 
359
  progress(1.0, "Complete!")
360
 
 
368
  f"{'='*72}",
369
  f" Multi-node convergence ({int(n_nodes)} nodes, {int(n_orderings)} orderings): {'PASS' if c else 'FAIL'}",
370
  f" Network partition healing ({int(n_partitions)} partitions): {'PASS' if p else 'FAIL'}",
371
+ f" Cross-strategy sweep ({summaries['strategy_sweep']['tested']}/{summaries['strategy_sweep']['total_strategies']} strategies): {'PASS' if sw else 'FAIL'}",
372
  f" Scalability benchmark: PASS",
373
  f"{'='*72}",
374
  ]
 
375
  if c and p and sw:
376
  report.append(f"\n >>> ALL EXPERIMENTS PASSED - CRDT COMPLIANCE VERIFIED <<<")
377
 
378
+ return "\n\n".join(all_logs) + "\n" + "\n".join(report), json.dumps(summaries, indent=2)
 
379
 
380
 
381
+ # ===== Gradio UI =====
382
 
383
  DESCRIPTION = """
384
  # CRDT-Merge Multi-Node Convergence Laboratory
 
389
 
390
  > **Patent Pending**: UK Application No. 2607132.4 | **Library**: [crdt-merge](https://pypi.org/project/crdt-merge/) v0.9.4
391
 
392
+ **Four experiments**: Multi-node convergence | Network partition & healing | All 26 strategies | Scalability benchmark
393
  """
394
 
395
  with gr.Blocks(title="CRDT-Merge Convergence Lab", theme=gr.themes.Default(primary_hue="slate", neutral_hue="slate")) as demo:
 
397
 
398
  with gr.Tabs():
399
  with gr.TabItem("Full Suite"):
400
+ gr.Markdown("### Run all four experiments tests all 26 merge strategies")
401
  with gr.Row():
402
  with gr.Column(scale=1):
403
  n_nodes = gr.Slider(3, 100, 30, step=1, label="Nodes")
404
  tensor_dim = gr.Slider(16, 512, 128, step=16, label="Tensor Dim (d x d)")
405
+ strategy = gr.Dropdown(ALL_STRATEGIES, value="weight_average", label="Primary Strategy")
406
  n_orderings = gr.Slider(2, 20, 5, step=1, label="Random Orderings")
407
  n_partitions = gr.Slider(2, 10, 3, step=1, label="Partitions")
408
  seed = gr.Number(42, label="Seed", precision=0)
409
+ skip_slow = gr.Checkbox(True, label="Skip evolutionary strategies (~2 min each on CPU)")
410
  run_btn = gr.Button("Run Full Suite", variant="primary", size="lg")
411
  with gr.Column(scale=2):
412
  out_log = gr.Textbox(label="Experiment Log", lines=35, max_lines=80)
413
  out_json = gr.Textbox(label="JSON", lines=10, max_lines=40)
414
+ run_btn.click(run_full_experiment, [n_nodes, tensor_dim, strategy, n_orderings, n_partitions, seed, skip_slow], [out_log, out_json])
415
 
416
  with gr.TabItem("Convergence"):
417
+ gr.Markdown("### N nodes merge in different random orderings — all must produce identical results")
418
  with gr.Row():
419
  with gr.Column(scale=1):
420
  c_n = gr.Slider(3, 100, 30, step=1, label="Nodes")
421
  c_d = gr.Slider(16, 512, 128, step=16, label="Tensor Dim")
422
+ c_s = gr.Dropdown(ALL_STRATEGIES, value="slerp", label="Strategy")
423
  c_o = gr.Slider(2, 20, 8, step=1, label="Orderings")
424
  c_seed = gr.Number(42, label="Seed", precision=0)
425
  c_btn = gr.Button("Run", variant="primary")
 
429
  c_btn.click(run_convergence_experiment, [c_n, c_d, c_s, c_o, c_seed], [c_log, c_json])
430
 
431
  with gr.TabItem("Partition & Healing"):
432
+ gr.Markdown("### Split nodes into isolated partitions, gossip internally, heal, verify convergence")
433
  with gr.Row():
434
  with gr.Column(scale=1):
435
  p_n = gr.Slider(6, 100, 30, step=1, label="Nodes")
436
  p_d = gr.Slider(16, 512, 128, step=16, label="Tensor Dim")
437
+ p_s = gr.Dropdown(ALL_STRATEGIES, value="ties", label="Strategy")
438
  p_p = gr.Slider(2, 10, 4, step=1, label="Partitions")
439
  p_seed = gr.Number(42, label="Seed", precision=0)
440
  p_btn = gr.Button("Run", variant="primary")
 
443
  p_json = gr.Textbox(label="JSON", lines=8)
444
  p_btn.click(run_partition_experiment, [p_n, p_d, p_s, p_p, p_seed], [p_log, p_json])
445
 
446
+ with gr.TabItem("All 26 Strategies"):
447
+ gr.Markdown("### Every merge strategy tested for convergence — 13 base-free + 13 base-required")
448
  with gr.Row():
449
  with gr.Column(scale=1):
450
  sw_n = gr.Slider(3, 30, 10, step=1, label="Nodes")
451
  sw_d = gr.Slider(16, 256, 64, step=16, label="Tensor Dim")
452
  sw_seed = gr.Number(42, label="Seed", precision=0)
453
+ sw_skip = gr.Checkbox(True, label="Skip evolutionary strategies (~2 min each)")
454
  sw_btn = gr.Button("Run Sweep", variant="primary")
455
  with gr.Column(scale=2):
456
  sw_log = gr.Textbox(label="Log", lines=30, max_lines=60)
457
  sw_json = gr.Textbox(label="JSON", lines=8)
458
+ sw_btn.click(run_strategy_sweep, [sw_n, sw_d, sw_seed, sw_skip], [sw_log, sw_json])
459
 
460
  with gr.TabItem("Scalability"):
461
+ gr.Markdown("### Measure convergence overhead from 2 to N nodes")
462
  with gr.Row():
463
  with gr.Column(scale=1):
464
  sc_m = gr.Slider(10, 100, 50, step=5, label="Max Nodes")
465
  sc_d = gr.Slider(16, 256, 64, step=16, label="Tensor Dim")
466
+ sc_s = gr.Dropdown(ALL_STRATEGIES, value="weight_average", label="Strategy")
467
  sc_seed = gr.Number(42, label="Seed", precision=0)
468
  sc_btn = gr.Button("Run Benchmark", variant="primary")
469
  with gr.Column(scale=2):
requirements.txt CHANGED
@@ -1,3 +1,3 @@
1
- crdt-merge[all]>=0.9.3
2
  numpy>=1.24.0
3
  gradio>=5.0.0,<6.0.0
 
1
+ crdt-merge[all]>=0.9.4
2
  numpy>=1.24.0
3
  gradio>=5.0.0,<6.0.0