From b1ffcd6975faec8ad0d45d0ffc8a20afef686058 Mon Sep 17 00:00:00 2001 From: Towsty Date: Fri, 11 Sep 2026 18:43:25 -0500 Subject: [PATCH] Retry confirmed Studio 2 host cleanup after failures --- scripts/studio2-purge.mjs | 7 +++- server/plugins/00-resume-studio2.ts | 3 +- server/utils/studio2/cleanup.ts | 50 +++++++++++++++++++++++++++++ server/utils/studio2/runner.ts | 18 +---------- tests/studio2-cleanup.test.mjs | 15 +++++++++ 5 files changed, 74 insertions(+), 19 deletions(-) create mode 100644 server/utils/studio2/cleanup.ts create mode 100644 tests/studio2-cleanup.test.mjs diff --git a/scripts/studio2-purge.mjs b/scripts/studio2-purge.mjs index 6d1b0c1..6e172f8 100644 --- a/scripts/studio2-purge.mjs +++ b/scripts/studio2-purge.mjs @@ -1,6 +1,11 @@ import { resolve, relative, isAbsolute } from 'node:path' import { existsSync, realpathSync, statSync, unlinkSync } from 'node:fs' -import { scopedFile } from '../shared/studio2/contracts.mjs' +// Keep the host cleanup endpoint independently deployable from the web application. +function scopedFile(file, prefix) { + const path=[file.subfolder,file.filename].map(s=>String(s||'').replaceAll('\\','/')).filter(Boolean).join('/') + const p=String(prefix||'').replaceAll('\\','/').replace(/\/$/,'') + return !!p && !path.startsWith('/') && !path.split('/').some(s=>s==='..'||s==='.') && path.startsWith(p+'/') +} /** Exact Studio 2 files only: no folder fallback, sibling deletion, or sweeping. */ export function purgeStudio2Files({prefix, files}, rootsForType) { diff --git a/server/plugins/00-resume-studio2.ts b/server/plugins/00-resume-studio2.ts index 979a314..0dc8522 100644 --- a/server/plugins/00-resume-studio2.ts +++ b/server/plugins/00-resume-studio2.ts @@ -1,2 +1,3 @@ +import { retryCleanup } from '../utils/studio2/cleanup'; import { resumeStudio2Jobs } from '../utils/studio2/runner'; -export default defineNitroPlugin(() => { resumeStudio2Jobs(); }); +export default defineNitroPlugin(() => { resumeStudio2Jobs(); const timer = setInterval(() => { void retryCleanup().catch(error => console.warn('[Studio 2 cleanup retry]', error.message)); }, 30000); timer.unref(); }); diff --git a/server/utils/studio2/cleanup.ts b/server/utils/studio2/cleanup.ts new file mode 100644 index 0000000..8daeb1b --- /dev/null +++ b/server/utils/studio2/cleanup.ts @@ -0,0 +1,50 @@ +import { scopedFile } from '~/shared/studio2/contracts.mjs'; +import { sharedGpuHeaders, withSharedGpuStart } from '../sharedGpu'; +import { records, saveRecord } from './store'; + +/** Called only after the output and handoff frame are durably saved. */ +export async function purge(r: any, includeCurrent = true) { + const config = useRuntimeConfig(), prefix = String(config.comfyFilenamePrefix); + const pending = r.cleanupPending || []; + const files = [...pending, ...(includeCurrent ? r.files || [] : [])].filter((file, i, all) => all.findIndex(f => f.filename === file.filename && f.subfolder === file.subfolder && f.type === file.type) === i); + r.cleanupPending = files; + r.purgeResult = 'Cleanup pending'; + saveRecord(r); + if (!files.length) return; + try { + if (files.some(file => !scopedFile(file, prefix))) throw new Error('Cleanup manifest is outside this studio'); + for (let offset = 0; offset < files.length; offset += 100) { + const batch = files.slice(offset, offset + 100); + const response = await fetch(String(config.comfyControlUrl).replace(/\/$/, '') + '/studio2/purge', { + method: 'POST', headers: { 'Content-Type': 'application/json', ...sharedGpuHeaders() }, + body: JSON.stringify({ prefix, files: batch }), signal: AbortSignal.timeout(12000) + }); + const result = await response.json() as any; + if (!response.ok || !result.ok || result.cleared !== batch.length) throw new Error(`Local cleanup was not confirmed (${response.status})`); + r.cleanupPending = files.slice(offset + batch.length); + saveRecord(r); + } + r.purgeResult = 'Cleared input + output'; + r.cleanupError = ''; + } catch (error: any) { + r.cleanupError = error.message || 'Local cleanup unavailable'; + console.warn('[Studio 2 cleanup]', r.id, r.cleanupError); + } + saveRecord(r); +} + +let retrying = false; +export async function retryCleanup() { + if (retrying) return; + retrying = true; + try { + const pending = records().filter(r => r.cleanupPending?.length && ['complete', 'failed', 'cancelled'].includes(r.state)); + if (!pending.length) return; + await withSharedGpuStart(async () => { + for (const r of pending) { + // Never add a later failed shot's unsaved files to the saved cleanup manifest. + await purge(r, false); + } + }, async () => {}); + } finally { retrying = false; } +} diff --git a/server/utils/studio2/runner.ts b/server/utils/studio2/runner.ts index 6c453ce..37cd235 100644 --- a/server/utils/studio2/runner.ts +++ b/server/utils/studio2/runner.ts @@ -1,3 +1,4 @@ +import { purge } from './cleanup'; import { watchProgress } from './progress'; import { fitStill } from './media'; import { resolveRequestSize } from './size'; @@ -46,23 +47,6 @@ async function upload(r: any, name: string, data: Buffer): Promise { saveRecord(r); return `${disk.subfolder}/${disk.filename}`; } -async function purge(r: any) { - const config = useRuntimeConfig(), prefix = String(config.comfyFilenamePrefix); - const files = (r.files as DiskFile[]).filter(file => scopedFile(file, prefix)); - r.purgeResult = 'Left on host'; - if (files.length !== r.files.length || !files.length) - return; - try { - const response = await fetch(String(config.comfyControlUrl).replace(/\/$/, '') + '/studio2/purge', { - method: 'POST', headers: { 'Content-Type': 'application/json', ...sharedGpuHeaders() }, - body: JSON.stringify({ prefix, files }), signal: AbortSignal.timeout(12000) - }); - const result = await response.json() as any; - if (response.ok && result.ok && result.cleared === files.length) - r.purgeResult = 'Cleared input + output'; - } - catch { /* Keep the durable library copy and report unconfirmed host cleanup. */ } -} async function prepareGraph(r: any) { const q = r.request, s = q.settings, prefix = String(useRuntimeConfig().comfyFilenamePrefix).replace(/\/$/, '') + `/studio2/${r.id}/${r.index}`; const load = async (id: string, label: string) => id ? upload(r, label, readFileSync(stillPath(r.owner, getStill(r.owner, id).id))) : ''; diff --git a/tests/studio2-cleanup.test.mjs b/tests/studio2-cleanup.test.mjs new file mode 100644 index 0000000..cfe9edc --- /dev/null +++ b/tests/studio2-cleanup.test.mjs @@ -0,0 +1,15 @@ +import test from 'node:test'; +import assert from 'node:assert/strict'; +import {readFileSync} from 'node:fs'; +import ts from 'typescript'; +import {scopedFile} from '../shared/studio2/contracts.mjs'; +const source=readFileSync(new URL('../server/utils/studio2/cleanup.ts',import.meta.url),'utf8').replace(/^import .*;\r?\n/gm,'').replace(/export /g,''); +function load(scope){const js=ts.transpileModule(source,{compilerOptions:{target:ts.ScriptTarget.ES2022}}).outputText;return new Function(...Object.keys(scope),js+';return {purge,retryCleanup}')( ...Object.values(scope));} +test('cleanup persists failures and retries only saved manifests after restart',async()=>{ + const saved=[];let online=false,calls=[]; + const file={filename:'hero.png',subfolder:'video/dev/studio2/job/0',type:'input'}; + const r={id:'job',state:'complete',files:[file]}; + const scope={scopedFile,useRuntimeConfig:()=>({comfyFilenamePrefix:'video/dev',comfyControlUrl:'http://local'}),sharedGpuHeaders:()=>({}),saveRecord:r=>saved.push(structuredClone(r)),records:()=>[r],withSharedGpuStart:async fn=>fn(),console:{warn(){}},fetch:async(url,opts)=>{calls.push(JSON.parse(opts.body));return {ok:online,status:online?200:404,json:async()=>({ok:online,cleared:1})}}}; + await load(scope).purge(r);assert.equal(r.cleanupPending.length,1);assert.match(r.cleanupError,/404/);assert.equal(saved[0].cleanupPending.length,1); + r.files=[{...file,filename:'unsaved.png'}];online=true;await load(scope).retryCleanup();assert.deepEqual(calls[1].files,[file]);assert.deepEqual(r.cleanupPending,[]);assert.equal(r.files[0].filename,'unsaved.png'); +});