import json import mimetypes import random import re import sys from uuid import uuid4 from curl_cffi import CurlMime, requests from perplexity.config import ( DEFAULT_HEADERS, ENDPOINT_AUTH_SESSION, ENDPOINT_AUTH_SIGNIN, ENDPOINT_SSE_ASK, ENDPOINT_UPLOAD_URL, MODEL_MAPPINGS, ) from perplexity.exceptions import FileUploadError from perplexity.utils import parse_nested_json_response, validate_search_params, validate_query_limits from perplexity.logger import get_logger from .emailnator import Emailnator logger = get_logger("async_client") class AsyncMixin: def __init__(self, *args, **kwargs): self.__storedargs = args, kwargs self.async_initialized = False async def __ainit__(self, *args, **kwargs): pass async def __initobj(self): assert not self.async_initialized self.async_initialized = True await self.__ainit__(*self.__storedargs[0], **self.__storedargs[1]) return self def __await__(self): return self.__initobj().__await__() class Client(AsyncMixin): async def __ainit__(self, cookies={}): self.session = requests.AsyncSession( headers=DEFAULT_HEADERS.copy(), cookies=cookies, impersonate="chrome", ) self.own = bool(cookies) self.copilot = 0 if not cookies else float("inf") self.file_upload = 0 if not cookies else float("inf") self.signin_regex = re.compile( r'"(https://www\.perplexity\.ai/api/auth/callback/email\?callbackUrl=.*?)"' ) self.timestamp = format(random.getrandbits(32), "08x") await self.session.get(ENDPOINT_AUTH_SESSION) async def create_account(self, cookies): while True: try: emailnator_cli = await Emailnator(cookies) resp = await self.session.post( ENDPOINT_AUTH_SIGNIN, data={ "email": emailnator_cli.email, "csrfToken": self.session.cookies.get_dict()["next-auth.csrf-token"].split( "%" )[0], "callbackUrl": "https://www.perplexity.ai/", "json": "true", }, ) if resp.ok: new_msgs = await emailnator_cli.reload( wait_for=lambda x: x["subject"] == "Sign in to Perplexity", timeout=20, ) if new_msgs: break else: logger.error(f"Account creation failed: {resp.status_code}") except Exception as e: logger.warning(f"Account creation attempt failed: {e}") msg = emailnator_cli.get(func=lambda x: x["subject"] == "Sign in to Perplexity") new_account_link = self.signin_regex.search( await emailnator_cli.open(msg["messageID"]) ).group(1) await self.session.get(new_account_link) self.copilot = 5 self.file_upload = 10 return True async def search( self, query, mode="auto", model=None, sources=["web"], files={}, stream=False, language="en-US", follow_up=None, incognito=False, ): validate_search_params(mode, model, sources, self.own) validate_query_limits(self.copilot, self.file_upload, mode, len(files)) self.copilot = ( self.copilot - 1 if mode in ["pro", "reasoning", "deep research"] else self.copilot ) self.file_upload = self.file_upload - len(files) if files else self.file_upload uploaded_files = [] for filename, file in files.items(): file_type = mimetypes.guess_type(filename)[0] file_upload_info = ( await self.session.post( ENDPOINT_UPLOAD_URL, params={"version": "2.18", "source": "default"}, json={ "content_type": file_type, "file_size": sys.getsizeof(file), "filename": filename, "force_image": False, "source": "default", }, ) ).json() mp = CurlMime() for key, value in file_upload_info["fields"].items(): mp.addpart(name=key, data=value) mp.addpart( name="file", content_type=file_type, filename=filename, data=file, ) upload_resp = await self.session.post(file_upload_info["s3_bucket_url"], multipart=mp) if not upload_resp.ok: raise FileUploadError(f"File upload failed: {upload_resp.status_code}") if "image/upload" in file_upload_info["s3_object_url"]: uploaded_url = re.sub( r"/private/s--.*?--/v\d+/user_uploads/", "/private/user_uploads/", upload_resp.json()["secure_url"], ) else: uploaded_url = file_upload_info["s3_object_url"] uploaded_files.append(uploaded_url) json_data = { "query_str": query, "params": { "attachments": ( uploaded_files + follow_up.get("attachments", []) if follow_up else uploaded_files ), "frontend_context_uuid": str(uuid4()), "frontend_uuid": str(uuid4()), "is_incognito": incognito, "language": language, "last_backend_uuid": (follow_up["backend_uuid"] if follow_up else None), "mode": "concise" if mode == "auto" else "copilot", "model_preference": MODEL_MAPPINGS[mode][model], "source": "default", "sources": sources, "version": "2.18", }, } resp = await self.session.post(ENDPOINT_SSE_ASK, json=json_data, stream=True) chunks = [] async def stream_response(resp): async for chunk in resp.aiter_lines(delimiter=b"\r\n\r\n"): content = chunk.decode("utf-8") if content.startswith("event: message\r\n"): try: content_json = json.loads(content[len("event: message\r\ndata: "):]) content_json = parse_nested_json_response(content_json) chunks.append(content_json) yield chunks[-1] except (json.JSONDecodeError, KeyError): continue elif content.startswith("event: end_of_stream\r\n"): return if stream: return stream_response(resp) async for chunk in resp.aiter_lines(delimiter=b"\r\n\r\n"): content = chunk.decode("utf-8") if content.startswith("event: message\r\n"): try: content_json = json.loads(content[len("event: message\r\ndata: "):]) content_json = parse_nested_json_response(content_json) chunks.append(content_json) except (json.JSONDecodeError, KeyError): continue elif content.startswith("event: end_of_stream\r\n"): return chunks[-1] if chunks else {}