47 lines
1.1 KiB
TypeScript
47 lines
1.1 KiB
TypeScript
export default defineEventHandler(async (event) => {
|
|
const id = getRouterParam(event, 'id')
|
|
let job = id ? getJob(id) : undefined
|
|
if (!job && id) {
|
|
const pending = readPendingJob(id)
|
|
if (pending) job = ensurePendingWatch(pending)
|
|
}
|
|
if (!job) {
|
|
throw createError({ statusCode: 404, statusMessage: 'Job not found' })
|
|
}
|
|
|
|
setResponseHeaders(event, {
|
|
'Cache-Control': 'no-cache, no-store, no-transform',
|
|
'X-Accel-Buffering': 'no',
|
|
Connection: 'keep-alive'
|
|
})
|
|
|
|
const stream = createEventStream(event)
|
|
const send = async (payload: unknown) => {
|
|
await stream.push(JSON.stringify(payload))
|
|
}
|
|
|
|
await send(jobSnapshot(job))
|
|
for (const past of job.events) {
|
|
await send(past)
|
|
}
|
|
|
|
const unsubscribe = subscribeJob(job, (payload) => {
|
|
void send(payload).then(() => {
|
|
if (payload.type === 'complete' || payload.type === 'error') {
|
|
void stream.close()
|
|
}
|
|
})
|
|
})
|
|
|
|
const ping = setInterval(() => {
|
|
void send(jobSnapshot(job))
|
|
}, 1000)
|
|
|
|
stream.onClosed(() => {
|
|
clearInterval(ping)
|
|
unsubscribe()
|
|
})
|
|
|
|
return stream.send()
|
|
})
|