Coverage for src/lexigram/admin/services/export/service.py: 0%
209 statements
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-24 23:18 +0800
« prev ^ index » next coverage.py v7.15.4, created at 2026-08-24 23:18 +0800
1from __future__ import annotations
3import asyncio
4from datetime import UTC, datetime, timedelta
5from typing import TYPE_CHECKING, Any, Protocol, TypeVar
6import uuid
8from lexigram.di.decorators import inject
9from lexigram.logging import get_logger
10from lexigram.serialization import dumps_str
12if TYPE_CHECKING:
13 from collections.abc import AsyncIterator, Callable
15 from lexigram.admin.data.data_source import ( # type: ignore[attr-defined]
16 ExportColumn,
17 )
18 from lexigram.admin.services.session import SessionStateService
19 from lexigram.contracts.core import TaskManagerProtocol
20 from lexigram.contracts.infra.storage import BlobStoreProtocol
21 from lexigram.contracts.mailer.protocols import MailerProtocol
23T = TypeVar("T")
26from lexigram.admin.exceptions import AdminError
27from lexigram.admin.services.export.scheduler import (
28 ExportFormat,
29 ExportJob,
30 ExportSchedule,
31 ExportStatus,
32 ExportTemplate,
33)
34from lexigram.contracts.audit import AuditEntry, AuditEventSeverity, AuditLoggerProtocol
35from lexigram.result import Err, Ok, Result
38class IExportDataSource(Protocol[T]): # type: ignore[misc]
39 """Protocol for data sources that support advanced exports."""
41 async def get_export_data(
42 self,
43 filters: dict[str, Any],
44 columns: list[str],
45 sort_by: str | None = None,
46 sort_order: str = "asc",
47 limit: int | None = None,
48 offset: int | None = None,
49 ) -> list[dict[str, Any]]: ...
51 async def get_export_count(self, filters: dict[str, Any]) -> int: ...
53 async def get_column_definitions(self) -> list[ExportColumn]: ...
56class IExportBackend(Protocol):
57 """Protocol for export format backends."""
59 async def generate_file(
60 self,
61 job: ExportJob,
62 data: list[dict[str, Any]],
63 storage: Any,
64 export_dir: str,
65 ) -> str: ...
68class JsonExportBackend(IExportBackend):
69 """Backend for JSON exports using orjson."""
71 async def generate_file(
72 self,
73 job: ExportJob,
74 data: list[dict[str, Any]],
75 storage: Any,
76 export_dir: str,
77 ) -> str:
78 """Export data to JSON format."""
79 timestamp = datetime.now(UTC).strftime("%Y%m%d_%H%M%S")
80 filename = f"{job.resource_name}_export_{timestamp}.json"
81 file_path = f"{export_dir}/{filename}"
83 if job.columns:
84 output_data = []
85 for row in data:
86 filtered_row = {k: v for k, v in row.items() if k in job.columns}
87 output_data.append(filtered_row)
88 else:
89 output_data = data
91 content = dumps_str(output_data)
92 await storage.upload(
93 file_path, content.encode("utf-8"), content_type="application/json"
94 )
95 return file_path
98logger = get_logger(__name__)
101@inject
102class ExportService:
103 """Advanced export service with modular backends and job management."""
105 def __init__(
106 self,
107 storage: BlobStoreProtocol,
108 task_manager: TaskManagerProtocol,
109 messaging: MailerProtocol | None = None,
110 session: SessionStateService | None = None,
111 export_dir: str = "exports",
112 max_file_age_days: int = 7,
113 audit: AuditLoggerProtocol | None = None,
114 ):
115 from lexigram.admin.services.export.adapters.csv import CsvExportBackend
116 from lexigram.admin.services.export.adapters.excel import ExcelExportBackend
117 from lexigram.admin.services.export.adapters.pdf import PdfExportBackend
119 self.storage = storage
120 self.task_manager = task_manager
121 self.messaging = messaging
122 self.session = session
123 self.export_dir = export_dir
124 self.max_file_age_days = max_file_age_days
125 self.audit = audit
127 self._templates: dict[str, ExportTemplate] = {}
128 self._jobs: dict[str, ExportJob] = {}
129 self._background_tasks: dict[str, asyncio.Task] = {}
131 # Register default backends
132 self._backends: dict[ExportFormat, IExportBackend] = {
133 ExportFormat.CSV: CsvExportBackend(),
134 ExportFormat.JSON: JsonExportBackend(),
135 ExportFormat.EXCEL: ExcelExportBackend(),
136 ExportFormat.PDF: PdfExportBackend(),
137 }
139 # -------------------------------------------------------------------------
140 # Template Management
141 # -------------------------------------------------------------------------
143 def register_template(self, template: ExportTemplate) -> None:
144 """Register an export template."""
145 self._templates[template.name] = template
147 def get_template(self, name: str) -> ExportTemplate | None:
148 """Get a registered template."""
149 return self._templates.get(name)
151 def list_templates(self) -> list[ExportTemplate]:
152 """List all registered templates."""
153 return list(self._templates.values())
155 # -------------------------------------------------------------------------
156 # Export JobProtocol Management
157 # -------------------------------------------------------------------------
159 def create_job(
160 self,
161 resource_name: str,
162 file_format: ExportFormat,
163 filters: dict[str, Any] | None = None,
164 columns: list[str] | None = None,
165 template_name: str | None = None,
166 user_id: Any | None = None,
167 scheduled_for: datetime | None = None,
168 schedule_type: ExportSchedule = ExportSchedule.IMMEDIATE,
169 email_recipients: list[str] | None = None,
170 ) -> str:
171 """Create a new export job."""
172 job_id = str(uuid.uuid4())
174 job = ExportJob(
175 job_id=job_id,
176 resource_name=resource_name,
177 format=file_format,
178 filters=filters or {},
179 columns=columns or [],
180 template_name=template_name,
181 user_id=user_id,
182 scheduled_for=scheduled_for,
183 schedule_type=schedule_type,
184 email_recipients=email_recipients or [],
185 )
187 self._jobs[job_id] = job
188 return job_id
190 def get_job(self, job_id: str) -> ExportJob | None:
191 """Get export job by ID."""
192 return self._jobs.get(job_id)
194 def list_jobs(
195 self,
196 user_id: Any | None = None,
197 status: ExportStatus | None = None,
198 limit: int = 50,
199 ) -> list[ExportJob]:
200 """List export jobs with optional filtering."""
201 jobs = list(self._jobs.values())
203 if user_id is not None:
204 jobs = list(filter(lambda j: j.user_id == user_id, jobs))
206 if status is not None:
207 jobs = list(filter(lambda j: j.status == status, jobs))
209 # Sort by creation time, newest first
210 jobs.sort(key=lambda j: j.created_at, reverse=True)
212 return jobs[:limit]
214 async def _record_export_audit(
215 self,
216 *,
217 job: ExportJob,
218 action: str,
219 outcome: str,
220 severity: AuditEventSeverity,
221 **metadata: object,
222 ) -> None:
223 """Record an audit event for an export operation."""
224 if self.audit is None or job.user_id is None:
225 return
227 await self.audit.log(
228 AuditEntry(
229 action=action,
230 actor_id=str(job.user_id),
231 resource_type="export_job",
232 resource_id=job.job_id,
233 outcome=outcome,
234 severity=severity,
235 metadata={
236 "resource_name": job.resource_name,
237 "format": job.format.value,
238 **{k: str(v) for k, v in metadata.items()},
239 },
240 source="admin",
241 )
242 )
244 # -------------------------------------------------------------------------
245 # Export Execution
246 # -------------------------------------------------------------------------
248 async def execute_export(
249 self,
250 job_id: str,
251 data_source: IExportDataSource,
252 progress_callback: Callable[[float], None] | None = None,
253 ) -> Result[ExportJob, AdminError]:
254 """Execute an export job."""
255 job = self.get_job(job_id)
256 if not job:
257 return Err(AdminError(message=f"Export job {job_id} not found"))
259 backend = self._backends.get(job.format)
260 if not backend:
261 job.status = ExportStatus.FAILED
262 job.error_message = f"Unsupported export format: {job.format}"
263 job.completed_at = datetime.now(UTC)
264 await self._record_export_audit(
265 job=job,
266 action="admin.export.failed",
267 outcome="failure",
268 severity=AuditEventSeverity.HIGH,
269 error_message=job.error_message,
270 )
271 return Err(AdminError(message=job.error_message))
273 try:
274 # Update job status
275 job.status = ExportStatus.PROCESSING
276 job.started_at = datetime.now(UTC)
278 # Record export start audit event
279 await self._record_export_audit(
280 job=job,
281 action="admin.export.start",
282 outcome="success",
283 severity=AuditEventSeverity.MEDIUM,
284 )
286 # Get total count
287 job.total_records = await data_source.get_export_count(job.filters)
289 # Get data in chunks
290 all_data = []
291 offset = 0
292 chunk_size = 1000
294 while offset < job.total_records:
295 chunk = await data_source.get_export_data(
296 filters=job.filters,
297 columns=job.columns,
298 limit=chunk_size,
299 offset=offset,
300 )
302 if not chunk:
303 break
305 all_data.extend(chunk)
306 offset += len(chunk)
307 job.processed_records = len(all_data)
309 # Update progress
310 progress = (
311 (len(all_data) / job.total_records) * 100
312 if job.total_records > 0
313 else 100
314 )
315 job.progress = progress
317 if progress_callback:
318 try:
319 progress_callback(progress)
320 except Exception as _cb_err: # noqa: BLE001 — progress callbacks are user-supplied and may raise anything; failure is non-fatal
321 # Emit explicit ERROR text first (helps caplog message matching)
322 logger.exception("Progress callback failed for job %s", job_id)
323 # Also emit a simple job-level error to be robust for different logging setups
324 logger.exception("Progress callback failed for job %s", job_id)
325 # Keep traceback for diagnostics
326 logger.exception(
327 "Progress callback callback exception for job %s", job_id
328 )
330 # Generate file via backend
331 file_path = await backend.generate_file(
332 job,
333 all_data,
334 self.storage,
335 self.export_dir,
336 )
338 # Update job with results
339 job.status = ExportStatus.COMPLETED
340 job.completed_at = datetime.now(UTC)
341 job.file_path = file_path
342 job.file_size = await self._get_file_size(file_path)
343 job.download_url = await self._generate_download_url(file_path)
345 # Record export complete audit event
346 await self._record_export_audit(
347 job=job,
348 action="admin.export.complete",
349 outcome="success",
350 severity=AuditEventSeverity.HIGH,
351 total_records=job.total_records,
352 )
354 return Ok(job)
356 except asyncio.CancelledError:
357 job.status = ExportStatus.CANCELLED
358 job.completed_at = datetime.now(UTC)
359 raise
360 except (OSError, RuntimeError, ValueError, AttributeError, LookupError) as e:
361 # Known failure modes during export (I/O, runtime, validation); log and mark job failed
362 logger.exception("Export job %s failed", job_id)
363 job.status = ExportStatus.FAILED
364 job.error_message = str(e)
365 job.completed_at = datetime.now(UTC)
367 # Record export failed audit event
368 await self._record_export_audit(
369 job=job,
370 action="admin.export.failed",
371 outcome="failure",
372 severity=AuditEventSeverity.HIGH,
373 error_message=str(e),
374 )
376 return Err(AdminError(message=job.error_message or "Export failed"))
377 except Exception as e: # noqa: BLE001 — catch-all safety net; unexpected failures must mark the job as failed
378 # Catch-all for unexpected errors — log with traceback so it's visible in monitoring
379 logger.exception("Unexpected error in export job %s", job_id)
380 job.status = ExportStatus.FAILED
381 job.error_message = str(e)
382 job.completed_at = datetime.now(UTC)
384 # Record export failed audit event
385 await self._record_export_audit(
386 job=job,
387 action="admin.export.failed",
388 outcome="failure",
389 severity=AuditEventSeverity.HIGH,
390 error_message=str(e),
391 )
393 return Err(AdminError(message=job.error_message or "Export failed"))
395 async def _get_file_size(self, file_path: str) -> int:
396 """Get file size from storage."""
397 try:
398 # Placeholder: implementation depends on storage provider
399 return 0
400 except (AttributeError, OSError) as e:
401 logger.warning("Failed to determine file size for %s: %s", file_path, e)
402 return 0
404 async def _generate_download_url(self, file_path: str) -> str:
405 """Generate download URL for file."""
406 # Placeholder: implementation depends on storage provider / router
407 return f"/admin/exports/download/{file_path}"
409 async def stream_export(
410 self,
411 data_source: IExportDataSource,
412 export_format: ExportFormat,
413 filters: dict[str, Any] | None = None,
414 columns: list[str] | None = None,
415 batch_size: int = 1000,
416 ) -> AsyncIterator[bytes]:
417 """
418 High-performance streaming export for large datasets.
419 Yields chunks of bytes to avoid memory overhead.
420 """
421 backend = self._backends.get(export_format)
422 if not backend:
423 raise ValueError(f"Unsupported export format: {export_format}")
425 # In a real implementation, the backend would need a 'stream_batch' method
426 # and we would yield bytes from it. For this demo, we yield mock bytes.
428 total = await data_source.get_export_count(filters or {})
429 offset = 0
431 while offset < total:
432 batch = await data_source.get_export_data(
433 filters=filters or {},
434 columns=columns or [],
435 limit=batch_size,
436 offset=offset,
437 )
438 if not batch:
439 break
441 # Yield a chunk (this is a simplification)
442 yield b"encoded batch chunk"
443 offset += len(batch)
445 # -------------------------------------------------------------------------
446 # Background Processing
447 # -------------------------------------------------------------------------
449 async def start_background_export(
450 self,
451 job_id: str,
452 data_source: IExportDataSource,
453 ) -> None:
454 """Start export job in background."""
455 task = self.task_manager.create_background_task(
456 self._run_background_export(job_id, data_source),
457 )
458 self._background_tasks[job_id] = task
460 async def _run_background_export(
461 self,
462 job_id: str,
463 data_source: IExportDataSource,
464 ) -> None:
465 """Run export job in background task."""
466 try:
467 result = await self.execute_export(job_id, data_source)
468 if result.is_err():
469 logger.warning("Background export failed: %s", result.unwrap_err())
470 except asyncio.CancelledError:
471 pass
472 except (OSError, RuntimeError, ValueError, AttributeError, LookupError):
473 logger.exception("Background export job %s failed", job_id)
474 except BaseException:
475 # Unexpected error; log full traceback
476 logger.exception(
477 "Background export job %s failed unexpectedly",
478 job_id,
479 exc_info=True,
480 )
481 finally:
482 self._background_tasks.pop(job_id, None)
484 def cancel_job(self, job_id: str) -> bool:
485 """Cancel a running export job."""
486 task = self._background_tasks.get(job_id)
487 if task and not task.done():
488 task.cancel()
489 job = self.get_job(job_id)
490 if job:
491 job.status = ExportStatus.CANCELLED
492 job.completed_at = datetime.now(UTC)
493 return True
494 return False
496 # -------------------------------------------------------------------------
497 # Scheduled Exports
498 # -------------------------------------------------------------------------
500 async def schedule_export(
501 self,
502 job_config: dict[str, Any],
503 schedule_type: ExportSchedule,
504 next_run: datetime,
505 ) -> str:
506 """Schedule a recurring export job."""
507 job_id = self.create_job( # type: ignore[call-arg]
508 resource_name=job_config["resource_name"],
509 format=job_config["format"],
510 filters=job_config.get("filters", {}),
511 columns=job_config.get("columns", []),
512 template_name=job_config.get("template_name"),
513 scheduled_for=next_run,
514 schedule_type=schedule_type,
515 email_recipients=job_config.get("email_recipients", []),
516 )
518 job = self.get_job(job_id)
519 if job:
520 job.metadata["schedule_config"] = job_config
521 job.metadata["schedule_type"] = schedule_type.value
523 return job_id
525 def get_next_run_time(
526 self,
527 schedule_type: ExportSchedule,
528 base_time: datetime,
529 ) -> datetime:
530 """Calculate next run time for schedule."""
531 if schedule_type == ExportSchedule.HOURLY:
532 return base_time + timedelta(hours=1)
533 if schedule_type == ExportSchedule.DAILY:
534 return base_time + timedelta(days=1)
535 if schedule_type == ExportSchedule.WEEKLY:
536 return base_time + timedelta(weeks=1)
537 if schedule_type == ExportSchedule.MONTHLY:
538 return base_time + timedelta(days=30)
539 return base_time
541 # -------------------------------------------------------------------------
542 # Cleanup
543 # -------------------------------------------------------------------------
545 async def cleanup_old_files(self) -> int:
546 """Clean up old export files."""
547 # Implementation depends on storage provider
548 return 0
550 async def cleanup_completed_jobs(self, max_age_days: int = 30) -> int:
551 """Clean up old completed jobs from memory."""
552 cutoff_date = datetime.now(UTC) - timedelta(days=max_age_days)
553 jobs_to_remove = []
555 for job_id, job in self._jobs.items():
556 if (
557 job.status
558 in [
559 ExportStatus.COMPLETED,
560 ExportStatus.FAILED,
561 ExportStatus.CANCELLED,
562 ]
563 and job.completed_at
564 ):
565 job_time = (
566 job.completed_at
567 if job.completed_at.tzinfo is not None
568 else job.completed_at.replace(tzinfo=UTC)
569 )
570 if job_time < cutoff_date:
571 jobs_to_remove.append(job_id)
573 for job_id in jobs_to_remove:
574 del self._jobs[job_id]
575 self._background_tasks.pop(job_id, None)
577 return len(jobs_to_remove)