#!/usr/bin/env python3 """ Proxy: OpenAI /v1/chat/completions -> CommandCode /alpha/generate Python port of server.js. Standard library only. All settings live in config.json next to this file. """ import http.client import json import os import platform import random import re import signal import socket import subprocess import sys import threading import time from datetime import datetime, timezone from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer BASE_DIR = os.path.dirname(os.path.abspath(__file__)) # -- config ------------------------------------------------------------------- DEFAULTS = { 'port': 3456, 'cc_version': '1.15.1', 'debug': False, 'upstream_host': 'api.commandcode.ai', 'upstream_path': '/alpha/generate', 'models_path': '/provider/v1/models', 'timeout_seconds': 300, 'log_file': 'proxy.log', 'max_body_bytes': 10 * 1024 * 1024, 'default_model': 'deepseek/deepseek-v4-pro', 'default_max_tokens': 32000, 'max_retries': 2, 'retry_max_delay_seconds': 60, 'fallback_models': [], } def load_config(): """Shallow-merge config.json over DEFAULTS. Missing file is fine, bad JSON is not.""" path = os.path.join(BASE_DIR, 'config.json') cfg = dict(DEFAULTS) if not os.path.exists(path): return cfg, 'config.json not found, using defaults' try: with open(path, 'r', encoding='utf-8') as f: loaded = json.load(f) except (ValueError, OSError) as e: sys.stderr.write(f'failed to read config.json: {e}\n') sys.exit(1) if not isinstance(loaded, dict): sys.stderr.write('config.json must contain a JSON object\n') sys.exit(1) for k, v in loaded.items(): if k in cfg: cfg[k] = v return cfg, None CONFIG, CONFIG_NOTE = load_config() PORT = int(CONFIG['port']) HOST = CONFIG['upstream_host'] PATH = CONFIG['upstream_path'] MODELS_PATH = CONFIG['models_path'] CC_VERSION = CONFIG['cc_version'] DEBUG = bool(CONFIG['debug']) TIMEOUT = float(CONFIG['timeout_seconds']) MAX_BODY = int(CONFIG['max_body_bytes']) DEFAULT_MODEL = CONFIG['default_model'] DEFAULT_MAX_TOKENS = int(CONFIG['default_max_tokens']) MAX_RETRIES = int(CONFIG['max_retries']) MAX_RETRY_DELAY = float(CONFIG['retry_max_delay_seconds']) FALLBACK_MODELS = CONFIG['fallback_models'] # -- secret redaction --------------------------------------------------------- # Upstream error bodies can echo the caller's credentials; they are both written to # proxy.log and forwarded to the client, so scrub them first. # Ported from pi-commandcode-provider src/overflow.ts. _REDACT = [ (re.compile(r'\bBearer\s+[A-Za-z0-9._~+/=-]+', re.I), 'Bearer [redacted]'), (re.compile(r'\b(?:user|cc)_[A-Za-z0-9_-]{8,}\b', re.I), '[redacted]'), (re.compile(r'([?&](?:api[-_ ]?key|apikey|access_token|refresh_token|token|secret|' r'password)=)[^&#\s]+', re.I), r'\1[redacted]'), (re.compile(r'\b(?:sk|rk|ghp|github_pat|xox[baprs])[-_A-Za-z0-9]{16,}\b' r'|\beyJ[A-Za-z0-9_-]{20,}\.[A-Za-z0-9_-]{10,}\.[A-Za-z0-9_-]{10,}\b'), '[redacted]'), ] _REDACT_KV = re.compile( r'\b(?:api[-_ ]?key|apikey|access[-_ ]?token|refresh[-_ ]?token|token|secret|password|' r'authorization)\s*[=:]\s*[^\s,;)"\']+', re.I) def redact(text): """Strip credentials from arbitrary text before logging or forwarding it.""" if not text: return text for pattern, repl in _REDACT: text = pattern.sub(repl, text) def _kv(m): s = m.group(0) i = min((s.find(c) for c in '=:' if c in s), default=-1) return s if i < 0 else s[:i + 1] + '[redacted]' return _REDACT_KV.sub(_kv, text) # -- logging ------------------------------------------------------------------ C = { 'reset': '\x1b[0m', 'cyan': '\x1b[36m', 'green': '\x1b[32m', 'yellow': '\x1b[33m', 'dim': '\x1b[2m', 'bold': '\x1b[1m', } LEVEL_COLORS = {'req': C['cyan'], 'upstream': C['green'], 'done': C['bold'], 'error': C['yellow']} _log_lock = threading.Lock() _log_file = open(os.path.join(BASE_DIR, CONFIG['log_file']), 'w', encoding='utf-8') def _now_iso(): now = datetime.now(timezone.utc) return now.strftime('%Y-%m-%dT%H:%M:%S.') + f'{now.microsecond // 1000:03d}Z' def write_log(level, msg): msg = redact(msg) ts = f'[{_now_iso()}]' c = LEVEL_COLORS.get(level, C['reset']) tag = f'{c}{level}{C["reset"]}' if level else '' full = f'{C["dim"]}{ts}{C["reset"]} ' + (f'[{tag}] ' if tag else '') + msg plain = f'{ts}' + (f' [{level}] ' if tag else ' ') + msg with _log_lock: sys.stdout.write(full + '\n') sys.stdout.flush() _log_file.write(plain + '\n') _log_file.flush() def log_req(msg): write_log('req', msg) def log_up(msg): write_log('upstream', msg) def log_done(msg): write_log('done', msg) def log_err(msg): write_log('error', msg) def log(msg): write_log('', msg) _ANSI = re.compile(r'\x1b\[[0-9;]*m') def banner(): def strip(s): return _ANSI.sub('', s) def pad(s, w): return s + ' ' * max(0, w - len(strip(s))) def L(s): return f'{C["bold"]}{s}{C["reset"]}' title = f'{C["cyan"]}{C["bold"]}Proxy CommandCode{C["reset"]}' sub = f'{C["dim"]}OpenAI → CommandCode /alpha/generate{C["reset"]}' rows = [ f'{L("Listening")} http://localhost:{PORT}', f'{L("Endpoint")} /v1/chat/completions', f'{L("Upstream")} {HOST}{PATH}', f'{L("CC Version")} {CC_VERSION}', f'{L("Debug")} ' + (f'{C["green"]}ON{C["reset"]}' if DEBUG else f'{C["yellow"]}OFF{C["reset"]}'), ] all_rows = [title, sub] + rows w = max(len(strip(s)) for s in all_rows) def box(s): return f'{C["cyan"]}│{C["reset"]} {pad(s, w)} {C["cyan"]}│{C["reset"]}' print(f'{C["cyan"]}┌{"─" * (w + 4)}┐{C["reset"]}') print(box(title)) print(box(sub)) print(f'{C["cyan"]}├{"─" * (w + 4)}┤{C["reset"]}') for r in rows: print(box(r)) print(f'{C["cyan"]}└{"─" * (w + 4)}┘{C["reset"]}') # -- transform ---------------------------------------------------------------- CORS = { 'Access-Control-Allow-Origin': '*', 'Access-Control-Allow-Headers': 'Content-Type, Authorization', } STATIC_CONFIG = { 'workingDir': '', 'date': '', 'environment': '', 'structure': [], 'isGitRepo': False, 'currentBranch': '', 'mainBranch': 'main', 'gitStatus': '', 'recentCommits': [], } ENVIRONMENT_INFO = f'{sys.platform}-{platform.machine()}, Python {platform.python_version()}' # Sampling knobs the OpenAI request may carry. Forwarded verbatim into params when present; # CommandCode ignores what it does not understand rather than rejecting the request. PASSTHROUGH_PARAMS = ('temperature', 'top_p', 'top_k', 'stop', 'seed', 'presence_penalty', 'frequency_penalty', 'reasoning_effort') def _js_str(v): """String(v) the way JavaScript does it for the values we actually see.""" if v is None: return 'null' if v is True: return 'true' if v is False: return 'false' if isinstance(v, str): return v return json.dumps(v, ensure_ascii=False) def _dumps(obj): return json.dumps(obj, separators=(',', ':'), ensure_ascii=False) def _today(): return datetime.now(timezone.utc).strftime('%Y-%m-%d') def _tool_input(arguments): """OpenAI streams tool arguments as a JSON string; CommandCode wants the object. Fragments are not always complete JSON, so fall back to the raw string. (pi-commandcode-provider `recordOrEmpty`, src/converters.ts:23) """ if isinstance(arguments, dict): return arguments if isinstance(arguments, str): try: parsed = json.loads(arguments) if isinstance(parsed, dict): return parsed except ValueError: pass return arguments def _image_part(url): """CommandCode expects {image: , mimeType}, not OpenAI's {url}. (pi-commandcode-provider `imageToCommandCode`, src/converters.ts:76) """ if not isinstance(url, str): return None m = re.match(r'data:([^;,]+);base64,', url) if m: return {'type': 'image', 'image': url, 'mimeType': m.group(1)} # Remote URLs are passed through as-is; CommandCode may or may not fetch them, # but rewriting them into a data URI would mean downloading on its behalf. return {'type': 'image', 'image': url} def _apply_tool_choice(tools, value): """OpenAI `tool_choice` -> (tools to send, params.tool_choice to send). Upstream accepts exactly one value, `{"type": "auto"}` — anything else answers `expected "auto" at "params.tool_choice.type"`, and a bare string answers `expected object, received string`. So only `auto` is expressible upstream; the rest are emulated by shaping the tool list, which is what actually constrains the model: none -> send no tools at all; a tool call becomes impossible (exact) {function: "X"} -> send only X; strongly biases toward it (not a guarantee) required -> not expressible, and not emulable; left to the model """ if value is None: return tools, None if isinstance(value, dict) and value.get('type') == 'function': name = (value.get('function') or {}).get('name') if name: picked = [t for t in tools if t.get('name') == name] return (picked or tools), None value = 'auto' if isinstance(value, dict): value = value.get('type') # already upstream-shaped, or OpenAI's object form if value == 'auto': return tools, {'type': 'auto'} if value == 'none': return [], None if value == 'required': return tools, None # upstream cannot force a call; the model decides log_err(f'[tool_choice] unrecognised, ignoring: {_dumps(value)[:120]}') return tools, None def _paired_tool_call_ids(src): """Tool call ids that have a matching result, and vice versa. An assistant tool_call with no tool result (or a result with no call) makes the upstream reject the whole request; editors truncate history and produce exactly that. (pi-commandcode-provider `completeToolCallIds`, src/converters.ts:162) """ calls, results = set(), set() for m in src: if m.get('role') == 'assistant': for tc in m.get('tool_calls') or []: if tc.get('id'): calls.add(tc['id']) elif m.get('role') == 'tool' and m.get('tool_call_id'): results.add(m['tool_call_id']) return calls & results def transform(oai_body): """OpenAI request body -> CommandCode /alpha/generate envelope (bytes).""" model = oai_body.get('model') or DEFAULT_MODEL system_parts = [] messages = [] src = oai_body.get('messages') or [] paired = _paired_tool_call_ids(src) # tool_call_id -> tool name, so tool results can be attributed tool_name_map = {} for m in src: if m.get('role') == 'assistant' and m.get('tool_calls'): for tc in m['tool_calls']: fn = tc.get('function') or {} if tc.get('id') and fn.get('name'): tool_name_map[tc['id']] = fn['name'] for m in src: role = m.get('role') content = m.get('content') if role == 'system': if isinstance(content, str): system_parts.append(content) elif isinstance(content, list): system_parts.append('\n'.join( p.get('text', '') for p in content if isinstance(p, dict) and p.get('type') == 'text')) else: system_parts.append(_js_str(content)) continue if role == 'tool': if m.get('tool_call_id') not in paired: continue # orphan result: upstream rejects the whole request if isinstance(content, str): out = {'type': 'text', 'value': content} elif content: out = content else: out = {'type': 'text', 'value': _js_str(content)} messages.append({'role': 'tool', 'content': [{ 'type': 'tool-result', 'toolCallId': m.get('tool_call_id'), 'toolName': tool_name_map.get(m.get('tool_call_id')) or 'unknown', 'output': out, }]}) continue if role == 'assistant': parts = [] if content: if isinstance(content, str): parts.append({'type': 'text', 'text': content}) elif isinstance(content, list): for p in content: if isinstance(p, dict) and p.get('type') == 'text': parts.append({'type': 'text', 'text': p.get('text')}) for tc in m.get('tool_calls') or []: if tc.get('type') == 'function' and tc.get('function'): if tc.get('id') not in paired: continue # unanswered call: same rejection risk as above parts.append({ 'type': 'tool-call', 'toolCallId': tc.get('id'), 'toolName': tc['function'].get('name'), 'input': _tool_input(tc['function'].get('arguments')), }) if parts: messages.append({'role': 'assistant', 'content': parts}) continue if isinstance(content, str): messages.append({'role': role, 'content': [{'type': 'text', 'text': content}]}) elif isinstance(content, list): parts = [] for p in content: if not isinstance(p, dict): continue if p.get('type') == 'text': parts.append({'type': 'text', 'text': p.get('text')}) elif p.get('type') == 'image_url': img = _image_part((p.get('image_url') or {}).get('url')) if img: parts.append(img) messages.append({'role': role, 'content': parts}) else: messages.append({'role': role, 'content': [{'type': 'text', 'text': _js_str(content)}]}) tools = [] for t in oai_body.get('tools') or []: fn = t.get('function') or {} tools.append({ 'type': 'function', 'name': fn.get('name') or t.get('name'), 'description': fn.get('description') or t.get('description') or '', 'input_schema': fn.get('parameters') or t.get('input_schema') or {'type': 'object', 'properties': {}}, }) tools, tool_choice = _apply_tool_choice(tools, oai_body.get('tool_choice')) system_text = '\n\n'.join(p for p in system_parts if p) params = {'model': model} if system_text: params['system'] = system_text params['messages'] = messages if tools: params['tools'] = tools if tool_choice: params['tool_choice'] = tool_choice params['max_tokens'] = oai_body.get('max_tokens') or DEFAULT_MAX_TOKENS # Always stream upstream, whatever the client asked for. `stream: false` makes the # endpoint answer "Proxy use detected. This endpoint only serves CLI." — the real CLI # never sends it (pi-commandcode-provider hardcodes stream: true, src/core.ts:495). # A non-streaming client is served by buffering the NDJSON here instead. params['stream'] = True # All eight confirmed accepted by /alpha/generate (probed individually); the endpoint # normalises most of them away, but none of them 400. for key in PASSTHROUGH_PARAMS: if oai_body.get(key) is not None: params[key] = oai_body[key] cfg = dict(STATIC_CONFIG) cfg['date'] = _today() cfg['environment'] = ENVIRONMENT_INFO # No threadId: measured A/B (stable vs random) showed identical cache growth (+640 # cached tokens either way), so upstream caching is content-based and the field only # adds a UUID-validation failure mode. return _dumps({ 'config': cfg, 'memory': None, 'taste': None, 'skills': None, 'permissionMode': 'standard', 'params': params, }).encode('utf-8') # -- upstream finish event ---------------------------------------------------- ZERO_USAGE = {'prompt_tokens': 0, 'completion_tokens': 0, 'total_tokens': 0} def usage_from_finish(evt): """CommandCode `finish` -> OpenAI `usage`. Live shape (confirmed from dump/): {"type":"finish","finishReason":"stop","rawFinishReason":"stop", "totalUsage":{"inputTokens":8233,"outputTokens":35,"totalTokens":8268, "inputTokenDetails":{"noCacheTokens":553,"cacheReadTokens":7680}, "outputTokenDetails":{"textTokens":35,"reasoningTokens":0}, "cachedInputTokens":7680}} """ tu = evt.get('totalUsage') if not isinstance(tu, dict): return None prompt = tu.get('inputTokens') or 0 completion = tu.get('outputTokens') or 0 usage = { 'prompt_tokens': prompt, 'completion_tokens': completion, 'total_tokens': tu.get('totalTokens') or (prompt + completion), } details = tu.get('inputTokenDetails') or {} cached = tu.get('cachedInputTokens') if cached is None: cached = details.get('cacheReadTokens') if cached is not None: usage['prompt_tokens_details'] = {'cached_tokens': cached} out_details = tu.get('outputTokenDetails') or {} if out_details.get('reasoningTokens') is not None: usage['completion_tokens_details'] = {'reasoning_tokens': out_details['reasoningTokens']} return usage def finish_reason_from(evt, had_tool_calls): """`rawFinishReason` is already OpenAI-spelled; `finishReason` uses dashes. (pi-commandcode-provider `mapFinishReason`, src/converters.ts:258) """ raw = evt.get('rawFinishReason') if raw in ('stop', 'length', 'tool_calls', 'content_filter', 'function_call'): return raw reason = evt.get('finishReason') if reason in ('tool-calls', 'tool_calls'): return 'tool_calls' if reason in ('length', 'max_tokens', 'max-tokens', 'max_output_tokens'): return 'length' if reason == 'stop': return 'tool_calls' if had_tool_calls else 'stop' return None def openai_error(message, code=None, etype='upstream_error'): """OpenAI-shaped error envelope. Clients read error.message; give them that field.""" return {'error': {'message': redact(str(message)), 'type': etype, 'param': None, 'code': code}} # -- retry -------------------------------------------------------------------- # Ported from pi-commandcode-provider src/core.ts:53-85. Only applied before any byte # reaches the client — retrying mid-stream would duplicate half a response. BASE_RETRY_DELAY = 0.5 def is_retryable(status): return status == 429 or 500 <= status < 600 def retry_delay(attempt, retry_after): """Honour Retry-After when present, else exponential backoff with jitter. Returns -1 when the server asks for longer than we are willing to wait. """ if retry_after: try: seconds = float(retry_after) except ValueError: seconds = None if seconds is not None and seconds >= 0: return -1 if seconds > MAX_RETRY_DELAY else seconds exponential = BASE_RETRY_DELAY * (2 ** attempt) return min(exponential + exponential * 0.2 * random.random(), MAX_RETRY_DELAY) def _close(conn): try: if conn: conn.close() except Exception: pass # -- handler ------------------------------------------------------------------ CLIENT_GONE = (BrokenPipeError, ConnectionResetError, ConnectionAbortedError, OSError) class ProxyHandler(BaseHTTPRequestHandler): protocol_version = 'HTTP/1.1' disable_nagle_algorithm = True server_version = 'proxy-commandcode' sys_version = '' timeout = TIMEOUT def log_message(self, fmt, *args): pass # our own logging only def handle_one_request(self): # A client that walks away mid-stream (editor cancels a completion) is routine; # the base class' trailing wfile.flush() would raise it up to socketserver. try: super().handle_one_request() except CLIENT_GONE: self.close_connection = True # -- response helpers -- def _cors_headers(self): for k, v in CORS.items(): self.send_header(k, v) def _send_bytes(self, status, body, content_type='application/json', extra=None): if self._headers_sent: return self._headers_sent = True try: self.send_response(status) self._cors_headers() for k, v in (extra or {}).items(): self.send_header(k, v) if body and content_type: self.send_header('Content-Type', content_type) self.send_header('Content-Length', str(len(body or b''))) self.end_headers() if body: self.wfile.write(body) except CLIENT_GONE: self.close_connection = True finally: self._log_response(status, content_type, len(body or b'')) def _send_json(self, status, obj, extra=None): self._send_bytes(status, _dumps(obj).encode('utf-8'), 'application/json', extra) def _log_response(self, status, content_type, nbytes): """One compact line per non-streaming response, only when debug is on.""" if not DEBUG: return ip = self.client_address[0] if self.client_address else '-' ms = int((time.time() - self._t0) * 1000) log_req(f'[detail] {self.command or "-"} {self.path or "-"} | {ip} | {status} | ' f'{content_type or "-"} | {nbytes} B | {ms}ms') def _debug_detail(self, oai, upstream): """Headers, parsed request body, and both params shapes, pretty-printed.""" if not DEBUG: return lines = ['[detail] request', ' -- headers ' + '-' * 40] for k, v in (self.headers.items() if self.headers else []): lines.append(f' {k}: {v}') lines.append(' -- request body ' + '-' * 36) lines.append(' ' + json.dumps(oai, indent=2, ensure_ascii=False).replace('\n', '\n ')) lines.append(' -- params (original OpenAI) ' + '-' * 25) orig = {k: v for k, v in oai.items() if k in ('model', 'max_tokens', 'stream', 'stream_options') or k in PASSTHROUGH_PARAMS} lines.append(' ' + json.dumps(orig, indent=2, ensure_ascii=False).replace('\n', '\n ')) lines.append(' -- params (transformed upstream) ' + '-' * 21) try: transformed = json.loads(upstream.decode('utf-8')).get('params', {}) except (ValueError, UnicodeDecodeError, AttributeError): transformed = {} lines.append(' ' + json.dumps(transformed, indent=2, ensure_ascii=False).replace('\n', '\n ')) log_req('\n'.join(lines)) def _begin_stream(self, content_type='text/event-stream'): """HTTP/1.1 chunked: BaseHTTPRequestHandler will not frame writes for us.""" if self._headers_sent: return self._headers_sent = True self._streaming = True try: self.send_response(200) self._cors_headers() self.send_header('Content-Type', content_type) self.send_header('Cache-Control', 'no-cache') self.send_header('Transfer-Encoding', 'chunked') self.end_headers() except CLIENT_GONE: self._client_gone = True self.close_connection = True def _write_chunk(self, data): if self._client_gone or not data: return try: self.wfile.write(f'{len(data):X}\r\n'.encode('ascii') + data + b'\r\n') self.wfile.flush() except CLIENT_GONE: self._client_gone = True self.close_connection = True def _write_sse(self, obj): self._write_chunk(f'data: {_dumps(obj)}\n\n'.encode('utf-8')) def _end_stream(self): if not self._streaming or self._stream_ended: return self._stream_ended = True if self._client_gone: return try: self.wfile.write(b'0\r\n\r\n') self.wfile.flush() except CLIENT_GONE: self._client_gone = True self.close_connection = True def _reset(self): self._headers_sent = False self._streaming = False self._stream_ended = False self._client_gone = False self._t0 = time.time() def _read_body(self): """Always consume the body, even on an early error — an undrained body would be parsed as the next request on this keep-alive connection. Returns None if the body is over the limit.""" if (self.headers.get('Transfer-Encoding') or '').lower() == 'chunked': return self._read_chunked_body() try: length = int(self.headers.get('Content-Length') or 0) except ValueError: length = 0 if length > MAX_BODY: self.close_connection = True return None try: return self.rfile.read(length) if length else b'' except CLIENT_GONE: self.close_connection = True return None def _read_chunked_body(self): """Not every client sends Content-Length; Node's http server decodes this for free.""" parts, total = [], 0 try: while True: size = int(self.rfile.readline(64).split(b';')[0].strip() or b'0', 16) if size == 0: while self.rfile.readline(65536).strip(): pass # trailers break total += size if total > MAX_BODY: self.close_connection = True return None parts.append(self.rfile.read(size)) self.rfile.read(2) # trailing CRLF except (ValueError,) + CLIENT_GONE: self.close_connection = True return b'' # malformed framing -> let it fail as a bad request return b''.join(parts) # -- routes -- def do_OPTIONS(self): self._reset() self._read_body() self._send_bytes(204, b'', None, { 'Access-Control-Allow-Methods': 'POST,GET,OPTIONS', 'Access-Control-Max-Age': '86400', }) def do_GET(self): self._reset() self._read_body() if self.path == '/health': self._send_json(200, {'status': 'ok'}) elif self.path.rstrip('/').endswith('/v1/models') or self.path.rstrip('/') == '/models': self._models() else: self._send_json(404, openai_error('POST /v1/chat/completions', code='not_found', etype='invalid_request_error')) def _models(self): """Many clients call /v1/models before they will talk to an endpoint at all. CommandCode's catalog already answers in OpenAI's {object:"list", data:[...]} shape, but it lives on the Pro-only /provider surface — so fall back to config.json's fallback_models when it refuses. """ auth = self.headers.get('Authorization') or '' try: conn = http.client.HTTPSConnection(HOST, timeout=30) conn.request('GET', MODELS_PATH, headers={ 'Authorization': auth, 'x-command-code-version': CC_VERSION}) r = conn.getresponse() body = r.read() status = r.status conn.close() except Exception as e: log_err(f'[models] {e}') status, body = 0, b'' log_up(f'{status or "ERR"} | models') if status == 200: self._send_bytes(200, body) return now = int(time.time()) self._send_json(200, {'object': 'list', 'data': [ {'id': m, 'object': 'model', 'created': now, 'owned_by': 'commandcode'} for m in FALLBACK_MODELS or [DEFAULT_MODEL]]}) def do_POST(self): self._reset() body = self._read_body() if body is None: self._send_bytes(413, b'{}') return if not self.path.startswith('/v1/chat/completions'): self._send_json(404, openai_error('POST /v1/chat/completions', code='not_found', etype='invalid_request_error')) return try: oai = json.loads(body.decode('utf-8')) except (ValueError, UnicodeDecodeError): self._send_json(400, openai_error('Invalid JSON', code='invalid_json', etype='invalid_request_error')) return model = oai.get('model') or '-' is_stream = oai.get('stream') is True ip = self.client_address[0] if self.client_address else '-' t0 = time.time() log_req(f'{model} | {ip} | {"stream" if is_stream else "sync"} | {len(body)} bytes') debug_transform_logs = ['[detail] request', ' -- headers (org) ' + '-' * 40] for k, v in (self.headers.items() if self.headers else []): debug_transform_logs.append(f' {k}: {v}') debug_transform_logs.append(' -- body (org) ' + '-' * 36) debug_transform_logs.append(' ' + json.dumps(oai, indent=2, ensure_ascii=False).replace('\n', '\n ')) try: upstream = transform(oai) except Exception: self._send_json(500, openai_error('Transform error', code='transform_error', etype='proxy_error')) return include_usage = bool((oai.get('stream_options') or {}).get('include_usage')) auth = self.headers.get('Authorization') or '' headers = { 'Content-Type': 'application/json', 'Content-Length': str(len(upstream)), 'Authorization': auth, 'x-command-code-version': CC_VERSION, 'x-cli-environment': 'production', 'x-project-slug': 'project', 'x-taste-learning': 'true', 'x-co-flag': 'false', } debug_transform_logs.append(' -- headers (xform) ' + '-' * 40) for k, v in (headers.items() if headers else []): debug_transform_logs.append(f' {k}: {v}') debug_transform_logs.append(' -- body (xform) ' + '-' * 36) transformed = json.loads(upstream.decode('utf-8')) debug_transform_logs.append(' ' + json.dumps(transformed, indent=2, ensure_ascii=False).replace('\n', '\n ')) log_req('\n'.join(debug_transform_logs)) for attempt in range(MAX_RETRIES + 1): conn, resp = None, None try: conn = http.client.HTTPSConnection(HOST, timeout=TIMEOUT) conn.request('POST', PATH, body=upstream, headers=headers) resp = conn.getresponse() except socket.timeout: _close(conn) if attempt < MAX_RETRIES: log_err(f'[upstream] timeout, retry {attempt + 1}/{MAX_RETRIES}') continue log_err('[upstream] timeout') self._send_json(504, openai_error('Upstream timed out', code='timeout')) return except Exception as e: _close(conn) if attempt < MAX_RETRIES: log_err(f'[upstream] {e}, retry {attempt + 1}/{MAX_RETRIES}') time.sleep(retry_delay(attempt, None)) continue log_err(f'[upstream] {e}') self._send_json(502, openai_error(e, code='upstream_unreachable')) return ok = 200 <= resp.status < 300 log_up(f'{resp.status} {"OK" if ok else "ERR"} | {model} | ' f'{int((time.time() - t0) * 1000)}ms') # 429/5xx are worth another go; nothing has reached the client yet. if is_retryable(resp.status) and attempt < MAX_RETRIES: wait = retry_delay(attempt, resp.getheader('Retry-After')) if wait >= 0: try: resp.read() except Exception: pass _close(conn) log_err(f'[upstream] {resp.status}, retry {attempt + 1}/{MAX_RETRIES} ' f'in {wait:.1f}s') time.sleep(wait) continue try: self._handle_upstream(resp, model, is_stream, t0, include_usage) finally: _close(conn) return # -- upstream response -- def _handle_upstream(self, resp, model, is_stream, t0, include_usage=False): if resp.status >= 400: try: body = resp.read().decode('utf-8', 'replace') except Exception: body = '' # CommandCode answers {"success":false,"error":{...}}; clients expect OpenAI's # {"error":{"message":...}}. Reshape, and redact any echoed credential. message, code = body, None try: parsed = json.loads(body) err = parsed.get('error') if isinstance(parsed, dict) else None if isinstance(err, dict): message = err.get('message') or body code = err.get('code') elif isinstance(err, str): message = err except ValueError: pass self._send_json(resp.status, openai_error(message[:2000], code=code)) return gen_id = 'chatcmpl-' + str(int(time.time() * 1000)) dump = None if DEBUG: os.makedirs(os.path.join(BASE_DIR, 'dump'), exist_ok=True) dump = open(os.path.join(BASE_DIR, 'dump', f'dump-{gen_id}.txt'), 'wb') log(f'[debug] dumping to dump/dump-{gen_id}.txt') try: if is_stream: self._stream_response(resp, model, gen_id, t0, dump, include_usage) else: self._buffer_response(resp, model, gen_id, t0, dump) except socket.timeout: log_err('[upstream] timeout') if self._streaming: self._end_stream() else: self._send_json(504, openai_error('Upstream timed out', code='timeout')) except CLIENT_GONE as e: log_err(f'[upstream] {e}') if self._streaming: self._end_stream() else: self._send_json(502, openai_error(e, code='upstream_error')) finally: if dump: dump.close() @staticmethod def _events(resp, dump): """CommandCode streams NDJSON; readline yields one event as soon as it lands. Tolerates SSE-style `data:` prefixes and `[DONE]` sentinels, which the upstream emits on some routes (pi-commandcode-provider `parseStreamEventLine`). """ for raw in resp: if dump: dump.write(raw) line = raw.strip() if not line or line.startswith(b':') or line.startswith(b'event:'): continue if line.startswith(b'data:'): line = line[5:].strip() if not line or line == b'[DONE]': continue try: yield json.loads(line.decode('utf-8')) except (ValueError, UnicodeDecodeError): continue def _buffer_response(self, resp, model, gen_id, t0, dump): full_text, full_reasoning, error_msg = '', '', '' tool_calls, tool_part = [], None usage, finish_reason = None, None for evt in self._events(resp, dump): kind = evt.get('type') if kind == 'error': err = evt.get('error') error_msg = (err or {}).get('message') if isinstance(err, dict) else None error_msg = error_msg or _dumps(err) elif kind == 'text-delta': full_text += evt.get('text') or '' elif kind == 'reasoning-delta': full_reasoning += evt.get('text') or '' elif kind == 'tool-input-start': tool_part = {'id': evt.get('id'), 'type': 'function', 'function': {'name': evt.get('toolName'), 'arguments': ''}} tool_calls.append(tool_part) elif kind == 'tool-input-delta': if evt.get('delta') and tool_part: tool_part['function']['arguments'] += evt['delta'] elif kind == 'tool-call': # authoritative parsed arguments; the deltas can be fragments for tc in tool_calls: if tc['id'] == evt.get('toolCallId') and isinstance(evt.get('input'), dict): tc['function']['arguments'] = _dumps(evt['input']) tool_part = None elif kind == 'tool-input-end': tool_part = None elif kind == 'finish': usage = usage_from_finish(evt) or usage finish_reason = finish_reason_from(evt, bool(tool_calls)) or finish_reason break # nothing meaningful follows the finish event if error_msg: log_err(f'[error] {model} | {error_msg}') self._send_json(502, openai_error(error_msg, code='context_length_exceeded')) return # Reasoning is NOT content. Substituting one for the other shows the model's private # chain-of-thought as the answer; clients read it from reasoning_content (DeepSeek) # or reasoning (vLLM/OpenRouter), so emit both and leave content to the real reply. msg = {'role': 'assistant', 'content': full_text or None, 'refusal': None, 'annotations': [], 'audio': None, 'function_call': None} if full_reasoning: msg['reasoning_content'] = full_reasoning msg['reasoning'] = full_reasoning if tool_calls: msg['tool_calls'] = tool_calls if finish_reason is None: finish_reason = 'tool_calls' if tool_calls else 'stop' self._send_json(200, { 'id': gen_id, 'object': 'chat.completion', 'created': int(time.time()), 'model': model, 'system_fingerprint': None, 'service_tier': None, 'choices': [{'index': 0, 'message': msg, 'logprobs': None, 'finish_reason': finish_reason}], 'usage': usage or dict(ZERO_USAGE), }) log_done(f'{model} | {len(full_text)} text / {len(full_reasoning)} reasoning / ' f'{len(tool_calls)} tools | {finish_reason} | ' f'{self._usage_note(usage)}{int((time.time() - t0) * 1000)}ms') @staticmethod def _usage_note(usage): if not usage: return '' cached = (usage.get('prompt_tokens_details') or {}).get('cached_tokens') note = f'{usage["prompt_tokens"]}in/{usage["completion_tokens"]}out' if cached: note += f' ({cached} cached)' return note + ' | ' def _stream_response(self, resp, model, gen_id, t0, dump, include_usage=False): tool_calls, tool_idx = [], 0 role_sent = False t_chars, r_chars = 0, 0 error_msg = '' usage, finish_reason = None, None def base(): return {'id': gen_id, 'object': 'chat.completion.chunk', 'created': int(time.time()), 'model': model} def write(chunk): self._begin_stream() self._write_sse(chunk) def delta(d, finish=None): write({**base(), 'choices': [{'index': 0, 'delta': d, 'finish_reason': finish}]}) def ensure_role(): nonlocal role_sent if not role_sent: role_sent = True delta({'role': 'assistant', 'content': ''}) for evt in self._events(resp, dump): kind = evt.get('type') if kind == 'error': err = evt.get('error') error_msg = (err or {}).get('message') if isinstance(err, dict) else None error_msg = error_msg or _dumps(err) self._begin_stream() self._write_sse(openai_error(error_msg)) self._end_stream() log(f'[error] {model} | {error_msg}') break if kind == 'text-start': # Only opens the message. Emitting a second role chunk here reads as a new # assistant turn, and clearing tool_calls would lose calls made earlier in # the same response (finish_reason would come back "stop"). ensure_role() elif kind == 'text-delta': if evt.get('text'): t_chars += len(evt['text']) delta({'content': evt['text']}) elif kind == 'reasoning-delta': if evt.get('text'): r_chars += len(evt['text']) ensure_role() delta({'reasoning_content': evt['text'], 'reasoning': evt['text']}) elif kind == 'tool-input-start': ensure_role() tool_idx = len(tool_calls) tool_calls.append({'id': evt.get('id'), 'name': evt.get('toolName')}) delta({'tool_calls': [{'index': tool_idx, 'id': evt.get('id'), 'type': 'function', 'function': {'name': evt.get('toolName'), 'arguments': ''}}]}) elif kind == 'tool-input-delta': if evt.get('delta') and tool_idx < len(tool_calls): delta({'tool_calls': [{'index': tool_idx, 'function': {'arguments': evt['delta']}}]}) elif kind == 'finish': usage = usage_from_finish(evt) or usage finish_reason = finish_reason_from(evt, bool(tool_calls)) or finish_reason break # nothing meaningful follows the finish event # skipped: start, start-step, text-end, reasoning-start/end, # tool-input-end, tool-call, finish-step, provider-metadata if error_msg: return # already reported on the wire reason = finish_reason or ('tool_calls' if tool_calls else 'stop') final = {**base(), 'choices': [{'index': 0, 'delta': {}, 'logprobs': None, 'finish_reason': reason}]} if include_usage: final['usage'] = None write(final) if include_usage: # OpenAI's include_usage sends one extra chunk with empty choices and the totals. write({**base(), 'choices': [], 'usage': usage or dict(ZERO_USAGE)}) self._write_chunk(b'data: [DONE]\n\n') self._end_stream() log(f'[done] {model} | {t_chars} text / {r_chars} reasoning / {len(tool_calls)} tools | ' f'{reason} | {self._usage_note(usage)}{int((time.time() - t0) * 1000)}ms') # -- start -------------------------------------------------------------------- def _pids_on_port(port): """PIDs listening on `port`, without shelling out to a pipeline.""" pids = [] if os.name == 'nt': out = subprocess.run(['netstat', '-ano'], capture_output=True, timeout=5).stdout for line in out.decode('utf-8', 'ignore').splitlines(): parts = line.split() if len(parts) < 5 or parts[3] != 'LISTENING': continue local = parts[1] if local.rsplit(':', 1)[-1] == str(port): pids.append(parts[4]) else: out = subprocess.run(['lsof', '-ti', f'tcp:{port}'], capture_output=True, timeout=5).stdout pids = [p for p in out.decode('utf-8', 'ignore').split() if p] return pids def kill_port(port): try: pids = _pids_on_port(port) except Exception: return # no netstat/lsof, or nothing listening for pid in pids: if pid == str(os.getpid()): continue log(f'killing existing process on port {port} (PID {pid})') try: if os.name == 'nt': subprocess.run(['taskkill', '/F', '/PID', pid], capture_output=True, timeout=5) else: os.kill(int(pid), signal.SIGKILL) except Exception: pass def main(): if os.name == 'nt': os.system('') # enable VT escape sequences in legacy consoles for stream in (sys.stdout, sys.stderr): try: stream.reconfigure(encoding='utf-8', errors='replace') except (AttributeError, ValueError): pass # already utf-8, or a stream that cannot be reconfigured banner() if CONFIG_NOTE: log(CONFIG_NOTE) log(f'=== proxy started (debug: {"ON" if DEBUG else "OFF"}) ===') kill_port(PORT) try: server = ThreadingHTTPServer(('', PORT), ProxyHandler) except OSError as e: if e.errno in (48, 98, 10048): # EADDRINUSE log_err(f'Port {PORT} still in use after kill attempt') sys.exit(1) raise server.daemon_threads = True server.timeout = TIMEOUT def shutdown(signum, frame): log('shutting down...') threading.Thread(target=server.shutdown, daemon=True).start() for signame in ('SIGINT', 'SIGTERM', 'SIGBREAK'): sig = getattr(signal, signame, None) if sig is None: continue try: signal.signal(sig, shutdown) except (OSError, ValueError): pass log(f'listening on http://localhost:{PORT}') try: server.serve_forever() finally: server.server_close() _log_file.close() if __name__ == '__main__': main()