Skip to content

Commit fb7a370

Browse files
Add hydrate and persist endpoints. (vllm-project#216)
## Summary Exposes the two halves of a stateful Responses turn as separate endpoints, so an external orchestrator (the llm-d coordinator) can make the inference call itself: - `POST /v1alpha/responses/hydrate`: expands `previous_response_id` into a stateless upstream request, plus a sealed context. - `POST /v1alpha/responses/persist`: takes that context and the model's response (a JSON body, or the SSE frames a streaming caller relayed), stores the turn, and returns the envelope carrying the stored `resp_` id. Both compose existing core operations (`rehydrate_conversation`, `upstream_request`, `decode_upstream`, `commit`), so they add no parsing, storage, or request building of their own. `decode_upstream` is the OpenAI JSON/SSE adapter; `commit` owns the shared validation and persistence. The in-process path is unchanged. The endpoints ship as a new crate and binary, `agentic-llm-d`, depending on `agentic-server-core` alone. The wire protocol (`SplitContext`, the sealing, the splittability check) lives there, so core stays unaware of a specific consumer. The context is HMAC-signed with an expiry and an audience, and the routes require a shared workload token in `x-agentic-workload-token`, leaving `Authorization` free for end-user identity. ## Test Plan `cargo test` - new `split_execution_integration.rs` covers the two-turn replay, the error paths, and the context round-trip. Existing suites pass unchanged. Verified on a cluster with a multi-turn conversation. Closes vllm-project#215. --------- Signed-off-by: Mohammad <mohammad.nassar@ibm.com> Co-authored-by: Francisco Javier Arceo <arceofrancisco@gmail.com>
1 parent 98325a8 commit fb7a370

20 files changed

Lines changed: 1152 additions & 43 deletions

File tree

‎Cargo.lock‎

Lines changed: 18 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -33,7 +33,7 @@ reqwest = { version = "0.12", default-features = false }
3333
rmcp-reqwest = { package = "reqwest", version = "0.13.2", default-features = false, features = ["json", "stream", "rustls"] }
3434
rmcp = { version = "1.8", default-features = false }
3535
serde = { version = "1", features = ["derive"] }
36-
serde_json = "1"
36+
serde_json = { version = "1", features = ["raw_value"] }
3737
thiserror = "2"
3838
tokio = { version = "1", features = ["full"] }
3939
tokio-util = "0.7"

‎Dockerfile‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -34,9 +34,11 @@ COPY crates ./crates
3434
RUN cargo clean \
3535
-p agentic-server-core \
3636
-p agentic-server \
37-
-p agentic-praxis && \
38-
cargo build --locked --release -p agentic-server && \
39-
install -Dm755 -s target/release/agentic-server /out/agentic-server
37+
-p agentic-praxis \
38+
-p agentic-llm-d && \
39+
cargo build --locked --release -p agentic-server -p agentic-llm-d && \
40+
install -Dm755 -s target/release/agentic-server /out/agentic-server && \
41+
install -Dm755 -s target/release/agentic-llm-d /out/agentic-llm-d
4042

4143
FROM debian:${DEBIAN_VERSION}-slim@${DEBIAN_IMAGE_DIGEST} AS runtime
4244

@@ -51,6 +53,7 @@ RUN apt-get update && \
5153
chmod g=u,g+s /var/lib/agentic-api
5254

5355
COPY --from=rust-build /out/agentic-server /usr/local/bin/agentic-server
56+
COPY --from=rust-build /out/agentic-llm-d /usr/local/bin/agentic-llm-d
5457
COPY --chmod=0755 docker-entrypoint.sh /usr/local/bin/docker-entrypoint.sh
5558

5659
ARG OCI_CREATED=""

‎crates/agentic-llm-d/Cargo.toml‎

Lines changed: 30 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,30 @@
1+
[package]
2+
name = "agentic-llm-d"
3+
description = "Backend mode for agentic-api: split-execution endpoints for the llm-d coordinator"
4+
version.workspace = true
5+
edition.workspace = true
6+
license.workspace = true
7+
repository.workspace = true
8+
publish = false
9+
10+
[dependencies]
11+
agentic-core.workspace = true
12+
axum.workspace = true
13+
clap.workspace = true
14+
jsonwebtoken.workspace = true
15+
serde.workspace = true
16+
serde_json.workspace = true
17+
thiserror.workspace = true
18+
tokio.workspace = true
19+
tokio-util.workspace = true
20+
tracing.workspace = true
21+
tracing-subscriber.workspace = true
22+
23+
[lints]
24+
workspace = true
25+
26+
[dev-dependencies]
27+
agentic-core = { workspace = true, features = [] }
28+
reqwest = { workspace = true, features = ["json"] }
29+
serde_json.workspace = true
30+
tokio = { workspace = true, features = ["macros", "rt-multi-thread"] }
Lines changed: 124 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,124 @@
1+
//! The wire protocol between `/v1alpha/responses/hydrate` and `.../persist`. Only
2+
//! this crate speaks it; core keeps the operations it composes.
3+
#![allow(clippy::result_large_err)] // `ExecutorError` is core's; boxing it is not ours to decide
4+
5+
use std::time::{Duration, SystemTime, UNIX_EPOCH};
6+
7+
use jsonwebtoken::{Algorithm, DecodingKey, EncodingKey, Header, Validation, decode, encode};
8+
use serde::{Deserialize, Serialize};
9+
use serde_json::value::RawValue;
10+
11+
use agentic_core::executor::request::RequestContext;
12+
use agentic_core::executor::{ExecutorError, ExecutorResult};
13+
use agentic_core::types::io::{ResponsesInput, ToolChoice};
14+
use agentic_core::types::request_response::RequestPayload;
15+
use agentic_core::types::tools::ResponsesTool;
16+
17+
/// What `hydrate` returns. Raw JSON: the caller forwards it uninterpreted.
18+
#[derive(Debug, Clone, Serialize, Deserialize)]
19+
pub struct Hydration {
20+
pub request: Box<RawValue>,
21+
/// Sealed: echo back to `persist` unchanged. Opaque to the caller.
22+
pub context: String,
23+
}
24+
25+
/// Wire form of a [`RequestContext`]. `enriched_request` and `new_input_items`
26+
/// are absent on purpose: both are rebuilt on return.
27+
#[derive(Debug, Clone, Serialize, Deserialize)]
28+
pub struct SplitContext {
29+
pub response_id: String,
30+
pub original_request: RequestPayload,
31+
/// Inherited from the continued turn; the request's own is rejected.
32+
pub conversation_id: Option<String>,
33+
/// The only part of `enriched_request` that `original_request` cannot supply.
34+
pub effective_tools: Option<Vec<ResponsesTool>>,
35+
pub effective_tool_choice: Option<ToolChoice>,
36+
}
37+
38+
impl From<RequestContext> for SplitContext {
39+
fn from(ctx: RequestContext) -> Self {
40+
Self {
41+
response_id: ctx.response_id,
42+
original_request: ctx.original_request,
43+
conversation_id: ctx.conversation_id,
44+
effective_tools: ctx.enriched_request.tools,
45+
effective_tool_choice: ctx.enriched_request.tool_choice,
46+
}
47+
}
48+
}
49+
50+
impl From<SplitContext> for RequestContext {
51+
fn from(wire: SplitContext) -> Self {
52+
let new_input_items = Vec::from(&wire.original_request.input);
53+
let mut enriched_request = wire.original_request.clone();
54+
enriched_request.previous_response_id = None;
55+
enriched_request.input = ResponsesInput::Items(new_input_items.clone());
56+
enriched_request.tools = wire.effective_tools;
57+
enriched_request.tool_choice = wire.effective_tool_choice;
58+
Self {
59+
original_request: wire.original_request,
60+
enriched_request,
61+
new_input_items,
62+
response_id: wire.response_id,
63+
conversation_id: wire.conversation_id,
64+
// Conversation mode is rejected, so there is no version to resume.
65+
conversation_version: None,
66+
}
67+
}
68+
}
69+
70+
/// Rejects requests needing state the in-process flow keeps between steps.
71+
///
72+
/// # Errors
73+
/// [`ExecutorError::InvalidRequest`] naming the feature that cannot be split.
74+
pub fn ensure_splittable(request: &RequestPayload) -> ExecutorResult<()> {
75+
if let Some(feature) = request.in_process_feature() {
76+
return Err(ExecutorError::InvalidRequest(format!(
77+
"{feature} is not supported for split execution"
78+
)));
79+
}
80+
Ok(())
81+
}
82+
83+
/// Only has to outlive one inference call, so generous rather than tuned.
84+
const CONTEXT_TTL: Duration = Duration::from_secs(600);
85+
const AUDIENCE: &str = "agentic-llm-d";
86+
87+
#[derive(Serialize, Deserialize)]
88+
struct SealedClaims {
89+
exp: u64,
90+
aud: String,
91+
ctx: SplitContext,
92+
}
93+
94+
/// Seals a context so `persist` can prove `hydrate` issued it.
95+
///
96+
/// # Errors
97+
/// [`ExecutorError::InvalidRequest`] if the token cannot be produced.
98+
pub fn seal(context: SplitContext, key: &[u8]) -> ExecutorResult<String> {
99+
let expires = SystemTime::now()
100+
.checked_add(CONTEXT_TTL)
101+
.and_then(|at| at.duration_since(UNIX_EPOCH).ok())
102+
.ok_or_else(|| ExecutorError::InvalidRequest("cannot compute context expiry".to_owned()))?;
103+
let claims = SealedClaims {
104+
exp: expires.as_secs(),
105+
aud: AUDIENCE.to_owned(),
106+
ctx: context,
107+
};
108+
encode(&Header::new(Algorithm::HS256), &claims, &EncodingKey::from_secret(key))
109+
.map_err(|error| ExecutorError::InvalidRequest(format!("cannot seal context: {error}")))
110+
}
111+
112+
/// Opens a sealed context, rejecting one that was tampered with or has expired.
113+
///
114+
/// # Errors
115+
/// [`ExecutorError::InvalidRequest`] for a bad signature, a wrong audience, or
116+
/// a context past its expiry.
117+
pub fn unseal(token: &str, key: &[u8]) -> ExecutorResult<SplitContext> {
118+
let mut validation = Validation::new(Algorithm::HS256);
119+
validation.set_audience(&[AUDIENCE]);
120+
validation.set_required_spec_claims(&["exp", "aud"]);
121+
decode::<SealedClaims>(token, &DecodingKey::from_secret(key), &validation)
122+
.map(|data| data.claims.ctx)
123+
.map_err(|error| ExecutorError::InvalidRequest(format!("context rejected: {error}")))
124+
}
Lines changed: 159 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,159 @@
1+
//! The endpoints and their axum glue. The split routes require a shared token;
2+
//! probes do not. Keep the listener cluster-internal regardless — the token
3+
//! authenticates the calling workload, not a tenant.
4+
5+
use std::time::Duration;
6+
7+
use axum::body::Body;
8+
use axum::extract::{Request, State};
9+
use axum::http::StatusCode;
10+
use axum::middleware::Next;
11+
use axum::response::{IntoResponse, Response};
12+
use serde::Deserialize;
13+
use serde::de::DeserializeOwned;
14+
use serde_json::value::RawValue;
15+
use tracing::warn;
16+
17+
use agentic_core::executor::request::RequestContext;
18+
19+
use agentic_core::executor::{
20+
ExecutorError, UpstreamBody, commit, decode_upstream, rehydrate_conversation, upstream_request,
21+
};
22+
use agentic_core::types::request_response::RequestPayload;
23+
24+
use crate::BackendState;
25+
use crate::context::{Hydration, ensure_splittable, seal, unseal};
26+
27+
const MAX_BODY_SIZE: usize = 10 * 1024 * 1024;
28+
/// The calling workload's shared secret.
29+
pub const WORKLOAD_TOKEN_HEADER: &str = "x-agentic-workload-token";
30+
/// Readiness means storage answers - llm-d owns the model fleet.
31+
const STORAGE_PROBE_TIMEOUT: Duration = Duration::from_secs(2);
32+
33+
/// Body of `POST /v1alpha/responses/persist`: the context, plus one response form.
34+
#[derive(Debug, Deserialize)]
35+
pub struct PersistRequest {
36+
context: String,
37+
response: Option<Box<RawValue>>,
38+
sse: Option<String>,
39+
}
40+
41+
/// Rejects any split-route call without the shared secret. The probes are
42+
/// layered separately and stay open.
43+
pub async fn require_token(State(state): State<BackendState>, request: Request, next: Next) -> Response {
44+
// Not `Authorization`: that stays free for the end user's token.
45+
let presented = request
46+
.headers()
47+
.get(WORKLOAD_TOKEN_HEADER)
48+
.and_then(|value| value.to_str().ok());
49+
match presented {
50+
Some(token) if token_matches(token, &state.api_token) => next.run(request).await,
51+
_ => json(
52+
StatusCode::UNAUTHORIZED,
53+
br#"{"error":{"type":"invalid_request_error","message":"missing or invalid bearer token"}}"#.to_vec(),
54+
),
55+
}
56+
}
57+
58+
/// No early return, so a wrong token takes the same time whatever byte differs.
59+
fn token_matches(presented: &str, expected: &str) -> bool {
60+
presented.len() == expected.len()
61+
&& presented
62+
.bytes()
63+
.zip(expected.bytes())
64+
.fold(0_u8, |differences, (a, b)| differences | (a ^ b))
65+
== 0
66+
}
67+
68+
pub async fn health() -> StatusCode {
69+
StatusCode::OK
70+
}
71+
72+
pub async fn ready(State(state): State<BackendState>) -> StatusCode {
73+
if state.exec_ctx.storage_ready(STORAGE_PROBE_TIMEOUT).await {
74+
StatusCode::OK
75+
} else {
76+
StatusCode::SERVICE_UNAVAILABLE
77+
}
78+
}
79+
80+
pub async fn hydrate(State(state): State<BackendState>, req: Request) -> Response {
81+
let payload: RequestPayload = match read_json(req.into_body()).await {
82+
Ok(payload) => payload,
83+
Err(response) => return response,
84+
};
85+
match build_hydration(payload, &state).await {
86+
Ok(hydration) => axum::Json(hydration).into_response(),
87+
Err(error) => error_response(error),
88+
}
89+
}
90+
91+
/// Rehydrates the turn and builds the request the caller forwards to a model.
92+
#[allow(clippy::result_large_err)] // `ExecutorError` is core's; boxing it is not ours to decide
93+
async fn build_hydration(
94+
request: RequestPayload,
95+
state: &BackendState,
96+
) -> agentic_core::executor::ExecutorResult<Hydration> {
97+
ensure_splittable(&request)?;
98+
let ctx = rehydrate_conversation(request, state.exec_ctx.as_ref()).await?;
99+
// Rehydration can restore a gateway-owned tool from the stored turn, so
100+
// check what will actually run.
101+
ensure_splittable(&ctx.enriched_request)?;
102+
let stream = ctx.original_request.stream;
103+
let request = RawValue::from_string(upstream_request(&ctx, stream)?).map_err(ExecutorError::JsonError)?;
104+
let context = seal(ctx.into(), &state.signing_key)?;
105+
Ok(Hydration { request, context })
106+
}
107+
108+
pub async fn persist(State(state): State<BackendState>, req: Request) -> Response {
109+
let PersistRequest { context, response, sse } = match read_json(req.into_body()).await {
110+
Ok(request) => request,
111+
Err(response) => return response,
112+
};
113+
// serde rejects `RawValue` in `untagged`, so "exactly one of" is checked here.
114+
let upstream = match (response.as_deref(), sse.as_deref()) {
115+
(Some(json), None) => UpstreamBody::Json(json.get()),
116+
(None, Some(sse)) => UpstreamBody::Sse(sse),
117+
_ => {
118+
let message = "exactly one of `response` or `sse` is required".to_owned();
119+
return error_response(ExecutorError::InvalidRequest(message));
120+
}
121+
};
122+
let context = match unseal(&context, &state.signing_key) {
123+
Ok(context) => context,
124+
Err(error) => return error_response(error),
125+
};
126+
let ctx = RequestContext::from(context);
127+
let stored = match decode_upstream(&ctx, upstream) {
128+
Ok(payload) => commit(ctx, payload, state.exec_ctx.as_ref()).await,
129+
Err(error) => Err(error),
130+
};
131+
match stored {
132+
Ok(payload) => axum::Json(payload).into_response(),
133+
Err(error) => error_response(error),
134+
}
135+
}
136+
137+
/// Renders an error with the status and envelope core defines.
138+
fn error_response(error: ExecutorError) -> Response {
139+
let status = error.http_status();
140+
warn!("backend error ({status}): {error}");
141+
json(status, error.into_response_body())
142+
}
143+
144+
#[allow(clippy::result_large_err)] // an axum `Response` is the idiomatic error here
145+
async fn read_json<T: DeserializeOwned>(body: Body) -> Result<T, Response> {
146+
let too_large = br#"{"error":{"type":"invalid_request_error","message":"request body too large"}}"#;
147+
let bytes = axum::body::to_bytes(body, MAX_BODY_SIZE)
148+
.await
149+
.map_err(|_| json(StatusCode::PAYLOAD_TOO_LARGE, too_large.to_vec()))?;
150+
serde_json::from_slice(&bytes).map_err(|error| error_response(ExecutorError::from(error)))
151+
}
152+
153+
fn json(status: StatusCode, body: Vec<u8>) -> Response {
154+
Response::builder()
155+
.status(status)
156+
.header("Content-Type", "application/json")
157+
.body(Body::from(body))
158+
.expect("valid response")
159+
}

0 commit comments

Comments
 (0)