38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144 | @flow(
name="process-s3-l1-olci",
flow_run_name="s3-l1-olci-from-{source_l0_run_id}",
)
async def process_s3l1_olci(
flow_params: Level1FlowParams | None = None,
l0_products: list[dict[str, Any]] | None = None,
source_l0_run_id: str = "manual", # pylint: disable=unused-argument
) -> list[dict[str, Any]]:
"""
Sentinel-3 OLCI L1 processing.
The input_products should have been processed before by L0.
``l0_products`` is the raw product list emitted by S3 L0. When supplied by
a Prefect Automation, it is converted here into the four (or more) processor inputs
expected by OLCI L1. ``source_l0_run_id`` provides a short upstream
reference used in the L1 flow-run name.
"""
mission = "3"
# how to use s3-l1-default-setting
flow_parameters = await (flow_params or Level1FlowParams()).resolve(mission)
if l0_products is not None:
orchestration_settings = await read_s3_orchestration_settings()
prepared_inputs = build_olci_l1_input_products(
l0_products,
orchestration_settings.s3_l0_output_collection,
)
flow_parameters.input_products = [FlowInputProduct.model_validate(product) for product in prepared_inputs]
get_run_logger().info(
"Built %d S3 L1 input product(s) from %d raw L0 product(s) received from Automation",
len(flow_parameters.input_products),
len(l0_products),
)
get_run_logger().info(f"Flow params: {flow_parameters}")
# Call DPR flow
products = await call_dpr_flow(
FlowEnvArgs(owner_id=flow_parameters.owner_identifier),
input_products=flow_parameters.input_products,
external_variables={
"start_datetime": flow_parameters.start_datetime,
"end_datetime": flow_parameters.end_datetime,
"satellite": flow_parameters.satellite,
},
dask_cluster_label=flow_parameters.dask_cluster_label,
processor_name=flow_parameters.processor_name,
processor_version=flow_parameters.processor_version,
pipeline=flow_parameters.pipeline,
unit=flow_parameters.unit,
priority=flow_parameters.priority,
processing_mode=flow_parameters.processing_mode,
workflow=flow_parameters.workflow,
generated_product_to_collection_identifier=flow_parameters.generated_product_to_collection_identifier or [],
auxiliary_product_to_collection_identifier=flow_parameters.auxiliary_product_to_collection_identifier or [],
)
flow_run_id = str(runtime.flow_run.id or "unknown")
# Trigger quicklooks independently for all published products.
emit_olci_quicklook_event(
level="l1",
owner_id=flow_parameters.owner_identifier,
products=products,
flow_run_id=flow_run_id,
)
input_products = [
{
"name": "S3OLCIL1",
"item_id": product["id"],
"collection_name": product["collection"],
}
for product in products
if product.get("properties", {}).get("product:type") == "S03OLCEFR"
]
if not input_products:
get_run_logger().warning("No published S03OLCEFR products; skipping the S3 L1 products-ready event")
return products
event_name = products_ready_event_name(mission="3", level="1")
emitted_event = emit_event(
event=event_name,
resource={
"prefect.resource.id": f"rs-python.s3-l1-result.{flow_run_id}",
"prefect.resource.name": "S3 OLCI L1 products",
},
related=[
{
"prefect.resource.id": f"prefect.flow-run.{flow_run_id}",
"prefect.resource.role": "flow-run",
},
],
payload={"flow_run_id": flow_run_id, "input_products": input_products},
)
if emitted_event is None:
get_run_logger().warning(
"Products-ready event was not emitted: event=%s, flow_run_id=%s",
event_name,
flow_run_id,
)
else:
get_run_logger().info(
"Emitted event=%s, event_id=%s, product_count=%d",
event_name,
emitted_event.id,
len(input_products),
)
return products
|