diff --git a/apps/api/app/ai_opportunity.py b/apps/api/app/ai_opportunity.py
new file mode 100644
index 0000000..3cec4fd
--- /dev/null
+++ b/apps/api/app/ai_opportunity.py
@@ -0,0 +1,157 @@
+"""Strict, evidence-grounded, review-only opportunity assessment normalization."""
+from __future__ import annotations
+
+from typing import Any, Callable
+
+ASSESSMENT_SCHEMA_VERSION = "opportunity-assessment-v3"
+DETERMINISTIC_ASSESSMENT_THRESHOLD = 70
+RECOMMENDATIONS = frozenset({"contact", "review", "low_priority", "do_not_contact", "insufficient_evidence"})
+PRIORITIES = frozenset({"high", "medium", "low"})
+WEBSITE_STATUSES = frozenset({"healthy", "outdated", "broken", "missing", "parked", "unknown"})
+DOMAIN_STATUSES = frozenset({"registered", "missing", "likely_available", "unknown"})
+CONTACT_TYPES = frozenset({"none", "general_business", "named_business", "free_mail", "unknown"})
+_ALLOWED_FIELDS = frozenset({
+ "opportunity_score", "confidence_score", "recommendation", "priority", "reasons", "missing_evidence",
+ "website_assessment", "domain_assessment", "contactability", "recommended_services",
+ "human_review_required", "evidence_references",
+})
+
+
+def _score(value: Any, *, confidence: bool = False) -> int:
+ if isinstance(value, bool) or not isinstance(value, (int, float)):
+ return 0
+ numeric = float(value)
+ if numeric != numeric or numeric in (float("inf"), float("-inf")):
+ return 0
+ if confidence and 0 <= numeric <= 1:
+ numeric *= 100
+ return max(0, min(100, int(round(numeric))))
+
+
+def _enum(value: Any, allowed: frozenset[str], default: str) -> str:
+ item = value.strip().lower() if isinstance(value, str) else ""
+ return item if item in allowed else default
+
+
+def _text_list(value: Any) -> list[str]:
+ if value is None:
+ return []
+ if not isinstance(value, list) or len(value) > 20 or any(not isinstance(item, str) for item in value):
+ raise ValueError("invalid_assessment_list")
+ result: list[str] = []
+ for item in value:
+ item = item.strip()
+ if not item or len(item) > 300 or item in result:
+ continue
+ result.append(item)
+ return result
+
+
+def _website(value: Any) -> dict[str, Any]:
+ default = {"status": "unknown", "broken": False, "outdated": False, "mobile_issue": False, "https_issue": False, "performance_issue": False}
+ if value is None:
+ return default
+ if not isinstance(value, dict) or set(value) - set(default):
+ raise ValueError("invalid_assessment_schema")
+ result = dict(default)
+ result["status"] = _enum(value.get("status"), WEBSITE_STATUSES, "unknown")
+ for key in set(default) - {"status"}:
+ if key in value:
+ if not isinstance(value[key], bool):
+ raise ValueError("invalid_assessment_schema")
+ result[key] = value[key]
+ return result
+
+
+def _domain(value: Any) -> dict[str, str]:
+ if value is None:
+ return {"status": "unknown"}
+ if not isinstance(value, dict) or set(value) != {"status"}:
+ raise ValueError("invalid_assessment_schema")
+ return {"status": _enum(value.get("status"), DOMAIN_STATUSES, "unknown")}
+
+
+def _contactability(value: Any) -> dict[str, Any]:
+ default = {"public_business_contact_found": False, "contact_type": "unknown"}
+ if value is None:
+ return default
+ if not isinstance(value, dict) or set(value) - set(default):
+ raise ValueError("invalid_assessment_schema")
+ result = dict(default)
+ if "public_business_contact_found" in value:
+ if not isinstance(value["public_business_contact_found"], bool):
+ raise ValueError("invalid_assessment_schema")
+ result["public_business_contact_found"] = value["public_business_contact_found"]
+ result["contact_type"] = _enum(value.get("contact_type"), CONTACT_TYPES, "unknown")
+ return result
+
+
+def normalize_assessment(raw: dict[str, Any] | None, known_evidence_ids: set[int], *, suppressed: bool = False) -> dict[str, Any]:
+ """Return exactly the assessment contract; reject invented evidence IDs."""
+ raw = {} if raw is None else raw
+ if not isinstance(raw, dict) or set(raw) - _ALLOWED_FIELDS:
+ raise ValueError("invalid_assessment_schema")
+ references = raw.get("evidence_references", [])
+ if not isinstance(references, list) or len(references) > 100:
+ raise ValueError("invalid_evidence_references")
+ evidence_references: list[int] = []
+ for reference in references:
+ if isinstance(reference, bool) or not isinstance(reference, int):
+ raise ValueError("invalid_evidence_reference")
+ if reference not in known_evidence_ids:
+ raise ValueError("unknown_evidence_reference")
+ if reference not in evidence_references:
+ evidence_references.append(reference)
+ evidence_references.sort()
+ confidence_score = _score(raw.get("confidence_score"), confidence=True)
+ recommendation = _enum(raw.get("recommendation"), RECOMMENDATIONS, "insufficient_evidence")
+ weak_evidence = len(evidence_references) < 2 or confidence_score < 70 or recommendation == "insufficient_evidence"
+ contactability = _contactability(raw.get("contactability"))
+ if suppressed:
+ recommendation = "do_not_contact"
+ contactability = {"public_business_contact_found": False, "contact_type": "none"}
+ return {
+ "opportunity_score": _score(raw.get("opportunity_score")),
+ "confidence_score": confidence_score,
+ "recommendation": recommendation,
+ "priority": _enum(raw.get("priority"), PRIORITIES, "low"),
+ "reasons": _text_list(raw.get("reasons")),
+ "missing_evidence": _text_list(raw.get("missing_evidence")),
+ "website_assessment": _website(raw.get("website_assessment")),
+ "domain_assessment": _domain(raw.get("domain_assessment")),
+ "contactability": contactability,
+ "recommended_services": _text_list(raw.get("recommended_services")),
+ "human_review_required": bool(suppressed or weak_evidence or raw.get("human_review_required", False)),
+ "evidence_references": evidence_references,
+ }
+
+
+def deterministic_assessment(business: dict[str, Any], evidence: list[dict[str, Any]]) -> dict[str, Any]:
+ try:
+ score = max(0, min(100, int(business.get("score", 0) or 0)))
+ except (TypeError, ValueError):
+ score = 0
+ references = [item["id"] for item in evidence if isinstance(item.get("id"), int)][:100]
+ has_website = bool(business.get("website") or business.get("website_domain"))
+ website_class = str(business.get("website_class", "")).lower()
+ website_status = "missing" if not has_website else website_class if website_class in WEBSITE_STATUSES else "unknown"
+ has_contact = bool(business.get("email") or business.get("phone"))
+ return {
+ "opportunity_score": score,
+ "confidence_score": min(95, 35 + 20 * len(references)),
+ "recommendation": "review" if references and score >= DETERMINISTIC_ASSESSMENT_THRESHOLD else "insufficient_evidence",
+ "priority": "high" if score >= 70 else "medium" if score >= 40 else "low",
+ "reasons": ["Stored evidence requires human review."] if references else [],
+ "missing_evidence": [item for item, present in (("website evidence", has_website), ("corroborating evidence", len(references) >= 2)) if not present],
+ "website_assessment": {"status": website_status, "broken": website_status == "broken", "outdated": website_status == "outdated", "mobile_issue": False, "https_issue": has_website and not str(business.get("website", "")).startswith("https://"), "performance_issue": False},
+ "domain_assessment": {"status": "registered" if business.get("website_domain") else "missing"},
+ "contactability": {"public_business_contact_found": has_contact, "contact_type": "general_business" if has_contact else "none"},
+ "recommended_services": ["website" if not has_website else "website_repair"],
+ "human_review_required": True,
+ "evidence_references": references,
+ }
+
+
+def assess_opportunity(business: dict[str, Any], evidence: list[dict[str, Any]], *, suppressed: bool = False, provider: Callable[[dict[str, Any], list[dict[str, Any]],], dict[str, Any]] | None = None) -> dict[str, Any]:
+ raw = provider(business, evidence) if provider else deterministic_assessment(business, evidence)
+ return normalize_assessment(raw, {item["id"] for item in evidence if isinstance(item.get("id"), int)}, suppressed=suppressed)
diff --git a/apps/api/app/ai_research.py b/apps/api/app/ai_research.py
index 33371e2..02d8883 100644
--- a/apps/api/app/ai_research.py
+++ b/apps/api/app/ai_research.py
@@ -66,7 +66,10 @@ def _config():
nous_key = credentials.get("step_api_key", credentials.get("nous_api_key", "")) if provider in STEPFUN_PROVIDER_IDS else credentials.get("nous_api_key", "")
return {"provider": provider, "model": row["model"], "nous_url": row["nous_base_url"], "nous_allowed": {urlparse(row["nous_base_url"]).hostname}, "nous_key": nous_key, "searxng_url": row["firecrawl_base_url"] if row["firecrawl_base_url"].startswith("http://searxng") else os.environ.get("SEARXNG_BASE_URL", "").strip(), "searxng_allowed": _hosts("SEARXNG_ALLOWED_HOSTS", "searxng"), "firecrawl_url": row["firecrawl_base_url"], "firecrawl_allowed": {urlparse(row["firecrawl_base_url"]).hostname}, "firecrawl_key": credentials.get("firecrawl_api_key", "")}
except Exception:
- return {"provider": "", "model": "", "endpoint": "", "allowed": set(), "api_key": ""}
+ # A removed/legacy database must not poison subsequent server or test
+ # contexts. Fall back to the explicit environment configuration, which
+ # is still validated fail-closed by _endpoint().
+ pass
provider = os.environ.get("AI_RESEARCH_PROVIDER", "").strip().lower()
# Nous uses its conventional key directly; no gateway or key translation is needed.
nous_key = os.environ.get("NOUS_API_KEY", "").strip()
diff --git a/apps/api/app/main.py b/apps/api/app/main.py
index 6f0b840..83368ca 100644
--- a/apps/api/app/main.py
+++ b/apps/api/app/main.py
@@ -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)
diff --git a/apps/api/app/sources.py b/apps/api/app/sources.py
index a0c286c..73b2a7d 100644
--- a/apps/api/app/sources.py
+++ b/apps/api/app/sources.py
@@ -7,7 +7,8 @@ Adapters never emit or persist credential values.
from __future__ import annotations
from dataclasses import dataclass, field
from typing import Any, Mapping, Protocol, Sequence
-import csv, io, random, time, json, re
+from html.parser import HTMLParser
+import csv, io, os, random, threading, time, json, re
from urllib.parse import urlencode, urlparse
from urllib.request import Request, urlopen
@@ -328,9 +329,22 @@ class ApprovedDirectorySource(GatedSource):
def discover(self, config, cursor=None, criteria=None, limits=None):
result=self.validate_config(config)
if not result.valid: raise ValueError(result.errors[0])
- provider=str(config["provider"]).lower(); query=str(config["query"]).strip(); limit=max(1,min(100,int(config.get("max_records",50))))
+ criteria = criteria if isinstance(criteria, Mapping) else {}
+ keyword_values = criteria.get("keywords", criteria.get("keyword", criteria.get("query", "")))
+ if isinstance(keyword_values, str): keyword_values = [keyword_values]
+ terms = [str(value).strip() for value in keyword_values[:10] if str(value).strip()] if isinstance(keyword_values, Sequence) and not isinstance(keyword_values, (bytes, bytearray, str)) else []
+ for key in ("category", "industry"):
+ value = str(criteria.get(key, "")).strip()
+ if value: terms.append(value)
+ query = " ".join(dict.fromkeys(terms))[:500]
+ if not query: raise ValueError("discovery_criteria_required")
+ location_parts = [str(criteria.get(key, "")).strip() for key in ("city", "location", "province", "country")]
+ area = next((value for value in location_parts if value), "South Africa")
+ limit_source = limits if isinstance(limits, Mapping) else {}
+ try: limit=max(1,min(100,int(limit_source.get("per_run_limit", limit_source.get("max_records",50)))))
+ except (TypeError, ValueError): raise ValueError("invalid_limits")
+ provider=str(config["provider"]).lower()
if provider == "openstreetmap":
- area=str(config.get("location", "South Africa")).strip()
terms=[token.lower() for token in re.findall(r"[A-Za-z0-9]{2,32}",query)[:5]]
variants=sorted({variant for term in terms for variant in (term,term[:-1] if term.endswith('s') and len(term)>3 else term)})
pattern="|".join(re.escape(term) for term in variants)
@@ -352,6 +366,158 @@ class ApprovedDirectorySource(GatedSource):
payload, _ = _HttpJsonSource()._get_json(index); records=[normalize_record({"name":str(x.get("url","")).split('/')[2] if '://' in str(x.get("url","")) else x.get("url", ""),"website":x.get("url","")}) for x in (payload if isinstance(payload,list) else [])]
return DiscoveryPage(records[:limit],metadata={"adapter":self.source_code,"provider":provider,"record_count":len(records)})
+class GoogleBrowserSearchBlocked(RuntimeError):
+ """Fail-closed result for a disabled, rate-limited, or Google-blocked fetch."""
+
+ code = "GOOGLE_BROWSER_BLOCKED"
+
+ def __init__(self, reason: str):
+ self.reason = reason
+ super().__init__(f"{self.code}:{reason}")
+
+ def as_dict(self) -> dict[str, str]:
+ return {"code": self.code, "reason": self.reason}
+
+
+class _GoogleVisibleResults(HTMLParser):
+ """Extract only human-visible heading links from public result HTML."""
+
+ def __init__(self):
+ super().__init__(convert_charrefs=True)
+ self._href = ""
+ self._depth = 0
+ self._parts: list[str] = []
+ self.results: list[tuple[str, str]] = []
+
+ def handle_starttag(self, tag, attrs):
+ attributes = dict(attrs)
+ if tag == "a" and not self._href:
+ self._href = str(attributes.get("href") or "")
+ if tag == "h3" and self._href:
+ self._depth = 1
+ self._parts = []
+ elif self._depth:
+ self._depth += 1
+
+ def handle_data(self, data):
+ if self._depth:
+ self._parts.append(data)
+
+ def handle_endtag(self, tag):
+ if not self._depth:
+ if tag == "a":
+ self._href = ""
+ return
+ self._depth -= 1
+ if tag != "h3" or self._depth:
+ return
+ title = " ".join("".join(self._parts).split())[:300]
+ href = self._href
+ self._href = ""
+ self._parts = []
+ if title and href:
+ self.results.append((title, href))
+
+
+class GoogleBrowserSearchSource(GatedSource):
+ """Experimental, feature-flagged public Google result-page connector.
+
+ This connector only requests the public result HTML. It does not use a
+ browser profile, JavaScript execution, login, proxy, CAPTCHA solver, or
+ alternative endpoint when Google blocks access.
+ """
+
+ kind = source_code = "google_browser_search"
+ display_name = "Google Browser Search (experimental)"
+ available = True
+ optional = True
+ requires_credentials = False
+ _last_request_at: float | None = None
+ _rate_lock = threading.Lock()
+ max_response_bytes = 512 * 1024
+ timeout = 10
+ max_results = 10
+
+ def validate_config(self, config):
+ result = super().validate_config(config)
+ if not result.valid:
+ return result
+ rate = config.get("rate_limit")
+ if not isinstance(rate, int) or isinstance(rate, bool) or not 1 <= rate <= 12:
+ return ValidationResult(False, ["rate_limit must be an integer from 1 to 12 requests per minute"])
+ return ValidationResult(True)
+
+ @staticmethod
+ def _query(criteria: Mapping[str, Any]) -> str:
+ if not isinstance(criteria, Mapping):
+ raise ValueError("criteria must be an object")
+ values: list[str] = []
+ keywords = criteria.get("keywords", criteria.get("keyword", criteria.get("query", "")))
+ if isinstance(keywords, str):
+ keywords = [keywords]
+ if isinstance(keywords, Sequence) and not isinstance(keywords, (bytes, bytearray, str)):
+ values.extend(str(value).strip() for value in keywords[:10] if str(value).strip())
+ for key in ("category", "industry", "city", "location", "province", "country"):
+ value = str(criteria.get(key, "")).strip()
+ if value:
+ values.append(value)
+ query = " ".join(values)
+ if not query or len(query) > 500 or any(char in query for char in "\r\n"):
+ raise ValueError("bounded discovery criteria are required")
+ return query
+
+ @staticmethod
+ def _blocked(html: str) -> bool:
+ lowered = html.lower()
+ markers = ("our systems have detected unusual traffic", "recaptcha", "captcha", "automated queries", "access denied", "sorry...")
+ return any(marker in lowered for marker in markers)
+
+ def discover(self, config, cursor=None, criteria=None, limits=None):
+ result = self.validate_config(config)
+ if not result.valid:
+ raise ValueError(result.errors[0])
+ if os.environ.get("GOOGLE_BROWSER_SEARCH_ENABLED", "").strip().lower() != "true":
+ raise GoogleBrowserSearchBlocked("feature_disabled")
+ query = self._query(criteria or {})
+ limits = limits if isinstance(limits, Mapping) else {}
+ try:
+ requested = int(limits.get("per_run_limit", limits.get("max_records", self.max_results)))
+ except (TypeError, ValueError):
+ raise ValueError("invalid_limits")
+ count = max(1, min(self.max_results, requested))
+ interval = 60.0 / int(config["rate_limit"])
+ with self._rate_lock:
+ now = time.monotonic()
+ if self._last_request_at is not None and now - self._last_request_at < interval:
+ raise GoogleBrowserSearchBlocked("rate_limited")
+ self.__class__._last_request_at = now
+ url = "https://www.google.com/search?" + urlencode({"q": query, "num": count, "hl": str(criteria.get("language", "en"))[:12] or "en"})
+ request = Request(url, headers={"User-Agent": "ProspectPlatform/0.1 public-search (no-login; experimental)", "Accept": "text/html,application/xhtml+xml"})
+ try:
+ with urlopen(request, timeout=self.timeout) as response:
+ html = response.read(self.max_response_bytes + 1)
+ except Exception as exc:
+ raise GoogleBrowserSearchBlocked("access_denied") from exc
+ if len(html) > self.max_response_bytes:
+ raise GoogleBrowserSearchBlocked("response_too_large")
+ text = html.decode("utf-8", "replace")
+ if self._blocked(text):
+ raise GoogleBrowserSearchBlocked("google_challenge_or_denial")
+ parser = _GoogleVisibleResults()
+ parser.feed(text)
+ records, seen = [], set()
+ for title, href in parser.results:
+ parsed = urlparse(href)
+ host = (parsed.hostname or "").lower()
+ if parsed.scheme not in {"http", "https"} or not host or host.endswith("google.com") or href in seen:
+ continue
+ seen.add(href)
+ records.append(normalize_record({"name": title, "website": href, "description": "Public Google search result"}))
+ if len(records) >= count:
+ break
+ return DiscoveryPage(records, metadata={"adapter": self.source_code, "experimental": True, "public_html_only": True, "record_count": len(records), "query": query})
+
+
class GooglePlacesSource(GatedSource):
kind = source_code = "google_places"
display_name = "Google Places"
@@ -408,7 +574,7 @@ def _gated(code, name):
BingLocalSource = _gated("bing_local", "Bing / approved local API")
PermittedSocialSource = _gated("permitted_social", "Permitted social")
-ADAPTERS = {x.source_code: x for x in (ManualSource, CsvSource, GooglePlacesSource, BingLocalSource, ApprovedDirectorySource, OpenStreetMapSource, WikidataSource, CommonCrawlSource, PublicWebsiteSource, PermittedSocialSource, CtLogsSource, DnsSource, RdapSource)}
+ADAPTERS = {x.source_code: x for x in (ManualSource, CsvSource, GooglePlacesSource, GoogleBrowserSearchSource, BingLocalSource, ApprovedDirectorySource, OpenStreetMapSource, WikidataSource, CommonCrawlSource, PublicWebsiteSource, PermittedSocialSource, CtLogsSource, DnsSource, RdapSource)}
# common aliases used by clients
ADAPTER_REGISTRY = ADAPTERS
diff --git a/apps/api/schema.sql b/apps/api/schema.sql
index e2e241c..bfea3af 100644
--- a/apps/api/schema.sql
+++ b/apps/api/schema.sql
@@ -136,7 +136,7 @@ CREATE INDEX IF NOT EXISTS idx_job_events_job_sequence ON job_events(job_id,sequ
-- Phase 5 source framework (additive-safe; credentials contain metadata only).
CREATE TABLE IF NOT EXISTS sources (
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,
diff --git a/apps/api/tests/test_ai_opportunity.py b/apps/api/tests/test_ai_opportunity.py
new file mode 100644
index 0000000..edef1cf
--- /dev/null
+++ b/apps/api/tests/test_ai_opportunity.py
@@ -0,0 +1,170 @@
+import json
+import os
+import sqlite3
+import threading
+import unittest
+from http.client import HTTPConnection
+from tempfile import TemporaryDirectory
+from unittest.mock import Mock
+
+from app.ai_opportunity import normalize_assessment
+from app.ai_research import configure_db as configure_ai_research_db
+from app.main import create_server, hash_password
+
+
+class OpportunityNormalizationTests(unittest.TestCase):
+ def test_normalizes_scores_enums_and_known_evidence_only(self):
+ provider = Mock(return_value={
+ "opportunity_score": 81,
+ "confidence_score": 0.82,
+ "recommendation": "contact",
+ "priority": "high",
+ "reasons": ["Two independent public listings corroborate the business."],
+ "missing_evidence": [],
+ "website_assessment": {"status": "healthy", "broken": False, "outdated": False, "mobile_issue": False, "https_issue": False, "performance_issue": False},
+ "domain_assessment": {"status": "registered"},
+ "contactability": {"public_business_contact_found": True, "contact_type": "general_business"},
+ "recommended_services": ["seo"],
+ "evidence_references": [2, 1, 2],
+ "human_review_required": False,
+ })
+ result = normalize_assessment(provider(), {1, 2})
+ self.assertEqual(result["opportunity_score"], 81)
+ self.assertEqual(result["confidence_score"], 82)
+ self.assertEqual(result["evidence_references"], [1, 2])
+ self.assertFalse(result["human_review_required"])
+ self.assertEqual(result["recommendation"], "contact")
+ self.assertEqual(result["reasons"], ["Two independent public listings corroborate the business."])
+ self.assertEqual(result["missing_evidence"], [])
+ self.assertEqual(result["website_assessment"]["status"], "healthy")
+ self.assertEqual(result["domain_assessment"]["status"], "registered")
+ self.assertEqual(result["contactability"]["contact_type"], "general_business")
+ self.assertEqual(result["recommended_services"], ["seo"])
+
+ def test_contact_recommendation_is_an_internal_no_send_recommendation(self):
+ payload = {"recommendation": "contact", "evidence_references": [1, 2], "confidence_score": 90}
+ self.assertEqual(normalize_assessment(payload, {1, 2})["recommendation"], "contact")
+
+ def test_unknown_values_default_and_weak_evidence_requires_human_review(self):
+ result = normalize_assessment({
+ "opportunity_score": "not-a-number",
+ "confidence_score": 20,
+ "recommendation": "email_them_now",
+ "priority": "urgent",
+ "evidence_references": [],
+ "human_review_required": False,
+ }, set())
+ self.assertEqual(result["opportunity_score"], 0)
+ self.assertEqual(result["confidence_score"], 20)
+ self.assertEqual(result["recommendation"], "insufficient_evidence")
+ self.assertEqual(result["priority"], "low")
+ self.assertTrue(result["human_review_required"])
+
+ def test_unknown_evidence_reference_is_rejected(self):
+ with self.assertRaisesRegex(ValueError, "unknown_evidence_reference"):
+ normalize_assessment({"evidence_references": [99]}, {1})
+
+ def test_suppression_overrides_provider_recommendation(self):
+ result = normalize_assessment({
+ "opportunity_score": 99,
+ "confidence_score": 99,
+ "recommendation": "contact",
+ "priority": "high",
+ "evidence_references": [1, 2],
+ "human_review_required": False,
+ }, {1, 2}, suppressed=True)
+ self.assertEqual(result["recommendation"], "do_not_contact")
+ self.assertTrue(result["human_review_required"])
+
+
+class OpportunityAssessmentApiTests(unittest.TestCase):
+ def setUp(self):
+ self.tmp = TemporaryDirectory()
+ self.old = {key: os.environ.get(key) for key in ("AI_PROVIDER", "BOOTSTRAP_ADMIN_EMAIL", "BOOTSTRAP_ADMIN_PASSWORD")}
+ os.environ.update({"AI_PROVIDER": "local", "BOOTSTRAP_ADMIN_EMAIL": "opportunity-owner@example.test", "BOOTSTRAP_ADMIN_PASSWORD": "password"})
+ self.server = create_server("127.0.0.1", 0, self.tmp.name + "/opportunity.db")
+ self.thread = threading.Thread(target=self.server.serve_forever, daemon=True); self.thread.start()
+ self.conn = HTTPConnection("127.0.0.1", self.server.server_port, timeout=3); self.cookie = None
+ self.request("POST", "/api/v1/auth/login", {"email": "opportunity-owner@example.test", "password": "password"})
+
+ def tearDown(self):
+ self.server.shutdown(); self.server.server_close(); self.thread.join(timeout=2)
+ configure_ai_research_db("")
+ self.tmp.cleanup()
+ for key, value in self.old.items():
+ if value is None: os.environ.pop(key, None)
+ else: os.environ[key] = value
+
+ def request(self, method, path, payload=None):
+ headers = {"Content-Type": "application/json"}
+ if self.cookie: headers["Cookie"] = self.cookie
+ self.conn.request(method, path, json.dumps(payload).encode() if payload is not None else None, headers)
+ response = self.conn.getresponse(); cookie = response.getheader("Set-Cookie")
+ if cookie: self.cookie = cookie.split(";", 1)[0]
+ return response.status, json.loads(response.read() or b"{}")
+
+ def selected_business_at_threshold(self):
+ status, business = self.request("POST", "/api/v1/businesses", {"name": "Selected Co"})
+ self.assertEqual(status, 201)
+ db = sqlite3.connect(self.tmp.name + "/opportunity.db")
+ db.execute("UPDATE businesses SET score=70 WHERE id=?", (business["id"],)); db.commit(); db.close()
+ return business["id"]
+
+ def test_manual_selected_business_after_threshold_returns_grounded_review_only_assessment(self):
+ bid = self.selected_business_at_threshold()
+ self.request("POST", f"/api/v1/businesses/{bid}/evidence", {"kind": "source", "url": "https://source.test", "claim": "Needs a modern website"})
+ status, result = self.request("POST", f"/api/v1/businesses/{bid}/ai/opportunity-assessment", {})
+ self.assertEqual(status, 201)
+ self.assertEqual(result["business_id"], bid)
+ self.assertEqual(result["assessment"]["evidence_references"], [1])
+ self.assertIn("opportunity_score", result["assessment"])
+ self.assertIn("confidence_score", result["assessment"])
+ self.assertTrue(result["assessment"]["human_review_required"])
+ self.assertFalse(result["network_send"])
+ self.assertFalse(result["automatic_outreach"])
+ status, runs = self.request("GET", "/api/v1/ai-runs")
+ self.assertEqual(status, 200)
+ self.assertEqual(runs["items"][0]["output"]["assessment"], result["assessment"])
+ self.assertEqual(len(runs["items"][0]["input_evidence_hashes"]), 1)
+
+ def test_manual_assessment_requires_deterministic_threshold(self):
+ status, business = self.request("POST", "/api/v1/businesses", {"name": "Below threshold"})
+ self.assertEqual(status, 201)
+ status, result = self.request("POST", f"/api/v1/businesses/{business['id']}/ai/opportunity-assessment", {})
+ self.assertEqual(status, 409)
+ self.assertEqual(result["error"], "deterministic_threshold_not_met")
+
+ def test_unconfigured_provider_fails_closed_without_assessment_output(self):
+ bid = self.selected_business_at_threshold()
+ with unittest.mock.patch("app.main.provider_status", return_value={"status": "not_configured", "provider": ""}):
+ status, result = self.request("POST", f"/api/v1/businesses/{bid}/ai/opportunity-assessment", {})
+ self.assertEqual(status, 409)
+ self.assertEqual(result["error"], "ai_provider_not_configured")
+ self.assertFalse(result["network_send"])
+ self.assertNotIn("assessment", result)
+
+ def test_suppressed_selected_business_assessment_is_do_not_contact(self):
+ bid = self.selected_business_at_threshold()
+ self.request("POST", f"/api/v1/businesses/{bid}/evidence", {"kind": "source", "url": "https://source.test", "claim": "Evidence"})
+ self.request("POST", "/api/v1/suppressions", {"kind": "domain", "value": "selected.test"})
+ db = sqlite3.connect(self.tmp.name + "/opportunity.db")
+ db.execute("UPDATE businesses SET website_domain='selected.test' WHERE id=?", (bid,)); db.commit(); db.close()
+ status, result = self.request("POST", f"/api/v1/businesses/{bid}/ai/opportunity-assessment", {})
+ self.assertEqual(status, 201)
+ self.assertEqual(result["assessment"]["recommendation"], "do_not_contact")
+ self.assertTrue(result["assessment"]["human_review_required"])
+
+ def test_manual_assessment_is_tenant_scoped(self):
+ bid = self.selected_business_at_threshold()
+ password_hash, password_salt = hash_password("other-password")
+ db = sqlite3.connect(self.tmp.name + "/opportunity.db")
+ db.execute("INSERT INTO organizations(id,name) VALUES(?,?)", ("other-tenant", "Other"))
+ db.execute("INSERT INTO users(organization_id,email,password_hash,password_salt,role) VALUES(?,?,?,?,?)", ("other-tenant", "other-opportunity@example.test", password_hash, password_salt, "owner"))
+ db.commit(); db.close()
+ self.cookie = None
+ self.assertEqual(self.request("POST", "/api/v1/auth/login", {"email": "other-opportunity@example.test", "password": "other-password"})[0], 200)
+ self.assertEqual(self.request("POST", f"/api/v1/businesses/{bid}/ai/opportunity-assessment", {})[0], 404)
+
+
+if __name__ == "__main__":
+ unittest.main()
diff --git a/apps/api/tests/test_google_browser_search.py b/apps/api/tests/test_google_browser_search.py
new file mode 100644
index 0000000..2b5b129
--- /dev/null
+++ b/apps/api/tests/test_google_browser_search.py
@@ -0,0 +1,59 @@
+import os
+import unittest
+from unittest.mock import patch
+
+from app.sources import GoogleBrowserSearchBlocked, GoogleBrowserSearchSource
+
+
+class GoogleBrowserSearchSourceTests(unittest.TestCase):
+ def setUp(self):
+ self.old_enabled = os.environ.pop("GOOGLE_BROWSER_SEARCH_ENABLED", None)
+
+ def tearDown(self):
+ GoogleBrowserSearchSource._last_request_at = None
+ if self.old_enabled is None:
+ os.environ.pop("GOOGLE_BROWSER_SEARCH_ENABLED", None)
+ else:
+ os.environ["GOOGLE_BROWSER_SEARCH_ENABLED"] = self.old_enabled
+
+ def test_disabled_feature_blocks_without_network_io(self):
+ source = GoogleBrowserSearchSource()
+ with patch("app.sources.urlopen") as network:
+ with self.assertRaises(GoogleBrowserSearchBlocked) as raised:
+ source.discover(
+ {"approved": True, "public_access": True, "terms_accepted": True, "rate_limit": 6},
+ criteria={"keywords": ["solar installers"], "city": "Cape Town"},
+ limits={"max_records": 5},
+ )
+ self.assertEqual(raised.exception.code, "GOOGLE_BROWSER_BLOCKED")
+ self.assertEqual(raised.exception.reason, "feature_disabled")
+ network.assert_not_called()
+
+ def test_enabled_fetch_builds_query_from_criteria_and_parses_visible_result_links(self):
+ os.environ["GOOGLE_BROWSER_SEARCH_ENABLED"] = "true"
+
+ class Response:
+ def read(self, _limit):
+ return b'
Acme Solar Settings '
+
+ def __enter__(self):
+ return self
+
+ def __exit__(self, *_):
+ return False
+
+ with patch("app.sources.urlopen", return_value=Response()) as network:
+ page = GoogleBrowserSearchSource().discover(
+ {"approved": True, "public_access": True, "terms_accepted": True, "rate_limit": 6},
+ criteria={"keywords": ["solar"], "city": "Cape Town"},
+ limits={"max_records": 4},
+ )
+ self.assertEqual(page.records, [{"name": "Acme Solar", "website": "https://acme.example/about", "email": "", "phone": "", "description": "Public Google search result", "location": ""}])
+ request = network.call_args.args[0]
+ self.assertIn("q=solar+Cape+Town", request.full_url)
+ self.assertNotIn("query", request.full_url)
+ self.assertTrue(page.metadata["public_html_only"])
+
+
+if __name__ == "__main__":
+ unittest.main()
diff --git a/apps/api/tests/test_scoped_discovery.py b/apps/api/tests/test_scoped_discovery.py
index 54bdf19..1c79a1d 100644
--- a/apps/api/tests/test_scoped_discovery.py
+++ b/apps/api/tests/test_scoped_discovery.py
@@ -84,6 +84,46 @@ class ScopedDiscoveryApiTests(unittest.TestCase):
run = self.request('GET', '/api/v1/discovery-runs')[1]['items'][0]
self.assertEqual(run['seed_urls'], []); self.assertEqual(run['result']['candidates'][0]['provenance']['mechanism'], 'criteria_search_provider')
+ def test_selected_disabled_source_is_rejected_before_job_creation(self):
+ status, source = self.request('POST', '/api/v1/sources', {
+ 'name': 'Disabled manual', 'kind': 'manual',
+ 'config': {'rows': [{'name': 'Not runnable'}]},
+ })
+ self.assertEqual(status, 201)
+ status, body = self.request('POST', '/api/v1/discovery', {
+ 'criteria': {'category': 'plumbers'}, 'source_ids': [source['id']],
+ 'idempotency_key': 'disabled-source',
+ })
+ self.assertEqual(status, 409)
+ self.assertEqual(body['error'], 'selected_source_not_ready')
+ self.assertEqual(self.request('GET', '/api/v1/discovery-runs')[1]['items'], [])
+
+ def test_source_dry_run_validates_enabled_source_without_persisting_candidates(self):
+ status, source = self.request('POST', '/api/v1/sources', {
+ 'name': 'Preview manual', 'kind': 'manual', 'enabled': True,
+ 'config': {'rows': [{'name': 'Preview Plumbing', 'website': 'https://preview.example.test'}]},
+ })
+ self.assertEqual(status, 201)
+ with patch('app.main.enrich_source_business') as enrich:
+ status, job = self.request('POST', '/api/v1/discovery', {
+ 'criteria': {'category': 'plumbers', 'city': 'Cape Town'},
+ 'source_ids': [source['id']], 'dry_run': True,
+ 'idempotency_key': 'source-dry-run', 'max_candidates': 5,
+ })
+ self.assertEqual(status, 202)
+ for _ in range(100):
+ _, current = self.request('GET', '/api/v1/jobs/' + str(job['id']))
+ if current['status'] in ('succeeded', 'failed'):
+ break
+ time.sleep(.02)
+ self.assertEqual(current['status'], 'succeeded')
+ self.assertFalse(enrich.called)
+ self.assertEqual(self.request('GET', '/api/v1/businesses')[1]['items'], [])
+ self.assertEqual(self.request('GET', '/api/v1/source-records')[1]['items'], [])
+ run = self.request('GET', '/api/v1/discovery-runs')[1]['items'][0]
+ self.assertTrue(run['dry_run'])
+ self.assertEqual(run['result_count'], 1)
+
def test_results_are_tenant_isolated(self):
ph, salt = hash_password('other-password')
db = sqlite3.connect(self.server.db_path); db.execute("INSERT INTO organizations VALUES ('other-tenant','Other',CURRENT_TIMESTAMP)"); db.execute("INSERT INTO users (organization_id,email,password_hash,password_salt,role) VALUES (?,?,?,?,?)", ('other-tenant','other@example.test',ph,salt,'owner')); db.commit(); db.close()
diff --git a/apps/api/tests/test_sources_phase5.py b/apps/api/tests/test_sources_phase5.py
index 66d96e4..9b91372 100644
--- a/apps/api/tests/test_sources_phase5.py
+++ b/apps/api/tests/test_sources_phase5.py
@@ -37,13 +37,14 @@ class SourceAdapterTests(unittest.TestCase):
def read(self, _): return b'{"elements":[{"tags":{"name":"Cape Plumber","craft":"plumber"}}]}'
def __enter__(self): return self
def __exit__(self, *_): return False
- config={'provider':'openstreetmap','query':'plumbers','location':'Cape Town','approved':True,'public_access':True,'terms_accepted':True,'rate_limit':1}
+ config={'provider':'openstreetmap','approved':True,'public_access':True,'terms_accepted':True,'rate_limit':1}
+ criteria={'keywords':['plumbers'],'city':'Cape Town'}
with patch('app.sources.urlopen',return_value=Response()) as request:
- page=ApprovedDirectorySource().discover(config)
+ result=ApprovedDirectorySource().discover(config,criteria=criteria,limits={'max_records':10})
query=request.call_args.args[0].data.decode()
self.assertIn('"craft"~"plumber|plumbers",i]',query)
self.assertIn('"shop"~"plumber|plumbers",i]',query)
- self.assertEqual(page.records[0]['name'],'Cape Plumber')
+ self.assertEqual(result.records[0]['name'],'Cape Plumber')
def test_legacy_source_kind_constraint_is_migrated(self):
with TemporaryDirectory() as tmp:
@@ -54,7 +55,8 @@ class SourceAdapterTests(unittest.TestCase):
db=connect(path)
self.assertEqual(db.execute("SELECT kind FROM sources WHERE name='Existing manual'").fetchone()[0],'manual')
db.execute("INSERT INTO sources(organization_id,name,kind,source_code) VALUES(?,?,?,?)",(ORGANIZATION_ID,'OpenStreetMap / Overpass · plumbers','approved_directory','openstreetmap'))
- db.commit(); self.assertEqual(db.execute("SELECT kind FROM sources WHERE source_code='openstreetmap'").fetchone()[0],'approved_directory'); db.close()
+ db.execute("INSERT INTO sources(organization_id,name,kind,source_code) VALUES(?,?,?,?)",(ORGANIZATION_ID,'Experimental Google','google_browser_search','google_browser_search'))
+ db.commit(); self.assertEqual(db.execute("SELECT kind FROM sources WHERE source_code='openstreetmap'").fetchone()[0],'approved_directory'); self.assertEqual(db.execute("SELECT kind FROM sources WHERE source_code='google_browser_search'").fetchone()[0],'google_browser_search'); db.close()
class SourceApiTests(unittest.TestCase):
def setUp(self):
@@ -177,6 +179,52 @@ class SourceApiTests(unittest.TestCase):
threading.Event().wait(.01)
self.assertEqual(current['status'],'failed'); self.assertEqual(current['error_code'],'SOURCE_DISABLED')
+ def test_google_browser_source_cannot_enable_without_runtime_feature_flag(self):
+ with patch.dict(os.environ, {'GOOGLE_BROWSER_SEARCH_ENABLED': 'false'}):
+ status, source = self.req('POST', '/api/v1/sources', {
+ 'name': 'Disabled Experimental Google', 'kind': 'google_browser_search',
+ 'config': {'approved': True, 'public_access': True, 'terms_accepted': True, 'rate_limit': 6},
+ })
+ self.assertEqual(status, 201)
+ status, body = self.req('PATCH', f"/api/v1/sources/{source['id']}", {'enabled': True})
+ self.assertEqual(status, 409)
+ self.assertEqual(body['error'], 'source_feature_disabled')
+ with patch.dict(os.environ, {'GOOGLE_BROWSER_SEARCH_ENABLED': 'false'}):
+ status, body = self.req('POST', '/api/v1/sources', {
+ 'name': 'Still Disabled Experimental Google', 'kind': 'google_browser_search', 'enabled': True,
+ 'config': {'approved': True, 'public_access': True, 'terms_accepted': True, 'rate_limit': 6},
+ })
+ self.assertEqual(status, 409)
+ self.assertEqual(body['error'], 'source_feature_disabled')
+
+ def test_google_browser_block_is_structured_and_stops_source_job(self):
+ class Response:
+ def read(self, _limit): return b"Our systems have detected unusual traffic from your computer network."
+ def __enter__(self): return self
+ def __exit__(self, *_): return False
+ with patch.dict(os.environ, {'GOOGLE_BROWSER_SEARCH_ENABLED': 'true'}):
+ status, source = self.req('POST', '/api/v1/sources', {
+ 'name': 'Experimental Google', 'kind': 'google_browser_search', 'enabled': True,
+ 'config': {'approved': True, 'public_access': True, 'terms_accepted': True, 'rate_limit': 6},
+ })
+ self.assertEqual(status, 201)
+ status, job = self.req('POST', '/api/v1/discovery', {
+ 'criteria': {'keywords': ['solar'], 'city': 'Cape Town'},
+ 'selected_adapters': ['google_browser_search'], 'idempotency_key': 'google-blocked', 'max_records': 2,
+ })
+ self.assertEqual(status, 202)
+ with patch('app.sources.urlopen', return_value=Response()) as network:
+ for _ in range(100):
+ _, current = self.req('GET', f"/api/v1/jobs/{job['id']}")
+ if current['status'] in ('succeeded', 'failed'):
+ break
+ threading.Event().wait(.01)
+ self.assertEqual(current['status'], 'failed')
+ self.assertEqual(current['error_code'], 'GOOGLE_BROWSER_BLOCKED')
+ self.assertEqual(network.call_count, 1)
+ events = self.req('GET', f"/api/v1/jobs/{job['id']}/events")[1]['items']
+ self.assertIn('GOOGLE_BROWSER_BLOCKED', [event.get('error_code') for event in events])
+
def test_source_discovery_persists_pipeline_and_is_idempotent(self):
status, source = self.req('POST', '/api/v1/sources', {
'name': 'Manual leads', 'kind': 'manual', 'enabled': True,
diff --git a/apps/web/app.js b/apps/web/app.js
index 17a65f7..8ac29eb 100644
--- a/apps/web/app.js
+++ b/apps/web/app.js
@@ -266,9 +266,7 @@
}
async function loadJobDetail(id) { selectedJobId = id; renderJobsList(); $('jobDetailPanel').innerHTML = 'Loading job detail…
'; try { const job = await jobsRequest(`/api/v1/jobs/${encodeURIComponent(id)}`); let events = job.events || job.event_timeline; if (!events) { const eventPayload = await jobsRequest(`/api/v1/jobs/${encodeURIComponent(id)}/events`); events = eventPayload.events || eventPayload.items || eventPayload; } renderJobDetail(job, events); } catch (error) { if (error.message !== 'unauthorized') $('jobDetailPanel').innerHTML = `Unable to load job ${esc(error.message)}
Try again `; } }
async function loadJobs({silent = false} = {}) { if (!silent) { jobMessage('Loading jobs…'); resetJobCounts(); } try { const payload = await jobsRequest('/api/v1/jobs'); jobs = Array.isArray(payload) ? payload : (payload.jobs || payload.items || []); renderJobCounts(payload); renderJobsList(); $('jobsUpdatedAt').textContent = `Updated ${new Date().toLocaleTimeString([], {hour:'2-digit', minute:'2-digit'})}`; jobMessage(''); const active = jobs.some(job => ['queued','running'].includes(jobStatus(job))); if (active && !jobPollTimer) jobPollTimer = setInterval(() => loadJobs({silent:true}), 5000); if (!active && jobPollTimer) { clearInterval(jobPollTimer); jobPollTimer = null; } if (selectedJobId) { const selected = jobs.find(job => String(job.id) === String(selectedJobId)); if (selected) await loadJobDetail(selectedJobId); } } catch (error) { jobs = []; resetJobCounts(); renderJobsList(); if (error.message !== 'unauthorized') jobMessage(error.message || 'Unable to load jobs.', true); } }
- async function startDemoJob() { if (!canManageJobs()) { jobMessage('Your role is not permitted to start jobs.', true); return; } const button = $('startDemoJobBtn'); button.disabled = true; jobMessage('Starting demo job…'); try { const job = await jobsRequest('/api/v1/jobs', {method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify({type:'noop', payload:{}, idempotency_key:`demo-${Date.now()}-${Math.random().toString(36).slice(2)}`})}); await loadJobs({silent:true}); if (job?.id) await loadJobDetail(job.id); jobMessage('Demo job started.'); } catch (error) { if (error.message !== 'unauthorized') jobMessage(error.message || 'Unable to start demo job.', true); } finally { button.disabled = !canManageJobs(); } }
- async function jobAction(action) { const job = jobs.find(item => String(item.id) === String(selectedJobId)); if (!job || !canManageJobs()) return; const endpointPath = action === 'cancel' ? `/api/v1/jobs/${encodeURIComponent(job.id)}/cancel` : `/api/v1/jobs/${encodeURIComponent(job.id)}/retry`; const label = action === 'cancel' ? 'cancel' : 'retry'; try { await jobsRequest(endpointPath, {method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify({})}); jobMessage(`Job ${label} requested.`); await loadJobs({silent:true}); await loadJobDetail(job.id); } catch (error) { if (error.message !== 'unauthorized') jobMessage(error.message || `Unable to ${label} job.`, true); } }
- function updateJobPermissions() { if ($('startDemoJobBtn')) $('startDemoJobBtn').disabled = !canManageJobs(); }
+ function updateJobPermissions() {}
let sources = [], selectedSourceId = null;
const sourceState = source => Boolean(source?.enabled ?? source?.active);
@@ -279,28 +277,32 @@
const sourceType = source => source.source_code || source.kind || source.type || 'unknown';
const sourceText = (source, keys, fallback='Not returned') => { for (const key of keys) if (source?.[key] !== undefined && source[key] !== null && source[key] !== '') return source[key]; return fallback; };
function sourceMessage(text, error = false) { const el = $('sourcesMessage'); if (el) { el.textContent = text || ''; el.className = `sources-message${error ? ' error' : ''}`; } }
- function renderSourceSelect() { const select = $('discoverySource'); if (select) select.innerHTML = `Select a source ${sources.filter(source=>!source.optional).map(s => `${esc(sourceLabel(s))} · ${esc(sourceStatus(s))} `).join('')}`; const multi=$('directDiscoverySources'); if(multi) multi.innerHTML=sources.filter(sourceState).map(s=>`${esc(sourceLabel(s))} · ${esc(sourceType(s))} `).join('') || 'No approved sources enabled '; }
- function renderSources() { renderSourceSelect(); const list = $('sourcesList'); if (!sources.length) { list.innerHTML = 'No registered or available sources returned by the workspace.
'; return; } list.innerHTML = sources.map(source => { const policy=sourceJson(source.policy_json||source.policy), quota=sourceJson(source.quota_json||source.quota), config=sourceJson(source.config_json||source.config), status=sourceStatus(source), configured=source.configured ?? (!source.optional && (source.approved || sourceType(source)==='manual'||sourceType(source)==='csv')), available=source.available ?? (!source.optional || Boolean(source.configured)), health=sourceText(source,['health_status','health'],'Not tested'), failures=sourceText(source,['consecutive_failures'],'0'), circuit=source.circuit_open===true||source.circuit_open===1?'Open':'Closed', credential=sourceText(source,['api_credential_status','credential_status'],source.optional?'Required / not configured':'Not required'), terms=sourceText(source,['terms_status','terms_reviewed'],policy.terms_accepted===true?'Accepted':policy.terms_url||config.terms_url?'Provided':'Not reviewed'), owner=sourceText(source,['owner','owner_name'],policy.owner||config.owner||'Not assigned'), rate=sourceText(source,['rate_limit','rate_limit_label'],policy.rate_limit||config.rate_limit||'Not set'), daily=sourceText(source,['daily_quota','daily_limit'],quota.daily_limit||'Not set'), lastHealth=sourceText(source,['last_health_at','last_checked_at','updated_at'],'Not checked'), success=sourceText(source,['last_success_at'],'No successful run'), error=sourceText(source,['last_error','error'],'None recorded'); return `${esc(sourceLabel(source))} ${esc(sourceType(source))}
${source.id ? (source.optional?'Optional adapter · configuration-gated':'Registered workspace source') : 'Adapter preview · register before testing or enabling'} ${esc(status==='unavailable'?'Unavailable':status)}
Integration ${esc(source.optional?'Optional adapter':'Registered')}
Configured ${configured?'Yes':'No'}
Available ${available?'Yes':'No'}
Enabled ${sourceState(source)?'Yes':'No'}
API credential ${esc(credential)}
Terms ${esc(terms)}
Owner ${esc(owner)}
Rate limit ${esc(rate)}
Daily quota ${esc(daily)}
Last health ${esc(lastHealth)} · ${esc(health)}
Success / error ${esc(success)} ${esc(error)}
Circuit ${esc(circuit)} · ${esc(failures)} failures Review Test ${source.id ? (status === 'enabled' ? 'Disable' : 'Enable') : 'Set up source'}
`; }).join(''); }
- function renderSourceRecords(items) { const list = $('sourceRecordsList'); if (!items.length) { list.innerHTML = 'No source records returned by the workspace.
'; return; } list.innerHTML = `Record Source Status Observed ${items.slice(0,25).map(record => `${esc(record.name || record.title || record.external_id || record.id || 'Unnamed record')} ${esc(record.source_name || record.source || 'Unknown source')} ${esc(record.status || 'Pending')} ${esc(record.observed_at || record.created_at || 'Time unavailable')} `).join('')}
`; }
- async function loadSources() { sourceMessage('Loading sources…'); $('sourcesList').innerHTML = 'Loading source registry…
'; $('sourceRecordsList').innerHTML = 'Loading source records…
'; try { const [sourcePayload, adapterPayload, recordPayload] = await Promise.all([jsonRequest('/api/v1/sources'), jsonRequest('/api/v1/sources/adapters'), jsonRequest('/api/v1/source-records?page_size=25')]); const configured=sourceItems(sourcePayload), registeredCodes=new Set(configured.map(sourceType)); const optional=sourceItems(adapterPayload).filter(adapter=>!registeredCodes.has(adapter.source_code)).map(adapter=>({...adapter,source_code:adapter.source_code,display_name:adapter.display_name,kind:adapter.source_code,optional:Boolean(adapter.optional),configured:false,available:Boolean(adapter.available),enabled:false})); sources=[...configured,...optional]; renderSources(); renderSourceRecords(sourceItems(recordPayload)); $('sourcesUpdatedAt').textContent = `Updated ${new Date().toLocaleTimeString([], {hour:'2-digit', minute:'2-digit'})}`; sourceMessage(sources.some(sourceState) ? '' : 'No live source is enabled. Configure and enable a ready source to begin discovery.'); } catch (error) { sources = []; renderSources(); renderSourceRecords([]); if (error.message !== 'unauthorized') sourceMessage(error.message || 'Unable to load sources.', true); } }
- async function saveSource(event) { event.preventDefault(); const form = event.currentTarget, fields = Object.fromEntries(new FormData(form).entries()); if (fields.source_type === 'csv' && !fields.csv_content.trim()) { message('sourceFormMessage', 'CSV content is required for a CSV source.', true); return; } const config = {url:fields.url, terms_url:fields.terms_url, owner:fields.owner, rate_limit:fields.rate_limit}; if (fields.source_type === 'csv') config.csv = fields.csv_content; else config.rows = []; try { await jsonRequest('/api/v1/sources', {method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify({name:fields.name, kind:fields.source_type, config, enabled:false})}); message('sourceFormMessage', 'Source saved. It remains disabled until explicitly enabled.'); form.reset(); $('sourceCsvField').hidden = true; await loadSources(); } catch (error) { if (error.message !== 'unauthorized') message('sourceFormMessage', error.message || 'Unable to save source.', true); } }
- function findRegisteredSource(id) { return sources.find(item => String(item.id) === String(id) && !item.optional) || sources.find(item => String(item.id) === String(id)) || sources.find(item => !item.optional && sourceType(item) === String(id)); }
- async function registerAdapterPreview(code) {
- const preview = sources.find(item => !item.id && sourceType(item) === String(code));
- if (!preview) { sourceMessage('This adapter preview is no longer available. Refresh the source registry.', true); return; }
- const type = sourceType(preview);
- if (type === 'google_places') { document.querySelector('.google-places-panel')?.scrollIntoView({behavior:'smooth',block:'center'}); sourceMessage('Google Places requires its protected setup form and an approved server-side key.'); return; }
- if (!['openstreetmap','wikidata','common_crawl','public_website','dns','rdap','ct_logs'].includes(type)) { sourceMessage('Use the source setup form to register this adapter before it can be tested or enabled.'); return; }
- const promptText = type === 'public_website' ? 'Enter a public website URL:' : type === 'dns' ? 'Enter a domain for DNS lookup:' : type === 'rdap' ? 'Enter a domain for RDAP lookup:' : type === 'ct_logs' ? 'Enter a domain for certificate-transparency lookup:' : 'Enter the business/category search query:';
- const value = window.prompt(promptText);
- if (!value || !value.trim()) return;
- const text = value.trim();
- const config = type === 'public_website' ? {urls:[text]} : type === 'dns' ? {domains:[text]} : type === 'rdap' || type === 'ct_logs' ? {domain:text} : {provider:type,query:text,location:(window.prompt('Enter the location (optional):','South Africa') || 'South Africa').trim(),approved:true,public_access:true,terms_accepted:true,rate_limit:1};
- try { const registered=await jsonRequest('/api/v1/sources', {method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify({name:`${sourceLabel(preview)} · ${text}${['openstreetmap','wikidata','common_crawl'].includes(type) ? ` · ${config.location}` : ''}`,kind:type,config})}); await loadSources(); sourceMessage(registered.created === false ? 'Existing source reloaded. Test it, then enable it once the health check succeeds.' : 'Source registered. Test it, then enable it once the health check succeeds.'); }
- catch (error) { if (error.message !== 'unauthorized') sourceMessage(error.message || 'Unable to register source.', true); }
+ function renderSourceSelect(){const multi=$('directDiscoverySources'),enabled=sources.filter(sourceState);if(multi)multi.innerHTML=enabled.map(s=>`${esc(sourceLabel(s))} `).join('')||'No enabled sources are available ';const help=$('directDiscoverySourceHelp');if(help)help.textContent=enabled.length?'Select the approved sources to use for this run.':'No enabled sources are available. Configure, test, and enable a source first.';}
+ function renderSources(){renderSourceSelect();const list=$('sourcesList');if(!sources.length){list.innerHTML='No registered sources returned by the workspace. Save configuration to add one.
';return;}list.innerHTML=sources.map(source=>{const config=sourceJson(source.config_json||source.config),policy=sourceJson(source.policy_json||source.policy),quota=sourceJson(source.quota_json||source.quota),owner=sourceText(source,['owner'],policy.owner||config.owner||'Not assigned'),terms=sourceText(source,['terms_status'],policy.terms_status||config.terms_status||'Pending'),rate=sourceText(source,['rate_limit'],policy.rate_limit||config.rate_limit||'Not set'),daily=sourceText(source,['daily_quota'],quota.daily_quota||config.daily_quota||'Not set'),health=sourceText(source,['health_status','health'],'Not tested'),ready=source.id&&terms==='approved'&&owner!=='Not assigned'&&rate!=='Not set'&&daily!=='Not set'&&/healthy|passed|success/i.test(health);return `${esc(sourceLabel(source))} ${esc(sourceType(source))} ${esc(sourceStatus(source))}
Registry configuration · no discovery criteria stored
Owner ${esc(owner)}
Terms ${esc(terms)}
Rate limit ${esc(rate)}
Daily quota ${esc(daily)}
Health ${esc(health)}
Status ${esc(sourceStatus(source))} Test ${sourceState(source)?'Disable':'Enable'}
`;}).join('');}
+ function renderSourceRecords(items){const list=$('sourceRecordsList');list.innerHTML=items.length?`${items.slice(0,25).map(record=>`${esc(record.name||record.title||record.id||'Unnamed record')} ${esc(record.source_name||record.source||'Unknown source')} ${esc(record.status||'Pending')} `).join('')}
`:'No source records returned by the workspace.
';}
+ async function loadSources(){sourceMessage('Loading sources…');try{const [sourcePayload,adapterPayload,recordPayload]=await Promise.all([jsonRequest('/api/v1/sources'),jsonRequest('/api/v1/sources/adapters'),jsonRequest('/api/v1/source-records?page_size=25')]);const configured=sourceItems(sourcePayload),registeredCodes=new Set(configured.map(sourceType)),optional=sourceItems(adapterPayload).filter(a=>!registeredCodes.has(a.source_code)).map(a=>({...a,kind:a.source_code,optional:true,configured:false,available:Boolean(a.available),enabled:false}));sources=[...configured,...optional];renderSources();renderSourceRecords(sourceItems(recordPayload));$('sourcesUpdatedAt').textContent=`Updated ${new Date().toLocaleTimeString([], {hour:'2-digit',minute:'2-digit'})}`;sourceMessage(sources.some(sourceState)?'':'No enabled source. Configure, test, and enable a source before discovery.');}catch(error){sources=[];renderSources();renderSourceRecords([]);if(error.message!=='unauthorized')sourceMessage(error.message||'Unable to load source registry.',true);}}
+ function setSourceButtons(busy,status=''){['sourceSaveBtn','sourceTestBtn','sourceEnableBtn'].forEach(id=>{const el=$(id);if(el)el.disabled=busy||(id!=='sourceSaveBtn'&&!selectedSourceId);});if(status)$('sourceStatus').textContent=status;}
+ async function saveSource(event){event.preventDefault();const form=event.currentTarget,fields=Object.fromEntries(new FormData(form).entries());if(fields.source_type==='csv'&&!fields.csv_content.trim()){message('sourceFormMessage','CSV content is required for a CSV source.',true);return;}let providerSettings={};try{providerSettings=fields.provider_settings.trim()?JSON.parse(fields.provider_settings):{};}catch{message('sourceFormMessage','Provider settings must be valid JSON.',true);return;}setSourceButtons(true,'Saving source configuration…');message('sourceFormMessage','Saving source configuration…');const config={owner:fields.owner,terms_url:fields.terms_url,terms_status:fields.terms_status,rate_limit:fields.rate_limit,daily_quota:Number(fields.daily_quota),provider_settings:fields.provider_settings};if(fields.source_type==='csv')config.csv=fields.csv_content;else config.rows=[];try{const saved=await jsonRequest('/api/v1/sources',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({name:fields.name,kind:fields.source_type,config,policy:{owner:fields.owner,terms_url:fields.terms_url,terms_status:fields.terms_status,rate_limit:fields.rate_limit},quota:{daily_quota:Number(fields.daily_quota)},enabled:false})});selectedSourceId=saved.id||saved.source?.id||null;message('sourceFormMessage','Source configuration saved. It remains disabled until a successful test and explicit enablement.');await loadSources();}catch(error){if(error.message!=='unauthorized')message('sourceFormMessage',error.message||'Unable to save source configuration.',true);}finally{setSourceButtons(false);}}
+ function findRegisteredSource(id){return sources.find(item=>String(item.id)===String(id)&&!item.optional)||sources.find(item=>String(item.id)===String(id));}
+ async function sourceAction(id, action) {
+ const source = findRegisteredSource(id || selectedSourceId);
+ if (!source?.id) { sourceMessage('Save source configuration before testing or enabling.', true); return; }
+ try {
+ if (action === 'test') {
+ sourceMessage('Testing source…');
+ await jsonRequest(`/api/v1/sources/${encodeURIComponent(source.id)}/test`, { method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify({}) });
+ sourceMessage('Source test completed. Review health before enabling.');
+ } else if (action === 'toggle') {
+ sourceMessage(sourceState(source) ? 'Disabling source…' : 'Enabling source…');
+ await jsonRequest(`/api/v1/sources/${encodeURIComponent(source.id)}`, { method:'PATCH', headers:{'Content-Type':'application/json'}, body:JSON.stringify({ enabled: !sourceState(source) }) });
+ sourceMessage(`Source ${sourceState(source) ? 'disabled' : 'enabled'} successfully.`);
+ }
+ await loadSources();
+ } catch (error) {
+ if (error.message !== 'unauthorized') sourceMessage(error.message || `Unable to ${action} source.`, true);
+ }
}
- async function sourceAction(id, action) { if (action === 'setup') return registerAdapterPreview(id); const source = findRegisteredSource(id); if (!source || !source.id) { sourceMessage('This source is an adapter preview, not a registered source.', true); return; } if (action === 'review') { selectedSourceId=source.id; document.querySelector(`[data-source-id="${CSS.escape(String(id))}"]`)?.scrollIntoView({behavior:'smooth',block:'center'}); sourceMessage('Source details are shown below. Review terms, owner, limits, health, and circuit state before enabling.'); return; } try { let notice=''; if (action === 'test') { await jsonRequest(`/api/v1/sources/${encodeURIComponent(source.id)}/test`, {method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify({})}); notice='Source test completed.'; } else { const enabled = sourceState(source); await jsonRequest(`/api/v1/sources/${encodeURIComponent(source.id)}`, {method:'PATCH', headers:{'Content-Type':'application/json'}, body:JSON.stringify({enabled:!enabled})}); notice=`Source ${enabled ? 'disabled' : 'enabled'} by the workspace.`; } await loadSources(); sourceMessage(notice); } catch (error) { if (error.message !== 'unauthorized') sourceMessage(error.message || `Unable to ${action} source.`, true); } }
- async function runDiscovery(dryRun) { const form = $('discoveryForm'), data = Object.fromEntries(new FormData(form).entries()); data.dry_run = Boolean(dryRun); if (!data.source_id || !data.query.trim()) { message('discoveryMessage', 'Select a source and enter a query.', true); return; } const source = sources.find(item => String(item.id) === String(data.source_id)); if (!dryRun && !sourceState(source)) { message('discoveryMessage', 'This source is disabled. Enable it only after review.', true); return; } const button = dryRun ? $('discoveryDryRunBtn') : $('discoveryRunBtn'); button.disabled = true; message('discoveryMessage', dryRun ? 'Validating query…' : 'Starting discovery…'); try { const query = await jsonRequest('/api/v1/discovery-queries', {method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify({source_id:Number(data.source_id), name:data.query.trim().slice(0,80), query:{text:data.query.trim()}, dry_run:data.dry_run})}); if (!dryRun) await jsonRequest(`/api/v1/discovery-queries/${encodeURIComponent(query.id)}/run`, {method:'POST', headers:{'Content-Type':'application/json'}, body:JSON.stringify({})}); message('discoveryMessage', dryRun ? 'Dry run completed; no discovery job was started.' : 'Discovery request accepted. Check Jobs for progress.'); } catch (error) { if (error.message !== 'unauthorized') message('discoveryMessage', error.message || 'Discovery request failed.', true); } finally { button.disabled = false; } }
+
let discoveryRuns = [], selectedDiscoveryRunId = null, discoveryPollTimer = null;
const discoveryStatus = run => String(run?.status || run?.state || 'queued').toLowerCase().replaceAll('_','-');
const discoveryLabel = state => ({queued:'Queued',running:'Running',succeeded:'Complete',failed:'Failed',cancelled:'Cancelled',partial:'Partial result',paused:'Paused',stale:'Stale'})[state] || state.replaceAll('-',' ');
@@ -316,7 +318,7 @@
const errors = run?.errors || result.errors || [];
$('directDiscoveryResultTitle').textContent = 'Run ' + run.id + ' · ' + (run.criteria?.keywords || run.criteria?.keyword || 'Discovery brief');
$('directDiscoveryResultStatus').textContent = discoveryLabel(status);
- const controls = 'Ⅱ Pause ▶ Resume Cancel Pause/resume require the job-control API.
';
+ const controls = 'Cancel Retry Cancel and retry use available job routes only.
';
const sourceRows = run.sources || run.source_health || [];
const health = sourceRows.length ? '
Source health ' + sourceRows.length + ' sources ' + sourceRows.map(source => '
' + esc(source.name || source.source_name || ('Source ' + (source.source_id || ''))) + ' ' + esc(source.status || source.health || 'Unknown') + ' ' + esc(source.result_count ?? source.results ?? 0) + ' results · ' + esc(source.error || 'No errors') + ' ').join('') + '
' : '
Source health Per-source health will appear when the API returns source telemetry.
';
const log = events.length ? events.map(event => '' + esc(event.message || event.type || 'Run event') + ' ' + esc(event.created_at || event.timestamp || '') + '
').join('') : 'No persisted events returned yet. ';
@@ -326,11 +328,11 @@
const partial = status === 'partial' || run.partial === true || result.partial === true;
el.innerHTML = controls + (partial ? 'Partial results Review provenance before using any candidate.
' : '') + (status === 'failed' ? 'Discovery failed ' + esc(job?.error?.message || job?.message || run.error || 'The job returned no safe result.') + '
' : '') + stats + health + '';
}
- async function discoveryAction(action,run){ if(!run?.id)return; if(action!=='cancel'&&!['pause','resume'].includes(action))return; if(action==='cancel'&&!window.confirm('Cancel this discovery run? The current partial result will remain reviewable.'))return; try{await jsonRequest(`/api/v1/discovery-runs/${encodeURIComponent(run.id)}/${action}`,{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({})});message('directDiscoveryMessage',`Discovery run ${action} requested.`);await loadDiscoveryRuns();}catch(error){if(error.message!=='unauthorized')message('directDiscoveryMessage',error.message,true);} }
+ async function discoveryAction(action,run){if(!run?.id)return;const path=action==='cancel'?`/api/v1/discovery-runs/${encodeURIComponent(run.id)}/cancel`:action==='retry'&&run.job_id?`/api/v1/jobs/${encodeURIComponent(run.job_id)}/retry`:null;if(!path){message('directDiscoveryMessage','This action is not available for the selected run.',true);return;}message('directDiscoveryMessage',action==='cancel'?'Cancelling discovery run…':'Retrying failed discovery job…');try{await jsonRequest(path,{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({})});message('directDiscoveryMessage',action==='cancel'?'Discovery run cancelled.':'Discovery job retry accepted.');await loadDiscoveryRuns();}catch(error){if(error.message!=='unauthorized')message('directDiscoveryMessage',error.message||`Unable to ${action} discovery run.`,true);}}
async function loadDiscoveryRun(run) { selectedDiscoveryRunId=run.id; renderDiscoveryResult(run); try { const job=await jsonRequest(`/api/v1/jobs/${encodeURIComponent(run.job_id)}`); renderDiscoveryResult(run,job); if(['queued','running','paused'].includes(discoveryStatus(job))){clearTimeout(discoveryPollTimer);discoveryPollTimer=setTimeout(loadDiscoveryRuns,3000);} } catch(error){if(error.message!=='unauthorized')renderDiscoveryResult(run,{status:'error',message:error.message});} }
function renderDiscoveryRuns() { const el=$('directDiscoveryRunsState'); if(!discoveryRuns.length){el.innerHTML='No discovery runs yet. Start with a market brief and approved sources.
';return;} el.innerHTML=discoveryRuns.map(run=>`Run ${esc(run.id)} ${esc(run.criteria?.keywords||run.criteria?.keyword||'Discovery brief')} · ${esc(run.criteria?.city||run.criteria?.location||'Global')} ${esc(discoveryLabel(discoveryStatus(run)))} ${esc(run.result_count??discoveryCandidates(run).length)} records `).join(''); }
async function loadDiscoveryRuns() { const state=$('directDiscoveryRunsState'); if(state)state.innerHTML='Loading discovery runs…
'; try { const payload=await jsonRequest('/api/v1/discovery-runs?page_size=50'); discoveryRuns=payloadItems(payload,['items','runs']); const jobsById=new Map(await Promise.all(discoveryRuns.map(async run=>{try{return [String(run.job_id),await jsonRequest(`/api/v1/jobs/${encodeURIComponent(run.job_id)}`)]}catch{return [String(run.job_id),{}]}}))); discoveryRuns=discoveryRuns.map(run=>({...run,status:jobsById.get(String(run.job_id))?.status||run.status||'queued'})); renderDiscoveryRuns(); const selected=discoveryRuns.find(run=>String(run.id)===String(selectedDiscoveryRunId))||discoveryRuns[0]; if(selected)loadDiscoveryRun(selected); } catch(error) { if(error.message!=='unauthorized'&&state)state.innerHTML=`Unable to load discovery runs ${esc(error.message)}
Try again `; } }
- async function submitDirectDiscovery(event) { event.preventDefault(); const form=event.currentTarget, button=$('directDiscoveryRunBtn'), fields=Object.fromEntries(new FormData(form).entries()), seeds=fields.seed_urls.split(/\r?\n|,/).map(value=>value.trim()).filter(Boolean), dryRun=$('directDiscoveryDryRun')?.checked; if(!fields.keywords.trim()){message('directDiscoveryMessage','Enter a category or keyword.',true);return;} if(seeds.length>5){message('directDiscoveryMessage','Use no more than 5 public seed URLs.',true);return;} if(dryRun){message('directDiscoveryMessage',`Plan validated: ${selectedSourceIds().length||'criteria-first'} source selection, max ${fields.max_candidates} records, ${fields.daily_limit}/day. No job created.`);return;} button.disabled=true; message('directDiscoveryMessage','Starting bounded discovery job…'); try { const criteria={keywords:fields.keywords.trim(),city:fields.city.trim(),province:fields.province.trim(),country:fields.country.trim()}; const payload={criteria,source_ids:selectedSourceIds(),daily_limit:Number(fields.daily_limit),schedule:fields.schedule,max_pages:Number(fields.max_pages),max_candidates:Number(fields.max_candidates),idempotency_key:`discovery-${Date.now()}-${Math.random().toString(36).slice(2)}`}; if(seeds.length)payload.seed_urls=seeds; const job=await jsonRequest('/api/v1/discovery',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify(payload)}); message('directDiscoveryMessage','Discovery accepted. Tracking job status and partial results below.'); await loadDiscoveryRuns(); const run=discoveryRuns.find(item=>String(item.job_id)===String(job.id)); if(run)loadDiscoveryRun(run); } catch(error) { if(error.message!=='unauthorized')message('directDiscoveryMessage',error.message||'Discovery request failed.',true); } finally {button.disabled=false;} }
+ async function submitDirectDiscovery(event){event.preventDefault();const form=event.currentTarget,button=$('directDiscoveryRunBtn'),fields=Object.fromEntries(new FormData(form).entries()),include=(fields.include_websites||'').split(/\r?\n|,/).map(v=>v.trim()).filter(Boolean),exclude=(fields.exclude_websites||'').split(/\r?\n|,/).map(v=>v.trim()).filter(Boolean),dryRun=$('directDiscoveryDryRun')?.checked;if(!fields.category.trim()||!fields.keywords.trim()){message('directDiscoveryMessage','Enter a category and keywords before starting discovery.',true);return;}if(!selectedSourceIds().length){message('directDiscoveryMessage','Select at least one enabled source before starting a discovery run.',true);return;}const criteria={category:fields.category.trim(),keywords:fields.keywords.trim(),city:fields.city.trim(),province:fields.province.trim(),country:fields.country.trim(),language:fields.language.trim()};const payload={criteria,source_ids:selectedSourceIds(),daily_limit:Number(fields.daily_limit),max_pages:Number(fields.max_pages),max_results:Number(fields.max_results),dry_run:Boolean(dryRun),website_analysis:{priorities:fields.priorities.trim(),include_websites:include,exclude_websites:exclude},idempotency_key:`discovery-${Date.now()}-${Math.random().toString(36).slice(2)}`};if(dryRun){message('directDiscoveryMessage',`Dry run validated: ${payload.source_ids.length} enabled sources, ${payload.max_results} max results. No job created.`);return;}button.disabled=true;message('directDiscoveryMessage','Starting bounded discovery job…');try{const job=await jsonRequest('/api/v1/discovery',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify(payload)});message('directDiscoveryMessage','Discovery accepted. Tracking job status and partial results below.');await loadDiscoveryRuns();const run=discoveryRuns.find(item=>String(item.job_id)===String(job.id));if(run)loadDiscoveryRun(run);}catch(error){if(error.message!=='unauthorized')message('directDiscoveryMessage',error.message||'Unable to start discovery run.',true);}finally{button.disabled=false;}}
function parseCsv(text){const lines=text.trim().split(/\r?\n/).filter(Boolean),cells=line=>line.match(/("[^"]*(?:""[^"]*)*"|[^,]+)(?=,|$)/g)?.map(x=>x.replace(/^"|"$/g,'').replaceAll('""','"'))||[];if(!lines.length)return[];const headers=cells(lines[0]);return lines.slice(1,11).map(l=>Object.fromEntries(cells(l).map((v,i)=>[headers[i]||`column_${i+1}`,v])));}
function renderCsv(rows){if(!rows.length){$('csvPreview').innerHTML='⊞ No data rows found
';return;}const h=Object.keys(rows[0]);$('csvPreview').className='csv-table';$('csvPreview').innerHTML=`${h.map(x=>`${esc(x)} `).join('')} ${rows.map(r=>`${h.map(x=>`${esc(r[x])} `).join('')} `).join('')}
Showing up to 10 rows · Preview only; nothing added yet. `;}
@@ -339,7 +341,7 @@
async function bootstrap(){let res;try{res=await fetch(endpoint('/api/v1/auth/me'),{credentials:'include'});}catch(e){showLogin('Unable to connect to the workspace. Try again.');return;}if(res.status===401){showLogin();return;}if(!res.ok){showLogin('Unable to verify the workspace session. Try again.');return;}let user;try{user=await res.json();}catch(e){showLogin('Unable to read the workspace session. Try again.');return;}showDashboard(user.user||user);const panelLoads=[loadData(),loadJobs(),loadSources(),loadDiscoveryRuns(),loadScorePanels(),loadSavedFilters(),loadReviewQueue(),loadCrmPipeline(),loadReports(),loadSuppressions(),loadProviderPolicy(),loadAiProviderSettings()];const results=await Promise.allSettled(panelLoads);const rejected=results.filter(result=>result.status==='rejected');if(rejected.length){console.error('Workspace dashboard panels failed',rejected);$('apiStatus').textContent='● Some panels need refresh';$('apiStatus').classList.remove('live');setExplorerState('Some workspace panels could not be loaded. Use Refresh to retry.',true);}}
document.addEventListener('submit',e=>{if(e.target.id==='contactForm')saveContact(e.target);if(e.target.id==='noteForm')saveNote(e.target);if(e.target.id==='pipelineForm')saveStage(e.target);});
document.addEventListener('click',e=>{if(e.target.id==='verifyBtn')verify();if(e.target.id==='recalculateScoreBtn')recalculateScore();if(e.target.id==='generateAiSuggestionBtn')generateAiSuggestion();if(e.target.id==='createOutreachDraftBtn')createOutreachDraft();const outreachButton=e.target.closest?.('[data-outreach-action]');if(outreachButton)approveOutreachDraft(outreachButton.dataset.outreachId);const aiButton=e.target.closest?.('[data-ai-action]');if(aiButton)aiAction(aiButton.dataset.aiAction,aiButton.dataset.aiId);const scoreEdit=e.target.closest?.('[data-score-edit]');if(scoreEdit)updateScoreRule(scoreEdit.dataset.scoreEdit,scoreEdit.dataset.scorePoints);if(e.target.id==='extractContactsBtn'&&selectedId)loadContactExtraction(selectedId,{extract:true});if(e.target.id==='refreshContactExtractionBtn'&&selectedId)loadContactExtraction(selectedId,{extract:true});if(e.target.id==='retryContactExtractionBtn'&&selectedId)loadContactExtraction(selectedId);if(e.target.id==='scanWebsiteBtn'&&selectedId)loadWebsiteScan(selectedId,{scan:true});if(e.target.id==='refreshWebsiteScanBtn'&&selectedId)loadWebsiteScan(selectedId);if(e.target.id==='retryWebsiteScanBtn'&&selectedId)loadWebsiteScan(selectedId);if(e.target.id==='runDomainCheckBtn')runDomainCheck();if(e.target.id==='retryDomainBtn'&&selectedId)loadDomainIntelligence(selectedId);const availability=e.target.closest?.('[data-domain-availability]');if(availability)checkDomainAvailability(availability.dataset.domain,availability);if(e.target.id==='retryDetailBtn'&&selectedId)loadDetail(selectedId);if(e.target.id==='retryDedupBtn'&&selectedId)loadMatchSuggestions(selectedId);if(e.target.id==='retryHistoryBtn'&&selectedId)loadMergeHistory(selectedId);if(e.target.id==='cancelMergeBtn'||e.target.id==='cancelMergeBtnSecondary')closeMergeDialog();if(e.target.id==='confirmMergeBtn')confirmMerge();const mergeButton=e.target.closest?.('[data-merge-target]');if(mergeButton)openMergeDialog(mergeButton.dataset.mergeTarget,mergeButton.dataset.mergeTargetName);const reverseButton=e.target.closest?.('[data-reverse-merge]');if(reverseButton)reverseMerge(reverseButton.dataset.reverseMerge);if(e.target.id==='retryJobDetailBtn'&&selectedJobId)loadJobDetail(selectedJobId);if(e.target.id==='cancelJobBtn')jobAction('cancel');if(e.target.id==='retryJobBtn')jobAction('retry');const row=e.target.closest?.('[data-job-id]');if(row)loadJobDetail(row.dataset.jobId);});
- $('loginForm').addEventListener('submit',login);$('logoutBtn').addEventListener('click',logout);$('directDiscoveryForm').addEventListener('submit',submitDirectDiscovery);$('directDiscoveryRefreshBtn').addEventListener('click',loadDiscoveryRuns);$('directDiscoveryRunsState').addEventListener('click',e=>{const row=e.target.closest?.('[data-discovery-run-id]');if(row){const run=discoveryRuns.find(item=>String(item.id)===String(row.dataset.discoveryRunId));if(run)loadDiscoveryRun(run);}if(e.target.id==='retryDiscoveryRunsBtn')loadDiscoveryRuns();});$('searchInput').addEventListener('input',()=>{page=1;renderRows();});['scoreFilter','statusFilter','websiteClassFilter','pipelineFilter'].forEach(id=>$(id).addEventListener('change',()=>{page=1;loadData();}));$('pageSize').addEventListener('change',e=>{pageSize=Number(e.target.value);page=1;loadData();});$('nextPageBtn').addEventListener('click',()=>{if(hasNextPage){page+=1;loadData();}});$('refreshBtn').addEventListener('click',loadData);$('jobsRefreshBtn').addEventListener('click',()=>loadJobs());$('startDemoJobBtn').addEventListener('click',startDemoJob);$('sourcesRefreshBtn').addEventListener('click',loadSources);$('sourceForm').addEventListener('submit',saveSource);$('sourceType').addEventListener('change',e=>{$('sourceCsvField').hidden=e.target.value!=='csv';});$('discoveryForm').addEventListener('submit',e=>{e.preventDefault();runDiscovery(true);});$('discoveryRunBtn').addEventListener('click',()=>runDiscovery(false));$('sourcesList').addEventListener('click',e=>{const button=e.target.closest?.('[data-source-action]');if(button)sourceAction(button.dataset.sourceId,button.dataset.sourceAction);});$('addForm').addEventListener('submit',addProspect);$('csvInput').addEventListener('change',e=>{const file=e.target.files[0];if(file){const reader=new FileReader();reader.onload=()=>renderCsv(parseCsv(reader.result));reader.readAsText(file);}});$('menuBtn').addEventListener('click',()=>document.querySelector('.sidebar').classList.toggle('open'));document.querySelectorAll('[data-scroll]').forEach(b=>b.addEventListener('click',()=>document.querySelector(b.dataset.scroll)?.scrollIntoView()));
+ $('loginForm').addEventListener('submit',login);$('logoutBtn').addEventListener('click',logout);$('directDiscoveryForm').addEventListener('submit',submitDirectDiscovery);$('directDiscoveryRefreshBtn').addEventListener('click',loadDiscoveryRuns);$('directDiscoveryRunsState').addEventListener('click',e=>{const row=e.target.closest?.('[data-discovery-run-id]');if(row){const run=discoveryRuns.find(item=>String(item.id)===String(row.dataset.discoveryRunId));if(run)loadDiscoveryRun(run);}if(e.target.id==='retryDiscoveryRunsBtn')loadDiscoveryRuns();});$('searchInput').addEventListener('input',()=>{page=1;renderRows();});['scoreFilter','statusFilter','websiteClassFilter','pipelineFilter'].forEach(id=>$(id).addEventListener('change',()=>{page=1;loadData();}));$('pageSize').addEventListener('change',e=>{pageSize=Number(e.target.value);page=1;loadData();});$('nextPageBtn').addEventListener('click',()=>{if(hasNextPage){page+=1;loadData();}});$('refreshBtn').addEventListener('click',loadData);$('jobsRefreshBtn').addEventListener('click',()=>loadJobs());$('sourcesRefreshBtn').addEventListener('click',loadSources);$('sourceForm').addEventListener('submit',saveSource);$('sourceType').addEventListener('change',e=>{$('sourceCsvField').hidden=e.target.value!=='csv';});$('sourceTestBtn').addEventListener('click',()=>sourceAction(selectedSourceId,'test'));$('sourceEnableBtn').addEventListener('click',()=>sourceAction(selectedSourceId,'toggle'));$('sourcesList').addEventListener('click',e=>{const button=e.target.closest?.('[data-source-action]');if(button)sourceAction(button.dataset.sourceId,button.dataset.sourceAction);});$('addForm').addEventListener('submit',addProspect);$('csvInput').addEventListener('change',e=>{const file=e.target.files[0];if(file){const reader=new FileReader();reader.onload=()=>renderCsv(parseCsv(reader.result));reader.readAsText(file);}});$('menuBtn').addEventListener('click',()=>document.querySelector('.sidebar').classList.toggle('open'));document.querySelectorAll('[data-scroll]').forEach(b=>b.addEventListener('click',()=>document.querySelector(b.dataset.scroll)?.scrollIntoView()));
$('savedFilterForm').addEventListener('submit',saveCurrentFilters);$('savedFilterSelect').addEventListener('change',e=>{ $('deleteSavedFilterBtn').disabled=!e.target.value;if(e.target.value)loadSavedFilter(e.target.value);});$('deleteSavedFilterBtn').addEventListener('click',deleteSavedFilter);$('reviewQueueState').addEventListener('change',e=>{const input=e.target.closest?.('[data-review-id]');if(input){if(input.checked)selectedReviewIds.add(String(input.dataset.reviewId));else selectedReviewIds.delete(String(input.dataset.reviewId));updateBulkState();}});$('selectAllReview').addEventListener('change',e=>{reviewQueue.forEach(p=>e.target.checked?selectedReviewIds.add(String(p.id)):selectedReviewIds.delete(String(p.id)));renderReviewQueue({items:reviewQueue,count:$('reviewQueueCount').textContent});});$('bulkVerifyBtn').addEventListener('click',()=>bulkReview('verify'));$('bulkRejectBtn').addEventListener('click',()=>bulkReview('reject'));$('reviewQueueState').addEventListener('click',e=>{if(e.target.id==='retryReviewQueueBtn')loadReviewQueue();});document.querySelectorAll('[data-dashboard-filter]').forEach(card=>card.addEventListener('click',()=>{const kind=card.dataset.dashboardFilter;if(kind==='review'||kind==='suppressed')applyFilters({status:kind});else if(kind==='high')applyFilters({score:'high'});else if(kind==='fresh')applyFilters({});else applyFilters({status:'all',score:'all'});}));
['sourceFilter','geographyFilter','categoryFilter','contactStatusFilter','scoreFilter','statusFilter','websiteClassFilter','pipelineFilter'].forEach(id=>$(id)?.addEventListener('change',()=>{page=1;loadData();}));
$('interactionForm').addEventListener('submit',saveInteraction);$('suppressionForm').addEventListener('submit',addSuppression);$('aiProviderForm').addEventListener('submit',saveAiProviderSettings);$('aiProvider').addEventListener('change',updateAiProviderForm);$('aiProviderRefreshBtn').addEventListener('click',loadAiProviderSettings);$('testAiProviderBtn').addEventListener('click',testAiProviderConnection);$('crmRefreshBtn').addEventListener('click',()=>{loadCrmPipeline();loadInteractions(selectedId);});$('reportsRefreshBtn').addEventListener('click',loadReports);$('suppressionRefreshBtn').addEventListener('click',loadSuppressions);$('pipelineViewToggle').addEventListener('click',()=>{crmListMode=!crmListMode;$('pipelineViewToggle').textContent=crmListMode?'▦ Board view':'☷ List view';$('pipelineViewToggle').setAttribute('aria-pressed',String(crmListMode));renderPipeline();});$('pipelineBoard').addEventListener('click',e=>{const selectButton=e.target.closest?.('[data-crm-select]');if(selectButton){selectedId=Number(selectButton.dataset.crmSelect);selectProspect(selectedId);loadInteractions(selectedId);$('crmActivity').scrollIntoView({behavior:'smooth',block:'start'});}const stageButton=e.target.closest?.('[data-crm-save-stage]');if(stageButton)saveCrmStage(stageButton.dataset.crmSaveStage);if(e.target.id==='retryCrmBtn')loadCrmPipeline();if(e.target.id==='retryInteractionsBtn'&&selectedId)loadInteractions(selectedId);});$('suppressionState').addEventListener('change',e=>{const input=e.target.closest?.('[data-suppression-id]');if(input){if(input.checked)selectedSuppressionIds.add(String(input.dataset.suppressionId));else selectedSuppressionIds.delete(String(input.dataset.suppressionId));updateSuppressionSelection();}});$('suppressionState').addEventListener('click',e=>{const button=e.target.closest?.('[data-remove-suppression]');if(button)removeSuppression(button.dataset.removeSuppression);if(e.target.id==='retrySuppressionsBtn')loadSuppressions();});$('selectAllSuppressions').addEventListener('change',e=>{suppressions.forEach(item=>e.target.checked?selectedSuppressionIds.add(String(item.id)):selectedSuppressionIds.delete(String(item.id)));renderSuppressions();});$('bulkReviewSuppressionsBtn').addEventListener('click',bulkReviewSuppressions);$('discoveryWorkspace').addEventListener('click',e=>{const action=e.target.closest?.('[data-run-action]');if(action){const run=discoveryRuns.find(item=>String(item.id)===String(selectedDiscoveryRunId));discoveryAction(action.dataset.runAction,run);}});
@@ -347,37 +349,6 @@ $('aiProviderForm').addEventListener('submit',saveAiProviderSettings);$('aiProvi
const viewMap = {dashboard:'dashboard', explorer:'prospects', add:'add', discoveryWorkspace:'discovery', jobs:'jobs', sources:'sources', crmPipeline:'crm', crmActivity:'crm', crmReports:'crm', suppressionCenter:'settings', outreachSettings:'settings', aiProviderSettings:'settings', scoreRules:'dashboard'};
function activateView(view) { document.body.dataset.view = view; document.querySelectorAll('.content > section').forEach(section => section.classList.toggle('view-active', viewMap[section.id] === view)); const labels = {dashboard:'Dashboard', prospects:'Prospect explorer', add:'Add prospects', discovery:'Discovery runs', jobs:'Jobs', sources:'Sources', crm:'CRM pipeline', settings:'Settings'}; const page = document.getElementById('topbarPage'); if (page) page.textContent = labels[view] || 'Workspace'; window.scrollTo({top:0, behavior:'smooth'}); }
document.body.dataset.view = 'dashboard';
- const sourceSetup = document.querySelector('.source-config-panel');
- if (sourceSetup) {
- const badge = sourceSetup.querySelector('.small-label');
- const copy = sourceSetup.querySelector('.source-setup-copy');
- if (badge) badge.textContent = 'Public and operator-controlled sources';
- if (copy) copy.textContent = 'Register bounded public sources or operator-controlled imports. New sources start disabled and must be tested before use.';
- sourceSetup.insertAdjacentHTML('afterend', `NO BILLING REQUIRED
Free public source Rate limited Search OpenStreetMap, Wikidata, or Common Crawl without a Google API key.
`);
- $('freeSourceForm')?.addEventListener('submit', async event => { event.preventDefault(); const fields=Object.fromEntries(new FormData(event.currentTarget).entries()), msg=$('freeSourceMessage'); msg.textContent='Registering source…'; msg.className='form-message'; try { await jsonRequest('/api/v1/sources',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({name:`${fields.provider} · ${fields.query.trim()} · ${(fields.location||'South Africa').trim()}`,kind:'approved_directory',config:{provider:fields.provider,query:fields.query.trim(),location:(fields.location||'South Africa').trim(),approved:true,public_access:true,terms_accepted:true,rate_limit:1},policy:{owner:'workspace operator',rate_limit:1,terms_url:'https://www.openstreetmap.org/copyright'}})}); msg.textContent='Source added disabled. Test it, then enable it below.'; event.currentTarget.reset(); event.currentTarget.querySelector('[name="location"]').value='South Africa'; await loadSources(); } catch(error) { msg.textContent=error.message||'Unable to add source.'; msg.className='form-message error'; } });
- sourceSetup.insertAdjacentHTML('afterend', `OPTIONAL PROVIDER
Google Places discovery Secure key storage Add a Google Places API key here. The key is sent directly to the server, encrypted, and never displayed again.
`);
- $('googlePlacesForm')?.addEventListener('submit', async event => { event.preventDefault(); const form=event.currentTarget, fields=Object.fromEntries(new FormData(form).entries()); const msg=$('googlePlacesMessage'); msg.textContent='Saving securely…'; msg.className='form-message'; try { const created=await jsonRequest('/api/v1/sources',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({name:'Google Places',kind:'google_places',config:{approved:true,public_access:true,terms_accepted:true,credential_ref:'google_places_api_key',rate_limit:1,query:fields.query.trim(),region_code:(fields.region_code||'ZA').toUpperCase()}})}); await jsonRequest(`/api/v1/sources/${encodeURIComponent(created.id)}/credentials`,{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({provider:'google_places',key_name:'api_key',api_key:fields.api_key})}); msg.textContent='Google Places is registered and the key is encrypted. Test it, then enable the source.'; form.reset(); form.querySelector('[name="region_code"]').value='ZA'; await loadSources(); } catch(error) { msg.textContent=error.message||'Unable to save Google Places settings.'; msg.className='form-message error'; } });
- }
- async function configureSourceFromUi(source, button) {
- const type = sourceType(source);
- const label = type === 'public_website' ? 'Enter a public website URL' : type === 'dns' ? 'Enter a domain for DNS lookup' : type === 'rdap' ? 'Enter a domain for RDAP lookup' : type === 'ct_logs' ? 'Enter a domain for certificate-transparency lookup' : type === 'openstreetmap' || type === 'wikidata' || type === 'common_crawl' ? 'Enter the business/category search query' : 'Enter source configuration';
- const value = window.prompt(label + ':');
- if (!value || !value.trim()) return;
- const text = value.trim();
- const config = type === 'public_website' ? {urls:[text]} : type === 'dns' ? {domains:[text]} : type === 'openstreetmap' || type === 'wikidata' || type === 'common_crawl' ? {provider:type,query:text,location:(window.prompt('Enter the location (optional):','South Africa') || 'South Africa').trim(),approved:true,public_access:true,terms_accepted:true,rate_limit:1} : {domain:text};
- button.disabled = true;
- try { await jsonRequest(`/api/v1/sources/${encodeURIComponent(source.id)}`, {method:'PATCH', headers:{'Content-Type':'application/json'}, body:JSON.stringify({config})}); await loadSources(); sourceMessage('Source configuration saved. Test it before enabling.'); }
- catch (error) { if (error.message !== 'unauthorized') sourceMessage(error.message || 'Unable to configure source.', true); }
- finally { button.disabled = false; }
- }
- $('sourcesList')?.addEventListener('click', event => {
- const actionButton = event.target.closest?.('[data-source-action="toggle"]');
- if (!actionButton) return;
- const source = findRegisteredSource(actionButton.dataset.sourceId);
- if (!source) return;
- const needsPromptConfig = ['public_website','dns','rdap','ct_logs','openstreetmap','wikidata','common_crawl'].includes(sourceType(source));
- if (source && !source.configured && source.available && needsPromptConfig) { event.preventDefault(); event.stopImmediatePropagation(); configureSourceFromUi(source, actionButton); }
- }, true);
document.querySelectorAll('.sidebar nav a, [data-scroll]').forEach(link => link.addEventListener('click', event => { const target = (link.getAttribute('href') || link.dataset.scroll || '').replace(/^#/, ''); const view = viewMap[target]; if (view) { event.preventDefault(); activateView(view); document.querySelector('.sidebar')?.classList.remove('open'); } }));
bootstrap();
})();
diff --git a/apps/web/asset-manifest.json b/apps/web/asset-manifest.json
index 85ae77c..dbdc6be 100644
--- a/apps/web/asset-manifest.json
+++ b/apps/web/asset-manifest.json
@@ -1,6 +1,6 @@
{
"schema": 1,
- "version": "phase-34",
+ "version": "phase-35",
"entrypoints": [
"config.js",
"app.js",
@@ -13,10 +13,10 @@
"healthz"
],
"integrity": {
- "config.js": "sha256-7792697d8640937cb0573d314897e8be96230b39418fc1b8a40b66c8138005df",
- "app.js": "sha256-662a5eb000a9ce90bc4a75c4cb5fbfbe097f5ebff5c06afa9452635872a60bf1",
- "styles.css": "sha256-ba90290ab11e82a6b2639dfd70d1e74502c1cacb2b26cf1db92b45beb67ac03f",
- "index.html": "sha256-3bd7fda6c92bdb3b65e61ac4457c855c99e5bbd82fcfc26959ce70d93b064aa9",
+ "config.js": "sha256-7310c02ad6a24ab7d1b49536bfea169ef28f7eb55ccbfe77cb8b0ee1918319ca",
+ "app.js": "sha256-7772e21dd629a56665264a714cf720b8a9cd20b6599e7d7aa5440211a6310651",
+ "styles.css": "sha256-c1b192a15de735819fbb7ec3cddc9574808ae366a8255b4f690a7a40bc21caba",
+ "index.html": "sha256-34b95523d465521a78845acc26fb6b9b8de84d3540c8b7acab224e2af28417a2",
"health.html": "sha256-c352a6f37aa24628cfc8d5709a70ff2d94181d192ed5fe00916cca9378a61d81",
"error.html": "sha256-f3cc28d2dfc9e8af112d0257c6ee47b9e9146e5da8590a754dd589bf43702aaf",
"healthz": "sha256-dc51b8c96c2d745df3bd5590d990230a482fd247123599548e0632fdbf97fc22"
diff --git a/apps/web/config.js b/apps/web/config.js
index f8d1883..89ba8af 100644
--- a/apps/web/config.js
+++ b/apps/web/config.js
@@ -1,5 +1,5 @@
/* Public, non-secret runtime configuration. Replace this file at deploy time if needed. */
window.__PROSPECT_CONFIG__ = Object.freeze({
apiBase: '',
- assetVersion: 'phase-34'
+ assetVersion: 'phase-35'
});
diff --git a/apps/web/index.html b/apps/web/index.html
index e80650e..2561b78 100644
--- a/apps/web/index.html
+++ b/apps/web/index.html
@@ -5,7 +5,7 @@
ProspectOS · Pipeline intelligence
-
+
@@ -58,7 +58,7 @@
SOURCE INTELLIGENCE
Discovery runs Build a criteria-first run across approved sources, monitor it live, and review every result with provenance before it enters your pipeline.
Bounded · review first
@@ -86,7 +86,7 @@
@@ -101,10 +101,7 @@
Discovery is disabled by default. No live source is enabled in this workspace. Enable a source only after its terms, owner, rate limit, and health have been reviewed.
-
+
REGISTRY STATUS
Configured & available integrations Configured sources are shown alongside optional adapters so unavailable capability is never mistaken for an active source.
Not loaded Sign in to load sources from the workspace.
RECENT OUTPUT
Recent source records
@@ -138,7 +135,7 @@
REVIEW REQUIRED
Confirm merge ×
This action is reversible. The merge will be recorded in history and can be reversed later.
Cancel Confirm merge
-
-
+
+