Files
LWLT-AIBOT/control-plane/src/input-attachment.ts
2026-08-31 15:42:11 +08:00

351 lines
13 KiB
TypeScript

import { lookup as dnsLookup } from 'node:dns/promises';
import { request as httpsRequest } from 'node:https';
import { isIP, type LookupFunction } from 'node:net';
import { basename } from 'node:path';
import { sha256Bytes } from './crypto.js';
import { diagnosticDurationMs } from './diagnostics.js';
const SHA256_PATTERN = /^[a-f0-9]{64}$/i;
const BASE64_PATTERN = /^[A-Za-z0-9+/_-]*={0,2}$/;
const ALLOWED_ROSTER_EXTENSIONS = new Set(['.xls', '.xlsx']);
const MAX_REDIRECTS = 2;
const DOWNLOAD_TIMEOUT_MS = 30_000;
export interface TaskInputAttachmentInput {
fileName: string;
contentType: string;
content: Buffer;
declaredSize?: number;
declaredSha256?: string;
source: 'manual' | 'agentbus';
}
export interface EncodedTaskInputAttachment {
name: string;
content_type?: string;
size?: number;
sha256?: string;
content_base64: string;
}
export interface AgentBusInputAttachmentReference {
name: string;
contentType: string;
size?: number;
sha256?: string;
url: string;
}
export type InputAttachmentDiagnostic = (
event: string,
metadata: Record<string, unknown>
) => void;
export class InputAttachmentError extends Error {
constructor(
public readonly code: string,
message: string,
public readonly details: Record<string, unknown> = {}
) {
super(message);
this.name = 'InputAttachmentError';
}
}
function emitAttachmentDiagnostic(
diagnostic: InputAttachmentDiagnostic | undefined,
event: string,
metadata: Record<string, unknown>
): void {
if (!diagnostic) return;
try {
diagnostic(event, metadata);
} catch {
// Diagnostics are deliberately non-authoritative and cannot interrupt
// validation, DNS pinning, downloading, or byte verification.
}
}
function extensionOf(fileName: string): string {
const lower = fileName.toLowerCase();
if (lower.endsWith('.xlsx')) return '.xlsx';
if (lower.endsWith('.xls')) return '.xls';
return '';
}
export function normalizeInputAttachmentFileName(value: unknown): string {
const normalized = basename(String(value ?? '').replace(/[\u0000-\u001f\u007f]/g, '_'))
.trim()
.slice(0, 200);
if (!normalized || !ALLOWED_ROSTER_EXTENSIONS.has(extensionOf(normalized))) {
throw new InputAttachmentError('roster_file_type_unsupported', '名单附件必须是 .xls 或 .xlsx 文件。');
}
return normalized;
}
function normalizeContentType(value: unknown): string {
return String(value ?? '')
.replace(/[\r\n]/g, '')
.trim()
.slice(0, 200) || 'application/octet-stream';
}
function normalizeDeclaredSize(value: unknown, maxBytes: number): number | undefined {
if (value === undefined || value === null || value === '') return undefined;
const size = Number(value);
if (!Number.isInteger(size) || size < 1) {
throw new InputAttachmentError('roster_file_size_invalid', '名单附件大小声明无效。');
}
if (size > maxBytes) {
throw new InputAttachmentError('roster_file_too_large', '名单附件超过大小限制。', { max_bytes: maxBytes });
}
return size;
}
function normalizeDeclaredSha256(value: unknown): string | undefined {
const normalized = String(value ?? '').trim().toLowerCase();
if (!normalized) return undefined;
if (!SHA256_PATTERN.test(normalized)) {
throw new InputAttachmentError('roster_file_sha256_invalid', '名单附件 SHA-256 声明无效。');
}
return normalized;
}
function validateAttachmentBytes(
content: Buffer,
maxBytes: number,
declaredSize?: number,
declaredSha256?: string
): void {
if (!content.byteLength) throw new InputAttachmentError('roster_file_empty', '名单附件为空。');
if (content.byteLength > maxBytes) {
throw new InputAttachmentError('roster_file_too_large', '名单附件超过大小限制。', { max_bytes: maxBytes });
}
if (declaredSize !== undefined && content.byteLength !== declaredSize) {
throw new InputAttachmentError('roster_file_size_mismatch', '名单附件实际大小与声明不一致。', {
declared_size: declaredSize,
actual_size: content.byteLength
});
}
const digest = sha256Bytes(content);
if (declaredSha256 && digest !== declaredSha256) {
throw new InputAttachmentError('roster_file_sha256_mismatch', '名单附件 SHA-256 校验失败。');
}
}
export function decodeInlineInputAttachment(
value: EncodedTaskInputAttachment,
maxBytes: number,
source: TaskInputAttachmentInput['source'] = 'manual'
): TaskInputAttachmentInput {
const fileName = normalizeInputAttachmentFileName(value.name);
const declaredSize = normalizeDeclaredSize(value.size, maxBytes);
const declaredSha256 = normalizeDeclaredSha256(value.sha256);
const encoded = String(value.content_base64 ?? '').replace(/\s+/g, '');
if (!encoded || encoded.length % 4 === 1 || !BASE64_PATTERN.test(encoded)) {
throw new InputAttachmentError('roster_file_base64_invalid', '名单附件内容不是合法 Base64。');
}
const normalizedBase64 = encoded.replace(/-/g, '+').replace(/_/g, '/');
const content = Buffer.from(normalizedBase64, 'base64');
validateAttachmentBytes(content, maxBytes, declaredSize, declaredSha256);
return {
fileName,
contentType: normalizeContentType(value.content_type),
content,
declaredSize,
declaredSha256,
source
};
}
export function validateAgentBusAttachmentUrl(value: unknown): URL {
let url: URL;
try {
url = new URL(String(value ?? ''));
} catch {
throw new InputAttachmentError('roster_attachment_url_invalid', '名单附件 URL 无效。');
}
if (url.protocol !== 'https:' || url.username || url.password) {
throw new InputAttachmentError('roster_attachment_url_unsafe', '名单附件必须使用不含用户名密码的 HTTPS URL。');
}
if (!url.hostname.replace(/^\[|\]$/g, '')) {
throw new InputAttachmentError('roster_attachment_url_invalid', '名单附件 URL 无效。');
}
return url;
}
export function parseAgentBusInputAttachment(
value: unknown,
maxBytes: number
): AgentBusInputAttachmentReference {
if (!value || typeof value !== 'object' || Array.isArray(value)) {
throw new InputAttachmentError('roster_attachment_metadata_invalid', '名单附件元数据无效。');
}
const record = value as Record<string, unknown>;
return {
name: normalizeInputAttachmentFileName(record.name),
contentType: normalizeContentType(record.content_type),
size: normalizeDeclaredSize(record.size, maxBytes),
sha256: normalizeDeclaredSha256(record.sha256),
url: validateAgentBusAttachmentUrl(record.url).toString()
};
}
export async function resolveAgentBusAttachmentAddresses(
hostname: string
): Promise<Array<{ address: string; family: number }>> {
if (isIP(hostname)) return [{ address: hostname, family: isIP(hostname) }];
let addresses: Array<{ address: string; family: number }>;
try {
addresses = await dnsLookup(hostname, { all: true, verbatim: true });
} catch {
throw new InputAttachmentError('roster_attachment_host_unresolved', '名单附件地址无法解析。');
}
if (!addresses.length) {
throw new InputAttachmentError('roster_attachment_host_unresolved', '名单附件地址无法解析。');
}
return addresses;
}
export function createPinnedAttachmentLookup(
address: { address: string; family: number }
): LookupFunction {
return (_hostname, options, callback) => {
if (options.all) {
callback(null, [address]);
return;
}
callback(null, address.address, address.family);
};
}
async function downloadPinnedHttps(
url: URL,
address: { address: string; family: number },
maxBytes: number
): Promise<{ statusCode: number; location?: string; contentType: string; content: Buffer }> {
return new Promise((resolve, reject) => {
const request = httpsRequest(url, {
method: 'GET',
headers: { Accept: 'application/vnd.ms-excel, application/vnd.openxmlformats-officedocument.spreadsheetml.sheet, application/octet-stream' },
timeout: DOWNLOAD_TIMEOUT_MS,
lookup: createPinnedAttachmentLookup(address)
}, (response) => {
const statusCode = Number(response.statusCode || 0);
const location = Array.isArray(response.headers.location) ? response.headers.location[0] : response.headers.location;
if (statusCode >= 300 && statusCode < 400) {
response.resume();
resolve({ statusCode, location, contentType: '', content: Buffer.alloc(0) });
return;
}
if (statusCode !== 200) {
response.resume();
reject(new InputAttachmentError('roster_attachment_download_failed', '名单附件下载失败。', { status_code: statusCode }));
return;
}
const declaredLength = Number(response.headers['content-length']);
if (Number.isFinite(declaredLength) && declaredLength > maxBytes) {
response.destroy();
reject(new InputAttachmentError('roster_file_too_large', '名单附件超过大小限制。', { max_bytes: maxBytes }));
return;
}
const chunks: Buffer[] = [];
let total = 0;
response.on('data', (chunk: Buffer | string) => {
const buffer = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk);
total += buffer.byteLength;
if (total > maxBytes) {
response.destroy(new InputAttachmentError('roster_file_too_large', '名单附件超过大小限制。', { max_bytes: maxBytes }));
return;
}
chunks.push(buffer);
});
response.once('end', () => resolve({
statusCode,
contentType: normalizeContentType(response.headers['content-type']),
content: Buffer.concat(chunks, total)
}));
response.once('error', reject);
});
request.once('timeout', () => request.destroy(new InputAttachmentError('roster_attachment_download_timeout', '名单附件下载超时。')));
request.once('error', reject);
request.end();
});
}
export async function downloadAgentBusInputAttachment(
reference: AgentBusInputAttachmentReference,
maxBytes: number,
diagnostic?: InputAttachmentDiagnostic
): Promise<TaskInputAttachmentInput> {
const startedAt = process.hrtime.bigint();
let redirectCount = 0;
emitAttachmentDiagnostic(diagnostic, 'download_started', {
file_extension: extensionOf(reference.name),
declared_size: reference.size ?? null,
declared_sha256_present: Boolean(reference.sha256),
max_bytes: maxBytes
});
try {
let url = validateAgentBusAttachmentUrl(reference.url);
for (let redirects = 0; redirects <= MAX_REDIRECTS; redirects += 1) {
redirectCount = redirects;
const dnsStartedAt = process.hrtime.bigint();
emitAttachmentDiagnostic(diagnostic, 'dns_started', { redirect_count: redirects });
const addresses = await resolveAgentBusAttachmentAddresses(url.hostname.replace(/^\[|\]$/g, ''));
emitAttachmentDiagnostic(diagnostic, 'dns_validated', {
redirect_count: redirects,
address_count: addresses.length,
address_families: [...new Set(addresses.map((address) => address.family))].sort(),
selected_address_family: addresses[0]?.family || null,
duration_ms: diagnosticDurationMs(dnsStartedAt)
});
const requestStartedAt = process.hrtime.bigint();
const response = await downloadPinnedHttps(url, addresses[0], maxBytes);
emitAttachmentDiagnostic(diagnostic, 'http_response', {
redirect_count: redirects,
status_code: response.statusCode,
response_bytes: response.content.byteLength,
duration_ms: diagnosticDurationMs(requestStartedAt)
});
if (response.statusCode >= 300 && response.statusCode < 400) {
if (!response.location || redirects === MAX_REDIRECTS) {
throw new InputAttachmentError('roster_attachment_redirect_invalid', '名单附件重定向无效或次数过多。');
}
emitAttachmentDiagnostic(diagnostic, 'redirect_followed', {
redirect_count: redirects + 1,
status_code: response.statusCode
});
url = validateAgentBusAttachmentUrl(new URL(response.location, url).toString());
continue;
}
validateAttachmentBytes(response.content, maxBytes, reference.size, reference.sha256);
emitAttachmentDiagnostic(diagnostic, 'download_completed', {
redirect_count: redirects,
byte_size: response.content.byteLength,
declared_size_match: reference.size === undefined || reference.size === response.content.byteLength,
declared_sha256_verified: Boolean(reference.sha256),
duration_ms: diagnosticDurationMs(startedAt)
});
return {
fileName: reference.name,
contentType: response.contentType === 'application/octet-stream' ? reference.contentType : response.contentType,
content: response.content,
declaredSize: reference.size,
declaredSha256: reference.sha256,
source: 'agentbus'
};
}
throw new InputAttachmentError('roster_attachment_download_failed', '名单附件下载失败。');
} catch (error) {
emitAttachmentDiagnostic(diagnostic, 'download_failed', {
redirect_count: redirectCount,
duration_ms: diagnosticDurationMs(startedAt),
error_code: error instanceof InputAttachmentError
? error.code
: 'roster_attachment_download_exception'
});
throw error;
}
}