|
| 1 | +import asyncio |
| 2 | +import json |
| 3 | +from typing import List |
| 4 | + |
| 5 | +import aiohttp |
| 6 | + |
| 7 | +from tensorrt_llm.logger import logger |
| 8 | + |
| 9 | +from ..pyexecutor.llm_request import * |
| 10 | +from ..pyexecutor.scheduler import ScheduledRequests |
| 11 | +from .drafter import Drafter |
| 12 | + |
| 13 | + |
| 14 | +class APIDrafter(Drafter): |
| 15 | + |
| 16 | + def __init__( |
| 17 | + self, |
| 18 | + spec_config: "ExternalAPIConfig", |
| 19 | + ): |
| 20 | + super().__init__() |
| 21 | + self.max_draft_len = spec_config.max_draft_len |
| 22 | + self.endpoint = spec_config.endpoint |
| 23 | + assert self.endpoint is not None, "API endpoint is required for external API speculative decoding." |
| 24 | + self.template = spec_config.template if spec_config.template is not None else {} |
| 25 | + self.response_field = spec_config.response_field if spec_config.response_field is not None else "draft_tokens" |
| 26 | + |
| 27 | + def single_draft_call(self): |
| 28 | + return True |
| 29 | + |
| 30 | + def get_nested_field_from_response(self, response: dict) -> List[int]: |
| 31 | + # Allows for nested fields in the response. |
| 32 | + # Example: "choices.0.message.content" |
| 33 | + # Returns the value of the nested field: response["choices"][0]["message"]["content"] |
| 34 | + keys = self.response_field.split(".") |
| 35 | + current = response |
| 36 | + |
| 37 | + for key in keys: |
| 38 | + try: |
| 39 | + if key.isdigit(): |
| 40 | + key = int(key) |
| 41 | + if isinstance(current, list) and 0 <= key < len(current): |
| 42 | + current = current[key] |
| 43 | + else: |
| 44 | + logger.warning( |
| 45 | + f"Response field {self.response_field} is invalid for response {response}. Index {key} is invalid." |
| 46 | + ) |
| 47 | + return [] |
| 48 | + else: |
| 49 | + if isinstance(current, dict) and key in current: |
| 50 | + current = current[key] |
| 51 | + else: |
| 52 | + logger.warning( |
| 53 | + f"Response field {self.response_field} is invalid for response {response}. Index {key} is invalid." |
| 54 | + ) |
| 55 | + return [] |
| 56 | + |
| 57 | + except (KeyError, ValueError, IndexError): |
| 58 | + logger.warning( |
| 59 | + f"Response field path is invalid: {self.response_field}") |
| 60 | + return [] |
| 61 | + |
| 62 | + if not isinstance(current, list): |
| 63 | + logger.warning( |
| 64 | + f"API response '{self.response_field}' must be a list. Got type: {type(current)}" |
| 65 | + ) |
| 66 | + return [] |
| 67 | + return current |
| 68 | + |
| 69 | + async def get_draft_tokens( |
| 70 | + self, |
| 71 | + prefix: list[int], |
| 72 | + request_id: int, |
| 73 | + end_id: int, |
| 74 | + max_sequence_length: int, |
| 75 | + ) -> List[int]: |
| 76 | + try: |
| 77 | + request_data = { |
| 78 | + "prefix": prefix, |
| 79 | + "request_id": request_id, |
| 80 | + "end_id": end_id, |
| 81 | + "max_sequence_length": max_sequence_length, |
| 82 | + } |
| 83 | + if self.template: |
| 84 | + request_data.update(self.template) |
| 85 | + |
| 86 | + async with aiohttp.ClientSession() as session: |
| 87 | + async with session.post( |
| 88 | + url=self.endpoint, |
| 89 | + json=request_data, |
| 90 | + headers={"Content-Type": "application/json"}, |
| 91 | + timeout=aiohttp.ClientTimeout(total=10), |
| 92 | + ) as response: |
| 93 | + |
| 94 | + # check for unsuccessful response |
| 95 | + if response.status != 200: |
| 96 | + logger.error( |
| 97 | + f"Failed to get draft tokens. API call failed for request {request_id} with status code {response.status}" |
| 98 | + ) |
| 99 | + return [] |
| 100 | + |
| 101 | + result = await response.json() |
| 102 | + draft_tokens = self.get_nested_field_from_response(result) |
| 103 | + if len(draft_tokens) > self.max_draft_len: |
| 104 | + draft_tokens = draft_tokens[:self.max_draft_len] |
| 105 | + logger.debug( |
| 106 | + f"Retrieved draft tokens for request {request_id}: {draft_tokens}" |
| 107 | + ) |
| 108 | + return draft_tokens |
| 109 | + |
| 110 | + except json.JSONDecodeError as e: |
| 111 | + logger.error( |
| 112 | + f"Failed to parse JSON response for request {request_id}: {e}") |
| 113 | + return [] |
| 114 | + |
| 115 | + except Exception as e: |
| 116 | + logger.error( |
| 117 | + f"Failed to get draft tokens. API call failed for request {request_id} with the following error: {e}" |
| 118 | + ) |
| 119 | + return [] |
| 120 | + |
| 121 | + async def async_prepare_draft_tokens( |
| 122 | + self, |
| 123 | + scheduled_requests: ScheduledRequests, |
| 124 | + resource_manager: None, |
| 125 | + ) -> None: |
| 126 | + # Sort by request_id when py_batch_idx is None as a fallback. |
| 127 | + # This happens in the disagg case: for a set of new requests, we draft |
| 128 | + # before forward_step, so py_batch_idx is not assigned. |
| 129 | + sorted_requests = sorted( |
| 130 | + scheduled_requests.generation_requests, |
| 131 | + key=lambda r: |
| 132 | + (r.py_batch_idx is None, r.py_batch_idx or r.request_id), |
| 133 | + ) |
| 134 | + |
| 135 | + tasks = [] |
| 136 | + for request in sorted_requests: |
| 137 | + # Add new token to a copy of the generated tokens to find new draft tokens |
| 138 | + prefix = list(request.get_tokens()[0]) # Get a copy |
| 139 | + task = self.get_draft_tokens( |
| 140 | + prefix, |
| 141 | + request.request_id, |
| 142 | + request.py_end_id, |
| 143 | + request.py_orig_prompt_len + request.py_max_new_tokens, |
| 144 | + ) |
| 145 | + tasks.append(task) |
| 146 | + |
| 147 | + try: |
| 148 | + all_draft_tokens = await asyncio.wait_for(asyncio.gather( |
| 149 | + *tasks, return_exceptions=True), |
| 150 | + timeout=10.0) |
| 151 | + except asyncio.TimeoutError: |
| 152 | + logger.error( |
| 153 | + f"Timeout occurred while getting draft tokens for batch of requests" |
| 154 | + ) |
| 155 | + all_draft_tokens = [[] for _ in tasks] |
| 156 | + |
| 157 | + for request, draft_tokens in zip(sorted_requests, all_draft_tokens): |
| 158 | + if isinstance(draft_tokens, Exception): |
| 159 | + logger.error( |
| 160 | + f"An exception occurred while getting draft tokens for request {request.request_id}. Set TLLM_LOG_LEVEL for more details." |
| 161 | + ) |
| 162 | + draft_tokens = [] |
| 163 | + elif len(draft_tokens) == 0: |
| 164 | + logger.error( |
| 165 | + f"Draft tokens could not be generated for request {request.request_id}. Set TLLM_LOG_LEVEL for more details." |
| 166 | + ) |
| 167 | + else: |
| 168 | + # Pad length to `self.max_draft_len` |
| 169 | + if len(draft_tokens) > 0: |
| 170 | + pad_length = self.max_draft_len - len(draft_tokens) |
| 171 | + draft_tokens.extend([request.py_end_id] * pad_length) |
| 172 | + |
| 173 | + request.py_draft_tokens = draft_tokens |
| 174 | + |
| 175 | + def prepare_draft_tokens( |
| 176 | + self, |
| 177 | + scheduled_requests: ScheduledRequests, |
| 178 | + resource_manager: None, |
| 179 | + ) -> None: |
| 180 | + asyncio.run( |
| 181 | + self.async_prepare_draft_tokens(scheduled_requests, |
| 182 | + resource_manager)) |
0 commit comments