diff --git a/k8s/welearn-datastack/templates/urlcollectors/cron-workflow.yaml b/k8s/welearn-datastack/templates/urlcollectors/cron-workflow.yaml index b17f2a8..bf0b94a 100644 --- a/k8s/welearn-datastack/templates/urlcollectors/cron-workflow.yaml +++ b/k8s/welearn-datastack/templates/urlcollectors/cron-workflow.yaml @@ -370,5 +370,16 @@ spec: name: {{ .name }} template: collect-world-bank-okr + - name: notre-environnement + templateRef: + name: {{ .name }} + template: collect-sitemap + arguments: + parameters: + - name: sitemap_url + value: https://www.notre-environnement.gouv.fr/sitemap.xml + - name: corpus_name + value: notre-environnement + {{- end }} {{- end }} diff --git a/k8s/welearn-datastack/templates/urlcollectors/workflow-template.yaml b/k8s/welearn-datastack/templates/urlcollectors/workflow-template.yaml index 93ef626..045d5fd 100644 --- a/k8s/welearn-datastack/templates/urlcollectors/workflow-template.yaml +++ b/k8s/welearn-datastack/templates/urlcollectors/workflow-template.yaml @@ -121,6 +121,7 @@ spec: volumeAttributes: secretName: {{ $.Values.common.azureShare.secret.name }} shareName: {{ $.Values.common.azureShare.name }} + - name: collect-wikipedia synchronization: semaphores: @@ -460,7 +461,6 @@ spec: secretName: {{ $.Values.common.azureShare.secret.name }} shareName: {{ $.Values.common.azureShare.name }} - - name: collect-world-bank-okr synchronization: semaphores: @@ -495,4 +495,48 @@ spec: volumeAttributes: secretName: {{ $.Values.common.azureShare.secret.name }} shareName: {{ $.Values.common.azureShare.name }} + + - name: collect-sitemap + synchronization: + semaphores: + - configMapKeyRef: + name: {{ .collectorSemaphore.configmapName }} + key: {{ .collectorSemaphore.standard.keyName }} + inputs: + parameters: + - name: sitemap_url + - name: corpus_name + container: + {{- with $.Values.image }} + image: {{ include "common.pods.image" (dict "root" $ "image" (dict "repository" .repository "path" .path "tag" .tag))}} + {{- end }} + args: + - python + - "-m" + - welearn_datastack.nodes_workflow.URLCollectors.node_sitemap_collect + envFrom: + - configMapRef: + name: {{ .name }} + env: + - name: SITEMAP_URL + value: {{ print "{{inputs.parameters.sitemap_url}}" | quote }} + - name: CORPUS_NAME + value: {{ print "{{inputs.parameters.corpus_name}}" | quote }} + volumeMounts: + - name: secrets + mountPath: "/secrets" + readOnly: true + - name: azure-share + mountPath: {{ $.Values.common.azureShare.mountPath }} + volumes: + - name: secrets + secret: + secretName: {{ .name }} + - name: azure-share + csi: + driver: file.csi.azure.com + readOnly: true + volumeAttributes: + secretName: {{ $.Values.common.azureShare.secret.name }} + shareName: {{ $.Values.common.azureShare.name }} {{- end }} diff --git a/tests/url_collector/test_sitemap_collector.py b/tests/url_collector/test_sitemap_collector.py new file mode 100644 index 0000000..3b9c22d --- /dev/null +++ b/tests/url_collector/test_sitemap_collector.py @@ -0,0 +1,256 @@ +from unittest import TestCase +from unittest.mock import MagicMock, call, patch + +from welearn_database.data.models import Corpus, WeLearnDocument + +from welearn_datastack.collectors.sitemap_collector import SiteMapURLCollector + + +class TestSitemapURLCollector(TestCase): + def setUp(self): + self.sitemap_index_content = """ + + + + https://www.example.org/post-sitemap.xml + 2026-07-15T08:48:24+00:00 + + + https://www.example.org/post-sitemap2.xml + 2025-01-16T10:33:52+00:00 + + + https://www.example.org/post-sitemap3.xml + 2026-07-15T08:48:24+00:00 + + + https://www.example.org/page-sitemap.xml + 2026-07-08T13:24:10+00:00 + + + https://www.example.org/formation-sitemap.xml + 2026-07-09T15:20:48+00:00 + + + """ + + self.sitemap_content = """ + + + + https://www.example.com/ + 2005-01-01 + monthly + 0.8 + + + + https://www.example.com/catalog?item=12&desc=vacation_hawaii + weekly + + + + https://www.example.com/catalog?item=73&desc=vacation_new_zealand + 2004-12-23 + weekly + + + + https://www.example.com/catalog?item=74&desc=vacation_newfoundland + 2004-12-23T18:00:15+00:00 + 0.3 + + + + https://www.example.com/catalog?item=83&desc=vacation_usa + 2004-11-23 + + + """ + self.corpus = Corpus(source_name="example", main_url="example.org") + self.collector = SiteMapURLCollector( + sitemap_url="https://example.org/sitemap.xml", corpus=self.corpus + ) + + def test__check_sitemap_index_is_index(self) -> None: + self.assertTrue(self.collector._is_sitemap_index(self.sitemap_index_content)) + + def test__check_sitemap_index_is_not_index(self) -> None: + self.assertFalse(self.collector._is_sitemap_index(self.sitemap_content)) + + def test__extract_url_from_regular_sitemap(self): + awaited_ret = [ + "https://www.example.com/", + "https://www.example.com/catalog?item=12&desc=vacation_hawaii", + "https://www.example.com/catalog?item=73&desc=vacation_new_zealand", + "https://www.example.com/catalog?item=74&desc=vacation_newfoundland", + "https://www.example.com/catalog?item=83&desc=vacation_usa", + ] + urls = self.collector._extract_urls(self.sitemap_content) + self.assertListEqual(urls, awaited_ret) + + def test__exxtract_url_from_sitemap_index(self): + awaited_ret = [ + "https://www.example.org/post-sitemap.xml", + "https://www.example.org/post-sitemap2.xml", + "https://www.example.org/post-sitemap3.xml", + "https://www.example.org/page-sitemap.xml", + "https://www.example.org/formation-sitemap.xml", + ] + urls = self.collector._extract_urls(self.sitemap_index_content) + self.assertListEqual(urls, awaited_ret) + + @patch( + "welearn_datastack.collectors.sitemap_collector.extracted_url_to_url_datastore" + ) + @patch("welearn_datastack.collectors.sitemap_collector.get_new_https_session") + def test_collect_simple_sitemap_not_index( + self, mock_get_session, mock_extract_datastore + ): + mock_session = MagicMock() + mock_get_session.return_value = mock_session + + sitemap_resp = MagicMock() + sitemap_resp.content = b"" + sitemap_resp.raise_for_status = MagicMock() + + pages_resp = MagicMock() + pages_resp.content = b"..." + pages_resp.raise_for_status = MagicMock() + + # Premier appel .get() -> sitemap racine, deuxième -> pages + mock_session.get.side_effect = [sitemap_resp, pages_resp] + + with ( + patch.object( + self.collector, "_is_sitemap_index", return_value=False + ) as mock_is_index, + patch.object( + self.collector, + "_extract_urls", + return_value=["https://example.com/page1"], + ) as mock_extract_urls, + ): + expected_docs = [MagicMock(spec=WeLearnDocument)] + mock_extract_datastore.return_value = expected_docs + + result = self.collector.collect() + + mock_get_session.assert_called_once() + self.assertEqual(mock_session.get.call_count, 2) + mock_session.get.assert_any_call("https://example.org/sitemap.xml") + mock_session.get.assert_any_call("https://example.org/sitemap.xml") + + sitemap_resp.raise_for_status.assert_called_once() + pages_resp.raise_for_status.assert_called_once() + + mock_is_index.assert_called_once_with(sitemap_resp.content.decode("utf-8")) + + self.assertEqual(mock_extract_urls.call_count, 1) + mock_extract_urls.assert_called_once_with(pages_resp.content.decode("utf-8")) + + mock_extract_datastore.assert_called_once_with( + urls=["https://example.com/page1"], corpus=self.corpus + ) + self.assertEqual(result, expected_docs) + + @patch( + "welearn_datastack.collectors.sitemap_collector.extracted_url_to_url_datastore" + ) + @patch("welearn_datastack.collectors.sitemap_collector.get_new_https_session") + def test_collect_sitemap_index_with_subsitemaps( + self, mock_get_session, mock_extract_datastore + ): + mock_session = MagicMock() + mock_get_session.return_value = mock_session + + root_resp = MagicMock(content=b"") + sub_resp = MagicMock(content=b"sub") + + for resp in (root_resp, sub_resp): + resp.raise_for_status = MagicMock() + + mock_session.get.side_effect = [root_resp, sub_resp] + + with ( + patch.object(self.collector, "_is_sitemap_index", return_value=True), + patch.object( + self.collector, + "_extract_urls", + side_effect=[ + [ + "https://example.com/sub-sitemap.xml" + ], # extraction depuis root (sous-sitemaps) + [ + "https://example.com/page1", + "https://example.com/page2", + ], # extraction pages + ], + ) as mock_extract_urls, + ): + expected_docs = ["doc1", "doc2"] + mock_extract_datastore.return_value = expected_docs + + result = self.collector.collect() + + self.assertEqual(mock_session.get.call_count, 2) + mock_session.get.assert_has_calls( + [ + call("https://example.org/sitemap.xml"), + call("https://example.com/sub-sitemap.xml"), + ] + ) + for resp in (root_resp, sub_resp): + resp.raise_for_status.assert_called_once() + + self.assertEqual(mock_extract_urls.call_count, 2) + mock_extract_datastore.assert_called_once_with( + urls=["https://example.com/page1", "https://example.com/page2"], + corpus=self.corpus, + ) + self.assertEqual(result, expected_docs) + + @patch("welearn_datastack.collectors.sitemap_collector.get_new_https_session") + def test_collect_raises_when_root_sitemap_http_error(self, mock_get_session): + mock_session = MagicMock() + mock_get_session.return_value = mock_session + + sitemap_resp = MagicMock() + sitemap_resp.raise_for_status.side_effect = Exception("HTTP 500") + mock_session.get.return_value = sitemap_resp + + with self.assertRaises(Exception): + self.collector.collect() + + mock_session.get.assert_called_once_with("https://example.org/sitemap.xml") + + @patch( + "welearn_datastack.collectors.sitemap_collector.extracted_url_to_url_datastore" + ) + @patch("welearn_datastack.collectors.sitemap_collector.get_new_https_session") + def test_collect_raises_when_subsitemap_http_error( + self, mock_get_session, mock_extract_datastore + ): + mock_session = MagicMock() + mock_get_session.return_value = mock_session + + root_resp = MagicMock(content=b"") + root_resp.raise_for_status = MagicMock() + + sub_resp = MagicMock() + sub_resp.raise_for_status.side_effect = Exception("HTTP 404") + + mock_session.get.side_effect = [root_resp, sub_resp] + + with ( + patch.object(self.collector, "_is_sitemap_index", return_value=True), + patch.object( + self.collector, + "_extract_urls", + return_value=["https://example.com/broken-sub.xml"], + ), + ): + with self.assertRaises(Exception): + self.collector.collect() + + mock_extract_datastore.assert_not_called() diff --git a/welearn_datastack/collectors/sitemap_collector.py b/welearn_datastack/collectors/sitemap_collector.py new file mode 100644 index 0000000..e1b67fc --- /dev/null +++ b/welearn_datastack/collectors/sitemap_collector.py @@ -0,0 +1,84 @@ +import logging +import os +from typing import List + +from welearn_database.data.models import Corpus, WeLearnDocument + +from welearn_datastack.collectors.helpers.feed_helpers import ( + extracted_url_to_url_datastore, +) +from welearn_datastack.data.url_collector import URLCollector +from welearn_datastack.modules.xml_extractor import XMLExtractor +from welearn_datastack.utils_.http_client_utils import get_new_https_session + +log_level: int = logging.getLevelName(os.getenv("LOG_LEVEL", "INFO")) +log_format: str = os.getenv( + "LOG_FORMAT", "[%(asctime)s][%(name)s][%(levelname)s] - %(message)s" +) + +if not isinstance(log_level, int): + raise ValueError("Log level is not recognized : '%s'", log_level) + +logging.basicConfig( + level=logging.getLevelName(log_level), + format=log_format, +) +logger = logging.getLogger(__name__) + + +class SiteMapURLCollector(URLCollector): + def __init__( + self, + sitemap_url: str, + corpus: Corpus, + ) -> None: + self.sitemap_url = sitemap_url + self.corpus = corpus + + @staticmethod + def _is_sitemap_index(sitemap_to_test: str) -> bool: + extractor = XMLExtractor(sitemap_to_test) + index = extractor.extract_content("sitemapindex") + + return bool(index) + + @staticmethod + def _extract_urls(sitemap_to_test: str) -> list[str]: + ret = [] + extractor = XMLExtractor(sitemap_to_test) + for loc in extractor.extract_content("loc"): + ret.append(loc.content) + return ret + + def collect(self) -> List[WeLearnDocument]: + logger.info("Start sitemap url collector") + http_client = get_new_https_session() + + sitemap_resp = http_client.get(self.sitemap_url) + sitemap_resp.raise_for_status() + sitemap_content = sitemap_resp.content.decode("utf-8") + + is_index = self._is_sitemap_index(sitemap_content) + logger.info("Sitemap is index ? : %s", is_index) + + sitemaps_urls = [] + if is_index: + logger.info("Sitemap %s is index", self.sitemap_url) + sitemaps_urls.extend(self._extract_urls(sitemap_content)) + else: + logger.info("Sitemap %s is not index", self.sitemap_url) + sitemaps_urls.append(self.sitemap_url) + + page_urls = [] + + for sm_url in sitemaps_urls: + logger.info("Get URLs from %s", sm_url) + pages_container = http_client.get(sm_url) + pages_container.raise_for_status() + pages_container_content = pages_container.content.decode("utf-8") + page_urls.extend(self._extract_urls(pages_container_content)) + + logger.info("We found %s urls", len(page_urls)) + ret = extracted_url_to_url_datastore(urls=page_urls, corpus=self.corpus) + + return ret diff --git a/welearn_datastack/nodes_workflow/URLCollectors/node_sitemap_collect.py b/welearn_datastack/nodes_workflow/URLCollectors/node_sitemap_collect.py new file mode 100644 index 0000000..46636cd --- /dev/null +++ b/welearn_datastack/nodes_workflow/URLCollectors/node_sitemap_collect.py @@ -0,0 +1,65 @@ +import logging +import os + +from dotenv import load_dotenv +from welearn_database.data.models import Corpus + +from welearn_datastack.collectors.sitemap_collector import SiteMapURLCollector +from welearn_datastack.nodes_workflow.URLCollectors.nodes_helpers.collect import ( + insert_urls, +) +from welearn_datastack.utils_.database_utils import create_db_session +from welearn_datastack.utils_.virtual_environement_utils import load_dotenv_local + +log_level: int = logging.getLevelName(os.getenv("LOG_LEVEL", "INFO")) +log_format: str = os.getenv( + "LOG_FORMAT", "[%(asctime)s][%(name)s][%(levelname)s] - %(message)s" +) + +if not isinstance(log_level, int): + raise ValueError("Log level is not recognized : '%s'", log_level) + +logging.basicConfig( + level=logging.getLevelName(log_level), + format=log_format, +) +logger = logging.getLogger(__name__) + + +if __name__ == "__main__": + load_dotenv() + logger.info("Sitemap collector starting...") + load_dotenv_local() + session = create_db_session() + sitemap_url = os.getenv("SITEMAP_URL") + + corpus_name = os.getenv("CORPUS_NAME") + + if not sitemap_url: + raise ValueError("SITEMAP_URL is not defined") + + if not corpus_name: + raise ValueError("CORPUS_NAME is not defined") + + corpus: Corpus | None = ( + session.query(Corpus).filter_by(source_name=corpus_name).one_or_none() + ) + + if corpus is None: + raise ValueError(f"Corpus {corpus_name} not found") + + sitemap_collector = SiteMapURLCollector( + sitemap_url=sitemap_url, + corpus=corpus, + ) + + urls = sitemap_collector.collect() + + logger.info("URLs retrieved : '%s'", len(urls)) + + insert_urls( + session=session, + urls=urls, + ) + + logger.info("Sitemap collector ended")