forcol,definitionin(("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 '{}'")):
# Phase 12 is additive-safe for databases created before CRM metadata existed.
fortable,additionsin{
"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")),
"sources":(("source_code","TEXT NOT NULL DEFAULT ''"),("display_name","TEXT NOT NULL DEFAULT ''"),("approved","INTEGER NOT NULL DEFAULT 0"),("policy_json","TEXT NOT NULL DEFAULT '{}'"),("quota_json","TEXT NOT NULL DEFAULT '{}'")),
"discovery_queries":(("selected_adapters_json","TEXT NOT NULL DEFAULT '[]'"),("location","TEXT NOT NULL DEFAULT ''"),("category","TEXT NOT NULL DEFAULT ''"),("max_records","INTEGER NOT NULL DEFAULT 100"),("daily_limit","INTEGER NOT NULL DEFAULT 1000"),("schedule","TEXT NOT NULL DEFAULT ''"),("dry_run","INTEGER NOT NULL DEFAULT 0"),("lifecycle","TEXT NOT NULL DEFAULT 'draft'")),
"discovery_runs":(("selected_adapters_json","TEXT NOT NULL DEFAULT '[]'"),("location","TEXT NOT NULL DEFAULT ''"),("category","TEXT NOT NULL DEFAULT ''"),("max_records","INTEGER NOT NULL DEFAULT 100"),("daily_limit","INTEGER NOT NULL DEFAULT 1000"),("schedule","TEXT NOT NULL DEFAULT ''"),("dry_run","INTEGER NOT NULL DEFAULT 0"),("lifecycle","TEXT NOT NULL DEFAULT 'draft'"),("paused_at","TEXT")),
"source_records":(("discovery_run_id","INTEGER"),("normalized_key","TEXT NOT NULL DEFAULT ''"),("provenance_json","TEXT NOT NULL DEFAULT '{}'"),("response_metadata_json","TEXT NOT NULL DEFAULT '{}'")),
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 <> ''")
fororganizationindb.execute("SELECT id FROM organizations").fetchall():
forruleinDEFAULT_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"]))
returndb.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()
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()
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=(",",":"))))
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()
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)ifclauseelse"")+" GROUP BY p.stage,p.status ORDER BY p.stage,p.status"
elifkind=="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)ifclauseelse"")+" 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)ifclauseelse"")+" GROUP BY activity_type"
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))
returnbool(db.execute("SELECT id FROM businesses WHERE id=? AND organization_id=? AND email=?",(bid,org,value)).fetchone()ordb.execute("SELECT id FROM contacts WHERE business_id=? AND organization_id=? AND email=?",(bid,org,value)).fetchone()ordb.execute("SELECT id FROM contact_extractions WHERE business_id=? AND organization_id=? AND kind='email' AND value=?",(bid,org,value)).fetchone())
returnbool(db.execute("SELECT id FROM businesses WHERE id=? AND organization_id=? AND phone=?",(bid,org,value)).fetchone()ordb.execute("SELECT id FROM contacts WHERE business_id=? AND organization_id=? AND phone=?",(bid,org,value)).fetchone()ordb.execute("SELECT id FROM contact_extractions WHERE business_id=? AND organization_id=? AND kind IN ('phone','whatsapp') AND value=?",(bid,org,value)).fetchone())
evidence=[dict(r)forrindb.execute("SELECT id,kind,url,claim,created_at FROM evidence WHERE business_id=? AND organization_id=? ORDER BY id",(bid,org))]
ifdb.execute("SELECT COUNT(*) FROM outreach_drafts WHERE organization_id=? AND created_at>=datetime('now','-1 day')",(org,)).fetchone()[0]>=self.OUTREACH_DAILY_CAP:returnself.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}forxinids},sort_keys=True),str(payload.get("legal_basis",""))[:100],int(bool(payload.get("consent_confirmed",False))),user["id"],key))
rows=db.execute("SELECT * FROM outreach_drafts WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?",(user["organization_id"],limit+1,offset)).fetchall()
evidence=[dict(r)forrindb.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"]))]
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()
returnself.send_json(200,self._draft_json(db.execute("SELECT * FROM outreach_drafts WHERE id=?",(did,)).fetchone()))
defapprove_outreach_draft(self,did,db,user):
row=db.execute("SELECT * FROM outreach_drafts WHERE id=? AND organization_id=?",(did,user["organization_id"])).fetchone()
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()
returnself.send_json(200,self._draft_json(db.execute("SELECT * FROM outreach_drafts WHERE id=?",(did,)).fetchone()))
defsend_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()
suppressed=is_suppressed({"email":row["target_value"]ifrow["target_kind"]=="email"else"","phone":row["target_value"]ifrow["target_kind"]!="email"else"","website_domain":business["website_domain"]ifbusinesselse""},[dict(r)forrindb.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1",(org,))])
ifcfganddb.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")
db.execute("INSERT INTO ai_remote_provider_configs(organization_id,provider,model,enabled,nous_base_url,firecrawl_base_url,credentials_ciphertext,credentials_fingerprint) VALUES(?,?,?,?,?,?,?,?) ON CONFLICT(organization_id) DO UPDATE SET provider=excluded.provider,model=excluded.model,enabled=excluded.enabled,nous_base_url=excluded.nous_base_url,firecrawl_base_url=excluded.firecrawl_base_url,credentials_ciphertext=excluded.credentials_ciphertext,credentials_fingerprint=excluded.credentials_fingerprint,updated_at=CURRENT_TIMESTAMP",(org,config["provider"],config["model"],int(config["enabled"]),config["nous_base_url"],config["firecrawl_base_url"],ciphertext,fingerprint))
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)))
scans=[dict(x)forxindb.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)forxindb.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)forxindb.execute("SELECT * FROM contacts WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?",(bid,org,limit))]
extracted=[dict(x)forxindb.execute("SELECT * FROM contact_extractions WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT ?",(bid,org,limit))]
scans=[dict(r)forrindb.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)forrindb.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)forrindb.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()]
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"ifproviderelse"",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
forsuggestioninoutput.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)))
rows=db.execute("SELECT * FROM ai_runs WHERE organization_id=? ORDER BY id DESC LIMIT ? OFFSET ?",(user["organization_id"],limit+1,offset)).fetchall()
items=[]
forrowinrows[:limit]:
suggestions=[json.loads(x[0])forxindb.execute("SELECT output_json FROM ai_suggestions WHERE ai_run_id=? AND organization_id=? ORDER BY id",(row["id"],user["organization_id"]))]
ifdecision=="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"]))
returnself.send_json(200,{"items":[dict(r)forrindb.execute("SELECT id,email,role,organization_id,created_at FROM users WHERE organization_id=? ORDER BY id",(org,))]})
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]}
returnself.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}}forkey,valueincounts.items()]})
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()
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"]))
exceptsqlite3.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()
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"]))
rows=db.execute("SELECT * FROM contact_extractions WHERE "+" AND ".join(where)+" ORDER BY id DESC LIMIT ? OFFSET ?",params+[limit+1,offset]).fetchall()
existing=db.execute("SELECT * FROM contact_extractions WHERE organization_id=? AND business_id=? AND extraction_key=? ORDER BY id",(org,bid,key)).fetchall()
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"]ifscanelseNone,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"]))
rows=db.execute("SELECT * FROM contact_extractions WHERE organization_id=? AND business_id=? AND extraction_key=? ORDER BY id",(org,bid,key)).fetchall()
row=db.execute("SELECT * FROM website_scans WHERE business_id=? AND organization_id=? ORDER BY id DESC LIMIT 1",(bid,user["organization_id"])).fetchone()
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()
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"]ifwebsiteelseNone,safe_url,result["classification"],json.dumps(result,sort_keys=True),cache_key,now.isoformat(),expires.isoformat()))
candidates=db.execute("SELECT domain FROM domain_candidates WHERE organization_id=? AND business_id=? ORDER BY rank,id",(user["organization_id"],bid)).fetchall()
db.execute("INSERT INTO discovery_runs(organization_id,job_id,selected_adapters_json,location,category,max_records,daily_limit,schedule,dry_run,lifecycle,criteria_json,seed_urls_json) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)",(user["organization_id"],job["id"],json.dumps(selected_adapters),str(payload.get("location","")),str(payload.get("category","")),int(payload.get("max_records",100)),int(payload.get("daily_limit",1000)),str(payload.get("schedule","")),int(bool(payload.get("dry_run",False))),"queued",json.dumps(criteria,sort_keys=True),json.dumps(seedsifnotcriteria_onlyelse[])))
events=[row_json(r)forrindb.execute("SELECT * FROM job_events WHERE job_id=? AND organization_id=? AND sequence>? ORDER BY sequence",(job["id"],org,after))]
events=[row_json(r)forrindb.execute("SELECT * FROM job_events WHERE job_id=? AND organization_id=? AND sequence>? ORDER BY sequence",(job["id"],org,after))]
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
row=db.execute("SELECT * FROM jobs WHERE organization_id=? AND idempotency_key=?",(user["organization_id"],key)).fetchone();returnself.send_json(200,job_json(row))
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));returnseq
ifjob["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();returnself.send_json(200,job_json(db.execute("SELECT * FROM jobs WHERE id=?",(jid,)).fetchone()))
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();returnself.send_json(200,job_json(db.execute("SELECT * FROM jobs WHERE id=?",(jid,)).fetchone()))
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"]))
deflist_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()
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))
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"])
returnself.send_json(201,item)
defupdate_score_rule(self,rid,payload,db,user):
row=db.execute("SELECT * FROM score_rules WHERE id=? AND organization_id=?",(rid,user["organization_id"])).fetchone()
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())
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"]}
contacts=[dict(r)forrindb.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)ifdrowelse{}
signals=signals_for_business(dict(business),website,contacts,domain,suppressed);self.ensure_score_rules(db,org);rules=[dict(r)forrindb.execute("SELECT * FROM score_rules WHERE organization_id=?",(org,))];result=evaluate_score(signals,rules)
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"]))
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"]forrindb.execute("SELECT priority_band,COUNT(*) count FROM score_history WHERE organization_id=? GROUP BY priority_band",(org,))}
returnself.send_json(200,{"organization_id":org,"businesses":row["businesses"],"average_score":round(row["average_score"],2),"high_priority":row["high_priority"]or0,"bands":bands,"history_count":db.execute("SELECT COUNT(*) FROM score_history WHERE organization_id=?",(org,)).fetchone()[0]})
ifq:where.append("(b.name LIKE ? OR b.website_domain LIKE ? OR b.email LIKE ?)");params+=[f"%{q}%"]*3
ifstage: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)or0)+(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]
ifaction=="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)
elifaction=="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)
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();returnself.send_json(200,{"ok":True,"id":int(bits[4])})
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()
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();returnself.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()))})
b=normalize_business(payload);suppressions=[dict(r)forrindb.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1",(org,))]
iffieldsanddb.execute("SELECT id FROM businesses WHERE organization_id=? AND ("+" OR ".join(f"{c}=?"forc,_infields)+")",[org]+[vfor_,vinfields]).fetchone():returnself.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();returnself.send_json(201,row_json(db.execute("SELECT * FROM businesses WHERE id=?",(cur.lastrowid,)).fetchone()))
try:db.execute("INSERT INTO suppressions(organization_id,kind,value,actor_user_id,active) VALUES(?,?,?,?,1)",(user["organization_id"],kind,value,user["id"]))
self.audit(db,user,"suppression.created",kind);db.commit();returnself.send_json(201,row_json(db.execute("SELECT * FROM suppressions WHERE organization_id=? AND kind=? AND value=?",(user["organization_id"],kind,value)).fetchone()))
suppressed=is_suppressed({"email":email,"phone":phone},[dict(r)forrindb.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")))orint(suppressed))
columns=CHILD_TABLES[table];db.execute(f"INSERT INTO {table}(business_id,organization_id,{','.join(columns)}) VALUES(?, ?, {','.join('?'for_incolumns)})",(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();returnself.send_json(201,row_json(db.execute(f"SELECT * FROM {table} WHERE id=?",(rid,)).fetchone()))
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"],keyorNone));rid=cur.lastrowid
self.audit(db,user,"pipeline.created",str(rid));db.commit();returnself.send_json(200ifgetattr(self,"command","")=="PATCH"else201,row_json(db.execute("SELECT * FROM pipeline_entries WHERE id=?",(rid,)).fetchone()))
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()
returnself.send_json(200,row_json(db.execute("SELECT * FROM pipeline_entries WHERE id=?",(eid,)).fetchone()))
prior=db.execute("SELECT * FROM interactions WHERE organization_id=? AND idempotency_key=?",(org,key)).fetchone()
ifprior:returnself.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"],keyorNone));self.audit(db,user,"interaction.created",str(cur.lastrowid));db.commit()
returnself.send_json(201,row_json(db.execute("SELECT * FROM interactions WHERE id=?",(cur.lastrowid,)).fetchone()))
defupdate_interaction(self,iid,payload,db,user):
org=user["organization_id"];row=db.execute("SELECT * FROM interactions WHERE id=? AND organization_id=?",(iid,org)).fetchone()
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();returnself.send_json(200,row_json(db.execute("SELECT * FROM interactions WHERE id=?",(iid,)).fetchone()))
defupdate_suppression(self,sid,payload,db,user):
row=db.execute("SELECT * FROM suppressions WHERE id=? AND organization_id=?",(sid,user["organization_id"])).fetchone()
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();returnself.send_json(200,row_json(db.execute("SELECT * FROM suppressions WHERE id=?",(sid,)).fetchone()))
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,"pipeline_stage.created",str(cur.lastrowid));db.commit();returnself.send_json(201,row_json(db.execute("SELECT * FROM pipeline_stages WHERE id=?",(cur.lastrowid,)).fetchone()))
defupdate_stage(self,sid,payload,db,user):
row=db.execute("SELECT * FROM pipeline_stages WHERE id=? AND organization_id=?",(sid,user["organization_id"])).fetchone()
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();returnself.send_json(200,row_json(db.execute("SELECT * FROM pipeline_stages WHERE id=?",(sid,)).fetchone()))
defdelete_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()
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();returnself.send_json(200,{"ok":True,"id":ident})
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),nowifverifiedelseNone,"verified"ifverifiedelse"pending",bid,user["organization_id"]));self.audit(db,user,"business.verified",str(verified));db.commit();row=self.business(db,bid,user["organization_id"]);returnself.send_json(200,row_json(row))
normalized=deduplicate_businesses([rforrinrowsifisinstance(r,dict)andstr(r.get("name","")).strip()]);suppressions=[dict(r)forrindb.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1",(org,))];existing=[row_json(r)forrindb.execute("SELECT * FROM businesses WHERE organization_id=?",(org,))];seen=set();accepted=[];suppressed=0;existing_keys={deduplication_key(x)forxinexisting}
returnself.send_json(200,{"organization_id":org,"items":[row_json(r)forrindb.execute(f"SELECT {cols} FROM sources WHERE organization_id=? ORDER BY id",(org,))]})
deflist_queries(self,db,org):
returnself.send_json(200,{"organization_id":org,"items":[row_json(r)forrindb.execute("SELECT * FROM discovery_queries WHERE organization_id=? ORDER BY id",(org,))]})
cur=db.execute("INSERT INTO sources(organization_id,name,kind,source_code,display_name,enabled,approved,config_json,policy_json,quota_json) VALUES(?,?,?,?,?,?,?,?,?,?)",(user['organization_id'],name,kind,kind,str(payload.get('display_name')oradapter_for(kind).display_name),int(bool(payload.get('enabled',False))),int(bool(payload.get('approved',config.get('approved',False)))),json.dumps(config,sort_keys=True),json.dumps(payload.get('policy',{}),sort_keys=True),json.dumps(payload.get('quota',{}),sort_keys=True)))
self.audit(db,user,'source.created',str(cur.lastrowid));db.commit();returnself.send_json(201,row_json(db.execute("SELECT id,organization_id,name,kind,source_code,display_name,enabled,approved,policy_json,quota_json,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()))
ifnotdb.execute("SELECT id FROM sources WHERE id=? AND organization_id=?",(sid,user['organization_id'])).fetchone():returnself.send_json(404,{"error":"not_found"})
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'ifvalueelse'source.disabled',str(sid));db.commit();returnself.send_json(200,row_json(db.execute("SELECT * FROM sources WHERE id=?",(sid,)).fetchone()))
ifnotdb.execute("SELECT id FROM sources WHERE id=? AND organization_id=?",(sid,user['organization_id'])).fetchone():returnself.send_json(404,{"error":"not_found"})
try:cur=db.execute("INSERT INTO discovery_queries(organization_id,source_id,name,query_json,selected_adapters_json,location,category,max_records,daily_limit,schedule,dry_run,lifecycle) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)",(user['organization_id'],sid,name,json.dumps(query,sort_keys=True),json.dumps(selected),location,category,max_records,daily_limit,schedule,int(bool(payload.get('dry_run',False))),str(payload.get('lifecycle','draft'))))
self.audit(db,user,'discovery_query.created',str(cur.lastrowid));db.commit();returnself.send_json(201,row_json(db.execute("SELECT * FROM discovery_queries WHERE id=?",(cur.lastrowid,)).fetchone()))
job=db.execute("SELECT * FROM jobs WHERE organization_id=? AND idempotency_key=?",(user["organization_id"],key)).fetchone()
ifnotdb.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,selected_adapters_json,location,category,max_records,daily_limit,schedule,dry_run,lifecycle,criteria_json,seed_urls_json) VALUES(?,?,?,?,?,?,?,?,?,?,?,?)",(user["organization_id"],job["id"],query["selected_adapters_json"],query["location"],query["category"],query["max_records"],query["daily_limit"],query["schedule"],query["dry_run"],"queued",query["query_json"],"[]"));db.commit()
ifok: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'
db.execute("INSERT INTO source_records(organization_id,source_id,content_hash,raw_json,normalized_json,normalized_key,source_url,provenance_json,query_context_json,cursor_json,rate_policy_json) VALUES(?,?,?,?,?,?,?,?,?,?,?)",(user['organization_id'],sid,digest,raw,json.dumps(normalize_record(record),sort_keys=True),normalized_key,str(payload.get('source_url','')),json.dumps({'adapter':source['kind']},sort_keys=True),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)))
record_id=db.execute("SELECT last_insert_rowid()").fetchone()[0];db.execute("INSERT OR IGNORE INTO enrichment_queue(organization_id,source_record_id) VALUES(?,?)",(user['organization_id'],record_id));inserted+=1
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();returnself.send_json(201ifinsertedelse200,{"inserted":inserted,"records":len(page.records)})
db.execute("UPDATE discovery_runs SET lifecycle=?,paused_at=CASE WHEN ?='paused' THEN CURRENT_TIMESTAMP ELSE paused_at END,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=?",(lifecycle,lifecycle,rid,user["organization_id"]))
ifaction=="cancel":db.execute("UPDATE jobs SET status='cancelled',completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=? AND organization_id=? AND status IN ('queued','running')",(run["job_id"],user["organization_id"]))
ifrows: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))
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"]))
suppressed=is_suppressed(b,[dict(r)forrindb.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1",(org,))])
ifsuppressedornotb["website_domain"]:continue
existing=db.execute("SELECT id FROM businesses WHERE organization_id=? AND website_domain=?",(org,b["website_domain"])).fetchone()
ifexisting: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"]))
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=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()
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
forpageincandidate["evidence"]:
db.execute("INSERT INTO evidence(business_id,organization_id,kind,url,claim) VALUES(?,?,?,?,?)",(bid,org,page["kind"],page["url"],page["claim"]))
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"]))
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"]))
linked=db.execute("SELECT kind FROM sources WHERE id=? AND organization_id=?",(query["source_id"],org)).fetchone()
selected=[linked["kind"]]iflinkedelse[]
placeholders=','.join('?'*len(selected))or"NULL"
sources=db.execute("SELECT * FROM sources WHERE organization_id=? AND enabled=1 AND circuit_open=0 AND (source_code IN ("+placeholders+") OR kind IN ("+placeholders+"))",[org]+list(selected)+list(selected)).fetchall()ifselectedelse[]
used_today=db.execute("SELECT COUNT(*) FROM source_records WHERE organization_id=? AND source_id=? AND date(created_at)=date('now')",(org,source["id"])).fetchone()[0]
ifused_today>=daily_limitorper_run_limit<=0:
handler.add_job_event(db,job["id"],org,"source.quota_exceeded",f"Quota reached for {source['kind']}",10,"SOURCE_QUOTA_EXCEEDED");continue
db.execute("UPDATE sources SET health_status='unhealthy',consecutive_failures=?,circuit_open=?,last_failure_at=CURRENT_TIMESTAMP,last_error=?,updated_at=CURRENT_TIMESTAMP WHERE id=?",(failures,int(circuit_is_open(failures)),str(exc)[:300],source["id"]))
cur=db.execute("INSERT INTO source_records(organization_id,source_id,discovery_query_id,discovery_run_id,content_hash,raw_json,normalized_json,normalized_key,source_url,provenance_json,query_context_json,response_metadata_json,processing_status) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?)",(org,source["id"],query["id"]ifqueryelseNone,run["id"]ifrunelseNone,digest,raw,norm,nkey,str(config.get("source_url","")),json.dumps({"adapter":source["kind"],"metadata":page.metadata},sort_keys=True),json.dumps(payload.get("criteria",{}),sort_keys=True),json.dumps(page.metadata,sort_keys=True),"raw"))
record_id=cur.lastrowid;total+=1
exceptsqlite3.IntegrityError:
existing=db.execute("SELECT id,processing_status FROM source_records WHERE organization_id=? AND source_id=? AND content_hash=?",(org,source["id"],digest)).fetchone()
suppressed=is_suppressed(normalized,[dict(r)forrindb.execute("SELECT kind,value FROM suppressions WHERE organization_id=? AND active=1",(org,))])
ifsuppressed:
ifrecord_id:db.execute("UPDATE source_records SET processing_status='skipped' WHERE id=? AND organization_id=?",(record_id,org))
handler.add_job_event(db,job["id"],org,"review.skipped",f"Suppressed source record {record_id}",40,"SUPPRESSED");continue
existing=db.execute("SELECT * FROM businesses WHERE organization_id=? AND ((website_domain<>'' AND website_domain=?) OR (email<>'' AND email=?) OR (phone<>'' AND phone=?) OR (name=? AND city=?)) ORDER BY id LIMIT 1",(org,normalized["website_domain"],normalized["email"],normalized["phone"],normalized["name"],normalized["city"])).fetchone()
ifexisting:bid=existing["id"];handler.add_job_event(db,job["id"],org,"business.matched",f"Matched business {bid}",50)
else:
scored=score_business(normalized);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,review_status) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)",(org,normalized["name"],normalized["website"],normalized["website_domain"],normalized["email"],normalized["phone"],normalized.get("description",""),normalized["province"],normalized["city"],normalized["suburb"],scored["score"],scored["score_version"],json.dumps(scored["factors"]),scored["website_class"],"pending"));bid=cur.lastrowid
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,"high"ifscored["score"]>=70else"medium"ifscored["score"]>=40else"low",scored["score_version"],json.dumps(scored["factors"]),json.dumps({"source":source["kind"]})))
handler.add_job_event(db,job["id"],org,"business.created",f"Created business {bid}",60)
ifnormalized["website_domain"]andnotdb.execute("SELECT 1 FROM domains WHERE organization_id=? AND business_id=? AND domain=?",(org,bid,normalized["website_domain"])).fetchone():db.execute("INSERT INTO domains(business_id,organization_id,domain,kind) VALUES(?,?,?,?)",(bid,org,normalized["website_domain"],"website"))
ifnormalized["website"]andnotdb.execute("SELECT 1 FROM websites WHERE organization_id=? AND business_id=? AND url=?",(org,bid,normalized["website"])).fetchone():db.execute("INSERT INTO websites(business_id,organization_id,url,website_class) VALUES(?,?,?,?)",(bid,org,normalized["website"],"business_site"))
ifnormalized["website"]ornormalized["website_domain"]:db.execute("INSERT OR IGNORE INTO evidence(business_id,organization_id,kind,url,claim) VALUES(?,?,?,?,?)",(bid,org,"source_record",normalized["website"],"Discovered by "+source["kind"]))
ifnormalized["email"]ornormalized["phone"]:
ifnotdb.execute("SELECT 1 FROM contacts WHERE organization_id=? AND business_id=? AND email=? AND phone=?",(org,bid,normalized["email"],normalized["phone"])).fetchone():db.execute("INSERT INTO contacts(business_id,organization_id,email,phone) VALUES(?,?,?,?)",(bid,org,normalized["email"],normalized["phone"]))
ifrecord_id:db.execute("UPDATE source_records SET processing_status='processed',normalized_json=?,normalized_key=? WHERE id=? AND organization_id=?",(norm,nkey,record_id,org));db.execute("INSERT OR IGNORE INTO enrichment_queue(organization_id,source_record_id,status) VALUES(?,?,?)",(org,record_id,"completed"));db.execute("UPDATE enrichment_queue SET status='completed',updated_at=CURRENT_TIMESTAMP WHERE source_record_id=? AND organization_id=?",(record_id,org))
handler.add_job_event(db,job["id"],org,"enrichment.queued",f"Enrichment complete for business {bid}",75);handler.add_job_event(db,job["id"],org,"review.queued",f"Business {bid} queued for review",90)
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
db.execute("UPDATE jobs SET status='succeeded',progress=100,completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));db.commit()
exceptExceptionasexc:
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()
db.execute("UPDATE jobs SET status='cancelled',completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));server_handler.add_job_event(db,jid,org,"cancelled","Discovery cancelled",job["progress"]);db.commit();continue
ifrun_stateandrun_state["lifecycle"]=="paused":
db.execute("UPDATE jobs SET status='queued',updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));server_handler.add_job_event(db,jid,org,"paused","Discovery paused",job["progress"]);db.commit();continue
db.execute("UPDATE jobs SET status='succeeded',progress=100,completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?",(jid,));db.commit()
exceptExceptionasexc:
db.execute("UPDATE jobs SET status='failed',error_code=?,completed_at=CURRENT_TIMESTAMP,updated_at=CURRENT_TIMESTAMP WHERE id=?",(str(exc)[:80]or"SOURCE_DISCOVERY_FAILED",jid));server_handler.add_job_event(db,jid,org,"failed","Source discovery failed",job["progress"],str(exc)[:80]);db.commit()
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()
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)
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)