Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
28 commits
Select commit Hold shift + click to select a range
db18e28
Update auth token check.
rhysrevans3 Jul 2, 2026
f2b0708
Use walrus.
rhysrevans3 Jul 2, 2026
1f48071
Add debug log to create item.
rhysrevans3 Jul 13, 2026
a424883
Add timeout setting for client.
rhysrevans3 Jul 13, 2026
a3fe1ba
Dump to item to json.
rhysrevans3 Jul 13, 2026
57a932f
Update CORDEX extension version.
rhysrevans3 Jul 14, 2026
5a90527
Update schema versions.
rhysrevans3 Jul 17, 2026
3b472fe
Merge pull request #52 from ESGF/integration
lukaszlacinski Jul 21, 2026
3c63b00
Increase read timeout for EGI checkin.
rhysrevans3 Jul 22, 2026
867c931
Add default timeout.
rhysrevans3 Jul 22, 2026
02843f6
Hard code timeout.
rhysrevans3 Jul 22, 2026
3707b2f
Add timeout default.
rhysrevans3 Jul 22, 2026
b8d2f86
Filp patch ordering.
rhysrevans3 Jul 24, 2026
03e822c
Update operation to partial item.
rhysrevans3 Jul 24, 2026
87e2053
Switch to pop.
rhysrevans3 Jul 24, 2026
361977a
Update entitlements get.
rhysrevans3 Jul 24, 2026
2b8b35e
Drop CMIP7 schema version to v1.2.8.
rhysrevans3 Jul 27, 2026
3148aed
Move validation to core utils.
rhysrevans3 Jul 29, 2026
073e0a5
Remove unused import.
rhysrevans3 Jul 29, 2026
99dd567
temporary remove core utils dependency.
rhysrevans3 Jul 29, 2026
b78d763
Add extra logging.
rhysrevans3 Aug 4, 2026
2c2b9d5
Merge branch 'main' of github.com:ESGF/stac-transaction-api into dev
rhysrevans3 Aug 4, 2026
ec81e2b
Install core-utils from pull.
rhysrevans3 Aug 4, 2026
7c64f20
Remove core utils requirement.
rhysrevans3 Aug 4, 2026
0f3ced8
Add aud check.
rhysrevans3 Aug 4, 2026
c3e43d4
Add patch role debug log.
rhysrevans3 Aug 4, 2026
c6a69dd
Adding exta error logging.
rhysrevans3 Aug 5, 2026
66bcf09
Update exception handling.
rhysrevans3 Aug 10, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1,027 changes: 21 additions & 1,006 deletions poetry.lock

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,7 @@ authors = [
]
readme = "README.md"
dependencies = [
"esgf-core-utils>=1.1.0",
# "esgf-core-utils>=1.2.0",
"fastapi>=0.114.0",
"jsonschema>=4.24.0",
"packaging>=26.2",
Expand Down
42 changes: 28 additions & 14 deletions src/authorizer/egi_authorizer.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,7 @@
from esgf_core_utils.models.auth import Authorizer
from esgf_core_utils.models.exceptions import InvalidTokenAudienceException
from esgf_core_utils.models.kafka.events import RequesterData
from fastapi import Request
from fastapi import HTTPException, Request
from starlette.middleware.base import BaseHTTPMiddleware

from settings import settings
Expand All @@ -17,8 +17,8 @@
Authorizer type: FastAPI Middleware
Event payload: Token
Token source: Authorization
Token validation: ^Bearer\s[^\s]+$ # noqa: W605
^Bearer\s[0-9A-Za-z]+$ for access tokens issued by Globus Auth (?) # noqa: W605
Token validation: ^Bearer\\s[^\\s]+$ # noqa: W605
^Bearer\\s[0-9A-Za-z]+$ for access tokens issued by Globus Auth (?) # noqa: W605
Authorization caching: 300 seconds
"""

Expand All @@ -40,25 +40,39 @@ async def dispatch(self, request: Request, call_next):
password=settings.client.client_secret,
)

async with httpx.AsyncClient(timeout=5.0, verify=False) as client:
async with httpx.AsyncClient(
timeout=httpx.Timeout(
10.0,
connect=settings.client.timeout.connect,
read=settings.client.timeout.read,
),
verify=False,
) as client:
logger.debug(
"Post request to %s",
settings.client.introspection_endpoint,
)
response = await client.post(
settings.client.introspection_endpoint,
headers={"Content-type": "application/x-www-form-urlencoded"},
data=f"token={request.headers.get('authorization')[7:]}",
auth=auth,
timeout=5,
)
response.raise_for_status()
if token := request.headers.get("authorization", "").removeprefix("Bearer "):
try:
response = await client.post(
settings.client.introspection_endpoint,
headers={"Content-type": "application/x-www-form-urlencoded"},
data=f"token={token}",
auth=auth,
)
response.raise_for_status()

except httpx.HTTPError as exc:
raise HTTPException(status_code=401) from exc

else:
raise HTTPException(status_code=401, detail="Missing or invalid bearer token")

token_info = response.json()

logger.debug("Token info: %s", token_info)

if request.headers["host"] not in [urlparse(aud).hostname for aud in token_info["aud"]]:
if "aud" in token_info and request.headers["host"] not in [urlparse(aud).hostname for aud in token_info["aud"]]:
raise InvalidTokenAudienceException(
token_audience=request.headers["host"],
expected_audience=", ".join(token_info["aud"]),
Expand All @@ -73,7 +87,7 @@ async def dispatch(self, request: Request, call_next):
),
)

authorizer.add(token_info["entitlements"])
authorizer.add(token_info.get("entitlements", []))

request.state.authorizer = authorizer

Expand Down
81 changes: 42 additions & 39 deletions src/client.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@
from datetime import datetime
from typing import Optional, Union


from esgf_core_utils.models.auth import Authorizer
from esgf_core_utils.models.exceptions import (
AuthorizationException,
Expand All @@ -26,27 +25,25 @@
Publisher,
)
from esgf_core_utils.models.kafka.producer import KafkaProducer
from fastapi import Request, Response, status
from stac_fastapi.extensions.transaction import BaseTransactionsClient
from stac_fastapi.extensions.transaction.request import PartialItem, PatchOperation
from stac_fastapi.types.stac import Collection
from stac_pydantic.item import Item
from pydantic import TypeAdapter

from utils import (
from esgf_core_utils.models.validation import (
evaluate_patch,
operation_to_partial_item,
patch_adapter,
validate_extensions,
validate_patch,
validate_post,
)
from fastapi import Request, Response, status
from stac_fastapi.extensions.transaction import BaseTransactionsClient
from stac_fastapi.extensions.transaction.request import PartialItem, PatchOperation
from stac_fastapi.types.stac import Collection
from stac_pydantic.item import Item

# Setup logger
# logger = logging.getLogger(__name__)
logger = logging.getLogger("uvicorn.error")
logger.setLevel(logging.DEBUG)

patch_adapter = TypeAdapter(PartialItem | list[PatchOperation])


class TransactionClient(BaseTransactionsClient):

Expand Down Expand Up @@ -94,6 +91,7 @@ async def create_item(
request: Request,
) -> Optional[Union[Item, Response, None]]:

logger.debug("CREATE REQUEST: %s", item.model_dump_json())
headers = request.headers

event_id = uuid.uuid4().hex
Expand All @@ -112,9 +110,12 @@ async def create_item(
except MissingPermissionException as exc:
raise AuthorizationException(instance=f"{request_id}:{event_id}") from exc

item_extensions = item.stac_extensions if item.stac_extensions else []
try:
item_extensions = validate_extensions(collection_id=collection_id, item_extensions=item_extensions)
item_extensions = validate_extensions(
collection_id=collection_id,
item_extensions=item.stac_extensions or [],
)

validate_post(
item_id=item.id,
item=item,
Expand All @@ -128,13 +129,10 @@ async def create_item(
UnexpectedExtensionException,
ExtensionBelowMinimumException,
) as exc:
rfc_exc = RFC9457Exception()
rfc_exc.status_code = 400
rfc_exc.type = exc.type
rfc_exc.title = exc.title
rfc_exc.detail = exc.detail
rfc_exc.instance = f"{request_id}:{event_id}"
raise rfc_exc from exc
logger.error("Error producing message for %s: %s", item.id, exc)
logger.error("Item model dump: %s", item.model_dump_json())
exc.instance = f"{request_id}:{event_id}"
raise

user_agent = headers.get("user-agent", "/").split("/")

Expand All @@ -146,7 +144,9 @@ async def create_item(

data = Data(type="STAC", payload=payload)

publisher = Publisher(package=user_agent[0], version=user_agent[1] if len(user_agent) > 1 else "")
publisher = Publisher(
package=user_agent[0], version=user_agent[1] if len(user_agent) > 1 else ""
)

metadata = Metadata(
auth=auth,
Expand All @@ -165,7 +165,8 @@ async def create_item(
)

except Exception as exc:
logger.error("Error producing message: %s", exc)
logger.error("Error producing message for %s: %s", item.id, exc)
logger.error("Item model dump: %s", item.model_dump_json())
raise UnknownException(instance=f"{request_id}:{event_id}") from exc

return Response(
Expand All @@ -191,17 +192,25 @@ async def patch_item(
) -> Item | Response | None:
logger.info("PATCH REQUEST: %s", patch)

item = operation_to_partial_item(collection_id=collection_id, operations=patch) if isinstance(patch, list) else patch
item = (
operation_to_partial_item(collection_id=collection_id, operations=patch)
if isinstance(patch, list)
else patch
)

headers = request.headers.get("headers", {})

event_id = uuid.uuid4().hex
request_id = headers.get("X-Request-ID", uuid.uuid4().hex)

role = evaluate_patch(patch)

logger.debug("PATCH ROLE: %s", role)

auth = self.authorize(
collection_id=collection_id,
item=item,
role="UPDATE",
role=role,
request=request,
request_id=request_id,
event_id=event_id,
Expand All @@ -210,7 +219,9 @@ async def patch_item(
item_extensions = item.stac_extensions if item.stac_extensions else []
try:

item_extensions = validate_extensions(collection_id=collection_id, item_extensions=item_extensions)
item_extensions = validate_extensions(
collection_id=collection_id, item_extensions=item_extensions
)

validate_patch(
item_id=item_id,
Expand All @@ -224,31 +235,23 @@ async def patch_item(
UnexpectedExtensionException,
ExtensionBelowMinimumException,
) as exc:
rfc_exc = RFC9457Exception()
rfc_exc.status_code = 400
rfc_exc.type = exc.type
rfc_exc.title = exc.title
rfc_exc.detail = exc.detail
rfc_exc.instance = f"{request_id}:{event_id}"
raise rfc_exc from exc
exc.instance = f"{request_id}:{event_id}"
raise

user_agent = headers.get("user-agent", "/").split("/")

if isinstance(patch, list):
patch_body = [op.model_dump() for op in patch]
else:
patch_body = item.model_dump()

payload = PatchPayload(
method="PATCH",
collection_id=collection_id,
item_id=item_id,
patch=patch_body,
patch=patch_adapter.dump_python(patch),
)

data = Data(type="STAC", payload=payload)

publisher = Publisher(package=user_agent[0], version=user_agent[1] if len(user_agent) > 1 else "")
publisher = Publisher(
package=user_agent[0], version=user_agent[1] if len(user_agent) > 1 else ""
)
metadata = Metadata(
auth=auth,
event_id=event_id,
Expand Down
71 changes: 1 addition & 70 deletions src/settings/__init__.py
Original file line number Diff line number Diff line change
@@ -1,82 +1,13 @@
import os
from typing import Literal
import re

from pydantic_settings import BaseSettings, SettingsConfigDict

if os.environ.get("TRANSACTION_AUTHORIZER") == "egi":
from settings.ceda import CEDAClientSettings as ClientSettings
else:
from settings.globus import GlobusClientSettings as ClientSettings

DEFAULT_EXTENSIONS = {
"CMIP6": {
"CMIP6": {
"regex": [r"https:\/\/esgf\.github\.io\/stac-transaction-api\/cmip6\/v[0-9]\.[0-9]\.[0-9]/schema\.json"],
"default": "https://esgf.github.io/stac-transaction-api/cmip6/v2.0.0/schema.json",
},
"alternate_assets": {
"regex": [r"https:\/\/stac-extensions\.github\.io\/alternate-assets\/v[0-9]\.[0-9]\.[0-9]\/schema\.json"],
"default": "https://stac-extensions.github.io/alternate-assets/v1.2.0/schema.json",
},
"file": {
"regex": [r"https:\/\/stac-extensions\.github\.io\/file\/v[0-9]\.[0-9]\.[0-9]/schema\.json"],
"default": "https://stac-extensions.github.io/file/v2.1.0/schema.json",
},
},
"CMIP6Plus": {
"CMIP6Plus": {
"regex": [r"https:\/\/esgf\.github\.io\/stac-transaction-api\/cmip6plus\/v[0-9]\.[0-9]\.[0-9]/schema\.json"],
"default": "https://esgf.github.io/stac-transaction-api/cmip6plus/v2.0.0/schema.json",
},
"alternate_assets": {
"regex": [r"https:\/\/stac-extensions\.github\.io\/alternate-assets\/v[0-9]\.[0-9]\.[0-9]\/schema\.json"],
"default": "https://stac-extensions.github.io/alternate-assets/v1.2.0/schema.json",
},
"file": {
"regex": [r"https:\/\/stac-extensions\.github\.io\/file\/v[0-9]\.[0-9]\.[0-9]/schema\.json"],
"default": "https://stac-extensions.github.io/file/v2.1.0/schema.json",
},
},
"CMIP7": {
"CMIP7": {
"regex": [r"https:\/\/esgf\.github\.io\/stac-transaction-api\/cmip7\/v[0-9]\.[0-9]\.[0-9]\/schema\.json"],
"default": "https://esgf.github.io/stac-transaction-api/cmip7/v1.2.1/schema.json",
},
"alternate_assets": {
"regex": [r"https:\/\/stac-extensions\.github\.io\/alternate-assets\/v[0-9]\.[0-9]\.[0-9]\/schema\.json"],
"default": "https://stac-extensions.github.io/alternate-assets/v1.2.0/schema.json",
},
"file": {
"regex": [r"https:\/\/stac-extensions\.github\.io\/file\/v[0-9]\.[0-9]\.[0-9]/schema\.json"],
"default": "https://stac-extensions.github.io/file/v2.1.0/schema.json",
},
},
"CORDEX-CMIP6": {
"CORDEX-CMIP6": {
"regex": [r"https:\/\/esgf\.github\.io\/stac-transaction-api\/cordex-cmip6\/v[0-9]\.[0-9]\.[0-9]/schema\.json"],
"default": "https://esgf.github.io/stac-transaction-api/cordex-cmip6/v1.2.0/schema.json",
},
"alternate_assets": {
"regex": [r"https:\/\/stac-extensions\.github\.io\/alternate-assets\/v[0-9]\.[0-9]\.[0-9]\/schema\.json"],
"default": "https://stac-extensions.github.io/alternate-assets/v1.2.0/schema.json",
},
"file": {
"regex": [r"https:\/\/stac-extensions\.github\.io\/file\/v[0-9]\.[0-9]\.[0-9]/schema\.json"],
"default": "https://stac-extensions.github.io/file/v2.1.0/schema.json",
},
},
}

VERSION_REGEX = re.compile(
r"/v("
r"(?P<major>0|[1-9]\d*)\."
r"(?P<minor>0|[1-9]\d*)\."
r"(?P<patch>0|[1-9]\d*)"
r"(?:-[0-9A-Za-z-]+(?:\.[0-9A-Za-z-]+)*)?"
r"(?:\+[0-9A-Za-z-]+(?:\.[0-9A-Za-z-]+)*)?"
r")/"
)


class Settings(BaseSettings):
"""
Expand Down
11 changes: 11 additions & 0 deletions src/settings/ceda.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,15 @@
from pydantic import BaseModel


class TimeoutSettings(BaseModel):
"""
Timeout settings
"""

connect: float = 5.0
read: float = 15.0


class CEDAClientSettings(BaseModel):
"""
CEDA settings
Expand All @@ -10,6 +19,8 @@ class CEDAClientSettings(BaseModel):
client_secret: str
token_url: str = "https://aai.egi.eu/auth/realms/egi/protocol/openid-connect/token"
introspection_endpoint: str = "https://aai.egi.eu/auth/realms/egi/protocol/openid-connect/token/introspect"
timeout: TimeoutSettings = TimeoutSettings()

regex: str = (
r"urn\:mace\:egi\.eu\:group\:esgf.vo.egi.eu\:(?P<type>[^:]*)\:(?P<id>[^:]*)"
r"(\:institution\:(?P<institution>[^:]*))?\:role=(?P<role>[^:]*)#aai\.egi\.eu"
Expand Down
17 changes: 0 additions & 17 deletions src/test_api.py

This file was deleted.

Loading