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

1from __future__ import annotations 

2 

3import asyncio 

4from datetime import UTC, datetime, timedelta 

5from typing import TYPE_CHECKING, Any, Protocol, TypeVar 

6import uuid 

7 

8from lexigram.di.decorators import inject 

9from lexigram.logging import get_logger 

10from lexigram.serialization import dumps_str 

11 

12if TYPE_CHECKING: 

13 from collections.abc import AsyncIterator, Callable 

14 

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 

22 

23T = TypeVar("T") 

24 

25 

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 

36 

37 

38class IExportDataSource(Protocol[T]): # type: ignore[misc] 

39 """Protocol for data sources that support advanced exports.""" 

40 

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]]: ... 

50 

51 async def get_export_count(self, filters: dict[str, Any]) -> int: ... 

52 

53 async def get_column_definitions(self) -> list[ExportColumn]: ... 

54 

55 

56class IExportBackend(Protocol): 

57 """Protocol for export format backends.""" 

58 

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: ... 

66 

67 

68class JsonExportBackend(IExportBackend): 

69 """Backend for JSON exports using orjson.""" 

70 

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}" 

82 

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 

90 

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 

96 

97 

98logger = get_logger(__name__) 

99 

100 

101@inject 

102class ExportService: 

103 """Advanced export service with modular backends and job management.""" 

104 

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 

118 

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 

126 

127 self._templates: dict[str, ExportTemplate] = {} 

128 self._jobs: dict[str, ExportJob] = {} 

129 self._background_tasks: dict[str, asyncio.Task] = {} 

130 

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 } 

138 

139 # ------------------------------------------------------------------------- 

140 # Template Management 

141 # ------------------------------------------------------------------------- 

142 

143 def register_template(self, template: ExportTemplate) -> None: 

144 """Register an export template.""" 

145 self._templates[template.name] = template 

146 

147 def get_template(self, name: str) -> ExportTemplate | None: 

148 """Get a registered template.""" 

149 return self._templates.get(name) 

150 

151 def list_templates(self) -> list[ExportTemplate]: 

152 """List all registered templates.""" 

153 return list(self._templates.values()) 

154 

155 # ------------------------------------------------------------------------- 

156 # Export JobProtocol Management 

157 # ------------------------------------------------------------------------- 

158 

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()) 

173 

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 ) 

186 

187 self._jobs[job_id] = job 

188 return job_id 

189 

190 def get_job(self, job_id: str) -> ExportJob | None: 

191 """Get export job by ID.""" 

192 return self._jobs.get(job_id) 

193 

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()) 

202 

203 if user_id is not None: 

204 jobs = list(filter(lambda j: j.user_id == user_id, jobs)) 

205 

206 if status is not None: 

207 jobs = list(filter(lambda j: j.status == status, jobs)) 

208 

209 # Sort by creation time, newest first 

210 jobs.sort(key=lambda j: j.created_at, reverse=True) 

211 

212 return jobs[:limit] 

213 

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 

226 

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 ) 

243 

244 # ------------------------------------------------------------------------- 

245 # Export Execution 

246 # ------------------------------------------------------------------------- 

247 

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")) 

258 

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)) 

272 

273 try: 

274 # Update job status 

275 job.status = ExportStatus.PROCESSING 

276 job.started_at = datetime.now(UTC) 

277 

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 ) 

285 

286 # Get total count 

287 job.total_records = await data_source.get_export_count(job.filters) 

288 

289 # Get data in chunks 

290 all_data = [] 

291 offset = 0 

292 chunk_size = 1000 

293 

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 ) 

301 

302 if not chunk: 

303 break 

304 

305 all_data.extend(chunk) 

306 offset += len(chunk) 

307 job.processed_records = len(all_data) 

308 

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 

316 

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 ) 

329 

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 ) 

337 

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) 

344 

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 ) 

353 

354 return Ok(job) 

355 

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) 

366 

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 ) 

375 

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) 

383 

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 ) 

392 

393 return Err(AdminError(message=job.error_message or "Export failed")) 

394 

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 

403 

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}" 

408 

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}") 

424 

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. 

427 

428 total = await data_source.get_export_count(filters or {}) 

429 offset = 0 

430 

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 

440 

441 # Yield a chunk (this is a simplification) 

442 yield b"encoded batch chunk" 

443 offset += len(batch) 

444 

445 # ------------------------------------------------------------------------- 

446 # Background Processing 

447 # ------------------------------------------------------------------------- 

448 

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 

459 

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) 

483 

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 

495 

496 # ------------------------------------------------------------------------- 

497 # Scheduled Exports 

498 # ------------------------------------------------------------------------- 

499 

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 ) 

517 

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 

522 

523 return job_id 

524 

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 

540 

541 # ------------------------------------------------------------------------- 

542 # Cleanup 

543 # ------------------------------------------------------------------------- 

544 

545 async def cleanup_old_files(self) -> int: 

546 """Clean up old export files.""" 

547 # Implementation depends on storage provider 

548 return 0 

549 

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 = [] 

554 

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) 

572 

573 for job_id in jobs_to_remove: 

574 del self._jobs[job_id] 

575 self._background_tasks.pop(job_id, None) 

576 

577 return len(jobs_to_remove)