@@ -32,11 +32,13 @@ type EventStreamBody = NodeReadableLike | ReadableStream<Uint8Array> | string;
3232const SSE_RETRY_INTERVAL = 5000 ;
3333const SSE_MAX_RETRIES = 100 ;
3434/**
35- * 等待任务时,状态接口连续几次回的不是任务(空响应体、没有 status)就报错。实测生产环境对不属于
36- * 当前身份的任务回 200 与空响应体;原先把它当成「还在跑」,白等到超时(600 秒),最后只剩一句
37- * 「The operation was aborted due to timeout」。留几次余地,偶发的一次空响应不至于打断等待。
35+ * 等待任务时,状态接口回的不是任务(空响应体、没有 status)最多再等这么久(秒)就报错。
36+ *
37+ * 实测生产环境对不属于当前身份的任务一直回 200 与空响应体,原先当成「还在跑」白等到超时(600 秒)。
38+ * 但刚建好的任务也会短暂地回空:2.5.1 连续三次(约 5 秒)就报错,测试环境上一个正常的任务因此
39+ * 失败。按时间留余地,刚建好时的空响应等得过去,一直回空的仍在半分钟内报错。
3840 */
39- const INVALID_TASK_RESPONSE_LIMIT = 3 ;
41+ const EMPTY_STATUS_GRACE_SECONDS = 30 ;
4042
4143/** The endpoint returned a snapshot, not a stream; polling can continue the wait. */
4244class EventStreamUnavailableError extends Error { }
@@ -272,7 +274,8 @@ export class TasksApi {
272274 // Fast path: task may already be terminal before SSE connects (common for
273275 // quick guest-mode tasks after an explicit start). Avoid hanging on an SSE
274276 // connection that never emits a terminal event for an already-finished task.
275- const current = await this . getSnapshot < T > ( taskId , options . signal , options . spaceId , options . pollInterval ?? DEFAULT_POLL_INTERVAL ) ;
277+ const grace = options . emptyStatusGrace ?? EMPTY_STATUS_GRACE_SECONDS ;
278+ const current = await this . getSnapshot < T > ( taskId , options . signal , options . spaceId , options . pollInterval ?? DEFAULT_POLL_INTERVAL , grace ) ;
276279 options . onProgress ?.( current ) ;
277280 if ( current . status === 'completed' ) {
278281 return current ;
@@ -285,7 +288,7 @@ export class TasksApi {
285288 if ( options . useEventStream !== false ) {
286289 return await this . waitWithEventStream < T > (
287290 taskId , remainingTimeout , options . onProgress , options . signal ,
288- options . pollInterval ?? DEFAULT_POLL_INTERVAL , options . spaceId
291+ options . pollInterval ?? DEFAULT_POLL_INTERVAL , options . spaceId , grace
289292 ) ;
290293 }
291294 return await this . waitWithPolling < T > (
@@ -294,7 +297,8 @@ export class TasksApi {
294297 options . pollInterval ?? DEFAULT_POLL_INTERVAL ,
295298 options . onProgress ,
296299 options . signal ,
297- options . spaceId
300+ options . spaceId ,
301+ grace
298302 ) ;
299303 }
300304
@@ -672,7 +676,8 @@ export class TasksApi {
672676 onProgress ?: ( task : DeckTask ) => void ,
673677 signal ?: AbortSignal ,
674678 pollInterval = DEFAULT_POLL_INTERVAL ,
675- spaceId ?: string
679+ spaceId ?: string ,
680+ emptyStatusGrace = EMPTY_STATUS_GRACE_SECONDS
676681 ) : Promise < DeckTask < T > > {
677682 return await new Promise < DeckTask < T > > ( ( resolve , reject ) => {
678683 throwIfAborted ( signal ) ;
@@ -706,7 +711,7 @@ export class TasksApi {
706711 fallbackStarted = true ;
707712 clearInterval ( timer ) ;
708713 cancel ?.( ) ;
709- this . waitWithPolling < T > ( taskId , remainingTimeout ( ) , pollInterval , onProgress , signal , spaceId )
714+ this . waitWithPolling < T > ( taskId , remainingTimeout ( ) , pollInterval , onProgress , signal , spaceId , emptyStatusGrace )
710715 . then ( ( task ) => {
711716 finish ( ( ) => resolve ( task ) ) ;
712717 } )
@@ -754,7 +759,8 @@ export class TasksApi {
754759 pollInterval : number ,
755760 onProgress ?: ( task : DeckTask ) => void ,
756761 signal ?: AbortSignal ,
757- spaceId ?: string
762+ spaceId ?: string ,
763+ emptyStatusGrace = EMPTY_STATUS_GRACE_SECONDS
758764 ) : Promise < DeckTask < T > > {
759765 const start = Date . now ( ) ;
760766 for ( ; ; ) {
@@ -763,7 +769,7 @@ export class TasksApi {
763769 throw new Error ( `Task ${ taskId } did not complete within ${ timeout } s` ) ;
764770 }
765771
766- const task = await this . getSnapshot < T > ( taskId , signal , spaceId , pollInterval ) ;
772+ const task = await this . getSnapshot < T > ( taskId , signal , spaceId , pollInterval , emptyStatusGrace ) ;
767773 onProgress ?.( task ) ;
768774
769775 if ( task . status === 'completed' ) {
@@ -777,13 +783,15 @@ export class TasksApi {
777783 }
778784 }
779785
780- /** 等待时取任务状态:回的不是任务就隔一个轮询间隔再取,连续 {@link INVALID_TASK_RESPONSE_LIMIT} 次就报错。 */
781- private async getSnapshot < T extends DeckTaskType > ( taskId : string , signal : AbortSignal | undefined , spaceId : string | undefined , pollInterval : number ) : Promise < DeckTask < T > > {
786+ /** 等待时取任务状态:回的不是任务就隔一个轮询间隔再取,超过宽限(秒)仍是空的就报错。 */
787+ private async getSnapshot < T extends DeckTaskType > ( taskId : string , signal : AbortSignal | undefined , spaceId : string | undefined , pollInterval : number , graceSeconds : number ) : Promise < DeckTask < T > > {
788+ const firstEmptyAt = Date . now ( ) ;
782789 for ( let attempt = 1 ; ; attempt += 1 ) {
783790 const task : unknown = await this . get < T > ( taskId , { signal, spaceId } ) ;
784791 if ( typeof task === 'object' && task !== null && typeof ( task as { status ?: unknown } ) . status === 'string' ) return task as DeckTask < T > ;
785- if ( attempt >= INVALID_TASK_RESPONSE_LIMIT ) {
786- throw new Error ( `The cloud returned no task status for ${ taskId } ${ attempt } times in a row (GET /tools/tasks/${ taskId } ); the task may belong to another account or space.` ) ;
792+ const waited = ( Date . now ( ) - firstEmptyAt ) / 1000 ;
793+ if ( waited >= graceSeconds ) {
794+ throw new Error ( `The cloud returned no task status for ${ taskId } for ${ Math . round ( waited ) } s (${ attempt } tries, GET /tools/tasks/${ taskId } ); the task may belong to another account or space.` ) ;
787795 }
788796 await delay ( pollInterval , signal ) ;
789797 }
0 commit comments