import { F as FiberEngine, W as WasmEngineRuntime, A as Async, i as RuntimeFiber, j as FiberEngineStats, k as Fiber, l as FiberId, m as FiberStatus, E as Exit, d as RuntimeEvent, n as WasmBridge, o as OpcodeProgram, p as FiberId$1, q as EngineEvent, r as RefId, N as NodeId, s as OpcodeNode, t as EngineStats, O as Option, b as Runtime, S as Scope, c as RingBufferOptions } from './effect-DbEMiMvv.js'; export { u as AbortablePromiseFinish, v as AbortablePromiseLabelStats, w as AbortablePromiseOptions, x as AbortablePromiseOutcome, y as AbortablePromiseStats, z as AbortablePromiseTimerWheel, C as AsyncRegisterRef, D as AsyncWithPromise, G as BoundedRingBuffer, H as CancelToken, I as Canceler, K as Cause, L as CausePrettyOptions, M as ContextNode, P as CustomHostAction, Q as DbHostAction, U as DecodeRef, V as DefaultHostExecutor, X as EngineKind, Y as EngineSelection, _ as EngineSelectionMode, $ as FiberContext, a0 as FiberEngineKind, a1 as FiberInfo, a2 as FiberReadyQueue, a3 as FiberReadyQueueOptions, a4 as FiberReadyQueueStats, a5 as FiberRunState, a6 as FlatMapRef, a7 as FoldFailureRef, a8 as FoldSuccessRef, a9 as HostAction, aa as HostActionKind, ab as HostActionResult, ac as HostExecutionContext, ad as HostExecutor, ae as HostRegistry, af as HostRegistryStats, ag as HttpHostAction, ah as Interrupted, ai as InterruptibilityMode, J as JSONValue, aj as Joiner, ak as LaneStatsData, al as None, am as NoopHooks, an as ProgramBuilder, ao as ProgramPatch, ap as PushStatus, aq as QueueHostAction, ar as ReadyQueueScheduleKind, as as RestoreInterruptibility, at as RingBuffer, au as RingBufferEngine, av as RingBufferStatsData, aw as RuntimeCapabilities, e as RuntimeEmitContext, ax as RuntimeEngineMode, f as RuntimeEventRecord, R as RuntimeHooks, a as RuntimeOptions, h as RuntimeRegistry, g as RuntimeSpanLink, ay as ScheduleResult, az as Scheduler, aA as SchedulerEngine, aB as SchedulerLaneMode, aC as SchedulerOptions, aD as SchedulerStats, aE as SchedulerStatsData, aF as ScopeId, aG as ScopeInfo, aH as Some, aI as SyncRef, aJ as Task, T as TraceContext, aK as WasmFiberEngine, aL as WasmFiberEngineOptions, Z as ZIO, aM as abortablePromiseStats, aN as acquireRelease, aO as async, aP as asyncCatchAll, aQ as asyncFail, aR as asyncFlatMap, aS as asyncFold, aT as asyncInterruptible, aU as asyncMap, aV as asyncMapError, aW as asyncSucceed, aX as asyncSync, aY as asyncTotal, aZ as catchAll, a_ as ctxExtend, a$ as ctxToObject, b0 as emptyContext, b1 as end, b2 as engineStats, b3 as fail, b4 as flatMap, b5 as fork, b6 as formatCause, b7 as fromPromiseAbortable, b8 as getBenchmarkBudget, b9 as getCurrentFiber, ba as globalScheduler, bb as inferCallerLaneFromStack, bc as interruptible, bd as isCause, be as laneTag, bf as linkAbortController, bg as makeBoundedRingBuffer, bh as makeCancelToken, bi as makeFiberReadyQueue, bj as makeRuntimeEventRecord, bk as map, bl as mapAsync, bm as mapError, bn as mapTryAsync, bo as none, bp as orElseOptional, bq as prettyCause, br as recordAbortablePromiseFinish, bs as recordAbortablePromiseStart, bt as resetAbortablePromiseStats, bu as runtimeCapabilities, bv as runtimeEventRecordContext, bw as runtimeForCaller, bx as sanitizeLaneKey, by as selectedEngineStats, bz as setAbortablePromisePerLabelTracking, bA as setBenchmarkBudget, bB as setCurrentFiber, bC as some, bD as succeed, bE as sync, bF as toPromise, bG as toPromiseByCaller, bH as uninterruptible, bI as uninterruptibleMask, bJ as unit, bK as unsafeGetCurrentRuntime, bL as unsafeRunAsync, bM as unsafeRunFoldWithEnv, bN as withAsyncPromise, bO as withCurrentFiber, bP as withScope, bQ as withScopeAsync } from './effect-DbEMiMvv.js'; export { ConfigLayerOptions, ConfigLayerSource, EffectEnvironment, EffectFailure, EffectSuccess, FiberRef, LogLevel, MakeRuntimeOptions, Ref, RetryPolicy, RetryState, RuntimeLayerEnv, RuntimeLayerOptions, RuntimeRecorder, RuntimeRecorderExplainOptions, RuntimeRecorderOptions, RuntimeRecorderStats, RuntimeService, Semaphore, SemaphoreStats, ShutdownConfig, ShutdownStats, SupervisedChildSpec, SupervisedChildStatus, SupervisedFiber, Supervisor, SupervisorConfig, SupervisorEscalation, SupervisorEvent, SupervisorRestartContext, SupervisorRestartMode, SupervisorRestartPolicy, SupervisorStrategy, TaggedError, TestClock, TestClockTimerSnapshot, TestRuntime, TestRuntimeOptions, TestScheduledTask, TestScheduler, TimeoutError, WorkerPool, WorkerPoolConfig, WorkerPoolError, WorkerPoolStats, assertCompletesWithin, assertFails, assertFailsWith, assertSucceeds, catchTag, catchTags, consoleJsonLogger, defineConfigLayer, delayedEffect, derivedRef, dumpAllFibers, fiberRefSnapshot, flakyEffect, getFiberRef, gracefulShutdown, joinSupervised, locallyFiberRef, locallyFiberRefWith, makeConfigLayer, makeFiberRef, makeRef, makeRuntime, makeRuntimeLayer, makeRuntimeRecorder, makeSemaphore, makeSupervisor, makeTestRuntime, makeWorkerPool, mapErrorTyped, modifyFiberRef, neverEffect, orElse, registerShutdownHooks, retry, retryN, retryWithBackoff, runEffect, runExit, runPromise, setFiberRef, sleep, supervise, tagError, timeout, unsafeGetFiberRef, unsafeSetFiberRef, updateFiberRef } from './core/index.js'; export { M as ManagedResource, R as Resource, R as ResourceDescriptor, S as Span, a as SpanContext, b as SpanEvent, c as SpanStatus, T as Tracer, d as TracerConfig, e as bracket, f as ensuring, m as makeResource, g as makeTracer, h as managed, i as managedAll, r as resource, j as resourceAll, k as resourceFromManaged, l as resourceSucceed, u as useManaged, n as useResource } from './tracing-DIAUf2x8.js'; export { B as BuiltLayer, b as CircuitBreaker, c as CircuitBreakerConfig, d as CircuitBreakerError, e as CircuitBreakerState, C as CircuitBreakerStats, L as Layer, a as LayerContext, f as LayerErrorOf, g as LayerInputOf, h as LayerOutputOf, i as LayerScope, M as MissingLayerServiceError, R as RuntimeClock, j as RuntimeClockEnv, k as RuntimeTimerId, l as Schedule, m as ScheduleDecision, n as ScheduleDriver, o as ScheduleDriverDecision, p as ScheduleDriverOptions, q as ScheduleDriverSnapshot, r as ScheduleObserver, s as ScheduleObserverEvent, t as ScheduleStepContext, S as ServiceTag, u as ServiceTagMap, v as ServicesOf, T as TestLayerProvider, w as andThenSchedule, x as buildLayer, y as composeAll, z as composeLayer, A as contramapSchedule, D as defineLayer, E as defineService, F as elapsed, G as exponential, H as fibonacci, I as fixed, J as forever, K as formatLayerError, N as getService, O as getServices, P as intersect, Q as jitter, U as jittered, V as jitteredSchedule, W as layer, X as layerEffect, Y as layerFail, Z as layerFrom, _ as layerFromContext, $ as layerSucceed, a0 as layerValue, a1 as linear, a2 as liveClock, a3 as makeCircuitBreaker, a4 as makeLayerScope, a5 as makeScheduleDriver, a6 as makeServiceTag, a7 as makeTestLayer, a8 as makeTestLayers, a9 as mapLayer, aa as mapSchedule, ab as maxDelay, ac as maxElapsed, ad as mergeAll, ae as mergeLayer, af as namedSchedule, ag as never, ah as once, ai as pollWithSchedule, aj as provide, ak as provideContext, al as provideLayer, am as provideLayerContext, an as recurs, ao as repeatWithSchedule, ap as repeatWithScheduleAlias, aq as retryWithSchedule, ar as retryWithScheduleAlias, as as runSchedule, at as runtimeClockFromEnv, au as scheduleDriver, av as serviceTag, aw as spaced, ax as takeSchedule, ay as tapDecision, az as union, aA as untilInput, aB as untilOutput, aC as upTo, aD as useService, aE as useServices, aF as whileInput, aG as whileOutput, aH as windowed } from './layer-CsNGeVee.js'; export { B as BrassEnv, C as Counter, E as EventBus, f as EventBusOptions, g as EventHandler, G as Gauge, H as Histogram, h as HistogramBuckets, I as InMemoryTracer, e as InMemoryTracerOptions, i as InMemoryTracerStats, a as MetricExemplar, b as MetricSnapshot, j as MetricType, k as MetricValue, M as MetricsRegistry, R as RuntimeSpan, l as RuntimeSpanEvent, d as RuntimeTraceIdGenerator, m as defaultTracer, n as makeMetrics, r as runtimeHooksToEventHandler } from './tracer-CS3yOXZx.js'; import { Z as ZStream } from './stream-B8c_UKZq.js'; export { C as Concat, E as Emit, a as Empty, F as Flatten, b as FromArray, c as FromPull, M as Managed, d as Merge, N as Normalize, R as ReadableStreamOptions, S as Scoped, e as assertNever, f as collectStream, g as concatStream, h as emitStream, i as emptyStream, j as flattenStream, k as foreachStream, l as fromArray, m as fromPull, n as managedStream, o as mapStream, p as merge, q as mergeStream, r as rangeStream, s as streamFromReadableStream, u as uncons, t as unwrapScoped, w as widenOpt, z as zip } from './stream-B8c_UKZq.js'; import './schema/index.js'; declare class JsFiberEngine implements FiberEngine { private readonly runtime; readonly kind: "ts"; private startedFibers; constructor(runtime: WasmEngineRuntime & any); fork(effect: Async, scopeId?: number): RuntimeFiber; stats(): FiberEngineStats; } type InternalFiberStatus = "queued" | "running" | "suspended" | "done" | "failed" | "interrupted"; declare class EngineFiberHandle implements Fiber { private readonly onScheduledStep; private readonly onInterrupt; private readonly onJoiner?; private readonly onQueued?; private readonly onScheduleDropped?; private readonly onScheduleRequest?; readonly id: FiberId; readonly runtime: WasmEngineRuntime & any; fiberContext: any; name?: string; scopeId?: number; parentFiberId?: number; lane?: string; private result; private readonly joiners; private readonly finalizers; private finalizersDrained; private internalStatus; private queued; constructor(id: FiberId, runtime: WasmEngineRuntime & any, onScheduledStep: (fiberId: FiberId) => void, onInterrupt: (fiberId: FiberId, reason: unknown) => void, onJoiner?: ((fiberId: FiberId) => void) | undefined, onQueued?: ((fiberId: FiberId) => void) | undefined, onScheduleDropped?: ((fiberId: FiberId, label: string) => void) | undefined, onScheduleRequest?: ((fiberId: FiberId, label: string) => "accepted" | "dropped") | undefined); status(): FiberStatus; engineStatus(): InternalFiberStatus; setEngineStatus(status: InternalFiberStatus): void; markDequeued(): void; join(cb: (exit: Exit) => void): void; interrupt(): void; addFinalizer(f: (exit: Exit) => void): void; schedule(tag?: string): void; private scheduleWithRuntime; emit(ev: RuntimeEvent): void; succeed(value: A): void; fail(error: E): void; die(defect: unknown): void; interrupted(): void; complete(exit: Exit): void; private runFinalizersOnce; } declare class WasmPackFiberBridge implements WasmBridge { readonly kind: "wasm"; readonly supportsBinary: boolean; readonly supportsZeroCopy: boolean; readonly supportsNoJsonMetrics: boolean; private readonly vm; private jsonEventCalls; private binaryEventCalls; private zeroCopyEventCalls; private eventsReceived; private maxEventsPerCall; private jsonPrograms; private binaryPrograms; private zeroCopyPrograms; private jsonPatches; private binaryPatches; private zeroCopyPatches; constructor(modulePath?: string); createFiber(program: OpcodeProgram): FiberId$1; poll(fiberId: FiberId$1): EngineEvent; driveBatch(fiberId: FiberId$1, budget: number): readonly EngineEvent[]; provideValue(fiberId: FiberId$1, valueRef: RefId): EngineEvent; provideValueBatch(fiberId: FiberId$1, valueRef: RefId, budget: number): readonly EngineEvent[]; provideError(fiberId: FiberId$1, errorRef: RefId): EngineEvent; provideErrorBatch(fiberId: FiberId$1, errorRef: RefId, budget: number): readonly EngineEvent[]; provideEffect(fiberId: FiberId$1, root: NodeId, nodes: OpcodeNode[]): EngineEvent; provideEffectBatch(fiberId: FiberId$1, root: NodeId, nodes: OpcodeNode[], budget: number): readonly EngineEvent[]; interrupt(fiberId: FiberId$1, reasonRef: RefId): EngineEvent; interruptBatch(fiberId: FiberId$1, reasonRef: RefId, budget: number): readonly EngineEvent[]; dropFiber(fiberId: FiberId$1): void; stats(): unknown; private assertStrictWasmHotPath; private decodeZeroCopy; private memory; private readU32; private readF64; private writeWords; private readMetricsSnapshot; } type WasmFiberRegistryStats = { readonly live: number; readonly queued: number; readonly running: number; readonly suspended: number; readonly done: number; readonly failed: number; readonly interrupted: number; readonly wakeQueueLen: number; readonly registered: number; readonly completed: number; readonly wakeups: number; readonly duplicateWakeups: number; readonly joins: number; }; type FiberRegistryStatus = "queued" | "running" | "suspended" | "done" | "failed" | "interrupted"; declare class WasmFiberRegistryBridge { private readonly registry; constructor(); registerFiber(fiberId: FiberId$1, parentId?: number, scopeId?: number): void; markQueued(fiberId: FiberId$1): void; markRunning(fiberId: FiberId$1): void; markSuspended(fiberId: FiberId$1): void; markDone(fiberId: FiberId$1, status: Exclude): number; dropFiber(fiberId: FiberId$1): void; addJoiner(fiberId: FiberId$1): void; wake(fiberId: FiberId$1): boolean; drainWakeup(): FiberId$1 | undefined; drainWakeups(): FiberId$1[]; wakeQueueLength(): number; stateOf(fiberId: FiberId$1): FiberRegistryStatus | "missing"; stats(): WasmFiberRegistryStats; } declare const ABI_VERSION = 1; declare const EVENT_WORDS = 5; declare const NONE_U32 = 4294967295; declare const enum OpcodeTagCode { Succeed = 0, Fail = 1, Sync = 2, Async = 3, FlatMap = 4, Fold = 5, Fork = 6, HostAction = 7 } declare const enum EventKindCode { Continue = 0, Done = 1, Failed = 2, Interrupted = 3, InvokeSync = 4, InvokeAsync = 5, InvokeFlatMap = 6, InvokeFoldFailure = 7, InvokeFoldSuccess = 8, InvokeFork = 9, InvokeHostAction = 10 } declare function encodeOpcodeProgram(program: OpcodeProgram): Uint32Array; declare function encodeOpcodeNodes(nodes: readonly OpcodeNode[]): Uint32Array; declare function decodeEvent(words: ArrayLike, offset?: number): EngineEvent; declare function decodeEventBatch(words: ArrayLike | null | undefined): EngineEvent[]; type StreamChunkEngine = "ts" | "wasm"; type StreamChunkOptions = { /** * ts: always use the TypeScript array chunker. * wasm: require BrassWasmChunkBuffer from wasm/pkg. * * Strict mode never falls back between engines. */ engine?: StreamChunkEngine; }; type StreamChunkStats = { len: number; maxChunkSize: number; emittedChunks: number; emittedItems: number; flushes: number; }; type Chunker = { readonly length: number; readonly maxChunkSize: number; push(value: A): boolean; isFull(): boolean; isEmpty(): boolean; takeChunk(): readonly A[]; clear(): void; stats(): EngineStats; }; declare function makeStreamChunker(chunkSize: number, options?: StreamChunkOptions): Chunker; /** * Re-chunk a stream so downstream operators receive arrays instead of single * items. This is the intended WASM boundary: pay the JS↔WASM crossing while * assembling chunks, then process bigger batches downstream. */ declare function chunks(input: ZStream, chunkSize: number, options?: StreamChunkOptions): ZStream; declare function mapChunks(input: ZStream, chunkSize: number, f: (chunk: readonly A[]) => readonly B[], options?: StreamChunkOptions): ZStream; declare function mapChunksEffect(chunkSize: number, f: (chunk: readonly A[]) => Async, options?: StreamChunkOptions): (input: ZStream) => ZStream; /** * ZPipeline-style transformer. * * A pipeline that consumes `In` and produces `Out`, potentially requiring `Rp` and failing with `Ep`. * When applied to a stream `ZStream`, the result is `ZStream`. */ type ZPipeline = (input: ZStream) => ZStream; /** Apply a pipeline to a stream (alias of `pipeline(stream)`). * * OPTIMIZATION: When the pipeline is a single pure operator (has PURE_PIPELINE_TAG) * and the stream can be drained synchronously, uses the fast fused path. * The FusedPipelineRepr is cached on the pipeline to avoid recalculation. */ declare function via(stream: ZStream, pipeline: ZPipeline): ZStream; /** Compose pipelines left-to-right (p1 >>> p2). */ declare function andThen(p1: ZPipeline, p2: ZPipeline): ZPipeline; /** Compose pipelines right-to-left (p2 <<< p1). */ declare function compose(p2: ZPipeline, p1: ZPipeline): ZPipeline; /** Identity pipeline. */ declare function identity(): ZPipeline; /** Map elements. */ declare function mapP(f: (a: A) => B): ZPipeline; /** Filter elements, preserving end/error. */ declare function filterP(pred: (a: A) => boolean): ZPipeline; /** * Filter-map (aka collectSome). * If `f(a)` returns None, the element is dropped. */ declare function filterMapP(f: (a: A) => Option): ZPipeline; /** Take at most N elements. */ declare function takeP(n: number): ZPipeline; /** Drop the first N elements. */ declare function dropP(n: number): ZPipeline; declare function mapEffectP(f: (a: A) => Async): ZPipeline; /** Tap each element with an effect, preserving the element. */ declare function tapEffectP(f: (a: A) => Async): ZPipeline; /** Re-chunk a stream into arrays of up to `chunkSize` elements. */ declare function chunksP(chunkSize: number, options?: StreamChunkOptions): ZPipeline; /** Apply one effect per chunk and flatten the returned chunk back to elements. */ declare function mapChunksEffectP(chunkSize: number, f: (chunk: readonly A[]) => Async, options?: StreamChunkOptions): ZPipeline; /** Buffer upstream using your existing queue-based buffer implementation. */ declare function bufferP(capacity: number, strategy?: "backpressure" | "dropping" | "sliding"): ZPipeline; /** * Group elements into arrays of size `n` (last chunk may be smaller). * Example: [1,2,3,4,5].grouped(2) => [1,2],[3,4],[5] */ declare function groupedP(n: number): ZPipeline; type Stream = ZStream & { readonly pipe: (pipeline: ZPipeline) => Stream; readonly map: (f: (value: A) => B) => Stream; readonly filter: (predicate: (value: A) => boolean) => Stream; readonly collect: (runtime: Runtime) => Promise; }; declare const Stream: Readonly<{ from: (values: readonly A[]) => Stream; empty: () => Stream; range: (start: number, end: number) => Stream; wrap: typeof asStream; }>; declare const Pipeline: Readonly<{ map: typeof mapP; filter: typeof filterP; }>; declare function asStream(stream: ZStream): Stream; declare function buffer(stream: ZStream<{} & R, E, A>, capacity: number, strategy?: "backpressure" | "dropping" | "sliding"): ZStream<{} & R, E, A>; /** * race(A, B): * - corre A y B en paralelo * - el primero que termine gana * - el otro es cancelado */ declare function race(left: Async, right: Async, parentScope: Scope): Async; /** * zipPar(A, B): * - corre A y B en paralelo * - si ambas terminan bien → éxito con (A, B) * - si una falla → cancelar todo y devolver fallo */ declare function zipPar(left: Async, right: Async, parentScope: Scope): Async; /** * collectAllPar: * - corre todos en paralelo * - si uno falla → cancela todos * - si todos terminan bien → devuelve array de resultados */ declare function collectAllPar(effects: ReadonlyArray>, parentScope: Scope): Async; declare function raceWith(left: Async, right: Async, parentScope: Scope, onLeft: (exit: Exit, rightFiber: Fiber, scope: Scope) => Async, onRight: (exit: Exit, leftFiber: Fiber, scope: Scope) => Async): Async; type Strategy = "backpressure" | "dropping" | "sliding"; type QueueClosed = { _tag: "QueueClosed"; }; type Queue = { offer: (a: A) => Async; take: () => Async; /** Offer multiple values in a single effect. Returns array of success flags. */ offerBatch: (values: readonly A[]) => Async; /** Take up to N values in a single effect. Returns available values (may be fewer than N). */ takeBatch: (n: number) => Async; size: () => number; shutdown: () => void; }; type QueueOptions = RingBufferOptions; declare function bounded(capacity: number, strategy?: Strategy, options?: QueueOptions): Async>; type HubStrategy = "BackPressure" | "Dropping" | "Sliding"; type HubClosed = { _tag: "HubClosed"; }; type Subscription = Queue & { unsubscribe: () => void; }; type Hub = { publish: (a: A) => Async; publishAll: (as: Iterable) => Async; subscribe: () => Async>; shutdown: () => Async; }; declare function makeHub(capacity: number, strategy?: HubStrategy): Hub; declare const broadcast: typeof makeHub; declare function broadcastToHub(stream: ZStream, hub: Hub): Async; declare function fromHub(hub: Hub): ZStream; /** Serialized representation of a single step in a fused pipeline */ type SerializedStep = { kind: "map"; fnSource: string; } | { kind: "filter"; predSource: string; } | { kind: "take"; n: number; } | { kind: "drop"; n: number; }; /** JSON-safe serialized representation of a fused pipeline */ type SerializedFusedPipeline = { readonly version: 1; readonly steps: SerializedStep[]; }; /** Enable or disable fusion globally. When disabled, andThen will not attempt fusion. */ declare function setFusionEnabled(enabled: boolean): void; /** Check if fusion is globally enabled. */ declare function isFusionEnabled(): boolean; /** Set verbose mode globally. When enabled, fusion decisions are logged to console. */ declare function setFusionVerbose(verbose: boolean): void; /** Check if verbose mode is globally enabled. */ declare function isFusionVerbose(): boolean; /** * Get stats from a fused pipeline, or null if the pipeline is not fused. * Works with pipelines that have been fused via `andThen` (have `_fusedSteps`) * or with `FusedPipelineRepr` objects directly. */ declare function getStats(pipeline: ZPipeline): FusedPipelineStats | null; declare const PURE_PIPELINE_TAG: unique symbol; /** Result of a fused step for an element */ type FuseResult = { readonly tag: "emit"; readonly value: A; } | { readonly tag: "skip"; } | { readonly tag: "halt"; }; /** Metadata of a step in the original pipeline */ type FusedStep = { readonly kind: "map"; } | { readonly kind: "filter"; } | { readonly kind: "take"; readonly n: number; } | { readonly kind: "drop"; readonly n: number; }; /** Stats of a fused pipeline */ type FusedPipelineStats = { readonly fusedSteps: number; readonly steps: readonly FusedStep[]; readonly hasTake: boolean; readonly hasDrop: boolean; }; /** Internal representation of a fused pipeline */ type FusedPipelineRepr = { readonly _tag: "FusedPipeline"; readonly step: (a: In, state: FuseState) => FuseResult; readonly initState: () => FuseState; readonly stats: FusedPipelineStats; }; /** Mutable state during execution (per-step counters for take/drop) */ type FuseState = { /** Per-step counters: each take/drop step gets its own independent counter */ counters: number[]; }; /** Fusion options */ type FusionOptions = { readonly enabled?: boolean; readonly verbose?: boolean; }; type PurePipelineMetadata = { readonly kind: "map" | "filter" | "take" | "drop"; readonly fn?: (a: In) => Out; readonly pred?: (a: In) => boolean; readonly n?: number; }; type PurePipelineTag = { readonly [PURE_PIPELINE_TAG]: PurePipelineMetadata; }; /** Creates a FuseState with per-step counters (one for each take/drop step) */ declare function initState(counterCount: number): () => FuseState; /** * Detects if a pipeline is fusable and returns the fused representation. * Returns null if the pipeline is not fusable (not pure or fusion disabled). */ declare function fuse(pipeline: ZPipeline, options?: FusionOptions): FusedPipelineRepr | null; /** * Applies a fused pipeline to an array of inputs synchronously. * This is the fastest execution path — a pure `for` loop with no effects, * no fibers, no scheduling overhead. O(n) with minimal constant factor. */ declare function runFusedArray(input: readonly In[], fused: FusedPipelineRepr): Out[]; /** * Applies a fused pipeline to a stream using a single pull loop. * No intermediate fibers are created between fused operators. * * OPTIMIZATION: If the input stream can be drained synchronously (e.g., fromArray), * the entire pipeline is executed as a pure synchronous loop — no effects at all. */ declare function applyFused(stream: ZStream, fused: FusedPipelineRepr): ZStream; /** * Serialize a fused pipeline to a JSON-safe representation. * Returns null if the pipeline is not fused (no _fusedSteps metadata). * * For map/filter steps, the function/predicate source is captured via `.toString()`. * Note: toString() has limitations with closures and minification, but is sufficient * for debugging/observability use cases. */ declare function serializeFusedPipeline(pipeline: ZPipeline): SerializedFusedPipeline | null; /** * Deserialize a serialized pipeline back to a functional pipeline. * Returns null if deserialization fails (e.g., invalid version, malformed fnSource). * * Reconstructs functions using `new Function(...)` with appropriate wrapping. * The resulting pipeline is functionally equivalent to the original for pure * (non-closure) functions. * * WARNING: Uses `new Function()` which has security implications similar to `eval`. * Only deserialize trusted serialized pipelines. */ declare function deserializeFusedPipeline(serialized: SerializedFusedPipeline): ZPipeline | null; /** * Throttles a stream to emit at most one element per `intervalMs`. * Elements arriving during the cooldown period are dropped. * * ```ts * const throttled = throttle(clickStream, 1000); // max 1 click per second * ``` */ declare function throttle(stream: ZStream, intervalMs: number): ZStream; /** * Debounces a stream: only emits an element after `delayMs` of silence. * If a new element arrives before the delay expires, the previous is dropped. * * Note: This is a simplified debounce that works by buffering the last element * and emitting it after a delay. For real-time use cases, consider using * the Hub-based approach with timers. * * ```ts * const debounced = debounce(inputStream, 300); // wait 300ms of silence * ``` */ declare function debounce(stream: ZStream, delayMs: number): ZStream; /** * Zips two streams together, pairing elements by position. * The resulting stream ends when either input stream ends. * * ```ts * const zipped = zip(numbersStream, lettersStream); * // [1, "a"], [2, "b"], [3, "c"], ... * ``` */ declare function zip(left: ZStream, right: ZStream): ZStream; /** * Zips two streams with a custom combiner function. * * ```ts * const summed = zipWith(xs, ys, (x, y) => x + y); * ``` */ declare function zipWith(left: ZStream, right: ZStream, f: (a: A, b: B) => C): ZStream; /** * Produces a stream of accumulated values using a reducer function. * Emits the initial value first, then each accumulated result. * * ```ts * const running = scan(numbersStream, 0, (acc, n) => acc + n); * // 0, 1, 3, 6, 10, ... (running sum) * ``` */ declare function scan(stream: ZStream, initial: B, f: (acc: B, a: A) => B): ZStream; /** * Interleaves two streams, alternating elements from each. * When one stream ends, remaining elements from the other are emitted. * * ```ts * const mixed = interleave(evens, odds); * // 0, 1, 2, 3, 4, 5, ... * ``` */ declare function interleave(left: ZStream, right: ZStream): ZStream; /** * Takes the first N elements from a stream. */ declare function take(stream: ZStream, n: number): ZStream; /** * Drops the first N elements from a stream. */ declare function drop(stream: ZStream, n: number): ZStream; export { ABI_VERSION, Async, EVENT_WORDS, EngineEvent, EngineFiberHandle, EngineStats, EventKindCode, Exit, Fiber, FiberEngine, FiberEngineStats, FiberId, FiberStatus, type FuseResult, type FuseState, type FusedPipelineRepr, type FusedPipelineStats, type FusedStep, type FusionOptions, type Hub, type HubClosed, type HubStrategy, type InternalFiberStatus, JsFiberEngine, NONE_U32, NodeId, OpcodeNode, OpcodeProgram, OpcodeTagCode, Option, PURE_PIPELINE_TAG, Pipeline, type PurePipelineMetadata, type PurePipelineTag, type Queue, type QueueClosed, type QueueOptions, RefId, RingBufferOptions, Runtime, RuntimeEvent, RuntimeFiber, Scope, type SerializedFusedPipeline, type SerializedStep, type Strategy, Stream, type StreamChunkEngine, type StreamChunkOptions, type StreamChunkStats, type Subscription, WasmBridge, WasmEngineRuntime, WasmFiberRegistryBridge, type WasmFiberRegistryStats, WasmPackFiberBridge, type ZPipeline, ZStream, andThen, applyFused, asStream, bounded, broadcast, broadcastToHub, buffer, bufferP, chunks, chunksP, collectAllPar, compose, debounce, decodeEvent, decodeEventBatch, deserializeFusedPipeline, dropP, drop as dropStream, encodeOpcodeNodes, encodeOpcodeProgram, filterMapP, filterP, fromHub, fuse, getStats, groupedP, identity, initState, interleave, isFusionEnabled, isFusionVerbose, makeHub, makeStreamChunker, mapChunks, mapChunksEffect, mapChunksEffectP, mapEffectP, mapP, race, raceWith, runFusedArray, scan, serializeFusedPipeline, setFusionEnabled, setFusionVerbose, takeP, take as takeStream, tapEffectP, throttle, via, zipPar, zip as zipStream, zipWith };