Coverage for .venv/lib/python3.13/site-packages/litellm/proxy/guardrails/guardrail_hooks/microsoft_purview/purview_dlp.py: 10%

251 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 12:01 +0000

1""" 

2Microsoft Purview DLP Guardrail for LiteLLM. 

3 

4Supports three modes: 

5- pre_call: Block sensitive data in prompts before they reach the LLM. 

6- post_call: Block sensitive data in LLM responses. 

7- logging_only: Log interactions to Purview for audit/compliance without blocking. 

8""" 

9 

10import asyncio 

11import threading 

12import uuid 

13from collections.abc import AsyncGenerator, AsyncIterable, Mapping, Sequence 

14from datetime import datetime 

15from typing import TYPE_CHECKING, Any, Final, cast 

16 

17import httpx 

18from fastapi import HTTPException 

19 

20from litellm._logging import verbose_proxy_logger 

21from litellm.integrations.custom_guardrail import ( 

22 CustomGuardrail, 

23 log_guardrail_information, 

24) 

25from litellm.types.guardrails import GuardrailEventHooks 

26from litellm.types.utils import ( 

27 Choices, 

28 GuardrailStatus, 

29 ModelResponse, 

30 ModelResponseStream, 

31 ResponsesAPIResponse, 

32 TextChoices, 

33 TextCompletionResponse, 

34) 

35 

36from .base import PurviewGuardrailBase 

37 

38if TYPE_CHECKING: 38 ↛ 39line 38 didn't jump to line 39 because the condition on line 38 was never true

39 from litellm.caching.dual_cache import DualCache 

40 from litellm.proxy._types import UserAPIKeyAuth 

41 from litellm.types.llms.openai import AllMessageValues 

42 from litellm.types.proxy.guardrails.guardrail_hooks.base import ( 

43 GuardrailConfigModel, 

44 ) 

45 from litellm.types.utils import ( 

46 CallTypesLiteral, 

47 LLMResponseTypes, 

48 ) 

49 

50 

51class MicrosoftPurviewDLPGuardrail(PurviewGuardrailBase, CustomGuardrail): 

52 """ 

53 Microsoft Purview DLP guardrail. 

54 

55 Evaluates prompts and responses against Microsoft Purview DLP policies 

56 via the Microsoft Graph ``processContent`` API. 

57 """ 

58 

59 def __init__( 

60 self, 

61 guardrail_name: str, 

62 tenant_id: str, 

63 client_id: str, 

64 client_secret: str, 

65 purview_app_name: str = "LiteLLM", 

66 user_id_field: str = "user_id", 

67 **kwargs: object, 

68 ): 

69 super().__init__( 

70 tenant_id=tenant_id, 

71 client_id=client_id, 

72 client_secret=client_secret, 

73 purview_app_name=purview_app_name, 

74 user_id_field=user_id_field, 

75 guardrail_name=guardrail_name, 

76 supported_event_hooks=list(self.get_supported_event_hooks()), 

77 **kwargs, 

78 ) 

79 self.guardrail_provider = "microsoft_purview" 

80 verbose_proxy_logger.info( 

81 "Initialized Microsoft Purview DLP Guardrail: %s", 

82 guardrail_name, 

83 ) 

84 

85 @staticmethod 

86 def get_config_model() -> type["GuardrailConfigModel"] | None: 

87 return None # Config model can be added later for UI support 

88 

89 @classmethod 

90 def get_supported_event_hooks(cls) -> list[GuardrailEventHooks]: 

91 return [ 

92 GuardrailEventHooks.pre_call, 

93 GuardrailEventHooks.post_call, 

94 GuardrailEventHooks.logging_only, 

95 ] 

96 

97 # ------------------------------------------------------------------ 

98 # Core DLP check 

99 # ------------------------------------------------------------------ 

100 

101 async def _check_content( 

102 self, 

103 user_id: str, 

104 text: str, 

105 activity: str, 

106 request_data: dict[str, Any], 

107 block_on_violation: bool = True, 

108 ) -> dict[str, object]: 

109 """Evaluate content against Purview DLP policies. 

110 

111 Args: 

112 user_id: Entra object ID. 

113 text: Content to evaluate. 

114 activity: ``"uploadText"`` or ``"downloadText"``. 

115 request_data: Original request dict (used for logging metadata). 

116 block_on_violation: If False, log only — do not raise. 

117 

118 Returns: 

119 The processContent response dict. 

120 """ 

121 start_time: Final = datetime.now() 

122 status: GuardrailStatus = "success" 

123 response: dict[str, object] = {} 

124 

125 try: 

126 etag, _ = await self._compute_protection_scopes(user_id) 

127 correlation_id: Final = request_data.get("litellm_call_id") or str(uuid.uuid4()) 

128 response = await self._process_content( 

129 user_id=user_id, 

130 text=text, 

131 activity=activity, 

132 etag=etag, 

133 correlation_id=correlation_id, 

134 ) 

135 

136 if self._should_block(response): 

137 status = "guardrail_intervened" 

138 except HTTPException: 

139 status = "guardrail_failed_to_respond" 

140 raise 

141 except httpx.HTTPStatusError as exc: 

142 # Preserve the upstream Graph API status code (e.g. 429, 503) so 

143 # callers can distinguish a transient infrastructure error from a 

144 # DLP policy block (signaled separately as HTTP 400 below) and can 

145 # implement retry-after handling on rate limits. 401/403 upstream 

146 # responses indicate a proxy-side credential / consent problem the 

147 # caller can do nothing about, so they are mapped to 502. 

148 status = "guardrail_failed_to_respond" 

149 if block_on_violation: 

150 upstream_status: Final = exc.response.status_code 

151 client_status: Final = 502 if upstream_status in (401, 403) else upstream_status 

152 headers: dict[str, str] | None = None 

153 retry_after: Final[str | None] = exc.response.headers.get("retry-after") 

154 if retry_after: 

155 headers = {"Retry-After": retry_after} 

156 raise HTTPException( 

157 status_code=client_status, 

158 detail={ 

159 "error": "Microsoft Purview DLP: upstream policy evaluation failed", 

160 "activity": activity, 

161 "upstream_status": upstream_status, 

162 "exception": str(exc), 

163 }, 

164 headers=headers, 

165 ) from exc 

166 verbose_proxy_logger.warning( 

167 "Purview DLP: API/network error in logging-only mode (not re-raised): %s", 

168 exc, 

169 ) 

170 except Exception as exc: 

171 status = "guardrail_failed_to_respond" 

172 if block_on_violation: 

173 raise HTTPException( 

174 status_code=400, 

175 detail={ 

176 "error": "Microsoft Purview DLP: upstream policy evaluation failed", 

177 "activity": activity, 

178 "exception": str(exc), 

179 }, 

180 ) from exc 

181 verbose_proxy_logger.warning( 

182 "Purview DLP: API/network error in logging-only mode (not re-raised): %s", 

183 exc, 

184 ) 

185 finally: 

186 end_time: Final = datetime.now() 

187 self.add_standard_logging_guardrail_information_to_request_data( 

188 guardrail_provider=self.guardrail_provider, 

189 guardrail_json_response=response, 

190 request_data=request_data, 

191 guardrail_status=status, 

192 start_time=start_time.timestamp(), 

193 end_time=end_time.timestamp(), 

194 duration=(end_time - start_time).total_seconds(), 

195 ) 

196 

197 if block_on_violation and status == "guardrail_intervened": 

198 raise HTTPException( 

199 status_code=400, 

200 detail={ 

201 "error": "Microsoft Purview DLP: Content blocked by policy", 

202 "activity": activity, 

203 }, 

204 ) 

205 

206 return response 

207 

208 @staticmethod 

209 def _extract_responses_api_function_call_args(result: object) -> list[str]: 

210 """Return tool-call argument strings from a ``ResponsesAPIResponse.output``. 

211 

212 ``ResponsesAPIResponse.output_text`` only aggregates ``output_text`` 

213 content blocks and ignores ``function_call`` items. Model-generated 

214 tool-call arguments can themselves contain sensitive data, so we 

215 extract them explicitly to keep DLP coverage consistent with the 

216 chat (``ModelResponse``) path. 

217 """ 

218 args: Final[list[str]] = [] 

219 output: Final[Sequence[object] | None] = getattr(result, "output", None) 

220 if not output: 

221 return args 

222 for item in output: 

223 if isinstance(item, dict): 

224 item_type = item.get("type") 

225 arguments = item.get("arguments") 

226 else: 

227 item_type = getattr(item, "type", None) 

228 arguments = getattr(item, "arguments", None) 

229 if item_type == "function_call" and isinstance(arguments, str): 

230 if arguments.strip(): 

231 args.append(arguments) 

232 return args 

233 

234 def _completion_response_text_parts(self, result: object) -> list[str]: 

235 """Collect non-empty text segments from chat, text completions, or responses API. 

236 

237 Includes assistant message content *and* model-generated tool-call 

238 arguments so that sensitive data returned inside function calls is not 

239 missed by the DLP scan. 

240 """ 

241 parts: Final[list[str]] = [] 

242 if isinstance(result, TextCompletionResponse) and result.choices: 

243 for text_choice in result.choices: 

244 if not isinstance(text_choice, TextChoices): 

245 continue 

246 raw = text_choice.get("text") 

247 if isinstance(raw, str) and raw.strip(): 

248 parts.append(raw) 

249 elif isinstance(result, ResponsesAPIResponse): 

250 text: Final = result.output_text 

251 if text and text.strip(): 

252 parts.append(text) 

253 # Include tool-call arguments from ``function_call`` output items 

254 # (``output_text`` ignores them). 

255 parts.extend(self._extract_responses_api_function_call_args(result)) 

256 elif isinstance(result, ModelResponse) and result.choices: 

257 for chat_choice in result.choices: 

258 if not isinstance(chat_choice, Choices): 

259 continue 

260 msg = chat_choice.message 

261 if msg is None: 

262 continue 

263 raw = msg.get("content") if isinstance(msg, dict) else getattr(msg, "content", None) 

264 if isinstance(raw, str) and raw.strip(): 

265 parts.append(raw) 

266 # Include tool-call arguments returned by the model 

267 parts.extend(self._extract_tool_call_args_from_message(msg)) 

268 return parts 

269 

270 def _assemble_responses_api_from_chunks(self, chunks: Sequence[object]) -> tuple[bool, ResponsesAPIResponse | None]: 

271 """Extract the final ``ResponsesAPIResponse`` from a buffered Responses API stream. 

272 

273 Returns a ``(is_responses_api_stream, assembled)`` tuple so the caller 

274 can distinguish "not a Responses API stream" (fall through to 

275 ``stream_chunk_builder``) from "Responses API stream but no final 

276 response event was received" (fail closed with an accurate error). 

277 When the stream is a Responses API stream the latest event carrying a 

278 ``ResponsesAPIResponse`` body is returned (``response.completed``, or 

279 ``response.failed`` / ``response.incomplete`` as fallbacks). 

280 """ 

281 looks_like_responses_api = False 

282 final: ResponsesAPIResponse | None = None 

283 for chunk in chunks: 

284 event_type = getattr(chunk, "type", None) 

285 if isinstance(event_type, str) and event_type.startswith("response."): 

286 looks_like_responses_api = True 

287 candidate = getattr(chunk, "response", None) 

288 if isinstance(candidate, ResponsesAPIResponse): 

289 final = candidate 

290 return looks_like_responses_api, final 

291 

292 def _responses_api_input_to_str(self, data: dict[str, Any], raise_on_failure: bool = False) -> str | None: 

293 """Extract DLP-scannable text from a Responses API request ``input`` field. 

294 

295 ``input`` may be a plain string or a list of input items (messages). In 

296 the latter case the items are converted to chat messages via the standard 

297 LiteLLM transformation and then concatenated by ``get_prompt_text_for_dlp``. 

298 

299 When ``raise_on_failure`` is True (blocking mode), a transformation error 

300 raises ``HTTPException`` so the request is fail-closed. In logging-only 

301 mode the error is swallowed and ``None`` is returned so audit attempts on 

302 the response side can still run. 

303 """ 

304 from litellm.responses.litellm_completion_transformation.transformation import ( 

305 LiteLLMCompletionResponsesConfig, 

306 ) 

307 

308 input_data: Final = data.get("input") 

309 if input_data is None and not data.get("instructions"): 

310 return None 

311 try: 

312 # Always transform via messages so ``instructions`` become a system message 

313 # (string ``input`` alone would skip instructions and bypass DLP). 

314 messages: Final = LiteLLMCompletionResponsesConfig.transform_responses_api_input_to_messages( 

315 input=input_data if input_data is not None else "", 

316 responses_api_request=data, 

317 ) 

318 return self.get_prompt_text_for_dlp(cast(list["AllMessageValues"], messages)) 

319 except Exception: 

320 verbose_proxy_logger.warning( 

321 "Purview DLP: failed to transform responses API input", 

322 exc_info=True, 

323 ) 

324 if raise_on_failure: 

325 raise HTTPException( 

326 status_code=400, 

327 detail={ 

328 "error": ( 

329 "Microsoft Purview DLP: Responses API input could " 

330 "not be transformed for DLP scanning in blocking mode" 

331 ), 

332 }, 

333 ) 

334 return None 

335 

336 # ------------------------------------------------------------------ 

337 # Identity resolution for blocking modes 

338 # ------------------------------------------------------------------ 

339 

340 def _resolve_user_id_for_blocking( 

341 self, 

342 data: Mapping[str, object], 

343 user_api_key_dict: "UserAPIKeyAuth", 

344 ) -> str: 

345 """Resolve user ID for blocking (pre_call / post_call) DLP hooks. 

346 

347 Uses only trusted proxy-authenticated sources (``_resolve_trusted_user_id``). 

348 Caller-supplied ``UserAPIKeyAuth.end_user_id`` (from request ``user``, 

349 ``metadata.user_id``, ``safety_identifier``, etc.) and 

350 ``metadata[user_id_field]`` are rejected (fail closed) because they can 

351 impersonate another Entra user's Purview policy. 

352 

353 Raises ``HTTPException`` when no API-key-bound ``user_id`` exists or when 

354 only caller-influenceable identity fields are available (fail closed). 

355 """ 

356 trusted_id: Final = self._resolve_trusted_user_id(data, user_api_key_dict) 

357 if trusted_id: 

358 return trusted_id 

359 

360 if self._resolve_user_id(data, user_api_key_dict): 

361 raise HTTPException( 

362 status_code=400, 

363 detail={ 

364 "error": ( 

365 "Microsoft Purview DLP: No proxy-authenticated user identity; " 

366 "bind user_id to the API key (caller-supplied metadata cannot " 

367 "be used for blocking DLP)" 

368 ), 

369 }, 

370 ) 

371 

372 raise HTTPException( 

373 status_code=400, 

374 detail={ 

375 "error": ( 

376 "Microsoft Purview DLP: No proxy-authenticated user identity; " 

377 "bind user_id to the API key for blocking DLP" 

378 ), 

379 }, 

380 ) 

381 

382 # ------------------------------------------------------------------ 

383 # Pre-call hook — DLP on prompts 

384 # ------------------------------------------------------------------ 

385 

386 @log_guardrail_information 

387 async def async_pre_call_hook( 

388 self, 

389 user_api_key_dict: "UserAPIKeyAuth", 

390 cache: "DualCache", 

391 data: dict[str, Any], 

392 call_type: "CallTypesLiteral", 

393 ) -> dict[str, object] | None: 

394 """Check user prompt against Purview DLP policies before LLM call.""" 

395 user_id: Final = self._resolve_user_id_for_blocking(data, user_api_key_dict) 

396 

397 prompt_text: str | None = None 

398 if call_type in ("responses", "aresponses"): 

399 # Route Responses API calls to the responses-specific extractor 

400 # before the generic ``messages`` branch. This mirrors 

401 # ``async_logging_hook`` and ensures ``instructions`` (system 

402 # prompt) content is included in the DLP scan, and prevents a 

403 # crafted ``messages`` key in the request from being scanned in 

404 # place of the actual ``input``. 

405 prompt_text = self._responses_api_input_to_str(data, raise_on_failure=True) 

406 elif call_type in ("text_completion", "atext_completion"): 

407 raw_prompt: Final = data.get("prompt") 

408 # Reject every token-id prompt shape Purview cannot evaluate — 

409 # flat ``list[int]`` (single prompt), ``list[list[int]]`` (multi-prompt 

410 # batches), and mixed lists that include any token-id sub-array. 

411 # Empty/whitespace-only strings also yield ``prompt_text is None`` but 

412 # contain no sensitive data and pass through harmlessly below. 

413 if self.is_token_id_prompt(raw_prompt): 

414 raise HTTPException( 

415 status_code=400, 

416 detail={ 

417 "error": ( 

418 "Microsoft Purview DLP: Token-id completion prompts " 

419 "cannot be scanned for DLP in blocking mode" 

420 ), 

421 }, 

422 ) 

423 prompt_text = self.completion_prompt_to_str(raw_prompt) 

424 else: 

425 messages: Final[list | None] = data.get("messages") 

426 if messages: 

427 prompt_text = self.get_prompt_text_for_dlp(cast(list["AllMessageValues"], messages)) 

428 

429 if not prompt_text: 

430 return data 

431 

432 await self._check_content( 

433 user_id=user_id, 

434 text=prompt_text, 

435 activity="uploadText", 

436 request_data=data, 

437 block_on_violation=True, 

438 ) 

439 return data 

440 

441 # ------------------------------------------------------------------ 

442 # Post-call hook — DLP on responses 

443 # ------------------------------------------------------------------ 

444 

445 @log_guardrail_information 

446 async def async_post_call_success_hook( 

447 self, 

448 data: dict, 

449 user_api_key_dict: "UserAPIKeyAuth", 

450 response: "LLMResponseTypes", 

451 ) -> "LLMResponseTypes": 

452 """Check LLM response against Purview DLP policies (non-streaming only). 

453 

454 Streaming responses are handled by ``async_post_call_streaming_iterator_hook`` 

455 which buffers all chunks before scanning. The proxy automatically skips 

456 this hook for requests that have a streaming iterator hook defined. 

457 """ 

458 user_id: Final = self._resolve_user_id_for_blocking(data, user_api_key_dict) 

459 

460 parts: Final = self._completion_response_text_parts(response) 

461 

462 if parts: 

463 combined: Final = "\n\n---\n\n".join(parts) 

464 await self._check_content( 

465 user_id=user_id, 

466 text=combined, 

467 activity="downloadText", 

468 request_data=data, 

469 block_on_violation=True, 

470 ) 

471 return response 

472 

473 async def async_post_call_streaming_iterator_hook( 

474 self, 

475 user_api_key_dict: "UserAPIKeyAuth", 

476 response: AsyncIterable[ModelResponseStream], 

477 request_data: dict, 

478 ) -> AsyncGenerator[ModelResponseStream, None]: 

479 """Check streaming LLM responses against Purview DLP policies. 

480 

481 All chunks are buffered before the DLP scan so that no content is 

482 delivered to the client if a policy violation is detected. After a 

483 clean scan the assembled response is re-yielded chunk-by-chunk via a 

484 ``MockResponseIterator`` so the caller receives normal streaming output. 

485 

486 The proxy automatically skips ``async_post_call_success_hook`` for 

487 guardrails that define this method, preventing duplicate scans. 

488 """ 

489 from litellm.llms.base_llm.base_model_iterator import MockResponseIterator 

490 from litellm.main import stream_chunk_builder 

491 

492 # Resolve user ID up-front so identity failures don't waste work 

493 # buffering and assembling the stream. 

494 user_id: Final = self._resolve_user_id_for_blocking(request_data, user_api_key_dict) 

495 

496 # Buffer the entire stream before any DLP scan. 

497 all_chunks: Final[list[ModelResponseStream]] = [] 

498 async for chunk in response: 

499 all_chunks.append(chunk) 

500 

501 # Responses API streams emit typed events (e.g. ``response.completed``) 

502 # whose final event carries the full ``ResponsesAPIResponse`` — these 

503 # are not understood by ``stream_chunk_builder`` (which is built for 

504 # chat/text-completion deltas). Detect and scan them via the same 

505 # ``_completion_response_text_parts`` path used by non-streaming. 

506 ( 

507 is_responses_api_stream, 

508 responses_api_assembled, 

509 ) = self._assemble_responses_api_from_chunks(all_chunks) 

510 if is_responses_api_stream: 

511 if responses_api_assembled is None: 

512 # Fail closed: Responses API events were seen but no final 

513 # ``response.completed`` / ``response.failed`` / 

514 # ``response.incomplete`` event carrying a ``ResponsesAPIResponse`` 

515 # body was received, so we cannot scan the content. 

516 raise HTTPException( 

517 status_code=400, 

518 detail={ 

519 "error": ( 

520 "Microsoft Purview DLP: Incomplete Responses API " 

521 "stream — no final response event received for " 

522 "DLP scanning; blocking response." 

523 ), 

524 }, 

525 ) 

526 parts = self._completion_response_text_parts(responses_api_assembled) 

527 if parts: 

528 combined = "\n\n---\n\n".join(parts) 

529 await self._check_content( 

530 user_id=user_id, 

531 text=combined, 

532 activity="downloadText", 

533 request_data=request_data, 

534 block_on_violation=True, 

535 ) 

536 for chunk in all_chunks: 

537 yield chunk 

538 return 

539 

540 assembled_response: Final = stream_chunk_builder(chunks=all_chunks) 

541 

542 if assembled_response is None and all_chunks: 

543 # Fail closed: stream_chunk_builder dropped all chunks, so we cannot 

544 # scan the content. Refuse to release the buffered chunks. 

545 raise HTTPException( 

546 status_code=400, 

547 detail={ 

548 "error": ( 

549 "Microsoft Purview DLP: Unable to assemble streamed response for scanning; blocking response." 

550 ), 

551 }, 

552 ) 

553 

554 if isinstance(assembled_response, (TextCompletionResponse, ResponsesAPIResponse)): 

555 parts = self._completion_response_text_parts(assembled_response) 

556 if parts: 

557 combined = "\n\n---\n\n".join(parts) 

558 await self._check_content( 

559 user_id=user_id, 

560 text=combined, 

561 activity="downloadText", 

562 request_data=request_data, 

563 block_on_violation=True, 

564 ) 

565 for chunk in all_chunks: 

566 yield chunk 

567 return 

568 

569 if not isinstance(assembled_response, ModelResponse): 

570 # Non-content response (e.g. embeddings) — pass through unchanged. 

571 for chunk in all_chunks: 

572 yield chunk 

573 return 

574 

575 parts = self._completion_response_text_parts(assembled_response) 

576 if parts: 

577 combined = "\n\n---\n\n".join(parts) 

578 # Raises HTTPException(400) on violation — no chunks are yielded. 

579 await self._check_content( 

580 user_id=user_id, 

581 text=combined, 

582 activity="downloadText", 

583 request_data=request_data, 

584 block_on_violation=True, 

585 ) 

586 

587 # DLP passed — re-yield chunks from the assembled chat response. 

588 mock_response: Final = MockResponseIterator(model_response=assembled_response) 

589 async for chunk in mock_response: 

590 yield chunk 

591 

592 # ------------------------------------------------------------------ 

593 # Logging-only hook — audit without blocking 

594 # ------------------------------------------------------------------ 

595 

596 def logging_hook(self, kwargs: dict, result: object, call_type: str) -> tuple[dict, object]: 

597 """Fire-and-forget async audit logging; returns original (kwargs, result) immediately. 

598 

599 In the proxy's async success path, litellm independently calls both 

600 ``logging_hook`` (sync) and ``async_logging_hook`` (async) for every 

601 ``CustomGuardrail`` callback. To avoid making two complete sets of 

602 Purview API calls per request, this sync hook is a no-op whenever an 

603 event loop is running — the framework's async path will invoke 

604 ``async_logging_hook`` directly. 

605 

606 For genuine sync-only call paths (no running event loop, so the async 

607 success handler will not fire either), schedule ``async_logging_hook`` 

608 on a short-lived background daemon thread so audit logging still runs 

609 without blocking the caller on two Graph API round-trips. 

610 """ 

611 

612 try: 

613 asyncio.get_running_loop() 

614 # Async context — let the framework's async success handler invoke 

615 # async_logging_hook to avoid duplicate Purview API calls. Log so 

616 # the deferral is observable if the framework ever stops dispatching 

617 # async_logging_hook on a given code path (otherwise audit silently 

618 # drops). 

619 verbose_proxy_logger.debug("Purview audit: deferring to async_logging_hook (running event loop detected)") 

620 return kwargs, result 

621 except RuntimeError: 

622 pass 

623 

624 async def _log_safe() -> None: 

625 try: 

626 await self.async_logging_hook(kwargs=kwargs, result=result, call_type=call_type) 

627 except Exception as exc: 

628 verbose_proxy_logger.error("Purview audit background logging error: %s", exc) 

629 

630 def _run_in_new_loop() -> None: 

631 new_loop: Final = asyncio.new_event_loop() 

632 try: 

633 asyncio.set_event_loop(new_loop) 

634 new_loop.run_until_complete(_log_safe()) 

635 finally: 

636 new_loop.close() 

637 asyncio.set_event_loop(None) 

638 

639 thread: Final = threading.Thread(target=_run_in_new_loop, daemon=True) 

640 thread.start() 

641 

642 return kwargs, result 

643 

644 async def async_logging_hook(self, kwargs: dict, result: object, call_type: str) -> tuple[dict, object]: 

645 """Send both prompt and response to Purview for audit logging. 

646 

647 Errors are logged but never raised — this mode is non-blocking. 

648 Each audit call (prompt and response) is wrapped in its own try/except 

649 so a failure on the first does not prevent the second from running. 

650 """ 

651 user_id: Final = self._resolve_user_id_from_logging_kwargs(kwargs) 

652 if not user_id: 

653 verbose_proxy_logger.debug("Purview audit: no user_id, skipping") 

654 return kwargs, result 

655 

656 # Log prompt (uploadText) 

657 try: 

658 prompt_text: str | None = None 

659 if call_type in ("responses", "aresponses"): 

660 # Responses API: route to the responses-specific extractor 

661 # before the generic ``messages`` branch. litellm's logging 

662 # pipeline stores the raw responses ``input`` (a string or a 

663 # list of input items) under ``model_call_details["messages"]`` 

664 # via ``function_setup``, which is NOT the chat message format 

665 # ``get_prompt_text_for_dlp`` expects. Use the original 

666 # ``input`` / ``instructions`` keys that ``pre_call`` and 

667 # ``update_environment_variables`` persist on the call details. 

668 prompt_text = self._responses_api_input_to_str(kwargs) 

669 elif call_type in ("text_completion", "atext_completion"): 

670 prompt_text = self.completion_prompt_to_str(kwargs.get("prompt")) 

671 else: 

672 messages: Final = kwargs.get("messages") 

673 if messages: 

674 prompt_text = self.get_prompt_text_for_dlp(cast(list["AllMessageValues"], messages)) 

675 

676 if prompt_text: 

677 await self._check_content( 

678 user_id=user_id, 

679 text=prompt_text, 

680 activity="uploadText", 

681 request_data=kwargs, 

682 block_on_violation=False, 

683 ) 

684 except Exception as e: 

685 verbose_proxy_logger.error("Purview audit logging error (prompt): %s", e) 

686 

687 # Log response (downloadText) — runs regardless of prompt audit outcome 

688 try: 

689 parts: Final = self._completion_response_text_parts(result) 

690 if parts: 

691 combined: Final = "\n\n---\n\n".join(parts) 

692 await self._check_content( 

693 user_id=user_id, 

694 text=combined, 

695 activity="downloadText", 

696 request_data=kwargs, 

697 block_on_violation=False, 

698 ) 

699 except Exception as e: 

700 verbose_proxy_logger.error("Purview audit logging error (response): %s", e) 

701 

702 return kwargs, result