Clove Deploy
v2: session-based multi-turn with backend_uuid follow-up, seamless fallback on session expiry
8445e7c
Raw
History Blame Contribute Delete
7.62 kB
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 {}