Repository navigation
[E2S 1.0] Established shared Pipeline / coupled model execution primitives - #1212
pzharrington wants to merge 12 commits into
Conversation
|
Auto-sync is disabled for ready for review pull requests in this repository. Workflows must be run manually. Contributors can view more details about this message here. |
|
Disclaimer: This is AI-generated, please review response for accuracy
|
| return fetch_data( | ||
| self.plan.source, | ||
| time=np.array([self.item.time]), | ||
| variable=np.array(requirement.variables), | ||
| lead_time=np.array(requirement.lead_offsets), | ||
| ) |
There was a problem hiding this comment.
If an FCN model is placed on CUDA, this call still fetches its initial condition on CPU. FCN uses that input with CUDA-resident normalization buffers without moving it first, so the first forecast step fails with a device-mismatch error. The existing deterministic workflow fetches data on the model's inference device.
| return fetch_data( | ||
| self.plan.source, | ||
| time=np.array([self.item.time]), | ||
| variable=np.array(requirement.variables), | ||
| lead_time=np.array(requirement.lead_offsets), | ||
| ) |
There was a problem hiding this comment.
When a source's numeric grid differs from the model's declared input grid, the fetched array goes straight to the iterator. FCN rejects the coordinate mismatch, so the forecast cannot start. The existing deterministic workflow maps the source field to the model's input grid before creating the iterator.
| single_lead = input_signature.sizes["lead_time"] == 1 | ||
| same_variables = _labels(input_signature, "variable") == _labels( | ||
| output_signature, "variable" | ||
| ) | ||
| capability = ( | ||
| CheckpointCapability.SNAPSHOTTABLE | ||
| if single_lead and same_variables | ||
| else CheckpointCapability.UNSUPPORTED | ||
| ) |
There was a problem hiding this comment.
Stochastic restart loses RNG state
StormCast has one input lead and matching input and output variables, so this check marks it as restartable. But its steps sample random latents, and the snapshot stores no RNG state. Resuming from that snapshot therefore follows a different random trajectory instead of continuing the interrupted forecast.
| model_type = f"{type(model).__module__}.{type(model).__qualname__}" | ||
| declaration = "|".join( | ||
| [ | ||
| model_type, | ||
| repr(tuple(transform.identity for transform in transforms)), | ||
| spec_fingerprint(spec), | ||
| ";".join( | ||
| f"{port.name}:{','.join(port.variables)}" | ||
| for port in self._output_ports.values() | ||
| ), | ||
| ] | ||
| ) | ||
| self._identity = hashlib.sha256(declaration.encode()).hexdigest()[:16] | ||
| self.compatibility = SnapshotCompatibility( | ||
| schema_version=SNAPSHOT_SCHEMA_VERSION, | ||
| earth2studio_version=earth2studio.__version__, | ||
| plan_identity=self._identity, | ||
| component_versions={name: model_type}, |
There was a problem hiding this comment.
Different weights share snapshot identity
Two instances of the same model class can have different weights yet receive the same identity here, because the hash includes the class and declarations but not the weights. A snapshot from one instance then passes the other's compatibility check, causing old model state to be advanced with different weights and producing an incorrect continuation.
| return {"forecast": x.where(x["lat"] > 0, 0.0)} | ||
|
|
||
|
|
||
| masked = SingleModelPlan(model, source, session=NorthernHemisphere) |
There was a problem hiding this comment.
| # See the License for the specific language governing permissions and | ||
| # limitations under the License. | ||
|
|
||
| """Forecast workflow helpers and simulation supervision. |
There was a problem hiding this comment.
for more information, see https://pre-commit.ci
for more information, see https://pre-commit.ci
Earth2Studio Pull Request
Description
Establishes the interfaces that are likely to be shared between coupled models and
Pipelineexecution so development of each can proceed in parallel. The end goal is to havePipelinebe the package's offering for general execution (handling distribution of work, I/O, and resume), and able to drive either more traditional single-model forecasts like inrun.deterministic, or a coupled model run. See theEXECUTION_CONTRACT_SPEC.mdfor full details. Current code is more draft/sketch, focus should be on the interface definitions.Checklist
Dependencies