1445 lines
136 KiB
Python
1445 lines
136 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, input_fingerprint, provider_status, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS
|
|
from app.discovery import discover as scoped_discover
|
|
from app.search_provider import provider_status as search_provider_status
|
|
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, input_fingerprint, provider_status, MAX_INPUT_ITEMS, MAX_OUTPUT_CHARS
|
|
from .discovery import discover as scoped_discover
|
|
from .search_provider import provider_status as search_provider_status
|
|
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 ai_provider_config(self, db, user, payload=None):
|
|
org = user["organization_id"]
|
|
row = db.execute("SELECT * FROM ai_provider_configs WHERE organization_id=?", (org,)).fetchone()
|
|
if payload is not None:
|
|
provider = str(payload.get("provider", "local")).strip().lower()
|
|
if not provider or len(provider) > 80: return self.send_json(400, {"error": "invalid_provider"})
|
|
reviewed = bool(payload.get("reviewed", False))
|
|
enabled = bool(payload.get("enabled", True))
|
|
# A remote provider cannot be enabled by configuration alone: this
|
|
# service has no reviewed transport abstraction and never sends data.
|
|
if provider not in {"local", "deterministic"} and (enabled or reviewed):
|
|
return self.send_json(409, {"error": "provider_not_reviewed", "network_enabled": False})
|
|
db.execute("INSERT INTO ai_provider_configs(organization_id,provider,enabled,reviewed) VALUES(?,?,?,?) ON CONFLICT(organization_id) DO UPDATE SET provider=excluded.provider,enabled=excluded.enabled,reviewed=excluded.reviewed,updated_at=CURRENT_TIMESTAMP", (org, provider, int(enabled), int(reviewed)))
|
|
self.audit(db, user, "ai_provider.updated", provider); db.commit()
|
|
row = db.execute("SELECT * FROM ai_provider_configs WHERE organization_id=?", (org,)).fetchone()
|
|
configured = row["provider"] if row and row["enabled"] else None
|
|
result = provider_status(configured if configured is not None else None)
|
|
result.update({"organization_id": org, "configured": bool(row), "enabled": bool(row and row["enabled"]), "reviewed": bool(row and row["reviewed"])})
|
|
return self.send_json(200, result)
|
|
|
|
def _ai_current_fingerprint(self, db, row):
|
|
if not row["business_id"]: return ""
|
|
try: limit = int(json.loads(row["prompt_metadata_json"] or "{}").get("request", {}).get("max_items", MAX_INPUT_ITEMS))
|
|
except (TypeError, ValueError): limit = MAX_INPUT_ITEMS
|
|
limit = max(1, min(limit, MAX_INPUT_ITEMS))
|
|
org, bid = row["organization_id"], row["business_id"]
|
|
business = self.business(db, bid, org)
|
|
if not business: return ""
|
|
scans = [dict(x) for x in db.execute("SELECT id,input_url,classification,result_json,scanned_at FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, limit))]
|
|
evidence = [dict(x) for x in db.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, limit))]
|
|
contacts = [dict(x) for x in db.execute("SELECT * FROM contacts WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, limit))]
|
|
extracted = [dict(x) for x in db.execute("SELECT * FROM contact_extractions WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, limit))]
|
|
return input_fingerprint(row_json(business), scans, contacts + extracted, evidence, limit)
|
|
|
|
def suggest_ai(self, bid, payload, db, user):
|
|
org = user["organization_id"]
|
|
business = self.business(db, bid, org)
|
|
if not business: return self.send_json(404, {"error": "not_found"})
|
|
business_dict = row_json(business)
|
|
suppressions = [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))]
|
|
contacts = [dict(r) for r in db.execute("SELECT * FROM contacts WHERE business_id=? AND organization_id=?", (bid, org))]
|
|
extracted = [dict(r) for r in db.execute("SELECT * FROM contact_extractions WHERE business_id=? AND organization_id=?", (bid, org))]
|
|
if is_suppressed(business_dict, suppressions) or any(bool(c.get("do_not_contact") or c.get("suppressed")) for c in contacts + extracted):
|
|
return self.send_json(409, {"error": "ai_blocked_suppressed_business", "business_id": bid})
|
|
try:
|
|
requested = int(payload.get("max_items", MAX_INPUT_ITEMS))
|
|
if requested < 1 or requested > MAX_INPUT_ITEMS: raise ValueError
|
|
except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_limits"})
|
|
scans = [dict(r) for r in db.execute("SELECT id,input_url,classification,result_json,scanned_at FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, requested)).fetchall()]
|
|
evidence = [dict(r) for r in db.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, requested)).fetchall()]
|
|
history = [dict(r) for r in db.execute("SELECT score,eligible,priority_band,score_version,explanations_json,created_at FROM score_history WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?", (bid, org, requested)).fetchall()]
|
|
status, provider, version, metadata = generate_ai(business_dict, scans[:requested], contacts[:requested] + extracted[:requested], evidence[:requested], history[:requested])
|
|
output = metadata.pop("output", {}) if status == "succeeded" else {"suggestions": [], "grounded": True, "claim_policy": "stored_evidence_only"}
|
|
hashes = metadata.get("evidence_hashes", [])
|
|
cur = db.execute("INSERT INTO ai_runs(organization_id,business_id,input_evidence_hashes_json,model,provider,version,prompt_metadata_json,data_minimization_json,status,approval_state,output_json,actor_user_id) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)", (org, bid, json.dumps(hashes, sort_keys=True), "local-deterministic" if provider else "", provider, version, json.dumps({"request": {"max_items": requested}, "output_limit": MAX_OUTPUT_CHARS}, sort_keys=True), json.dumps(metadata, sort_keys=True), status, "pending", json.dumps(output, sort_keys=True), user["id"]))
|
|
run_id = cur.lastrowid
|
|
for suggestion in output.get("suggestions", []):
|
|
db.execute("INSERT INTO ai_suggestions(ai_run_id,organization_id,business_id,suggestion_type,citations_json,output_json) VALUES(?,?,?,?,?,?)", (run_id, org, bid, str(suggestion.get("type", ""))[:80], json.dumps(suggestion.get("citations", []), sort_keys=True), json.dumps(suggestion, sort_keys=True)))
|
|
self.audit(db, user, "ai.suggested", str(run_id)); db.commit()
|
|
row = db.execute("SELECT * FROM ai_runs WHERE id=? AND organization_id=?", (run_id, org)).fetchone()
|
|
return self.send_json(201, self._ai_run_json(row, output.get("suggestions", [])))
|
|
|
|
def list_ai_runs(self, db, user, query):
|
|
try:
|
|
limit = int((query.get("page_size") or [50])[0]); offset = max(0, int((query.get("offset") or [0])[0]))
|
|
if limit < 1 or limit > 100: raise ValueError
|
|
except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"})
|
|
rows = db.execute("SELECT * FROM ai_runs WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?", (user["organization_id"], limit + 1, offset)).fetchall()
|
|
items = []
|
|
for row in rows[:limit]:
|
|
suggestions = [json.loads(x[0]) for x in db.execute("SELECT output_json FROM ai_suggestions WHERE ai_run_id=? AND organization_id=? ORDER BY id", (row["id"], user["organization_id"]))]
|
|
items.append(self._ai_run_json(row, suggestions))
|
|
return self.send_json(200, {"organization_id": user["organization_id"], "items": items, "limit": limit, "offset": offset, "has_more": len(rows) > limit})
|
|
|
|
def decide_ai(self, run_id, decision, db, user):
|
|
row = db.execute("SELECT * FROM ai_runs WHERE id=? AND organization_id=?", (run_id, user["organization_id"])).fetchone()
|
|
if not row: return self.send_json(404, {"error": "not_found"})
|
|
if decision == "approve" and row["status"] != "succeeded": return self.send_json(409, {"error": "run_not_approvable"})
|
|
if row["approval_state"] != "pending": return self.send_json(409, {"error": "already_decided"})
|
|
if decision == "approve":
|
|
try: metadata = json.loads(row["data_minimization_json"] or "{}")
|
|
except (TypeError, ValueError): metadata = {}
|
|
expected = metadata.get("input_fingerprint")
|
|
if expected and expected != self._ai_current_fingerprint(db, row):
|
|
return self.send_json(409, {"error": "ai_run_stale", "approval_state": "pending", "reason": "source_hash_changed"})
|
|
now = datetime.now(timezone.utc).replace(microsecond=0).isoformat()
|
|
if decision == "approve": db.execute("UPDATE ai_runs SET approval_state='approved',approved_at=? WHERE id=? AND organization_id=?", (now, run_id, user["organization_id"]))
|
|
else: db.execute("UPDATE ai_runs SET approval_state='rejected',rejected_at=? WHERE id=? AND organization_id=?", (now, run_id, user["organization_id"]))
|
|
self.audit(db, user, "ai." + decision, str(run_id)); db.commit()
|
|
return self.send_json(200, self._ai_run_json(db.execute("SELECT * FROM ai_runs WHERE id=?", (run_id,)).fetchone()))
|
|
|
|
def do_GET(self):
|
|
parsed=urlparse(self.path); path=parsed.path.rstrip("/")
|
|
if path=="/api/v1/health/live": return self.send_json(200,{"status":"ok","organization_id":ORGANIZATION_ID,"outreach_enabled":False})
|
|
if path=="/api/v1/health/ready":
|
|
try:
|
|
check=connect(self.server.db_path); check.execute("SELECT 1"); check.close()
|
|
return self.send_json(200,{"status":"ok","ready":True,"database":"ok","outreach_enabled":False})
|
|
except (sqlite3.Error, OSError):
|
|
return self.send_json(503,{"status":"not_ready","ready":False,"database":"error","outreach_enabled":False})
|
|
db=self.db()
|
|
try:
|
|
user=self.require_auth(db)
|
|
if not user:return
|
|
org=user["organization_id"]
|
|
if path=="/api/v1/auth/me": return self.send_json(200,{"id":user["id"],"email":user["email"],"role":user["role"],"organization_id":org})
|
|
if path=="/api/v1/admin/users":
|
|
if user["role"] not in {"owner","admin"}: return self.send_json(403,{"error":"forbidden"})
|
|
return self.send_json(200,{"items":[dict(r) for r in db.execute("SELECT id,email,role,organization_id,created_at FROM users WHERE organization_id=? ORDER BY id",(org,))]})
|
|
if path=="/api/v1/dashboard/summary":
|
|
row=db.execute("SELECT COUNT(*) businesses,COALESCE(AVG(score),0) average_score FROM businesses WHERE organization_id=?",(org,)).fetchone()
|
|
counts={"new":db.execute("SELECT COUNT(*) FROM businesses WHERE organization_id=? AND created_at>=datetime('now','-7 days')",(org,)).fetchone()[0],"hot":db.execute("SELECT COUNT(*) FROM businesses WHERE organization_id=? AND score>=70 AND merge_status='active'",(org,)).fetchone()[0],"review":db.execute("SELECT COUNT(*) FROM businesses WHERE organization_id=? AND review_status='pending'",(org,)).fetchone()[0],"source_health":db.execute("SELECT COUNT(*) FROM sources WHERE organization_id=? AND health_status IN ('healthy','unhealthy')",(org,)).fetchone()[0],"active_jobs":db.execute("SELECT COUNT(*) FROM jobs WHERE organization_id=? AND status IN ('queued','running')",(org,)).fetchone()[0]}
|
|
clickable={key:{"count":value,"filter":{"dashboard_filter":key}} for key,value in counts.items()}
|
|
return self.send_json(200,{"organization_id":org,"businesses":row["businesses"],"average_score":round(row["average_score"],2),"suppressed":db.execute("SELECT COUNT(*) FROM suppressions WHERE organization_id=?",(org,)).fetchone()[0],"counts":counts,"clickable_filters":clickable,"quick_filters":[{"key":key,"count":value,"filter":{"dashboard_filter":key}} for key,value in counts.items()]})
|
|
if path=="/api/v1/saved-filters": return self.list_saved_filters(db,user)
|
|
if path=="/api/v1/review-queue": return self.review_queue(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/score-rules": return self.list_score_rules(db,org)
|
|
if path=="/api/v1/scoring/summary": return self.scoring_summary(db,org)
|
|
if path=="/api/v1/businesses": return self.list_businesses(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/merge-history": return self.list_merge_history(db,org)
|
|
if path=="/api/v1/sources": return self.list_sources(db,org)
|
|
if path=="/api/v1/discovery-queries": return self.list_queries(db,org)
|
|
if path=="/api/v1/discovery-runs": return self.list_discovery_runs(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/discovery/provider-status": return self.send_json(200, search_provider_status())
|
|
if path=="/api/v1/source-records": return self.list_source_records(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/jobs": return self.list_jobs(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/domain-checks": return self.list_domain_checks(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/website-scans": return self.list_website_scans(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/contact-extractions": return self.list_contact_extractions(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/pipeline-entries": return self.list_pipeline(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/pipeline-stages": return self.list_stages(db,org)
|
|
if path=="/api/v1/outcomes": return self.send_json(200,{"items":sorted(self.OUTCOMES)})
|
|
if path=="/api/v1/interactions": return self.list_interactions(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/suppressions": return self.list_suppressions(db,org,parse_qs(parsed.query))
|
|
if path=="/api/v1/ai-runs": return self.list_ai_runs(db,user,parse_qs(parsed.query))
|
|
if path=="/api/v1/ai/provider-config": return self.ai_provider_config(db,user)
|
|
if path=="/api/v1/outreach/drafts": return self.list_outreach_drafts(db,user,parse_qs(parsed.query))
|
|
if path=="/api/v1/outreach/provider-config": return self.provider_config(db,user)
|
|
if path in ("/api/v1/reports/pipeline","/api/v1/reports/outcomes","/api/v1/reports/activity"): return self.report(db,org,path.rsplit('/',1)[1],parse_qs(parsed.query))
|
|
if path.startswith("/api/v1/jobs/"): return self.get_job_route(db,org,path,parse_qs(parsed.query))
|
|
if path.startswith("/api/v1/businesses/"):
|
|
bits=path.split("/"); ident=bits[4] if len(bits)>4 else ""
|
|
if not ident.isdigit(): return self.send_json(404,{"error":"not_found"})
|
|
row=self.business(db,int(ident),org)
|
|
if not row:return self.send_json(404,{"error":"not_found"})
|
|
if len(bits)==7 and bits[5:]==["websites","scan"]: return self.get_latest_website_scan(int(ident),db,user)
|
|
if len(bits)==7 and bits[5:]==["domains","check"]: return self.get_domain_check(int(ident),db,user,parse_qs(parsed.query))
|
|
if len(bits)==7 and bits[5:]==["contacts","extract"]: return self.send_json(405,{"error":"method_not_allowed"})
|
|
if len(bits)==6 and bits[5]=="domain-candidates": return self.list_domain_candidates(int(ident),db,org)
|
|
if len(bits)==6 and bits[5]=="matches": return self.matches(int(ident),db,org)
|
|
payload=row_json(row); payload.update(self.nested(db,int(ident),org)); return self.send_json(200,payload)
|
|
return self.send_json(404,{"error":"not_found"})
|
|
finally: db.close()
|
|
def _domain_result(self, row, cache_hit=False):
|
|
try: result=json.loads(row["result_json"] or "{}")
|
|
except (TypeError,ValueError): result={}
|
|
result.update({"id":row["id"],"business_id":row["business_id"],"domain":row["domain"],"status":row["status"],"cache_hit":cache_hit,"checked_at":row["checked_at"],"cache_expires_at":row["cache_expires_at"]})
|
|
return result
|
|
|
|
def _check_domain(self, bid, domain, db, user):
|
|
normalized=normalize_registrable_domain(domain)
|
|
if normalized == "unknown": return self.send_json(400,{"error":"unsupported_domain","status":"unknown"})
|
|
now=datetime.now(timezone.utc).replace(microsecond=0).isoformat()
|
|
cache_key=normalized+":a_aaaa:v1"
|
|
cached=db.execute("SELECT * FROM domain_checks WHERE organization_id=? AND business_id=? AND cache_key=? AND cache_expires_at>? ORDER BY id DESC LIMIT 1",(user["organization_id"],bid,cache_key,now)).fetchone()
|
|
if cached:return self.send_json(200,self._domain_result(cached,True))
|
|
result=resolve_domain(normalized)
|
|
try: cur=db.execute("INSERT INTO domain_checks(organization_id,business_id,domain,status,result_json,cache_key,checked_at,cache_expires_at) VALUES(?,?,?,?,?,?,?,?)",(user["organization_id"],bid,normalized,result["status"],json.dumps(result,sort_keys=True),cache_key,result["checked_at"],result["cache_expires_at"]))
|
|
except sqlite3.IntegrityError:
|
|
existing=db.execute("SELECT id FROM domain_checks WHERE organization_id=? AND business_id=? AND cache_key=?",(user["organization_id"],bid,cache_key)).fetchone()
|
|
if not existing: return self.send_json(409,{"error":"domain_check_conflict"})
|
|
db.execute("UPDATE domain_checks SET domain=?,status=?,result_json=?,checked_at=?,cache_expires_at=? WHERE id=?",(normalized,result["status"],json.dumps(result,sort_keys=True),result["checked_at"],result["cache_expires_at"],existing["id"]))
|
|
self.audit(db,user,"domain.checked",f"{bid}:{normalized}:{result['status']}"); db.commit()
|
|
return self.send_json(200,self._domain_result(db.execute("SELECT * FROM domain_checks WHERE id=?",(existing["id"],)).fetchone()))
|
|
self.audit(db,user,"domain.checked",f"{bid}:{normalized}:{result['status']}"); db.commit()
|
|
return self.send_json(200,self._domain_result(db.execute("SELECT * FROM domain_checks WHERE id=?",(cur.lastrowid,)).fetchone()))
|
|
|
|
def post_domain_check(self,bid,payload,db,user):
|
|
if not self.business(db,bid,user["organization_id"]): return self.send_json(404,{"error":"not_found"})
|
|
domain=payload.get("domain") or self.business(db,bid,user["organization_id"])["website_domain"]
|
|
if not str(domain).strip(): return self.send_json(400,{"error":"domain_required"})
|
|
return self._check_domain(bid,str(domain),db,user)
|
|
|
|
def get_domain_check(self,bid,db,user,query):
|
|
domain=(query.get("domain") or [""])[0]
|
|
if not domain:return self.send_json(400,{"error":"domain_required"})
|
|
return self._check_domain(bid,domain,db,user)
|
|
|
|
def list_domain_checks(self,db,org,query):
|
|
try: limit=max(1,min(int((query.get("page_size") or [50])[0]),100)); offset=max(0,int((query.get("offset") or [0])[0]))
|
|
except (ValueError,TypeError): return self.send_json(400,{"error":"invalid_pagination"})
|
|
rows=db.execute("SELECT * FROM domain_checks WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?",(org,limit+1,offset)).fetchall()
|
|
return self.send_json(200,{"organization_id":org,"items":[self._domain_result(r) for r in rows[:limit]],"limit":limit,"offset":offset,"has_more":len(rows)>limit})
|
|
|
|
def _website_scan_result(self, row, cache_hit=False):
|
|
result = json.loads(row["result_json"] or "{}")
|
|
result.update({"id": row["id"], "business_id": row["business_id"], "website_id": row["website_id"], "input_url": row["input_url"], "classification": row["classification"], "scanned_at": row["scanned_at"], "cache_expires_at": row["cache_expires_at"], "cache_hit": cache_hit})
|
|
return result
|
|
|
|
def list_website_scans(self, db, org, query):
|
|
try:
|
|
limit = int((query.get("page_size") or [50])[0]); offset = max(0, int((query.get("offset") or [0])[0]))
|
|
if limit < 1 or limit > WEBSITE_SCAN_PAGE_SIZE: raise ValueError
|
|
except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"})
|
|
params = [org]; where = ["organization_id=?"]
|
|
if (query.get("business_id") or [""])[0].isdigit(): where.append("business_id=?"); params.append(int(query["business_id"][0]))
|
|
rows = db.execute("SELECT * FROM website_scans WHERE " + " AND ".join(where) + " ORDER BY id DESC LIMIT ? OFFSET ?", params + [limit + 1, offset]).fetchall()
|
|
return self.send_json(200, {"organization_id": org, "items": [self._website_scan_result(r) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit})
|
|
|
|
def _contact_json(self, row):
|
|
result = row_json(row)
|
|
for key in ("public_business", "suppressed", "do_not_contact"):
|
|
if key in result: result[key] = bool(result[key])
|
|
return result
|
|
|
|
def list_contact_extractions(self, db, org, query):
|
|
try:
|
|
limit = int((query.get("page_size") or [50])[0]); offset = max(0, int((query.get("offset") or [0])[0]))
|
|
if limit < 1 or limit > CONTACT_EXTRACTION_PAGE_SIZE: raise ValueError
|
|
except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"})
|
|
params = [org]; where = ["organization_id=?"]
|
|
business_id = (query.get("business_id") or [""])[0]
|
|
if business_id.isdigit(): where.append("business_id=?"); params.append(int(business_id))
|
|
rows = db.execute("SELECT * FROM contact_extractions WHERE " + " AND ".join(where) + " ORDER BY id DESC LIMIT ? OFFSET ?", params + [limit + 1, offset]).fetchall()
|
|
return self.send_json(200, {"organization_id": org, "items": [self._contact_json(r) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit})
|
|
|
|
def extract_business_contacts(self, bid, payload, db, user):
|
|
org = user["organization_id"]; business = self.business(db, bid, org)
|
|
if not business: return self.send_json(404, {"error": "not_found"})
|
|
scan_id = payload.get("website_scan_id", payload.get("scan_id"))
|
|
scan = None
|
|
if scan_id is not None:
|
|
if not isinstance(scan_id, int): return self.send_json(400, {"error": "invalid_scan"})
|
|
scan = db.execute("SELECT * FROM website_scans WHERE id=? AND business_id=? AND organization_id=?", (scan_id, bid, org)).fetchone()
|
|
if not scan: return self.send_json(404, {"error": "scan_not_found"})
|
|
try: stored = json.loads(scan["result_json"] or "{}")
|
|
except (TypeError, ValueError): stored = {}
|
|
source_url = str(payload.get("source_url") or scan["input_url"]).strip()
|
|
source_html = payload.get("html") if isinstance(payload.get("html"), str) else stored.get("html")
|
|
if source_html is None: return self.send_json(409, {"error": "scan_html_unavailable"})
|
|
if source_url != scan["input_url"] and source_url != (stored.get("final_url") or ""): return self.send_json(400, {"error": "source_not_approved"})
|
|
else:
|
|
source_url = str(payload.get("source_url") or "").strip()
|
|
source_html = payload.get("html")
|
|
if not source_url or not isinstance(source_html, str): return self.send_json(400, {"error": "approved_scan_required"})
|
|
parsed = urlparse(source_url)
|
|
official_hosts = {str(business["website_domain"]).lower().strip(".")}
|
|
if business["website"]:
|
|
official_hosts.add((urlparse(business["website"]).hostname or "").lower().strip("."))
|
|
if parsed.scheme not in {"http", "https"} or not parsed.hostname or parsed.hostname.lower().strip(".") not in official_hosts:
|
|
return self.send_json(400, {"error": "source_not_approved"})
|
|
if len(source_html.encode("utf-8")) > MAX_HTML_BYTES: return self.send_json(413, {"error": "html_too_large"})
|
|
try: requested_limit = int(payload.get("limit", MAX_RESULTS))
|
|
except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_limits"})
|
|
if requested_limit < 1 or requested_limit > MAX_RESULTS: return self.send_json(400, {"error": "invalid_limits"})
|
|
key = str(payload.get("idempotency_key") or hashlib.sha256((str(scan_id or "") + source_url + source_html).encode()).hexdigest())[:200]
|
|
existing = db.execute("SELECT * FROM contact_extractions WHERE organization_id=? AND business_id=? AND extraction_key=? ORDER BY id", (org, bid, key)).fetchall()
|
|
if existing: return self.send_json(200, {"business_id": bid, "extraction_key": key, "items": [self._contact_json(r) for r in existing], "idempotent": True})
|
|
suppressions = [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))]
|
|
try: found = extract_contacts(source_html, source_url, suppressions=suppressions, max_results=requested_limit)
|
|
except ValueError as exc: return self.send_json(400, {"error": str(exc)})
|
|
for item in found:
|
|
db.execute("INSERT INTO contact_extractions(organization_id,business_id,website_scan_id,extraction_key,kind,value,label,classification,confidence,source_url,public_business,mx_status,suppressed,do_not_contact,provenance) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (org,bid,scan["id"] if scan else None,key,item["kind"],item["value"],item["label"],item["classification"],item["confidence"],item["source_url"],int(item["public_business"]),item["mx_status"],int(item["suppressed"]),int(item["do_not_contact"]),item["provenance"]))
|
|
self.audit(db, user, "contacts.extracted", f"{bid}:{len(found)}:{key}"); db.commit()
|
|
rows = db.execute("SELECT * FROM contact_extractions WHERE organization_id=? AND business_id=? AND extraction_key=? ORDER BY id", (org,bid,key)).fetchall()
|
|
return self.send_json(201, {"business_id": bid, "extraction_key": key, "items": [self._contact_json(r) for r in rows], "idempotent": False})
|
|
|
|
def get_latest_website_scan(self, bid, db, user):
|
|
row = db.execute("SELECT * FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT 1", (bid, user["organization_id"])).fetchone()
|
|
if not row: return self.send_json(404, {"error": "scan_not_found"})
|
|
return self.send_json(200, self._website_scan_result(row, False))
|
|
|
|
def scan_business_website(self, bid, payload, db, user):
|
|
org = user["organization_id"]; business = self.business(db, bid, org)
|
|
if not business: return self.send_json(404, {"error": "not_found"})
|
|
requested = str(payload.get("url", "")).strip() if isinstance(payload, dict) else ""
|
|
if not requested:
|
|
child = db.execute("SELECT * FROM websites WHERE business_id=? AND organization_id=? ORDER BY id LIMIT 1", (bid, org)).fetchone()
|
|
requested = (child["url"] if child else business["website"]) or ""
|
|
try: safe_url = validate_url(requested)
|
|
except ValueError as exc:
|
|
self.audit(db, user, "website.scan.rejected", f"{bid}:{str(exc)}"); db.commit()
|
|
return self.send_json(400, {"error": "unsafe_url", "reason": str(exc)})
|
|
website = db.execute("SELECT id FROM websites WHERE business_id=? AND organization_id=? AND url=? ORDER BY id LIMIT 1", (bid, org, safe_url)).fetchone()
|
|
cache_key = hashlib.sha256(safe_url.encode()).hexdigest(); now = datetime.now(timezone.utc).replace(microsecond=0); expires = now + timedelta(seconds=WEBSITE_SCAN_CACHE_SECONDS)
|
|
cached = db.execute("SELECT * FROM website_scans WHERE organization_id=? AND business_id=? AND cache_key=? AND cache_expires_at>? ORDER BY id DESC LIMIT 1", (org, bid, cache_key, now.isoformat())).fetchone()
|
|
if cached:
|
|
self.audit(db, user, "website.scan.cache_hit", f"{bid}:{safe_url}"); db.commit()
|
|
return self.send_json(200, self._website_scan_result(cached, True))
|
|
result = scan_website(safe_url); result["business_id"] = bid
|
|
cur = db.execute("INSERT INTO website_scans(organization_id,business_id,website_id,input_url,classification,result_json,cache_key,scanned_at,cache_expires_at) VALUES(?,?,?,?,?,?,?,?,?)", (org, bid, website["id"] if website else None, safe_url, result["classification"], json.dumps(result, sort_keys=True), cache_key, now.isoformat(), expires.isoformat()))
|
|
self.audit(db, user, "website.scanned", f"{bid}:{result['classification']}"); db.commit()
|
|
return self.send_json(201, self._website_scan_result(db.execute("SELECT * FROM website_scans WHERE id=?", (cur.lastrowid,)).fetchone()))
|
|
|
|
def list_domain_candidates(self,bid,db,org):
|
|
business=self.business(db,bid,org)
|
|
if not business:return self.send_json(404,{"error":"not_found"})
|
|
rows=db.execute("SELECT * FROM domain_candidates WHERE organization_id=? AND business_id=? ORDER BY rank,id",(org,bid)).fetchall()
|
|
if not rows:
|
|
generated=generate_candidate_domains(business["name"],business["description"],business["city"])
|
|
for rank,domain in enumerate(generated,1):
|
|
db.execute("INSERT OR IGNORE INTO domain_candidates(organization_id,business_id,domain,rank) VALUES(?,?,?,?)",(org,bid,domain,rank))
|
|
db.commit(); rows=db.execute("SELECT * FROM domain_candidates WHERE organization_id=? AND business_id=? ORDER BY rank,id",(org,bid)).fetchall()
|
|
return self.send_json(200,{"business_id":bid,"items":[row_json(r) for r in rows]})
|
|
|
|
def check_availability(self,bid,payload,db,user):
|
|
if not self.business(db,bid,user["organization_id"]):return self.send_json(404,{"error":"not_found"})
|
|
candidates=db.execute("SELECT domain FROM domain_candidates WHERE organization_id=? AND business_id=? ORDER BY rank,id",(user["organization_id"],bid)).fetchall()
|
|
domains=[r["domain"] for r in candidates]
|
|
if isinstance(payload.get("domains"),list): domains=[normalize_registrable_domain(x) for x in payload["domains"] if normalize_registrable_domain(x)!="unknown"][:20]
|
|
self.audit(db,user,"domain.availability.checked",str(bid)); db.commit()
|
|
return self.send_json(200,{"business_id":bid,"status":"unknown","reason":"not_configured","provider_configured":False,"items":[{"domain":d,"status":"unknown","reason":"not_configured"} for d in domains]})
|
|
|
|
def list_jobs(self, db, org, query):
|
|
try: limit=max(1,min(int(query.get("page_size",[50])[0]),JOB_PAGE_SIZE)); offset=max(0,int(query.get("offset",[0])[0]))
|
|
except (ValueError, TypeError): return self.send_json(400,{"error":"invalid_pagination"})
|
|
rows=db.execute("SELECT * FROM jobs WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?",(org,limit+1,offset)).fetchall(); more=len(rows)>limit
|
|
return self.send_json(200,{"organization_id":org,"items":[job_json(r) for r in rows[:limit]],"limit":limit,"offset":offset,"has_more":more})
|
|
|
|
def _discovery_run_json(self, row):
|
|
item = row_json(row)
|
|
for field, default in (("criteria_json", {}), ("seed_urls_json", []), ("result_json", {})):
|
|
try: item[field[:-5]] = json.loads(item.pop(field) or json.dumps(default))
|
|
except (TypeError, ValueError): item[field[:-5]] = default
|
|
return item
|
|
|
|
def list_discovery_runs(self, db, org, query):
|
|
try: limit = max(1, min(int((query.get("page_size") or [50])[0]), 100)); offset = max(0, int((query.get("offset") or [0])[0]))
|
|
except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"})
|
|
rows = db.execute("SELECT * FROM discovery_runs WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?", (org, limit + 1, offset)).fetchall()
|
|
return self.send_json(200, {"organization_id": org, "items": [self._discovery_run_json(r) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit})
|
|
|
|
def create_scoped_discovery(self, payload, db, user):
|
|
criteria = payload.get("criteria", {}); seeds = payload.get("seed_urls")
|
|
criteria_only = seeds is None
|
|
if not isinstance(criteria, dict): return self.send_json(400, {"error": "invalid_criteria"})
|
|
if not criteria_only and (not isinstance(seeds, list) or not seeds): return self.send_json(400, {"error": "seed_urls_required"})
|
|
if criteria_only:
|
|
status = search_provider_status()
|
|
if status["status"] != "ready": return self.send_json(503, {"error": status["status"], "provider": status["provider"]})
|
|
try:
|
|
if not criteria_only and len(seeds) > 5: raise ValueError("invalid_criteria")
|
|
if len(json.dumps(criteria).encode()) > 8192: raise ValueError("invalid_criteria")
|
|
if not isinstance(criteria.get("keywords", criteria.get("keyword", [])), (list, str)): raise ValueError("invalid_criteria")
|
|
max_pages = int(payload.get("max_pages", 20)); max_candidates = int(payload.get("max_candidates", 50))
|
|
if not 1 <= max_pages <= 20 or not 1 <= max_candidates <= 50: raise ValueError("invalid_limits")
|
|
except (ValueError, TypeError) as exc: return self.send_json(400, {"error": str(exc) or "invalid_criteria"})
|
|
if not criteria_only:
|
|
try:
|
|
for url in seeds: validate_url(url)
|
|
except (ValueError, TypeError):
|
|
return self.send_json(400, {"error": "unsafe_seed_url"})
|
|
key = str(payload.get("idempotency_key", "")).strip()
|
|
if not key or len(key) > 200: return self.send_json(400, {"error": "invalid_idempotency_key"})
|
|
job_payload = {"criteria": criteria, "seed_urls": seeds if not criteria_only else None, "max_pages": max_pages, "max_candidates": max_candidates}
|
|
result = self.create_job({"type": "scoped_discovery", "payload": job_payload, "idempotency_key": key, "_accepted": True, "_defer_wakeup": True}, db, user)
|
|
# create_job has already committed; read its id from the response is not available,
|
|
# so resolve by the tenant-scoped idempotency key.
|
|
job = db.execute("SELECT * FROM jobs WHERE organization_id=? AND idempotency_key=?", (user["organization_id"], key)).fetchone()
|
|
if not db.execute("SELECT id FROM discovery_runs WHERE organization_id=? AND job_id=?", (user["organization_id"], job["id"])).fetchone():
|
|
db.execute("INSERT INTO discovery_runs(organization_id,job_id,criteria_json,seed_urls_json) VALUES(?,?,?,?)", (user["organization_id"], job["id"], json.dumps(criteria, sort_keys=True), json.dumps(seeds if not criteria_only else [])))
|
|
self.audit(db, user, "discovery.created", str(job["id"])); db.commit()
|
|
getattr(self.server, "job_wakeup", threading.Event()).set()
|
|
return result
|
|
|
|
def get_job_route(self, db, org, path, query):
|
|
bits=path.split("/")
|
|
if len(bits)<5 or not bits[4].isdigit(): return self.send_json(404,{"error":"not_found"})
|
|
job=db.execute("SELECT * FROM jobs WHERE id=? AND organization_id=?",(int(bits[4]),org)).fetchone()
|
|
if not job:return self.send_json(404,{"error":"not_found"})
|
|
if len(bits)==5:return self.send_json(200,job_json(job))
|
|
if len(bits)==6 and bits[5]=="events":
|
|
try: after=max(0,int(query.get("after",[0])[0]))
|
|
except (ValueError,TypeError): return self.send_json(400,{"error":"invalid_sequence"})
|
|
events=[row_json(r) for r in db.execute("SELECT * FROM job_events WHERE job_id=? AND organization_id=? AND sequence>? ORDER BY sequence",(job["id"],org,after))]
|
|
return self.send_json(200,{"items":events,"after":after})
|
|
if len(bits)==7 and bits[5]=="events" and bits[6]=="stream":
|
|
try: after=max(0,int(query.get("after",[0])[0]))
|
|
except (ValueError,TypeError): return self.send_json(400,{"error":"invalid_sequence"})
|
|
events=[row_json(r) for r in db.execute("SELECT * FROM job_events WHERE job_id=? AND organization_id=? AND sequence>? ORDER BY sequence",(job["id"],org,after))]
|
|
body=b"".join((b"event: "+str(e["event_type"]).encode()+b"\\ndata: "+json.dumps(e,sort_keys=True).encode()+b"\\n\\n") for e in events)
|
|
self.send_response(200);self.send_header("Content-Type","text/event-stream");self.send_header("Cache-Control","no-cache");self.send_header("Content-Length",str(len(body)));self.end_headers();self.wfile.write(body);return
|
|
return self.send_json(404,{"error":"not_found"})
|
|
|
|
def create_job(self, payload, db, user):
|
|
kind=str(payload.get("type","")).strip(); key=str(payload.get("idempotency_key","")).strip(); data=payload.get("payload",{})
|
|
if kind not in JOB_TYPES:return self.send_json(400,{"error":"invalid_job_type"})
|
|
if not key or len(key)>200 or not isinstance(data,dict):return self.send_json(400,{"error":"invalid_job_request"})
|
|
safe=json.dumps(redact(data),sort_keys=True,separators=(",",":")); max_attempts=max(1,min(int(payload.get("max_attempts",3)),5)) if str(payload.get("max_attempts",3)).isdigit() else 3
|
|
try:
|
|
cur=db.execute("INSERT INTO jobs(organization_id,idempotency_key,type,payload,max_attempts) VALUES(?,?,?,?,?)",(user["organization_id"],key,kind,safe,max_attempts)); jid=cur.lastrowid
|
|
self.add_job_event(db,jid,user["organization_id"],"queued","Job queued",0); self.audit(db,user,"job.created",str(jid)); db.commit()
|
|
if not payload.get("_defer_wakeup"): getattr(self.server,"job_wakeup",threading.Event()).set()
|
|
return self.send_json(202 if payload.get("_accepted") else 201,job_json(db.execute("SELECT * FROM jobs WHERE id=?",(jid,)).fetchone()))
|
|
except sqlite3.IntegrityError:
|
|
row=db.execute("SELECT * FROM jobs WHERE organization_id=? AND idempotency_key=?",(user["organization_id"],key)).fetchone(); return self.send_json(200,job_json(row))
|
|
|
|
def add_job_event(self, db, jid, org, event_type, message, progress=0, error_code=None):
|
|
seq=db.execute("SELECT COALESCE(MAX(sequence),0)+1 FROM job_events WHERE job_id=?",(jid,)).fetchone()[0]
|
|
db.execute("INSERT INTO job_events(job_id,organization_id,sequence,event_type,message,progress,error_code) VALUES(?,?,?,?,?,?,?)",(jid,org,seq,event_type,str(message)[:500],progress,error_code)); return seq
|
|
|
|
def job_action(self, db, user, path):
|
|
bits=path.split("/"); jid=int(bits[4]) if len(bits)>4 and bits[4].isdigit() else -1; action=bits[5] if len(bits)>5 else ""
|
|
job=db.execute("SELECT * FROM jobs WHERE id=? AND organization_id=?",(jid,user["organization_id"])).fetchone()
|
|
if not job:return self.send_json(404,{"error":"not_found"})
|
|
if action=="cancel":
|
|
if job["status"] in ("queued","running"): db.execute("UPDATE jobs SET status='cancelled',completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));self.add_job_event(db,jid,user["organization_id"],"cancelled","Job cancelled",job["progress"])
|
|
self.audit(db,user,"job.cancelled",str(jid));db.commit();return self.send_json(200,job_json(db.execute("SELECT * FROM jobs WHERE id=?",(jid,)).fetchone()))
|
|
if action=="retry":
|
|
if job["status"]!="failed":return self.send_json(409,{"error":"job_not_failed"})
|
|
db.execute("UPDATE jobs SET status='queued',error_code=NULL,completed_at=NULL,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));self.add_job_event(db,jid,user["organization_id"],"retry","Job retry queued",job["progress"]);self.audit(db,user,"job.retried",str(jid));db.commit();getattr(self.server,"job_wakeup",threading.Event()).set();return self.send_json(200,job_json(db.execute("SELECT * FROM jobs WHERE id=?",(jid,)).fetchone()))
|
|
return self.send_json(404,{"error":"not_found"})
|
|
|
|
def ensure_score_rules(self, db, org):
|
|
for rule in DEFAULT_RULES:
|
|
db.execute("INSERT OR IGNORE INTO score_rules(organization_id,code,name,description,condition_json,points,max_applications,enabled,version) VALUES(?,?,?,?,?,?,?,?,?)", (org, rule["code"], rule["name"], rule["description"], json.dumps(rule["condition_json"], sort_keys=True), rule["points"], rule["max_applications"], rule["enabled"], rule["version"]))
|
|
|
|
def list_score_rules(self, db, org):
|
|
self.ensure_score_rules(db, org); db.commit()
|
|
rows = db.execute("SELECT * FROM score_rules WHERE organization_id=? ORDER BY code,id", (org,)).fetchall()
|
|
items = []
|
|
for row in rows:
|
|
item = row_json(row)
|
|
try: item["condition_json"] = json.loads(item["condition_json"])
|
|
except (TypeError, ValueError): item["condition_json"] = {}
|
|
item["enabled"] = bool(item["enabled"]); items.append(item)
|
|
return self.send_json(200, {"organization_id": org, "items": items})
|
|
|
|
def create_score_rule(self, payload, db, user):
|
|
code = str(payload.get("code", "")).strip(); name = str(payload.get("name", "")).strip(); condition = payload.get("condition_json", payload.get("condition", {}))
|
|
try: points = int(payload.get("points", 0)); maximum = int(payload.get("max_applications", 1)); version = int(payload.get("version", 1))
|
|
except (TypeError, ValueError): return self.send_json(400, {"error": "invalid_rule"})
|
|
if not code or not name or not isinstance(condition, dict) or maximum < 1 or version < 1 or points < -100 or points > 100: return self.send_json(400, {"error": "invalid_rule"})
|
|
try:
|
|
cur = db.execute("INSERT INTO score_rules(organization_id,code,name,description,condition_json,points,max_applications,enabled,version) VALUES(?,?,?,?,?,?,?,?,?)", (user["organization_id"], code, name, str(payload.get("description", "")), json.dumps(condition, sort_keys=True), points, maximum, int(bool(payload.get("enabled", True))), version))
|
|
except sqlite3.IntegrityError: return self.send_json(409, {"error": "duplicate_rule"})
|
|
self.audit(db, user, "score_rule.created", code); db.commit()
|
|
row = db.execute("SELECT * FROM score_rules WHERE id=?", (cur.lastrowid,)).fetchone(); item = row_json(row); item["condition_json"] = condition; item["enabled"] = bool(item["enabled"])
|
|
return self.send_json(201, item)
|
|
|
|
def update_score_rule(self, rid, payload, db, user):
|
|
row = db.execute("SELECT * FROM score_rules WHERE id=? AND organization_id=?", (rid, user["organization_id"])).fetchone()
|
|
if not row: return self.send_json(404, {"error": "not_found"})
|
|
allowed = {"name", "description", "condition_json", "points", "max_applications", "enabled", "version"}; values = {k: payload[k] for k in allowed if k in payload}
|
|
if not values: return self.send_json(400, {"error": "no_changes"})
|
|
if "condition_json" in values and not isinstance(values["condition_json"], dict): return self.send_json(400, {"error": "invalid_rule"})
|
|
if "points" in values:
|
|
try: values["points"] = int(values["points"])
|
|
except (TypeError, ValueError): return self.send_json(400, {"error": "invalid_rule"})
|
|
columns=[]; params=[]
|
|
for key, value in values.items(): columns.append(key + "=?"); params.append(json.dumps(value, sort_keys=True) if key == "condition_json" else (int(bool(value)) if key == "enabled" else value))
|
|
params += [rid, user["organization_id"]]; db.execute("UPDATE score_rules SET " + ",".join(columns) + ",updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?", params); self.audit(db, user, "score_rule.updated", str(rid)); db.commit()
|
|
item = row_json(db.execute("SELECT * FROM score_rules WHERE id=?", (rid,)).fetchone())
|
|
try: item["condition_json"] = json.loads(item["condition_json"])
|
|
except (TypeError, ValueError): item["condition_json"] = {}
|
|
item["enabled"] = bool(item["enabled"]); return self.send_json(200, item)
|
|
|
|
def recalculate_score(self, bid, payload, db, user):
|
|
org = user["organization_id"]; business = self.business(db, bid, org)
|
|
if not business: return self.send_json(404, {"error": "not_found"})
|
|
scans = db.execute("SELECT result_json,classification,scanned_at FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT 1", (bid, org)).fetchone(); website = {"classification": business["website_class"]}
|
|
if scans:
|
|
try: website.update(json.loads(scans["result_json"] or "{}"))
|
|
except (TypeError, ValueError): pass
|
|
website["classification"] = scans["classification"]
|
|
contacts = [dict(r) for r in db.execute("SELECT public_business,suppressed,do_not_contact FROM contact_extractions WHERE business_id=? AND organization_id=?", (bid, org))]
|
|
drow = db.execute("SELECT status,result_json,checked_at FROM domain_checks WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT 1", (bid, org)).fetchone(); domain = dict(drow) if drow else {}
|
|
if drow:
|
|
try: domain.update(json.loads(drow["result_json"] or "{}"))
|
|
except (TypeError, ValueError): pass
|
|
suppressed = is_suppressed(dict(business), [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))])
|
|
signals = signals_for_business(dict(business), website, contacts, domain, suppressed); self.ensure_score_rules(db, org); rules = [dict(r) for r in db.execute("SELECT * FROM score_rules WHERE organization_id=?", (org,))]; result = evaluate_score(signals, rules)
|
|
override_score = payload.get("override_score"); override_eligible = payload.get("override_eligible")
|
|
if override_score is not None or override_eligible is not None:
|
|
reason = str(payload.get("override_reason", "")).strip()
|
|
if not reason: return self.send_json(400, {"error": "override_reason_required"})
|
|
if override_score is not None: result["score"] = max(0, min(100, int(override_score)))
|
|
if override_eligible is not None and not suppressed: result["eligible"] = bool(override_eligible)
|
|
result["priority_band"] = "ineligible" if not result["eligible"] else ("high" if result["score"] >= 70 else "medium" if result["score"] >= 40 else "low")
|
|
db.execute("UPDATE businesses SET score=?,score_version=?,score_factors=?,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?", (result["score"], SCORE_VERSION, json.dumps(result["explanations"], sort_keys=True), bid, org))
|
|
cur=db.execute("INSERT INTO score_history(organization_id,business_id,score,eligible,priority_band,score_version,explanations_json,signals_json,override_score,override_eligible,override_reason,actor_user_id) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)", (org,bid,result["score"],int(result["eligible"]),result["priority_band"],SCORE_VERSION,json.dumps(result["explanations"],sort_keys=True),json.dumps(signals,sort_keys=True),override_score,override_eligible,payload.get("override_reason"),user["id"]))
|
|
self.audit(db,user,"business.score_recalculated",f"{bid}:{result['score']}"); db.commit(); result.update({"business_id": bid, "history_id": cur.lastrowid}); return self.send_json(200, result)
|
|
|
|
def scoring_summary(self, db, org):
|
|
row=db.execute("SELECT COUNT(*) businesses,COALESCE(AVG(score),0) average_score,SUM(CASE WHEN score>=70 THEN 1 ELSE 0 END) high_priority FROM businesses WHERE organization_id=? AND merge_status='active'",(org,)).fetchone()
|
|
bands={r["priority_band"]:r["count"] for r in db.execute("SELECT priority_band,COUNT(*) count FROM score_history WHERE organization_id=? GROUP BY priority_band",(org,))}
|
|
return self.send_json(200,{"organization_id":org,"businesses":row["businesses"],"average_score":round(row["average_score"],2),"high_priority":row["high_priority"] or 0,"bands":bands,"history_count":db.execute("SELECT COUNT(*) FROM score_history WHERE organization_id=?",(org,)).fetchone()[0]})
|
|
|
|
def list_businesses(self,db,org,query):
|
|
def number(name, default=None):
|
|
raw=query.get(name,[None])[0]
|
|
if raw is None:return default
|
|
try:return int(raw)
|
|
except ValueError: raise ValueError
|
|
try:
|
|
page=number("page",1); size=number("page_size",50); score=number("score_min",None); cursor=number("cursor",0)
|
|
except ValueError:return self.send_json(400,{"error":"invalid_pagination"})
|
|
if page<1 or size<1 or size>100 or cursor<0:return self.send_json(400,{"error":"invalid_pagination"})
|
|
params=[org]; where=["b.organization_id=?"]; q=query.get("q",[""])[0].strip(); website_class=query.get("website_class",[""])[0].strip(); stage=query.get("pipeline_stage",[""])[0].strip()
|
|
if score is not None: where.append("b.score>=?"); params.append(score)
|
|
if website_class: where.append("b.website_class=?"); params.append(website_class)
|
|
if q: where.append("(b.name LIKE ? OR b.website_domain LIKE ? OR b.email LIKE ?)"); params += [f"%{q}%"]*3
|
|
if stage: where.append("EXISTS (SELECT 1 FROM pipeline_entries p WHERE p.business_id=b.id AND p.organization_id=b.organization_id AND p.stage=?)"); params.append(stage)
|
|
offset=(number("cursor",0) or 0)+(page-1)*size
|
|
rows=db.execute("SELECT b.* FROM businesses b WHERE "+" AND ".join(where)+" ORDER BY b.score DESC,b.id LIMIT ? OFFSET ?",params+[size+1,offset]).fetchall(); more=len(rows)>size; rows=rows[:size]
|
|
return self.send_json(200,{"organization_id":org,"items":[row_json(r) for r in rows],"page":page,"page_size":size,"next_cursor":str(offset+size) if more else None})
|
|
def bulk_review(self, payload, db, user):
|
|
ids = payload.get("ids", payload.get("business_ids")); action = str(payload.get("action", "")).strip().lower()
|
|
if not isinstance(ids, list) or not ids or len(ids) > 100 or any(not isinstance(i, int) or i < 1 for i in ids) or len(set(ids)) != len(ids): return self.send_json(400, {"error": "invalid_bulk_ids"})
|
|
if action not in {"verify", "reject", "assign"}: return self.send_json(400, {"error": "invalid_bulk_action"})
|
|
if action == "assign":
|
|
assignee = str(payload.get("assigned_to", payload.get("assignee", ""))).strip()
|
|
if not assignee or len(assignee) > 120: return self.send_json(400, {"error": "assignee_required"})
|
|
else: assignee = ""
|
|
org = user["organization_id"]; marks = ",".join("?" for _ in ids)
|
|
rows = db.execute("SELECT id FROM businesses WHERE organization_id=? AND id IN (" + marks + ")", [org] + ids).fetchall()
|
|
if len(rows) != len(ids): return self.send_json(404, {"error": "not_found"})
|
|
try:
|
|
if action == "verify": db.execute("UPDATE businesses SET verified=1,verified_at=CURRENT_TIMESTAMP,review_status='verified',updated_at=CURRENT_TIMESTAMP WHERE organization_id=? AND id IN (" + marks + ")", [org] + ids)
|
|
elif action == "reject": db.execute("UPDATE businesses SET verified=0,review_status='rejected',updated_at=CURRENT_TIMESTAMP WHERE organization_id=? AND id IN (" + marks + ")", [org] + ids)
|
|
else: db.execute("UPDATE businesses SET assigned_to=?,review_status='assigned',updated_at=CURRENT_TIMESTAMP WHERE organization_id=? AND id IN (" + marks + ")", [assignee, org] + ids)
|
|
self.audit(db, user, "businesses.bulk_review", json.dumps({"action": action, "ids": ids, "assigned_to": assignee}, sort_keys=True))
|
|
db.commit()
|
|
except Exception:
|
|
db.rollback(); raise
|
|
return self.send_json(200, {"action": action, "updated": ids, "count": len(ids)})
|
|
|
|
def do_POST(self):
|
|
path=urlparse(self.path).path.rstrip("/")
|
|
if path=="/api/v1/auth/login":return self.login(self.read_json())
|
|
db=self.db()
|
|
try:
|
|
user=self.require_auth(db)
|
|
if not user:return
|
|
if path=="/api/v1/auth/logout":
|
|
c=SimpleCookie();c.load(self.headers.get("Cookie",""));t=c.get("session");
|
|
if t:db.execute("DELETE FROM sessions WHERE token_hash=?",(hashlib.sha256(t.value.encode()).hexdigest(),))
|
|
self.audit(db,user,"logout");db.commit();return self.send_json(200,{"ok":True},{"Set-Cookie":self.auth_cookie("",0)})
|
|
if path=="/api/v1/jobs":
|
|
if user["role"] not in MUTATING_ROLES:return self.send_json(403,{"error":"forbidden"})
|
|
return self.create_job(self.read_json(),db,user)
|
|
if user["role"] not in MUTATING_ROLES:return self.send_json(403,{"error":"forbidden"})
|
|
payload=self.read_json(); org=user["organization_id"]
|
|
if path=="/api/v1/saved-filters": return self.save_filter(payload,db,user)
|
|
if path=="/api/v1/outreach/provider-config": return self.provider_config(db,user,payload)
|
|
if path=="/api/v1/ai/provider-config": return self.ai_provider_config(db,user,payload)
|
|
if path=="/api/v1/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()
|