Skip to content
Merged
Show file tree
Hide file tree
Changes from 3 commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .rust-file-sizes.json
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,7 @@
"crates/agentic-server-core/src/types/request_response.rs": 582,
"crates/agentic-server-core/src/types/tools/params.rs": 537,
"crates/agentic-server/src/agentic_process.rs": 557,
"crates/agentic-server/src/auth.rs": 799,
"crates/agentic-server/src/auth.rs": 791,
"crates/agentic-server/src/handler/websocket/responses.rs": 841
},
"exceptions": {},
Expand Down
9 changes: 8 additions & 1 deletion ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -871,7 +871,8 @@ round that omits `usage` still reports the hidden rounds' counters.
here.
- **`types/`** — the conversion layer from those raw rows into business types, via
`From`/`TryFrom` impls: `ConversationData`/`ConversationSnapshot`, `ResponseData`/
`ResponseMetadata` (parses the JSON metadata column into a typed struct),
`ResponseMetadata` (parses the JSON metadata column into a typed struct, including an optional
terminal `ResponsePayload` snapshot for GET retrieval),
`InOutItem` (parses an `Item.data` JSON blob back into a typed `InputItem` or
`OutputItem`), and `StorageError`. `InOutItem::into_input_items` turns a full
history into the `Vec<InputItem>` used for continuation processing: stored
Expand Down Expand Up @@ -901,6 +902,12 @@ round that omits `usage` still reports the hidden rounds' counters.
them directly for fixtures — that's expected and fine; production code paths should
not.)

Stored Responses snapshots are written in the same transaction as response history. Retrieval goes through
`ResponseHandler::retrieve`, independently of upstream availability. Continuation checkpoints omit the
snapshot to avoid retaining a duplicate response; they continue to use canonical history and effective
settings. Legacy history-only records remain usable for continuation, but GET retrieval reports a conflict
rather than fabricating status, usage, or output.

### `tool/` — the tool framework

Wire shapes for tool declarations live in `types::tools` (see above); this module owns
Expand Down
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

11 changes: 11 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -65,13 +65,24 @@ flowchart LR
| Endpoint | Description | Status |
| --- | --- | --- |
| `POST /v1/responses` | OpenAI-compatible Responses API with state, tools, and streaming | ✅ |
| `GET /v1/responses/{response_id}` | Retrieve a locally stored response | ✅ |
| `GET /v1/responses` | WebSocket transport for the Responses API | ✅ |
| `POST /v1/conversations` | Conversation management | ✅ |
| `GET /v1/models` | Model listing proxied from vLLM | ✅ |
| `GET /health` · `GET /ready` | Liveness and readiness probes | ✅ |
| Messages API | Anthropic-style stateful messages on shared primitives | 🚧 Planned |
| Interactions API | Higher-level agentic workflow surface | ⏳ Planned |

Responses created with `store: true` retain a terminal snapshot for retrieval, including status, usage,
and this turn's output. Retrieval does not call the upstream model. Unknown IDs return a JSON `404`;
older records created before snapshot storage (or through history-only APIs) return `409` because their
original response cannot be reconstructed faithfully. Request-scoped MCP credentials are stripped from
stored tool definitions. Responses created with `store: false` do not have retrievable snapshots.

Retrieval requires a valid OIDC bearer token when OIDC is enabled. Otherwise, when `OPENAI_API_KEY` is
nonempty, callers must send that key in `Authorization: Bearer <key>`; missing or invalid credentials
return `401` before storage is read. With neither configured, retrieval allows unauthenticated access.

## 🚀 Quickstart

### Agentic API CLI
Expand Down
2 changes: 2 additions & 0 deletions crates/agentic-server-core/src/executor/compaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -481,6 +481,7 @@ pub async fn compact_response(
ctx,
tool_search_metadata,
Vec::new(),
None,
&exec_ctx.conv_handler,
&exec_ctx.resp_handler,
)
Expand Down Expand Up @@ -1266,6 +1267,7 @@ mod tests {
previous_response_id: None,
effective_tools: None,
tool_search_loaded_tools: None,
response_snapshot: None,
effective_tool_choice: crate::ToolChoice::Auto,
effective_instructions: None,
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -132,6 +132,7 @@ impl ConversationHandler {
previous_response_id: ctx.original_request.previous_response_id.take(),
effective_tools: ctx.enriched_request.tools.take(),
tool_search_loaded_tools: None,
response_snapshot: None,
effective_tool_choice: ctx.enriched_request.tool_choice.take().unwrap_or_default(),
effective_instructions: ctx.enriched_request.instructions.take(),
};
Expand Down
30 changes: 28 additions & 2 deletions crates/agentic-server-core/src/executor/modes/response.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@

use crate::storage::{InOutItem, ResponseData, ResponseMetadata, ResponseStore};
use crate::types::io::OutputItem;
use crate::types::request_response::ResponsePayload;

use crate::executor::error::{ExecutorError, ExecutorResult};
use crate::executor::request::RequestContext;
Expand All @@ -18,6 +19,24 @@ impl ResponseHandler {
Self { store }
}

/// Retrieves the terminal payload of a response created with `store=true`.
///
/// # Errors
/// Returns a storage error for missing IDs or unavailable storage, and a conflict
/// for older records that contain continuation history but no response snapshot.
pub async fn retrieve(&self, response_id: &str) -> ExecutorResult<ResponsePayload> {
let stored = self.store.get(response_id).await?;
stored
.metadata
.response_snapshot
.map(|snapshot| *snapshot)
.ok_or_else(|| {
ExecutorError::Conflict(
"stored response has no retrievable payload; create a new response with store=true".into(),
)
})
}

/// Retrieves the stored response for `previous_response_id`.
///
/// Reads `previous_response_id` from `ctx.original_request`.
Expand All @@ -31,7 +50,10 @@ impl ResponseHandler {
.previous_response_id
.as_deref()
.ok_or_else(|| ExecutorError::InvalidRequest("previous_response_id is required for get".into()))?;
self.store.get(prev_id).await.map_err(ExecutorError::Storage)
let mut stored = self.store.get(prev_id).await?;
// Continuation checkpoints retain history and settings, not a duplicate payload.
stored.metadata.response_snapshot = None;
Ok(stored)
}

/// Validates that the response for `previous_response_id` exists.
Expand Down Expand Up @@ -75,6 +97,7 @@ impl ResponseHandler {
previous_response_id: ctx.original_request.previous_response_id.take(),
effective_tools: ctx.enriched_request.tools.take(),
tool_search_loaded_tools: None,
response_snapshot: None,
effective_tool_choice: ctx.enriched_request.tool_choice.take().unwrap_or_default(),
effective_instructions: ctx.enriched_request.instructions.take(),
};
Expand All @@ -87,7 +110,7 @@ impl ResponseHandler {
&self,
mut ctx: RequestContext,
output_items: Vec<OutputItem>,
metadata: ResponseMetadata,
mut metadata: ResponseMetadata,
) -> ExecutorResult<()> {
let continuation = ctx.continuation.take();
let write_durable = continuation.is_none() || ctx.original_request.store;
Expand All @@ -105,6 +128,7 @@ impl ResponseHandler {
.map(|(_, item)| InOutItem::Output(item)),
);

let snapshot = metadata.response_snapshot.take();
let checkpoint = continuation
.as_ref()
.map(|lease| {
Expand All @@ -118,6 +142,8 @@ impl ResponseHandler {
})
.transpose()?;

metadata.response_snapshot = snapshot;

// A stored child of a transient parent needs its complete canonical
// checkpoint, not a database reference to a response that was never stored.
let transient_parent = continuation
Expand Down
29 changes: 26 additions & 3 deletions crates/agentic-server-core/src/executor/persist.rs
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,7 @@ pub async fn persist_response(

let (ctx, tool_search_state) = prepare_request_tools(ctx, &conv_handler, &resp_handler).await?;
let tool_search_metadata = tool_search_state.map(ToolSearchState::into_public_metadata);
persist_prepared_turn(ctx, tool_search_metadata, payload.output, &conv_handler, &resp_handler).await
persist_prepared_response(payload, ctx, tool_search_metadata, conv_handler, resp_handler).await
}

async fn persist_prepared_response(
Expand All @@ -88,7 +88,20 @@ async fn persist_prepared_response(
return Ok(());
}

persist_prepared_turn(ctx, tool_search_metadata, payload.output, &conv_handler, &resp_handler).await
let (output_items, snapshot) = if ctx.original_request.store {
(payload.output.clone(), Some(Box::new(payload)))
} else {
(payload.output, None)
};
persist_prepared_turn(
ctx,
tool_search_metadata,
output_items,
snapshot,
&conv_handler,
&resp_handler,
)
.await
}

/// Persists one completed turn with the handler selected by its explicit conversation discriminator.
Expand All @@ -103,13 +116,22 @@ pub async fn persist_turn(
) -> ExecutorResult<()> {
let (ctx, tool_search_state) = prepare_request_tools(ctx, conv_handler, resp_handler).await?;
let tool_search_metadata = tool_search_state.map(ToolSearchState::into_public_metadata);
persist_prepared_turn(ctx, tool_search_metadata, output_items, conv_handler, resp_handler).await
persist_prepared_turn(
ctx,
tool_search_metadata,
output_items,
None,
conv_handler,
resp_handler,
)
.await
}

pub(crate) async fn persist_prepared_turn(
mut ctx: RequestContext,
tool_search_metadata: Option<ToolSearchMetadata>,
output_items: Vec<OutputItem>,
response_snapshot: Option<Box<ResponsePayload>>,
conv_handler: &ConversationHandler,
resp_handler: &ResponseHandler,
) -> ExecutorResult<()> {
Expand All @@ -118,6 +140,7 @@ pub(crate) async fn persist_prepared_turn(
previous_response_id: ctx.original_request.previous_response_id.take(),
effective_tools: ctx.enriched_request.tools.take(),
tool_search_loaded_tools: None,
response_snapshot,
effective_tool_choice: ctx.enriched_request.tool_choice.take().unwrap_or_default(),
effective_instructions: ctx.enriched_request.instructions.take(),
};
Expand Down
30 changes: 30 additions & 0 deletions crates/agentic-server-core/src/storage/types/response.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,12 +7,16 @@ use serde::{Deserialize, Serialize};
use super::super::models::Response as StorageDbResponse;
use super::errors::StorageError;
use crate::types::io::ToolChoice;
use crate::types::request_response::ResponsePayload;
use crate::types::tools::ResponsesTool;
use crate::utils::common::serialize_to_string;

/// Response metadata with effective configuration.
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
pub struct ResponseMetadata {
/// Exact terminal Responses payload, absent for legacy and non-Responses records.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub response_snapshot: Option<Box<ResponsePayload>>,
pub model: String,
pub previous_response_id: Option<String>,
pub effective_tools: Option<Vec<ResponsesTool>>,
Expand Down Expand Up @@ -73,6 +77,13 @@ impl TryFrom<&ResponseMetadata> for String {
tool.sanitize_for_persistence();
}
}
if let Some(snapshot) = persisted.response_snapshot.as_mut() {
if let Some(tools) = snapshot.tools.as_mut() {
for tool in tools {
tool.sanitize_for_persistence();
}
}
}
serialize_to_string(&persisted).map_err(StorageError::Serialization)
}
}
Expand Down Expand Up @@ -128,6 +139,7 @@ mod tests {
previous_response_id: Some("resp_1".to_string()),
effective_tools: None,
tool_search_loaded_tools: None,
response_snapshot: None,
effective_tool_choice: ToolChoice::Auto,
effective_instructions: Some("be helpful".to_string()),
};
Expand Down Expand Up @@ -164,13 +176,31 @@ mod tests {
}))
.expect("discovered MCP tool"),
});
let snapshot = ResponsePayload {
id: "resp_snapshot".into(),
object: "response".into(),
created_at: 123,
model: "test-model".into(),
status: "completed".into(),
output: Vec::new(),
usage: None,
incomplete_details: None,
error: None,
previous_response_id: None,
conversation_id: None,
instructions: None,
tools: Some(vec![tool.clone()]),
tool_choice: None,
};
let metadata = ResponseMetadata {
effective_tools: Some(vec![tool]),
tool_search_loaded_tools: None,
response_snapshot: Some(Box::new(snapshot)),
..ResponseMetadata::default()
};

let serialized = String::try_from(&metadata).expect("serialization failed");
assert!(!serialized.contains("secret"));
let serialized_value: serde_json::Value =
serde_json::from_str(&serialized).expect("serialized response metadata");
assert!(
Expand Down
1 change: 1 addition & 0 deletions crates/agentic-server/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ opentelemetry_sdk.workspace = true
reqwest = { workspace = true, default-features = false, features = ["rustls-tls"] }
serde.workspace = true
serde_json.workspace = true
subtle = "2.6"
thiserror.workspace = true
tempfile = "3"
tokio.workspace = true
Expand Down
16 changes: 13 additions & 3 deletions crates/agentic-server/src/app.rs
Original file line number Diff line number Diff line change
Expand Up @@ -15,11 +15,12 @@ use tower_http::cors::{AllowOrigin, Any, CorsLayer};
use agentic_core::executor::ExecutionContext;
use agentic_core::proxy::ProxyState;

use crate::auth::api_key::require_api_key;
use crate::auth::{ANTHROPIC_COUNT_TOKENS_PATH, ANTHROPIC_MESSAGES_PATH, OidcAuthenticator, require_oidc};
use crate::handler::{
compact_response, count_tokens, create_conversation, create_item, delete_conversation, delete_item, health,
list_items, messages, models, ready, responses, responses_ws_with_auth, retrieve_conversation, retrieve_item,
update_conversation,
retrieve_response, update_conversation,
};
use crate::model_capabilities::ModelCapabilities;
use crate::telemetry::http::{HttpMetrics, track_request};
Expand Down Expand Up @@ -252,7 +253,8 @@ pub struct AppState {
/// Whether `/ready` should omit the upstream health check.
pub skip_llm_ready_check: bool,
/// Server-configured API key; used as fallback when the request carries no
/// `Authorization` header on the executor path.
/// `Authorization` header on the executor path. Also authenticates local response
/// retrieval when OIDC is disabled.
pub openai_api_key: Option<String>,
/// Configured per-model input-modality overrides applied to the Codex model catalog.
pub model_capabilities: Arc<ModelCapabilities>,
Expand Down Expand Up @@ -284,6 +286,13 @@ pub fn build_router_with_auth(
} else {
public_routes
};
// Retrieval reads local storage, so it cannot rely on upstream credential validation.
let mut retrieval_route = get(retrieve_response);
if authenticator.is_none()
&& let Some(key) = state.openai_api_key.as_ref().filter(|key| !key.is_empty())
{
retrieval_route = retrieval_route.route_layer(middleware::from_fn_with_state(key.clone(), require_api_key));
}
let protected_routes = Router::new()
.route("/v1/conversations", post(create_conversation))
.route(
Expand All @@ -304,7 +313,8 @@ pub fn build_router_with_auth(
.route(ANTHROPIC_MESSAGES_PATH, post(messages))
.route(ANTHROPIC_COUNT_TOKENS_PATH, post(count_tokens))
.route("/v1/responses", post(responses).get(responses_ws_with_auth))
.route("/v1/responses/compact", post(compact_response));
.route("/v1/responses/compact", post(compact_response))
.route("/v1/responses/{response_id}", retrieval_route);
let protected_routes = match authenticator {
Some(authenticator) => {
protected_routes.route_layer(middleware::from_fn_with_state(authenticator, require_oidc))
Expand Down
16 changes: 4 additions & 12 deletions crates/agentic-server/src/auth.rs
Original file line number Diff line number Diff line change
@@ -1,3 +1,7 @@
pub(crate) mod api_key;

use api_key::bearer_token;

use std::collections::HashMap;
use std::sync::Arc;
use std::time::Duration;
Expand Down Expand Up @@ -384,18 +388,6 @@ pub async fn require_oidc(
}
}

fn bearer_token(headers: &axum::http::HeaderMap) -> Option<&str> {
headers
.get(header::AUTHORIZATION)?
.to_str()
.ok()?
.split_once(' ')
.and_then(|(scheme, token)| {
let token = token.trim();
(scheme.eq_ignore_ascii_case("bearer") && !token.is_empty()).then_some(token)
})
}

#[derive(Clone, Copy)]
enum AuthErrorFormat {
OpenAi,
Expand Down
Loading
Loading