Guides

Chunked upload and import

Overview

This script is a customer-facing Python client for importing an app archive into the OCP Environments Manager. It selects the transport automatically:

  • Files below the threshold (default 90 MiB) → use the direct /import endpoint with no upload-session overhead.

  • Files at/above the threshold (or --multipart) → use a resumable chunked upload: open a session, upload parts with bounded concurrency and per-part integrity checks, commit, then poll the background import to completion.

The chunked path always runs asynchronously: S3 assembly completes synchronously (so missing-part or conflict errors surface immediately), then the import runs in the background and is polled until it reaches a terminal state.

Requirements

  • Python 3.10+

  • httpx - pip install httpx

Authentication

For authentication, choose one of these options (in this order of preference):

  • Option A - Personal Access Token (PAT)

Variable

Meaning

OCP_PAT

Personal access token (or --pat)

Sends header: X-OCP-PERSONAL-ACCESS-TOKEN: <value>

  • Option B - Bearer token

Variable

Meaning

OCP_TOKEN

Bearer access token copied from the browser / Bruno (or --token)

Option C - Keycloak password grant

This option is used when neither PAT nor token is set.

Variable

Meaning

OCP_AUTH_SERVER_URL

Keycloak base, for example, https://host

OCP_REALM

Keycloak realm

OCP_CLIENT_ID

OIDC client ID

OCP_CLIENT_SECRET

OIDC client secret

OCP_USERNAME

Username

OCP_PASSWORD

Password

OCP_GRANT_TYPE

Optional, defaults to password

Token endpoint: {OCP_AUTH_SERVER_URL}/auth/realms/{OCP_REALM}/protocol/openid-connect/token

Do not commit or paste real credentials or tokens into tickets, pages, or logs. Use environment variables only.

API base URL

OCP_API_URL (or --api-url) - base URL up to the service prefix, for example: https://host/envs-manager.

Usage

Bash
# PAT auth
export OCP_API_URL=https://host/envs-manager
export OCP_PAT=...
python ocp_app_upload.py app.zip --target-group my-group

# Bearer token auth
export OCP_TOKEN=...
python ocp_app_upload.py app.zip --target-group my-group

# Large file (>= 90 MiB) automatically uses the chunked path
python ocp_app_upload.py big_app.zip --target-group my-group

# Force the chunked path regardless of size
python ocp_app_upload.py app.zip --target-group my-group --multipart

# Tag import
python ocp_app_upload.py app.zip --import-type tag --app-id 123

Key options

Option

Default

Purpose

--pat

-

Personal access token (or OCP_PAT)

--token

-

Bearer token (or OCP_TOKEN)

--threshold

90 MiB

Bytes; at/above this use the chunked path

--chunk-size

8 MiB

Part size; must be ≥ 5 MiB for multi-part

--multipart

off

Force the chunked path regardless of size

--timeout

300s

HTTP timeout

Behavior

  • Concurrency: 4 concurrent part uploads, bounded by a semaphore.

  • Integrity: per-part Content-MD5; whole-file SHA-256 sent at commit.

  • Retry: transient failures (429/500/502/503/504, transport/timeout) → exponential backoff + jitter, up to 5 attempts per part. 400/401/403/413/415/422 surface immediately, no retry.

  • Resume: a 422 at commit with missing_parts re-uploads only the gaps then retries commit (up to 3 times).

  • Cleanup: on any fatal failure, the session is aborted (DELETE /upload-sessions/{id}) to release S3 + Redis resources.

  • S3 limits: non-final parts must be ≥ 5 MiB; max 10000 parts.

  • Commit: always asynchronous - returns a result_id that is polled to FINISHED or FAILED.

Script

Save the file as ocp_app_upload.py.

Python
#!/usr/bin/env python3
"""
OCP Environments Manager — app import uploader.

Transparently imports an app archive into the Environments Manager:

  * Files BELOW the threshold (default 90 MiB) go straight to the existing
    direct import endpoint, with no upload-session overhead.
  * Files AT/ABOVE the threshold (or with --multipart) are uploaded in
    resumable chunks: a session is opened, parts are uploaded concurrently
    with per-part integrity checks, the upload is committed, and the
    background import is polled to a terminal state.

Authentication (checked in this order):
  * Personal Access Token:  --pat <PAT>        (or env OCP_PAT)
  * Bearer token:           --token <JWT>      (or env OCP_TOKEN)
  * Keycloak password grant: the OCP_* auth env vars below

Requires: Python 3.10+, httpx  (pip install httpx)

Environment variables
---------------------
API:
  OCP_API_URL          Base URL up to the service prefix, e.g.
                       https://host/envs-manager   (CLI: --api-url)

Auth — option A (Personal Access Token):
  OCP_PAT              Personal access token; sends X-OCP-PERSONAL-ACCESS-TOKEN.

Auth — option B (paste a bearer token):
  OCP_TOKEN            A bearer access token copied from the browser/Bruno.

Auth — option C (Keycloak password grant; used when PAT and token are unset):
  OCP_AUTH_SERVER_URL  Keycloak base, e.g. https://host       (the {server})
  OCP_REALM            Keycloak realm
  OCP_CLIENT_ID        OIDC client id
  OCP_CLIENT_SECRET    OIDC client secret
  OCP_USERNAME         Username
  OCP_PASSWORD         Password
  OCP_GRANT_TYPE       Optional, defaults to "password"

The token endpoint called is:
  {OCP_AUTH_SERVER_URL}/auth/realms/{OCP_REALM}/protocol/openid-connect/token
"""

# Standalone CLI: stdout/stderr are the interface, so `print` is intentional.
# ruff: noqa: T201
from __future__ import annotations

import argparse
import asyncio
import base64
import hashlib
import math
import os
import sys
from dataclasses import dataclass
from pathlib import Path

try:
    import httpx
except ModuleNotFoundError:  # surfaced with a friendly hint in main()
    httpx = None

# --- Defaults --------------------------------------------------------------

DEFAULT_THRESHOLD = 90 * 1024 * 1024  # 90 MiB — switch to chunked at/above
DEFAULT_CHUNK_SIZE = 8 * 1024 * 1024  # matches server UPLOAD_DEFAULT_CHUNK_SIZE
SERVER_MIN_CHUNK = 5 * 1024 * 1024  # S3 floor; non-final parts must be >= this
PART_CONCURRENCY = 4  # concurrent part uploads (AC: 4)
PART_READ_CHUNK = 1024 * 1024  # streaming read granularity for hashing
MAX_PART_RETRIES = 5  # per-part transient-failure retries
MAX_COMMIT_RETRIES = 3  # commit retries after re-uploading missing parts
POLL_INTERVAL_S = 3.0  # async result polling cadence
POLL_TIMEOUT_S = 3600.0  # give up polling after this long

# HTTP statuses that must surface immediately, never retried (AC #3).
NO_RETRY_STATUSES = {400, 401, 403, 413, 415, 422}
# Transient statuses worth retrying.
RETRY_STATUSES = {429, 500, 502, 503, 504}


class UploadError(RuntimeError):
    """Fatal, non-retryable client error."""


# --- Auth ------------------------------------------------------------------


async def fetch_token_keycloak(cfg: AuthConfig) -> str:
    """Keycloak OIDC password grant (Python port of the Bruno script)."""
    url = (
        f"{cfg.auth_server_url}/auth/realms/{cfg.realm}"
        "/protocol/openid-connect/token"
    )
    data = {
        "grant_type": cfg.grant_type,
        "client_id": cfg.client_id,
        "client_secret": cfg.client_secret,
        "username": cfg.username,
        "password": cfg.password,
    }
    async with httpx.AsyncClient(timeout=30.0) as client:
        resp = await client.post(
            url,
            data=data,
            headers={"Content-Type": "application/x-www-form-urlencoded"},
        )
    if resp.status_code != 200:
        raise UploadError(
            f"Keycloak token request failed ({resp.status_code}): {resp.text}"
        )
    body = resp.json()
    token = body.get("access_token")
    if not token:
        raise UploadError("Keycloak response did not contain an access_token")
    return token


@dataclass
class AuthConfig:
    auth_server_url: str | None = None
    realm: str | None = None
    grant_type: str = "password"
    client_id: str | None = None
    client_secret: str | None = None
    username: str | None = None
    password: str | None = None

    @classmethod
    def from_env(cls) -> AuthConfig:
        return cls(
            auth_server_url=os.environ.get("OCP_AUTH_SERVER_URL"),
            realm=os.environ.get("OCP_REALM"),
            grant_type=os.environ.get("OCP_GRANT_TYPE", "password"),
            client_id=os.environ.get("OCP_CLIENT_ID"),
            client_secret=os.environ.get("OCP_CLIENT_SECRET"),
            username=os.environ.get("OCP_USERNAME"),
            password=os.environ.get("OCP_PASSWORD"),
        )

    def is_complete(self) -> bool:
        return all(
            [
                self.auth_server_url,
                self.realm,
                self.client_id,
                self.client_secret,
                self.username,
                self.password,
            ]
        )


async def resolve_auth_headers(args: argparse.Namespace) -> dict[str, str]:
    """Return auth headers for all requests.

    Precedence: PAT > bearer token > Keycloak password grant.
    """
    pat = args.pat or os.environ.get("OCP_PAT")
    if pat:
        return {"X-OCP-PERSONAL-ACCESS-TOKEN": pat}

    token = args.token or os.environ.get("OCP_TOKEN")
    if token:
        return {"Authorization": f"Bearer {token}"}

    auth = AuthConfig.from_env()
    if not auth.is_complete():
        raise UploadError(
            "No auth provided. Set OCP_PAT, OCP_TOKEN, or all of "
            "OCP_AUTH_SERVER_URL / OCP_REALM / OCP_CLIENT_ID / "
            "OCP_CLIENT_SECRET / OCP_USERNAME / OCP_PASSWORD."
        )
    token = await fetch_token_keycloak(auth)
    return {"Authorization": f"Bearer {token}"}


# --- Helpers ---------------------------------------------------------------


def _b64_md5(data: bytes) -> str:
    """Base64-encoded MD5 of bytes (S3 Content-MD5, RFC 1864)."""
    return base64.b64encode(hashlib.md5(data).digest()).decode()


def _sha256_hex(path: Path) -> str:
    """Whole-file SHA-256 (hex) computed by streaming the file once."""
    h = hashlib.sha256()
    with path.open("rb") as fh:
        for block in iter(lambda: fh.read(PART_READ_CHUNK), b""):
            h.update(block)
    return h.hexdigest()


def _read_part(path: Path, index: int, chunk_size: int) -> bytes:
    """Read the 0-based part `index` of `chunk_size` bytes from `path`."""
    with path.open("rb") as fh:
        fh.seek(index * chunk_size)
        return fh.read(chunk_size)


def _v1(base: str, suffix: str) -> str:
    return f"{base.rstrip('/')}/api/v1{suffix}"


def _raise_for_fatal(resp: httpx.Response, context: str) -> None:
    """Raise UploadError for statuses that must never be retried."""
    if resp.status_code in NO_RETRY_STATUSES:
        raise UploadError(f"{context}: {resp.status_code} {resp.text}".strip())


def _backoff_delay(attempt: int) -> float:
    """Exponential backoff with jitter. attempt is 0-based."""
    import secrets

    base = min(2.0**attempt, 30.0)
    return base + secrets.SystemRandom().uniform(0, 1.0)


# --- Direct import (below threshold) --------------------------------------


async def direct_import(
    client: httpx.AsyncClient,
    base: str,
    path: Path,
    args: argparse.Namespace,
) -> httpx.Response:
    """POST the whole file to the existing direct import endpoint."""
    files = {"file": (path.name, path.open("rb"), "application/zip")}
    if args.import_type == "tag":
        if not args.app_id:
            raise UploadError("--app-id is required for a tag import")
        url = _v1(base, f"/apps/{args.app_id}/tags/import")
        data = {"deployment_type": args.deployment_type}
        if args.vb_profile_id:
            data["vb_profile_id"] = args.vb_profile_id
    else:
        url = _v1(base, "/apps/import")
        data = {
            "target_group": args.target_group,
            "deployment_type": args.deployment_type,
            "include_vars": str(args.include_vars).lower(),
            "include_nlu": str(args.include_nlu).lower(),
        }
        if args.target_app_name:
            data["target_app_name"] = args.target_app_name
        if args.vb_profile_id:
            data["vb_profile_id"] = args.vb_profile_id

    resp = await client.post(url, data=data, files=files)
    _raise_for_fatal(resp, "Direct import")
    if resp.status_code >= 400:
        raise UploadError(
            f"Direct import failed: {resp.status_code} {resp.text}"
        )
    return resp


# --- Chunked upload (at/above threshold) ----------------------------------


async def init_session(
    client: httpx.AsyncClient,
    base: str,
    path: Path,
    size: int,
    total_parts: int,
    args: argparse.Namespace,
) -> dict:
    """Open an upload session and return the server's session parameters.

    total_parts is declared here and is binding: the server validates commit
    coverage against it, so the caller must slice the file into exactly this
    many parts using the same chunk size it derived total_parts from.
    """
    payload = {
        "filename": path.name,
        "expected_size": size,
        "total_parts": total_parts,
        "import_type": args.import_type,
        "target_group": args.target_group,
        "include_vars": args.include_vars,
        "include_nlu": args.include_nlu,
    }
    if args.app_id:
        payload["app_id"] = args.app_id
    if args.deployment_type:
        payload["deployment_type"] = args.deployment_type
    if args.target_app_name:
        payload["target_app_name"] = args.target_app_name
    if args.vb_profile_id:
        payload["vb_profile_id"] = args.vb_profile_id

    resp = await client.post(_v1(base, "/upload-sessions"), json=payload)
    _raise_for_fatal(resp, "Init session")
    if resp.status_code != 201:
        raise UploadError(
            f"Init session failed: {resp.status_code} {resp.text}"
        )
    return resp.json()["data"]


async def upload_part(
    client: httpx.AsyncClient,
    base: str,
    session_id: str,
    path: Path,
    part_number: int,
    chunk_size: int,
    sem: asyncio.Semaphore,
) -> dict:
    """Upload a single part with bounded concurrency and transient retry."""
    body = await asyncio.to_thread(
        _read_part, path, part_number - 1, chunk_size
    )
    content_md5 = _b64_md5(body)
    url = _v1(base, f"/upload-sessions/{session_id}/parts/{part_number}")

    async with sem:
        last_exc: Exception | None = None
        for attempt in range(MAX_PART_RETRIES):
            try:
                resp = await client.put(
                    url,
                    files={"chunk": (f"part-{part_number}", body)},
                    data={"content_md5": content_md5},
                )
                _raise_for_fatal(resp, f"Part {part_number}")
                if resp.status_code == 200:
                    data = resp.json()["data"]
                    return {"part_number": part_number, "etag": data["etag"]}
                if resp.status_code not in RETRY_STATUSES:
                    raise UploadError(
                        f"Part {part_number} failed: "
                        f"{resp.status_code} {resp.text}"
                    )
                last_exc = UploadError(
                    f"Part {part_number}: {resp.status_code}"
                )
            except (httpx.TransportError, httpx.TimeoutException) as exc:
                last_exc = exc
            await asyncio.sleep(_backoff_delay(attempt))
        raise UploadError(f"Part {part_number} exhausted retries: {last_exc}")


async def upload_parts(
    client: httpx.AsyncClient,
    base: str,
    session_id: str,
    path: Path,
    part_numbers: list[int],
    chunk_size: int,
) -> list[dict]:
    """Upload the given part numbers concurrently (bounded by the semaphore)."""
    sem = asyncio.Semaphore(PART_CONCURRENCY)
    tasks = [
        upload_part(client, base, session_id, path, n, chunk_size, sem)
        for n in part_numbers
    ]
    return await asyncio.gather(*tasks)


async def commit_session(
    client: httpx.AsyncClient,
    base: str,
    session_id: str,
    parts: list[dict],
    checksum_sha256: str,
    path: Path,
    chunk_size: int,
) -> httpx.Response:
    """Commit the upload, re-uploading missing parts on a 422 and retrying.

    The server always returns 202; the caller is responsible for polling the
    returned result_id to a terminal state.
    """
    url = _v1(base, f"/upload-sessions/{session_id}/commit")
    by_number = {p["part_number"]: p for p in parts}

    for _attempt in range(MAX_COMMIT_RETRIES):
        payload = {
            "parts": sorted(by_number.values(), key=lambda p: p["part_number"]),
            "checksum_sha256": checksum_sha256,
        }
        resp = await client.post(url, json=payload)

        if resp.status_code == 202:
            return resp

        # 422 with missing_parts -> re-upload only the gaps, then retry commit.
        if resp.status_code == 422:
            missing = _missing_parts(resp)
            if missing:
                print(f"  commit: re-uploading missing parts {missing}")
                fresh = await upload_parts(
                    client, base, session_id, path, missing, chunk_size
                )
                for p in fresh:
                    by_number[p["part_number"]] = p
                continue
            # 422 without missing_parts (e.g. checksum) is fatal.
            raise UploadError(
                f"Commit rejected: {resp.status_code} {resp.text}"
            )

        _raise_for_fatal(resp, "Commit")
        raise UploadError(f"Commit failed: {resp.status_code} {resp.text}")

    raise UploadError("Commit exhausted retries after re-uploading parts")


def _missing_parts(resp: httpx.Response) -> list[int]:
    try:
        for err in resp.json().get("errors", []):
            meta = err.get("meta") or {}
            if "missing_parts" in meta:
                return list(meta["missing_parts"])
    except ValueError:
        pass
    return []


async def chunked_upload(
    client: httpx.AsyncClient,
    base: str,
    path: Path,
    size: int,
    args: argparse.Namespace,
) -> httpx.Response:
    """Full chunked path: init -> parts -> commit, with abort on fatal error.

    The client's chunk size is authoritative: total_parts is computed from it
    and declared at init, and the file is sliced by the same size so the parts
    line up with the declared count.
    """
    chunk_size = args.chunk_size
    total_parts = max(1, math.ceil(size / chunk_size))

    # Non-final parts must meet the S3 5 MiB floor; a single-part upload is
    # exempt. Fail fast with a clear message instead of a server 400 mid-upload.
    if total_parts > 1 and chunk_size < SERVER_MIN_CHUNK:
        raise UploadError(
            f"--chunk-size {chunk_size} is below the {SERVER_MIN_CHUNK}-byte "
            "minimum for multi-part uploads (5 MiB)"
        )
    if total_parts > 10_000:
        raise UploadError(
            f"{total_parts} parts exceeds the 10000 limit; raise --chunk-size"
        )

    init = await init_session(client, base, path, size, total_parts, args)
    session_id = init["session_id"]
    print(
        f"  session {session_id}: {total_parts} part(s) of {chunk_size} bytes"
    )

    try:
        checksum = await asyncio.to_thread(_sha256_hex, path)
        parts = await upload_parts(
            client,
            base,
            session_id,
            path,
            list(range(1, total_parts + 1)),
            chunk_size,
        )
        return await commit_session(
            client, base, session_id, parts, checksum, path, chunk_size
        )
    except BaseException:
        # On any fatal failure release server-side resources (S3 + Redis).
        await _abort_quietly(client, base, session_id)
        raise


async def _abort_quietly(
    client: httpx.AsyncClient, base: str, session_id: str
) -> None:
    try:
        await client.delete(_v1(base, f"/upload-sessions/{session_id}"))
        print(f"  aborted session {session_id}")
    except Exception as exc:  # best-effort cleanup
        print(f"  warning: abort of {session_id} failed: {exc}")


# --- Async result polling --------------------------------------------------


async def poll_result(
    client: httpx.AsyncClient, base: str, result_id: str
) -> dict:
    """Poll /results/{id} until FINISHED or FAILED (or timeout)."""
    url = _v1(base, f"/results/{result_id}")
    waited = 0.0
    while waited < POLL_TIMEOUT_S:
        resp = await client.get(url)
        _raise_for_fatal(resp, "Poll result")
        if resp.status_code == 200:
            data = resp.json()["data"]
            status = data.get("status")
            if status in ("FINISHED", "FAILED"):
                return data
        await asyncio.sleep(POLL_INTERVAL_S)
        waited += POLL_INTERVAL_S
    raise UploadError(f"Result {result_id} did not finish within timeout")


# --- Orchestration ---------------------------------------------------------


async def run(args: argparse.Namespace) -> int:
    path = Path(args.file)
    if not path.is_file():
        raise UploadError(f"File not found: {path}")
    size = path.stat().st_size
    base = (args.api_url or os.environ.get("OCP_API_URL") or "").rstrip("/")
    if not base:
        raise UploadError("API base URL required (--api-url or OCP_API_URL)")

    headers = await resolve_auth_headers(args)
    use_chunked = args.multipart or size >= args.threshold
    mode = "chunked" if use_chunked else "direct"
    why = (
        "forced --multipart"
        if args.multipart
        else f"threshold {args.threshold}"
    )
    print(f"{path.name}: {size} bytes -> {mode} ({why})")

    async with httpx.AsyncClient(
        headers=headers, timeout=args.timeout
    ) as client:
        if mode == "direct":
            resp = await direct_import(client, base, path, args)
            report = resp.json().get("data")
            print(f"  done ({resp.status_code}): {report}")
            return 0

        # Chunked: commit always returns 202; poll to a terminal state.
        resp = await chunked_upload(client, base, path, size, args)
        result_id = resp.json()["data"]["result_id"]
        print(f"  result_id={result_id}; polling...")
        result = await poll_result(client, base, result_id)
        print(f"  result: {result.get('status')}")
        return 0 if result.get("status") == "FINISHED" else 1


def parse_args(argv: list[str]) -> argparse.Namespace:
    p = argparse.ArgumentParser(
        description="Import an app archive into OCP Environments Manager "
        "(direct for small files, chunked for large)."
    )
    p.add_argument("file", help="Path to the app .zip archive")
    p.add_argument(
        "--api-url", help="API base incl. service prefix (or OCP_API_URL)"
    )
    p.add_argument("--pat", help="Personal access token (or OCP_PAT env)")
    p.add_argument("--token", help="Bearer token (or OCP_TOKEN env)")

    p.add_argument("--import-type", choices=["app", "tag"], default="app")
    p.add_argument(
        "--target-group", help="Target group (required for app import)"
    )
    p.add_argument("--app-id", help="App id (required for tag import)")
    p.add_argument(
        "--deployment-type",
        default="production",
        choices=["production", "testing"],
    )
    p.add_argument("--target-app-name", help="Override imported app name")
    p.add_argument("--vb-profile-id", help="Voice biometrics profile id")
    p.add_argument("--include-vars", action="store_true", default=True)
    p.add_argument(
        "--no-include-vars", dest="include_vars", action="store_false"
    )
    p.add_argument("--include-nlu", action="store_true", default=False)

    p.add_argument(
        "--multipart",
        action="store_true",
        help="Force the chunked/multipart path regardless of file size",
    )
    p.add_argument(
        "--threshold",
        type=int,
        default=DEFAULT_THRESHOLD,
        help=f"Bytes; at/above this use chunked (default {DEFAULT_THRESHOLD})",
    )
    p.add_argument(
        "--chunk-size",
        type=int,
        default=DEFAULT_CHUNK_SIZE,
        help="Part size in bytes; >=5 MiB for multi-part (default 8 MiB)",
    )
    p.add_argument("--timeout", type=float, default=300.0)

    args = p.parse_args(argv)
    if args.import_type == "app" and not args.target_group:
        p.error("--target-group is required for an app import")
    if args.import_type == "tag" and not args.app_id:
        p.error("--app-id is required for a tag import")
    return args


def main() -> int:
    args = parse_args(sys.argv[1:])
    if httpx is None:
        print(
            "ERROR: this script requires httpx (pip install httpx)",
            file=sys.stderr,
        )
        return 1
    try:
        return asyncio.run(run(args))
    except UploadError as exc:
        print(f"ERROR: {exc}", file=sys.stderr)
        return 1
    except KeyboardInterrupt:
        print("interrupted", file=sys.stderr)
        return 130


if __name__ == "__main__":
    raise SystemExit(main())