Merge pull request #138 from christianlouis/copilot/fix-quality-gate-issues
Fix SonarQube Quality Gate: eliminate OAuth duplication, secure logging, add path validation
This commit is contained in:
+46
-9
@@ -1,10 +1,11 @@
|
|||||||
"""
|
"""
|
||||||
Common utilities for API routes
|
Common utilities for API routes
|
||||||
"""
|
"""
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
import os
|
import os
|
||||||
from sqlalchemy.orm import Session
|
from pathlib import Path
|
||||||
from fastapi import Depends
|
from fastapi import HTTPException, status
|
||||||
|
|
||||||
from app.database import SessionLocal
|
from app.database import SessionLocal
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
@@ -12,6 +13,7 @@ from app.config import settings
|
|||||||
# Set up logging
|
# Set up logging
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
def get_db():
|
def get_db():
|
||||||
"""Database dependency injection for routes"""
|
"""Database dependency injection for routes"""
|
||||||
db = SessionLocal()
|
db = SessionLocal()
|
||||||
@@ -20,17 +22,52 @@ def get_db():
|
|||||||
finally:
|
finally:
|
||||||
db.close()
|
db.close()
|
||||||
|
|
||||||
|
|
||||||
def resolve_file_path(file_path: str, subfolder: str = None) -> str:
|
def resolve_file_path(file_path: str, subfolder: str = None) -> str:
|
||||||
"""
|
"""
|
||||||
Resolves a file path to an absolute path.
|
Resolves a file path to an absolute path with path traversal protection.
|
||||||
If the path is not absolute, it will be joined with the workdir path.
|
If the path is not absolute, it will be joined with the workdir path.
|
||||||
Optionally, can include a subfolder like 'processed'.
|
Optionally, can include a subfolder like 'processed'.
|
||||||
|
|
||||||
Returns the absolute file path.
|
Security: Validates that the resolved path stays within the workdir
|
||||||
|
to prevent path traversal attacks (e.g., ../../etc/passwd).
|
||||||
|
|
||||||
|
Args:
|
||||||
|
file_path: The file path to resolve
|
||||||
|
subfolder: Optional subfolder within workdir
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
The validated absolute file path
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
HTTPException: If the path attempts to escape the workdir
|
||||||
"""
|
"""
|
||||||
|
# Get the workdir as the security boundary
|
||||||
|
workdir = Path(settings.workdir).resolve()
|
||||||
|
|
||||||
|
# Build the base directory (workdir or workdir/subfolder)
|
||||||
|
if subfolder:
|
||||||
|
base_dir = workdir / subfolder
|
||||||
|
else:
|
||||||
|
base_dir = workdir
|
||||||
|
|
||||||
|
# Resolve the file path
|
||||||
if not os.path.isabs(file_path):
|
if not os.path.isabs(file_path):
|
||||||
if subfolder:
|
# Relative path: join with base_dir
|
||||||
file_path = os.path.join(settings.workdir, subfolder, file_path)
|
resolved_path = (base_dir / file_path).resolve()
|
||||||
else:
|
else:
|
||||||
file_path = os.path.join(settings.workdir, file_path)
|
# Absolute path: resolve as-is
|
||||||
return file_path
|
resolved_path = Path(file_path).resolve()
|
||||||
|
|
||||||
|
# Ensure the resolved path is within workdir (path traversal protection)
|
||||||
|
# This checks both relative and absolute paths against workdir
|
||||||
|
try:
|
||||||
|
resolved_path.relative_to(workdir)
|
||||||
|
except ValueError:
|
||||||
|
# Path is outside the workdir - potential path traversal attack
|
||||||
|
logger.warning(f"Path traversal attempt detected: {file_path} -> {resolved_path}")
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=status.HTTP_400_BAD_REQUEST, detail="Invalid file path: path traversal not allowed"
|
||||||
|
)
|
||||||
|
|
||||||
|
return str(resolved_path)
|
||||||
|
|||||||
+40
-117
@@ -1,22 +1,22 @@
|
|||||||
"""
|
"""
|
||||||
Dropbox API endpoints
|
Dropbox API endpoints
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from fastapi import APIRouter, Request, HTTPException, status, Form
|
from fastapi import APIRouter, Request, HTTPException, status, Form
|
||||||
import logging
|
import logging
|
||||||
import os
|
import os
|
||||||
import requests
|
import requests
|
||||||
import json
|
|
||||||
from datetime import datetime, timedelta
|
|
||||||
from typing import Optional
|
|
||||||
|
|
||||||
from app.auth import require_login
|
from app.auth import require_login
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
|
from app.utils.oauth_helper import exchange_oauth_token
|
||||||
|
|
||||||
# Set up logging
|
# Set up logging
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
router = APIRouter()
|
router = APIRouter()
|
||||||
|
|
||||||
|
|
||||||
@router.post("/dropbox/exchange-token")
|
@router.post("/dropbox/exchange-token")
|
||||||
@require_login
|
@require_login
|
||||||
async def exchange_dropbox_token(
|
async def exchange_dropbox_token(
|
||||||
@@ -25,90 +25,33 @@ async def exchange_dropbox_token(
|
|||||||
client_secret: str = Form(...),
|
client_secret: str = Form(...),
|
||||||
redirect_uri: str = Form(...),
|
redirect_uri: str = Form(...),
|
||||||
code: str = Form(...),
|
code: str = Form(...),
|
||||||
folder_path: str = Form(None)
|
folder_path: str = Form(None),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Exchange an authorization code for a refresh token from Dropbox.
|
Exchange an authorization code for a refresh token from Dropbox.
|
||||||
This is done on the server to avoid exposing client secret in the browser.
|
This is done on the server to avoid exposing client secret in the browser.
|
||||||
"""
|
"""
|
||||||
try:
|
# Prepare the token request
|
||||||
logger.info("Starting Dropbox token exchange process")
|
token_url = "https://api.dropboxapi.com/oauth2/token"
|
||||||
|
|
||||||
# Prepare the token request
|
payload = {
|
||||||
token_url = "https://api.dropboxapi.com/oauth2/token"
|
"client_id": client_id,
|
||||||
|
"client_secret": client_secret,
|
||||||
|
"code": code,
|
||||||
|
"redirect_uri": redirect_uri,
|
||||||
|
"grant_type": "authorization_code",
|
||||||
|
}
|
||||||
|
|
||||||
payload = {
|
# Use shared OAuth helper (handles secure logging and error handling)
|
||||||
'client_id': client_id,
|
token_data = exchange_oauth_token(provider_name="Dropbox", token_url=token_url, payload=payload)
|
||||||
'client_secret': client_secret,
|
|
||||||
'code': code,
|
|
||||||
'redirect_uri': redirect_uri,
|
|
||||||
'grant_type': 'authorization_code'
|
|
||||||
}
|
|
||||||
|
|
||||||
# Log request details (excluding secret)
|
# Return just what's needed by the frontend
|
||||||
safe_payload = payload.copy()
|
return {
|
||||||
safe_payload['client_secret'] = '[REDACTED]'
|
"refresh_token": token_data["refresh_token"],
|
||||||
safe_payload['code'] = f"{code[:5]}...{code[-5:]}" if len(code) > 10 else '[REDACTED]'
|
"access_token": token_data["access_token"],
|
||||||
logger.info(f"Token exchange request payload: {safe_payload}")
|
"expires_in": token_data.get("expires_in", 14400),
|
||||||
|
}
|
||||||
|
|
||||||
# Make the token request
|
|
||||||
logger.info("Sending POST request to Dropbox for token exchange")
|
|
||||||
response = requests.post(token_url, data=payload, timeout=settings.http_request_timeout)
|
|
||||||
|
|
||||||
# Check if the request was successful
|
|
||||||
logger.info(f"Token exchange response status: {response.status_code}")
|
|
||||||
|
|
||||||
if response.status_code != 200:
|
|
||||||
# Log the error response for debugging
|
|
||||||
try:
|
|
||||||
error_json = response.json()
|
|
||||||
logger.error(f"Token exchange failed with status {response.status_code}: {error_json}")
|
|
||||||
error_detail = error_json
|
|
||||||
except Exception as json_err:
|
|
||||||
logger.error(f"Failed to parse error response as JSON: {str(json_err)}")
|
|
||||||
logger.error(f"Raw response content: {response.content[:500]}") # Limit log size
|
|
||||||
error_detail = {"error": "Unknown error", "raw_content_snippet": str(response.content[:100])}
|
|
||||||
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_400_BAD_REQUEST,
|
|
||||||
detail=f"Token exchange failed: {error_detail}"
|
|
||||||
)
|
|
||||||
|
|
||||||
# Return the token response
|
|
||||||
token_data = response.json()
|
|
||||||
|
|
||||||
# Validate the token response
|
|
||||||
if "refresh_token" not in token_data:
|
|
||||||
logger.error(f"Dropbox returned success but no refresh token found in response: {token_data.keys()}")
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_502_BAD_GATEWAY,
|
|
||||||
detail="Dropbox OAuth server returned success but no refresh token was included"
|
|
||||||
)
|
|
||||||
|
|
||||||
# Calculate token length for logging
|
|
||||||
refresh_token_length = len(token_data.get("refresh_token", ""))
|
|
||||||
access_token_length = len(token_data.get("access_token", ""))
|
|
||||||
|
|
||||||
logger.info(f"Successfully exchanged authorization code for Dropbox tokens. "
|
|
||||||
f"Refresh token length: {refresh_token_length}, "
|
|
||||||
f"Access token length: {access_token_length}")
|
|
||||||
|
|
||||||
# Return just what's needed by the frontend
|
|
||||||
return {
|
|
||||||
"refresh_token": token_data["refresh_token"],
|
|
||||||
"access_token": token_data["access_token"],
|
|
||||||
"expires_in": token_data.get("expires_in", 14400)
|
|
||||||
}
|
|
||||||
|
|
||||||
except HTTPException:
|
|
||||||
# Re-raise HTTP exceptions as they already have appropriate status codes
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
logger.exception(f"Unexpected error during Dropbox token exchange: {str(e)}")
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
|
||||||
detail=f"Failed to exchange token: {str(e)}"
|
|
||||||
)
|
|
||||||
|
|
||||||
@router.post("/dropbox/update-settings")
|
@router.post("/dropbox/update-settings")
|
||||||
@require_login
|
@require_login
|
||||||
@@ -117,7 +60,7 @@ async def update_dropbox_settings(
|
|||||||
app_key: str = Form(None),
|
app_key: str = Form(None),
|
||||||
app_secret: str = Form(None),
|
app_secret: str = Form(None),
|
||||||
refresh_token: str = Form(...),
|
refresh_token: str = Form(...),
|
||||||
folder_path: str = Form(None)
|
folder_path: str = Form(None),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Update Dropbox settings in memory
|
Update Dropbox settings in memory
|
||||||
@@ -144,18 +87,15 @@ async def update_dropbox_settings(
|
|||||||
|
|
||||||
# Test token validity would be here, but we'll skip it for now
|
# Test token validity would be here, but we'll skip it for now
|
||||||
|
|
||||||
return {
|
return {"status": "success", "message": "Dropbox settings have been updated in memory"}
|
||||||
"status": "success",
|
|
||||||
"message": "Dropbox settings have been updated in memory"
|
|
||||||
}
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Unexpected error updating Dropbox settings: {str(e)}")
|
logger.exception(f"Unexpected error updating Dropbox settings: {str(e)}")
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=f"Failed to update Dropbox settings: {str(e)}"
|
||||||
detail=f"Failed to update Dropbox settings: {str(e)}"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@router.get("/dropbox/test-token")
|
@router.get("/dropbox/test-token")
|
||||||
@require_login
|
@require_login
|
||||||
async def test_dropbox_token(request: Request):
|
async def test_dropbox_token(request: Request):
|
||||||
@@ -167,17 +107,14 @@ async def test_dropbox_token(request: Request):
|
|||||||
|
|
||||||
if not settings.dropbox_refresh_token or not settings.dropbox_app_key or not settings.dropbox_app_secret:
|
if not settings.dropbox_refresh_token or not settings.dropbox_app_key or not settings.dropbox_app_secret:
|
||||||
logger.warning("Dropbox credentials not fully configured")
|
logger.warning("Dropbox credentials not fully configured")
|
||||||
return {
|
return {"status": "error", "message": "Dropbox credentials are not fully configured"}
|
||||||
"status": "error",
|
|
||||||
"message": "Dropbox credentials are not fully configured"
|
|
||||||
}
|
|
||||||
|
|
||||||
# Check token validity by getting current account info
|
# Check token validity by getting current account info
|
||||||
headers = {"Authorization": f"Bearer {settings.dropbox_refresh_token}"}
|
headers = {"Authorization": f"Bearer {settings.dropbox_refresh_token}"}
|
||||||
response = requests.post(
|
response = requests.post(
|
||||||
"https://api.dropboxapi.com/2/users/get_current_account",
|
"https://api.dropboxapi.com/2/users/get_current_account",
|
||||||
headers=headers,
|
headers=headers,
|
||||||
timeout=settings.http_request_timeout
|
timeout=settings.http_request_timeout,
|
||||||
)
|
)
|
||||||
|
|
||||||
# If token is invalid, try refreshing it
|
# If token is invalid, try refreshing it
|
||||||
@@ -190,18 +127,14 @@ async def test_dropbox_token(request: Request):
|
|||||||
"grant_type": "refresh_token",
|
"grant_type": "refresh_token",
|
||||||
"refresh_token": settings.dropbox_refresh_token,
|
"refresh_token": settings.dropbox_refresh_token,
|
||||||
"client_id": settings.dropbox_app_key,
|
"client_id": settings.dropbox_app_key,
|
||||||
"client_secret": settings.dropbox_app_secret
|
"client_secret": settings.dropbox_app_secret,
|
||||||
}
|
}
|
||||||
|
|
||||||
refresh_response = requests.post(refresh_url, data=refresh_data, timeout=settings.http_request_timeout)
|
refresh_response = requests.post(refresh_url, data=refresh_data, timeout=settings.http_request_timeout)
|
||||||
|
|
||||||
if refresh_response.status_code != 200:
|
if refresh_response.status_code != 200:
|
||||||
logger.error(f"Failed to refresh Dropbox token: {refresh_response.text}")
|
logger.error(f"Failed to refresh Dropbox token: {refresh_response.text}")
|
||||||
return {
|
return {"status": "error", "message": "Refresh token has expired or is invalid", "needs_reauth": True}
|
||||||
"status": "error",
|
|
||||||
"message": "Refresh token has expired or is invalid",
|
|
||||||
"needs_reauth": True
|
|
||||||
}
|
|
||||||
|
|
||||||
token_info = refresh_response.json()
|
token_info = refresh_response.json()
|
||||||
access_token = token_info.get("access_token")
|
access_token = token_info.get("access_token")
|
||||||
@@ -211,14 +144,14 @@ async def test_dropbox_token(request: Request):
|
|||||||
response = requests.post(
|
response = requests.post(
|
||||||
"https://api.dropboxapi.com/2/users/get_current_account",
|
"https://api.dropboxapi.com/2/users/get_current_account",
|
||||||
headers=headers,
|
headers=headers,
|
||||||
timeout=settings.http_request_timeout
|
timeout=settings.http_request_timeout,
|
||||||
)
|
)
|
||||||
|
|
||||||
if response.status_code != 200:
|
if response.status_code != 200:
|
||||||
logger.error(f"Dropbox token test failed: {response.status_code} {response.text}")
|
logger.error(f"Dropbox token test failed: {response.status_code} {response.text}")
|
||||||
return {
|
return {
|
||||||
"status": "error",
|
"status": "error",
|
||||||
"message": f"Token validation failed with status {response.status_code}: {response.text}"
|
"message": f"Token validation failed with status {response.status_code}: {response.text}",
|
||||||
}
|
}
|
||||||
|
|
||||||
# Get account info
|
# Get account info
|
||||||
@@ -227,27 +160,22 @@ async def test_dropbox_token(request: Request):
|
|||||||
account_name = account_info.get("name", {}).get("display_name", "Unknown user")
|
account_name = account_info.get("name", {}).get("display_name", "Unknown user")
|
||||||
|
|
||||||
# Dropbox refresh tokens don't expire, but we should note that in our response
|
# Dropbox refresh tokens don't expire, but we should note that in our response
|
||||||
token_info = {
|
token_info = {"expires_in_human": "Never expires (perpetual token)", "is_perpetual": True}
|
||||||
"expires_in_human": "Never expires (perpetual token)",
|
|
||||||
"is_perpetual": True
|
|
||||||
}
|
|
||||||
|
|
||||||
logger.info(f"Successfully connected to Dropbox as {account_email}")
|
logger.info(f"Successfully connected to Dropbox as {account_email}")
|
||||||
|
|
||||||
return {
|
return {
|
||||||
"status": "success",
|
"status": "success",
|
||||||
"message": f"Dropbox connection successful",
|
"message": "Dropbox connection successful",
|
||||||
"account": account_email,
|
"account": account_email,
|
||||||
"account_name": account_name,
|
"account_name": account_name,
|
||||||
"token_info": token_info
|
"token_info": token_info,
|
||||||
}
|
}
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Unexpected error testing Dropbox token: {str(e)}")
|
logger.exception(f"Unexpected error testing Dropbox token: {str(e)}")
|
||||||
return {
|
return {"status": "error", "message": f"Connection error: {str(e)}"}
|
||||||
"status": "error",
|
|
||||||
"message": f"Connection error: {str(e)}"
|
|
||||||
}
|
|
||||||
|
|
||||||
@router.post("/dropbox/save-settings")
|
@router.post("/dropbox/save-settings")
|
||||||
@require_login
|
@require_login
|
||||||
@@ -256,7 +184,7 @@ async def save_dropbox_settings(
|
|||||||
app_key: str = Form(None),
|
app_key: str = Form(None),
|
||||||
app_secret: str = Form(None),
|
app_secret: str = Form(None),
|
||||||
refresh_token: str = Form(...),
|
refresh_token: str = Form(...),
|
||||||
folder_path: str = Form(None)
|
folder_path: str = Form(None),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Save Dropbox settings to the .env file
|
Save Dropbox settings to the .env file
|
||||||
@@ -268,8 +196,7 @@ async def save_dropbox_settings(
|
|||||||
if not os.path.exists(env_path):
|
if not os.path.exists(env_path):
|
||||||
logger.error(f".env file not found at {env_path}")
|
logger.error(f".env file not found at {env_path}")
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail="Could not find .env file to update"
|
||||||
detail="Could not find .env file to update"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.info(f"Updating Dropbox settings in {env_path}")
|
logger.info(f"Updating Dropbox settings in {env_path}")
|
||||||
@@ -329,16 +256,12 @@ async def save_dropbox_settings(
|
|||||||
|
|
||||||
logger.info("Successfully updated Dropbox settings")
|
logger.info("Successfully updated Dropbox settings")
|
||||||
|
|
||||||
return {
|
return {"status": "success", "message": "Dropbox settings have been saved"}
|
||||||
"status": "success",
|
|
||||||
"message": "Dropbox settings have been saved"
|
|
||||||
}
|
|
||||||
|
|
||||||
except HTTPException:
|
except HTTPException:
|
||||||
raise
|
raise
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Unexpected error saving Dropbox settings: {str(e)}")
|
logger.exception(f"Unexpected error saving Dropbox settings: {str(e)}")
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=f"Failed to save Dropbox settings: {str(e)}"
|
||||||
detail=f"Failed to save Dropbox settings: {str(e)}"
|
|
||||||
)
|
)
|
||||||
|
|||||||
+69
-145
@@ -1,22 +1,23 @@
|
|||||||
"""
|
"""
|
||||||
Google Drive API endpoints
|
Google Drive API endpoints
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from fastapi import APIRouter, Request, HTTPException, status, Form
|
from fastapi import APIRouter, Request, HTTPException, status, Form
|
||||||
import logging
|
import logging
|
||||||
import os
|
import os
|
||||||
import requests
|
|
||||||
import json
|
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime
|
||||||
|
|
||||||
from app.auth import require_login
|
from app.auth import require_login
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
|
from app.utils.oauth_helper import exchange_oauth_token
|
||||||
|
|
||||||
# Set up logging
|
# Set up logging
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
router = APIRouter()
|
router = APIRouter()
|
||||||
|
|
||||||
|
|
||||||
@router.post("/google-drive/exchange-token")
|
@router.post("/google-drive/exchange-token")
|
||||||
@require_login
|
@require_login
|
||||||
async def exchange_google_drive_token(
|
async def exchange_google_drive_token(
|
||||||
@@ -25,90 +26,33 @@ async def exchange_google_drive_token(
|
|||||||
client_secret: str = Form(...),
|
client_secret: str = Form(...),
|
||||||
redirect_uri: str = Form(...),
|
redirect_uri: str = Form(...),
|
||||||
code: str = Form(...),
|
code: str = Form(...),
|
||||||
folder_id: Optional[str] = Form(None)
|
folder_id: Optional[str] = Form(None),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Exchange an authorization code for refresh and access tokens from Google.
|
Exchange an authorization code for refresh and access tokens from Google.
|
||||||
This is done on the server to avoid exposing client secret in the browser.
|
This is done on the server to avoid exposing client secret in the browser.
|
||||||
"""
|
"""
|
||||||
try:
|
# Prepare the token request
|
||||||
logger.info("Starting Google Drive token exchange process")
|
token_url = "https://oauth2.googleapis.com/token"
|
||||||
|
|
||||||
# Prepare the token request
|
payload = {
|
||||||
token_url = "https://oauth2.googleapis.com/token"
|
"client_id": client_id,
|
||||||
|
"client_secret": client_secret,
|
||||||
|
"code": code,
|
||||||
|
"redirect_uri": redirect_uri,
|
||||||
|
"grant_type": "authorization_code",
|
||||||
|
}
|
||||||
|
|
||||||
payload = {
|
# Use shared OAuth helper (handles secure logging and error handling)
|
||||||
'client_id': client_id,
|
token_data = exchange_oauth_token(provider_name="Google Drive", token_url=token_url, payload=payload)
|
||||||
'client_secret': client_secret,
|
|
||||||
'code': code,
|
|
||||||
'redirect_uri': redirect_uri,
|
|
||||||
'grant_type': 'authorization_code'
|
|
||||||
}
|
|
||||||
|
|
||||||
# Log request details (excluding secret)
|
# Return just what's needed by the frontend
|
||||||
safe_payload = payload.copy()
|
return {
|
||||||
safe_payload['client_secret'] = '[REDACTED]'
|
"refresh_token": token_data["refresh_token"],
|
||||||
safe_payload['code'] = f"{code[:5]}...{code[-5:]}" if len(code) > 10 else '[REDACTED]'
|
"access_token": token_data["access_token"],
|
||||||
logger.info(f"Token exchange request payload: {safe_payload}")
|
"expires_in": token_data.get("expires_in", 3600),
|
||||||
|
}
|
||||||
|
|
||||||
# Make the token request
|
|
||||||
logger.info("Sending POST request to Google for token exchange")
|
|
||||||
response = requests.post(token_url, data=payload, timeout=settings.http_request_timeout)
|
|
||||||
|
|
||||||
# Check if the request was successful
|
|
||||||
logger.info(f"Token exchange response status: {response.status_code}")
|
|
||||||
|
|
||||||
if response.status_code != 200:
|
|
||||||
# Log the error response for debugging
|
|
||||||
try:
|
|
||||||
error_json = response.json()
|
|
||||||
logger.error(f"Token exchange failed with status {response.status_code}: {error_json}")
|
|
||||||
error_detail = error_json
|
|
||||||
except Exception as json_err:
|
|
||||||
logger.error(f"Failed to parse error response as JSON: {str(json_err)}")
|
|
||||||
logger.error(f"Raw response content: {response.content[:500]}") # Limit log size
|
|
||||||
error_detail = {"error": "Unknown error", "raw_content_snippet": str(response.content[:100])}
|
|
||||||
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_400_BAD_REQUEST,
|
|
||||||
detail=f"Token exchange failed: {error_detail}"
|
|
||||||
)
|
|
||||||
|
|
||||||
# Return the token response
|
|
||||||
token_data = response.json()
|
|
||||||
|
|
||||||
# Validate the token response
|
|
||||||
if "refresh_token" not in token_data:
|
|
||||||
logger.error(f"Google returned success but no refresh token found in response: {token_data.keys()}")
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_502_BAD_GATEWAY,
|
|
||||||
detail="Google OAuth server returned success but no refresh token was included"
|
|
||||||
)
|
|
||||||
|
|
||||||
# Calculate token length for logging
|
|
||||||
refresh_token_length = len(token_data.get("refresh_token", ""))
|
|
||||||
access_token_length = len(token_data.get("access_token", ""))
|
|
||||||
|
|
||||||
logger.info(f"Successfully exchanged authorization code for Google Drive tokens. "
|
|
||||||
f"Refresh token length: {refresh_token_length}, "
|
|
||||||
f"Access token length: {access_token_length}")
|
|
||||||
|
|
||||||
# Return just what's needed by the frontend
|
|
||||||
return {
|
|
||||||
"refresh_token": token_data["refresh_token"],
|
|
||||||
"access_token": token_data["access_token"],
|
|
||||||
"expires_in": token_data.get("expires_in", 3600)
|
|
||||||
}
|
|
||||||
|
|
||||||
except HTTPException:
|
|
||||||
# Re-raise HTTP exceptions as they already have appropriate status codes
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
logger.exception(f"Unexpected error during Google Drive token exchange: {str(e)}")
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
|
||||||
detail=f"Failed to exchange token: {str(e)}"
|
|
||||||
)
|
|
||||||
|
|
||||||
@router.post("/google-drive/update-settings")
|
@router.post("/google-drive/update-settings")
|
||||||
@require_login
|
@require_login
|
||||||
@@ -118,7 +62,7 @@ async def update_google_drive_settings(
|
|||||||
client_secret: str = Form(None),
|
client_secret: str = Form(None),
|
||||||
refresh_token: str = Form(...),
|
refresh_token: str = Form(...),
|
||||||
folder_id: str = Form(None),
|
folder_id: str = Form(None),
|
||||||
use_oauth: str = Form("true")
|
use_oauth: str = Form("true"),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Update Google Drive settings in memory
|
Update Google Drive settings in memory
|
||||||
@@ -150,18 +94,16 @@ async def update_google_drive_settings(
|
|||||||
settings.google_drive_use_oauth = use_oauth_bool
|
settings.google_drive_use_oauth = use_oauth_bool
|
||||||
logger.info(f"Updated GOOGLE_DRIVE_USE_OAUTH in memory to {use_oauth_bool}")
|
logger.info(f"Updated GOOGLE_DRIVE_USE_OAUTH in memory to {use_oauth_bool}")
|
||||||
|
|
||||||
return {
|
return {"status": "success", "message": "Google Drive settings have been updated in memory"}
|
||||||
"status": "success",
|
|
||||||
"message": "Google Drive settings have been updated in memory"
|
|
||||||
}
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Unexpected error updating Google Drive settings: {str(e)}")
|
logger.exception(f"Unexpected error updating Google Drive settings: {str(e)}")
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
||||||
detail=f"Failed to update Google Drive settings: {str(e)}"
|
detail=f"Failed to update Google Drive settings: {str(e)}",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@router.get("/google-drive/test-token")
|
@router.get("/google-drive/test-token")
|
||||||
@require_login
|
@require_login
|
||||||
async def test_google_drive_token(request: Request):
|
async def test_google_drive_token(request: Request):
|
||||||
@@ -175,15 +117,14 @@ async def test_google_drive_token(request: Request):
|
|||||||
logger.info("Testing Google Drive token validity")
|
logger.info("Testing Google Drive token validity")
|
||||||
|
|
||||||
# Check if OAuth is enabled and configured
|
# Check if OAuth is enabled and configured
|
||||||
if getattr(settings, 'google_drive_use_oauth', False):
|
if getattr(settings, "google_drive_use_oauth", False):
|
||||||
if not (settings.google_drive_client_id and
|
if not (
|
||||||
settings.google_drive_client_secret and
|
settings.google_drive_client_id
|
||||||
settings.google_drive_refresh_token):
|
and settings.google_drive_client_secret
|
||||||
|
and settings.google_drive_refresh_token
|
||||||
|
):
|
||||||
logger.warning("Google Drive OAuth credentials not fully configured")
|
logger.warning("Google Drive OAuth credentials not fully configured")
|
||||||
return {
|
return {"status": "error", "message": "Google Drive OAuth credentials are not fully configured"}
|
||||||
"status": "error",
|
|
||||||
"message": "Google Drive OAuth credentials are not fully configured"
|
|
||||||
}
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# Test OAuth connection
|
# Test OAuth connection
|
||||||
@@ -198,7 +139,7 @@ async def test_google_drive_token(request: Request):
|
|||||||
refresh_token=settings.google_drive_refresh_token,
|
refresh_token=settings.google_drive_refresh_token,
|
||||||
token_uri="https://oauth2.googleapis.com/token",
|
token_uri="https://oauth2.googleapis.com/token",
|
||||||
client_id=settings.google_drive_client_id,
|
client_id=settings.google_drive_client_id,
|
||||||
client_secret=settings.google_drive_client_secret
|
client_secret=settings.google_drive_client_secret,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Force a refresh to update the token expiration
|
# Force a refresh to update the token expiration
|
||||||
@@ -207,14 +148,14 @@ async def test_google_drive_token(request: Request):
|
|||||||
|
|
||||||
# Get token expiration info
|
# Get token expiration info
|
||||||
expiration_info = {}
|
expiration_info = {}
|
||||||
if hasattr(credentials, 'expiry') and credentials.expiry:
|
if hasattr(credentials, "expiry") and credentials.expiry:
|
||||||
now = datetime.now()
|
now = datetime.now()
|
||||||
expiry = credentials.expiry
|
expiry = credentials.expiry
|
||||||
time_left = expiry - now
|
time_left = expiry - now
|
||||||
expiration_info = {
|
expiration_info = {
|
||||||
"expires_at": expiry.isoformat(),
|
"expires_at": expiry.isoformat(),
|
||||||
"expires_in_seconds": max(0, int(time_left.total_seconds())),
|
"expires_in_seconds": max(0, int(time_left.total_seconds())),
|
||||||
"expires_in_human": format_time_remaining(time_left)
|
"expires_in_human": format_time_remaining(time_left),
|
||||||
}
|
}
|
||||||
|
|
||||||
# Test basic API operation
|
# Test basic API operation
|
||||||
@@ -228,7 +169,7 @@ async def test_google_drive_token(request: Request):
|
|||||||
"message": f"OAuth token is valid! Connected as {user_email}",
|
"message": f"OAuth token is valid! Connected as {user_email}",
|
||||||
"account": user_email,
|
"account": user_email,
|
||||||
"auth_type": "oauth",
|
"auth_type": "oauth",
|
||||||
"token_info": expiration_info
|
"token_info": expiration_info,
|
||||||
}
|
}
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
error_msg = str(e)
|
error_msg = str(e)
|
||||||
@@ -239,20 +180,14 @@ async def test_google_drive_token(request: Request):
|
|||||||
return {
|
return {
|
||||||
"status": "error",
|
"status": "error",
|
||||||
"message": f"OAuth token validation failed: {error_msg}",
|
"message": f"OAuth token validation failed: {error_msg}",
|
||||||
"needs_reauth": True
|
"needs_reauth": True,
|
||||||
}
|
}
|
||||||
return {
|
return {"status": "error", "message": f"Connection error: {error_msg}"}
|
||||||
"status": "error",
|
|
||||||
"message": f"Connection error: {error_msg}"
|
|
||||||
}
|
|
||||||
else:
|
else:
|
||||||
# Test service account connection
|
# Test service account connection
|
||||||
if not settings.google_drive_credentials_json:
|
if not settings.google_drive_credentials_json:
|
||||||
logger.warning("Google Drive service account credentials not configured")
|
logger.warning("Google Drive service account credentials not configured")
|
||||||
return {
|
return {"status": "error", "message": "Google Drive service account credentials are not configured"}
|
||||||
"status": "error",
|
|
||||||
"message": "Google Drive service account credentials are not configured"
|
|
||||||
}
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
service = get_google_drive_service()
|
service = get_google_drive_service()
|
||||||
@@ -260,7 +195,7 @@ async def test_google_drive_token(request: Request):
|
|||||||
|
|
||||||
# For service accounts, try to show the delegated user if available
|
# For service accounts, try to show the delegated user if available
|
||||||
user_email = about.get("user", {}).get("emailAddress", "Unknown")
|
user_email = about.get("user", {}).get("emailAddress", "Unknown")
|
||||||
delegated_user = getattr(settings, 'google_drive_delegate_to', None)
|
delegated_user = getattr(settings, "google_drive_delegate_to", None)
|
||||||
|
|
||||||
if delegated_user:
|
if delegated_user:
|
||||||
user_display = f"{user_email} (delegating as {delegated_user})"
|
user_display = f"{user_email} (delegating as {delegated_user})"
|
||||||
@@ -273,22 +208,17 @@ async def test_google_drive_token(request: Request):
|
|||||||
"status": "success",
|
"status": "success",
|
||||||
"message": f"Service account is valid! Connected as {user_display}",
|
"message": f"Service account is valid! Connected as {user_display}",
|
||||||
"account": user_email,
|
"account": user_email,
|
||||||
"auth_type": "service_account"
|
"auth_type": "service_account",
|
||||||
}
|
}
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
error_msg = str(e)
|
error_msg = str(e)
|
||||||
logger.error(f"Google Drive service account test failed: {error_msg}")
|
logger.error(f"Google Drive service account test failed: {error_msg}")
|
||||||
return {
|
return {"status": "error", "message": f"Service account validation failed: {error_msg}"}
|
||||||
"status": "error",
|
|
||||||
"message": f"Service account validation failed: {error_msg}"
|
|
||||||
}
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception("Unexpected error testing Google Drive token")
|
logger.exception("Unexpected error testing Google Drive token")
|
||||||
return {
|
return {"status": "error", "message": f"Unexpected error: {str(e)}"}
|
||||||
"status": "error",
|
|
||||||
"message": f"Unexpected error: {str(e)}"
|
|
||||||
}
|
|
||||||
|
|
||||||
@router.get("/google-drive/get-token-info")
|
@router.get("/google-drive/get-token-info")
|
||||||
@require_login
|
@require_login
|
||||||
@@ -302,21 +232,20 @@ async def get_google_drive_token_info(request: Request):
|
|||||||
logger.info("Getting Google Drive token information")
|
logger.info("Getting Google Drive token information")
|
||||||
|
|
||||||
# Check if OAuth is enabled and configured
|
# Check if OAuth is enabled and configured
|
||||||
if not getattr(settings, 'google_drive_use_oauth', False):
|
if not getattr(settings, "google_drive_use_oauth", False):
|
||||||
logger.warning("OAuth is not enabled, using service account instead")
|
logger.warning("OAuth is not enabled, using service account instead")
|
||||||
return {
|
return {
|
||||||
"status": "error",
|
"status": "error",
|
||||||
"message": "OAuth is not enabled. Service accounts don't support user-facing features like folder picker."
|
"message": "OAuth is not enabled. Service accounts don't support user-facing features.",
|
||||||
}
|
}
|
||||||
|
|
||||||
if not (settings.google_drive_client_id and
|
if not (
|
||||||
settings.google_drive_client_secret and
|
settings.google_drive_client_id
|
||||||
settings.google_drive_refresh_token):
|
and settings.google_drive_client_secret
|
||||||
|
and settings.google_drive_refresh_token
|
||||||
|
):
|
||||||
logger.warning("Google Drive OAuth credentials not fully configured")
|
logger.warning("Google Drive OAuth credentials not fully configured")
|
||||||
return {
|
return {"status": "error", "message": "Google Drive OAuth credentials are not fully configured"}
|
||||||
"status": "error",
|
|
||||||
"message": "Google Drive OAuth credentials are not fully configured"
|
|
||||||
}
|
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# Get credentials and access token
|
# Get credentials and access token
|
||||||
@@ -328,7 +257,7 @@ async def get_google_drive_token_info(request: Request):
|
|||||||
refresh_token=settings.google_drive_refresh_token,
|
refresh_token=settings.google_drive_refresh_token,
|
||||||
token_uri="https://oauth2.googleapis.com/token",
|
token_uri="https://oauth2.googleapis.com/token",
|
||||||
client_id=settings.google_drive_client_id,
|
client_id=settings.google_drive_client_id,
|
||||||
client_secret=settings.google_drive_client_secret
|
client_secret=settings.google_drive_client_secret,
|
||||||
)
|
)
|
||||||
|
|
||||||
# Force a refresh to get a fresh access token
|
# Force a refresh to get a fresh access token
|
||||||
@@ -337,14 +266,14 @@ async def get_google_drive_token_info(request: Request):
|
|||||||
|
|
||||||
# Get token expiration info
|
# Get token expiration info
|
||||||
expiration_info = {}
|
expiration_info = {}
|
||||||
if hasattr(credentials, 'expiry') and credentials.expiry:
|
if hasattr(credentials, "expiry") and credentials.expiry:
|
||||||
now = datetime.now()
|
now = datetime.now()
|
||||||
expiry = credentials.expiry
|
expiry = credentials.expiry
|
||||||
time_left = expiry - now
|
time_left = expiry - now
|
||||||
expiration_info = {
|
expiration_info = {
|
||||||
"expires_at": expiry.isoformat(),
|
"expires_at": expiry.isoformat(),
|
||||||
"expires_in_seconds": max(0, int(time_left.total_seconds())),
|
"expires_in_seconds": max(0, int(time_left.total_seconds())),
|
||||||
"expires_in_human": format_time_remaining(time_left)
|
"expires_in_human": format_time_remaining(time_left),
|
||||||
}
|
}
|
||||||
|
|
||||||
# Return the token info
|
# Return the token info
|
||||||
@@ -354,7 +283,7 @@ async def get_google_drive_token_info(request: Request):
|
|||||||
"status": "success",
|
"status": "success",
|
||||||
"message": "Access token successfully retrieved",
|
"message": "Access token successfully retrieved",
|
||||||
"access_token": credentials.token,
|
"access_token": credentials.token,
|
||||||
"token_info": expiration_info
|
"token_info": expiration_info,
|
||||||
}
|
}
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
error_msg = str(e)
|
error_msg = str(e)
|
||||||
@@ -365,19 +294,14 @@ async def get_google_drive_token_info(request: Request):
|
|||||||
return {
|
return {
|
||||||
"status": "error",
|
"status": "error",
|
||||||
"message": f"OAuth token retrieval failed: {error_msg}",
|
"message": f"OAuth token retrieval failed: {error_msg}",
|
||||||
"needs_reauth": True
|
"needs_reauth": True,
|
||||||
}
|
}
|
||||||
return {
|
return {"status": "error", "message": f"Token retrieval error: {error_msg}"}
|
||||||
"status": "error",
|
|
||||||
"message": f"Token retrieval error: {error_msg}"
|
|
||||||
}
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception("Unexpected error getting Google Drive token info")
|
logger.exception("Unexpected error getting Google Drive token info")
|
||||||
return {
|
return {"status": "error", "message": f"Unexpected error: {str(e)}"}
|
||||||
"status": "error",
|
|
||||||
"message": f"Unexpected error: {str(e)}"
|
|
||||||
}
|
|
||||||
|
|
||||||
def format_time_remaining(time_delta):
|
def format_time_remaining(time_delta):
|
||||||
"""Format a timedelta into a human-readable string."""
|
"""Format a timedelta into a human-readable string."""
|
||||||
@@ -398,6 +322,7 @@ def format_time_remaining(time_delta):
|
|||||||
|
|
||||||
return ", ".join(parts)
|
return ", ".join(parts)
|
||||||
|
|
||||||
|
|
||||||
@router.post("/google-drive/save-settings")
|
@router.post("/google-drive/save-settings")
|
||||||
@require_login
|
@require_login
|
||||||
async def save_dropbox_settings(
|
async def save_dropbox_settings(
|
||||||
@@ -406,7 +331,7 @@ async def save_dropbox_settings(
|
|||||||
client_secret: str = Form(None),
|
client_secret: str = Form(None),
|
||||||
refresh_token: str = Form(...),
|
refresh_token: str = Form(...),
|
||||||
folder_id: str = Form(None),
|
folder_id: str = Form(None),
|
||||||
use_oauth: str = Form("true")
|
use_oauth: str = Form("true"),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Save Google Drive settings to the .env file
|
Save Google Drive settings to the .env file
|
||||||
@@ -419,9 +344,7 @@ async def save_dropbox_settings(
|
|||||||
use_oauth_bool = use_oauth.lower() in ("true", "1", "yes", "y", "t")
|
use_oauth_bool = use_oauth.lower() in ("true", "1", "yes", "y", "t")
|
||||||
|
|
||||||
# Define settings to update
|
# Define settings to update
|
||||||
drive_settings = {
|
drive_settings = {"GOOGLE_DRIVE_USE_OAUTH": str(use_oauth_bool).lower()}
|
||||||
"GOOGLE_DRIVE_USE_OAUTH": str(use_oauth_bool).lower()
|
|
||||||
}
|
|
||||||
|
|
||||||
# Only update these if provided
|
# Only update these if provided
|
||||||
if use_oauth_bool:
|
if use_oauth_bool:
|
||||||
@@ -475,7 +398,9 @@ async def save_dropbox_settings(
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.warning(f"Failed to update .env file: {str(e)}, but will continue with in-memory update")
|
logger.warning(f"Failed to update .env file: {str(e)}, but will continue with in-memory update")
|
||||||
else:
|
else:
|
||||||
logger.warning(f".env file not found at {env_path}, skipping file update but continuing with in-memory update")
|
logger.warning(
|
||||||
|
f".env file not found at {env_path}, skipping file update but continuing with in-memory update"
|
||||||
|
)
|
||||||
|
|
||||||
# Update the settings in memory (this always happens)
|
# Update the settings in memory (this always happens)
|
||||||
if refresh_token:
|
if refresh_token:
|
||||||
@@ -495,12 +420,11 @@ async def save_dropbox_settings(
|
|||||||
return {
|
return {
|
||||||
"status": "success",
|
"status": "success",
|
||||||
"message": "Google Drive settings have been saved",
|
"message": "Google Drive settings have been saved",
|
||||||
"in_memory_only": not os.path.exists(env_path)
|
"in_memory_only": not os.path.exists(env_path),
|
||||||
}
|
}
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Unexpected error saving Google Drive settings: {str(e)}")
|
logger.exception(f"Unexpected error saving Google Drive settings: {str(e)}")
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=f"Failed to save Google Drive settings: {str(e)}"
|
||||||
detail=f"Failed to save Google Drive settings: {str(e)}"
|
|
||||||
)
|
)
|
||||||
|
|||||||
+57
-135
@@ -1,22 +1,23 @@
|
|||||||
"""
|
"""
|
||||||
OneDrive API endpoints
|
OneDrive API endpoints
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from fastapi import APIRouter, Request, HTTPException, status, Form
|
from fastapi import APIRouter, Request, HTTPException, status, Form
|
||||||
import logging
|
import logging
|
||||||
import os
|
import os
|
||||||
import requests
|
import requests
|
||||||
import json
|
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta
|
||||||
from typing import Optional
|
|
||||||
|
|
||||||
from app.auth import require_login
|
from app.auth import require_login
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
|
from app.utils.oauth_helper import exchange_oauth_token
|
||||||
|
|
||||||
# Set up logging
|
# Set up logging
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
router = APIRouter()
|
router = APIRouter()
|
||||||
|
|
||||||
|
|
||||||
@router.post("/onedrive/exchange-token")
|
@router.post("/onedrive/exchange-token")
|
||||||
@require_login
|
@require_login
|
||||||
async def exchange_onedrive_token(
|
async def exchange_onedrive_token(
|
||||||
@@ -25,91 +26,30 @@ async def exchange_onedrive_token(
|
|||||||
client_secret: str = Form(...),
|
client_secret: str = Form(...),
|
||||||
redirect_uri: str = Form(...),
|
redirect_uri: str = Form(...),
|
||||||
code: str = Form(...),
|
code: str = Form(...),
|
||||||
tenant_id: str = Form(...)
|
tenant_id: str = Form(...),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Exchange an authorization code for a refresh token.
|
Exchange an authorization code for a refresh token.
|
||||||
This is done on the server to avoid exposing client secret in the browser.
|
This is done on the server to avoid exposing client secret in the browser.
|
||||||
"""
|
"""
|
||||||
try:
|
# Prepare the token request
|
||||||
logger.info(f"Starting OneDrive token exchange process with tenant_id: {tenant_id}")
|
token_url = f"https://login.microsoftonline.com/{tenant_id}/oauth2/v2.0/token"
|
||||||
|
|
||||||
# Prepare the token request
|
payload = {
|
||||||
token_url = f"https://login.microsoftonline.com/{tenant_id}/oauth2/v2.0/token"
|
"client_id": client_id,
|
||||||
logger.info(f"Using token URL: {token_url}")
|
"scope": "https://graph.microsoft.com/.default offline_access",
|
||||||
|
"code": code,
|
||||||
|
"redirect_uri": redirect_uri,
|
||||||
|
"grant_type": "authorization_code",
|
||||||
|
"client_secret": client_secret,
|
||||||
|
}
|
||||||
|
|
||||||
payload = {
|
# Use shared OAuth helper (handles secure logging and error handling)
|
||||||
'client_id': client_id,
|
token_data = exchange_oauth_token(provider_name="OneDrive", token_url=token_url, payload=payload)
|
||||||
'scope': 'https://graph.microsoft.com/.default offline_access',
|
|
||||||
'code': code,
|
|
||||||
'redirect_uri': redirect_uri,
|
|
||||||
'grant_type': 'authorization_code',
|
|
||||||
'client_secret': client_secret
|
|
||||||
}
|
|
||||||
|
|
||||||
# Log request details (excluding secret)
|
# Return just what's needed by the frontend
|
||||||
safe_payload = payload.copy()
|
return {"refresh_token": token_data["refresh_token"], "expires_in": token_data.get("expires_in", 3600)}
|
||||||
safe_payload['client_secret'] = '[REDACTED]'
|
|
||||||
safe_payload['code'] = f"{code[:5]}...{code[-5:]}" if len(code) > 10 else '[REDACTED]'
|
|
||||||
logger.info(f"Token exchange request payload: {safe_payload}")
|
|
||||||
|
|
||||||
# Make the token request
|
|
||||||
logger.info("Sending POST request to Microsoft for token exchange")
|
|
||||||
response = requests.post(token_url, data=payload, timeout=settings.http_request_timeout)
|
|
||||||
|
|
||||||
# Check if the request was successful
|
|
||||||
logger.info(f"Token exchange response status: {response.status_code}")
|
|
||||||
|
|
||||||
if response.status_code != 200:
|
|
||||||
# Log the error response for debugging
|
|
||||||
try:
|
|
||||||
error_json = response.json()
|
|
||||||
logger.error(f"Token exchange failed with status {response.status_code}: {error_json}")
|
|
||||||
error_detail = error_json
|
|
||||||
except Exception as json_err:
|
|
||||||
logger.error(f"Failed to parse error response as JSON: {str(json_err)}")
|
|
||||||
logger.error(f"Raw response content: {response.content[:500]}") # Limit log size
|
|
||||||
error_detail = {"error": "Unknown error", "raw_content_snippet": str(response.content[:100])}
|
|
||||||
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_400_BAD_REQUEST,
|
|
||||||
detail=f"Token exchange failed: {error_detail}"
|
|
||||||
)
|
|
||||||
|
|
||||||
# Return the token response
|
|
||||||
token_data = response.json()
|
|
||||||
|
|
||||||
# Validate the token response
|
|
||||||
if "refresh_token" not in token_data:
|
|
||||||
logger.error(f"Microsoft returned success but no refresh token found in response: {token_data.keys()}")
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_502_BAD_GATEWAY,
|
|
||||||
detail="Microsoft OAuth server returned success but no refresh token was included"
|
|
||||||
)
|
|
||||||
|
|
||||||
# Calculate token length for logging
|
|
||||||
refresh_token_length = len(token_data.get("refresh_token", ""))
|
|
||||||
access_token_length = len(token_data.get("access_token", ""))
|
|
||||||
|
|
||||||
logger.info(f"Successfully exchanged authorization code for OneDrive tokens. "
|
|
||||||
f"Refresh token length: {refresh_token_length}, "
|
|
||||||
f"Access token length: {access_token_length}")
|
|
||||||
|
|
||||||
# Return just what's needed by the frontend
|
|
||||||
return {
|
|
||||||
"refresh_token": token_data["refresh_token"],
|
|
||||||
"expires_in": token_data.get("expires_in", 3600)
|
|
||||||
}
|
|
||||||
|
|
||||||
except HTTPException:
|
|
||||||
# Re-raise HTTP exceptions as they already have appropriate status codes
|
|
||||||
raise
|
|
||||||
except Exception as e:
|
|
||||||
logger.exception(f"Unexpected error during OneDrive token exchange: {str(e)}")
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
|
||||||
detail=f"Failed to exchange token: {str(e)}"
|
|
||||||
)
|
|
||||||
|
|
||||||
@router.get("/onedrive/test-token")
|
@router.get("/onedrive/test-token")
|
||||||
@require_login
|
@require_login
|
||||||
@@ -120,12 +60,13 @@ async def test_onedrive_token(request: Request):
|
|||||||
try:
|
try:
|
||||||
logger.info("Testing OneDrive token validity")
|
logger.info("Testing OneDrive token validity")
|
||||||
|
|
||||||
if not settings.onedrive_refresh_token or not settings.onedrive_client_id or not settings.onedrive_client_secret:
|
if (
|
||||||
|
not settings.onedrive_refresh_token
|
||||||
|
or not settings.onedrive_client_id
|
||||||
|
or not settings.onedrive_client_secret
|
||||||
|
):
|
||||||
logger.warning("OneDrive credentials not fully configured")
|
logger.warning("OneDrive credentials not fully configured")
|
||||||
return {
|
return {"status": "error", "message": "OneDrive credentials are not fully configured"}
|
||||||
"status": "error",
|
|
||||||
"message": "OneDrive credentials are not fully configured"
|
|
||||||
}
|
|
||||||
|
|
||||||
# Refresh token to get a new access token and expiration info
|
# Refresh token to get a new access token and expiration info
|
||||||
tenant_id = settings.onedrive_tenant_id or "common"
|
tenant_id = settings.onedrive_tenant_id or "common"
|
||||||
@@ -136,18 +77,14 @@ async def test_onedrive_token(request: Request):
|
|||||||
"client_secret": settings.onedrive_client_secret,
|
"client_secret": settings.onedrive_client_secret,
|
||||||
"refresh_token": settings.onedrive_refresh_token,
|
"refresh_token": settings.onedrive_refresh_token,
|
||||||
"grant_type": "refresh_token",
|
"grant_type": "refresh_token",
|
||||||
"scope": "offline_access Files.ReadWrite"
|
"scope": "offline_access Files.ReadWrite",
|
||||||
}
|
}
|
||||||
|
|
||||||
response = requests.post(token_url, data=refresh_data, timeout=settings.http_request_timeout)
|
response = requests.post(token_url, data=refresh_data, timeout=settings.http_request_timeout)
|
||||||
|
|
||||||
if response.status_code != 200:
|
if response.status_code != 200:
|
||||||
logger.error(f"Failed to refresh OneDrive token: {response.text}")
|
logger.error(f"Failed to refresh OneDrive token: {response.text}")
|
||||||
return {
|
return {"status": "error", "message": "Refresh token has expired or is invalid", "needs_reauth": True}
|
||||||
"status": "error",
|
|
||||||
"message": "Refresh token has expired or is invalid",
|
|
||||||
"needs_reauth": True
|
|
||||||
}
|
|
||||||
|
|
||||||
token_data = response.json()
|
token_data = response.json()
|
||||||
access_token = token_data.get("access_token")
|
access_token = token_data.get("access_token")
|
||||||
@@ -199,7 +136,7 @@ async def test_onedrive_token(request: Request):
|
|||||||
logger.error(f"OneDrive token test failed: {user_response.status_code} {user_response.text}")
|
logger.error(f"OneDrive token test failed: {user_response.status_code} {user_response.text}")
|
||||||
return {
|
return {
|
||||||
"status": "error",
|
"status": "error",
|
||||||
"message": f"Token validation failed with status {user_response.status_code}: {user_response.text}"
|
"message": f"Token validation failed with status {user_response.status_code}: {user_response.text}",
|
||||||
}
|
}
|
||||||
|
|
||||||
# Get user info
|
# Get user info
|
||||||
@@ -217,25 +154,23 @@ async def test_onedrive_token(request: Request):
|
|||||||
"expires_at": expiry_time.isoformat(),
|
"expires_at": expiry_time.isoformat(),
|
||||||
"expires_in_seconds": expires_in,
|
"expires_in_seconds": expires_in,
|
||||||
"expires_in_human": format_time_remaining(time_left),
|
"expires_in_human": format_time_remaining(time_left),
|
||||||
"refresh_token_validity": "Refresh token is valid for 90 days of inactivity"
|
"refresh_token_validity": "Refresh token is valid for 90 days of inactivity",
|
||||||
}
|
}
|
||||||
|
|
||||||
logger.info(f"Successfully connected to OneDrive as {email}")
|
logger.info(f"Successfully connected to OneDrive as {email}")
|
||||||
|
|
||||||
return {
|
return {
|
||||||
"status": "success",
|
"status": "success",
|
||||||
"message": f"OneDrive connection successful",
|
"message": "OneDrive connection successful",
|
||||||
"account": email,
|
"account": email,
|
||||||
"account_name": display_name,
|
"account_name": display_name,
|
||||||
"token_info": token_info
|
"token_info": token_info,
|
||||||
}
|
}
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Unexpected error testing OneDrive token: {str(e)}")
|
logger.exception(f"Unexpected error testing OneDrive token: {str(e)}")
|
||||||
return {
|
return {"status": "error", "message": f"Connection error: {str(e)}"}
|
||||||
"status": "error",
|
|
||||||
"message": f"Connection error: {str(e)}"
|
|
||||||
}
|
|
||||||
|
|
||||||
def format_time_remaining(time_delta):
|
def format_time_remaining(time_delta):
|
||||||
"""Format a timedelta into a human-readable string."""
|
"""Format a timedelta into a human-readable string."""
|
||||||
@@ -256,6 +191,7 @@ def format_time_remaining(time_delta):
|
|||||||
|
|
||||||
return ", ".join(parts)
|
return ", ".join(parts)
|
||||||
|
|
||||||
|
|
||||||
@router.post("/onedrive/save-settings")
|
@router.post("/onedrive/save-settings")
|
||||||
@require_login
|
@require_login
|
||||||
async def save_onedrive_settings(
|
async def save_onedrive_settings(
|
||||||
@@ -264,7 +200,7 @@ async def save_onedrive_settings(
|
|||||||
client_secret: str = Form(None),
|
client_secret: str = Form(None),
|
||||||
refresh_token: str = Form(...),
|
refresh_token: str = Form(...),
|
||||||
tenant_id: str = Form("common"),
|
tenant_id: str = Form("common"),
|
||||||
folder_path: str = Form(None)
|
folder_path: str = Form(None),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Save OneDrive settings to the .env file
|
Save OneDrive settings to the .env file
|
||||||
@@ -276,8 +212,7 @@ async def save_onedrive_settings(
|
|||||||
if not os.path.exists(env_path):
|
if not os.path.exists(env_path):
|
||||||
logger.error(f".env file not found at {env_path}")
|
logger.error(f".env file not found at {env_path}")
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail="Could not find .env file to update"
|
||||||
detail="Could not find .env file to update"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
logger.info(f"Updating OneDrive settings in {env_path}")
|
logger.info(f"Updating OneDrive settings in {env_path}")
|
||||||
@@ -341,20 +276,17 @@ async def save_onedrive_settings(
|
|||||||
|
|
||||||
logger.info("Successfully updated OneDrive settings")
|
logger.info("Successfully updated OneDrive settings")
|
||||||
|
|
||||||
return {
|
return {"status": "success", "message": "OneDrive settings have been saved"}
|
||||||
"status": "success",
|
|
||||||
"message": "OneDrive settings have been saved"
|
|
||||||
}
|
|
||||||
|
|
||||||
except HTTPException:
|
except HTTPException:
|
||||||
raise
|
raise
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Unexpected error saving OneDrive settings: {str(e)}")
|
logger.exception(f"Unexpected error saving OneDrive settings: {str(e)}")
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=f"Failed to save OneDrive settings: {str(e)}"
|
||||||
detail=f"Failed to save OneDrive settings: {str(e)}"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@router.post("/onedrive/update-settings")
|
@router.post("/onedrive/update-settings")
|
||||||
@require_login
|
@require_login
|
||||||
async def update_onedrive_settings(
|
async def update_onedrive_settings(
|
||||||
@@ -363,7 +295,7 @@ async def update_onedrive_settings(
|
|||||||
client_secret: str = Form(None),
|
client_secret: str = Form(None),
|
||||||
refresh_token: str = Form(...),
|
refresh_token: str = Form(...),
|
||||||
tenant_id: str = Form("common"),
|
tenant_id: str = Form("common"),
|
||||||
folder_path: str = Form(None)
|
folder_path: str = Form(None),
|
||||||
):
|
):
|
||||||
"""
|
"""
|
||||||
Update OneDrive settings in memory (without modifying .env file)
|
Update OneDrive settings in memory (without modifying .env file)
|
||||||
@@ -395,27 +327,22 @@ async def update_onedrive_settings(
|
|||||||
# Test the token to make sure it works
|
# Test the token to make sure it works
|
||||||
try:
|
try:
|
||||||
from app.tasks.upload_to_onedrive import get_onedrive_token
|
from app.tasks.upload_to_onedrive import get_onedrive_token
|
||||||
access_token = get_onedrive_token()
|
|
||||||
|
get_onedrive_token() # Test that token can be retrieved
|
||||||
logger.info("Successfully tested OneDrive token")
|
logger.info("Successfully tested OneDrive token")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"Token test failed after updating settings: {str(e)}")
|
logger.error(f"Token test failed after updating settings: {str(e)}")
|
||||||
return {
|
return {"status": "warning", "message": "Settings updated but token test failed: " + str(e)}
|
||||||
"status": "warning",
|
|
||||||
"message": "Settings updated but token test failed: " + str(e)
|
|
||||||
}
|
|
||||||
|
|
||||||
return {
|
return {"status": "success", "message": "OneDrive settings have been updated in memory"}
|
||||||
"status": "success",
|
|
||||||
"message": "OneDrive settings have been updated in memory"
|
|
||||||
}
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception(f"Unexpected error updating OneDrive settings: {str(e)}")
|
logger.exception(f"Unexpected error updating OneDrive settings: {str(e)}")
|
||||||
raise HTTPException(
|
raise HTTPException(
|
||||||
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR,
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=f"Failed to update OneDrive settings: {str(e)}"
|
||||||
detail=f"Failed to update OneDrive settings: {str(e)}"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
@router.get("/onedrive/get-full-config")
|
@router.get("/onedrive/get-full-config")
|
||||||
@require_login
|
@require_login
|
||||||
async def get_onedrive_full_config(request: Request):
|
async def get_onedrive_full_config(request: Request):
|
||||||
@@ -429,26 +356,21 @@ async def get_onedrive_full_config(request: Request):
|
|||||||
"client_secret": settings.onedrive_client_secret or "",
|
"client_secret": settings.onedrive_client_secret or "",
|
||||||
"tenant_id": settings.onedrive_tenant_id or "common",
|
"tenant_id": settings.onedrive_tenant_id or "common",
|
||||||
"refresh_token": settings.onedrive_refresh_token or "",
|
"refresh_token": settings.onedrive_refresh_token or "",
|
||||||
"folder_path": settings.onedrive_folder_path or "Documents/Uploads"
|
"folder_path": settings.onedrive_folder_path or "Documents/Uploads",
|
||||||
}
|
}
|
||||||
|
|
||||||
# Generate environment variable format
|
# Generate environment variable format
|
||||||
env_format = "\n".join([
|
env_format = "\n".join(
|
||||||
f"ONEDRIVE_CLIENT_ID={config['client_id']}",
|
[
|
||||||
f"ONEDRIVE_CLIENT_SECRET={config['client_secret']}",
|
f"ONEDRIVE_CLIENT_ID={config['client_id']}",
|
||||||
f"ONEDRIVE_TENANT_ID={config['tenant_id']}",
|
f"ONEDRIVE_CLIENT_SECRET={config['client_secret']}",
|
||||||
f"ONEDRIVE_REFRESH_TOKEN={config['refresh_token']}",
|
f"ONEDRIVE_TENANT_ID={config['tenant_id']}",
|
||||||
f"ONEDRIVE_FOLDER_PATH={config['folder_path']}"
|
f"ONEDRIVE_REFRESH_TOKEN={config['refresh_token']}",
|
||||||
])
|
f"ONEDRIVE_FOLDER_PATH={config['folder_path']}",
|
||||||
|
]
|
||||||
|
)
|
||||||
|
|
||||||
return {
|
return {"status": "success", "config": config, "env_format": env_format}
|
||||||
"status": "success",
|
|
||||||
"config": config,
|
|
||||||
"env_format": env_format
|
|
||||||
}
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.exception("Error getting OneDrive configuration")
|
logger.exception("Error getting OneDrive configuration")
|
||||||
return {
|
return {"status": "error", "message": str(e)}
|
||||||
"status": "error",
|
|
||||||
"message": str(e)
|
|
||||||
}
|
|
||||||
|
|||||||
@@ -2,16 +2,14 @@
|
|||||||
|
|
||||||
import os
|
import os
|
||||||
import subprocess
|
import subprocess
|
||||||
import json
|
|
||||||
import tempfile
|
|
||||||
import logging
|
import logging
|
||||||
from pathlib import Path
|
|
||||||
from app.config import settings
|
from app.config import settings
|
||||||
from app.tasks.retry_config import BaseTaskWithRetry
|
from app.tasks.retry_config import BaseTaskWithRetry
|
||||||
from app.celery_app import celery
|
from app.celery_app import celery
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
@celery.task(base=BaseTaskWithRetry)
|
@celery.task(base=BaseTaskWithRetry)
|
||||||
def upload_with_rclone(file_path: str, destination: str):
|
def upload_with_rclone(file_path: str, destination: str):
|
||||||
"""
|
"""
|
||||||
@@ -28,6 +26,18 @@ def upload_with_rclone(file_path: str, destination: str):
|
|||||||
# Extract filename
|
# Extract filename
|
||||||
filename = os.path.basename(file_path)
|
filename = os.path.basename(file_path)
|
||||||
|
|
||||||
|
# Validate destination format to prevent command injection
|
||||||
|
if ":" not in destination:
|
||||||
|
raise ValueError(f"Invalid destination format: {destination}. Expected format: remote:path")
|
||||||
|
|
||||||
|
# Split and validate destination components
|
||||||
|
remote, remote_path = destination.split(":", 1)
|
||||||
|
|
||||||
|
# Validate remote name to prevent command injection
|
||||||
|
# Must start with alphanumeric, can contain alphanumeric, underscore, hyphen
|
||||||
|
if not remote or not remote[0].isalnum() or not all(c.isalnum() or c in ("_", "-") for c in remote):
|
||||||
|
raise ValueError(f"Invalid remote name: {remote}")
|
||||||
|
|
||||||
# Check if rclone is installed and config exists
|
# Check if rclone is installed and config exists
|
||||||
rclone_config_path = os.path.join(settings.workdir, "rclone.conf")
|
rclone_config_path = os.path.join(settings.workdir, "rclone.conf")
|
||||||
if not os.path.exists(rclone_config_path):
|
if not os.path.exists(rclone_config_path):
|
||||||
@@ -36,31 +46,13 @@ def upload_with_rclone(file_path: str, destination: str):
|
|||||||
raise ValueError(error_msg)
|
raise ValueError(error_msg)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# Split destination into remote and path
|
|
||||||
if ":" not in destination:
|
|
||||||
raise ValueError(f"Invalid destination format: {destination}. Expected format: remote:path")
|
|
||||||
|
|
||||||
remote, remote_path = destination.split(":", 1)
|
|
||||||
|
|
||||||
# Ensure the remote path exists (create folders if needed)
|
# Ensure the remote path exists (create folders if needed)
|
||||||
mkdir_cmd = [
|
mkdir_cmd = ["rclone", "mkdir", "--config", rclone_config_path, destination]
|
||||||
"rclone",
|
|
||||||
"mkdir",
|
|
||||||
"--config", rclone_config_path,
|
|
||||||
destination
|
|
||||||
]
|
|
||||||
|
|
||||||
subprocess.run(mkdir_cmd, check=True, capture_output=True)
|
subprocess.run(mkdir_cmd, check=True, capture_output=True)
|
||||||
|
|
||||||
# Construct the upload command
|
# Construct the upload command
|
||||||
upload_cmd = [
|
upload_cmd = ["rclone", "copy", "--config", rclone_config_path, file_path, destination, "--progress"]
|
||||||
"rclone",
|
|
||||||
"copy",
|
|
||||||
"--config", rclone_config_path,
|
|
||||||
file_path,
|
|
||||||
destination,
|
|
||||||
"--progress"
|
|
||||||
]
|
|
||||||
|
|
||||||
# Execute the upload command
|
# Execute the upload command
|
||||||
result = subprocess.run(upload_cmd, check=True, capture_output=True, text=True)
|
result = subprocess.run(upload_cmd, check=True, capture_output=True, text=True)
|
||||||
@@ -69,38 +61,29 @@ def upload_with_rclone(file_path: str, destination: str):
|
|||||||
if result.returncode == 0:
|
if result.returncode == 0:
|
||||||
# Try to get a public link if possible
|
# Try to get a public link if possible
|
||||||
try:
|
try:
|
||||||
link_cmd = [
|
link_cmd = ["rclone", "link", "--config", rclone_config_path, f"{destination}/{filename}"]
|
||||||
"rclone",
|
|
||||||
"link",
|
|
||||||
"--config", rclone_config_path,
|
|
||||||
f"{destination}/{filename}"
|
|
||||||
]
|
|
||||||
link_result = subprocess.run(link_cmd, capture_output=True, text=True)
|
link_result = subprocess.run(link_cmd, capture_output=True, text=True)
|
||||||
public_url = link_result.stdout.strip() if link_result.returncode == 0 else None
|
public_url = link_result.stdout.strip() if link_result.returncode == 0 else None
|
||||||
except Exception:
|
except (subprocess.SubprocessError, OSError) as e:
|
||||||
|
logger.warning(f"Failed to get public link for {filename}: {str(e)}")
|
||||||
public_url = None
|
public_url = None
|
||||||
|
|
||||||
logger.info(f"Successfully uploaded {filename} to {destination}")
|
logger.info(f"Successfully uploaded {filename} to {destination}")
|
||||||
return {
|
return {"status": "Completed", "file": file_path, "destination": destination, "public_url": public_url}
|
||||||
"status": "Completed",
|
|
||||||
"file": file_path,
|
|
||||||
"destination": destination,
|
|
||||||
"public_url": public_url
|
|
||||||
}
|
|
||||||
else:
|
else:
|
||||||
error_msg = f"Failed to upload {filename} to {destination}: {result.stderr}"
|
error_msg = f"Failed to upload {filename} to {destination}: {result.stderr}"
|
||||||
logger.error(error_msg)
|
logger.error(error_msg)
|
||||||
raise Exception(error_msg)
|
raise RuntimeError(error_msg)
|
||||||
|
|
||||||
except subprocess.CalledProcessError as e:
|
except subprocess.CalledProcessError as e:
|
||||||
error_msg = f"Rclone error: {e.stderr.decode('utf-8') if hasattr(e.stderr, 'decode') else e.stderr}"
|
error_msg = f"Rclone error: {e.stderr.decode('utf-8') if hasattr(e.stderr, 'decode') else e.stderr}"
|
||||||
logger.error(error_msg)
|
logger.error(error_msg)
|
||||||
raise Exception(error_msg)
|
raise RuntimeError(error_msg) from e
|
||||||
|
|
||||||
except Exception as e:
|
except (OSError, ValueError) as e:
|
||||||
error_msg = f"Error uploading {filename} to {destination}: {str(e)}"
|
error_msg = f"Error uploading {filename} to {destination}: {str(e)}"
|
||||||
logger.error(error_msg)
|
logger.error(error_msg)
|
||||||
raise Exception(error_msg)
|
raise RuntimeError(error_msg) from e
|
||||||
|
|
||||||
|
|
||||||
@celery.task(base=BaseTaskWithRetry)
|
@celery.task(base=BaseTaskWithRetry)
|
||||||
@@ -134,7 +117,7 @@ def send_to_all_rclone_destinations(file_path: str):
|
|||||||
# Target directories for each remote (from settings)
|
# Target directories for each remote (from settings)
|
||||||
remote_paths = {}
|
remote_paths = {}
|
||||||
for remote in remotes:
|
for remote in remotes:
|
||||||
remote_name = remote.rstrip(':')
|
remote_name = remote.rstrip(":")
|
||||||
path_setting_name = f"rclone_{remote_name}_path"
|
path_setting_name = f"rclone_{remote_name}_path"
|
||||||
if hasattr(settings, path_setting_name) and getattr(settings, path_setting_name):
|
if hasattr(settings, path_setting_name) and getattr(settings, path_setting_name):
|
||||||
remote_paths[remote] = getattr(settings, path_setting_name)
|
remote_paths[remote] = getattr(settings, path_setting_name)
|
||||||
@@ -146,24 +129,20 @@ def send_to_all_rclone_destinations(file_path: str):
|
|||||||
results = {}
|
results = {}
|
||||||
for remote, path in remote_paths.items():
|
for remote, path in remote_paths.items():
|
||||||
full_destination = f"{remote}{path}"
|
full_destination = f"{remote}{path}"
|
||||||
if path and not path.endswith('/'):
|
if path and not path.endswith("/"):
|
||||||
full_destination += '/'
|
full_destination += "/"
|
||||||
|
|
||||||
logger.info(f"Queueing {file_path} for upload to {full_destination}")
|
logger.info(f"Queueing {file_path} for upload to {full_destination}")
|
||||||
task = upload_with_rclone.delay(file_path, full_destination)
|
task = upload_with_rclone.delay(file_path, full_destination)
|
||||||
results[f"rclone_{remote.rstrip(':')}_task_id"] = task.id
|
results[f"rclone_{remote.rstrip(':')}_task_id"] = task.id
|
||||||
|
|
||||||
return {
|
return {"status": "Queued", "file_path": file_path, "tasks": results}
|
||||||
"status": "Queued",
|
|
||||||
"file_path": file_path,
|
|
||||||
"tasks": results
|
|
||||||
}
|
|
||||||
else:
|
else:
|
||||||
error_msg = f"Failed to list rclone remotes: {result.stderr}"
|
error_msg = f"Failed to list rclone remotes: {result.stderr}"
|
||||||
logger.error(error_msg)
|
logger.error(error_msg)
|
||||||
raise Exception(error_msg)
|
raise RuntimeError(error_msg)
|
||||||
|
|
||||||
except Exception as e:
|
except (subprocess.SubprocessError, OSError) as e:
|
||||||
error_msg = f"Error setting up rclone uploads for {filename}: {str(e)}"
|
error_msg = f"Error setting up rclone uploads for {filename}: {str(e)}"
|
||||||
logger.error(error_msg)
|
logger.error(error_msg)
|
||||||
raise Exception(error_msg)
|
raise RuntimeError(error_msg) from e
|
||||||
|
|||||||
@@ -0,0 +1,103 @@
|
|||||||
|
"""
|
||||||
|
OAuth helper utilities for token exchange operations.
|
||||||
|
Shared across multiple OAuth providers to reduce code duplication.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import logging
|
||||||
|
from typing import Dict, Any, Optional
|
||||||
|
import requests
|
||||||
|
from fastapi import HTTPException, status
|
||||||
|
|
||||||
|
from app.config import settings
|
||||||
|
|
||||||
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
|
def exchange_oauth_token(
|
||||||
|
provider_name: str, token_url: str, payload: Dict[str, str], timeout: Optional[int] = None
|
||||||
|
) -> Dict[str, Any]:
|
||||||
|
"""
|
||||||
|
Exchange an authorization code for tokens from an OAuth provider.
|
||||||
|
|
||||||
|
This function handles the common OAuth token exchange flow across multiple providers
|
||||||
|
(OneDrive, Google Drive, Dropbox) with proper error handling and secure logging.
|
||||||
|
|
||||||
|
Args:
|
||||||
|
provider_name: Name of the OAuth provider (for logging)
|
||||||
|
token_url: OAuth token endpoint URL
|
||||||
|
payload: Request payload containing client credentials and auth code
|
||||||
|
timeout: Request timeout in seconds (defaults to settings.http_request_timeout)
|
||||||
|
|
||||||
|
Returns:
|
||||||
|
Dict containing the token response from the provider
|
||||||
|
|
||||||
|
Raises:
|
||||||
|
HTTPException: If token exchange fails or response is invalid
|
||||||
|
"""
|
||||||
|
if timeout is None:
|
||||||
|
timeout = settings.http_request_timeout
|
||||||
|
|
||||||
|
try:
|
||||||
|
logger.info(f"Starting {provider_name} token exchange process")
|
||||||
|
|
||||||
|
# SECURITY: Never log sensitive data - only log non-sensitive metadata
|
||||||
|
safe_info = {
|
||||||
|
"provider": provider_name,
|
||||||
|
"token_url": token_url,
|
||||||
|
"grant_type": payload.get("grant_type", "unknown"),
|
||||||
|
}
|
||||||
|
logger.info(f"Token exchange request: {safe_info}")
|
||||||
|
|
||||||
|
# Make the token request
|
||||||
|
logger.info(f"Sending POST request to {provider_name} for token exchange")
|
||||||
|
response = requests.post(token_url, data=payload, timeout=timeout)
|
||||||
|
|
||||||
|
# Check if the request was successful
|
||||||
|
logger.info(f"Token exchange response status: {response.status_code}")
|
||||||
|
|
||||||
|
if response.status_code != 200:
|
||||||
|
# Log the error response for debugging (without sensitive data)
|
||||||
|
try:
|
||||||
|
error_json = response.json()
|
||||||
|
# Extract only error type, not full details which may contain sensitive info
|
||||||
|
error_type = error_json.get("error", "unknown_error")
|
||||||
|
logger.error(f"Token exchange failed with status {response.status_code}: {error_type}")
|
||||||
|
error_detail = {"error": error_type, "error_description": error_json.get("error_description", "")}
|
||||||
|
except (ValueError, requests.exceptions.JSONDecodeError) as json_err:
|
||||||
|
logger.error(f"Failed to parse error response as JSON: {str(json_err)}")
|
||||||
|
error_detail = {"error": "Unknown error", "status_code": response.status_code}
|
||||||
|
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=status.HTTP_400_BAD_REQUEST, detail=f"Token exchange failed: {error_detail}"
|
||||||
|
)
|
||||||
|
|
||||||
|
# Parse the token response
|
||||||
|
token_data = response.json()
|
||||||
|
|
||||||
|
# Validate the token response
|
||||||
|
if "refresh_token" not in token_data:
|
||||||
|
logger.error(f"{provider_name} returned success but no refresh_token found in response")
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=status.HTTP_502_BAD_GATEWAY,
|
||||||
|
detail=f"{provider_name} OAuth server returned success but no refresh token was included",
|
||||||
|
)
|
||||||
|
|
||||||
|
# Log success with non-sensitive metadata only
|
||||||
|
logger.info(f"Successfully exchanged authorization code for {provider_name} tokens")
|
||||||
|
|
||||||
|
return token_data
|
||||||
|
|
||||||
|
except HTTPException:
|
||||||
|
# Re-raise HTTP exceptions as they already have appropriate status codes
|
||||||
|
raise
|
||||||
|
except requests.exceptions.RequestException as e:
|
||||||
|
logger.exception(f"Network error during {provider_name} token exchange: {str(e)}")
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
|
||||||
|
detail=f"Failed to connect to {provider_name} OAuth service: {str(e)}",
|
||||||
|
)
|
||||||
|
except Exception as e:
|
||||||
|
logger.exception(f"Unexpected error during {provider_name} token exchange: {str(e)}")
|
||||||
|
raise HTTPException(
|
||||||
|
status_code=status.HTTP_500_INTERNAL_SERVER_ERROR, detail=f"Failed to exchange token: {str(e)}"
|
||||||
|
)
|
||||||
Reference in New Issue
Block a user