import { WebCompressor } from '@cloudpss/zstd'; import { encode as encodeUbjson } from '@cloudpss/ubjson/rxjs'; import { of, type OperatorFunction, Observable, tap } from 'rxjs'; import { DEFAULT_LEVEL } from '@cloudpss/zstd/config'; const MAX_DEBUG_SIZE = 4 * 1024 * 1024; const MAX_BUFFER_SIZE = 4 * 1024 * 1024; /** 压缩结果 */ interface EncodeResult { /** 类型 */ contentType: string; /** 数据 */ data: BodyInit; /** Ubjson 原始大小 */ rawSize: number; /** Ubjson 压缩大小 */ compressedSize: number; } /** 压缩 */ function compress(): OperatorFunction, Uint8Array> { return (source) => { return new Observable((subscriber) => { const controller: TransformStreamDefaultController> = { desiredSize: 0, enqueue(chunk) { if (!chunk) return; subscriber.next(chunk); }, error(reason) { subscriber.error(reason); }, terminate() { subscriber.complete(); }, }; const compressor = new WebCompressor(DEFAULT_LEVEL); Promise.resolve() .then(async () => { await Promise.resolve(compressor.start(controller)); source.subscribe({ next(value) { Promise.resolve(compressor.transform(value, controller)).catch((err: Error) => controller.error(err), ); }, complete() { Promise.resolve(compressor.flush(controller)) .then(() => controller.terminate()) .catch((err: Error) => controller.error(err)); }, error(err: Error) { controller.error(err); }, }); }) .catch((err) => controller.error(err)); }); }; } /** 编码并压缩 */ async function encodeImpl(input: unknown): Promise { let rawSize = 0; const data = await new Promise((resolve, reject) => { const blobs: Blob[] = []; let bufferedSize = 0; const result: Array> = []; of(input) .pipe( encodeUbjson(), tap((chunk) => (rawSize += chunk.byteLength)), compress(), ) .subscribe({ next(chunk) { result.push(chunk); bufferedSize += chunk.byteLength; if (bufferedSize > MAX_BUFFER_SIZE) { blobs.push(new Blob(result.splice(0), { type: `application/ubjson+zstd` })); bufferedSize = 0; } }, complete() { if (result.length > 0) { blobs.push(new Blob(result)); } if (blobs.length === 1) { resolve(blobs[0]!); } resolve(new Blob(blobs, { type: `application/ubjson+zstd` })); }, error(err: Error) { reject(err); }, }); }); return { contentType: `application/ubjson+zstd`, data, rawSize, compressedSize: data?.size, }; } /** 压缩 */ export async function encodeDev(input: unknown): Promise { const result = await encodeImpl(input); if (result.rawSize < MAX_DEBUG_SIZE) { const json = JSON.stringify(input); return { contentType: `application/json;charset=utf-8`, data: json, rawSize: result.rawSize, compressedSize: result.compressedSize, }; } return result; } /** 压缩 */ export async function encode(input: unknown): Promise { return await encodeImpl(input); }