Files
MarketingTool/apps/api/app/main.py
T

1393 lines
131 KiB
Python

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, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS
from app.discovery import discover as scoped_discover
from app.config import load_config
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, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS
from .discovery import discover as scoped_discover
from .config import load_config
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 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], 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"})
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/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/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")
if not isinstance(criteria, dict) or not isinstance(seeds, list) or not seeds: return self.send_json(400, {"error": "seed_urls_required"})
try:
if len(seeds) > 5 or 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"})
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, "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)))
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/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()
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()