import json, os, sqlite3, threading, unittest from http.client import HTTPConnection from tempfile import TemporaryDirectory from app.main import create_server from app.sources import CsvSource, ManualSource class SourceAdapterTests(unittest.TestCase): def test_csv_adapter_is_deterministic_and_normalizes(self): src = CsvSource() a = src.discover({'csv': 'Name,Website,Email\n Acme ,https://acme.test,a@acme.test\n'}) b = src.discover({'csv': 'Name,Website,Email\n Acme ,https://acme.test,a@acme.test\n'}) self.assertEqual(a.records, b.records) self.assertEqual(a.records[0]['name'], 'Acme') self.assertEqual(a.records[0]['email'], 'a@acme.test') def test_manual_validation_rejects_secret_fields(self): result = ManualSource().validate({'rows': [{'name': 'x', 'api_key': 'secret'}]}) self.assertFalse(result.valid) self.assertIn('secret', result.errors[0].lower()) class SourceApiTests(unittest.TestCase): def setUp(self): self.tmp=TemporaryDirectory(); os.environ['BOOTSTRAP_ADMIN_EMAIL']='owner@example.test'; os.environ['BOOTSTRAP_ADMIN_PASSWORD']='development-password' self.server=create_server('127.0.0.1',0,self.tmp.name+'/x.db'); self.thread=threading.Thread(target=self.server.serve_forever,daemon=True); self.thread.start(); self.c=HTTPConnection('127.0.0.1',self.server.server_port); self.cookie=None self.req('POST','/api/v1/auth/login',{'email':'owner@example.test','password':'development-password'}) def tearDown(self): self.server.shutdown(); self.server.server_close(); self.thread.join(2); self.tmp.cleanup() def req(self,m,p,x=None,cookie=True): body=json.dumps(x).encode() if x is not None else None; h={'Content-Type':'application/json'} if body else {}; if cookie and self.cookie:h['Cookie']=self.cookie self.c.request(m,p,body,h); r=self.c.getresponse(); sc=r.getheader('Set-Cookie'); if sc:self.cookie=sc.split(';',1)[0] raw=r.read(); return r.status,json.loads(raw or b'{}') def test_source_lifecycle_ingestion_idempotency_health_and_disabled(self): s, source=self.req('POST','/api/v1/sources',{'name':'Import','kind':'manual','config':{}}); self.assertEqual(s,201) sid=source['id']; self.assertFalse(source['enabled']) self.assertEqual(self.req('PATCH',f'/api/v1/sources/{sid}',{'enabled':True})[0],200) payload={'rows':[{'name':'Acme','website':'https://acme.test'}], 'source_url':'file://import.csv','query_context':{'q':'test'}} self.assertEqual(self.req('POST',f'/api/v1/sources/{sid}/ingest',payload)[0],201) self.assertEqual(self.req('POST',f'/api/v1/sources/{sid}/ingest',payload)[0],200) self.assertEqual(len(self.req('GET','/api/v1/source-records')[1]['items']),1) self.assertEqual(self.req('POST',f'/api/v1/sources/{sid}/test')[0],200) self.assertEqual(self.req('PATCH',f'/api/v1/sources/{sid}',{'enabled':False})[0],200) self.assertEqual(self.req('POST',f'/api/v1/sources/{sid}/ingest',payload)[0],409) db=sqlite3.connect(self.tmp.name+'/x.db'); self.assertTrue(db.execute("select 1 from audit_log where action='source.disabled'").fetchone()); db.close() def test_queries_enqueue_and_records_are_tenant_scoped(self): _,source=self.req('POST','/api/v1/sources',{'name':'CSV','kind':'csv','enabled':True}) _,q=self.req('POST','/api/v1/discovery-queries',{'source_id':source['id'],'name':'q','query':{'csv':'name\nA'}}) status,job=self.req('POST',f"/api/v1/discovery-queries/{q['id']}/run",{}) self.assertEqual(status,202); self.assertEqual(job['type'],'source_discovery') for _ in range(100): _, current=self.req('GET',f"/api/v1/jobs/{job['id']}") if current['status'] in ('succeeded','failed'): break threading.Event().wait(.01) self.assertEqual(current['status'],'succeeded') self.assertEqual(len(self.req('GET','/api/v1/businesses')[1]['items']),1) self.assertEqual(self.req('GET','/api/v1/source-records')[1]['items'][0]['discovery_query_id'],q['id']) self.assertEqual(self.req('GET','/api/v1/source-records?page_size=101')[0],400) def test_sources_require_auth(self): self.cookie=None; self.assertEqual(self.req('GET','/api/v1/sources',cookie=False)[0],401) def test_fresh_schema_accepts_optional_source_kind_fail_closed(self): status, source = self.req('POST', '/api/v1/sources', {'name': 'RDAP', 'kind': 'rdap', 'config': {}}) self.assertEqual(status, 201) self.assertEqual(source['kind'], 'rdap') self.assertEqual(self.req('POST', f"/api/v1/sources/{source['id']}/test", {})[0], 200) status, health = self.req('GET', f"/api/v1/sources/{source['id']}/health") self.assertEqual(status, 200) self.assertFalse(health['configured']) def test_source_discovery_persists_pipeline_and_is_idempotent(self): status, source = self.req('POST', '/api/v1/sources', { 'name': 'Manual leads', 'kind': 'manual', 'enabled': True, 'config': {'rows': [{'name': 'Acme Solar', 'website': 'https://acme.test', 'email': 'hello@acme.test', 'phone': '011 555 0100', 'description': 'solar installers', 'location': 'Cape Town'}]}}) self.assertEqual(status, 201) payload = {'criteria': {'keywords': ['solar']}, 'selected_adapters': ['manual'], 'idempotency_key': 'source-run-1', 'max_records': 10} status, job = self.req('POST', '/api/v1/discovery', payload) self.assertEqual(status, 202) for _ in range(100): _, current = self.req('GET', f"/api/v1/jobs/{job['id']}") if current['status'] in ('succeeded', 'failed'): break threading.Event().wait(.01) self.assertEqual(current['status'], 'succeeded') events = self.req('GET', f"/api/v1/jobs/{job['id']}/events")[1]['items'] event_types = [event['event_type'] for event in events] for stage in ('source.started', 'source.raw_persisted', 'source.normalized', 'business.created', 'enrichment.queued', 'review.queued', 'discovery.completed'): self.assertIn(stage, event_types) businesses = self.req('GET', '/api/v1/businesses')[1]['items'] self.assertEqual(len(businesses), 1) detail = self.req('GET', f"/api/v1/businesses/{businesses[0]['id']}")[1] self.assertTrue(detail['domains']); self.assertTrue(detail['websites']); self.assertTrue(detail['evidence']) self.assertTrue(detail['contacts']); self.assertEqual(detail['review_status'], 'pending') self.assertEqual(detail['score_version'], 'opportunity-v1') self.assertTrue(all({'code', 'name', 'points', 'version'} <= set(item) for item in detail['score_factors'])) db = sqlite3.connect(self.tmp.name + '/x.db') self.assertEqual(db.execute('SELECT processing_status FROM source_records').fetchone()[0], 'processed') self.assertEqual(db.execute('SELECT status FROM enrichment_queue').fetchone()[0], 'completed') db.close() status, second = self.req('POST', '/api/v1/discovery', {**payload, 'idempotency_key': 'source-run-2'}) self.assertEqual(status, 202) for _ in range(100): _, current = self.req('GET', f"/api/v1/jobs/{second['id']}") if current['status'] in ('succeeded', 'failed'): break threading.Event().wait(.01) self.assertEqual(current['status'], 'succeeded') self.assertEqual(len(self.req('GET', '/api/v1/businesses')[1]['items']), 1) self.assertEqual(self.req('GET', '/api/v1/source-records')[1]['items'].__len__(), 1) def test_source_limits_and_run_lifecycle_are_enforced(self): status, source = self.req('POST', '/api/v1/sources', {'name': 'Limited', 'kind': 'manual', 'enabled': True, 'config': {'rows': [{'name': 'A'}, {'name': 'B'}]}, 'quota': {'daily_limit': 1, 'per_run_limit': 1}}) self.assertEqual(status, 201) status, job = self.req('POST', '/api/v1/discovery', {'criteria': {}, 'selected_adapters': ['manual'], 'idempotency_key': 'limited-1', 'max_records': 10, 'daily_limit': 1}) self.assertEqual(status, 202) runs = self.req('GET', '/api/v1/discovery-runs')[1]['items']; rid = runs[0]['id'] self.assertEqual(self.req('POST', f'/api/v1/discovery-runs/{rid}/pause', {})[0], 200) self.assertEqual(self.req('POST', f'/api/v1/discovery-runs/{rid}/resume', {})[0], 200) self.assertEqual(self.req('POST', f'/api/v1/discovery-runs/{rid}/cancel', {})[0], 200) db = sqlite3.connect(self.tmp.name + '/x.db') self.assertEqual(db.execute("SELECT lifecycle FROM discovery_runs WHERE id=?", (rid,)).fetchone()[0], 'cancelled') self.assertIn(db.execute("SELECT status FROM jobs WHERE id=?", (job['id'],)).fetchone()[0], ('cancelled', 'succeeded')) db.close() if __name__=='__main__': unittest.main()