rebuild source discovery workflows
CI / compose (push) Failing after 5m46s

This commit is contained in:
Marco0300
2026-09-04 20:52:37 +02:00
parent 3a9b553440
commit 6aefbdc1f3
18 changed files with 853 additions and 106 deletions
+81 -10
View File
@@ -9,12 +9,13 @@ from urllib.request import Request, urlopen
if __package__ in (None, ""):
sys.path.insert(0, str(Path(__file__).resolve().parents[1]))
from app.domain import deduplication_key, deduplicate_businesses, is_suppressed, normalize_business, normalize_domain, normalize_phone, match_businesses
from app.sources import adapter_for, contains_secret, available_adapters, normalize_record, circuit_is_open, DISCOVERY_CRITERIA_FIELDS
from app.sources import adapter_for, contains_secret, available_adapters, normalize_record, circuit_is_open, DISCOVERY_CRITERIA_FIELDS, GoogleBrowserSearchBlocked
from app.domain_intelligence import normalize_registrable_domain, resolve_domain, generate_candidate_domains
from app.website_scanner import scan_website, validate_url
from app.contact_extractor import extract_contacts, MAX_HTML_BYTES, MAX_RESULTS
from app.scoring import DEFAULT_RULES, signals_for_business, evaluate_score, score_business_opportunity, SCORE_VERSION
from app.ai_assistance import generate as generate_ai, input_fingerprint, provider_status, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS
from app.ai_assistance import generate as generate_ai, input_fingerprint, evidence_hashes, provider_status, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS
from app.ai_opportunity import assess_opportunity, DETERMINISTIC_ASSESSMENT_THRESHOLD, ASSESSMENT_SCHEMA_VERSION
from app.discovery import discover as scoped_discover
from app.ai_research import provider_status as ai_research_provider_status, configure_db as configure_ai_research_db, validate_criteria as validate_ai_research_criteria, AIResearchConfigError
from app.search_provider import provider_status as search_provider_status
@@ -22,12 +23,13 @@ if __package__ in (None, ""):
from app.provider_config import validate_payload as validate_remote_provider, encrypt as encrypt_provider_secret, decrypt as decrypt_provider_secret, safe_status as remote_provider_status, test_connectivity as test_remote_connectivity
else:
from .domain import deduplication_key, deduplicate_businesses, is_suppressed, normalize_business, normalize_domain, normalize_phone, match_businesses
from .sources import adapter_for, contains_secret, available_adapters, normalize_record, circuit_is_open, DISCOVERY_CRITERIA_FIELDS
from .sources import adapter_for, contains_secret, available_adapters, normalize_record, circuit_is_open, DISCOVERY_CRITERIA_FIELDS, GoogleBrowserSearchBlocked
from .domain_intelligence import normalize_registrable_domain, resolve_domain, generate_candidate_domains
from .website_scanner import scan_website, validate_url
from .contact_extractor import extract_contacts, MAX_HTML_BYTES, MAX_RESULTS
from .scoring import DEFAULT_RULES, signals_for_business, evaluate_score, score_business_opportunity, SCORE_VERSION
from .ai_assistance import generate as generate_ai, input_fingerprint, provider_status, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS
from .ai_assistance import generate as generate_ai, input_fingerprint, evidence_hashes, provider_status, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS
from .ai_opportunity import assess_opportunity, DETERMINISTIC_ASSESSMENT_THRESHOLD, ASSESSMENT_SCHEMA_VERSION
from .discovery import discover as scoped_discover
from .ai_research import provider_status as ai_research_provider_status, configure_db as configure_ai_research_db, validate_criteria as validate_ai_research_criteria, AIResearchConfigError
from .search_provider import provider_status as search_provider_status
@@ -98,13 +100,13 @@ def _initialize_database(db_path: str) -> sqlite3.Connection:
db.execute("INSERT OR IGNORE INTO organizations (id,name) VALUES (?,?)", (ORGANIZATION_ID, "Demo organization"))
source_sql_row=db.execute("SELECT sql FROM sqlite_master WHERE type='table' AND name='sources'").fetchone()
source_sql=(source_sql_row[0] or '') if source_sql_row else ''
if "CHECK(kind IN ('csv','manual'))" in source_sql:
if "google_browser_search" not in source_sql:
# SQLite cannot change foreign-key enforcement during a transaction.
db.commit(); db.execute("PRAGMA foreign_keys=OFF")
db.executescript("""
CREATE TABLE sources_rebuilt (
id INTEGER PRIMARY KEY AUTOINCREMENT, organization_id TEXT NOT NULL REFERENCES organizations(id),
name TEXT NOT NULL, kind TEXT NOT NULL CHECK(kind IN ('csv','manual','google_places','bing_local','approved_directory','public_website','permitted_social','ct_logs','dns','rdap')), source_code TEXT NOT NULL DEFAULT '', display_name TEXT NOT NULL DEFAULT '', enabled INTEGER NOT NULL DEFAULT 0, approved INTEGER NOT NULL DEFAULT 0,
name TEXT NOT NULL, kind TEXT NOT NULL CHECK(kind IN ('csv','manual','google_places','google_browser_search','bing_local','approved_directory','public_website','permitted_social','ct_logs','dns','rdap')), source_code TEXT NOT NULL DEFAULT '', display_name TEXT NOT NULL DEFAULT '', enabled INTEGER NOT NULL DEFAULT 0, approved INTEGER NOT NULL DEFAULT 0,
config_json TEXT NOT NULL DEFAULT '{}', policy_json TEXT NOT NULL DEFAULT '{}', quota_json TEXT NOT NULL DEFAULT '{}', health_status TEXT NOT NULL DEFAULT 'unknown',
consecutive_failures INTEGER NOT NULL DEFAULT 0, circuit_open INTEGER NOT NULL DEFAULT 0,
last_success_at TEXT, last_failure_at TEXT, last_error TEXT, created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP,
@@ -652,6 +654,50 @@ class ApiHandler(BaseHTTPRequestHandler):
row = db.execute("SELECT * FROM ai_runs WHERE id=? AND organization_id=?", (run_id, org)).fetchone()
return self.send_json(201, self._ai_run_json(row, output.get("suggestions", [])))
def assess_ai_opportunity(self, bid, payload, db, user):
"""Assess one operator-selected, tenant-scoped business; never contact it."""
org = user["organization_id"]
business = self.business(db, bid, org)
if not business:
return self.send_json(404, {"error": "not_found"})
configured_provider = provider_status()
if configured_provider["status"] != "ready":
return self.send_json(409, {"error": "ai_provider_not_configured", "business_id": bid,
"provider": configured_provider["provider"], "network_send": False,
"automatic_outreach": False})
try:
deterministic_score = int(business["score"] or 0)
except (TypeError, ValueError):
deterministic_score = 0
if deterministic_score < DETERMINISTIC_ASSESSMENT_THRESHOLD:
return self.send_json(409, {"error": "deterministic_threshold_not_met", "business_id": bid,
"score": deterministic_score, "threshold": DETERMINISTIC_ASSESSMENT_THRESHOLD,
"network_send": False})
if payload:
return self.send_json(400, {"error": "manual_selection_only"})
evidence = [dict(row) for row in db.execute(
"SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id LIMIT ?",
(bid, org, MAX_INPUT_ITEMS)
)]
suppressions = [dict(row) for row in db.execute(
"SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,)
)]
suppressed = is_suppressed(dict(business), suppressions)
assessment = assess_opportunity(row_json(business), evidence, suppressed=suppressed)
fingerprint = input_fingerprint(row_json(business), [], [], evidence)
metadata = {"assessment_type": "manual_selected_business", "deterministic_threshold": DETERMINISTIC_ASSESSMENT_THRESHOLD,
"input_fingerprint": fingerprint, "network_send": False, "automatic_outreach": False}
cur = db.execute(
"INSERT INTO ai_runs(organization_id,business_id,input_evidence_hashes_json,model,provider,version,prompt_metadata_json,data_minimization_json,status,approval_state,output_json,actor_user_id) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)",
(org, bid, json.dumps(evidence_hashes(evidence), sort_keys=True), "local-deterministic", "local", ASSESSMENT_SCHEMA_VERSION,
json.dumps({"request": "manual_selected_business"}, sort_keys=True), json.dumps(metadata, sort_keys=True),
"succeeded", "pending", json.dumps({"assessment": assessment}, sort_keys=True), user["id"])
)
self.audit(db, user, "ai.opportunity_assessed", str(cur.lastrowid)); db.commit()
return self.send_json(201, {"id": cur.lastrowid, "business_id": bid, "assessment": assessment,
"human_review_required": assessment["human_review_required"], "network_send": False,
"automatic_outreach": False})
def list_ai_runs(self, db, user, query):
try:
limit = int((query.get("page_size") or [50])[0]); offset = max(0, int((query.get("offset") or [0])[0]))
@@ -938,6 +984,18 @@ class ApiHandler(BaseHTTPRequestHandler):
criteria = payload.get("criteria", {}); seeds = payload.get("seed_urls")
selected_adapters = resolve_selected_source_codes(db, user["organization_id"], payload.get("selected_adapters", payload.get("source_ids", payload.get("sources", []))))
source_mode = bool(selected_adapters) and seeds is None
if source_mode:
marks = ",".join("?" for _ in selected_adapters)
rows = db.execute("SELECT * FROM sources WHERE organization_id=? AND (source_code IN (" + marks + ") OR kind IN (" + marks + "))", [user["organization_id"]] + selected_adapters + selected_adapters).fetchall()
ready_codes = set()
for row in rows:
code = row["source_code"] or row["kind"]
try: configuration = json.loads(row["config_json"] or "{}")
except (TypeError, ValueError): configuration = {}
if row["enabled"] and not row["circuit_open"] and adapter_for(code).validate_config(configuration).valid:
ready_codes.add(code)
if set(selected_adapters) - ready_codes:
return self.send_json(409, {"error": "selected_source_not_ready", "details": sorted(set(selected_adapters) - ready_codes)})
criteria_only = seeds is None and not source_mode
if not isinstance(criteria, dict): return self.send_json(400, {"error": "invalid_criteria"})
if not criteria_only and not source_mode and (not isinstance(seeds, list) or not seeds): return self.send_json(400, {"error": "seed_urls_required"})
@@ -1160,6 +1218,7 @@ class ApiHandler(BaseHTTPRequestHandler):
if path=="/api/v1/admin/ai-provider-config/test": return self.remote_ai_provider_config(db,user,connectivity=True)
if path=="/api/v1/businesses/bulk-review": return self.bulk_review(payload,db,user)
bits_ai=path.split("/")
if len(bits_ai)==7 and bits_ai[:4]==["","api","v1","businesses"] and bits_ai[5]=="ai" and bits_ai[6]=="opportunity-assessment": return self.assess_ai_opportunity(int(bits_ai[4]) if bits_ai[4].isdigit() else -1,payload,db,user)
if len(bits_ai)==7 and bits_ai[:4]==["","api","v1","businesses"] and bits_ai[5]=="ai" and bits_ai[6]=="suggest": return self.suggest_ai(int(bits_ai[4]) if bits_ai[4].isdigit() else -1,payload,db,user)
if len(bits_ai)==6 and bits_ai[:4]==["","api","v1","ai-runs"] and bits_ai[4].isdigit() and bits_ai[5] in {"approve","reject"}: return self.decide_ai(int(bits_ai[4]),bits_ai[5],db,user)
if path=="/api/v1/score-rules": return self.create_score_rule(payload,db,user)
@@ -1444,7 +1503,8 @@ class ApiHandler(BaseHTTPRequestHandler):
provider_kinds={'openstreetmap','wikidata','common_crawl'}
kind='approved_directory' if requested in provider_kinds else requested
source_code=str(config.get('provider','')).strip().lower() if requested == 'approved_directory' and str(config.get('provider','')).strip().lower() in {'openstreetmap','wikidata','common_crawl'} else requested
if not name or requested not in ('csv','manual','google_places','bing_local','approved_directory','openstreetmap','wikidata','common_crawl','public_website','permitted_social','ct_logs','dns','rdap') or not isinstance(config,dict):return self.send_json(400,{"error":"invalid_source"})
if not name or requested not in ('csv','manual','google_places','google_browser_search','bing_local','approved_directory','openstreetmap','wikidata','common_crawl','public_website','permitted_social','ct_logs','dns','rdap') or not isinstance(config,dict):return self.send_json(400,{"error":"invalid_source"})
if requested == 'google_browser_search' and bool(payload.get('enabled', False)) and os.environ.get('GOOGLE_BROWSER_SEARCH_ENABLED', '').strip().lower() != 'true': return self.send_json(409,{"error":"source_feature_disabled","detail":"GOOGLE_BROWSER_SEARCH_ENABLED=true is required"})
if contains_secret(config):return self.send_json(400,{"error":"secret_not_permitted"})
if any(field in config for field in DISCOVERY_CRITERIA_FIELDS):return self.send_json(400,{"error":"source_configuration_contains_criteria"})
try:
@@ -1465,6 +1525,7 @@ class ApiHandler(BaseHTTPRequestHandler):
if 'config' in payload:
config=payload.get('config')
if not isinstance(config,dict) or contains_secret(config): return self.send_json(400,{"error":"invalid_source_config"})
if any(field in config for field in DISCOVERY_CRITERIA_FIELDS): return self.send_json(400,{"error":"source_configuration_contains_criteria"})
adapter=adapter_for(source['source_code'] or source['kind']); validation=adapter.validate_config(config)
if not validation.valid: return self.send_json(400,{"error":"invalid_source_config","details":validation.errors})
db.execute("UPDATE sources SET config_json=?,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",(json.dumps(config,sort_keys=True),sid,user['organization_id']))
@@ -1478,6 +1539,7 @@ class ApiHandler(BaseHTTPRequestHandler):
except (TypeError,ValueError): config={}
metadata=next((item for item in available_adapters() if item['source_code']==(source['source_code'] or source['kind'])), {})
if not metadata.get('available', False): return self.send_json(409,{"error":"source_unavailable"})
if (source['source_code'] or source['kind']) == 'google_browser_search' and os.environ.get('GOOGLE_BROWSER_SEARCH_ENABLED', '').strip().lower() != 'true': return self.send_json(409,{"error":"source_feature_disabled","detail":"GOOGLE_BROWSER_SEARCH_ENABLED=true is required"})
validation=adapter.validate_config(config)
# A blank manual source is a deliberate staging point: the query or
# ingest payload can provide rows later. Other adapters must be ready
@@ -1676,10 +1738,19 @@ def _run_source_discovery(db, job, handler):
page=adapter_for(source["source_code"] or source["kind"]).discover(source_config_with_credentials(db, source, config), criteria=criteria, limits=limits)
except Exception as exc:
blocked_count+=1; failures=int(source["consecutive_failures"])+1
db.execute("UPDATE sources SET health_status='unhealthy',consecutive_failures=?,circuit_open=?,last_failure_at=CURRENT_TIMESTAMP,last_error=?,updated_at=CURRENT_TIMESTAMP WHERE id=?",(failures,int(circuit_is_open(failures)),str(exc)[:300],source["id"]))
detail=str(exc).lower(); code="SOURCE_NETWORK_ERROR" if any(token in detail for token in ("urlopen", "gaierror", "timed out", "temporary failure", "network is unreachable")) else "SOURCE_EXECUTION_FAILED"
handler.add_job_event(db,job["id"],org,"source.blocked",f"{source['display_name'] or source['kind']} could not be reached",10,code); continue
google_blocked=isinstance(exc, GoogleBrowserSearchBlocked)
db.execute("UPDATE sources SET health_status=?,consecutive_failures=?,circuit_open=?,last_failure_at=CURRENT_TIMESTAMP,last_error=?,updated_at=CURRENT_TIMESTAMP WHERE id=?",("blocked" if google_blocked else "unhealthy",failures,int(google_blocked or circuit_is_open(failures)),str(exc)[:300],source["id"]))
detail=str(exc).lower(); code="GOOGLE_BROWSER_BLOCKED" if google_blocked else ("SOURCE_NETWORK_ERROR" if any(token in detail for token in ("urlopen", "gaierror", "timed out", "temporary failure", "network is unreachable")) else "SOURCE_EXECUTION_FAILED")
handler.add_job_event(db,job["id"],org,"source.blocked",f"{source['display_name'] or source['kind']} could not be reached",10,code)
if google_blocked:
if run: db.execute("UPDATE discovery_runs SET lifecycle='failed',updated_at=CURRENT_TIMESTAMP WHERE id=?",(run["id"],))
raise RuntimeError(code)
continue
limit=min(max_records-total, per_run_limit, max(0,daily_limit-used_today))
if bool(payload.get("dry_run", False)):
preview_count=len(page.records[:limit]); total += preview_count
handler.add_job_event(db,job["id"],org,"source.preview",f"Validated {source['display_name'] or source['kind']}: {preview_count} candidate(s) would be collected",50)
continue
for record in page.records[:limit]:
raw=json.dumps(record,sort_keys=True,separators=(",",":")); digest=hashlib.sha256(raw.encode()).hexdigest(); normalized=normalize_business(normalize_record(record)); norm=json.dumps(normalized,sort_keys=True); nkey=hashlib.sha256(norm.encode()).hexdigest()
handler.add_job_event(db,job["id"],org,"source.raw_persisted",f"Persisting {source['kind']} record",20)