Orchestration🔗
Generated from the source docstrings by mkdocstrings at build time.
See the API reference overview for the other groups.
WorkflowManager 🔗
Orchestrates workflow execution by: 1. Setting up environment and credentials 2. Managing authentication tokens 3. Loading entities to process 4. Instantiating extractors 5. Executing Prefect flows with prepared inputs
eds_token
property
🔗
EDS access token for this env, exchanged on first use and cached.
Without this every LRTS extractor exchanges its own token at construction, so a workflow with several LRTS steps performs N exchanges for a credential that is valid 24 hours. Held here, one exchange serves the whole run.
Returns:
| Name | Type | Description |
|---|---|---|
tuple |
|
|
|
expiry margin. |
Raises:
| Type | Description |
|---|---|
ValueError
|
No EDS API token configured for this environment. Raised on first use, not at construction, so unrelated workflows are unaffected. |
__init__ 🔗
__init__(env: str = None, log_level: str = 'INFO', log_to_console: bool = True, log_to_console_only: bool = False, project_root: str = None, output_result_dir: str | None = None, partial_result_dir: str | None = None, cache_dir: str | None = None, storage: Literal['auto', 'local', 's3'] = 'auto', eds_env: str | None = None)
Initialize workflow manager with configuration.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
env
|
str
|
Environment to use ('prod' or 'preprod'). If None, reads from ENVIRONMENT env var. |
None
|
log_level
|
str
|
Logging level (TRACE, DEBUG, INFO, WARNING, ERROR) |
'INFO'
|
log_to_console
|
bool
|
Whether to output logs to console (stderr). |
True
|
log_to_console_only
|
bool
|
When True, skip the file-rotation log sink and
emit only to the console. The env var
|
False
|
project_root
|
str
|
Override for workspace root (where results/, partials/, etc. are created). If None, resolves automatically via pyproject.toml lookup. |
None
|
output_result_dir
|
str | None
|
Override the final-results write path. Defaults
to |
None
|
partial_result_dir
|
str | None
|
Override the partial-results write path. Same
semantics as |
None
|
cache_dir
|
str | None
|
Override the cache directory. Defaults to
|
None
|
eds_env
|
str | None
|
Override the environment for the EarthDaily Data
Studio provider, which LRTS uses. Defaults to This is an escape hatch, not part of the normal flow. It exists only because an account may be entitled to LRTS on one environment and not the other — a back-end entitlement matter, nothing to do with the client. Leave it unset unless you are deliberately crossing environments. It selects both the EDS token and the LRTS data host, which must agree: they are separate identity deployments and a token from one is refused by the other. |
None
|
storage
|
Literal['auto', 'local', 's3']
|
High-level storage mode. Notebook-friendly toggle that sits between the env var and the explicit kwargs. Values:
Explicit |
'auto'
|
refresh_token_if_needed 🔗
Refresh token if expired or close to expiration.
refresh_eds_token_if_needed 🔗
Return a valid EDS token, exchanging one only when needed.
The EDS counterpart of :meth:refresh_token_if_needed. Both providers
now expose the same "give me a usable token" call, which is what
@requires_token -> BaseExtractor.ensure_token_valid() drives.
get_token 🔗
Return (token, expires_at) for an identity provider.
The single entry point extractors use, so adding a provider does not mean touching every call site.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
provider
|
str
|
|
'geosys'
|
initialize_eds 🔗
Exchange the EDS token now rather than on first use.
Opt-in fail-fast, the same role storage="s3" plays for the S3 client:
call it when you would rather a missing or wrong-environment EDS token
surface at startup than midway through a long run.
load_workflow 🔗
Load and validate a multi-step workflow YAML configuration.
Supports two YAML formats: - Legacy flat format: {defaults, analytics: {name: {module, class, ...}}} - New step format: {workflow: {name, settings, steps: [{name, depends_on, ...}]}}
Stores parsed config in self.workflow_cfg. Builds step dependency graph and validates references.
.. note::
Breaking change — input_from field. Steps now accept an
optional input_from field that selects the entity_list source
("original" or another step's name). When omitted, it defaults
to the first depends_on if present, else "original".
Previously, every step received the entity_list passed to
run_workflow() regardless of depends_on. Existing YAMLs that
rely on the old behavior must add input_from: "original"
explicitly. load_workflow() validates input_from references
at load time (must be "original" or a step in an earlier
execution level) and raises ValueError otherwise.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
config_path
|
str
|
Path to the workflow YAML file. |
required |
Returns:
| Name | Type | Description |
|---|---|---|
dict
|
The parsed config — unwrapped for the new step format. The YAML is |
|
dict
|
|
|
dict
|
the INNER dict, so the keys to index are |
|
dict
|
|
|
dict
|
The shape depends on the input format, which is the trap: |
|
dict
|
=================== ==================================== ============== |
|
dict
|
YAML format Returns |
|
dict
|
=================== ==================================== ============== |
|
dict
|
new ( |
|
dict
|
legacy ( |
|
dict
|
=================== ==================================== ============== |
|
dict
|
Reaching for |
|
dict
|
file therefore raises a bare |
|
dict
|
nowhere near the cause. The same object is also stored on the manager |
|
as |
dict
|
attr: |
dict
|
read-only as :attr: |
Examples:
Indexing the returned config (new step format)::
cfg = manager.load_workflow("configuration/workflow.yml")
cfg["name"] # workflow name
cfg["settings"] # settings block
cfg["steps"] # list of step dicts
cfg is manager.workflow_cfg # True — same object
NOT this — it raises KeyError: 'workflow'::
cfg["workflow"]["steps"] # ✗ the wrapper is already stripped
If you genuinely need the wrapped form (round-tripping the YAML, say), use the raw view rather than re-wrapping by hand::
manager.workflow_cfg_raw["workflow"]["steps"]
run_workflow 🔗
run_workflow(entity_list: DataFrame = None, run_prefix: str | None = None, disabled_steps: list[str] | None = None, *, generate_report: bool = False, report_options: dict[str, Any] | None = None) -> dict[str, dict]
Execute all workflow steps respecting dependency order and transforms.
.. note::
Breaking change — per-step input resolution. Each step now
receives an entity_list determined by its input_from field
(or, when omitted, the first depends_on, falling back to
"original"). Previously every step received the entity_list
argument unchanged. Steps that auto-chain from an upstream no
longer need a no-op use_upstream_entities transform; if a step
still requires the original entity_list despite having
depends_on, declare input_from: "original" in the YAML.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
entity_list
|
DataFrame
|
Input entities. If None, uses self.sfd_list. |
None
|
run_prefix
|
str | None
|
Optional per-run tag inserted between the workflow-level
|
None
|
disabled_steps
|
list[str] | None
|
Optional list of step names to disable for this run only.
Equivalent to setting |
None
|
generate_report
|
bool
|
When |
False
|
report_options
|
dict[str, Any] | None
|
Optional dict forwarded to the reporter. Accepted keys split between the constructor and the run-context:
Unknown keys raise |
None
|
Returns:
| Type | Description |
|---|---|
dict[str, dict]
|
Dict mapping step_name -> {results_df, global_errors, failed_ids} |
Flow per step
- Resolve entity_list source via
input_from(or its default) - Run transform (if defined) — reshape entity_list into new entity_list
- Apply condition/filter (if defined) — subset entities
- Instantiate extractor (if defined) — dynamic import from YAML module/class
- Setup parameters — call setup_method with merged settings + step params
- Run extraction — call run_method with entity_list
Transform-only steps (no extractor) store the transform output as results_df.
dag_export_config 🔗
Return a Plotly config dict to pass to fig.show() / write_html.
Bakes in a custom PNG filename and 2× scale so the modebar download
button produces a publication-ready image instead of newplot.png.
export_workflow_dag_html 🔗
export_workflow_dag_html(path: str, *, mode: str = 'results', reporter=None, filename_stem: str | None = None, **kwargs) -> str
Write a standalone HTML file of the workflow DAG with export config baked in.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
str
|
Output path (local or fsspec-compatible URI). |
required |
mode
|
str
|
|
'results'
|
reporter
|
Optional |
None
|
|
filename_stem
|
str | None
|
Stem for the PNG-download filename surfaced by the
modebar button. Defaults to |
None
|
**kwargs
|
Forwarded to the underlying |
{}
|
Returns:
| Type | Description |
|---|---|
str
|
The path written to. |
visualize_workflow 🔗
visualize_workflow(*, height: int = 560, h_spacing: float = 2.6, v_spacing: float = 1.5, max_value_len: int = 70)
Return an interactive Plotly DAG visualization of the loaded workflow.
Nodes are placed left→right by topological level. Color encodes role:
transform + extractor (orange), extractor only (green), transform-only (blue).
Hover reveals extractor / transform setup + run params, dependencies, and
the resolved input_from.
.. note::
Breaking change. Earlier versions of this method returned a plain
text DAG (str). It now returns a plotly.graph_objects.Figure.
Callers that did print(manager.visualize_workflow()) should switch
to manager.visualize_workflow().show() (or simply leave it as the
last expression of a notebook cell to auto-render).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
height
|
int
|
Figure height in pixels. |
560
|
h_spacing
|
float
|
Horizontal spacing between execution levels. |
2.6
|
v_spacing
|
float
|
Vertical spacing between sibling steps within a level. |
1.5
|
max_value_len
|
int
|
Max length for hover param values before truncation. |
70
|
Returns:
| Type | Description |
|---|---|
|
plotly.graph_objects.Figure — call |
visualize_workflow_results 🔗
visualize_workflow_results(reporter=None, *, height: int = 560, h_spacing: float = 2.6, v_spacing: float = 1.5)
Return a Plotly DAG colored by run results (rows / errors / skipped).
Mirrors :meth:visualize_workflow layout but nodes are colored by
outcome instead of role, and each node carries an extraction-KPI line
(rows • errors). Hover shows per-step run details:
rows, errors, failed_ids sample, skip reason, cache stats.
Data source: a :class:~earthdaily.agriculture.reporting.WorkflowRunReporter.
If reporter is None, one is built on-the-fly from
self.workflow_results (so a notebook user who didn't construct
a reporter explicitly still gets a chart). The on-the-fly reporter
is not stored on the manager.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
reporter
|
Optional pre-populated reporter
(e.g. |
None
|
|
height
|
int
|
Figure height in pixels. |
560
|
h_spacing
|
float
|
Horizontal spacing between execution levels. |
2.6
|
v_spacing
|
float
|
Vertical spacing between sibling steps within a level. |
1.5
|
Returns:
| Type | Description |
|---|---|
|
plotly.graph_objects.Figure — call |
|
|
|
|
|
custom-filename PNG download button) to render in a notebook. |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
if no workflow is loaded, or no reporter was passed
and |
select_steps_to_run 🔗
Render an ipywidgets-based step selector with a checkbox per step and a "Run workflow" button.
Toggling a checkbox mutates self.step_map[name]["enabled"] in place,
so the chosen subset is honored by run_workflow() (or any other path
that reads the flag). Clicking "Run workflow" calls run_workflow()
with the given entity_list and run_prefix.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
entity_list
|
DataFrame
|
Entities to pass to |
None
|
run_prefix
|
str | None
|
Optional |
None
|
Returns:
| Type | Description |
|---|---|
|
ipywidgets.VBox — display in a Jupyter notebook to interact. |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
if no workflow is loaded. |
ImportError
|
if |
interactive_workflow 🔗
Return a clickable Plotly DAG of the loaded workflow.
Same layout as visualize_workflow() but returns a
plotly.graph_objects.FigureWidget. Clicking a node toggles its
enabled state in self.step_map (and recolors it to grey).
After toggling the desired subset, call run_workflow() to execute
only the still-enabled steps. State persists on the manager between
interactive_workflow() calls and into run_workflow().
Note: cascade-skip of downstream steps is enforced at runtime by the existing empty-upstream branch — this widget does not pre-shade descendants. Disabled-node colour is the only visual cue.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
**kwargs
|
Forwarded to |
{}
|
Returns:
| Type | Description |
|---|---|
|
ipywidgets.HBox wrapping a |
|
|
the wrapper makes the chart fill the cell width (FigureWidget on |
|
|
its own defaults to ~700px in VS Code / classic notebook). Clicks |
|
|
still fire on the underlying widget; display in a Jupyter |
|
|
notebook to interact (clicks won't fire on a static |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
if no workflow is loaded. |
ImportError
|
if |
inspect_workflow 🔗
Return a structured summary of the loaded workflow for programmatic inspection.
load_seasonfields 🔗
load_seasonfields(sowing_date_gte=None, sowing_date_lte=None, crop_id=None, start_date=None, end_date=None, farm_name=None, fields=None)
Load seasonfields from EarthDaily platform using EntityManager with optional filters. Delegates to EntityManager.load_seasonfields().
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
farm_name
|
str | list[str]
|
Filter by farm name (Field.Farm.Name).
A single name ( |
None
|
fields
|
str
|
Comma-separated list of fields to retrieve (e.g. "id,name,geometry,sowingDate,crop.code,acreage,customerExternalId"). If None, uses the EntityManager default. |
None
|
load_target_dates 🔗
Load target dates from a CSV file for specific dates extraction mode.
The CSV should have a column 'target_date' with dates in YYYY-MM-DD format.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
file_path
|
str
|
Path to CSV file with target dates |
required |
Returns:
| Name | Type | Description |
|---|---|---|
list |
List of target dates as strings |
Example CSV format
target_date 2024-12-02 2024-12-15 2024-12-21
load_sfd_list 🔗
Load seasonfields from a file (SHP, CSV, etc.) Automatically detects and parses target_dates column if present.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
file_path
|
str
|
Path to the file |
required |
geometry_col
|
str
|
Name of geometry column (mainly for CSV files) |
'geometry'
|
crs
|
str
|
Coordinate reference system |
'EPSG:4326'
|
CSV format with specific dates
id,name,target_dates 2anvgxb,K-5-BE-378,2024-12-02|2024-12-15|2024-12-21 x37bdzx,K-5-BE-385,2024-12-02|2024-12-15|2024-12-21
load_seasonfields_batch 🔗
Load seasonfields using batch processing for improved performance. Delegates to EntityManager.load_seasonfields_batch().
BaseExtractor 🔗
Base class for all analytic extractors (Coverage, InSeason, Emergence, etc.) Handles token management, environment setup, shared config, and logging.
Results export format is controlled by self.export_format — "csv"
(default) or "parquet". Set it via config["export_format"] (forwarded
by WorkflowManager from workflow.yml settings.export_format), the
EDAGRO_EXPORT_FORMAT env var, or by assigning extractor.export_format
directly. Parquet requires pyarrow (already a dependency); error files are
always written as CSV.
get_contextualized_logger 🔗
Returns a logger with additional context for specific operations.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
context
|
str
|
Context identifier (e.g., 'API_CALL', 'BULK_PROCESSING') |
required |
Returns:
| Name | Type | Description |
|---|---|---|
logger |
Contextualized loguru logger |
set_column_mapping 🔗
Update column mapping to match user's DataFrame column names.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
mapping
|
dict
|
Maps internal field names to user column names. Example: {"id": "entity_id", "geometry": "wkt", "crop": "crop_type"} |
required |
validate_column_mapping 🔗
Validate that mapped columns actually exist in the entity DataFrame. Call this after set_column_mapping to catch mismatches early.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
entity_list
|
DataFrame
|
The DataFrame that will be used for extraction. |
required |
Returns:
| Name | Type | Description |
|---|---|---|
bool |
True if all mapped columns are found (or have canonical fallbacks). |
validate_entity 🔗
Validate that an entity has all required fields before making an API call. Uses column mapping to resolve field names, with clear error messages.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
entity_data
|
dict | Series
|
Entity data to validate. |
required |
required_fields
|
list[str]
|
List of internal field names required (e.g., ['geometry', 'crop']). |
required |
context
|
str
|
Context string for error messages (e.g., 'LAI extraction'). |
''
|
Returns:
| Name | Type | Description |
|---|---|---|
dict |
{'valid': bool, 'missing': list[str], 'details': str} |
normalize_date
staticmethod
🔗
Normalize a date value to YYYY-MM-DD string format. Handles ISO timestamps (e.g. '2025-04-01T00:00:00'), pd.Timestamp, and datetime objects.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
value
|
Date value (str, pd.Timestamp, datetime, or None). |
required |
Returns:
| Type | Description |
|---|---|
|
str or None: Date in YYYY-MM-DD format, or the original value if not a recognized date. |
get_entity_value 🔗
Resolve an entity field value using the column mapping. Automatically normalizes date fields to YYYY-MM-DD format.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
row
|
dict | Series
|
Entity data row. |
required |
key
|
str
|
Internal/canonical field name (e.g., 'id', 'geometry', 'crop'). |
required |
default
|
Default value if the mapped column is not found. |
None
|
Returns:
| Type | Description |
|---|---|
|
The value from the row at the mapped column name. |
has_entity_field 🔗
Check if an entity field exists in the row using the column mapping.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
row
|
dict | Series
|
Entity data row. |
required |
key
|
str
|
Internal/canonical field name (e.g., 'id', 'geometry'). |
required |
Returns:
| Name | Type | Description |
|---|---|---|
bool |
bool
|
True if the mapped column exists in the row. |
get_mapped_column 🔗
Get the user-facing column name for an internal field.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
key
|
str
|
Internal field name (e.g., 'id', 'geometry'). |
required |
Returns:
| Name | Type | Description |
|---|---|---|
str |
str
|
The mapped column name. |
normalize_entity_row 🔗
Convert a row with user column names to a dict with canonical/internal names. Useful for passing normalized data to API call wrappers.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
row
|
dict | Series
|
Entity data row with user column names. |
required |
Returns:
| Name | Type | Description |
|---|---|---|
dict |
dict
|
Row data with internal/canonical keys. |
validate_api_response 🔗
Validate an API response in format_*_json methods. Returns an empty DataFrame if the response is empty/None, or None if valid.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
response
|
Raw API response (list, dict, or None) |
required | |
entity_id
|
Entity identifier for logging |
required | |
data_type
|
str
|
Label for the data type (e.g., 'coverage', 'weather') |
'data'
|
Returns:
| Name | Type | Description |
|---|---|---|
|
pd.DataFrame: Empty DataFrame if response is empty/None |
||
None |
If response is valid (caller should continue processing) |
configure_output 🔗
Configure output DataFrame formatting. All parameters are optional — by default everything is kept as-is (backward compatible).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
output_mapping
|
dict
|
Rename columns in output. {current_name: desired_name} Example: {"crop.code": "crop", "sowingDate": "sowing_date"} |
None
|
exclude_columns
|
list
|
Columns to drop from output. Replaces the default (["geometry"]) — include "geometry" in your list if you want to keep the default geometry exclusion. Example: ["field.farm.grower.firstname", "geometry"] |
None
|
output_columns
|
list
|
Whitelist of columns to keep (applied after rename). If set, only these columns appear in the output. Example: ["id", "name", "emergence_date", "confirmation_status"] |
None
|
apply_output_format 🔗
Apply output formatting to a results DataFrame. Operations are applied in order: rename → exclude → select. Returns the original DataFrame unchanged if no output config is set.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
df
|
DataFrame
|
Results DataFrame to format. |
required |
Returns:
| Type | Description |
|---|---|
|
pd.DataFrame: Formatted DataFrame. |
apply_cache_setting 🔗
Apply cache enable/disable override from setup methods. Call this from any extractor's setup_*_parameters() method.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
use_cache
|
bool | None
|
True to enable, False to disable, None to keep the current setting (from config). |
None
|
When cache_dir is a remote URI (s3://, gs://, ...), this method
silently keeps the cache disabled regardless of use_cache=True.
The construction-time warning already explained why; per-call setup
methods don't repeat it. Disabling (use_cache=False) always works.
clear_cache 🔗
Delete the cache file for this extractor (and optional parameter set).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
params
|
dict
|
If provided, clears only the cache for these parameters. If None, clears the default (un-parameterized) cache file. |
None
|
cache_info 🔗
Return information about the current cache state.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
params
|
dict
|
Extraction parameters (for cache filename). |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
dict |
dict
|
Cache statistics (path, exists, records, oldest, newest). |
clear_spatial_cache 🔗
Delete the geohash-keyed cache file for this extractor (and parameter set).
spatial_cache_info 🔗
Cells, rows and date span currently held in the geohash-keyed cache.
is_token_expired 🔗
Check if the bearer token is expired or about to expire.
ensure_token_valid 🔗
Ensure the current token is valid. - If part of a workflow, delegate refresh to workflow manager. - If standalone, refresh locally using get_new_token method.
get_new_token 🔗
Default implementation of token refresh logic. Uses EDAuthenticator to get a new token based on environment. Child classes can override this if they need custom token logic.
save_map_file 🔗
Save one raster/map file, and mirror it to output_uri when set.
Shared by the postprocess="file" writers in FLM, Difference and
Zoning. Routes through :mod:earthdaily.agriculture.core._fs, so
save_path may itself be a remote URI (s3://, gs://, …) and
not just a local directory — the plain open(..., "wb") these
writers used before treated an s3:// prefix as a directory name.
When self.output_uri is set the same bytes are written a second
time to <output_uri>/<filename> as the durable copy. A failure on
that write is raised, not swallowed: the caller reports it as a normal
per-entity failure. A raster recorded in the manifest but living only
on a runner is worse than one that failed loudly.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
save_path
|
Directory (local) or prefix (remote URI) for the working copy. |
required | |
filename
|
File name, already sanitized by the caller. |
required | |
data
|
bytes
|
File contents. |
required |
Returns:
| Name | Type | Description |
|---|---|---|
str |
str
|
The path of the working copy — what callers put in |
str
|
|
|
str
|
local so the analysis step can keep opening files off disk. |