Skip to content
Draft
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
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.

1 change: 1 addition & 0 deletions crates/alien-ai-gateway/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ aws-credential-types = { workspace = true }
aws-smithy-eventstream = { workspace = true }
aws-smithy-types = { workspace = true }
http = { workspace = true }
uuid = { workspace = true, features = ["v4"] }

[dev-dependencies]
httpmock = { workspace = true }
Expand Down
10 changes: 8 additions & 2 deletions crates/alien-ai-gateway/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,14 +10,20 @@ mod config;
mod creds;
mod error;
mod router;
mod usage;
pub use config::{bindings_from_env, bindings_from_env_map, route_from_remote_ai_lease};
pub use creds::{
AmbientCred, AnthropicApiKeyCred, AwsSigV4Cred, BearerTokenCred, OpenAiApiKeyCred,
};
pub use error::{ErrorData, Result};
pub use router::{
build_router, build_router_with_availability, route_from_direct_anthropic,
route_from_direct_openai, AvailableModels, GatewayRoute, GatewayTarget,
build_router, build_router_with_availability, build_router_with_availability_and_observer,
build_router_with_observer, route_from_direct_anthropic, route_from_direct_openai,
AvailableModels, GatewayRoute, GatewayTarget,
};
pub use usage::{
parse_ai_token_usage, AiTokenUsage, AiUsageClientApi, AiUsageEvent, AiUsageObserver,
AiUsageOutcome, AiUsageProvider,
};

use std::net::{Ipv4Addr, SocketAddr};
Expand Down
138 changes: 126 additions & 12 deletions crates/alien-ai-gateway/src/router/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,9 @@ use serde_json::{json, Value};

use crate::creds::{AmbientCred, AnthropicApiKeyCred, OpenAiApiKeyCred};
use crate::error::{ErrorData, Result};
use crate::usage::{
observe_response, AiUsageClientApi, AiUsageContext, AiUsageObserver, AiUsageProvider,
};

mod bedrock;
mod eventstream;
Expand Down Expand Up @@ -141,6 +144,7 @@ struct AppState {
/// Account-specific, read-only control-plane observations supplied by the
/// hosted route resolver. `None` keeps embedded gateways catalog-only.
available_models: Option<AvailableModels>,
usage_observer: Option<Arc<dyn AiUsageObserver>>,
}

/// Available public model IDs keyed by binding name.
Expand All @@ -150,7 +154,15 @@ pub type AvailableModels = HashMap<String, HashSet<String>>;
/// `POST /<name>/v1/chat/completions` (OpenAI), `POST /<name>/v1/messages`
/// (Anthropic), and `GET /<name>/v1/models`.
pub fn build_router(routes: Vec<GatewayRoute>) -> Router {
build_router_inner(routes, None)
build_router_inner(routes, None, None)
}

/// Build a router that reports completed requests to a non-blocking observer.
pub fn build_router_with_observer(
routes: Vec<GatewayRoute>,
usage_observer: Arc<dyn AiUsageObserver>,
) -> Router {
build_router_inner(routes, None, Some(usage_observer))
}

/// Build a router whose model listing and inference paths are restricted by a
Expand All @@ -159,19 +171,30 @@ pub fn build_router_with_availability(
routes: Vec<GatewayRoute>,
available_models: AvailableModels,
) -> Router {
build_router_inner(routes, Some(available_models))
build_router_inner(routes, Some(available_models), None)
}

/// Build a hosted router with both bounded model availability and usage observation.
pub fn build_router_with_availability_and_observer(
routes: Vec<GatewayRoute>,
available_models: AvailableModels,
usage_observer: Arc<dyn AiUsageObserver>,
) -> Router {
build_router_inner(routes, Some(available_models), Some(usage_observer))
}

fn build_router_inner(
routes: Vec<GatewayRoute>,
available_models: Option<AvailableModels>,
usage_observer: Option<Arc<dyn AiUsageObserver>>,
) -> Router {
let routes: HashMap<String, GatewayRoute> =
routes.into_iter().map(|r| (r.name.clone(), r)).collect();
let state = Arc::new(AppState {
routes,
client: reqwest::Client::new(),
available_models,
usage_observer,
});
Router::new()
.route(
Expand Down Expand Up @@ -275,6 +298,23 @@ async fn forward_response(upstream: reqwest::Response) -> Result<Response> {
})
}

fn usage_client_api(client_api: ClientApi) -> AiUsageClientApi {
match client_api {
ClientApi::OpenAiChatCompletions => AiUsageClientApi::OpenAiChatCompletions,
ClientApi::OpenAiResponses => AiUsageClientApi::OpenAiResponses,
ClientApi::AnthropicMessages => AiUsageClientApi::AnthropicMessages,
}
}

fn cloud_usage_provider(cloud: Platform) -> AiUsageProvider {
match cloud {
Platform::Aws => AiUsageProvider::AwsBedrock,
Platform::Gcp => AiUsageProvider::GcpVertex,
Platform::Azure => AiUsageProvider::AzureFoundry,
_ => unreachable!("AI cloud routes are available only on AWS, GCP, and Azure"),
}
}

/// Build a JSON POST to `url`, sign it with the ambient credential for `service`,
/// and execute it. The handlers differ only in URL, signing service, body, and any
/// protocol-required header, so the build + sign + execute + upstream-error
Expand Down Expand Up @@ -365,7 +405,24 @@ async fn proxy(
message: format!("direct Anthropic supports only /{binding}/v1/messages"),
}));
}
return proxy_direct_anthropic(&state.client, route, payload, &model, &headers).await;
let provider_model = ai_catalog::resolve_direct_anthropic(&model)
.map(|resolved| resolved.upstream_id)
.unwrap_or(model.as_str());
let descriptor = AiUsageContext::new(
&binding,
AiUsageProvider::Anthropic,
&model,
provider_model,
usage_client_api(client_api),
None,
);
let response =
proxy_direct_anthropic(&state.client, route, payload, &model, &headers).await?;
return Ok(observe_response(
response,
state.usage_observer.as_ref(),
descriptor,
));
}
GatewayTarget::DirectOpenAi => {
ensure_model_available(&state, &binding, &model)?;
Expand All @@ -376,14 +433,27 @@ async fn proxy(
),
}));
}
return proxy_direct_openai(
let descriptor = AiUsageContext::new(
&binding,
AiUsageProvider::OpenAi,
&model,
&model,
usage_client_api(client_api),
None,
);
let response = proxy_direct_openai(
&state.client,
route,
payload,
&model,
"/v1/chat/completions",
)
.await;
return Ok(observe_response(
response?,
state.usage_observer.as_ref(),
descriptor,
));
}
GatewayTarget::Cloud(_) => {}
}
Expand All @@ -401,6 +471,15 @@ async fn proxy(
})?;
ensure_model_available(&state, &binding, &model)?;

let descriptor = AiUsageContext::new(
&binding,
cloud_usage_provider(cloud),
&model,
cm.upstream_id,
usage_client_api(client_api),
route.region.clone(),
);

if !cm.client_apis.contains(&client_api) {
let expected_path = match cm.client_apis.first() {
Some(ClientApi::OpenAiChatCompletions) => "v1/chat/completions",
Expand All @@ -419,21 +498,38 @@ async fn proxy(
// endpoint: the model id travels in the URL and the streamed reply is AWS
// event-stream framing, so it needs its own request/response shape.
if cloud == Platform::Aws && cm.provider_api == ProviderApi::Anthropic {
return proxy_bedrock_anthropic(&state.client, route, cm.upstream_id, payload, &headers)
.await;
let response =
proxy_bedrock_anthropic(&state.client, route, cm.upstream_id, payload, &headers)
.await?;
return Ok(observe_response(
response,
state.usage_observer.as_ref(),
descriptor,
));
}
// GCP serves Claude through Vertex rawPredict: the model id travels in the URL
// and streaming is chosen by the URL verb, but the reply is native Anthropic
// JSON/SSE — no decoder needed, unlike Bedrock.
if cloud == Platform::Gcp && cm.provider_api == ProviderApi::Anthropic {
return proxy_vertex_anthropic(&state.client, route, cm.upstream_id, payload, &headers)
.await;
let response =
proxy_vertex_anthropic(&state.client, route, cm.upstream_id, payload, &headers).await?;
return Ok(observe_response(
response,
state.usage_observer.as_ref(),
descriptor,
));
}
// Azure serves Claude through Foundry's Anthropic endpoint: standard Messages
// in both directions, on the `/anthropic/v1` path with the version header.
if cloud == Platform::Azure && cm.provider_api == ProviderApi::Anthropic {
return proxy_foundry_anthropic(&state.client, route, cm.upstream_id, payload, &headers)
.await;
let response =
proxy_foundry_anthropic(&state.client, route, cm.upstream_id, payload, &headers)
.await?;
return Ok(observe_response(
response,
state.usage_observer.as_ref(),
descriptor,
));
}

payload["model"] = Value::String(cm.upstream_id.to_string());
Expand All @@ -456,7 +552,12 @@ async fn proxy(
)
.await?;

forward_response(upstream).await
let response = forward_response(upstream).await?;
Ok(observe_response(
response,
state.usage_observer.as_ref(),
descriptor,
))
}

/// Proxy an OpenAI Responses request (`POST /<name>/v1/responses`, used by Codex).
Expand Down Expand Up @@ -509,6 +610,14 @@ async fn proxy_responses(
binding: binding.clone(),
})
})?;
let descriptor = AiUsageContext::new(
&binding,
cloud_usage_provider(cloud),
&model,
target.upstream_id,
AiUsageClientApi::OpenAiResponses,
route.region.clone(),
);

payload["model"] = Value::String(target.upstream_id.to_string());
let upstream_body =
Expand Down Expand Up @@ -538,7 +647,12 @@ async fn proxy_responses(
)
.await?;

forward_response(upstream).await
let response = forward_response(upstream).await?;
Ok(observe_response(
response,
state.usage_observer.as_ref(),
descriptor,
))
}

/// `GET /<name>/v1/models`: the qualified catalog, intersected with the bounded
Expand Down
Loading
Loading