(res)
| 168 | event_loop = None |
| 169 | |
| 170 | def process_res(res): |
| 171 | client_id = res.client_id |
| 172 | nonlocal event_loop |
| 173 | nonlocal async_queues |
| 174 | |
| 175 | queue = self._results[client_id].queue |
| 176 | if isinstance(queue, _SyncQueue): |
| 177 | queue.put_nowait(res) |
| 178 | async_queues.append(queue) |
| 179 | # all the loops are identical |
| 180 | event_loop = event_loop or queue.loop |
| 181 | else: |
| 182 | queue.put(res) |
| 183 | |
| 184 | # FIXME: Add type annotations and make 'res' type more homogeneous (e.g. |
| 185 | # include PostprocWorker.Output in is_llm_response and unify is_final APIs). |
| 186 | if (is_llm_response(res) and res.result.is_final) or isinstance( |
| 187 | res, |
| 188 | ErrorResponse) or (isinstance(res, PostprocWorker.Output) |
| 189 | and res.is_final): |
| 190 | self._results.pop(client_id) |
| 191 | |
| 192 | res = res if isinstance(res, list) else [res] |
| 193 |
nothing calls this directly
no test coverage detected