Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions config.json.example
Original file line number Diff line number Diff line change
Expand Up @@ -34,8 +34,8 @@
},

"pch": {
"parallel_downloads": 8,
"parallel_parsers": 40
"parallel_downloads": 1,
"parallel_parsers": 8
},

"openintel": {
Expand Down
217 changes: 78 additions & 139 deletions iyp/crawlers/pch/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,24 +8,24 @@
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

from iyp import (AddressValueError, BaseCrawler, CacheHandler,
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].
Expand All @@ -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),
Expand All @@ -91,175 +93,112 @@ 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."""
for status, resp, resp_name in self.fetch_urls([(name, url)]):
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
Expand Down
Loading