Repository navigation
feat: add SerpAPI dataset source - #312
Conversation
houtanb
left a comment
There was a problem hiding this comment.
Looks great. Some small changes here. Need to look more in depth later, but this is something for now
7bc07fa to
f0a795e
Compare
|
I still need more time to dig in, but probably useful for you to address the above first. I'd also recommend:
|
f0a795e to
2975f83
Compare
|
@nikbpetrov |
b4ff05b to
42a0e92
Compare
There was a problem hiding this comment.
(This was posted by Fable 5.1 in my name)
@pythoryn I went through each comment below and elaborated on them, noting what seemed important.
Full review of the SerpAPI source at 42a0e92. Well-structured and thoroughly tested; nothing here is a hard blocker.
Main risk (1 comment): a systematic failure in one category (e.g. a payload shape change for one engine) is only logged per id. Every question in that category silently retires from curation and, after the fallback window, resolves to NaN. Only the all-requests-failed case raises.
Maintainability (3 comments): per-category behavior (date filling, snapshot fallback, flight handling) is hardcoded in four places instead of on the spec, with one latent divergence between update() and _resolve(); growth-threshold arithmetic is implemented twice; the update job forks three _source_io helpers in a way that will drop fetch_datetime if anyone later unifies them.
Minor (5 comments): secret access pattern, a test that pins private call shapes, an O(n²) loop in _resolve, a dead explanation string, and an unused browser_url() helper.
Generated by Claude Code
| snapshot_names = { | ||
| "amazon_minimum_product_price", | ||
| "walmart_food_drink_price", | ||
| "youtube_channel_max_video_views", | ||
| } | ||
| assert all( | ||
| name in snapshot_names | ||
| for name, spec in QUESTION_SPECS.items() | ||
| if spec["engine"] in TODAY_ONLY_ENGINES | ||
| ), "Every today-only question category must have snapshot resolution support." | ||
| dfr = dfr[["id", "date", "value"]].copy() | ||
| numeric = pd.to_numeric(dfr["value"], errors="coerce") | ||
| dfr["value"] = numeric.where(np.isfinite(numeric), np.nan) | ||
| # Finance histories retain raw observations and missing values. Fill only | ||
| # for resolution, after all saved and refreshed observations were merged. | ||
| finance = dfr["id"].str.startswith("google_finance_ticker_price__") |
There was a problem hiding this comment.
Category behavior is hardcoded in four places instead of on the spec. update() decides date-filling from spec.get("fill_missing_dates") (line 240), but _resolve() decides it from the google_finance_ticker_price__ id prefix here. They agree today only because google_finance is the sole spec with the flag. A second spec that sets fill_missing_dates would get a filled baseline in the question bank and an unfilled resolution baseline, so the forecasted and scored events diverge.
The same pattern shows up in the snapshot_names set + assert here, spec["engine"] == "google_finance" in fetch() (line 126), and name == "flight_departure_delay" in fetch() and update() (lines 148, 302). Adding a sixth category means touching all of them.
Suggest deriving these from the spec so the "declarative request specification" in the module docstring holds:
filled_ids = {id for id, (name, spec, _) in configured.items() if spec.get("fill_missing_dates")}
snapshot_names = {name for name, spec in QUESTION_SPECS.items() if spec["engine"] in TODAY_ONLY_ENGINES}If retired categories must keep resolving after their spec is removed (per the comment above), keep a small RETIRED_SNAPSHOT_NAMES constant and union it in, rather than freezing the live set here.
Generated by Claude Code
There was a problem hiding this comment.
I got a more helpful comment for this:
update() checks the spec flag to decide whether to forward-fill. Good. Nothing to change_resolve()does not check the flag. It checks whether the question id starts with the literal stringgoogle_finance_ticker_price__.fetch()checksspec["engine"] == "google_finance"to decide whether to parse a price history.fetch()andupdate()checkname == "flight_departure_delay"for the flight-specific logic._resolve()has a hand-typed set of the three snapshot category names, plus an assert that the set matches the specs.
There was a problem hiding this comment.
ok yes, resolution should follow the same configuration as update. Done. I’ve kept the specific finance and flight delay branches because they handle different response structures and recovery rules.
| return frame if not frame.empty else pd.DataFrame(columns=constants.QUESTION_FILE_COLUMNS) | ||
|
|
||
|
|
||
| def upload_resolution_files(resolution_files: dict[str, pd.DataFrame]) -> None: |
There was a problem hiding this comment.
Forked orchestration helpers. This job reimplements question-bank load, per-id history download, and resolution upload rather than using _source_io. I understand why: the shared upload helper slices to [id, date, value] and would drop fetch_datetime, and the shared download swallows non-404 errors. But that makes the divergence a trap. Anyone who later "cleans this up" by switching to _source_io.upload_resolution_files silently loses the timestamps that guard against stale replays.
Two options, either is fine:
- Extend
_source_iowith what serpapi needs (anextra_columnspassthrough on upload, a strict download that re-raises anything but 404) and use it here. Other sources would benefit from the strict download too. - Keep the fork but add a short comment on each helper saying which shared-helper behavior it deliberately differs from and why.
Generated by Claude Code
There was a problem hiding this comment.
@nikbpetrov is authoratative here, but my take is that these should go into _source_io.py. Snippets from Claude that would clean up src/orchestration/func_serpapi_update/main.py
def upload_resolution_files(source: str, resolution_files: dict[str, pd.DataFrame], *, extra_columns: Iterable[str] = ()) -> None:
...
columns = ["id", "date", "value", *extra_columns]
df[columns].to_json(
...and
def load_existing_resolution_files(
source: str,
ids: Iterable[str] | None = None,
*,
strict: bool = False,
) -> dict[str, pd.DataFrame]:
...
if os.path.exists(local_filename):
os.remove(local_filename) # playbook §2.5: never read a stale /tmp copy
if strict:
try:
gcp.storage.download(bucket_name=..., filename=remote_path, local_filename=local_filename)
except NotFound:
continue
else:
gcp.storage.download_no_error_message_on_404(...)There was a problem hiding this comment.
yes, I was agonizing about this and thought initially about modifying _source_io.py but felt the branch and PR should stay within the bounds of the serpapi branch. Hence, my decision to place upload_resolution_files in src/orchestration/func_serpapi_update/main.py. But since you don't mind I have moved this into _source_io.py, adding optional extra columns and strict history downloads. SerpAPI retains its timestamps and download-error safeguards, while existing callers retain their defaults.
| keys.get_secret.assert_called_once_with("API_KEY_SERPAPI") | ||
| assert source.api_key == "test-key" | ||
| write.assert_called_once_with("serpapi", fetched) |
There was a problem hiding this comment.
Minor: these assertions pin the private call shape of keys.get_secret and _source_io.write_fetch_output, which AGENTS.md asks us not to do. Switching to keys.API_KEY_SERPAPI (see comment on the entrypoint) would break this test without changing behavior. The source.api_key == "test-key" assertion plus checking what was written (e.g. the frame passed to the upload, or the uploaded file's contents) covers the contract.
Generated by Claude Code
There was a problem hiding this comment.
Agreed. Fixed by previous secret-access edit.
42a0e92 to
83d24d8
Compare
| fetch_datetime: Series[str] | ||
|
|
||
|
|
||
| class SerpapiFetchFrame(ResolutionFrame): |
There was a problem hiding this comment.
FetchFrame inheriting from ResolutionFrame is, to put it mildly, not expected. I would prefer you repeated the cols/col definitions as opposed to this
There was a problem hiding this comment.
You are right. I changed SerpapiFetchFrame to inherit directly from pa.DataFrameModel, with explicit column definitions and validation settings.
| ) | ||
| if failures and failures == requests_attempted: | ||
| raise RuntimeError("All SerpAPI requests failed; retaining the previous fetch output.") | ||
| return pd.DataFrame(rows, columns=list(SerpapiFetchFrame.to_schema().columns)) |
There was a problem hiding this comment.
pd.DataFrame(rows) should do the job, as is the case w/ other sources, no need for columns=... stuff
There was a problem hiding this comment.
Note that the explicit columns also preserve the fetch schema when no questions are eligible (for example if all configured IDs are nullified). Without them, the empty dataframe fails pandera validation. Note also that kalshi fetch also explicitly supplies column. I have kept this and add a short comment explaining why.
There was a problem hiding this comment.
Good catch for kalshi - it's bad. Openeed #319
There was a problem hiding this comment.
no need to change this for now then
| @@ -0,0 +1,14 @@ | |||
| """SerpAPI metadata and stable question-ID semantics.""" | |||
There was a problem hiding this comment.
if i understand, the sole purpose of this is more of a path as the core refactor is ongoing; specifically, eventually the naive and dummy baselines job should be able to use the sources directly so importing from helpers/
if that's the intent behind this file, then i'd actually prefer to see this as a separate commit, like tmp(serpapi) or patch(serpapi) or sth so that after/during the refactor it gets reverted/updated but @houtanb might have another pref
There was a problem hiding this comment.
Indeed, these are hangovers until the refactor is complete.
No strong feelings on the separate commit, especially if pressed for time given the urgency of merging this code. It will be cleaned up anyway once we move away from all of those helpers/<source> files
There was a problem hiding this comment.
Yes, this supports the existing helper imports used by curation and the baseline forecaster. It also contains percentage-threshold functions shared by the baseline and source resolver, so those will need to be relocated when this module is removed. Given @houtanb’s response, I'll keep the current commit structure and leave that migration to the broader refactor.
| for day, price in sorted(prices.items()) | ||
| ) | ||
| continue | ||
| if name == "flight_departure_delay": |
There was a problem hiding this comment.
ok, changed to elif
| raise ValueError("Expected a JSON object from SerpAPI.") | ||
| if "error" in data: | ||
| raise ValueError(data["error"]) | ||
| if spec["engine"] == "google_finance": |
There was a problem hiding this comment.
actuall,y should the engine-specific code even be in fetch? ideally it'd be a lot neater, similar to the other soruces so it's auditable w/o knowledge of serapi-specifics
There was a problem hiding this comment.
Ok. I moved the engine-specific parsing into a helper so fetch() is easier to follow.
| """Append observations to history and refresh the reusable question bank. | ||
|
|
||
| Args: | ||
| dfq (DataFrame[QuestionFrame]): Existing questions. |
There was a problem hiding this comment.
just noticed this in other sources too but I dont think the Args' definition should be re-specified in the docstring -- no need to fix as this is a sources-wide issue
There was a problem hiding this comment.
Understood. I’ll keep the existing style for this PR.
| "amazon_minimum_product_price": { | ||
| "question_template": ( | ||
| "Will the lowest returned listed price in USD for '{product}' on Amazon.com, " | ||
| "across all conditions (including used), be higher on {resolution_date} " |
There was a problem hiding this comment.
all product conditions?
There was a problem hiding this comment.
ok. changed to “all product conditions” for clarity
| ) | ||
| result = source.update(dfq, dff, existing_resolution_files=existing) | ||
| if result.resolution_files: | ||
| upload_resolution_files(result.resolution_files) |
There was a problem hiding this comment.
ditto as above - why not:
_source_io.upload_resolution_files(SOURCE, result.resolution_files)
There was a problem hiding this comment.
okay i kinda see - it limts df to 3 cols which is not the case for serapi - but i'd recommend patching the source_io func rather than reinventing one here
There was a problem hiding this comment.
Yes, SerpAPI now uses the shared uploader, with fetch_datetime retained
|
@pythoryn wtf i didnt click submit on my original comments so i am posting them onyl now will add to some of them as some seem to be addressed |
83d24d8 to
03d6a60
Compare
6734fc1 to
5f5aa33
Compare
698e63f to
196188e
Compare
Adds a SerpAPI dataset source with 150 configured questions covering Amazon
prices, Walmart food and drink prices, and flight departure delays.
Includes daily collection, persistent observation histories, question curation,
resolution, and Prophet baseline forecasts.
Daily snapshots allow a seven-day fallback for missing comparison dates,
with baselines selected separately for each horizon.
departure. Resolution compares the target delay with the preceding 14-day
median, counting early/on-time departures as zero and omitting missing days.
Horizons lacking required measurements remain unresolved and unscored.
Deployment prerequisite: Before the first nightly run that includes SerpAPI, create an empty 'serpapi_questions.jsonl'
at the root of the production question-bank bucket:
Closes #49