Polyweave

@polyweave/loop-utils (1.1.0)

Published 2026-07-17 22:23:31 +00:00 by Dvorak

Installation

@polyweave:registry=https://hub.kl1.tenere.ai/api/packages/Polyweave/npm/
npm install @polyweave/loop-utils@1.1.0
"@polyweave/loop-utils": "1.1.0"

About this package

@polyweave/loop-utils

Shared utilities for self-improving optimization loops — Pareto-front computation, dataset reservoir sampling, experience library CRUD, sequential iteration helpers, async/sync/frame invokers, and Polyweave state-path read/write.

Synopsis

import { computeParetoFront, sampleDataset, applyExperienceOperations } from '@polyweave/loop-utils';

const front = computeParetoFront(candidates, 'paretoScores');

const sample = sampleDataset(dataset, 10);

const next = applyExperienceOperations(library, [
  { operation: 'ADD', entry: { text: 'new experience' } },
  { operation: 'UPDATE', id: 'exp-1', content: 'updated text' },
  { operation: 'DELETE', id: 'exp-2' }
]);

Install

npm install @polyweave/loop-utils

Why

Self-improving optimization loops (training, GRPO, DPO, sieve-based RL) need non-trivial iteration control, multi-objective selection, experience library management, dataset sampling, and state-path plumbing. These utilities are shared across @polyweave/training, @polyweave/optimization, @polyweave/self-improving, and @polyweave/nursery. Extracted to avoid 4 copies of computeParetoFront.

API

From utils.js

runSeries(items, iterator, done)

Sequentially iterates items using a callback-based iterator.

Signature: runSeries(items: Array<*>, iterator: Function, done: Function): void

Name Type Required Default Description
items Array Yes Items to iterate.
iterator (item, next) => void Yes Called for each item. Calls next(err?) to proceed.
done (err?) => void Yes Callback after all items complete, or on first error.

Returns: void. done is always called exactly once.

Throws: Does not throw. Errors are passed to done(err).


runSeriesCollect(items, iterator, done)

Like runSeries but collects results into an array.

Signature: runSeriesCollect(items: Array<*>, iterator: Function, done: Function): void

Name Type Required Default Description
items Array Yes Items to iterate.
iterator (item, next) => void Yes Called per item. Calls next(err, result?).
done (err?, results?) => void Yes Final callback. Receives (null, results) on success, (err) on error.

Returns: void.

Throws: Does not throw. Errors are passed to done(err).


runParallelCollect(items, iterator, done)

Runs an iterator over every item concurrently and collects results in input order. First error wins; remaining callbacks become no-ops.

Signature: runParallelCollect(items: Array<*>, iterator: Function, done: Function): void

Name Type Required Default Description
items Array Yes Items to iterate.
iterator (item, next) => void Yes Called per item. Calls next(err, result?).
done (err?, results?) => void Yes Final callback. Receives (null, results) on success, (err) on error.

Returns: void. done is called exactly once.

Throws: Does not throw. Errors are passed to done(err).


runParallelCollectLimited(items, iterator, options, done)

Bounded-concurrency variant of runParallelCollect. Throttles via an injected scheduler (preferred, dependency injection) or auto-creates one from options.concurrency.

Signature: runParallelCollectLimited(items: Array<*>, iterator: Function, options: Object, done: Function): void

Name Type Required Default Description
items Array Yes Items to iterate.
iterator (item, next) => void Yes Called per item. Calls next(err, result?).
options Object Yes Concurrency options.
options.scheduler { schedule(task) => void } No Inject a pre-built scheduler (e.g. createSharedScheduler from @polyweave/pool). Preferred over concurrency.
options.concurrency number No 0 If positive and scheduler is omitted, a createSharedScheduler({ concurrency }) is created automatically.
options.delayMs number No 0 Per-item delay before invocation.
done (err?, results?) => void Yes Final callback. Receives (null, results) on success, (err) on error.

Behavior:

  • If neither scheduler nor a positive concurrency is supplied, falls back to runParallelCollect (unbounded).
  • Results are stored at results[index], preserving input order.

Returns: void. done is called exactly once.

Throws: Does not throw. Errors are passed to done(err).


invokeMaybeSync(fn, args, done)

Invokes a function that may be sync (returns a value) or async (uses a callback). Resolves via callback.

Signature: invokeMaybeSync(fn: Function, args: Array<*>, done: Function): void

Name Type Required Default Description
fn Function Yes A function fn(...args, callback).
args Array Yes Arguments to pass (the callback is appended).
done (err?, result?) => void Yes Receives result from sync return or async callback.

Behavior: Calls fn(...args, finish). If fn returns a non-undefined value before finish is called, that value is used as the result. Otherwise waits for finish(err, result). Only the first resolution is used (subsequent calls to finish are no-ops).

Throws: Does not throw.


invokeMaybeFrame(handler, payload, done, options)

Invokes a function or frame handler (duplex), routing META frames and resolving TEXT, END, TOOL_RESULT, or META response streams.

Signature: invokeMaybeFrame(handler: Function|IFrameDuplex, payload: *, done: Function, options?: Object): void

Name Type Required Default Description
handler Function | IFrameDuplex Yes A callback-style function or a frame handler duplex (with sink + source).
payload * Yes Payload passed to the handler. For frame handlers, becomes content.payload in the META frame.
done (err?, result?) => void Yes Receives the resolved result.
options Object No {} Options
options.action string No 'self-improving' The META frame action name.
options.resolveFromFrame string No undefined If 'tool_result', resolves with the content of the matching TOOL_RESULT frame.
options.toolName string No undefined When resolveFromFrame: 'tool_result', filters for a specific tool name.

Behavior:

  • If handler is a function: calls it via invokeMaybeSync(handler, [payload], done).
  • If handler is a frame handler duplex: writes a META frame, then listens for completion on the source. Resolution path:
    • TEXT frames accumulate into a buffer.
    • END frame: resolves with accumulated text buffer, or if resolveFromFrame: 'tool_result', with { text, toolResults }.
    • TOOL_RESULT frame: appends to toolResults array. If resolveFromFrame: 'tool_result', resolves immediately with { text, toolResults }.
    • META frame: resolves with the meta result. If resolveFromFrame: 'tool_result', returns { result, toolResults }.
    • ERROR frame: rejects via done(error).

Throws: If handler is neither a function nor a frame handler, calls done(new Error('Handler must be a function or frame handler')).


sampleDataset(dataset, size, rng)

Fisher-Yates shuffle reservoir sampling.

Signature: sampleDataset(dataset: Array<*>, size: number, rng?: Function): Array<*>

Name Type Required Default Description
dataset Array Yes Source array to sample from.
size number Yes Number of items to return.
rng () => number No Math.random Random number generator returning [0, 1).

Returns: Array<*> — Sampled array.

Behavior:

  • If dataset.length === 0: returns [].
  • If size === 0 (exactly zero): returns [].
  • If size >= dataset.length: returns a shallow copy (dataset.slice()).
  • Otherwise: Fisher-Yates shuffles a copy, returns the first size items.

Throws: Does not throw.


normalizeExperienceLibrary(library)

Shallow-copies experience entries from an array.

Signature: normalizeExperienceLibrary(library: Array<*>|*): Array<Object>

Name Type Required Default Description
library Array | * Yes Anything. If not an array, treated as [].

Returns: Array<Object> — Each element is a shallow spread copy ({ ...entry }).

Throws: Does not throw.


applyExperienceOperations(library, operations)

Applies ADD, UPDATE, DELETE operations to an experience library.

Signature: applyExperienceOperations(library: Array<*>|*, operations: Array<Object>|*): Array<Object>

Name Type Required Default Description
library Array | * Yes Source library. Falsey values are normalized to [].
operations Array | * Yes Operations list. Falsey values are treated as [].

Operation shapes:

Operation Fields Effect
ADD { operation: 'ADD', entry: { id?, text, metadata? } } Adds entry if text exists and id doesn't already exist. Auto-generates id if missing.
UPDATE { operation: 'UPDATE', id, content } Updates entry text by id. If id not found, creates new entry with the content as text.
DELETE { operation: 'DELETE', id } Removes entry by id.

Also recognizes op as alias for operation (checked via op.operation || op.op).

Returns: Array<Object> — The modified library (new array).

Throws: Does not throw. Invalid operations (null, undefined, missing text for ADD) are silently skipped.


computeParetoFront(candidates, scoreKey)

Returns the Pareto-optimal front from a candidate array using multi-objective dominance.

Signature: computeParetoFront(candidates: Array<Object>, scoreKey?: string): Array<Object>

Name Type Required Default Description
candidates Array<Object> Yes Array of candidate objects.
scoreKey string No 'paretoScores' Property name for the numeric score vector on each candidate.

Returns: Array<Object> — Candidates that are not dominated by any other candidate.

Dominance rule: A dominates B iff all scores in A >= B and at least one is strictly >. Candidates missing the score key are excluded from the front (not compared as dominators, but skipped when checking if they are dominated).

Complexity: O(n^2) comparisons.

Throws: Does not throw.


dominates(scoresA, scoresB)

Pareto dominance test between two score vectors.

Signature: dominates(scoresA: number[], scoresB: number[]): boolean

Name Type Required Default Description
scoresA number[] Yes First score vector.
scoresB number[] Yes Second score vector. Must be same length as scoresA.

Returns: booleantrue iff scoresA[i] >= scoresB[i] for every objective and scoresA[i] > scoresB[i] for at least one objective.

Throws: Does not throw. Inputs must be non-empty same-length arrays of numbers (callers validate upstream).


From stateUtils.js

readStateValue(store, sectionId, pathValue, fallbackPath)

Reads a nested value from a Polyweave document store section.

Signature: readStateValue(store: Object, sectionId: string, pathValue: string|string[], fallbackPath?: string|string[]): *

Name Type Required Default Description
store Object Yes Must have a getSection(sectionId) method.
sectionId string Yes Section identifier.
pathValue string | string[] Yes Dot-separated string or array path. If falsely/empty and fallbackPath is provided, uses fallback.
fallbackPath string | string[] No 'data' Fallback path if pathValue is empty/falsy.

Returns: * — The value at the path, or undefined if store is invalid, getSection returns null/missing, section.data is missing, or any intermediate path segment is missing.

Path normalization: String paths are split on '.'. Array paths are used directly. Fallback defaults to ['data'].

Throws: Does not throw.


writeStateValue(store, sectionId, pathValue, fallbackPath, value)

Writes a nested value to a Polyweave document store via DOC_UPDATE frame.

Signature: writeStateValue(store: Object, sectionId: string, pathValue: string|string[], fallbackPath?: string|string[], value: *): void

Name Type Required Default Description
store Object Yes Must have a sink.write(frame) method.
sectionId string Yes Section identifier.
pathValue string | string[] Yes Dot-separated string or array path.
fallbackPath string | string[] No 'data' Fallback path.
value * Yes Value to write.

Returns: void. If store is invalid or store.sink.write is not a function, silently returns.

Frame emitted: { type: 'DOC_UPDATE', content: { sectionId, op: 'set', path, value } }.

Throws: Does not throw.

Tests

npm test

Covers: Pareto-front computation with dominance edge cases, reservoir sampling (size=0, size>=len, empty input), experience library operations (ADD, UPDATE, DELETE, idempotent ADD), sequential iteration (runSeries, runSeriesCollect), invokeMaybeSync dual-mode, invokeMaybeFrame with function and duplex handlers, state value read with dot-path and array-path, and state value write via DOC_UPDATE frame.

Requirements

  • Node.js 22+
  • ESM only
  • Peer: @polyweave/core
  • Dependencies: @polyweave/pool (for runParallelCollectLimited auto-scheduler)

Dependencies

Dependencies

ID Version
@polyweave/pool *

Development Dependencies

ID Version
@rigor/core *

Peer Dependencies

ID Version
@polyweave/core *

Keywords

polyweave loop optimization utilities
Details
npm
2026-07-17 22:23:31 +00:00
59
John Dvorak
SEE LICENSE IN LICENSE
latest
7.7 KiB
Assets (1)
Versions (4) View all
1.1.0 2026-07-17
1.0.5 2026-07-11
1.0.2 2026-07-10
1.0.1 2026-07-10