Coverage for src/lexigram/admin/tasks/bulk_operations.py: 71%

106 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-21 14:56 +0800

1"""Background tasks for lexigram-admin. 

2 

3Provides task definitions for long-running operations like 

4bulk exports, imports, and scheduled maintenance. 

5 

6Integrates with lexigram-tasks for background job processing. 

7""" 

8 

9from __future__ import annotations 

10 

11from dataclasses import dataclass, field 

12from enum import StrEnum 

13from typing import Any 

14 

15from lexigram.logging import get_logger 

16 

17logger = get_logger(__name__) 

18 

19 

20class AdminTaskType(StrEnum): 

21 """Types of admin background tasks.""" 

22 

23 BULK_EXPORT = "admin.bulk_export" 

24 BULK_IMPORT = "admin.bulk_import" 

25 BULK_UPDATE = "admin.bulk_update" 

26 BULK_DELETE = "admin.bulk_delete" 

27 REPORT_GENERATION = "admin.report_generation" 

28 DATA_CLEANUP = "admin.data_cleanup" 

29 CACHE_WARM = "admin.cache_warm" 

30 INDEX_REBUILD = "admin.index_rebuild" 

31 

32 

33@dataclass 

34class TaskProgress: 

35 """Progress information for a task.""" 

36 

37 total: int = 0 

38 completed: int = 0 

39 failed: int = 0 

40 current_item: str | None = None 

41 message: str = "" 

42 

43 @property 

44 def percentage(self) -> int: 

45 if self.total == 0: 

46 return 0 

47 return int((self.completed / self.total) * 100) 

48 

49 def to_dict(self) -> dict[str, Any]: 

50 return { 

51 "total": self.total, 

52 "completed": self.completed, 

53 "failed": self.failed, 

54 "percentage": self.percentage, 

55 "current_item": self.current_item, 

56 "message": self.message, 

57 } 

58 

59 

60@dataclass 

61class AdminTaskResult: 

62 """Result of an admin task.""" 

63 

64 success: bool 

65 message: str = "" 

66 data: dict[str, Any] = field(default_factory=dict) 

67 errors: list[str] = field(default_factory=list) 

68 warnings: list[str] = field(default_factory=list) 

69 

70 def to_dict(self) -> dict[str, Any]: 

71 return { 

72 "success": self.success, 

73 "message": self.message, 

74 "data": self.data, 

75 "errors": self.errors, 

76 "warnings": self.warnings, 

77 } 

78 

79 

80# ========== Bulk Export Task ========== 

81 

82 

83async def bulk_export_task( 

84 resource_type: str, 

85 query: dict[str, Any], 

86 file_format: str, 

87 columns: list[str] | None = None, 

88 notify_email: str | None = None, 

89 user_id: str | None = None, 

90 export_dir: str | None = None, 

91) -> AdminTaskResult: 

92 """Export resources in background. 

93 

94 Args: 

95 resource_type: Type of resource to export 

96 query: Query parameters for filtering 

97 file_format: Export format (csv, json, xlsx) 

98 columns: Columns to include (None = all) 

99 notify_email: Email to notify when complete 

100 user_id: ID of user who initiated export 

101 export_dir: Directory to store export files (defaults to system temp) 

102 

103 Returns: 

104 AdminTaskResult with export file path 

105 """ 

106 from pathlib import Path 

107 import tempfile 

108 from uuid import uuid4 

109 

110 export_id = str(uuid4()) 

111 if export_dir: 

112 export_dir_path = Path(export_dir) 

113 else: 

114 # Use system temp directory for security 

115 temp_dir = Path(tempfile.gettempdir()) / "lexigram_exports" 

116 temp_dir.mkdir(parents=True, exist_ok=True) 

117 export_dir_path = temp_dir 

118 

119 export_dir_path.mkdir(parents=True, exist_ok=True) 

120 

121 try: 

122 # Get data source (would be injected in real implementation) 

123 # For now, return a placeholder result 

124 logger.info( 

125 "Starting export: resource=%s format=%s user=%s", 

126 resource_type, 

127 format, 

128 user_id, 

129 ) 

130 

131 # Placeholder: actual implementation would: 

132 # 1. Get data source from container 

133 # 2. Build query from params 

134 # 3. Stream rows to file 

135 # 4. Send notification 

136 

137 export_path = export_dir_path / f"{export_id}.{format}" 

138 

139 # Write placeholder file 

140 import aiofiles 

141 

142 async with aiofiles.open(export_path, "w") as f: 

143 await f.write(f"Export of {resource_type}\n") 

144 

145 logger.info("Export completed: %s", export_path) 

146 

147 return AdminTaskResult( 

148 success=True, 

149 message="Export completed successfully", 

150 data={ 

151 "export_id": export_id, 

152 "file_path": str(export_path), 

153 "format": format, 

154 }, 

155 ) 

156 

157 except (ValueError, ConnectionError, TimeoutError, OSError, KeyError) as e: 

158 logger.exception("Export failed") 

159 return AdminTaskResult( 

160 success=False, 

161 message=str(e), 

162 errors=[str(e)], 

163 ) 

164 

165 

166# ========== Bulk Import Task ========== 

167 

168 

169async def bulk_import_task( 

170 resource_type: str, 

171 file_path: str, 

172 file_format: str, 

173 mapping: dict[str, str] | None = None, 

174 on_duplicate: str = "skip", 

175 user_id: str | None = None, 

176) -> AdminTaskResult: 

177 """Import resources from file in background. 

178 

179 Args: 

180 resource_type: Type of resource to import 

181 file_path: Path to import file 

182 file_format: File format (csv, json, xlsx) 

183 mapping: Column mapping (file_col -> model_field) 

184 on_duplicate: How to handle duplicates (skip, update, error) 

185 user_id: ID of user who initiated import 

186 

187 Returns: 

188 AdminTaskResult with import statistics 

189 """ 

190 from pathlib import Path 

191 

192 try: 

193 logger.info( 

194 "Starting import: resource=%s file=%s user=%s", 

195 resource_type, 

196 file_path, 

197 user_id, 

198 ) 

199 

200 # Validate file exists 

201 path = Path(file_path) 

202 if not path.exists(): 

203 return AdminTaskResult( 

204 success=False, 

205 message=f"File not found: {file_path}", 

206 errors=["File not found"], 

207 ) 

208 

209 # Placeholder: actual implementation would: 

210 # 1. Parse file based on format 

211 # 2. Validate and transform rows 

212 # 3. Create/update records 

213 # 4. Track progress 

214 

215 created = 0 

216 updated = 0 

217 skipped = 0 

218 errors: list[str] = [] 

219 

220 logger.info( 

221 "Import completed: created=%d updated=%d skipped=%d", 

222 created, 

223 updated, 

224 skipped, 

225 ) 

226 

227 return AdminTaskResult( 

228 success=True, 

229 message=f"Import completed: {created} created, {updated} updated, {skipped} skipped", 

230 data={ 

231 "created": created, 

232 "updated": updated, 

233 "skipped": skipped, 

234 "error_count": len(errors), 

235 }, 

236 errors=errors, 

237 ) 

238 

239 except (ValueError, ConnectionError, TimeoutError, OSError, KeyError) as e: 

240 return AdminTaskResult( 

241 success=False, 

242 message=str(e), 

243 errors=[str(e)], 

244 ) 

245 

246 

247# ========== Bulk Update Task ========== 

248 

249 

250async def bulk_update_task( 

251 resource_type: str, 

252 resource_ids: list[Any], 

253 data: dict[str, Any], 

254 user_id: str | None = None, 

255) -> AdminTaskResult: 

256 """Bulk update resources in background. 

257 

258 Args: 

259 resource_type: Type of resource 

260 resource_ids: IDs to update 

261 data: Data to apply to all records 

262 user_id: ID of user who initiated 

263 

264 Returns: 

265 AdminTaskResult with update statistics 

266 """ 

267 try: 

268 logger.info( 

269 "Starting bulk update: resource=%s count=%d user=%s", 

270 resource_type, 

271 len(resource_ids), 

272 user_id, 

273 ) 

274 

275 updated = 0 

276 failed = 0 

277 errors: list[str] = [] 

278 

279 # Placeholder: actual implementation would update each record 

280 

281 return AdminTaskResult( 

282 success=True, 

283 message=f"Bulk update completed: {updated} updated, {failed} failed", 

284 data={ 

285 "updated": updated, 

286 "failed": failed, 

287 }, 

288 errors=errors, 

289 ) 

290 

291 except (ValueError, ConnectionError, TimeoutError, OSError, KeyError) as e: 

292 return AdminTaskResult( 

293 success=False, 

294 message=str(e), 

295 errors=[str(e)], 

296 ) 

297 

298 

299# ========== Bulk Delete Task ========== 

300 

301 

302async def bulk_delete_task( 

303 resource_type: str, 

304 resource_ids: list[Any], 

305 soft_delete: bool = True, 

306 user_id: str | None = None, 

307) -> AdminTaskResult: 

308 """Bulk delete resources in background. 

309 

310 Args: 

311 resource_type: Type of resource 

312 resource_ids: IDs to delete 

313 soft_delete: Whether to soft delete 

314 user_id: ID of user who initiated 

315 

316 Returns: 

317 AdminTaskResult with delete statistics 

318 """ 

319 try: 

320 logger.info( 

321 "Starting bulk delete: resource=%s count=%d soft=%s user=%s", 

322 resource_type, 

323 len(resource_ids), 

324 soft_delete, 

325 user_id, 

326 ) 

327 

328 deleted = 0 

329 failed = 0 

330 

331 # Placeholder: actual implementation would delete each record 

332 

333 return AdminTaskResult( 

334 success=True, 

335 message=f"Bulk delete completed: {deleted} deleted, {failed} failed", 

336 data={ 

337 "deleted": deleted, 

338 "failed": failed, 

339 }, 

340 ) 

341 

342 except (ValueError, ConnectionError, TimeoutError, OSError, KeyError) as e: 

343 return AdminTaskResult( 

344 success=False, 

345 message=str(e), 

346 errors=[str(e)], 

347 ) 

348 

349 

350# ========== Scheduled Tasks ========== 

351 

352 

353async def data_cleanup_task( 

354 resource_type: str | None = None, 

355 older_than_days: int = 90, 

356) -> AdminTaskResult: 

357 """Clean up old data (soft-deleted, expired, etc). 

358 

359 Can be scheduled to run periodically. 

360 """ 

361 try: 

362 logger.info( 

363 "Starting data cleanup: resource=%s older_than=%d days", 

364 resource_type or "all", 

365 older_than_days, 

366 ) 

367 

368 cleaned = 0 

369 

370 return AdminTaskResult( 

371 success=True, 

372 message=f"Cleaned up {cleaned} records", 

373 data={"cleaned": cleaned}, 

374 ) 

375 

376 except (ValueError, ConnectionError, TimeoutError, OSError, KeyError) as e: 

377 return AdminTaskResult( 

378 success=False, 

379 message=str(e), 

380 errors=[str(e)], 

381 ) 

382 

383 

384async def cache_warm_task( 

385 resource_types: list[str] | None = None, 

386) -> AdminTaskResult: 

387 """Warm caches for frequently accessed resources. 

388 

389 Can be scheduled to run after deployments or periodically. 

390 """ 

391 try: 

392 logger.info("Starting cache warm: resources=%s", resource_types or "all") 

393 

394 warmed = 0 

395 

396 return AdminTaskResult( 

397 success=True, 

398 message=f"Warmed {warmed} cache entries", 

399 data={"warmed": warmed}, 

400 ) 

401 

402 except (ValueError, ConnectionError, TimeoutError, OSError, KeyError) as e: 

403 return AdminTaskResult( 

404 success=False, 

405 message=str(e), 

406 errors=[str(e)], 

407 )