@@ -679,11 +679,23 @@ func (p *Processor) processModelAsync(
679679 delete (pending , resp .RequestID )
680680 }
681681
682+ // Best-effort: tell the dispatcher to drop still-pending requests before
683+ // dispatch. Use a detached timeout — requestAbortCtx is often already
684+ // cancelled when we reach this path.
685+ if len (pending ) > 0 {
686+ cancelCtx , cancelFn := context .WithTimeout (context .Background (), 5 * time .Second )
687+ if err := asyncClient .Cancel (cancelCtx ); err != nil {
688+ logger .Error (err , "Failed to cancel pending async requests" , "pendingCount" , len (pending ))
689+ }
690+ cancelFn ()
691+ }
692+
682693 // Drain submitted-but-uncollected requests as errors so that
683- // output_lines + error_lines == total_requests.
694+ // output_lines + error_lines == total_requests. Error code follows the
695+ // same abort-reason priority as drainAndFinalize.
696+ errCode , errMsg := uncollectedPendingError (mainCtx , sloCtx , userCancelCtx , modelErr )
684697 for _ , pr := range pending {
685- out := newErrorOutputLine (pr .batchReqID , pr .customID ,
686- string (batch_types .ErrCodeBatchExpired ), "result not collected before deadline" )
698+ out := newErrorOutputLine (pr .batchReqID , pr .customID , string (errCode ), errMsg )
687699 lineBytes , err := json .Marshal (out )
688700 if err != nil {
689701 return fmt .Errorf ("marshal uncollected error line: %w" , err )
@@ -699,6 +711,30 @@ func (p *Processor) processModelAsync(
699711 inputFile , entries [submitCount :], writers , progress , modelErr , logger , len (entries ), 0 )
700712}
701713
714+ // uncollectedPendingError selects the error code/message for submitted-but-
715+ // uncollected async requests, matching drainAndFinalize's abort-reason order.
716+ // Context parameter order mirrors drainAndFinalize (minus requestAbortCtx).
717+ func uncollectedPendingError (
718+ mainCtx , sloCtx , userCancelCtx context.Context ,
719+ modelErr error ,
720+ ) (batch_types.BatchErrorCode , string ) {
721+ switch {
722+ case errors .Is (sloCtx .Err (), context .DeadlineExceeded ):
723+ return batch_types .ErrCodeBatchExpired , batch_types .ErrCodeBatchExpired .Message ()
724+ case userCancelCtx .Err () != nil :
725+ return batch_types .ErrCodeBatchCancelled , batch_types .ErrCodeBatchCancelled .Message ()
726+ case modelErr != nil :
727+ return batch_types .ErrCodeBatchFailed , batch_types .ErrCodeBatchFailed .Message ()
728+ case mainCtx .Err () != nil :
729+ // SIGTERM: leave terminalization to orphan recovery, but still record
730+ // lines so completed+failed == total for this model pass.
731+ return batch_types .ErrCodeBatchFailed , batch_types .ErrCodeBatchFailed .Message ()
732+ default :
733+ // Sibling abort or collect interrupted without a classified reason.
734+ return batch_types .ErrCodeBatchFailed , batch_types .ErrCodeBatchFailed .Message ()
735+ }
736+ }
737+
702738// drainAndFinalize drains undispatched entries based on termination reason and
703739// returns the appropriate sentinel error. Shared by processModel and processModelAsync.
704740func (p * Processor ) drainAndFinalize (
0 commit comments