From fd01cce72724f8df29c693a568beb333a94d7439 Mon Sep 17 00:00:00 2001 From: Malte Tashiro Date: Fri, 10 Jul 2026 08:28:45 +0000 Subject: [PATCH] Update PCH crawlers PCH is back online with an updated website that provides an API to list files. While this greatly simplifies the code, the website is also pretty weak so add additional checks for retries and reduce default parallelization to 1 query at a time. --- config.json.example | 4 +- iyp/crawlers/pch/__init__.py | 217 +++++++++++++---------------------- 2 files changed, 80 insertions(+), 141 deletions(-) diff --git a/config.json.example b/config.json.example index c83079d7..82f21609 100644 --- a/config.json.example +++ b/config.json.example @@ -34,8 +34,8 @@ }, "pch": { - "parallel_downloads": 8, - "parallel_parsers": 40 + "parallel_downloads": 1, + "parallel_parsers": 8 }, "openintel": { diff --git a/iyp/crawlers/pch/__init__.py b/iyp/crawlers/pch/__init__.py index f05fd61e..82a2ccfb 100644 --- a/iyp/crawlers/pch/__init__.py +++ b/iyp/crawlers/pch/__init__.py @@ -8,10 +8,7 @@ from multiprocessing import Pool from typing import Iterable, Tuple -from bs4 import BeautifulSoup -from bs4.element import ResultSet -from requests.adapters import HTTPAdapter, Response -from requests.exceptions import ChunkedEncodingError +from requests.adapters import HTTPAdapter from requests_futures.sessions import FuturesSession from urllib3.util.retry import Retry @@ -19,13 +16,16 @@ DataNotAvailableError) from iyp.crawlers.pch.show_bgp_parser import ShowBGPParser -PARALLEL_DOWNLOADS = 8 +PARALLEL_DOWNLOADS = 1 PARALLEL_PARSERS = 8 if os.path.exists('config.json'): config = json.load(open('config.json', 'r')) PARALLEL_DOWNLOADS = config['pch']['parallel_downloads'] PARALLEL_PARSERS = config['pch']['parallel_parsers'] +COLLECTOR_LIST_URL_FMT = 'http://downloads.pch.net/api/files/Routing_Data/IPv{af}_daily_snapshots/%Y/%m/' +FILE_FMT = os.path.join(COLLECTOR_LIST_URL_FMT, '{collector}/{collector}-ipv{af}_bgp_routes.%Y.%m.%d.gz') + class RoutingSnapshotCrawler(BaseCrawler): """Crawler for PCH route collector data[0]. @@ -50,22 +50,24 @@ def __init__(self, organization: str, url: str, name: str, af: int): logging.error(f'Invalid address family: {af}') raise AddressValueError(f'Invalid address family: {af}') self.MAX_LOOKBACK = timedelta(days=7) + # self.curr_date = datetime.now(tz=timezone.utc) + self.curr_date = datetime(2026, 7, 3, tzinfo=timezone.utc) + self.max_lookback_dt = self.curr_date - self.MAX_LOOKBACK + self.latest_available_date = None self.af = af - if self.af == 4: - self.file_format = '{collector}-ipv4_bgp_routes.{year}.{month:02d}.{day:02d}.gz' - else: - self.file_format = '{collector}-ipv6_bgp_routes.{year}.{month:02d}.{day:02d}.gz' self.parser = ShowBGPParser(self.af) - cache_file_prefix = f'CACHED.{datetime.now().strftime("%Y%m%d")}.v{self.af}.' + cache_file_prefix = f'CACHED.{self.curr_date.strftime("%Y%m%d")}.v{self.af}.' self.cache_handler = CacheHandler(self.get_tmp_dir(), cache_file_prefix) self.collector_files = dict() - self.collector_site_url = str() + self.collector_urls = dict() self.__initialize_session() super().__init__(organization, url, name) + self.reference['reference_url_data'] = self.curr_date.strftime(COLLECTOR_LIST_URL_FMT.format(af=self.af)) self.reference['reference_url_info'] = 'https://www.pch.net/resources/Routing_Data/' def __initialize_session(self) -> None: self.session = FuturesSession(max_workers=PARALLEL_DOWNLOADS) + self.session.headers['User-Agent'] = 'Internet Yellow Pages - admin@ihr.live' retry = Retry( backoff_factor=0.1, status_forcelist=(429, 500, 502, 503, 504), @@ -91,10 +93,10 @@ def fetch_urls(self, urls: list) -> Iterable: try: resp = query.result() yield resp.ok, resp.content, name - except ChunkedEncodingError as e: + except Exception as e: logging.error(f'Failed to retrieve data for {query}') logging.error(e) - return False, str(), name + yield False, str(), name def fetch_url(self, url: str, name: str = str()) -> Tuple[bool, str, str]: """Helper function for single URL.""" @@ -102,164 +104,101 @@ def fetch_url(self, url: str, name: str = str()) -> Tuple[bool, str, str]: return status, resp, resp_name return False, str(), str() - def fetch_collector_site(self) -> str: - """Fetch the HTML code of the collector site for the current month. + def fetch_and_parse_collector_urls(self, date: datetime) -> None: + """Fetch the list of collectors available on the specified date. + + Only the year and month components of the date are used, since collectors are + listed by month. Then create the direct file URL by using the mtime value + specified for each collector. - If the site does not yet exist, check the previous month, as long as it is - within the lookback interval. + This function propagates self.collector_urls with the URL to the latest file per + collector and sets self.latest_available_date to the newest mtime. In + particular, collectors with existing entries in self.collector_urls are not + overwritten so this function can be called for multiple dates. """ - logging.info('Fetching list of collectors.') - today = datetime.now(tz=timezone.utc) - self.collector_site_url = self.url + today.strftime('%Y/%m/') - resp = self.session.get(self.collector_site_url).result() - if resp.ok: - return resp.text - logging.warning(f'Failed to retrieve collector site from: {self.collector_site_url}') - curr_month = today.month - lookback = today - self.MAX_LOOKBACK - if lookback.month == curr_month: - logging.error('Failed to find current data.') - raise DataNotAvailableError('Failed to find current data.') - self.collector_site_url = self.url + today.strftime('%Y/%m/') - resp: Response = self.session.get(self.collector_site_url).result() - if resp.ok: - return resp.text - logging.warning(f'Failed to retrieve collector site from: {self.collector_site_url}') - logging.error('Failed to find current data.') - raise DataNotAvailableError('Failed to find current data.') - - @staticmethod - def filter_route_collector_links(links: ResultSet) -> list: - """Extract route collector names from HTML code.""" - collector_names = list() - for a in links: - if 'href' not in a.attrs or not a['href'].startswith('route-collector'): + status, resp, _ = self.fetch_url(date.strftime(COLLECTOR_LIST_URL_FMT.format(af=self.af))) + if not status: + return + try: + collector_list = json.loads(resp) + except json.JSONDecodeError as e: + logging.error(f'Failed to decode collector list: {e}') + return + for entry in collector_list: + collector = entry['name'] + if collector in self.collector_urls: continue - collector_names.append(a['href'].rstrip('/')) - return collector_names - - def make_url(self, collector_name: str, date: datetime) -> str: - """Create file URLs based on the template and gathered route collector names.""" - file_name = self.file_format.format(collector=collector_name, - year=date.year, - month=date.month, - day=date.day) - file_url = f'{self.url}{date.strftime("%Y/%m/")}{collector_name}/{file_name}' - return file_url - - def probe_latest_set(self, collector_name: str) -> datetime: - """Find the date of the latest available dataset for the specified collector. + mtime = entry['mtime'] + try: + mtime_dt = datetime.strptime(mtime, '%a, %d %b %Y %H:%M:%S GMT').replace(tzinfo=timezone.utc) + except ValueError as e: + logging.warning(f'Failed to parse mtime "{mtime}" for collector {collector}: {e}') + continue + if mtime_dt < self.max_lookback_dt: + print(f'Ignoring collector {collector} due to stale entry: {mtime_dt.isoformat()}') + continue + self.collector_urls[collector] = mtime_dt.strftime(FILE_FMT.format(af=self.af, collector=collector)) + if self.latest_available_date is None or self.latest_available_date < mtime_dt: + self.latest_available_date = mtime_dt - Start with the current date and look up to MAX_LOOKBACK days into the past if no - current data is found. + def get_collector_urls(self): + """Get the latest file URL per collector. - Return None if no data is found within the valid interval. + This wrapper function handles the edge case when the lookback interval overlaps + a month boundary and we need to check two directories for collectors. """ - logging.info('Probing latest available dataset.') - curr_date = datetime.now(tz=timezone.utc) - max_lookback = curr_date - self.MAX_LOOKBACK - while curr_date >= max_lookback: - file_name = self.file_format.format(collector=collector_name, - year=curr_date.year, - month=curr_date.month, - day=curr_date.day) - probe_url = f'{self.url}{curr_date.strftime("%Y/%m/")}{collector_name}/{file_name}' - resp = self.session.head(probe_url).result() - if resp.status_code == 200: - logging.info(f'Latest available dataset: {curr_date.strftime("%Y-%m-%d")}') - return curr_date - curr_date -= timedelta(days=1) - logging.error('Failed to find current data.') - raise DataNotAvailableError('Failed to find current data.') + self.fetch_and_parse_collector_urls(self.curr_date) + if self.curr_date.month != self.max_lookback_dt.month: + self.fetch_and_parse_collector_urls(self.max_lookback_dt) def fetch(self) -> None: """Fetch and cache all data. First get a list of collector names and their associated files. Then fetch the - files in parallel. If some files are not available for the current date, try - fetching older data as long as it is within the lookback interval. + files in parallel. All downloaded files are cached, so if this process is restarted, only files that are not in the cache are fetched, the rest is loaded from cache. Return True if there was an error during the fetching process, else False. """ - collector_names_name = 'collectors' tmp_dir = self.get_tmp_dir() if not os.path.exists(tmp_dir): tmp_dir = self.create_tmp_dir() - # Get a list of collector names - if self.cache_handler.cached_object_exists(collector_names_name): - self.collector_site_url, collector_names = self.cache_handler.load_cached_object(collector_names_name) - else: - collector_site = self.fetch_collector_site() - soup = BeautifulSoup(collector_site, features='html.parser') - links = soup.find_all('a') - collector_names = self.filter_route_collector_links(links) - self.cache_handler.save_cached_object(collector_names_name, (self.collector_site_url, collector_names)) - self.reference['reference_url_data'] = self.collector_site_url - - # Get the date of the latest available dataset based on the - # specified beacon collector. - # This may be not the best method if only the beacon collector - # is missing the most up-to-date data for some reason, but - # generally this prevents a lot of requests to non-existing - # files (one per collector) if the data for the current date - # is not yet available for all collectors. - # The chosen beacon collector has been very stable so far. - beacon_collector = 'route-collector.ams.pch.net' - latest_available_date = self.probe_latest_set(beacon_collector) - self.reference['reference_time_modification'] = latest_available_date.replace(hour=0, - minute=0, - second=0, - microsecond=0) - curr_date = datetime.now(tz=timezone.utc) - max_lookback = curr_date - self.MAX_LOOKBACK + self.get_collector_urls() + if not self.collector_urls: + raise DataNotAvailableError('Failed to find valid collectors.') + + self.reference['reference_time_modification'] = self.latest_available_date # Build list of URLs for files that are not yet cached, and # load existing files from cache. to_fetch = list() - for collector_name in collector_names: + for collector_name in self.collector_urls: if self.cache_handler.cached_object_exists(collector_name): collector_file = self.cache_handler.load_cached_object(collector_name) self.collector_files[collector_name] = collector_file else: - collector_url = self.make_url(collector_name, latest_available_date) - to_fetch.append((collector_name, collector_url)) + to_fetch.append((collector_name, self.collector_urls[collector_name])) # Fetch remaining files from PCH. - if to_fetch: - logging.info(f'{len(self.collector_files)}/{len(collector_names)} collector files in cache, fetching ' - f'{len(to_fetch)}') - - # If some collectors do not have current data available, - # try again until the max lookback window is reached. - failed_fetches = list() - while to_fetch and latest_available_date >= max_lookback: - failed_fetches = list() - for ok, content, name in self.fetch_urls(to_fetch): - if not ok: - failed_fetches.append(name) - continue - # Files are compressed with gzip. - content = gzip.decompress(content).decode('utf-8') - self.collector_files[name] = content - self.cache_handler.save_cached_object(name, content) - - # Create new URLs for collectors that have failed. - to_fetch = list() - latest_available_date -= timedelta(days=1) - for collector_name in failed_fetches: - collector_url = self.make_url(collector_name, latest_available_date) - to_fetch.append((collector_name, collector_url)) - if to_fetch: - logging.info(f'Retrying fetch for {len(to_fetch)} collectors for date ' - f'{latest_available_date.strftime("%Y-%m-%d")}') - if failed_fetches: - # Max lookback reached. - logging.warning(f'Failed to find current data for {len(failed_fetches)} collectors: {failed_fetches}') + attempt = 1 + while to_fetch and attempt <= 10: + logging.info(f' Attempt {attempt}: {len(self.collector_files)}/{len(self.collector_urls)} collector files ' + f'in cache, fetching {len(to_fetch)}') + for ok, content, name in self.fetch_urls(to_fetch): + if not ok: + continue + # Files are compressed with gzip. + content = gzip.decompress(content).decode('utf-8') + self.collector_files[name] = content + self.cache_handler.save_cached_object(name, content) + + missing_collectors = set(self.collector_urls.keys()) - self.collector_files.keys() + to_fetch = [(collector, self.collector_urls[collector]) for collector in missing_collectors] + attempt += 1 def run(self) -> None: """Fetch data from PCH, parse the files, and push nodes and relationships to the