from __future__ import annotations import argparse, hashlib, json, os, re, secrets, sqlite3, sys, threading, time from datetime import datetime, timedelta, timezone from http.cookies import SimpleCookie from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from pathlib import Path from urllib.parse import parse_qs, urlparse 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, score_business, normalize_domain, normalize_phone, match_businesses from app.sources import adapter_for, contains_secret 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_VERSION from app.ai_assistance import generate as generate_ai, input_fingerprint, provider_status, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS 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 from app.config import load_config 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, score_business, normalize_domain, normalize_phone, match_businesses from .sources import adapter_for, contains_secret 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_VERSION from .ai_assistance import generate as generate_ai, input_fingerprint, provider_status, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS 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 from .config import load_config from .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 ORGANIZATION_ID = "demo-tenant" SCHEMA = Path(__file__).resolve().parents[1] / "schema.sql" SESSION_DAYS = 7 PBKDF2_ITERATIONS = 300_000 MUTATING_ROLES = {"owner", "admin", "researcher"} JOB_TYPES = {"noop", "prospect_recalculate", "source_discovery", "domain_check", "scoped_discovery"} JOB_PAGE_SIZE = 100 WEBSITE_SCAN_PAGE_SIZE = 100 WEBSITE_SCAN_CACHE_SECONDS = 3600 CONTACT_EXTRACTION_PAGE_SIZE = 100 SECRET_KEYS = {"password", "passwd", "secret", "token", "api_key", "apikey", "authorization", "credential", "private_key"} CHILD_TABLES = {"contacts": ("name", "email", "phone", "title", "do_not_contact"), "domains": ("domain", "kind"), "websites": ("url", "website_class"), "evidence": ("kind", "url", "claim"), "notes": ("body",)} def redact(value): if isinstance(value, dict): return {k: ("[REDACTED]" if str(k).lower() in SECRET_KEYS or any(s in str(k).lower() for s in ("password", "token", "secret", "api_key")) else redact(v)) for k,v in value.items()} if isinstance(value, list): return [redact(v) for v in value[:100]] if isinstance(value, str): return value[:2000] return value def job_json(row): result = row_json(row) try: result["payload"] = json.loads(result.get("payload") or "{}") except (ValueError, TypeError): result["payload"] = {} return result def hash_password(password: str, salt: bytes | None = None) -> tuple[str, str]: salt = salt or secrets.token_bytes(16); return hashlib.pbkdf2_hmac("sha256", password.encode(), salt, PBKDF2_ITERATIONS).hex(), salt.hex() def verify_password(password, encoded_hash, encoded_salt): try: return secrets.compare_digest(hashlib.pbkdf2_hmac("sha256", password.encode(), bytes.fromhex(encoded_salt), PBKDF2_ITERATIONS).hex(), encoded_hash) except (TypeError, ValueError): return False def connect(db_path: str) -> sqlite3.Connection: db = sqlite3.connect(db_path); db.row_factory = sqlite3.Row; db.execute("PRAGMA foreign_keys = ON"); db.executescript(SCHEMA.read_text()) # Upgrade databases created by Phase 1/2 without destroying data. cols = {r[1] for r in db.execute("PRAGMA table_info(businesses)")} for col, definition in (("verified", "INTEGER NOT NULL DEFAULT 0"), ("verified_at", "TEXT"), ("updated_at", "TEXT"), ("province", "TEXT NOT NULL DEFAULT ''"), ("city", "TEXT NOT NULL DEFAULT ''"), ("suburb", "TEXT NOT NULL DEFAULT ''"), ("merge_status", "TEXT NOT NULL DEFAULT 'active'"), ("merged_into_id", "INTEGER"), ("review_status", "TEXT NOT NULL DEFAULT 'pending'"), ("assigned_to", "TEXT NOT NULL DEFAULT ''"), ("review_metadata_json", "TEXT NOT NULL DEFAULT '{}'")): if col not in cols: db.execute(f"ALTER TABLE businesses ADD COLUMN {col} {definition}") db.execute("UPDATE businesses SET updated_at=COALESCE(updated_at,created_at) WHERE updated_at IS NULL") # Phase 12 is additive-safe for databases created before CRM metadata existed. for table, additions in { "pipeline_entries": (("notes", "TEXT NOT NULL DEFAULT ''"), ("next_action", "TEXT NOT NULL DEFAULT ''"), ("follow_up_at", "TEXT"), ("actor_user_id", "INTEGER"), ("idempotency_key", "TEXT"), ("version", "INTEGER NOT NULL DEFAULT 1")), "interactions": (("outcome", "TEXT NOT NULL DEFAULT 'other'"), ("notes", "TEXT NOT NULL DEFAULT ''"), ("next_action", "TEXT NOT NULL DEFAULT ''"), ("follow_up_at", "TEXT"), ("actor_user_id", "INTEGER"), ("idempotency_key", "TEXT")), "suppressions": (("active", "INTEGER NOT NULL DEFAULT 1"), ("updated_at", "TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP"), ("actor_user_id", "INTEGER")), }.items(): existing = {r[1] for r in db.execute(f"PRAGMA table_info({table})")} for col, definition in additions: if col not in existing: db.execute(f"ALTER TABLE {table} ADD COLUMN {col} {definition}") db.execute("CREATE UNIQUE INDEX IF NOT EXISTS uq_pipeline_idempotency ON pipeline_entries(organization_id,idempotency_key) WHERE idempotency_key IS NOT NULL AND idempotency_key <> ''") db.execute("CREATE UNIQUE INDEX IF NOT EXISTS uq_interaction_idempotency ON interactions(organization_id,idempotency_key) WHERE idempotency_key IS NOT NULL AND idempotency_key <> ''") db.execute("INSERT OR IGNORE INTO organizations (id,name) VALUES (?,?)", (ORGANIZATION_ID, "Demo organization")) for organization in db.execute("SELECT id FROM organizations").fetchall(): for rule in DEFAULT_RULES: db.execute("INSERT OR IGNORE INTO score_rules(organization_id,code,name,description,condition_json,points,max_applications,enabled,version) VALUES(?,?,?,?,?,?,?,?,?)", (organization["id"], rule["code"], rule["name"], rule["description"], json.dumps(rule["condition_json"], sort_keys=True), rule["points"], rule["max_applications"], rule["enabled"], rule["version"])) email, password = os.environ.get("BOOTSTRAP_ADMIN_EMAIL"), os.environ.get("BOOTSTRAP_ADMIN_PASSWORD") if email and password and not db.execute("SELECT id FROM users WHERE email=?", (email.strip().lower(),)).fetchone(): ph, salt = hash_password(password); db.execute("INSERT INTO users (organization_id,email,password_hash,password_salt,role) VALUES (?,?,?,?,?)", (ORGANIZATION_ID,email.strip().lower(),ph,salt,"owner")) db.commit(); return db def safe_value(value): if isinstance(value, bytes): return value.decode("utf-8", "replace") return value def row_json(row): result = {k: safe_value(v) for k, v in dict(row).items()} if "score_factors" in result: try: result["score_factors"] = json.loads(result["score_factors"] or "[]") except (TypeError, ValueError): result["score_factors"] = [] for key in ("verified", "do_not_contact"): if key in result: result[key] = bool(result[key]) return result class ApiHandler(BaseHTTPRequestHandler): server_version = "ProspectPlatform/0.1" def send_json(self, status, payload, extra_headers=None): body = json.dumps(payload, sort_keys=True, default=str).encode(); self.send_response(status); self.send_header("Content-Type","application/json; charset=utf-8"); self.send_header("Cache-Control","no-store, private"); self.send_header("Pragma","no-cache"); self.send_header("Vary","Cookie, Origin"); self.send_header("Access-Control-Allow-Origin",os.environ.get("CORS_ORIGINS","http://localhost:8080")); self.send_header("Access-Control-Allow-Credentials","true"); self.send_header("Access-Control-Allow-Methods","GET, POST, PATCH, OPTIONS"); self.send_header("Access-Control-Allow-Headers","Content-Type") for k,v in (extra_headers or {}).items(): self.send_header(k,v) self.send_header("Content-Length",str(len(body))); self.end_headers(); self.wfile.write(body) def read_json(self): try: value=json.loads(self.rfile.read(int(self.headers.get("Content-Length","0"))) or b"{}"); return value if isinstance(value,dict) else {} except (ValueError,json.JSONDecodeError): return {} def db(self): return connect(getattr(self.server,"db_path")) def do_OPTIONS(self): self.send_response(204); self.send_header("Access-Control-Allow-Methods","GET, POST, PATCH, OPTIONS"); self.end_headers() def session_user(self, db): cookie=SimpleCookie(); cookie.load(self.headers.get("Cookie","")); token=cookie.get("session") if not token: return None now=datetime.now(timezone.utc).replace(microsecond=0).isoformat(); h=hashlib.sha256(token.value.encode()).hexdigest() return db.execute("SELECT u.id,u.email,u.role,u.organization_id FROM sessions s JOIN users u ON u.id=s.user_id WHERE s.token_hash=? AND s.expires_at>?",(h,now)).fetchone() def require_auth(self,db): user=self.session_user(db) if not user: self.send_json(401,{"error":"unauthorized"}); return None return user def auth_cookie(self,token,max_age): return f"session={token}; Max-Age={max_age}; Path=/; HttpOnly; SameSite=Lax" def audit(self, db, user, action, details=""): db.execute("INSERT INTO audit_log (organization_id,user_id,action,details) VALUES (?,?,?,?)",(user["organization_id"],user["id"],action,details)) def business(self, db, ident, org): return db.execute("SELECT * FROM businesses WHERE id=? AND organization_id=?",(ident,org)).fetchone() def nested(self, db, bid, org): result={"contacts":[],"contact_extractions":[],"domains":[],"websites":[],"evidence":[],"pipeline":[],"interactions":[],"notes":[]} tables={"contacts":"contacts","contact_extractions":"contact_extractions","domains":"domains","websites":"websites","evidence":"evidence","pipeline":"pipeline_entries","interactions":"interactions","notes":"notes"} for key, table in tables.items(): result[key]=[row_json(r) for r in db.execute(f"SELECT * FROM {table} WHERE business_id=? AND organization_id=? ORDER BY id",(bid,org))] return result def _bounded_filters(self, value): def walk(item, depth=0): if depth > 4: raise ValueError("filters_too_deep") if isinstance(item, dict): if len(item) > 30: raise ValueError("filters_too_large") return {str(k)[:80]: walk(v, depth + 1) for k, v in item.items()} if isinstance(item, list): if len(item) > 50: raise ValueError("filters_too_large") return [walk(v, depth + 1) for v in item] if isinstance(item, str): if len(item) > 500: raise ValueError("filter_value_too_large") return item if item is None or isinstance(item, (bool, int, float)): return item raise ValueError("invalid_filters") if not isinstance(value, dict): raise ValueError("invalid_filters") result = walk(value) if len(json.dumps(result, separators=(",", ":"), ensure_ascii=False).encode()) > 8192: raise ValueError("filters_too_large") return result def _saved_filter_json(self, row): item = row_json(row) try: item["filters"] = json.loads(item.pop("filters_json") or "{}") except (TypeError, ValueError): item["filters"] = {} return item def list_saved_filters(self, db, user): rows = db.execute("SELECT * FROM saved_filters WHERE organization_id=? AND user_id=? ORDER BY updated_at DESC,id DESC", (user["organization_id"], user["id"])).fetchall() return self.send_json(200, {"organization_id": user["organization_id"], "items": [self._saved_filter_json(r) for r in rows]}) def save_filter(self, payload, db, user): name = str(payload.get("name", "")).strip() raw = payload.get("filters", payload.get("filter", payload.get("filters_json", {}))) if not name or len(name) > 120: return self.send_json(400, {"error": "invalid_saved_filter"}) try: filters = self._bounded_filters(raw) except ValueError as exc: return self.send_json(400, {"error": str(exc)}) try: cur = db.execute("INSERT INTO saved_filters(organization_id,user_id,name,filters_json) VALUES(?,?,?,?)", (user["organization_id"], user["id"], name, json.dumps(filters, sort_keys=True, separators=(",", ":")))) except sqlite3.IntegrityError: return self.send_json(409, {"error": "duplicate_saved_filter"}) self.audit(db, user, "saved_filter.created", str(cur.lastrowid)); db.commit() return self.send_json(201, self._saved_filter_json(db.execute("SELECT * FROM saved_filters WHERE id=?", (cur.lastrowid,)).fetchone())) def update_saved_filter(self, fid, payload, db, user): row = db.execute("SELECT * FROM saved_filters WHERE id=? AND organization_id=? AND user_id=?", (fid, user["organization_id"], user["id"])).fetchone() if not row: return self.send_json(404, {"error": "not_found"}) fields = []; values = [] if "name" in payload: name = str(payload["name"]).strip() if not name or len(name) > 120: return self.send_json(400, {"error": "invalid_saved_filter"}) fields.append("name=?"); values.append(name) if any(k in payload for k in ("filters", "filter", "filters_json")): try: filters = self._bounded_filters(payload.get("filters", payload.get("filter", payload.get("filters_json")))) except ValueError as exc: return self.send_json(400, {"error": str(exc)}) fields.append("filters_json=?"); values.append(json.dumps(filters, sort_keys=True, separators=(",", ":"))) if not fields: return self.send_json(400, {"error": "no_changes"}) values += [fid, user["organization_id"], user["id"]] try: db.execute("UPDATE saved_filters SET " + ",".join(fields) + ",updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=? AND user_id=?", values) except sqlite3.IntegrityError: return self.send_json(409, {"error": "duplicate_saved_filter"}) self.audit(db, user, "saved_filter.updated", str(fid)); db.commit() return self.send_json(200, self._saved_filter_json(db.execute("SELECT * FROM saved_filters WHERE id=?", (fid,)).fetchone())) def delete_saved_filter(self, fid, db, user): row = db.execute("SELECT id FROM saved_filters WHERE id=? AND organization_id=? AND user_id=?", (fid, user["organization_id"], user["id"])).fetchone() if not row: return self.send_json(404, {"error": "not_found"}) db.execute("DELETE FROM saved_filters WHERE id=? AND organization_id=? AND user_id=?", (fid, user["organization_id"], user["id"])) self.audit(db, user, "saved_filter.deleted", str(fid)); db.commit() return self.send_json(200, {"ok": True, "id": fid}) def _review_item(self, row, db, org): item = row_json(row) suppressed = is_suppressed(dict(row), [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))]) merged = row["merge_status"] == "merged" item.update({"suppressed": bool(suppressed), "merged": bool(merged), "outreach_eligible": not suppressed and not merged and row["review_status"] not in ("rejected",), "review_flags": [x for x, yes in (("suppressed", suppressed), ("merged", merged), ("rejected", row["review_status"] == "rejected")) if yes]}) return item def review_queue(self, db, org, query): try: limit = int((query.get("page_size") or [50])[0]); offset = int((query.get("offset") or [0])[0]) if limit < 1 or limit > 100 or offset < 0: raise ValueError except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"}) where, params = ["b.organization_id=?"], [org] status = (query.get("review_status") or query.get("status") or [""])[0].strip() if status: if status not in {"pending", "verified", "rejected", "assigned"}: return self.send_json(400, {"error": "invalid_review_status"}) where.append("b.review_status=?"); params.append(status) for key, op in (("score_min", ">="), ("score_max", "<=")): raw = (query.get(key) or [""])[0] if raw: try: value = int(raw) except ValueError: return self.send_json(400, {"error": "invalid_score"}) if value < 0 or value > 100: return self.send_json(400, {"error": "invalid_score"}) where.append("b.score" + op + "?"); params.append(value) priority = (query.get("priority") or query.get("priority_band") or [""])[0].strip() if priority: bounds = {"high": (70, 100), "medium": (40, 69), "low": (0, 39)} if priority not in bounds: return self.send_json(400, {"error": "invalid_priority"}) where += ["b.score BETWEEN ? AND ?"]; params += list(bounds[priority]) website = (query.get("website_state") or query.get("website_class") or [""])[0].strip() if website: where.append("b.website_class=?"); params.append(website) freshness = (query.get("freshness_days") or [""])[0] if freshness: try: days = int(freshness) except ValueError: return self.send_json(400, {"error": "invalid_freshness"}) if days < 0 or days > 3650: return self.send_json(400, {"error": "invalid_freshness"}) where.append("b.updated_at >= datetime('now', ?)"); params.append(f"-{days} days") source_health = (query.get("source_health") or [""])[0].strip() if source_health: where.append("EXISTS (SELECT 1 FROM sources s WHERE s.organization_id=b.organization_id AND s.health_status=?)"); params.append(source_health) rows = db.execute("SELECT b.* FROM businesses b WHERE " + " AND ".join(where) + " ORDER BY b.score DESC,b.updated_at DESC,b.id DESC LIMIT ? OFFSET ?", params + [limit + 1, offset]).fetchall() return self.send_json(200, {"organization_id": org, "items": [self._review_item(r, db, org) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit}) OUTCOMES = {"connected", "no_answer", "left_message", "meeting_booked", "meeting_held", "qualified", "disqualified", "won", "lost", "other"} STAGES = {"new", "contacted", "qualified", "proposal", "negotiation", "won", "lost"} def _crm_suppressed(self, db, org, business): return is_suppressed(dict(business), [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))]) def _date_where(self, query, column, params): start = (query.get("from") or query.get("start") or [""])[0] end = (query.get("to") or query.get("end") or [""])[0] for value in (start, end): if value and (len(value) > 30 or not re.match(r"^\\d{4}-\\d{2}-\\d{2}(?:T[^ ]*)?$", value)): raise ValueError("invalid_date") if start and end: try: left=datetime.fromisoformat(start.replace("Z","+00:00")).date(); right=datetime.fromisoformat(end.replace("Z","+00:00")).date() if (right-left).days > 366: raise ValueError("date_range_too_large") except ValueError as exc: if str(exc) == "date_range_too_large": raise raise ValueError("invalid_date_range") if start: params.append(start); clause = f"{column}>=?" else: clause = "" if end: params.append(end); clause += (" AND " if clause else "") + f"{column}<=?" return clause def list_pipeline(self, db, org, query): params=[org]; where=["p.organization_id=?"] bid=(query.get("business_id") or [""])[0] if bid.isdigit(): where.append("p.business_id=?"); params.append(int(bid)) stage=(query.get("stage") or [""])[0] if stage: where.append("p.stage=?"); params.append(stage) rows=db.execute("SELECT p.* FROM pipeline_entries p WHERE " + " AND ".join(where) + " ORDER BY p.updated_at DESC,p.id DESC",params).fetchall() return self.send_json(200,{"organization_id":org,"items":[row_json(r) for r in rows]}) def list_stages(self, db, org): rows=db.execute("SELECT * FROM pipeline_stages WHERE organization_id=? AND active=1 ORDER BY position,id",(org,)).fetchall() if not rows: return self.send_json(200,{"organization_id":org,"items":[{"name":x,"position":i,"active":True} for i,x in enumerate(("new","contacted","qualified","proposal","negotiation","won","lost"))]}) return self.send_json(200,{"organization_id":org,"items":[row_json(r) for r in rows]}) def list_interactions(self, db, org, query): params=[org]; where=["i.organization_id=?"] bid=(query.get("business_id") or [""])[0] if bid.isdigit(): where.append("i.business_id=?"); params.append(int(bid)) outcome=(query.get("outcome") or [""])[0] if outcome: where.append("i.outcome=?"); params.append(outcome) rows=db.execute("SELECT i.* FROM interactions i WHERE " + " AND ".join(where) + " ORDER BY i.created_at DESC,i.id DESC",params).fetchall() return self.send_json(200,{"organization_id":org,"items":[row_json(r) for r in rows]}) def list_suppressions(self, db, org, query): rows=db.execute("SELECT * FROM suppressions WHERE organization_id=? ORDER BY id DESC",(org,)).fetchall() return self.send_json(200,{"organization_id":org,"items":[row_json(r) for r in rows]}) def report(self, db, org, kind, query): try: if kind == "pipeline": params=[org]; clause=self._date_where(query,"p.created_at",params); sql="SELECT p.stage, p.status, COUNT(*) count FROM pipeline_entries p WHERE p.organization_id=?" + ((" AND "+clause) if clause else "") + " GROUP BY p.stage,p.status ORDER BY p.stage,p.status" elif kind == "outcomes": params=[org]; clause=self._date_where(query,"i.created_at",params); sql="SELECT i.outcome, COUNT(*) count FROM interactions i WHERE i.organization_id=?" + ((" AND "+clause) if clause else "") + " GROUP BY i.outcome ORDER BY i.outcome" else: params=[org]; clause=self._date_where(query,"activity_at",params); sql="SELECT activity_type, COUNT(*) count FROM (SELECT 'pipeline' activity_type, created_at activity_at FROM pipeline_entries WHERE organization_id=? UNION ALL SELECT 'interaction',created_at FROM interactions WHERE organization_id=?) WHERE 1=1" + ((" AND "+clause) if clause else "") + " GROUP BY activity_type" params=[org,org] + params[1:] rows=db.execute(sql,params).fetchall(); return self.send_json(200,{"organization_id":org,"items":[dict(r) for r in rows]}) except ValueError as exc: return self.send_json(400,{"error":str(exc)}) def _ai_run_json(self, row, suggestions=None): item = row_json(row) for key in ("input_evidence_hashes_json", "prompt_metadata_json", "data_minimization_json", "output_json"): source = item.pop(key, "{}" if key != "input_evidence_hashes_json" else "[]") try: item[key[:-5] if key.endswith("_json") else key] = json.loads(source or ("{}" if key != "input_evidence_hashes_json" else "[]")) except (TypeError, ValueError): item[key[:-5] if key.endswith("_json") else key] = {} if key != "input_evidence_hashes_json" else [] if suggestions is not None: item["suggestions"] = suggestions return item # Phase 14: outreach is a reviewable preparation workflow only. There is # deliberately no provider client in this service and send never performs # network I/O. OUTREACH_PROVIDERS = {"smtp", "sendgrid", "twilio", "whatsapp"} OUTREACH_KINDS = {"email", "phone", "whatsapp"} OUTREACH_DAILY_CAP = 100 OUTREACH_BATCH_CAP = 25 def _provider_json(self, row, org): if not row: return {"organization_id": org, "provider": "", "enabled": False, "policy": {"consent_required": True}, "daily_cap": self.OUTREACH_DAILY_CAP, "batch_cap": self.OUTREACH_BATCH_CAP} item = {"id": row["id"], "organization_id": org, "provider": row["provider"], "enabled": bool(row["enabled"]), "daily_cap": row["daily_cap"], "batch_cap": row["batch_cap"]} try: item["policy"] = json.loads(row["policy_json"] or "{}") except (TypeError, ValueError): item["policy"] = {"consent_required": True} return item def provider_config(self, db, user, payload=None): org = user["organization_id"] row = db.execute("SELECT * FROM outreach_provider_configs WHERE organization_id=?", (org,)).fetchone() if payload is None: return self.send_json(200, self._provider_json(row, org)) provider = str(payload.get("provider", row["provider"] if row else "")).strip().lower() if provider and provider not in self.OUTREACH_PROVIDERS: return self.send_json(400, {"error": "invalid_provider"}) enabled = bool(payload.get("enabled", bool(row["enabled"]) if row else False)) policy = payload.get("policy", payload.get("legal_policy", {} if not row else None)) if policy is None: try: policy = json.loads(row["policy_json"] or "{}") except (TypeError, ValueError): policy = {"consent_required": True} if not isinstance(policy, dict) or len(policy) > 20: return self.send_json(400, {"error": "invalid_policy"}) policy = {str(k)[:80]: v for k, v in policy.items()} policy.setdefault("consent_required", True) try: daily = int(payload.get("daily_cap", row["daily_cap"] if row else self.OUTREACH_DAILY_CAP)); batch = int(payload.get("batch_cap", row["batch_cap"] if row else self.OUTREACH_BATCH_CAP)) except (TypeError, ValueError): return self.send_json(400, {"error": "invalid_limits"}) if daily < 1 or daily > self.OUTREACH_DAILY_CAP or batch < 1 or batch > self.OUTREACH_BATCH_CAP: return self.send_json(400, {"error": "invalid_limits"}) secret = payload.get("secret", None) fingerprint = row["secret_fingerprint"] if row else "" if secret is not None: if not isinstance(secret, str) or not secret or len(secret) > 4096: return self.send_json(400, {"error": "invalid_secret"}) fingerprint = hashlib.sha256(secret.encode()).hexdigest() if enabled and (not provider or not fingerprint): return self.send_json(400, {"error": "provider_credentials_required"}) if row: db.execute("UPDATE outreach_provider_configs SET provider=?,enabled=?,secret_fingerprint=?,policy_json=?,daily_cap=?,batch_cap=?,updated_at=CURRENT_TIMESTAMP WHERE organization_id=?", (provider, int(enabled), fingerprint, json.dumps(policy, sort_keys=True), daily, batch, org)) else: db.execute("INSERT INTO outreach_provider_configs(organization_id,provider,enabled,secret_fingerprint,policy_json,daily_cap,batch_cap) VALUES(?,?,?,?,?,?,?)", (org, provider, int(enabled), fingerprint, json.dumps(policy, sort_keys=True), daily, batch)) self.audit(db, user, "outreach.provider_config.updated", provider or "disabled"); db.commit() return self.send_json(200, self._provider_json(db.execute("SELECT * FROM outreach_provider_configs WHERE organization_id=?", (org,)).fetchone(), org)) def _draft_json(self, row): item = row_json(row) for field, default in (("template_json", {}), ("citations_json", []), ("provenance_json", {})): key = field[:-5] try: item[key] = json.loads(item.pop(field) or json.dumps(default)) except (TypeError, ValueError): item[key] = default item["target_verified"] = bool(item.get("target_verified")); item["consent_confirmed"] = bool(item.get("consent_confirmed")) item["target"] = {"kind": item.pop("target_kind"), "value": item.pop("target_value"), "verified": item["target_verified"]} return item def _template(self, text, business, evidence): pattern = re.compile(r"{{\s*([^}]+?)\s*}}") used = [] def replace(match): token = match.group(1).strip() if token == "business.name": return str(business["name"]) m = re.fullmatch(r"evidence\.(\d+)\.claim", token) if m: index = int(m.group(1)); if index < 1 or index > len(evidence): raise ValueError("unsupported_template_variable") used.append(evidence[index - 1]["id"]); return str(evidence[index - 1]["claim"]) raise ValueError("unsupported_template_variable") return pattern.sub(replace, text), sorted(set(used)) def _target_exists(self, db, bid, org, kind, value): if kind == "email": return bool(db.execute("SELECT id FROM businesses WHERE id=? AND organization_id=? AND email=?", (bid, org, value)).fetchone() or db.execute("SELECT id FROM contacts WHERE business_id=? AND organization_id=? AND email=?", (bid, org, value)).fetchone() or db.execute("SELECT id FROM contact_extractions WHERE business_id=? AND organization_id=? AND kind='email' AND value=?", (bid, org, value)).fetchone()) return bool(db.execute("SELECT id FROM businesses WHERE id=? AND organization_id=? AND phone=?", (bid, org, value)).fetchone() or db.execute("SELECT id FROM contacts WHERE business_id=? AND organization_id=? AND phone=?", (bid, org, value)).fetchone() or db.execute("SELECT id FROM contact_extractions WHERE business_id=? AND organization_id=? AND kind IN ('phone','whatsapp') AND value=?", (bid, org, value)).fetchone()) def create_outreach_draft(self, bid, payload, db, user): org = user["organization_id"]; business = self.business(db, bid, org) if not business: return self.send_json(404, {"error": "not_found"}) target = payload.get("target", {}); kind = str(target.get("kind", "email")).lower(); value = str(target.get("value", "")).strip().lower() if kind not in self.OUTREACH_KINDS or not value or len(value) > 320 or not isinstance(target.get("verified", False), bool): return self.send_json(400, {"error": "invalid_target"}) subject, body = str(payload.get("subject", "")).strip(), str(payload.get("body", "")).strip() if not subject or not body or len(subject) > 500 or len(body) > 10000: return self.send_json(400, {"error": "invalid_draft"}) if contains_secret(payload): return self.send_json(400, {"error": "secret_not_permitted"}) try: evidence = [dict(r) for r in db.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id", (bid, org))] subject, sids = self._template(subject, business, evidence); body, bids = self._template(body, business, evidence); ids = sorted(set(sids + bids)) except ValueError as exc: return self.send_json(400, {"error": str(exc)}) key = str(payload.get("idempotency_key", "")).strip() or hashlib.sha256(json.dumps({"business_id": bid, "target": target, "subject": subject, "body": body}, sort_keys=True).encode()).hexdigest() if len(key) > 200: return self.send_json(400, {"error": "invalid_idempotency_key"}) existing = db.execute("SELECT * FROM outreach_drafts WHERE organization_id=? AND idempotency_key=?", (org, key)).fetchone() if existing: return self.send_json(200, self._draft_json(existing)) if db.execute("SELECT COUNT(*) FROM outreach_drafts WHERE organization_id=? AND created_at>=datetime('now','-1 day')", (org,)).fetchone()[0] >= self.OUTREACH_DAILY_CAP: return self.send_json(429, {"error": "outreach_daily_cap"}) cur = db.execute("INSERT INTO outreach_drafts(organization_id,business_id,target_kind,target_value,target_verified,subject,body,template_json,citations_json,provenance_json,legal_basis,consent_confirmed,actor_user_id,idempotency_key) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (org, bid, kind, value, int(target["verified"]), subject, body, json.dumps({"variables": ids}, sort_keys=True), json.dumps(ids), json.dumps({str(x): {"type": "evidence", "id": x} for x in ids}, sort_keys=True), str(payload.get("legal_basis", ""))[:100], int(bool(payload.get("consent_confirmed", False))), user["id"], key)) self.audit(db, user, "outreach_draft.created", str(cur.lastrowid)); db.commit() return self.send_json(201, self._draft_json(db.execute("SELECT * FROM outreach_drafts WHERE id=?", (cur.lastrowid,)).fetchone())) def list_outreach_drafts(self, db, user, query): try: limit = int((query.get("page_size") or [self.OUTREACH_BATCH_CAP])[0]); offset = int((query.get("offset") or [0])[0]) except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"}) if limit < 1 or limit > self.OUTREACH_BATCH_CAP or offset < 0: return self.send_json(400, {"error": "invalid_pagination"}) rows = db.execute("SELECT * FROM outreach_drafts WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?", (user["organization_id"], limit + 1, offset)).fetchall() return self.send_json(200, {"organization_id": user["organization_id"], "items": [self._draft_json(r) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit}) def update_outreach_draft(self, did, payload, db, user): row = db.execute("SELECT * FROM outreach_drafts WHERE id=? AND organization_id=?", (did, user["organization_id"])).fetchone() if not row: return self.send_json(404, {"error": "not_found"}) if row["status"] in ("approved", "sent"): return self.send_json(409, {"error": "draft_locked"}) fields, values = [], [] subject, body = row["subject"], row["body"] if "subject" in payload: subject = str(payload["subject"]).strip() if "body" in payload: body = str(payload["body"]).strip() if ("subject" in payload and (not subject or len(subject) > 500)) or ("body" in payload and (not body or len(body) > 10000)): return self.send_json(400, {"error": "invalid_draft"}) if "subject" in payload or "body" in payload: business = self.business(db, row["business_id"], user["organization_id"]) evidence = [dict(r) for r in db.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id", (row["business_id"], user["organization_id"]))] try: subject, sids = self._template(subject, business, evidence); body, bids = self._template(body, business, evidence) except ValueError as exc: return self.send_json(400, {"error": str(exc)}) ids = sorted(set(sids + bids)) fields += ["subject=?", "body=?", "template_json=?", "citations_json=?", "provenance_json=?"] values += [subject, body, json.dumps({"variables": ids}, sort_keys=True), json.dumps(ids), json.dumps({str(x): {"type": "evidence", "id": x} for x in ids}, sort_keys=True)] if not fields: return self.send_json(400, {"error": "no_changes"}) fields.append("updated_at=CURRENT_TIMESTAMP"); values += [did, user["organization_id"]] db.execute("UPDATE outreach_drafts SET " + ",".join(fields) + " WHERE id=? AND organization_id=?", values); self.audit(db, user, "outreach_draft.updated", str(did)); db.commit() return self.send_json(200, self._draft_json(db.execute("SELECT * FROM outreach_drafts WHERE id=?", (did,)).fetchone())) def approve_outreach_draft(self, did, db, user): row = db.execute("SELECT * FROM outreach_drafts WHERE id=? AND organization_id=?", (did, user["organization_id"])).fetchone() if not row: return self.send_json(404, {"error": "not_found"}) if row["status"] != "pending_review": return self.send_json(409, {"error": "draft_not_reviewable"}) db.execute("UPDATE outreach_drafts SET status='approved',approved_by=?,approved_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?", (user["id"], did, user["organization_id"])); self.audit(db, user, "outreach_draft.approved", str(did)); db.commit() return self.send_json(200, self._draft_json(db.execute("SELECT * FROM outreach_drafts WHERE id=?", (did,)).fetchone())) def send_outreach_draft(self, did, db, user): org = user["organization_id"]; row = db.execute("SELECT * FROM outreach_drafts WHERE id=? AND organization_id=?", (did, org)).fetchone() if not row: return self.send_json(404, {"error": "not_found"}) cfg = db.execute("SELECT * FROM outreach_provider_configs WHERE organization_id=?", (org,)).fetchone(); reasons = [] if not cfg or not cfg["enabled"] or not cfg["provider"] or not cfg["secret_fingerprint"]: reasons.append("provider") business = self.business(db, row["business_id"], org) if not row["target_verified"] or not self._target_exists(db, row["business_id"], org, row["target_kind"], row["target_value"]): reasons.append("verified_target") suppressed = is_suppressed({"email": row["target_value"] if row["target_kind"] == "email" else "", "phone": row["target_value"] if row["target_kind"] != "email" else "", "website_domain": business["website_domain"] if business else ""}, [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))]) if suppressed: reasons.append("suppressed") policy = {} if cfg: try: policy = json.loads(cfg["policy_json"] or "{}") except (TypeError, ValueError): policy = {"consent_required": True} if policy.get("consent_required", True) and (not row["consent_confirmed"] or not row["legal_basis"]): reasons.append("consent_or_legal_policy") if row["status"] != "approved": reasons.append("approved_draft") if cfg and db.execute("SELECT COUNT(*) FROM outreach_drafts WHERE organization_id=? AND status='sent' AND sent_at>=datetime('now','-1 day')", (org,)).fetchone()[0] >= cfg["daily_cap"]: reasons.append("daily_cap") if reasons: status = "not_configured" if "provider" in reasons and (not cfg or not cfg["enabled"] or not cfg["provider"] or not cfg["secret_fingerprint"]) else "blocked"; self.audit(db, user, "outreach_draft.send_blocked", f"{did}:{','.join(reasons)}"); db.commit() return self.send_json(409, {"status": status, "blocked_reasons": reasons, "network_send": False, "id": did}) reasons.append("network_send_disabled"); self.audit(db, user, "outreach_draft.send_blocked", f"{did}:network_send_disabled"); db.commit() return self.send_json(409, {"status": "blocked", "blocked_reasons": reasons, "network_send": False, "id": did}) def remote_ai_provider_config(self, db, user, payload=None, connectivity=False): org = user["organization_id"] row = db.execute("SELECT * FROM ai_remote_provider_configs WHERE organization_id=?", (org,)).fetchone() if connectivity: if user["role"] not in {"owner", "admin"}: return self.send_json(403, {"error": "forbidden"}) return self.send_json(200, test_remote_connectivity(row)) if payload is not None: if user["role"] not in {"owner", "admin"}: return self.send_json(403, {"error": "forbidden"}) try: config = validate_remote_provider(payload) credentials = dict(config["credentials"]) if row: try: old = json.loads(decrypt_provider_secret(row["credentials_ciphertext"])) if row["credentials_ciphertext"] else {} except Exception: old = {} for name in ("nous_api_key", "firecrawl_api_key"): if name not in credentials and name in old: credentials[name] = old[name] if config["enabled"] and any(name not in credentials for name in ("nous_api_key", "firecrawl_api_key")): return self.send_json(400, {"error": "provider_credentials_required"}) ciphertext = encrypt_provider_secret(json.dumps(credentials, sort_keys=True)) if credentials else "" fingerprint = hashlib.sha256(json.dumps(credentials, sort_keys=True).encode()).hexdigest() if credentials else "" except (ValueError, TypeError, RuntimeError) as exc: return self.send_json(400, {"error": str(exc)}) db.execute("INSERT INTO ai_remote_provider_configs(organization_id,provider,model,enabled,nous_base_url,firecrawl_base_url,credentials_ciphertext,credentials_fingerprint) VALUES(?,?,?,?,?,?,?,?) ON CONFLICT(organization_id) DO UPDATE SET provider=excluded.provider,model=excluded.model,enabled=excluded.enabled,nous_base_url=excluded.nous_base_url,firecrawl_base_url=excluded.firecrawl_base_url,credentials_ciphertext=excluded.credentials_ciphertext,credentials_fingerprint=excluded.credentials_fingerprint,updated_at=CURRENT_TIMESTAMP", (org, config["provider"], config["model"], int(config["enabled"]), config["nous_base_url"], config["firecrawl_base_url"], ciphertext, fingerprint)) self.audit(db, user, "ai.remote_provider.updated", config["provider"]); db.commit() row = db.execute("SELECT * FROM ai_remote_provider_configs WHERE organization_id=?", (org,)).fetchone() return self.send_json(200, dict(remote_provider_status(row), organization_id=org)) def ai_provider_config(self, db, user, payload=None): remote = db.execute("SELECT * FROM ai_remote_provider_configs WHERE organization_id=?", (user["organization_id"],)).fetchone() if remote or (payload and str(payload.get("provider", "")).lower() in {"nous_portal", "nous_portal_web_research"}): return self.remote_ai_provider_config(db, user, payload) org = user["organization_id"] row = db.execute("SELECT * FROM ai_provider_configs WHERE organization_id=?", (org,)).fetchone() if payload is not None: provider = str(payload.get("provider", "local")).strip().lower() if not provider or len(provider) > 80: return self.send_json(400, {"error": "invalid_provider"}) reviewed = bool(payload.get("reviewed", False)) enabled = bool(payload.get("enabled", True)) # A remote provider cannot be enabled by configuration alone: this # service has no reviewed transport abstraction and never sends data. if provider not in {"local", "deterministic"} and (enabled or reviewed): return self.send_json(409, {"error": "provider_not_reviewed", "network_enabled": False}) db.execute("INSERT INTO ai_provider_configs(organization_id,provider,enabled,reviewed) VALUES(?,?,?,?) ON CONFLICT(organization_id) DO UPDATE SET provider=excluded.provider,enabled=excluded.enabled,reviewed=excluded.reviewed,updated_at=CURRENT_TIMESTAMP", (org, provider, int(enabled), int(reviewed))) self.audit(db, user, "ai_provider.updated", provider); db.commit() row = db.execute("SELECT * FROM ai_provider_configs WHERE organization_id=?", (org,)).fetchone() configured = row["provider"] if row and row["enabled"] else None result = provider_status(configured if configured is not None else None) result.update({"organization_id": org, "configured": bool(row), "enabled": bool(row and row["enabled"]), "reviewed": bool(row and row["reviewed"])}) return self.send_json(200, result) def _ai_current_fingerprint(self, db, row): if not row["business_id"]: return "" try: limit = int(json.loads(row["prompt_metadata_json"] or "{}").get("request", {}).get("max_items", MAX_INPUT_ITEMS)) except (TypeError, ValueError): limit = MAX_INPUT_ITEMS limit = max(1, min(limit, MAX_INPUT_ITEMS)) org, bid = row["organization_id"], row["business_id"] business = self.business(db, bid, org) if not business: return "" scans = [dict(x) for x in db.execute("SELECT id,input_url,classification,result_json,scanned_at FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, limit))] evidence = [dict(x) for x in db.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, limit))] contacts = [dict(x) for x in db.execute("SELECT * FROM contacts WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, limit))] extracted = [dict(x) for x in db.execute("SELECT * FROM contact_extractions WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, limit))] return input_fingerprint(row_json(business), scans, contacts + extracted, evidence, limit) def suggest_ai(self, bid, payload, db, user): org = user["organization_id"] business = self.business(db, bid, org) if not business: return self.send_json(404, {"error": "not_found"}) business_dict = row_json(business) suppressions = [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))] contacts = [dict(r) for r in db.execute("SELECT * FROM contacts WHERE business_id=? AND organization_id=?", (bid, org))] extracted = [dict(r) for r in db.execute("SELECT * FROM contact_extractions WHERE business_id=? AND organization_id=?", (bid, org))] if is_suppressed(business_dict, suppressions) or any(bool(c.get("do_not_contact") or c.get("suppressed")) for c in contacts + extracted): return self.send_json(409, {"error": "ai_blocked_suppressed_business", "business_id": bid}) try: requested = int(payload.get("max_items", MAX_INPUT_ITEMS)) if requested < 1 or requested > MAX_INPUT_ITEMS: raise ValueError except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_limits"}) scans = [dict(r) for r in db.execute("SELECT id,input_url,classification,result_json,scanned_at FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, requested)).fetchall()] evidence = [dict(r) for r in db.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, requested)).fetchall()] history = [dict(r) for r in db.execute("SELECT score,eligible,priority_band,score_version,explanations_json,created_at FROM score_history WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, requested)).fetchall()] status, provider, version, metadata = generate_ai(business_dict, scans[:requested], contacts[:requested] + extracted[:requested], evidence[:requested], history[:requested]) output = metadata.pop("output", {}) if status == "succeeded" else {"suggestions": [], "grounded": True, "claim_policy": "stored_evidence_only"} hashes = metadata.get("evidence_hashes", []) 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(hashes, sort_keys=True), "local-deterministic" if provider else "", provider, version, json.dumps({"request": {"max_items": requested}, "output_limit": MAX_OUTPUT_CHARS}, sort_keys=True), json.dumps(metadata, sort_keys=True), status, "pending", json.dumps(output, sort_keys=True), user["id"])) run_id = cur.lastrowid for suggestion in output.get("suggestions", []): db.execute("INSERT INTO ai_suggestions(ai_run_id,organization_id,business_id,suggestion_type,citations_json,output_json) VALUES(?,?,?,?,?,?)", (run_id, org, bid, str(suggestion.get("type", ""))[:80], json.dumps(suggestion.get("citations", []), sort_keys=True), json.dumps(suggestion, sort_keys=True))) self.audit(db, user, "ai.suggested", str(run_id)); db.commit() 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 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])) if limit < 1 or limit > 100: raise ValueError except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"}) rows = db.execute("SELECT * FROM ai_runs WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?", (user["organization_id"], limit + 1, offset)).fetchall() items = [] for row in rows[:limit]: suggestions = [json.loads(x[0]) for x in db.execute("SELECT output_json FROM ai_suggestions WHERE ai_run_id=? AND organization_id=? ORDER BY id", (row["id"], user["organization_id"]))] items.append(self._ai_run_json(row, suggestions)) return self.send_json(200, {"organization_id": user["organization_id"], "items": items, "limit": limit, "offset": offset, "has_more": len(rows) > limit}) def decide_ai(self, run_id, decision, db, user): row = db.execute("SELECT * FROM ai_runs WHERE id=? AND organization_id=?", (run_id, user["organization_id"])).fetchone() if not row: return self.send_json(404, {"error": "not_found"}) if decision == "approve" and row["status"] != "succeeded": return self.send_json(409, {"error": "run_not_approvable"}) if row["approval_state"] != "pending": return self.send_json(409, {"error": "already_decided"}) if decision == "approve": try: metadata = json.loads(row["data_minimization_json"] or "{}") except (TypeError, ValueError): metadata = {} expected = metadata.get("input_fingerprint") if expected and expected != self._ai_current_fingerprint(db, row): return self.send_json(409, {"error": "ai_run_stale", "approval_state": "pending", "reason": "source_hash_changed"}) now = datetime.now(timezone.utc).replace(microsecond=0).isoformat() if decision == "approve": db.execute("UPDATE ai_runs SET approval_state='approved',approved_at=? WHERE id=? AND organization_id=?", (now, run_id, user["organization_id"])) else: db.execute("UPDATE ai_runs SET approval_state='rejected',rejected_at=? WHERE id=? AND organization_id=?", (now, run_id, user["organization_id"])) self.audit(db, user, "ai." + decision, str(run_id)); db.commit() return self.send_json(200, self._ai_run_json(db.execute("SELECT * FROM ai_runs WHERE id=?", (run_id,)).fetchone())) def do_GET(self): parsed=urlparse(self.path); path=parsed.path.rstrip("/") if path=="/api/v1/health/live": return self.send_json(200,{"status":"ok","organization_id":ORGANIZATION_ID,"outreach_enabled":False}) if path=="/api/v1/health/ready": try: check=connect(self.server.db_path); check.execute("SELECT 1"); check.close() return self.send_json(200,{"status":"ok","ready":True,"database":"ok","outreach_enabled":False}) except (sqlite3.Error, OSError): return self.send_json(503,{"status":"not_ready","ready":False,"database":"error","outreach_enabled":False}) db=self.db() try: user=self.require_auth(db) if not user:return org=user["organization_id"] if path=="/api/v1/auth/me": return self.send_json(200,{"id":user["id"],"email":user["email"],"role":user["role"],"organization_id":org}) if path=="/api/v1/admin/users": if user["role"] not in {"owner","admin"}: return self.send_json(403,{"error":"forbidden"}) return self.send_json(200,{"items":[dict(r) for r in db.execute("SELECT id,email,role,organization_id,created_at FROM users WHERE organization_id=? ORDER BY id",(org,))]}) if path=="/api/v1/dashboard/summary": row=db.execute("SELECT COUNT(*) businesses,COALESCE(AVG(score),0) average_score FROM businesses WHERE organization_id=?",(org,)).fetchone() counts={"new":db.execute("SELECT COUNT(*) FROM businesses WHERE organization_id=? AND created_at>=datetime('now','-7 days')",(org,)).fetchone()[0],"hot":db.execute("SELECT COUNT(*) FROM businesses WHERE organization_id=? AND score>=70 AND merge_status='active'",(org,)).fetchone()[0],"review":db.execute("SELECT COUNT(*) FROM businesses WHERE organization_id=? AND review_status='pending'",(org,)).fetchone()[0],"source_health":db.execute("SELECT COUNT(*) FROM sources WHERE organization_id=? AND health_status IN ('healthy','unhealthy')",(org,)).fetchone()[0],"active_jobs":db.execute("SELECT COUNT(*) FROM jobs WHERE organization_id=? AND status IN ('queued','running')",(org,)).fetchone()[0]} clickable={key:{"count":value,"filter":{"dashboard_filter":key}} for key,value in counts.items()} return self.send_json(200,{"organization_id":org,"businesses":row["businesses"],"average_score":round(row["average_score"],2),"suppressed":db.execute("SELECT COUNT(*) FROM suppressions WHERE organization_id=?",(org,)).fetchone()[0],"counts":counts,"clickable_filters":clickable,"quick_filters":[{"key":key,"count":value,"filter":{"dashboard_filter":key}} for key,value in counts.items()]}) if path=="/api/v1/saved-filters": return self.list_saved_filters(db,user) if path=="/api/v1/review-queue": return self.review_queue(db,org,parse_qs(parsed.query)) if path=="/api/v1/score-rules": return self.list_score_rules(db,org) if path=="/api/v1/scoring/summary": return self.scoring_summary(db,org) if path=="/api/v1/businesses": return self.list_businesses(db,org,parse_qs(parsed.query)) if path=="/api/v1/merge-history": return self.list_merge_history(db,org) if path=="/api/v1/sources": return self.list_sources(db,org) if path=="/api/v1/discovery-queries": return self.list_queries(db,org) if path=="/api/v1/discovery-runs": return self.list_discovery_runs(db,org,parse_qs(parsed.query)) if path=="/api/v1/discovery/provider-status": return self.send_json(200, ai_research_provider_status()) if path=="/api/v1/discovery/ai-provider-status": return self.send_json(200, ai_research_provider_status()) if path=="/api/v1/source-records": return self.list_source_records(db,org,parse_qs(parsed.query)) if path=="/api/v1/jobs": return self.list_jobs(db,org,parse_qs(parsed.query)) if path=="/api/v1/domain-checks": return self.list_domain_checks(db,org,parse_qs(parsed.query)) if path=="/api/v1/website-scans": return self.list_website_scans(db,org,parse_qs(parsed.query)) if path=="/api/v1/contact-extractions": return self.list_contact_extractions(db,org,parse_qs(parsed.query)) if path=="/api/v1/pipeline-entries": return self.list_pipeline(db,org,parse_qs(parsed.query)) if path=="/api/v1/pipeline-stages": return self.list_stages(db,org) if path=="/api/v1/outcomes": return self.send_json(200,{"items":sorted(self.OUTCOMES)}) if path=="/api/v1/interactions": return self.list_interactions(db,org,parse_qs(parsed.query)) if path=="/api/v1/suppressions": return self.list_suppressions(db,org,parse_qs(parsed.query)) if path=="/api/v1/ai-runs": return self.list_ai_runs(db,user,parse_qs(parsed.query)) if path=="/api/v1/ai/provider-config": return self.ai_provider_config(db,user) if path=="/api/v1/admin/ai-provider-config": return self.remote_ai_provider_config(db,user) if path=="/api/v1/outreach/drafts": return self.list_outreach_drafts(db,user,parse_qs(parsed.query)) if path=="/api/v1/outreach/provider-config": return self.provider_config(db,user) if path in ("/api/v1/reports/pipeline","/api/v1/reports/outcomes","/api/v1/reports/activity"): return self.report(db,org,path.rsplit('/',1)[1],parse_qs(parsed.query)) if path.startswith("/api/v1/jobs/"): return self.get_job_route(db,org,path,parse_qs(parsed.query)) if path.startswith("/api/v1/businesses/"): bits=path.split("/"); ident=bits[4] if len(bits)>4 else "" if not ident.isdigit(): return self.send_json(404,{"error":"not_found"}) row=self.business(db,int(ident),org) if not row:return self.send_json(404,{"error":"not_found"}) if len(bits)==7 and bits[5:]==["websites","scan"]: return self.get_latest_website_scan(int(ident),db,user) if len(bits)==7 and bits[5:]==["domains","check"]: return self.get_domain_check(int(ident),db,user,parse_qs(parsed.query)) if len(bits)==7 and bits[5:]==["contacts","extract"]: return self.send_json(405,{"error":"method_not_allowed"}) if len(bits)==6 and bits[5]=="domain-candidates": return self.list_domain_candidates(int(ident),db,org) if len(bits)==6 and bits[5]=="matches": return self.matches(int(ident),db,org) payload=row_json(row); payload.update(self.nested(db,int(ident),org)); return self.send_json(200,payload) return self.send_json(404,{"error":"not_found"}) finally: db.close() def _domain_result(self, row, cache_hit=False): try: result=json.loads(row["result_json"] or "{}") except (TypeError,ValueError): result={} result.update({"id":row["id"],"business_id":row["business_id"],"domain":row["domain"],"status":row["status"],"cache_hit":cache_hit,"checked_at":row["checked_at"],"cache_expires_at":row["cache_expires_at"]}) return result def _check_domain(self, bid, domain, db, user): normalized=normalize_registrable_domain(domain) if normalized == "unknown": return self.send_json(400,{"error":"unsupported_domain","status":"unknown"}) now=datetime.now(timezone.utc).replace(microsecond=0).isoformat() cache_key=normalized+":a_aaaa:v1" cached=db.execute("SELECT * FROM domain_checks WHERE organization_id=? AND business_id=? AND cache_key=? AND cache_expires_at>? ORDER BY id DESC LIMIT 1",(user["organization_id"],bid,cache_key,now)).fetchone() if cached:return self.send_json(200,self._domain_result(cached,True)) result=resolve_domain(normalized) try: cur=db.execute("INSERT INTO domain_checks(organization_id,business_id,domain,status,result_json,cache_key,checked_at,cache_expires_at) VALUES(?,?,?,?,?,?,?,?)",(user["organization_id"],bid,normalized,result["status"],json.dumps(result,sort_keys=True),cache_key,result["checked_at"],result["cache_expires_at"])) except sqlite3.IntegrityError: existing=db.execute("SELECT id FROM domain_checks WHERE organization_id=? AND business_id=? AND cache_key=?",(user["organization_id"],bid,cache_key)).fetchone() if not existing: return self.send_json(409,{"error":"domain_check_conflict"}) db.execute("UPDATE domain_checks SET domain=?,status=?,result_json=?,checked_at=?,cache_expires_at=? WHERE id=?",(normalized,result["status"],json.dumps(result,sort_keys=True),result["checked_at"],result["cache_expires_at"],existing["id"])) self.audit(db,user,"domain.checked",f"{bid}:{normalized}:{result['status']}"); db.commit() return self.send_json(200,self._domain_result(db.execute("SELECT * FROM domain_checks WHERE id=?",(existing["id"],)).fetchone())) self.audit(db,user,"domain.checked",f"{bid}:{normalized}:{result['status']}"); db.commit() return self.send_json(200,self._domain_result(db.execute("SELECT * FROM domain_checks WHERE id=?",(cur.lastrowid,)).fetchone())) def post_domain_check(self,bid,payload,db,user): if not self.business(db,bid,user["organization_id"]): return self.send_json(404,{"error":"not_found"}) domain=payload.get("domain") or self.business(db,bid,user["organization_id"])["website_domain"] if not str(domain).strip(): return self.send_json(400,{"error":"domain_required"}) return self._check_domain(bid,str(domain),db,user) def get_domain_check(self,bid,db,user,query): domain=(query.get("domain") or [""])[0] if not domain:return self.send_json(400,{"error":"domain_required"}) return self._check_domain(bid,domain,db,user) def list_domain_checks(self,db,org,query): try: limit=max(1,min(int((query.get("page_size") or [50])[0]),100)); offset=max(0,int((query.get("offset") or [0])[0])) except (ValueError,TypeError): return self.send_json(400,{"error":"invalid_pagination"}) rows=db.execute("SELECT * FROM domain_checks WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?",(org,limit+1,offset)).fetchall() return self.send_json(200,{"organization_id":org,"items":[self._domain_result(r) for r in rows[:limit]],"limit":limit,"offset":offset,"has_more":len(rows)>limit}) def _website_scan_result(self, row, cache_hit=False): result = json.loads(row["result_json"] or "{}") result.update({"id": row["id"], "business_id": row["business_id"], "website_id": row["website_id"], "input_url": row["input_url"], "classification": row["classification"], "scanned_at": row["scanned_at"], "cache_expires_at": row["cache_expires_at"], "cache_hit": cache_hit}) return result def list_website_scans(self, db, org, query): try: limit = int((query.get("page_size") or [50])[0]); offset = max(0, int((query.get("offset") or [0])[0])) if limit < 1 or limit > WEBSITE_SCAN_PAGE_SIZE: raise ValueError except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"}) params = [org]; where = ["organization_id=?"] if (query.get("business_id") or [""])[0].isdigit(): where.append("business_id=?"); params.append(int(query["business_id"][0])) rows = db.execute("SELECT * FROM website_scans WHERE " + " AND ".join(where) + " ORDER BY id DESC LIMIT ? OFFSET ?", params + [limit + 1, offset]).fetchall() return self.send_json(200, {"organization_id": org, "items": [self._website_scan_result(r) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit}) def _contact_json(self, row): result = row_json(row) for key in ("public_business", "suppressed", "do_not_contact"): if key in result: result[key] = bool(result[key]) return result def list_contact_extractions(self, db, org, query): try: limit = int((query.get("page_size") or [50])[0]); offset = max(0, int((query.get("offset") or [0])[0])) if limit < 1 or limit > CONTACT_EXTRACTION_PAGE_SIZE: raise ValueError except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"}) params = [org]; where = ["organization_id=?"] business_id = (query.get("business_id") or [""])[0] if business_id.isdigit(): where.append("business_id=?"); params.append(int(business_id)) rows = db.execute("SELECT * FROM contact_extractions WHERE " + " AND ".join(where) + " ORDER BY id DESC LIMIT ? OFFSET ?", params + [limit + 1, offset]).fetchall() return self.send_json(200, {"organization_id": org, "items": [self._contact_json(r) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit}) def extract_business_contacts(self, bid, payload, db, user): org = user["organization_id"]; business = self.business(db, bid, org) if not business: return self.send_json(404, {"error": "not_found"}) scan_id = payload.get("website_scan_id", payload.get("scan_id")) scan = None if scan_id is not None: if not isinstance(scan_id, int): return self.send_json(400, {"error": "invalid_scan"}) scan = db.execute("SELECT * FROM website_scans WHERE id=? AND business_id=? AND organization_id=?", (scan_id, bid, org)).fetchone() if not scan: return self.send_json(404, {"error": "scan_not_found"}) try: stored = json.loads(scan["result_json"] or "{}") except (TypeError, ValueError): stored = {} source_url = str(payload.get("source_url") or scan["input_url"]).strip() source_html = payload.get("html") if isinstance(payload.get("html"), str) else stored.get("html") if source_html is None: return self.send_json(409, {"error": "scan_html_unavailable"}) if source_url != scan["input_url"] and source_url != (stored.get("final_url") or ""): return self.send_json(400, {"error": "source_not_approved"}) else: source_url = str(payload.get("source_url") or "").strip() source_html = payload.get("html") if not source_url or not isinstance(source_html, str): return self.send_json(400, {"error": "approved_scan_required"}) parsed = urlparse(source_url) official_hosts = {str(business["website_domain"]).lower().strip(".")} if business["website"]: official_hosts.add((urlparse(business["website"]).hostname or "").lower().strip(".")) if parsed.scheme not in {"http", "https"} or not parsed.hostname or parsed.hostname.lower().strip(".") not in official_hosts: return self.send_json(400, {"error": "source_not_approved"}) if len(source_html.encode("utf-8")) > MAX_HTML_BYTES: return self.send_json(413, {"error": "html_too_large"}) try: requested_limit = int(payload.get("limit", MAX_RESULTS)) except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_limits"}) if requested_limit < 1 or requested_limit > MAX_RESULTS: return self.send_json(400, {"error": "invalid_limits"}) key = str(payload.get("idempotency_key") or hashlib.sha256((str(scan_id or "") + source_url + source_html).encode()).hexdigest())[:200] existing = db.execute("SELECT * FROM contact_extractions WHERE organization_id=? AND business_id=? AND extraction_key=? ORDER BY id", (org, bid, key)).fetchall() if existing: return self.send_json(200, {"business_id": bid, "extraction_key": key, "items": [self._contact_json(r) for r in existing], "idempotent": True}) suppressions = [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))] try: found = extract_contacts(source_html, source_url, suppressions=suppressions, max_results=requested_limit) except ValueError as exc: return self.send_json(400, {"error": str(exc)}) for item in found: db.execute("INSERT INTO contact_extractions(organization_id,business_id,website_scan_id,extraction_key,kind,value,label,classification,confidence,source_url,public_business,mx_status,suppressed,do_not_contact,provenance) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (org,bid,scan["id"] if scan else None,key,item["kind"],item["value"],item["label"],item["classification"],item["confidence"],item["source_url"],int(item["public_business"]),item["mx_status"],int(item["suppressed"]),int(item["do_not_contact"]),item["provenance"])) self.audit(db, user, "contacts.extracted", f"{bid}:{len(found)}:{key}"); db.commit() rows = db.execute("SELECT * FROM contact_extractions WHERE organization_id=? AND business_id=? AND extraction_key=? ORDER BY id", (org,bid,key)).fetchall() return self.send_json(201, {"business_id": bid, "extraction_key": key, "items": [self._contact_json(r) for r in rows], "idempotent": False}) def get_latest_website_scan(self, bid, db, user): row = db.execute("SELECT * FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT 1", (bid, user["organization_id"])).fetchone() if not row: return self.send_json(404, {"error": "scan_not_found"}) return self.send_json(200, self._website_scan_result(row, False)) def scan_business_website(self, bid, payload, db, user): org = user["organization_id"]; business = self.business(db, bid, org) if not business: return self.send_json(404, {"error": "not_found"}) requested = str(payload.get("url", "")).strip() if isinstance(payload, dict) else "" if not requested: child = db.execute("SELECT * FROM websites WHERE business_id=? AND organization_id=? ORDER BY id LIMIT 1", (bid, org)).fetchone() requested = (child["url"] if child else business["website"]) or "" try: safe_url = validate_url(requested) except ValueError as exc: self.audit(db, user, "website.scan.rejected", f"{bid}:{str(exc)}"); db.commit() return self.send_json(400, {"error": "unsafe_url", "reason": str(exc)}) website = db.execute("SELECT id FROM websites WHERE business_id=? AND organization_id=? AND url=? ORDER BY id LIMIT 1", (bid, org, safe_url)).fetchone() cache_key = hashlib.sha256(safe_url.encode()).hexdigest(); now = datetime.now(timezone.utc).replace(microsecond=0); expires = now + timedelta(seconds=WEBSITE_SCAN_CACHE_SECONDS) cached = db.execute("SELECT * FROM website_scans WHERE organization_id=? AND business_id=? AND cache_key=? AND cache_expires_at>? ORDER BY id DESC LIMIT 1", (org, bid, cache_key, now.isoformat())).fetchone() if cached: self.audit(db, user, "website.scan.cache_hit", f"{bid}:{safe_url}"); db.commit() return self.send_json(200, self._website_scan_result(cached, True)) result = scan_website(safe_url); result["business_id"] = bid cur = db.execute("INSERT INTO website_scans(organization_id,business_id,website_id,input_url,classification,result_json,cache_key,scanned_at,cache_expires_at) VALUES(?,?,?,?,?,?,?,?,?)", (org, bid, website["id"] if website else None, safe_url, result["classification"], json.dumps(result, sort_keys=True), cache_key, now.isoformat(), expires.isoformat())) self.audit(db, user, "website.scanned", f"{bid}:{result['classification']}"); db.commit() return self.send_json(201, self._website_scan_result(db.execute("SELECT * FROM website_scans WHERE id=?", (cur.lastrowid,)).fetchone())) def list_domain_candidates(self,bid,db,org): business=self.business(db,bid,org) if not business:return self.send_json(404,{"error":"not_found"}) rows=db.execute("SELECT * FROM domain_candidates WHERE organization_id=? AND business_id=? ORDER BY rank,id",(org,bid)).fetchall() if not rows: generated=generate_candidate_domains(business["name"],business["description"],business["city"]) for rank,domain in enumerate(generated,1): db.execute("INSERT OR IGNORE INTO domain_candidates(organization_id,business_id,domain,rank) VALUES(?,?,?,?)",(org,bid,domain,rank)) db.commit(); rows=db.execute("SELECT * FROM domain_candidates WHERE organization_id=? AND business_id=? ORDER BY rank,id",(org,bid)).fetchall() return self.send_json(200,{"business_id":bid,"items":[row_json(r) for r in rows]}) def check_availability(self,bid,payload,db,user): if not self.business(db,bid,user["organization_id"]):return self.send_json(404,{"error":"not_found"}) candidates=db.execute("SELECT domain FROM domain_candidates WHERE organization_id=? AND business_id=? ORDER BY rank,id",(user["organization_id"],bid)).fetchall() domains=[r["domain"] for r in candidates] if isinstance(payload.get("domains"),list): domains=[normalize_registrable_domain(x) for x in payload["domains"] if normalize_registrable_domain(x)!="unknown"][:20] self.audit(db,user,"domain.availability.checked",str(bid)); db.commit() return self.send_json(200,{"business_id":bid,"status":"unknown","reason":"not_configured","provider_configured":False,"items":[{"domain":d,"status":"unknown","reason":"not_configured"} for d in domains]}) def list_jobs(self, db, org, query): try: limit=max(1,min(int(query.get("page_size",[50])[0]),JOB_PAGE_SIZE)); offset=max(0,int(query.get("offset",[0])[0])) except (ValueError, TypeError): return self.send_json(400,{"error":"invalid_pagination"}) rows=db.execute("SELECT * FROM jobs WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?",(org,limit+1,offset)).fetchall(); more=len(rows)>limit return self.send_json(200,{"organization_id":org,"items":[job_json(r) for r in rows[:limit]],"limit":limit,"offset":offset,"has_more":more}) def _discovery_run_json(self, row): item = row_json(row) for field, default in (("criteria_json", {}), ("seed_urls_json", []), ("result_json", {})): try: item[field[:-5]] = json.loads(item.pop(field) or json.dumps(default)) except (TypeError, ValueError): item[field[:-5]] = default return item def list_discovery_runs(self, db, org, query): try: limit = max(1, min(int((query.get("page_size") or [50])[0]), 100)); offset = max(0, int((query.get("offset") or [0])[0])) except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"}) rows = db.execute("SELECT * FROM discovery_runs WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?", (org, limit + 1, offset)).fetchall() return self.send_json(200, {"organization_id": org, "items": [self._discovery_run_json(r) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit}) def create_scoped_discovery(self, payload, db, user): criteria = payload.get("criteria", {}); seeds = payload.get("seed_urls") criteria_only = seeds is None if not isinstance(criteria, dict): return self.send_json(400, {"error": "invalid_criteria"}) if not criteria_only and (not isinstance(seeds, list) or not seeds): return self.send_json(400, {"error": "seed_urls_required"}) if criteria_only: try: validate_ai_research_criteria(criteria) except AIResearchConfigError as exc: if str(exc) in {"prompt_injection_rejected", "criteria_too_large"}: return self.send_json(400, {"error": str(exc)}) status = ai_research_provider_status() # Legacy SEARCH_PROVIDER_* may pass only during migration; the # primary status and endpoint remain AI research. legacy = search_provider_status() if status["status"] != "ready" and legacy["status"] != "ready": return self.send_json(503, {"error": status["status"], "provider": status["provider"]}) try: if not criteria_only and len(seeds) > 5: raise ValueError("invalid_criteria") if len(json.dumps(criteria).encode()) > 8192: raise ValueError("invalid_criteria") if not isinstance(criteria.get("keywords", criteria.get("keyword", [])), (list, str)): raise ValueError("invalid_criteria") max_pages = int(payload.get("max_pages", 20)); max_candidates = int(payload.get("max_candidates", 50)) if not 1 <= max_pages <= 20 or not 1 <= max_candidates <= 50: raise ValueError("invalid_limits") except (ValueError, TypeError) as exc: return self.send_json(400, {"error": str(exc) or "invalid_criteria"}) if not criteria_only: try: for url in seeds: validate_url(url) except (ValueError, TypeError): return self.send_json(400, {"error": "unsafe_seed_url"}) key = str(payload.get("idempotency_key", "")).strip() if not key or len(key) > 200: return self.send_json(400, {"error": "invalid_idempotency_key"}) job_payload = {"criteria": criteria, "seed_urls": seeds if not criteria_only else None, "max_pages": max_pages, "max_candidates": max_candidates} result = self.create_job({"type": "scoped_discovery", "payload": job_payload, "idempotency_key": key, "_accepted": True, "_defer_wakeup": True}, db, user) # create_job has already committed; read its id from the response is not available, # so resolve by the tenant-scoped idempotency key. job = db.execute("SELECT * FROM jobs WHERE organization_id=? AND idempotency_key=?", (user["organization_id"], key)).fetchone() if not db.execute("SELECT id FROM discovery_runs WHERE organization_id=? AND job_id=?", (user["organization_id"], job["id"])).fetchone(): db.execute("INSERT INTO discovery_runs(organization_id,job_id,criteria_json,seed_urls_json) VALUES(?,?,?,?)", (user["organization_id"], job["id"], json.dumps(criteria, sort_keys=True), json.dumps(seeds if not criteria_only else []))) self.audit(db, user, "discovery.created", str(job["id"])); db.commit() getattr(self.server, "job_wakeup", threading.Event()).set() return result def get_job_route(self, db, org, path, query): bits=path.split("/") if len(bits)<5 or not bits[4].isdigit(): return self.send_json(404,{"error":"not_found"}) job=db.execute("SELECT * FROM jobs WHERE id=? AND organization_id=?",(int(bits[4]),org)).fetchone() if not job:return self.send_json(404,{"error":"not_found"}) if len(bits)==5:return self.send_json(200,job_json(job)) if len(bits)==6 and bits[5]=="events": try: after=max(0,int(query.get("after",[0])[0])) except (ValueError,TypeError): return self.send_json(400,{"error":"invalid_sequence"}) events=[row_json(r) for r in db.execute("SELECT * FROM job_events WHERE job_id=? AND organization_id=? AND sequence>? ORDER BY sequence",(job["id"],org,after))] return self.send_json(200,{"items":events,"after":after}) if len(bits)==7 and bits[5]=="events" and bits[6]=="stream": try: after=max(0,int(query.get("after",[0])[0])) except (ValueError,TypeError): return self.send_json(400,{"error":"invalid_sequence"}) events=[row_json(r) for r in db.execute("SELECT * FROM job_events WHERE job_id=? AND organization_id=? AND sequence>? ORDER BY sequence",(job["id"],org,after))] body=b"".join((b"event: "+str(e["event_type"]).encode()+b"\\ndata: "+json.dumps(e,sort_keys=True).encode()+b"\\n\\n") for e in events) self.send_response(200);self.send_header("Content-Type","text/event-stream");self.send_header("Cache-Control","no-cache");self.send_header("Content-Length",str(len(body)));self.end_headers();self.wfile.write(body);return return self.send_json(404,{"error":"not_found"}) def create_job(self, payload, db, user): kind=str(payload.get("type","")).strip(); key=str(payload.get("idempotency_key","")).strip(); data=payload.get("payload",{}) if kind not in JOB_TYPES:return self.send_json(400,{"error":"invalid_job_type"}) if not key or len(key)>200 or not isinstance(data,dict):return self.send_json(400,{"error":"invalid_job_request"}) safe=json.dumps(redact(data),sort_keys=True,separators=(",",":")); max_attempts=max(1,min(int(payload.get("max_attempts",3)),5)) if str(payload.get("max_attempts",3)).isdigit() else 3 try: cur=db.execute("INSERT INTO jobs(organization_id,idempotency_key,type,payload,max_attempts) VALUES(?,?,?,?,?)",(user["organization_id"],key,kind,safe,max_attempts)); jid=cur.lastrowid self.add_job_event(db,jid,user["organization_id"],"queued","Job queued",0); self.audit(db,user,"job.created",str(jid)); db.commit() if not payload.get("_defer_wakeup"): getattr(self.server,"job_wakeup",threading.Event()).set() return self.send_json(202 if payload.get("_accepted") else 201,job_json(db.execute("SELECT * FROM jobs WHERE id=?",(jid,)).fetchone())) except sqlite3.IntegrityError: row=db.execute("SELECT * FROM jobs WHERE organization_id=? AND idempotency_key=?",(user["organization_id"],key)).fetchone(); return self.send_json(200,job_json(row)) def add_job_event(self, db, jid, org, event_type, message, progress=0, error_code=None): seq=db.execute("SELECT COALESCE(MAX(sequence),0)+1 FROM job_events WHERE job_id=?",(jid,)).fetchone()[0] db.execute("INSERT INTO job_events(job_id,organization_id,sequence,event_type,message,progress,error_code) VALUES(?,?,?,?,?,?,?)",(jid,org,seq,event_type,str(message)[:500],progress,error_code)); return seq def job_action(self, db, user, path): bits=path.split("/"); jid=int(bits[4]) if len(bits)>4 and bits[4].isdigit() else -1; action=bits[5] if len(bits)>5 else "" job=db.execute("SELECT * FROM jobs WHERE id=? AND organization_id=?",(jid,user["organization_id"])).fetchone() if not job:return self.send_json(404,{"error":"not_found"}) if action=="cancel": if job["status"] in ("queued","running"): db.execute("UPDATE jobs SET status='cancelled',completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));self.add_job_event(db,jid,user["organization_id"],"cancelled","Job cancelled",job["progress"]) self.audit(db,user,"job.cancelled",str(jid));db.commit();return self.send_json(200,job_json(db.execute("SELECT * FROM jobs WHERE id=?",(jid,)).fetchone())) if action=="retry": if job["status"]!="failed":return self.send_json(409,{"error":"job_not_failed"}) db.execute("UPDATE jobs SET status='queued',error_code=NULL,completed_at=NULL,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));self.add_job_event(db,jid,user["organization_id"],"retry","Job retry queued",job["progress"]);self.audit(db,user,"job.retried",str(jid));db.commit();getattr(self.server,"job_wakeup",threading.Event()).set();return self.send_json(200,job_json(db.execute("SELECT * FROM jobs WHERE id=?",(jid,)).fetchone())) return self.send_json(404,{"error":"not_found"}) def ensure_score_rules(self, db, org): for rule in DEFAULT_RULES: db.execute("INSERT OR IGNORE INTO score_rules(organization_id,code,name,description,condition_json,points,max_applications,enabled,version) VALUES(?,?,?,?,?,?,?,?,?)", (org, rule["code"], rule["name"], rule["description"], json.dumps(rule["condition_json"], sort_keys=True), rule["points"], rule["max_applications"], rule["enabled"], rule["version"])) def list_score_rules(self, db, org): self.ensure_score_rules(db, org); db.commit() rows = db.execute("SELECT * FROM score_rules WHERE organization_id=? ORDER BY code,id", (org,)).fetchall() items = [] for row in rows: item = row_json(row) try: item["condition_json"] = json.loads(item["condition_json"]) except (TypeError, ValueError): item["condition_json"] = {} item["enabled"] = bool(item["enabled"]); items.append(item) return self.send_json(200, {"organization_id": org, "items": items}) def create_score_rule(self, payload, db, user): code = str(payload.get("code", "")).strip(); name = str(payload.get("name", "")).strip(); condition = payload.get("condition_json", payload.get("condition", {})) try: points = int(payload.get("points", 0)); maximum = int(payload.get("max_applications", 1)); version = int(payload.get("version", 1)) except (TypeError, ValueError): return self.send_json(400, {"error": "invalid_rule"}) if not code or not name or not isinstance(condition, dict) or maximum < 1 or version < 1 or points < -100 or points > 100: return self.send_json(400, {"error": "invalid_rule"}) try: cur = db.execute("INSERT INTO score_rules(organization_id,code,name,description,condition_json,points,max_applications,enabled,version) VALUES(?,?,?,?,?,?,?,?,?)", (user["organization_id"], code, name, str(payload.get("description", "")), json.dumps(condition, sort_keys=True), points, maximum, int(bool(payload.get("enabled", True))), version)) except sqlite3.IntegrityError: return self.send_json(409, {"error": "duplicate_rule"}) self.audit(db, user, "score_rule.created", code); db.commit() row = db.execute("SELECT * FROM score_rules WHERE id=?", (cur.lastrowid,)).fetchone(); item = row_json(row); item["condition_json"] = condition; item["enabled"] = bool(item["enabled"]) return self.send_json(201, item) def update_score_rule(self, rid, payload, db, user): row = db.execute("SELECT * FROM score_rules WHERE id=? AND organization_id=?", (rid, user["organization_id"])).fetchone() if not row: return self.send_json(404, {"error": "not_found"}) allowed = {"name", "description", "condition_json", "points", "max_applications", "enabled", "version"}; values = {k: payload[k] for k in allowed if k in payload} if not values: return self.send_json(400, {"error": "no_changes"}) if "condition_json" in values and not isinstance(values["condition_json"], dict): return self.send_json(400, {"error": "invalid_rule"}) if "points" in values: try: values["points"] = int(values["points"]) except (TypeError, ValueError): return self.send_json(400, {"error": "invalid_rule"}) columns=[]; params=[] for key, value in values.items(): columns.append(key + "=?"); params.append(json.dumps(value, sort_keys=True) if key == "condition_json" else (int(bool(value)) if key == "enabled" else value)) params += [rid, user["organization_id"]]; db.execute("UPDATE score_rules SET " + ",".join(columns) + ",updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?", params); self.audit(db, user, "score_rule.updated", str(rid)); db.commit() item = row_json(db.execute("SELECT * FROM score_rules WHERE id=?", (rid,)).fetchone()) try: item["condition_json"] = json.loads(item["condition_json"]) except (TypeError, ValueError): item["condition_json"] = {} item["enabled"] = bool(item["enabled"]); return self.send_json(200, item) def recalculate_score(self, bid, payload, db, user): org = user["organization_id"]; business = self.business(db, bid, org) if not business: return self.send_json(404, {"error": "not_found"}) scans = db.execute("SELECT result_json,classification,scanned_at FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT 1", (bid, org)).fetchone(); website = {"classification": business["website_class"]} if scans: try: website.update(json.loads(scans["result_json"] or "{}")) except (TypeError, ValueError): pass website["classification"] = scans["classification"] contacts = [dict(r) for r in db.execute("SELECT public_business,suppressed,do_not_contact FROM contact_extractions WHERE business_id=? AND organization_id=?", (bid, org))] drow = db.execute("SELECT status,result_json,checked_at FROM domain_checks WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT 1", (bid, org)).fetchone(); domain = dict(drow) if drow else {} if drow: try: domain.update(json.loads(drow["result_json"] or "{}")) except (TypeError, ValueError): pass suppressed = is_suppressed(dict(business), [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))]) signals = signals_for_business(dict(business), website, contacts, domain, suppressed); self.ensure_score_rules(db, org); rules = [dict(r) for r in db.execute("SELECT * FROM score_rules WHERE organization_id=?", (org,))]; result = evaluate_score(signals, rules) override_score = payload.get("override_score"); override_eligible = payload.get("override_eligible") if override_score is not None or override_eligible is not None: reason = str(payload.get("override_reason", "")).strip() if not reason: return self.send_json(400, {"error": "override_reason_required"}) if override_score is not None: result["score"] = max(0, min(100, int(override_score))) if override_eligible is not None and not suppressed: result["eligible"] = bool(override_eligible) result["priority_band"] = "ineligible" if not result["eligible"] else ("high" if result["score"] >= 70 else "medium" if result["score"] >= 40 else "low") db.execute("UPDATE businesses SET score=?,score_version=?,score_factors=?,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?", (result["score"], SCORE_VERSION, json.dumps(result["explanations"], sort_keys=True), bid, org)) cur=db.execute("INSERT INTO score_history(organization_id,business_id,score,eligible,priority_band,score_version,explanations_json,signals_json,override_score,override_eligible,override_reason,actor_user_id) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)", (org,bid,result["score"],int(result["eligible"]),result["priority_band"],SCORE_VERSION,json.dumps(result["explanations"],sort_keys=True),json.dumps(signals,sort_keys=True),override_score,override_eligible,payload.get("override_reason"),user["id"])) self.audit(db,user,"business.score_recalculated",f"{bid}:{result['score']}"); db.commit(); result.update({"business_id": bid, "history_id": cur.lastrowid}); return self.send_json(200, result) def scoring_summary(self, db, org): row=db.execute("SELECT COUNT(*) businesses,COALESCE(AVG(score),0) average_score,SUM(CASE WHEN score>=70 THEN 1 ELSE 0 END) high_priority FROM businesses WHERE organization_id=? AND merge_status='active'",(org,)).fetchone() bands={r["priority_band"]:r["count"] for r in db.execute("SELECT priority_band,COUNT(*) count FROM score_history WHERE organization_id=? GROUP BY priority_band",(org,))} return self.send_json(200,{"organization_id":org,"businesses":row["businesses"],"average_score":round(row["average_score"],2),"high_priority":row["high_priority"] or 0,"bands":bands,"history_count":db.execute("SELECT COUNT(*) FROM score_history WHERE organization_id=?",(org,)).fetchone()[0]}) def list_businesses(self,db,org,query): def number(name, default=None): raw=query.get(name,[None])[0] if raw is None:return default try:return int(raw) except ValueError: raise ValueError try: page=number("page",1); size=number("page_size",50); score=number("score_min",None); cursor=number("cursor",0) except ValueError:return self.send_json(400,{"error":"invalid_pagination"}) if page<1 or size<1 or size>100 or cursor<0:return self.send_json(400,{"error":"invalid_pagination"}) params=[org]; where=["b.organization_id=?"]; q=query.get("q",[""])[0].strip(); website_class=query.get("website_class",[""])[0].strip(); stage=query.get("pipeline_stage",[""])[0].strip() if score is not None: where.append("b.score>=?"); params.append(score) if website_class: where.append("b.website_class=?"); params.append(website_class) if q: where.append("(b.name LIKE ? OR b.website_domain LIKE ? OR b.email LIKE ?)"); params += [f"%{q}%"]*3 if stage: where.append("EXISTS (SELECT 1 FROM pipeline_entries p WHERE p.business_id=b.id AND p.organization_id=b.organization_id AND p.stage=?)"); params.append(stage) offset=(number("cursor",0) or 0)+(page-1)*size rows=db.execute("SELECT b.* FROM businesses b WHERE "+" AND ".join(where)+" ORDER BY b.score DESC,b.id LIMIT ? OFFSET ?",params+[size+1,offset]).fetchall(); more=len(rows)>size; rows=rows[:size] return self.send_json(200,{"organization_id":org,"items":[row_json(r) for r in rows],"page":page,"page_size":size,"next_cursor":str(offset+size) if more else None}) def bulk_review(self, payload, db, user): ids = payload.get("ids", payload.get("business_ids")); action = str(payload.get("action", "")).strip().lower() if not isinstance(ids, list) or not ids or len(ids) > 100 or any(not isinstance(i, int) or i < 1 for i in ids) or len(set(ids)) != len(ids): return self.send_json(400, {"error": "invalid_bulk_ids"}) if action not in {"verify", "reject", "assign"}: return self.send_json(400, {"error": "invalid_bulk_action"}) if action == "assign": assignee = str(payload.get("assigned_to", payload.get("assignee", ""))).strip() if not assignee or len(assignee) > 120: return self.send_json(400, {"error": "assignee_required"}) else: assignee = "" org = user["organization_id"]; marks = ",".join("?" for _ in ids) rows = db.execute("SELECT id FROM businesses WHERE organization_id=? AND id IN (" + marks + ")", [org] + ids).fetchall() if len(rows) != len(ids): return self.send_json(404, {"error": "not_found"}) try: if action == "verify": db.execute("UPDATE businesses SET verified=1,verified_at=CURRENT_TIMESTAMP,review_status='verified',updated_at=CURRENT_TIMESTAMP WHERE organization_id=? AND id IN (" + marks + ")", [org] + ids) elif action == "reject": db.execute("UPDATE businesses SET verified=0,review_status='rejected',updated_at=CURRENT_TIMESTAMP WHERE organization_id=? AND id IN (" + marks + ")", [org] + ids) else: db.execute("UPDATE businesses SET assigned_to=?,review_status='assigned',updated_at=CURRENT_TIMESTAMP WHERE organization_id=? AND id IN (" + marks + ")", [assignee, org] + ids) self.audit(db, user, "businesses.bulk_review", json.dumps({"action": action, "ids": ids, "assigned_to": assignee}, sort_keys=True)) db.commit() except Exception: db.rollback(); raise return self.send_json(200, {"action": action, "updated": ids, "count": len(ids)}) def do_POST(self): path=urlparse(self.path).path.rstrip("/") if path=="/api/v1/auth/login":return self.login(self.read_json()) db=self.db() try: user=self.require_auth(db) if not user:return if path=="/api/v1/auth/logout": c=SimpleCookie();c.load(self.headers.get("Cookie",""));t=c.get("session"); if t:db.execute("DELETE FROM sessions WHERE token_hash=?",(hashlib.sha256(t.value.encode()).hexdigest(),)) self.audit(db,user,"logout");db.commit();return self.send_json(200,{"ok":True},{"Set-Cookie":self.auth_cookie("",0)}) if path=="/api/v1/jobs": if user["role"] not in MUTATING_ROLES:return self.send_json(403,{"error":"forbidden"}) return self.create_job(self.read_json(),db,user) if user["role"] not in MUTATING_ROLES:return self.send_json(403,{"error":"forbidden"}) payload=self.read_json(); org=user["organization_id"] if path=="/api/v1/saved-filters": return self.save_filter(payload,db,user) if path=="/api/v1/outreach/provider-config": return self.provider_config(db,user,payload) if path=="/api/v1/ai/provider-config": return self.ai_provider_config(db,user,payload) if path=="/api/v1/admin/ai-provider-config": return self.remote_ai_provider_config(db,user,payload) 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]=="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) if len(path.split("/"))==7 and path.split("/")[3:]==["businesses",path.split("/")[4],"score","recalculate"]: return self.recalculate_score(int(path.split("/")[4]) if path.split("/")[4].isdigit() else -1,payload,db,user) if path.startswith("/api/v1/jobs/"): return self.job_action(db,user,path) if path=="/api/v1/businesses":return self.create_business(payload,db,user) if len(path.split("/"))==7 and path.split("/")[3:6]==["businesses",path.split("/")[4],"outreach"] and path.split("/")[6]=="drafts": return self.create_outreach_draft(int(path.split("/")[4]) if path.split("/")[4].isdigit() else -1,payload,db,user) bits_outreach=path.split("/") if len(bits_outreach)==7 and bits_outreach[:4]==["","api","v1","outreach"] and bits_outreach[4]=="drafts" and bits_outreach[5].isdigit() and bits_outreach[6] in {"approve","send"}: return self.approve_outreach_draft(int(bits_outreach[5]),db,user) if bits_outreach[6]=="approve" else self.send_outreach_draft(int(bits_outreach[5]),db,user) if path=="/api/v1/sources":return self.create_source(payload,db,user) if path=="/api/v1/discovery":return self.create_scoped_discovery(payload,db,user) if path=="/api/v1/discovery-queries":return self.create_query(payload,db,user) if path=="/api/v1/suppressions":return self.create_suppression(payload,db,user) if path=="/api/v1/suppressions/import":return self.import_suppressions(payload,db,user) if path=="/api/v1/pipeline-stages":return self.create_stage(payload,db,user) if path=="/api/v1/interactions": return self.send_json(400,{"error":"business_id_required"}) if path=="/api/v1/contact-extractions": return self.send_json(405,{"error":"method_not_allowed"}) if path=="/api/v1/imports/preview":return self.preview_import(payload,db,org) if len(path.split("/"))==7 and path.split("/")[3:6]==["businesses",path.split("/")[4],"websites"] and path.split("/")[6]=="scan": return self.scan_business_website(int(path.split("/")[4]) if path.split("/")[4].isdigit() else -1,payload,db,user) if len(path.split("/"))==7 and path.split("/")[3:6]==["businesses",path.split("/")[4],"contacts"] and path.split("/")[6]=="extract": return self.extract_business_contacts(int(path.split("/")[4]) if path.split("/")[4].isdigit() else -1,payload,db,user) if len(path.split("/"))==7 and path.split("/")[3:6]==["businesses",path.split("/")[4],"domains"] and path.split("/")[6]=="check": return self.post_domain_check(int(path.split("/")[4]) if path.split("/")[4].isdigit() else -1,payload,db,user) if len(path.split("/"))==7 and path.split("/")[3:6]==["businesses",path.split("/")[4],"domain-candidates"] and path.split("/")[6]=="check-availability": return self.check_availability(int(path.split("/")[4]) if path.split("/")[4].isdigit() else -1,payload,db,user) if path.startswith("/api/v1/merge-history/") and path.endswith("/reverse"): ident=path.split("/")[4] return self.reverse_merge(int(ident) if ident.isdigit() else -1,db,user) bits=path.split("/") if len(bits)==6 and bits[3] == "sources" and bits[4].isdigit() and bits[5] in {"test","ingest"}: return self.test_source(int(bits[4]),db,user) if bits[5]=="test" else self.ingest_source(int(bits[4]),payload,db,user) if len(bits)==6 and bits[3] == "discovery-queries" and bits[4].isdigit() and bits[5]=="run": return self.run_query(int(bits[4]),db,user) if len(bits)==7 and bits[:4]==["","api","v1","businesses"] and bits[5] in CHILD_TABLES and bits[6]=="": pass if len(bits)==6 and bits[:4]==["","api","v1","businesses"] and bits[5] == "pipeline": return self.update_pipeline(int(bits[4]) if bits[4].isdigit() else -1,payload,db,user) if len(bits)==6 and bits[:4]==["","api","v1","businesses"] and bits[5] == "interactions": return self.create_interaction(int(bits[4]) if bits[4].isdigit() else -1,payload,db,user) if len(bits)==6 and bits[:4]==["","api","v1","businesses"] and bits[5] in CHILD_TABLES:return self.create_child(int(bits[4]) if bits[4].isdigit() else -1,bits[5],payload,db,user) if len(bits)==6 and bits[:4]==["","api","v1","businesses"] and bits[5]=="verify":return self.verify_business(int(bits[4]) if bits[4].isdigit() else -1,payload,db,user) if len(bits)==6 and bits[:4]==["","api","v1","businesses"] and bits[5]=="merge":return self.merge_business(int(bits[4]) if bits[4].isdigit() else -1,payload,db,user) return self.send_json(404,{"error":"not_found"}) finally:db.close() def do_PATCH(self): path=urlparse(self.path).path.rstrip("/"); db=self.db() try: user=self.require_auth(db) if not user:return if user["role"] not in MUTATING_ROLES:return self.send_json(403,{"error":"forbidden"}) bits=path.split("/") if len(bits)==5 and bits[:4]==["","api","v1","pipeline-stages"] and bits[4].isdigit(): return self.update_stage(int(bits[4]),self.read_json(),db,user) if len(bits)==5 and bits[:4]==["","api","v1","suppressions"] and bits[4].isdigit(): return self.update_suppression(int(bits[4]),self.read_json(),db,user) if len(bits)==5 and bits[:4]==["","api","v1","pipeline-entries"] and bits[4].isdigit(): return self.update_pipeline_entry(int(bits[4]),self.read_json(),db,user) if len(bits)==5 and bits[:4]==["","api","v1","interactions"] and bits[4].isdigit(): return self.update_interaction(int(bits[4]),self.read_json(),db,user) if len(bits)==5 and bits[:4]==["","api","v1","saved-filters"] and bits[4].isdigit(): return self.update_saved_filter(int(bits[4]),self.read_json(),db,user) if len(bits)==5 and bits[:4]==["","api","v1","score-rules"] and bits[4].isdigit(): return self.update_score_rule(int(bits[4]),self.read_json(),db,user) if len(bits)==5 and bits[:4]==["","api","v1","outreach"] and bits[4]=="provider-config": return self.provider_config(db,user,self.read_json()) if len(bits)==5 and bits[:4]==["","api","v1","sources"] and bits[4].isdigit(): return self.update_source(int(bits[4]),self.read_json(),db,user) if len(bits)==6 and bits[:4]==["","api","v1","outreach"] and bits[4]=="drafts" and bits[5].isdigit(): return self.update_outreach_draft(int(bits[5]),self.read_json(),db,user) if len(bits)==6 and bits[:4]==["","api","v1","businesses"] and bits[5]=="pipeline":return self.update_pipeline(int(bits[4]) if bits[4].isdigit() else -1,self.read_json(),db,user) return self.send_json(404,{"error":"not_found"}) finally:db.close() def do_DELETE(self): path=urlparse(self.path).path.rstrip("/"); db=self.db() try: user=self.require_auth(db) if not user:return if user["role"] not in MUTATING_ROLES:return self.send_json(403,{"error":"forbidden"}) bits=path.split("/") if len(bits)==5 and bits[:4]==["","api","v1","pipeline-stages"] and bits[4].isdigit(): return self.delete_crm_item("pipeline_stages",int(bits[4]),db,user) if len(bits)==5 and bits[:4]==["","api","v1","pipeline-entries"] and bits[4].isdigit(): return self.delete_crm_item("pipeline_entries",int(bits[4]),db,user) if len(bits)==5 and bits[:4]==["","api","v1","interactions"] and bits[4].isdigit(): return self.delete_crm_item("interactions",int(bits[4]),db,user) if len(bits)==5 and bits[:4]==["","api","v1","suppressions"] and bits[4].isdigit(): row=db.execute("SELECT id FROM suppressions WHERE id=? AND organization_id=?",(int(bits[4]),user["organization_id"])).fetchone() if not row:return self.send_json(404,{"error":"not_found"}) db.execute("DELETE FROM suppressions WHERE id=? AND organization_id=?",(int(bits[4]),user["organization_id"]));self.audit(db,user,"suppression.deleted",bits[4]);db.commit();return self.send_json(200,{"ok":True,"id":int(bits[4])}) if len(bits)==5 and bits[:4]==["","api","v1","saved-filters"] and bits[4].isdigit(): return self.delete_saved_filter(int(bits[4]),db,user) return self.send_json(404,{"error":"not_found"}) finally: db.close() def login(self,payload): db=self.db(); email=str(payload.get("email"," ")).strip().lower(); password=str(payload.get("password","")); user=db.execute("SELECT * FROM users WHERE email=?",(email,)).fetchone() try: if not user or not verify_password(password,user["password_hash"],user["password_salt"]):return self.send_json(401,{"error":"invalid_credentials"}) token=secrets.token_urlsafe(32); expires=datetime.now(timezone.utc)+timedelta(days=SESSION_DAYS);db.execute("INSERT INTO sessions(user_id,token_hash,expires_at) VALUES(?,?,?)",(user["id"],hashlib.sha256(token.encode()).hexdigest(),expires.replace(microsecond=0).isoformat()));self.audit(db,user,"login");db.commit();return self.send_json(200,{"id":user["id"],"email":user["email"],"role":user["role"],"organization_id":user["organization_id"]},{"Set-Cookie":self.auth_cookie(token,int(timedelta(days=SESSION_DAYS).total_seconds()))}) finally:db.close() def create_business(self,payload,db,user): org=user["organization_id"] if not str(payload.get("name","")).strip():return self.send_json(400,{"error":"name_required"}) b=normalize_business(payload); suppressions=[dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1",(org,))] if is_suppressed(b,suppressions):return self.send_json(409,{"error":"suppressed"}) fields=[(c,b[c]) for c in ("website_domain","email","phone") if b[c]] if fields and db.execute("SELECT id FROM businesses WHERE organization_id=? AND ("+" OR ".join(f"{c}=?" for c,_ in fields)+")",[org]+[v for _,v in fields]).fetchone():return self.send_json(409,{"error":"duplicate"}) scored=score_business(b);cur=db.execute("INSERT INTO businesses(organization_id,name,website,website_domain,email,phone,description,province,city,suburb,score,score_version,score_factors,website_class) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)",(org,b["name"],b["website"],b["website_domain"],b["email"],b["phone"],str(b.get("description","")),b["province"],b["city"],b["suburb"],scored["score"],scored["score_version"],json.dumps(scored["factors"]),scored["website_class"])); self.audit(db,user,"business.created",str(cur.lastrowid));db.commit();return self.send_json(201,row_json(db.execute("SELECT * FROM businesses WHERE id=?",(cur.lastrowid,)).fetchone())) def create_suppression(self,payload,db,user): kind,value=payload.get("kind"),str(payload.get("value","")).strip().lower() if kind not in {"email","domain","phone"} or not value:return self.send_json(400,{"error":"invalid_suppression"}) try:db.execute("INSERT INTO suppressions(organization_id,kind,value,actor_user_id,active) VALUES(?,?,?,?,1)",(user["organization_id"],kind,value,user["id"])) except sqlite3.IntegrityError:pass self.audit(db,user,"suppression.created",kind);db.commit();return self.send_json(201,row_json(db.execute("SELECT * FROM suppressions WHERE organization_id=? AND kind=? AND value=?",(user["organization_id"],kind,value)).fetchone())) def child_business(self,db,bid,user):return self.business(db,bid,user["organization_id"]) def create_child(self,bid,table,payload,db,user): if not self.child_business(db,bid,user):return self.send_json(404,{"error":"not_found"}) if table=="contacts": email=str(payload.get("email","")).strip().lower(); phone=normalize_phone(payload.get("phone")); if email and not re.match(r"^[^@\s]+@[^@\s]+\.[^@\s]+$",email):return self.send_json(400,{"error":"invalid_contact"}) suppressed=is_suppressed({"email":email,"phone":phone},[dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1",(user["organization_id"],))]); values=(str(payload.get("name","")).strip(),email,phone,str(payload.get("title","")).strip(),int(bool(payload.get("do_not_contact"))) or int(suppressed)) elif table=="domains": value=normalize_domain(payload.get("domain")); if not value:return self.send_json(400,{"error":"invalid_domain"}) values=(value,str(payload.get("kind","other")).strip() or "other") elif table=="websites": value=str(payload.get("url","")).strip(); if not urlparse(value).scheme or not urlparse(value).netloc:return self.send_json(400,{"error":"invalid_website"}) values=(value,str(payload.get("website_class","business_site")).strip() or "business_site") elif table=="evidence": if not str(payload.get("kind","")).strip():return self.send_json(400,{"error":"invalid_evidence"}) values=(str(payload["kind"]).strip(),str(payload.get("url","")).strip(),str(payload.get("claim","")).strip()) else: if not str(payload.get("body","")).strip():return self.send_json(400,{"error":"body_required"}) values=(str(payload["body"]).strip(),) columns=CHILD_TABLES[table]; db.execute(f"INSERT INTO {table}(business_id,organization_id,{','.join(columns)}) VALUES(?, ?, {','.join('?' for _ in columns)})",(bid,user["organization_id"])+values); rid=db.execute("SELECT last_insert_rowid()").fetchone()[0];self.audit(db,user,f"{table}.created",str(rid));db.commit();return self.send_json(201,row_json(db.execute(f"SELECT * FROM {table} WHERE id=?",(rid,)).fetchone())) def update_pipeline(self,bid,payload,db,user): if not self.child_business(db,bid,user): return self.send_json(404,{"error":"not_found"}) stage=str(payload.get("stage","")).strip(); status=str(payload.get("status","active")).strip() or "active" if not stage or stage not in self.STAGES or status not in {"active","won","lost","paused"}: return self.send_json(400,{"error":"invalid_pipeline"}) key=str(payload.get("idempotency_key","")).strip(); org=user["organization_id"] if key: prior=db.execute("SELECT * FROM pipeline_entries WHERE organization_id=? AND idempotency_key=?",(org,key)).fetchone() if prior:return self.send_json(200,row_json(prior)) suppressed=self._crm_suppressed(db,org,self.business(db,bid,org)) if suppressed:return self.send_json(409,{"error":"do_not_contact","outreach_disabled":True}) cur=db.execute("INSERT INTO pipeline_entries(business_id,organization_id,stage,status,notes,next_action,follow_up_at,actor_user_id,idempotency_key) VALUES(?,?,?,?,?,?,?,?,?)",(bid,org,stage,status,str(payload.get("notes",payload.get("body","")))[:5000],str(payload.get("next_action",""))[:500],payload.get("follow_up_at"),user["id"],key or None)); rid=cur.lastrowid self.audit(db,user,"pipeline.created",str(rid)); db.commit(); return self.send_json(200 if getattr(self,"command","")=="PATCH" else 201,row_json(db.execute("SELECT * FROM pipeline_entries WHERE id=?",(rid,)).fetchone())) def update_pipeline_entry(self,eid,payload,db,user): org=user["organization_id"]; row=db.execute("SELECT * FROM pipeline_entries WHERE id=? AND organization_id=?",(eid,org)).fetchone() if not row:return self.send_json(404,{"error":"not_found"}) if payload.get("expected_updated_at") and payload["expected_updated_at"] != row["updated_at"]: return self.send_json(409,{"error":"conflict"}) values={k:payload[k] for k in ("stage","status","notes","next_action","follow_up_at") if k in payload} if "stage" in values and values["stage"] not in self.STAGES:return self.send_json(400,{"error":"invalid_stage"}) if "status" in values and values["status"] not in {"active","won","lost","paused"}:return self.send_json(400,{"error":"invalid_status"}) if not values:return self.send_json(400,{"error":"no_changes"}) cols=[]; args=[] for k,v in values.items():cols.append(k+"=?");args.append(v) args += [eid,org]; db.execute("UPDATE pipeline_entries SET "+",".join(cols)+",version=version+1,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",args); self.audit(db,user,"pipeline.updated",str(eid));db.commit() return self.send_json(200,row_json(db.execute("SELECT * FROM pipeline_entries WHERE id=?",(eid,)).fetchone())) def create_interaction(self,bid,payload,db,user): org=user["organization_id"]; business=self.business(db,bid,org) if not business:return self.send_json(404,{"error":"not_found"}) outcome=str(payload.get("outcome","other")).strip().lower(); kind=str(payload.get("kind","")).strip().lower() if not kind or outcome not in self.OUTCOMES:return self.send_json(400,{"error":"invalid_interaction"}) if self._crm_suppressed(db,org,business):return self.send_json(409,{"error":"do_not_contact","outreach_disabled":True}) key=str(payload.get("idempotency_key","")).strip() if key: prior=db.execute("SELECT * FROM interactions WHERE organization_id=? AND idempotency_key=?",(org,key)).fetchone() if prior:return self.send_json(200,row_json(prior)) cur=db.execute("INSERT INTO interactions(business_id,organization_id,kind,body,outcome,notes,next_action,follow_up_at,actor_user_id,idempotency_key) VALUES(?,?,?,?,?,?,?,?,?,?)",(bid,org,kind,str(payload.get("body",payload.get("notes","")))[:5000],outcome,str(payload.get("notes",""))[:5000],str(payload.get("next_action",""))[:500],payload.get("follow_up_at"),user["id"],key or None)); self.audit(db,user,"interaction.created",str(cur.lastrowid));db.commit() return self.send_json(201,row_json(db.execute("SELECT * FROM interactions WHERE id=?",(cur.lastrowid,)).fetchone())) def update_interaction(self,iid,payload,db,user): org=user["organization_id"]; row=db.execute("SELECT * FROM interactions WHERE id=? AND organization_id=?",(iid,org)).fetchone() if not row:return self.send_json(404,{"error":"not_found"}) values={k:payload[k] for k in ("kind","body","outcome","notes","next_action","follow_up_at") if k in payload} if "outcome" in values and values["outcome"] not in self.OUTCOMES:return self.send_json(400,{"error":"invalid_outcome"}) if not values:return self.send_json(400,{"error":"no_changes"}) cols=[];args=[] for k,v in values.items():cols.append(k+"=?");args.append(v) args += [iid,org];db.execute("UPDATE interactions SET "+",".join(cols)+",updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",args);self.audit(db,user,"interaction.updated",str(iid));db.commit();return self.send_json(200,row_json(db.execute("SELECT * FROM interactions WHERE id=?",(iid,)).fetchone())) def update_suppression(self,sid,payload,db,user): row=db.execute("SELECT * FROM suppressions WHERE id=? AND organization_id=?",(sid,user["organization_id"])).fetchone() if not row:return self.send_json(404,{"error":"not_found"}) if "active" not in payload:return self.send_json(400,{"error":"active_required"}) db.execute("UPDATE suppressions SET active=?,actor_user_id=?,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",(int(bool(payload["active"])),user["id"],sid,user["organization_id"]));self.audit(db,user,"suppression.updated",str(sid));db.commit();return self.send_json(200,row_json(db.execute("SELECT * FROM suppressions WHERE id=?",(sid,)).fetchone())) def import_suppressions(self,payload,db,user): items=payload.get("items",payload.get("suppressions")); if not isinstance(items,list) or len(items)>1000:return self.send_json(400,{"error":"invalid_suppression_import"}) imported=0 for item in items: if not isinstance(item,dict) or item.get("kind") not in {"email","domain","phone"} or not str(item.get("value","")).strip():return self.send_json(400,{"error":"invalid_suppression"}) kind=item["kind"];value=str(item["value"]).strip().lower(); cur=db.execute("INSERT OR IGNORE INTO suppressions(organization_id,kind,value,actor_user_id) VALUES(?,?,?,?)",(user["organization_id"],kind,value,user["id"]));imported += cur.rowcount self.audit(db,user,"suppressions.imported",str(imported));db.commit();return self.send_json(201,{"imported":imported,"received":len(items)}) def create_stage(self,payload,db,user): name=str(payload.get("name","")).strip().lower() if not name or len(name)>80 or not re.match(r"^[a-z0-9_-]+$",name):return self.send_json(400,{"error":"invalid_stage"}) try: position=int(payload.get("position",0)) except (TypeError,ValueError):return self.send_json(400,{"error":"invalid_stage"}) try:cur=db.execute("INSERT INTO pipeline_stages(organization_id,name,position) VALUES(?,?,?)",(user["organization_id"],name,position)) except sqlite3.IntegrityError:return self.send_json(409,{"error":"duplicate_stage"}) self.audit(db,user,"pipeline_stage.created",str(cur.lastrowid));db.commit();return self.send_json(201,row_json(db.execute("SELECT * FROM pipeline_stages WHERE id=?",(cur.lastrowid,)).fetchone())) def update_stage(self,sid,payload,db,user): row=db.execute("SELECT * FROM pipeline_stages WHERE id=? AND organization_id=?",(sid,user["organization_id"])).fetchone() if not row:return self.send_json(404,{"error":"not_found"}) values={k:payload[k] for k in ("name","position","active") if k in payload} if not values:return self.send_json(400,{"error":"no_changes"}) if "name" in values and (not isinstance(values["name"],str) or not re.match(r"^[a-z0-9_-]+$",values["name"])):return self.send_json(400,{"error":"invalid_stage"}) cols=[];args=[] for k,v in values.items():cols.append(k+"=?");args.append(int(bool(v)) if k=="active" else v) args += [sid,user["organization_id"]];db.execute("UPDATE pipeline_stages SET "+",".join(cols)+",updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",args);self.audit(db,user,"pipeline_stage.updated",str(sid));db.commit();return self.send_json(200,row_json(db.execute("SELECT * FROM pipeline_stages WHERE id=?",(sid,)).fetchone())) def delete_crm_item(self,table,ident,db,user): row=db.execute(f"SELECT id FROM {table} WHERE id=? AND organization_id=?",(ident,user["organization_id"])).fetchone() if not row:return self.send_json(404,{"error":"not_found"}) db.execute(f"DELETE FROM {table} WHERE id=? AND organization_id=?",(ident,user["organization_id"]));self.audit(db,user,table+".deleted",str(ident));db.commit();return self.send_json(200,{"ok":True,"id":ident}) def verify_business(self,bid,payload,db,user): if not self.child_business(db,bid,user):return self.send_json(404,{"error":"not_found"}) verified=bool(payload.get("verified",True)); now=datetime.now(timezone.utc).replace(microsecond=0).isoformat();db.execute("UPDATE businesses SET verified=?,verified_at=?,review_status=?,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",(int(verified),now if verified else None,"verified" if verified else "pending",bid,user["organization_id"]));self.audit(db,user,"business.verified",str(verified));db.commit();row=self.business(db,bid,user["organization_id"]);return self.send_json(200,row_json(row)) def preview_import(self,payload,db,org): rows=payload.get("rows",[]) if not isinstance(rows,list):return self.send_json(400,{"error":"rows_required"}) normalized=deduplicate_businesses([r for r in rows if isinstance(r,dict) and str(r.get("name","")).strip()]); suppressions=[dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1",(org,))];existing=[row_json(r) for r in db.execute("SELECT * FROM businesses WHERE organization_id=?",(org,))];seen=set();accepted=[];suppressed=0;existing_keys={deduplication_key(x) for x in existing} for b in normalized: key=deduplication_key(b) if is_suppressed(b,suppressions):suppressed+=1 elif key in existing_keys or key in seen:continue else:seen.add(key);accepted.append(b) return self.send_json(200,{"accepted":len(accepted),"duplicates":len(rows)-len(normalized)+len(normalized)-len(accepted)-suppressed,"suppressed":suppressed,"rows":accepted}) def list_sources(self,db,org): cols='id,organization_id,name,kind,enabled,health_status,consecutive_failures,circuit_open,last_success_at,last_failure_at,last_error,created_at,updated_at' return self.send_json(200,{"organization_id":org,"items":[row_json(r) for r in db.execute(f"SELECT {cols} FROM sources WHERE organization_id=? ORDER BY id",(org,))]}) def list_queries(self,db,org): return self.send_json(200,{"organization_id":org,"items":[row_json(r) for r in db.execute("SELECT * FROM discovery_queries WHERE organization_id=? ORDER BY id",(org,))]}) def list_source_records(self,db,org,q): try: limit=int(q.get('page_size',[50])[0]); offset=max(0,int(q.get('offset',[0])[0])) if limit<1 or limit>100: raise ValueError except (ValueError,TypeError): return self.send_json(400,{"error":"invalid_pagination"}) rows=db.execute("SELECT * FROM source_records WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?",(org,limit+1,offset)).fetchall(); out=[] for r in rows[:limit]: x=row_json(r) for k in ('raw_json','normalized_json','query_context_json','cursor_json','rate_policy_json'): try:x[k]=json.loads(x[k]) except (ValueError,TypeError):pass out.append(x) return self.send_json(200,{"organization_id":org,"items":out,"limit":limit,"offset":offset,"has_more":len(rows)>limit}) def create_source(self,payload,db,user): name=str(payload.get('name','')).strip(); kind=str(payload.get('kind','')).strip().lower(); config=payload.get('config',{}) if not name or kind not in ('csv','manual') or not isinstance(config,dict):return self.send_json(400,{"error":"invalid_source"}) if contains_secret(config):return self.send_json(400,{"error":"secret_not_permitted"}) try: validation=adapter_for(kind).validate(config) if config and not validation.valid:return self.send_json(400,{"error":"invalid_source_config","details":validation.errors}) cur=db.execute("INSERT INTO sources(organization_id,name,kind,enabled,config_json) VALUES(?,?,?,?,?)",(user['organization_id'],name,kind,int(bool(payload.get('enabled',False))),json.dumps(config,sort_keys=True))) except sqlite3.IntegrityError:return self.send_json(409,{"error":"duplicate_source"}) self.audit(db,user,'source.created',str(cur.lastrowid));db.commit();return self.send_json(201,row_json(db.execute("SELECT id,organization_id,name,kind,enabled,health_status,consecutive_failures,circuit_open,last_success_at,last_failure_at,last_error,created_at,updated_at FROM sources WHERE id=?",(cur.lastrowid,)).fetchone())) def update_source(self,sid,payload,db,user): if not db.execute("SELECT id FROM sources WHERE id=? AND organization_id=?",(sid,user['organization_id'])).fetchone():return self.send_json(404,{"error":"not_found"}) if 'enabled' not in payload:return self.send_json(400,{"error":"enabled_required"}) value=int(bool(payload['enabled']));db.execute("UPDATE sources SET enabled=?,updated_at=CURRENT_TIMESTAMP WHERE id=?",(value,sid));self.audit(db,user,'source.enabled' if value else 'source.disabled',str(sid));db.commit();return self.send_json(200,row_json(db.execute("SELECT * FROM sources WHERE id=?",(sid,)).fetchone())) def create_query(self,payload,db,user): sid=payload.get('source_id');name=str(payload.get('name','')).strip();query=payload.get('query',{}) if not isinstance(sid,int) or not name or not isinstance(query,dict) or contains_secret(query):return self.send_json(400,{"error":"invalid_query"}) if not db.execute("SELECT id FROM sources WHERE id=? AND organization_id=?",(sid,user['organization_id'])).fetchone():return self.send_json(404,{"error":"not_found"}) try:cur=db.execute("INSERT INTO discovery_queries(organization_id,source_id,name,query_json) VALUES(?,?,?,?)",(user['organization_id'],sid,name,json.dumps(query,sort_keys=True))) except sqlite3.IntegrityError:return self.send_json(409,{"error":"duplicate_query"}) self.audit(db,user,'discovery_query.created',str(cur.lastrowid));db.commit();return self.send_json(201,row_json(db.execute("SELECT * FROM discovery_queries WHERE id=?",(cur.lastrowid,)).fetchone())) def run_query(self,qid,db,user): if not db.execute("SELECT id FROM discovery_queries WHERE id=? AND organization_id=?",(qid,user['organization_id'])).fetchone():return self.send_json(404,{"error":"not_found"}) return self.create_job({"type":"source_discovery","_accepted":True,"payload":{"discovery_query_id":qid},"idempotency_key":f"discovery-query-{qid}-{int(time.time())}"},db,user) def test_source(self,sid,db,user): source=db.execute("SELECT * FROM sources WHERE id=? AND organization_id=?",(sid,user['organization_id'])).fetchone() if not source:return self.send_json(404,{"error":"not_found"}) try: result=adapter_for(source['kind']).validate(json.loads(source['config_json'])); ok=result.valid; error='; '.join(result.errors) if not ok else None except Exception as exc:ok=False;error=str(exc)[:300] if ok:db.execute("UPDATE sources SET health_status='healthy',consecutive_failures=0,circuit_open=0,last_success_at=CURRENT_TIMESTAMP,last_error=NULL WHERE id=?",(sid,));action='source.test.succeeded' else:db.execute("UPDATE sources SET health_status='unhealthy',consecutive_failures=consecutive_failures+1,circuit_open=CASE WHEN consecutive_failures+1>=3 THEN 1 ELSE circuit_open END,last_failure_at=CURRENT_TIMESTAMP,last_error=? WHERE id=?",(error,sid));action='source.test.failed' self.audit(db,user,action,str(sid));db.commit();return self.send_json(200,{"ok":ok,"errors":[] if ok else [error]}) def ingest_source(self,sid,payload,db,user): source=db.execute("SELECT * FROM sources WHERE id=? AND organization_id=?",(sid,user['organization_id'])).fetchone() if not source:return self.send_json(404,{"error":"not_found"}) if not source['enabled'] or source['circuit_open']:return self.send_json(409,{"error":"source_disabled"}) if contains_secret(payload):return self.send_json(400,{"error":"secret_not_permitted"}) config={k:v for k,v in payload.items() if k not in ('source_url','query_context','cursor','rate_policy')} if len(json.dumps(config).encode())>5*1024*1024:return self.send_json(400,{"error":"ingest_limits"}) if isinstance(config.get('rows'),list) and (len(config['rows'])>1000 or any(not isinstance(r,dict) or len(r)>50 or any(len(str(v))>10000 for v in r.values()) for r in config['rows'])):return self.send_json(400,{"error":"ingest_limits"}) if isinstance(config.get('csv'),str) and config['csv'].count('\n')>1001:return self.send_json(400,{"error":"ingest_limits"}) try:page=adapter_for(source['kind']).discover(config) except (ValueError,KeyError) as exc:return self.send_json(400,{"error":"invalid_ingest","detail":str(exc)}) inserted=0 for record in page.records[:1000]: raw=json.dumps(record,sort_keys=True,separators=(',',':'));digest=hashlib.sha256(raw.encode()).hexdigest() try:db.execute("INSERT INTO source_records(organization_id,source_id,content_hash,raw_json,normalized_json,source_url,query_context_json,cursor_json,rate_policy_json) VALUES(?,?,?,?,?,?,?,?,?)",(user['organization_id'],sid,digest,raw,raw,str(payload.get('source_url','')),json.dumps(payload.get('query_context',{}),sort_keys=True),json.dumps(payload.get('cursor',{}),sort_keys=True),json.dumps(payload.get('rate_policy',{}),sort_keys=True)));inserted+=1 except sqlite3.IntegrityError:pass db.execute("UPDATE sources SET health_status='healthy',consecutive_failures=0,last_success_at=CURRENT_TIMESTAMP,last_error=NULL WHERE id=?",(sid,));self.audit(db,user,'source.ingested',f'{sid}:{inserted}');db.commit();return self.send_json(201 if inserted else 200,{"inserted":inserted,"records":len(page.records)}) def matches(self,bid,db,org): source=self.business(db,bid,org) if not source:return self.send_json(404,{"error":"not_found"}) try: threshold=float(parse_qs(urlparse(self.path).query).get("threshold",["0.72"])[0]) except ValueError:return self.send_json(400,{"error":"invalid_threshold"}) rows=[row_json(r) for r in db.execute("SELECT * FROM businesses WHERE organization_id=? AND id<>? AND merge_status='active' ORDER BY id",(org,bid))] return self.send_json(200,{"business_id":bid,"threshold":threshold,"items":match_businesses(row_json(source),rows,threshold)}) def list_merge_history(self,db,org): rows=[row_json(r) for r in db.execute("SELECT * FROM merge_history WHERE organization_id=? ORDER BY id DESC",(org,))] for row in rows: for key in ("source_snapshot_json","child_reassignment_json"): try: row[key]=json.loads(row[key]) except (TypeError,ValueError): pass return self.send_json(200,{"organization_id":org,"items":rows}) def merge_business(self,bid,payload,db,user): org=user["organization_id"]; target_id=payload.get("target_business_id",payload.get("target_id")) if not isinstance(target_id,int) or target_id==bid:return self.send_json(400,{"error":"target_required"}) source=self.business(db,bid,org); target=self.business(db,target_id,org) if not source or not target:return self.send_json(404,{"error":"not_found"}) if source["merge_status"] != "active":return self.send_json(409,{"error":"source_already_merged"}) snapshot={"business":row_json(source),"children":{}} child_meta={} for table in ("business_identifiers","contacts","domains","websites","evidence","pipeline_entries","interactions","notes"): rows=[row_json(r) for r in db.execute(f"SELECT * FROM {table} WHERE business_id=? AND organization_id=? ORDER BY id",(bid,org))] snapshot["children"][table]=rows; child_meta[table]={"ids":[r["id"] for r in rows],"count":len(rows)} if rows: db.execute(f"UPDATE {table} SET business_id=? WHERE business_id=? AND organization_id=?",(target_id,bid,org)) cur=db.execute("INSERT INTO merge_history(organization_id,source_business_id,target_business_id,source_snapshot_json,child_reassignment_json,actor_user_id) VALUES(?,?,?,?,?,?)",(org,bid,target_id,json.dumps(snapshot,sort_keys=True),json.dumps(child_meta,sort_keys=True),user["id"])) db.execute("UPDATE businesses SET merge_status='merged',merged_into_id=?,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",(target_id,bid,org)) self.audit(db,user,"business.merged",f"{bid}->{target_id}");db.commit() return self.send_json(200,{"merge_history_id":cur.lastrowid,"source_business_id":bid,"target_business_id":target_id,"status":"merged","reassigned":child_meta}) def reverse_merge(self,hid,db,user): row=db.execute("SELECT * FROM merge_history WHERE id=? AND organization_id=?",(hid,user["organization_id"])).fetchone() if not row:return self.send_json(404,{"error":"not_found"}) if not row["reversible"]:return self.send_json(409,{"error":"merge_not_reversible"}) source=self.business(db,row["source_business_id"],user["organization_id"]); target=self.business(db,row["target_business_id"],user["organization_id"]) if not source or not target:return self.send_json(409,{"error":"business_missing"}) snapshot=json.loads(row["source_snapshot_json"]); ids=json.loads(row["child_reassignment_json"]) for table, meta in ids.items(): if not meta.get("ids"):continue marks=",".join("?" for _ in meta["ids"]) db.execute(f"UPDATE {table} SET business_id=? WHERE business_id=? AND organization_id=? AND id IN ({marks})",[source["id"],target["id"],user["organization_id"]]+meta["ids"]) db.execute("UPDATE businesses SET merge_status='active',merged_into_id=NULL,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",(source["id"],user["organization_id"])) db.execute("UPDATE merge_history SET reversible=0,reversed_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",(hid,user["organization_id"])) self.audit(db,user,"business.merge_reversed",str(hid));db.commit() return self.send_json(200,{"id":hid,"status":"reversed","source_business_id":source["id"],"target_business_id":target["id"]}) def log_message(self,*_):pass def _run_scoped_discovery(db, job, handler): payload = json.loads(job["payload"] or "{}") result = scoped_discover(payload.get("criteria", {}), payload.get("seed_urls"), max_pages=payload.get("max_pages", 20), max_candidates=payload.get("max_candidates", 50)) org = job["organization_id"]; persisted = [] for candidate in result["candidates"]: b = normalize_business(candidate) suppressed = is_suppressed(b, [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))]) if suppressed or not b["website_domain"]: continue existing = db.execute("SELECT id FROM businesses WHERE organization_id=? AND website_domain=?", (org, b["website_domain"])).fetchone() if existing: bid = existing["id"] else: scored = score_business(b) cur = db.execute("INSERT INTO businesses(organization_id,name,website,website_domain,email,phone,description,province,city,suburb,score,score_version,score_factors,website_class) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (org,b["name"],b["website"],b["website_domain"],b["email"],b["phone"],b.get("description", ""),b["province"],b["city"],b["suburb"],scored["score"],scored["score_version"],json.dumps(scored["factors"]),scored["website_class"])) bid = cur.lastrowid priority = "high" if scored["score"] >= 70 else "medium" if scored["score"] >= 40 else "low" db.execute("INSERT INTO score_history(organization_id,business_id,score,eligible,priority_band,score_version,explanations_json,signals_json) VALUES(?,?,?,?,?,?,?,?)", (org,bid,scored["score"],1,priority,scored["score_version"],json.dumps(scored["factors"]),json.dumps({"source": "scoped_discovery"}, sort_keys=True))) scan_ids = {} for page in candidate.get("pages", []): scan_key = hashlib.sha256((org + ":" + page["url"]).encode()).hexdigest() scan = db.execute("SELECT id FROM website_scans WHERE organization_id=? AND business_id=? AND cache_key=? ORDER BY id DESC LIMIT 1", (org, bid, scan_key)).fetchone() if scan: scan_ids[page["url"]] = scan["id"]; continue scan_result = {"input_url": page["url"], "final_url": page.get("final_url"), "status": page.get("status"), "title": page.get("title", ""), "headings": page.get("headings", []), "html": page.get("html", ""), "provenance": "scoped_discovery"} cur_scan = db.execute("INSERT INTO website_scans(organization_id,business_id,website_id,input_url,classification,result_json,cache_key,scanned_at,cache_expires_at) VALUES(?,?,?,?,?,?,?,?,?)", (org,bid,None,page["url"],"healthy",json.dumps(scan_result, sort_keys=True),scan_key,datetime.now(timezone.utc).replace(microsecond=0).isoformat(),None)) scan_ids[page["url"]] = cur_scan.lastrowid for page in candidate["evidence"]: db.execute("INSERT INTO evidence(business_id,organization_id,kind,url,claim) VALUES(?,?,?,?,?)", (bid, org, page["kind"], page["url"], page["claim"])) for contact in candidate["contacts"]: key = hashlib.sha256((str(job["id"]) + contact["source_url"] + contact["kind"] + contact["value"]).encode()).hexdigest() db.execute("INSERT OR IGNORE INTO contact_extractions(organization_id,business_id,website_scan_id,extraction_key,kind,value,label,classification,confidence,source_url,public_business,mx_status,suppressed,do_not_contact,provenance) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (org,bid,scan_ids.get(contact["source_url"]),key,contact["kind"],contact["value"],contact["label"],contact["classification"],contact["confidence"],contact["source_url"],1,"unknown",int(contact["suppressed"]),int(contact["do_not_contact"]),contact["provenance"])) persisted.append({"business_id": bid, "domain": b["website_domain"], "provenance": candidate["provenance"]}) safe_result = dict(result); safe_result["candidates"] = persisted db.execute("UPDATE discovery_runs SET result_json=?,result_count=?,updated_at=CURRENT_TIMESTAMP WHERE organization_id=? AND job_id=?", (json.dumps(safe_result, sort_keys=True), len(persisted), org, job["id"])) handler.add_job_event(db, job["id"], org, "discovery.completed", f"Persisted {len(persisted)} candidates", 100) def _job_worker(server): while not server.job_stop.is_set(): db=connect(server.db_path) try: job=db.execute("SELECT * FROM jobs WHERE status='queued' ORDER BY id LIMIT 1").fetchone() if not job: db.close(); server.job_wakeup.wait(.1); server.job_wakeup.clear(); continue changed=db.execute("UPDATE jobs SET status='running',attempts=attempts+1,started_at=COALESCE(started_at,CURRENT_TIMESTAMP),updated_at=CURRENT_TIMESTAMP WHERE id=? AND status='queued'",(job["id"],)).rowcount if not changed: db.close(); continue db.commit(); org=job["organization_id"]; jid=job["id"]; server_handler=object.__new__(ApiHandler) server_handler.add_job_event(db,jid,org,"started","Job started",0); db.commit() try: payload=json.loads(job["payload"] or "{}") except ValueError: payload={} if job["type"] == "scoped_discovery": try: _run_scoped_discovery(db, job, server_handler) db.execute("UPDATE jobs SET status='succeeded',progress=100,completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?", (jid,)); db.commit() except Exception as exc: db.execute("UPDATE jobs SET status='failed',error_code=?,completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?", (str(exc)[:80] or "DISCOVERY_FAILED", jid)); server_handler.add_job_event(db,jid,org,"failed","Discovery failed",job["progress"],str(exc)[:80]); db.commit() continue try: steps=1 if job["type"]=="noop" else max(1,min(int(payload.get("steps",5)),20)) except (ValueError,TypeError): steps=5 cancelled=False for i in range(steps): time.sleep(.01) fresh=db.execute("SELECT status FROM jobs WHERE id=?",(jid,)).fetchone() if not fresh or fresh["status"]=="cancelled": cancelled=True; break progress=int((i+1)*100/steps); db.execute("UPDATE jobs SET progress=?,updated_at=CURRENT_TIMESTAMP WHERE id=? AND status='running'",(progress,jid)); server_handler.add_job_event(db,jid,org,"progress",f"Job progress {progress}%",progress); db.commit() if cancelled: continue if payload.get("force_fail") or (payload.get("fail_once") and job["attempts"] == 0): db.execute("UPDATE jobs SET status='failed',error_code='DEMO_FAILURE',completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));server_handler.add_job_event(db,jid,org,"failed","Job failed",job["progress"],"DEMO_FAILURE") else: db.execute("UPDATE jobs SET status='succeeded',progress=100,completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));server_handler.add_job_event(db,jid,org,"succeeded","Job completed",100) db.commit() finally: db.close() def create_server(host="127.0.0.1",port=8000,db_path="prospects.db"): load_config() configure_ai_research_db(db_path) server=ThreadingHTTPServer((host,port),ApiHandler);server.db_path=db_path;connect(db_path).close();server.job_stop=threading.Event();server.job_wakeup=threading.Event();server.job_thread=threading.Thread(target=_job_worker,args=(server,),daemon=True);server.job_thread.start() original_close=server.server_close def close(): server.job_stop.set();server.job_wakeup.set();server.job_thread.join(timeout=2);original_close() server.server_close=close return server if __name__=="__main__": parser=argparse.ArgumentParser();parser.add_argument("--host",default="127.0.0.1");parser.add_argument("--port",type=int,default=int(os.environ.get("PROSPECT_API_PORT","8000")));parser.add_argument("--db",default=os.environ.get("PROSPECT_API_DB","prospects.db"));args=parser.parse_args();server=create_server(args.host,args.port,args.db);print(f"Prospect API listening on http://{args.host}:{args.port}",flush=True) try:server.serve_forever() except KeyboardInterrupt:pass finally:server.server_close()