mirror of
https://github.com/THU-MAIC/OpenMAIC.git
synced 2026-10-02 01:15:18 +08:00
* feat(render): add admission and per-task resource budgets * fix(ci): isolate patched producer test configuration * fix(render): bind portable resource builds and installed validation * fix(render): preserve deadline errors and service exit ownership * style(render): format verified lifecycle changes * fix(render): retain cleanup evidence and enforce task pid limits * refactor(render): narrow resource patch and acceptance documentation Revert cleanup adapters in two unreferenced Producer helpers. Remove historical runtime claims from the patched README, retain the current Linux NOT_RUN qualification, and map remaining B4 evidence to concrete owner and platform boundaries. Validation: fixed-source patch application and all 44 source hashes pass; 43 retained patch sections and both dependency locks are unchanged. Source/package tests: 17 passed. Cleanup/context/supervisor tests: 30 passed. No native build, VM, installed Linux execution, or export was run. * test(render): cover remaining B4 lifecycle acceptance paths Extend the existing installed runner with guardian unlink failure, overlapping task reference settlement, and successful new-supervisor service after explicit platform takeover. Keep original per-task and outer limits; only the overlap case reserves two task budgets through the existing Producer API. Local validation: syntax, formatting, lint, two evidence regressions, and fixed-source patch/hash verification pass. Installed Linux cases remain NOT_RUN: the first SSH preflight timed out before authentication or any remote command ran. No runtime implementation or dependency version changed. * docs(render): keep runtime qualification tied to candidate evidence Keep the validation document as a reproducible contract rather than a stale NOT_RUN snapshot. Current native and Producer/B4 evidence passed for 6be78252; HTTP remains unexecuted after an adapter input-preflight failure. Runtime status and remaining work are tracked in the PR. * fix(render): harden resource startup and simplify owner transport * docs(render): clarify post-commit cleanup settlement * style(render): format outer resource fence files * fix(render-service): harden resource-mode closeout * fix(render): bind resource cleanup to trusted directory identity Use nonrecursive stage cleanup and preserve quarantine when project identity cannot be verified. Bind mount, publication and cleanup to verified directory descriptors, contain recovery bookkeeping errors, and cover helper cancellation and closeout gates. --------- Co-authored-by: wyuc <wang-yc24@mails.tsinghua.edu.cn>
237 lines
8.6 KiB
TypeScript
237 lines
8.6 KiB
TypeScript
import { access, mkdtemp, rm } from 'node:fs/promises';
|
|
import { tmpdir } from 'node:os';
|
|
import { join } from 'node:path';
|
|
import { afterEach, describe, expect, it } from 'vitest';
|
|
import type { JobStore } from '../src/job-store.js';
|
|
import { RenderCoordinator } from '../src/render-coordinator.js';
|
|
import type { RenderExecutor } from '../src/render-executor.js';
|
|
import type {
|
|
RenderExecutionRequest,
|
|
RenderExecutionResult,
|
|
RenderJobRecord,
|
|
} from '../src/types.js';
|
|
import { createMemoryArtifactStore, createMemoryJobStore } from './support/fakes.js';
|
|
|
|
const scratch: string[] = [];
|
|
|
|
afterEach(async () => {
|
|
await Promise.all(scratch.splice(0).map((path) => rm(path, { recursive: true, force: true })));
|
|
});
|
|
|
|
class FakeExecutor implements RenderExecutor {
|
|
readonly requests: RenderExecutionRequest[] = [];
|
|
|
|
constructor(
|
|
private readonly handler: (request: RenderExecutionRequest) => Promise<RenderExecutionResult>,
|
|
) {}
|
|
|
|
async execute(request: RenderExecutionRequest): Promise<RenderExecutionResult> {
|
|
this.requests.push(request);
|
|
return this.handler(request);
|
|
}
|
|
}
|
|
|
|
async function projectDir(): Promise<string> {
|
|
const path = await mkdtemp(join(tmpdir(), 'render-coordinator-'));
|
|
scratch.push(path);
|
|
return path;
|
|
}
|
|
|
|
async function waitForJob(
|
|
jobs: JobStore,
|
|
id: string,
|
|
predicate: (job: RenderJobRecord) => boolean,
|
|
): Promise<RenderJobRecord> {
|
|
for (let attempt = 0; attempt < 100; attempt += 1) {
|
|
const job = await jobs.get(id);
|
|
if (job && predicate(job)) return job;
|
|
await new Promise((resolve) => setTimeout(resolve, 5));
|
|
}
|
|
throw new Error(`Timed out waiting for job ${id}`);
|
|
}
|
|
|
|
/**
|
|
* Project cleanup runs after the job reaches its terminal status, so a bare
|
|
* `access` right after `waitForJob` races the removal.
|
|
*/
|
|
async function waitForCleanup(path: string): Promise<void> {
|
|
for (let attempt = 0; attempt < 100; attempt += 1) {
|
|
try {
|
|
await access(path);
|
|
} catch {
|
|
return;
|
|
}
|
|
await new Promise((resolve) => setTimeout(resolve, 5));
|
|
}
|
|
throw new Error(`Timed out waiting for ${path} to be removed`);
|
|
}
|
|
|
|
const renderOptions = { fps: 30, quality: 'standard', format: 'mp4' } as const;
|
|
|
|
describe('RenderCoordinator through the RenderExecutor seam', () => {
|
|
it('keeps a dispatched video queued until the shared execution slot is acquired', async () => {
|
|
const jobs = createMemoryJobStore();
|
|
const artifacts = createMemoryArtifactStore();
|
|
let releasePreview!: () => void;
|
|
const previewParked = new Promise<void>((resolve) => {
|
|
releasePreview = resolve;
|
|
});
|
|
let finishVideo!: () => void;
|
|
const videoParked = new Promise<void>((resolve) => {
|
|
finishVideo = resolve;
|
|
});
|
|
const executor = new FakeExecutor(async () => {
|
|
await videoParked;
|
|
return { status: 'succeeded' };
|
|
});
|
|
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store, {
|
|
maxConcurrency: 1,
|
|
});
|
|
const preview = coordinator.tryRunWithExecutionSlot(() => previewParked);
|
|
expect(preview).toBeDefined();
|
|
|
|
const dir = await projectDir();
|
|
const id = await coordinator.submit(coordinator.reserve('video-user'), dir, renderOptions);
|
|
await Promise.resolve();
|
|
expect(await jobs.get(id)).toMatchObject({ status: 'queued', currentStage: 'queued' });
|
|
expect(executor.requests).toHaveLength(0);
|
|
|
|
releasePreview();
|
|
await preview;
|
|
await waitForJob(jobs, id, () => executor.requests.length === 1);
|
|
expect(await jobs.get(id)).toMatchObject({ status: 'running', currentStage: 'preparing' });
|
|
finishVideo();
|
|
await waitForJob(jobs, id, (job) => job.status === 'succeeded');
|
|
});
|
|
|
|
it('persists normalized progress, performance, and the artifact on success', async () => {
|
|
const jobs = createMemoryJobStore();
|
|
const artifacts = createMemoryArtifactStore();
|
|
const performance = {
|
|
totalElapsedMs: 800,
|
|
stages: { captureMs: 600 },
|
|
workers: 1,
|
|
totalFrames: 30,
|
|
captureMode: 'beginframe',
|
|
};
|
|
const executor = new FakeExecutor(async (request) => {
|
|
await request.onProgress({
|
|
progress: 0.5,
|
|
stage: 'capturing',
|
|
framesRendered: 15,
|
|
totalFrames: 30,
|
|
});
|
|
return { status: 'succeeded', performance };
|
|
});
|
|
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store, {
|
|
jobDeadlineMs: 12_345,
|
|
});
|
|
const dir = await projectDir();
|
|
const id = await coordinator.submit(coordinator.reserve('alice'), dir, renderOptions);
|
|
|
|
const job = await waitForJob(jobs, id, (current) => current.status === 'succeeded');
|
|
expect(executor.requests).toHaveLength(1);
|
|
expect(executor.requests[0].deadlineMs).toBe(12_345);
|
|
expect(job).toMatchObject({
|
|
status: 'succeeded',
|
|
progress: 1,
|
|
currentStage: 'complete',
|
|
framesRendered: 15,
|
|
totalFrames: 30,
|
|
performance,
|
|
});
|
|
expect(artifacts.paths.get(id)).toBe(join(dir, 'output.mp4'));
|
|
});
|
|
|
|
it('routes running-job cancellation through the executor signal and cleans up', async () => {
|
|
const jobs = createMemoryJobStore();
|
|
const artifacts = createMemoryArtifactStore();
|
|
const executor = new FakeExecutor(
|
|
(request) =>
|
|
new Promise((resolve) => {
|
|
request.signal.addEventListener('abort', () => {
|
|
resolve({
|
|
status: 'cancelled',
|
|
failure: { code: 'cancelled', message: 'Render cancelled' },
|
|
});
|
|
});
|
|
}),
|
|
);
|
|
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store);
|
|
const dir = await projectDir();
|
|
const id = await coordinator.submit(coordinator.reserve('bob'), dir, renderOptions);
|
|
await waitForJob(jobs, id, () => executor.requests.length === 1);
|
|
|
|
expect(await coordinator.cancel(id)).toBe(true);
|
|
const job = await waitForJob(jobs, id, (current) => current.status === 'cancelled');
|
|
expect(job.failure).toEqual({ code: 'cancelled', message: 'Render cancelled' });
|
|
await waitForCleanup(dir);
|
|
});
|
|
|
|
it('keeps default-executor late success cancellable when no publication was committed', async () => {
|
|
const jobs = createMemoryJobStore();
|
|
const artifacts = createMemoryArtifactStore();
|
|
let settle!: (result: RenderExecutionResult) => void;
|
|
const executor = new FakeExecutor(
|
|
() =>
|
|
new Promise((resolve) => {
|
|
settle = resolve;
|
|
}),
|
|
);
|
|
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store);
|
|
const dir = await projectDir();
|
|
const id = await coordinator.submit(coordinator.reserve('late-default'), dir, renderOptions);
|
|
await waitForJob(jobs, id, () => executor.requests.length === 1);
|
|
|
|
expect(await coordinator.cancel(id)).toBe(true);
|
|
settle({ status: 'succeeded' });
|
|
|
|
const job = await waitForJob(jobs, id, (current) => current.status === 'cancelled');
|
|
expect(job.failure).toEqual({ code: 'cancelled', message: 'Render cancelled' });
|
|
expect(artifacts.paths.has(id)).toBe(false);
|
|
await waitForCleanup(dir);
|
|
});
|
|
|
|
it('keeps deadline failure classification from a replaceable executor', async () => {
|
|
const jobs = createMemoryJobStore();
|
|
const artifacts = createMemoryArtifactStore();
|
|
const executor = new FakeExecutor(async () => ({
|
|
status: 'failed',
|
|
failure: { code: 'deadline_exceeded', message: 'Render exceeded the deadline' },
|
|
}));
|
|
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store, {
|
|
jobDeadlineMs: 42,
|
|
});
|
|
const dir = await projectDir();
|
|
const id = await coordinator.submit(coordinator.reserve('carol'), dir, renderOptions);
|
|
|
|
const job = await waitForJob(jobs, id, (current) => current.status === 'failed');
|
|
expect(executor.requests[0].deadlineMs).toBe(42);
|
|
expect(job).toMatchObject({
|
|
status: 'failed',
|
|
error: 'Render exceeded the deadline',
|
|
failure: { code: 'deadline_exceeded' },
|
|
});
|
|
expect(artifacts.paths.has(id)).toBe(false);
|
|
await waitForCleanup(dir);
|
|
});
|
|
|
|
it('classifies unexpected executor errors and still performs cleanup', async () => {
|
|
const jobs = createMemoryJobStore();
|
|
const artifacts = createMemoryArtifactStore();
|
|
const executor = new FakeExecutor(async () => {
|
|
throw new Error('executor unavailable');
|
|
});
|
|
const coordinator = new RenderCoordinator(executor, jobs, artifacts.store);
|
|
const dir = await projectDir();
|
|
const id = await coordinator.submit(coordinator.reserve('dana'), dir, renderOptions);
|
|
|
|
const job = await waitForJob(jobs, id, (current) => current.status === 'failed');
|
|
expect(job).toMatchObject({
|
|
error: 'executor unavailable',
|
|
failure: { code: 'execution_failed', message: 'executor unavailable' },
|
|
});
|
|
await waitForCleanup(dir);
|
|
});
|
|
});
|