✨(backend) use a celery task for eml import

Unify the import implementations to use celery task (eml, mbox, imap)
Same returns for all import tasks.
Update openapi schema.
This commit is contained in:
Sabrina Demagny
2025-06-11 23:53:48 +02:00
parent 55c3f2a081
commit 761784e91c
7 changed files with 263 additions and 153 deletions
+5 -37
View File
@@ -1218,7 +1218,7 @@
"/api/v1.0/import/file/": {
"post": {
"operationId": "import_file_create",
"description": "\n Import messages by uploading an EML or MBOX file.\n \n - For EML files: Import is processed synchronously and returns immediately\n - For MBOX files: Import is processed asynchronously and returns a task ID\n \n The file must be a valid EML or MBOX format. The recipient mailbox must exist\n and the user must have access to it.\n ",
"description": "\n Import messages by uploading an EML or MBOX file.\n \n The import is processed asynchronously and returns a task ID for tracking.\n The file must be a valid EML or MBOX format. The recipient mailbox must exist\n and the user must have access to it.\n ",
"parameters": [
{
"in": "query",
@@ -1259,30 +1259,6 @@
}
],
"responses": {
"200": {
"content": {
"application/json": {
"schema": {
"type": "object",
"properties": {
"type": {
"type": "string",
"description": "Type of import (eml)"
},
"status": {
"type": "string",
"description": "Status of the import"
},
"message": {
"type": "string",
"description": "Success message"
}
}
}
}
},
"description": "EML file import completed successfully"
},
"202": {
"content": {
"application/json": {
@@ -1291,21 +1267,17 @@
"properties": {
"task_id": {
"type": "string",
"description": "Celery task ID for tracking import progress"
"description": "Task ID for tracking the import"
},
"type": {
"type": "string",
"description": "Type of import (mbox)"
},
"status": {
"type": "string",
"description": "Initial status of the import task"
"description": "Type of import (eml or mbox)"
}
}
}
}
},
"description": "MBOX file import started. Returns Celery task ID for tracking."
"description": "Import started. Returns Celery task ID for tracking."
},
"400": {
"description": "Invalid input data or file format"
@@ -1434,15 +1406,11 @@
"properties": {
"task_id": {
"type": "string",
"description": "Celery task ID for tracking import progress"
"description": "Task ID for tracking the import"
},
"type": {
"type": "string",
"description": "Type of import (imap)"
},
"status": {
"type": "string",
"description": "Initial status of the import task"
}
}
}
+19 -24
View File
@@ -33,23 +33,19 @@ class ImportViewSet(viewsets.ViewSet):
@extend_schema(
request=ImportFileSerializer,
responses={
200: OpenApiResponse(
description="EML file import completed successfully",
response={
"type": "object",
"properties": {
"task_id": {"type": "string", "description": "Task ID for tracking the import"},
"type": {"type": "string", "description": "Type of import (eml)"},
},
},
),
202: OpenApiResponse(
description="MBOX file import started. Returns Celery task ID for tracking.",
description="Import started. Returns Celery task ID for tracking.",
response={
"type": "object",
"properties": {
"task_id": {"type": "string", "description": "Task ID for tracking the import"},
"type": {"type": "string", "description": "Type of import (mbox)"},
"task_id": {
"type": "string",
"description": "Task ID for tracking the import",
},
"type": {
"type": "string",
"description": "Type of import (eml or mbox)",
},
},
},
),
@@ -62,9 +58,7 @@ class ImportViewSet(viewsets.ViewSet):
description="""
Import messages by uploading an EML or MBOX file.
- For EML files: Import is processed synchronously and returns immediately
- For MBOX files: Import is processed asynchronously and returns a task ID
The import is processed asynchronously and returns a task ID for tracking.
The file must be a valid EML or MBOX format. The recipient mailbox must exist
and the user must have access to it.
""",
@@ -103,12 +97,7 @@ class ImportViewSet(viewsets.ViewSet):
if not success:
return Response(response_data, status=status.HTTP_403_FORBIDDEN)
return Response(
response_data,
status=status.HTTP_202_ACCEPTED
if response_data["type"] == "mbox"
else status.HTTP_200_OK,
)
return Response(response_data, status=status.HTTP_202_ACCEPTED)
@extend_schema(
request=ImportIMAPSerializer,
@@ -118,8 +107,14 @@ class ImportViewSet(viewsets.ViewSet):
response={
"type": "object",
"properties": {
"task_id": {"type": "string", "description": "Task ID for tracking the import"},
"type": {"type": "string", "description": "Type of import (imap)"},
"task_id": {
"type": "string",
"description": "Task ID for tracking the import",
},
"type": {
"type": "string",
"description": "Type of import (imap)",
},
},
},
),
+19 -22
View File
@@ -7,10 +7,12 @@ from django.contrib import messages
from django.core.files.uploadedfile import UploadedFile
from django.http import HttpRequest
from core.mda.inbound import deliver_inbound_message
from core.mda.rfc5322 import parse_email_message
from core.models import Mailbox
from core.tasks import import_imap_messages_task, process_mbox_file_task
from core.tasks import (
import_imap_messages_task,
process_eml_file_task,
process_mbox_file_task,
)
logger = logging.getLogger(__name__)
@@ -54,26 +56,21 @@ class ImportService:
"This may take a while. You can check the status in the Celery task monitor.",
)
return True, response_data
else:
# Process EML file synchronously
parsed_email = parse_email_message(file_content)
success = deliver_inbound_message(
str(recipient), parsed_email, file_content, is_import=True
)
response_data = {"success": success, "type": "eml"}
elif file.name.endswith(".eml"):
# Process EML file asynchronously
task = process_eml_file_task.delay(file_content, str(recipient.id))
response_data = {"task_id": task.id, "type": "eml"}
if request:
if success:
messages.success(
request,
f"Successfully processed EML file: {file.name} for recipient {recipient}",
)
else:
messages.error(
request,
f"Failed to process EML file: {file.name} for recipient {recipient}",
)
return success, response_data
messages.info(
request,
f"Started processing EML file: {file.name} for recipient {recipient}. "
"This may take a while. You can check the status in the Celery task monitor.",
)
return True, response_data
else:
return False, {
"detail": "Invalid file format. Only EML and MBOX files are supported."
}
except Exception as e:
logger.exception("Error processing file: %s", e)
if request:
+76 -1
View File
@@ -335,7 +335,13 @@ def process_mbox_file_task(
)
failure_count += 1
return success_count, failure_count
return {
"status": "completed",
"total_messages": len(messages),
"success_count": success_count,
"failure_count": failure_count,
"type": "mbox",
}
def split_mbox_file(content: bytes) -> List[bytes]:
@@ -482,9 +488,78 @@ def import_imap_messages_task(
"total_messages": total_messages,
"success_count": success_count,
"failure_count": failure_count,
"type": "imap",
}
except Exception as e:
logger.exception("Error in import_imap_messages_task: %s", e)
self.update_state(state="FAILURE", meta={"status": "failed", "error": str(e)})
raise
@celery_app.task(bind=True)
def process_eml_file_task(
self, file_content: bytes, recipient_id: str
) -> Dict[str, Any]:
"""
Process an EML file asynchronously.
Args:
file_content: The content of the EML file
recipient_id: The UUID of the recipient mailbox
Returns:
Dictionary with import statistics
"""
try:
recipient = Mailbox.objects.get(id=recipient_id)
except Mailbox.DoesNotExist:
logger.error("Recipient mailbox %s not found", recipient_id)
return {
"status": "failed",
"total_messages": 0,
"success_count": 0,
"failure_count": 0,
"type": "eml",
"error": "Recipient mailbox not found",
}
try:
# Parse the email message
parsed_email = parse_email_message(file_content)
# Deliver the message
success = deliver_inbound_message(
str(recipient), parsed_email, file_content, is_import=True
)
if success:
return {
"status": "completed",
"total_messages": 1,
"success_count": 1,
"failure_count": 0,
"type": "eml",
}
return {
"status": "failed",
"total_messages": 1,
"success_count": 0,
"failure_count": 1,
"type": "eml",
"error": "Failed to deliver message",
}
except Exception as e:
logger.exception(
"Error processing EML file for recipient %s: %s",
recipient_id,
e,
)
self.update_state(state="FAILURE", meta={"status": "failed", "error": str(e)})
return {
"status": "failed",
"total_messages": 1,
"success_count": 0,
"failure_count": 1,
"type": "eml",
"error": str(e),
}
@@ -69,8 +69,7 @@ def test_import_eml_file(api_client, user, mailbox, eml_file_path):
{"import_file": f, "recipient": str(mailbox.id)},
format="multipart",
)
assert response.status_code == 200
assert response.data["success"] is True
assert response.status_code == 202
assert response.data["type"] == "eml"
assert Message.objects.count() == 1
message = Message.objects.first()
@@ -2,6 +2,7 @@
# pylint: disable=redefined-outer-name, unused-argument, no-value-for-parameter
import datetime
from unittest.mock import patch
from django.core.files.uploadedfile import SimpleUploadedFile
from django.urls import reverse
@@ -10,7 +11,7 @@ import pytest
from core import factories
from core.models import Mailbox, MailDomain, Message, Thread
from core.tasks import process_mbox_file_task
from core.tasks import process_eml_file_task, process_mbox_file_task
@pytest.fixture
@@ -82,42 +83,62 @@ def test_import_eml_file(admin_client, eml_file, mailbox):
"""Test submitting the import form with a valid EML file."""
url = reverse("admin:core_message_import_messages")
# Create a test EML file
eml_file = SimpleUploadedFile("test.eml", eml_file, content_type="message/rfc822")
# Submit the form
response = admin_client.post(
url, {"import_file": eml_file, "recipient": mailbox.id}, follow=True
# Create a SimpleUploadedFile from the bytes content
test_file = SimpleUploadedFile(
"test.eml",
eml_file, # eml_file is already bytes
content_type="message/rfc822",
)
# Check response
assert response.status_code == 200
assert (
f"Successfully processed EML file: test.eml for recipient {mailbox}"
in response.content.decode()
)
# check that the message was created
assert Message.objects.count() == 1
message = Message.objects.first()
assert message.subject == "Mon mail avec joli pj"
assert message.attachments.count() == 1
assert message.sender.email == "sender@example.com"
assert message.recipients.get().contact.email == "recipient@example.com"
assert message.sent_at == message.thread.messaged_at
assert message.sent_at == (
datetime.datetime(2025, 5, 26, 20, 13, 44, tzinfo=datetime.timezone.utc)
)
with patch("core.tasks.process_eml_file_task.delay") as mock_task:
mock_task.return_value.id = "fake-task-id"
# Submit the form
response = admin_client.post(
url, {"import_file": test_file, "recipient": mailbox.id}, follow=True
)
# Check response
assert response.status_code == 200
assert (
f"Started processing EML file: test.eml for recipient {mailbox}"
in response.content.decode()
)
mock_task.assert_called_once()
# Run the task synchronously for testing
result = process_eml_file_task(
file_content=eml_file, recipient_id=str(mailbox.id)
)
assert result["status"] == "completed"
assert result["type"] == "eml"
assert result["total_messages"] == 1
assert result["success_count"] == 1
assert result["failure_count"] == 0
# check that the message was created
assert Message.objects.count() == 1
message = Message.objects.first()
assert message.subject == "Mon mail avec joli pj"
assert message.attachments.count() == 1
assert message.sender.email == "sender@example.com"
assert message.recipients.get().contact.email == "recipient@example.com"
assert message.sent_at == message.thread.messaged_at
assert message.sent_at == (
datetime.datetime(2025, 5, 26, 20, 13, 44, tzinfo=datetime.timezone.utc)
)
@pytest.mark.django_db
def test_process_mbox_file_task(mailbox, mbox_file):
"""Test the Celery task that processes MBOX files."""
# Run the task synchronously for testing
success_count, failure_count = process_mbox_file_task(
result = process_mbox_file_task(
file_content=mbox_file, recipient_id=str(mailbox.id)
)
assert success_count == 3 # Three messages in the test MBOX file
assert failure_count == 0
assert result["status"] == "completed"
assert result["type"] == "mbox"
assert result["total_messages"] == 3 # Three messages in the test MBOX file
assert result["success_count"] == 3
assert result["failure_count"] == 0
# Verify messages were created
assert Message.objects.count() == 3
@@ -1,8 +1,9 @@
"""Tests for the ImportService class."""
import datetime
from unittest.mock import MagicMock, patch
from unittest.mock import patch
from django.contrib.messages.storage.fallback import FallbackStorage
from django.core.files.uploadedfile import SimpleUploadedFile
from django.http import HttpRequest
@@ -44,6 +45,18 @@ def mailbox(domain):
return Mailbox.objects.create(local_part="test", domain=domain)
@pytest.fixture
def mock_request():
"""Create a mock request object with messages framework support."""
request = HttpRequest()
request.user = None
# Set up messages framework
request.session = "session"
messages = FallbackStorage(request)
request._messages = messages
return request
@pytest.fixture
def eml_file():
"""Get test eml file from test data."""
@@ -60,27 +73,37 @@ def mbox_file():
)
@pytest.fixture
def mock_request():
"""Create a mock request object."""
request = MagicMock(spec=HttpRequest)
request._messages = MagicMock()
return request
@pytest.mark.django_db
def test_import_file_eml_by_superuser(admin_user, mailbox, eml_file, mock_request):
"""Test successful EML file import for superuser."""
success, response_data = ImportService.import_file(
file=eml_file,
recipient=mailbox,
user=admin_user,
request=mock_request,
)
with patch("core.tasks.process_eml_file_task.delay") as mock_task:
mock_task.return_value.id = "fake-task-id"
success, response_data = ImportService.import_file(
file=eml_file,
recipient=mailbox,
user=admin_user,
request=mock_request,
)
assert success is True
assert response_data["type"] == "eml"
assert response_data["success"] is True
assert success is True
assert response_data["type"] == "eml"
assert response_data["task_id"] == "fake-task-id"
mock_task.assert_called_once()
@pytest.mark.django_db
def test_import_file_eml_by_superuser_sync(admin_user, mailbox, eml_file):
# Run the task synchronously for testing
from core.tasks import process_eml_file_task
result = process_eml_file_task(
file_content=eml_file.read(), recipient_id=str(mailbox.id)
)
assert result["status"] == "completed"
assert result["type"] == "eml"
assert result["total_messages"] == 1
assert result["success_count"] == 1
assert result["failure_count"] == 0
assert Message.objects.count() == 1
message = Message.objects.first()
@@ -95,21 +118,41 @@ def test_import_file_eml_by_superuser(admin_user, mailbox, eml_file, mock_reques
@pytest.mark.django_db
def test_import_file_eml_by_user_with_access(user, mailbox, eml_file, mock_request):
def test_import_file_eml_by_user_with_access_task(
user, mailbox, eml_file, mock_request
):
"""Test successful EML file import by user with access on mailbox."""
# Add access to mailbox
mailbox.accesses.create(user=user, role=MailboxRoleChoices.ADMIN)
success, response_data = ImportService.import_file(
file=eml_file,
recipient=mailbox,
user=user,
request=mock_request,
)
with patch("core.tasks.process_eml_file_task.delay") as mock_task:
mock_task.return_value.id = "fake-task-id"
success, response_data = ImportService.import_file(
file=eml_file,
recipient=mailbox,
user=user,
request=mock_request,
)
assert success is True
assert response_data["type"] == "eml"
assert response_data["success"] is True
assert success is True
assert response_data["type"] == "eml"
assert response_data["task_id"] == "fake-task-id"
mock_task.assert_called_once()
@pytest.mark.django_db
def test_import_file_eml_by_user_with_access_sync(user, mailbox, eml_file):
# Run the task synchronously for testing
from core.tasks import process_eml_file_task
result = process_eml_file_task(
file_content=eml_file.read(), recipient_id=str(mailbox.id)
)
assert result["status"] == "completed"
assert result["type"] == "eml"
assert result["total_messages"] == 1
assert result["success_count"] == 1
assert result["failure_count"] == 0
assert Message.objects.count() == 1
message = Message.objects.first()
@@ -210,21 +253,32 @@ def test_import_file_no_access(user, domain, eml_file, mock_request):
assert Message.objects.count() == 0
@pytest.mark.django_db
def test_import_file_invalid_file(admin_user, mailbox, mock_request):
"""Test import with invalid file type."""
"""Test import with an invalid file."""
# Create an invalid file (not EML or MBOX)
invalid_file = SimpleUploadedFile(
"test.txt", b"Not an email file", content_type="text/plain"
"test.txt", b"Invalid file content", content_type="text/plain"
)
success, response_data = ImportService.import_file(
file=invalid_file,
recipient=mailbox,
user=admin_user,
request=mock_request,
)
with patch("core.tasks.process_eml_file_task.delay") as mock_task:
# The task should not be called for invalid files
mock_task.assert_not_called()
assert success is False
assert Message.objects.count() == 0
success, response_data = ImportService.import_file(
file=invalid_file,
recipient=mailbox,
user=admin_user,
request=mock_request,
)
assert success is False
assert "detail" in response_data
assert (
"Invalid file format. Only EML and MBOX files are supported."
in response_data["detail"]
)
assert Message.objects.count() == 0
def test_import_imap_by_superuser(admin_user, mailbox, mock_request):
@@ -250,10 +304,11 @@ def test_import_imap_by_superuser(admin_user, mailbox, mock_request):
mock_task.assert_called_once()
def test_import_imap_by_user_with_access(user, mailbox, mock_request):
@pytest.mark.parametrize("role", [MailboxRoleChoices.ADMIN, MailboxRoleChoices.EDITOR])
def test_import_imap_by_user_with_access(user, mailbox, mock_request, role):
"""Test successful IMAP import by user with access on mailbox."""
# Add access to mailbox
mailbox.accesses.create(user=user, role=MailboxRoleChoices.ADMIN)
mailbox.accesses.create(user=user, role=role)
with patch("core.tasks.import_imap_messages_task.delay") as mock_task:
mock_task.return_value.id = "fake-task-id"