File size: 2,909 Bytes
a5d718a | 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 | var __defProp = Object.defineProperty;
var __getOwnPropDesc = Object.getOwnPropertyDescriptor;
var __getOwnPropNames = Object.getOwnPropertyNames;
var __hasOwnProp = Object.prototype.hasOwnProperty;
var __export = (target, all) => {
for (var name in all)
__defProp(target, name, { get: all[name], enumerable: true });
};
var __copyProps = (to, from, except, desc) => {
if (from && typeof from === "object" || typeof from === "function") {
for (let key of __getOwnPropNames(from))
if (!__hasOwnProp.call(to, key) && key !== except)
__defProp(to, key, { get: () => from[key], enumerable: !(desc = __getOwnPropDesc(from, key)) || desc.enumerable });
}
return to;
};
var __toCommonJS = (mod) => __copyProps(__defProp({}, "__esModule", { value: true }), mod);
var stream_exports = {};
__export(stream_exports, {
StreamingApi: () => StreamingApi
});
module.exports = __toCommonJS(stream_exports);
class StreamingApi {
writer;
encoder;
writable;
abortSubscribers = [];
responseReadable;
/**
* Whether the stream has been aborted.
*/
aborted = false;
/**
* Whether the stream has been closed normally.
*/
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 the stream.
* You can call this method when stream is aborted by external event.
*/
abort() {
if (!this.aborted) {
this.aborted = true;
this.abortSubscribers.forEach((subscriber) => subscriber());
}
}
}
// Annotate the CommonJS export names for ESM import in node:
0 && (module.exports = {
StreamingApi
});
|