diff --git a/scripts/build_context_bundle.py b/scripts/build_context_bundle.py
index ddeccea..17a1408 100644
--- a/scripts/build_context_bundle.py
+++ b/scripts/build_context_bundle.py
@@ -20,6 +20,7 @@
write_context_bundle,
)
from research_signal_context_pipelines.research_context_adapter import ResearchContextAdapter # noqa: E402
+from research_signal_context_pipelines.price_history import parse_price_date # noqa: E402
from research_signal_context_pipelines.theme_universe import ( # noqa: E402
build_theme_context,
load_symbol_theme_exposure,
@@ -70,11 +71,18 @@ def main() -> int:
web_research_context = None
if args.web_research_sources:
- web_research_context = ResearchContextAdapter(
+ research_adapter = ResearchContextAdapter(
Path(args.web_research_sources),
timeout_seconds=args.web_research_timeout,
max_entries=args.web_research_max_entries,
- ).build_context(pit_timestamp=generated_at)
+ )
+ end_date = parse_price_date(args.end_date) if args.end_date else None
+ if end_date is not None and end_date < generated_at.date():
+ # Calendar-date marker only, not an asserted market-close timestamp.
+ historical_cutoff = dt.datetime.combine(end_date, dt.time.min, tzinfo=dt.timezone.utc)
+ web_research_context = research_adapter.build_context(pit_timestamp=historical_cutoff)
+ else:
+ web_research_context = research_adapter.build_context()
try:
bundle = build_context_from_source(
diff --git a/src/research_signal_context_pipelines/research_context_adapter.py b/src/research_signal_context_pipelines/research_context_adapter.py
index e3127fe..7fe2b17 100644
--- a/src/research_signal_context_pipelines/research_context_adapter.py
+++ b/src/research_signal_context_pipelines/research_context_adapter.py
@@ -252,12 +252,18 @@ def __init__(
self.max_entries = int(max_entries)
def build_context(self, *, pit_timestamp: dt.datetime | None = None) -> dict[str, Any]:
- fetched_at = _isoformat_utc(pit_timestamp)
+ """Fetch current observations; past cutoffs require saved, then-visible data."""
+ started_at = dt.datetime.now(dt.timezone.utc)
+ cutoff = dt.datetime.fromisoformat(_isoformat_utc(pit_timestamp)) if pit_timestamp is not None else None
context: dict[str, Any] = {
- "pit_timestamp": fetched_at,
+ "pit_timestamp": _isoformat_utc(cutoff or started_at),
"research_sources": [],
+ "source_count": 0,
"warnings": [],
}
+ if cutoff is not None and cutoff < started_at:
+ context["warnings"].append("historical live-web research is unsupported; saved observations visible at the requested cutoff are required")
+ return context
if not self.sources_path.exists():
context["warnings"].append(f"research sources file not found: {self.sources_path}")
return context
@@ -285,10 +291,15 @@ def build_context(self, *, pit_timestamp: dt.datetime | None = None) -> dict[str
with urlopen(request, timeout=self.timeout_seconds) as response: # noqa: S310 - operator-controlled research fetch.
body = response.read()
content_type = response.headers.get_content_type() if hasattr(response.headers, "get_content_type") else None
+ fetched_time = dt.datetime.now(dt.timezone.utc)
+ fetched_at = _isoformat_utc(fetched_time)
except (OSError, URLError, TimeoutError, ValueError) as exc:
warnings.append(f"failed to fetch research source {source.url}: {type(exc).__name__}: {exc}")
continue
+ if cutoff is not None and fetched_time > cutoff:
+ warnings.append("skipped research response fetched after the requested cutoff")
+ continue
kind = _source_kind(source, content_type=content_type, body=body)
try:
entries = _feed_entries(body, source_url=source.url, fetched_at=fetched_at) if kind == "rss" else [_html_entry(body, source_url=source.url, fetched_at=fetched_at)]
@@ -299,8 +310,18 @@ def build_context(self, *, pit_timestamp: dt.datetime | None = None) -> dict[str
for entry in entries:
if len(collected) >= self.max_entries:
break
+ published_at = entry.get("published_at")
+ if not published_at:
+ warnings.append("skipped research entry with missing or invalid publication time")
+ continue
+ if dt.datetime.fromisoformat(published_at) > fetched_time:
+ warnings.append("skipped research entry with publication time after its observed fetch time")
+ continue
collected.append(entry)
+ if cutoff is None:
+ # Current-query cutoff is completion, not a timestamp before I/O.
+ context["pit_timestamp"] = _isoformat_utc(None)
context["research_sources"] = collected
context["source_count"] = len(collected)
return context
diff --git a/tests/test_research_context_adapter.py b/tests/test_research_context_adapter.py
index e8f7c76..4c72a77 100644
--- a/tests/test_research_context_adapter.py
+++ b/tests/test_research_context_adapter.py
@@ -2,9 +2,47 @@
import datetime as dt
import json
+import importlib.util
+import sys
+
+import pytest
from pathlib import Path
from research_signal_context_pipelines.research_context_adapter import ResearchContextAdapter
+from research_signal_context_pipelines import research_context_adapter as adapter_module
+
+
+@pytest.fixture(autouse=True)
+def clock(monkeypatch):
+ real_datetime = dt.datetime
+
+ class Clock(real_datetime):
+ current = real_datetime(2026, 9, 5, 12, tzinfo=dt.timezone.utc)
+ observed = []
+
+ @classmethod
+ def now(cls, tz=None):
+ value = cls.current
+ cls.observed.append(value)
+ cls.current += dt.timedelta(milliseconds=1)
+ return value.astimezone(tz) if tz else value.replace(tzinfo=None)
+
+ monkeypatch.setattr(adapter_module.dt, "datetime", Clock)
+ return Clock
+
+
+def _config(tmp_path, kind="rss"):
+ path = tmp_path / "synthetic_sources.json"
+ path.write_text(json.dumps({"whitelist": ["https://synthetic.example"],
+ "sources": [{"url": "https://synthetic.example/source", "kind": kind}]}))
+ return path
+
+
+def _rss(*dates):
+ return "" + "".join(
+ f"- synthetic-{index}{date}
"
+ for index, date in enumerate(dates)
+ ) + ""
class _FakeHeaders:
@@ -33,7 +71,7 @@ def __exit__(self, exc_type, exc, tb) -> bool:
return False
-def test_research_context_adapter_respects_whitelist_and_max_entries(tmp_path: Path, monkeypatch) -> None:
+def test_research_context_adapter_respects_whitelist_and_max_entries(tmp_path: Path, monkeypatch, clock) -> None:
config_path = tmp_path / "web_research.json"
config_path.write_text(
json.dumps(
@@ -75,11 +113,9 @@ def fake_urlopen(request, timeout):
monkeypatch.setattr("research_signal_context_pipelines.research_context_adapter.urlopen", fake_urlopen)
- context = ResearchContextAdapter(config_path, timeout_seconds=3.0, max_entries=1).build_context(
- pit_timestamp=dt.datetime(2026, 1, 1, tzinfo=dt.timezone.utc)
- )
+ context = ResearchContextAdapter(config_path, timeout_seconds=3.0, max_entries=1).build_context()
- assert context["pit_timestamp"] == "2026-01-01T00:00:00Z"
+ assert context["pit_timestamp"] == clock.observed[-1].isoformat().replace("+00:00", "Z")
assert context["source_count"] == 1
assert [item["title"] for item in context["research_sources"]] == ["First signal"]
assert context["research_sources"][0]["url"] == "https://allowed.example/a"
@@ -103,15 +139,13 @@ def fake_urlopen(request, timeout): # pragma: no cover - should not be called
monkeypatch.setattr("research_signal_context_pipelines.research_context_adapter.urlopen", fake_urlopen)
- context = ResearchContextAdapter(config_path, timeout_seconds=3.0, max_entries=3).build_context(
- pit_timestamp=dt.datetime(2026, 1, 1, tzinfo=dt.timezone.utc)
- )
+ context = ResearchContextAdapter(config_path, timeout_seconds=3.0, max_entries=3).build_context()
assert context["research_sources"] == []
assert any("skipped non-whitelisted source" in warning for warning in context["warnings"])
-def test_research_context_adapter_extracts_html_news(tmp_path: Path, monkeypatch) -> None:
+def test_research_context_adapter_extracts_html_news(tmp_path: Path, monkeypatch, clock) -> None:
config_path = tmp_path / "web_research.json"
config_path.write_text(
json.dumps(
@@ -143,9 +177,7 @@ def fake_urlopen(request, timeout):
monkeypatch.setattr("research_signal_context_pipelines.research_context_adapter.urlopen", fake_urlopen)
- context = ResearchContextAdapter(config_path, timeout_seconds=5.0, max_entries=3).build_context(
- pit_timestamp=dt.datetime(2026, 1, 1, tzinfo=dt.timezone.utc)
- )
+ context = ResearchContextAdapter(config_path, timeout_seconds=5.0, max_entries=3).build_context()
assert context["source_count"] == 1
assert context["research_sources"] == [
@@ -156,6 +188,117 @@ def fake_urlopen(request, timeout):
"url": "https://news.example/story",
"source_url": "https://news.example/story",
"source_type": "news",
- "fetched_at": "2026-01-01T00:00:00Z",
+ "fetched_at": (clock.observed[0] + dt.timedelta(milliseconds=1)).isoformat().replace("+00:00", "Z"),
}
]
+
+
+@pytest.mark.parametrize("cutoff", [dt.datetime(2026, 9, 4, tzinfo=dt.timezone.utc), dt.datetime(2026, 9, 4)])
+def test_historical_request_is_empty_unsupported_and_does_not_fetch(tmp_path, monkeypatch, cutoff):
+ calls = []
+ def fetch(*args, **kwargs):
+ calls.append(args)
+ return _FakeResponse(_rss("2026-09-03T00:00:00Z"), "application/rss+xml")
+ monkeypatch.setattr(adapter_module, "urlopen", fetch)
+ context = ResearchContextAdapter(_config(tmp_path)).build_context(pit_timestamp=cutoff)
+ assert calls == []
+ assert context["research_sources"] == []
+ assert context["source_count"] == 0
+ assert context["pit_timestamp"] == "2026-09-04T00:00:00Z"
+ assert any("unsupported" in warning and "saved" in warning for warning in context["warnings"])
+
+
+def test_current_fetch_timestamps_are_observed_after_response_not_query_start(tmp_path, monkeypatch, clock):
+ class DelayedResponse(_FakeResponse):
+ def read(self):
+ clock.current += dt.timedelta(seconds=1)
+ return super().read()
+
+ monkeypatch.setattr(adapter_module, "urlopen", lambda *args, **kwargs: DelayedResponse(
+ _rss("2026-09-05T11:59:00Z"), "application/rss+xml"))
+ context = ResearchContextAdapter(_config(tmp_path)).build_context()
+ assert context["source_count"] == 1
+ entry = context["research_sources"][0]
+ fetched = dt.datetime.fromisoformat(entry["fetched_at"].replace("Z", "+00:00"))
+ cutoff = dt.datetime.fromisoformat(context["pit_timestamp"].replace("Z", "+00:00"))
+ assert clock.observed[0] < fetched <= cutoff
+ assert fetched == clock.observed[1]
+
+
+@pytest.mark.parametrize("kind", ["rss", "news"])
+@pytest.mark.parametrize("published", ["2026-09-06T00:00:00Z", "invalid", ""])
+def test_future_and_unknown_publication_times_are_not_pit_evidence(tmp_path, monkeypatch, kind, published):
+ body = _rss(published) if kind == "rss" else f'
synthetic'
+ monkeypatch.setattr(adapter_module, "urlopen", lambda *args, **kwargs: _FakeResponse(body, "text/html"))
+ context = ResearchContextAdapter(_config(tmp_path, kind)).build_context()
+ assert context["research_sources"] == []
+ assert context["source_count"] == 0
+ assert context["warnings"]
+
+
+def test_future_feed_item_does_not_consume_entry_limit(tmp_path, monkeypatch):
+ monkeypatch.setattr(adapter_module, "urlopen", lambda *args, **kwargs: _FakeResponse(
+ _rss("2026-09-06T00:00:00Z", "2026-09-04T00:00:00Z"), "application/rss+xml"))
+ context = ResearchContextAdapter(_config(tmp_path), max_entries=1).build_context()
+ assert [entry["title"] for entry in context["research_sources"]] == ["synthetic-1"]
+
+
+@pytest.mark.parametrize("end_date,current", [(None, True), ("2026-09-05", True), ("2026-09-06", True), ("2026-09-04", False)])
+def test_cli_routes_current_and_explicit_historical_web_without_price_or_artifact_io(tmp_path, monkeypatch, clock, end_date, current):
+ script = Path(__file__).resolve().parents[1] / "scripts/build_context_bundle.py"
+ spec = importlib.util.spec_from_file_location("synthetic_build_context_cli", script)
+ cli = importlib.util.module_from_spec(spec)
+ spec.loader.exec_module(cli)
+ calls, captured, written = [], {}, []
+
+ def fetch(*args, **kwargs):
+ calls.append(True)
+ return _FakeResponse(_rss("2026-09-04T00:00:00Z"), "application/rss+xml")
+
+ def prices(**kwargs):
+ captured.update(kwargs)
+ return {"as_of": end_date or "2026-09-05", "universe": ["SPY"], "web_research": kwargs["web_research_context"]}
+
+ monkeypatch.setattr(adapter_module, "urlopen", fetch)
+ monkeypatch.setattr(cli, "build_context_from_source", prices)
+ monkeypatch.setattr(cli, "write_context_bundle", lambda bundle, path: written.append(bundle))
+ argv = [str(script), "--symbols", "SPY", "--no-theme-context", "--web-research-sources", str(_config(tmp_path))]
+ if end_date:
+ argv.extend(["--end-date", end_date])
+ monkeypatch.setattr(sys, "argv", argv)
+ assert cli.main() == 0
+ web = captured["web_research_context"]
+ assert len(calls) == int(current)
+ assert web["source_count"] == int(current)
+ assert captured["generated_at"] == clock.observed[0]
+ assert len(written) == 1
+ if current:
+ assert dt.datetime.fromisoformat(web["research_sources"][0]["fetched_at"].replace("Z", "+00:00")) > captured["generated_at"]
+ else:
+ assert web["warnings"]
+ assert web["pit_timestamp"].startswith("2026-09-04")
+
+
+def test_explicit_future_cutoff_does_not_replace_actual_fetch_time(tmp_path, monkeypatch):
+ monkeypatch.setattr(adapter_module, "urlopen", lambda *args, **kwargs: _FakeResponse(
+ _rss("2026-09-04T00:00:00Z"), "application/rss+xml"))
+ cutoff = dt.datetime(2026, 9, 6, tzinfo=dt.timezone.utc)
+ context = ResearchContextAdapter(_config(tmp_path)).build_context(pit_timestamp=cutoff)
+ assert context["source_count"] == 1
+ assert context["pit_timestamp"] == "2026-09-06T00:00:00Z"
+ assert context["research_sources"][0]["fetched_at"].startswith("2026-09-05T12:00:00.")
+
+
+def test_response_finishing_after_explicit_cutoff_is_not_backdated(tmp_path, monkeypatch, clock):
+ class DelayedResponse(_FakeResponse):
+ def read(self):
+ clock.current += dt.timedelta(seconds=2)
+ return super().read()
+
+ monkeypatch.setattr(adapter_module, "urlopen", lambda *args, **kwargs: DelayedResponse(
+ _rss("2026-09-04T00:00:00Z"), "application/rss+xml"))
+ cutoff = clock.current + dt.timedelta(seconds=1)
+ context = ResearchContextAdapter(_config(tmp_path)).build_context(pit_timestamp=cutoff)
+ assert context["source_count"] == 0
+ assert context["pit_timestamp"] == "2026-09-05T12:00:01Z"
+ assert any("fetched after" in warning for warning in context["warnings"])