Skip to content

Logistics API

jukto.interfaces.logistics.Parcel dataclass

Source code in src/jukto/interfaces/logistics.py
@dataclass
class Parcel:
    invoice_id: str
    recipient_name: str
    recipient_phone: str
    recipient_address: str
    cod_amount: Decimal | float | int
    note: Optional[str] = None
    # Extended fields (useful for Pathao, RedX, etc.)
    recipient_city: Optional[int] = None
    recipient_zone: Optional[int] = None
    recipient_area: Optional[int] = None
    store_id: Optional[int] = None
    item_weight: float = 0.5
    item_quantity: int = 1
    item_type: int = 2  # 1: Document, 2: Parcel
    delivery_type: int = 48  # 48: Normal Delivery (Pathao default)
    # Independent goods value; never inferred from COD. RedX-only routing fields.
    declared_value: Optional[Decimal | float | int] = None
    redx_delivery_area: Optional[str] = None
    redx_delivery_area_id: Optional[int] = None
    redx_pickup_store_id: Optional[int] = None
    pathao_options: Optional[PathaoOptions] = None
    redx_options: Optional[RedXOptions] = None

jukto.interfaces.logistics.BaseLogisticsProvider

Bases: TransportLifecycle, ABC

Source code in src/jukto/interfaces/logistics.py
class BaseLogisticsProvider(TransportLifecycle, ABC):
    provider_name: str
    """The Contract: All logistics providers MUST implement these methods."""

    def __init__(
        self,
        api_key: str,
        secret_key: Optional[str] = None,
        base_url: Optional[str] = None,
        timeout: float | httpx.Timeout = 30.0,
        *,
        allow_insecure_http: bool = False,
        http_client: Optional[httpx.Client] = None,
        limits: Optional[httpx.Limits] = None,
    ) -> None:
        if base_url is not None:
            validate_endpoint(base_url, allow_insecure_http)
        self.allow_insecure_http = allow_insecure_http
        self.api_key = api_key
        self.secret_key = secret_key
        self.base_url = base_url
        self.timeout = timeout
        self._transport = SyncTransport(timeout=timeout, limits=limits, http_client=http_client,
                                        allow_insecure_http=allow_insecure_http)

    @abstractmethod
    def create_order(self, parcel: Parcel) -> dict[str, Any]:
        """Creates a new shipment/parcel."""
        pass

    @abstractmethod
    def check_status(self, consignment_id: str) -> dict[str, Any]:
        """Checks the status of a specific parcel."""
        pass

    @property
    def capabilities(self) -> ProviderCapabilities:
        return ProviderCapabilities(self.provider_name, shipment_creation=True, shipment_status=True,
                                    declared_value=self.provider_name == 'RedX',
                                    geography_options=self.provider_name in ('Pathao','RedX'))

    def create_shipment(self, parcel: Parcel) -> ShipmentCreationResult:
        raw = self.create_order(parcel)
        return shipment_creation(self.provider_name, raw)

    def check_shipment_status(self, consignment_id: str) -> ShipmentStatusResult:
        raw = self._status_payload(consignment_id)
        return shipment_status(self.provider_name, consignment_id, raw)

    def _status_payload(self, consignment_id: str) -> dict[str, Any]:
        return self.check_status(consignment_id)

provider_name instance-attribute

The Contract: All logistics providers MUST implement these methods.

check_status(consignment_id) abstractmethod

Checks the status of a specific parcel.

Source code in src/jukto/interfaces/logistics.py
@abstractmethod
def check_status(self, consignment_id: str) -> dict[str, Any]:
    """Checks the status of a specific parcel."""
    pass

create_order(parcel) abstractmethod

Creates a new shipment/parcel.

Source code in src/jukto/interfaces/logistics.py
@abstractmethod
def create_order(self, parcel: Parcel) -> dict[str, Any]:
    """Creates a new shipment/parcel."""
    pass

jukto.options.PathaoOptions dataclass

Source code in src/jukto/options.py
@dataclass(frozen=True)
class PathaoOptions:
    store_id: Optional[int] = None
    recipient_city: Optional[int] = None
    recipient_zone: Optional[int] = None
    recipient_area: Optional[int] = None
    item_type: int = 2
    delivery_type: int = 48

jukto.options.RedXOptions dataclass

Source code in src/jukto/options.py
@dataclass(frozen=True)
class RedXOptions:
    delivery_area: str
    delivery_area_id: int
    pickup_store_id: Optional[int] = None

jukto.options.ProviderCapabilities dataclass

Source code in src/jukto/options.py
@dataclass(frozen=True)
class ProviderCapabilities:
    provider: str
    shipment_creation: bool = False
    shipment_status: bool = False
    sms_submission: bool = False
    per_recipient_submission: bool = False
    declared_value: bool = False
    geography_options: bool = False

jukto.results.ShipmentCreationResult dataclass

Source code in src/jukto/results.py
@dataclass(frozen=True)
class ShipmentCreationResult:
    provider: str
    provider_id: Optional[str]
    status: Outcome
    provider_status: Any = field(default=None, repr=False)
    raw: Any = field(default=None, repr=False)

    @property
    def reconciliation_required(self) -> bool:
        return self.status == Outcome.UNKNOWN

jukto.results.ShipmentStatusResult dataclass

Source code in src/jukto/results.py
@dataclass(frozen=True)
class ShipmentStatusResult:
    provider: str
    provider_id: str
    status: Outcome
    provider_status: Any = field(default=None, repr=False)
    raw: Any = field(default=None, repr=False)

jukto.providers.logistics.steadfast.SteadfastClient

Bases: SteadfastContract, BaseLogisticsProvider

Adapter for Steadfast Courier API.

Source code in src/jukto/providers/logistics/steadfast.py
class SteadfastClient(SteadfastContract, BaseLogisticsProvider):
    """Adapter for Steadfast Courier API."""

    provider_name = "Steadfast"

    DEFAULT_URL = "https://portal.packzy.com/api/v1"

    def __init__(
        self,
        api_key: str,
        secret_key: str,
        base_url: Optional[str] = None,
        timeout: float | httpx.Timeout = 30.0,
        *,
        allow_insecure_http: bool = False,
        http_client: Optional[httpx.Client] = None,
        limits: Optional[httpx.Limits] = None,
    ) -> None:
        super().__init__(
            api_key=api_key,
            secret_key=secret_key,
            base_url=base_url or self.DEFAULT_URL,
            timeout=timeout,
            allow_insecure_http=allow_insecure_http,
            http_client=http_client,
            limits=limits,
        )
        self.headers = {
            "Api-Key": self.api_key,
            "Secret-Key": self.secret_key,
            "Content-Type": "application/json",
        }


    def _request(self, method: str, endpoint: str, **kwargs: Any) -> dict[str, Any]:
        """Execute HTTP request with robust error handling and debug logging."""
        validate_endpoint(endpoint, self.allow_insecure_http)
        kwargs["follow_redirects"] = False

        logger.debug('provider=Steadfast operation=request event=request')

        try:
            response = self._transport.request(
                method=method,
                url=endpoint,
                headers=self.headers,
                **kwargs,
            )
            logger.debug('provider=Steadfast operation=request event=response status=%d', response.status_code)
            response.raise_for_status()
            return self._inspect_response(response, method, endpoint)
        except httpx.TimeoutException as exc:
            logger.error('provider=Steadfast operation=request event=error')
            raise ProviderTimeoutError(
                'Request to Steadfast timed out: [details withheld]',
                provider="Steadfast",
            ) from None
        except httpx.HTTPStatusError as exc:
            self._handle_http_status_error(exc)
        except httpx.RequestError as exc:
            logger.error('provider=Steadfast operation=request event=error')
            raise ProviderConnectionError(
                'Failed to connect to Steadfast: [details withheld]',
                provider="Steadfast",
            ) from None

    @operation
    def create_order(self, parcel: Parcel) -> dict[str, Any]:
        endpoint, payload = self._prepare_order(parcel)
        return self._request("POST", endpoint, json=payload)

    @operation
    def check_status(self, consignment_id: str) -> dict[str, Any]:
        endpoint = f"{self.base_url}/status_by_cid/{consignment_id}"
        return self._request("GET", endpoint)

jukto.providers.logistics.pathao.PathaoClient

Bases: PathaoContract, BaseLogisticsProvider

Adapter for Pathao Courier API.

Implements automatic OAuth2 token issue, caching, and token refresh.

Source code in src/jukto/providers/logistics/pathao.py
class PathaoClient(PathaoContract, BaseLogisticsProvider):
    """Adapter for Pathao Courier API.

    Implements automatic OAuth2 token issue, caching, and token refresh.
    """

    provider_name = "Pathao"

    PRODUCTION_URL = "https://api-hermes.pathao.com"
    SANDBOX_URL = "https://courier-api-sandbox.pathao.com"

    def __init__(
        self,
        client_id: str,
        client_secret: str,
        username: str,
        password: str,
        store_id: Optional[int] = None,
        base_url: Optional[str] = None,
        sandbox: bool = False,
        timeout: float | httpx.Timeout = 30.0,
        *,
        allow_insecure_http: bool = False,
        http_client: Optional[httpx.Client] = None,
        limits: Optional[httpx.Limits] = None,
        retry_unauthorized_orders: bool = False,
    ) -> None:
        self.client_id = client_id
        self.client_secret = client_secret
        self.username = username
        self.password = password
        self.store_id = store_id
        self.sandbox = sandbox

        resolved_url = base_url or (self.SANDBOX_URL if sandbox else self.PRODUCTION_URL)
        super().__init__(
            api_key=client_id,
            secret_key=client_secret,
            base_url=resolved_url,
            timeout=timeout,
            allow_insecure_http=allow_insecure_http,
            http_client=http_client,
            limits=limits,
        )

        self.access_token: Optional[str] = None
        self.refresh_token: Optional[str] = None
        self.token_expiry: float = 0.0
        self.retry_unauthorized_orders = retry_unauthorized_orders
        self._token_lock = RLock()
        self._token_generation = 0
        self._refresh_margin = 60.0


    def _exchange_token(self, refresh: bool) -> str:
        """Called under the token lock; publish only a fully validated response."""
        endpoint = f"{self.base_url}/aladdin/api/v1/issue-token"
        payload = self._token_payload(refresh)
        started = time.monotonic()
        logger.debug('provider=Pathao operation=exchange_token event=auth')
        try:
            response = self._transport.request("POST", endpoint, json=payload,
                headers={"Accept": "application/json", "Content-Type": "application/json"})
            # RFC 6749 section 5.2: only this explicit rejection permits fallback.
            body = error_body(response) if response.status_code == 400 else None
            if refresh and isinstance(body, dict) and body.get("error") == "invalid_grant":
                return self._exchange_token(False)
            response.raise_for_status()
            return self._publish_token(response, refresh, started)
        except httpx.TimeoutException:
            raise ProviderTimeoutError('Pathao token request timed out: [details withheld]', provider="Pathao") from None
        except httpx.HTTPStatusError as exc:
            self._handle_http_status_error(exc)
            raise
        except httpx.RequestError:
            raise ProviderConnectionError('Failed to connect to Pathao auth server: [details withheld]', provider="Pathao") from None

    def _issue_token(self) -> str:
        with self._token_lock:
            return self._exchange_token(False)

    def _refresh_access_token(self) -> str:
        with self._token_lock:
            return self._exchange_token(bool(self.refresh_token))

    def _get_token_snapshot(self, rejected_generation: Optional[int] = None) -> tuple[str, int]:
        self._transport.ensure_open()
        # Recheck after acquiring the lock. A stale 401 cannot invalidate a newer token.
        with self._token_lock:
            rejected = rejected_generation == self._token_generation
            if self.access_token and not rejected and time.monotonic() < self.token_expiry - self._refresh_margin:
                return self.access_token, self._token_generation
            token = self._exchange_token(bool(self.refresh_token))
            return token, self._token_generation

    def get_token(self) -> str:
        """Return a token from this instance's process-local, synchronized cache."""
        return self._get_token_snapshot()[0]


    def _request(self, method: str, endpoint: str, **kwargs: Any) -> dict[str, Any]:
        """Replay reads once on 401; shipment replay requires explicit opt-in."""
        validate_endpoint(endpoint, self.allow_insecure_http)
        kwargs["follow_redirects"] = False
        token, generation = self._get_token_snapshot()
        headers = {"Accept": "application/json", "Content-Type": "application/json"}
        headers.update(kwargs.pop("headers", {}))
        replay = method.upper() == "GET" or (method.upper() == "POST" and endpoint.endswith("/orders") and self.retry_unauthorized_orders)
        for attempt in range(2):
            headers["Authorization"] = f"Bearer {token}"
            logger.debug('provider=Pathao operation=request event=request')
            try:
                response = self._transport.request(method=method, url=endpoint, headers=headers, **kwargs)
                logger.debug('provider=Pathao operation=request event=response status=%d', response.status_code)
                response.raise_for_status()
                return self._inspect_response(response, endpoint)
            except httpx.TimeoutException:
                raise ProviderTimeoutError('Pathao request timed out: [details withheld]', provider="Pathao") from None
            except httpx.HTTPStatusError as exc:
                if attempt == 0 and replay and exc.response.status_code == 401:
                    token, generation = self._get_token_snapshot(rejected_generation=generation)
                    continue
                self._handle_http_status_error(exc)
                raise
            except httpx.RequestError:
                raise ProviderConnectionError('Failed to connect to Pathao: [details withheld]', provider="Pathao") from None
        raise AssertionError("unreachable retry state")

    @operation
    def create_order(self, parcel: Parcel) -> dict[str, Any]:
        endpoint, payload = self._prepare_order(parcel)
        return self._request("POST", endpoint, json=payload)

    @operation
    def check_status(self, consignment_id: str) -> dict[str, Any]:
        endpoint = f"{self.base_url}/aladdin/api/v1/orders/{consignment_id}/info"
        return self._request("GET", endpoint)

    @operation
    def get_stores(self) -> dict[str, Any]:
        """Fetch merchant store locations registered with Pathao."""
        endpoint = f"{self.base_url}/aladdin/api/v1/stores"
        return self._request("GET", endpoint)

get_stores()

Fetch merchant store locations registered with Pathao.

Source code in src/jukto/providers/logistics/pathao.py
@operation
def get_stores(self) -> dict[str, Any]:
    """Fetch merchant store locations registered with Pathao."""
    endpoint = f"{self.base_url}/aladdin/api/v1/stores"
    return self._request("GET", endpoint)

get_token()

Return a token from this instance's process-local, synchronized cache.

Source code in src/jukto/providers/logistics/pathao.py
def get_token(self) -> str:
    """Return a token from this instance's process-local, synchronized cache."""
    return self._get_token_snapshot()[0]

jukto.providers.logistics.redx.RedXClient

Bases: RedXContract, BaseLogisticsProvider

Adapter for RedX Logistics API.

Source code in src/jukto/providers/logistics/redx.py
class RedXClient(RedXContract, BaseLogisticsProvider):
    """Adapter for RedX Logistics API."""

    provider_name = "RedX"

    PRODUCTION_URL = "https://openapi.redx.com.bd/v1.0.0-beta"
    SANDBOX_URL = "https://sandbox.redx.com.bd/v1.0.0-beta"

    def __init__(
        self,
        api_key: str,
        base_url: Optional[str] = None,
        sandbox: bool = False,
        timeout: float | httpx.Timeout = 30.0,
        *,
        allow_insecure_http: bool = False,
        http_client: Optional[httpx.Client] = None,
        limits: Optional[httpx.Limits] = None,
    ) -> None:
        resolved_url = base_url or (self.SANDBOX_URL if sandbox else self.PRODUCTION_URL)
        super().__init__(
            api_key=api_key,
            secret_key=None,
            base_url=resolved_url,
            timeout=timeout,
            allow_insecure_http=allow_insecure_http,
            http_client=http_client,
            limits=limits,
        )
        self.headers = {
            "API-ACCESS-TOKEN": f"Bearer {self.api_key}",
            "Content-Type": "application/json",
            "Accept": "application/json",
        }


    def _request(self, method: str, endpoint: str, **kwargs: Any) -> dict[str, Any]:
        validate_endpoint(endpoint, self.allow_insecure_http)
        kwargs["follow_redirects"] = False

        logger.debug('provider=RedX operation=request event=request')

        try:
            response = self._transport.request(
                method=method,
                url=endpoint,
                headers=self.headers,
                **kwargs,
            )
            logger.debug('provider=RedX operation=request event=response status=%d', response.status_code)
            response.raise_for_status()
            return self._inspect_response(response, method, endpoint)
        except httpx.TimeoutException as exc:
            logger.error('provider=RedX operation=request event=timeout')
            raise ProviderTimeoutError(
                'Request to RedX timed out: [details withheld]',
                provider="RedX",
            ) from None
        except httpx.HTTPStatusError as exc:
            self._handle_http_status_error(exc)
            raise
        except httpx.RequestError as exc:
            logger.error('provider=RedX operation=request event=error')
            raise ProviderConnectionError(
                'Failed to connect to RedX: [details withheld]',
                provider="RedX",
            ) from None

    @operation
    def create_order(self, parcel: Parcel) -> dict[str, Any]:
        endpoint, payload = self._prepare_order(parcel)
        return self._request("POST", endpoint, json=payload)

    @operation
    def check_status(self, consignment_id: str) -> dict[str, Any]:
        endpoint = f"{self.base_url}/parcels/tracking/{consignment_id}"
        return self._request("GET", endpoint)

    @operation
    def _status_payload(self, consignment_id: str) -> dict[str, Any]:
        return self._request("GET", f"{self.base_url}/parcel/info/{consignment_id}")

Native async (0.1.0a1 alpha)

jukto.interfaces.asynchronous.AsyncBaseLogisticsProvider

Bases: AsyncTransportLifecycle, ABC

Source code in src/jukto/interfaces/asynchronous.py
class AsyncBaseLogisticsProvider(AsyncTransportLifecycle, ABC):
    provider_name: str
    """The Contract: All logistics providers MUST implement these methods."""

    def __init__(
        self,
        api_key: str,
        secret_key: Optional[str] = None,
        base_url: Optional[str] = None,
        timeout: float | httpx.Timeout = 30.0,
        *,
        allow_insecure_http: bool = False,
        http_client: Optional[httpx.AsyncClient] = None,
        limits: Optional[httpx.Limits] = None,
    ) -> None:
        if base_url is not None:
            validate_endpoint(base_url, allow_insecure_http)
        self.allow_insecure_http = allow_insecure_http
        self.api_key = api_key
        self.secret_key = secret_key
        self.base_url = base_url
        self.timeout = timeout
        self._transport = AsyncTransport(timeout=timeout, limits=limits, http_client=http_client,
                                        allow_insecure_http=allow_insecure_http)

    @abstractmethod
    async def create_order(self, parcel: Parcel) -> dict[str, Any]:
        """Creates a new shipment/parcel."""
        pass

    @abstractmethod
    async def check_status(self, consignment_id: str) -> dict[str, Any]:
        """Checks the status of a specific parcel."""
        pass

    @property
    def capabilities(self) -> ProviderCapabilities:
        return ProviderCapabilities(self.provider_name, shipment_creation=True, shipment_status=True,
                                    declared_value=self.provider_name == 'RedX',
                                    geography_options=self.provider_name in ('Pathao','RedX'))

    async def create_shipment(self, parcel: Parcel) -> ShipmentCreationResult:
        raw = await self.create_order(parcel)
        return shipment_creation(self.provider_name, raw)

    async def check_shipment_status(self, consignment_id: str) -> ShipmentStatusResult:
        raw = await self._status_payload(consignment_id)
        return shipment_status(self.provider_name, consignment_id, raw)

    async def _status_payload(self, consignment_id: str) -> dict[str, Any]:
        return await self.check_status(consignment_id)

provider_name instance-attribute

The Contract: All logistics providers MUST implement these methods.

check_status(consignment_id) abstractmethod async

Checks the status of a specific parcel.

Source code in src/jukto/interfaces/asynchronous.py
@abstractmethod
async def check_status(self, consignment_id: str) -> dict[str, Any]:
    """Checks the status of a specific parcel."""
    pass

create_order(parcel) abstractmethod async

Creates a new shipment/parcel.

Source code in src/jukto/interfaces/asynchronous.py
@abstractmethod
async def create_order(self, parcel: Parcel) -> dict[str, Any]:
    """Creates a new shipment/parcel."""
    pass

jukto.providers.logistics.async_steadfast.AsyncSteadfastClient

Bases: SteadfastContract, AsyncBaseLogisticsProvider

Adapter for Steadfast Courier API.

Source code in src/jukto/providers/logistics/async_steadfast.py
class AsyncSteadfastClient(SteadfastContract, AsyncBaseLogisticsProvider):
    """Adapter for Steadfast Courier API."""

    provider_name = "Steadfast"

    DEFAULT_URL = "https://portal.packzy.com/api/v1"

    def __init__(
        self,
        api_key: str,
        secret_key: str,
        base_url: Optional[str] = None,
        timeout: float | httpx.Timeout = 30.0,
        *,
        allow_insecure_http: bool = False,
        http_client: Optional[httpx.AsyncClient] = None,
        limits: Optional[httpx.Limits] = None,
    ) -> None:
        super().__init__(
            api_key=api_key,
            secret_key=secret_key,
            base_url=base_url or self.DEFAULT_URL,
            timeout=timeout,
            allow_insecure_http=allow_insecure_http,
            http_client=http_client,
            limits=limits,
        )
        self.headers = {
            "Api-Key": self.api_key,
            "Secret-Key": self.secret_key,
            "Content-Type": "application/json",
        }


    async def _request(self, method: str, endpoint: str, **kwargs: Any) -> dict[str, Any]:
        """Execute HTTP request with robust error handling and debug logging."""
        validate_endpoint(endpoint, self.allow_insecure_http)
        kwargs["follow_redirects"] = False

        logger.debug('provider=Steadfast operation=request event=request')

        try:
            response = await self._transport.request(
                method=method,
                url=endpoint,
                headers=self.headers,
                **kwargs,
            )
            logger.debug('provider=Steadfast operation=request event=response status=%d', response.status_code)
            response.raise_for_status()
            return self._inspect_response(response, method, endpoint)
        except httpx.TimeoutException as exc:
            logger.error('provider=Steadfast operation=request event=error')
            raise ProviderTimeoutError(
                'Request to Steadfast timed out: [details withheld]',
                provider="Steadfast",
            ) from None
        except httpx.HTTPStatusError as exc:
            self._handle_http_status_error(exc)
        except httpx.RequestError as exc:
            logger.error('provider=Steadfast operation=request event=error')
            raise ProviderConnectionError(
                'Failed to connect to Steadfast: [details withheld]',
                provider="Steadfast",
            ) from None

    @operation
    async def create_order(self, parcel: Parcel) -> dict[str, Any]:
        endpoint, payload = self._prepare_order(parcel)
        return await self._request("POST", endpoint, json=payload)

    @operation
    async def check_status(self, consignment_id: str) -> dict[str, Any]:
        endpoint = f"{self.base_url}/status_by_cid/{consignment_id}"
        return await self._request("GET", endpoint)

jukto.providers.logistics.async_pathao.AsyncPathaoClient

Bases: PathaoContract, AsyncBaseLogisticsProvider

Adapter for Pathao Courier API.

Implements automatic OAuth2 token issue, caching, and token refresh.

Source code in src/jukto/providers/logistics/async_pathao.py
class AsyncPathaoClient(PathaoContract, AsyncBaseLogisticsProvider):
    """Adapter for Pathao Courier API.

    Implements automatic OAuth2 token issue, caching, and token refresh.
    """

    provider_name = "Pathao"

    PRODUCTION_URL = "https://api-hermes.pathao.com"
    SANDBOX_URL = "https://courier-api-sandbox.pathao.com"

    def __init__(
        self,
        client_id: str,
        client_secret: str,
        username: str,
        password: str,
        store_id: Optional[int] = None,
        base_url: Optional[str] = None,
        sandbox: bool = False,
        timeout: float | httpx.Timeout = 30.0,
        *,
        allow_insecure_http: bool = False,
        http_client: Optional[httpx.AsyncClient] = None,
        limits: Optional[httpx.Limits] = None,
        retry_unauthorized_orders: bool = False,
    ) -> None:
        self.client_id = client_id
        self.client_secret = client_secret
        self.username = username
        self.password = password
        self.store_id = store_id
        self.sandbox = sandbox

        resolved_url = base_url or (self.SANDBOX_URL if sandbox else self.PRODUCTION_URL)
        super().__init__(
            api_key=client_id,
            secret_key=client_secret,
            base_url=resolved_url,
            timeout=timeout,
            allow_insecure_http=allow_insecure_http,
            http_client=http_client,
            limits=limits,
        )

        self.access_token: Optional[str] = None
        self.refresh_token: Optional[str] = None
        self.token_expiry: float = 0.0
        self.retry_unauthorized_orders = retry_unauthorized_orders
        self._token_lock = Lock()
        self._token_generation = 0
        self._refresh_margin = 60.0


    async def _exchange_token(self, refresh: bool) -> str:
        """Called under the token lock; publish only a fully validated response."""
        endpoint = f"{self.base_url}/aladdin/api/v1/issue-token"
        payload = self._token_payload(refresh)
        started = time.monotonic()
        logger.debug('provider=Pathao operation=exchange_token event=auth')
        try:
            response = await self._transport.request("POST", endpoint, json=payload,
                headers={"Accept": "application/json", "Content-Type": "application/json"})
            # RFC 6749 section 5.2: only this explicit rejection permits fallback.
            body = error_body(response) if response.status_code == 400 else None
            if refresh and isinstance(body, dict) and body.get("error") == "invalid_grant":
                return await self._exchange_token(False)
            response.raise_for_status()
            return self._publish_token(response, refresh, started)
        except httpx.TimeoutException:
            raise ProviderTimeoutError('Pathao token request timed out: [details withheld]', provider="Pathao") from None
        except httpx.HTTPStatusError as exc:
            self._handle_http_status_error(exc)
            raise
        except httpx.RequestError:
            raise ProviderConnectionError('Failed to connect to Pathao auth server: [details withheld]', provider="Pathao") from None

    async def _issue_token(self) -> str:
        async with self._token_lock:
            return await self._exchange_token(False)

    async def _refresh_access_token(self) -> str:
        async with self._token_lock:
            return await self._exchange_token(bool(self.refresh_token))

    async def _get_token_snapshot(self, rejected_generation: Optional[int] = None) -> tuple[str, int]:
        self._transport.ensure_open()
        # Recheck after acquiring the lock. A stale 401 cannot invalidate a newer token.
        async with self._token_lock:
            rejected = rejected_generation == self._token_generation
            if self.access_token and not rejected and time.monotonic() < self.token_expiry - self._refresh_margin:
                return self.access_token, self._token_generation
            token = await self._exchange_token(bool(self.refresh_token))
            return token, self._token_generation

    async def get_token(self) -> str:
        """Return a token from this instance's process-local, synchronized cache."""
        return (await self._get_token_snapshot())[0]


    async def _request(self, method: str, endpoint: str, **kwargs: Any) -> dict[str, Any]:
        """Replay reads once on 401; shipment replay requires explicit opt-in."""
        validate_endpoint(endpoint, self.allow_insecure_http)
        kwargs["follow_redirects"] = False
        token, generation = await self._get_token_snapshot()
        headers = {"Accept": "application/json", "Content-Type": "application/json"}
        headers.update(kwargs.pop("headers", {}))
        replay = method.upper() == "GET" or (method.upper() == "POST" and endpoint.endswith("/orders") and self.retry_unauthorized_orders)
        for attempt in range(2):
            headers["Authorization"] = f"Bearer {token}"
            logger.debug('provider=Pathao operation=request event=request')
            try:
                response = await self._transport.request(method=method, url=endpoint, headers=headers, **kwargs)
                logger.debug('provider=Pathao operation=request event=response status=%d', response.status_code)
                response.raise_for_status()
                return self._inspect_response(response, endpoint)
            except httpx.TimeoutException:
                raise ProviderTimeoutError('Pathao request timed out: [details withheld]', provider="Pathao") from None
            except httpx.HTTPStatusError as exc:
                if attempt == 0 and replay and exc.response.status_code == 401:
                    token, generation = await self._get_token_snapshot(rejected_generation=generation)
                    continue
                self._handle_http_status_error(exc)
                raise
            except httpx.RequestError:
                raise ProviderConnectionError('Failed to connect to Pathao: [details withheld]', provider="Pathao") from None
        raise AssertionError("unreachable retry state")

    @operation
    async def create_order(self, parcel: Parcel) -> dict[str, Any]:
        endpoint, payload = self._prepare_order(parcel)
        return await self._request("POST", endpoint, json=payload)

    @operation
    async def check_status(self, consignment_id: str) -> dict[str, Any]:
        endpoint = f"{self.base_url}/aladdin/api/v1/orders/{consignment_id}/info"
        return await self._request("GET", endpoint)

    @operation
    async def get_stores(self) -> dict[str, Any]:
        """Fetch merchant store locations registered with Pathao."""
        endpoint = f"{self.base_url}/aladdin/api/v1/stores"
        return await self._request("GET", endpoint)

get_stores() async

Fetch merchant store locations registered with Pathao.

Source code in src/jukto/providers/logistics/async_pathao.py
@operation
async def get_stores(self) -> dict[str, Any]:
    """Fetch merchant store locations registered with Pathao."""
    endpoint = f"{self.base_url}/aladdin/api/v1/stores"
    return await self._request("GET", endpoint)

get_token() async

Return a token from this instance's process-local, synchronized cache.

Source code in src/jukto/providers/logistics/async_pathao.py
async def get_token(self) -> str:
    """Return a token from this instance's process-local, synchronized cache."""
    return (await self._get_token_snapshot())[0]

jukto.providers.logistics.async_redx.AsyncRedXClient

Bases: RedXContract, AsyncBaseLogisticsProvider

Adapter for RedX Logistics API.

Source code in src/jukto/providers/logistics/async_redx.py
class AsyncRedXClient(RedXContract, AsyncBaseLogisticsProvider):
    """Adapter for RedX Logistics API."""

    provider_name = "RedX"

    PRODUCTION_URL = "https://openapi.redx.com.bd/v1.0.0-beta"
    SANDBOX_URL = "https://sandbox.redx.com.bd/v1.0.0-beta"

    def __init__(
        self,
        api_key: str,
        base_url: Optional[str] = None,
        sandbox: bool = False,
        timeout: float | httpx.Timeout = 30.0,
        *,
        allow_insecure_http: bool = False,
        http_client: Optional[httpx.AsyncClient] = None,
        limits: Optional[httpx.Limits] = None,
    ) -> None:
        resolved_url = base_url or (self.SANDBOX_URL if sandbox else self.PRODUCTION_URL)
        super().__init__(
            api_key=api_key,
            secret_key=None,
            base_url=resolved_url,
            timeout=timeout,
            allow_insecure_http=allow_insecure_http,
            http_client=http_client,
            limits=limits,
        )
        self.headers = {
            "API-ACCESS-TOKEN": f"Bearer {self.api_key}",
            "Content-Type": "application/json",
            "Accept": "application/json",
        }


    async def _request(self, method: str, endpoint: str, **kwargs: Any) -> dict[str, Any]:
        validate_endpoint(endpoint, self.allow_insecure_http)
        kwargs["follow_redirects"] = False

        logger.debug('provider=RedX operation=request event=request')

        try:
            response = await self._transport.request(
                method=method,
                url=endpoint,
                headers=self.headers,
                **kwargs,
            )
            logger.debug('provider=RedX operation=request event=response status=%d', response.status_code)
            response.raise_for_status()
            return self._inspect_response(response, method, endpoint)
        except httpx.TimeoutException as exc:
            logger.error('provider=RedX operation=request event=timeout')
            raise ProviderTimeoutError(
                'Request to RedX timed out: [details withheld]',
                provider="RedX",
            ) from None
        except httpx.HTTPStatusError as exc:
            self._handle_http_status_error(exc)
            raise
        except httpx.RequestError as exc:
            logger.error('provider=RedX operation=request event=error')
            raise ProviderConnectionError(
                'Failed to connect to RedX: [details withheld]',
                provider="RedX",
            ) from None

    @operation
    async def create_order(self, parcel: Parcel) -> dict[str, Any]:
        endpoint, payload = self._prepare_order(parcel)
        return await self._request("POST", endpoint, json=payload)

    @operation
    async def check_status(self, consignment_id: str) -> dict[str, Any]:
        endpoint = f"{self.base_url}/parcels/tracking/{consignment_id}"
        return await self._request("GET", endpoint)

    @operation
    async def _status_payload(self, consignment_id: str) -> dict[str, Any]:
        return await self._request("GET", f"{self.base_url}/parcel/info/{consignment_id}")