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
« 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
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
25from crawler.factory import crawler_factory
26from crawler.models import Source
27from crawler.utils import get_all_cols
29request_interval = getattr(settings, "REQUESTS_INTERVAL", 3)
30chunk_size_collections = settings.CHUNK_SIZE_COLLECTIONS
31chunk_size_sources = settings.CHUNK_SIZE_SOURCES
33if TYPE_CHECKING:
34 from history.model_data import HistoryEventDict
35 from ptf.model_data import IssueData
37tracer = trace.get_tracer(__name__)
38logger = logging.getLogger(__name__)
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)
53def download_pdf(obj: Article, remove_first_page: bool, only_new: bool, pause_function=time.sleep):
54 collection = obj.get_collection()
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
61 extlink = qs.first()
62 href = extlink.location
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
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")
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)
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)
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)
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 )
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")
130 qs = obj.datastream_set.filter(mimetype="application/pdf")
132 if qs:
133 datastream = qs.first()
134 else:
135 datastream = DataStream()
136 datastream.resource = obj
138 if datastream.location != pdf_filename:
139 datastream.location = pdf_filename
140 datastream.save()
141 return disk_location
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")
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)
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")
180 event_dict: "HistoryEventDict" = {
181 "type": "download-sources",
182 "pid": "all sources",
183 "col": None,
184 "status": HistoryEventStatus.PENDING,
185 }
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)
191 source_count = sources.count()
192 logger.info(f"Downloading article PDFs for {source_count} sources")
194 self.set_progress(current=1, total=source_count)
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)
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
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)
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 )
333 result = self.wait_child(promise, propagate=True)
334 results.append(result)
335 increment_progress.delay(self.request.id)
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)
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 }
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
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]
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
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
400 task.set_progress(current=0, total=len(issue_list), col=colid)
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)
413 event_dict["status"] = HistoryEventStatus.OK
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"])
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
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
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
479 # TODO here filtrage par période: quand start_year, end_year remplaceront year dans les conteneurs
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)
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
513 else:
514 event_dict["status"] = HistoryEventStatus.WARNING
515 event_dict["message"] = f"No article to download for the collection: {colid}"
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)
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)
539 return crawler.crawl_collection()
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)
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)
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 = {}
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
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
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 }
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)}
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