feat: add resource limits management and network client integration

- Introduced a new module for managing process resource limits, specifically for raising the open file limit on Unix systems.
- Added a `NetworkClients` struct to handle reusable outbound HTTP clients, improving network request management.
- Updated various components, including `ControlService` and `CursorProxy`, to utilize the new network client structure for better client handling.
- Enhanced the API router to accept network clients, ensuring consistent client usage across different services.
- Added tests to validate the integration of network clients and resource limits functionality.
This commit is contained in:
leokun
2026-09-01 21:03:02 +08:00
parent 76417e005b
commit 2dad593263
13 changed files with 203 additions and 32 deletions
+5 -2
View File
@@ -27,8 +27,11 @@ use crate::{
Result,
};
pub fn router(registry: TransportRegistry) -> Result<Router> {
let proxy = CursorProxy::cursor(registry.store().clone())?;
pub fn router(
registry: TransportRegistry,
clients: crate::network::NetworkClients,
) -> Result<Router> {
let proxy = CursorProxy::cursor(clients);
let knowledge = knowledge::KnowledgeService::managed()?;
Ok(router_with_proxy(registry, proxy, knowledge))
}
+6 -15
View File
@@ -14,8 +14,7 @@ pub const UPSTREAM_URL_HEADER: &str = "x-server-upstream-url";
#[derive(Clone)]
pub struct CursorProxy {
client: Option<reqwest::Client>,
store: Option<crate::store::Store>,
clients: crate::network::NetworkClients,
upstream: String,
}
@@ -47,23 +46,15 @@ impl BufferedResponse {
}
impl CursorProxy {
pub fn cursor(store: crate::store::Store) -> Result<Self> {
Ok(Self {
client: None,
store: Some(store),
pub fn cursor(clients: crate::network::NetworkClients) -> Self {
Self {
clients,
upstream: CURSOR_UPSTREAM.into(),
})
}
}
async fn client(&self) -> Result<reqwest::Client> {
match (&self.client, &self.store) {
(Some(client), _) => Ok(client.clone()),
(_, Some(store)) => Ok(crate::network::client_builder(store)
.await?
.redirect(reqwest::redirect::Policy::none())
.build()?),
_ => unreachable!("Cursor proxy always has a client or store"),
}
self.clients.cursor_client().await
}
}
+3 -3
View File
@@ -1,7 +1,7 @@
//! Builds the top-level server router.
use crate::{cursor::transport::TransportRegistry, Result};
use crate::{cursor::transport::TransportRegistry, network::NetworkClients, Result};
pub fn router(registry: TransportRegistry) -> Result<axum::Router> {
super::cursor::router(registry)
pub fn router(registry: TransportRegistry, clients: NetworkClients) -> Result<axum::Router> {
super::cursor::router(registry, clients)
}
+10 -3
View File
@@ -44,9 +44,11 @@ impl App {
plugin_runtime.clone(),
config.app_version.clone(),
)?;
let clients = crate::network::NetworkClients::new(store.clone());
let provider = std::sync::Arc::new(ProviderRouter::new(
store.clone(),
plugins.clone(),
clients.clone(),
config.provider_request_timeout,
config.provider_stream_idle_timeout,
));
@@ -58,10 +60,15 @@ impl App {
plugins.clone(),
crate::config::managed_data_dir()?.join("rules"),
);
let control =
control::ControlService::new(store.clone(), provider, plugin_runtime, plugins)?;
let control = control::ControlService::new(
store.clone(),
provider,
plugin_runtime,
plugins,
clients.clone(),
)?;
let harness = control.cursor_harness().clone();
let mut router = api::router(registry.clone())?;
let mut router = api::router(registry.clone(), clients)?;
router = match &config.console {
Some(ConsoleSource::Directory(directory)) => {
router.merge(control::web_router(control.clone(), directory))
+9 -4
View File
@@ -41,6 +41,7 @@ pub struct ControlService {
provider: Arc<dyn Provider>,
plugin_runtime: PluginRuntime,
plugins: PluginRegistry,
clients: crate::network::NetworkClients,
model_tests: Arc<Mutex<BTreeMap<String, CancellationToken>>>,
}
@@ -151,6 +152,7 @@ impl ControlService {
provider: Arc<dyn Provider>,
plugin_runtime: PluginRuntime,
plugins: PluginRegistry,
clients: crate::network::NetworkClients,
) -> Result<Self> {
Ok(Self {
cursor_harness: CursorHarness::new(store.clone())?,
@@ -158,6 +160,7 @@ impl ControlService {
provider,
plugin_runtime,
plugins,
clients,
model_tests: Arc::new(Mutex::new(BTreeMap::new())),
})
}
@@ -256,7 +259,7 @@ impl ControlService {
disabled_ad_ids: Option<&str>,
language: &str,
) -> Result<AdRuntime> {
let client = crate::network::client(&self.store).await?;
let client = self.clients.default_client().await?;
let installation_id = self.store.installation_id().await?;
let mut request = client
.get(ADS_ENDPOINT)
@@ -281,7 +284,7 @@ impl ControlService {
}
pub(super) async fn dismiss_ad(&self, ad_id: &str, input: &AdDismissalInput) -> Result<()> {
let client = crate::network::client(&self.store).await?;
let client = self.clients.default_client().await?;
let installation_id = self.store.installation_id().await?;
let mut endpoint = Url::parse(ADS_ENDPOINT).map_err(|error| {
Error::Config(format!("advertisement endpoint is invalid: {error}"))
@@ -510,7 +513,7 @@ impl ControlService {
}
pub async fn discover_models(&self, input: &ModelDiscoveryInput) -> Result<DiscoveredModels> {
let client = crate::network::client(&self.store).await?;
let client = self.clients.default_client().await?;
let base_url = crate::model::normalize_request_url(&input.base_url)?;
discover_models_from_endpoint(
&client,
@@ -677,7 +680,9 @@ impl ControlService {
}
pub async fn set_proxy_settings(&self, settings: ProxySettingsInput) -> Result<ProxySettings> {
self.store.set_proxy_settings(settings).await
let settings = self.store.set_proxy_settings(settings).await?;
self.clients.invalidate().await;
Ok(settings)
}
pub async fn tab_settings(&self) -> Result<TabSettings> {
+87 -2
View File
@@ -1,8 +1,93 @@
//! Provides shared network client and transport configuration.
//! Outbound HTTP clients configured from persisted application proxy settings.
//! Owns reusable outbound HTTP clients configured from persisted proxy settings.
use std::{sync::Arc, time::Duration};
use tokio::sync::RwLock;
use crate::{store::Store, Result};
#[derive(Clone)]
pub struct NetworkClients {
store: Store,
cache: Arc<RwLock<ClientCache>>,
}
#[derive(Default)]
struct ClientCache {
default: Option<reqwest::Client>,
cursor: Option<reqwest::Client>,
provider: Option<(Duration, reqwest::Client)>,
}
impl NetworkClients {
pub fn new(store: Store) -> Self {
Self {
store,
cache: Arc::new(RwLock::new(ClientCache::default())),
}
}
pub async fn default_client(&self) -> Result<reqwest::Client> {
if let Some(client) = self.cache.read().await.default.clone() {
return Ok(client);
}
let mut cache = self.cache.write().await;
if let Some(client) = cache.default.clone() {
return Ok(client);
}
let client = client_builder(&self.store).await?.build()?;
cache.default = Some(client.clone());
Ok(client)
}
pub async fn cursor_client(&self) -> Result<reqwest::Client> {
if let Some(client) = self.cache.read().await.cursor.clone() {
return Ok(client);
}
let mut cache = self.cache.write().await;
if let Some(client) = cache.cursor.clone() {
return Ok(client);
}
let client = client_builder(&self.store)
.await?
.redirect(reqwest::redirect::Policy::none())
.build()?;
cache.cursor = Some(client.clone());
Ok(client)
}
pub async fn provider_client(&self, timeout: Duration) -> Result<reqwest::Client> {
if let Some((_, client)) = self
.cache
.read()
.await
.provider
.as_ref()
.filter(|(cached_timeout, _)| *cached_timeout == timeout)
{
return Ok(client.clone());
}
let mut cache = self.cache.write().await;
if let Some((_, client)) = cache
.provider
.as_ref()
.filter(|(cached_timeout, _)| *cached_timeout == timeout)
{
return Ok(client.clone());
}
let client = client_builder(&self.store)
.await?
.timeout(timeout)
.build()?;
cache.provider = Some((timeout, client.clone()));
Ok(client)
}
pub async fn invalidate(&self) {
*self.cache.write().await = ClientCache::default();
}
}
pub async fn client_builder(store: &Store) -> Result<reqwest::ClientBuilder> {
let settings = store.proxy_settings_secret().await?;
// Use the platform TLS stack for compatibility with provider gateways that
+5 -1
View File
@@ -21,6 +21,7 @@ use super::{
pub struct ProviderRouter {
store: Store,
plugins: PluginRegistry,
clients: crate::network::NetworkClients,
request_timeout: Duration,
stream_idle_timeout: Duration,
}
@@ -29,12 +30,14 @@ impl ProviderRouter {
pub fn new(
store: Store,
plugins: PluginRegistry,
clients: crate::network::NetworkClients,
request_timeout: Duration,
stream_idle_timeout: Duration,
) -> Self {
Self {
store,
plugins,
clients,
request_timeout,
stream_idle_timeout,
}
@@ -49,6 +52,7 @@ impl Provider for ProviderRouter {
) -> ProviderStream {
let store = self.store.clone();
let plugins = self.plugins.clone();
let clients = self.clients.clone();
let request_timeout = self.request_timeout;
let stream_idle_timeout = self.stream_idle_timeout;
Box::pin(try_stream! {
@@ -91,7 +95,7 @@ impl Provider for ProviderRouter {
request_timeout,
allowed_body_fields: None,
};
let client = crate::network::client_builder(&store).await?.timeout(request_timeout).build()?;
let client = clients.provider_client(request_timeout).await?;
let provider = build_observed(&config, recorder.clone(), client)?;
(recorder, guard, provider.stream(routed, cancellation.clone()))
};
+2 -2
View File
@@ -106,7 +106,7 @@ async fn decode<M: Message + Default>(response: Response<Body>) -> M {
#[tokio::test]
async fn offline_crud_round_trip_persists_markdown() {
let (_store_dir, store) = fixtures::temp_store().await;
let upstream = CursorProxy::cursor(store).unwrap();
let upstream = CursorProxy::cursor(cursor_server::network::NetworkClients::new(store));
let rules_dir = tempfile::tempdir().unwrap();
let rules_root = rules_dir.path().join("rules");
let service = KnowledgeService::with_root(rules_root.clone()).unwrap();
@@ -193,7 +193,7 @@ async fn offline_crud_round_trip_persists_markdown() {
#[tokio::test]
async fn updating_missing_rule_reports_failure() {
let (_store_dir, store) = fixtures::temp_store().await;
let upstream = CursorProxy::cursor(store).unwrap();
let upstream = CursorProxy::cursor(cursor_server::network::NetworkClients::new(store));
let rules_dir = tempfile::tempdir().unwrap();
let service = KnowledgeService::with_root(rules_dir.path().join("rules")).unwrap();