diff --git a/scraper.py b/scraper.py index 3bf78ad..470c926 100644 --- a/scraper.py +++ b/scraper.py @@ -61,36 +61,44 @@ def _collect_normalized_pages( since_ts: Optional[int], ) -> Tuple[List[Dict[str, Any]], List[str], Optional[str], int]: """ - Collect normalized items across pages until we have at least `target` items when counting: - - all new (not in DB), plus - - duplicates we can top up from DB later. + Collect pages while counting ONLY *new* posts toward `target`. + Keep dup IDs to optionally top-up from DB after paging. Returns: - normalized_new -> list of normalized, not in DB - dup_ids_ordered -> list of IG IDs already in DB (for top-up) - final_cursor -> next_max_id to persist if we traversed older pages - pages_walked + normalized_new -> normalized items not in DB (in order found) + dup_ids_ordered -> ig_post_id list already in DB (in order found) + final_cursor -> last next_max_id seen (persist to continue older) + pages_walked -> how many API pages we pulled """ normalized_new: List[Dict[str, Any]] = [] dup_ids_ordered: List[str] = [] pages = 0 - # Phase A: NEWEST first + # Start from stored cursor if present to keep drilling older across runs + max_id = start_cursor newest_cursor = None - max_id = None - while (len(normalized_new) + len(dup_ids_ordered) < target) and pages < MAX_PAGES: - per_req = 12 # always ask 12 for fewer roundtrips + final_cursor: Optional[str] = None + seen_cursors: Set[Optional[str]] = set() + + while len(normalized_new) < target and pages < MAX_PAGES: + per_req = 12 # BoxAPI cap per request try: page, _log, nxt = get_media_page(username=username, count=per_req, max_id=max_id) except BoxAPIError as e: - dprint(f"BoxAPI error (A): {e}") + dprint(f"BoxAPI error: {e}") break pages += 1 - dprint(f"[A] Page#{pages}: request max_id={max_id!r} → got {len(page)} items, next_max_id={nxt!r}") + dprint(f"[PAGE]#{pages}: request max_id={max_id!r} → got {len(page)} items, next_max_id={nxt!r}") if newest_cursor is None: newest_cursor = nxt + if not page: + dprint("No items → stopping.") + break + norm = _normalize_page(page) + + # optional time filter if since_ts is not None: try: s = int(since_ts) @@ -103,61 +111,35 @@ def _collect_normalized_pages( new_here = [n for n in norm if n.get("ig_post_id") and str(n["ig_post_id"]) not in already] dups_here = [str(n["ig_post_id"]) for n in norm if n.get("ig_post_id") and str(n["ig_post_id"]) in already] + # Append in discovery order normalized_new.extend(new_here) dup_ids_ordered.extend(dups_here) - dprint(f" normalized={len(norm)} new={len(new_here)} dups={len(dups_here)} " - f"total_collected={len(normalized_new)+len(dup_ids_ordered)}") + dprint( + f" normalized={len(norm)} new={len(new_here)} dups={len(dups_here)} " + f"new_collected_total={len(normalized_new)}/{target}" + ) - if len(normalized_new) >= target or (len(normalized_new) + len(dup_ids_ordered)) >= target: - # enough candidates to build exact N - return normalized_new, dup_ids_ordered, None, pages - - if not nxt: - break - max_id = nxt - - # Phase B: BACKFILL older (use stored cursor if provided, else continue from newest_cursor) - backfill_cursor = start_cursor or newest_cursor - final_cursor = None - while (len(normalized_new) + len(dup_ids_ordered) < target) and backfill_cursor and pages < MAX_PAGES: - per_req = 12 - try: - page, _log, nxt = get_media_page(username=username, count=per_req, max_id=backfill_cursor) - except BoxAPIError as e: - dprint(f"BoxAPI error (B): {e}") + # Stop only when we have enough NEW items + if len(normalized_new) >= target: + final_cursor = nxt # where to resume next time (older side) break - pages += 1 - dprint(f"[B] Page#{pages}: request max_id={backfill_cursor!r} → got {len(page)} items, next_max_id={nxt!r}") - final_cursor = nxt # we advanced older pages - - norm = _normalize_page(page) - if since_ts is not None: - try: - s = int(since_ts) - norm = [n for n in norm if (n.get("taken_at") or 0) > s] - except Exception: - pass - - ids = [str(n.get("ig_post_id")) for n in norm if n.get("ig_post_id")] - already = _db_existing_ids(db, ids) - new_here = [n for n in norm if n.get("ig_post_id") and str(n["ig_post_id"]) not in already] - dups_here = [str(n["ig_post_id"]) for n in norm if n.get("ig_post_id") and str(n["ig_post_id"]) in already] - - normalized_new.extend(new_here) - dup_ids_ordered.extend(dups_here) - - dprint(f" normalized={len(norm)} new={len(new_here)} dups={len(dups_here)} " - f"total_collected={len(normalized_new)+len(dup_ids_ordered)}") - - if len(normalized_new) >= target or (len(normalized_new) + len(dup_ids_ordered)) >= target: - return normalized_new, dup_ids_ordered, final_cursor, pages - + # Advance to older page; if none, we're done if not nxt: dprint("No next_max_id → end of feed") + final_cursor = None break - backfill_cursor = nxt + + # Loop guard: avoid cursor cycles + if nxt in seen_cursors: + dprint("Repeating cursor detected → stopping.") + final_cursor = nxt + break + + seen_cursors.add(nxt) + final_cursor = nxt + max_id = nxt # CRITICAL: go older return normalized_new, dup_ids_ordered, final_cursor, pages