2323 lead_row ,
2424 normalize_domain_from_url ,
2525)
26- from leadgen .pipeline .checkpoints import CheckpointWriter
2726from leadgen .pipeline .dedupe import dedupe_leads
27+ from leadgen .pipeline .event_log import EventLogWriter
2828from leadgen .scraping .discovery import discover_contact_like_urls , homepage_url
2929from leadgen .scraping .extractor import extract_contacts
3030from leadgen .scraping .fetcher import Fetcher
@@ -109,27 +109,27 @@ def __init__(
109109 self ._scrape_delay_seconds = scrape_delay_seconds
110110 self ._max_output_leads = max_output_leads
111111 self ._sleep = sleeper
112- self ._checkpoints = CheckpointWriter (checkpoint_dir ) if checkpoint_dir else None
112+ self ._event_log = EventLogWriter (checkpoint_dir ) if checkpoint_dir else None
113113
114114 def run (self , * , queries_path : Path , output_csv_path : Path ) -> PipelineOutput :
115115 queries_file = load_queries (queries_path )
116116
117117 logger .info ("Loaded %s queries from %s" , len (queries_file .queries ), queries_path )
118- if self ._checkpoints :
119- self ._checkpoints .write_json (
120- "queries.validated" ,
121- queries_file .model_dump (),
122- )
118+ if self ._event_log :
119+ self ._event_log .append ("queries.validated" , queries_file .model_dump ())
123120 domain_contexts : dict [str , _DomainContext ] = {}
124121 leads : list [Lead ] = []
125122 for query in queries_file .queries :
126123 logger .info ("Searching query_id=%s" , query .id )
127124 results = self ._search .search (query )
128125 logger .info ("Search results query_id=%s count=%s" , query .id , len (results ))
129- if self ._checkpoints :
130- self ._checkpoints .write_json (
131- f"search_results.{ query .id } " ,
132- [r .model_dump () | {"domain" : r .domain } for r in results ],
126+ if self ._event_log :
127+ self ._event_log .append (
128+ "search.results" ,
129+ {
130+ "query_id" : query .id ,
131+ "results" : [r .model_dump () | {"domain" : r .domain } for r in results ],
132+ },
133133 )
134134 leads .extend (
135135 self ._scrape_results (
@@ -139,19 +139,24 @@ def run(self, *, queries_path: Path, output_csv_path: Path) -> PipelineOutput:
139139 domain_contexts = domain_contexts ,
140140 )
141141 )
142- if self ._checkpoints :
143- self ._checkpoints .write_json (
144- f"leads.scraped.{ query .id } " ,
145- [lead .model_dump () for lead in leads if lead .source_query_id == query .id ],
142+ if self ._event_log :
143+ self ._event_log .append (
144+ "leads.scraped" ,
145+ {
146+ "query_id" : query .id ,
147+ "leads" : [
148+ lead .model_dump () for lead in leads if lead .source_query_id == query .id
149+ ],
150+ },
146151 )
147152
148153 logger .info ("Collected %s raw leads (pre-dedupe)" , len (leads ))
149- if self ._checkpoints :
150- self ._checkpoints . write_json ("leads.raw" , [lead .model_dump () for lead in leads ])
154+ if self ._event_log :
155+ self ._event_log . append ("leads.raw" , [lead .model_dump () for lead in leads ])
151156 deduped = dedupe_leads (leads )
152157 logger .info ("Deduped leads by domain: %s -> %s" , len (leads ), len (deduped ))
153- if self ._checkpoints :
154- self ._checkpoints . write_json ("leads.deduped" , [lead .model_dump () for lead in deduped ])
158+ if self ._event_log :
159+ self ._event_log . append ("leads.deduped" , [lead .model_dump () for lead in deduped ])
155160
156161 filtered = [
157162 lead for lead in deduped if (is_valid_email (lead .email ) or is_valid_phone (lead .phone ))
@@ -166,9 +171,9 @@ def run(self, *, queries_path: Path, output_csv_path: Path) -> PipelineOutput:
166171 )
167172 )
168173 limited = filtered [: self ._max_output_leads ]
169- if self ._checkpoints :
170- self ._checkpoints . write_json ("leads.filtered" , [lead .model_dump () for lead in filtered ])
171- self ._checkpoints . write_json ("leads.final" , [lead .model_dump () for lead in limited ])
174+ if self ._event_log :
175+ self ._event_log . append ("leads.filtered" , [lead .model_dump () for lead in filtered ])
176+ self ._event_log . append ("leads.final" , [lead .model_dump () for lead in limited ])
172177
173178 written = _write_csv (output_csv_path , limited )
174179 logger .info ("Wrote output CSV: %s" , output_csv_path )
0 commit comments