Changes committed
This commit is contained in:
0
docengine/app/tasks/__init__.py
Normal file
0
docengine/app/tasks/__init__.py
Normal file
94
docengine/app/tasks/document_tasks.py
Normal file
94
docengine/app/tasks/document_tasks.py
Normal file
@@ -0,0 +1,94 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import uuid
|
||||
|
||||
from app.core.logging_config import get_logger
|
||||
from app.workers.celery_app import celery_app
|
||||
from app.core.database import get_db_context
|
||||
from app.services.document_service import DocumentProcessingService
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
@celery_app.task(
|
||||
name="app.tasks.document_tasks.process_document_task",
|
||||
bind=True,
|
||||
max_retries=3,
|
||||
default_retry_delay=60,
|
||||
acks_late=True,
|
||||
)
|
||||
def process_document_task(self, document_id: str) -> dict: # noqa: ANN001
|
||||
"""Celery task to process a document asynchronously."""
|
||||
logger.info("task_started", task_id=self.request.id, document_id=document_id)
|
||||
|
||||
try:
|
||||
doc_uuid = uuid.UUID(document_id)
|
||||
with get_db_context() as db:
|
||||
service = DocumentProcessingService(db)
|
||||
document = service.process_document(doc_uuid)
|
||||
logger.info(
|
||||
"task_completed",
|
||||
task_id=self.request.id,
|
||||
document_id=document_id,
|
||||
status=document.status,
|
||||
)
|
||||
return {
|
||||
"document_id": document_id,
|
||||
"status": document.status,
|
||||
"page_count": document.page_count,
|
||||
}
|
||||
except Exception as exc:
|
||||
logger.exception(
|
||||
"task_failed",
|
||||
task_id=self.request.id,
|
||||
document_id=document_id,
|
||||
error=str(exc),
|
||||
retry=self.request.retries,
|
||||
)
|
||||
raise self.retry(exc=exc)
|
||||
|
||||
|
||||
@celery_app.task(
|
||||
name="app.tasks.document_tasks.match_document_task",
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=30,
|
||||
)
|
||||
def match_document_task(
|
||||
self, # noqa: ANN001
|
||||
document_id: str,
|
||||
min_confidence: float = 0.5,
|
||||
max_results: int = 5,
|
||||
) -> dict:
|
||||
"""Celery task to match a document against templates."""
|
||||
logger.info("match_task_started", task_id=self.request.id, document_id=document_id)
|
||||
|
||||
try:
|
||||
doc_uuid = uuid.UUID(document_id)
|
||||
with get_db_context() as db:
|
||||
from app.services.matching_service import MatchingService
|
||||
service = MatchingService(db)
|
||||
matches = service.match_document(
|
||||
document_id=doc_uuid,
|
||||
min_confidence=min_confidence,
|
||||
max_results=max_results,
|
||||
)
|
||||
logger.info(
|
||||
"match_task_completed",
|
||||
task_id=self.request.id,
|
||||
document_id=document_id,
|
||||
matches=len(matches),
|
||||
)
|
||||
return {
|
||||
"document_id": document_id,
|
||||
"matches": len(matches),
|
||||
"best_score": matches[0].confidence_score if matches else 0.0,
|
||||
}
|
||||
except Exception as exc:
|
||||
logger.exception(
|
||||
"match_task_failed",
|
||||
task_id=self.request.id,
|
||||
document_id=document_id,
|
||||
error=str(exc),
|
||||
)
|
||||
raise self.retry(exc=exc)
|
||||
46
docengine/app/tasks/maintenance_tasks.py
Normal file
46
docengine/app/tasks/maintenance_tasks.py
Normal file
@@ -0,0 +1,46 @@
|
||||
from __future__ import annotations
|
||||
|
||||
from app.core.logging_config import get_logger
|
||||
from app.workers.celery_app import celery_app
|
||||
from app.core.database import get_db_context
|
||||
|
||||
logger = get_logger(__name__)
|
||||
|
||||
|
||||
@celery_app.task(
|
||||
name="app.tasks.maintenance_tasks.cleanup_expired_tokens",
|
||||
bind=True,
|
||||
)
|
||||
def cleanup_expired_tokens(self) -> dict: # noqa: ANN001
|
||||
"""Cleanup expired and revoked refresh tokens."""
|
||||
logger.info("cleanup_tokens_started", task_id=self.request.id)
|
||||
|
||||
try:
|
||||
with get_db_context() as db:
|
||||
from app.repositories.user_repository import RefreshTokenRepository
|
||||
repo = RefreshTokenRepository(db)
|
||||
count = repo.cleanup_expired_tokens()
|
||||
logger.info("cleanup_tokens_completed", removed=count)
|
||||
return {"removed_tokens": count}
|
||||
except Exception as exc:
|
||||
logger.exception("cleanup_tokens_failed", error=str(exc))
|
||||
return {"error": str(exc)}
|
||||
|
||||
|
||||
@celery_app.task(
|
||||
name="app.tasks.maintenance_tasks.cleanup_temp_storage",
|
||||
bind=True,
|
||||
)
|
||||
def cleanup_temp_storage(self) -> dict: # noqa: ANN001
|
||||
"""Cleanup temporary storage files."""
|
||||
logger.info("cleanup_temp_started", task_id=self.request.id)
|
||||
|
||||
try:
|
||||
from app.storage.provider import get_storage_provider
|
||||
storage = get_storage_provider()
|
||||
count = storage.cleanup_temp()
|
||||
logger.info("cleanup_temp_completed", removed=count)
|
||||
return {"removed_files": count}
|
||||
except Exception as exc:
|
||||
logger.exception("cleanup_temp_failed", error=str(exc))
|
||||
return {"error": str(exc)}
|
||||
Reference in New Issue
Block a user