67 lines
2.3 KiB
JavaScript
67 lines
2.3 KiB
JavaScript
"use strict";
|
|
Object.defineProperty(exports, "__esModule", { value: true });
|
|
exports.EventStreamCodec = void 0;
|
|
const crc32_1 = require("@aws-crypto/crc32");
|
|
const HeaderMarshaller_1 = require("./HeaderMarshaller");
|
|
const splitMessage_1 = require("./splitMessage");
|
|
class EventStreamCodec {
|
|
constructor(toUtf8, fromUtf8) {
|
|
this.headerMarshaller = new HeaderMarshaller_1.HeaderMarshaller(toUtf8, fromUtf8);
|
|
this.messageBuffer = [];
|
|
this.isEndOfStream = false;
|
|
}
|
|
feed(message) {
|
|
this.messageBuffer.push(this.decode(message));
|
|
}
|
|
endOfStream() {
|
|
this.isEndOfStream = true;
|
|
}
|
|
getMessage() {
|
|
const message = this.messageBuffer.pop();
|
|
const isEndOfStream = this.isEndOfStream;
|
|
return {
|
|
getMessage() {
|
|
return message;
|
|
},
|
|
isEndOfStream() {
|
|
return isEndOfStream;
|
|
},
|
|
};
|
|
}
|
|
getAvailableMessages() {
|
|
const messages = this.messageBuffer;
|
|
this.messageBuffer = [];
|
|
const isEndOfStream = this.isEndOfStream;
|
|
return {
|
|
getMessages() {
|
|
return messages;
|
|
},
|
|
isEndOfStream() {
|
|
return isEndOfStream;
|
|
},
|
|
};
|
|
}
|
|
encode({ headers: rawHeaders, body }) {
|
|
const headers = this.headerMarshaller.format(rawHeaders);
|
|
const length = headers.byteLength + body.byteLength + 16;
|
|
const out = new Uint8Array(length);
|
|
const view = new DataView(out.buffer, out.byteOffset, out.byteLength);
|
|
const checksum = new crc32_1.Crc32();
|
|
view.setUint32(0, length, false);
|
|
view.setUint32(4, headers.byteLength, false);
|
|
view.setUint32(8, checksum.update(out.subarray(0, 8)).digest(), false);
|
|
out.set(headers, 12);
|
|
out.set(body, headers.byteLength + 12);
|
|
view.setUint32(length - 4, checksum.update(out.subarray(8, length - 4)).digest(), false);
|
|
return out;
|
|
}
|
|
decode(message) {
|
|
const { headers, body } = (0, splitMessage_1.splitMessage)(message);
|
|
return { headers: this.headerMarshaller.parse(headers), body };
|
|
}
|
|
formatHeaders(rawHeaders) {
|
|
return this.headerMarshaller.format(rawHeaders);
|
|
}
|
|
}
|
|
exports.EventStreamCodec = EventStreamCodec;
|