-
Notifications
You must be signed in to change notification settings - Fork 3k
Expand file tree
/
Copy pathAsyncScheduler.ts
More file actions
52 lines (45 loc) · 1.63 KB
/
Copy pathAsyncScheduler.ts
File metadata and controls
52 lines (45 loc) · 1.63 KB
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
import { Scheduler } from '../Scheduler';
import { Action } from './Action';
import { AsyncAction } from './AsyncAction';
import { reportUnhandledError } from '../util/reportUnhandledError';
import { TimerHandle } from './timerHandle';
export class AsyncScheduler extends Scheduler {
public actions: Array<AsyncAction<any>> = [];
/**
* A flag to indicate whether the Scheduler is currently executing a batch of
* queued actions.
* @internal
*/
public _active: boolean = false;
/**
* An internal ID used to track the latest asynchronous task such as those
* coming from `setTimeout`, `setInterval`, `requestAnimationFrame`, and
* others.
* @internal
*/
public _scheduled: TimerHandle | undefined;
constructor(SchedulerAction: typeof Action, now: () => number = Scheduler.now) {
super(SchedulerAction, now);
}
public flush(action: AsyncAction<any>): void {
const { actions } = this;
if (this._active) {
actions.push(action);
return;
}
let error: any;
this._active = true;
do {
if ((error = action.execute(action.state, action.delay))) {
// Report the error asynchronously so it doesn't tear down the
// synchronous subscriber chain (e.g. observeOn(queueScheduler)).
// The erroring action already unsubscribed itself in _execute().
// Continue flushing remaining actions — they are independent
// operations that should not be affected by a sibling's error.
reportUnhandledError(error);
error = null;
}
} while ((action = actions.shift()!)); // exhaust the scheduler queue
this._active = false;
}
}