This is an automated email from the ASF dual-hosted git repository.
yzeng1618 pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new ce5a6ba4cc [Feature][CLI] Add bedrock-mantle provider for
OpenAI-family Bedrock models (#11548)
ce5a6ba4cc is described below
commit ce5a6ba4cc86c674d5c1ca1b93ab9f5fce9ca425
Author: SEZ <[email protected]>
AuthorDate: Fri Jul 31 09:47:56 2026 +0800
[Feature][CLI] Add bedrock-mantle provider for OpenAI-family Bedrock models
(#11548)
Co-authored-by: SEZ9 <[email protected]>
---
docs/en/ai-cli/quickstart.md | 45 ++-
docs/zh/ai-cli/quickstart.md | 40 ++-
seatunnel-cli/README.md | 25 +-
seatunnel-cli/env.example.sh | 3 +
seatunnel-cli/pyproject.toml | 5 +-
seatunnel-cli/seatunnel_cli/cli.py | 32 +-
seatunnel-cli/seatunnel_cli/llm_provider.py | 308 +++++++++++++++-
seatunnel-cli/tests/test_cli_provider_routing.py | 85 +++++
.../tests/test_llm_provider_bedrock_mantle.py | 388 +++++++++++++++++++++
9 files changed, 917 insertions(+), 14 deletions(-)
diff --git a/docs/en/ai-cli/quickstart.md b/docs/en/ai-cli/quickstart.md
index 27b858106c..a1ea10f88f 100644
--- a/docs/en/ai-cli/quickstart.md
+++ b/docs/en/ai-cli/quickstart.md
@@ -45,7 +45,8 @@ seatunnel --init # interactive provider setup
export AI_PROVIDER=bedrock
export AWS_REGION=us-east-1
-# Option A2: OpenAI-family Bedrock models (Responses-API-only, e.g.
gpt-5.6-terra)
+# Option A2: OpenAI-family Bedrock models (bedrock-mantle) — see the
+# dedicated section below for the full contract
export AI_PROVIDER=bedrock-mantle
export OPENAI_MODEL='openai.gpt-5.6-terra'
@@ -59,6 +60,48 @@ export OPENAI_API_KEY=sk-...
# export OPENAI_BASE_URL=https://... # Azure OpenAI, DeepSeek, local vLLM,
...
```
+### bedrock-mantle: OpenAI-family models on Bedrock
+
+Some OpenAI models on Bedrock (e.g. `openai.gpt-5.6-terra`,
`openai.gpt-5.6-sol`)
+are not in the Bedrock foundation-model catalog and only support the OpenAI
+**Responses API** on the dedicated `bedrock-mantle` endpoint — the regular
+`bedrock` provider (Converse API) and the `openai` provider (Chat Completions)
+cannot reach them. Use the `bedrock-mantle` provider:
+
+```bash
+# 1. Install the provider extra (openai SDK >= 2.45 + AWS token generator)
+pip install -e ".[bedrock-mantle]"
+
+# 2. Configure — AWS credentials only, no OpenAI account or API key needed
+export AI_PROVIDER=bedrock-mantle
+export AWS_REGION=us-east-1 # us-east-1 / us-east-2 /
us-west-2
+export OPENAI_MODEL='openai.gpt-5.6-terra' # default if unset
+# export OPENAI_SMALL_FAST_MODEL='openai.gpt-5.6-terra'
+
+# 3. Generate as usual
+seatunnel "Sync MySQL users table to S3 Parquet"
+```
+
+Provider contract:
+
+- **Endpoint**: `https://bedrock-mantle.{region}.api.aws/openai/v1` — the
+ model-specific `openai/v1` path required by these models (the generic `v1`
+ Responses path rejects them).
+- **Auth**: a short-term bearer token is derived automatically from your AWS
+ credentials (profile, env vars, or IAM role) via
`aws-bedrock-token-generator`
+ and refreshed every 30 minutes. No long-lived key is stored anywhere.
+- **Data retention**: every request is sent with `store=false`, so Bedrock does
+ not retain your prompts or generated configs server-side (the service default
+ would otherwise keep them for 30 days).
+- **Parameters**: these models reject `temperature`; the provider never sends
+ it, so any configured temperature value is not applied.
+- **Errors**: truncated (`incomplete`), failed, and refused responses raise an
+ explicit error instead of being returned as a normal answer.
+
+The provider fully supports the CLI's internal tool-calling loop (connector
+lookups during planning) and multi-turn sessions, including replay of the
+model's reasoning output between tool calls.
+
API keys are read from environment variables only — they are never written to
config files.
## Generate Your First Pipeline
diff --git a/docs/zh/ai-cli/quickstart.md b/docs/zh/ai-cli/quickstart.md
index 1b73f11906..d7b722e982 100644
--- a/docs/zh/ai-cli/quickstart.md
+++ b/docs/zh/ai-cli/quickstart.md
@@ -45,7 +45,7 @@ seatunnel --init # 交互式配置提供商
export AI_PROVIDER=bedrock
export AWS_REGION=us-east-1
-# 方式 A2:Bedrock 上的 OpenAI 系模型(仅支持 Responses API,如 gpt-5.6-terra)
+# 方式 A2:Bedrock 上的 OpenAI 系模型(bedrock-mantle,完整说明见下方专节)
export AI_PROVIDER=bedrock-mantle
export OPENAI_MODEL='openai.gpt-5.6-terra'
@@ -59,6 +59,44 @@ export OPENAI_API_KEY=sk-...
# export OPENAI_BASE_URL=https://... # Azure OpenAI、DeepSeek、本地 vLLM 等
```
+### bedrock-mantle:Bedrock 上的 OpenAI 系模型
+
+Bedrock 上的部分 OpenAI 模型(如 `openai.gpt-5.6-terra`、`openai.gpt-5.6-sol`)
+不在 Bedrock 基础模型目录中,只支持专用 `bedrock-mantle` 端点上的 OpenAI
+**Responses API**——常规 `bedrock` 提供商(Converse API)和 `openai` 提供商
+(Chat Completions)都无法调用它们。请使用 `bedrock-mantle` 提供商:
+
+```bash
+# 1. 安装提供商依赖(openai SDK >= 2.45 + AWS token 生成器)
+pip install -e ".[bedrock-mantle]"
+
+# 2. 配置——只需 AWS 凭证,不需要 OpenAI 账号或 API key
+export AI_PROVIDER=bedrock-mantle
+export AWS_REGION=us-east-1 # us-east-1 / us-east-2 /
us-west-2
+export OPENAI_MODEL='openai.gpt-5.6-terra' # 不设置时的默认值
+# export OPENAI_SMALL_FAST_MODEL='openai.gpt-5.6-terra'
+
+# 3. 正常生成
+seatunnel "把 MySQL 的 users 表同步到 S3,Parquet 格式"
+```
+
+提供商契约:
+
+- **端点**:`https://bedrock-mantle.{region}.api.aws/openai/v1`——这类模型
+ 要求的专属 `openai/v1` 路径(通用的 `v1` Responses 路径会拒绝这些模型)。
+- **认证**:通过 `aws-bedrock-token-generator` 从你的 AWS 凭证(profile、
+ 环境变量或 IAM 角色)自动派生短期 bearer token,每 30 分钟自动轮换,
+ 不在任何地方存储长期密钥。
+- **数据留存**:所有请求携带 `store=false`,Bedrock 不会在服务端留存你的
+ 提示词和生成的配置(服务默认行为是保留 30 天)。
+- **参数**:这类模型不接受 `temperature`,提供商不会发送该参数,配置的
+ temperature 值不会生效。
+- **错误处理**:截断(`incomplete`)、失败和拒答的响应会抛出显式错误,
+ 而不是伪装成正常结果返回。
+
+该提供商完整支持 CLI 内部的工具调用循环(规划阶段的连接器查询)和多轮
+会话,包括工具调用之间模型推理输出(reasoning)的保留与回放。
+
API 密钥只从环境变量读取——绝不写入任何配置文件。
## 生成第一条管道
diff --git a/seatunnel-cli/README.md b/seatunnel-cli/README.md
index dc285199be..de890b0964 100644
--- a/seatunnel-cli/README.md
+++ b/seatunnel-cli/README.md
@@ -111,6 +111,29 @@ Converse responses when Bedrock returns them.
> model for the rest of the session, so subsequent calls skip the parameter.
> For these models any configured temperature value is not applied.
+#### Option A2: AWS Bedrock — OpenAI-family models (bedrock-mantle)
+
+Some OpenAI models on Bedrock (e.g. `openai.gpt-5.6-terra`) are not in the
+foundation-model catalog and only support the OpenAI Responses API on the
+dedicated `bedrock-mantle` endpoint. Use the `bedrock-mantle` provider for
+these:
+
+```bash
+export AI_PROVIDER=bedrock-mantle
+export AWS_REGION=us-east-1
+export OPENAI_MODEL='openai.gpt-5.6-terra'
+
+# Requires: pip install -e ".[bedrock-mantle]"
+# Auth: a short-term bearer token is derived automatically from your AWS
+# credentials (aws-bedrock-token-generator) and refreshed every 30 minutes.
+```
+
+These models do not accept the `temperature` parameter; the provider omits it.
+
+All requests are sent with `store=false`, so Bedrock does not retain your
+prompts or responses server-side (the service default would otherwise keep
+them for 30 days).
+
#### Option B: Anthropic API
```bash
@@ -184,7 +207,7 @@ When the engine is running, the CLI operates in **cluster
mode** with live conne
| Variable | Required | Default | Description |
|----------|----------|---------|-------------|
-| `AI_PROVIDER` | No | `bedrock` | LLM provider: `bedrock`, `anthropic`, or
`openai` |
+| `AI_PROVIDER` | No | `bedrock` | LLM provider: `bedrock`, `bedrock-mantle`,
`anthropic`, or `openai` |
| `AWS_REGION` | Bedrock | `us-east-1` | AWS region for Bedrock |
| `ANTHROPIC_API_KEY` | Anthropic | -- | Anthropic API key |
| `OPENAI_API_KEY` | OpenAI | -- | OpenAI API key |
diff --git a/seatunnel-cli/env.example.sh b/seatunnel-cli/env.example.sh
old mode 100644
new mode 100755
index 180f2d9161..c684257df4
--- a/seatunnel-cli/env.example.sh
+++ b/seatunnel-cli/env.example.sh
@@ -26,6 +26,9 @@
# export AI_PROVIDER=anthropic # Option A
# export AI_PROVIDER=openai # Option B
# export AI_PROVIDER=bedrock # Option C
+# export AI_PROVIDER=bedrock-mantle # Option C2: OpenAI-family models on
Bedrock
+# # (GPT-5.6 Terra/Sol; needs
".[bedrock-mantle]" extra;
+# # model via OPENAI_MODEL, e.g.
openai.gpt-5.6-terra)
# ─── Option A: Anthropic API (AI_PROVIDER=anthropic) ───
# export ANTHROPIC_API_KEY=sk-ant-...
diff --git a/seatunnel-cli/pyproject.toml b/seatunnel-cli/pyproject.toml
index 51d8055538..e88f56efdb 100644
--- a/seatunnel-cli/pyproject.toml
+++ b/seatunnel-cli/pyproject.toml
@@ -39,8 +39,9 @@ seatunnel = "seatunnel_cli.cli:main"
bedrock = ["boto3>=1.34.0"]
anthropic = ["anthropic>=0.42.0"]
openai = ["openai>=1.0.0"]
-all = ["boto3>=1.34.0", "anthropic>=0.42.0", "openai>=1.0.0"]
-dev = ["pytest", "black", "ruff", "boto3>=1.34.0", "anthropic>=0.42.0",
"openai>=1.0.0"]
+bedrock-mantle = ["boto3>=1.34.0", "openai>=2.45.0",
"aws-bedrock-token-generator>=1.0.0"]
+all = ["boto3>=1.34.0", "anthropic>=0.42.0", "openai>=2.45.0",
"aws-bedrock-token-generator>=1.0.0"]
+dev = ["pytest", "black", "ruff", "boto3>=1.34.0", "anthropic>=0.42.0",
"openai>=2.45.0", "aws-bedrock-token-generator>=1.0.0"]
[tool.setuptools.packages.find]
include = ["seatunnel_cli*"]
diff --git a/seatunnel-cli/seatunnel_cli/cli.py
b/seatunnel-cli/seatunnel_cli/cli.py
index 6be82c7f70..61988744af 100644
--- a/seatunnel-cli/seatunnel_cli/cli.py
+++ b/seatunnel-cli/seatunnel_cli/cli.py
@@ -511,19 +511,23 @@ class SeaTunnelCLI:
console.print(" [bold]3[/bold]. bedrock — AWS Bedrock (Claude
via AWS)")
console.print(" Requires: AWS credentials (aws configure / env
vars / IAM role)")
console.print(" Docs: https://docs.aws.amazon.com/bedrock/\n")
+ console.print(" [bold]4[/bold]. bedrock-mantle — OpenAI-family
models on Bedrock (GPT-5.6 Terra/Sol)")
+ console.print(" Requires: AWS credentials + pip install
\".[bedrock-mantle]\"")
+ console.print(" Note: Responses-API-only models on the
bedrock-mantle endpoint\n")
try:
- choice = pt_prompt(" Enter your choice (1/2/3): ").strip().lower()
+ choice = pt_prompt(" Enter your choice (1/2/3/4):
").strip().lower()
except (EOFError, KeyboardInterrupt):
console.print("\n Setup cancelled.", style="warning")
return
- choice_map = {"1": "anthropic", "2": "openai", "3": "bedrock"}
+ choice_map = {"1": "anthropic", "2": "openai", "3": "bedrock",
+ "4": "bedrock-mantle"}
choice = choice_map.get(choice, choice)
if not choice or choice not in _PROVIDERS:
console.print(
- f" [error]Invalid choice: '{choice}'. Please enter 1, 2, or
3.[/error]"
+ f" [error]Invalid choice: '{choice}'. Please enter 1, 2, 3,
or 4.[/error]"
)
return
@@ -599,8 +603,14 @@ class SeaTunnelCLI:
config.setdefault("settings", {})["openai_base_url"] =
base_url
console.print(f" Base URL set: [bold]{base_url}[/bold]")
- elif choice == "bedrock":
+ elif choice in ("bedrock", "bedrock-mantle"):
console.print(" AWS Bedrock requires AWS credentials.\n")
+ if choice == "bedrock-mantle":
+ console.print(
+ " [dim]bedrock-mantle also requires the optional
extra:[/dim]\n"
+ " pip install \".[bedrock-mantle]\" "
+ "[dim](openai SDK + aws-bedrock-token-generator)[/dim]\n"
+ )
console.print(
" [dim]Options:[/dim]\n"
" - aws configure (interactive setup)\n"
@@ -655,6 +665,11 @@ class SeaTunnelCLI:
default_fast = "gpt-4o-mini"
model_env = "OPENAI_MODEL"
fast_env = "OPENAI_SMALL_FAST_MODEL"
+ elif choice == "bedrock-mantle":
+ default_model = "openai.gpt-5.6-terra"
+ default_fast = "openai.gpt-5.6-terra"
+ model_env = "OPENAI_MODEL"
+ fast_env = "OPENAI_SMALL_FAST_MODEL"
else: # bedrock
default_model = "us.anthropic.claude-sonnet-4-20250514-v1:0"
default_fast = "us.anthropic.claude-haiku-4-5-20251001-v1:0"
@@ -1512,7 +1527,7 @@ def main():
)
parser.add_argument(
"--provider",
- choices=["bedrock", "anthropic", "openai"],
+ choices=["bedrock", "bedrock-mantle", "anthropic", "openai"],
help="LLM provider (overrides AI_PROVIDER env var and config.json)",
)
parser.add_argument(
@@ -1563,15 +1578,18 @@ def main():
# Override provider if specified via CLI flags
if args.provider:
os.environ["AI_PROVIDER"] = args.provider
+ # Providers speaking the OpenAI protocol read OPENAI_MODEL*;
+ # bedrock/anthropic read ANTHROPIC_MODEL*.
+ _OPENAI_FAMILY = ("openai", "bedrock-mantle")
if args.model:
provider = os.environ.get("AI_PROVIDER", "").lower()
- if provider == "openai":
+ if provider in _OPENAI_FAMILY:
os.environ["OPENAI_MODEL"] = args.model
else:
os.environ["ANTHROPIC_MODEL"] = args.model
if args.fast_model:
provider = os.environ.get("AI_PROVIDER", "").lower()
- if provider == "openai":
+ if provider in _OPENAI_FAMILY:
os.environ["OPENAI_SMALL_FAST_MODEL"] = args.fast_model
else:
os.environ["ANTHROPIC_SMALL_FAST_MODEL"] = args.fast_model
diff --git a/seatunnel-cli/seatunnel_cli/llm_provider.py
b/seatunnel-cli/seatunnel_cli/llm_provider.py
index 8d547666bc..29a5c13384 100644
--- a/seatunnel-cli/seatunnel_cli/llm_provider.py
+++ b/seatunnel-cli/seatunnel_cli/llm_provider.py
@@ -277,6 +277,10 @@ class LLMProvider(abc.ABC):
reasoning_text["text"] += event.get("text", "")
if event.get("signature"):
reasoning_text["signature"] += event["signature"]
+ elif etype == "mantle_responses_item":
+ flush_text()
+ flush_model_state()
+ content.append({"mantleResponsesItem": event.get("item", {})})
elif etype == "tool_start":
flush_text()
flush_model_state()
@@ -947,6 +951,305 @@ class OpenAIProvider(LLMProvider):
return None
+# ─── Bedrock Mantle Provider (OpenAI Responses API) ───
+
+class BedrockMantleProvider(LLMProvider):
+ """AWS Bedrock 'mantle' endpoint provider for OpenAI-family models.
+
+ Some OpenAI models on Bedrock (e.g. openai.gpt-5.6-terra) are NOT in the
+ foundation-model catalog and reject Converse/InvokeModel/ChatCompletions.
+ They are served only via the OpenAI Responses API on the bedrock-mantle
+ endpoint: https://bedrock-mantle.{region}.api.aws/openai/v1
+
+ Auth uses a short-term bearer token derived from the caller's SigV4
+ credentials (aws-bedrock-token-generator); tokens are refreshed
+ automatically. These models do not accept `temperature` — it is omitted.
+ """
+
+ TOKEN_TTL_SECONDS = 1800 # regenerate well within the 12h token validity
+
+ def __init__(self):
+ try:
+ import openai # noqa: F401
+ from aws_bedrock_token_generator import provide_token # noqa: F401
+ except ImportError as e:
+ raise ImportError(
+ "bedrock-mantle provider requires: "
+ "pip install openai aws-bedrock-token-generator"
+ ) from e
+
+ self._region = os.environ.get(
+ "AWS_REGION", os.environ.get("AWS_DEFAULT_REGION", "us-east-1"))
+ self._model_id = os.environ.get("OPENAI_MODEL", "openai.gpt-5.6-terra")
+ self._fast_model_id = os.environ.get("OPENAI_SMALL_FAST_MODEL",
self._model_id)
+ self._client = None
+ self._token_born = 0.0
+
+ def _get_client(self):
+ """Return an OpenAI client for the bedrock-mantle endpoint.
+
+ Auth uses a short-term bearer token derived from the caller's AWS
+ credentials; the client (and token) is rebuilt after TOKEN_TTL_SECONDS.
+ The base URL must be the model-specific ``openai/v1`` path — the
+ generic ``v1`` path rejects these models (see the GPT-5.6 model card).
+ """
+ import time as _time
+ import openai
+ from aws_bedrock_token_generator import provide_token
+ if self._client is None or _time.monotonic() - self._token_born >
self.TOKEN_TTL_SECONDS:
+ token = provide_token(region=self._region)
+ self._client = openai.OpenAI(
+ api_key=token,
+
base_url=f"https://bedrock-mantle.{self._region}.api.aws/openai/v1",
+ )
+ self._token_born = _time.monotonic()
+ return self._client
+
+ @property
+ def provider_name(self) -> str:
+ return "bedrock-mantle"
+
+ @property
+ def model_id(self) -> str:
+ return self._model_id
+
+ @property
+ def fast_model_id(self) -> str:
+ return self._fast_model_id
+
+ # ── format conversion ──
+
+ @staticmethod
+ def _to_responses_input(messages: list[dict]) -> list[dict]:
+ """Convert internal (Converse-shaped) messages to Responses API
items."""
+ items: list[dict] = []
+ for msg in messages:
+ role = msg["role"]
+ for block in msg.get("content", []):
+ if "text" in block:
+ items.append({"role": role, "content": block["text"]})
+ elif "toolUse" in block:
+ tu = block["toolUse"]
+ items.append({
+ "type": "function_call",
+ "call_id": tu["toolUseId"],
+ "name": tu["name"],
+ "arguments": json.dumps(tu.get("input", {})),
+ })
+ elif "toolResult" in block:
+ tr = block["toolResult"]
+ text_parts = [c["text"] for c in tr.get("content", []) if
"text" in c]
+ items.append({
+ "type": "function_call_output",
+ "call_id": tr["toolUseId"],
+ "output": "\n".join(text_parts),
+ })
+ elif "mantleResponsesItem" in block:
+ # Opaque Responses API output item (e.g. reasoning)
+ # captured verbatim from a previous turn; replayed
+ # in-order so multi-step tool loops keep their context,
+ # as the AWS tool-calling guide requires.
+ item = block["mantleResponsesItem"]
+ if item:
+ items.append(item)
+ return items
+
+ @staticmethod
+ def _to_responses_tools(tools: list[dict]) -> list[dict]:
+ result = []
+ for tool in tools:
+ spec = tool.get("toolSpec", {})
+ result.append({
+ "type": "function",
+ "name": spec["name"],
+ "description": spec.get("description", ""),
+ "parameters": spec.get("inputSchema", {}).get("json", {}),
+ })
+ return result
+
+ @staticmethod
+ def _dump_item(item) -> dict:
+ """Serialize a Responses output item to a replayable plain dict."""
+ try:
+ return item.model_dump(exclude_none=True)
+ except AttributeError:
+ return dict(item) if isinstance(item, dict) else {}
+
+ @staticmethod
+ def _raise_on_terminal_status(response) -> None:
+ """Fail loudly on non-completed terminal states.
+
+ Bedrock Mantle reports truncation/filtering/model errors in-band via
+ ``status`` + ``incomplete_details``/``error`` rather than transport
+ errors; treating those as success would hand truncated configs to
+ the agent loop as if they were complete.
+ """
+ status = getattr(response, "status", None)
+ if status in (None, "completed"):
+ return
+ detail = ""
+ incomplete = getattr(response, "incomplete_details", None)
+ if incomplete is not None:
+ detail = f" ({getattr(incomplete, 'reason', '') or incomplete})"
+ error = getattr(response, "error", None)
+ if error is not None:
+ detail += f" error: {getattr(error, 'message', '') or error}"
+ raise RuntimeError(
+ f"Bedrock Mantle response ended with status '{status}'{detail}")
+
+ def chat(
+ self,
+ messages: list[dict],
+ system: str = "",
+ model: str | None = None,
+ temperature: float = 0.3,
+ max_tokens: int = 4096,
+ tools: list[dict] | None = None,
+ ) -> dict:
+ """Send a non-streaming Responses API request.
+
+ Requests are sent with ``store=False`` so Bedrock does not retain
+ prompt/response data server-side (its default is 30-day retention).
+ Reasoning output items are preserved verbatim in the returned
+ history so subsequent turns can replay them. Non-``completed``
+ terminal statuses and refusals raise ``RuntimeError`` instead of
+ being returned as an apparently-successful message.
+ """
+ client = self._get_client()
+ kwargs = {
+ "model": model or self._model_id,
+ "input": self._to_responses_input(messages),
+ "max_output_tokens": max_tokens,
+ "store": False,
+ }
+ if system:
+ kwargs["instructions"] = system
+ if tools:
+ kwargs["tools"] = self._to_responses_tools(tools)
+
+ response = client.responses.create(**kwargs)
+ self._raise_on_terminal_status(response)
+
+ content: list[dict] = []
+ has_tool_use = False
+ for item in response.output:
+ itype = getattr(item, "type", "")
+ if itype == "message":
+ for part in getattr(item, "content", []) or []:
+ if getattr(part, "type", "") == "refusal":
+ refusal = getattr(part, "refusal", "") or ""
+ raise RuntimeError(
+ f"Bedrock Mantle model refused the request:
{refusal}")
+ text = getattr(part, "text", None)
+ if text:
+ content.append({"text": text})
+ elif itype == "function_call":
+ has_tool_use = True
+ try:
+ parsed = json.loads(item.arguments or "{}")
+ except (json.JSONDecodeError, TypeError):
+ parsed = {}
+ content.append({
+ "toolUse": {
+ "toolUseId": item.call_id,
+ "name": item.name,
+ "input": parsed,
+ }
+ })
+ elif itype == "reasoning":
+ content.append({
+ "mantleResponsesItem": self._dump_item(item),
+ })
+
+ return {
+ "output": {"message": {"role": "assistant", "content": content}},
+ "stopReason": "tool_use" if has_tool_use else "end_turn",
+ }
+
+ def chat_stream(
+ self,
+ messages: list[dict],
+ system: str = "",
+ model: str | None = None,
+ temperature: float = 0.3,
+ max_tokens: int = 4096,
+ tools: list[dict] | None = None,
+ ) -> Generator[dict, None, None]:
+ """Stream a Responses API request as internal events.
+
+ Same contracts as :meth:`chat`: ``store=False`` on every request,
+ reasoning output items forwarded for replay, and terminal
+ failed/incomplete/error events raised as ``RuntimeError`` rather
+ than silently ending the stream (collect_stream would otherwise
+ default the missing stop to a successful ``end_turn``).
+ """
+ client = self._get_client()
+ kwargs = {
+ "model": model or self._model_id,
+ "input": self._to_responses_input(messages),
+ "max_output_tokens": max_tokens,
+ "store": False,
+ "stream": True,
+ }
+ if system:
+ kwargs["instructions"] = system
+ if tools:
+ kwargs["tools"] = self._to_responses_tools(tools)
+
+ stream = client.responses.create(**kwargs)
+ current_tool_id = None
+ saw_tool_use = False
+ for event in stream:
+ etype = getattr(event, "type", "")
+ if etype == "response.output_text.delta":
+ yield {"type": "text_delta", "text": getattr(event, "delta",
"")}
+ elif etype == "response.output_item.added":
+ item = getattr(event, "item", None)
+ if item is not None and getattr(item, "type", "") ==
"function_call":
+ saw_tool_use = True
+ current_tool_id = getattr(item, "call_id", "") or ""
+ yield {
+ "type": "tool_start",
+ "tool_use_id": current_tool_id,
+ "name": getattr(item, "name", "") or "",
+ }
+ elif etype == "response.function_call_arguments.delta":
+ yield {
+ "type": "tool_input_delta",
+ "tool_use_id": current_tool_id or "",
+ "delta": getattr(event, "delta", ""),
+ }
+ elif etype == "response.output_item.done":
+ item = getattr(event, "item", None)
+ item_type = getattr(item, "type", "") if item is not None else
""
+ if item_type == "function_call" and current_tool_id:
+ yield {"type": "tool_stop", "tool_use_id": current_tool_id}
+ current_tool_id = None
+ elif item_type == "reasoning":
+ yield {
+ "type": "mantle_responses_item",
+ "item": self._dump_item(item),
+ }
+ elif etype == "response.refusal.done":
+ refusal = getattr(event, "refusal", "") or ""
+ raise RuntimeError(
+ f"Bedrock Mantle model refused the request: {refusal}")
+ elif etype in ("response.failed", "response.incomplete",
+ "response.error", "error"):
+ resp = getattr(event, "response", None)
+ self._raise_on_terminal_status(resp) if resp is not None \
+ else None
+ message = getattr(event, "message", "") or etype
+ raise RuntimeError(
+ f"Bedrock Mantle stream ended abnormally: {message}")
+ elif etype == "response.completed":
+ self._raise_on_terminal_status(getattr(event, "response",
None))
+ yield {
+ "type": "message_stop",
+ "stop_reason": "tool_use" if saw_tool_use else "end_turn",
+ }
+
+
# ─── Config file ───
@@ -1006,6 +1309,7 @@ def _auto_detect_provider() -> str | None:
_PROVIDERS = {
"bedrock": BedrockProvider,
+ "bedrock-mantle": BedrockMantleProvider,
"anthropic": AnthropicProvider,
"openai": OpenAIProvider,
}
@@ -1051,12 +1355,12 @@ def create_provider(provider: str | None = None) ->
LLMProvider:
if name and "models" in config:
model_config = config["models"].get(name, {})
if model_config.get("model") and not os.environ.get("ANTHROPIC_MODEL")
and not os.environ.get("OPENAI_MODEL"):
- if name == "openai":
+ if name in ("openai", "bedrock-mantle"):
os.environ.setdefault("OPENAI_MODEL", model_config["model"])
else:
os.environ.setdefault("ANTHROPIC_MODEL", model_config["model"])
if model_config.get("fast_model") and not
os.environ.get("ANTHROPIC_SMALL_FAST_MODEL") and not
os.environ.get("OPENAI_SMALL_FAST_MODEL"):
- if name == "openai":
+ if name in ("openai", "bedrock-mantle"):
os.environ.setdefault("OPENAI_SMALL_FAST_MODEL",
model_config["fast_model"])
else:
os.environ.setdefault("ANTHROPIC_SMALL_FAST_MODEL",
model_config["fast_model"])
diff --git a/seatunnel-cli/tests/test_cli_provider_routing.py
b/seatunnel-cli/tests/test_cli_provider_routing.py
new file mode 100644
index 0000000000..39531fdbbc
--- /dev/null
+++ b/seatunnel-cli/tests/test_cli_provider_routing.py
@@ -0,0 +1,85 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+"""Regression tests for CLI --provider/--model routing of the
+bedrock-mantle provider family (issue: mantle models were routed to
+ANTHROPIC_MODEL and the provider was missing from argparse choices)."""
+
+import os
+import sys
+from unittest import mock
+
+import pytest
+
+
+def _run_main_until_provider(argv):
+ """Run cli.main() with argv, stopping right after env routing."""
+ from seatunnel_cli import cli
+
+ captured = {}
+
+ class _Stop(Exception):
+ pass
+
+ def fake_console(*a, **k):
+ # capture env state at the point the CLI would build the console
+ captured["AI_PROVIDER"] = os.environ.get("AI_PROVIDER")
+ captured["OPENAI_MODEL"] = os.environ.get("OPENAI_MODEL")
+ captured["ANTHROPIC_MODEL"] = os.environ.get("ANTHROPIC_MODEL")
+ raise _Stop()
+
+ with mock.patch.object(sys, "argv", ["seatunnel"] + argv), \
+ mock.patch.object(cli, "Console", side_effect=fake_console), \
+ pytest.raises(_Stop):
+ cli.main()
+ return captured
+
+
[email protected](autouse=True)
+def _clean_env():
+ saved = {k: os.environ.pop(k, None) for k in
+ ("AI_PROVIDER", "OPENAI_MODEL", "ANTHROPIC_MODEL",
+ "OPENAI_SMALL_FAST_MODEL", "ANTHROPIC_SMALL_FAST_MODEL")}
+ yield
+ for k, v in saved.items():
+ if v is None:
+ os.environ.pop(k, None)
+ else:
+ os.environ[k] = v
+
+
+def test_bedrock_mantle_accepted_by_argparse_and_routes_openai_model():
+ captured = _run_main_until_provider(
+ ["--provider", "bedrock-mantle", "--model", "openai.gpt-5.6-sol",
"hi"])
+ assert captured["AI_PROVIDER"] == "bedrock-mantle"
+ assert captured["OPENAI_MODEL"] == "openai.gpt-5.6-sol"
+ assert captured["ANTHROPIC_MODEL"] is None
+
+
+def test_bedrock_still_routes_anthropic_model():
+ captured = _run_main_until_provider(
+ ["--provider", "bedrock", "--model", "us.anthropic.claude-sonnet-5",
"hi"])
+ assert captured["ANTHROPIC_MODEL"] == "us.anthropic.claude-sonnet-5"
+ assert captured["OPENAI_MODEL"] is None
+
+
+def test_unknown_provider_rejected():
+ from seatunnel_cli import cli
+ with mock.patch.object(sys, "argv",
+ ["seatunnel", "--provider", "nonsense", "hi"]), \
+ pytest.raises(SystemExit):
+ cli.main()
diff --git a/seatunnel-cli/tests/test_llm_provider_bedrock_mantle.py
b/seatunnel-cli/tests/test_llm_provider_bedrock_mantle.py
new file mode 100644
index 0000000000..2b67c05b95
--- /dev/null
+++ b/seatunnel-cli/tests/test_llm_provider_bedrock_mantle.py
@@ -0,0 +1,388 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements. See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License. You may obtain a copy of the License at
+#
+# http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+"""Tests for BedrockMantleProvider message/tool format conversion and the
+chat/chat_stream response mapping (Responses API <-> internal format)."""
+
+from types import SimpleNamespace
+from unittest import mock
+
+import pytest
+
+from seatunnel_cli.llm_provider import BedrockMantleProvider
+
+
+def _make_provider():
+ provider = BedrockMantleProvider.__new__(BedrockMantleProvider)
+ provider._region = "us-east-1"
+ provider._model_id = "openai.gpt-5.6-terra"
+ provider._fast_model_id = "openai.gpt-5.6-terra"
+ provider._client = mock.MagicMock()
+ provider._token_born = float("inf") # never refresh in tests
+ return provider
+
+
+# ── input conversion: internal (Converse-shaped) -> Responses API items ──
+
+def test_to_responses_input_text_and_tool_roundtrip():
+ messages = [
+ {"role": "user", "content": [{"text": "hi"}]},
+ {"role": "assistant", "content": [
+ {"text": "let me check"},
+ {"toolUse": {"toolUseId": "call_1", "name": "calc",
+ "input": {"expr": "6*7"}}},
+ ]},
+ {"role": "user", "content": [
+ {"toolResult": {"toolUseId": "call_1",
+ "content": [{"text": "42"}]}},
+ ]},
+ ]
+ items = BedrockMantleProvider._to_responses_input(messages)
+ assert items[0] == {"role": "user", "content": "hi"}
+ assert items[1] == {"role": "assistant", "content": "let me check"}
+ assert items[2]["type"] == "function_call"
+ assert items[2]["call_id"] == "call_1"
+ assert items[2]["name"] == "calc"
+ assert items[3] == {"type": "function_call_output",
+ "call_id": "call_1", "output": "42"}
+
+
+def test_to_responses_tools_conversion():
+ tools = [{"toolSpec": {
+ "name": "calc", "description": "calculate",
+ "inputSchema": {"json": {"type": "object",
+ "properties": {"expr": {"type": "string"}}}},
+ }}]
+ converted = BedrockMantleProvider._to_responses_tools(tools)
+ assert converted == [{
+ "type": "function", "name": "calc", "description": "calculate",
+ "parameters": {"type": "object",
+ "properties": {"expr": {"type": "string"}}},
+ }]
+
+
+# ── chat: Responses output -> internal format ──
+
+def _resp(output_items):
+ return SimpleNamespace(output=output_items)
+
+
+def _message_item(text):
+ return SimpleNamespace(type="message",
+ content=[SimpleNamespace(text=text)])
+
+
+def _function_call_item(call_id, name, arguments):
+ return SimpleNamespace(type="function_call", call_id=call_id,
+ name=name, arguments=arguments)
+
+
+def test_chat_maps_text_response():
+ provider = _make_provider()
+ provider._client.responses.create.return_value = _resp(
+ [_message_item("OK")])
+ resp = provider.chat([{"role": "user", "content": [{"text": "hi"}]}])
+ assert resp["stopReason"] == "end_turn"
+ assert resp["output"]["message"]["content"] == [{"text": "OK"}]
+ # temperature must never be sent (unsupported by these models)
+ kwargs = provider._client.responses.create.call_args.kwargs
+ assert "temperature" not in kwargs
+
+
+def test_chat_maps_tool_call_and_bad_json_input():
+ provider = _make_provider()
+ provider._client.responses.create.return_value = _resp([
+ _function_call_item("call_9", "calc", '{"expr": "6*7"}'),
+ _function_call_item("call_x", "calc", "NOT JSON"),
+ ])
+ resp = provider.chat([{"role": "user", "content": [{"text": "hi"}]}])
+ assert resp["stopReason"] == "tool_use"
+ tool_uses = [b["toolUse"] for b in resp["output"]["message"]["content"]]
+ assert tool_uses[0] == {"toolUseId": "call_9", "name": "calc",
+ "input": {"expr": "6*7"}}
+ assert tool_uses[1]["input"] == {} # malformed arguments degrade to {}
+
+
+def test_chat_passes_system_as_instructions():
+ provider = _make_provider()
+ provider._client.responses.create.return_value = _resp(
+ [_message_item("OK")])
+ provider.chat([{"role": "user", "content": [{"text": "hi"}]}],
+ system="be brief")
+ kwargs = provider._client.responses.create.call_args.kwargs
+ assert kwargs["instructions"] == "be brief"
+
+
+# ── chat_stream: Responses stream events -> internal events ──
+
+def _ev(type_, **attrs):
+ return SimpleNamespace(type=type_, **attrs)
+
+
+def test_chat_stream_maps_text_and_tool_events():
+ provider = _make_provider()
+ fc_item = SimpleNamespace(type="function_call", call_id="call_1",
+ name="calc")
+ provider._client.responses.create.return_value = iter([
+ _ev("response.output_text.delta", delta="hel"),
+ _ev("response.output_text.delta", delta="lo"),
+ _ev("response.output_item.added", item=fc_item),
+ _ev("response.function_call_arguments.delta", delta='{"expr":'),
+ _ev("response.function_call_arguments.delta", delta='"6*7"}'),
+ _ev("response.output_item.done", item=fc_item),
+ _ev("response.completed"),
+ ])
+ events = list(provider.chat_stream(
+ [{"role": "user", "content": [{"text": "hi"}]}]))
+ types = [e["type"] for e in events]
+ assert types == ["text_delta", "text_delta", "tool_start",
+ "tool_input_delta", "tool_input_delta", "tool_stop",
+ "message_stop"]
+ assert events[2] == {"type": "tool_start", "tool_use_id": "call_1",
+ "name": "calc"}
+ assert events[-1]["stop_reason"] == "tool_use"
+
+
+def test_chat_stream_plain_text_ends_with_end_turn():
+ provider = _make_provider()
+ provider._client.responses.create.return_value = iter([
+ _ev("response.output_text.delta", delta="OK"),
+ _ev("response.completed"),
+ ])
+ events = list(provider.chat_stream(
+ [{"role": "user", "content": [{"text": "hi"}]}]))
+ assert events[-1] == {"type": "message_stop", "stop_reason": "end_turn"}
+
+
+# ── stream events integrate with LLMProvider.collect_stream ──
+
+def test_stream_events_collect_to_internal_response():
+ provider = _make_provider()
+ fc_item = SimpleNamespace(type="function_call", call_id="call_1",
+ name="calc")
+ provider._client.responses.create.return_value = iter([
+ _ev("response.output_item.added", item=fc_item),
+ _ev("response.function_call_arguments.delta", delta='{"expr": "1"}'),
+ _ev("response.output_item.done", item=fc_item),
+ _ev("response.completed"),
+ ])
+ events = list(provider.chat_stream(
+ [{"role": "user", "content": [{"text": "hi"}]}]))
+ resp = BedrockMantleProvider.collect_stream(events)
+ assert resp["stopReason"] == "tool_use"
+ tool_uses = BedrockMantleProvider.extract_tool_use(resp)
+ assert tool_uses == [{"toolUseId": "call_1", "name": "calc",
+ "input": {"expr": "1"}}]
+
+
+# ── endpoint construction ──
+
+def test_client_uses_documented_mantle_openai_path():
+ """These models are served on the `openai/v1` path of the bedrock-mantle
+ endpoint — NOT the generic `v1` path used by other Responses-API models.
+ See the AWS model card for gpt-5.6-terra ("available on the
+ openai/v1/responses path ... different from the v1/responses path").
+ Empirically, the generic /v1 path rejects these models with
+ "does not support the '/v1/responses' API"."""
+ provider = BedrockMantleProvider.__new__(BedrockMantleProvider)
+ provider._region = "eu-west-3"
+ provider._model_id = "m"
+ provider._fast_model_id = "m"
+ provider._client = None
+ provider._token_born = 0.0
+
+ fake_openai = mock.MagicMock()
+ fake_generator = mock.MagicMock()
+ fake_generator.provide_token.return_value = "tok"
+ with mock.patch.dict("sys.modules", {
+ "openai": fake_openai,
+ "aws_bedrock_token_generator": fake_generator,
+ }):
+ provider._get_client()
+ kwargs = fake_openai.OpenAI.call_args.kwargs
+ assert kwargs["base_url"] == \
+ "https://bedrock-mantle.eu-west-3.api.aws/openai/v1"
+ assert kwargs["api_key"] == "tok"
+
+
+# ── token refresh ──
+
+def test_client_refreshes_after_ttl():
+ provider = BedrockMantleProvider.__new__(BedrockMantleProvider)
+ provider._region = "us-east-1"
+ provider._model_id = "m"
+ provider._fast_model_id = "m"
+ provider._client = None
+ provider._token_born = 0.0
+
+ fake_openai = mock.MagicMock()
+ fake_generator = mock.MagicMock()
+ fake_generator.provide_token.return_value = "tok"
+ with mock.patch.dict("sys.modules", {
+ "openai": fake_openai,
+ "aws_bedrock_token_generator": fake_generator,
+ }):
+ provider._get_client()
+ assert fake_generator.provide_token.called
+ first_client = provider._client
+ # within TTL: same client reused
+ provider._get_client()
+ assert provider._client is first_client
+ # expire TTL: new token requested
+ provider._token_born = -10_000.0
+ provider._get_client()
+ assert fake_generator.provide_token.call_count == 2
+
+# ── store=False on every request (Bedrock retains data by default) ──
+
+def test_chat_sends_store_false():
+ provider = _make_provider()
+ provider._client.responses.create.return_value = _resp(
+ [_message_item("OK")])
+ provider.chat([{"role": "user", "content": [{"text": "hi"}]}])
+ assert provider._client.responses.create.call_args.kwargs["store"] is False
+
+
+def test_chat_stream_sends_store_false():
+ provider = _make_provider()
+ provider._client.responses.create.return_value = iter([
+ _ev("response.completed"),
+ ])
+ list(provider.chat_stream([{"role": "user", "content": [{"text": "hi"}]}]))
+ assert provider._client.responses.create.call_args.kwargs["store"] is False
+
+
+# ── reasoning items are preserved and replayed in order ──
+
+def _reasoning_item():
+ item = mock.MagicMock()
+ item.type = "reasoning"
+ item.model_dump.return_value = {
+ "type": "reasoning", "id": "rs_1",
+ "summary": [], "content": None,
+ }
+ return item
+
+
+def test_chat_preserves_reasoning_for_replay():
+ provider = _make_provider()
+ provider._client.responses.create.return_value = _resp([
+ _reasoning_item(),
+ _function_call_item("call_1", "calc", '{"expr": "1"}'),
+ ])
+ resp = provider.chat([{"role": "user", "content": [{"text": "hi"}]}])
+ blocks = resp["output"]["message"]["content"]
+ assert blocks[0] == {"mantleResponsesItem": {
+ "type": "reasoning", "id": "rs_1", "summary": [], "content": None}}
+ assert "toolUse" in blocks[1]
+
+ # replay: history containing the reasoning block converts back verbatim,
+ # in order, before the function_call item
+ history = [
+ {"role": "user", "content": [{"text": "hi"}]},
+ resp["output"]["message"],
+ {"role": "user", "content": [{"toolResult": {
+ "toolUseId": "call_1", "content": [{"text": "1"}]}}]},
+ ]
+ items = BedrockMantleProvider._to_responses_input(history)
+ assert items[1] == {"type": "reasoning", "id": "rs_1",
+ "summary": [], "content": None}
+ assert items[2]["type"] == "function_call"
+ assert items[3]["type"] == "function_call_output"
+
+
+def test_chat_stream_forwards_reasoning_items():
+ provider = _make_provider()
+ r_item = mock.MagicMock()
+ r_item.type = "reasoning"
+ r_item.model_dump.return_value = {"type": "reasoning", "id": "rs_2"}
+ provider._client.responses.create.return_value = iter([
+ _ev("response.output_item.added", item=r_item),
+ _ev("response.output_item.done", item=r_item),
+ _ev("response.completed"),
+ ])
+ events = list(provider.chat_stream(
+ [{"role": "user", "content": [{"text": "hi"}]}]))
+ assert {"type": "mantle_responses_item",
+ "item": {"type": "reasoning", "id": "rs_2"}} in events
+ # and collect_stream lands it in history
+ resp = BedrockMantleProvider.collect_stream(events)
+ assert {"mantleResponsesItem": {"type": "reasoning", "id": "rs_2"}} \
+ in resp["output"]["message"]["content"]
+
+
+# ── terminal states must not become successful end_turns ──
+
+def test_chat_incomplete_status_raises():
+ provider = _make_provider()
+ response = _resp([_message_item("truncated par")])
+ response.status = "incomplete"
+ response.incomplete_details = mock.MagicMock(reason="max_output_tokens")
+ response.error = None
+ provider._client.responses.create.return_value = response
+ with pytest.raises(RuntimeError, match="incomplete.*max_output_tokens"):
+ provider.chat([{"role": "user", "content": [{"text": "hi"}]}])
+
+
+def test_chat_refusal_raises():
+ provider = _make_provider()
+ refusal_part = mock.MagicMock()
+ refusal_part.type = "refusal"
+ refusal_part.refusal = "cannot help with that"
+ msg = mock.MagicMock()
+ msg.type = "message"
+ msg.content = [refusal_part]
+ provider._client.responses.create.return_value = _resp([msg])
+ with pytest.raises(RuntimeError, match="refused"):
+ provider.chat([{"role": "user", "content": [{"text": "hi"}]}])
+
+
+def test_chat_stream_failed_event_raises():
+ provider = _make_provider()
+ provider._client.responses.create.return_value = iter([
+ _ev("response.output_text.delta", delta="par"),
+ _ev("response.failed", response=None, message="internal model error"),
+ ])
+ with pytest.raises(RuntimeError, match="abnormally"):
+ list(provider.chat_stream(
+ [{"role": "user", "content": [{"text": "hi"}]}]))
+
+# ── dependency contract: every extra bundling this provider must satisfy it ──
+
+def test_extras_bundling_mantle_share_responses_capable_floor():
+ """The provider calls client.responses.create, which needs a modern
+ openai SDK. Any extra that installs this provider (bedrock-mantle, all,
+ dev) must therefore pin the same floor — a lower one would install a
+ broken provider (review finding)."""
+ import re
+ from pathlib import Path
+ pyproject = Path(__file__).parent.parent / "pyproject.toml"
+ extras = {}
+ for line in pyproject.read_text().splitlines():
+ m = re.match(r'^([\w-]+)\s*=\s*\[(.*)\]', line.strip())
+ if m:
+ extras[m.group(1)] = m.group(2)
+ for extra in ("bedrock-mantle", "all", "dev"):
+ assert extra in extras, f"extra '{extra}' missing"
+ deps = extras[extra]
+ m = re.search(r'openai>=([\d.]+)', deps)
+ assert m, f"extra '{extra}' has no openai floor"
+ version = tuple(int(x) for x in m.group(1).split("."))
+ assert version >= (2, 45, 0), \
+ f"extra '{extra}' allows openai {m.group(1)} < 2.45.0"
+ assert "aws-bedrock-token-generator" in deps, \
+ f"extra '{extra}' missing aws-bedrock-token-generator"