add gated outreach preparation
This commit is contained in:
@@ -307,6 +307,167 @@ class ApiHandler(BaseHTTPRequestHandler):
|
||||
if suggestions is not None: item["suggestions"] = suggestions
|
||||
return item
|
||||
|
||||
# Phase 14: outreach is a reviewable preparation workflow only. There is
|
||||
# deliberately no provider client in this service and send never performs
|
||||
# network I/O.
|
||||
OUTREACH_PROVIDERS = {"smtp", "sendgrid", "twilio", "whatsapp"}
|
||||
OUTREACH_KINDS = {"email", "phone", "whatsapp"}
|
||||
OUTREACH_DAILY_CAP = 100
|
||||
OUTREACH_BATCH_CAP = 25
|
||||
|
||||
def _provider_json(self, row, org):
|
||||
if not row:
|
||||
return {"organization_id": org, "provider": "", "enabled": False,
|
||||
"policy": {"consent_required": True}, "daily_cap": self.OUTREACH_DAILY_CAP,
|
||||
"batch_cap": self.OUTREACH_BATCH_CAP}
|
||||
item = {"id": row["id"], "organization_id": org, "provider": row["provider"],
|
||||
"enabled": bool(row["enabled"]), "daily_cap": row["daily_cap"],
|
||||
"batch_cap": row["batch_cap"]}
|
||||
try: item["policy"] = json.loads(row["policy_json"] or "{}")
|
||||
except (TypeError, ValueError): item["policy"] = {"consent_required": True}
|
||||
return item
|
||||
|
||||
def provider_config(self, db, user, payload=None):
|
||||
org = user["organization_id"]
|
||||
row = db.execute("SELECT * FROM outreach_provider_configs WHERE organization_id=?", (org,)).fetchone()
|
||||
if payload is None:
|
||||
return self.send_json(200, self._provider_json(row, org))
|
||||
provider = str(payload.get("provider", row["provider"] if row else "")).strip().lower()
|
||||
if provider and provider not in self.OUTREACH_PROVIDERS: return self.send_json(400, {"error": "invalid_provider"})
|
||||
enabled = bool(payload.get("enabled", bool(row["enabled"]) if row else False))
|
||||
policy = payload.get("policy", payload.get("legal_policy", {} if not row else None))
|
||||
if policy is None:
|
||||
try: policy = json.loads(row["policy_json"] or "{}")
|
||||
except (TypeError, ValueError): policy = {"consent_required": True}
|
||||
if not isinstance(policy, dict) or len(policy) > 20: return self.send_json(400, {"error": "invalid_policy"})
|
||||
policy = {str(k)[:80]: v for k, v in policy.items()}
|
||||
policy.setdefault("consent_required", True)
|
||||
try:
|
||||
daily = int(payload.get("daily_cap", row["daily_cap"] if row else self.OUTREACH_DAILY_CAP)); batch = int(payload.get("batch_cap", row["batch_cap"] if row else self.OUTREACH_BATCH_CAP))
|
||||
except (TypeError, ValueError): return self.send_json(400, {"error": "invalid_limits"})
|
||||
if daily < 1 or daily > self.OUTREACH_DAILY_CAP or batch < 1 or batch > self.OUTREACH_BATCH_CAP: return self.send_json(400, {"error": "invalid_limits"})
|
||||
secret = payload.get("secret", None)
|
||||
fingerprint = row["secret_fingerprint"] if row else ""
|
||||
if secret is not None:
|
||||
if not isinstance(secret, str) or not secret or len(secret) > 4096: return self.send_json(400, {"error": "invalid_secret"})
|
||||
fingerprint = hashlib.sha256(secret.encode()).hexdigest()
|
||||
if enabled and (not provider or not fingerprint): return self.send_json(400, {"error": "provider_credentials_required"})
|
||||
if row:
|
||||
db.execute("UPDATE outreach_provider_configs SET provider=?,enabled=?,secret_fingerprint=?,policy_json=?,daily_cap=?,batch_cap=?,updated_at=CURRENT_TIMESTAMP WHERE organization_id=?", (provider, int(enabled), fingerprint, json.dumps(policy, sort_keys=True), daily, batch, org))
|
||||
else:
|
||||
db.execute("INSERT INTO outreach_provider_configs(organization_id,provider,enabled,secret_fingerprint,policy_json,daily_cap,batch_cap) VALUES(?,?,?,?,?,?,?)", (org, provider, int(enabled), fingerprint, json.dumps(policy, sort_keys=True), daily, batch))
|
||||
self.audit(db, user, "outreach.provider_config.updated", provider or "disabled"); db.commit()
|
||||
return self.send_json(200, self._provider_json(db.execute("SELECT * FROM outreach_provider_configs WHERE organization_id=?", (org,)).fetchone(), org))
|
||||
|
||||
def _draft_json(self, row):
|
||||
item = row_json(row)
|
||||
for field, default in (("template_json", {}), ("citations_json", []), ("provenance_json", {})):
|
||||
key = field[:-5]
|
||||
try: item[key] = json.loads(item.pop(field) or json.dumps(default))
|
||||
except (TypeError, ValueError): item[key] = default
|
||||
item["target_verified"] = bool(item.get("target_verified")); item["consent_confirmed"] = bool(item.get("consent_confirmed"))
|
||||
item["target"] = {"kind": item.pop("target_kind"), "value": item.pop("target_value"), "verified": item["target_verified"]}
|
||||
return item
|
||||
|
||||
def _template(self, text, business, evidence):
|
||||
pattern = re.compile(r"{{\s*([^}]+?)\s*}}")
|
||||
used = []
|
||||
def replace(match):
|
||||
token = match.group(1).strip()
|
||||
if token == "business.name": return str(business["name"])
|
||||
m = re.fullmatch(r"evidence\.(\d+)\.claim", token)
|
||||
if m:
|
||||
index = int(m.group(1));
|
||||
if index < 1 or index > len(evidence): raise ValueError("unsupported_template_variable")
|
||||
used.append(evidence[index - 1]["id"]); return str(evidence[index - 1]["claim"])
|
||||
raise ValueError("unsupported_template_variable")
|
||||
return pattern.sub(replace, text), sorted(set(used))
|
||||
|
||||
def _target_exists(self, db, bid, org, kind, value):
|
||||
if kind == "email":
|
||||
return bool(db.execute("SELECT id FROM businesses WHERE id=? AND organization_id=? AND email=?", (bid, org, value)).fetchone() or db.execute("SELECT id FROM contacts WHERE business_id=? AND organization_id=? AND email=?", (bid, org, value)).fetchone() or db.execute("SELECT id FROM contact_extractions WHERE business_id=? AND organization_id=? AND kind='email' AND value=?", (bid, org, value)).fetchone())
|
||||
return bool(db.execute("SELECT id FROM businesses WHERE id=? AND organization_id=? AND phone=?", (bid, org, value)).fetchone() or db.execute("SELECT id FROM contacts WHERE business_id=? AND organization_id=? AND phone=?", (bid, org, value)).fetchone() or db.execute("SELECT id FROM contact_extractions WHERE business_id=? AND organization_id=? AND kind IN ('phone','whatsapp') AND value=?", (bid, org, value)).fetchone())
|
||||
|
||||
def create_outreach_draft(self, bid, payload, db, user):
|
||||
org = user["organization_id"]; business = self.business(db, bid, org)
|
||||
if not business: return self.send_json(404, {"error": "not_found"})
|
||||
target = payload.get("target", {}); kind = str(target.get("kind", "email")).lower(); value = str(target.get("value", "")).strip().lower()
|
||||
if kind not in self.OUTREACH_KINDS or not value or len(value) > 320 or not isinstance(target.get("verified", False), bool): return self.send_json(400, {"error": "invalid_target"})
|
||||
subject, body = str(payload.get("subject", "")).strip(), str(payload.get("body", "")).strip()
|
||||
if not subject or not body or len(subject) > 500 or len(body) > 10000: return self.send_json(400, {"error": "invalid_draft"})
|
||||
if contains_secret(payload): return self.send_json(400, {"error": "secret_not_permitted"})
|
||||
try:
|
||||
evidence = [dict(r) for r in db.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id", (bid, org))]
|
||||
subject, sids = self._template(subject, business, evidence); body, bids = self._template(body, business, evidence); ids = sorted(set(sids + bids))
|
||||
except ValueError as exc: return self.send_json(400, {"error": str(exc)})
|
||||
key = str(payload.get("idempotency_key", "")).strip() or hashlib.sha256(json.dumps({"business_id": bid, "target": target, "subject": subject, "body": body}, sort_keys=True).encode()).hexdigest()
|
||||
if len(key) > 200: return self.send_json(400, {"error": "invalid_idempotency_key"})
|
||||
existing = db.execute("SELECT * FROM outreach_drafts WHERE organization_id=? AND idempotency_key=?", (org, key)).fetchone()
|
||||
if existing: return self.send_json(200, self._draft_json(existing))
|
||||
if db.execute("SELECT COUNT(*) FROM outreach_drafts WHERE organization_id=? AND created_at>=datetime('now','-1 day')", (org,)).fetchone()[0] >= self.OUTREACH_DAILY_CAP: return self.send_json(429, {"error": "outreach_daily_cap"})
|
||||
cur = db.execute("INSERT INTO outreach_drafts(organization_id,business_id,target_kind,target_value,target_verified,subject,body,template_json,citations_json,provenance_json,legal_basis,consent_confirmed,actor_user_id,idempotency_key) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (org, bid, kind, value, int(target["verified"]), subject, body, json.dumps({"variables": ids}, sort_keys=True), json.dumps(ids), json.dumps({str(x): {"type": "evidence", "id": x} for x in ids}, sort_keys=True), str(payload.get("legal_basis", ""))[:100], int(bool(payload.get("consent_confirmed", False))), user["id"], key))
|
||||
self.audit(db, user, "outreach_draft.created", str(cur.lastrowid)); db.commit()
|
||||
return self.send_json(201, self._draft_json(db.execute("SELECT * FROM outreach_drafts WHERE id=?", (cur.lastrowid,)).fetchone()))
|
||||
|
||||
def list_outreach_drafts(self, db, user, query):
|
||||
try: limit = int((query.get("page_size") or [self.OUTREACH_BATCH_CAP])[0]); offset = int((query.get("offset") or [0])[0])
|
||||
except (ValueError, TypeError): return self.send_json(400, {"error": "invalid_pagination"})
|
||||
if limit < 1 or limit > self.OUTREACH_BATCH_CAP or offset < 0: return self.send_json(400, {"error": "invalid_pagination"})
|
||||
rows = db.execute("SELECT * FROM outreach_drafts WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?", (user["organization_id"], limit + 1, offset)).fetchall()
|
||||
return self.send_json(200, {"organization_id": user["organization_id"], "items": [self._draft_json(r) for r in rows[:limit]], "limit": limit, "offset": offset, "has_more": len(rows) > limit})
|
||||
|
||||
def update_outreach_draft(self, did, payload, db, user):
|
||||
row = db.execute("SELECT * FROM outreach_drafts WHERE id=? AND organization_id=?", (did, user["organization_id"])).fetchone()
|
||||
if not row: return self.send_json(404, {"error": "not_found"})
|
||||
if row["status"] in ("approved", "sent"): return self.send_json(409, {"error": "draft_locked"})
|
||||
fields, values = [], []
|
||||
subject, body = row["subject"], row["body"]
|
||||
if "subject" in payload: subject = str(payload["subject"]).strip()
|
||||
if "body" in payload: body = str(payload["body"]).strip()
|
||||
if ("subject" in payload and (not subject or len(subject) > 500)) or ("body" in payload and (not body or len(body) > 10000)): return self.send_json(400, {"error": "invalid_draft"})
|
||||
if "subject" in payload or "body" in payload:
|
||||
business = self.business(db, row["business_id"], user["organization_id"])
|
||||
evidence = [dict(r) for r in db.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id", (row["business_id"], user["organization_id"]))]
|
||||
try:
|
||||
subject, sids = self._template(subject, business, evidence); body, bids = self._template(body, business, evidence)
|
||||
except ValueError as exc: return self.send_json(400, {"error": str(exc)})
|
||||
ids = sorted(set(sids + bids))
|
||||
fields += ["subject=?", "body=?", "template_json=?", "citations_json=?", "provenance_json=?"]
|
||||
values += [subject, body, json.dumps({"variables": ids}, sort_keys=True), json.dumps(ids), json.dumps({str(x): {"type": "evidence", "id": x} for x in ids}, sort_keys=True)]
|
||||
if not fields: return self.send_json(400, {"error": "no_changes"})
|
||||
fields.append("updated_at=CURRENT_TIMESTAMP"); values += [did, user["organization_id"]]
|
||||
db.execute("UPDATE outreach_drafts SET " + ",".join(fields) + " WHERE id=? AND organization_id=?", values); self.audit(db, user, "outreach_draft.updated", str(did)); db.commit()
|
||||
return self.send_json(200, self._draft_json(db.execute("SELECT * FROM outreach_drafts WHERE id=?", (did,)).fetchone()))
|
||||
|
||||
def approve_outreach_draft(self, did, db, user):
|
||||
row = db.execute("SELECT * FROM outreach_drafts WHERE id=? AND organization_id=?", (did, user["organization_id"])).fetchone()
|
||||
if not row: return self.send_json(404, {"error": "not_found"})
|
||||
if row["status"] != "pending_review": return self.send_json(409, {"error": "draft_not_reviewable"})
|
||||
db.execute("UPDATE outreach_drafts SET status='approved',approved_by=?,approved_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?", (user["id"], did, user["organization_id"])); self.audit(db, user, "outreach_draft.approved", str(did)); db.commit()
|
||||
return self.send_json(200, self._draft_json(db.execute("SELECT * FROM outreach_drafts WHERE id=?", (did,)).fetchone()))
|
||||
|
||||
def send_outreach_draft(self, did, db, user):
|
||||
org = user["organization_id"]; row = db.execute("SELECT * FROM outreach_drafts WHERE id=? AND organization_id=?", (did, org)).fetchone()
|
||||
if not row: return self.send_json(404, {"error": "not_found"})
|
||||
cfg = db.execute("SELECT * FROM outreach_provider_configs WHERE organization_id=?", (org,)).fetchone(); reasons = []
|
||||
if not cfg or not cfg["enabled"] or not cfg["provider"] or not cfg["secret_fingerprint"]: reasons.append("provider")
|
||||
business = self.business(db, row["business_id"], org)
|
||||
if not row["target_verified"] or not self._target_exists(db, row["business_id"], org, row["target_kind"], row["target_value"]): reasons.append("verified_target")
|
||||
suppressed = is_suppressed({"email": row["target_value"] if row["target_kind"] == "email" else "", "phone": row["target_value"] if row["target_kind"] != "email" else "", "website_domain": business["website_domain"] if business else ""}, [dict(r) for r in db.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1", (org,))])
|
||||
if suppressed: reasons.append("suppressed")
|
||||
policy = {}
|
||||
if cfg:
|
||||
try: policy = json.loads(cfg["policy_json"] or "{}")
|
||||
except (TypeError, ValueError): policy = {"consent_required": True}
|
||||
if policy.get("consent_required", True) and (not row["consent_confirmed"] or not row["legal_basis"]): reasons.append("consent_or_legal_policy")
|
||||
if row["status"] != "approved": reasons.append("approved_draft")
|
||||
if cfg and db.execute("SELECT COUNT(*) FROM outreach_drafts WHERE organization_id=? AND status='sent' AND sent_at>=datetime('now','-1 day')", (org,)).fetchone()[0] >= cfg["daily_cap"]: reasons.append("daily_cap")
|
||||
if reasons:
|
||||
status = "not_configured" if "provider" in reasons and (not cfg or not cfg["enabled"] or not cfg["provider"] or not cfg["secret_fingerprint"]) else "blocked"; self.audit(db, user, "outreach_draft.send_blocked", f"{did}:{','.join(reasons)}"); db.commit()
|
||||
return self.send_json(409, {"status": status, "blocked_reasons": reasons, "network_send": False, "id": did})
|
||||
reasons.append("network_send_disabled"); self.audit(db, user, "outreach_draft.send_blocked", f"{did}:network_send_disabled"); db.commit()
|
||||
return self.send_json(409, {"status": "blocked", "blocked_reasons": reasons, "network_send": False, "id": did})
|
||||
|
||||
def suggest_ai(self, bid, payload, db, user):
|
||||
org = user["organization_id"]
|
||||
business = self.business(db, bid, org)
|
||||
@@ -394,6 +555,8 @@ class ApiHandler(BaseHTTPRequestHandler):
|
||||
if path=="/api/v1/interactions": return self.list_interactions(db,org,parse_qs(parsed.query))
|
||||
if path=="/api/v1/suppressions": return self.list_suppressions(db,org,parse_qs(parsed.query))
|
||||
if path=="/api/v1/ai-runs": return self.list_ai_runs(db,user,parse_qs(parsed.query))
|
||||
if path=="/api/v1/outreach/drafts": return self.list_outreach_drafts(db,user,parse_qs(parsed.query))
|
||||
if path=="/api/v1/outreach/provider-config": return self.provider_config(db,user)
|
||||
if path in ("/api/v1/reports/pipeline","/api/v1/reports/outcomes","/api/v1/reports/activity"): return self.report(db,org,path.rsplit('/',1)[1],parse_qs(parsed.query))
|
||||
if path.startswith("/api/v1/jobs/"): return self.get_job_route(db,org,path,parse_qs(parsed.query))
|
||||
if path.startswith("/api/v1/businesses/"):
|
||||
@@ -753,6 +916,7 @@ class ApiHandler(BaseHTTPRequestHandler):
|
||||
if user["role"] not in MUTATING_ROLES:return self.send_json(403,{"error":"forbidden"})
|
||||
payload=self.read_json(); org=user["organization_id"]
|
||||
if path=="/api/v1/saved-filters": return self.save_filter(payload,db,user)
|
||||
if path=="/api/v1/outreach/provider-config": return self.provider_config(db,user,payload)
|
||||
if path=="/api/v1/businesses/bulk-review": return self.bulk_review(payload,db,user)
|
||||
bits_ai=path.split("/")
|
||||
if len(bits_ai)==7 and bits_ai[:4]==["","api","v1","businesses"] and bits_ai[5]=="ai" and bits_ai[6]=="suggest": return self.suggest_ai(int(bits_ai[4]) if bits_ai[4].isdigit() else -1,payload,db,user)
|
||||
@@ -762,6 +926,9 @@ class ApiHandler(BaseHTTPRequestHandler):
|
||||
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-queries":return self.create_query(payload,db,user)
|
||||
if path=="/api/v1/suppressions":return self.create_suppression(payload,db,user)
|
||||
@@ -802,7 +969,9 @@ class ApiHandler(BaseHTTPRequestHandler):
|
||||
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()
|
||||
|
||||
Reference in New Issue
Block a user