diff --git a/data/fred/IPMAN.parquet b/data/fred/IPMAN.parquet new file mode 100644 index 000000000..0c2bd4e4e Binary files /dev/null and b/data/fred/IPMAN.parquet differ diff --git a/implementations/manufacturing_stress_forecasting/README.md b/implementations/manufacturing_stress_forecasting/README.md new file mode 100644 index 000000000..c31f8b881 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/README.md @@ -0,0 +1,76 @@ +# Manufacturing stress forecasting — minimal IPMAN MVP + +This implementation asks one Track 1 question: + +> Given information available at a monthly forecast origin, what is the +> probability that U.S. manufacturing will be under stress three months later? + +The model deliberately uses only five explanatory variables: trailing +1-, 3-, and 6-month IPMAN changes, the effective federal funds rate, and the +10-year minus 2-year Treasury yield spread. The small panel keeps the first +multivariate experiment interpretable. + +## Target + +A month is labelled `1` (stress) when IPMAN has declined by at least 2% over +its preceding three months; otherwise it is `0`. The threshold is a provisional +version-1 definition and should be reviewed visually before expanding the +project. + +The forecast made at month `t` predicts the stress label at `t + 3 months`. +That distinction makes this forecasting rather than current-state detection. + +## Predictors + +- `HistoricalFrequencyPredictor`: the visible historical stress rate. +- `ManufacturingStressLogisticPredictor`: fit-at-origin logistic regression on + the five IPMAN/rate variables. +- `manufacturing_stress_analyst`: a structured LLM predictor receiving the same + five cutoff-safe signals plus recent IPMAN history and historical base rates. + +All predictors return `BinaryForecast` probabilities; backtested predictors are scored with Brier score. + +## Data and cutoff assumptions + +`FREDAdapter` caches `IPMAN`, `DFF`, `DGS10`, and `DGS2` under `data/fred/`. +IPMAN is conservatively treated as available one month after its reference +month. Daily rate observations are treated as available on the next business +day and collapsed to their final monthly observation. The standard FRED API +does not provide full point-in-time vintages, so historical observations may +still contain later revisions; a production study should use ALFRED vintages. + +## Run + +From the repository root, put a personal FRED key in `.env` or export it: + +```bash +export FRED_API_KEY="..." +``` + +Populate the cache and inspect the registered series: + +```bash +uv run python scripts/fetch_manufacturing_stress.py +``` + +Run the small backtest: + +```bash +uv run --directory implementations python -m manufacturing_stress_forecasting.run_smoke +``` + +The output prints one mean Brier score per predictor; lower is better. The +logistic model should be compared against historical frequency, not judged in +isolation. + +Run one current forecast, including the structured agent: + +```bash +uv run --directory implementations python -m manufacturing_stress_forecasting.run_agent_prediction +``` + +## Next steps + +1. Plot IPMAN and the derived stress months; confirm or revise the 2% threshold. +2. Compare the five-variable logistic score with the earlier IPMAN-only result. +3. Backtest the agent only after the deterministic model is stable. diff --git a/implementations/manufacturing_stress_forecasting/__init__.py b/implementations/manufacturing_stress_forecasting/__init__.py new file mode 100644 index 000000000..e4cbe5929 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/__init__.py @@ -0,0 +1,16 @@ +"""Minimal IPMAN-based manufacturing-stress forecasting use case.""" + +from manufacturing_stress_forecasting.data import ( + IPMAN_SERIES_ID, + STRESS_SERIES_ID, + build_manufacturing_stress_service, +) +from manufacturing_stress_forecasting.predictors import ManufacturingStressLogisticPredictor + + +__all__ = [ + "IPMAN_SERIES_ID", + "STRESS_SERIES_ID", + "ManufacturingStressLogisticPredictor", + "build_manufacturing_stress_service", +] diff --git a/implementations/manufacturing_stress_forecasting/analyst_agent/__init__.py b/implementations/manufacturing_stress_forecasting/analyst_agent/__init__.py new file mode 100644 index 000000000..4d704e53c --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/analyst_agent/__init__.py @@ -0,0 +1,14 @@ +"""Quantitative-only manufacturing-stress analyst agent.""" + +from manufacturing_stress_forecasting.analyst_agent.agent import ( + ManufacturingStressPromptBuilder, + build_manufacturing_stress_agent_config, + build_manufacturing_stress_agent_predictor, +) + + +__all__ = [ + "ManufacturingStressPromptBuilder", + "build_manufacturing_stress_agent_config", + "build_manufacturing_stress_agent_predictor", +] diff --git a/implementations/manufacturing_stress_forecasting/analyst_agent/agent.py b/implementations/manufacturing_stress_forecasting/analyst_agent/agent.py new file mode 100644 index 000000000..1f3c7e3b9 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/analyst_agent/agent.py @@ -0,0 +1,163 @@ +"""Quantitative-only ADK agent for binary manufacturing-stress forecasts.""" + +from __future__ import annotations + +import json +from typing import Any + +import pandas as pd +from aieng.forecasting.data.context import ForecastContext +from aieng.forecasting.evaluation.task import ForecastingTask +from aieng.forecasting.methods.agentic import ( + AgentPredictor, + DiscreteAgentForecastOutput, + build_adk_agent, +) +from aieng.forecasting.methods.agentic.agent_factory import AgentConfig +from aieng.forecasting.models import LITE_MODEL +from manufacturing_stress_forecasting.data import IPMAN_SERIES_ID +from manufacturing_stress_forecasting.features import ( + FEATURE_SERIES_IDS, + IPMAN_FEATURE_SERIES_IDS, + MACRO_FEATURE_SERIES_IDS, + build_feature_snapshot, +) +from manufacturing_stress_forecasting.targets import ( + DEFAULT_LOOKBACK_MONTHS, + DEFAULT_STRESS_THRESHOLD_PCT, +) +from pydantic import BaseModel, Field + + +def _build_instruction() -> str: + schema = DiscreteAgentForecastOutput.prompt_schema_json() + return ( + "## Role\n\n" + "You are a cautious U.S. manufacturing-cycle analyst. Estimate the probability that the " + "binary IPMAN stress event in the supplied task resolves to 1 at the specified forecast date.\n\n" + "## Rules\n\n" + "1. Use only the JSON payload. Do not use remembered events or facts after `as_of`.\n" + "2. Start from the supplied historical base rate, then adjust using the five supplied signals.\n" + "3. Treat negative IPMAN momentum, a restrictive fed funds rate, and an inverted 10Y-2Y spread " + "as possible evidence for stress; explain how the signals interact.\n" + "4. Do not double-count correlated signals or turn a weak signal into certainty.\n" + "5. `probability` means P(stress=1), not confidence in your explanation.\n" + "6. Give a concise rationale, identify both supporting and countervailing evidence, and remain calibrated.\n" + "7. Use `direction_bias='down'` when signals point toward manufacturing stress, `up` when they point " + "away from stress, and `neutral` when mixed.\n\n" + "## Output\n\n" + "Return exactly one JSON object matching this structure, with no markdown fence or preamble:\n\n" + schema + ) + + +class ManufacturingStressPromptBuilder(BaseModel): + """Serialize cutoff-safe IPMAN evidence into the agent's prompt.""" + + model_config = {"extra": "forbid"} + + recent_history_months: int = Field(default=24, ge=6, le=120) + trailing_base_rate_months: int = Field(default=60, ge=12, le=240) + + def __call__(self, *, task: ForecastingTask, context: ForecastContext) -> str: + """Build one structured, cutoff-safe forecast payload.""" + if task.payload_type != "binary" or len(task.horizons) != 1: + raise ValueError("ManufacturingStressPromptBuilder requires one binary forecast horizon.") + + as_of = pd.Timestamp(context.as_of) + offset = pd.tseries.frequencies.to_offset(task.frequency) + forecast_date = as_of + offset * task.horizons[0] + ipman = context.get_series(IPMAN_SERIES_ID).sort_values("timestamp") + target = context.get_series(task.target_series_id).sort_values("timestamp") + feature_frames = {series_id: context.get_series(series_id) for series_id in FEATURE_SERIES_IDS} + current_signals = build_feature_snapshot(as_of, feature_frames) + current_ipman_signals = ( + {series_id: current_signals[series_id] for series_id in IPMAN_FEATURE_SERIES_IDS} + if current_signals is not None + else None + ) + current_macro_signals = ( + {series_id: current_signals[series_id] for series_id in MACRO_FEATURE_SERIES_IDS} + if current_signals is not None + else None + ) + + target_values = target["value"].astype(float) + trailing_values = target_values.tail(self.trailing_base_rate_months) + recent_ipman = [ + { + "reference_month": str(pd.Timestamp(timestamp).date()), + "value": float(value), + "released_at": str(pd.Timestamp(released_at).date()), + } + for timestamp, value, released_at in zip( + ipman["timestamp"].tail(self.recent_history_months), + ipman["value"].tail(self.recent_history_months), + ipman["released_at"].tail(self.recent_history_months), + strict=True, + ) + ] + + payload: dict[str, Any] = { + "task": { + "task_id": task.task_id, + "question": task.description, + "horizon_months": task.horizons[0], + }, + "as_of": str(as_of.date()), + "forecast_date": str(forecast_date.date()), + "target_definition": { + "event": "manufacturing stress", + "stress_value": 1, + "no_stress_value": 0, + "lookback_months": DEFAULT_LOOKBACK_MONTHS, + "threshold_pct": DEFAULT_STRESS_THRESHOLD_PCT, + "rule": ("stress=1 when trailing IPMAN percentage change is less than or equal to threshold_pct"), + }, + "current_ipman_signals_pct": current_ipman_signals, + "current_macro_signals": current_macro_signals, + "historical_stress": { + "n_visible_months": len(target_values), + "all_history_base_rate": float(target_values.mean()) if len(target_values) else None, + "trailing_window_months": self.trailing_base_rate_months, + "trailing_base_rate": float(trailing_values.mean()) if len(trailing_values) else None, + }, + "recent_ipman": recent_ipman, + } + return json.dumps(payload, indent=2) + + +def build_manufacturing_stress_agent_config(model: str = LITE_MODEL) -> AgentConfig: + """Build the tool-free manufacturing analyst configuration.""" + return AgentConfig( + name="manufacturing_stress_analyst", + model=model, + instruction=_build_instruction(), + temperature=0.1, + seed=42, + max_output_tokens=2_048, + ) + + +def build_manufacturing_stress_agent_predictor( + config: AgentConfig | None = None, +) -> AgentPredictor: + """Wrap the analyst in the standard binary AgentPredictor contract.""" + return AgentPredictor( + agent_config=config or build_manufacturing_stress_agent_config(), + prompt_builder=ManufacturingStressPromptBuilder(), + output_schema=DiscreteAgentForecastOutput, + ) + + +def __getattr__(name: str) -> Any: + """Expose a schema-free root agent for ``adk run`` and ``adk web``.""" + if name == "root_agent": + return build_adk_agent(build_manufacturing_stress_agent_config()) + raise AttributeError(f"module {__name__!r} has no attribute {name!r}") + + +__all__ = [ + "ManufacturingStressPromptBuilder", + "build_manufacturing_stress_agent_config", + "build_manufacturing_stress_agent_predictor", +] diff --git a/implementations/manufacturing_stress_forecasting/data.py b/implementations/manufacturing_stress_forecasting/data.py new file mode 100644 index 000000000..5a6f184d5 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/data.py @@ -0,0 +1,140 @@ +"""FRED data service for the five-variable manufacturing-stress MVP.""" + +from __future__ import annotations + +from pathlib import Path + +from aieng.forecasting.data import DataService, SeriesMetadata +from aieng.forecasting.data.adapters import FREDAdapter +from aieng.forecasting.data.features import StaticFrameAdapter +from manufacturing_stress_forecasting.features import ( + FEATURE_PERIODS, + FED_FUNDS_SERIES_ID, + YIELD_CURVE_SERIES_ID, + apply_conservative_monthly_release_lag, + build_ipman_feature_frames, + build_macro_feature_frames, +) +from manufacturing_stress_forecasting.targets import ( + DEFAULT_LOOKBACK_MONTHS, + DEFAULT_STRESS_THRESHOLD_PCT, + derive_manufacturing_stress_labels, +) + + +IPMAN_FRED_ID = "IPMAN" +FED_FUNDS_FRED_ID = "DFF" +TREASURY_10Y_FRED_ID = "DGS10" +TREASURY_2Y_FRED_ID = "DGS2" + +IPMAN_SERIES_ID = "ipman_us_manufacturing_production" +STRESS_SERIES_ID = "manufacturing_stress" + +_REPO_ROOT = Path(__file__).resolve().parents[2] +DEFAULT_FRED_CACHE_DIR = _REPO_ROOT / "data" / "fred" + + +def build_manufacturing_stress_service( + *, + cache_dir: str | Path = DEFAULT_FRED_CACHE_DIR, + refresh: bool = False, + release_lag_months: int = 1, + stress_lookback_months: int = DEFAULT_LOOKBACK_MONTHS, + stress_threshold_pct: float = DEFAULT_STRESS_THRESHOLD_PCT, +) -> DataService: + """Build a service containing the target and five cutoff-aware input features.""" + raw_ipman = FREDAdapter(IPMAN_FRED_ID, cache_dir=cache_dir, refresh=refresh).fetch() + raw_fed_funds = FREDAdapter(FED_FUNDS_FRED_ID, cache_dir=cache_dir, refresh=refresh).fetch() + raw_treasury_10y = FREDAdapter(TREASURY_10Y_FRED_ID, cache_dir=cache_dir, refresh=refresh).fetch() + raw_treasury_2y = FREDAdapter(TREASURY_2Y_FRED_ID, cache_dir=cache_dir, refresh=refresh).fetch() + + ipman = apply_conservative_monthly_release_lag(raw_ipman, months=release_lag_months) + ipman_feature_frames = build_ipman_feature_frames(ipman) + macro_feature_frames = build_macro_feature_frames( + raw_fed_funds, + raw_treasury_10y, + raw_treasury_2y, + ) + stress = derive_manufacturing_stress_labels( + ipman, + lookback_months=stress_lookback_months, + threshold_pct=stress_threshold_pct, + ) + + service = DataService() + service.register( + IPMAN_SERIES_ID, + StaticFrameAdapter(ipman), + SeriesMetadata( + series_id=IPMAN_SERIES_ID, + description="U.S. manufacturing industrial production index (IPMAN)", + source="FRED (IPMAN)", + units="Index", + frequency="MS", + ), + ) + + for series_id, frame in ipman_feature_frames.items(): + periods = FEATURE_PERIODS[series_id] + service.register( + series_id, + StaticFrameAdapter(frame), + SeriesMetadata( + series_id=series_id, + description=f"Trailing {periods}-month percentage change in IPMAN", + source="Derived from FRED IPMAN", + units="Percent", + frequency="MS", + ), + ) + + service.register( + FED_FUNDS_SERIES_ID, + StaticFrameAdapter(macro_feature_frames[FED_FUNDS_SERIES_ID]), + SeriesMetadata( + series_id=FED_FUNDS_SERIES_ID, + description="Month-end effective federal funds rate", + source="FRED (DFF), derived monthly", + units="Percent", + frequency="MS", + ), + ) + service.register( + YIELD_CURVE_SERIES_ID, + StaticFrameAdapter(macro_feature_frames[YIELD_CURVE_SERIES_ID]), + SeriesMetadata( + series_id=YIELD_CURVE_SERIES_ID, + description="Month-end 10-year minus 2-year Treasury yield spread", + source="FRED (DGS10 minus DGS2), derived monthly", + units="Percentage points", + frequency="MS", + ), + ) + + service.register( + STRESS_SERIES_ID, + StaticFrameAdapter(stress), + SeriesMetadata( + series_id=STRESS_SERIES_ID, + description=( + "Binary U.S. manufacturing stress label: 1 when trailing " + f"{stress_lookback_months}-month IPMAN change is at or below {stress_threshold_pct:.1f}%" + ), + source="Derived from FRED IPMAN", + units="Binary event (0=no stress, 1=stress)", + frequency="MS", + ), + ) + return service + + +__all__ = [ + "DEFAULT_FRED_CACHE_DIR", + "FED_FUNDS_FRED_ID", + "IPMAN_FRED_ID", + "IPMAN_SERIES_ID", + "STRESS_SERIES_ID", + "TREASURY_10Y_FRED_ID", + "TREASURY_2Y_FRED_ID", + "build_manufacturing_stress_service", +] diff --git a/implementations/manufacturing_stress_forecasting/features.py b/implementations/manufacturing_stress_forecasting/features.py new file mode 100644 index 000000000..3599a2106 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/features.py @@ -0,0 +1,138 @@ +"""Leak-safe monthly features for the manufacturing-stress MVP.""" + +from __future__ import annotations + +from collections.abc import Sequence + +import pandas as pd +from aieng.forecasting.data.features import canonical_three_col + + +IPMAN_CHANGE_1M_SERIES_ID = "ipman_change_1m_pct" +IPMAN_CHANGE_3M_SERIES_ID = "ipman_change_3m_pct" +IPMAN_CHANGE_6M_SERIES_ID = "ipman_change_6m_pct" +FED_FUNDS_SERIES_ID = "fed_funds_rate_pct" +YIELD_CURVE_SERIES_ID = "treasury_10y_minus_2y_pct_points" + +FEATURE_PERIODS: dict[str, int] = { + IPMAN_CHANGE_1M_SERIES_ID: 1, + IPMAN_CHANGE_3M_SERIES_ID: 3, + IPMAN_CHANGE_6M_SERIES_ID: 6, +} +IPMAN_FEATURE_SERIES_IDS: tuple[str, ...] = tuple(FEATURE_PERIODS) +MACRO_FEATURE_SERIES_IDS: tuple[str, ...] = ( + FED_FUNDS_SERIES_ID, + YIELD_CURVE_SERIES_ID, +) +FEATURE_SERIES_IDS: tuple[str, ...] = IPMAN_FEATURE_SERIES_IDS + MACRO_FEATURE_SERIES_IDS + + +def apply_conservative_monthly_release_lag(frame: pd.DataFrame, months: int = 1) -> pd.DataFrame: + """Stamp observations as available ``months`` after their reference month. + + The standard FRED adapter uses ``released_at = timestamp`` because it does + not retrieve release vintages. For this monthly prototype, one month is a + deliberately conservative approximation that prevents a month-start + forecast from seeing that same month's completed production observation. + """ + if months < 0: + raise ValueError(f"months must be non-negative; got {months}") + out = frame.copy() + out["released_at"] = pd.to_datetime(out["timestamp"]) + pd.offsets.MonthBegin(months) + return canonical_three_col(out) + + +def percent_change_feature(ipman: pd.DataFrame, periods: int) -> pd.DataFrame: + """Return the trailing ``periods``-month IPMAN percentage change.""" + if periods < 1: + raise ValueError(f"periods must be positive; got {periods}") + out = ipman.copy().sort_values("timestamp").reset_index(drop=True) + out["value"] = out["value"].pct_change(periods=periods, fill_method=None) * 100.0 + return canonical_three_col(out) + + +def build_ipman_feature_frames(ipman: pd.DataFrame) -> dict[str, pd.DataFrame]: + """Build the three IPMAN momentum features used by the model.""" + return {series_id: percent_change_feature(ipman, periods) for series_id, periods in FEATURE_PERIODS.items()} + + +def monthly_last_observation(frame: pd.DataFrame) -> pd.DataFrame: + """Collapse a daily canonical series to its final observation each month.""" + out = canonical_three_col(frame) + out["month"] = out["timestamp"].dt.to_period("M") + out = out.sort_values(["month", "timestamp"]).groupby("month", as_index=False).tail(1) + out["timestamp"] = out["month"].dt.to_timestamp() + return canonical_three_col(out) + + +def build_macro_feature_frames( + fed_funds: pd.DataFrame, + treasury_10y: pd.DataFrame, + treasury_2y: pd.DataFrame, +) -> dict[str, pd.DataFrame]: + """Build monthly fed-funds and 10Y-minus-2Y rate features. + + Daily FRED observations are treated as available on the next business day. + The model then uses the final published observation associated with each + calendar month. This keeps the monthly feature panel small while preserving + an honest ``released_at`` cutoff. + """ + fed = canonical_three_col(fed_funds) + fed["released_at"] = fed["timestamp"] + pd.offsets.BDay(1) + fed_monthly = monthly_last_observation(fed) + + ten_year = canonical_three_col(treasury_10y).rename( + columns={"value": "value_10y", "released_at": "released_at_10y"} + ) + two_year = canonical_three_col(treasury_2y).rename(columns={"value": "value_2y", "released_at": "released_at_2y"}) + spread = pd.merge(ten_year, two_year, on="timestamp", how="inner") + spread["value"] = spread["value_10y"] - spread["value_2y"] + spread["released_at"] = spread[["released_at_10y", "released_at_2y"]].max(axis=1) + pd.offsets.BDay(1) + spread_monthly = monthly_last_observation(spread[["timestamp", "value", "released_at"]]) + + return { + FED_FUNDS_SERIES_ID: fed_monthly, + YIELD_CURVE_SERIES_ID: spread_monthly, + } + + +def build_feature_snapshot( + origin: pd.Timestamp, + feature_frames: dict[str, pd.DataFrame], + *, + series_ids: Sequence[str] = FEATURE_SERIES_IDS, +) -> dict[str, float] | None: + """Return the latest feature values that were published by ``origin``. + + Filtering on ``released_at`` here is essential when reconstructing older + training examples: the surrounding ``ForecastContext`` protects the current + forecast origin, while this function recreates the stricter cutoff at each + past origin used to train the fit-at-origin logistic model. + """ + snapshot: dict[str, float] = {} + for series_id in series_ids: + frame = feature_frames[series_id] + visible = frame[pd.to_datetime(frame["released_at"]) <= origin] + if visible.empty: + return None + snapshot[series_id] = float(visible.sort_values("timestamp")["value"].iloc[-1]) + return snapshot + + +__all__ = [ + "FED_FUNDS_SERIES_ID", + "FEATURE_PERIODS", + "FEATURE_SERIES_IDS", + "IPMAN_FEATURE_SERIES_IDS", + "IPMAN_CHANGE_1M_SERIES_ID", + "IPMAN_CHANGE_3M_SERIES_ID", + "IPMAN_CHANGE_6M_SERIES_ID", + "MACRO_FEATURE_SERIES_IDS", + "YIELD_CURVE_SERIES_ID", + "apply_conservative_monthly_release_lag", + "build_feature_snapshot", + "build_ipman_feature_frames", + "build_macro_feature_frames", + "monthly_last_observation", + "percent_change_feature", +] diff --git a/implementations/manufacturing_stress_forecasting/predictors/__init__.py b/implementations/manufacturing_stress_forecasting/predictors/__init__.py new file mode 100644 index 000000000..a9dad38c7 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/predictors/__init__.py @@ -0,0 +1,6 @@ +"""Predictors for the manufacturing-stress implementation.""" + +from manufacturing_stress_forecasting.predictors.logistic import ManufacturingStressLogisticPredictor + + +__all__ = ["ManufacturingStressLogisticPredictor"] diff --git a/implementations/manufacturing_stress_forecasting/predictors/logistic.py b/implementations/manufacturing_stress_forecasting/predictors/logistic.py new file mode 100644 index 000000000..bfea5de23 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/predictors/logistic.py @@ -0,0 +1,107 @@ +"""Five-variable logistic baseline for three-month-ahead manufacturing stress.""" + +from __future__ import annotations + +from datetime import datetime, timezone + +import numpy as np +import pandas as pd +from aieng.forecasting.data.context import ForecastContext +from aieng.forecasting.evaluation.prediction import BinaryForecast, Prediction +from aieng.forecasting.evaluation.predictor import Predictor +from aieng.forecasting.evaluation.task import ForecastingTask +from manufacturing_stress_forecasting.features import FEATURE_SERIES_IDS, build_feature_snapshot + + +class ManufacturingStressLogisticPredictor(Predictor): + """Forecast manufacturing stress from three IPMAN and two rate signals. + + The model is rebuilt at every backtest origin. For each resolved historical + outcome at month ``r``, its feature vector is reconstructed at ``r - lead`` + so the training examples obey the same three-month forecast horizon as the + current prediction. + """ + + def __init__(self, *, regularization_c: float = 1.0, min_training_examples: int = 24) -> None: + self._c = regularization_c + self._min_training_examples = min_training_examples + + @property + def predictor_id(self) -> str: + """Return the stable artifact identifier.""" + return "manufacturing_stress_logistic_ipman_rates" + + def predict(self, task: ForecastingTask, context: ForecastContext) -> list[Prediction]: + """Fit on visible history and return one binary stress probability.""" + if task.payload_type != "binary": + raise ValueError(f"{type(self).__name__} requires payload_type='binary'.") + if len(task.horizons) != 1: + raise ValueError(f"{type(self).__name__} supports exactly one horizon; got {task.horizons}.") + + as_of = pd.Timestamp(context.as_of) + target = context.get_series(task.target_series_id) + feature_frames = {series_id: context.get_series(series_id) for series_id in FEATURE_SERIES_IDS} + lead = pd.tseries.frequencies.to_offset(task.frequency) * task.horizons[0] + + rows, outcomes = self._training_data(target, feature_frames, lead) + current = build_feature_snapshot(as_of, feature_frames) + payload, model_metadata = self._fit_and_predict(rows, outcomes, current) + + return [ + Prediction( + predictor_id=self.predictor_id, + task_id=task.task_id, + issued_at=datetime.now(tz=timezone.utc).replace(tzinfo=None), + as_of=context.as_of, + forecast_date=(as_of + lead).to_pydatetime(), + payload=payload, + metadata={"n_train": len(outcomes), **model_metadata}, + ) + ] + + def _training_data( + self, + target: pd.DataFrame, + feature_frames: dict[str, pd.DataFrame], + lead: pd.DateOffset, + ) -> tuple[list[list[float]], list[float]]: + rows: list[list[float]] = [] + outcomes: list[float] = [] + for resolution_date, outcome in zip(target["timestamp"], target["value"], strict=True): + past_origin = pd.Timestamp(resolution_date) - lead + snapshot = build_feature_snapshot(past_origin, feature_frames) + if snapshot is None: + continue + rows.append([snapshot[series_id] for series_id in FEATURE_SERIES_IDS]) + outcomes.append(float(outcome)) + return rows, outcomes + + def _fit_and_predict( + self, + rows: list[list[float]], + outcomes: list[float], + current: dict[str, float] | None, + ) -> tuple[BinaryForecast, dict[str, object]]: + base_rate = float(np.mean(outcomes)) if outcomes else 0.1 + if current is None: + return BinaryForecast(probability=base_rate), {"model": "base_rate_fallback"} + if len(outcomes) < self._min_training_examples or len(set(outcomes)) < 2: + return BinaryForecast(probability=base_rate), {"model": "base_rate_fallback"} + + from sklearn.linear_model import LogisticRegression # noqa: PLC0415 + from sklearn.pipeline import make_pipeline # noqa: PLC0415 + from sklearn.preprocessing import StandardScaler # noqa: PLC0415 + + model = make_pipeline(StandardScaler(), LogisticRegression(C=self._c, max_iter=1000)) + model.fit(np.asarray(rows), np.asarray(outcomes)) + current_row = np.asarray([[current[series_id] for series_id in FEATURE_SERIES_IDS]]) + probability = float(model.predict_proba(current_row)[0, 1]) + coefficients = model.named_steps["logisticregression"].coef_[0] + return BinaryForecast(probability=probability), { + "model": "logistic_regression", + "features": dict(zip(FEATURE_SERIES_IDS, (float(value) for value in current_row[0]), strict=True)), + "coefficients": dict(zip(FEATURE_SERIES_IDS, (float(value) for value in coefficients), strict=True)), + } + + +__all__ = ["ManufacturingStressLogisticPredictor"] diff --git a/implementations/manufacturing_stress_forecasting/run_agent_prediction.py b/implementations/manufacturing_stress_forecasting/run_agent_prediction.py new file mode 100644 index 000000000..5a7245df1 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/run_agent_prediction.py @@ -0,0 +1,54 @@ +"""Run one current, quantitative-only manufacturing-stress agent forecast.""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pandas as pd +import yaml +from aieng.forecasting.evaluation import BacktestSpec +from aieng.forecasting.methods import HistoricalFrequencyPredictor +from manufacturing_stress_forecasting.analyst_agent import build_manufacturing_stress_agent_predictor +from manufacturing_stress_forecasting.data import ( + IPMAN_SERIES_ID, + build_manufacturing_stress_service, +) +from manufacturing_stress_forecasting.predictors import ManufacturingStressLogisticPredictor + + +SPEC_PATH = Path(__file__).resolve().parent / "specs" / "manufacturing_stress_smoke.yaml" + + +def main() -> None: + """Forecast from the most recent cached IPMAN release date.""" + with SPEC_PATH.open() as file: + task = BacktestSpec.model_validate(yaml.safe_load(file)).task + + service = build_manufacturing_stress_service() + full_ipman = service.get_series(IPMAN_SERIES_ID, as_of=pd.Timestamp("2100-01-01").to_pydatetime()) + as_of = pd.Timestamp(full_ipman["released_at"].max()) + context = service.context(as_of=as_of.to_pydatetime()) + + predictors = [ + HistoricalFrequencyPredictor(), + ManufacturingStressLogisticPredictor(), + build_manufacturing_stress_agent_predictor(), + ] + print(f"Forecast origin: {as_of.date()}") + print( + f"Latest visible IPMAN reference month: {pd.Timestamp(context.get_series(IPMAN_SERIES_ID)['timestamp'].max()).date()}" + ) + for predictor in predictors: + prediction = predictor.predict(task, context)[0] + output = { + "predictor_id": prediction.predictor_id, + "forecast_date": str(pd.Timestamp(prediction.forecast_date).date()), + "stress_probability": prediction.payload.probability, + "metadata": prediction.metadata, + } + print(json.dumps(output, indent=2)) + + +if __name__ == "__main__": + main() diff --git a/implementations/manufacturing_stress_forecasting/run_smoke.py b/implementations/manufacturing_stress_forecasting/run_smoke.py new file mode 100644 index 000000000..fdef561b8 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/run_smoke.py @@ -0,0 +1,28 @@ +"""Run the two-predictor manufacturing-stress smoke backtest.""" + +from pathlib import Path + +import yaml +from aieng.forecasting.evaluation import BacktestSpec, backtest +from aieng.forecasting.methods import HistoricalFrequencyPredictor +from manufacturing_stress_forecasting.data import build_manufacturing_stress_service +from manufacturing_stress_forecasting.predictors import ManufacturingStressLogisticPredictor + + +SPEC_PATH = Path(__file__).resolve().parent / "specs" / "manufacturing_stress_smoke.yaml" + + +def main() -> None: + """Load data, run both baselines, and print their mean Brier scores.""" + with SPEC_PATH.open() as file: + spec = BacktestSpec.model_validate(yaml.safe_load(file)) + + service = build_manufacturing_stress_service() + predictors = [HistoricalFrequencyPredictor(), ManufacturingStressLogisticPredictor()] + for predictor in predictors: + result = backtest(predictor=predictor, spec=spec, data_service=service) + print(f"{predictor.predictor_id}: {result.mean_score:.4f} mean {result.metric}") + + +if __name__ == "__main__": + main() diff --git a/implementations/manufacturing_stress_forecasting/specs/manufacturing_stress_smoke.yaml b/implementations/manufacturing_stress_forecasting/specs/manufacturing_stress_smoke.yaml new file mode 100644 index 000000000..39cc9e5cc --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/specs/manufacturing_stress_smoke.yaml @@ -0,0 +1,21 @@ +# Initial Track 1 experiment: one binary target, one three-month horizon. + +description: >- + Development backtest for forecasting whether U.S. manufacturing will be + under IPMAN-defined stress three months after each forecast origin. + +task: + task_id: manufacturing_stress_3m + target_series_id: manufacturing_stress + horizons: [3] + frequency: MS + payload_type: binary + description: >- + Probability that U.S. manufacturing will be under stress three months + ahead. A resolved month is stressed when IPMAN has declined by at least + 2 percent over its preceding three months. + +start: "2018-01-01" +end: "2024-12-01" +stride: 3 +warmup: 60 diff --git a/implementations/manufacturing_stress_forecasting/targets.py b/implementations/manufacturing_stress_forecasting/targets.py new file mode 100644 index 000000000..d2db91aa7 --- /dev/null +++ b/implementations/manufacturing_stress_forecasting/targets.py @@ -0,0 +1,42 @@ +"""Deterministic manufacturing-stress target construction.""" + +from __future__ import annotations + +import pandas as pd +from aieng.forecasting.data.features import canonical_three_col + + +DEFAULT_LOOKBACK_MONTHS = 3 +DEFAULT_STRESS_THRESHOLD_PCT = -2.0 + + +def derive_manufacturing_stress_labels( + ipman: pd.DataFrame, + *, + lookback_months: int = DEFAULT_LOOKBACK_MONTHS, + threshold_pct: float = DEFAULT_STRESS_THRESHOLD_PCT, +) -> pd.DataFrame: + """Create a monthly 0/1 target from trailing IPMAN deterioration. + + A month is labelled stressed when IPMAN has fallen by at least + ``abs(threshold_pct)`` percent over the preceding ``lookback_months``. + The label inherits the current IPMAN observation's ``released_at`` date, + because it cannot be known before that observation is published. + """ + if lookback_months < 1: + raise ValueError(f"lookback_months must be positive; got {lookback_months}") + if threshold_pct >= 0: + raise ValueError(f"threshold_pct must be negative; got {threshold_pct}") + + out = ipman.copy().sort_values("timestamp").reset_index(drop=True) + deterioration = out["value"].pct_change(periods=lookback_months, fill_method=None) * 100.0 + out["value"] = (deterioration <= threshold_pct).astype(float) + out.loc[deterioration.isna(), "value"] = float("nan") + return canonical_three_col(out) + + +__all__ = [ + "DEFAULT_LOOKBACK_MONTHS", + "DEFAULT_STRESS_THRESHOLD_PCT", + "derive_manufacturing_stress_labels", +] diff --git a/implementations/tests/manufacturing_stress_forecasting/test_agent_prompt.py b/implementations/tests/manufacturing_stress_forecasting/test_agent_prompt.py new file mode 100644 index 000000000..d3220fe39 --- /dev/null +++ b/implementations/tests/manufacturing_stress_forecasting/test_agent_prompt.py @@ -0,0 +1,74 @@ +"""Tests for the manufacturing-stress agent payload.""" + +import json +from datetime import datetime + +import pandas as pd +from aieng.forecasting.data import DataService, SeriesMetadata +from aieng.forecasting.data.features import StaticFrameAdapter +from aieng.forecasting.evaluation import ForecastingTask +from manufacturing_stress_forecasting.analyst_agent import ManufacturingStressPromptBuilder +from manufacturing_stress_forecasting.data import IPMAN_SERIES_ID, STRESS_SERIES_ID +from manufacturing_stress_forecasting.features import FEATURE_SERIES_IDS + + +def _frame(values: list[float], *, future_value: float | None = None) -> pd.DataFrame: + dates = pd.date_range("2020-01-01", periods=len(values), freq="MS") + frame = pd.DataFrame({"timestamp": dates, "value": values, "released_at": dates}) + if future_value is not None: + frame = pd.concat( + [ + frame, + pd.DataFrame( + { + "timestamp": [pd.Timestamp("2020-07-01")], + "value": [future_value], + "released_at": [pd.Timestamp("2020-08-01")], + } + ), + ], + ignore_index=True, + ) + return frame + + +def test_prompt_uses_only_cutoff_visible_evidence() -> None: + service = DataService() + metadata = lambda series_id: SeriesMetadata( # noqa: E731 + series_id=series_id, + description=series_id, + source="test", + units="test", + frequency="MS", + ) + service.register( + IPMAN_SERIES_ID, + StaticFrameAdapter(_frame([100, 101, 102, 103, 104, 105], future_value=999)), + metadata(IPMAN_SERIES_ID), + ) + service.register( + STRESS_SERIES_ID, StaticFrameAdapter(_frame([0, 0, 1, 0, 0, 0], future_value=1)), metadata(STRESS_SERIES_ID) + ) + for index, series_id in enumerate(FEATURE_SERIES_IDS): + service.register( + series_id, StaticFrameAdapter(_frame([float(index)] * 6, future_value=999)), metadata(series_id) + ) + + task = ForecastingTask( + task_id="manufacturing_stress_3m", + target_series_id=STRESS_SERIES_ID, + horizons=[3], + frequency="MS", + payload_type="binary", + description="Will manufacturing be stressed three months ahead?", + ) + prompt = ManufacturingStressPromptBuilder()(task=task, context=service.context(datetime(2020, 6, 1))) + payload = json.loads(prompt) + + assert payload["as_of"] == "2020-06-01" + assert payload["forecast_date"] == "2020-09-01" + assert payload["recent_ipman"][-1]["value"] == 105.0 + assert len(payload["current_ipman_signals_pct"]) == 3 + assert len(payload["current_macro_signals"]) == 2 + assert 999.0 not in payload["current_ipman_signals_pct"].values() + assert 999.0 not in payload["current_macro_signals"].values() diff --git a/implementations/tests/manufacturing_stress_forecasting/test_targets_and_features.py b/implementations/tests/manufacturing_stress_forecasting/test_targets_and_features.py new file mode 100644 index 000000000..3b94d08ea --- /dev/null +++ b/implementations/tests/manufacturing_stress_forecasting/test_targets_and_features.py @@ -0,0 +1,76 @@ +"""Focused tests for label construction and historical feature cutoffs.""" + +import pandas as pd +import pytest +from manufacturing_stress_forecasting.features import ( + FEATURE_SERIES_IDS, + FED_FUNDS_SERIES_ID, + YIELD_CURVE_SERIES_ID, + build_feature_snapshot, + build_macro_feature_frames, +) +from manufacturing_stress_forecasting.targets import derive_manufacturing_stress_labels + + +def test_stress_label_uses_trailing_three_month_decline() -> None: + dates = pd.date_range("2024-01-01", periods=5, freq="MS") + ipman = pd.DataFrame( + { + "timestamp": dates, + "value": [100.0, 100.0, 100.0, 100.0, 97.9], + "released_at": dates + pd.offsets.MonthBegin(1), + } + ) + + labels = derive_manufacturing_stress_labels(ipman) + + assert labels["value"].tolist() == [0.0, 1.0] + assert labels["timestamp"].tolist() == [pd.Timestamp("2024-04-01"), pd.Timestamp("2024-05-01")] + assert labels["released_at"].tolist() == [pd.Timestamp("2024-05-01"), pd.Timestamp("2024-06-01")] + + +def test_feature_snapshot_ignores_values_released_after_origin() -> None: + origin = pd.Timestamp("2024-03-01") + frames: dict[str, pd.DataFrame] = {} + for index, series_id in enumerate(FEATURE_SERIES_IDS): + frames[series_id] = pd.DataFrame( + { + "timestamp": [pd.Timestamp("2024-01-01"), pd.Timestamp("2024-02-01")], + "value": [float(index), 999.0], + "released_at": [pd.Timestamp("2024-02-01"), pd.Timestamp("2024-04-01")], + } + ) + + snapshot = build_feature_snapshot(origin, frames) + + assert snapshot is not None + for index, series_id in enumerate(FEATURE_SERIES_IDS): + assert snapshot[series_id] == pytest.approx(float(index)) + + +def test_macro_features_use_month_end_values_with_next_business_day_release() -> None: + timestamps = pd.to_datetime(["2024-01-30", "2024-01-31", "2024-02-28", "2024-02-29"]) + + def frame(values: list[float]) -> pd.DataFrame: + return pd.DataFrame( + { + "timestamp": timestamps, + "value": values, + "released_at": timestamps, + } + ) + + features = build_macro_feature_frames( + frame([5.30, 5.31, 5.32, 5.33]), + frame([4.00, 4.10, 4.20, 4.30]), + frame([4.40, 4.50, 4.55, 4.60]), + ) + + fed = features[FED_FUNDS_SERIES_ID] + spread = features[YIELD_CURVE_SERIES_ID] + + assert fed["timestamp"].tolist() == [pd.Timestamp("2024-01-01"), pd.Timestamp("2024-02-01")] + assert fed["value"].tolist() == pytest.approx([5.31, 5.33]) + assert fed["released_at"].tolist() == [pd.Timestamp("2024-02-01"), pd.Timestamp("2024-03-01")] + assert spread["value"].tolist() == pytest.approx([-0.40, -0.30]) + assert spread["released_at"].tolist() == [pd.Timestamp("2024-02-01"), pd.Timestamp("2024-03-01")] diff --git a/scripts/fetch_manufacturing_stress.py b/scripts/fetch_manufacturing_stress.py new file mode 100644 index 000000000..240114b33 --- /dev/null +++ b/scripts/fetch_manufacturing_stress.py @@ -0,0 +1,28 @@ +"""Fetch/cache IPMAN and rate inputs, then print the registered-series summary.""" + +from __future__ import annotations + +import sys +from pathlib import Path + + +REPO_ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(REPO_ROOT)) +sys.path.insert(0, str(REPO_ROOT / "implementations")) + +from dotenv import load_dotenv + + +load_dotenv(REPO_ROOT / ".env", override=False) + +from manufacturing_stress_forecasting.data import build_manufacturing_stress_service + + +def main() -> None: + """Populate the IPMAN/rates cache and report the derived series.""" + service = build_manufacturing_stress_service() + print(service.summary().to_string(index=False)) + + +if __name__ == "__main__": + main() diff --git a/uv.lock b/uv.lock index 7314a4a59..a1ba2aec6 100644 --- a/uv.lock +++ b/uv.lock @@ -1,5 +1,5 @@ version = 1 -revision = 3 +revision = 2 requires-python = ">=3.12, <4.0" resolution-markers = [ "python_full_version >= '3.15' and sys_platform == 'win32'", @@ -70,7 +70,7 @@ dev = [ [package.metadata] requires-dist = [ - { name = "agentic-forecasting-implementations", editable = "implementations" }, + { name = "agentic-forecasting-implementations", virtual = "implementations" }, { name = "ipykernel", specifier = ">=7.3.0" }, { name = "jupyter", specifier = ">=1.1.1" }, { name = "lightgbm", specifier = ">=4.0.0" }, @@ -97,7 +97,7 @@ dev = [ [[package]] name = "agentic-forecasting-implementations" version = "0.1.0" -source = { editable = "implementations" } +source = { virtual = "implementations" } dependencies = [ { name = "aieng-forecasting", extra = ["agentic", "documents", "llm", "numerical"] }, { name = "beautifulsoup4" },