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
/importendpoint 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 |
|---|---|
|
|
Personal access token (or |
Sends header: X-OCP-PERSONAL-ACCESS-TOKEN: <value>
-
Option B - Bearer token
|
Variable |
Meaning |
|---|---|
|
|
Bearer access token copied from the browser / Bruno (or |
Option C - Keycloak password grant
This option is used when neither PAT nor token is set.
|
Variable |
Meaning |
|---|---|
|
|
Keycloak base, for example, |
|
|
Keycloak realm |
|
|
OIDC client ID |
|
|
OIDC client secret |
|
|
Username |
|
|
Password |
|
|
Optional, defaults to |
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
# 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 |
|---|---|---|
|
|
- |
Personal access token (or |
|
|
- |
Bearer token (or |
|
|
90 MiB |
Bytes; at/above this use the chunked path |
|
|
8 MiB |
Part size; must be ≥ 5 MiB for multi-part |
|
|
off |
Force the chunked path regardless of size |
|
|
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/422surface immediately, no retry. -
Resume: a
422at commit withmissing_partsre-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_idthat is polled toFINISHEDorFAILED.
Script
Save the file as ocp_app_upload.py.
#!/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())