-
Notifications
You must be signed in to change notification settings - Fork 2
Expand file tree
/
Copy pathgateway_request.py
More file actions
945 lines (820 loc) · 33.8 KB
/
Copy pathgateway_request.py
File metadata and controls
945 lines (820 loc) · 33.8 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
"""Inbound request-boundary helpers for the CodexHub Gateway.
Owns request decoding and model extraction, Official reasoning sanitization,
browser-context detection, local Gateway authorization, request-context
headers, response-header filtering, reasoning-effort validation, and the
request-time Vision Proxy adapter factory.
"""
from __future__ import annotations
import gzip
import hmac
import io
import json
import zlib
from collections.abc import Mapping
import re
from typing import Any
import gateway_events
import proxy_telemetry
import vision_proxy
from catalog import canonical_model_id
import gateway_catalog_runtime as _catalog
import maintained_catalog
import gateway_stream_semantics as _stream_semantics
import gateway_settings
from gateway_transport import get_header as _get_header, header_items as _header_items
from route_primitives import (
RouteProtocol,
)
try:
import zstandard
except ImportError: # pragma: no cover - optional dependency on older Python installs.
zstandard = None
DECODE_ERRORS = (OSError, zlib.error) + ((zstandard.ZstdError,) if zstandard is not None else ())
OFFICIAL_ULTRA_REASONING_MODELS = {"gpt-5.6-sol", "gpt-5.6-terra"}
OFFICIAL_ALIAS_PREFIX = "openai/"
UNSUPPORTED_REASONING_MODEL_PREFIXES = () # retained for tests; catalog drives drop now
_UPSTREAM_TO_PROVIDER = {
"ollama_cloud": "ollama-cloud",
"volcengine": "volc",
"minimax_cn": "minimax-cn",
"kimi": "kimi",
"commandcode": "commandcode",
"opencode_go": "opencode-go",
"xai": "xai",
}
OFFICIAL_ENCRYPTED_CONTENT_PREFIX = "gAAAA"
BROWSER_CONTEXT_MARKERS = (
"# in app browser",
"# browser comments",
"browser visual feedback",
)
BROWSER_CURRENT_URL_RE = re.compile(
r"(?im)^\s*(?:current\s+url|current\s+browser\s+url|browser\s+url|url)\s*:\s*https?://\S+"
)
BROWSER_CONTEXT_GUIDANCE_SENTINEL = "Codex browser context detected."
BROWSER_CONTEXT_GUIDANCE = (
BROWSER_CONTEXT_GUIDANCE_SENTINEL
+ "\nRequired browser-control workflow:\n"
"- Load and follow the browser:control-in-app-browser skill before saying browser control is unavailable.\n"
'- For OpenAI/Codex native discovery, use tool_search with query "node_repl js" if mcp__node_repl.js is not already visible.\n'
"- Browser control is unavailable only when that search does not return mcp__node_repl.js, or when mcp__node_repl.js reports no in-app browser session.\n"
"- If executable alias mcp__node_repl__js is visible, use it directly to bootstrap browser-client.mjs and select the iab browser.\n"
'- In a CLI/no-browser environment, report "browser session unavailable"; do not report "browser tool not exposed".'
)
EMBEDDED_MODEL_RE = re.compile(rb'"model"\s*:\s*"(?:[^"\\]|\\.)+"')
FORM_MODEL_RE = re.compile(rb'name="model"(?:\r?\n[^\r\n]*)*\r?\n\r?\n([^\r\n]+)')
HOP_BY_HOP_RESPONSE_HEADERS = {
"connection",
"keep-alive",
"proxy-authenticate",
"proxy-authorization",
"server",
"te",
"trailer",
"trailers",
"transfer-encoding",
"upgrade",
}
def decoded_request_body(body: bytes, content_encoding: str | None = None) -> tuple[bytes, bool, str | None]:
if not content_encoding:
return body, False, None
encoding = content_encoding.lower()
try:
if "gzip" in encoding:
return gzip.decompress(body), True, None
if "deflate" in encoding:
return zlib.decompress(body), True, None
if "zstd" in encoding:
if zstandard is None:
return body, False, "zstandard module is not available"
with zstandard.ZstdDecompressor().stream_reader(io.BytesIO(body)) as reader:
return reader.read(), True, None
except DECODE_ERRORS as exc:
return body, False, f"{type(exc).__name__}: {exc}"
return body, False, None
def _decode_json_string_token(token: bytes) -> str | None:
try:
value = json.loads(token.decode("utf-8"))
except (UnicodeDecodeError, json.JSONDecodeError):
return None
return value if isinstance(value, str) and value.strip() else None
def try_extract_model(body: bytes, content_encoding: str | None = None) -> str | None:
scan_body, _, _ = decoded_request_body(body, content_encoding)
try:
payload = json.loads(scan_body.decode("utf-8-sig"))
except (UnicodeDecodeError, json.JSONDecodeError):
payload = None
if isinstance(payload, dict):
model = payload.get("model")
return model if isinstance(model, str) and model.strip() else None
form_match = FORM_MODEL_RE.search(scan_body)
if form_match:
try:
form_model = form_match.group(1).strip().decode("utf-8")
except UnicodeDecodeError:
form_model = ""
if form_model:
return form_model
for match in EMBEDDED_MODEL_RE.finditer(scan_body):
token = match.group(0).split(b":", 1)[1].strip()
model = _decode_json_string_token(token)
if model:
return model
return None
def extract_model(body: bytes) -> str:
model = try_extract_model(body)
if model:
return model
try:
payload = json.loads(body.decode("utf-8-sig"))
except (UnicodeDecodeError, json.JSONDecodeError) as exc:
raise ValueError("request body must include a string model") from exc
if not isinstance(payload, dict):
raise ValueError("request body must be a JSON object")
raise ValueError("request body must include a string model")
def _looks_like_official_encrypted_content(value: Any) -> bool:
return isinstance(value, str) and value.startswith(OFFICIAL_ENCRYPTED_CONTENT_PREFIX)
def sanitize_official_reasoning_items(value: Any) -> bool:
changed = False
if isinstance(value, list):
for item in value:
if _sanitize_official_reasoning_items(item):
changed = True
return changed
if not isinstance(value, dict):
return False
if value.get("type") == "reasoning" and "encrypted_content" in value:
if not _looks_like_official_encrypted_content(value.get("encrypted_content")):
value.pop("encrypted_content", None)
changed = True
for item in value.values():
if _sanitize_official_reasoning_items(item):
changed = True
return changed
_sanitize_official_reasoning_items = sanitize_official_reasoning_items
def sanitize_official_input_reasoning_items(payload: dict[str, Any]) -> tuple[bool, dict[str, int]]:
"""Remove non-portable reasoning references at the Official input boundary.
Codex stores third-party Responses reasoning items in the shared task
history. An Official ``store=false`` request cannot resolve those
provider-local IDs, so forwarding the whole item turns a model switch into
a permanent 404/reconnect loop. Only a self-contained Official encrypted
item is portable; every other reasoning item is dropped as a whole. The
walk is deliberately limited to input containers so response metadata,
tool schemas, and ordinary transcript fields remain untouched.
"""
counts = {
"removed_non_portable": 0,
"kept_official_encrypted": 0,
}
def sanitize_input_list(items: list[Any]) -> tuple[list[Any], bool]:
changed = False
rewritten: list[Any] = []
for item in items:
if isinstance(item, dict) and item.get("type") == "reasoning":
encrypted_content = item.get("encrypted_content")
if _looks_like_official_encrypted_content(encrypted_content):
counts["kept_official_encrypted"] += 1
rewritten.append(item)
else:
counts["removed_non_portable"] += 1
changed = True
continue
if isinstance(item, list):
nested_items, nested_changed = sanitize_input_list(item)
if nested_changed:
item = nested_items
changed = True
elif isinstance(item, dict) and isinstance(item.get("input"), list):
nested_items, nested_changed = sanitize_input_list(item["input"])
if nested_changed:
item = {**item, "input": nested_items}
changed = True
rewritten.append(item)
return rewritten, changed
input_items = payload.get("input")
if not isinstance(input_items, list):
return False, counts
rewritten_items, changed = sanitize_input_list(input_items)
if changed:
payload["input"] = rewritten_items
return changed, counts
_sanitize_official_input_reasoning_items = sanitize_official_input_reasoning_items
def strip_reasoning_encrypted_content(value: Any) -> bool:
changed = False
if isinstance(value, list):
for item in value:
if _strip_reasoning_encrypted_content(item):
changed = True
return changed
if not isinstance(value, dict):
return False
if value.get("type") == "reasoning" and "encrypted_content" in value:
value.pop("encrypted_content", None)
changed = True
for item in value.values():
if _strip_reasoning_encrypted_content(item):
changed = True
return changed
_strip_reasoning_encrypted_content = strip_reasoning_encrypted_content
def _third_party_reasoning_part_text(part: Any) -> str | None:
if isinstance(part, str) and part:
return part
if not isinstance(part, dict):
return None
for key in ("text", "summary"):
value = part.get(key)
if isinstance(value, str) and value:
return value
return None
def _third_party_reasoning_summary_parts(value: Any) -> list[dict[str, str]]:
if value is None:
return []
if isinstance(value, str) and value:
return [{"type": "summary_text", "text": value}]
if not isinstance(value, list):
return []
parts: list[dict[str, str]] = []
for item in value:
text = _third_party_reasoning_part_text(item)
if text:
parts.append({"type": "summary_text", "text": text})
return parts
_THIRD_PARTY_ENCRYPTED_AGENT_MESSAGE_PLACEHOLDER = (
"[Official encrypted agent_message unavailable]"
)
def _rewrite_third_party_encrypted_agent_message(
item: dict[str, Any],
) -> dict[str, Any] | None:
"""Drop Official ciphertext from collaboration handoff history.
Third-party Responses endpoints cannot decrypt ``encrypted_content``.
Keeping it for later fail-closed kills compact and the child turn; the
ciphertext is also useless as model input. Preserve any plaintext parts.
Replace an encrypted-only item with a bounded developer placeholder.
"""
if item.get("type") != "agent_message":
return None
content = item.get("content")
if not isinstance(content, list):
return None
kept: list[Any] = []
saw_encrypted = False
for part in content:
if isinstance(part, Mapping) and part.get("type") == "encrypted_content":
saw_encrypted = True
continue
kept.append(part)
if not saw_encrypted:
return None
if kept:
rewritten = dict(item)
rewritten["content"] = kept
return rewritten
return {
"type": "message",
"role": "developer",
"content": _THIRD_PARTY_ENCRYPTED_AGENT_MESSAGE_PLACEHOLDER,
}
def sanitize_third_party_reasoning_items(
value: Any,
*,
preserve_collaboration_agent_message_encryption: bool = False,
inside_collaboration_agent_message: bool = False,
) -> bool:
# Codex App history stores Official/xAI reasoning blobs that other
# Responses providers cannot decrypt, plus content: null which xAI rejects
# as "invalid type: null, expected a sequence".
changed = False
if isinstance(value, list):
index = 0
while index < len(value):
item = value[index]
if (
isinstance(item, dict)
and not preserve_collaboration_agent_message_encryption
):
rewritten_agent_message = _rewrite_third_party_encrypted_agent_message(item)
if rewritten_agent_message is not None:
value[index] = rewritten_agent_message
item = rewritten_agent_message
changed = True
if (
isinstance(item, dict)
and item.get("type") in {"encrypted_content", "item_reference"}
and not (
preserve_collaboration_agent_message_encryption
and inside_collaboration_agent_message
and item.get("type") == "encrypted_content"
)
):
del value[index]
changed = True
continue
if sanitize_third_party_reasoning_items(
item,
preserve_collaboration_agent_message_encryption=preserve_collaboration_agent_message_encryption,
inside_collaboration_agent_message=inside_collaboration_agent_message,
):
changed = True
if (
isinstance(item, dict)
and item.get("type") == "reasoning"
and not item.get("summary")
and item.get("content") == []
and "id" not in item
):
del value[index]
changed = True
continue
index += 1
return changed
if not isinstance(value, dict):
return False
if isinstance(value.get("previous_response_id"), str):
value.pop("previous_response_id", None)
changed = True
include = value.get("include")
if isinstance(include, list):
filtered = [item for item in include if item != "reasoning.encrypted_content"]
if len(filtered) != len(include):
if filtered:
value["include"] = filtered
else:
value.pop("include", None)
changed = True
if value.get("type") == "reasoning":
if "id" in value:
value.pop("id", None)
changed = True
if "encrypted_content" in value:
value.pop("encrypted_content", None)
changed = True
summary_parts = _third_party_reasoning_summary_parts(value.get("summary"))
content_parts = _third_party_reasoning_summary_parts(value.get("content"))
if not summary_parts:
summary_parts = content_parts
if value.get("content") != []:
value["content"] = []
changed = True
if value.get("summary") != summary_parts:
value["summary"] = summary_parts
changed = True
nested_inside_agent_message = (
inside_collaboration_agent_message
or (
preserve_collaboration_agent_message_encryption
and value.get("type") == "agent_message"
)
)
for key, nested in list(value.items()):
if (
isinstance(nested, dict)
and not preserve_collaboration_agent_message_encryption
):
rewritten_agent_message = _rewrite_third_party_encrypted_agent_message(nested)
if rewritten_agent_message is not None:
value[key] = rewritten_agent_message
nested = rewritten_agent_message
changed = True
if sanitize_third_party_reasoning_items(
nested,
preserve_collaboration_agent_message_encryption=preserve_collaboration_agent_message_encryption,
inside_collaboration_agent_message=nested_inside_agent_message,
):
changed = True
return changed
_sanitize_third_party_reasoning_items = sanitize_third_party_reasoning_items
def has_browser_context_signal(value: Any) -> bool:
for fragment in _stream_semantics.collect_text_fragments(value):
lowered = fragment.lower()
if any(marker in lowered for marker in BROWSER_CONTEXT_MARKERS):
return True
if BROWSER_CURRENT_URL_RE.search(fragment):
return True
return False
_has_browser_context_signal = has_browser_context_signal
def _has_browser_context_guidance(value: Any) -> bool:
return any(
BROWSER_CONTEXT_GUIDANCE_SENTINEL in fragment
for fragment in _stream_semantics.collect_text_fragments(value)
)
def reasoning_param_is_unsupported(upstream_name: Any, requested_model: Any, upstream_model: Any) -> bool:
if upstream_name == "official":
return False
controls = _thinking_controls_for(upstream_name, requested_model, upstream_model)
return bool(controls and controls.drop_reasoning_effort)
def _thinking_controls_for(
upstream_name: Any,
requested_model: Any,
upstream_model: Any,
*,
effort: str | None = None,
) -> maintained_catalog.ThinkingPayload | None:
provider_id = _UPSTREAM_TO_PROVIDER.get(str(upstream_name or ""))
if not provider_id:
return None
wire = ""
for model in (upstream_model, requested_model):
if isinstance(model, str) and model.strip():
wire = canonical_model_id(model)
break
if not wire:
return None
return maintained_catalog.thinking_payload(provider_id, wire, effort=effort)
def apply_maintained_thinking_controls(
payload: dict[str, Any],
upstream_name: Any,
requested_model: Any,
upstream_model: Any,
) -> bool:
"""Strip unsupported effort grades and attach vendor thinking JSON."""
inbound_effort = None
inbound_reasoning = payload.get("reasoning")
if isinstance(inbound_reasoning, dict) and isinstance(inbound_reasoning.get("effort"), str):
inbound_effort = inbound_reasoning["effort"]
elif isinstance(inbound_reasoning, str):
inbound_effort = inbound_reasoning
controls = _thinking_controls_for(
upstream_name, requested_model, upstream_model, effort=inbound_effort
)
if controls is None:
return False
changed = False
inbound_thinking = payload.get("thinking")
thinking_enabled = None
if isinstance(inbound_thinking, dict):
thinking_type = str(inbound_thinking.get("type") or "").strip().lower()
if thinking_type == "disabled":
thinking_enabled = False
elif thinking_type in {"enabled", "adaptive"}:
thinking_enabled = True
if thinking_enabled is not None:
controls = maintained_catalog.thinking_payload(
_UPSTREAM_TO_PROVIDER.get(str(upstream_name or "")),
canonical_model_id(str(upstream_model or requested_model or "")),
effort=None,
thinking_enabled=thinking_enabled,
)
if controls.drop_reasoning_effort:
for key in ("reasoning", "reasoning_effort"):
if key in payload:
del payload[key]
changed = True
template = payload.get("chat_template_kwargs")
if isinstance(template, dict) and "reasoning_effort" in template:
template.pop("reasoning_effort", None)
if not template:
payload.pop("chat_template_kwargs", None)
changed = True
if controls.thinking is not None and not isinstance(inbound_thinking, dict):
payload["thinking"] = dict(controls.thinking)
changed = True
if (
controls.reasoning_effort
and not controls.drop_reasoning_effort
and isinstance(payload.get("reasoning"), dict)
and payload["reasoning"].get("effort") != controls.reasoning_effort
):
payload["reasoning"]["effort"] = controls.reasoning_effort
changed = True
return changed
_reasoning_param_is_unsupported = reasoning_param_is_unsupported
_apply_maintained_thinking_controls = apply_maintained_thinking_controls
def _request_carries_reasoning_control(payload: Mapping[str, Any]) -> bool:
effort = payload.get("reasoning_effort")
if isinstance(effort, str) and effort:
return True
reasoning = payload.get("reasoning")
if isinstance(reasoning, str) and reasoning:
return True
if isinstance(reasoning, Mapping) and reasoning:
return True
template_kwargs = payload.get("chat_template_kwargs")
if isinstance(template_kwargs, Mapping):
template_effort = template_kwargs.get("reasoning_effort")
if isinstance(template_effort, str) and template_effort:
return True
return False
def reasoning_policy_for_request(
inbound_payload: Any,
upstream: Mapping[str, Any] | None,
model: str | None,
) -> str | None:
if not isinstance(inbound_payload, Mapping) or not isinstance(upstream, Mapping):
return None
if _request_carries_reasoning_control(inbound_payload):
return "explicit"
levels = upstream.get("supported_reasoning_levels")
if not levels and model:
candidate = _catalog.generated_catalog_by_slug().get(
_catalog.catalog_identity_slug(canonical_model_id(model))
)
if isinstance(candidate, Mapping):
levels = candidate.get("supported_reasoning_levels")
if levels:
return "provider-default"
return None
_reasoning_policy_for_request = reasoning_policy_for_request
def validate_reasoning_effort_for_upstream(
payload: Any,
upstream: Mapping[str, Any],
model: str | None,
) -> None:
if not isinstance(payload, Mapping):
return
requested_efforts = [payload.get("reasoning_effort")]
reasoning = payload.get("reasoning")
if isinstance(reasoning, Mapping):
requested_efforts.append(reasoning.get("effort"))
elif isinstance(reasoning, str):
requested_efforts.append(reasoning)
template_kwargs = payload.get("chat_template_kwargs")
if isinstance(template_kwargs, Mapping):
requested_efforts.append(template_kwargs.get("reasoning_effort"))
is_ultra = any(
isinstance(effort, str) and effort.strip().lower() == "ultra" for effort in requested_efforts
)
if not is_ultra:
return
is_official = upstream.get("name") == "official" and upstream.get("auth") == "codex_auth"
model_id = canonical_model_id(model or "").lower()
if model_id.startswith(OFFICIAL_ALIAS_PREFIX):
model_id = model_id[len(OFFICIAL_ALIAS_PREFIX) :]
if is_official and model_id in OFFICIAL_ULTRA_REASONING_MODELS:
return
if is_official:
raise ValueError(
"reasoning effort 'ultra' is supported only for gpt-5.6-sol and gpt-5.6-terra"
)
raise ValueError("reasoning effort 'ultra' is not supported for third-party models")
_validate_reasoning_effort_for_upstream = validate_reasoning_effort_for_upstream
def _bearer_token(headers: Mapping[str, str] | Any) -> str | None:
auth_header = _get_header(headers, "Authorization")
if not auth_header:
return None
value = auth_header.strip()
if not value:
return None
if value.lower().startswith("bearer "):
return value[7:].strip() or None
return value
def local_request_authorized(
headers: Mapping[str, str] | Any,
request_context: Mapping[str, str],
) -> bool:
expected_key = gateway_settings.gateway_client_key()
if expected_key is None:
return True
token = _bearer_token(headers)
claude_gateway_key = _get_header(headers, "x-codexhub-gateway-key")
return bool(
(token and hmac.compare_digest(token, expected_key))
or (
claude_gateway_key
and hmac.compare_digest(claude_gateway_key.strip(), expected_key)
)
)
_local_request_authorized = local_request_authorized
def _truthy_probe_value(value: str | None) -> bool:
return isinstance(value, str) and value.strip().lower() in {"1", "true", "yes", "on"}
def raw_provider_probe_requested(headers: Mapping[str, str] | Any, path: str) -> bool:
from urllib.parse import parse_qs, urlsplit
if _truthy_probe_value(_get_header(headers, "X-CodexHub-Raw-Provider-Probe")):
return True
try:
query_values = parse_qs(urlsplit(path).query, keep_blank_values=True)
except ValueError:
return False
return any(_truthy_probe_value(value) for value in query_values.get("raw_provider_probe", []))
def _header_tokens(headers: Mapping[str, str] | Any, name: str) -> set[str]:
value = _get_header(headers, name)
if not value:
return set()
return {token.strip().lower() for token in value.split(",") if token.strip()}
def is_websocket_upgrade(headers: Mapping[str, str] | Any) -> bool:
upgrade = _get_header(headers, "Upgrade")
if not upgrade or upgrade.lower() != "websocket":
return False
return "upgrade" in _header_tokens(headers, "Connection")
_is_websocket_upgrade = is_websocket_upgrade
def websocket_probe_frame_metadata(frame: Any) -> dict[str, Any]:
metadata: dict[str, Any] = {
"direction": "client_to_proxy",
"opcode": int(frame.opcode),
"fin": bool(frame.fin),
"payload_length": len(frame.payload),
"appears_json": False,
"json_top_level_keys": [],
}
if frame.opcode == 0x8:
metadata["close_code"] = int.from_bytes(frame.payload[:2], "big") if len(frame.payload) >= 2 else None
metadata["close_reason_length"] = max(0, len(frame.payload) - 2)
return metadata
if frame.opcode not in {0x1, 0x2}:
return metadata
try:
payload = json.loads(frame.payload.decode("utf-8-sig"))
except (UnicodeDecodeError, json.JSONDecodeError):
return metadata
metadata["appears_json"] = True
if isinstance(payload, Mapping):
metadata["json_top_level_keys"] = proxy_telemetry.protocol_field_names(payload)
metadata["unknown_json_key_count"] = len(payload) - len(metadata["json_top_level_keys"])
return metadata
_websocket_probe_frame_metadata = websocket_probe_frame_metadata
def request_context_from_headers(headers: Mapping[str, str] | Any) -> dict[str, str]:
context: dict[str, str] = {}
direct_headers = {
"x-codex-turn-id": "turn_id",
"x-codex-thread-id": "thread_id",
"x-codex-session-id": "session_id",
"x-codex-window-id": "window_id",
"x-codex-client-id": "client_id",
"x-request-id": "client_request_id",
"x-query-id": "query_id",
"x-session-id": "session_id",
"x-zcode-trace-id": "trace_id",
}
for header_name, field_name in direct_headers.items():
value = _get_header(headers, header_name)
if value:
context[field_name] = value[:200]
if field_name == "client_id":
context["client_inference_source"] = "header"
for header_name in ("x-codex-client-metadata", "x-codex-metadata"):
value = _get_header(headers, header_name)
if not value:
continue
try:
metadata = json.loads(value)
except json.JSONDecodeError:
continue
if not isinstance(metadata, dict):
continue
for key in (
"client_id",
"session_id",
"thread_id",
"turn_id",
"window_id",
"request_kind",
"thread_source",
):
item = metadata.get(key)
if isinstance(item, str) and item and key not in context:
context[key] = item[:200]
if key == "client_id":
context["client_inference_source"] = "metadata"
# Claude's explicit session header takes precedence over generic session hints.
# Preserve its source so another client using the same value cannot share tools.
claude_session = (_get_header(headers, "x-claude-code-session-id") or "").strip()
if claude_session:
context["session_id"] = claude_session[:200]
context["session_source"] = "claude-code"
user_agent = _get_header(headers, "User-Agent")
if user_agent:
context["user_agent_hash"] = proxy_telemetry.telemetry_hmac(
gateway_events.RUNTIME_CODEX_DIR,
b"user-agent",
user_agent[:500].encode("utf-8", errors="ignore"),
)
if "client_id" not in context:
inferred = _infer_client_id(user_agent)
if inferred:
context["client_id"] = inferred
context["client_inference_source"] = "user_agent"
context.setdefault("client_id", "unknown")
context.setdefault("client_inference_source", "unknown")
return context
def _infer_client_id(user_agent: str | None) -> str | None:
if not user_agent:
return None
value = user_agent.lower()
if "opencode" in value:
return "opencode"
if "zcode" in value:
return "zcode"
if "omp" in value:
return "omp"
if "codex desktop/" in value or "codex-app" in value:
return "codex-app"
return None
def request_observability_with_prefix(fields: Mapping[str, Any], prefix: str) -> dict[str, Any]:
renamed: dict[str, Any] = {}
for key, value in fields.items():
if key == "request_body_hmac":
renamed[f"{prefix}_request_body_hmac"] = value
elif key == "request_body_hmac_skipped":
renamed[f"{prefix}_request_body_hmac_skipped"] = value
elif key == "request_prefix_hmac":
renamed[f"{prefix}_request_prefix_hmac"] = value
elif key == "prefix_bytes":
renamed[f"{prefix}_prefix_bytes"] = value
elif key == "prompt_cache_key_hash":
renamed[f"{prefix}_prompt_cache_key_hash"] = value
elif key == "prompt_cache_key_state":
renamed[f"{prefix}_prompt_cache_key_state"] = value
elif key == "body_bytes":
renamed[f"{prefix}_body_bytes"] = value
elif key == "body_sha256":
renamed[f"{prefix}_body_sha256"] = value
return renamed
_request_observability_with_prefix = request_observability_with_prefix
def is_event_stream(headers: Mapping[str, str] | Any) -> bool:
content_type = _get_header(headers, "Content-Type")
if content_type and "text/event-stream" in content_type.lower():
return True
if content_type and "json" in content_type.lower():
return False
# Some upstreams (e.g. chatgpt.com/backend-api/codex) return SSE without
# an explicit Content-Type header but do signal chunked transfer.
transfer_encoding = _get_header(headers, "Transfer-Encoding")
return bool(transfer_encoding and "chunked" in transfer_encoding.lower())
_is_event_stream = is_event_stream
UNSET_CONTENT_ENCODING = object()
_UNSET_CONTENT_ENCODING = UNSET_CONTENT_ENCODING
def filtered_response_headers(
headers: Mapping[str, str] | Any,
is_event_stream: bool,
content_length: int | None = None,
content_type: str | None = None,
content_encoding: str | None | object = _UNSET_CONTENT_ENCODING,
) -> list[tuple[str, str]]:
outgoing: list[tuple[str, str]] = []
for key, value in _header_items(headers):
lowered = key.lower()
if lowered in HOP_BY_HOP_RESPONSE_HEADERS:
continue
if lowered == "content-length" and (is_event_stream or content_length is not None):
continue
if lowered == "content-type" and content_type is not None:
continue
if lowered == "content-encoding" and content_encoding is not _UNSET_CONTENT_ENCODING:
continue
outgoing.append((key, value))
if content_type is not None:
outgoing.append(("Content-Type", content_type))
if content_length is not None:
outgoing.append(("Content-Length", str(content_length)))
if content_encoding is not _UNSET_CONTENT_ENCODING and isinstance(content_encoding, str) and content_encoding:
outgoing.append(("Content-Encoding", content_encoding))
return outgoing
_filtered_response_headers = filtered_response_headers
def json_response_bytes(payload: dict[str, Any]) -> bytes:
return json.dumps(payload, ensure_ascii=True).encode("utf-8")
_json_response_bytes = json_response_bytes
OFFICIAL_IMAGE_POST_PATHS: Mapping[str, str] = {
"/v1/images/generations": "/images/generations",
"/v1/images/edits": "/images/edits",
"/v1/images/variations": "/images/variations",
}
def official_image_upstream_path(path: str) -> str | None:
"""Return the Official Images suffix for a Gateway POST, or None."""
return OFFICIAL_IMAGE_POST_PATHS.get(path)
def provider_scoped_path(path: str, endpoint_suffix: str) -> str | None:
from urllib.parse import unquote
prefix = "/v1/providers/"
suffix = "/" + endpoint_suffix.strip("/")
if not path.startswith(prefix) or not path.endswith(suffix):
return None
provider_part = path[len(prefix) : -len(suffix)]
if not provider_part or "/" in provider_part:
return None
provider = unquote(provider_part).strip()
return provider or None
def provider_scoped_route_model(model_id: str | None, provider_hint: str | None) -> str | None:
if not model_id:
return None
slug = canonical_model_id(str(model_id))
if not slug or not provider_hint:
return slug
provider = canonical_model_id(str(provider_hint))
if not provider or slug.startswith(f"{provider}/"):
return slug
return f"{provider}/{slug}"
def value_contains_image(value: Any) -> bool:
return vision_proxy.value_contains_image(value)
_value_contains_image = value_contains_image
def enforce_text_only_image_boundary(
payload: dict[str, Any],
*,
inbound_format: str,
target_model: str | None,
target_upstream: Mapping[str, Any],
vision_plan: Any,
event_context: Mapping[str, Any] | None = None,
progress_callback: Any = None,
) -> bool:
from gateway_errors import ImageProxyError
try:
inbound_protocol = RouteProtocol(inbound_format)
except ValueError as exc:
raise ImageProxyError("Vision Proxy received an unsupported inbound protocol") from exc
return vision_proxy.enforce_text_only_boundary(
payload,
inbound_protocol=inbound_protocol,
target_model=target_model,
target_upstream=target_upstream,
vision_plan=vision_plan,
event_context=event_context,
progress_callback=progress_callback,
)