import { NOOP_SUBSCRIBER, NOOP_SUBSCRIPTION, Source, Subscriber, Subscription } from "../_core"
import { Transaction } from "../_tx"
// TODO: flatten to constants
export enum EventType {
NEXT = 1,
ERROR = 2,
END = 3,
}
export abstract class Operator implements Subscriber, Source, Subscription {
public readonly weight: number
protected sink: Subscriber = NOOP_SUBSCRIBER
protected active: boolean = false
protected ord: number = -1
protected subs: Subscription = NOOP_SUBSCRIPTION
constructor(protected source: Source) {
this.weight = source.weight
}
public subscribe(subscriber: Subscriber, order: number): Subscription {
this.init(subscriber, order, this.source.subscribe(this, order))
return this
}
public activate(initialNeeded: boolean): void {
this.subs.activate(initialNeeded)
}
public reorder(order: number): void {
this.ord = order
this.subs.reorder(order)
}
public dispose(): void {
const { subs } = this
this.sink = NOOP_SUBSCRIBER
this.subs = NOOP_SUBSCRIPTION
this.active = false
this.ord = -1
subs.dispose()
}
public abstract next(tx: Transaction, val: A): void
public error(tx: Transaction, err: Error): void {
this.sink.error(tx, err)
}
public end(tx: Transaction): void {
this.sink.end(tx)
}
protected init(subscriber: Subscriber, order: number, subscription: Subscription): void {
this.active = true
this.sink = subscriber
this.ord = order
this.subs = subscription
}
}
export interface PipeSubscriber {
pipedNext(sender: Pipe, tx: Transaction, val: T): void
pipedError(sender: Pipe, tx: Transaction, err: Error): void
pipedEnd(sender: Pipe, tx: Transaction): void
}
export class Pipe implements Subscriber {
constructor(public s: PipeSubscriber) {}
public next(tx: Transaction, val: T): void {
this.s.pipedNext(this, tx, val)
}
public error(tx: Transaction, err: Error): void {
this.s.pipedError(this, tx, err)
}
public end(tx: Transaction): void {
this.s.pipedEnd(this, tx)
}
}
export class LinkedPipe extends Pipe {
constructor(
s: PipeSubscriber,
public h: LinkedPipe | null,
public t: LinkedPipe | null,
) {
super(s)
}
}
export class LinkedPipeList {
public size: number
private h: LinkedPipe | null
constructor(subscribers: Array>) {
this.size = subscribers.length
if (subscribers.length === 0) {
this.h = null
} else {
let tail = (this.h = new LinkedPipe(subscribers[0], null, null))
for (let i = 0; i < this.size; i++) {
tail = tail.t = new LinkedPipe(subscribers[i], tail, null)
}
}
}
public head(): LinkedPipe | null {
return this.h
}
public append(subscriber: PipeSubscriber): LinkedPipe {
const node = new LinkedPipe(subscriber, this.h, null)
this.h === null && (this.h = node)
++this.size
return node
}
public remove(node: LinkedPipe): void {
node.h !== null ? (node.h.t = node.t) : (this.h = node.t)
node.t !== null ? (node.t.h = node.h) : void 0
--this.size
}
public clear(): void {
this.h = null
this.size = 0
}
}
export class Identity extends Operator {
public next(tx: Transaction, val: T): void {
this.sink.next(tx, val)
}
}