File size: 2,195 Bytes
d58afdb
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
// 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(() => {});
  }
}