Files
wehub-resource-sync 6ede33ccdb
Build and Push Docker Images / create_manifest (web, surfsense-web, , cpu) (push) Has been cancelled
Build and Push Docker Images / finalize_release (push) Has been cancelled
Obsidian Plugin Lint / lint (push) Has been cancelled
Build and Push Docker Images / build (./surfsense_web, cpu, ./surfsense_web/Dockerfile, web, surfsense-web, ubuntu-24.04-arm, linux/arm64, arm64, , runner, false, cpu) (push) Has been cancelled
Build and Push Docker Images / build (./surfsense_web, cpu, ./surfsense_web/Dockerfile, web, surfsense-web, ubuntu-latest, linux/amd64, amd64, , runner, false, cpu) (push) Has been cancelled
Build and Push Docker Images / compute_version (push) Has been cancelled
Build and Push Docker Images / build (./surfsense_backend, cpu, ./surfsense_backend/Dockerfile, backend, surfsense-backend, ubuntu-24.04-arm, linux/arm64, arm64, , production, false, cpu) (push) Has been cancelled
Build and Push Docker Images / build (./surfsense_backend, cpu, ./surfsense_backend/Dockerfile, backend, surfsense-backend, ubuntu-latest, linux/amd64, amd64, , production, false, cpu) (push) Has been cancelled
Build and Push Docker Images / build (./surfsense_backend, cu126, ./surfsense_backend/Dockerfile, backend, surfsense-backend, ubuntu-24.04-arm, linux/arm64, arm64, -cuda126, production, true, cuda126) (push) Has been cancelled
Build and Push Docker Images / build (./surfsense_backend, cu126, ./surfsense_backend/Dockerfile, backend, surfsense-backend, ubuntu-latest, linux/amd64, amd64, -cuda126, production, true, cuda126) (push) Has been cancelled
Build and Push Docker Images / build (./surfsense_backend, cu128, ./surfsense_backend/Dockerfile, backend, surfsense-backend, ubuntu-24.04-arm, linux/arm64, arm64, -cuda, production, true, cuda) (push) Has been cancelled
Build and Push Docker Images / build (./surfsense_backend, cu128, ./surfsense_backend/Dockerfile, backend, surfsense-backend, ubuntu-latest, linux/amd64, amd64, -cuda, production, true, cuda) (push) Has been cancelled
Build and Push Docker Images / verify_digests (push) Has been cancelled
Build and Push Docker Images / create_manifest (backend, surfsense-backend, , cpu) (push) Has been cancelled
Build and Push Docker Images / create_manifest (backend, surfsense-backend, -cuda, cuda) (push) Has been cancelled
Build and Push Docker Images / create_manifest (backend, surfsense-backend, -cuda126, cuda126) (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 13:33:44 +08:00

499 lines
18 KiB
Python

"""
Google Gmail Connector Module | Google OAuth Credentials | Gmail API
A module for retrieving emails from Gmail using Google OAuth credentials.
Allows fetching emails from Gmail mailbox using Google OAuth credentials.
"""
import base64
import json
import logging
from typing import Any
from bs4 import BeautifulSoup
from google.auth.transport.requests import Request
from google.oauth2.credentials import Credentials
from googleapiclient.discovery import build
from markdownify import markdownify as md
from sqlalchemy.ext.asyncio import AsyncSession
from sqlalchemy.future import select
from sqlalchemy.orm.attributes import flag_modified
from app.db import (
SearchSourceConnector,
SearchSourceConnectorType,
)
logger = logging.getLogger(__name__)
def fetch_google_user_email(credentials: Credentials) -> str | None:
"""
Fetch user email from Gmail API using Google credentials.
Uses the Gmail users.getProfile endpoint which returns the authenticated
user's email address.
Args:
credentials: Google OAuth Credentials object (not encrypted)
Returns:
User's email address or None if fetch fails
"""
try:
service = build("gmail", "v1", credentials=credentials)
profile = service.users().getProfile(userId="me").execute()
email = profile.get("emailAddress")
if email:
logger.debug(f"Fetched Google user email: {email}")
return email
return None
except Exception as e:
logger.warning(f"Error fetching Google user email: {e!s}")
return None
class GoogleGmailConnector:
"""Class for retrieving emails from Gmail using Google OAuth credentials."""
def __init__(
self,
credentials: Credentials,
session: AsyncSession,
user_id: str,
connector_id: int | None = None,
):
"""
Initialize the GoogleGmailConnector class.
Args:
credentials: Google OAuth Credentials object
session: Database session for updating connector
user_id: User ID (kept for backward compatibility)
connector_id: Optional connector ID for direct updates
"""
self._credentials = credentials
self._session = session
self._user_id = user_id
self._connector_id = connector_id
self.service = None
async def _get_credentials(
self,
) -> Credentials:
"""
Get valid Google OAuth credentials.
Supports both native OAuth (with refresh_token) and Composio-sourced
credentials (with refresh_handler). For Composio credentials, validation
and DB persistence are skipped since Composio manages its own tokens.
"""
has_standard_refresh = bool(self._credentials.refresh_token)
if has_standard_refresh and not all(
[self._credentials.client_id, self._credentials.client_secret]
):
raise ValueError(
"Google OAuth credentials (client_id, client_secret) must be set"
)
if self._credentials and not self._credentials.expired:
return self._credentials
if has_standard_refresh:
self._credentials = Credentials(
token=self._credentials.token,
refresh_token=self._credentials.refresh_token,
token_uri=self._credentials.token_uri,
client_id=self._credentials.client_id,
client_secret=self._credentials.client_secret,
scopes=self._credentials.scopes,
expiry=self._credentials.expiry,
)
if self._credentials.expired or not self._credentials.valid:
try:
self._credentials.refresh(Request())
# Only persist refreshed token for native OAuth (Composio manages its own)
if has_standard_refresh and self._session:
if self._connector_id:
result = await self._session.execute(
select(SearchSourceConnector).filter(
SearchSourceConnector.id == self._connector_id
)
)
else:
result = await self._session.execute(
select(SearchSourceConnector).filter(
SearchSourceConnector.user_id == self._user_id,
SearchSourceConnector.connector_type
== SearchSourceConnectorType.GOOGLE_GMAIL_CONNECTOR,
)
)
connector = result.scalars().first()
if connector is None:
raise RuntimeError(
"GMAIL connector not found; cannot persist refreshed token."
)
from app.config import config
from app.utils.oauth_security import TokenEncryption
creds_dict = json.loads(self._credentials.to_json())
token_encrypted = connector.config.get("_token_encrypted", False)
if token_encrypted and config.SECRET_KEY:
token_encryption = TokenEncryption(config.SECRET_KEY)
if creds_dict.get("token"):
creds_dict["token"] = token_encryption.encrypt_token(
creds_dict["token"]
)
if creds_dict.get("refresh_token"):
creds_dict["refresh_token"] = (
token_encryption.encrypt_token(
creds_dict["refresh_token"]
)
)
if creds_dict.get("client_secret"):
creds_dict["client_secret"] = (
token_encryption.encrypt_token(
creds_dict["client_secret"]
)
)
creds_dict["_token_encrypted"] = True
connector.config = creds_dict
flag_modified(connector, "config")
await self._session.commit()
except Exception as e:
error_str = str(e)
if (
"invalid_grant" in error_str.lower()
or "token has been expired or revoked" in error_str.lower()
):
raise Exception(
"Gmail authentication failed. Please re-authenticate."
) from e
raise Exception(
f"Failed to refresh Google OAuth credentials: {e!s}"
) from e
return self._credentials
async def _get_service(self):
"""
Get the Gmail service instance using Google OAuth credentials.
Returns:
Gmail service instance
Raises:
ValueError: If credentials have not been set
Exception: If service creation fails
"""
if self.service:
return self.service
try:
credentials = await self._get_credentials()
self.service = build("gmail", "v1", credentials=credentials)
return self.service
except Exception as e:
error_str = str(e)
# If the error already contains a user-friendly re-authentication message, preserve it
if (
"re-authenticate" in error_str.lower()
or "expired or been revoked" in error_str.lower()
or "authentication failed" in error_str.lower()
):
raise Exception(error_str) from e
raise Exception(f"Failed to create Gmail service: {e!s}") from e
async def get_user_profile(self) -> tuple[dict[str, Any], str | None]:
"""
Fetch user's Gmail profile information.
Returns:
Tuple containing (profile dict, error message or None)
"""
try:
service = await self._get_service()
profile = service.users().getProfile(userId="me").execute()
return {
"email_address": profile.get("emailAddress"),
"messages_total": profile.get("messagesTotal", 0),
"threads_total": profile.get("threadsTotal", 0),
"history_id": profile.get("historyId"),
}, None
except Exception as e:
return {}, f"Error fetching user profile: {e!s}"
async def get_messages_list(
self,
max_results: int = 100,
query: str = "",
label_ids: list[str] | None = None,
include_spam_trash: bool = False,
) -> tuple[list[dict[str, Any]], str | None]:
"""
Fetch list of messages from Gmail.
Args:
max_results: Maximum number of messages to fetch (default: 100)
query: Gmail search query (e.g., "is:unread", "from:example@gmail.com")
label_ids: List of label IDs to filter by
include_spam_trash: Whether to include spam and trash
Returns:
Tuple containing (messages list, error message or None)
"""
try:
service = await self._get_service()
# Build request parameters
request_params = {
"userId": "me",
"maxResults": max_results,
"includeSpamTrash": include_spam_trash,
}
if query:
request_params["q"] = query
if label_ids:
request_params["labelIds"] = label_ids
# Get messages list
result = service.users().messages().list(**request_params).execute()
messages = result.get("messages", [])
return messages, None
except Exception as e:
error_str = str(e)
# If the error already contains a user-friendly re-authentication message, preserve it
if (
"re-authenticate" in error_str.lower()
or "expired or been revoked" in error_str.lower()
or "authentication failed" in error_str.lower()
):
return [], error_str
return [], f"Error fetching messages list: {e!s}"
async def get_message_details(
self, message_id: str
) -> tuple[dict[str, Any], str | None]:
"""
Fetch detailed information for a specific message.
Args:
message_id: The ID of the message to fetch
Returns:
Tuple containing (message details dict, error message or None)
"""
try:
service = await self._get_service()
# Get full message details
message = (
service.users()
.messages()
.get(userId="me", id=message_id, format="full")
.execute()
)
return message, None
except Exception as e:
return {}, f"Error fetching message details: {e!s}"
async def get_recent_messages(
self,
max_results: int = 50,
start_date: str | None = None,
end_date: str | None = None,
) -> tuple[list[dict[str, Any]], str | None]:
"""
Fetch recent messages from Gmail within specified date range.
Args:
max_results: Maximum number of messages to fetch (default: 50)
start_date: Start date in YYYY-MM-DD format (default: 3 days ago)
end_date: End date in YYYY-MM-DD format (default: today)
Returns:
Tuple containing (messages list with details, error message or None)
"""
try:
from datetime import datetime, timedelta
# Normalize date values - handle "undefined" strings from frontend
# This prevents "time data 'undefined' does not match format" errors
if start_date == "undefined" or start_date == "":
start_date = None
if end_date == "undefined" or end_date == "":
end_date = None
# Build date query
query_parts = []
if start_date:
# Parse start_date from YYYY-MM-DD to Gmail query format YYYY/MM/DD
start_dt = datetime.strptime(start_date, "%Y-%m-%d")
start_query = start_dt.strftime("%Y/%m/%d")
query_parts.append(f"after:{start_query}")
else:
# Default to 3 days ago
cutoff_date = datetime.now() - timedelta(days=3)
date_query = cutoff_date.strftime("%Y/%m/%d")
query_parts.append(f"after:{date_query}")
if end_date:
# Parse end_date from YYYY-MM-DD to Gmail query format YYYY/MM/DD
end_dt = datetime.strptime(end_date, "%Y-%m-%d")
end_query = end_dt.strftime("%Y/%m/%d")
query_parts.append(f"before:{end_query}")
query = " ".join(query_parts)
# Get messages list
messages_list, error = await self.get_messages_list(
max_results=max_results, query=query
)
if error:
return [], error
# Get detailed information for each message
detailed_messages = []
for msg in messages_list:
message_details, detail_error = await self.get_message_details(
msg["id"]
)
if detail_error:
continue # Skip messages that can't be fetched
detailed_messages.append(message_details)
return detailed_messages, None
except Exception as e:
return [], f"Error fetching recent messages: {e!s}"
@staticmethod
def _html_to_markdown(html: str) -> str:
"""Convert HTML (especially email layouts with nested tables) to clean markdown."""
soup = BeautifulSoup(html, "html.parser")
for tag in soup.find_all(["style", "script", "img"]):
tag.decompose()
for tag in soup.find_all(
["table", "thead", "tbody", "tfoot", "tr", "td", "th"]
):
tag.unwrap()
return md(str(soup)).strip()
def extract_message_text(self, message: dict[str, Any]) -> str:
"""
Extract text content from a Gmail message.
Args:
message: Gmail message object
Returns:
Extracted text content
"""
def get_message_parts(payload):
"""Recursively extract message parts."""
parts = []
if "parts" in payload:
for part in payload["parts"]:
parts.extend(get_message_parts(part))
else:
parts.append(payload)
return parts
try:
payload = message.get("payload", {})
parts = get_message_parts(payload)
text_content = ""
for part in parts:
mime_type = part.get("mimeType", "")
body = part.get("body", {})
data = body.get("data", "")
if mime_type == "text/plain" and data:
# Decode base64 content
decoded_data = base64.urlsafe_b64decode(data + "===").decode(
"utf-8", errors="ignore"
)
text_content += decoded_data + "\n"
elif mime_type == "text/html" and data and not text_content:
decoded_data = base64.urlsafe_b64decode(data + "===").decode(
"utf-8", errors="ignore"
)
text_content = self._html_to_markdown(decoded_data)
return text_content.strip()
except Exception as e:
return f"Error extracting message text: {e!s}"
def format_message_to_markdown(self, message: dict[str, Any]) -> str:
"""
Format a Gmail message to markdown.
Args:
message: Message object from Gmail API
Returns:
Formatted markdown string
"""
try:
# Extract basic message information
message_id = message.get("id", "")
thread_id = message.get("threadId", "")
label_ids = message.get("labelIds", [])
# Extract headers
payload = message.get("payload", {})
headers = payload.get("headers", [])
# Parse headers into a dict
header_dict = {}
for header in headers:
name = header.get("name", "").lower()
value = header.get("value", "")
header_dict[name] = value
# Extract key information
subject = header_dict.get("subject", "No Subject")
from_email = header_dict.get("from", "Unknown Sender")
to_email = header_dict.get("to", "Unknown Recipient")
date_str = header_dict.get("date", "Unknown Date")
# Extract message content
message_text = self.extract_message_text(message)
# Build markdown content
markdown_content = f"# {subject}\n\n"
# Add message details
markdown_content += f"**From:** {from_email}\n"
markdown_content += f"**To:** {to_email}\n"
markdown_content += f"**Date:** {date_str}\n"
if label_ids:
markdown_content += f"**Labels:** {', '.join(label_ids)}\n"
markdown_content += "\n"
# Add message content
if message_text:
markdown_content += f"## Message Content\n\n{message_text}\n\n"
# Add message metadata
markdown_content += "## Message Details\n\n"
markdown_content += f"- **Message ID:** {message_id}\n"
markdown_content += f"- **Thread ID:** {thread_id}\n"
# Add snippet if available
snippet = message.get("snippet", "")
if snippet:
markdown_content += f"- **Snippet:** {snippet}\n"
return markdown_content
except Exception as e:
return f"Error formatting message to markdown: {e!s}"