|
|
| var StreamingApi = class {
|
| writer;
|
| encoder;
|
| writable;
|
| abortSubscribers = [];
|
| responseReadable;
|
| |
| |
|
|
| aborted = false;
|
| |
| |
|
|
| closed = false;
|
| constructor(writable, _readable) {
|
| this.writable = writable;
|
| this.writer = writable.getWriter();
|
| this.encoder = new TextEncoder();
|
| const reader = _readable.getReader();
|
| this.abortSubscribers.push(async () => {
|
| await reader.cancel();
|
| });
|
| this.responseReadable = new ReadableStream({
|
| async pull(controller) {
|
| const { done, value } = await reader.read();
|
| done ? controller.close() : controller.enqueue(value);
|
| },
|
| cancel: () => {
|
| if (!this.closed) {
|
| this.abort();
|
| }
|
| }
|
| });
|
| }
|
| async write(input) {
|
| try {
|
| if (typeof input === "string") {
|
| input = this.encoder.encode(input);
|
| }
|
| await this.writer.write(input);
|
| } catch {
|
| }
|
| return this;
|
| }
|
| async writeln(input) {
|
| await this.write(input + "\n");
|
| return this;
|
| }
|
| sleep(ms) {
|
| return new Promise((res) => setTimeout(res, ms));
|
| }
|
| async close() {
|
| this.closed = true;
|
| try {
|
| await this.writer.close();
|
| } catch {
|
| }
|
| }
|
| async pipe(body) {
|
| this.writer.releaseLock();
|
| await body.pipeTo(this.writable, { preventClose: true });
|
| this.writer = this.writable.getWriter();
|
| }
|
| onAbort(listener) {
|
| this.abortSubscribers.push(listener);
|
| }
|
| |
| |
| |
|
|
| abort() {
|
| if (!this.aborted) {
|
| this.aborted = true;
|
| this.abortSubscribers.forEach((subscriber) => subscriber());
|
| }
|
| }
|
| };
|
| export {
|
| StreamingApi
|
| };
|
|
|