BackPressureStrategy<A>Represents the back-pressure strategy for bounded PubSub values.
When to use
Use to preserve every message for current subscribers when a bounded custom
PubSub should make publishers wait for capacity instead of dropping or
evicting messages.
Details
Publishers wait when the PubSub is at capacity, so all current subscribers
can receive every published message.
Gotchas
A slow subscriber can slow down publishers and other subscribers.
export class class BackPressureStrategy<in out A>class BackPressureStrategy {
publishers: MutableList.MutableList<readonly [A, Deferred.Deferred<boolean>, boolean]>;
shutdown: Effect.Effect<void, never, never>;
handleSurplus: (pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>, elements: Iterable<A>, isShutdown: MutableRef.MutableRef<boolean>) => Effect.Effect<boolean>;
onPubSubEmptySpaceUnsafe: (pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>) => void;
completePollersUnsafe: (pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>, subscription: PubSub.BackingSubscription<A>, pollers: MutableList.MutableList<Deferred.Deferred<A>>) => void;
completeSubscribersUnsafe: (pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>) => void;
offerUnsafe: (elements: Iterable<A>, deferred: Deferred.Deferred<boolean>) => void;
removeUnsafe: (deferred: Deferred.Deferred<boolean>) => void;
}
Represents the back-pressure strategy for bounded PubSub values.
When to use
Use to preserve every message for current subscribers when a bounded custom
PubSub should make publishers wait for capacity instead of dropping or
evicting messages.
Details
Publishers wait when the PubSub is at capacity, so all current subscribers
can receive every published message.
Gotchas
A slow subscriber can slow down publishers and other subscribers.
BackPressureStrategy<in out function (type parameter) A in BackPressureStrategy<in out A>A> implements PubSub.interface PubSub<in out A>.Strategy<in out A>Strategy interface defining how PubSub handles backpressure and message distribution.
Strategy<function (type parameter) A in BackPressureStrategy<in out A>A> {
BackPressureStrategy<in out A>.publishers: MutableList.MutableList<readonly [A, Deferred.Deferred<boolean>, boolean]>(property) BackPressureStrategy<in out A>.publishers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
publishers: import MutableListMutableList.type MutableList.MutableList = /*unresolved*/ anyMutableList<
readonly [function (type parameter) A in BackPressureStrategy<in out A>A, import DeferredDeferred.type Deferred.Deferred = /*unresolved*/ anyDeferred<boolean>, boolean]
> = import MutableListMutableList.make()
get BackPressureStrategy<in out A>.shutdown: Effect.Effect<void, never, never>(getter) BackPressureStrategy<in out A>.shutdown: {
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
Describes any finalization logic associated with this strategy.
shutdown(): import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<void> {
return import EffectEffect.withFiber((fiber: Fiber<unknown, unknown>(parameter) fiber: {
id: number;
currentOpCount: number;
getRef: <A>(ref: Context.Reference<A>) => A;
context: Context.Context<never>;
setContext: (context: Context.Context<never>) => void;
currentScheduler: Scheduler;
currentDispatcher: SchedulerDispatcher;
currentSpan: AnySpan | undefined;
currentLogLevel: LogLevel;
minimumLogLevel: LogLevel;
currentStackFrame: StackFrame | undefined;
maxOpsBeforeYield: number;
currentPreventYield: boolean;
addObserver: (cb: (exit: Exit<A, E>) => void) => () => void;
interruptUnsafe: (fiberId?: number | undefined, annotations?: Context.Context<never> | undefined) => void;
pollUnsafe: () => Exit<A, E> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
fiber) =>
import EffectEffect.forEach(
import MutableListMutableList.takeAll(this.BackPressureStrategy<in out A>.publishers: MutableList.MutableList<readonly [A, Deferred.Deferred<boolean>, boolean]>(property) BackPressureStrategy<in out A>.publishers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
publishers),
([_: any_, deferred: Deferred.Deferred<boolean, never>(parameter) deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred, last: anylast]) => last: anylast ? import DeferredDeferred.interruptWith(deferred: Deferred.Deferred<boolean, never>(parameter) deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred, fiber: Fiber<unknown, unknown>(parameter) fiber: {
id: number;
currentOpCount: number;
getRef: <A>(ref: Context.Reference<A>) => A;
context: Context.Context<never>;
setContext: (context: Context.Context<never>) => void;
currentScheduler: Scheduler;
currentDispatcher: SchedulerDispatcher;
currentSpan: AnySpan | undefined;
currentLogLevel: LogLevel;
minimumLogLevel: LogLevel;
currentStackFrame: StackFrame | undefined;
maxOpsBeforeYield: number;
currentPreventYield: boolean;
addObserver: (cb: (exit: Exit<A, E>) => void) => () => void;
interruptUnsafe: (fiberId?: number | undefined, annotations?: Context.Context<never> | undefined) => void;
pollUnsafe: () => Exit<A, E> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
fiber.id) : import EffectEffect.void,
{ concurrency: stringconcurrency: "unbounded", discard: booleandiscard: true }
)
)
}
function BackPressureStrategy(pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>, elements: Iterable<A>, isShutdown: MutableRef.MutableRef<boolean>): Effect.Effect<boolean>Describes how publishers should signal to subscribers that they are
waiting for space to become available in the PubSub.
handleSurplus(
pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub: PubSub.interface PubSub<in out A>.Atomic<in out A>Low-level atomic PubSub interface that handles the core message storage and retrieval.
Atomic<function (type parameter) A in BackPressureStrategy<in out A>A>,
subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers: PubSub.type PubSub<in out A>.Subscribers<A> = Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A>>>>Tracks the pollers currently waiting on each backing subscription.
Details
This type is part of the low-level PubSub.Strategy contract. Most
application code should use subscribe, take, and the other PubSub
operations instead of manipulating subscriber maps directly.
Subscribers<function (type parameter) A in BackPressureStrategy<in out A>A>,
elements: Iterable<A>elements: interface Iterable<T, TReturn = any, TNext = any>Iterable<function (type parameter) A in BackPressureStrategy<in out A>A>,
isShutdown: MutableRef.MutableRef<boolean>(parameter) isShutdown: {
current: T;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
isShutdown: import MutableRefMutableRef.type MutableRef.MutableRef = /*unresolved*/ anyMutableRef<boolean>
): import EffectEffect.type Effect.Effect = /*unresolved*/ anyEffect<boolean> {
return import EffectEffect.suspend(() => {
const const deferred: Deferred.Deferred<
boolean,
never
>
const deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred = import DeferredDeferred.makeUnsafe<boolean>()
this.function BackPressureStrategy(elements: Iterable<A>, deferred: Deferred.Deferred<boolean>): voidofferUnsafe(elements: Iterable<A>elements, const deferred: Deferred.Deferred<
boolean,
never
>
const deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred)
this.BackPressureStrategy<in out A>.onPubSubEmptySpaceUnsafe(pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>): voidDescribes how subscribers should signal to publishers waiting for space
to become available in the PubSub that space may be available.
onPubSubEmptySpaceUnsafe(pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub, subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers)
this.BackPressureStrategy<in out A>.completeSubscribersUnsafe(pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>): voidDescribes how publishers should signal to subscribers waiting for
additional values from the PubSub that new values are available.
completeSubscribersUnsafe(pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub, subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers)
return (import MutableRefMutableRef.get(isShutdown: MutableRef.MutableRef<boolean>(parameter) isShutdown: {
current: T;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
toString: () => string;
toJSON: () => unknown;
}
isShutdown) ? import EffectEffect.interrupt : import DeferredDeferred.await(const deferred: Deferred.Deferred<
boolean,
never
>
const deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred)).pipe(
import EffectEffect.onInterrupt(() => {
this.function BackPressureStrategy(deferred: Deferred.Deferred<boolean>): voidremoveUnsafe(const deferred: Deferred.Deferred<
boolean,
never
>
const deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred)
return import EffectEffect.void
})
)
})
}
BackPressureStrategy<in out A>.onPubSubEmptySpaceUnsafe(pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>): voidDescribes how subscribers should signal to publishers waiting for space
to become available in the PubSub that space may be available.
onPubSubEmptySpaceUnsafe(
pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub: PubSub.interface PubSub<in out A>.Atomic<in out A>Low-level atomic PubSub interface that handles the core message storage and retrieval.
Atomic<function (type parameter) A in BackPressureStrategy<in out A>A>,
subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers: PubSub.type PubSub<in out A>.Subscribers<A> = Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A>>>>Tracks the pollers currently waiting on each backing subscription.
Details
This type is part of the low-level PubSub.Strategy contract. Most
application code should use subscribe, take, and the other PubSub
operations instead of manipulating subscriber maps directly.
Subscribers<function (type parameter) A in BackPressureStrategy<in out A>A>
): void {
let let keepPolling: booleankeepPolling = true
while (let keepPolling: booleankeepPolling && !pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub.PubSub<in out A>.Atomic<A>.isFull(): booleanisFull()) {
const const publisher: anypublisher = import MutableListMutableList.take(this.BackPressureStrategy<in out A>.publishers: MutableList.MutableList<readonly [A, Deferred.Deferred<boolean>, boolean]>(property) BackPressureStrategy<in out A>.publishers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
publishers)
if (const publisher: anypublisher === import MutableListMutableList.Empty) {
let keepPolling: booleankeepPolling = false
} else {
const [const value: anyvalue, const deferred: Deferred.Deferred<
boolean,
never
>
const deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred] = const publisher: anyconst publisher: {
0: A;
1: Deferred.Deferred<boolean, never>;
2: boolean;
length: 3;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<boolean | A | Deferred.Deferred<boolean, never>>>): Array<boolean | A | Deferred.Deferred<boolean, never>>; (...items: Array<boolean | A | Deferred.Deferred<boolean, never> | ConcatArray<boolean | A | Deferre…;
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
indexOf: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => number;
lastIndexOf: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => number;
every: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: boolean |…;
some: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => void, thisArg?: any) => void;
map: (callbackfn: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): Array<S>; (predicate: (value: boolean | A | Deferre…;
reduce: { (callbackfn: (previousValue: boolean | A | Deferred.Deferred<boolean, never>, currentValue: boolean | A | Deferred.Deferred<boolean, never>, currentIndex: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => …;
reduceRight: { (callbackfn: (previousValue: boolean | A | Deferred.Deferred<boolean, never>, currentValue: boolean | A | Deferred.Deferred<boolean, never>, currentIndex: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => …;
find: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, obj: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): S | undefined; (predicate: (value: boolean | A | Defe…;
findIndex: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, obj: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, boolean | A | Deferred.Deferred<boolean, never>]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<boolean | A | Deferred.Deferred<boolean, never>>;
includes: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: Array<boolean | A | Deferred.Deferred<boolean, never>>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => boolean | A | Deferred.Deferred<boolean, never> | undefined;
findLast: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): S | undefined; (predicate: (value: boolean | A | De…;
findLastIndex: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => number;
toReversed: () => Array<boolean | A | Deferred.Deferred<boolean, never>>;
toSorted: (compareFn?: ((a: boolean | A | Deferred.Deferred<boolean, never>, b: boolean | A | Deferred.Deferred<boolean, never>) => number) | undefined) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<boolean | A | Deferred.Deferred<boolean, never>>): Array<boolean | A | Deferred.Deferred<boolean, never>>; (start: number, deleteCount?: number): Array<boolean | A | Deferred.Deferred<…;
with: (index: number, value: boolean | A | Deferred.Deferred<boolean, never>) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
}
publisher
const const published: booleanpublished = pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub.PubSub<in out A>.Atomic<A>.publish(value: A): booleanpublish(const value: anyvalue)
if (const published: booleanpublished && const publisher: anyconst publisher: {
0: A;
1: Deferred.Deferred<boolean, never>;
2: boolean;
length: 3;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<boolean | A | Deferred.Deferred<boolean, never>>>): Array<boolean | A | Deferred.Deferred<boolean, never>>; (...items: Array<boolean | A | Deferred.Deferred<boolean, never> | ConcatArray<boolean | A | Deferre…;
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
indexOf: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => number;
lastIndexOf: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => number;
every: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: boolean |…;
some: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => void, thisArg?: any) => void;
map: (callbackfn: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): Array<S>; (predicate: (value: boolean | A | Deferre…;
reduce: { (callbackfn: (previousValue: boolean | A | Deferred.Deferred<boolean, never>, currentValue: boolean | A | Deferred.Deferred<boolean, never>, currentIndex: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => …;
reduceRight: { (callbackfn: (previousValue: boolean | A | Deferred.Deferred<boolean, never>, currentValue: boolean | A | Deferred.Deferred<boolean, never>, currentIndex: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => …;
find: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, obj: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): S | undefined; (predicate: (value: boolean | A | Defe…;
findIndex: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, obj: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, boolean | A | Deferred.Deferred<boolean, never>]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<boolean | A | Deferred.Deferred<boolean, never>>;
includes: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: Array<boolean | A | Deferred.Deferred<boolean, never>>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => boolean | A | Deferred.Deferred<boolean, never> | undefined;
findLast: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): S | undefined; (predicate: (value: boolean | A | De…;
findLastIndex: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => number;
toReversed: () => Array<boolean | A | Deferred.Deferred<boolean, never>>;
toSorted: (compareFn?: ((a: boolean | A | Deferred.Deferred<boolean, never>, b: boolean | A | Deferred.Deferred<boolean, never>) => number) | undefined) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<boolean | A | Deferred.Deferred<boolean, never>>): Array<boolean | A | Deferred.Deferred<boolean, never>>; (start: number, deleteCount?: number): Array<boolean | A | Deferred.Deferred<…;
with: (index: number, value: boolean | A | Deferred.Deferred<boolean, never>) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
}
publisher[2]) {
import DeferredDeferred.doneUnsafe(const deferred: Deferred.Deferred<
boolean,
never
>
const deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred, import ExitExit.succeed(true))
} else if (!const published: booleanpublished) {
import MutableListMutableList.prepend(this.BackPressureStrategy<in out A>.publishers: MutableList.MutableList<readonly [A, Deferred.Deferred<boolean>, boolean]>(property) BackPressureStrategy<in out A>.publishers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
publishers, const publisher: anyconst publisher: {
0: A;
1: Deferred.Deferred<boolean, never>;
2: boolean;
length: 3;
toString: () => string;
toLocaleString: { (): string; (locales: string | string[], options?: Intl.NumberFormatOptions & Intl.DateTimeFormatOptions): string };
concat: { (...items: Array<ConcatArray<boolean | A | Deferred.Deferred<boolean, never>>>): Array<boolean | A | Deferred.Deferred<boolean, never>>; (...items: Array<boolean | A | Deferred.Deferred<boolean, never> | ConcatArray<boolean | A | Deferre…;
join: (separator?: string) => string;
slice: (start?: number, end?: number) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
indexOf: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => number;
lastIndexOf: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => number;
every: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): this is readonly S[]; (predicate: (value: boolean |…;
some: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => boolean;
forEach: (callbackfn: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => void, thisArg?: any) => void;
map: (callbackfn: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => U, thisArg?: any) => Array<U>;
filter: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): Array<S>; (predicate: (value: boolean | A | Deferre…;
reduce: { (callbackfn: (previousValue: boolean | A | Deferred.Deferred<boolean, never>, currentValue: boolean | A | Deferred.Deferred<boolean, never>, currentIndex: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => …;
reduceRight: { (callbackfn: (previousValue: boolean | A | Deferred.Deferred<boolean, never>, currentValue: boolean | A | Deferred.Deferred<boolean, never>, currentIndex: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => …;
find: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, obj: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): S | undefined; (predicate: (value: boolean | A | Defe…;
findIndex: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, obj: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => number;
entries: () => ArrayIterator<[number, boolean | A | Deferred.Deferred<boolean, never>]>;
keys: () => ArrayIterator<number>;
values: () => ArrayIterator<boolean | A | Deferred.Deferred<boolean, never>>;
includes: (searchElement: boolean | A | Deferred.Deferred<boolean, never>, fromIndex?: number) => boolean;
flatMap: (callback: (this: This, value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: Array<boolean | A | Deferred.Deferred<boolean, never>>) => U | ReadonlyArray<U>, thisArg?: This | undefined) => Array<U>;
flat: (this: A, depth?: D | undefined) => Array<FlatArray<A, D>>;
at: (index: number) => boolean | A | Deferred.Deferred<boolean, never> | undefined;
findLast: { (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => value is S, thisArg?: any): S | undefined; (predicate: (value: boolean | A | De…;
findLastIndex: (predicate: (value: boolean | A | Deferred.Deferred<boolean, never>, index: number, array: ReadonlyArray<boolean | A | Deferred.Deferred<boolean, never>>) => unknown, thisArg?: any) => number;
toReversed: () => Array<boolean | A | Deferred.Deferred<boolean, never>>;
toSorted: (compareFn?: ((a: boolean | A | Deferred.Deferred<boolean, never>, b: boolean | A | Deferred.Deferred<boolean, never>) => number) | undefined) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
toSpliced: { (start: number, deleteCount: number, ...items: Array<boolean | A | Deferred.Deferred<boolean, never>>): Array<boolean | A | Deferred.Deferred<boolean, never>>; (start: number, deleteCount?: number): Array<boolean | A | Deferred.Deferred<…;
with: (index: number, value: boolean | A | Deferred.Deferred<boolean, never>) => Array<boolean | A | Deferred.Deferred<boolean, never>>;
}
publisher)
}
this.BackPressureStrategy<in out A>.completeSubscribersUnsafe(pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>): voidDescribes how publishers should signal to subscribers waiting for
additional values from the PubSub that new values are available.
completeSubscribersUnsafe(pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub, subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers)
}
}
}
function BackPressureStrategy(pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>, subscription: PubSub.BackingSubscription<A>, pollers: MutableList.MutableList<Deferred.Deferred<A>>): voidDescribes how subscribers waiting for additional values from the PubSub
should take those values and signal to publishers that they are no
longer waiting for additional values.
completePollersUnsafe(
pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub: PubSub.interface PubSub<in out A>.Atomic<in out A>Low-level atomic PubSub interface that handles the core message storage and retrieval.
Atomic<function (type parameter) A in BackPressureStrategy<in out A>A>,
subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers: PubSub.type PubSub<in out A>.Subscribers<A> = Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A>>>>Tracks the pollers currently waiting on each backing subscription.
Details
This type is part of the low-level PubSub.Strategy contract. Most
application code should use subscribe, take, and the other PubSub
operations instead of manipulating subscriber maps directly.
Subscribers<function (type parameter) A in BackPressureStrategy<in out A>A>,
subscription: PubSub.BackingSubscription<A>(parameter) subscription: {
isEmpty: () => boolean;
size: () => number;
poll: () => typeof MutableList.Empty | A;
pollUpTo: (n: number) => Array<A>;
unsubscribe: () => void;
}
subscription: PubSub.interface PubSub<in out A>.BackingSubscription<out A>Low-level subscription interface that handles message polling for individual subscribers.
BackingSubscription<function (type parameter) A in BackPressureStrategy<in out A>A>,
pollers: MutableList.MutableList<
Deferred.Deferred<A>
>
(parameter) pollers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
pollers: import MutableListMutableList.type MutableList.MutableList = /*unresolved*/ anyMutableList<import DeferredDeferred.type Deferred.Deferred = /*unresolved*/ anyDeferred<function (type parameter) A in BackPressureStrategy<in out A>A>>
): void {
return const strategyCompletePollersUnsafe: <A>(
strategy: PubSub.Strategy<A>,
pubsub: PubSub.Atomic<A>,
subscribers: PubSub.Subscribers<A>,
subscription: PubSub.BackingSubscription<A>,
pollers: MutableList.MutableList<
Deferred.Deferred<A>
>
) => void
strategyCompletePollersUnsafe(this, pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub, subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers, subscription: PubSub.BackingSubscription<A>(parameter) subscription: {
isEmpty: () => boolean;
size: () => number;
poll: () => typeof MutableList.Empty | A;
pollUpTo: (n: number) => Array<A>;
unsubscribe: () => void;
}
subscription, pollers: MutableList.MutableList<
Deferred.Deferred<A>
>
(parameter) pollers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
pollers)
}
BackPressureStrategy<in out A>.completeSubscribersUnsafe(pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>): voidDescribes how publishers should signal to subscribers waiting for
additional values from the PubSub that new values are available.
completeSubscribersUnsafe(pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub: PubSub.interface PubSub<in out A>.Atomic<in out A>Low-level atomic PubSub interface that handles the core message storage and retrieval.
Atomic<function (type parameter) A in BackPressureStrategy<in out A>A>, subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers: PubSub.type PubSub<in out A>.Subscribers<A> = Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A>>>>Tracks the pollers currently waiting on each backing subscription.
Details
This type is part of the low-level PubSub.Strategy contract. Most
application code should use subscribe, take, and the other PubSub
operations instead of manipulating subscriber maps directly.
Subscribers<function (type parameter) A in BackPressureStrategy<in out A>A>): void {
return const strategyCompleteSubscribersUnsafe: <A>(strategy: PubSub<in out A>.Strategy<A>, pubsub: PubSub.Atomic<A>, subscribers: PubSub.Subscribers<A>) => voidstrategyCompleteSubscribersUnsafe(this, pubsub: PubSub.Atomic<A>(parameter) pubsub: {
capacity: number;
isEmpty: () => boolean;
isFull: () => boolean;
size: () => number;
publish: (value: A) => boolean;
publishAll: (elements: Iterable<A>) => Array<A>;
slide: () => void;
subscribe: () => PubSub.BackingSubscription<A>;
replayWindow: () => PubSub.ReplayWindow<A>;
}
pubsub, subscribers: PubSub.Subscribers<A>(parameter) subscribers: {
clear: () => void;
delete: (key: PubSub.BackingSubscription<A>) => boolean;
forEach: (callbackfn: (value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>, key: PubSub.BackingSubscription<A>, map: Map<PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>) => void, thisArg?: any)…;
get: (key: PubSub.BackingSubscription<A>) => Set<MutableList.MutableList<Deferred.Deferred<A, never>>> | undefined;
has: (key: PubSub.BackingSubscription<A>) => boolean;
set: (key: PubSub.BackingSubscription<A>, value: Set<MutableList.MutableList<Deferred.Deferred<A, never>>>) => PubSub.Subscribers<A>;
size: number;
entries: () => MapIterator<[PubSub.BackingSubscription<A>, Set<MutableList.MutableList<Deferred.Deferred<A, never>>>]>;
keys: () => MapIterator<PubSub.BackingSubscription<A>>;
values: () => MapIterator<Set<MutableList.MutableList<Deferred.Deferred<A, never>>>>;
}
subscribers)
}
private function BackPressureStrategy(elements: Iterable<A>, deferred: Deferred.Deferred<boolean>): voidofferUnsafe(elements: Iterable<A>elements: interface Iterable<T, TReturn = any, TNext = any>Iterable<function (type parameter) A in BackPressureStrategy<in out A>A>, deferred: Deferred.Deferred<boolean>(parameter) deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred: import DeferredDeferred.type Deferred.Deferred = /*unresolved*/ anyDeferred<boolean>): void {
const const iterator: Iterator<A, any, any>iterator = elements: Iterable<A>elements[var Symbol: SymbolConstructorSymbol.SymbolConstructor.iterator: typeof Symbol.iteratorA method that returns the default iterator for an object. Called by the semantics of the
for-of statement.
iterator]()
let let next: IteratorResult<A, any>next: type IteratorResult<T, TReturn = any> =
| IteratorYieldResult<T>
| IteratorReturnResult<TReturn>
IteratorResult<function (type parameter) A in BackPressureStrategy<in out A>A> = const iterator: Iterator<A, any, any>iterator.Iterator<A, any, any>.next(...[value]: [] | [any]): IteratorResult<A, any>next()
if (!let next: IteratorResult<A, any>next.done?: boolean | undefineddone) {
// oxlint-disable-next-line no-constant-condition
while (1) {
const const value: in out Avalue = let next: IteratorYieldResult<A>next.IteratorYieldResult<A>.value: in out Avalue
let next: IteratorResult<A, any>next = const iterator: Iterator<A, any, any>iterator.Iterator<A, any, any>.next(...[value]: [] | [any]): IteratorResult<A, any>next()
if (let next: IteratorResult<A, any>next.done?: boolean | undefineddone) {
import MutableListMutableList.append(this.BackPressureStrategy<in out A>.publishers: MutableList.MutableList<readonly [A, Deferred.Deferred<boolean>, boolean]>(property) BackPressureStrategy<in out A>.publishers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
publishers, [const value: in out Avalue, deferred: Deferred.Deferred<boolean>(parameter) deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred, true])
break
}
import MutableListMutableList.append(this.BackPressureStrategy<in out A>.publishers: MutableList.MutableList<readonly [A, Deferred.Deferred<boolean>, boolean]>(property) BackPressureStrategy<in out A>.publishers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
publishers, [const value: in out Avalue, deferred: Deferred.Deferred<boolean>(parameter) deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred, false])
}
}
}
function BackPressureStrategy(deferred: Deferred.Deferred<boolean>): voidremoveUnsafe(deferred: Deferred.Deferred<boolean>(parameter) deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred: import DeferredDeferred.type Deferred.Deferred = /*unresolved*/ anyDeferred<boolean>): void {
import MutableListMutableList.filter(this.BackPressureStrategy<in out A>.publishers: MutableList.MutableList<readonly [A, Deferred.Deferred<boolean>, boolean]>(property) BackPressureStrategy<in out A>.publishers: {
head: MutableList.Bucket<A> | undefined;
tail: MutableList.Bucket<A> | undefined;
length: number;
}
publishers, ([_: any_, d: Deferred.Deferred<boolean, never>(parameter) d: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
d]) => d: Deferred.Deferred<boolean, never>(parameter) d: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
d !== deferred: Deferred.Deferred<boolean>(parameter) deferred: {
effect: Effect<A, E>;
resumes: Array<(effect: Effect<A, E>) => void> | undefined;
pipe: { <A>(this: A): A; <A, B = never>(this: A, ab: (_: A) => B): B; <A, B = never, C = never>(this: A, ab: (_: A) => B, bc: (_: B) => C): C; <A, B = never, C = never, D = never>(this: A, ab: (_: A) => B, bc: (_: B) => C, cd: (_: C) => D): D; <…;
}
deferred)
}
}