decision-studio / static /queue-client.js
Xunzhuo's picture
Align Studio and Tetris with HTTP queue results
d58afdb
Raw History Blame Contribute Delete
2.2 kB
// Studio's bounded queue transport. A completed job is validated before use.
export async function requestJSON(url, options = {}) {
const response = await fetch(url, options);
const body = await response.json();
if (!response.ok) throw Error(typeof body.detail === 'string' ? body.detail : 'The server rejected this request.');
return body;
}
function pollPause(signal) {
return new Promise((resolve, reject) => {
const cancel = () => { clearTimeout(timer); reject(new DOMException('Detached', 'AbortError')); };
const timer = setTimeout(() => { signal.removeEventListener('abort', cancel); resolve(); }, 600);
signal.addEventListener('abort', cancel, {once:true});
if(signal.aborted) cancel();
});
}
const jobURL = (id, model) => '/api/jobs/' + encodeURIComponent(id) + '?model=' + encodeURIComponent(model);
export async function queuedPrediction(payload, current, info, {onProgress, validateCompletion, pause=pollPause}) {
let id;
const validJob = job => {if(job.model !== payload.model || job.manifest_sha256 !== info.manifest_sha256) throw Error('The job belongs to a different model.');};
try {
// Keep only the acceptance response alive on switch so its job can be
// cancelled by ID. A detached prediction is never rendered.
const accepted = await requestJSON('/api/jobs', {method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify(payload), signal:AbortSignal.timeout(45000)});
validJob(accepted); id = accepted.id;
if(current.signal.aborted) throw new DOMException('Detached', 'AbortError');
for (;;) {
const job = await requestJSON(jobURL(id, payload.model), {signal:current.signal, cache:'no-store'});
validJob(job);
if (job.status === 'succeeded') {
validateCompletion(job.result);
return job.result;
}
if (!['queued','running'].includes(job.status)) throw Error(job.detail || 'The request did not finish.');
if (!current.signal.aborted) onProgress(job.status);
await pause(current.signal);
}
} finally {
if (id && current.signal.aborted) fetch(jobURL(id, payload.model), {method:'DELETE', keepalive:true}).catch(() => {});
}
}