Press n or j to go to the next uncovered block, b, p or k for the previous block.
| 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 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 | 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x 1x | import { Observable, Subject, firstValueFrom } from "rxjs";
import { take } from "rxjs/operators";
import { DurableSocket } from "./durable-socket";
import { ResettableReplaySubject } from "./resettable-subject";
export interface RPCChannel {
received: Observable<string>;
/**
* This optional event signifies that the ongoing state of the channel has been lost,
* and any outstanding requests are no longer resolvable. An example of such
* an event might be the connection being lost.
*
* The value of the observable is the error message when state has been lost.
*/
stateLost?: Observable<string>;
/**
* This optional event signifies that the channel is ready for communications.
*
* If provided, then this event must fire as soon as the channel is ready for use
* (even if the subscription occurs after the channel becomes ready for use).
*
* If the channel supports re-establishment, then subscribing to this
* observable while re-establishment is occuring must not cause an event until
* the channel is re-established.
*/
ready?: Observable<void>;
send(message: string);
close?();
}
/**
* A channel which operates within the local process.
*/
export class LocalChannel implements RPCChannel {
private _received = new Subject<string>();
get received() { return this._received.asObservable(); }
send(message: string) {
this.otherChannel._received.next(message);
}
private otherChannel: LocalChannel;
static makePair(): [ LocalChannel, LocalChannel ] {
let a = new LocalChannel();
let b = new LocalChannel();
a.otherChannel = b;
b.otherChannel = a;
return [a, b];
}
}
/**
* A channel that operates over a WebSocket or RTCDataChannel.
*/
export class SocketChannel implements RPCChannel {
constructor(readonly socket: WebSocket | RTCDataChannel) {
if (socket.readyState === WebSocket.OPEN) {
this.markReady();
} else {
socket.addEventListener('open', () => this.markReady());
}
socket.addEventListener('message', (ev: MessageEvent<any>) => this._received.next(ev.data));
socket.addEventListener('close', () => this.stateWasLost(`Disconnected permanently`));
socket.addEventListener('error', () => this.stateWasLost(`Disconnected permanently`));
}
private _ready = new ResettableReplaySubject<void>(1);
get ready() { return this._ready.asObservable(); }
markReady() {
this._ready.next();
}
markNotReady() {
this._ready.reset();
}
private _stateLost = new Subject<string>();
private _stateLost$ = this._stateLost.asObservable();
get stateLost() { return this._stateLost$; }
private _received = new Subject<string>();
private _received$ = this._received.asObservable();
get received() { return this._received$; }
protected stateWasLost(errorMessage: string) {
this._stateLost.next(errorMessage);
}
async send(message: any) {
await firstValueFrom(this.ready);
this.socket.send(message)
}
close() {
this.socket.close();
}
}
/**
* A channel that operates on a DurableSocket.
*/
export class DurableSocketChannel extends SocketChannel {
constructor(socket: DurableSocket) {
super(socket);
socket.addEventListener('lost', () => (this.markNotReady(), this.stateWasLost(`Connection lost`)));
socket.addEventListener('restore', () => this.markReady());
}
readonly socket: DurableSocket;
}
/**
* A channel that operates via window-to-window (or frame-to-frame) postMessage.
*/
export class WindowChannel implements RPCChannel {
constructor(private remoteWindow: Window, origin?: string) {
window.addEventListener('message', this.handler = ev => {
if (origin && ev.origin !== origin)
return;
this._received.next(ev.data)
});
}
private handler;
private _received = new Subject<string>();
get received() { return this._received.asObservable(); }
send(message: any) {
this.remoteWindow.postMessage(message, '*');
}
close() {
window.removeEventListener('message', this.handler);
}
}
/**
* @unstable Caution: This API may change in minor or patch releases until it is marked as stable.
*/
interface PostMessageTarget {
addEventListener(event: 'message', callback: (ev: MessageEvent) => void);
removeEventListener(event: 'message', callback: (ev: MessageEvent) => void);
postMessage(data: string);
}
/**
* A channel that operates via window-to-window (or frame-to-frame) postMessage.
* @unstable Caution: This API may change in minor or patch releases until it is marked as stable.
*/
export class PostMessageChannel implements RPCChannel {
constructor(private remote: PostMessageTarget, readonly requiredOrigin?: string) {
remote.addEventListener('message', this.handler);
}
private handler = (ev: MessageEvent) => {
if (this.requiredOrigin && ev.origin !== this.requiredOrigin)
return;
this._received.next(ev.data)
};
private _received = new Subject<string>();
get received() { return this._received.asObservable(); }
send(message: any) {
this.remote.postMessage(message);
}
close() {
this.remote.removeEventListener('message', this.handler);
}
} |