Skip to content
Merged
Show file tree
Hide file tree
Changes from all 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 @@ -12,7 +12,7 @@
"crates/agentic-server-core/src/types/request_response.rs": 518,
"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": 784
},
"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 @@ -469,6 +469,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 @@ -1277,6 +1278,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 @@ -305,6 +305,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,7 +116,15 @@ 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
}

#[tracing::instrument(name = "agentic.persist", skip_all, fields(
Expand All @@ -113,6 +134,7 @@ 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 @@ -121,6 +143,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