plugin runtime

This commit is contained in:
leookun
2026-08-30 13:35:03 +08:00
parent e5bfcdd202
commit 1609b57433
20 changed files with 1417 additions and 26 deletions
+7
View File
@@ -4,6 +4,7 @@ mod calls;
mod harness;
mod models;
mod overview;
mod plugins;
mod service;
mod settings;
@@ -136,6 +137,12 @@ pub fn api_router(service: ControlService) -> Router {
)
.route("/__byok-api__/api/llm-calls", get(calls::list))
.route("/__byok-api__/api/llm-calls/{call_id}", get(calls::detail))
.route(
"/__byok-api__/api/plugins/runtime",
get(plugins::runtime_status)
.post(plugins::initialize_runtime)
.delete(plugins::cancel_runtime_initialization),
)
.route(
"/__byok-api__/api/settings/observability",
get(settings::get).put(settings::update),
+24
View File
@@ -0,0 +1,24 @@
//! Exposes plugin runtime initialization and status endpoints.
use axum::{extract::State, Json};
use crate::{plugin::PluginRuntimeStatus, Result};
use super::ControlService;
pub async fn runtime_status(
State(service): State<ControlService>,
) -> Result<Json<PluginRuntimeStatus>> {
Ok(Json(service.plugin_runtime_status()))
}
pub async fn initialize_runtime(
State(service): State<ControlService>,
) -> Result<Json<PluginRuntimeStatus>> {
Ok(Json(service.initialize_plugin_runtime()))
}
pub async fn cancel_runtime_initialization(
State(service): State<ControlService>,
) -> Result<Json<PluginRuntimeStatus>> {
Ok(Json(service.cancel_plugin_runtime_initialization()))
}
+15
View File
@@ -25,6 +25,7 @@ use crate::{
ModelRequest, ModelSpec, ModelType, Overview, ProjectedContent, ProjectedMessage,
PromptSpec, ProviderType, Role,
},
plugin::{PluginRuntime, PluginRuntimeStatus},
provider::{is_valid_response_event, ModelEvent, Provider},
store::{
DesktopSettings, PortSettings, ProxySettings, ProxySettingsInput, StatisticsStorage, Store,
@@ -38,6 +39,7 @@ pub struct ControlService {
store: Store,
cursor_harness: CursorHarness,
provider: Arc<dyn Provider>,
plugin_runtime: PluginRuntime,
model_tests: Arc<Mutex<BTreeMap<String, CancellationToken>>>,
}
@@ -148,6 +150,7 @@ impl ControlService {
cursor_harness: CursorHarness::new(store.clone())?,
store,
provider,
plugin_runtime: PluginRuntime::managed()?,
model_tests: Arc::new(Mutex::new(BTreeMap::new())),
})
}
@@ -156,6 +159,18 @@ impl ControlService {
&self.cursor_harness
}
pub fn plugin_runtime_status(&self) -> PluginRuntimeStatus {
self.plugin_runtime.status()
}
pub fn initialize_plugin_runtime(&self) -> PluginRuntimeStatus {
self.plugin_runtime.initialize(self.store.clone())
}
pub fn cancel_plugin_runtime_initialization(&self) -> PluginRuntimeStatus {
self.plugin_runtime.cancel_initialization()
}
pub(super) async fn ads(
&self,
disabled_ad_ids: Option<&str>,
+1
View File
@@ -8,6 +8,7 @@ pub mod error;
pub mod local_app;
pub mod model;
pub mod network;
pub mod plugin;
pub mod provider;
pub mod run;
pub mod search;
+88
View File
@@ -0,0 +1,88 @@
//! Maps supported desktop platforms to pinned Deno release assets.
pub(super) const DENO_VERSION: &str = "2.9.6";
#[derive(Clone, Copy, Debug)]
pub(super) struct RuntimeAsset {
pub target: &'static str,
pub sha256: &'static str,
}
impl RuntimeAsset {
pub fn current() -> Option<Self> {
Self::for_platform(std::env::consts::OS, std::env::consts::ARCH)
}
pub(super) fn for_platform(os: &str, arch: &str) -> Option<Self> {
let (target, sha256) = match (os, arch) {
("macos", "aarch64") => (
"aarch64-apple-darwin",
"213a2f304f04d3c9cb5220669afad138f60a5aab1fe80962abdeb8f35807a472",
),
("macos", "x86_64") => (
"x86_64-apple-darwin",
"7d4524b82bcc557fe020a1a5b56956ed42b992ae5b28026e8ad5d17329533f5f",
),
("windows", "aarch64") => (
"aarch64-pc-windows-msvc",
"acb014afe2299847764e232b4993e162e3946cdeec36603e3f1a0b548cd1ea55",
),
("windows", "x86_64") => (
"x86_64-pc-windows-msvc",
"15e5300b0ba3c3695a7621d90160a746ec9e710228cee639afa9d580f6e3cd11",
),
("linux", "aarch64") => (
"aarch64-unknown-linux-gnu",
"9a46afc6c392c7cd2ff71a31558935545b46408d0e87f7a86908c712721c046e",
),
("linux", "x86_64") => (
"x86_64-unknown-linux-gnu",
"394f07f4da2bebe6ce6f1e7ce0fa16429b29b08c35e3fac3fe25972676dff4b2",
),
_ => return None,
};
Some(Self { target, sha256 })
}
pub fn archive_name(self) -> String {
format!("deno-{}.zip", self.target)
}
pub fn download_url(self) -> String {
format!(
"https://github.com/denoland/deno/releases/download/v{DENO_VERSION}/{}",
self.archive_name()
)
}
pub fn executable_name(self) -> &'static str {
if self.target.contains("windows") {
"deno.exe"
} else {
"deno"
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn maps_every_supported_desktop_target() {
let cases = [
("macos", "aarch64", "aarch64-apple-darwin"),
("macos", "x86_64", "x86_64-apple-darwin"),
("windows", "aarch64", "aarch64-pc-windows-msvc"),
("windows", "x86_64", "x86_64-pc-windows-msvc"),
("linux", "aarch64", "aarch64-unknown-linux-gnu"),
("linux", "x86_64", "x86_64-unknown-linux-gnu"),
];
for (os, arch, expected) in cases {
assert_eq!(
RuntimeAsset::for_platform(os, arch).unwrap().target,
expected
);
}
assert!(RuntimeAsset::for_platform("linux", "x86").is_none());
}
}
+268
View File
@@ -0,0 +1,268 @@
//! Downloads, verifies, extracts, and validates a pinned Deno runtime.
use std::{
io,
path::{Path, PathBuf},
process::Stdio,
time::Duration,
};
use futures_util::StreamExt;
use sha2::{Digest, Sha256};
use tokio::io::AsyncWriteExt;
use tokio_util::sync::CancellationToken;
use super::{
asset::{RuntimeAsset, DENO_VERSION},
runtime::PluginRuntimePhase,
};
use crate::{network, store::Store, Error, Result};
const DOWNLOAD_TIMEOUT: Duration = Duration::from_secs(10 * 60);
const VALIDATION_TIMEOUT: Duration = Duration::from_secs(15);
const MAX_ARCHIVE_BYTES: u64 = 128 * 1024 * 1024;
pub(super) async fn install(
root: &Path,
store: &Store,
asset: RuntimeAsset,
cancellation: CancellationToken,
on_progress: impl Fn(PluginRuntimePhase, u64, Option<u64>),
) -> Result<()> {
let paths = RuntimePaths::new(root, asset);
tokio::fs::create_dir_all(&paths.download_dir).await?;
tokio::fs::create_dir_all(&paths.install_dir).await?;
remove_if_exists(&paths.archive).await?;
remove_if_exists(&paths.executable_staging).await?;
remove_if_exists(&paths.ready_marker).await?;
let result = download_and_install(store, asset, &paths, &cancellation, &on_progress).await;
if result.is_err() {
let _ = remove_if_exists(&paths.archive).await;
let _ = remove_if_exists(&paths.executable_staging).await;
let _ = remove_if_exists(&paths.executable).await;
let _ = remove_if_exists(&paths.ready_marker).await;
}
result
}
pub(super) fn runtime_complete(root: &Path, asset: RuntimeAsset) -> bool {
let paths = RuntimePaths::new(root, asset);
paths.executable.is_file() && paths.ready_marker.is_file()
}
async fn download_and_install(
store: &Store,
asset: RuntimeAsset,
paths: &RuntimePaths,
cancellation: &CancellationToken,
on_progress: &impl Fn(PluginRuntimePhase, u64, Option<u64>),
) -> Result<()> {
ensure_not_cancelled(cancellation)?;
let client = network::client(store).await?;
let response = tokio::select! {
_ = cancellation.cancelled() => return Err(Error::Cancelled),
response = client
.get(asset.download_url())
.timeout(DOWNLOAD_TIMEOUT)
.send() => response?,
}
.error_for_status()?;
let total_bytes = response.content_length();
if total_bytes.is_some_and(|size| size > MAX_ARCHIVE_BYTES) {
return Err(Error::Config(
"Deno runtime archive is larger than allowed".into(),
));
}
on_progress(PluginRuntimePhase::Downloading, 0, total_bytes);
let mut archive = tokio::fs::File::create(&paths.archive).await?;
let mut hasher = Sha256::new();
let mut downloaded_bytes = 0_u64;
let mut stream = response.bytes_stream();
loop {
let next = tokio::select! {
_ = cancellation.cancelled() => return Err(Error::Cancelled),
next = stream.next() => next,
};
let Some(chunk) = next else { break };
let chunk = chunk?;
downloaded_bytes = downloaded_bytes.saturating_add(chunk.len() as u64);
if downloaded_bytes > MAX_ARCHIVE_BYTES {
return Err(Error::Config(
"Deno runtime archive is larger than allowed".into(),
));
}
archive.write_all(&chunk).await?;
hasher.update(&chunk);
on_progress(
PluginRuntimePhase::Downloading,
downloaded_bytes,
total_bytes,
);
}
archive.flush().await?;
archive.sync_all().await?;
drop(archive);
ensure_not_cancelled(cancellation)?;
on_progress(PluginRuntimePhase::Verifying, downloaded_bytes, total_bytes);
let actual_hash = hex::encode(hasher.finalize());
if actual_hash != asset.sha256 {
return Err(Error::Config(format!(
"Deno runtime checksum mismatch: expected {}, received {actual_hash}",
asset.sha256
)));
}
on_progress(
PluginRuntimePhase::Installing,
downloaded_bytes,
total_bytes,
);
let archive_path = paths.archive.clone();
let staging_path = paths.executable_staging.clone();
let executable_name = asset.executable_name();
tokio::task::spawn_blocking(move || {
extract_runtime_archive(&archive_path, &staging_path, executable_name)
})
.await
.map_err(|error| Error::Config(format!("Deno extraction task failed: {error}")))??;
ensure_not_cancelled(cancellation)?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
tokio::fs::set_permissions(
&paths.executable_staging,
std::fs::Permissions::from_mode(0o700),
)
.await?;
}
remove_if_exists(&paths.executable).await?;
tokio::fs::rename(&paths.executable_staging, &paths.executable).await?;
ensure_not_cancelled(cancellation)?;
on_progress(
PluginRuntimePhase::Validating,
downloaded_bytes,
total_bytes,
);
validate_runtime(&paths.executable, cancellation).await?;
ensure_not_cancelled(cancellation)?;
tokio::fs::write(&paths.ready_marker, format!("deno {DENO_VERSION}\n")).await?;
remove_if_exists(&paths.archive).await?;
tracing::info!(
version = DENO_VERSION,
target = asset.target,
path = %paths.executable.display(),
"plugin runtime initialized"
);
Ok(())
}
struct RuntimePaths {
download_dir: PathBuf,
install_dir: PathBuf,
archive: PathBuf,
executable: PathBuf,
executable_staging: PathBuf,
ready_marker: PathBuf,
}
impl RuntimePaths {
fn new(root: &Path, asset: RuntimeAsset) -> Self {
let download_dir = root.join(".downloads");
let install_dir = root
.join("deno")
.join(format!("v{DENO_VERSION}"))
.join(asset.target);
let executable = install_dir.join(asset.executable_name());
Self {
archive: download_dir.join(format!("{}.part", asset.archive_name())),
executable_staging: install_dir.join(format!("{}.part", asset.executable_name())),
ready_marker: install_dir.join(".ready"),
download_dir,
install_dir,
executable,
}
}
}
fn extract_runtime_archive(archive: &Path, output: &Path, executable_name: &str) -> Result<()> {
let file = std::fs::File::open(archive)?;
let mut archive = zip::ZipArchive::new(file)
.map_err(|error| Error::Config(format!("invalid Deno runtime archive: {error}")))?;
let mut executable = archive
.by_name(executable_name)
.map_err(|error| Error::Config(format!("Deno executable missing from archive: {error}")))?;
let mut destination = std::fs::File::create(output)?;
io::copy(&mut executable, &mut destination)?;
destination.sync_all()?;
Ok(())
}
async fn validate_runtime(executable: &Path, cancellation: &CancellationToken) -> Result<()> {
let mut command = tokio::process::Command::new(executable);
command
.arg("--version")
.stdin(Stdio::null())
.stderr(Stdio::piped())
.stdout(Stdio::piped())
.kill_on_drop(true);
let output = tokio::select! {
_ = cancellation.cancelled() => return Err(Error::Cancelled),
result = tokio::time::timeout(VALIDATION_TIMEOUT, command.output()) => {
result.map_err(|_| Error::Config("Deno runtime validation timed out".into()))??
}
};
if !output.status.success() {
return Err(Error::Config(format!(
"Deno runtime validation failed: {}",
String::from_utf8_lossy(&output.stderr).trim()
)));
}
let expected = format!("deno {DENO_VERSION}");
let stdout = String::from_utf8_lossy(&output.stdout);
let version_line = stdout.lines().next().unwrap_or_default().trim();
if version_line != expected && !version_line.starts_with(&format!("{expected} ")) {
return Err(Error::Config(format!(
"unexpected Deno runtime version: {version_line}"
)));
}
Ok(())
}
fn ensure_not_cancelled(cancellation: &CancellationToken) -> Result<()> {
if cancellation.is_cancelled() {
Err(Error::Cancelled)
} else {
Ok(())
}
}
async fn remove_if_exists(path: &Path) -> Result<()> {
match tokio::fs::remove_file(path).await {
Ok(()) => Ok(()),
Err(error) if error.kind() == io::ErrorKind::NotFound => Ok(()),
Err(error) => Err(error.into()),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn uses_versioned_runtime_directory() {
let root = PathBuf::from("/tmp/plugin-runtime");
let asset = super::super::asset::RuntimeAsset::for_platform("macos", "aarch64").unwrap();
let paths = RuntimePaths::new(&root, asset);
assert_eq!(
paths.executable,
root.join("deno")
.join(format!("v{DENO_VERSION}"))
.join(asset.target)
.join("deno")
);
}
}
+6
View File
@@ -0,0 +1,6 @@
//! Owns plugin runtime installation and lifecycle infrastructure.
mod asset;
mod installation;
mod runtime;
pub use runtime::{PluginRuntime, PluginRuntimePhase, PluginRuntimeState, PluginRuntimeStatus};
+221
View File
@@ -0,0 +1,221 @@
//! Tracks Deno runtime readiness and coordinates one initialization task.
use std::{
path::PathBuf,
sync::{
atomic::{AtomicBool, Ordering},
Arc,
},
};
use parking_lot::{Mutex, RwLock};
use serde::Serialize;
use tokio_util::sync::CancellationToken;
use super::{
asset::{RuntimeAsset, DENO_VERSION},
installation,
};
use crate::{config, store::Store, Error, Result};
#[derive(Clone)]
pub struct PluginRuntime {
inner: Arc<PluginRuntimeInner>,
}
struct PluginRuntimeInner {
root: PathBuf,
asset: Option<RuntimeAsset>,
status: RwLock<PluginRuntimeStatus>,
initializing: AtomicBool,
cancellation: Mutex<Option<CancellationToken>>,
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum PluginRuntimeState {
Uninitialized,
Initializing,
Ready,
Failed,
Unsupported,
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum PluginRuntimePhase {
Checking,
Downloading,
Verifying,
Installing,
Validating,
}
#[derive(Clone, Debug, Serialize, PartialEq, Eq)]
pub struct PluginRuntimeStatus {
pub state: PluginRuntimeState,
pub version: String,
pub target: Option<String>,
pub phase: Option<PluginRuntimePhase>,
pub downloaded_bytes: u64,
pub total_bytes: Option<u64>,
pub error: Option<String>,
}
impl PluginRuntimeStatus {
fn uninitialized(asset: RuntimeAsset) -> Self {
Self::new(PluginRuntimeState::Uninitialized, Some(asset))
}
fn ready(asset: RuntimeAsset) -> Self {
Self::new(PluginRuntimeState::Ready, Some(asset))
}
fn unsupported() -> Self {
Self {
state: PluginRuntimeState::Unsupported,
version: DENO_VERSION.into(),
target: None,
phase: None,
downloaded_bytes: 0,
total_bytes: None,
error: Some(format!(
"unsupported platform: {}/{}",
std::env::consts::OS,
std::env::consts::ARCH
)),
}
}
fn new(state: PluginRuntimeState, asset: Option<RuntimeAsset>) -> Self {
Self {
state,
version: DENO_VERSION.into(),
target: asset.map(|value| value.target.into()),
phase: None,
downloaded_bytes: 0,
total_bytes: None,
error: None,
}
}
}
impl PluginRuntime {
pub fn managed() -> Result<Self> {
Self::new(config::managed_data_dir()?.join("plugins").join("runtime"))
}
fn new(root: PathBuf) -> Result<Self> {
std::fs::create_dir_all(&root)?;
#[cfg(unix)]
{
use std::os::unix::fs::PermissionsExt;
std::fs::set_permissions(&root, std::fs::Permissions::from_mode(0o700))?;
}
let asset = RuntimeAsset::current();
let status = match asset {
Some(asset) if installation::runtime_complete(&root, asset) => {
PluginRuntimeStatus::ready(asset)
}
Some(asset) => PluginRuntimeStatus::uninitialized(asset),
None => PluginRuntimeStatus::unsupported(),
};
Ok(Self {
inner: Arc::new(PluginRuntimeInner {
root,
asset,
status: RwLock::new(status),
initializing: AtomicBool::new(false),
cancellation: Mutex::new(None),
}),
})
}
pub fn status(&self) -> PluginRuntimeStatus {
let mut status = self.inner.status.write();
if status.state == PluginRuntimeState::Ready {
if let Some(asset) = self.inner.asset {
if !installation::runtime_complete(&self.inner.root, asset) {
*status = PluginRuntimeStatus::uninitialized(asset);
}
}
}
status.clone()
}
pub fn initialize(&self, store: Store) -> PluginRuntimeStatus {
let Some(asset) = self.inner.asset else {
return self.status();
};
if installation::runtime_complete(&self.inner.root, asset) {
let ready = PluginRuntimeStatus::ready(asset);
*self.inner.status.write() = ready.clone();
return ready;
}
if self
.inner
.initializing
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return self.status();
}
let mut initializing =
PluginRuntimeStatus::new(PluginRuntimeState::Initializing, Some(asset));
initializing.phase = Some(PluginRuntimePhase::Checking);
*self.inner.status.write() = initializing.clone();
let cancellation = CancellationToken::new();
*self.inner.cancellation.lock() = Some(cancellation.clone());
let runtime = self.clone();
tokio::spawn(async move {
let result = installation::install(
&runtime.inner.root,
&store,
asset,
cancellation,
|phase, downloaded, total| {
runtime.update_progress(phase, downloaded, total);
},
)
.await;
let status = match result {
Ok(()) => PluginRuntimeStatus::ready(asset),
Err(Error::Cancelled) => PluginRuntimeStatus::uninitialized(asset),
Err(error) => {
tracing::error!(%error, target = asset.target, "plugin runtime initialization failed");
let mut failed =
PluginRuntimeStatus::new(PluginRuntimeState::Failed, Some(asset));
failed.error = Some("plugin runtime initialization failed".into());
failed
}
};
*runtime.inner.status.write() = status;
runtime.inner.cancellation.lock().take();
runtime.inner.initializing.store(false, Ordering::Release);
});
initializing
}
pub fn cancel_initialization(&self) -> PluginRuntimeStatus {
if let Some(cancellation) = self.inner.cancellation.lock().as_ref() {
cancellation.cancel();
}
self.status()
}
fn update_progress(
&self,
phase: PluginRuntimePhase,
downloaded_bytes: u64,
total_bytes: Option<u64>,
) {
let mut status = self.inner.status.write();
status.state = PluginRuntimeState::Initializing;
status.phase = Some(phase);
status.downloaded_bytes = downloaded_bytes;
status.total_bytes = total_bytes;
status.error = None;
}
}