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
Generated
+1
View File
@@ -1176,6 +1176,7 @@ version = "0.1.5"
dependencies = [ dependencies = [
"axum", "axum",
"cursor-server", "cursor-server",
"libc",
"rfd", "rfd",
"serde", "serde",
"serde_json", "serde_json",
+1
View File
@@ -14,6 +14,7 @@ tauri-build = { version = "2", features = [] }
[dependencies] [dependencies]
axum = "0.8" axum = "0.8"
cursor-server = { path = "../../../server" } cursor-server = { path = "../../../server" }
libc = "0.2"
rfd = "0.15" rfd = "0.15"
serde = { version = "1", features = ["derive"] } serde = { version = "1", features = ["derive"] }
serde_json = "1" serde_json = "1"
+17
View File
@@ -149,6 +149,23 @@ pub fn run() -> ExitCode {
return ExitCode::FAILURE; return ExitCode::FAILURE;
} }
}; };
#[cfg(unix)]
{
let open_file_limit = match crate::resource_limits::raise_open_file_limit() {
Ok(limit) => limit,
Err(error) => {
diagnostics.report_fatal(&error);
return ExitCode::FAILURE;
}
};
tracing::info!(
requested = crate::resource_limits::REQUESTED_OPEN_FILE_LIMIT,
previous = open_file_limit.previous,
effective = open_file_limit.effective,
hard = open_file_limit.hard,
"open file limit configured"
);
}
tracing::info!( tracing::info!(
version = env!("CARGO_PKG_VERSION"), version = env!("CARGO_PKG_VERSION"),
os = std::env::consts::OS, os = std::env::consts::OS,
+1
View File
@@ -1,6 +1,7 @@
mod desktop; mod desktop;
#[cfg(not(dev))] #[cfg(not(dev))]
mod frontend; mod frontend;
mod resource_limits;
mod startup; mod startup;
mod tray; mod tray;
@@ -0,0 +1,56 @@
//! Configures process resource limits before the desktop runtime starts.
#[cfg(unix)]
use std::io;
#[cfg(unix)]
pub(crate) const REQUESTED_OPEN_FILE_LIMIT: u64 = 65_536;
#[cfg(unix)]
pub(crate) struct OpenFileLimit {
pub(crate) previous: u64,
pub(crate) effective: u64,
pub(crate) hard: u64,
}
#[cfg(unix)]
pub(crate) fn raise_open_file_limit() -> io::Result<OpenFileLimit> {
let mut limits = libc::rlimit {
rlim_cur: 0,
rlim_max: 0,
};
// SAFETY: `limits` points to writable memory for one `rlimit` value.
if unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut limits) } != 0 {
return Err(io::Error::last_os_error());
}
let previous = limits.rlim_cur;
let target = limits
.rlim_max
.min(REQUESTED_OPEN_FILE_LIMIT as libc::rlim_t);
if previous < target {
let requested = libc::rlimit {
rlim_cur: target,
rlim_max: limits.rlim_max,
};
// SAFETY: `requested` is a valid `rlimit` value and does not raise the hard limit.
if unsafe { libc::setrlimit(libc::RLIMIT_NOFILE, &requested) } != 0 {
return Err(io::Error::last_os_error());
}
}
let mut effective = libc::rlimit {
rlim_cur: 0,
rlim_max: 0,
};
// SAFETY: `effective` points to writable memory for one `rlimit` value.
if unsafe { libc::getrlimit(libc::RLIMIT_NOFILE, &mut effective) } != 0 {
return Err(io::Error::last_os_error());
}
Ok(OpenFileLimit {
previous: previous as u64,
effective: effective.rlim_cur as u64,
hard: effective.rlim_max as u64,
})
}
+5 -2
View File
@@ -27,8 +27,11 @@ use crate::{
Result, Result,
}; };
pub fn router(registry: TransportRegistry) -> Result<Router> { pub fn router(
let proxy = CursorProxy::cursor(registry.store().clone())?; registry: TransportRegistry,
clients: crate::network::NetworkClients,
) -> Result<Router> {
let proxy = CursorProxy::cursor(clients);
let knowledge = knowledge::KnowledgeService::managed()?; let knowledge = knowledge::KnowledgeService::managed()?;
Ok(router_with_proxy(registry, proxy, knowledge)) 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)] #[derive(Clone)]
pub struct CursorProxy { pub struct CursorProxy {
client: Option<reqwest::Client>, clients: crate::network::NetworkClients,
store: Option<crate::store::Store>,
upstream: String, upstream: String,
} }
@@ -47,23 +46,15 @@ impl BufferedResponse {
} }
impl CursorProxy { impl CursorProxy {
pub fn cursor(store: crate::store::Store) -> Result<Self> { pub fn cursor(clients: crate::network::NetworkClients) -> Self {
Ok(Self { Self {
client: None, clients,
store: Some(store),
upstream: CURSOR_UPSTREAM.into(), upstream: CURSOR_UPSTREAM.into(),
}) }
} }
async fn client(&self) -> Result<reqwest::Client> { async fn client(&self) -> Result<reqwest::Client> {
match (&self.client, &self.store) { self.clients.cursor_client().await
(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"),
}
} }
} }
+3 -3
View File
@@ -1,7 +1,7 @@
//! Builds the top-level server router. //! 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> { pub fn router(registry: TransportRegistry, clients: NetworkClients) -> Result<axum::Router> {
super::cursor::router(registry) super::cursor::router(registry, clients)
} }
+10 -3
View File
@@ -44,9 +44,11 @@ impl App {
plugin_runtime.clone(), plugin_runtime.clone(),
config.app_version.clone(), config.app_version.clone(),
)?; )?;
let clients = crate::network::NetworkClients::new(store.clone());
let provider = std::sync::Arc::new(ProviderRouter::new( let provider = std::sync::Arc::new(ProviderRouter::new(
store.clone(), store.clone(),
plugins.clone(), plugins.clone(),
clients.clone(),
config.provider_request_timeout, config.provider_request_timeout,
config.provider_stream_idle_timeout, config.provider_stream_idle_timeout,
)); ));
@@ -58,10 +60,15 @@ impl App {
plugins.clone(), plugins.clone(),
crate::config::managed_data_dir()?.join("rules"), crate::config::managed_data_dir()?.join("rules"),
); );
let control = let control = control::ControlService::new(
control::ControlService::new(store.clone(), provider, plugin_runtime, plugins)?; store.clone(),
provider,
plugin_runtime,
plugins,
clients.clone(),
)?;
let harness = control.cursor_harness().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 { router = match &config.console {
Some(ConsoleSource::Directory(directory)) => { Some(ConsoleSource::Directory(directory)) => {
router.merge(control::web_router(control.clone(), directory)) router.merge(control::web_router(control.clone(), directory))
+9 -4
View File
@@ -41,6 +41,7 @@ pub struct ControlService {
provider: Arc<dyn Provider>, provider: Arc<dyn Provider>,
plugin_runtime: PluginRuntime, plugin_runtime: PluginRuntime,
plugins: PluginRegistry, plugins: PluginRegistry,
clients: crate::network::NetworkClients,
model_tests: Arc<Mutex<BTreeMap<String, CancellationToken>>>, model_tests: Arc<Mutex<BTreeMap<String, CancellationToken>>>,
} }
@@ -151,6 +152,7 @@ impl ControlService {
provider: Arc<dyn Provider>, provider: Arc<dyn Provider>,
plugin_runtime: PluginRuntime, plugin_runtime: PluginRuntime,
plugins: PluginRegistry, plugins: PluginRegistry,
clients: crate::network::NetworkClients,
) -> Result<Self> { ) -> Result<Self> {
Ok(Self { Ok(Self {
cursor_harness: CursorHarness::new(store.clone())?, cursor_harness: CursorHarness::new(store.clone())?,
@@ -158,6 +160,7 @@ impl ControlService {
provider, provider,
plugin_runtime, plugin_runtime,
plugins, plugins,
clients,
model_tests: Arc::new(Mutex::new(BTreeMap::new())), model_tests: Arc::new(Mutex::new(BTreeMap::new())),
}) })
} }
@@ -256,7 +259,7 @@ impl ControlService {
disabled_ad_ids: Option<&str>, disabled_ad_ids: Option<&str>,
language: &str, language: &str,
) -> Result<AdRuntime> { ) -> 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 installation_id = self.store.installation_id().await?;
let mut request = client let mut request = client
.get(ADS_ENDPOINT) .get(ADS_ENDPOINT)
@@ -281,7 +284,7 @@ impl ControlService {
} }
pub(super) async fn dismiss_ad(&self, ad_id: &str, input: &AdDismissalInput) -> Result<()> { 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 installation_id = self.store.installation_id().await?;
let mut endpoint = Url::parse(ADS_ENDPOINT).map_err(|error| { let mut endpoint = Url::parse(ADS_ENDPOINT).map_err(|error| {
Error::Config(format!("advertisement endpoint is invalid: {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> { 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)?; let base_url = crate::model::normalize_request_url(&input.base_url)?;
discover_models_from_endpoint( discover_models_from_endpoint(
&client, &client,
@@ -677,7 +680,9 @@ impl ControlService {
} }
pub async fn set_proxy_settings(&self, settings: ProxySettingsInput) -> Result<ProxySettings> { 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> { pub async fn tab_settings(&self) -> Result<TabSettings> {
+87 -2
View File
@@ -1,8 +1,93 @@
//! Provides shared network client and transport configuration. //! Owns reusable outbound HTTP clients configured from persisted proxy settings.
//! Outbound HTTP clients configured from persisted application proxy settings.
use std::{sync::Arc, time::Duration};
use tokio::sync::RwLock;
use crate::{store::Store, Result}; 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> { pub async fn client_builder(store: &Store) -> Result<reqwest::ClientBuilder> {
let settings = store.proxy_settings_secret().await?; let settings = store.proxy_settings_secret().await?;
// Use the platform TLS stack for compatibility with provider gateways that // Use the platform TLS stack for compatibility with provider gateways that
+5 -1
View File
@@ -21,6 +21,7 @@ use super::{
pub struct ProviderRouter { pub struct ProviderRouter {
store: Store, store: Store,
plugins: PluginRegistry, plugins: PluginRegistry,
clients: crate::network::NetworkClients,
request_timeout: Duration, request_timeout: Duration,
stream_idle_timeout: Duration, stream_idle_timeout: Duration,
} }
@@ -29,12 +30,14 @@ impl ProviderRouter {
pub fn new( pub fn new(
store: Store, store: Store,
plugins: PluginRegistry, plugins: PluginRegistry,
clients: crate::network::NetworkClients,
request_timeout: Duration, request_timeout: Duration,
stream_idle_timeout: Duration, stream_idle_timeout: Duration,
) -> Self { ) -> Self {
Self { Self {
store, store,
plugins, plugins,
clients,
request_timeout, request_timeout,
stream_idle_timeout, stream_idle_timeout,
} }
@@ -49,6 +52,7 @@ impl Provider for ProviderRouter {
) -> ProviderStream { ) -> ProviderStream {
let store = self.store.clone(); let store = self.store.clone();
let plugins = self.plugins.clone(); let plugins = self.plugins.clone();
let clients = self.clients.clone();
let request_timeout = self.request_timeout; let request_timeout = self.request_timeout;
let stream_idle_timeout = self.stream_idle_timeout; let stream_idle_timeout = self.stream_idle_timeout;
Box::pin(try_stream! { Box::pin(try_stream! {
@@ -91,7 +95,7 @@ impl Provider for ProviderRouter {
request_timeout, request_timeout,
allowed_body_fields: None, 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)?; let provider = build_observed(&config, recorder.clone(), client)?;
(recorder, guard, provider.stream(routed, cancellation.clone())) (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] #[tokio::test]
async fn offline_crud_round_trip_persists_markdown() { async fn offline_crud_round_trip_persists_markdown() {
let (_store_dir, store) = fixtures::temp_store().await; 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_dir = tempfile::tempdir().unwrap();
let rules_root = rules_dir.path().join("rules"); let rules_root = rules_dir.path().join("rules");
let service = KnowledgeService::with_root(rules_root.clone()).unwrap(); let service = KnowledgeService::with_root(rules_root.clone()).unwrap();
@@ -193,7 +193,7 @@ async fn offline_crud_round_trip_persists_markdown() {
#[tokio::test] #[tokio::test]
async fn updating_missing_rule_reports_failure() { async fn updating_missing_rule_reports_failure() {
let (_store_dir, store) = fixtures::temp_store().await; 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_dir = tempfile::tempdir().unwrap();
let service = KnowledgeService::with_root(rules_dir.path().join("rules")).unwrap(); let service = KnowledgeService::with_root(rules_dir.path().join("rules")).unwrap();