From 8ed1ccb3dd97194c9745441d3b1bc173b2724d06 Mon Sep 17 00:00:00 2001 From: Devin Ivy Date: Tue, 7 Apr 2026 01:19:41 -0400 Subject: [PATCH] feat: add examples and freeze module registry on first task dispatch Restructure examples into separate directories so moroutine definitions are isolated modules (preventing fork bombs from top-level side effects). Freeze per-module mo() registrations once a task is dispatched, erroring on any late dynamic mo() calls that would cause ID mismatches in workers. Co-Authored-By: Claude Opus 4.6 (1M context) --- examples/non-blocking/fibonacci.ts | 10 +++++++++ examples/non-blocking/main.ts | 30 +++++++++++++++++++++++++++ examples/parallel-batch/heavy-work.ts | 8 +++++++ examples/parallel-batch/main.ts | 30 +++++++++++++++++++++++++++ examples/primes/is-prime.ts | 11 ++++++++++ examples/primes/main.ts | 19 +++++++++++++++++ examples/tsconfig.json | 4 ++++ src/execute.ts | 3 +++ src/mo.ts | 10 ++++++++- src/registry.ts | 10 +++++++++ test/fixtures/freeze.ts | 7 +++++++ test/freeze.test.ts | 12 +++++++++++ 12 files changed, 153 insertions(+), 1 deletion(-) create mode 100644 examples/non-blocking/fibonacci.ts create mode 100644 examples/non-blocking/main.ts create mode 100644 examples/parallel-batch/heavy-work.ts create mode 100644 examples/parallel-batch/main.ts create mode 100644 examples/primes/is-prime.ts create mode 100644 examples/primes/main.ts create mode 100644 examples/tsconfig.json create mode 100644 test/fixtures/freeze.ts create mode 100644 test/freeze.test.ts diff --git a/examples/non-blocking/fibonacci.ts b/examples/non-blocking/fibonacci.ts new file mode 100644 index 0000000..f185f62 --- /dev/null +++ b/examples/non-blocking/fibonacci.ts @@ -0,0 +1,10 @@ +import { mo } from '../../src/index.ts'; + +export const fibonacci = mo(import.meta, (n: number): number => { + // Deliberately naive recursive fibonacci — CPU-heavy but bounded + function fib(x: number): number { + if (x <= 1) return x; + return fib(x - 1) + fib(x - 2); + } + return fib(n); +}); diff --git a/examples/non-blocking/main.ts b/examples/non-blocking/main.ts new file mode 100644 index 0000000..320f7df --- /dev/null +++ b/examples/non-blocking/main.ts @@ -0,0 +1,30 @@ +// Demonstrates that the main thread event loop stays responsive +// while heavy work runs on worker threads. +// +// Run: node --experimental-strip-types examples/non-blocking/main.ts + +import { workerPool } from '../../src/index.ts'; +import { fibonacci } from './fibonacci.ts'; + +// Tick a counter on the main thread to prove it's not blocked +let ticks = 0; +const interval = setInterval(() => { + ticks++; + process.stdout.write(`\r main thread tick #${ticks}`); +}, 100); + +console.log('Computing fibonacci(42) on a worker pool...'); +console.log('Meanwhile, the main thread keeps ticking:\n'); + +const pool = workerPool(2); +const start = performance.now(); +const [a, b] = await Promise.all([ + pool(fibonacci(42)), + pool(fibonacci(41)), +]); +const elapsed = (performance.now() - start).toFixed(0); +pool[Symbol.dispose](); + +clearInterval(interval); +console.log(`\n\nResults: fib(42)=${a}, fib(41)=${b}`); +console.log(`Computed in ${elapsed}ms with ${ticks} main-thread ticks (not blocked!)`); diff --git a/examples/parallel-batch/heavy-work.ts b/examples/parallel-batch/heavy-work.ts new file mode 100644 index 0000000..24f7af9 --- /dev/null +++ b/examples/parallel-batch/heavy-work.ts @@ -0,0 +1,8 @@ +import { mo } from '../../src/index.ts'; + +export const heavyWork = mo(import.meta, (item: number): number => { + // Simulate CPU-bound work (~50ms per item) + const start = Date.now(); + while (Date.now() - start < 50) { /* busy wait */ } + return item * item; +}); diff --git a/examples/parallel-batch/main.ts b/examples/parallel-batch/main.ts new file mode 100644 index 0000000..de1c680 --- /dev/null +++ b/examples/parallel-batch/main.ts @@ -0,0 +1,30 @@ +// Process a batch of work items in parallel using a worker pool. +// Compares sequential (dedicated worker) vs parallel (pool) execution. +// +// Run: node --experimental-strip-types examples/parallel-batch/main.ts + +import { workerPool } from '../../src/index.ts'; +import { heavyWork } from './heavy-work.ts'; + +const items = Array.from({ length: 20 }, (_, i) => i + 1); + +// Sequential: each call awaits on the dedicated worker +console.log('Sequential (dedicated worker)...'); +const seqStart = performance.now(); +const seqResults = []; +for (const item of items) { + seqResults.push(await heavyWork(item)); +} +const seqTime = (performance.now() - seqStart).toFixed(0); +console.log(` ${items.length} items in ${seqTime}ms\n`); + +// Parallel: distribute across a pool of 4 workers +console.log('Parallel (worker pool, 4 workers)...'); +const pool = workerPool(4); +const parStart = performance.now(); +const parResults = await Promise.all(items.map((item) => pool(heavyWork(item)))); +const parTime = (performance.now() - parStart).toFixed(0); +pool[Symbol.dispose](); +console.log(` ${items.length} items in ${parTime}ms\n`); + +console.log(`Speedup: ~${(Number(seqTime) / Number(parTime)).toFixed(1)}x`); diff --git a/examples/primes/is-prime.ts b/examples/primes/is-prime.ts new file mode 100644 index 0000000..cb510dc --- /dev/null +++ b/examples/primes/is-prime.ts @@ -0,0 +1,11 @@ +import { mo } from '../../src/index.ts'; + +export const isPrime = mo(import.meta, (n: number): boolean => { + if (n < 2) return false; + if (n < 4) return true; + if (n % 2 === 0 || n % 3 === 0) return false; + for (let i = 5; i * i <= n; i += 6) { + if (n % i === 0 || n % (i + 2) === 0) return false; + } + return true; +}); diff --git a/examples/primes/main.ts b/examples/primes/main.ts new file mode 100644 index 0000000..2d11db9 --- /dev/null +++ b/examples/primes/main.ts @@ -0,0 +1,19 @@ +// CPU-bound prime checking offloaded to a dedicated worker thread. +// +// Run: node --experimental-strip-types examples/primes/main.ts + +import { isPrime } from './is-prime.ts'; + +const candidates = [ + 999_999_937, + 1_000_000_007, + 1_000_000_009, + 1_000_000_021, + 1_000_000_033, + 999_999_938, +]; + +for (const n of candidates) { + const prime = await isPrime(n); + console.log(`${n} is ${prime ? 'prime' : 'not prime'}`); +} diff --git a/examples/tsconfig.json b/examples/tsconfig.json new file mode 100644 index 0000000..379a994 --- /dev/null +++ b/examples/tsconfig.json @@ -0,0 +1,4 @@ +{ + "extends": "../tsconfig.json", + "include": ["."] +} diff --git a/src/execute.ts b/src/execute.ts index cdccb3e..39599e6 100644 --- a/src/execute.ts +++ b/src/execute.ts @@ -1,4 +1,5 @@ import type { Worker } from 'node:worker_threads'; +import { freezeModule } from './registry.ts'; let nextCallId = 0; const pending = new Map void; reject: (reason: any) => void }>(); @@ -17,6 +18,8 @@ export function setupWorker(worker: Worker): void { } export function execute(worker: Worker, id: string, args: unknown[]): Promise { + const url = id.slice(0, id.lastIndexOf('#')); + freezeModule(url); const callId = nextCallId++; return new Promise((resolve, reject) => { pending.set(callId, { resolve, reject }); diff --git a/src/mo.ts b/src/mo.ts index feea66e..020dbd9 100644 --- a/src/mo.ts +++ b/src/mo.ts @@ -1,4 +1,4 @@ -import { registry } from './registry.ts'; +import { registry, isModuleFrozen } from './registry.ts'; import { Task } from './task.ts'; const counters = new Map(); @@ -8,6 +8,14 @@ export function mo( fn: (...args: A) => R, ): (...args: A) => Task { const url = importMeta.url; + + if (isModuleFrozen(url)) { + throw new Error( + `Cannot call mo() for ${url} after a task from this module has been dispatched. ` + + 'All mo() calls must happen at module load time.', + ); + } + const index = counters.get(url) ?? 0; counters.set(url, index + 1); const id = `${url}#${index}`; diff --git a/src/registry.ts b/src/registry.ts index 2431bb0..5739585 100644 --- a/src/registry.ts +++ b/src/registry.ts @@ -1 +1,11 @@ export const registry = new Map any>(); + +const frozen = new Set(); + +export function freezeModule(url: string): void { + frozen.add(url); +} + +export function isModuleFrozen(url: string): boolean { + return frozen.has(url); +} diff --git a/test/fixtures/freeze.ts b/test/fixtures/freeze.ts new file mode 100644 index 0000000..eb5bd5b --- /dev/null +++ b/test/fixtures/freeze.ts @@ -0,0 +1,7 @@ +import { mo } from 'moroutine'; + +export const double = mo(import.meta, (x: number) => 2 * x); + +export function lateRegistration() { + return mo(import.meta, (x: number) => x + 1); +} diff --git a/test/freeze.test.ts b/test/freeze.test.ts new file mode 100644 index 0000000..2628e8d --- /dev/null +++ b/test/freeze.test.ts @@ -0,0 +1,12 @@ +import { describe, it } from 'node:test'; +import assert from 'node:assert/strict'; +import { double, lateRegistration } from './fixtures/freeze.ts'; + +describe('module freezing', () => { + it('throws when mo() is called after a task from that module has been dispatched', async () => { + await double(2); + assert.throws(lateRegistration, { + message: /Cannot call mo\(\).*after a task from this module has been dispatched/, + }); + }); +}); -- 2.51.2