diff --git a/.env.sample b/.env.sample index 896f455..167b6da 100644 --- a/.env.sample +++ b/.env.sample @@ -4,7 +4,7 @@ DATAROBOT_ENDPOINT=https://app.datarobot.com/api/v2 DATAROBOT_API_TOKEN= # JointFM compatibility metadata for the selected deployment. -JOINTFM_SCHEMA_VERSION=v4 +JOINTFM_SCHEMA_VERSION=v5 # Optional drift-detection pin; the SDK discovers the model version from /healthz when unset: # JOINTFM_MODEL_VERSION=jointfm-inference:0.2.0+ckpt.fin-2026-05-22 diff --git a/README.md b/README.md index 2b1e7f0..658e355 100644 --- a/README.md +++ b/README.md @@ -10,7 +10,7 @@ The SDK targets the DataRobot-hosted unstructured prediction route and the same - Import namespace: `jointfm_client` - Supported Python: `>=3.11` - Current SDK package version: `0.8.0` -- Current JointFM service schema: `schema_version="v4"` +- Current JointFM service schema: `schema_version="v5"` The public API shape is a synchronous low-level `JointFMClient` with `health()`, `health_instances()`, and `predict(payload)` methods plus high-level `forecast(...)`, `forecast_mean(...)`, `forecast_samples(...)`, and `forecast_quantiles(...)` helpers. The SDK is not a proxy service; callers use it as a local Python library that talks to the hosted or local JointFM endpoint. @@ -50,7 +50,7 @@ Example deployment configuration: deployment: datarobot_endpoint: https://app.datarobot.com/api/v2 datarobot_api_token: - schema_version: v4 + schema_version: v5 deployment_id: # Optional model-version pin; the SDK discovers it from /healthz when unset: # model_version: jointfm-inference:0.3.0+ckpt.fin-2026-05-22 @@ -68,7 +68,7 @@ Equivalent `.env` deployment configuration: ```dotenv DATAROBOT_ENDPOINT=https://app.datarobot.com/api/v2 DATAROBOT_API_TOKEN= -JOINTFM_SCHEMA_VERSION=v4 +JOINTFM_SCHEMA_VERSION=v5 JOINTFM_DEPLOYMENT_ID= # Optional drift-detection pin; the SDK discovers the model version from /healthz when unset: # JOINTFM_MODEL_VERSION=jointfm-inference:0.3.0+ckpt.fin-2026-05-22 @@ -78,7 +78,7 @@ Equivalent local REST configuration for a service started from the `joint` repos ```dotenv JOINTFM_LOCAL_BASE_URL=http://127.0.0.1:8080 -JOINTFM_SCHEMA_VERSION=v4 +JOINTFM_SCHEMA_VERSION=v5 # Optional drift-detection pin; the SDK discovers the model version from /healthz when unset: # JOINTFM_MODEL_VERSION=jointfm-inference:0.3.0+ckpt.fin_i504_o63_f0_t10_h16l16_mam7_af_t3r1_cnn_k3l4_hpst_h16l2_studentt_m4cr2df8skew ``` @@ -167,51 +167,47 @@ The bootstrap helper resolves the nearest src-layout Python project root, switch The current forecast request contract is: -- `schema_version`: exactly `"v4"`, configured as `JOINTFM_SCHEMA_VERSION` for `from_env()` clients +- `schema_version`: exactly `"v5"`, configured as `JOINTFM_SCHEMA_VERSION` for `from_env()` clients - `model_version`: exact model version advertised by `/healthz` or otherwise selected by the caller. Optional for `from_env()` clients: when `JOINTFM_MODEL_VERSION` is unset the SDK reads it from `/healthz` on first use; when set it acts as a drift-detection pin -- `query_mode`: `"forecast"` for the unconditional forecast, or `"condition"` for a conditional query at one future position; the high-level helpers set it from whether a `condition` block was passed +- `query_mode`: `"forecast"` for the unconditional forecast, or `"condition"` for the forecast given conditions on some columns; the high-level helpers set it from whether a `condition` was passed - `return_mode`: one of `"mean"`, `"samples"`, or `"quantiles"` - `time_index_mode`: one of `"ordinal"`, `"continuous_float"`, or `"absolute_datetime"` - `time_column`: required for `"absolute_datetime"`, and used for ordered ordinal or continuous histories when supplied - `query_times`: non-empty future forecast times only - `requested_columns`: optional column names or integer column indices, with duplicates rejected - `n_samples`: positive sample count for sampled forecasts and quantile estimation. When `return_mode="samples"` exceeds the `max_sample_count` advertised by the deployment's health metadata, `forecast_samples(...)` splits the request into capped prediction batches up front and returns one merged `SampleForecastResult`. -- `condition`: required with `query_mode="condition"` and forbidden otherwise. A `ConditionBlock` naming one future position by its index into `query_times` and one condition per column, see below. +- `condition`: required with `query_mode="condition"` and forbidden otherwise. One `EqualityCondition` or `IntervalCondition`, or a list of them, each covering the future positions its `query_time_indices` names; sent on the wire as the `conditions` list, see below. ### Conditional Queries -The `condition` query mode asks for the model's joint distribution at one future position *given* something about some of its columns at that same position. Conditioning relates columns to each other within one position and never across horizons, so the block names the position once and the response describes that position alone: `outputs.query_times` carries exactly one entry however many `query_times` the request listed. +The `condition` query mode asks for the forecast *given* something about some of its columns. Each condition names the future positions it covers by index into `query_times` through `query_time_indices`, and covers every position when that is `None`. Two conditions on the same column must cover disjoint positions; conditions on different columns may share positions, which is how one position mixes both kinds: -A column carries at most one condition, of either kind: +- `EqualityCondition(column, value, query_time_indices=None)` pins the column to a finite value. It stays nameable in `requested_columns` and reads back the value the request supplied, so a scenario answer lines up column for column with an unconditioned one. +- `IntervalCondition(column, lower=None, upper=None, query_time_indices=None)` confines the column to a range; `None` leaves that side open, and at least one side must be bounded. The column stays readable, and what comes back is its distribution inside the range. -- `EqualityCondition(column, value)` pins the column to a finite value. It stays nameable in `requested_columns` and reads back the value the request supplied, so a scenario answer lines up column for column with an unconditioned one. -- `IntervalCondition(column, lower=None, upper=None)` confines the column to a range; `None` leaves that side open, and at least one side must be bounded. The column stays readable, and what comes back is its distribution inside the range. - -Every column without a condition is a read-out column, and at least one must remain. `requested_columns` chooses what the response carries independently of that, defaulting to every declared column in declared order. Pass the block to `forecast(...)`, `forecast_mean(...)`, `forecast_samples(...)`, or `forecast_quantiles(...)`: +The response answers every entry of `query_times`. A position no condition covers carries the unconditioned forecast, and that is exact: the model draws each future position from its own joint, independently of the others, so a condition relates columns to each other within a position and never reaches another one. At every covered position at least one column must stay unconditioned, because what the position reads out is the conditional distribution of those columns. `requested_columns` chooses what the response carries independently of that, defaulting to every declared column in declared order. Pass one condition or a list of them to `forecast(...)`, `forecast_mean(...)`, `forecast_samples(...)`, `forecast_quantiles(...)`, or `forecast_log_prob(...)`: ```python -from jointfm_client import ConditionBlock, EqualityCondition, IntervalCondition +from jointfm_client import EqualityCondition, IntervalCondition -block = ConditionBlock( - query_time_index=0, - conditions=[ - EqualityCondition(column="equity_index_level", value=4780.0), - IntervalCondition(column="treasury_10y_yield", lower=0.041, upper=0.045), - ], -) result = client.forecast_mean( history, query_times=query_times, requested_columns=["portfolio_nav", "realized_volatility"], columns=plan.columns, - condition=block, + condition=[ + EqualityCondition(column="equity_index_level", value=4780.0), + IntervalCondition( + column="treasury_10y_yield", lower=0.041, upper=0.045, query_time_indices=[0] + ), + ], ) print(result.plausibility) ``` -Whether a deployment can condition depends on the mounted checkpoint's head. `/healthz` advertises `condition` in `supported_query_modes` and the kinds it answers in `supported_condition_kinds` (empty when the mode is absent). The client checks that advertisement before sending, so a deployment that cannot condition is refused with `UnsupportedServiceContractError` rather than after a paid round trip; `require_condition_support(metadata, block)` exposes the same check. +Whether a deployment can condition depends on the mounted checkpoint's head. `/healthz` advertises `condition` in `supported_query_modes` and the kinds it answers in `supported_condition_kinds` (empty when the mode is absent). The client checks that advertisement before sending, so a deployment that cannot condition is refused with `UnsupportedServiceContractError` rather than after a paid round trip; `require_condition_support(metadata, condition)` exposes the same check. -A condition response carries a `plausibility` block: `equality_log_density` is the log density the model assigns to the pinned values and `region_log_probability` the log probability it gives the interval region, each `None` when the request carried no condition of that kind. They separate a confident answer from one conditioned on something the model finds implausible; the service reports them and never refuses on them. `diagnostics.condition_draws` counts the draws behind sampled outputs, and `diagnostics.interval_estimator` reports the numerical accounting (`points`, `effective_sample_size`) when more than one column carries an interval and the region probability had to be estimated. +A condition response carries a `plausibility` block: `equality_log_density` is the log density the model assigns to the pinned values and `region_log_probability` the log probability it gives the interval region *given the pinned values at the same position*, so the two add to the plausibility of the whole condition set, each `None` when the request carried no condition of that kind. Over several covered positions each is the sum of the per-position values, exact because positions are independent. They separate a confident answer from one conditioned on something the model finds implausible; the service reports them and never refuses on them. `diagnostics.condition_draws` counts the draws behind sampled outputs, and `diagnostics.interval_estimator` reports the numerical accounting (`points`, `effective_sample_size`) when some position bounds more than one column and its region probability had to be estimated; across several estimated positions it describes the one with the smallest effective sample size. Column descriptors support the server fields `name`, `modality`, `role`, `nullable`, `vocabulary_size`, `level_count`, `mapping`, `lower_bound`, `upper_bound`, `time_value_kind`, `time_value_scale_seconds`, `time_value_use_local_normalized_time`, `time_value_calendar_id`, and `time_value_timezone`. @@ -227,7 +223,7 @@ Successful forecast responses preserve `schema_version`, `image_version`, `model ```json { - "schema_version": "v4", + "schema_version": "v5", "errors": [ { "code": "VALIDATION_ERROR", @@ -242,7 +238,7 @@ Known error codes are `VALIDATION_ERROR`, `UNSUPPORTED_HEAD_QUERY_COMBINATION`, ## Compatibility Policy -The SDK supports only `schema_version="v4"`. `validate_service_metadata()` checks `/healthz` metadata and raises typed compatibility errors before prediction if the service advertises a different schema, an unexpected model version, mode capabilities outside the recorded service contract, or an unsupported `decoding_strategy`. Return modes and time-index modes must match the SDK's lists exactly. Query modes and condition kinds are derived by the service from the mounted head, so a deployment may advertise fewer of them than the SDK knows; it must advertise at least one query mode and nothing the SDK does not know. +The SDK supports only `schema_version="v5"`. `validate_service_metadata()` checks `/healthz` metadata and raises typed compatibility errors before prediction if the service advertises a different schema, an unexpected model version, mode capabilities outside the recorded service contract, or an unsupported `decoding_strategy`. Return modes and time-index modes must match the SDK's lists exactly. Query modes and condition kinds are derived by the service from the mounted head, so a deployment may advertise fewer of them than the SDK knows; it must advertise at least one query mode and nothing the SDK does not know. Callers should pass an expected `model_version` when they already know which deployment artifact they intend to use. A mismatch is treated as a hard compatibility error rather than silently downgrading, guessing, or retrying another model. @@ -284,7 +280,7 @@ Create `.env` from `.env.sample` or set the same values in your shell. A hosted ```dotenv DATAROBOT_ENDPOINT=https://app.datarobot.com/api/v2 DATAROBOT_API_TOKEN= -JOINTFM_SCHEMA_VERSION=v4 +JOINTFM_SCHEMA_VERSION=v5 JOINTFM_DEPLOYMENT_ID= # Or: JOINTFM_DEPLOYMENT_IDS=chevron-id,research-id # Optional drift-detection pin; the SDK discovers the model version from /healthz when unset: diff --git a/config.sample.yaml b/config.sample.yaml index a49d28e..dd45a56 100644 --- a/config.sample.yaml +++ b/config.sample.yaml @@ -52,7 +52,7 @@ transport: - X-DataRobot-Execution-ID user_agent_header: User-Agent forecast: - schema_version: v4 + schema_version: v5 query_mode: forecast return_mode: mean time_index_mode: ordinal diff --git a/docs/api-reference.md b/docs/api-reference.md index efc7f82..e434d76 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -6,7 +6,7 @@ This reference covers the supported public Python surface exported by `jointfm_c | Name | Purpose | | --- | --- | -| `JointFMClient` | Synchronous client for hosted or local JointFM endpoints. Use `from_env()` for `.env` and `config.yaml` backed hosted settings, `health()` for consensus typed service metadata, `health_instances()` for per-deployment probe results and pooled sample topology, `predict(payload)` for low-level JSON prediction, `forecast(...)` for validated tabular forecasts, and the `forecast_mean(...)`, `forecast_samples(...)`, `forecast_quantiles(...)`, and `forecast_log_prob(...)` convenience methods for typed forecast results. `forecast_log_prob(...)` is the one that asks about values the caller already holds: it takes `query_rows`, one observed row per entry of `query_times` carrying every declared column, and returns their log density under the model's joint. Each forecast method accepts `condition=ConditionBlock(...)`, which switches the request to the `condition` query mode and, before anything is sent, checks that the deployment's health metadata advertises the mode and every condition kind the block uses. `health()` probes `GET /healthz` for local deployments and POSTs `{"request_type": "health"}` to `predict_url` for hosted DataRobot deployments because the DataRobot deployment gateway only proxies the unstructured prediction route. `feature_importance(...)` runs permutation feature importance: one baseline `forecast_samples` call plus one per shuffled feature column, returning a list of `{"feature", "mean", "distance"}` dicts, each holding that feature's absolute forecast-mean shift and centered squared 2-Wasserstein distance indexed by target and horizon. | +| `JointFMClient` | Synchronous client for hosted or local JointFM endpoints. Use `from_env()` for `.env` and `config.yaml` backed hosted settings, `health()` for consensus typed service metadata, `health_instances()` for per-deployment probe results and pooled sample topology, `predict(payload)` for low-level JSON prediction, `forecast(...)` for validated tabular forecasts, and the `forecast_mean(...)`, `forecast_samples(...)`, `forecast_quantiles(...)`, and `forecast_log_prob(...)` convenience methods for typed forecast results. `forecast_log_prob(...)` is the one that asks about values the caller already holds: it takes `query_rows`, one observed row per entry of `query_times` carrying every declared column, and returns their log density under the model's joint. Each forecast method accepts `condition=`, one `EqualityCondition` or `IntervalCondition` or a list of them, which switches the request to the `condition` query mode and, before anything is sent, checks that the deployment's health metadata advertises the mode and every condition kind the conditions use. `health()` probes `GET /healthz` for local deployments and POSTs `{"request_type": "health"}` to `predict_url` for hosted DataRobot deployments because the DataRobot deployment gateway only proxies the unstructured prediction route. `feature_importance(...)` runs permutation feature importance: one baseline `forecast_samples` call plus one per shuffled feature column, returning a list of `{"feature", "mean", "distance"}` dicts, each holding that feature's absolute forecast-mean shift and centered squared 2-Wasserstein distance indexed by target and horizon. | `JointFMClient.from_env()` loads `config.yaml`, optional `.env` values, and process environment variables. `JointFMClient.health(cache=True)` caches health metadata only when requested. `JointFMClient.health_instances()` returns the same probe as a `HealthInstances` object: one `InstanceHealth` per configured deployment (including failures), `max_sample_count` as the sum of reachable caps (overall parallel capacity), and `topology` / `topology_label` grouping those caps (unavailable peers are listed but excluded from the sum and topology). Each endpoint's health payload describes only that endpoint; the client aggregates by calling each configured peer. `health()` still exposes the minimum reachable `max_sample_count`, which is the sample-batch cap used by forecast helpers. `JointFMClient.predict(payload)` requires `payload["model_version"]`; high-level forecast helpers resolve the configured model version when the caller does not pass one explicitly. When `forecast_samples(...)` requests an explicit `n_samples`, the client learns the deployment's `max_sample_count` from health metadata before the first prediction, splits oversized requests into capped prediction batches, and returns one merged `SampleForecastResult`. Clients configured without a reachable health route fall back to discovering the cap from the structured service error. @@ -17,12 +17,12 @@ This reference covers the supported public Python surface exported by `jointfm_c | `ColumnSpec` | Describes one modeled request column. Fields are `name`, `modality`, `role`, `nullable`, `vocabulary_size`, `level_count`, `mapping`, `lower_bound`, `upper_bound`, `time_value_kind`, `time_value_scale_seconds`, `time_value_use_local_normalized_time`, `time_value_calendar_id`, and `time_value_timezone`. | | `DataFrameSchema` | Describes tabular history layout. Fields are `columns`, `time_index_mode`, `time_column`, `time_scale_seconds`, `use_local_normalized_time`, `calendar_id`, and `timezone`. | | `ForecastRequestMetadata` | Holds `schema_version`, `model_version`, `query_mode` (`forecast` or `condition`), and `return_mode` for one forecast request. | -| `ForecastRequest` | Validated request object that combines metadata, schema, history rows, query times, requested columns, sample or quantile controls, `seed`, an optional `condition` block, and the optional `query_rows` a scored request supplies, then emits a JSON-compatible payload with `to_payload()`. `query_mode="condition"` and a `condition` block must appear together, as must `return_mode="log_prob"` and `query_rows`, and both are validated against the request's own schema and `query_times` before any round trip. | -| `EqualityCondition` | One column of the conditioned position pinned to a finite `value`. The column may be named in `requested_columns` like any other, and reads back the pinned value the request supplied. | -| `IntervalCondition` | One column of the conditioned position confined to `[lower, upper]`; `None` leaves a side open, at least one side must be bounded, and `lower < upper`. The column stays readable and the response describes its distribution inside the range. | -| `ConditionBlock` | Every condition of one request: `query_time_index` into the request's `query_times` and a sequence of `EqualityCondition` or `IntervalCondition` values, at most one per column. `kinds`, `pinned_columns`, and `conditioned_columns` expose what the capability gate and the request validation need. | -| `ConditionPlausibility` | What the model thinks of the conditions it was given: `equality_log_density` of the pinned values and `region_log_probability` of the interval region, each `None` when the request carried no condition of that kind. Reported by the service and never refused on. | -| `IntervalEstimator` | Numerical accounting (`points`, `effective_sample_size`) behind a region probability that had to be estimated, which happens when more than one column carries an interval condition. | +| `ForecastRequest` | Validated request object that combines metadata, schema, history rows, query times, requested columns, sample or quantile controls, `seed`, an optional `condition` (one condition or a list), and the optional `query_rows` a scored request supplies, then emits a JSON-compatible payload with `to_payload()`. `query_mode="condition"` and a `condition` must appear together, as must `return_mode="log_prob"` and `query_rows`, and both are validated against the request's own schema and `query_times` before any round trip. | +| `EqualityCondition` | One column pinned to a finite `value` at the positions `query_time_indices` names (indices into `query_times`, `None` for every position). The column may be named in `requested_columns` like any other, and reads back the pinned value the request supplied. | +| `IntervalCondition` | One column confined to `[lower, upper]` at the positions `query_time_indices` names (`None` for every position); `None` leaves a side open, at least one side must be bounded, and `lower < upper`. The column stays readable and the response describes its distribution inside the range. | +| `Condition` | Type alias for `EqualityCondition \| IntervalCondition`. Every `condition=` parameter takes one of them or a sequence of them. | +| `ConditionPlausibility` | What the model thinks of the conditions it was given: `equality_log_density` of the pinned values and `region_log_probability` of the interval region given the pinned values at the same position, each `None` when the request carried no condition of that kind. Reported by the service and never refused on. | +| `IntervalEstimator` | Numerical accounting (`points`, `effective_sample_size`) behind a region probability that had to be estimated, which happens when some position bounds more than one column. Across several estimated positions it describes the one with the smallest effective sample size. | | `HealthMetadata` | Typed service-health payload with service status, schema and model versions, checkpoint metadata, device, head, `decoding_strategy`, advertised query modes, `supported_condition_kinds` (empty when the deployment cannot condition), return modes, time-index modes, time-index encoding, `max_sample_count`, and an optional `data_generation` block carrying advertised capacity limits. The container exposes it on `GET /healthz` for direct local access and as the response to `POST {"request_type": "health"}` on the unstructured prediction route for DataRobot-hosted deployments. Each endpoint reports only its own capabilities. | | `InstanceHealth` | One configured deployment's probe outcome: `deployment_id`, optional `metadata` (`HealthMetadata` when reachable), and optional `error` when the peer was skipped. | | `HealthInstances` | Client aggregation of `health_instances()`: `instances` (one `InstanceHealth` per configured ID), `max_sample_count` (sum of reachable caps = overall parallel capacity), `topology` as `(count, cap)` pairs sorted by descending cap, and `topology_label` such as `2x5000` or `1x7000, 1x3000`. Unavailable peers stay in `instances` but are omitted from the sum and topology. | @@ -80,9 +80,9 @@ All SDK-specific exceptions inherit from `JointFMError`. | `JointFMHTTPStatusError` | The service returns an HTTP error status. | | `JointFMServiceError` | A response body contains non-empty JointFM `errors`, including the case where HTTP status unexpectedly succeeded. | | `JointFMCompatibilityError` | Base class for fail-fast service compatibility failures. | -| `UnsupportedSchemaVersionError` | The service or response advertises a schema version other than `v4`. | +| `UnsupportedSchemaVersionError` | The service or response advertises a schema version other than `v5`. | | `UnsupportedModelVersionError` | The service or response model version differs from the configured or requested version. | -| `UnsupportedServiceContractError` | The service-health payload advertises mode capabilities or a `decoding_strategy` outside the recorded service contract, or a condition request targets a deployment that does not advertise the `condition` mode or one of the block's condition kinds. | +| `UnsupportedServiceContractError` | The service-health payload advertises mode capabilities or a `decoding_strategy` outside the recorded service contract, or a condition request targets a deployment that does not advertise the `condition` mode or one of the request's condition kinds. | ## Public Functions @@ -106,7 +106,8 @@ All SDK-specific exceptions inherit from `JointFMError`. | `build_datarobot_prediction_headers(api_token)` | Build hosted prediction headers: bearer authorization, broad accept header, and JSON content type. | | `build_forecast_payload(...)` | Build a validated JSON-compatible forecast payload from explicit schema, history rows, query times, and return-mode controls. | | `validate_service_metadata(metadata, expected_model_version=None)` | Validate the service-health metadata against the supported schema version, the expected model when supplied, advertised mode capabilities, and a supported `decoding_strategy`. Return and time-index modes must match the SDK's lists exactly; `supported_query_modes` (non-empty) and `supported_condition_kinds` (may be empty) must be subsets of what the SDK knows, because the service derives them from the mounted head. | -| `require_condition_support(metadata, block)` | Raise `UnsupportedServiceContractError` when `HealthMetadata` does not advertise the `condition` query mode or one of the kinds the `ConditionBlock` uses. The forecast helpers call it before sending a condition request. | +| `require_condition_support(metadata, condition)` | Raise `UnsupportedServiceContractError` when `HealthMetadata` does not advertise the `condition` query mode or one of the kinds `condition` uses. The forecast helpers call it before sending a condition request. | +| `resolve_conditions(condition)` | Normalize one condition or a sequence of them into a tuple, raising `ValueError` for an empty sequence, an entry that is not a condition, or two conditions on one column at a shared position (`query_time_indices=None` overlaps every other condition on its column). | | `infer_column_specs_from_dataframe(frame, ...)` | Infer ordered `ColumnSpec` objects from a pandas `DataFrame` and explicit role, modality, mapping, nullability, time-value, and bounds hints. | | `dataframe_to_history_rows(frame, schema)` | Convert a pandas `DataFrame` into server-compatible `history_rows`. | | `arrays_to_history_rows(values, columns=..., ...)` | Convert a two-dimensional NumPy-like array plus column metadata into `history_rows`. | @@ -126,7 +127,7 @@ All SDK-specific exceptions inherit from `JointFMError`. | --- | --- | --- | | `DATAROBOT_ENDPOINT` | Hosted calls | HTTPS DataRobot API v2 endpoint, normalized without a trailing slash and required to end in `/api/v2`. | | `DATAROBOT_API_TOKEN` | Hosted calls | Non-empty, whitespace-free API token used in the hosted bearer authorization header. | -| `JOINTFM_SCHEMA_VERSION` | Hosted calls | Request schema pin. The SDK supports only `v4`. | +| `JOINTFM_SCHEMA_VERSION` | Hosted calls | Request schema pin. The SDK supports only `v5`. | | `JOINTFM_MODEL_VERSION` | Hosted calls | Exact JointFM deployment model version expected from the service-health payload and prediction responses. | | `JOINTFM_DEPLOYMENT_ID` | One selector | Deployment ID used to build hosted health and prediction URLs. | | `JOINTFM_DEPLOYMENT_IDS` | One selector | Comma-separated hosted deployment IDs for round-robin load balancing (at least two unique IDs). Mutually exclusive with other selectors. Peers must share `model_version` and `checkpoint_version`. `health()` uses the minimum reachable `max_sample_count` as the sample-batch cap; `health_instances()` sums reachable caps for overall parallel capacity and reports topology. | @@ -157,9 +158,9 @@ The string literals are exposed as `PREDICT_REQUEST_TYPE`, `HEALTH_REQUEST_TYPE` | Field | Required | Description | | --- | --- | --- | | `request_type` | Optional | One of `"predict"` (default) or `"health"`. Forecast requests omit this field or set it to `"predict"`. | -| `schema_version` | Yes | Must be `"v4"`. | +| `schema_version` | Yes | Must be `"v5"`. | | `model_version` | Yes | Exact deployed model version expected by the caller. | -| `query_mode` | Yes | `"forecast"` for the unconditional forecast or `"condition"` for a conditional query at one future position. | +| `query_mode` | Yes | `"forecast"` for the unconditional forecast or `"condition"` for the forecast given the request's `conditions`. | | `return_mode` | Yes | One of `"mean"`, `"samples"`, `"quantiles"`, or `"log_prob"`, covered by the `forecast_mean`, `forecast_samples`, `forecast_quantiles`, and `forecast_log_prob` helpers. | | `time_index_mode` | Yes | One of `"ordinal"`, `"continuous_float"`, or `"absolute_datetime"`. | | `columns` | Yes | Non-empty array of column descriptors for modeled columns. | @@ -171,7 +172,7 @@ The string literals are exposed as `PREDICT_REQUEST_TYPE`, `HEALTH_REQUEST_TYPE` | `n_samples` | Samples and quantiles controls | Positive sample count when sampling controls are needed. Oversized sample forecasts are batched automatically against the cap advertised in health metadata. | | `quantiles` | Quantiles mode | Quantile levels in `(0, 1)`, required for `return_mode="quantiles"`. | | `seed` | Optional | Integer random seed for reproducible stochastic outputs. | -| `condition` | With `query_mode="condition"` | Object with `query_time_index` (index into `query_times`) and `conditions`, a list of `{"column", "kind": "equality", "value"}` or `{"column", "kind": "interval", "lower", "upper"}` entries with `null` for an open bound. At most one condition per column and at least one column left unconditioned. Forbidden with any other query mode. | +| `conditions` | With `query_mode="condition"` | Non-empty list of `{"column", "kind": "equality", "value", "query_time_indices"}` or `{"column", "kind": "interval", "lower", "upper", "query_time_indices"}` entries, `null` marking an open bound. `query_time_indices` is a list of distinct indices into `query_times`, or `null` for every position. Two conditions on one column must cover disjoint positions, and every covered position must leave at least one column unconditioned. Forbidden with any other query mode. | | `time_scale_seconds` | Optional | Positive scale for continuous time indexes. | | `use_local_normalized_time` | Optional | Whether the service should use local normalized time features. | | `calendar_id` | Optional | Calendar identifier, defaulting to `pandas-default`. | @@ -200,14 +201,14 @@ The string literals are exposed as `PREDICT_REQUEST_TYPE`, `HEALTH_REQUEST_TYPE` | Field | Description | | --- | --- | -| `schema_version` | Response schema, expected to be `"v4"`. | +| `schema_version` | Response schema, expected to be `"v5"`. | | `image_version` | Service image version that produced the response. | | `model_version` | Model version that produced the response. | | `checkpoint_version` | Checkpoint version that produced the response. | | `head` | Forecast head used by the service. | | `query_mode` | Response query mode, matching the request: `"forecast"` or `"condition"`. | | `return_mode` | Response return mode matching the request. | -| `outputs.query_times` | Forecast horizon values preserved from the request. A condition response carries only the conditioned position. | +| `outputs.query_times` | Forecast horizon values preserved from the request, on condition responses too: a position no condition covers carries the unconditioned forecast. | | `outputs.requested_columns` | Output columns in response order. | | `outputs.mean` | Mean values with axis order `(horizon, column)` when `return_mode="mean"`. | | `outputs.samples` | Sample values with axis order `(sample, horizon, column)` when `return_mode="samples"`. | @@ -217,8 +218,8 @@ The string literals are exposed as `PREDICT_REQUEST_TYPE`, `HEALTH_REQUEST_TYPE` | `diagnostics.seed` | Optional seed used by the service. | | `outputs.log_prob` | `log_prob` responses only: per-horizon `values` and `nll_values` with their `total`, `mean`, `nll_total`, and `nll_mean` summaries. | | `diagnostics.condition_draws` | Condition responses only: number of draws behind sampled outputs. Batched sample requests report the merged count. | -| `diagnostics.interval_estimator` | Condition responses only, and only when more than one column carries an interval: `points` and `effective_sample_size` of the numerical region-probability estimate. | -| `plausibility` | `null` on forecast responses. On condition responses an object with `equality_log_density` (log density of the pinned values) and `region_log_probability` (log probability of the interval region), each `null` when the request carried no condition of that kind. | +| `diagnostics.interval_estimator` | Condition responses only, and only when some position bounds more than one column: `points` and `effective_sample_size` of the numerical region-probability estimate, describing the estimated position with the smallest effective sample size. | +| `plausibility` | `null` on forecast responses. On condition responses an object with `equality_log_density` (log density of the pinned values) and `region_log_probability` (log probability of the interval region, conditional on the pinned values at the same position), each `null` when the request carried no condition of that kind. Over several covered positions each is the sum of the per-position values. | | `errors` | Structured service errors. Non-empty arrays raise typed SDK exceptions. | ### Health Metadata @@ -226,7 +227,7 @@ The string literals are exposed as `PREDICT_REQUEST_TYPE`, `HEALTH_REQUEST_TYPE` | Field | Description | | --- | --- | | `status` | Service status string. | -| `schema_version` | Advertised schema version. The SDK requires `v4`. | +| `schema_version` | Advertised schema version. The SDK requires `v5`. | | `image_version` | Running service image version. | | `model_version` | Running model version. | | `checkpoint_version` | Loaded checkpoint version. | @@ -235,7 +236,7 @@ The string literals are exposed as `PREDICT_REQUEST_TYPE`, `HEALTH_REQUEST_TYPE` | `head` | Active forecast head. | | `decoding_strategy` | Horizon decoding mode advertised by the mounted model. Must be one of `SUPPORTED_DECODING_STRATEGIES`: `parallel_dense`, `parallel_scalable`, or `autoregressive`. Parallel strategies decode every horizon in one pass; `autoregressive` rolls horizons sequentially. | | `supported_query_modes` | Non-empty subset of the SDK's query modes (`forecast`, `condition`); the service derives it from the mounted head. | -| `supported_condition_kinds` | Subset of the SDK's condition kinds (`equality`, `interval`), empty when `condition` is not advertised. Condition requests are refused locally when the block uses a kind that is missing here. | +| `supported_condition_kinds` | Subset of the SDK's condition kinds (`equality`, `interval`), empty when `condition` is not advertised. Condition requests are refused locally when a condition uses a kind that is missing here. | | `supported_return_modes` | Must match the SDK's return modes (`mean`, `samples`, `quantiles`, `log_prob`). | | `supported_time_index_modes` | Must match the SDK's time-index modes. | | `time_index_encoding` | Time-index encoding advertised by the service. | diff --git a/notebooks/forecast_condition.ipynb b/notebooks/forecast_condition.ipynb index 4fac89e..7a13611 100644 --- a/notebooks/forecast_condition.ipynb +++ b/notebooks/forecast_condition.ipynb @@ -49,14 +49,14 @@ "# Conditional Forecast\n", "Forecast USD portfolio NAV and risk from 100 positive float daily observations: equity index level, 10-year Treasury yield, and one EUR/USD FX rate. Instead of the unconditional forecast, ask the model a *what-if* question about the first future step: what does it expect for portfolio NAV and realized volatility **given** something about the other columns at that same step.\n", "\n", - "The `condition` query mode answers this in closed form from the model's joint distribution at one future position. A request names that position once, by its index into `query_times`, and attaches one condition per column it wants to fix:\n", + "The `condition` query mode answers this in closed form from the model's joint distribution at each future position. Every condition names the positions it covers by index into `query_times` (`query_time_indices`), and covers every position when it names none. `condition=` takes one condition or a list of them; conditions on different columns may cover the same position, while two conditions on the same column must cover disjoint positions:\n", "\n", "- An **equality condition** pins a column to a value (`EqualityCondition`). The pinned column leaves the read-out set, because reading it back would only repeat the request.\n", "- An **interval condition** confines a column to a range whose bounds may be open on either side (`IntervalCondition`). The column stays readable: what comes back is its distribution inside the range.\n", "\n", "The cell below asks for the same forecast twice, once without the block and once with it, because a conditional mean is only readable next to the unconditional one: what the what-if is worth is the shift between them, not the level of either.\n", "\n", - "Every column without a condition is a read-out column, and the response describes the conditional distribution of those columns at the conditioned position only, so `outputs.query_times` has exactly one entry however many `query_times` the request carried.\n", + "Every column without a condition is a read-out column. The response answers every entry of `query_times`, and a position no condition covers carries the unconditioned forecast. That is exact rather than a placeholder: the model draws each future position from its own joint, independently of the others, so a condition at one step says nothing about another. The comparison below shows it — the shift is zero at every step except the conditioned one.\n", "\n", "Every read-out the service serves reads that same conditional, and the sections below take it in turn: as a mean beside the unconditional one, as coherent joint draws, as quantiles inside a bounded band, and as the log density of values you already hold. The last two sections use the plausibility numbers to rank candidate what-ifs against each other.\n", "\n", @@ -78,7 +78,6 @@ "import pandas as pd\n", "\n", "from jointfm_client import (\n", - " ConditionBlock,\n", " EqualityCondition,\n", " JointFMClient,\n", " plan_forecast_columns,\n", @@ -94,9 +93,7 @@ "EQUITY_RALLY = 1.02\n", "EXPECTED_COLUMNS = FEATURE_COLUMNS + TARGET_COLUMNS\n", "QUERY_TIMES = list(range(INPUT_STEPS, INPUT_STEPS + OUTPUT_HORIZONS))\n", - "# Every column a request may read back: an equality condition removes its own\n", - "# column from the read-out set, and nothing else does.\n", - "READABLE_COLUMNS = [name for name in EXPECTED_COLUMNS if name != PINNED_COLUMN]\n", + "CONDITIONED_TIME = QUERY_TIMES[CONDITIONED_STEP]\n", "\n", "history = pd.read_csv(HISTORY_PATH, dtype=float)\n", "if list(history.columns) != EXPECTED_COLUMNS:\n", @@ -125,9 +122,10 @@ "\n", "last_equity_level = float(history[PINNED_COLUMN].iloc[-1])\n", "rallied_equity_level = last_equity_level * EQUITY_RALLY\n", - "rally = ConditionBlock(\n", - " query_time_index=CONDITIONED_STEP,\n", - " conditions=[EqualityCondition(column=PINNED_COLUMN, value=rallied_equity_level)],\n", + "rally = EqualityCondition(\n", + " column=PINNED_COLUMN,\n", + " value=rallied_equity_level,\n", + " query_time_indices=[CONDITIONED_STEP],\n", ")\n", "baseline = client.forecast_mean(\n", " history,\n", @@ -144,22 +142,18 @@ " seed=7,\n", " condition=rally,\n", ")\n", - "if result.query_times != (QUERY_TIMES[CONDITIONED_STEP],):\n", - " raise ValueError(\n", - " f\"Expected the conditioned position alone, got {result.query_times!r}\"\n", - " )\n", + "if result.query_times != tuple(QUERY_TIMES):\n", + " raise ValueError(f\"Expected every query time, got {result.query_times!r}\")\n", "if result.plausibility is None:\n", " raise ValueError(\"A condition response must carry its plausibility block\")\n", "print(\"log density of the pinned value:\", result.plausibility.equality_log_density)\n", "forecast = result.to_pandas_tidy()\n", - "expected_forecast_rows = len(plan.requested_columns)\n", + "expected_forecast_rows = len(plan.requested_columns) * len(QUERY_TIMES)\n", "if len(forecast) != expected_forecast_rows:\n", " raise ValueError(\n", " f\"Expected {expected_forecast_rows} forecast rows, got {len(forecast)}\"\n", " )\n", "\n", - "# The conditional answers one position, so joining on it keeps exactly the rows\n", - "# the two requests have in common.\n", "comparison = baseline.to_pandas_tidy().merge(\n", " forecast,\n", " on=[\"query_time\", \"requested_column\"],\n", @@ -168,6 +162,12 @@ "comparison[\"shift\"] = (\n", " comparison[\"value_given_the_rally\"] - comparison[\"value_unconditional\"]\n", ")\n", + "# Positions are drawn independently, so the pin can only move its own step.\n", + "moved_elsewhere = comparison[\n", + " (comparison[\"query_time\"] != CONDITIONED_TIME) & (comparison[\"shift\"] != 0.0)\n", + "]\n", + "if not moved_elsewhere.empty:\n", + " raise ValueError(f\"A step the rally does not cover moved:\\n{moved_elsewhere}\")\n", "comparison" ] }, @@ -188,6 +188,55 @@ { "cell_type": "markdown", "id": "5", + "metadata": { + "id": "condition-every-step-description", + "language": "markdown" + }, + "source": [ + "## Conditioning every step\n", + "A condition that names no `query_time_indices` covers every entry of `query_times`. Below, the equity index is held at the rallied level at all ten steps, and each step answers under that pin.\n", + "\n", + "Read it as ten what-ifs, one per step, not as a path along which the index stays put. Each step is conditioned on its own joint, and nothing links one step's joint to the next, so a pin at step 3 cannot inform step 4 — the same independence that left the uncovered steps unchanged above. What does add up across steps is the plausibility: `equality_log_density` is the sum of the ten per-step log densities, exactly, because the steps are independent. The same rule reaches `region_log_probability`, which sums over the steps an interval covers." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "6", + "metadata": { + "id": "condition-every-step", + "language": "python" + }, + "outputs": [], + "source": [ + "held_rally = EqualityCondition(column=PINNED_COLUMN, value=rallied_equity_level)\n", + "held = client.forecast_mean(\n", + " history,\n", + " query_times=QUERY_TIMES,\n", + " requested_columns=plan.requested_columns,\n", + " columns=plan.columns,\n", + " seed=7,\n", + " condition=held_rally,\n", + ")\n", + "if held.query_times != tuple(QUERY_TIMES):\n", + " raise ValueError(f\"Expected every query time, got {held.query_times!r}\")\n", + "if held.plausibility is None:\n", + " raise ValueError(\"A condition response must carry its plausibility block\")\n", + "print(\n", + " \"log density of the held level over every step:\",\n", + " held.plausibility.equality_log_density,\n", + ")\n", + "held_path = baseline.to_pandas_wide().merge(\n", + " held.to_pandas_wide(),\n", + " on=\"query_time\",\n", + " suffixes=(\"_unconditional\", \"_held\"),\n", + ")\n", + "held_path" + ] + }, + { + "cell_type": "markdown", + "id": "7", "metadata": { "id": "condition-draws-description", "language": "markdown" @@ -196,13 +245,13 @@ "## Scenarios, not summaries\n", "A mean answers *where*, and a band answers *how wide*, but a portfolio question usually needs whole futures: draws. Every row below is one coherent joint scenario, because the service draws the read-out columns together from the same conditional mixture — the NAV and the volatility in one row belong to each other, so a function of several columns can be evaluated row by row. A quantile table cannot answer that: the 90th percentile of NAV and the 90th percentile of volatility need not describe any single future.\n", "\n", - "`diagnostics.condition_draws` reports how many draws stand behind the answer. Oversized sample requests are split into batches against the deployment's advertised cap and merged locally, and the merged count is what this field reports." + "The table shows the conditioned step; the other steps come back too, drawn unconditionally. `diagnostics.condition_draws` reports how many draws stand behind the answer. Oversized sample requests are split into batches against the deployment's advertised cap and merged locally, and the merged count is what this field reports." ] }, { "cell_type": "code", "execution_count": null, - "id": "6", + "id": "8", "metadata": { "id": "condition-draws", "language": "python" @@ -220,36 +269,35 @@ " seed=7,\n", " condition=rally,\n", ")\n", - "if scenarios.query_times != (QUERY_TIMES[CONDITIONED_STEP],):\n", - " raise ValueError(\n", - " f\"Expected the conditioned position alone, got {scenarios.query_times!r}\"\n", - " )\n", + "if scenarios.query_times != tuple(QUERY_TIMES):\n", + " raise ValueError(f\"Expected every query time, got {scenarios.query_times!r}\")\n", "print(\"draws behind the answer:\", scenarios.diagnostics.condition_draws)\n", - "scenarios.to_pandas_wide()" + "draws = scenarios.to_pandas_wide()\n", + "draws[draws[\"query_time\"] == CONDITIONED_TIME].reset_index(drop=True)" ] }, { "cell_type": "markdown", - "id": "7", + "id": "9", "metadata": { "id": "condition-interval-description", "language": "markdown" }, "source": [ "## Interval condition\n", - "Now confine the 10-year yield to a band around its last observed value instead of pinning it, and pin the equity index at the same time: a request may mix both kinds across the columns of one position. The response then carries `region_log_probability`, and the yield column itself stays readable because its distribution inside the band is a genuine answer.\n", + "Now confine the 10-year yield to a band around its last observed value instead of pinning it, and pin the equity index at the same time: conditions on different columns may cover the same step, which is how one step mixes both kinds. The response then carries `region_log_probability`, and the yield column itself stays readable because its distribution inside the band is a genuine answer.\n", "\n", "That number is the log probability of the band **given the pinned equity index**, not the band's own probability, because the two condition kinds compose in a fixed order and the band is measured on the distribution the pin has already reduced. The two numbers therefore chain rather than describe separate things, and adding them gives the plausibility of the whole request — `log(density(pin) * P(band | pin))` — with no independence assumed anywhere. A request carrying only interval conditions has nothing to combine, and its `region_log_probability` alone is the joint probability of everything it asked about.\n", "\n", "With one interval column, as here, the region probability is exact — one difference of distribution functions per mixture component — and `diagnostics.interval_estimator` stays empty, because there is no estimate to characterize. Bounding a second column makes the region a box with no closed form, which the service estimates numerically and then reports the accounting for — the next section does exactly that.\n", "\n", - "A band is also where the quantile read-out earns its place over the mean: what the banded column comes back with is a distribution inside its own range, so the cell asserts every quantile of it lands there." + "A band is also where the quantile read-out earns its place over the mean: what the banded column comes back with is a distribution inside its own range, so the cell asserts every quantile of it lands there at the step the band covers." ] }, { "cell_type": "code", "execution_count": null, - "id": "8", + "id": "10", "metadata": { "id": "condition-interval-example", "language": "python" @@ -268,15 +316,12 @@ "yield_lower = last_yield - YIELD_BAND_HALF_WIDTH\n", "yield_upper = last_yield + YIELD_BAND_HALF_WIDTH\n", "yield_band = IntervalCondition(\n", - " column=\"treasury_10y_yield\", lower=yield_lower, upper=yield_upper\n", - ")\n", - "rally_with_yield_band = ConditionBlock(\n", - " query_time_index=CONDITIONED_STEP,\n", - " conditions=[\n", - " EqualityCondition(column=PINNED_COLUMN, value=rallied_equity_level),\n", - " yield_band,\n", - " ],\n", + " column=\"treasury_10y_yield\",\n", + " lower=yield_lower,\n", + " upper=yield_upper,\n", + " query_time_indices=[CONDITIONED_STEP],\n", ")\n", + "rally_with_yield_band = [rally, yield_band]\n", "banded = client.forecast_quantiles(\n", " history,\n", " query_times=QUERY_TIMES,\n", @@ -292,7 +337,10 @@ "print(\"log probability of the yield band:\", banded.plausibility.region_log_probability)\n", "if banded.diagnostics.interval_estimator is not None:\n", " raise ValueError(\"A single interval column is exact and reports no estimator\")\n", - "band = banded.to_pandas_wide()\n", + "quantile_frame = banded.to_pandas_wide()\n", + "band = quantile_frame[quantile_frame[\"query_time\"] == CONDITIONED_TIME].reset_index(\n", + " drop=True\n", + ")\n", "outside_band = band[\n", " (band[\"treasury_10y_yield\"] < yield_lower - BAND_TOLERANCE)\n", " | (band[\"treasury_10y_yield\"] > yield_upper + BAND_TOLERANCE)\n", @@ -304,7 +352,7 @@ }, { "cell_type": "markdown", - "id": "9", + "id": "11", "metadata": { "id": "condition-box-description", "language": "markdown" @@ -313,13 +361,13 @@ "## When the region has to be estimated\n", "Bounding the FX rate as well turns the region into a box. That has no closed form, so the service estimates its probability with a quasi-random walk and only then reports the accounting in `diagnostics.interval_estimator`: `points` is how many quasi-random points went into the estimate and `effective_sample_size` is Kish's effective sample size of the weights behind them. A small effective sample size against a large point count means a deep-tail box where few points carry the answer. The cell asserts both sides of that rule — absent for the single band above, present here.\n", "\n", - "Drawing from a boxed conditional also shows what the box does to the draws. Each row pairs one truncated draw of the bounded columns with a read-out drawn from *that* draw's conditional, so the rows are exact draws from the joint conditional rather than separately summarized margins, and every one of them lands inside both bands." + "Drawing from a boxed conditional also shows what the box does to the draws. Each row pairs one truncated draw of the bounded columns with a read-out drawn from *that* draw's conditional, so the rows are exact draws from the joint conditional rather than separately summarized margins, and every one of them lands inside both bands at the step they cover." ] }, { "cell_type": "code", "execution_count": null, - "id": "10", + "id": "12", "metadata": { "id": "condition-box", "language": "python" @@ -331,15 +379,13 @@ "last_fx_rate = float(history[\"eur_usd_rate\"].iloc[-1])\n", "fx_lower = last_fx_rate - FX_BAND_HALF_WIDTH\n", "fx_upper = last_fx_rate + FX_BAND_HALF_WIDTH\n", - "fx_band = IntervalCondition(column=\"eur_usd_rate\", lower=fx_lower, upper=fx_upper)\n", - "rally_with_two_bands = ConditionBlock(\n", - " query_time_index=CONDITIONED_STEP,\n", - " conditions=[\n", - " EqualityCondition(column=PINNED_COLUMN, value=rallied_equity_level),\n", - " yield_band,\n", - " fx_band,\n", - " ],\n", + "fx_band = IntervalCondition(\n", + " column=\"eur_usd_rate\",\n", + " lower=fx_lower,\n", + " upper=fx_upper,\n", + " query_time_indices=[CONDITIONED_STEP],\n", ")\n", + "rally_with_two_bands = [rally, yield_band, fx_band]\n", "boxed = client.forecast_samples(\n", " history,\n", " query_times=QUERY_TIMES,\n", @@ -362,7 +408,10 @@ ")\n", "print(\"estimator points:\", estimator.points)\n", "print(\"effective sample size:\", estimator.effective_sample_size)\n", - "box_draws = boxed.to_pandas_wide()\n", + "boxed_frame = boxed.to_pandas_wide()\n", + "box_draws = boxed_frame[boxed_frame[\"query_time\"] == CONDITIONED_TIME].reset_index(\n", + " drop=True\n", + ")\n", "outside_box = box_draws[\n", " (box_draws[\"treasury_10y_yield\"] < yield_lower - BAND_TOLERANCE)\n", " | (box_draws[\"treasury_10y_yield\"] > yield_upper + BAND_TOLERANCE)\n", @@ -376,7 +425,7 @@ }, { "cell_type": "markdown", - "id": "11", + "id": "13", "metadata": { "id": "condition-ranking-description", "language": "markdown" @@ -391,7 +440,7 @@ { "cell_type": "code", "execution_count": null, - "id": "12", + "id": "14", "metadata": { "id": "condition-ranking", "language": "python" @@ -409,14 +458,18 @@ " requested_columns=plan.requested_columns,\n", " columns=plan.columns,\n", " seed=7,\n", - " condition=ConditionBlock(\n", - " query_time_index=CONDITIONED_STEP,\n", - " conditions=[EqualityCondition(column=PINNED_COLUMN, value=candidate_level)],\n", + " condition=EqualityCondition(\n", + " column=PINNED_COLUMN,\n", + " value=candidate_level,\n", + " query_time_indices=[CONDITIONED_STEP],\n", " ),\n", " )\n", " if candidate_result.plausibility is None:\n", " raise ValueError(\"A condition response must carry its plausibility block\")\n", - " conditional_mean = candidate_result.to_pandas_wide().iloc[0]\n", + " candidate_frame = candidate_result.to_pandas_wide()\n", + " conditional_mean = candidate_frame[\n", + " candidate_frame[\"query_time\"] == CONDITIONED_TIME\n", + " ].iloc[0]\n", " rankings.append(\n", " {\n", " \"rally\": candidate,\n", @@ -435,7 +488,7 @@ }, { "cell_type": "markdown", - "id": "13", + "id": "15", "metadata": { "id": "condition-score-description", "language": "markdown" @@ -444,7 +497,7 @@ "## Scoring an outcome you already have\n", "Every read-out so far answers *what does the model expect*. `log_prob` asks the opposite — how plausible are values I already hold — and is the only return mode where the caller supplies the future instead of receiving it. Under a condition it scores those values against the conditional, which is how two candidate outcomes of the same what-if are compared.\n", "\n", - "Three rules follow from scoring a joint rather than a projection of one, and the client enforces all three before anything is sent. `query_rows` carries one observed row per entry of `query_times`, each with a value for every declared column. `requested_columns` must name every declared column in declared order — a narrower projection would score a different distribution than the caller means to ask about — and a condition excuses nothing, so omitting the field is the natural spelling. The deployment refuses a row that contradicts the condition, a pinned column given another value or a banded column outside its range, instead of scoring it.\n", + "Three rules follow from scoring a joint rather than a projection of one, and the client enforces all three before anything is sent. `query_rows` carries one observed row per entry of `query_times`, each with a value for every declared column. `requested_columns` must name every declared column in declared order — a narrower projection would score a different distribution than the caller means to ask about — and a condition excuses nothing, so omitting the field is the natural spelling. The deployment refuses a row that contradicts a condition covering its step, a pinned column given another value or a banded column outside its range, instead of scoring it.\n", "\n", "The scores below are log densities of whole rows, so the unit caveat from the plausibility section applies unchanged: read their difference, not their level." ] @@ -452,7 +505,7 @@ { "cell_type": "code", "execution_count": null, - "id": "14", + "id": "16", "metadata": { "id": "condition-score", "language": "python" @@ -481,7 +534,6 @@ " history,\n", " query_times=SCORED_QUERY_TIMES,\n", " query_rows=pd.DataFrame([outcome], columns=EXPECTED_COLUMNS),\n", - " requested_columns=READABLE_COLUMNS,\n", " columns=plan.columns,\n", " seed=7,\n", " condition=rally,\n", diff --git a/src/jointfm_client/__init__.py b/src/jointfm_client/__init__.py index ffe9a94..77a9900 100644 --- a/src/jointfm_client/__init__.py +++ b/src/jointfm_client/__init__.py @@ -69,13 +69,14 @@ PACKAGE_VERSION, PREDICT_REQUEST_TYPE, SCHEMA_VERSION, - ConditionBlock, + Condition, ConditionKind, ConditionPlausibility, EqualityCondition, IntervalCondition, IntervalEstimator, require_condition_support, + resolve_conditions, MeanForecastResult, QuantileForecast, QuantileForecastResult, @@ -224,13 +225,14 @@ "PREDICT_REQUEST_TYPE", "RetryConfig", "SCHEMA_VERSION", - "ConditionBlock", + "Condition", "ConditionKind", "ConditionPlausibility", "EqualityCondition", "IntervalCondition", "IntervalEstimator", "require_condition_support", + "resolve_conditions", "SampleForecastResult", "SUPPORTED_COLUMN_MODALITIES", "SUPPORTED_COLUMN_ROLES", diff --git a/src/jointfm_client/adapters.py b/src/jointfm_client/adapters.py index d1aa57d..8872431 100644 --- a/src/jointfm_client/adapters.py +++ b/src/jointfm_client/adapters.py @@ -24,7 +24,7 @@ from jointfm_client.configuration import DEFAULT_FORECAST_SCHEMA_VERSION from jointfm_client.contract import ( DEFAULT_CALENDAR_ID, - ConditionBlock, + Condition, ColumnModality, ColumnRole, ColumnSpec, @@ -304,14 +304,15 @@ def build_forecast_payload_from_dataframe( time_value_columns: Sequence[str] | Mapping[str, TimeValueKind] | None = None, nullable_columns: Sequence[str] | None = None, bounds: ColumnBounds | None = None, - condition: ConditionBlock | None = None, + condition: Condition | Sequence[Condition] | None = None, query_rows: Any | None = None, ) -> dict[str, Any]: """Build a validated forecast payload from a pandas ``DataFrame``. - Passing ``condition`` makes this a conditioning request: the payload then - carries ``query_mode='condition'`` and the block, and the deployment - answers the conditional at the one future position the block names. + Passing ``condition`` — one condition or a list of them — makes this a + conditioning request: the payload then carries ``query_mode='condition'`` + and the conditions, and the deployment answers every future position under + the conditions covering it. ``query_rows`` carries the *observed* values at ``query_times`` that ``return_mode='log_prob'`` scores — a ``DataFrame`` shaped like ``frame``, diff --git a/src/jointfm_client/client.py b/src/jointfm_client/client.py index 8a6a5bd..7c08bd7 100644 --- a/src/jointfm_client/client.py +++ b/src/jointfm_client/client.py @@ -36,7 +36,7 @@ load_configuration, ) from jointfm_client.contract import ( - ConditionBlock, + Condition, DEFAULT_CALENDAR_ID, HEALTH_REQUEST_TYPE, SCHEMA_VERSION, @@ -313,15 +313,21 @@ def forecast( nullable_columns: Sequence[str] | None = None, bounds: Mapping[str, tuple[float | int | None, float | int | None]] | None = None, - condition: ConditionBlock | None = None, + condition: Condition | Sequence[Condition] | None = None, query_rows: Any | None = None, ) -> ForecastResponse: """Build and submit a forecast request from tabular history inputs. - Passing ``condition`` asks the deployment for the conditional at the one - future position the block names, instead of the unconditional forecast. - The deployment's advertised capability is checked first, so a deployment - that cannot condition is refused here rather than after a round trip. + Passing ``condition`` — one ``EqualityCondition`` or + ``IntervalCondition``, or a list of them — asks the deployment for the + forecast given those conditions. Each condition covers the positions its + ``query_time_indices`` names, every position when that is ``None``, and + two conditions on one column must cover disjoint positions. The answer + still covers every entry of ``query_times``: a position no condition + covers carries the unconditioned forecast, which is exact because the + deployment draws each position independently. The deployment's + advertised capability is checked first, so a deployment that cannot + condition is refused here rather than after a round trip. ``query_rows`` carries the observed values at ``query_times`` that ``return_mode='log_prob'`` scores, in the same shape as ``history``. @@ -417,12 +423,12 @@ def forecast_mean( requested_columns: Sequence[str | int] | None = None, model_version: str | None = None, seed: int | None = None, - condition: ConditionBlock | None = None, + condition: Condition | Sequence[Condition] | None = None, ) -> MeanForecastResult: """Forecast mean values through the shared forecast validation path. - With ``condition`` the mean is the conditional mean at the one future - position the block names; see :meth:`forecast`. + With ``condition`` the mean at each covered position is the conditional + mean there; see :meth:`forecast`. """ return cast( MeanForecastResult, @@ -454,12 +460,12 @@ def forecast_samples( model_version: str | None = None, n_samples: int | None = None, seed: int | None = None, - condition: ConditionBlock | None = None, + condition: Condition | Sequence[Condition] | None = None, ) -> SampleForecastResult: """Forecast sample paths through the shared forecast validation path. - With ``condition`` the draws come from the conditional at the one future - position the block names; see :meth:`forecast`. + With ``condition`` the draws at each covered position come from the + conditional there; see :meth:`forecast`. """ return cast( SampleForecastResult, @@ -493,12 +499,12 @@ def forecast_quantiles( n_samples: int | None = None, quantiles: Sequence[float | int] | None = None, seed: int | None = None, - condition: ConditionBlock | None = None, + condition: Condition | Sequence[Condition] | None = None, ) -> QuantileForecastResult: """Forecast quantiles through the shared forecast validation path. - With ``condition`` the quantiles describe the conditional at the one - future position the block names; see :meth:`forecast`. + With ``condition`` the quantiles at each covered position describe the + conditional there; see :meth:`forecast`. """ return cast( QuantileForecastResult, @@ -532,7 +538,7 @@ def forecast_log_prob( requested_columns: Sequence[str | int] | None = None, model_version: str | None = None, seed: int | None = None, - condition: ConditionBlock | None = None, + condition: Condition | Sequence[Condition] | None = None, ) -> LogProbResult: """Score observed future values through the shared forecast validation path. @@ -544,9 +550,9 @@ def forecast_log_prob( declared column in declared order, which is what omitting it already resolves to. - With ``condition`` the score is taken under the conditional at the one - future position the block names, and the service refuses rows that - contradict the condition instead of scoring them; see :meth:`forecast`. + With ``condition`` each covered row is scored under the conditional at + its position, and the service refuses a row that contradicts a + condition covering it instead of scoring it; see :meth:`forecast`. """ return cast( LogProbResult, @@ -721,7 +727,7 @@ def _forecast_payload_from_rows( use_local_normalized_time: bool, calendar_id: str, timezone: str | None, - condition: ConditionBlock | None = None, + condition: Condition | Sequence[Condition] | None = None, query_rows: Any | None = None, ) -> dict[str, Any]: if schema is None: diff --git a/src/jointfm_client/contract.py b/src/jointfm_client/contract.py index 2b028b4..598875e 100644 --- a/src/jointfm_client/contract.py +++ b/src/jointfm_client/contract.py @@ -37,7 +37,7 @@ # pyproject.toml the way a hand-maintained literal here did. PACKAGE_VERSION: Final = importlib.metadata.version(DISTRIBUTION_NAME) -SCHEMA_VERSION: Final = "v4" +SCHEMA_VERSION: Final = "v5" # The service sums and averages its log densities in the model's own precision, # so its summaries differ from a recomputation in the last few digits. DERIVED_SCORE_TOLERANCE: Final = 1e-6 @@ -311,19 +311,37 @@ def to_payload(self) -> dict[str, Any]: return payload +def _require_query_time_indices(value: Any) -> tuple[int, ...] | None: + """Validate the positions one condition covers, ``None`` meaning every one.""" + if value is None: + return None + indices = _require_sequence(value, field="condition.query_time_indices") + for index in indices: + if isinstance(index, bool) or not isinstance(index, int): + raise ValueError("condition.query_time_indices must hold integers") + if index < 0: + raise ValueError("condition.query_time_indices must not be negative") + if len(set(indices)) != len(indices): + raise ValueError("condition.query_time_indices names a position more than once") + return tuple(indices) + + @dataclass(frozen=True, slots=True) class EqualityCondition: - """One column of the conditioned position pinned to a value. - - An equality condition is an event of probability zero under a continuous - head, so the deployment answers it analytically rather than by filtering - draws. It fixes the column's value, so a projection that names the column - reads that value back rather than a model output — which is what lets a - scenario answer line up column for column with an unconditioned one. + """One column pinned to a value at the future positions it covers. + + ``query_time_indices`` names those positions by index into the request's + ``query_times``; ``None`` covers every one of them. An equality condition + is an event of probability zero under a continuous head, so the deployment + answers it analytically rather than by filtering draws. It fixes the + column's value, so a projection that names the column reads that value back + at every covered position rather than a model output — which is what lets + a scenario answer line up column for column with an unconditioned one. """ column: str value: float + query_time_indices: Sequence[int] | None = None def __post_init__(self) -> None: """Reject a pin the service would refuse anyway.""" @@ -332,25 +350,37 @@ def __post_init__(self) -> None: raise ValueError("condition value must be a number") if not math.isfinite(float(self.value)): raise ValueError("condition value must be finite") + object.__setattr__( + self, + "query_time_indices", + _require_query_time_indices(self.query_time_indices), + ) def to_payload(self) -> dict[str, Any]: """Return this condition as JSON-compatible payload fields.""" - return {"column": self.column, "kind": "equality", "value": float(self.value)} + return { + "column": self.column, + "kind": "equality", + "value": float(self.value), + "query_time_indices": _query_time_indices_payload(self), + } @dataclass(frozen=True, slots=True) class IntervalCondition: - """One column of the conditioned position confined to a range. + """One column confined to a range at the future positions it covers. - ``None`` on either side leaves it open. The region carries positive mass, - so the deployment answers it by composition on the reduced mixture, and the - column stays readable: what it returns is its distribution inside the - region. + ``None`` on either side leaves the range open, and ``query_time_indices`` + of ``None`` covers every requested position. The region carries positive + mass, so the deployment answers it by composition on the reduced mixture, + and the column stays readable: what it returns is its distribution inside + the region. """ column: str lower: float | None = None upper: float | None = None + query_time_indices: Sequence[int] | None = None def __post_init__(self) -> None: """Reject a range that bounds nothing or bounds it backwards.""" @@ -374,6 +404,11 @@ def __post_init__(self) -> None: raise ValueError( f"condition needs lower < upper, got [{self.lower}, {self.upper}]" ) + object.__setattr__( + self, + "query_time_indices", + _require_query_time_indices(self.query_time_indices), + ) def to_payload(self) -> dict[str, Any]: """Return this condition as JSON-compatible payload fields.""" @@ -382,78 +417,85 @@ def to_payload(self) -> dict[str, Any]: "kind": "interval", "lower": None if self.lower is None else float(self.lower), "upper": None if self.upper is None else float(self.upper), + "query_time_indices": _query_time_indices_payload(self), } -@dataclass(frozen=True, slots=True) -class ConditionBlock: - """Every condition of one request, against one future position. +Condition: TypeAlias = EqualityCondition | IntervalCondition - The position is named once, by its index into the request's - ``query_times``, because conditioning relates variables to each other - within a position and never across horizons. A statement spanning - horizons is out of reach of this mode and cannot be expressed here. - """ - query_time_index: int - conditions: Sequence[EqualityCondition | IntervalCondition] +def _query_time_indices_payload(condition: Condition) -> list[int] | None: + """Serialize a condition's positions, ``None`` staying the wire's ``null``.""" + if condition.query_time_indices is None: + return None + return list(condition.query_time_indices) - def __post_init__(self) -> None: - """Validate the block independently of any schema or deployment.""" - if isinstance(self.query_time_index, bool) or not isinstance( - self.query_time_index, int - ): - raise ValueError("query_time_index must be an integer") - if self.query_time_index < 0: - raise ValueError("query_time_index must not be negative") - conditions = _require_sequence(self.conditions, field="conditions") - if not conditions: - raise ValueError("a condition block must carry at least one condition") - seen: set[str] = set() - for index, condition in enumerate(conditions): - if not isinstance(condition, (EqualityCondition, IntervalCondition)): - raise ValueError( - f"conditions[{index}] must be an EqualityCondition or an " - "IntervalCondition" - ) - if condition.column in seen: - raise ValueError( - f"column {condition.column!r} carries more than one condition" - ) - seen.add(condition.column) - @property - def pinned_columns(self) -> tuple[str, ...]: - """Columns an equality condition fixes, whose answer is the request's own value.""" - return tuple( - condition.column - for condition in self.conditions - if isinstance(condition, EqualityCondition) - ) +def _covers(condition: Condition, position: int) -> bool: + """Return whether ``condition`` applies at future ``position``.""" + return ( + condition.query_time_indices is None or position in condition.query_time_indices + ) - @property - def conditioned_columns(self) -> tuple[str, ...]: - """Every column this block conditions, of either kind.""" - return tuple(condition.column for condition in self.conditions) - @property - def kinds(self) -> tuple[ConditionKind, ...]: - """The condition kinds this block uses, which a deployment must advertise.""" - kinds: list[ConditionKind] = [] - for condition in self.conditions: - kind: ConditionKind = ( - "equality" if isinstance(condition, EqualityCondition) else "interval" +def resolve_conditions( + condition: Condition | Sequence[Condition], +) -> tuple[Condition, ...]: + """Normalize one condition or a list of them into a validated tuple. + + Conditions on different columns may cover the same positions, which is how + one position mixes a pin and a band. Two conditions on the *same* column + must cover disjoint positions, because a column carries at most one + condition per position; ``query_time_indices=None`` covers every position, + so it overlaps any other condition on its column. + + Raises: + ValueError: when the list is empty, holds something other than an + ``EqualityCondition`` or ``IntervalCondition``, or puts two + conditions on one column at a shared position. + """ + conditions: tuple[Any, ...] = ( + (condition,) + if isinstance(condition, (EqualityCondition, IntervalCondition)) + else tuple(_require_sequence(condition, field="condition")) + ) + for index, entry in enumerate(conditions): + if not isinstance(entry, (EqualityCondition, IntervalCondition)): + raise ValueError( + f"condition[{index}] must be an EqualityCondition or an " + "IntervalCondition" ) - if kind not in kinds: - kinds.append(kind) - return tuple(kinds) + for later_index, later in enumerate(conditions): + for earlier_index, earlier in enumerate(conditions[:later_index]): + if earlier.column != later.column: + continue + if earlier.query_time_indices is None or later.query_time_indices is None: + overlap = "every position" + else: + shared = sorted( + set(earlier.query_time_indices) & set(later.query_time_indices) + ) + if not shared: + continue + overlap = f"query_time_indices {shared}" + raise ValueError( + f"column {later.column!r} carries more than one condition at " + f"{overlap}: condition[{later_index}] overlaps " + f"condition[{earlier_index}]" + ) + return cast(tuple[Condition, ...], conditions) - def to_payload(self) -> dict[str, Any]: - """Return the block as JSON-compatible payload fields.""" - return { - "query_time_index": self.query_time_index, - "conditions": [condition.to_payload() for condition in self.conditions], - } + +def condition_kinds(conditions: Sequence[Condition]) -> tuple[ConditionKind, ...]: + """Return the condition kinds ``conditions`` use, which a deployment must advertise.""" + kinds: list[ConditionKind] = [] + for condition in conditions: + kind: ConditionKind = ( + "equality" if isinstance(condition, EqualityCondition) else "interval" + ) + if kind not in kinds: + kinds.append(kind) + return tuple(kinds) @dataclass(frozen=True, slots=True) @@ -506,7 +548,7 @@ class ForecastRequest: quantiles: Sequence[float | int] | None = None seed: int | None = None query_row_ids: Sequence[int] | None = None - condition: ConditionBlock | None = None + condition: Condition | Sequence[Condition] | None = None query_rows: Sequence[Mapping[str, Any]] | None = None def __post_init__(self) -> None: @@ -531,7 +573,7 @@ def __post_init__(self) -> None: if (self.metadata.query_mode == "condition") != (self.condition is not None): raise ValueError( - "query_mode='condition' and a condition block go together: one " + "query_mode='condition' and a condition go together: one " "without the other means a different request than the caller wrote" ) if self.condition is not None: @@ -553,34 +595,50 @@ def __post_init__(self) -> None: self._validate_query_rows() def _validate_condition(self) -> None: - """Check the condition block against this request's schema and horizons. + """Check the conditions against this request's schema and positions. The deployment rejects all of this too, before executing anything; the - point of repeating it here is that a caller learns which column name it - got wrong without paying for a round trip. + point of repeating it here is that a caller learns which column name or + position it got wrong without paying for a round trip. """ - block = self.condition - assert block is not None + assert self.condition is not None + conditions = resolve_conditions(self.condition) declared = {column.name for column in self.schema.columns} - unknown = [name for name in block.conditioned_columns if name not in declared] + unknown = [ + condition.column + for condition in conditions + if condition.column not in declared + ] if unknown: raise ValueError(f"condition references undeclared columns {unknown}") horizon_count = len(_require_sequence(self.query_times, field="query_times")) - if block.query_time_index >= horizon_count: - raise ValueError( - f"condition.query_time_index {block.query_time_index} is outside the " - f"{horizon_count} requested future positions" - ) - if len(set(block.conditioned_columns)) >= len(declared): - if len(block.pinned_columns) == len(set(block.conditioned_columns)): + for condition in conditions: + outside = [ + index + for index in condition.query_time_indices or () + if index >= horizon_count + ] + if outside: raise ValueError( - "pinning every column leaves nothing to read out; the joint " - "density of those values is what return_mode='log_prob' answers" + f"condition on {condition.column!r} names query_time_indices " + f"{outside} outside the {horizon_count} requested future positions" + ) + for position in range(horizon_count): + covering = [ + condition for condition in conditions if _covers(condition, position) + ] + if len(covering) < len(declared): + continue + if all(isinstance(condition, EqualityCondition) for condition in covering): + raise ValueError( + f"pinning every column at query_time_indices position {position} " + "leaves nothing to read out; the joint density of those values " + "is what return_mode='log_prob' answers" ) raise ValueError( - "a request must leave at least one column unconditioned: what it " - "reads out is the conditional distribution of the columns it does " - "not condition" + "a request must leave at least one column unconditioned at " + f"query_time_indices position {position}: what it reads out is " + "the conditional distribution of the columns it does not condition" ) def _validate_query_rows(self) -> None: @@ -590,7 +648,7 @@ def _validate_query_rows(self) -> None: must carry every declared column and ``requested_columns`` must name every one of them in declared order: a narrower or reordered projection would score a different distribution than the caller believes it asked - about. A ``condition`` block changes nothing here — its columns are + about. A ``condition`` changes nothing here — its columns are declared columns like any other — and omitting the field already resolves to exactly this projection. @@ -663,7 +721,10 @@ def to_payload(self) -> dict[str, Any]: if self.seed is not None: payload["seed"] = self.seed if self.condition is not None: - payload["condition"] = self.condition.to_payload() + payload["conditions"] = [ + condition.to_payload() + for condition in resolve_conditions(self.condition) + ] if self.query_rows is not None: payload["query_rows"] = _serialize_rows(self.query_rows, field="query_rows") return payload @@ -850,9 +911,12 @@ def from_payload(cls, payload: Mapping[str, Any]) -> Self: class IntervalEstimator: """How a multi-column region's probability was estimated. - Present only when the interval block spans more than one column, which is - the only case where a component's box probability has to be estimated; one + Present only when some position bounds more than one column, which is the + only case where a component's box probability has to be estimated; one interval column is exact, and a number that never varies would say nothing. + When several positions are estimated, the accounting describes the one with + the smallest effective sample size, because that position bounds how far + the whole response can be trusted. """ points: int @@ -883,14 +947,17 @@ class ConditionPlausibility: kind. The service reports them and never refuses on them, so acting on them is the caller's decision. - Each number is joint over its whole condition block rather than one value + Each number is joint over every condition of its kind rather than one value per column: ``equality_log_density`` covers every pinned value at once, and ``region_log_probability`` covers the whole box at once. Do not rebuild either by multiplying per-column numbers, which would assume the columns are - independent and so discard what a joint model is for. + independent and so discard what a joint model is for. Across positions the + numbers *are* sums: the deployment draws each future position from its own + joint, independently of the others, so the value for a request covering + several positions is the sum of the per-position values, exactly. ``region_log_probability`` is measured *after* the pinned columns are - applied, so on a request carrying both kinds it is the log probability of + applied, so at a position carrying both kinds it is the log probability of the box **given the pins**, not the box's own. The two therefore chain, and their sum is the plausibility of the whole condition set:: @@ -1628,7 +1695,7 @@ def build_forecast_payload( seed: int | None = None, schema_version: str = SCHEMA_VERSION, query_mode: QueryMode = "forecast", - condition: ConditionBlock | None = None, + condition: Condition | Sequence[Condition] | None = None, query_rows: Sequence[Mapping[str, Any]] | None = None, ) -> dict[str, Any]: """Build a validated JSON-compatible forecast request payload.""" @@ -1653,7 +1720,7 @@ def build_forecast_payload( def require_condition_support( metadata: HealthMetadata, - block: ConditionBlock, + condition: Condition | Sequence[Condition], ) -> None: """Refuse a conditioning request the mounted deployment has not advertised. @@ -1664,7 +1731,7 @@ def require_condition_support( Raises: UnsupportedServiceContractError: when the deployment serves no - ``condition`` mode, or none of the condition kinds this block uses. + ``condition`` mode, or not every condition kind ``condition`` uses. """ if "condition" not in metadata.supported_query_modes: raise UnsupportedServiceContractError( @@ -1672,7 +1739,9 @@ def require_condition_support( f"queries; it advertises {list(metadata.supported_query_modes)}" ) missing = [ - kind for kind in block.kinds if kind not in metadata.supported_condition_kinds + kind + for kind in condition_kinds(resolve_conditions(condition)) + if kind not in metadata.supported_condition_kinds ] if missing: raise UnsupportedServiceContractError( @@ -2064,18 +2133,6 @@ def _forecast_response_expectations( ): requested_columns = tuple(requested_column_values) - condition_value = request_payload.get("condition") - if condition_value is not None and query_times is not None: - condition_block = _require_mapping( - condition_value, field="request_payload.condition" - ) - conditioned_index = condition_block.get("query_time_index") - if isinstance(conditioned_index, int) and not isinstance( - conditioned_index, bool - ): - if 0 <= conditioned_index < len(query_times): - query_times = (query_times[conditioned_index],) - query_mode_value = request_payload.get("query_mode") return_mode_value = request_payload.get("return_mode") quantiles_value = request_payload.get("quantiles") diff --git a/tests/fixtures/condition_interval_response.json b/tests/fixtures/condition_interval_response.json index f217332..d076318 100644 --- a/tests/fixtures/condition_interval_response.json +++ b/tests/fixtures/condition_interval_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", @@ -7,9 +7,10 @@ "query_mode": "condition", "return_mode": "mean", "outputs": { - "query_times": [3], + "query_times": [2, 3], "requested_columns": ["target"], "mean": [ + [12.0], [14.25] ], "samples": null, @@ -21,7 +22,7 @@ }, "diagnostics": { "history_rows": 2, - "horizon_count": 1, + "horizon_count": 2, "seed": 7, "condition_draws": 500, "interval_estimator": { diff --git a/tests/fixtures/condition_log_prob_request.json b/tests/fixtures/condition_log_prob_request.json index b6ef705..c3e5c74 100644 --- a/tests/fixtures/condition_log_prob_request.json +++ b/tests/fixtures/condition_log_prob_request.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "query_mode": "condition", "return_mode": "log_prob", @@ -27,16 +27,14 @@ ], "query_times": [2, 3], "requested_columns": ["driver", "target"], - "condition": { - "query_time_index": 1, - "conditions": [ - { - "column": "driver", - "kind": "equality", - "value": 1.5 - } - ] - }, + "conditions": [ + { + "column": "driver", + "kind": "equality", + "value": 1.5, + "query_time_indices": [1] + } + ], "query_rows": [ { "driver": 1.2, diff --git a/tests/fixtures/condition_log_prob_response.json b/tests/fixtures/condition_log_prob_response.json index 04513ea..a78c15c 100644 --- a/tests/fixtures/condition_log_prob_response.json +++ b/tests/fixtures/condition_log_prob_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", @@ -7,18 +7,18 @@ "query_mode": "condition", "return_mode": "log_prob", "outputs": { - "query_times": [3], + "query_times": [2, 3], "requested_columns": ["driver", "target"], "mean": null, "samples": null, "quantiles": null, "log_prob": { - "values": [-1.75], - "nll_values": [1.75], - "total": -1.75, - "mean": -1.75, - "nll_total": 1.75, - "nll_mean": 1.75 + "values": [-2.5, -1.75], + "nll_values": [2.5, 1.75], + "total": -4.25, + "mean": -2.125, + "nll_total": 4.25, + "nll_mean": 2.125 } }, "plausibility": { @@ -27,7 +27,7 @@ }, "diagnostics": { "history_rows": 2, - "horizon_count": 1, + "horizon_count": 2, "seed": 7, "condition_draws": 500, "interval_estimator": null diff --git a/tests/fixtures/condition_mean_request.json b/tests/fixtures/condition_mean_request.json index f506b20..ecd9608 100644 --- a/tests/fixtures/condition_mean_request.json +++ b/tests/fixtures/condition_mean_request.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "query_mode": "condition", "return_mode": "mean", @@ -27,15 +27,13 @@ ], "query_times": [2, 3], "requested_columns": ["target"], - "condition": { - "query_time_index": 1, - "conditions": [ - { - "column": "driver", - "kind": "equality", - "value": 1.5 - } - ] - }, + "conditions": [ + { + "column": "driver", + "kind": "equality", + "value": 1.5, + "query_time_indices": [1] + } + ], "seed": 7 } diff --git a/tests/fixtures/condition_mean_response.json b/tests/fixtures/condition_mean_response.json index 0c06ed6..db2d7ca 100644 --- a/tests/fixtures/condition_mean_response.json +++ b/tests/fixtures/condition_mean_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", @@ -7,9 +7,10 @@ "query_mode": "condition", "return_mode": "mean", "outputs": { - "query_times": [3], + "query_times": [2, 3], "requested_columns": ["target"], "mean": [ + [12.0], [13.5] ], "samples": null, @@ -21,7 +22,7 @@ }, "diagnostics": { "history_rows": 2, - "horizon_count": 1, + "horizon_count": 2, "seed": 7, "condition_draws": 500, "interval_estimator": null diff --git a/tests/fixtures/forecast_log_prob_request.json b/tests/fixtures/forecast_log_prob_request.json index 6ba470f..ccd069d 100644 --- a/tests/fixtures/forecast_log_prob_request.json +++ b/tests/fixtures/forecast_log_prob_request.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "query_mode": "forecast", "return_mode": "log_prob", diff --git a/tests/fixtures/forecast_log_prob_response.json b/tests/fixtures/forecast_log_prob_response.json index ab9ee00..548207b 100644 --- a/tests/fixtures/forecast_log_prob_response.json +++ b/tests/fixtures/forecast_log_prob_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", diff --git a/tests/fixtures/forecast_mean_request.json b/tests/fixtures/forecast_mean_request.json index 17c29ff..18d43ac 100644 --- a/tests/fixtures/forecast_mean_request.json +++ b/tests/fixtures/forecast_mean_request.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "query_mode": "forecast", "return_mode": "mean", diff --git a/tests/fixtures/forecast_mean_response.json b/tests/fixtures/forecast_mean_response.json index d289237..dab2ecd 100644 --- a/tests/fixtures/forecast_mean_response.json +++ b/tests/fixtures/forecast_mean_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", diff --git a/tests/fixtures/forecast_quantiles_request.json b/tests/fixtures/forecast_quantiles_request.json index 32442ba..a1551df 100644 --- a/tests/fixtures/forecast_quantiles_request.json +++ b/tests/fixtures/forecast_quantiles_request.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "query_mode": "forecast", "return_mode": "quantiles", diff --git a/tests/fixtures/forecast_quantiles_response.json b/tests/fixtures/forecast_quantiles_response.json index cd49f12..c61c0da 100644 --- a/tests/fixtures/forecast_quantiles_response.json +++ b/tests/fixtures/forecast_quantiles_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", diff --git a/tests/fixtures/forecast_samples_request.json b/tests/fixtures/forecast_samples_request.json index 9203269..7b06d74 100644 --- a/tests/fixtures/forecast_samples_request.json +++ b/tests/fixtures/forecast_samples_request.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "query_mode": "forecast", "return_mode": "samples", diff --git a/tests/fixtures/forecast_samples_response.json b/tests/fixtures/forecast_samples_response.json index 29cff5c..45a10b4 100644 --- a/tests/fixtures/forecast_samples_response.json +++ b/tests/fixtures/forecast_samples_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", diff --git a/tests/fixtures/health_metadata.json b/tests/fixtures/health_metadata.json index 804fdce..836b1ec 100644 --- a/tests/fixtures/health_metadata.json +++ b/tests/fixtures/health_metadata.json @@ -1,6 +1,6 @@ { "status": "ok", - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", diff --git a/tests/fixtures/input_size_exceeded_response.json b/tests/fixtures/input_size_exceeded_response.json index 2e1c7c1..ab5be10 100644 --- a/tests/fixtures/input_size_exceeded_response.json +++ b/tests/fixtures/input_size_exceeded_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "errors": [ { "code": "INPUT_SIZE_EXCEEDED", diff --git a/tests/fixtures/model_version_mismatch_response.json b/tests/fixtures/model_version_mismatch_response.json index 07f5363..321c7fa 100644 --- a/tests/fixtures/model_version_mismatch_response.json +++ b/tests/fixtures/model_version_mismatch_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "errors": [ { "code": "MODEL_VERSION_MISMATCH", diff --git a/tests/fixtures/schema_version_mismatch_response.json b/tests/fixtures/schema_version_mismatch_response.json index d5cffa9..5cd59f7 100644 --- a/tests/fixtures/schema_version_mismatch_response.json +++ b/tests/fixtures/schema_version_mismatch_response.json @@ -1,9 +1,9 @@ { - "schema_version": "v3", + "schema_version": "v4", "errors": [ { "code": "SCHEMA_VERSION_MISMATCH", - "message": "Unsupported schema_version: expected 'v4', got 'v3'", + "message": "Unsupported schema_version: expected 'v5', got 'v4'", "field": "schema_version", "detail": {} } diff --git a/tests/fixtures/validation_error_response.json b/tests/fixtures/validation_error_response.json index 30d6160..226105d 100644 --- a/tests/fixtures/validation_error_response.json +++ b/tests/fixtures/validation_error_response.json @@ -1,5 +1,5 @@ { - "schema_version": "v4", + "schema_version": "v5", "errors": [ { "code": "VALIDATION_ERROR", diff --git a/tests/test_cli.py b/tests/test_cli.py index 006392a..1cc1d98 100644 --- a/tests/test_cli.py +++ b/tests/test_cli.py @@ -48,7 +48,7 @@ class FakeHealthClient: "deployment-id/predictionsUnstructured" ), deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", deployment_id="deployment-id", ) @@ -58,7 +58,7 @@ def health(self, *, cache: bool = False, refresh: bool = False) -> HealthMetadat del cache, refresh return HealthMetadata( status="ok", - schema_version="v4", + schema_version="v5", image_version="0.3.0", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", checkpoint_version="sdk-test", @@ -172,7 +172,7 @@ def test_predict_command_writes_response_file(monkeypatch, tmp_path: Path) -> No request_file.write_text( json.dumps( { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", } ), diff --git a/tests/test_condition_mode.py b/tests/test_condition_mode.py index 7a1e25c..50557aa 100644 --- a/tests/test_condition_mode.py +++ b/tests/test_condition_mode.py @@ -15,8 +15,9 @@ """Tests for the ``condition`` query mode on the client side. Three surfaces are covered, and the split matters. The condition objects -validate what is wrong independently of any deployment, so a caller learns it -without a round trip. The capability gate refuses what *this* deployment has +validate what is wrong independently of any deployment — including two +conditions on one column at a shared position — so a caller learns it without a +round trip. The capability gate refuses what *this* deployment has not advertised, which is the only way to know before paying for a request. The response parser reads back the two numbers the service reports and never refuses on, so acting on them stays the caller's decision. @@ -31,7 +32,7 @@ from jointfm_client import ( ColumnSpec, - ConditionBlock, + Condition, ConditionPlausibility, DataFrameSchema, EqualityCondition, @@ -46,9 +47,10 @@ build_forecast_payload, build_forecast_payload_from_dataframe, require_condition_support, + resolve_conditions, validate_service_metadata, ) -from jointfm_client.contract import QueryMode +from jointfm_client.contract import QueryMode, condition_kinds from jointfm_client.exceptions import UnsupportedServiceContractError _MODEL_VERSION = "jointfm-inference:0.3.0+ckpt.sdk-test" @@ -67,12 +69,12 @@ def _schema() -> DataFrameSchema: def _request( - block: ConditionBlock | None, + condition: Condition | list[Condition] | None, *, requested_columns: tuple[str, ...] | None = ("target",), query_mode: QueryMode = "condition", ) -> ForecastRequest: - """Build one condition request against the two-column schema.""" + """Build one condition request against the three-column schema.""" return ForecastRequest( metadata=ForecastRequestMetadata( model_version=_MODEL_VERSION, @@ -85,7 +87,7 @@ def _request( ), query_times=(2, 3), requested_columns=requested_columns, - condition=block, + condition=condition, ) @@ -97,7 +99,7 @@ def _health( """Build one advertisement without going through a transport.""" return HealthMetadata( status="ok", - schema_version="v4", + schema_version="v5", image_version="0.3.0", model_version=_MODEL_VERSION, checkpoint_version="sdk-test", @@ -136,108 +138,184 @@ def test_an_interval_condition_rejects_a_range_that_bounds_nothing( IntervalCondition(column="driver", lower=lower, upper=upper) -def test_a_block_rejects_two_conditions_on_one_column() -> None: - """A column carries at most one condition, of either kind.""" - with pytest.raises(ValueError, match="more than one condition"): - ConditionBlock( - query_time_index=0, - conditions=( - EqualityCondition(column="driver", value=1.0), - IntervalCondition(column="driver", lower=0.0, upper=1.0), - ), +@pytest.mark.parametrize( + ("positions", "message"), + [ + ([], "must not be empty"), + ([0, 0], "more than once"), + ([-1], "must not be negative"), + ([True], "must hold integers"), + ("0", "JSON array"), + ], + ids=["empty", "duplicate", "negative", "bool", "string"], +) +def test_a_condition_rejects_positions_that_name_nothing_usable( + positions: Any, message: str +) -> None: + """``None`` covers every position; an explicit list must name real, distinct ones.""" + with pytest.raises(ValueError, match=message): + EqualityCondition(column="driver", value=1.0, query_time_indices=positions) + + +def test_a_condition_keeps_its_positions_as_an_immutable_tuple() -> None: + """A frozen condition must not change because the caller's list did.""" + positions = [0, 2] + condition = EqualityCondition( + column="driver", value=1.0, query_time_indices=positions + ) + positions.append(1) + + assert condition.query_time_indices == (0, 2) + + +def test_one_condition_and_a_list_of_them_resolve_alike() -> None: + """A single condition is the one-element list, so callers need not wrap it.""" + pin = EqualityCondition(column="driver", value=1.0) + + assert resolve_conditions(pin) == (pin,) + assert resolve_conditions([pin]) == (pin,) + + +@pytest.mark.parametrize( + ("first", "second", "overlap"), + [ + ((0,), (0, 1), r"query_time_indices \[0\]"), + (None, (1,), "every position"), + (None, None, "every position"), + ], + ids=["explicit", "null_against_explicit", "null_against_null"], +) +def test_two_conditions_on_one_column_at_a_shared_position_are_refused( + first: tuple[int, ...] | None, second: tuple[int, ...] | None, overlap: str +) -> None: + """A column carries one condition per position, whatever the two kinds are.""" + with pytest.raises( + ValueError, + match=rf"'driver' carries more than one condition at {overlap}: " + r"condition\[1\] overlaps condition\[0\]", + ): + resolve_conditions( + [ + EqualityCondition(column="driver", value=1.0, query_time_indices=first), + IntervalCondition( + column="driver", lower=0.0, upper=2.0, query_time_indices=second + ), + ] ) -def test_a_block_reports_the_kinds_a_deployment_must_advertise() -> None: - """The gate needs the kinds, and each kind appears once however many columns use it.""" - block = ConditionBlock( - query_time_index=0, - conditions=( +def test_one_column_may_carry_conditions_at_disjoint_positions() -> None: + """Disjoint positions never meet, so each keeps its own condition.""" + conditions = resolve_conditions( + [ + EqualityCondition(column="driver", value=1.0, query_time_indices=[0]), + EqualityCondition(column="driver", value=2.0, query_time_indices=[1]), + ] + ) + + assert len(conditions) == 2 + + +def test_different_columns_may_share_every_position() -> None: + """A pin and a band on different columns meet at a position; that is the point.""" + conditions = resolve_conditions( + [ EqualityCondition(column="driver", value=1.0), - IntervalCondition(column="target", lower=0.0, upper=None), - ), + IntervalCondition(column="hedge", lower=0.0, upper=None), + ] ) - assert block.kinds == ("equality", "interval") - assert block.pinned_columns == ("driver",) - assert block.conditioned_columns == ("driver", "target") + assert condition_kinds(conditions) == ("equality", "interval") -def test_the_mode_and_the_block_must_agree() -> None: +@pytest.mark.parametrize( + ("condition", "message"), + [([], "must not be empty"), (["driver"], "must be an EqualityCondition")], + ids=["empty", "not_a_condition"], +) +def test_a_condition_list_must_hold_conditions(condition: Any, message: str) -> None: + """An empty list conditions nothing, and only the two condition kinds exist.""" + with pytest.raises(ValueError, match=message): + resolve_conditions(condition) + + +def test_the_mode_and_the_condition_must_agree() -> None: """One without the other means a different request than the caller wrote.""" - block = ConditionBlock( - query_time_index=0, - conditions=(EqualityCondition(column="driver", value=1.0),), - ) + pin = EqualityCondition(column="driver", value=1.0) with pytest.raises(ValueError, match="go together"): - _request(block, query_mode="forecast") + _request(pin, query_mode="forecast") with pytest.raises(ValueError, match="go together"): _request(None) @pytest.mark.parametrize( - ("block", "requested_columns", "message"), + ("condition", "requested_columns", "message"), [ ( - ConditionBlock( - query_time_index=0, - conditions=(EqualityCondition(column="absent", value=1.0),), - ), + EqualityCondition(column="absent", value=1.0), ("target",), "undeclared columns", ), ( - ConditionBlock( - query_time_index=9, - conditions=(EqualityCondition(column="driver", value=1.0),), - ), + EqualityCondition(column="driver", value=1.0, query_time_indices=[1, 9]), ("target",), - "outside the 2 requested future positions", + r"query_time_indices \[9\] outside the 2 requested future positions", ), ( - ConditionBlock( - query_time_index=0, - conditions=( - EqualityCondition(column="driver", value=1.0), - EqualityCondition(column="hedge", value=0.5), - EqualityCondition(column="target", value=2.0), - ), - ), + [ + EqualityCondition(column="driver", value=1.0), + EqualityCondition(column="hedge", value=0.5), + EqualityCondition(column="target", value=2.0, query_time_indices=[1]), + ], None, - "log_prob", + "pinning every column at query_time_indices position 1.*log_prob", ), ( - ConditionBlock( - query_time_index=0, - conditions=( - EqualityCondition(column="driver", value=1.0), - EqualityCondition(column="hedge", value=0.5), - IntervalCondition(column="target", lower=0.0, upper=None), + [ + EqualityCondition(column="driver", value=1.0, query_time_indices=[0]), + EqualityCondition(column="hedge", value=0.5, query_time_indices=[0]), + IntervalCondition( + column="target", lower=0.0, upper=None, query_time_indices=[0] ), - ), + ], None, - "at least one column unconditioned", + "at least one column unconditioned at query_time_indices position 0", ), ], + ids=["undeclared", "outside", "all_pinned", "all_conditioned"], ) def test_a_request_is_checked_against_its_own_schema_before_any_round_trip( - block: ConditionBlock, requested_columns: tuple[str, ...] | None, message: str + condition: Condition | list[Condition], + requested_columns: tuple[str, ...] | None, + message: str, ) -> None: - """The caller learns which column name it got wrong without paying for a request.""" + """The caller learns which column or position it got wrong without a request.""" with pytest.raises(ValueError, match=message): - _request(block, requested_columns=requested_columns) + _request(condition, requested_columns=requested_columns) + + +def test_every_column_may_be_conditioned_somewhere_if_never_all_at_once() -> None: + """The read-out rule holds per position, not over the request as a whole.""" + request = _request( + [ + EqualityCondition(column="driver", value=1.0), + EqualityCondition(column="hedge", value=0.5, query_time_indices=[0]), + EqualityCondition(column="target", value=2.0, query_time_indices=[1]), + ], + requested_columns=None, + ) + + assert len(request.to_payload()["conditions"]) == 3 def test_an_interval_column_may_still_be_read_out() -> None: """An interval fixes a region, so the column's distribution inside it is an answer.""" - block = ConditionBlock( - query_time_index=0, - conditions=(IntervalCondition(column="driver", lower=0.5, upper=1.5),), + request = _request( + IntervalCondition(column="driver", lower=0.5, upper=1.5), + requested_columns=("driver", "target"), ) - request = _request(block, requested_columns=("driver", "target")) - assert request.to_payload()["requested_columns"] == ["driver", "target"] @@ -247,13 +325,11 @@ def test_a_pinned_column_may_be_read_out_in_any_position() -> None: Its answer is the request's own value, which is what lets a scenario frame be compared against an unconditioned one column for column. """ - block = ConditionBlock( - query_time_index=0, - conditions=(EqualityCondition(column="driver", value=1.0),), + request = _request( + EqualityCondition(column="driver", value=1.0), + requested_columns=("target", "driver"), ) - request = _request(block, requested_columns=("target", "driver")) - assert request.to_payload()["requested_columns"] == ["target", "driver"] @@ -263,18 +339,15 @@ def test_omitting_the_projection_states_nothing_on_the_wire() -> None: A client-side default would be a second copy of the rule, free to drift from the one the deployment actually applies. """ - block = ConditionBlock( - query_time_index=0, - conditions=(EqualityCondition(column="driver", value=1.0),), - ) - - payload = _request(block, requested_columns=None).to_payload() + payload = _request( + EqualityCondition(column="driver", value=1.0), requested_columns=None + ).to_payload() assert "requested_columns" not in payload -def test_the_payload_carries_the_block_the_service_parses() -> None: - """The wire form names its position once and each condition names its kind.""" +def test_the_payload_carries_the_conditions_the_service_parses() -> None: + """Each condition names its kind and its positions, ``null`` covering them all.""" payload = build_forecast_payload( model_version=_MODEL_VERSION, schema=_schema(), @@ -282,51 +355,56 @@ def test_the_payload_carries_the_block_the_service_parses() -> None: query_times=(2, 3), requested_columns=("target",), query_mode="condition", - condition=ConditionBlock( - query_time_index=1, - conditions=( - EqualityCondition(column="driver", value=1.5), - IntervalCondition(column="target", lower=None, upper=12.0), + condition=[ + EqualityCondition(column="driver", value=1.5), + IntervalCondition( + column="target", lower=None, upper=12.0, query_time_indices=[1] ), - ), + ], ) assert payload["query_mode"] == "condition" - assert payload["condition"] == { - "query_time_index": 1, - "conditions": [ - {"column": "driver", "kind": "equality", "value": 1.5}, - {"column": "target", "kind": "interval", "lower": None, "upper": 12.0}, - ], - } + assert "condition" not in payload + assert payload["conditions"] == [ + { + "column": "driver", + "kind": "equality", + "value": 1.5, + "query_time_indices": None, + }, + { + "column": "target", + "kind": "interval", + "lower": None, + "upper": 12.0, + "query_time_indices": [1], + }, + ] def test_the_gate_refuses_a_deployment_that_advertises_no_condition_mode() -> None: """Discovery before the request is the point: the error must not cost a round trip.""" - block = ConditionBlock( - query_time_index=0, - conditions=(EqualityCondition(column="driver", value=1.0),), - ) + pin = EqualityCondition(column="driver", value=1.0) with pytest.raises( UnsupportedServiceContractError, match="does not serve condition" ): require_condition_support( - _health(query_modes=("forecast",), condition_kinds=()), block + _health(query_modes=("forecast",), condition_kinds=()), pin ) def test_the_gate_refuses_a_kind_the_deployment_does_not_answer() -> None: - """A deployment may serve one kind before the other; the block says which it needs.""" - block = ConditionBlock( - query_time_index=0, - conditions=(IntervalCondition(column="driver", lower=0.0, upper=1.0),), - ) + """A deployment may serve one kind before the other; the conditions say which.""" + conditions = [ + EqualityCondition(column="driver", value=1.0), + IntervalCondition(column="hedge", lower=0.0, upper=1.0), + ] with pytest.raises(UnsupportedServiceContractError, match=r"kinds \['interval'\]"): - require_condition_support(_health(condition_kinds=("equality",)), block) + require_condition_support(_health(condition_kinds=("equality",)), conditions) - require_condition_support(_health(), block) + require_condition_support(_health(), conditions) def test_the_equality_response_reads_back_its_plausibility( @@ -348,20 +426,37 @@ def test_the_equality_response_reads_back_its_plausibility( assert response.diagnostics.interval_estimator is None -def test_the_response_describes_the_conditioned_position_alone( +def test_the_response_answers_every_requested_position( json_fixture_loader: Callable[[str], dict[str, Any]], ) -> None: - """The request asked for two future rows; a condition answers about one.""" + """The condition covers one of two rows, and the answer still carries both.""" request_payload = json_fixture_loader("condition_mean_request") assert request_payload["query_times"] == [2, 3] + assert request_payload["conditions"][0]["query_time_indices"] == [1] response = ForecastResponse.from_payload( json_fixture_loader("condition_mean_response"), request_payload=request_payload, ) - assert response.query_times == (3,) - assert response.diagnostics.horizon_count == 1 + assert response.query_times == (2, 3) + assert response.diagnostics.horizon_count == 2 + + +def test_a_response_narrowed_to_the_covered_position_is_refused( + json_fixture_loader: Callable[[str], dict[str, Any]], +) -> None: + """A deployment answering fewer positions than asked must fail, not parse.""" + response_payload = json_fixture_loader("condition_mean_response") + response_payload["outputs"]["query_times"] = [3] + response_payload["outputs"]["mean"] = [[13.5]] + response_payload["diagnostics"]["horizon_count"] = 1 + + with pytest.raises(ValueError, match="query_times"): + ForecastResponse.from_payload( + response_payload, + request_payload=json_fixture_loader("condition_mean_request"), + ) def test_an_interval_response_reads_back_its_estimator_accuracy( @@ -418,7 +513,7 @@ def post_json(self, url: str, payload: Mapping[str, Any]) -> Mapping[str, Any]: assert isinstance(sample_count, int) start = sum(cast(int, earlier["n_samples"]) for earlier in self.payloads[:-1]) return { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": _MODEL_VERSION, "checkpoint_version": "sdk-test", @@ -426,10 +521,12 @@ def post_json(self, url: str, payload: Mapping[str, Any]) -> Mapping[str, Any]: "query_mode": "condition", "return_mode": "samples", "outputs": { - "query_times": [3], + "query_times": [2, 3], "requested_columns": ["target"], "mean": None, - "samples": [[[float(start + index)]] for index in range(sample_count)], + "samples": [ + [[-1.0], [float(start + index)]] for index in range(sample_count) + ], "quantiles": None, }, "plausibility": { @@ -438,7 +535,7 @@ def post_json(self, url: str, payload: Mapping[str, Any]) -> Mapping[str, Any]: }, "diagnostics": { "history_rows": 2, - "horizon_count": 1, + "horizon_count": 2, "seed": payload.get("seed"), "condition_draws": sample_count, "interval_estimator": None, @@ -456,7 +553,7 @@ def _health_payload( """Build one health advertisement as the service serializes it.""" return { "status": "ok", - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": _MODEL_VERSION, "checkpoint_version": "sdk-test", @@ -490,16 +587,13 @@ def _client(transport: _ConditionTransport) -> JointFMClient: {"driver": 1.0, "hedge": 0.5, "target": 10.0}, {"driver": 1.1, "hedge": 0.6, "target": 11.0}, ) -_PIN_DRIVER = ConditionBlock( - query_time_index=1, - conditions=(EqualityCondition(column="driver", value=1.5),), -) +_PIN_DRIVER = EqualityCondition(column="driver", value=1.5, query_time_indices=[1]) def test_the_client_sends_the_condition_and_reads_the_answer_back( json_fixture_loader: Callable[[str], dict[str, Any]], ) -> None: - """The typed helper carries the block onto the wire and types what comes back.""" + """The typed helper carries the condition onto the wire and types what comes back.""" transport = _ConditionTransport( health_payload=_health_payload( query_modes=["forecast", "condition"], @@ -520,13 +614,13 @@ def test_the_client_sends_the_condition_and_reads_the_answer_back( assert isinstance(result, MeanForecastResult) assert result.query_mode == "condition" - assert result.query_times == (3,) - assert result.mean == ((13.5,),) + assert result.query_times == (2, 3) + assert result.mean == ((12.0,), (13.5,)) assert result.plausibility == ConditionPlausibility(equality_log_density=-1.27) assert len(transport.payloads) == 1 sent = transport.payloads[0] assert sent["query_mode"] == "condition" - assert sent["condition"] == _PIN_DRIVER.to_payload() + assert sent["conditions"] == [_PIN_DRIVER.to_payload()] def test_the_client_refuses_before_posting_when_the_deployment_cannot_condition() -> ( @@ -552,7 +646,7 @@ def test_the_client_refuses_before_posting_when_the_deployment_cannot_condition( def test_batched_condition_samples_merge_into_one_conditional_answer() -> None: - """Every batch repeats the same block; the merge recounts the draws and keeps the plausibility.""" + """Every batch repeats the same conditions; the merge recounts draws, keeps plausibility.""" transport = _ConditionTransport( health_payload=_health_payload( query_modes=["forecast", "condition"], @@ -573,20 +667,24 @@ def test_batched_condition_samples_merge_into_one_conditional_answer() -> None: ) assert isinstance(result, SampleForecastResult) - assert result.samples == (((0.0,),), ((1.0,),), ((2.0,),)) - assert result.query_times == (3,) + assert result.samples == ( + ((-1.0,), (0.0,)), + ((-1.0,), (1.0,)), + ((-1.0,), (2.0,)), + ) + assert result.query_times == (2, 3) assert result.diagnostics.condition_draws == 3 assert result.plausibility == ConditionPlausibility(equality_log_density=-1.27) assert [payload["n_samples"] for payload in transport.payloads] == [2, 1] assert all( payload["query_mode"] == "condition" - and payload["condition"] == _PIN_DRIVER.to_payload() + and payload["conditions"] == [_PIN_DRIVER.to_payload()] for payload in transport.payloads ) def test_the_dataframe_adapter_builds_a_condition_request() -> None: - """A pandas caller passes the block and gets the condition envelope.""" + """A pandas caller passes the conditions and gets the condition envelope.""" pandas = pytest.importorskip("pandas") frame = pandas.DataFrame(list(_HISTORY_ROWS)) @@ -601,7 +699,7 @@ def test_the_dataframe_adapter_builds_a_condition_request() -> None: ) assert payload["query_mode"] == "condition" - assert payload["condition"] == _PIN_DRIVER.to_payload() + assert payload["conditions"] == [_PIN_DRIVER.to_payload()] assert payload["requested_columns"] == ["target"] diff --git a/tests/test_configuration.py b/tests/test_configuration.py index 3a6fe7b..88344fe 100644 --- a/tests/test_configuration.py +++ b/tests/test_configuration.py @@ -47,7 +47,7 @@ class _HealthTransport: _METADATA: dict[str, object] = { "status": "ok", - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.yaml", "checkpoint_version": "yaml", @@ -147,7 +147,7 @@ def test_load_settings_layers_config_below_dotenv_and_environment( "datarobot_endpoint": "https://app.datarobot.com/api/v2", "datarobot_api_token": "yaml-token", "deployment_id": "yaml-deployment-id", - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.yaml", } }, @@ -191,7 +191,7 @@ def test_client_from_env_uses_transport_defaults_from_config( "datarobot_endpoint": "https://app.datarobot.com/api/v2", "datarobot_api_token": "yaml-token", "deployment_id": "yaml-deployment-id", - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.yaml", }, "transport": { diff --git a/tests/test_contract.py b/tests/test_contract.py index 9d80ad3..30ee0a1 100644 --- a/tests/test_contract.py +++ b/tests/test_contract.py @@ -76,7 +76,7 @@ def _health_metadata() -> dict[str, object]: """Health metadata.""" return { "status": "ok", - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", @@ -111,7 +111,7 @@ def test_package_identity_contract() -> None: assert DISTRIBUTION_NAME == "jointfm-client" assert IMPORT_NAMESPACE == "jointfm_client" assert FIRST_SUPPORTED_PYTHON_VERSION == "3.11" - assert SCHEMA_VERSION == "v4" + assert SCHEMA_VERSION == "v5" def test_package_version_matches_installed_distribution() -> None: @@ -267,7 +267,7 @@ def test_forecast_payload_matches_service_contract_without_mutating_inputs() -> ) assert payload == { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "query_mode": "forecast", "return_mode": "quantiles", @@ -362,7 +362,7 @@ def test_dataframe_payload_matches_service_forecast_request_shape() -> None: ) assert payload == { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "query_mode": "forecast", "return_mode": "mean", @@ -777,7 +777,7 @@ def test_health_and_response_models_parse_current_payloads() -> None: health = HealthMetadata.from_payload(_health_metadata()) response = ForecastResponse.from_payload( { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", @@ -840,7 +840,7 @@ def test_forecast_result_conversion_helpers_cover_mean_samples_and_quantiles() - """Forecast result conversion helpers cover mean samples and quantiles.""" mean_result = ForecastResponse.from_payload( { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", @@ -860,7 +860,7 @@ def test_forecast_result_conversion_helpers_cover_mean_samples_and_quantiles() - ) sample_result = ForecastResponse.from_payload( { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", @@ -883,7 +883,7 @@ def test_forecast_result_conversion_helpers_cover_mean_samples_and_quantiles() - ) quantile_result = ForecastResponse.from_payload( { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", @@ -957,7 +957,7 @@ def test_forecast_result_conversion_helpers_cover_mean_samples_and_quantiles() - def test_forecast_response_validates_request_scoped_shapes() -> None: """Forecast response validates request scoped shapes.""" request_payload = { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "query_mode": "forecast", "return_mode": "samples", @@ -966,7 +966,7 @@ def test_forecast_response_validates_request_scoped_shapes() -> None: "n_samples": 3, } response_payload = { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", @@ -996,7 +996,7 @@ def test_forecast_response_raises_typed_error_for_success_payload_errors() -> No with pytest.raises(JointFMServiceError) as exc_info: ForecastResponse.from_payload( { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", @@ -1162,7 +1162,7 @@ def post_json(self, url: str, payload: Mapping[str, Any]) -> Mapping[str, Any]: def _forecast_response_payload(*, return_mode: str) -> dict[str, object]: """Forecast response payload.""" return { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", diff --git a/tests/test_contract_models.py b/tests/test_contract_models.py index af48ec8..7352acd 100644 --- a/tests/test_contract_models.py +++ b/tests/test_contract_models.py @@ -76,7 +76,7 @@ def test_request_models_serialize_direct_payloads_without_mutating_inputs() -> N ) assert metadata.to_payload() == { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "query_mode": "forecast", "return_mode": "quantiles", @@ -91,7 +91,7 @@ def test_request_models_serialize_direct_payloads_without_mutating_inputs() -> N "timezone": "UTC", } assert request.to_payload() == { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "query_mode": "forecast", "return_mode": "quantiles", @@ -303,7 +303,7 @@ def test_response_models_reject_direct_validation_edges() -> None: def test_forecast_response_rejects_request_scoped_metadata_mismatches() -> None: """Forecast response rejects request scoped metadata mismatches.""" request_payload = { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "query_mode": "forecast", "return_mode": "mean", @@ -354,7 +354,7 @@ def test_forecast_response_rejects_request_scoped_metadata_mismatches() -> None: def test_forecast_response_rejects_sample_bound_violations() -> None: """Forecast response rejects sample bound violations.""" request_payload = { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "query_mode": "forecast", "return_mode": "samples", @@ -397,7 +397,7 @@ def test_forecast_response_rejects_sample_bound_violations() -> None: def _mean_response_payload() -> dict[str, Any]: """Mean response payload.""" return { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", @@ -419,7 +419,7 @@ def _mean_response_payload() -> dict[str, Any]: def _sample_response_payload() -> dict[str, Any]: """Sample response payload.""" return { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.smoke-1", "checkpoint_version": "smoke-1", diff --git a/tests/test_feature_importance.py b/tests/test_feature_importance.py index c706175..273d84a 100644 --- a/tests/test_feature_importance.py +++ b/tests/test_feature_importance.py @@ -47,7 +47,7 @@ def post_json(self, url: str, payload: Mapping[str, Any]) -> Mapping[str, Any]: self.payloads.append(dict(payload)) samples = self.sample_batches[len(self.payloads) - 1] return { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": _MODEL_VERSION, "checkpoint_version": "sdk-test", diff --git a/tests/test_log_prob_mode.py b/tests/test_log_prob_mode.py index 9c0cd34..2715b53 100644 --- a/tests/test_log_prob_mode.py +++ b/tests/test_log_prob_mode.py @@ -35,7 +35,7 @@ from jointfm_client import ( ColumnSpec, - ConditionBlock, + Condition, ConditionPlausibility, DataFrameSchema, EqualityCondition, @@ -61,10 +61,7 @@ {"driver": 1.2, "target": 12.0}, {"driver": 1.5, "target": 13.0}, ) -_PIN_DRIVER = ConditionBlock( - query_time_index=1, - conditions=(EqualityCondition(column="driver", value=1.5),), -) +_PIN_DRIVER = EqualityCondition(column="driver", value=1.5, query_time_indices=[1]) def _schema() -> DataFrameSchema: @@ -83,7 +80,7 @@ def _request( return_mode: ReturnMode = "log_prob", query_rows: Any = _QUERY_ROWS, requested_columns: tuple[str, ...] | None = ("driver", "target"), - condition: ConditionBlock | None = None, + condition: Condition | None = None, ) -> ForecastRequest: """Build one scoring request against the two-column schema.""" return ForecastRequest( @@ -218,7 +215,7 @@ def test_a_score_may_leave_out_no_column_at_all_under_a_condition() -> None: def test_the_payload_carries_the_rows_the_service_scores( json_fixture_loader: Callable[[str], dict[str, Any]], fixture_name: str, - condition: ConditionBlock | None, + condition: Condition | None, requested_columns: tuple[str, ...], query_rows: tuple[Mapping[str, Any], ...], ) -> None: @@ -324,7 +321,7 @@ def test_the_client_scores_observed_rows_and_types_the_answer( def test_the_client_scores_under_a_condition( json_fixture_loader: Callable[[str], dict[str, Any]], ) -> None: - """A conditioned score answers at the conditioned position and reports its pin.""" + """A conditioned score answers every position and reports its pin.""" health_payload = json_fixture_loader("health_metadata") health_payload["supported_query_modes"] = ["forecast", "condition"] health_payload["supported_condition_kinds"] = ["equality", "interval"] @@ -345,8 +342,8 @@ def test_the_client_scores_under_a_condition( ) assert isinstance(result, LogProbResult) - assert result.query_times == (3,) - assert result.log_prob.values == (-1.75,) + assert result.query_times == (2, 3) + assert result.log_prob.values == (-2.5, -1.75) assert result.plausibility == ConditionPlausibility(equality_log_density=-1.27) sent = transport.payloads[0] assert sent["query_mode"] == "condition" diff --git a/tests/test_pool.py b/tests/test_pool.py index 9de17bb..7615efd 100644 --- a/tests/test_pool.py +++ b/tests/test_pool.py @@ -56,7 +56,7 @@ def _health( """Health.""" return { "status": "ok", - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": model_version, "checkpoint_version": checkpoint_version, @@ -134,7 +134,7 @@ def test_pool_retries_next_instance_on_470() -> None: pool = _pool(transport=_Transport(fail_ids=frozenset({"a"}))) assert pool.next_instance().deployment_id == "a" assert pool.next_instance().deployment_id == "b" - assert pool.post_json({"schema_version": "v4"}) == { + assert pool.post_json({"schema_version": "v5"}) == { "ok": True, "deployment_id": "b", } @@ -144,7 +144,7 @@ def test_pool_raises_when_all_instances_unavailable() -> None: """Pool raises when all instances unavailable.""" pool = _pool(transport=_Transport(fail_ids=frozenset({"a", "b"}))) with pytest.raises(JointFMHTTPStatusError, match="unavailable"): - pool.post_json({"schema_version": "v4"}) + pool.post_json({"schema_version": "v5"}) def test_pool_health_rejects_mismatch_and_aligns_sample_cap() -> None: @@ -275,7 +275,7 @@ def post_json(self, url: str, payload: Mapping[str, Any]) -> Mapping[str, Any]: executor.submit( pool.post_json_to, pool.instance_at(index), - {"schema_version": "v4"}, + {"schema_version": "v5"}, ) for index in range(2) ] @@ -299,7 +299,7 @@ def test_pool_health_routes_only_reachable_peers() -> None: assert pool.instance_at(0).deployment_id == "b" assert pool.instance_at(1).deployment_id == "b" assert pool.next_instance().deployment_id == "b" - assert pool.post_json({"schema_version": "v4"})["deployment_id"] == "b" + assert pool.post_json({"schema_version": "v5"})["deployment_id"] == "b" def test_pool_failover_retries_health_excluded_peer() -> None: @@ -313,7 +313,7 @@ def test_pool_failover_retries_health_excluded_peer() -> None: assert pool.instance_at(0).deployment_id == "b" transport.fail_ids = frozenset({"b"}) - assert pool.post_json({"schema_version": "v4"})["deployment_id"] == "a" + assert pool.post_json({"schema_version": "v5"})["deployment_id"] == "a" assert pool.instance_at(0).deployment_id == "a" @@ -333,7 +333,7 @@ def test_pool_health_skips_incompatible_peer_when_another_matches_pin() -> None: assert metadata.model_version == pinned assert pool.instance_at(0).deployment_id == "a" assert pool.instance_at(1).deployment_id == "a" - assert pool.post_json({"schema_version": "v4"})["deployment_id"] == "a" + assert pool.post_json({"schema_version": "v5"})["deployment_id"] == "a" def test_pool_cooldown_restores_peer_after_transient_failure( @@ -345,7 +345,7 @@ def test_pool_cooldown_restores_peer_after_transient_failure( transport = _Transport(fail_ids=frozenset({"a"})) pool = _pool(transport=transport, peer_cooldown_seconds=10.0) - assert pool.post_json({"schema_version": "v4"})["deployment_id"] == "b" + assert pool.post_json({"schema_version": "v5"})["deployment_id"] == "b" assert pool.instance_at(0).deployment_id == "b" assert pool.instance_at(1).deployment_id == "b" @@ -355,4 +355,4 @@ def test_pool_cooldown_restores_peer_after_transient_failure( clock["now"] = 110.0 assert {pool.instance_at(i).deployment_id for i in range(2)} == {"a", "b"} - assert pool.post_json({"schema_version": "v4"})["deployment_id"] in {"a", "b"} + assert pool.post_json({"schema_version": "v5"})["deployment_id"] in {"a", "b"} diff --git a/tests/test_settings.py b/tests/test_settings.py index 13632c0..2769036 100644 --- a/tests/test_settings.py +++ b/tests/test_settings.py @@ -54,7 +54,7 @@ def _hosted_env(**overrides: str) -> dict[str, str]: DATAROBOT_ENDPOINT_ENV: "https://app.datarobot.com/api/v2/", DATAROBOT_API_TOKEN_ENV: "secret-token", JOINTFM_DEPLOYMENT_ID_ENV: "deployment-id", - JOINTFM_SCHEMA_VERSION_ENV: "v4", + JOINTFM_SCHEMA_VERSION_ENV: "v5", JOINTFM_MODEL_VERSION_ENV: "jointfm-inference:0.3.0+ckpt.sdk-test", } env.update(overrides) @@ -66,7 +66,7 @@ def test_load_settings_from_environment_with_deployment_id_builds_hosted_url() - settings = load_settings(env=_hosted_env(), dotenv_path=None) assert settings.datarobot_endpoint == "https://app.datarobot.com/api/v2" - assert settings.schema_version == "v4" + assert settings.schema_version == "v5" assert settings.model_version == "jointfm-inference:0.3.0+ckpt.sdk-test" assert settings.deployment_id == "deployment-id" assert settings.predict_url == ( @@ -82,7 +82,7 @@ def test_load_settings_with_local_service_base_url_builds_direct_urls() -> None: settings = load_settings( env={ JOINTFM_LOCAL_BASE_URL_ENV: "http://127.0.0.1:8080/", - JOINTFM_SCHEMA_VERSION_ENV: "v4", + JOINTFM_SCHEMA_VERSION_ENV: "v5", JOINTFM_MODEL_VERSION_ENV: "jointfm-inference:0.3.0+ckpt.local-test", }, dotenv_path=None, @@ -158,7 +158,7 @@ def test_load_settings_reads_dotenv_without_overriding_environment(tmp_path) -> "DATAROBOT_ENDPOINT=https://app.datarobot.com/api/v2", "DATAROBOT_API_TOKEN=file-token", "JOINTFM_DEPLOYMENT_ID=file-deployment-id", - "JOINTFM_SCHEMA_VERSION=v4", + "JOINTFM_SCHEMA_VERSION=v5", "JOINTFM_MODEL_VERSION=jointfm-inference:0.3.0+ckpt.sdk-test", ] ), @@ -209,7 +209,7 @@ def test_load_settings_rejects_missing_credentials_without_defaults() -> None: env={ DATAROBOT_API_TOKEN_ENV: "secret-token", JOINTFM_DEPLOYMENT_ID_ENV: "deployment-id", - JOINTFM_SCHEMA_VERSION_ENV: "v4", + JOINTFM_SCHEMA_VERSION_ENV: "v5", JOINTFM_MODEL_VERSION_ENV: "jointfm-inference:0.3.0+ckpt.sdk-test", }, dotenv_path=None, @@ -220,7 +220,7 @@ def test_load_settings_rejects_missing_credentials_without_defaults() -> None: env={ DATAROBOT_ENDPOINT_ENV: "https://app.datarobot.com/api/v2", JOINTFM_DEPLOYMENT_ID_ENV: "deployment-id", - JOINTFM_SCHEMA_VERSION_ENV: "v4", + JOINTFM_SCHEMA_VERSION_ENV: "v5", JOINTFM_MODEL_VERSION_ENV: "jointfm-inference:0.3.0+ckpt.sdk-test", }, dotenv_path=None, diff --git a/tests/test_surfaces.py b/tests/test_surfaces.py index bfd2a26..4f79ca6 100644 --- a/tests/test_surfaces.py +++ b/tests/test_surfaces.py @@ -70,7 +70,7 @@ def test_hosted_surface_uses_datarobot_routes_and_auth_headers( health_url=predict_url, predict_url=predict_url, deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version=request_payload["model_version"], deployment_id="deployment-id", ) @@ -209,7 +209,7 @@ def test_hosted_surface_auto_discovers_model_version_when_settings_unpinned( health_url=predict_url, predict_url=predict_url, deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", deployment_id="deployment-id", ) assert settings.model_version is None diff --git a/tests/test_transport.py b/tests/test_transport.py index c5fd1af..d291f99 100644 --- a/tests/test_transport.py +++ b/tests/test_transport.py @@ -154,7 +154,7 @@ def _health_payload( """Health payload.""" return { "status": "ok", - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": model_version, "checkpoint_version": checkpoint_version, @@ -190,7 +190,7 @@ def _forecast_response_payload(*, return_mode: str = "mean") -> dict[str, object else None, } return { - "schema_version": "v4", + "schema_version": "v5", "image_version": "0.3.0", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "checkpoint_version": "sdk-test", @@ -234,7 +234,7 @@ def test_transport_posts_json_with_headers_timeout_and_user_agent() -> None: session.mount("https://", adapter) result = transport.post_json( - "https://example.com/predict", {"schema_version": "v4"} + "https://example.com/predict", {"schema_version": "v5"} ) assert result == {"ok": True} @@ -248,7 +248,7 @@ def test_transport_posts_json_with_headers_timeout_and_user_agent() -> None: assert adapter.kwargs[0]["timeout"] == (1.5, 2.5) request_body = request.body assert isinstance(request_body, bytes) - assert json.loads(request_body.decode("utf-8")) == {"schema_version": "v4"} + assert json.loads(request_body.decode("utf-8")) == {"schema_version": "v5"} def test_transport_from_settings_attaches_hosted_auth_headers_and_closes_session() -> ( @@ -265,7 +265,7 @@ def test_transport_from_settings_attaches_hosted_auth_headers_and_closes_session "deployment-id/predictionsUnstructured" ), deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", deployment_id="deployment-id", ) @@ -293,7 +293,7 @@ def test_transport_from_local_settings_omits_hosted_auth_headers() -> None: health_url="http://127.0.0.1:8080/healthz", predict_url="http://127.0.0.1:8080/predict", deployment_selector="local_service", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.local-test", local_base_url="http://127.0.0.1:8080", ) @@ -319,7 +319,7 @@ def test_transport_retries_retryable_server_responses() -> None: retry_config=JointFMRetryConfig(max_attempts=2) ) - result = transport.post_json(_server_url(server), {"schema_version": "v4"}) + result = transport.post_json(_server_url(server), {"schema_version": "v5"}) assert result == {"ok": True} assert handler.request_count == 2 @@ -351,7 +351,7 @@ def test_transport_retries_html_bodied_gateway_errors() -> None: ), ) - result = transport.post_json(_server_url(server), {"schema_version": "v4"}) + result = transport.post_json(_server_url(server), {"schema_version": "v5"}) assert result == {"ok": True} assert handler.request_count == 2 @@ -380,7 +380,7 @@ def test_transport_raises_status_error_when_gateway_html_persists() -> None: ) with pytest.raises(JointFMHTTPStatusError) as exc_info: - transport.post_json(_server_url(server), {"schema_version": "v4"}) + transport.post_json(_server_url(server), {"schema_version": "v5"}) assert exc_info.value.status_code == HTTPStatus.BAD_GATEWAY assert "502 Bad Gateway" in exc_info.value.response_body_excerpt @@ -444,7 +444,7 @@ def request(self, *args: Any, **kwargs: Any) -> requests.Response: ) result = transport.post_json( - "https://example.com/predict", {"schema_version": "v4"} + "https://example.com/predict", {"schema_version": "v5"} ) assert result == {"ok": True} @@ -463,7 +463,7 @@ def test_status_error_carries_parsed_retry_after_seconds() -> None: ) with pytest.raises(JointFMHTTPStatusError) as exc_info: - transport.post_json(_server_url(server), {"schema_version": "v4"}) + transport.post_json(_server_url(server), {"schema_version": "v5"}) assert exc_info.value.retry_after_seconds == 0.5 assert handler.request_count == 1 @@ -481,7 +481,7 @@ def test_transport_does_not_retry_validation_errors() -> None: ) with pytest.raises(JointFMHTTPStatusError) as exc_info: - transport.post_json(_server_url(server), {"schema_version": "v4"}) + transport.post_json(_server_url(server), {"schema_version": "v5"}) assert exc_info.value.status_code == HTTPStatus.BAD_REQUEST assert exc_info.value.datarobot_request_id == "request-id-1" @@ -507,7 +507,7 @@ def test_transport_rejects_non_json_serializable_payloads() -> None: with pytest.raises(JointFMRequestEncodingError, match="JSON-serializable"): transport.post_json( "https://example.com/predict", - {"schema_version": "v4", "bad": object()}, + {"schema_version": "v5", "bad": object()}, ) @@ -572,14 +572,14 @@ def test_client_predict_uses_configured_transport_and_settings() -> None: health_url="https://app.datarobot.com/api/v2/deployments/deployment-id/healthz", predict_url="https://app.datarobot.com/api/v2/deployments/deployment-id/predictionsUnstructured", deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", deployment_id="deployment-id", ) transport = RecordingTransport() client = JointFMClient(settings=settings, transport=transport) payload = { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", } @@ -598,7 +598,7 @@ def test_client_health_returns_typed_metadata_and_caches_only_when_requested() - health_url="https://app.datarobot.com/api/v2/deployments/deployment-id/predictionsUnstructured", predict_url="https://app.datarobot.com/api/v2/deployments/deployment-id/predictionsUnstructured", deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", deployment_id="deployment-id", ) @@ -629,7 +629,7 @@ def test_client_health_instances_returns_one_entry_for_single_endpoint() -> None health_url="https://app.datarobot.com/api/v2/deployments/deployment-id/predictionsUnstructured", predict_url="https://app.datarobot.com/api/v2/deployments/deployment-id/predictionsUnstructured", deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", deployment_id="deployment-id", ) @@ -691,7 +691,7 @@ def capture_transport( "DATAROBOT_ENDPOINT": "https://app.datarobot.com/api/v2", "DATAROBOT_API_TOKEN": "secret-token", "JOINTFM_DEPLOYMENT_ID": "deployment-id", - "JOINTFM_SCHEMA_VERSION": "v4", + "JOINTFM_SCHEMA_VERSION": "v5", "JOINTFM_MODEL_VERSION": "jointfm-inference:0.3.0+ckpt.sdk-test", }, dotenv_path=None, @@ -713,7 +713,7 @@ def test_client_health_rejects_cached_model_mismatch() -> None: health_url="https://app.datarobot.com/api/v2/deployments/deployment-id/predictionsUnstructured", predict_url="https://app.datarobot.com/api/v2/deployments/deployment-id/predictionsUnstructured", deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", deployment_id="deployment-id", ) @@ -739,7 +739,7 @@ def test_client_hosted_health_posts_request_type_health_to_predict_url() -> None health_url=predict_url, predict_url=predict_url, deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", deployment_id="deployment-id", ) @@ -763,7 +763,7 @@ def test_client_local_health_keeps_get_request_to_healthz_route() -> None: health_url="http://127.0.0.1:8080/healthz", predict_url="http://127.0.0.1:8080/predict", deployment_selector="local_service", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", local_base_url="http://127.0.0.1:8080", ) @@ -805,7 +805,7 @@ def test_client_forecast_builds_payload_from_rows_and_returns_typed_response() - assert result.outputs.mean == ((12.0,),) assert transport.predict_url == "http://localhost:8080/predict" assert transport.payload == { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", "query_mode": "forecast", "return_mode": "mean", @@ -980,7 +980,7 @@ def post_json(self, url: str, payload: Mapping[str, Any]) -> Mapping[str, Any]: health_url="http://127.0.0.1:8080/healthz", predict_url="http://127.0.0.1:8080/predict", deployment_selector="local_service", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", local_base_url="http://127.0.0.1:8080", ) @@ -1032,13 +1032,13 @@ def test_client_predict_raises_typed_service_error_for_success_payload_errors() health_url="https://app.datarobot.com/api/v2/deployments/deployment-id/healthz", predict_url="https://app.datarobot.com/api/v2/deployments/deployment-id/predictionsUnstructured", deployment_selector="deployment_id", - schema_version="v4", + schema_version="v5", model_version="jointfm-inference:0.3.0+ckpt.sdk-test", deployment_id="deployment-id", ) transport = RecordingTransport() transport.predict_payload = { - "schema_version": "v4", + "schema_version": "v5", "errors": [ { "code": "VALIDATION_ERROR", @@ -1052,7 +1052,7 @@ def test_client_predict_raises_typed_service_error_for_success_payload_errors() with pytest.raises(JointFMServiceError) as exc_info: client.predict( { - "schema_version": "v4", + "schema_version": "v5", "model_version": settings.model_version, } ) @@ -1141,7 +1141,7 @@ def do_POST(self) -> None: payload = {"ok": True} else: payload = { - "schema_version": "v4", + "schema_version": "v5", "errors": [ { "code": "VALIDATION_ERROR", @@ -1225,7 +1225,7 @@ def _pool_settings(primary: str, backup: str) -> JointFMSettings: health_url=primary, predict_url=primary, deployment_selector="deployment_ids", - schema_version="v4", + schema_version="v5", instances=( JointFMInstanceSettings(deployment_id="primary-id", predict_url=primary), JointFMInstanceSettings(deployment_id="backup-id", predict_url=backup), @@ -1267,7 +1267,7 @@ def post_json(self, url: str, payload: Mapping[str, Any]) -> Mapping[str, Any]: transport = PoolTransport() client = JointFMClient(settings=settings, transport=transport) payload = { - "schema_version": "v4", + "schema_version": "v5", "model_version": "jointfm-inference:0.3.0+ckpt.sdk-test", } client.predict(payload)