Files
curriculum-project-hub/hub/src/capability/docmindClient.ts
T
hongjr03 97a99cd381 fix(hub): raise Docmind OSS upload timeout past httpx 3s default
SubmitDocParserJobAdvance uploads PDFs to Aliyun OSS via tea/httpx, which
defaults readTimeout to 3000ms when RuntimeOptions is empty. Multi-MB
teacher PDFs on para silo failed with ReadTimeout(3000) before the job
could start. Set connectTimeout=15s and readTimeout=5m.
2026-07-30 11:16:56 +08:00

279 lines
10 KiB
TypeScript

/**
* ADR-0027: Alibaba Cloud Document Mind (docmind) client.
*
* Uses the official @alicloud/docmind-api20220711 SDK to call the document
* parsing (large model version) API. The API is asynchronous:
* 1. SubmitDocParserJobAdvance — upload local file as a stream, get job id
* 2. QueryDocParserStatus — poll until data.status === "success";
* the markdown output URL is returned in outputFormatResult
* 3. Download markdown from OSS, extract and download referenced images
*
* Pricing (2025-07, aliyun docmind):
* - 图文文档基础链路: 0.02元/页
* - 图文文档增强链路 (含公式 LaTeX): 0.04元/页
* - 视频: 0.002元/秒
* - 音频: 0.00035元/秒
*/
import $DocmindClient, {
SubmitDocParserJobAdvanceRequest,
QueryDocParserStatusRequest,
} from "@alicloud/docmind-api20220711";
import { RuntimeOptions } from "@alicloud/tea-util";
import { createReadStream, type ReadStream } from "node:fs";
import { once } from "node:events";
import { basename } from "node:path";
import type { DocmindCapabilitySecretPayload } from "./types.js";
/** A single extracted image downloaded from the markdown's OSS image URLs. */
export interface DocmindExtractedImage {
readonly filename: string;
readonly data: Uint8Array;
}
/** The structured result of parsing one document. */
export interface DocmindParseResult {
readonly markdown: string;
readonly images: readonly DocmindExtractedImage[];
readonly pageCount: number;
readonly costUsd: number | null;
readonly requestId: string | null;
}
export interface DocmindParseOptions {
readonly inputFilePath: string;
}
export interface CapabilityProviderClient {
parse(credential: DocmindCapabilitySecretPayload, options: DocmindParseOptions): Promise<DocmindParseResult>;
}
export class DocmindClientError extends Error {
constructor(
message: string,
readonly code: "docmind_unreachable" | "docmind_rejected" | "docmind_invalid_response" | "docmind_no_output" | "docmind_timeout",
readonly upstreamStatus?: number,
) {
super(message);
this.name = "DocmindClientError";
}
}
/** Cost per page in USD (0.04 CNY/page ≈ 0.0056 USD). */
const COST_PER_PAGE_USD = 0.0056;
const POLL_INTERVAL_MS = 10_000;
const POLL_TIMEOUT_MS = 5 * 60_000;
/**
* httpx (tea transport) defaults read/connect timeout to 3000ms when unset.
* SubmitDocParserJobAdvance uploads the PDF to OSS; multi-MB files routinely
* exceed 3s on the silo host (production: ReadTimeout(3000) on
* docmind-api-cn-hangzhou.oss-cn-hangzhou.aliyuncs.com).
*/
export const DOCMIND_CONNECT_TIMEOUT_MS = 15_000;
/** Allow slow / large PDF OSS uploads up to the same bound as job polling. */
export const DOCMIND_READ_TIMEOUT_MS = POLL_TIMEOUT_MS;
/** RuntimeOptions for Docmind SDK calls that may upload or wait on the wire. */
export function createDocmindRuntimeOptions(): RuntimeOptions {
return new RuntimeOptions({
connectTimeout: DOCMIND_CONNECT_TIMEOUT_MS,
readTimeout: DOCMIND_READ_TIMEOUT_MS,
});
}
type DocmindConfig = ConstructorParameters<typeof $DocmindClient.default>[0];
export class AliyunDocmindClient implements CapabilityProviderClient {
async parse(credential: DocmindCapabilitySecretPayload, options: DocmindParseOptions): Promise<DocmindParseResult> {
// Open first so missing local inputs fail closed before touching the SDK.
// Unhandled createReadStream('error') previously crashed the Hub process.
const fileName = basename(options.inputFilePath);
const fileStream = await openLocalFileStream(options.inputFilePath);
const config: DocmindConfig = {
endpoint: credential.endpoint,
accessKeyId: credential.accessKeyId,
accessKeySecret: credential.accessKeySecret,
type: "access_key",
regionId: "cn-hangzhou",
} as DocmindConfig;
const client = new $DocmindClient.default(config);
// Submit job with local file as a ReadStream (not a Buffer — the SDK
// serializes Buffers as JSON {type:"Buffer",data:[...]} which the API
// can't read; a Stream is uploaded as multipart form data).
const advanceRequest = new SubmitDocParserJobAdvanceRequest({
fileUrlObject: fileStream,
fileName,
outputFormat: ["markdown"],
formulaEnhancement: true,
});
const runtime = createDocmindRuntimeOptions();
let submitResponse;
try {
submitResponse = await client.submitDocParserJobAdvance(advanceRequest, runtime);
} catch (e) {
fileStream.destroy();
throw new DocmindClientError(
e instanceof Error ? e.message : String(e),
"docmind_unreachable",
);
}
const jobId = submitResponse.body?.data?.id;
if (jobId === undefined || jobId === null || jobId === "") {
throw new DocmindClientError("docmind submit returned no job id", "docmind_invalid_response");
}
// 2. Poll QueryDocParserStatus until data.status === "success" or "fail".
const deadline = Date.now() + POLL_TIMEOUT_MS;
let status = "";
let markdownUrl: string | null = null;
let pageCount = 0;
while (Date.now() < deadline) {
await sleep(POLL_INTERVAL_MS);
const statusReq = new QueryDocParserStatusRequest({ id: jobId });
const statusResponse = await client.queryDocParserStatus(statusReq);
const body = statusResponse.body as {
data?: {
status?: string;
pageCountEstimate?: number;
outputFormatResult?: Array<{ outputFileUrl?: string; outputType?: string }>;
};
};
status = body.data?.status ?? "";
if (status === "success") {
const mdResult = body.data?.outputFormatResult?.find((r) => r.outputType === "markdown");
markdownUrl = mdResult?.outputFileUrl ?? null;
pageCount = body.data?.pageCountEstimate ?? 0;
break;
}
if (status === "fail") {
throw new DocmindClientError(`docmind job ${jobId} failed`, "docmind_rejected");
}
}
if (status !== "success") {
throw new DocmindClientError(`docmind job ${jobId} timed out (status: ${status})`, "docmind_timeout");
}
if (markdownUrl === null) {
throw new DocmindClientError("docmind returned no markdown output URL", "docmind_no_output");
}
// 3. Download the markdown file from OSS.
let markdown: string;
try {
const mdResp = await fetch(markdownUrl);
if (!mdResp.ok) {
throw new DocmindClientError(`failed to download markdown: HTTP ${mdResp.status}`, "docmind_unreachable");
}
markdown = await mdResp.text();
} catch (e) {
if (e instanceof DocmindClientError) throw e;
throw new DocmindClientError(
e instanceof Error ? e.message : String(e),
"docmind_unreachable",
);
}
if (markdown === "") {
throw new DocmindClientError("docmind returned empty markdown", "docmind_no_output");
}
// 4. Download images referenced in the markdown (OSS URLs).
// Markdown contains ![filename](http://...oss.../image.png?...) entries.
// We rewrite them to local relative paths and download the images.
const images: DocmindExtractedImage[] = [];
const imageRefRegex = /!\[([^\]]*)\]\((https?:\/\/[^)]+)\)/g;
const localMarkdown = markdown.replace(imageRefRegex, (match, altText: string, url: string) => {
const idx = images.findIndex((img) => img.filename === extractFilename(altText, url, images.length));
if (idx >= 0) {
const img = images[idx]!;
return `![${altText}](${img.filename})`;
}
return match;
});
// Collect all image URLs first, then download.
const imageUrls: Array<{ url: string; altText: string }> = [];
let match: RegExpExecArray | null;
const collectRegex = /!\[([^\]]*)\]\((https?:\/\/[^)]+)\)/g;
while ((match = collectRegex.exec(markdown)) !== null) {
imageUrls.push({ url: match[2]!, altText: match[1]! });
}
for (let i = 0; i < imageUrls.length; i++) {
const { url, altText } = imageUrls[i]!;
const filename = extractFilename(altText, url, i);
try {
const imgResp = await fetch(url);
if (!imgResp.ok) continue;
const data = new Uint8Array(await imgResp.arrayBuffer());
images.push({ filename, data });
} catch {
// Best-effort: skip images that fail to download.
}
}
// Rewrite markdown with local image paths.
let finalMarkdown = markdown;
let imageIdx = 0;
finalMarkdown = finalMarkdown.replace(imageRefRegex, (match, altText: string, _url: string) => {
if (imageIdx < images.length) {
const img = images[imageIdx]!;
imageIdx++;
return `![${altText}](${img.filename})`;
}
return match;
});
if (pageCount === 0) {
pageCount = Math.max(1, images.length);
}
const costUsd = (pageCount > 0 ? pageCount : 1) * COST_PER_PAGE_USD;
return { markdown: finalMarkdown, images, pageCount, costUsd, requestId: jobId };
}
}
function sleep(ms: number): Promise<void> {
return new Promise((resolve) => {
setTimeout(resolve, ms);
});
}
/** Derive a clean filename from the alt text or URL. */
function extractFilename(altText: string, url: string, index: number): string {
// Try alt text first (docmind often puts the original filename).
if (altText !== "" && altText.length < 100) {
const cleaned = altText.replace(/[^a-zA-Z0-9._-]/g, "_");
if (cleaned.length > 0) return cleaned;
}
// Fall back to URL path.
const urlPath = new URL(url).pathname;
const base = basename(urlPath);
if (base !== "" && base !== "/") return base;
return `image_${index + 1}.png`;
}
/**
* Open a local file as a ReadStream only after the fd is successfully open.
* createReadStream() emits asynchronous 'error' for missing paths; without a
* listener that becomes an unhandled EventEmitter error and exits Node.
*/
async function openLocalFileStream(path: string): Promise<ReadStream> {
const stream = createReadStream(path);
try {
await once(stream, "open");
} catch (error) {
stream.destroy();
const err = error instanceof Error ? error : new Error(String(error));
const code = (err as NodeJS.ErrnoException).code;
if (code === "ENOENT") {
throw new DocmindClientError(`input file not found: ${path}`, "docmind_rejected");
}
throw new DocmindClientError(err.message, "docmind_unreachable");
}
// After open, residual stream errors must not become unhandled and crash Hub.
stream.on("error", () => {
// The Aliyun SDK / destroy path owns consumption failures after open.
});
return stream;
}