From 69615402da5fa4e6ede3519a4294a7f10bf69beb Mon Sep 17 00:00:00 2001 From: Regan-Koopmans Date: Wed, 30 Sep 2026 09:45:44 +0200 Subject: [PATCH] AGT-38: Add Async Processing API SDK/CLI Add a client for the Async Processing API, served from services.sentinel-hub.com. The API has two operations: submit a request and get its status. Finished requests return 404. - AsyncProcessingClient: create_request, get_request, wait - AsyncProcessingAPI sync wrapper, exposed as Planet().async_processing - async_processing_request: builders for input, output, data sources, responses, and S3/GCS buckets - `planet async-processing request|create|get|wait`, with --deployment for the EU and US hosts - CLI tutorial, SDK reference entries, auth overview update - Unit tests for the builders; respx tests for client and CLI Requires OAuth2 auth. Planet API keys are not accepted by this API. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/auth/auth-overview.md | 8 +- docs/cli/cli-async-processing.md | 166 +++++++++ docs/python/sdk-reference.md | 8 + mkdocs.yml | 1 + planet/__init__.py | 6 +- planet/async_processing_request.py | 286 +++++++++++++++ planet/cli/async_processing.py | 341 ++++++++++++++++++ planet/cli/cli.py | 3 +- planet/clients/__init__.py | 3 + planet/clients/async_processing.py | 183 ++++++++++ planet/http.py | 5 +- planet/sync/async_processing.py | 78 ++++ planet/sync/client.py | 4 + .../integration/test_async_processing_api.py | 179 +++++++++ .../integration/test_async_processing_cli.py | 329 +++++++++++++++++ tests/unit/test_async_processing_request.py | 215 +++++++++++ 16 files changed, 1807 insertions(+), 8 deletions(-) create mode 100644 docs/cli/cli-async-processing.md create mode 100644 planet/async_processing_request.py create mode 100644 planet/cli/async_processing.py create mode 100644 planet/clients/async_processing.py create mode 100644 planet/sync/async_processing.py create mode 100644 tests/integration/test_async_processing_api.py create mode 100644 tests/integration/test_async_processing_cli.py create mode 100644 tests/unit/test_async_processing_request.py diff --git a/docs/auth/auth-overview.md b/docs/auth/auth-overview.md index c8148e03d..c82875f98 100644 --- a/docs/auth/auth-overview.md +++ b/docs/auth/auth-overview.md @@ -80,10 +80,10 @@ user-interactive and machine-to-machine use cases, as described in this guide. 🚧 OAuth2 machine-to-machine (M2M) access tokens are currently available for use with `services.sentinel-hub.com` APIs. Work to support `api.planet.com` is - ongoing. It should also be noted that at this time no API clients for - `services.sentinel-hub.com` APIs have been incorporated into this SDK. - The SDK may still be used to obtain and manage M2M access tokens to - support external applications. + ongoing. The SDK's Async Processing API client + (`planet async-processing`) is served from `services.sentinel-hub.com` + and accepts both user and M2M access tokens. The SDK may also be used to + obtain and manage M2M access tokens to support external applications. ### Planet API Keys Planet API keys are simple fixed strings that may be presented by the client diff --git a/docs/cli/cli-async-processing.md b/docs/cli/cli-async-processing.md new file mode 100644 index 000000000..72d0e3537 --- /dev/null +++ b/docs/cli/cli-async-processing.md @@ -0,0 +1,166 @@ +--- +title: CLI for Async Processing API Tutorial +--- + +## Introduction +The `planet async-processing` command submits and tracks requests to the [Async Processing API](https://docs.planet.com/develop/apis/async-processing/). The API runs large Process API requests in the background. Outputs may be up to 10000 pixels in each dimension. Results go to your S3 or GCS bucket, not back to the CLI. + +The API is served from `services.sentinel-hub.com`. It needs an OAuth2 login. Planet API keys do not work. + +```sh +planet auth login +``` + +For automated jobs, log in with an M2M client: + +```sh +planet auth login --auth-client-id --auth-client-secret +``` + +## Core Workflow + +### Write an Evalscript +An evalscript tells the service how to turn input bands into output pixels. Save this NDVI script as `ndvi.js`: + +```js +//VERSION=3 +function setup() { + return { + input: ["B04", "B08"], + output: { bands: 1, sampleType: "FLOAT32" } + }; +} +function evaluatePixel(s) { + return [(s.B08 - s.B04) / (s.B08 + s.B04)]; +} +``` + +### Generate a Request +`planet async-processing request` builds the request JSON. It sends nothing. + +```sh +planet async-processing request \ + --collection sentinel-2-l2a \ + --bbox 12.44,41.87,12.54,41.93 \ + --time-from 2024-06-01T00:00:00Z \ + --time-to 2024-06-30T23:59:59Z \ + --max-cloud-coverage 20 \ + --evalscript ndvi.js \ + --width 2500 --height 2500 \ + --delivery s3://my-bucket/ndvi \ + --iam-role-arn arn:aws:iam::123456789012:role/planet \ + > request.json +``` + +Give the area as `--bbox`, `--geometry` or both. Coordinates are WGS84 longitude/latitude unless `--crs` says otherwise. Give the size as `--width`/`--height` in pixels, or `--resx`/`--resy` in CRS units. + +`--geometry` accepts a GeoJSON geometry, Feature or FeatureCollection. It may be a string, a file, or `-` for stdin. + +Each `--response IDENTIFIER:FORMAT` adds an output file. The identifier must match an output `id` in the evalscript's `setup()`, or be `userdata`. The default is `default:image/tiff`. + +### Delivery +Results are written to `//`. To choose a different layout, use `` and `` placeholders in the URL: + +```sh +--delivery 's3://my-bucket/ndvi//' +``` + +S3 accepts an IAM role (recommended) or an access key pair: + +```sh +--delivery s3://my-bucket/ndvi --iam-role-arn arn:aws:iam::123456789012:role/planet +--delivery s3://my-bucket/ndvi --aws-access-key-id AKIA... --aws-secret-access-key ... +``` + +Add `--aws-region` if the bucket is in a different region from the deployment. + +GCS needs a service account key file. The CLI base64-encodes it for you: + +```sh +--delivery gs://my-bucket/ndvi --gcs-credentials key.json +``` + +See the [API documentation](https://docs.planet.com/develop/apis/async-processing/) for the bucket permissions the service needs. + +### Stored Evalscripts +To use an evalscript kept in your bucket, pass `--evalscript-url` instead of `--evalscript`. It uses the same credential options as `--delivery`: + +```sh +--evalscript-url s3://my-bucket/scripts/ndvi.js +``` + +### Submit +```sh +planet async-processing create request.json +``` + +```json +{"id": "7d9a1c2e-0000-4000-8000-000000000000", "status": "RUNNING"} +``` + +`create` also reads a JSON string, or `-` for stdin. Generate and submit in one step: + +```sh +planet async-processing request ... | planet async-processing create - +``` + +The API rejects invalid requests here. It also rejects a request when you have reached your concurrent request limit. + +### Track +`get` shows the status of a running request: + +```sh +planet async-processing get 7d9a1c2e-0000-4000-8000-000000000000 +``` + +`wait` polls until the request stops running: + +```sh +planet async-processing wait 7d9a1c2e-0000-4000-8000-000000000000 +``` + +The API reports only running requests. When a request finishes, successfully or not, it disappears. `get` then fails with a not-running message and `wait` exits 0. An unknown ID behaves the same way. + +A finished request does not always mean success. Check the delivery bucket. Results are in `//`. Failures write `error.json`. The service also stores a copy of the request there, with processing cost added after the run. + +Submit, wait, and list the results: + +```sh +id=$(planet async-processing create request.json | jq -r .id) +planet async-processing wait "$id" +aws s3 ls "s3://my-bucket/ndvi/$id/" +``` + +## Deployments +Input data must be hosted on the deployment the request goes to. The default is `aws-eu-central-1`. Use `--deployment` for others: + +```sh +planet async-processing --deployment aws-us-west-2 create request.json +``` + +## Python +The same operations are available in the SDK: + +```python +from planet import Planet, async_processing_request as apr + +pl = Planet() +request = apr.build_request( + input=apr.process_input( + data=[apr.data_source('sentinel-2-l2a', + time_from='2024-06-01T00:00:00Z', + time_to='2024-06-30T23:59:59Z')], + bbox=[12.44, 41.87, 12.54, 41.93]), + output=apr.process_output( + delivery=apr.s3_bucket('s3://my-bucket/ndvi', + iam_role_arn='arn:aws:iam::123456789012:role/planet'), + width=2500, + height=2500, + responses=[apr.response('default', 'image/tiff')]), + evalscript=open('ndvi.js').read()) + +req = pl.async_processing.create_request(request) +pl.async_processing.wait(req['id']) +``` + +Use `planet.AsyncProcessingClient` for the async interface, and pass `base_url` to select a deployment. diff --git a/docs/python/sdk-reference.md b/docs/python/sdk-reference.md index ba7046032..f92c0ec6c 100644 --- a/docs/python/sdk-reference.md +++ b/docs/python/sdk-reference.md @@ -54,6 +54,14 @@ title: Python SDK API Reference rendering: show_root_full_path: false +## ::: planet.AsyncProcessingClient + rendering: + show_root_full_path: false + +## ::: planet.async_processing_request + rendering: + show_root_full_path: false + ## ::: planet.Planet rendering: show_root_full_path: false diff --git a/mkdocs.yml b/mkdocs.yml index b69e5289e..a3a3f1a37 100644 --- a/mkdocs.yml +++ b/mkdocs.yml @@ -88,6 +88,7 @@ nav: - cli/cli-subscriptions.md - cli/cli-destinations.md - cli/cli-quota.md + - cli/cli-async-processing.md - cli/cli-tips-tricks.md - cli/cli-reference.md - "Python": diff --git a/planet/__init__.py b/planet/__init__.py index e80e928e8..c60b35692 100644 --- a/planet/__init__.py +++ b/planet/__init__.py @@ -13,15 +13,17 @@ # See the License for the specific language governing permissions and # limitations under the License. from .http import Session -from . import data_filter, order_request, reporting, subscription_request +from . import async_processing_request, data_filter, order_request, reporting, subscription_request from .__version__ import __version__ # NOQA from .auth import Auth from .auth_builtins import PlanetOAuthScopes -from .clients import DataClient, DestinationsClient, FeaturesClient, MosaicsClient, OrdersClient, QuotaClient, SubscriptionsClient # NOQA +from .clients import AsyncProcessingClient, DataClient, DestinationsClient, FeaturesClient, MosaicsClient, OrdersClient, QuotaClient, SubscriptionsClient # NOQA from .io import collect from .sync import Planet __all__ = [ + 'AsyncProcessingClient', + 'async_processing_request', 'Auth', 'PlanetOAuthScopes', 'collect', diff --git a/planet/async_processing_request.py b/planet/async_processing_request.py new file mode 100644 index 000000000..68b6e67e8 --- /dev/null +++ b/planet/async_processing_request.py @@ -0,0 +1,286 @@ +# Copyright 2026 Planet Labs PBC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); you may not +# use this file except in compliance with the License. You may obtain a copy of +# the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations under +# the License. +"""Functionality for preparing Async Processing API requests. + +A request has three parts: `input` (where and what data), `output` (size, +formats and delivery bucket) and an evalscript (how to process the data). + +Example: + ```python + from planet import async_processing_request as apr + + request = apr.build_request( + input=apr.process_input( + data=[apr.data_source('sentinel-2-l2a', + time_from='2022-06-20T00:00:00Z', + time_to='2022-06-30T23:59:59Z')], + bbox=[426000, 3960000, 462000, 3994000], + crs='http://www.opengis.net/def/crs/EPSG/0/32633'), + output=apr.process_output( + delivery=apr.s3_bucket('s3://my-bucket/ndvi', + iam_role_arn='arn:aws:iam::1:role/sh'), + resx=10, + resy=10), + evalscript=open('ndvi.js').read()) + ``` +""" +from typing import Any, Dict, List, Optional + +from planet.exceptions import ClientError + +DEFAULT_CRS = 'http://www.opengis.net/def/crs/OGC/1.3/CRS84' + + +def build_request(input: dict, + output: dict, + evalscript: Optional[str] = None, + evalscript_reference: Optional[dict] = None) -> dict: + """Prepare an Async Processing API request. + + Exactly one of `evalscript` and `evalscript_reference` must be given. + + Parameters: + input: Bounds and data sources. See + [planet.async_processing_request.process_input][]. + output: Output size, responses and delivery. See + [planet.async_processing_request.process_output][]. + evalscript: The evalscript source. + evalscript_reference: Location of a stored evalscript. See + [planet.async_processing_request.s3_bucket][] and + [planet.async_processing_request.gs_bucket][]. + + Returns: + dict: the request body. + + Raises: + planet.exceptions.ClientError: If neither or both evalscript + options are given. + """ + if (evalscript is None) == (evalscript_reference is None): + raise ClientError( + 'Exactly one of evalscript and evalscript_reference is required.') + + request: Dict[str, Any] = {'input': input, 'output': output} + if evalscript is not None: + request['evalscript'] = evalscript + else: + request['evalscriptReference'] = evalscript_reference + return request + + +def process_input(data: List[dict], + bbox: Optional[List[float]] = None, + geometry: Optional[dict] = None, + crs: Optional[str] = None) -> dict: + """Prepare the `input` part of a request. + + Give `bbox`, `geometry` or both. With both, the image covers the bbox and + data is rendered only inside the geometry. + + Parameters: + data: Data sources. See + [planet.async_processing_request.data_source][]. + bbox: `[minx, miny, maxx, maxy]` in `crs`. + geometry: GeoJSON geometry in `crs`. + crs: CRS URI of `bbox` and `geometry`, e.g. + `http://www.opengis.net/def/crs/EPSG/0/32633`. The server + defaults to WGS84 (CRS84). + + Raises: + planet.exceptions.ClientError: If neither bbox nor geometry is given, + bbox does not have four values, or data is empty. + """ + if bbox is None and geometry is None: + raise ClientError('One or both of bbox and geometry is required.') + if bbox is not None and len(bbox) != 4: + raise ClientError('bbox must have four values: minx,miny,maxx,maxy.') + if not data: + raise ClientError('At least one data source is required.') + + bounds: Dict[str, Any] = {} + if bbox is not None: + bounds['bbox'] = list(bbox) + if geometry is not None: + bounds['geometry'] = geometry + if crs is not None: + bounds['properties'] = {'crs': crs} + return {'bounds': bounds, 'data': list(data)} + + +def data_source(type: str, + time_from: Optional[str] = None, + time_to: Optional[str] = None, + max_cloud_coverage: Optional[float] = None, + mosaicking_order: Optional[str] = None, + id: Optional[str] = None, + processing: Optional[dict] = None) -> dict: + """Prepare one entry of `input.data`. + + Parameters: + type: Data collection type, e.g. `sentinel-2-l2a`, or `byoc-` + for a collection you own. + time_from: Start of the time range, ISO 8601, e.g. + `2022-06-20T00:00:00Z`. + time_to: End of the time range, ISO 8601. + max_cloud_coverage: Maximum scene cloud cover, 0 to 100. + mosaicking_order: `mostRecent`, `leastRecent` or `leastCC`. + id: Identifier used to reference this source in the evalscript when + the request fuses several sources. + processing: Extra processing options, e.g. + `{'upsampling': 'BICUBIC'}`. + """ + source: Dict[str, Any] = {'type': type} + if id is not None: + source['id'] = id + + data_filter: Dict[str, Any] = {} + if time_from is not None or time_to is not None: + time_range = {} + if time_from is not None: + time_range['from'] = time_from + if time_to is not None: + time_range['to'] = time_to + data_filter['timeRange'] = time_range + if max_cloud_coverage is not None: + data_filter['maxCloudCoverage'] = max_cloud_coverage + if mosaicking_order is not None: + data_filter['mosaickingOrder'] = mosaicking_order + if data_filter: + source['dataFilter'] = data_filter + + if processing: + source['processing'] = processing + return source + + +def process_output(delivery: dict, + width: Optional[int] = None, + height: Optional[int] = None, + resx: Optional[float] = None, + resy: Optional[float] = None, + responses: Optional[List[dict]] = None) -> dict: + """Prepare the `output` part of a request. + + Give `width` and `height` in pixels, or `resx` and `resy` in CRS units, + not both. Each dimension may be at most 10000 pixels. + + Parameters: + delivery: Destination bucket. See + [planet.async_processing_request.s3_bucket][] and + [planet.async_processing_request.gs_bucket][]. + width: Image width in pixels. + height: Image height in pixels. + resx: Horizontal resolution in CRS units. + resy: Vertical resolution in CRS units. + responses: Output files. See + [planet.async_processing_request.response][]. The server defaults + to one PNG named `default`. + + Raises: + planet.exceptions.ClientError: If size and resolution are both given, + or either is given incompletely. + """ + size = (width, height) + res = (resx, resy) + if None not in size and None not in res: + raise ClientError('Give width/height or resx/resy, not both.') + for pair, names in ((size, 'width and height'), (res, 'resx and resy')): + if (pair[0] is None) != (pair[1] is None): + raise ClientError(f'{names} must be given together.') + + output: Dict[str, Any] = {} + if width is not None: + output['width'] = width + output['height'] = height + if resx is not None: + output['resx'] = resx + output['resy'] = resy + if responses: + output['responses'] = list(responses) + output['delivery'] = delivery + return output + + +def response(identifier: str = 'default', + format_type: str = 'image/tiff') -> dict: + """Prepare one entry of `output.responses`. + + Parameters: + identifier: Must match an output `id` in the evalscript's `setup()`, + or be `userdata`. + format_type: `image/tiff`, `image/png`, `image/jpeg` or + `application/json`. + """ + return {'identifier': identifier, 'format': {'type': format_type}} + + +def s3_bucket(url: str, + iam_role_arn: Optional[str] = None, + access_key: Optional[str] = None, + secret_access_key: Optional[str] = None, + region: Optional[str] = None) -> dict: + """Amazon S3 location, for delivery or an evalscript reference. + + Give `iam_role_arn` (recommended) or `access_key` and + `secret_access_key`. + + For delivery, results go to `//` unless `url` + contains the `` or `` placeholders. + + Parameters: + url: `s3://bucket/prefix`. + iam_role_arn: Role the service assumes to access the bucket. + access_key: AWS access key ID. + secret_access_key: AWS secret access key. + region: Bucket region, if it differs from the deployment's region. + + Raises: + planet.exceptions.ClientError: If url is not an s3:// URL or + credentials are missing or incomplete. + """ + if not url.startswith('s3://'): + raise ClientError(f'S3 url must start with s3://; got {url!r}.') + has_keys = access_key is not None or secret_access_key is not None + if iam_role_arn is None and not has_keys: + raise ClientError('S3 access requires iam_role_arn or access_key and ' + 'secret_access_key.') + if has_keys and (access_key is None or secret_access_key is None): + raise ClientError( + 'access_key and secret_access_key must be given together.') + + info = {'url': url} + if iam_role_arn is not None: + info['iamRoleARN'] = iam_role_arn + if access_key is not None and secret_access_key is not None: + info['accessKey'] = access_key + info['secretAccessKey'] = secret_access_key + if region is not None: + info['region'] = region + return {'s3': info} + + +def gs_bucket(url: str, credentials: str) -> dict: + """Google Cloud Storage location, for delivery or an evalscript + reference. + + Parameters: + url: `gs://bucket/prefix`. + credentials: Base64-encoded service account key JSON. + + Raises: + planet.exceptions.ClientError: If url is not a gs:// URL. + """ + if not url.startswith('gs://'): + raise ClientError(f'GCS url must start with gs://; got {url!r}.') + return {'gs': {'url': url, 'credentials': credentials}} diff --git a/planet/cli/async_processing.py b/planet/cli/async_processing.py new file mode 100644 index 000000000..3badcbcfe --- /dev/null +++ b/planet/cli/async_processing.py @@ -0,0 +1,341 @@ +# Copyright 2026 Planet Labs PBC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); you may not +# use this file except in compliance with the License. You may obtain a copy of +# the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations under +# the License. +"""Async Processing API CLI""" +import base64 +from contextlib import asynccontextmanager +from typing import List, Tuple + +import click + +from planet import async_processing_request as apr +from planet import exceptions, geojson +from planet.cli.io import echo_json +from planet.clients.async_processing import (DEFAULT_DEPLOYMENT, + DEPLOYMENT_URLS, + AsyncProcessingClient) + +from .cmds import command +from .session import CliSession +from . import types + +FORMATS = ['image/tiff', 'image/png', 'image/jpeg', 'application/json'] + +NOT_RUNNING = ('Request {} is not running. It has finished or does not ' + 'exist. Check the delivery bucket for results or error.json.') + + +@asynccontextmanager +async def async_processing_client(ctx): + async with CliSession(ctx) as sess: + cl = AsyncProcessingClient(sess, base_url=ctx.obj['BASE_URL']) + yield cl + + +@click.group(name='async-processing') # type: ignore +@click.pass_context +@click.option('--deployment', + type=click.Choice(list(DEPLOYMENT_URLS)), + default=DEFAULT_DEPLOYMENT, + show_default=True, + help='Deployment to send requests to. Input data must be ' + 'hosted on the same deployment.') +@click.option('-u', + '--base-url', + default=None, + help='Assign custom base Async Processing API URL. ' + 'Overrides --deployment.') +def async_processing(ctx, deployment, base_url): + """Commands for the Async Processing API. + + Async Processing runs large Process API requests in the background and + delivers results to an S3 or GCS bucket. Outputs may be up to 10000 + pixels in each dimension. + + Requires an OAuth2 login (`planet auth login`). Planet API keys are not + accepted. + + Typical workflow: + + \b + planet async-processing request --collection sentinel-2-l2a \\ + --bbox 12.44,41.87,12.54,41.93 --time-from 2024-06-01T00:00:00Z \\ + --time-to 2024-06-30T23:59:59Z --evalscript ndvi.js \\ + --width 2500 --height 2500 --delivery s3://my-bucket/ndvi \\ + --iam-role-arn arn:aws:iam::123456789012:role/planet > request.json + planet async-processing create request.json + planet async-processing wait + """ + ctx.obj['BASE_URL'] = base_url or DEPLOYMENT_URLS[deployment] + + +def _parse_responses(values: Tuple[str, ...]) -> List[dict]: + """Parse repeated `--response IDENTIFIER:FORMAT` options.""" + responses = [] + for value in values: + identifier, sep, fmt = value.partition(':') + if not sep or not identifier or fmt not in FORMATS: + raise click.BadParameter( + f'expected IDENTIFIER:FORMAT with FORMAT one of ' + f'{", ".join(FORMATS)}; got {value!r}.', + param_hint='--response') + responses.append(apr.response(identifier, fmt)) + return responses + + +def _bucket(url, + iam_role_arn, + access_key, + secret_access_key, + region, + gcs_credentials) -> dict: + """Build an S3 or GCS bucket description from CLI options.""" + if url.startswith('gs://'): + if gcs_credentials is None: + raise click.UsageError(f'--gcs-credentials is required for {url}.') + return apr.gs_bucket(url, gcs_credentials) + return apr.s3_bucket(url, + iam_role_arn=iam_role_arn, + access_key=access_key, + secret_access_key=secret_access_key, + region=region) + + +@command(async_processing, name='request') +@click.option('--collection', + required=True, + help='Data collection type, e.g. sentinel-2-l2a, ' + 'sentinel-1-grd, landsat-ot-l2, or byoc-.') +@click.option('--time-from', + help='Start of the time range, ISO 8601, e.g. ' + '2024-06-01T00:00:00Z.') +@click.option('--time-to', help='End of the time range, ISO 8601.') +@click.option('--max-cloud-coverage', + type=click.FloatRange(0, 100), + help='Maximum scene cloud cover, 0 to 100.') +@click.option('--mosaicking-order', + type=click.Choice(['mostRecent', 'leastRecent', 'leastCC']), + help='Order in which scenes are mosaicked.') +@click.option('--bbox', + type=types.CommaSeparatedFloat(), + help='Bounding box as minx,miny,maxx,maxy in --crs.') +@click.option('--geometry', + type=types.JSON(), + help='GeoJSON geometry, Feature or FeatureCollection in --crs. ' + 'A JSON string, filename, or - for stdin.') +@click.option('--crs', + help='CRS URI of --bbox and --geometry, e.g. ' + 'http://www.opengis.net/def/crs/EPSG/0/32633. ' + 'Defaults to WGS84 longitude/latitude.') +@click.option('--evalscript', + type=click.File('r'), + help='Evalscript file, or - for stdin.') +@click.option('--evalscript-url', + help='s3:// or gs:// URL of a stored evalscript. Uses the ' + 'same credential options as --delivery.') +@click.option('--width', type=click.IntRange(1, 10000), help='Width in px.') +@click.option('--height', type=click.IntRange(1, 10000), help='Height in px.') +@click.option('--resx', + type=float, + help='Horizontal resolution in --crs units. Use with --resy ' + 'instead of --width/--height.') +@click.option('--resy', type=float, help='Vertical resolution in --crs units.') +@click.option('--response', + 'responses', + multiple=True, + default=['default:image/tiff'], + show_default=True, + metavar='IDENTIFIER:FORMAT', + help='Output file. IDENTIFIER matches an output id in the ' + 'evalscript setup(), or userdata. FORMAT is one of ' + f'{", ".join(FORMATS)}. May be repeated.') +@click.option('--delivery', + required=True, + help='s3://bucket/prefix or gs://bucket/prefix. Results are ' + 'written to // unless the URL contains ' + ' or placeholders.') +@click.option('--iam-role-arn', + help='S3: IAM role for the service to assume. Recommended.') +@click.option('--aws-access-key-id', help='S3: access key ID.') +@click.option('--aws-secret-access-key', help='S3: secret access key.') +@click.option('--aws-region', + help='S3: bucket region, if it differs from the deployment.') +@click.option('--gcs-credentials', + type=click.File('rb'), + help='GCS: service account key JSON file.') +async def request(ctx, + collection, + time_from, + time_to, + max_cloud_coverage, + mosaicking_order, + bbox, + geometry, + crs, + evalscript, + evalscript_url, + width, + height, + resx, + resy, + responses, + delivery, + iam_role_arn, + aws_access_key_id, + aws_secret_access_key, + aws_region, + gcs_credentials, + pretty): + """Generate an async processing request. + + Prints the request JSON. Pass it to `planet async-processing create`. + No request is sent to the API. + + Give --bbox, --geometry or both. Give --width/--height or --resx/--resy. + Give --evalscript or --evalscript-url. + + S3 delivery needs --iam-role-arn, or --aws-access-key-id and + --aws-secret-access-key. GCS delivery needs --gcs-credentials. + + Example: + + \b + planet async-processing request --collection sentinel-2-l2a \\ + --bbox 12.44,41.87,12.54,41.93 --time-from 2024-06-01T00:00:00Z \\ + --time-to 2024-06-30T23:59:59Z --evalscript ndvi.js \\ + --width 2500 --height 2500 --delivery s3://my-bucket/ndvi \\ + --iam-role-arn arn:aws:iam::123456789012:role/planet + """ + if (evalscript is None) == (evalscript_url is None): + raise click.UsageError( + 'Give exactly one of --evalscript and --evalscript-url.') + + if geometry is not None: + geometry = geojson.geom_from_geojson(geometry) + + # Read once: delivery and the evalscript reference may both use it. + gcs_creds = None + if gcs_credentials is not None: + gcs_creds = base64.b64encode(gcs_credentials.read()).decode() + + def bucket(url): + return _bucket(url, + iam_role_arn, + aws_access_key_id, + aws_secret_access_key, + aws_region, + gcs_creds) + + body = apr.build_request( + input=apr.process_input(data=[ + apr.data_source(collection, + time_from=time_from, + time_to=time_to, + max_cloud_coverage=max_cloud_coverage, + mosaicking_order=mosaicking_order) + ], + bbox=bbox, + geometry=geometry, + crs=crs), + output=apr.process_output(delivery=bucket(delivery), + width=width, + height=height, + resx=resx, + resy=resy, + responses=_parse_responses(responses)), + evalscript=evalscript.read() if evalscript else None, + evalscript_reference=bucket(evalscript_url) + if evalscript_url else None) + echo_json(body, pretty) + + +@command(async_processing, name='create') +@click.argument('request', type=types.JSON()) +async def create(ctx, request, pretty): + """Submit an async processing request. + + REQUEST is the request JSON: a JSON string, a filename, or - for stdin. + Generate one with `planet async-processing request`. + + Prints the request ID and status. Invalid requests fail here. Errors + during processing are written to error.json in the delivery bucket. + + Example: + + planet async-processing create request.json + """ + async with async_processing_client(ctx) as cl: + result = await cl.create_request(request) + echo_json(result, pretty) + + +@command(async_processing, name='get') +@click.argument('request_id') +async def get(ctx, request_id, pretty): + """Get the status of a running request. + + The API reports only running requests. Once a request finishes, + successfully or not, this command fails with a not-running message. + Check the delivery bucket for results or error.json. + + Example: + + planet async-processing get 7d9a1c2e-0000-4000-8000-000000000000 + """ + async with async_processing_client(ctx) as cl: + try: + result = await cl.get_request(request_id) + except exceptions.MissingResource: + raise click.ClickException(NOT_RUNNING.format(request_id)) + echo_json(result, pretty) + + +@command(async_processing, name='wait') +@click.argument('request_id') +@click.option('--delay', + type=click.IntRange(min=0), + default=10, + show_default=True, + help='Time (in seconds) between polls.') +@click.option('--max-attempts', + type=click.IntRange(min=0), + default=360, + show_default=True, + help='Maximum number of polls. Set to zero for no limit.') +async def wait(ctx, request_id, delay, max_attempts, pretty): + """Wait until a request is no longer running. + + Polls the request status until the API stops reporting it, then exits + with status 0. Fails if --max-attempts is reached first. + + Completion does not mean success. Check the delivery bucket for results + or error.json. An unknown request ID also returns at once, since the + API reports it the same way as a finished request. + + Example: + + planet async-processing wait 7d9a1c2e-0000-4000-8000-000000000000 + """ + quiet = ctx.obj['QUIET'] + + def report(status): + if not quiet: + click.echo(f'{request_id}: {status}', err=True) + + async with async_processing_client(ctx) as cl: + await cl.wait(request_id, + delay=delay, + max_attempts=max_attempts, + callback=report) + if not quiet: + click.echo(NOT_RUNNING.format(request_id), err=True) diff --git a/planet/cli/cli.py b/planet/cli/cli.py index d678ddd00..8d15b9338 100644 --- a/planet/cli/cli.py +++ b/planet/cli/cli.py @@ -22,7 +22,7 @@ import planet from planet.cli import mosaics -from . import auth, cmds, collect, data, destinations, orders, quota, subscriptions, features +from . import async_processing, auth, cmds, collect, data, destinations, orders, quota, subscriptions, features LOGGER = logging.getLogger(__name__) @@ -123,6 +123,7 @@ def _configure_logging(verbosity): main.add_command(cmd=planet_auth_utils.cmd_plauth_embedded, name="plauth") # type: ignore +main.add_command(async_processing.async_processing) # type: ignore main.add_command(auth.cmd_auth) # type: ignore main.add_command(data.data) # type: ignore main.add_command(orders.orders) # type: ignore diff --git a/planet/clients/__init__.py b/planet/clients/__init__.py index d831dd1fa..34a791fab 100644 --- a/planet/clients/__init__.py +++ b/planet/clients/__init__.py @@ -12,6 +12,7 @@ # WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. # See the License for the specific language governing permissions and # limitations under the License. +from .async_processing import AsyncProcessingClient from .data import DataClient from .destinations import DestinationsClient from .features import FeaturesClient @@ -21,6 +22,7 @@ from .subscriptions import SubscriptionsClient __all__ = [ + 'AsyncProcessingClient', 'DataClient', 'DestinationsClient', 'FeaturesClient', @@ -32,6 +34,7 @@ # Organize client classes by their module name to allow lookup. _client_directory = { + 'async_processing': AsyncProcessingClient, 'data': DataClient, 'destinations': DestinationsClient, 'features': FeaturesClient, diff --git a/planet/clients/async_processing.py b/planet/clients/async_processing.py new file mode 100644 index 000000000..addf8f8a7 --- /dev/null +++ b/planet/clients/async_processing.py @@ -0,0 +1,183 @@ +# Copyright 2026 Planet Labs PBC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); you may not +# use this file except in compliance with the License. You may obtain a copy of +# the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations under +# the License. +"""Planet Async Processing API Python client.""" + +import asyncio +import logging +import time +from typing import Callable, Dict, Optional + +from planet.clients.base import _BaseClient +from planet.exceptions import APIError, ClientError, MissingResource +from planet.http import Session + +# Async Processing is served per deployment. Data must come from the same +# deployment the request is sent to. +DEPLOYMENT_URLS: Dict[str, str] = { + 'aws-eu-central-1': 'https://services.sentinel-hub.com/async/v1', + 'aws-us-west-2': 'https://services-uswest2.sentinel-hub.com/async/v1', +} +DEFAULT_DEPLOYMENT = 'aws-eu-central-1' +BASE_URL = DEPLOYMENT_URLS[DEFAULT_DEPLOYMENT] + +# The only status the API reports. Finished requests return 404. +RUNNING = 'RUNNING' + +LOGGER = logging.getLogger(__name__) + + +class AsyncProcessingClient(_BaseClient): + """Asynchronous Async Processing API client. + + The Async Processing API runs large Process API requests in the + background and writes the results to your S3 or GCS bucket. + + Authenticate with an OAuth2 user or M2M session (`planet auth login`). + Planet API keys are not accepted by this API. + + For more information, see + https://docs.planet.com/develop/apis/async-processing/ + + Example: + ```python + >>> import asyncio + >>> from planet import Session + >>> + >>> async def main(): + ... async with Session() as sess: + ... cl = sess.client('async_processing') + ... # use client here + ... + >>> asyncio.run(main()) + ``` + """ + + def __init__(self, + session: Session, + base_url: Optional[str] = None) -> None: + """ + Parameters: + session: Open session connected to server. + base_url: The base URL to use. Defaults to the AWS EU + (Frankfurt) deployment. See `DEPLOYMENT_URLS` for others. + """ + super().__init__(session, base_url or BASE_URL) + self._process_url = f'{self._base_url}/process' + + async def create_request(self, request: dict) -> dict: + """Submit an async processing request. + + Errors in the request are raised here. Errors during processing are + written to `error.json` in the delivery bucket. + + Parameters: + request: The request body. See + [planet.async_processing_request.build_request][]. + + Returns: + dict: `id` and `status` of the submitted request. + + Raises: + planet.exceptions.APIError: On API error, including + `TooManyRequests` when the concurrent request limit is + reached. + """ + try: + resp = await self._session.request(method='POST', + url=self._process_url, + json=request) + except APIError: + raise + except ClientError: # pragma: no cover + raise + return resp.json() + + async def get_request(self, request_id: str) -> dict: + """Get the status of a running request. + + The API reports only running requests. Once a request finishes, + successfully or not, this raises `MissingResource`. Check the + delivery bucket for results or `error.json`. + + Parameters: + request_id: ID of the request. + + Returns: + dict: `id` and `status` of the request. + + Raises: + planet.exceptions.MissingResource: If the request has finished + or does not exist. + planet.exceptions.APIError: On other API errors. + planet.exceptions.ClientError: If request_id is empty. + """ + if not request_id: + raise ClientError('Must provide a request id.') + + url = f'{self._process_url}/{request_id}' + try: + resp = await self._session.request(method='GET', url=url) + except APIError: + raise + except ClientError: # pragma: no cover + raise + return resp.json() + + async def wait(self, + request_id: str, + delay: int = 10, + max_attempts: int = 360, + callback: Optional[Callable[[str], None]] = None) -> None: + """Wait until a request is no longer running. + + Polls the request status every `delay` seconds until the API stops + reporting it. The API cannot tell a finished request from an unknown + ID, so an ID that never existed also returns at once. Completion + does not mean success: check the delivery bucket for results or + `error.json`. + + Parameters: + request_id: ID of the request. + delay: Seconds between polls. + max_attempts: Maximum number of polls. Set to zero for no limit. + callback: Called with the status after each poll that finds the + request running. + + Raises: + planet.exceptions.APIError: On API error. + planet.exceptions.ClientError: If request_id is empty or + max_attempts is reached while the request is running. + """ + num_attempts = 0 + while not max_attempts or num_attempts < max_attempts: + t = time.time() + try: + status = (await self.get_request(request_id))['status'] + except MissingResource: + LOGGER.debug(f'{request_id} is no longer running') + return + + LOGGER.debug(status) + if callback: + callback(status) + + num_attempts += 1 + await asyncio.sleep(max(delay - (time.time() - t), 0)) + + raise ClientError( + f'Maximum number of attempts ({max_attempts}) reached. ' + f'Request {request_id} is still running.') + + +__all__ = ['AsyncProcessingClient'] diff --git a/planet/http.py b/planet/http.py index f971bb8d2..bc5fae9c3 100644 --- a/planet/http.py +++ b/planet/http.py @@ -462,7 +462,10 @@ async def stream( await response.aclose() def client(self, - name: Literal['data', 'orders', 'subscriptions'], + name: Literal['async_processing', + 'data', + 'orders', + 'subscriptions'], base_url: Optional[str] = None) -> object: """Get a client by its module name. diff --git a/planet/sync/async_processing.py b/planet/sync/async_processing.py new file mode 100644 index 000000000..2bc16168c --- /dev/null +++ b/planet/sync/async_processing.py @@ -0,0 +1,78 @@ +# Copyright 2026 Planet Labs PBC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); you may not +# use this file except in compliance with the License. You may obtain a copy of +# the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations under +# the License. +"""Synchronous Planet Async Processing API client.""" + +from typing import Any, Callable, Dict, Optional + +from planet.clients.async_processing import AsyncProcessingClient +from planet.http import Session + + +class AsyncProcessingAPI: + """Async Processing API client. + + Example: + ```python + >>> from planet import Planet + >>> + >>> pl = Planet() + >>> req = pl.async_processing.create_request(request) + >>> pl.async_processing.wait(req['id']) + ``` + """ + + _client: AsyncProcessingClient + + def __init__(self, + session: Session, + base_url: Optional[str] = None) -> None: + """ + Parameters: + session: Open session connected to server. + base_url: The base URL to use. Defaults to the AWS EU + (Frankfurt) deployment. + """ + self._client = AsyncProcessingClient(session, base_url) + + def create_request(self, request: dict) -> Dict[str, Any]: + """Submit an async processing request. + + See [planet.clients.async_processing.AsyncProcessingClient.create_request][] + for details. + """ + return self._client._call_sync(self._client.create_request(request)) + + def get_request(self, request_id: str) -> Dict[str, Any]: + """Get the status of a running request. + + See [planet.clients.async_processing.AsyncProcessingClient.get_request][] + for details. + """ + return self._client._call_sync(self._client.get_request(request_id)) + + def wait(self, + request_id: str, + delay: int = 10, + max_attempts: int = 360, + callback: Optional[Callable[[str], None]] = None) -> None: + """Wait until a request is no longer running. + + See [planet.clients.async_processing.AsyncProcessingClient.wait][] + for details. + """ + return self._client._call_sync( + self._client.wait(request_id, + delay=delay, + max_attempts=max_attempts, + callback=callback)) diff --git a/planet/sync/client.py b/planet/sync/client.py index 499e6fe6b..77060970f 100644 --- a/planet/sync/client.py +++ b/planet/sync/client.py @@ -1,5 +1,6 @@ from typing import Optional +from .async_processing import AsyncProcessingAPI from .features import FeaturesAPI from .data import DataAPI from .destinations import DestinationsAPI @@ -20,6 +21,7 @@ class Planet: Members: + - `async_processing`: Async Processing API. - `data`: for interacting with the Planet Data API. - `destinations`: Destinations API. - `orders`: Orders API. @@ -69,3 +71,5 @@ def __init__(self, self.features = FeaturesAPI(self._session, f"{planet_base}/features/v1/ogc/my/") self.quota = QuotaAPI(self._session, f"{planet_base}/account/v1") + # Served from services.sentinel-hub.com, not the Planet base URL. + self.async_processing = AsyncProcessingAPI(self._session) diff --git a/tests/integration/test_async_processing_api.py b/tests/integration/test_async_processing_api.py new file mode 100644 index 000000000..6e4a9b5e1 --- /dev/null +++ b/tests/integration/test_async_processing_api.py @@ -0,0 +1,179 @@ +# Copyright 2026 Planet Labs PBC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); you may not +# use this file except in compliance with the License. You may obtain a copy of +# the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations under +# the License. +"""Tests of the Planet Async Processing API client.""" +from http import HTTPStatus +import json + +import httpx +import pytest +import respx + +from planet import AsyncProcessingClient, Planet, Session +from planet.auth import Auth +from planet.clients.async_processing import BASE_URL, DEPLOYMENT_URLS +from planet.exceptions import (BadQuery, + ClientError, + MissingResource, + TooManyRequests) +from planet.sync.async_processing import AsyncProcessingAPI + +pytestmark = pytest.mark.anyio # noqa + +# Simulated host/path for testing purposes. Not a real subdomain. +TEST_URL = "http://test.planet.com/async/v1" +PROCESS_URL = f"{TEST_URL}/process" + +REQUEST_ID = "7d9a1c2e-0000-4000-8000-000000000000" +RUNNING = {"id": REQUEST_ID, "status": "RUNNING"} +NOT_FOUND = {"error": {"status": 404, "reason": "Not Found"}} + +REQUEST = { + "input": { + "bounds": { + "bbox": [1, 2, 3, 4] + }, "data": [{ + "type": "sentinel-2-l2a" + }] + }, + "output": { + "width": 10, + "height": 10, + "delivery": { + "s3": { + "url": "s3://b/p", "iamRoleARN": "arn" + } + } + }, + "evalscript": "//VERSION=3", +} + +test_session = Session(auth=Auth.from_key(key="test")) +cl_async = AsyncProcessingClient(test_session, base_url=TEST_URL) +cl_sync = AsyncProcessingAPI(test_session, base_url=TEST_URL) + + +def status_route(): + return respx.get(f"{PROCESS_URL}/{REQUEST_ID}") + + +def test_default_base_url(): + assert BASE_URL == "https://services.sentinel-hub.com/async/v1" + assert AsyncProcessingClient(test_session)._process_url == \ + f"{BASE_URL}/process" + assert set(DEPLOYMENT_URLS) == {"aws-eu-central-1", "aws-us-west-2"} + + +def test_session_client_lookup(): + assert isinstance(test_session.client("async_processing"), + AsyncProcessingClient) + + +def test_planet_sync_client_uses_sentinel_hub(): + pl = Planet(session=test_session, base_url="http://test.planet.com") + assert pl.async_processing._client._base_url == BASE_URL + + +@respx.mock +async def test_create_request(): + respx.post(PROCESS_URL).return_value = httpx.Response(HTTPStatus.OK, + json=RUNNING) + assert await cl_async.create_request(REQUEST) == RUNNING + assert json.loads(respx.calls.last.request.content) == REQUEST + + +@respx.mock +def test_create_request_sync(): + respx.post(PROCESS_URL).return_value = httpx.Response(HTTPStatus.OK, + json=RUNNING) + assert cl_sync.create_request(REQUEST) == RUNNING + + +@respx.mock +async def test_create_request_bad_request(): + respx.post(PROCESS_URL).return_value = httpx.Response( + HTTPStatus.BAD_REQUEST, json={"error": { + "message": "bad bbox" + }}) + with pytest.raises(BadQuery): + await cl_async.create_request(REQUEST) + + +@respx.mock +async def test_create_request_concurrency_limit(): + # The session retries 429s; make it give up at once. + respx.post(PROCESS_URL).return_value = httpx.Response( + HTTPStatus.TOO_MANY_REQUESTS, json={}) + sess = Session(auth=Auth.from_key(key="test")) + sess.max_retries = 0 + cl = AsyncProcessingClient(sess, base_url=TEST_URL) + with pytest.raises(TooManyRequests): + await cl.create_request(REQUEST) + + +@respx.mock +async def test_get_request(): + status_route().return_value = httpx.Response(HTTPStatus.OK, json=RUNNING) + assert await cl_async.get_request(REQUEST_ID) == RUNNING + + +@respx.mock +def test_get_request_sync(): + status_route().return_value = httpx.Response(HTTPStatus.OK, json=RUNNING) + assert cl_sync.get_request(REQUEST_ID) == RUNNING + + +@respx.mock +async def test_get_request_finished(): + status_route().return_value = httpx.Response(HTTPStatus.NOT_FOUND, + json=NOT_FOUND) + with pytest.raises(MissingResource): + await cl_async.get_request(REQUEST_ID) + + +async def test_get_request_empty_id(): + with pytest.raises(ClientError): + await cl_async.get_request("") + + +@respx.mock +async def test_wait(): + route = status_route() + route.side_effect = [ + httpx.Response(HTTPStatus.OK, json=RUNNING), + httpx.Response(HTTPStatus.OK, json=RUNNING), + httpx.Response(HTTPStatus.NOT_FOUND, json=NOT_FOUND), + ] + seen = [] + assert await cl_async.wait(REQUEST_ID, delay=0, + callback=seen.append) is None + assert seen == ["RUNNING", "RUNNING"] + assert route.call_count == 3 + + +@respx.mock +def test_wait_sync(): + status_route().side_effect = [ + httpx.Response(HTTPStatus.OK, json=RUNNING), + httpx.Response(HTTPStatus.NOT_FOUND, json=NOT_FOUND), + ] + assert cl_sync.wait(REQUEST_ID, delay=0) is None + + +@respx.mock +async def test_wait_max_attempts(): + route = status_route() + route.return_value = httpx.Response(HTTPStatus.OK, json=RUNNING) + with pytest.raises(ClientError, match="still running"): + await cl_async.wait(REQUEST_ID, delay=0, max_attempts=2) + assert route.call_count == 2 diff --git a/tests/integration/test_async_processing_cli.py b/tests/integration/test_async_processing_cli.py new file mode 100644 index 000000000..b291fe0d6 --- /dev/null +++ b/tests/integration/test_async_processing_cli.py @@ -0,0 +1,329 @@ +# Copyright 2026 Planet Labs PBC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); you may not +# use this file except in compliance with the License. You may obtain a copy of +# the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations under +# the License. +"""Tests of the planet async-processing CLI.""" +import base64 +from http import HTTPStatus +import json + +import httpx +import pytest +import respx +from click.testing import CliRunner + +from planet.cli import cli + +from tests.integration.test_async_processing_api import ( + NOT_FOUND, + PROCESS_URL, + REQUEST, + REQUEST_ID, + RUNNING, + TEST_URL, +) + +REQUIRED = [ + "--collection", + "sentinel-2-l2a", + "--bbox", + "1,2,3,4", + "--width", + "10", + "--height", + "10", + "--delivery", + "s3://b/p", + "--iam-role-arn", + "arn", +] + + +def invoke(*args, input=None, base_url=TEST_URL): + runner = CliRunner() + group = ["async-processing"] + if base_url: + group += ["--base-url", base_url] + return runner.invoke(cli.main, args=group + list(args), input=input) + + +@pytest.fixture +def evalscript(tmp_path): + path = tmp_path / "script.js" + path.write_text("//VERSION=3") + return str(path) + + +def test_help_lists_commands(): + result = invoke("--help", base_url=None) + assert result.exit_code == 0 + for cmd in ("request", "create", "get", "wait"): + assert cmd in result.output + + +@pytest.mark.parametrize("cmd", ["request", "create", "get", "wait"]) +def test_subcommand_help_has_example(cmd): + result = invoke(cmd, "--help", base_url=None) + assert result.exit_code == 0 + assert "Example:" in result.output + + +def test_request_minimal(evalscript): + result = invoke("request", *REQUIRED, "--evalscript", evalscript) + assert result.exit_code == 0, result.output + assert json.loads(result.output) == { + "input": { + "bounds": { + "bbox": [1.0, 2.0, 3.0, 4.0] + }, + "data": [{ + "type": "sentinel-2-l2a" + }] + }, + "output": { + "width": 10, + "height": 10, + "responses": [{ + "identifier": "default", "format": { + "type": "image/tiff" + } + }], + "delivery": { + "s3": { + "url": "s3://b/p", "iamRoleARN": "arn" + } + } + }, + "evalscript": "//VERSION=3" + } + + +def test_request_full(evalscript): + geom = { + "type": "Feature", + "properties": {}, + "geometry": { + "type": "Point", "coordinates": [1, 2] + } + } + result = invoke("request", + "--collection", + "sentinel-2-l2a", + "--time-from", + "2024-01-01T00:00:00Z", + "--time-to", + "2024-02-01T00:00:00Z", + "--max-cloud-coverage", + "20", + "--mosaicking-order", + "leastCC", + "--geometry", + json.dumps(geom), + "--crs", + "http://www.opengis.net/def/crs/EPSG/0/4326", + "--resx", + "10", + "--resy", + "10", + "--response", + "default:image/tiff", + "--response", + "userdata:application/json", + "--delivery", + "s3://b//", + "--aws-access-key-id", + "a", + "--aws-secret-access-key", + "s", + "--aws-region", + "us-east-1", + "--evalscript", + evalscript) + assert result.exit_code == 0, result.output + req = json.loads(result.output) + assert req["input"]["bounds"] == { + "geometry": { + "type": "Point", "coordinates": [1, 2] + }, + "properties": { + "crs": "http://www.opengis.net/def/crs/EPSG/0/4326" + } + } + assert req["input"]["data"][0]["dataFilter"] == { + "timeRange": { + "from": "2024-01-01T00:00:00Z", "to": "2024-02-01T00:00:00Z" + }, + "maxCloudCoverage": 20.0, + "mosaickingOrder": "leastCC" + } + assert req["output"]["resx"] == 10.0 + assert [r["identifier"] + for r in req["output"]["responses"]] == ["default", "userdata"] + assert req["output"]["delivery"] == { + "s3": { + "url": "s3://b//", + "accessKey": "a", + "secretAccessKey": "s", + "region": "us-east-1" + } + } + + +def test_request_gcs_and_evalscript_url(tmp_path): + creds = tmp_path / "key.json" + creds.write_bytes(b'{"type": "service_account"}') + args = [a if a != "s3://b/p" else "gs://b/p" for a in REQUIRED] + result = invoke("request", + *args, + "--evalscript-url", + "gs://b/script.js", + "--gcs-credentials", + str(creds)) + assert result.exit_code == 0, result.output + req = json.loads(result.output) + expected = base64.b64encode(creds.read_bytes()).decode() + assert req["output"]["delivery"] == { + "gs": { + "url": "gs://b/p", "credentials": expected + } + } + assert req["evalscriptReference"] == { + "gs": { + "url": "gs://b/script.js", "credentials": expected + } + } + assert "evalscript" not in req + + +@pytest.mark.parametrize( + "extra, message", + [ + ([], "exactly one of --evalscript and --evalscript-url"), + (["--evalscript-url", "s3://b/e.js", "--evalscript", "EVALSCRIPT"], + "exactly one of --evalscript and --evalscript-url"), + (["--evalscript", "EVALSCRIPT", "--response", "default"], + "IDENTIFIER:FORMAT"), + (["--evalscript", "EVALSCRIPT", "--response", "default:image/gif"], + "IDENTIFIER:FORMAT"), + (["--evalscript", "EVALSCRIPT", "--resx", "10", "--resy", "10"], + "not both"), + ]) +def test_request_invalid(evalscript, extra, message): + extra = [evalscript if a == "EVALSCRIPT" else a for a in extra] + result = invoke("request", *REQUIRED, *extra) + assert result.exit_code != 0 + assert message in result.output + + +def test_request_gcs_requires_credentials(evalscript): + args = [a if a != "s3://b/p" else "gs://b/p" for a in REQUIRED] + result = invoke("request", *args, "--evalscript", evalscript) + assert result.exit_code != 0 + assert "--gcs-credentials is required" in result.output + + +def test_request_s3_requires_credentials(evalscript): + args = REQUIRED[:-2] + result = invoke("request", *args, "--evalscript", evalscript) + assert result.exit_code != 0 + assert "iam_role_arn" in result.output + + +@respx.mock +def test_create(): + respx.post(PROCESS_URL).return_value = httpx.Response(HTTPStatus.OK, + json=RUNNING) + result = invoke("create", json.dumps(REQUEST)) + assert result.exit_code == 0, result.output + assert json.loads(result.output) == RUNNING + assert json.loads(respx.calls.last.request.content) == REQUEST + + +@respx.mock +def test_create_stdin(): + respx.post(PROCESS_URL).return_value = httpx.Response(HTTPStatus.OK, + json=RUNNING) + result = invoke("create", "-", input=json.dumps(REQUEST)) + assert result.exit_code == 0, result.output + assert json.loads(result.output) == RUNNING + + +@respx.mock +def test_create_bad_request(): + respx.post(PROCESS_URL).return_value = httpx.Response( + HTTPStatus.BAD_REQUEST, json={"error": { + "message": "bad bbox" + }}) + result = invoke("create", json.dumps(REQUEST)) + assert result.exit_code == 1 + assert "bad bbox" in result.output + + +@respx.mock +def test_get(): + respx.get(f"{PROCESS_URL}/{REQUEST_ID}").return_value = httpx.Response( + HTTPStatus.OK, json=RUNNING) + result = invoke("get", REQUEST_ID) + assert result.exit_code == 0, result.output + assert json.loads(result.output) == RUNNING + + +@respx.mock +def test_get_finished(): + respx.get(f"{PROCESS_URL}/{REQUEST_ID}").return_value = httpx.Response( + HTTPStatus.NOT_FOUND, json=NOT_FOUND) + result = invoke("get", REQUEST_ID) + assert result.exit_code == 1 + assert "is not running" in result.output + assert "error.json" in result.output + + +@respx.mock +def test_wait(): + route = respx.get(f"{PROCESS_URL}/{REQUEST_ID}") + route.side_effect = [ + httpx.Response(HTTPStatus.OK, json=RUNNING), + httpx.Response(HTTPStatus.NOT_FOUND, json=NOT_FOUND), + ] + result = invoke("wait", REQUEST_ID, "--delay", "0") + assert result.exit_code == 0, result.output + assert f"{REQUEST_ID}: RUNNING" in result.output + assert "is not running" in result.output + assert route.call_count == 2 + + +@respx.mock +def test_wait_max_attempts(): + respx.get(f"{PROCESS_URL}/{REQUEST_ID}").return_value = httpx.Response( + HTTPStatus.OK, json=RUNNING) + result = invoke("wait", REQUEST_ID, "--delay", "0", "--max-attempts", "1") + assert result.exit_code == 1 + assert "still running" in result.output + + +@pytest.mark.parametrize( + "deployment, host", + [ + ("aws-eu-central-1", "services.sentinel-hub.com"), + ("aws-us-west-2", "services-uswest2.sentinel-hub.com"), + ]) +@respx.mock +def test_deployment_selects_host(deployment, host): + route = respx.get(f"https://{host}/async/v1/process/{REQUEST_ID}").mock( + return_value=httpx.Response(HTTPStatus.OK, json=RUNNING)) + result = invoke("--deployment", + deployment, + "get", + REQUEST_ID, + base_url=None) + assert result.exit_code == 0, result.output + assert route.called diff --git a/tests/unit/test_async_processing_request.py b/tests/unit/test_async_processing_request.py new file mode 100644 index 000000000..ea7d27b74 --- /dev/null +++ b/tests/unit/test_async_processing_request.py @@ -0,0 +1,215 @@ +# Copyright 2026 Planet Labs PBC. +# +# Licensed under the Apache License, Version 2.0 (the "License"); you may not +# use this file except in compliance with the License. You may obtain a copy of +# the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +# WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +# License for the specific language governing permissions and limitations under +# the License. +import pytest + +from planet import async_processing_request as apr +from planet.exceptions import ClientError + +BUCKET = {'s3': {'url': 's3://b/p', 'iamRoleARN': 'arn:role'}} +GEOM = {'type': 'Point', 'coordinates': [1, 2]} + + +def test_build_request_evalscript(): + req = apr.build_request({'i': 1}, {'o': 1}, evalscript='//VERSION=3') + assert req == { + 'input': { + 'i': 1 + }, 'output': { + 'o': 1 + }, 'evalscript': '//VERSION=3' + } + + +def test_build_request_evalscript_reference(): + req = apr.build_request({}, {}, evalscript_reference=BUCKET) + assert req['evalscriptReference'] == BUCKET + assert 'evalscript' not in req + + +@pytest.mark.parametrize('kwargs', + [{}, { + 'evalscript': 'x', 'evalscript_reference': BUCKET + }]) +def test_build_request_evalscript_exclusive(kwargs): + with pytest.raises(ClientError): + apr.build_request({}, {}, **kwargs) + + +def test_process_input_bbox_geometry_crs(): + src = apr.data_source('sentinel-2-l2a') + inp = apr.process_input([src], + bbox=[1, 2, 3, 4], + geometry=GEOM, + crs='EPSG') + assert inp == { + 'bounds': { + 'bbox': [1, 2, 3, 4], + 'geometry': GEOM, + 'properties': { + 'crs': 'EPSG' + } + }, + 'data': [src] + } + + +def test_process_input_geometry_only(): + inp = apr.process_input([{'type': 't'}], geometry=GEOM) + assert inp['bounds'] == {'geometry': GEOM} + + +@pytest.mark.parametrize('kwargs', + [ + { + 'data': [{ + 'type': 't' + }] + }, + { + 'data': [{ + 'type': 't' + }], 'bbox': [1, 2, 3] + }, + { + 'data': [], 'bbox': [1, 2, 3, 4] + }, + ]) +def test_process_input_invalid(kwargs): + with pytest.raises(ClientError): + apr.process_input(**kwargs) + + +def test_data_source_minimal(): + assert apr.data_source('sentinel-1-grd') == {'type': 'sentinel-1-grd'} + + +def test_data_source_full(): + src = apr.data_source('sentinel-2-l2a', + time_from='2024-01-01T00:00:00Z', + time_to='2024-02-01T00:00:00Z', + max_cloud_coverage=20, + mosaicking_order='leastCC', + id='s2', + processing={'upsampling': 'BICUBIC'}) + assert src == { + 'type': 'sentinel-2-l2a', + 'id': 's2', + 'dataFilter': { + 'timeRange': { + 'from': '2024-01-01T00:00:00Z', 'to': '2024-02-01T00:00:00Z' + }, + 'maxCloudCoverage': 20, + 'mosaickingOrder': 'leastCC' + }, + 'processing': { + 'upsampling': 'BICUBIC' + } + } + + +def test_data_source_open_time_range(): + src = apr.data_source('t', time_to='2024-02-01T00:00:00Z') + assert src['dataFilter'] == {'timeRange': {'to': '2024-02-01T00:00:00Z'}} + + +def test_process_output_size(): + out = apr.process_output(BUCKET, + width=100, + height=200, + responses=[apr.response()]) + assert out == { + 'width': 100, + 'height': 200, + 'responses': [{ + 'identifier': 'default', 'format': { + 'type': 'image/tiff' + } + }], + 'delivery': BUCKET + } + + +def test_process_output_resolution(): + out = apr.process_output(BUCKET, resx=10, resy=10) + assert out == {'resx': 10, 'resy': 10, 'delivery': BUCKET} + + +@pytest.mark.parametrize('kwargs', + [ + { + 'width': 1, 'height': 1, 'resx': 1, 'resy': 1 + }, + { + 'width': 1 + }, + { + 'resy': 1 + }, + ]) +def test_process_output_invalid(kwargs): + with pytest.raises(ClientError): + apr.process_output(BUCKET, **kwargs) + + +def test_response(): + assert apr.response('userdata', 'application/json') == { + 'identifier': 'userdata', 'format': { + 'type': 'application/json' + } + } + + +def test_s3_bucket_role(): + assert apr.s3_bucket('s3://b/p', iam_role_arn='arn', region='us-east-1') \ + == {'s3': {'url': 's3://b/p', 'iamRoleARN': 'arn', + 'region': 'us-east-1'}} + + +def test_s3_bucket_keys(): + assert apr.s3_bucket('s3://b', access_key='a', secret_access_key='s') \ + == {'s3': {'url': 's3://b', 'accessKey': 'a', + 'secretAccessKey': 's'}} + + +@pytest.mark.parametrize('kwargs', + [ + { + 'url': 'gs://b', 'iam_role_arn': 'arn' + }, + { + 'url': 's3://b' + }, + { + 'url': 's3://b', 'access_key': 'a' + }, + { + 'url': 's3://b', 'secret_access_key': 's' + }, + ]) +def test_s3_bucket_invalid(kwargs): + with pytest.raises(ClientError): + apr.s3_bucket(**kwargs) + + +def test_gs_bucket(): + assert apr.gs_bucket('gs://b/p', 'Y3JlZHM=') == { + 'gs': { + 'url': 'gs://b/p', 'credentials': 'Y3JlZHM=' + } + } + + +def test_gs_bucket_invalid(): + with pytest.raises(ClientError): + apr.gs_bucket('s3://b', 'c')