/
BufferNode.ts
89 lines (81 loc) · 2.62 KB
/
BufferNode.ts
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
import { DataFrame } from '../../data/DataFrame';
import { GraphOptions, PullOptions } from '../../graph/options';
import { Node, NodeOptions } from '../../Node';
import { DataFrameService } from '../../service';
/**
* @category Flow shape
*/
export class BufferNode<InOut extends DataFrame> extends Node<InOut, InOut> {
protected service: DataFrameService<InOut>;
protected options: BufferOptions;
constructor(options?: BufferOptions) {
super(options);
this.on('pull', this.onPull.bind(this));
this.on('push', this.onPush.bind(this));
this.on('build', this._initService.bind(this));
}
private _initService(): Promise<void> {
return new Promise((resolve) => {
if (!this.service) {
this.service = this.model.findDataService(
this.options.service || DataFrame,
) as unknown as DataFrameService<InOut>;
}
resolve();
});
}
public onPull(options?: PullOptions): Promise<void> {
return new Promise((resolve, reject) => {
this.shift()
.then((frame) => {
if (frame) {
this.outlets.forEach((outlet) => outlet.push(frame, options as GraphOptions));
}
resolve();
})
.catch(reject);
});
}
public onPush(frame: InOut): Promise<void> {
return new Promise<void>((resolve, reject) => {
this.service
.insertFrame(frame)
.then(() => {
resolve();
})
.catch(reject);
});
}
protected next(): Promise<InOut> {
return new Promise((resolve, reject) => {
this.service
.findOne(
{},
{
sort: [['createdTimestamp', 1]],
},
)
.then(resolve)
.catch(reject);
});
}
protected shift(): Promise<InOut> {
return new Promise((resolve, reject) => {
let result: InOut;
this.next()
.then((frame: InOut) => {
if (frame) {
result = frame;
return this.service.delete(frame.uid);
} else {
resolve(undefined);
}
})
.then(() => resolve(result))
.catch(reject);
});
}
}
export interface BufferOptions extends NodeOptions {
service?: string;
}