Coverage for src/crawler/tasks.py: 0%

337 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-09-23 14:47 +0000

1import logging 

2import os 

3import subprocess 

4import time 

5import traceback 

6from os import path 

7from typing import TYPE_CHECKING 

8 

9import pypdf 

10import requests 

11from celery import chain, group, shared_task 

12from django.conf import settings 

13from django.contrib.auth.models import User 

14from django.db.models import Q 

15from history.model_data import HistoryEventStatus 

16from history.utils import insert_history_event 

17from opentelemetry import trace 

18from ptf.display import resolver 

19from ptf.models import Article, Collection, Container 

20from ptf.models.classes.datastream import DataStream 

21from task import TaskAborted 

22from task.custom_task import PtfAbortableTask, TaskWithProgress 

23from task.tasks import increment_progress 

24 

25from crawler.factory import crawler_factory 

26from crawler.models import Source 

27from crawler.utils import get_all_cols 

28 

29request_interval = getattr(settings, "REQUESTS_INTERVAL", 3) 

30chunk_size_collections = settings.CHUNK_SIZE_COLLECTIONS 

31chunk_size_sources = settings.CHUNK_SIZE_SOURCES 

32 

33if TYPE_CHECKING: 

34 from history.model_data import HistoryEventDict 

35 from ptf.model_data import IssueData 

36 

37tracer = trace.get_tracer(__name__) 

38logger = logging.getLogger(__name__) 

39 

40 

41def get(href, retries=0): 

42 try: 

43 r = requests.get(href, timeout=10.0, verify=False) 

44 return r 

45 except (requests.ConnectionError, requests.ConnectTimeout) as e: 

46 if retries >= 3: 

47 raise e 

48 logger.info("Retry query %s", href) 

49 time.sleep(60) 

50 return get(href, retries + 1) 

51 

52 

53def download_pdf(obj: Article, remove_first_page: bool, only_new: bool, pause_function=time.sleep): 

54 collection = obj.get_collection() 

55 

56 qs = obj.extlink_set.filter(rel="article-pdf").all() 

57 if not qs: 

58 logger.warning(f"No PDF link for {obj.pid}") 

59 return 

60 

61 extlink = qs.first() 

62 href = extlink.location 

63 

64 if hasattr(obj, "my_container"): 

65 container_id = obj.my_container.pid 

66 obj_id = obj.pid.replace("/", "_") 

67 else: 

68 # Download a PDF of a book 

69 container_id = obj.pid 

70 obj_id = None 

71 

72 if href.find("http") == 0: 

73 disk_location = resolver.get_disk_location( 

74 settings.RESOURCES_ROOT, 

75 collection.pid, 

76 "pdf", 

77 container_id=container_id, 

78 article_id=obj_id, 

79 do_create_folder=True, 

80 ) 

81 logger.info("disk_location: %s", disk_location) 

82 if not (os.path.isfile(disk_location) and only_new): 

83 # Download unless the file is already present AND only_new = True 

84 logger.info("tempo: %s", request_interval) 

85 pause_function(request_interval) 

86 r = get(href) 

87 logger.info("http response status: %s", r.status_code) 

88 if len(r.text) > 10 and r.text[0:21] == "<!DOCTYPE html PUBLIC": 

89 print(f"{obj.doi} has an embargo, no PDF") 

90 

91 else: 

92 # Remove front page 

93 if remove_first_page: 

94 temp_location = os.path.join(settings.TEMP_FOLDER, "file.pdf") 

95 with open(temp_location, "wb") as f_: 

96 f_.write(r.content) 

97 

98 pdf_reader = pypdf.PdfReader(temp_location) 

99 pdf_writer = pypdf.PdfWriter() 

100 for page in range(len(pdf_reader.pages)): 

101 current_page = pdf_reader.pages[page] 

102 if page > 0: 

103 pdf_writer.add_page(current_page) 

104 

105 with open(disk_location, "wb") as f_: 

106 pdf_writer.write(f_) 

107 logger.info("PDF file saved on disk, %s", disk_location) 

108 else: 

109 with open(disk_location, "wb") as f_: 

110 f_.write(r.content) 

111 logger.info("PDF file saved on disk, %s", disk_location) 

112 

113 else: 

114 disk_location = resolver.get_disk_location( 

115 settings.RESOURCES_ROOT, 

116 collection.pid, 

117 "pdf", 

118 container_id=container_id, 

119 article_id=obj_id, 

120 do_create_folder=False, 

121 ) 

122 

123 if not os.path.isfile(disk_location): 

124 logger.debug(f"Directory does not yet exists : {disk_location}") 

125 new_location = resolver.get_relative_folder( 

126 collection.pid, container_id=container_id, article_id=obj_id 

127 ) 

128 pdf_filename = os.path.join(new_location, obj.pid + ".pdf") 

129 

130 qs = obj.datastream_set.filter(mimetype="application/pdf") 

131 

132 if qs: 

133 datastream = qs.first() 

134 else: 

135 datastream = DataStream() 

136 datastream.resource = obj 

137 

138 if datastream.location != pdf_filename: 

139 datastream.location = pdf_filename 

140 datastream.save() 

141 return disk_location 

142 

143 

144def crawl_sources( 

145 user_name, 

146 only_new=False, 

147 period: tuple[int, int] = (0, 9999), 

148 number: tuple[int, int] = (0, 99999), 

149): 

150 logger.info("Start crawling all sources") 

151 sources = Source.objects.exclude(domain="NUMDAM") 

152 

153 # we launch the source crawlings concurrently 

154 for source in sources: 

155 collections = ( 

156 Collection.objects.filter(content__origin__source=source) 

157 .distinct() 

158 .order_by("pid") 

159 .values("pid") 

160 ) 

161 colids = [col["pid"] for col in collections] 

162 crawl_source.delay(colids, source.domain, user_name, only_new, period, number) 

163 

164 

165@shared_task( 

166 name="crawler.tasks.download_sources", 

167 bind=True, 

168 queue="coordinator", 

169 base=TaskWithProgress, 

170) 

171def download_sources( 

172 self: "TaskWithProgress", 

173 user_name, 

174 only_new=False, 

175 period: tuple[int, int] = (0, 9999), 

176 number: tuple[int, int] = (0, 99999), 

177): 

178 logger.info("Start downloading PDF from all sources") 

179 

180 event_dict: "HistoryEventDict" = { 

181 "type": "download-sources", 

182 "pid": "all sources", 

183 "col": None, 

184 "status": HistoryEventStatus.PENDING, 

185 } 

186 

187 sources = Source.objects.exclude(domain__in=["NUMDAM", "GEODESIC"]) 

188 # source_liste = ["SCHOLASTICA", "BMMS", "PTM", "ARSIA", "IPB"] 

189 # sources = Source.objects.filter(domain__in=source_liste) 

190 

191 source_count = sources.count() 

192 logger.info(f"Downloading article PDFs for {source_count} sources") 

193 

194 self.set_progress(current=1, total=source_count) 

195 

196 # we launch the sources downloading concurrently 

197 task_chains = [] 

198 for num, source in enumerate(sources): 

199 logger.info(f"sources: {source.domain}") 

200 collections = ( 

201 Collection.objects.filter(content__origin__source=source) 

202 .distinct() 

203 .order_by("pid") 

204 .values("pid") 

205 ) 

206 colids = [col["pid"] for col in collections] 

207 if len(colids) > 0: 

208 source_task = download_source.si( 

209 colids, source.domain, user_name, only_new, period, number 

210 ) 

211 increment_task = increment_progress.si(self.request.id) 

212 task_chain = chain(source_task, increment_task) 

213 task_chains.append(task_chain) 

214 logger.info("Source %s : launch the PDF download.", source.domain) 

215 if len(task_chains) == chunk_size_sources or num == source_count - 1: 

216 try: 

217 task_group = group(*task_chains).delay() 

218 self.wait_child(task_group, propagate=False) 

219 task_chains = [] 

220 except TaskAborted: 

221 logger.info("Download all sources ABORTED") 

222 event_dict["status"] = HistoryEventStatus.ERROR 

223 event_dict["message"] = "Task aborted by user" 

224 insert_history_event(event_dict) 

225 logger.error(event_dict["message"]) 

226 raise 

227 event_dict["status"] = HistoryEventStatus.OK 

228 logger.info("Download all source FINISHED") 

229 insert_history_event(event_dict) 

230 

231 

232@shared_task( 

233 name="crawler.tasks.download_source", 

234 bind=True, 

235 queue="coordinator", 

236 base=TaskWithProgress, 

237) 

238def download_source( 

239 self: "TaskWithProgress", 

240 colids: list, 

241 source_domain: str, 

242 user_name, 

243 only_new=False, 

244 period: tuple[int, int] = (0, 9999), 

245 number: tuple[int, int] = (0, 99999), 

246): 

247 # TODO comment gérer le parent_id de HistoryEvent des collections pour qu'il pointe vers celui des sources ? 

248 event_dict: "HistoryEventDict" = { 

249 "type": "download-source", 

250 "pid": source_domain, 

251 "col": None, 

252 "status": HistoryEventStatus.PENDING, 

253 } 

254 logger.info("Start downloading the PDFs for the source: %s", source_domain) 

255 logger.info("Number of collections: %s", len(colids)) 

256 try: 

257 self.set_progress(current=1, total=len(colids), col=source_domain) 

258 results = [] 

259 for col in colids: 

260 logger.info("Start dowloading the pdf from the collection: %s", col) 

261 promise = download_collection.delay( 

262 col, source_domain, user_name, only_new, period, number 

263 ) 

264 result = self.wait_child(promise, propagate=True) 

265 logger.info(f"after wait_child: {col}") 

266 # result = promise.get(disable_sync_subtasks=False, propagate=True) 

267 results.append(result) 

268 increment_progress.delay(self.request.id) 

269 exceptions = [result for result in results if isinstance(result, Exception)] 

270 if len(exceptions) > 0: 

271 raise ExceptionGroup("Encountered errors while processing subtasks", exceptions) 

272 event_dict["status"] = HistoryEventStatus.OK 

273 # check the collection status 

274 

275 except TaskAborted: 

276 event_dict["status"] = HistoryEventStatus.ERROR 

277 event_dict["message"] = "Task aborted by user" 

278 raise 

279 except Exception: 

280 event_dict["status"] = HistoryEventStatus.ERROR 

281 event_dict["message"] = traceback.format_exc() 

282 logger.error(event_dict["message"]) 

283 promise.abort() 

284 raise 

285 finally: 

286 logger.info("Download one source finished") 

287 # TODO fetch the article numbers 

288 total_files = 0 

289 total_articles = 0 

290 for colid in colids: 

291 file_count, article_count = articles_count(colid, source_domain) 

292 total_articles += article_count 

293 total_files += file_count 

294 if total_files != total_articles: 

295 event_dict["status"] = HistoryEventStatus.WARNING 

296 event_dict["message"] = ( 

297 f"{total_articles} articles à télécharger \n{total_files} articles sur disque." 

298 ) 

299 insert_history_event(event_dict) 

300 

301 

302@shared_task( 

303 name="crawler.tasks.crawl_source", 

304 bind=True, 

305 queue="coordinator", 

306 base=TaskWithProgress, 

307) 

308def crawl_source( 

309 self: "TaskWithProgress", 

310 colids: list, 

311 source_domain: str, 

312 user_name, 

313 only_new=False, 

314 period: tuple[int, int] = (0, 9999), 

315 number: tuple[int, int] = (0, 99999), 

316): 

317 event_dict: "HistoryEventDict" = { 

318 "type": "import-source", 

319 "pid": "import all", 

320 "col": None, 

321 "status": HistoryEventStatus.PENDING, 

322 } 

323 logger.info("Start crawling the source: %s", source_domain) 

324 try: 

325 self.set_progress(current=0, total=len(colids), col=source_domain) 

326 results = [] 

327 for col in colids: 

328 logger.info("Start crawling the collection: %s", col) 

329 promise = crawl_collection.delay( 

330 col, source_domain, user_name, only_new, period, number 

331 ) 

332 

333 result = self.wait_child(promise, propagate=True) 

334 results.append(result) 

335 increment_progress.delay(self.request.id) 

336 

337 event_dict["status"] = HistoryEventStatus.OK 

338 exceptions = [result for result in results if isinstance(result, Exception)] 

339 if len(exceptions) > 0: 

340 raise ExceptionGroup("Encountered errors while processing subtasks", exceptions) 

341 except Exception: 

342 event_dict["status"] = HistoryEventStatus.ERROR 

343 event_dict["message"] = traceback.format_exc() 

344 logger.error(event_dict["message"]) 

345 raise 

346 finally: 

347 insert_history_event(event_dict) 

348 

349 

350@shared_task( 

351 name="crawler.tasks.crawl_collection", 

352 bind=True, 

353 queue="coordinator", 

354 base=TaskWithProgress, 

355) 

356@tracer.start_as_current_span("CrawlCollectionTask.do") 

357def crawl_collection( 

358 task: "TaskWithProgress", 

359 colid: str, 

360 source_domain: str, 

361 user_name, 

362 only_new=False, 

363 period: tuple[int, int] = (0, 9999), 

364 number: tuple[int, int] = (0, 99999), 

365): 

366 event_dict: "HistoryEventDict" = { 

367 "type": "import-collection", 

368 "pid": f"{colid}-{source_domain}", 

369 "col": None, 

370 "source": source_domain, 

371 "status": HistoryEventStatus.PENDING, 

372 } 

373 

374 try: 

375 logger.debug("craw_collection task") 

376 user = User.objects.get(username=user_name) 

377 collection = Collection.objects.filter(pid=colid).first() 

378 event_dict["col"] = collection 

379 

380 event_dict["userid"] = user.pk 

381 all_cols = get_all_cols() 

382 col_data = all_cols[colid] 

383 url = col_data["sources"][source_domain] 

384 

385 issue_list = crawl_issue_list(source_domain, colid, url, user_name) 

386 if not issue_list: 

387 event_dict["status"] = HistoryEventStatus.WARNING 

388 event_dict["message"] = f"No issue to import for the collection: {colid}" 

389 return 

390 

391 issue_list = filter_issues(colid, issue_list, period, number, event_dict, only_new) 

392 if not issue_list: 

393 event_dict["message"] = ( 

394 f"No issue to import with the selection for the collection: {colid}" 

395 ) 

396 event_dict["status"] = HistoryEventStatus.OK 

397 logger.debug(event_dict["message"]) 

398 return 

399 

400 task.set_progress(current=0, total=len(issue_list), col=colid) 

401 

402 logger.info( 

403 "%s issues to process for the source: %s and the collection: %s", 

404 len(issue_list), 

405 source_domain, 

406 colid, 

407 ) 

408 for issue in issue_list.values(): 

409 promise = crawl_issue.delay(issue, source_domain, colid, url, user_name) 

410 task.wait_child(promise) 

411 increment_progress.delay(task.request.id) 

412 

413 event_dict["status"] = HistoryEventStatus.OK 

414 

415 except BaseException: 

416 event_dict["status"] = HistoryEventStatus.ERROR 

417 event_dict["message"] = traceback.format_exc() 

418 logger.error(event_dict["message"]) 

419 raise 

420 finally: 

421 insert_history_event(event_dict) 

422 logger.info("history event inserted %s", event_dict["pid"]) 

423 

424 

425def articles_count(colid, source_domain): 

426 """compute the articles number on the disk and the downloaded articles""" 

427 articles = Article.objects.order_by("pid") 

428 articles = Article.objects.filter( 

429 Q(my_container__origin__source__domain=source_domain), 

430 Q(my_container__my_collection__pid=colid) 

431 | Q(my_container__my_collection__parent__pid=colid), 

432 ) 

433 articles_count = articles.count() 

434 disk_path = settings.RESOURCES_ROOT 

435 disk_path = path.join(disk_path, colid) 

436 folders = subprocess.run(["find", disk_path, "-type", "f"], capture_output=True) 

437 file_count = len(str(folders.stdout).split("\\n")) - 1 

438 return file_count, articles_count 

439 

440 

441@shared_task( 

442 name="crawler.tasks.download_collection", 

443 bind=True, 

444 queue="coordinator_child", 

445 base=TaskWithProgress, 

446) 

447@tracer.start_as_current_span("DownloadCollectionTask.do") 

448def download_collection( 

449 self: "TaskWithProgress", 

450 colid: str, 

451 source_domain: str, 

452 user_name, 

453 only_new=False, 

454 period: tuple[int, int] = (0, 9999), 

455 number: tuple[int, int] = (0, 99999), 

456): 

457 logger.debug("download_collection task") 

458 user = User.objects.get(username=user_name) 

459 collection = Collection.objects.filter(pid=colid).first() 

460 event_dict: "HistoryEventDict" = { 

461 "type": "download-collection", 

462 "pid": f"{colid}", 

463 "col": collection, 

464 "source": source_domain, 

465 "status": HistoryEventStatus.PENDING, 

466 } 

467 try: 

468 event_dict["userid"] = user.pk 

469 

470 articles = Article.objects.order_by("pid") 

471 articles = Article.objects.filter( 

472 Q(my_container__origin__source__domain=source_domain), 

473 Q(my_container__my_collection__pid=colid) 

474 | Q(my_container__my_collection__parent__pid=colid), 

475 ) 

476 articles_count = articles.count() 

477 disk_path = settings.RESOURCES_ROOT 

478 

479 # TODO here filtrage par période: quand start_year, end_year remplaceront year dans les conteneurs 

480 

481 if articles_count > 0: 

482 logger.info( 

483 f"Downloading PDF articles: {articles_count} articles to process for the source: {source_domain} and the collection: {colid}" 

484 ) 

485 self.set_progress(current=1, total=articles_count, col=colid) 

486 

487 download_tasks = [] 

488 for num, article in enumerate(articles): 

489 if self.is_aborted(): 

490 raise TaskAborted 

491 promise = download_article.si(article, source_domain, colid, only_new) 

492 download_tasks.append(promise) 

493 download_tasks.append(increment_progress.si(self.request.id)) 

494 if num % chunk_size_collections == 0: 

495 logger.info(f"Chunk {num}") 

496 batch_tasks = chain(*download_tasks).delay() 

497 self.wait_child(batch_tasks, propagate=True) 

498 download_tasks = [] 

499 if download_tasks != []: 

500 batch_tasks = chain(*download_tasks).delay() 

501 self.wait_child(batch_tasks, propagate=True) 

502 # Compare files on disk/file to download 

503 disk_path = path.join(disk_path, colid) 

504 folders = subprocess.run(["find", disk_path, "-type", "f"], capture_output=True) 

505 file_count = len(str(folders.stdout).split("\\n")) - 1 

506 event_dict["status"] = HistoryEventStatus.OK 

507 event_dict["message"] = ( 

508 f"{articles_count} articles à télécharger \n{file_count} articles sur disque." 

509 ) 

510 if file_count != articles_count: 

511 event_dict["status"] = HistoryEventStatus.WARNING 

512 

513 else: 

514 event_dict["status"] = HistoryEventStatus.WARNING 

515 event_dict["message"] = f"No article to download for the collection: {colid}" 

516 

517 except TaskAborted: 

518 event_dict["status"] = HistoryEventStatus.ERROR 

519 event_dict["message"] = "Task aborted by user" 

520 # insert_history_event(event_dict) 

521 raise 

522 except BaseException: 

523 event_dict["status"] = HistoryEventStatus.ERROR 

524 event_dict["message"] = traceback.format_exc() 

525 # insert_history_event(event_dict) 

526 raise 

527 finally: 

528 insert_history_event(event_dict) 

529 

530 

531@shared_task( 

532 name="crawler.tasks.crawl_issue_list", bind=True, base=TaskWithProgress, queue="executor" 

533) 

534def crawl_issue_list( 

535 self: "TaskWithProgress", source_domain: str, colid: str, url: str, username: str 

536): 

537 crawler = crawler_factory(source_domain, colid, username, pause_function=self.wait) 

538 

539 return crawler.crawl_collection() 

540 

541 

542@shared_task(name="crawler.tasks.crawl_issue", queue="executor", bind=True, base=TaskWithProgress) 

543def crawl_issue( 

544 self: "TaskWithProgress", 

545 issue: "IssueData", 

546 source_domain: str, 

547 colid: str, 

548 url: str, 

549 username: str, 

550): 

551 crawler = crawler_factory(source_domain, colid, username, pause_function=self.wait) 

552 logger.info("crawl issue: %s, %s", source_domain, colid) 

553 crawler.crawl_issue(issue) 

554 

555 

556@shared_task( 

557 name="crawler.tasks.download_article", queue="executor", bind=True, base=PtfAbortableTask 

558) 

559def download_article(self, article: Article, source_domain: str, colid: str, only_new: bool): 

560 logger.info("Download article: %s, %s", source_domain, colid) 

561 download_pdf(article, False, only_new, pause_function=self.wait) 

562 

563 

564def filter_issues( 

565 colid: str, 

566 issues: "dict[str, IssueData]", 

567 period: tuple[int, int] = (0, 9999), 

568 number: tuple[int, int] = (0, 99999), 

569 event_dict=None, 

570 only_new=False, 

571): 

572 if event_dict is None: 

573 event_dict = {} 

574 

575 def is_year_in_range(year): 

576 try: 

577 return period[0] <= int(year) <= period[1] 

578 except ValueError: 

579 event_dict["status"] = HistoryEventStatus.ERROR 

580 insert_history_event(event_dict) 

581 logger.error("Missing the year property for issues in the collection %s", colid) 

582 return False 

583 

584 def is_number_in_range(n): 

585 try: 

586 # n can be "1-2" 

587 if "-" in n or "–" in n: 

588 n = n.split("-" if "-" in n else "–") 

589 return number[0] <= int(n[0]) <= int(n[1]) <= number[1] 

590 # n can be "3" 

591 return number[0] <= int(n) <= number[1] 

592 except ValueError: 

593 # issue.number is not an integer, nor an integer range 

594 event_dict["status"] = HistoryEventStatus.ERROR 

595 insert_history_event(event_dict) 

596 logger.warning( 

597 "The number property for issues in the collection %s is not defined or is not a number", 

598 colid, 

599 ) 

600 return False 

601 

602 if period == (0, 9999) and number != (0, 99999): 

603 # no filter expected for the period, let's filter only on the volume number 

604 # we select the issue if the issue.number is "" (not always defined) 

605 issues = { 

606 pid: issue 

607 for pid, issue in issues.items() 

608 if is_number_in_range(issue.number) or issue.number == "" 

609 } 

610 

611 if number == (0, 99999) and period != (0, 9999): 

612 # no filter expected for the volume numbers, let's filter only on the period 

613 issues = {pid: issue for pid, issue in issues.items() if is_year_in_range(issue.fyear)} 

614 

615 if only_new: 

616 all_containers = Container.objects.filter(my_collection__pid=colid).all() 

617 issues = { 

618 pid: issue 

619 for pid, issue in issues.items() 

620 if not any(container.pid == pid for container in all_containers) 

621 } 

622 logger.info("issues filtered: %s", issues) 

623 return issues