From 49787f5150c319a2100880f14b238b0cef0d026c Mon Sep 17 00:00:00 2001 From: Pieterjan Montens Date: Mon, 20 Apr 2026 21:38:33 +0200 Subject: [PATCH] :technologist: Updates done by claude --- py_agent/agent | 15 +- py_agent/classes/db.py | 305 +++++++++++++++++-- py_agent/fetch/__init__.py | 0 py_agent/fetch/commands.py | 109 +++++++ py_agent/parse/__init__.py | 0 py_agent/parse/commands.py | 213 ++++++++++++++ py_agent/parse/page.py | 585 +++++++++++++++++++++++++++++++++++++ py_agent/pyproject.toml | 2 +- 8 files changed, 1199 insertions(+), 30 deletions(-) create mode 100644 py_agent/fetch/__init__.py create mode 100644 py_agent/fetch/commands.py create mode 100644 py_agent/parse/__init__.py create mode 100644 py_agent/parse/commands.py create mode 100644 py_agent/parse/page.py diff --git a/py_agent/agent b/py_agent/agent index 7dc4433..909314f 100755 --- a/py_agent/agent +++ b/py_agent/agent @@ -10,15 +10,18 @@ logger.addHandler(logging.StreamHandler()) from scrape import commands as com_scrape -from debug import commands as com_debug +from fetch import commands as com_fetch +from parse import commands as com_parse +from debug import commands as com_debug # Agent is a simple CLI interface to multiple sub-parts -# needed to manage ETAAMB. Each module has its own CLI -# Commands. +# needed to manage ETAAMB. Each module has its own CLI commands. # # Modules: -# COM_DEBUG: Debugging and testing commands -# COM_SCRAPE: Scraping options +# COM_SCRAPE : Scrape numac IDs from the MB summary pages +# COM_FETCH : Fetch raw FR/NL article HTML for each numac +# COM_PARSE : Parse raw pages and populate structured DB tables +# COM_DEBUG : Debugging and testing helpers # # Support classes: # DB : database operations @@ -32,6 +35,8 @@ def init(): init.add_command(com_scrape.version) init.add_command(com_scrape.last_date) init.add_command(com_scrape.get_numacs) +init.add_command(com_fetch.fetch_raws) +init.add_command(com_parse.parse_raws) init.add_command(com_debug.config) init.add_command(com_debug.test_db) diff --git a/py_agent/classes/db.py b/py_agent/classes/db.py index fb81c8b..4f0af57 100644 --- a/py_agent/classes/db.py +++ b/py_agent/classes/db.py @@ -3,28 +3,60 @@ import logging logger = logging.getLogger(__name__) +# ── Versioning constants ─────────────────────────────────────────────────────── +RAW_VERSION = 1 +RAW_SOURCE_VERSION = '052024' +DOC_VERSION = 18 + + def get_config(): return { - 'DB_HOST': os.getenv('DB_HOST'), - 'DB_PORT': int(os.getenv('DB_PORT')), - 'DB_USER': os.getenv('DB_USER'), - 'DB_PASSWORD': os.getenv('DB_PASSWORD'), - 'DB_DATA': os.getenv('DB_DATA'), + 'DB_HOST': os.getenv('DB_HOST'), + 'DB_PORT': int(os.getenv('DB_PORT')), + 'DB_USER': os.getenv('DB_USER'), + 'DB_PASSWORD': os.getenv('DB_PASSWORD'), + 'DB_DATA': os.getenv('DB_DATA'), } + class obj: def __init__(self, config): self.config = config - self.conn = None; + self.conn = None + + # ── Internals ────────────────────────────────────────────────────────────── + + def ensure(self): + if not self.conn: + self.connect() + + def connect(self): + self.conn = pymysql.connect( + host=self.config['DB_HOST'], + port=self.config['DB_PORT'], + user=self.config['DB_USER'], + password=self.config['DB_PASSWORD'], + database=self.config['DB_DATA'], + cursorclass=pymysql.cursors.DictCursor, + ) + + def query(self, q): + self.ensure() + with self.conn.cursor() as cursor: + cursor.execute(q) + res = cursor.fetchall() + return res + + # ── Existing (scrape) ────────────────────────────────────────────────────── def test(self): try: self.query('SELECT COUNT(*) FROM done_dates') - return True; + return True except Exception as e: logger.exception(e) - return False; + return False def store_numacs(self, numacs, date_obj): self.ensure() @@ -37,23 +69,248 @@ def store_numacs(self, numacs, date_obj): for numac in numacs: cursor.execute(sql, (numac, date_str, 2)) - def ensure(self): - if not self.conn: - self.connect() + # ── Fetch module ─────────────────────────────────────────────────────────── - def query(self, q): + def get_numacs_to_fetch(self, limit=250, force=False): + """Return list of {doc_id, date} dicts that still need raw pages fetched.""" self.ensure() with self.conn.cursor() as cursor: - cursor.execute(q) - res = cursor.fetchall() - return res + if force: + cursor.execute( + "SELECT doc_id, date FROM raw_ids LIMIT %s", + (limit,), + ) + else: + cursor.execute( + """ + SELECT doc_id, date FROM raw_ids + LEFT JOIN raw_pages ON raw_ids.doc_id = raw_pages.numac + WHERE raw_pages.version != %s + UNION + SELECT doc_id, date FROM raw_ids + LEFT JOIN raw_pages ON raw_ids.doc_id = raw_pages.numac + WHERE raw_pages.numac IS NULL + LIMIT %s + """, + (RAW_VERSION, limit), + ) + return cursor.fetchall() - def connect(self): - self.conn = pymysql.connect( - host=self.config['DB_HOST'], - port=self.config['DB_PORT'], - user=self.config['DB_USER'], - password=self.config['DB_PASSWORD'], - database=self.config['DB_DATA'], - cursorclass=pymysql.cursors.DictCursor - ) + def get_numac_dates(self, numacs): + """Return {doc_id, date} for a specific list of numacs (for targeted fetch).""" + if not numacs: + return [] + self.ensure() + placeholders = ','.join(['%s'] * len(numacs)) + with self.conn.cursor() as cursor: + cursor.execute( + f"SELECT doc_id, date FROM raw_ids WHERE doc_id IN ({placeholders})", + numacs, + ) + return cursor.fetchall() + + def store_raw_page(self, numac, pub_date, raw_fr, raw_nl, force=False): + """Insert or conditionally overwrite a raw page record.""" + self.ensure() + if force: + sql = """ + INSERT INTO raw_pages (numac, pub_date, raw_fr, raw_nl, version, raw_source_version) + VALUES (%s, %s, %s, %s, %s, %s) + ON DUPLICATE KEY UPDATE + pub_date = VALUES(pub_date), + raw_fr = VALUES(raw_fr), + raw_nl = VALUES(raw_nl), + version = VALUES(version), + raw_source_version = VALUES(raw_source_version) + """ + else: + sql = """ + INSERT INTO raw_pages (numac, pub_date, raw_fr, raw_nl, version, raw_source_version) + VALUES (%s, %s, %s, %s, %s, %s) + ON DUPLICATE KEY UPDATE id = id + """ + with self.conn.cursor() as cursor: + cursor.execute(sql, (numac, pub_date, raw_fr, raw_nl, RAW_VERSION, RAW_SOURCE_VERSION)) + self.conn.commit() + + # ── Parse module ────────────────────────────────────────────────────────── + + def get_numacs_to_parse(self, limit=1000, force=False): + """Return list of numac strings that still need parsing.""" + self.ensure() + with self.conn.cursor() as cursor: + if force: + cursor.execute( + "SELECT numac FROM raw_pages WHERE raw_source_version = %s LIMIT %s", + (RAW_SOURCE_VERSION, limit), + ) + else: + cursor.execute( + """ + SELECT raw_pages.numac FROM raw_pages + LEFT JOIN docs ON raw_pages.numac = docs.numac + WHERE docs.numac IS NULL + AND raw_pages.raw_source_version = %s + UNION + SELECT raw_pages.numac FROM raw_pages + LEFT JOIN docs ON raw_pages.numac = docs.numac + WHERE docs.version != %s + AND raw_pages.raw_source_version = %s + LIMIT %s + """, + (RAW_SOURCE_VERSION, DOC_VERSION, RAW_SOURCE_VERSION, limit), + ) + return [row['numac'] for row in cursor.fetchall()] + + def get_raw_page(self, numac): + """Return {pub_date, raw_fr, raw_nl} for one numac, or None.""" + self.ensure() + with self.conn.cursor() as cursor: + cursor.execute( + "SELECT pub_date, raw_fr, raw_nl FROM raw_pages WHERE numac = %s", + (numac,), + ) + return cursor.fetchone() + + def get_or_create_source(self, src_nl, src_fr): + """Upsert a source pair and return its id.""" + self.ensure() + with self.conn.cursor() as cursor: + cursor.execute( + "SELECT id FROM sources WHERE source_nl = %s AND source_fr = %s", + (src_nl, src_fr), + ) + row = cursor.fetchone() + if row: + return row['id'] + cursor.execute( + "INSERT IGNORE INTO sources (source_nl, source_fr) VALUES (%s, %s)", + (src_nl, src_fr), + ) + self.conn.commit() + with self.conn.cursor() as cursor: + cursor.execute( + "SELECT id FROM sources WHERE source_nl = %s AND source_fr = %s", + (src_nl, src_fr), + ) + row = cursor.fetchone() + return row['id'] if row else None + + def get_or_create_type(self, type_nl, type_fr): + """Upsert a type pair and return its id.""" + self.ensure() + # When only one language has a real type value + if (type_nl == 'notype') != (type_fr == 'notype'): + lang = 'fr' if type_nl == 'notype' else 'nl' + value = type_fr if type_nl == 'notype' else type_nl + with self.conn.cursor() as cursor: + cursor.execute(f"SELECT id FROM types WHERE type_{lang} = %s", (value,)) + row = cursor.fetchone() + if row: + return row['id'] + else: + with self.conn.cursor() as cursor: + cursor.execute( + "SELECT id FROM types WHERE type_nl = %s AND type_fr = %s", + (type_nl, type_fr), + ) + row = cursor.fetchone() + if row: + return row['id'] + + with self.conn.cursor() as cursor: + cursor.execute( + "INSERT IGNORE INTO types (type_nl, type_fr) VALUES (%s, %s)", + (type_nl, type_fr), + ) + self.conn.commit() + with self.conn.cursor() as cursor: + cursor.execute( + "SELECT id FROM types WHERE type_nl = %s AND type_fr = %s", + (type_nl, type_fr), + ) + row = cursor.fetchone() + return row['id'] if row else None + + def store_doc(self, data, type_id, source_id): + """Upsert the main doc record. Always overwrites on conflict.""" + self.ensure() + sql = """ + INSERT INTO docs + (numac, pub_date, prom_date, type, source, version, anonymise, + eli_type_fr, eli_type_nl, chrono_id, + chamber_id, senate_id, chamber_leg, senate_leg) + VALUES (%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s,%s) + ON DUPLICATE KEY UPDATE + pub_date = VALUES(pub_date), + prom_date = VALUES(prom_date), + type = VALUES(type), + source = VALUES(source), + version = VALUES(version), + eli_type_fr = VALUES(eli_type_fr), + eli_type_nl = VALUES(eli_type_nl), + chrono_id = VALUES(chrono_id), + chamber_id = VALUES(chamber_id), + senate_id = VALUES(senate_id), + chamber_leg = VALUES(chamber_leg), + senate_leg = VALUES(senate_leg) + """ + with self.conn.cursor() as cursor: + cursor.execute(sql, ( + data['numac'], data['pub_date'], data['prom_date'], + type_id, source_id, DOC_VERSION, + 1 if data.get('anonymise') else 0, + data.get('eli_type_fr'), data.get('eli_type_nl'), + data.get('chrono_id'), + data.get('chamber_id'), data.get('senate_id'), + data.get('chamber_leg'), data.get('senate_leg'), + )) + self.conn.commit() + + def store_text(self, numac, lang, raw, pure): + """Upsert plain-text content for one language.""" + self.ensure() + sql = """ + INSERT INTO text (numac, ln, raw, pure, length) + VALUES (%s, %s, %s, %s, %s) + ON DUPLICATE KEY UPDATE + raw = VALUES(raw), pure = VALUES(pure), length = VALUES(length) + """ + with self.conn.cursor() as cursor: + cursor.execute(sql, (numac, lang, raw, pure, len(pure))) + self.conn.commit() + + def store_title(self, numac, lang, raw, pure): + """Upsert a title for one language.""" + self.ensure() + sql = """ + INSERT INTO titles (numac, ln, raw, pure) + VALUES (%s, %s, %s, %s) + ON DUPLICATE KEY UPDATE raw = VALUES(raw), pure = VALUES(pure) + """ + with self.conn.cursor() as cursor: + cursor.execute(sql, (numac, lang, raw, pure)) + self.conn.commit() + + def add_doc_language(self, numac, lang): + """Append a language code to docs.languages via CONCAT_WS.""" + self.ensure() + with self.conn.cursor() as cursor: + cursor.execute( + "UPDATE docs SET languages = CONCAT_WS(',', languages, %s) WHERE numac = %s", + (lang, numac), + ) + self.conn.commit() + + def store_doc_links(self, numac, chrono, eli, pdf): + """Upsert document external links.""" + self.ensure() + sql = """ + INSERT INTO doc_links (numac, chrono, eli, pdf) + VALUES (%s, %s, %s, %s) + ON DUPLICATE KEY UPDATE + chrono = VALUES(chrono), eli = VALUES(eli), pdf = VALUES(pdf) + """ + with self.conn.cursor() as cursor: + cursor.execute(sql, (numac, chrono, eli, pdf)) + self.conn.commit() diff --git a/py_agent/fetch/__init__.py b/py_agent/fetch/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/py_agent/fetch/commands.py b/py_agent/fetch/commands.py new file mode 100644 index 0000000..16595cd --- /dev/null +++ b/py_agent/fetch/commands.py @@ -0,0 +1,109 @@ +import asyncio +import random +import click +import aiohttp +import logging +from classes import db + +logger = logging.getLogger(__name__) + +_ROOT = 'https://www.ejustice.just.fgov.be' +_HEADERS = {'User-Agent': 'etaamb-agent/2.0'} +_TIMEOUT = aiohttp.ClientTimeout(total=15) + + +# ── HTTP helpers ─────────────────────────────────────────────────────────────── + +def _article_url(numac, pub_date, lang): + lg_txt = 'f' if lang == 'fr' else 'n' + return ( + f"{_ROOT}/cgi/article.pl" + f"?language={lang}&sum_date={pub_date}" + f"&lg_txt={lg_txt}&caller=sum&numac_search={numac}&view_numac=" + ) + + +async def _fetch(session, url): + async with session.get(url, headers=_HEADERS, timeout=_TIMEOUT) as resp: + resp.raise_for_status() + return await resp.text() + + +# ── Per-document coroutine ───────────────────────────────────────────────────── + +async def _process(session, sem, numac, pub_date, sleep_min, sleep_max, + dry_run, force, db_obj): + async with sem: + logger.info("Fetching %s (%s)…", numac, pub_date) + try: + page_fr, page_nl = await asyncio.gather( + _fetch(session, _article_url(numac, pub_date, 'fr')), + _fetch(session, _article_url(numac, pub_date, 'nl')), + ) + except Exception as exc: + logger.error("Failed to fetch %s: %s", numac, exc) + return False + + if not dry_run: + db_obj.store_raw_page(numac, pub_date, page_fr, page_nl, force=force) + + await asyncio.sleep(random.randint(sleep_min, sleep_max) / 1000) + return True + + +# ── Async runner ─────────────────────────────────────────────────────────────── + +async def _run(rows, concurrency, sleep_min, sleep_max, dry_run, force, db_obj): + sem = asyncio.Semaphore(concurrency) + async with aiohttp.ClientSession() as session: + tasks = [ + _process(session, sem, row['doc_id'], str(row['date']), + sleep_min, sleep_max, dry_run, force, db_obj) + for row in rows + ] + results = await asyncio.gather(*tasks) + return sum(1 for r in results if r) + + +# ── CLI command ──────────────────────────────────────────────────────────────── + +@click.command() +@click.option('-r', '--dry-run', is_flag=True, default=False, + help='Fetch pages but do not write to DB.') +@click.option('-n', '--numac', 'numacs', multiple=True, + help='Target a specific numac (repeatable: -n X -n Y).') +@click.option('-l', '--limit', default=250, show_default=True, + help='Max documents to fetch per run.') +@click.option('-f', '--force', is_flag=True, default=False, + help='Re-fetch and overwrite pages already stored.') +@click.option('-c', '--concurrency', default=5, show_default=True, + help='Number of concurrent HTTP requests.') +@click.option('--sleep-min', default=30, show_default=True, + help='Minimum polite sleep between requests (ms).') +@click.option('--sleep-max', default=80, show_default=True, + help='Maximum polite sleep between requests (ms).') +def fetch_raws(dry_run, numacs, limit, force, concurrency, sleep_min, sleep_max): + """Fetch raw FR/NL article pages from ejustice and store them.""" + db_obj = db.obj(db.get_config()) + + if numacs: + rows = db_obj.get_numac_dates(list(numacs)) + if not rows: + click.echo("None of the provided numacs were found in raw_ids.") + return + missing = set(numacs) - {str(r['doc_id']) for r in rows} + if missing: + click.echo(f"\tWarning: not found in raw_ids: {', '.join(missing)}") + else: + rows = db_obj.get_numacs_to_fetch(limit=limit, force=force) + + click.echo( + f"Fetching {len(rows)} document(s) " + f"[dry-run={dry_run}, force={force}, concurrency={concurrency}]" + ) + if not rows: + click.echo("Nothing to fetch.") + return + + done = asyncio.run(_run(rows, concurrency, sleep_min, sleep_max, dry_run, force, db_obj)) + click.echo(f"Done — {done}/{len(rows)} fetched successfully.") diff --git a/py_agent/parse/__init__.py b/py_agent/parse/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/py_agent/parse/commands.py b/py_agent/parse/commands.py new file mode 100644 index 0000000..6ecf3d9 --- /dev/null +++ b/py_agent/parse/commands.py @@ -0,0 +1,213 @@ +""" +parse/commands.py — Python port of parseRaws.pl. + +Uses asyncio + ThreadPoolExecutor so each concurrent worker gets its own DB +connection (mirrors Perl's Parallel::ForkManager fork-per-doc model). +aiohttp is NOT used here; all I/O is DB-bound and runs in threads. +""" + +import asyncio +import re +import click +import logging +from concurrent.futures import ThreadPoolExecutor +from bs4 import BeautifulSoup + +from classes import db +from .page import Page, normalize_text + +logger = logging.getLogger(__name__) + + +# ── Data extraction ──────────────────────────────────────────────────────────── + +def _make_data_object(raw_fr, raw_nl, numac, pub_date): + """Build the parsed data dict for one document (mirrors makeDataObject).""" + plain_fr = BeautifulSoup(raw_fr, 'html.parser').get_text() + plain_nl = BeautifulSoup(raw_nl, 'html.parser').get_text() + + page_fr = Page(raw_fr, numac, pub_date, 'fr') + page_nl = Page(raw_nl, numac, pub_date, 'nl') + + # Promulgation date: prefer FR, fall back to NL, default to null date + prom_fr = page_fr.get_prom_date() + prom_nl = page_nl.get_prom_date() + if prom_fr != '--': + prom_date = prom_fr + elif prom_nl != '--': + prom_date = prom_nl + else: + prom_date = '0000-00-00' + + raw_title_fr = page_fr.get_title() + raw_title_nl = page_nl.get_title() + type_fr = page_fr.get_type() + type_nl = page_nl.get_type() + + # Cross-language type reconciliation: pick the more specific classification + nl_idx = page_nl.type_index(type_nl) + fr_idx = page_fr.type_index(type_fr) + if nl_idx != fr_idx and raw_title_fr and raw_title_nl: + if (fr_idx < nl_idx and fr_idx != 0) or nl_idx == 0: + new_nl = page_nl.get_type_by_index(fr_idx) + new_fr = type_fr + else: + new_fr = page_fr.get_type_by_index(nl_idx) + new_nl = type_nl + logger.info( + "TYPE CORR %s || NL(%s):%s => %s, FR(%s):%s => %s", + numac, nl_idx, type_nl, new_nl, fr_idx, type_fr, new_fr, + ) + type_fr, type_nl = new_fr, new_nl + + if type_fr == type_nl == 'document': + logger.debug("DOCUMENT %s classified as generic document", numac) + + norm_title_fr = normalize_text(raw_title_fr) + norm_title_nl = normalize_text(raw_title_nl) + + # Anonymisation flag + _anon_fr = r'(loi accordant des naturalisations|relative aux noms et pr|reserve de recrutement|article 770|exercer la profession de detective prive|recrutement|etrangers)' + _anon_nl = r'(betreffende de namen en voornamen|wet die naturalisaties verleent|samenstelling van een wervingreserve|artikel 770|het beroep van prive|aanwerving|vreemdelingen)' + anonymise = bool( + re.search(_anon_fr, norm_title_fr, re.I) or + re.search(_anon_nl, norm_title_nl, re.I) + ) + if anonymise: + logger.info("ANON %s anonymise bit will be set", numac) + + chamber_id, chamber_leg = page_fr.get_chamber_data() + senate_id, senate_leg = page_fr.get_senate_data() + + return { + 'numac': numac, + 'pub_date': pub_date, + 'prom_date': prom_date, + 'raw_fr': plain_fr, + 'norm_fr': normalize_text(plain_fr), + 'raw_nl': plain_nl, + 'norm_nl': normalize_text(plain_nl), + 'raw_title_fr': raw_title_fr, + 'norm_title_fr':norm_title_fr, + 'raw_title_nl': raw_title_nl, + 'norm_title_nl':norm_title_nl, + 'type_fr': type_fr, + 'type_nl': type_nl, + 'source_fr': page_fr.get_source(), + 'source_nl': page_nl.get_source(), + 'chrono': page_fr.get_chrono(), + 'eli': page_fr.get_eli(), + 'pdf': page_fr.get_pdf(), + 'eli_type_fr': page_fr.get_eli_type(), + 'eli_type_nl': page_nl.get_eli_type(), + 'chrono_id': page_fr.get_chrono_data(), + 'chamber_id': chamber_id, + 'chamber_leg': chamber_leg, + 'senate_id': senate_id, + 'senate_leg': senate_leg, + 'anonymise': anonymise, + } + + +def _store_data(db_obj, data): + """Persist all parsed data for one document.""" + type_id = db_obj.get_or_create_type(data['type_nl'], data['type_fr']) + source_id = db_obj.get_or_create_source(data['source_nl'], data['source_fr']) + + db_obj.store_doc(data, type_id, source_id) + db_obj.store_text(data['numac'], 'fr', data['raw_fr'], data['norm_fr']) + db_obj.store_text(data['numac'], 'nl', data['raw_nl'], data['norm_nl']) + + if data['raw_title_fr']: + db_obj.store_title(data['numac'], 'fr', data['raw_title_fr'], data['norm_title_fr']) + db_obj.add_doc_language(data['numac'], 'fr') + + if data['raw_title_nl']: + db_obj.store_title(data['numac'], 'nl', data['raw_title_nl'], data['norm_title_nl']) + db_obj.add_doc_language(data['numac'], 'nl') + + db_obj.store_doc_links( + data['numac'], + data.get('chrono'), data.get('eli'), data.get('pdf'), + ) + + +# ── Per-document worker (runs in thread pool) ────────────────────────────────── + +def _process_one(config, numac, dry_run): + """Fetch raw pages from DB, parse, and store results. + Each call creates its own DB connection — safe for ThreadPoolExecutor.""" + thread_db = db.obj(config) + raw = thread_db.get_raw_page(numac) + if not raw: + logger.warning("No raw data found for %s", numac) + return False + + try: + data = _make_data_object( + raw['raw_fr'], raw['raw_nl'], + numac, str(raw['pub_date']), + ) + except Exception as exc: + logger.error("Parse failed for %s: %s", numac, exc) + return False + + if not dry_run: + _store_data(thread_db, data) + + logger.info("Parsed %s (type FR: %s / NL: %s)", numac, data['type_fr'], data['type_nl']) + return True + + +# ── Async runner ─────────────────────────────────────────────────────────────── + +async def _run(numacs, concurrency, config, dry_run): + sem = asyncio.Semaphore(concurrency) + loop = asyncio.get_event_loop() + + async def bounded(numac): + async with sem: + return await loop.run_in_executor(None, _process_one, config, numac, dry_run) + + results = await asyncio.gather(*[bounded(n) for n in numacs], return_exceptions=True) + + errors = [r for r in results if isinstance(r, Exception)] + for exc in errors: + logger.error("Worker exception: %s", exc) + + return sum(1 for r in results if r is True) + + +# ── CLI command ──────────────────────────────────────────────────────────────── + +@click.command() +@click.option('-r', '--dry-run', is_flag=True, default=False, + help='Parse documents but do not write to DB.') +@click.option('-n', '--numac', 'numacs', multiple=True, + help='Target a specific numac (repeatable: -n X -n Y).') +@click.option('-l', '--limit', default=1000, show_default=True, + help='Max documents to parse per run.') +@click.option('-f', '--force', is_flag=True, default=False, + help='Re-parse docs whose version already matches (skip version check).') +@click.option('-c', '--concurrency', default=4, show_default=True, + help='Number of parallel parse workers.') +def parse_raws(dry_run, numacs, limit, force, concurrency): + """Parse stored raw pages and populate docs/text/titles tables.""" + config = db.get_config() + db_obj = db.obj(config) + + if numacs: + to_process = list(numacs) + else: + to_process = db_obj.get_numacs_to_parse(limit=limit, force=force) + + click.echo( + f"Parsing {len(to_process)} document(s) " + f"[dry-run={dry_run}, force={force}, concurrency={concurrency}]" + ) + if not to_process: + click.echo("Nothing to parse.") + return + + done = asyncio.run(_run(to_process, concurrency, config, dry_run)) + click.echo(f"Done — {done}/{len(to_process)} parsed successfully.") diff --git a/py_agent/parse/page.py b/py_agent/parse/page.py new file mode 100644 index 0000000..7e2458f --- /dev/null +++ b/py_agent/parse/page.py @@ -0,0 +1,585 @@ +""" +page.py — Python port of Page.pm (Autocorrect dependency removed). + +Two public normalize functions are exported for use in commands.py: + - normalize_title(txt) mirrors Page.pm's normalize() — for type matching + - normalize_text(txt) mirrors parseRaws.pl normalize() — for DB storage +""" + +import re +import unicodedata +from bs4 import BeautifulSoup +import logging + +logger = logging.getLogger(__name__) + + +# ── Text normalization ───────────────────────────────────────────────────────── + +def _strip_accents(txt): + nfkd = unicodedata.normalize('NFD', txt) + return ''.join(c for c in nfkd if not unicodedata.combining(c)) + + +def normalize_title(txt): + """Strip accents, lowercase, remove punctuation (except . ; : - and space). + Used internally by Page.get_type() — mirrors Page.pm::normalize().""" + if not txt: + return '' + txt = _strip_accents(txt) + txt = txt.lower() + txt = re.sub(r'\s', ' ', txt) + txt = re.sub(r'[^A-Za-z0-9.;:\- ]', '', txt) + return txt + + +def normalize_text(txt): + """Strip accents, lowercase, remove elisions + punctuation. + Used when storing normalized text to DB — mirrors parseRaws.pl::normalize().""" + if not txt: + return '' + txt = _strip_accents(txt) + txt = txt.lower() + txt = re.sub(r'\s', ' ', txt) + txt = re.sub(r"\w'", '', txt) + txt = re.sub(r'[^A-Za-z0-9.;:\- ]', '', txt) + return txt + + +def _norm_html(html): + """Produce the normalized HTML string used for regex matching inside Page. + Mirrors the processing done in Perl's Page::new: + unac_string + lc + strip + fix ministrieel typo.""" + if not html: + return '' + html = re.sub(r'', '', html, flags=re.IGNORECASE) + html = html.replace('ministrieel', 'ministerieel') + return _strip_accents(html).lower() + + +# ── Type-matching helpers ────────────────────────────────────────────────────── + +_DATE_PREFIX = ( + r'(?:' + r'\s*(?:\d+/){2}\d+\s*\.?' # 01/01/2024. + r'|\s*\d+(?:er)?\s*\w+\s*\d+\s*\.?' # 1er janvier 2024. + r'|[^.]*\.\s*' # anything then dot + r'|\s*' # just whitespace + r')' +) + + +def _matches_at_start(title, mask): + """Match mask at the start of title, possibly after a date/preamble prefix.""" + pattern = rf'^(?:{mask})|^{_DATE_PREFIX}-?\s*[^a-z]{{0,}}{mask}' + return bool(re.search(pattern, title, re.IGNORECASE)) + + +def _matches_anywhere(title, mask): + return bool(re.search(rf'\W{mask}\W', title, re.IGNORECASE)) + + +def _matches_source(source, source_mask): + return bool(re.match(rf'^\s*{source_mask}\s*$', source, re.IGNORECASE)) + + +def _matches_doc_mask(title, doc_mask): + if not doc_mask: + return False + return bool(re.search(rf'\W{doc_mask}\W|^{doc_mask}\W', title, re.IGNORECASE)) + + +# ── Classification tables (faithfully ported from Page.pm) ─────────────────── + +SOURCE_TYPES = { + r'(ministere de la|service public federal) justice|(ministerie van|overheidsdienst) justitie': [ + {'mask': r'rechterlijke orde bekendmaking|ordre judiciaire publication', + 'fr': "publication de l'ordre judiciaire", + 'nl': 'bekendmaking van de juridische orde'}, + {'mask': r'jugement|vonnis', + 'fr': 'jugement', 'nl': 'vonnis'}, + {'mask': r'organisation dun concours de recrutement|vergelijkend wervingsexamen|oproep(.*)kandidaten|appel(.*)candidats', + 'fr': 'recrutement', 'nl': 'aanwerving'}, + {'mask': r'rechterlijke orde|ordre judiciaire', + 'fr': "document concernant l'ordre judiciaire", + 'nl': 'document betreffende de rechterlijke ordre'}, + {'mask': r'services judiciaires|gerechtelijke diensten', + 'fr': 'document concernant les services judiciaires', + 'nl': 'document betreffende de gerechtelijke diensten'}, + ], + r'(ministere de la|service public federal) affaires economiques|(ministerie van|overheidsdienst) economische zaken': [ + {'mask': r'liste|lijst', 'fr': 'liste', 'nl': 'lijst'}, + {'mask': r'prix.*gaz nature|prijzen.*aardgas', 'fr': 'fixation des prix', 'nl': 'vaststelling van de prijzen'}, + {'mask': r'enregistremen.*normes|registratie.*normen', 'fr': 'enregistrement de normes belges', 'nl': 'registratie van belgische normen'}, + {'mask': r'demande.*concession|concessieaanvraag', 'fr': 'demande de concession', 'nl': 'concessieaanvraag'}, + {'mask': r'mededeling|information', 'fr': 'communication', 'nl': 'mededeling'}, + {'mask': r'conseil de la concurrence.*avis|raad voor de mededinging.*kennisgeving', 'fr': 'avis du conseil de la concurrence', 'nl': 'kennisgeving van de raad voor de mededinging'}, + {'mask': r'conseil de la concurrence.*decision|raad voor de mededinging.*beslissing', 'fr': 'décision du conseil de la concurrence', 'nl': 'beslissing van de raad voor de mededinging'}, + {'mask': r'formule i.b.t|formule e.i.l', 'fr': 'formule i.b.t.', 'nl': 'formule e.i.l.'}, + {'mask': r'interdiction de mise sur le marche|verbod tot het in de handel brengen', 'fr': 'interdiction de mise sur le marche', 'nl': 'verbod tot het in de handel brengen'}, + {'mask': r'actes dapprobation|akten tot goedkeuring', 'fr': 'approbations', 'nl': 'goedkeuringen'}, + {'mask': r'indice des prix a la (consommation|production)|indexcijfer.*(produktieprijzen|consumptieprijzen)', 'fr': 'indice des prix', 'nl': 'indexcijfer van de prijzen'}, + ], + r"ministere.*emploi et du travail|ministerie van tewerkstelling en arbeid|service public federal emploi, travail et concertation sociale|federale overheidsdienst werkgelegenheid, arbeid en sociaal overleg": [ + {'mask': r'reglement general.*du travail|reglement.*arbeidsbescherming', 'fr': 'règlement général de la protection du travail', 'nl': 'algemeen reglement voor de arbeidsbescherming'}, + {'mask': r'neerlegging van collectieve arbeidsovereenkomsten|depot de conventions collectives de travail', 'fr': 'dépôt de conventions collectives de travail', 'nl': 'neerlegging van collectieve arbeidsovereenkomsten'}, + ], + r"(ministere de la|service public federal) affaires sociales.*environnement|(ministerie van|overheidsdienst) sociale zaken.*leefmilieu": [ + {'mask': r'directive|richtlijn', 'fr': 'directive', 'nl': 'richtlijn'}, + {'mask': r'liste|lijst', 'fr': 'liste', 'nl': 'lijst'}, + {'mask': r'rijksinstituut.*invaliditeitsverzekering|institut.*maladie.invalidite', 'fr': "document de l'institut national d'assurance maladie-invalidité", 'nl': ' document van het rijksinstituut voor ziekte- en invaliditeitsverzekering'}, + ], + r'^(services du premier ministre|diensten van de eerste minister)$': [ + {'mask': '', 'fr': 'document des services du premier ministre', 'nl': 'document van de diensten van de eerste minister'}, + ], + r'(ministere de la|service public federal) finances|(ministerie van|overheidsdienst) financien': [ + {'mask': r'lotenlening|emprunt a lots', 'fr': 'emprunt à lots', 'nl': 'lotenlening'}, + {'mask': r'domaines publications? prescrites? par|bekendmaking(en)? voorgeschreven bij', 'fr': 'publication', 'nl': 'bekendmaking'}, + {'mask': r'loterie nationale|nationale loterij|lucky bingo', 'fr': 'document concernant la loterie nationale', 'nl': 'document betreffende de nationale loterij'}, + {'mask': r'monnaie royale de belgique|koninklijke munt van belgie', 'fr': 'document concernant la monnaie royale de belgique', 'nl': 'document betreffende de koninklijke munt van belgie'}, + {'mask': r'(mise en competition|incompetitiestelling)', 'fr': 'mise en compétition', 'nl': 'incompetitiestelling'}, + {'mask': r'administration du cadastre|administratie van het kadaster', 'fr': 'document - cadastre', 'nl': 'document - kadaster'}, + {'mask': r'administration de la t.v.a.|administratie van de btw', 'fr': 'document - t.v.a.', 'nl': 'document - btw'}, + {'mask': r'administration (generale )?de la tresorerie|(algemene )?administratie (der|van de) thesaurie', 'fr': 'document - trésorerie', 'nl': 'document - thesaurie'}, + {'mask': r'administration (generale )?de la fiscalite|(algemene )?administratie (der|van de) fiscaliteit', 'fr': 'document - fiscalité', 'nl': 'document - fiscaliteit'}, + {'mask': r'decision du president|beslissing van de voorzitter', 'fr': 'décision', 'nl': 'beslissing'}, + ], + r"(ministere de la|service public federal) communications et de l'infrastructure|ministerie van verkeer en infrastructuur": [ + {'mask': r'loterie nationale|nationale loterij|lucky bingo', 'fr': 'document concernant la loterie nationale', 'nl': 'document betreffende de nationale loterij'}, + ], + r'(ministere.*|service public federal )interieur|(ministerie van|overheidsdienst) binnenlandse zaken': [ + {'mask': r"intrekking van de vergunning|abrogation de l'?autorisation", 'fr': "abrogation d'autorisation", 'nl': 'intrekking van vergunning'}, + {'mask': r'autorisations? (dexploiter une|dexercer)|vergunning(en)? (tot het exploiteren|om het beroep)', 'fr': 'autorisation', 'nl': 'vergunning'}, + {'mask': r'gendarmerie\.|rijkswacht\.', 'fr': 'document concernant la gendarmerie', 'nl': 'document betreffende de rijkswacht'}, + {'mask': r"elections? communales? d(u|e)", 'fr': 'document concernant les elections communales', 'nl': ''}, + {'mask': r"conseil detat. - avis|raad van state. - bericht", 'fr': "avis du conseil d'état", 'nl': 'bericht van de raad van state'}, + {'mask': r'arretes? concernant les provinces|besluit(en)? betreffende de provincies', 'fr': 'arrêté royal', 'nl': 'koninklijk besluit'}, + ], + r'service public federal economie, p.m.e., classes moyennes et energie|federale overheidsdienst economie, k.m.o., middenstand en energie': [ + {'mask': r'autorisation individuelle|individuele vergunning', 'fr': 'autorisation', 'nl': 'vergunning'}, + {'mask': r'federaal ontwikkelingsplan|plan de developpement federal', 'fr': 'plan de développement', 'nl': 'ontwikkelingsplan'}, + {'mask': r'conseil de la concurrence.*avis|raad voor de mededinging.*kennisgeving', 'fr': 'avis du conseil de la concurrence', 'nl': 'kennisgeving van de raad voor de mededinging'}, + {'mask': r'liste|lijst', 'fr': 'liste', 'nl': 'lijst'}, + {'mask': r'prix.*gaz nature|prijzen.*aardgas', 'fr': 'fixation des prix', 'nl': 'vaststelling van de prijzen'}, + {'mask': r'enregistremen.*normes|registratie.*normen', 'fr': 'enregistrement de normes belges', 'nl': 'registratie van belgische normen'}, + {'mask': r'demande.*concession|concessieaanvraag', 'fr': 'demande de concession', 'nl': 'concessieaanvraag'}, + {'mask': r'mededeling|information', 'fr': 'communication', 'nl': 'mededeling'}, + {'mask': r'conseil de la concurrence.*decision|raad voor de mededinging.*beslissing', 'fr': 'décision du conseil de la concurrence', 'nl': 'beslissing van de raad voor de mededinging'}, + {'mask': r'formule i.b.t|formule e.i.l', 'fr': 'formule i.b.t.', 'nl': 'formule e.i.l.'}, + {'mask': r'interdiction de mise sur le marche|verbod tot het in de handel brengen', 'fr': 'interdiction de mise sur le marché', 'nl': 'verbod tot het in de handel brengen'}, + {'mask': r'actes dapprobation|akten tot goedkeuring', 'fr': 'approbations', 'nl': 'goedkeuringen'}, + {'mask': r'indice des prix a la (consommation|production)|indexcijfer.*(produktieprijzen|consumptieprijzen)', 'fr': 'indice des prix', 'nl': 'indexcijfer van de prijzen'}, + ], + r"cour d'arbitrage|arbitragehof|cour constitutionnelle|grondwettelijk hof": [ + {'mask': r"extrait de larret|uittreksel uit arrest", 'fr': "extrait d'un arrêt", 'nl': 'arrest uittreksel'}, + {'mask': r'arrest|arret', 'fr': 'arrêt de la cour constitutionelle', 'nl': 'arrest van het grondwettelijk hof'}, + ], + r'ministere de la communaute francaise': [ + {'mask': r'jury de promotion charge de classer les candidats', 'fr': 'communique du ministere de la communaute francaise', 'nl': ''}, + ], + r'ministere.*bruxelles-capitale|ministerie.*brussels hoofdstedelijk gewest': [ + {'mask': r'arretes concernant la ville.*bruxelles|besluiten betreffende de stad.*brussel|arretes concernant.*villes.*communes|besluiten betreffende de provincies steden en gemeenten', + 'fr': 'arrêtés concernant bruxelles', 'nl': 'besluiten betreffende brussel'}, + ], + r'ministere de la communaute germanophone|ministerie van de duitstalige gemeenschap': [ + {'mask': r'arrete du (ministre|gouvernement)|besluit van de (minister|regering)', 'fr': 'arrêté de la communauté germanophone', 'nl': 'besluit van de duitstalige gemeenschap'}, + ], + r'ministere de la communaute flamande|ministerie van de vlaamse gemeenschap': [ + {'mask': r'societe publique des dechets|openbare afvalstoffenmaatschappij', 'fr': 'document du ministere de la communauté flamande', 'nl': 'document van het ministeriel van de vlaamse gemeenschapscommissie'}, + ], + r'ministere de la region wallonne': [ + {'mask': r'tresorerie situation mensuelle du tresor', 'fr': 'situation mensuelle du trésor', 'nl': ''}, + ], + r'parlement de la region de bruxelles-capitale|brussels hoofdstedelijk parlement': [ + {'mask': r'seances plenieres ordre du jour|plenaire vergaderingen agenda', 'fr': 'ordre du jour des séances plénières', 'nl': 'agenda van de plenaire vergaderingen'}, + ], + r'commission bancaire et financiere|commissie voor.*financiewezen': [ + {'mask': r'autorisation.*cession|toestemming.*overdracht', 'fr': 'document de la commission bancaire et financière', 'nl': 'document van de commissie voor het bank- en financiewezen'}, + ], + r'selor.*bureau de selection|selor.*overheid': [ + {'mask': '', 'fr': 'communication du selor', 'nl': 'bericht van selor'}, + ], + r'pouvoir judiciaire|rechterlijke macht': [ + {'mask': r'aanwijzing|designation|designe|aangewezen', 'fr': "désignation dans l'ordre judiciaire", 'nl': 'aanwijzing in de rechterlijke orde'}, + {'mask': r'ordonnances?|beschikking(en)?', 'fr': 'ordonnance', 'nl': 'beschikking'}, + ], + r'banque nationale de belgique|nationale bank van belgie': [ + {'mask': r'notification|mededeling', 'fr': 'notification de la banque nationale', 'nl': 'mededeling van de nationale bank'}, + {'mask': r'autorisation.*cession|toestemming.*overdracht', 'fr': 'document de la banque nationale', 'nl': 'document van de nationale bank'}, + ], + r'service public federal mobilite et transports|federale overheidsdienst mobiliteit en vervoer': [ + {'mask': r'mobilite et securite routiere|mobiliteit en verkeersveiligheid', 'fr': 'document concernant la mobilité et la sécurite routière', 'nl': 'document betreffende mobiliteit en verkeersveiligheid'}, + ], + r'commission communautaire flamande de la region de bruxelles-capitale|vlaamse gemeenschapscommissie van het brussels hoofdstedelijk gewest': [ + {'mask': r'reglement|verordening', 'fr': 'règlement de la commission communautaire flamande', 'nl': 'verordening van de vlaamse gemeenschapscommissie'}, + ], +} + +TYPES = { + 'fr': [ + ('erratum', ['errat(?:um|a)']), + ('proces-verbal', ['proces-verbal']), + ("arrêt du conseil d'état", [r"conseil d(?:')?etat.? (?:- )?annulation", + r"annulation par le conseil d'?etat", + r"conseil d(?:')?etat.? (?:- )?arret", + r"conseil d(?:')?etat (?:- )?suspension partielle", + r"conseil d(?:')?etat.? (?:- )?suspension"]), + ('arrêt de la cour constitutionelle', [r'arret .*(?:- )?questions? prejudicielles?', + r'arret .*(?:- )?recours en annulation partielle', + r'arret .*(?:- )?recours en annulation', + r'arret .*(?:- )?demande de suspension']), + ("ordonnance de la cour d'appel", [r"cour d(?:')?appel .*(?:- )?ordonnance"]), + ('ordonnance du tribunal de commerce',[r'tribunal d(?:e|u) commerce .*(?:- )?ordonnance']), + ('ordonnance du tribunal du travail', [r'tribunal d(?:e|u) travail .*(?:- )?ordonnance']), + ('arrêté de la commission communautaire francaise', [r'arrete.*du college.*commission.*francaise']), + ('arrêté du gouvernement de la communauté germanophone', [r'arrete du gouvernement.*germanophone']), + ('arrêté de la commission communautaire commune', [r'arrete.*du college reuni.*commission communautaire commune']), + ('arrêté du gouvernement flamand', [r'arrete du gouvernement flamand']), + ('arrêté du gouvernement wallon', [r'arrete du gouvernement wallon']), + ('arrêté du gouvernement de la région de bruxelles-capitale', [r'arrete du gouvernement de la region de bruxelles-capitale']), + ('arrêté du gouvernement de la communauté francaise', [r'arrete du gouvernement de la communaute francaise']), + ('circulaire coordonnée', [r'circulaire coordonnee']), + ('arrêtés concernant les membres des commissions paritaires', [r'arretes concernant les membres des commissions paritaires']), + ("règlement d'application", [r"reglement d(?:')?application"]), + ('arrêté de la commission bancaire et financière', [r'arrete de la commission bancaire et financiere']), + ("décision du comité de ministres de l'union économique benelux", [r'decision du comite de ministres de lunion economique benelux']), + ('enquete(s) publique(s)', [r'enquete(?:s)? publique(?:s)?']), + ('accord international', [r'accord entre le royaume', r'convention entre la belgique', r'convention internationale']), + ('ordre du jour des séances plénières', [r'seances? plenieres? ordre du jour']), + ('agrément', [r"agrement d'?un expert", r'agrements?.*experts', r'agrements? en tant que', r"agrements?.*es ecoles", r'agrement.*laboratoires?']), + ('enregistrement', [r'enregistrements? en tant']), + ('approbations', [r'actes? dapprobation']), + ('code judiciaire', [r'code judiciaire']), + ('code pénal', [r'code penal']), + ('code des sociétes', [r'code des societes']), + ('code civil', [r'code civil']), + ('constitution reserve de recrutement', [r'(?:constitue|constitution) .* reserve de recrutement', r'examen de recrutement']), + ("communiqué d'etat", [r"communique d'?etat"]), + ('consulats étrangers', [r'consulats etrangers']), + ("règlement d'ordre interieur", [r"reglement d(?:')?ordre interieur"]), + ('arrêté royal', [r'arrete royal', r'par arretes royaux']), + ('arrêté ministeriel', [r'arretes? ministeriels?']), + ('demande de concession', [r'demande de concession']), + ('décision ministerielle', [r'decision (?:de la|du) ministre']), + ('accord de cooperation', [r'accord de cooperation']), + ('modification de la constitution', [r'modification a la constitution', r'revision de la constitution']), + ("arrêté d'application", [r"arrete d'application"]), + ('règlement spécial', [r'reglement special']), + ('décret spécial', [r'decret special']), + ('décret-programme', [r'decret-programme']), + ('décret', [r'decret']), + ('loi-programme', [r'loi-programme']), + ('arrêté-loi', [r'arrete-?loi']), + ("changement d'adresse", [r'changement dadresse']), + ('arrêt', [r'arret']), + ('ordonnance', [r'ordonnance']), + ('rapport', [r'rapport']), + ('circulaire', [r'circulaire']), + ('loi', [r'loi']), + ('nominations', [r'personnel.*nominations?']), + ('avis', [r'avis']), + ('jugement', [r'jugement']), + ('avenant', [r'avenant']), + ('règlement', [r'reglement']), + ('ratification', [r'ratification']), + ('adhésion', [r'adhesion']), + ('autorisation', [r'autorisation']), + ('liste', [r'liste des']), + ('protocole', [r'protocole']), + ("vacance d'emploi", [r"vacance(?:s) d'emploi(?:s)"]), + ('communication', [r'communication']), + ('recrutement', [r'recrutement']), + ('plan de secteur', [r'plan de secteur']), + ('remise de lettres de créance', [r'remise de lettres de creance']), + ('composition', [r'composition']), + ("journal officiel des communautés européennes", [r'journal officiel des communautes europeennes']), + ('indices du prix', [r'indices du prix']), + ("publication de l'ordre judiciaire", []), + ("document concernant l'ordre judiciaire", []), + ('document concernant les services judiciaires', []), + ('nomination par arrete royal', []), + ], + 'nl': [ + ('erratum', ['errat(?:um|a)']), + ('proces-verbaal', ['proces-verbaal']), + ('arrest van de raad van state', [r'raad van state.? (?:- )?vernietiging', + r'vernietiging door de raad van state', + r'raad van state.? (?:- )?arrest', + r'raad van state (?:- )?gedeeltelijke schorsing', + r'raad van state.? (?:- )?schorsing']), + ('arrest van het grondwettelijk hof', [r'arrest .*(?:- )?prejudiciele vraa?g(?:en)?', + r'arrest .*(?:- )?beroep(?:en)? tot gedeeltelijke vernietiging', + r'arrest .*(?:- )?beroep(?:en)? tot vernietiging', + r'arrest .*(?:- )?vordering(?:en)? tot schorsing']), + ('beschikking van het hof van beroep', [r'hof van beroep .*(?:- )?beschikking']), + ('beschikking van de rechtbank van koophandel', [r'rechtbank van koophandel .*(?:- )?beschikking']), + ('beschikking van de arbeidsrechtbank', [r'arbeidsrechtbank .*(?:- )?beschikking']), + ('besluit van de franse gemeenschapscommissie', [r'besluit.*van het college.*franse gemeenschapscommissie']), + ('besluit van de regering van de duitstalige gemeenschap', [r'besluit van de regering van de duitstalige gemeenschap']), + ('besluit van de gemeenschappelijke gemeenschapscommissie', [r'besluit.*van.*gemeenschappelijke gemeenschapscommissie']), + ('besluit van de vlaamse regering', [r'besluit van de vlaamse regering']), + ('besluit van de waalse regering', [r'besluit van de waalse regering']), + ('besluit van de brusselse hoofdstedelijke regering', [r'besluit van de brusselse hoofdstedelijke regering']), + ('besluit van de regering van de franse gemeenschap', [r'besluit (?:van de regering )?van de franse gemeenschap']), + ('gecoordineerde omzendbrief', [r'gecoordineerde omzendbrief']), + ('besluiten betreffende de leden van de paritaire comites', [r'besluiten betreffende de leden .* paritaire comites']), + ('toepassingsreglement', [r'toepassingsreglement']), + ('besluit van de commissie voor het bank- en financiewezen', [r'besluit van de commissie .*financiewezen']), + ('beschikking van het comite van ministers van de benelux economische unie', [r'beschikking van het comite van ministers van de benelux economische unie']), + ('publicatie(s) ter kritiek', [r'publicatie(?:s)? ter kritiek']), + ('internationale overeenkomst', [r'overeenkomst tussen het koninkrijk', r'overeenkomst tussen belgie', r'internationa(?:al|le) verdrag']), + ('agenda van de plenaire vergaderingen', [r'plenaire vergadering(?:en)? agenda']), + ('erkenning', [r'erkenning van een deskundige', r'erkenning.*deskundigen', r'erkenning(?:en)? als', r'erkenning.*scholen', r'erkenning.*laboratori(?:a|um)']), + ('registratie', [r'registratie als']), + ('goedkeuringen', [r'akten? tot goedkeuring']), + ('gerechtelijk wetboek', [r'gerechtelijk wetboek']), + ('strafwetboek', [r'strafwetboek']), + ('wetboek van vennootschappen', [r'wetboek van vennootschappen']), + ('burgerlijk wetboek', [r'burgerlijk wetboek']), + ('samenstelling wervingsreserve', [r'wervings?reserve', r'aanleggen van een personeelsreserve', r'wervingsexamen']), + ('mededeling van de Staat', [r'mededeling van de staat']), + ('buitenlandse consulaten', [r'buitenlandse consulaten']), + ('huishoudelijk reglement', [r'huishoudelijk reglement']), + ('koninklijk besluit', [r'koninklijk besluit', r'bij koninklijke besluiten']), + ('ministerieel besluit', [r'ministeriee?le? besluit(?:en)?']), + ('concessieaanvraag', [r'concessieaanvraag']), + ('ministeriele beslissing', [r'beslissing van de minister']), + ('samenwerkingsakkoord', [r'samenwerkingsakkoord']), + ('wijziging aan de grondwet', [r'wijziging aan de grondwet', r'herziening van de grondwet']), + ('toepassingsbesluit', [r'toepassingsbesluit']), + ('bijzonder reglement', [r'bijzonder reglement']), + ('bijzonder decreet', [r'bijzonder decreet']), + ('programmadecreet', [r'programmadecreet']), + ('decreet', [r'decreet']), + ('programmawet', [r'programmawet']), + ('besluitwet', [r'besluitwet']), + ('adreswijziging', [r'adreswijziging']), + ('arrest', [r'arrest']), + ('beschikking', [r'ordonnantie', r'beschikking']), + ('verslag', [r'verslag']), + ('omzendbrief', [r'omzendbrief|circulaire']), + ('wet', [r'wet']), + ('benoemingen', [r'personeel.*benoeming(?:en)?']), + ('bericht', [r'bericht', r'advies']), + ('vonnis', [r'vonnis']), + ('bijakte', [r'bijakte']), + ('overeenkomst', [r'overeenkomst']), + ('bekrachtiging', [r'bekrachtiging']), + ('toetreding', [r'toetreding']), + ('vergunning', [r'verordening', r'vergunning']), + ('lijst', [r'lijst van']), + ('protocol', [r'protocol']), + ('vacante bettreking', [r'vacante bettreking(?:en)']), + ('mededeling', [r'mededeling']), + ('aanwerving', [r'(?:aan)?werving']), + ('gewestplan', [r'gewestplan']), + ('overhandiging van geloofsbrieven',[r'overhandiging van geloofsbrieven']), + ('samenstelling', [r'samenstelling']), + ('publicatieblad van de europese gemeenschappen', [r'publicatieblad van de europese gemeenschappen']), + ('indexcijfers', [r'indexcijfers']), + ('bekendmaking van de juridische orde', []), + ('document betreffende de rechterlijke ordre', []), + ('document betreffende de gerechtelijke diensten', []), + ('benoeming door koninklijk besluit', []), + ], +} + + +# ── Page class ───────────────────────────────────────────────────────────────── + +class Page: + """ + Parses a raw HTML page from ejustice.just.fgov.be and exposes document + metadata via get_*() methods. Drop-in Python replacement for Page.pm. + """ + + def __init__(self, html, numac, pub_date, lang): + self.raw = html + self.numac = numac + self.pub_date = pub_date + self.lang = lang + self.soup = BeautifulSoup(html, 'html.parser') + # Normalized HTML: unaccented + lowercase (used for regex matching) + self.content = _norm_html(html) + + # ── Public getters ───────────────────────────────────────────────────────── + + def get_source(self): + if 'no article available with such references' in self.content: + return 'nosource' + h1 = self.soup.find('h1', class_='page__title') + if h1: + span = h1.find('span') + if span: + return span.get_text(strip=True) + return 'nosource' + + def get_title(self): + final_title = '' + intro = self.soup.find(class_='intro-text') + + if intro: + text = intro.get_text(strip=True) + # Strip leading date/preamble: "21 JUIN 2024. -" or "21/06/2024 -" + m = re.match(r'^(?:[^.]+\.)?\s*-?\s*(.+?)(?:\s*-\s*\(.*\))?\s*$', text, re.DOTALL) + if m and m.group(1).strip(): + final_title = m.group(1).strip() + + if not final_title: + intro_date = self.soup.find(class_='intro-date') + main_el = self.soup.find(attrs={'role': 'main'}) + sub_title = intro_date.get_text(strip=True) if intro_date else '' + text_begin = main_el.get_text()[:100] if main_el else '' + if sub_title: + final_title = f"{sub_title}(...) {text_begin}(...)" + + if 'SSFetch Text failed:' in self.raw: + logger.warning("SSFetch failed for %s (%s)", self.numac, self.lang) + elif not final_title and 'No article available with such references' not in self.raw: + logger.debug("Could not match title for %s (%s)", self.numac, self.lang) + + final_title = re.sub(r'\([^)]+\)', '', final_title) + return final_title.strip() + + def get_prom_date(self): + # 1) Try ELI URI: /eli/type/yyyy/mm/dd/ + m = re.search(r'/eli/[a-z0-9_\-./]+/(\d{4})/(\d{2})/(\d{2})/', self.content, re.IGNORECASE) + if m: + return f"{m.group(1)}-{m.group(2)}-{m.group(3)}" + + # 2) Try date in intro-text: "21 juin 2024." + _months = { + 'janvier': 1, 'fevrier': 2, 'mars': 3, 'avril': 4, 'mai': 5, 'juin': 6, + 'juillet': 7, 'aout': 8, 'septembre': 9, 'octobre': 10, 'novembre': 11, 'decembre': 12, + 'januari': 1, 'februari': 2, 'maart': 3, 'april': 4, 'mei': 5, 'juni': 6, + 'juli': 7, 'augustus': 8, 'september': 9, 'oktober': 10, 'november': 11, 'december': 12, + } + m = re.search(r'"intro-text">\s*(\d{1,2})\s+(\w+)\s+(\d{4})\.', self.content, re.IGNORECASE) + if m: + month = m.group(2).lower() + if month in _months: + return f"{m.group(3)}-{_months[month]:02}-{int(m.group(1)):02}" + return '--' + + def get_type(self): + title = normalize_title(self.get_title()) + title = re.sub(r'\s{2,}', ' ', title) + + if not title.strip(): + logger.debug("NO TITLE — default type for %s", self.numac) + return 'document' + + # Special case: vacancy + if re.search( + r'vacante betrekking(?:en)?|openstaande plaats|openstaande betrekking|vacature' + r'|oproep.*kandidaten|te begeven betrekkingen|recrutering' + r"|places? vacantes?|emplois? vacants?|vacance.+emploi|vacance.+mandat" + r"|appel aux candidats|emplois a conferer|mandat vacant", + title, + ): + return "vacance d'emploi" if self.lang == 'fr' else 'vacante bettreking' + + # Special case: royal appointment + if re.search( + r'benoeming(?:en)?.*bij +koninklijke? +besluit(?:en)?' + r'|nominations?.*par +arretes? +roya(?:l|ux)', + title, + ): + return 'benoeming door koninklijk besluit' if self.lang == 'nl' else 'nomination par arrete royal' + + # Source-driven classification + source = _strip_accents(self.get_source()).lower() + for source_mask, items in SOURCE_TYPES.items(): + if _matches_source(source, source_mask): + for item in items: + if _matches_doc_mask(title, item['mask']): + return item[self.lang] + + # General type matching — beginning of title + for term, masks in TYPES[self.lang]: + for mask in masks: + if _matches_at_start(title, mask): + return term + + # General type matching — anywhere in title + for term, masks in TYPES[self.lang]: + for mask in masks: + if _matches_anywhere(title, mask): + logger.debug("LONGF %s found '%s' anywhere in title", self.numac, term) + return term + + logger.debug("NO TYPE — no type found for %s (%s)", self.numac, self.lang) + return 'document' + + def get_eli(self): + if 'no article available with such references' in self.content: + return None + m = re.search(r'(/eli/[a-z0-9/_\-./]+)', self.raw, re.IGNORECASE) + return m.group(1) if m else None + + def get_eli_type(self): + if 'no article available with such references' in self.content: + return None + m = re.search(r'/eli/([a-z0-9_\-.]+)/[a-z0-9/_\-.]+', self.raw, re.IGNORECASE) + return m.group(1) if m else None + + def get_pdf(self): + if 'no article available with such references' in self.content: + return None + m = re.search(r'href="(/mopdf[^"]+)"', self.raw, re.IGNORECASE) + return m.group(1).strip() if m else None + + def get_chrono(self): + if 'no article available with such references' in self.content: + return None + m = re.search(r'"(http://reflex\.raadvst-consetat\.be/[^"]+)"', self.raw, re.IGNORECASE) + return m.group(1).strip() if m else None + + def get_chrono_data(self): + if 'no article available with such references' in self.content: + return None + m = re.search( + r'reflex\.raadvst-consetat\.be/reflex/\?page=chrono&c=detail_get&d=detail&docid=(\d+)&tab=chrono', + self.raw, re.IGNORECASE, + ) + return m.group(1) if m else None + + def get_chamber_data(self): + if 'no article available with such references' in self.content: + return (None, None) + m = re.search( + r'www\.dekamer\.be/kvvcr/showpage\.cfm\?section=flwb&language=\w{2}' + r'&cfm=/site/wwwcfm/flwb/flwbn\.cfm\?&dossierID=(\d+)&legislat=(\d+)', + self.raw, re.IGNORECASE, + ) + return (m.group(1), m.group(2)) if m else (None, None) + + def get_senate_data(self): + if 'no article available with such references' in self.content: + return (None, None) + m = re.search( + r'www\.senaa?t\.be/www/\?MIval=dossier&LEG=(\d+)&NR=(\d+)&LANG=', + self.raw, re.IGNORECASE, + ) + # Perl returns ($2, $1) i.e. (NR, LEG) + return (m.group(2), m.group(1)) if m else (None, None) + + # ── Type index helpers (used for cross-language reconciliation) ──────────── + + def type_index(self, value): + """Return 1-based index of value in TYPES[lang], or 0 for 'document'.""" + if value == 'document': + return 0 + for i, (term, _) in enumerate(TYPES[self.lang]): + if value == term: + return i + 1 + return 0 + + def get_type_by_index(self, index): + """Return the type term at 1-based index, or 'document'.""" + if index == 0 or index == 100: + return 'document' + types = TYPES[self.lang] + if 1 <= index <= len(types): + return types[index - 1][0] + return 'document' diff --git a/py_agent/pyproject.toml b/py_agent/pyproject.toml index c7e94d2..df5c5a4 100644 --- a/py_agent/pyproject.toml +++ b/py_agent/pyproject.toml @@ -11,7 +11,7 @@ beautifulsoup4 = "^4.12.3" requests = "^2.32.3" click = "^8.1.7" pymysql = "^1.1.1" - +aiohttp = "^3.9" [build-system] requires = ["poetry-core"]