Merge branch 'main' into copilot/fix-mailbox-parameter-editing
This commit is contained in:
@@ -11,7 +11,7 @@ import logging
|
||||
|
||||
from app.core.database import get_db
|
||||
from app.core.deps import get_current_active_user
|
||||
from app.core.security import encrypt_credential
|
||||
from app.core.security import encrypt_credential, decrypt_credential
|
||||
from app.core.config import settings
|
||||
from app.models.database_models import User, GmailCredential
|
||||
from app.models.schemas import (
|
||||
@@ -22,7 +22,7 @@ from app.models.schemas import (
|
||||
GmailAuthorizeResponse,
|
||||
GmailCallbackRequest,
|
||||
)
|
||||
from app.services.gmail_service import GmailService, GMAIL_SCOPES
|
||||
from app.services.gmail_service import GmailService, GmailInjectionError, GMAIL_SCOPES
|
||||
|
||||
router = APIRouter()
|
||||
logger = logging.getLogger(__name__)
|
||||
@@ -319,6 +319,77 @@ async def get_gmail_authorize_url(
|
||||
return GmailAuthorizeResponse(authorization_url=url)
|
||||
|
||||
|
||||
@router.post("/gmail/debug-email", status_code=status.HTTP_200_OK)
|
||||
async def send_gmail_debug_email(
|
||||
current_user: User = Depends(get_current_active_user),
|
||||
db: AsyncSession = Depends(get_db),
|
||||
):
|
||||
"""
|
||||
Inject a debug/test email into the current user's Gmail inbox.
|
||||
|
||||
The message appears to have been sent by christian@docuelevate.org,
|
||||
carries today's date in the subject, and is tagged with the custom
|
||||
labels "test" and "imported" as well as placed in the inbox.
|
||||
|
||||
Useful for verifying that Gmail API delivery is working end-to-end
|
||||
without requiring an active mail-account polling cycle.
|
||||
"""
|
||||
result = await db.execute(
|
||||
select(GmailCredential).where(
|
||||
GmailCredential.user_id == current_user.id,
|
||||
GmailCredential.is_valid == True, # noqa: E712
|
||||
)
|
||||
)
|
||||
credential = result.scalar_one_or_none()
|
||||
|
||||
if not credential:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_400_BAD_REQUEST,
|
||||
detail="No valid Gmail credentials found. Connect Gmail first.",
|
||||
)
|
||||
|
||||
access_token = decrypt_credential(credential.encrypted_access_token) # type: ignore[arg-type]
|
||||
refresh_token = (
|
||||
decrypt_credential(credential.encrypted_refresh_token) # type: ignore[arg-type]
|
||||
if credential.encrypted_refresh_token
|
||||
else None
|
||||
)
|
||||
|
||||
gmail_service = GmailService(
|
||||
access_token=access_token,
|
||||
refresh_token=refresh_token,
|
||||
client_id=settings.GOOGLE_CLIENT_ID,
|
||||
client_secret=settings.GOOGLE_CLIENT_SECRET,
|
||||
)
|
||||
|
||||
try:
|
||||
inject_result = await gmail_service.inject_debug_email(
|
||||
recipient_email=credential.gmail_email, # type: ignore[arg-type]
|
||||
)
|
||||
except GmailInjectionError as exc:
|
||||
raise HTTPException(
|
||||
status_code=status.HTTP_502_BAD_GATEWAY,
|
||||
detail=f"Gmail injection failed: {exc}",
|
||||
)
|
||||
|
||||
# Persist refreshed token if the google-auth library renewed it
|
||||
refreshed = gmail_service.get_refreshed_token()
|
||||
if refreshed:
|
||||
credential.encrypted_access_token = encrypt_credential( # type: ignore[assignment]
|
||||
refreshed["access_token"]
|
||||
)
|
||||
if refreshed.get("expiry"):
|
||||
credential.token_expiry = refreshed["expiry"] # type: ignore[assignment]
|
||||
await db.commit()
|
||||
|
||||
return {
|
||||
"message": "Debug email injected successfully",
|
||||
"message_id": inject_result.get("message_id"),
|
||||
"thread_id": inject_result.get("thread_id"),
|
||||
"label_ids": inject_result.get("label_ids", []),
|
||||
}
|
||||
|
||||
|
||||
@router.post(
|
||||
"/gmail/callback",
|
||||
response_model=GmailCredentialResponse,
|
||||
|
||||
@@ -9,6 +9,10 @@ This is preferred over SMTP forwarding as it doesn't modify the email.
|
||||
import asyncio
|
||||
import base64
|
||||
import logging
|
||||
import textwrap
|
||||
from datetime import datetime, timezone
|
||||
from email.mime.text import MIMEText
|
||||
from email.utils import format_datetime
|
||||
from typing import Optional, Dict, Any
|
||||
|
||||
from google.oauth2.credentials import Credentials
|
||||
@@ -184,6 +188,119 @@ class GmailService:
|
||||
logger.error(f"Failed to get Gmail email address: {e}")
|
||||
return None
|
||||
|
||||
async def get_or_create_label(self, name: str) -> str:
|
||||
"""
|
||||
Return the Gmail label ID for a label with the given name.
|
||||
|
||||
Lists the user's existing labels and returns the ID of the first
|
||||
match (case-insensitive). If no matching label is found, a new
|
||||
label is created and its ID is returned.
|
||||
|
||||
Args:
|
||||
name: Human-readable label name (e.g. "test", "imported").
|
||||
|
||||
Returns:
|
||||
Gmail label ID string (e.g. "Label_1234567890").
|
||||
|
||||
Raises:
|
||||
GmailInjectionError: If the Gmail API call fails.
|
||||
"""
|
||||
loop = asyncio.get_event_loop()
|
||||
|
||||
try:
|
||||
labels_resp = await loop.run_in_executor(
|
||||
None,
|
||||
lambda: self.service.users().labels().list(userId="me").execute(),
|
||||
)
|
||||
for label in labels_resp.get("labels", []):
|
||||
if label.get("name", "").lower() == name.lower():
|
||||
return label["id"]
|
||||
|
||||
# Label not found – create it
|
||||
created = await loop.run_in_executor(
|
||||
None,
|
||||
lambda: self.service.users()
|
||||
.labels()
|
||||
.create(userId="me", body={"name": name})
|
||||
.execute(),
|
||||
)
|
||||
logger.info(f"Created Gmail label '{name}' with id={created['id']}")
|
||||
return created["id"]
|
||||
|
||||
except HttpError as e:
|
||||
error_msg = f"Gmail API error while managing label '{name}': {e.reason if hasattr(e, 'reason') else str(e)}"
|
||||
logger.error(error_msg)
|
||||
raise GmailInjectionError(error_msg)
|
||||
except Exception as e:
|
||||
error_msg = f"Failed to get/create Gmail label '{name}': {str(e)}"
|
||||
logger.error(error_msg)
|
||||
raise GmailInjectionError(error_msg)
|
||||
|
||||
async def inject_debug_email(
|
||||
self,
|
||||
recipient_email: str,
|
||||
) -> Dict[str, Any]:
|
||||
"""
|
||||
Inject a debug/test email into the user's Gmail inbox.
|
||||
|
||||
The message is made to appear as if it was sent by
|
||||
christian@docuelevate.org on the current date. It is placed in
|
||||
the inbox and tagged with the custom labels "test" and "imported"
|
||||
so it is easy to identify and clean up.
|
||||
|
||||
Args:
|
||||
recipient_email: The Gmail address to deliver the message to
|
||||
(the authenticated user's address).
|
||||
|
||||
Returns:
|
||||
Dict with message_id, thread_id, and label_ids.
|
||||
|
||||
Raises:
|
||||
GmailInjectionError: If injection or label management fails.
|
||||
"""
|
||||
now = datetime.now(timezone.utc)
|
||||
date_str = now.strftime("%d %B %Y") # e.g. "25 March 2026"
|
||||
|
||||
subject = f"Test Import – {date_str}"
|
||||
|
||||
body = textwrap.dedent(f"""\
|
||||
Hi there,
|
||||
|
||||
This is an automated test message injected via the Gmail API to
|
||||
confirm that the import pipeline is working correctly.
|
||||
|
||||
Date: {date_str}
|
||||
Source: DocuElevate Integration Test
|
||||
|
||||
If you can see this message in your inbox it means that Gmail API
|
||||
delivery is functioning as expected. Feel free to delete it.
|
||||
|
||||
Best regards,
|
||||
Christian Loris
|
||||
DocuElevate
|
||||
""")
|
||||
|
||||
msg = MIMEText(body, "plain", "utf-8")
|
||||
msg["From"] = "Christian Loris <christian@docuelevate.org>"
|
||||
msg["To"] = recipient_email
|
||||
msg["Subject"] = subject
|
||||
msg["Date"] = format_datetime(now)
|
||||
msg["Message-ID"] = f"<debug-{now.strftime('%Y%m%d%H%M%S')}@docuelevate.org>"
|
||||
|
||||
raw_bytes = msg.as_bytes()
|
||||
|
||||
# Resolve label IDs (create labels if they don't exist yet)
|
||||
test_label_id = await self.get_or_create_label("test")
|
||||
imported_label_id = await self.get_or_create_label("imported")
|
||||
|
||||
label_ids = ["INBOX", test_label_id, imported_label_id]
|
||||
|
||||
return await self.inject_email(
|
||||
raw_email=raw_bytes,
|
||||
label_ids=label_ids,
|
||||
source_account_name="debug",
|
||||
)
|
||||
|
||||
def get_refreshed_token(self) -> Optional[Dict[str, Any]]:
|
||||
"""
|
||||
Return the current access token and expiry if the token was refreshed
|
||||
|
||||
@@ -8,7 +8,7 @@ from celery import Task
|
||||
import logging
|
||||
|
||||
from app.workers.celery_app import celery_app
|
||||
from app.core.database import async_session_maker
|
||||
from app.core.database import async_session_maker, engine
|
||||
from app.core.security import decrypt_credential, encrypt_credential
|
||||
from app.models.database_models import (
|
||||
MailAccount,
|
||||
@@ -34,8 +34,20 @@ class AsyncTask(Task):
|
||||
|
||||
def __call__(self, *args, **kwargs):
|
||||
"""Run async task in event loop"""
|
||||
# Use asyncio.run() for better event loop management
|
||||
return asyncio.run(self.run(*args, **kwargs))
|
||||
|
||||
async def _run():
|
||||
try:
|
||||
return await self.run(*args, **kwargs)
|
||||
finally:
|
||||
# Dispose the connection pool before the event loop closes.
|
||||
# Each asyncio.run() creates a fresh event loop; if pooled
|
||||
# asyncpg connections are still open when the loop is torn
|
||||
# down, asyncpg raises "Exception terminating connection".
|
||||
# Disposing the engine here closes those connections cleanly
|
||||
# inside the same loop, before asyncio.run() shuts it down.
|
||||
await engine.dispose()
|
||||
|
||||
return asyncio.run(_run())
|
||||
|
||||
|
||||
@celery_app.task(base=AsyncTask, name="app.workers.tasks.process_mail_account")
|
||||
|
||||
@@ -11,6 +11,7 @@ psycopg2-binary==2.9.11
|
||||
asyncpg==0.31.0
|
||||
|
||||
# Authentication
|
||||
|
||||
python-jose[cryptography]==3.5.0 # Updated: Fixed algorithm confusion with OpenSSH ECDSA keys (was 3.3.0)
|
||||
bcrypt==4.3.0
|
||||
python-multipart==0.0.22 # Updated: Fixed multiple vulnerabilities (was 0.0.6)
|
||||
|
||||
Reference in New Issue
Block a user