#!/usr/bin/env node
'use strict';
// בדיקות יחידה ל-Ξ/DAG (הרכבת-משימות) ול-SkyLattice-Φ (זרימה בזמן-אמת).
// הכול דטרמיניסטי (שעון-וירטואלי), בלי רשת ובלי טיימרים אמיתיים.
process.env.SC_TEST = '1';
const t = require('./dump-test.js').__test;
const assert = require('node:assert');

let pass = 0;
const ok = (name, cond) => { assert.ok(cond, name); console.log('PASS ' + name); pass++; };
const solve = (program) => Promise.resolve(t.xiSolve(program));

(async () => {
  // ═══════════ Ξ/DAG: הרכבת-משימות ═══════════
  ok('DAG refs found in nested args', JSON.stringify(t.xiRefs({ vars: { c: { $ref: 'a', path: 'count' }, d: { $ref: 'b' } } }).sort()) === JSON.stringify(['a', 'b']));
  ok('DAG subst replaces $ref with upstream (by path)', t.xiSubst({ c: { $ref: 'a', path: 'count' } }, { a: { count: 25 } }).c === 25);

  const dag = {
    nodes: {
      a: { op: 'primes', args: { upTo: 100 } },                                  // count=25
      b: { op: 'formula', args: { expr: 'c*2', vars: { c: { $ref: 'a', path: 'count' } } } }, // 50
      c: { op: 'factorize', args: { n: { $ref: 'b' } } },                        // factorize(50)=[2,5,5]
    },
    output: 'c',
  };
  const order = t.xiDagPlan(dag);
  ok('DAG topological order respects deps (a<b<c)', order.indexOf('a') < order.indexOf('b') && order.indexOf('b') < order.indexOf('c'));

  const ev = await t.xiDagEval(dag, solve);
  ok('DAG pipes output→input correctly (final = factorize(50))', JSON.stringify(ev.output) === JSON.stringify([2, 5, 5]));
  ok('DAG intermediate results captured', ev.results.a.count === 25 && ev.results.b === 50);
  ok('DAG deterministic dagDigest', ev.dagDigest === (await t.xiDagEval(dag, solve)).dagDigest);
  ok('DAG id deterministic + content-addressed', t.xiDagId(dag) === t.xiDagId(JSON.parse(JSON.stringify(dag))));

  let cyc = false; try { t.xiDagPlan({ nodes: { x: { op: 'formula', args: { vars: { v: { $ref: 'y' } } } }, y: { op: 'formula', args: { vars: { v: { $ref: 'x' } } } } } }); } catch { cyc = true; }
  ok('DAG detects cycles (rejects non-DAG)', cyc);
  let miss = false; try { t.xiDagPlan({ nodes: { a: { op: 'formula', args: { vars: { v: { $ref: 'ghost' } } } } } }); } catch { miss = true; }
  ok('DAG rejects reference to missing node', miss);

  // ── value nodes + runtime input override + in-run memo (computed live site) ──
  const paramDag = {
    nodes: {
      p: { value: 10 },                                                        // קלט מוזרק
      out: { op: 'formula', args: { expr: 'x*x', vars: { x: { $ref: 'p' } } } },
    },
    output: 'out',
  };
  const e10 = await t.xiDagEval(paramDag, solve);
  ok('DAG value node feeds computation (10→100)', e10.output === 100);
  // דריסת-קלט: אותה צנרת, קלט אחר → תוצאה אחרת (בסיס לאתר-חי מחושב פר-בקשה)
  const overridden = { ...paramDag, nodes: { ...paramDag.nodes, p: { value: 7 } } };
  const e7 = await t.xiDagEval(overridden, solve);
  ok('DAG input override changes result (7→49)', e7.output === 49);
  ok('DAG value node reported as non-compute source', e10.sources.p === 'value');

  // מזכר תוך-ריצה: אותו jobId מופיע פעמיים → פותר פעם אחת בלבד
  let solveCalls = 0;
  const counting = (program) => { solveCalls++; return Promise.resolve(t.xiSolve(program)); };
  const dupDag = {
    nodes: {
      a: { op: 'factorize', args: { n: 360 } },
      b: { op: 'factorize', args: { n: 360 } },                                // תוכנית זהה ל-a
      out: { coalesce: [{ $ref: 'a' }, { $ref: 'b' }] },
    },
    output: 'out',
  };
  await t.xiDagEval(dupDag, counting);
  ok('DAG in-run memo dedupes identical jobs (1 solve, not 2)', solveCalls === 1);

  // ═══════════ Φ: זרימה בזמן-אמת עם שעון-וירטואלי ═══════════
  const clk = t.clock ? null : null; // (clock דורש מופע Cloud; משתמשים ישירות במחלקות)
  const Clock = t.Clock, Stream = t.Stream, Channel = t.Channel;

  // תדר → map/filter
  const c1 = new Clock();
  const evens = Stream.of(c1, (e) => e).filter((v) => v % 2 === 0).map((v) => v * 10);
  const got = [];
  evens.subscribe((v) => got.push(v));
  for (let i = 0; i < 6; i++) c1.tick(); // epochs 1..6
  ok('Φ stream map/filter over frequency ticks', JSON.stringify(got) === JSON.stringify([20, 40, 60]));

  // throttle (decimation): 1 מכל 3
  const c2 = new Clock();
  const thr = Stream.of(c2, (e) => e).throttle(3);
  const th = []; thr.subscribe((v) => th.push(v));
  for (let i = 0; i < 9; i++) c2.tick();
  ok('Φ throttle emits 1 of every 3 epochs', JSON.stringify(th) === JSON.stringify([3, 6, 9]));

  // scan (הצטברות מצב)
  const c3 = new Clock();
  const sums = Stream.of(c3, () => 1).scan((a, b) => a + b, 0);
  let lastSum = 0; sums.subscribe((v) => (lastSum = v));
  for (let i = 0; i < 5; i++) c3.tick();
  ok('Φ scan accumulates running total', lastSum === 5);

  // formula transform על זרם
  const c4 = new Clock();
  const fx = Stream.of(c4, (e) => e).formula('x*x + 1');
  const fxo = []; fx.subscribe((v) => fxo.push(v));
  for (let i = 0; i < 3; i++) c4.tick();
  ok('Φ formula transform (x*x+1)', JSON.stringify(fxo) === JSON.stringify([2, 5, 10]));

  // ═══════════ מחסום-סינכרון: תוצאות שלמות בלבד ═══════════
  const sA = new Stream(), sB = new Stream();
  const zipped = Stream.zipComplete([sA, sB]);
  const tuples = []; zipped.subscribe((tup, e) => tuples.push({ e, tup }));
  sA.push('a1', 1);              // epoch 1 חלקי — לא נפלט
  sB.push('b2', 2);              // epoch 2 חלקי — לא נפלט
  ok('Φ barrier holds partial epochs (no emit yet)', tuples.length === 0);
  sB.push('b1', 1);              // epoch 1 הושלם → נפלט
  ok('Φ barrier emits only COMPLETE tuple for epoch', tuples.length === 1 && JSON.stringify(tuples[0]) === JSON.stringify({ e: 1, tup: ['a1', 'b1'] }));
  sA.push('a2', 2);              // epoch 2 הושלם → נפלט
  ok('Φ barrier completes later epoch independently', tuples.length === 2 && tuples[1].tup[0] === 'a2' && tuples[1].tup[1] === 'b2');

  // merge לפי נתיבים (lanes)
  const m1 = new Stream(), m2 = new Stream();
  const merged = Stream.merge([m1, m2]);
  const ml = []; merged.subscribe((v) => ml.push(v));
  m1.push('x', 1); m2.push('y', 1);
  ok('Φ merge tags lane index', JSON.stringify(ml) === JSON.stringify([{ lane: 0, value: 'x' }, { lane: 1, value: 'y' }]));

  // ═══════════ Channel (CSP): send/recv/backpressure/close ═══════════
  const ch = new Channel(2);
  ok('Φ channel send within capacity', ch.send(1) === true && ch.send(2) === true);
  ok('Φ channel backpressure when full', ch.send(3) === false);
  const r1 = await ch.recv();
  ok('Φ channel recv FIFO', r1.value === 1 && r1.done === false);
  const ch2 = new Channel();
  const pending = ch2.recv();          // ממתין לפני send
  ch2.send('later');
  const rr = await pending;
  ok('Φ channel recv resolves on later send', rr.value === 'later');
  ch2.close();
  const done = await ch2.recv();
  ok('Φ channel signals done after close', done.done === true);

  console.log('\nALL ' + pass + ' FLUX/DAG/SITE TESTS PASSED');
})().catch((e) => { console.error('FAIL:', e.message); process.exit(1); });
