Implement MCP automation and approvals.

This commit is contained in:
2026-07-19 15:10:43 +02:00
parent fdfae200a0
commit fb5cae2131
35 changed files with 3654 additions and 80 deletions

View File

@@ -17,8 +17,9 @@ mod tests {
use crate::db::schema::{
ai_catalog_meta, ai_model_modalities, ai_models, ai_providers, chat_conversations,
chat_messages, db_notifications, dismissed_duplicate_pairs, embedding_keys,
generated_file_hashes, import_definitions, media, media_translations, post_links,
post_media, post_translations, posts, projects, scripts, settings, tags, templates,
generated_file_hashes, import_definitions, mcp_proposals, media, media_translations,
post_links, post_media, post_translations, posts, projects, scripts, settings, tags,
templates,
};
use diesel::prelude::*;
use diesel_migrations::MigrationHarness;
@@ -35,7 +36,7 @@ mod tests {
let applied = db
.conn()
.with_migrations(|conn| conn.applied_migrations().unwrap().len());
assert_eq!(applied, 3);
assert_eq!(applied, 4);
}
#[test]
@@ -75,6 +76,7 @@ mod tests {
embedding_keys::table,
dismissed_duplicate_pairs::table,
import_definitions::table,
mcp_proposals::table,
db_notifications::table,
);
}

View File

@@ -0,0 +1,153 @@
use diesel::prelude::*;
use crate::db::DbConnection;
use crate::db::schema::mcp_proposals;
use crate::model::{McpProposal, ProposalStatus};
pub fn insert_proposal(conn: &DbConnection, proposal: &McpProposal) -> QueryResult<()> {
conn.with(|connection| {
diesel::insert_into(mcp_proposals::table)
.values(proposal)
.execute(connection)
.map(|_| ())
})
}
pub fn get_proposal(conn: &DbConnection, id: &str) -> QueryResult<McpProposal> {
conn.with(|connection| {
mcp_proposals::table
.filter(mcp_proposals::id.eq(id))
.select(McpProposal::as_select())
.first(connection)
})
}
pub fn list_proposals(conn: &DbConnection, project_id: &str) -> QueryResult<Vec<McpProposal>> {
conn.with(|connection| {
mcp_proposals::table
.filter(mcp_proposals::project_id.eq(project_id))
.order(mcp_proposals::created_at.desc())
.select(McpProposal::as_select())
.load(connection)
})
}
pub fn list_pending_proposals(
conn: &DbConnection,
project_id: &str,
) -> QueryResult<Vec<McpProposal>> {
conn.with(|connection| {
mcp_proposals::table
.filter(mcp_proposals::project_id.eq(project_id))
.filter(mcp_proposals::status.eq(ProposalStatus::Pending))
.order(mcp_proposals::created_at.asc())
.select(McpProposal::as_select())
.load(connection)
})
}
pub fn expire_pending(conn: &DbConnection, now: i64) -> QueryResult<usize> {
conn.with(|connection| {
diesel::update(
mcp_proposals::table
.filter(mcp_proposals::status.eq(ProposalStatus::Pending))
.filter(mcp_proposals::expires_at.le(now)),
)
.set((
mcp_proposals::status.eq(ProposalStatus::Expired),
mcp_proposals::resolved_at.eq(Some(now)),
mcp_proposals::result.eq(Some("{\"message\":\"expired\"}".to_string())),
))
.execute(connection)
})
}
pub fn claim_pending(conn: &DbConnection, id: &str, now: i64) -> QueryResult<bool> {
conn.with(|connection| {
diesel::update(
mcp_proposals::table
.filter(mcp_proposals::id.eq(id))
.filter(mcp_proposals::status.eq(ProposalStatus::Pending))
.filter(mcp_proposals::expires_at.gt(now)),
)
.set(mcp_proposals::status.eq(ProposalStatus::Executing))
.execute(connection)
.map(|changed| changed == 1)
})
}
pub fn resolve_claimed(
conn: &DbConnection,
id: &str,
status: ProposalStatus,
result: &str,
resolved_at: i64,
) -> QueryResult<bool> {
debug_assert!(matches!(
status,
ProposalStatus::Accepted | ProposalStatus::Rejected
));
conn.with(|connection| {
diesel::update(
mcp_proposals::table
.filter(mcp_proposals::id.eq(id))
.filter(mcp_proposals::status.eq(ProposalStatus::Executing)),
)
.set((
mcp_proposals::status.eq(status),
mcp_proposals::result.eq(Some(result.to_string())),
mcp_proposals::resolved_at.eq(Some(resolved_at)),
))
.execute(connection)
.map(|changed| changed == 1)
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::db::Database;
use crate::db::queries::project::{insert_project, make_test_project};
use crate::model::ProposalKind;
fn setup() -> Database {
let db = Database::open_in_memory().unwrap();
db.migrate().unwrap();
insert_project(db.conn(), &make_test_project("p1", "blog")).unwrap();
db
}
fn proposal(id: &str, expires_at: i64) -> McpProposal {
McpProposal {
id: id.into(),
project_id: "p1".into(),
kind: ProposalKind::DraftPost,
status: ProposalStatus::Pending,
entity_id: None,
data: "{}".into(),
result: None,
created_at: 1,
expires_at,
resolved_at: None,
}
}
#[test]
fn lifecycle_claims_once_and_expires_pending_rows() {
let db = setup();
insert_proposal(db.conn(), &proposal("p1", 10)).unwrap();
insert_proposal(db.conn(), &proposal("p2", 1)).unwrap();
assert_eq!(expire_pending(db.conn(), 5).unwrap(), 1);
assert!(claim_pending(db.conn(), "p1", 5).unwrap());
assert!(!claim_pending(db.conn(), "p1", 5).unwrap());
assert!(resolve_claimed(db.conn(), "p1", ProposalStatus::Accepted, "{}", 6).unwrap());
assert_eq!(
get_proposal(db.conn(), "p1").unwrap().status,
ProposalStatus::Accepted
);
assert_eq!(
get_proposal(db.conn(), "p2").unwrap().status,
ProposalStatus::Expired
);
}
}

View File

@@ -1,6 +1,7 @@
pub mod db_notification;
pub mod generated_file_hash;
pub mod import_definition;
pub mod mcp_proposal;
pub mod media;
pub mod media_translation;
pub mod post;

View File

@@ -141,6 +141,21 @@ diesel::table! {
}
}
diesel::table! {
mcp_proposals (id) {
id -> Text,
project_id -> Text,
kind -> Text,
status -> Text,
entity_id -> Nullable<Text>,
data -> Text,
result -> Nullable<Text>,
created_at -> BigInt,
expires_at -> BigInt,
resolved_at -> Nullable<BigInt>,
}
}
diesel::table! {
media (id) {
id -> Text,
@@ -319,6 +334,7 @@ diesel::joinable!(chat_messages -> chat_conversations (conversation_id));
diesel::joinable!(dismissed_duplicate_pairs -> projects (project_id));
diesel::joinable!(generated_file_hashes -> projects (project_id));
diesel::joinable!(import_definitions -> projects (project_id));
diesel::joinable!(mcp_proposals -> projects (project_id));
diesel::joinable!(media -> projects (project_id));
diesel::joinable!(media_translations -> media (translation_for));
diesel::joinable!(media_translations -> projects (project_id));
@@ -344,6 +360,7 @@ diesel::allow_tables_to_appear_in_same_query!(
embedding_keys,
generated_file_hashes,
import_definitions,
mcp_proposals,
media,
media_translations,
post_links,

View File

@@ -5,8 +5,8 @@ use diesel::sql_types::{Integer, Text};
use diesel::sqlite::{Sqlite, SqliteValue};
use crate::model::{
NotificationAction, NotificationEntity, PostStatus, ScriptKind, ScriptStatus, TemplateKind,
TemplateStatus,
NotificationAction, NotificationEntity, PostStatus, ProposalKind, ProposalStatus, ScriptKind,
ScriptStatus, TemplateKind, TemplateStatus,
};
#[derive(Debug, AsExpression, FromSqlRow)]
@@ -93,3 +93,5 @@ text_enum_sql!(ScriptKind);
text_enum_sql!(ScriptStatus);
text_enum_sql!(NotificationEntity);
text_enum_sql!(NotificationAction);
text_enum_sql!(ProposalKind);
text_enum_sql!(ProposalStatus);

View File

@@ -0,0 +1,182 @@
use std::path::{Path, PathBuf};
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value, json};
use crate::engine::{EngineError, EngineResult};
const SERVER_NAME: &str = "bDS";
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum McpAgent {
ClaudeCode,
GithubCopilot,
}
impl McpAgent {
pub const fn all() -> [Self; 2] {
[Self::ClaudeCode, Self::GithubCopilot]
}
pub const fn label(self) -> &'static str {
match self {
Self::ClaudeCode => "Claude Code",
Self::GithubCopilot => "GitHub Copilot",
}
}
pub const fn as_str(self) -> &'static str {
match self {
Self::ClaudeCode => "claude_code",
Self::GithubCopilot => "github_copilot",
}
}
}
pub fn agent_config_path(agent: McpAgent, home_dir: &Path) -> PathBuf {
match agent {
McpAgent::ClaudeCode => home_dir.join(".claude.json"),
McpAgent::GithubCopilot => {
#[cfg(target_os = "macos")]
let path = home_dir.join("Library/Application Support/Code/User/mcp.json");
#[cfg(target_os = "windows")]
let path = home_dir.join("AppData/Roaming/Code/User/mcp.json");
#[cfg(not(any(target_os = "macos", target_os = "windows")))]
let path = home_dir.join(".config/Code/User/mcp.json");
path
}
}
}
pub fn packaged_mcp_executable() -> EngineResult<PathBuf> {
let executable = std::env::current_exe()?;
let sibling = executable.with_file_name(if cfg!(windows) {
"bds-mcp.exe"
} else {
"bds-mcp"
});
if sibling.is_file() {
Ok(sibling)
} else {
Err(EngineError::NotFound(format!(
"packaged MCP executable {}",
sibling.display()
)))
}
}
pub fn is_agent_configured(agent: McpAgent, home_dir: &Path) -> bool {
read_config(&agent_config_path(agent, home_dir))
.ok()
.is_some_and(|config| {
server_map(&config, agent).is_some_and(|servers| servers.contains_key(SERVER_NAME))
})
}
pub fn install_agent_config(
agent: McpAgent,
home_dir: &Path,
executable: &Path,
) -> EngineResult<PathBuf> {
if !executable.is_file() {
return Err(EngineError::NotFound(format!(
"MCP executable {}",
executable.display()
)));
}
let path = agent_config_path(agent, home_dir);
let mut config = read_config(&path)?;
let key = server_key(agent);
let servers = config
.as_object_mut()
.ok_or_else(|| {
EngineError::Validation(format!("{} must contain a JSON object", path.display()))
})?
.entry(key)
.or_insert_with(|| Value::Object(Map::new()))
.as_object_mut()
.ok_or_else(|| {
EngineError::Validation(format!("{key} in {} must be an object", path.display()))
})?;
let executable = executable.to_string_lossy();
let server = match agent {
McpAgent::ClaudeCode => json!({"command": executable, "args": []}),
McpAgent::GithubCopilot => {
json!({"type": "stdio", "command": executable, "args": []})
}
};
servers.insert(SERVER_NAME.into(), server);
write_config(&path, &config)?;
Ok(path)
}
pub fn remove_agent_config(agent: McpAgent, home_dir: &Path) -> EngineResult<PathBuf> {
let path = agent_config_path(agent, home_dir);
let mut config = read_config(&path)?;
if let Some(servers) = config
.as_object_mut()
.and_then(|object| object.get_mut(server_key(agent)))
.and_then(Value::as_object_mut)
{
servers.remove(SERVER_NAME);
}
write_config(&path, &config)?;
Ok(path)
}
fn server_key(agent: McpAgent) -> &'static str {
match agent {
McpAgent::ClaudeCode => "mcpServers",
McpAgent::GithubCopilot => "servers",
}
}
fn server_map(config: &Value, agent: McpAgent) -> Option<&Map<String, Value>> {
config.get(server_key(agent)).and_then(Value::as_object)
}
fn read_config(path: &Path) -> EngineResult<Value> {
match std::fs::read_to_string(path) {
Ok(source) => serde_json::from_str(&source)
.map_err(|error| EngineError::Parse(format!("{}: {error}", path.display()))),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(json!({})),
Err(error) => Err(error.into()),
}
}
fn write_config(path: &Path, config: &Value) -> EngineResult<()> {
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent)?;
}
let source = serde_json::to_string_pretty(config)?;
crate::util::atomic_write_str(path, &source)?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn install_and_remove_preserve_unrelated_config_without_secrets() {
let root = tempfile::tempdir().unwrap();
let executable = root.path().join("bds-mcp");
std::fs::write(&executable, "binary").unwrap();
for agent in McpAgent::all() {
let path = agent_config_path(agent, root.path());
if let Some(parent) = path.parent() {
std::fs::create_dir_all(parent).unwrap();
}
std::fs::write(&path, r#"{"unrelated":{"token":"kept"}}"#).unwrap();
install_agent_config(agent, root.path(), &executable).unwrap();
assert!(is_agent_configured(agent, root.path()));
let source = std::fs::read_to_string(&path).unwrap();
assert!(source.contains("kept"));
assert!(!source.contains("api_key"));
remove_agent_config(agent, root.path()).unwrap();
assert!(!is_agent_configured(agent, root.path()));
assert!(std::fs::read_to_string(path).unwrap().contains("kept"));
}
}
}

View File

@@ -0,0 +1,234 @@
use std::net::{Ipv4Addr, SocketAddr};
use std::path::PathBuf;
use std::thread::JoinHandle;
use axum::Router;
use axum::extract::State;
use axum::http::{HeaderMap, HeaderValue, StatusCode, header};
use axum::response::{IntoResponse, Response};
use axum::routing::post;
use serde_json::Value;
use crate::engine::{EngineError, EngineResult};
use super::McpContext;
use super::protocol::{MCP_PROTOCOL_VERSION, error, handle_rpc};
pub struct McpHttpServer {
address: SocketAddr,
shutdown: Option<tokio::sync::oneshot::Sender<()>>,
thread: Option<JoinHandle<()>>,
}
impl McpHttpServer {
pub fn start(database_path: PathBuf, port: u16) -> EngineResult<Self> {
McpContext::new(database_path.clone()).prepare()?;
let listener = std::net::TcpListener::bind((Ipv4Addr::LOCALHOST, port))?;
listener.set_nonblocking(true)?;
let address = listener.local_addr()?;
let (shutdown, shutdown_rx) = tokio::sync::oneshot::channel();
let thread = std::thread::Builder::new()
.name("bds-mcp-http".into())
.spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("MCP Tokio runtime");
runtime.block_on(async move {
let listener =
tokio::net::TcpListener::from_std(listener).expect("MCP loopback listener");
let context = McpContext::new(database_path);
let router = Router::new()
.route("/mcp", post(post_mcp).options(options_mcp))
.with_state(context);
let _ = axum::serve(listener, router)
.with_graceful_shutdown(async {
let _ = shutdown_rx.await;
})
.await;
});
})?;
Ok(Self {
address,
shutdown: Some(shutdown),
thread: Some(thread),
})
}
pub fn address(&self) -> SocketAddr {
self.address
}
pub fn endpoint(&self) -> String {
format!("http://{}/mcp", self.address)
}
pub fn stop(mut self) -> EngineResult<()> {
self.shutdown.take();
if let Some(thread) = self.thread.take() {
thread
.join()
.map_err(|_| EngineError::Parse("MCP server thread panicked".into()))?;
}
Ok(())
}
}
impl Drop for McpHttpServer {
fn drop(&mut self) {
self.shutdown.take();
if let Some(thread) = self.thread.take() {
let _ = thread.join();
}
}
}
async fn post_mcp(
State(context): State<McpContext>,
headers: HeaderMap,
body: axum::body::Bytes,
) -> Response {
if let Err((status, message)) = validate_http_headers(&headers) {
return with_cors((status, message).into_response(), &headers);
}
let request = match serde_json::from_slice::<Value>(&body) {
Ok(request) => request,
Err(_) => {
return with_cors(
(
StatusCode::BAD_REQUEST,
axum::Json(error(Value::Null, -32700, "Parse error")),
)
.into_response(),
&headers,
);
}
};
let response = match handle_rpc(&context, &request) {
Some(response) => (StatusCode::OK, axum::Json(response)).into_response(),
None => StatusCode::ACCEPTED.into_response(),
};
with_cors(response, &headers)
}
async fn options_mcp(headers: HeaderMap) -> Response {
if let Err((status, message)) = validate_origin_and_host(&headers) {
return with_cors((status, message).into_response(), &headers);
}
with_cors(StatusCode::NO_CONTENT.into_response(), &headers)
}
fn validate_http_headers(headers: &HeaderMap) -> Result<(), (StatusCode, &'static str)> {
validate_origin_and_host(headers)?;
if let Some(version) = headers
.get("mcp-protocol-version")
.and_then(|value| value.to_str().ok())
&& !["2025-03-26", MCP_PROTOCOL_VERSION].contains(&version)
{
return Err((StatusCode::BAD_REQUEST, "Unsupported MCP protocol version"));
}
if headers
.get(header::CONTENT_TYPE)
.and_then(|value| value.to_str().ok())
.is_none_or(|value| !value.starts_with("application/json"))
{
return Err((
StatusCode::UNSUPPORTED_MEDIA_TYPE,
"Expected application/json",
));
}
if headers
.get(header::ACCEPT)
.and_then(|value| value.to_str().ok())
.is_none_or(|value| {
!value.contains("application/json") || !value.contains("text/event-stream")
})
{
return Err((
StatusCode::NOT_ACCEPTABLE,
"Accept must include application/json and text/event-stream",
));
}
Ok(())
}
fn validate_origin_and_host(headers: &HeaderMap) -> Result<(), (StatusCode, &'static str)> {
let host = headers
.get(header::HOST)
.and_then(|value| value.to_str().ok())
.unwrap_or_default();
let host_name = host
.strip_prefix('[')
.and_then(|value| value.split_once(']').map(|(host, _)| host))
.unwrap_or_else(|| host.split(':').next().unwrap_or_default());
if !["localhost", "127.0.0.1", "::1"].contains(&host_name) {
return Err((StatusCode::FORBIDDEN, "Forbidden host"));
}
if let Some(origin) = headers
.get(header::ORIGIN)
.and_then(|value| value.to_str().ok())
{
let local = url::Url::parse(origin).ok().is_some_and(|origin| {
matches!(origin.host_str(), Some("localhost" | "127.0.0.1" | "::1"))
});
if !local {
return Err((StatusCode::FORBIDDEN, "Forbidden origin"));
}
}
Ok(())
}
fn with_cors(mut response: Response, request_headers: &HeaderMap) -> Response {
let headers = response.headers_mut();
let origin = request_headers
.get(header::ORIGIN)
.filter(|value| {
value.to_str().ok().is_some_and(|origin| {
url::Url::parse(origin).ok().is_some_and(|origin| {
matches!(origin.host_str(), Some("localhost" | "127.0.0.1" | "::1"))
})
})
})
.cloned()
.unwrap_or_else(|| HeaderValue::from_static("http://127.0.0.1"));
headers.insert(header::ACCESS_CONTROL_ALLOW_ORIGIN, origin);
headers.insert(
header::ACCESS_CONTROL_ALLOW_METHODS,
HeaderValue::from_static("POST, OPTIONS"),
);
headers.insert(
header::ACCESS_CONTROL_ALLOW_HEADERS,
HeaderValue::from_static("content-type, accept, origin, mcp-protocol-version"),
);
headers.insert(header::VARY, HeaderValue::from_static("Origin"));
response
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn header_validation_rejects_dns_rebinding_and_remote_origins() {
let mut headers = HeaderMap::new();
headers.insert(header::HOST, HeaderValue::from_static("attacker.example"));
assert!(validate_origin_and_host(&headers).is_err());
headers.insert(header::HOST, HeaderValue::from_static("127.0.0.1:4124"));
headers.insert(
header::ORIGIN,
HeaderValue::from_static("https://attacker.example"),
);
assert!(validate_origin_and_host(&headers).is_err());
headers.insert(
header::ORIGIN,
HeaderValue::from_static("http://localhost:3000"),
);
assert!(validate_origin_and_host(&headers).is_ok());
}
#[test]
fn router_only_allows_post_and_options() {
assert_eq!(axum::http::Method::POST.as_str(), "POST");
assert_eq!(axum::http::Method::OPTIONS.as_str(), "OPTIONS");
}
}

View File

@@ -0,0 +1,379 @@
mod agent_config;
mod http;
mod protocol;
mod resources;
mod tools;
use std::path::{Path, PathBuf};
use serde_json::Value;
use uuid::Uuid;
use crate::db::queries::{mcp_proposal as proposal_q, project as project_q};
use crate::db::{Database, DbConnection};
use crate::engine::{EngineError, EngineResult, cli_sync, domain_events};
use crate::model::DomainEvent;
use crate::util::now_unix_ms;
pub use crate::model::{McpProposal, ProposalKind, ProposalStatus};
pub use agent_config::{
McpAgent, agent_config_path, install_agent_config, is_agent_configured,
packaged_mcp_executable, remove_agent_config,
};
pub use http::McpHttpServer;
pub use protocol::{MCP_PROTOCOL_VERSION, handle_rpc};
pub use resources::ResourceContent;
pub const DEFAULT_HTTP_PORT: u16 = 4124;
pub const PROPOSAL_TTL_MS: i64 = 30 * 60 * 1_000;
pub const PROPOSALS_EVENT_KEY: &str = "mcp.proposals";
#[derive(Debug, Clone)]
pub struct McpContext {
database_path: PathBuf,
}
impl McpContext {
pub fn new(database_path: PathBuf) -> Self {
Self { database_path }
}
pub fn database_path(&self) -> &Path {
&self.database_path
}
/// Prepare shared storage once when a transport starts. Individual
/// stateless read requests never run migrations or repair derived state.
pub fn prepare(&self) -> EngineResult<()> {
let db = Database::open(&self.database_path)?;
db.migrate()
.map_err(|error| EngineError::Parse(error.to_string()))?;
crate::engine::search::prepare_search_index(db.conn())?;
Ok(())
}
pub fn list_resources(&self) -> Vec<Value> {
resources::list()
}
pub fn list_resource_templates(&self) -> Vec<Value> {
resources::templates()
}
pub fn read_resource(&self, uri: &str) -> EngineResult<resources::ResourceContent> {
let db = self.open_database()?;
resources::read(db.conn(), uri)
}
pub fn list_tools(&self) -> Vec<Value> {
tools::list()
}
pub fn call_tool(&self, name: &str, params: Value) -> EngineResult<Value> {
let db = self.open_database()?;
tools::call(db.conn(), name, params)
}
pub(crate) fn open_database(&self) -> EngineResult<Database> {
Ok(Database::open(&self.database_path)?)
}
}
pub fn get_proposal(conn: &DbConnection, proposal_id: &str) -> EngineResult<McpProposal> {
proposal_q::get_proposal(conn, proposal_id)
.map_err(|_| EngineError::NotFound(format!("MCP proposal {proposal_id}")))
}
pub fn list_proposals(conn: &DbConnection, project_id: &str) -> EngineResult<Vec<McpProposal>> {
expire_proposals(conn)?;
Ok(proposal_q::list_proposals(conn, project_id)?)
}
pub fn list_pending_proposals(
conn: &DbConnection,
project_id: &str,
) -> EngineResult<Vec<McpProposal>> {
expire_proposals(conn)?;
Ok(proposal_q::list_pending_proposals(conn, project_id)?)
}
pub fn expire_proposals(conn: &DbConnection) -> EngineResult<usize> {
let expired = proposal_q::expire_pending(conn, now_unix_ms())?;
if expired > 0 {
notify_proposals_changed();
}
Ok(expired)
}
pub(crate) fn create_proposal(
conn: &DbConnection,
kind: ProposalKind,
project_id: &str,
entity_id: Option<&str>,
data: &Value,
) -> EngineResult<McpProposal> {
expire_proposals(conn)?;
let now = now_unix_ms();
let proposal = McpProposal {
id: Uuid::new_v4().to_string(),
project_id: project_id.to_string(),
kind,
status: ProposalStatus::Pending,
entity_id: entity_id.map(str::to_string),
data: serde_json::to_string(data)?,
result: None,
created_at: now,
expires_at: now + PROPOSAL_TTL_MS,
resolved_at: None,
};
proposal_q::insert_proposal(conn, &proposal)?;
cli_sync::record_cli_event(
conn,
&DomainEvent::SettingsChanged {
project_id: None,
key: PROPOSALS_EVENT_KEY.to_string(),
},
)?;
Ok(proposal)
}
pub fn accept_proposal(
conn: &DbConnection,
data_dir: &Path,
proposal_id: &str,
) -> EngineResult<McpProposal> {
resolve_proposal(conn, data_dir, proposal_id, true)
}
pub fn reject_proposal(
conn: &DbConnection,
data_dir: &Path,
proposal_id: &str,
) -> EngineResult<McpProposal> {
resolve_proposal(conn, data_dir, proposal_id, false)
}
fn resolve_proposal(
conn: &DbConnection,
data_dir: &Path,
proposal_id: &str,
accept: bool,
) -> EngineResult<McpProposal> {
expire_proposals(conn)?;
conn.begin_savepoint()?;
let outcome = (|| {
let now = now_unix_ms();
if !proposal_q::claim_pending(conn, proposal_id, now)? {
let current = get_proposal(conn, proposal_id)?;
return Err(EngineError::Conflict(format!(
"MCP proposal {} is {}",
current.id,
current.status.as_str()
)));
}
let proposal = get_proposal(conn, proposal_id)?;
let result = if accept {
execute_proposal(conn, data_dir, &proposal)?
} else {
serde_json::json!({"message": "rejected"})
};
let status = if accept {
ProposalStatus::Accepted
} else {
ProposalStatus::Rejected
};
if !proposal_q::resolve_claimed(
conn,
proposal_id,
status,
&serde_json::to_string(&result)?,
now_unix_ms(),
)? {
return Err(EngineError::Conflict(format!(
"MCP proposal {proposal_id} was resolved concurrently"
)));
}
get_proposal(conn, proposal_id)
})();
match outcome {
Ok(proposal) => {
conn.release_savepoint()?;
notify_proposals_changed();
Ok(proposal)
}
Err(error) => {
let _ = conn.rollback_savepoint();
Err(error)
}
}
}
fn execute_proposal(
conn: &DbConnection,
data_dir: &Path,
proposal: &McpProposal,
) -> EngineResult<Value> {
let data: Value = serde_json::from_str(&proposal.data)?;
match proposal.kind {
ProposalKind::DraftPost => {
let post = crate::engine::post::create_post(
conn,
data_dir,
&proposal.project_id,
required_string(&data, "title")?,
Some(required_string(&data, "content")?),
string_array(&data, "tags"),
string_array(&data, "categories"),
optional_string(&data, "author"),
optional_string(&data, "language"),
None,
)?;
let post = if let Some(excerpt) = optional_string(&data, "excerpt") {
crate::engine::post::update_post(
conn,
data_dir,
&post.id,
None,
None,
Some(Some(excerpt)),
None,
None,
None,
None,
None,
None,
None,
)?
} else {
post
};
let post = crate::engine::post::publish_post(conn, data_dir, &post.id)?;
Ok(serde_json::to_value(post)?)
}
ProposalKind::ProposeScript => {
let kind = required_string(&data, "kind")?
.parse()
.map_err(EngineError::Validation)?;
let script = crate::engine::script::create_script(
conn,
&proposal.project_id,
required_string(&data, "title")?,
kind,
required_string(&data, "content")?,
optional_string(&data, "entrypoint"),
)?;
let script = crate::engine::script::publish_script(conn, data_dir, &script.id)?;
Ok(serde_json::to_value(script)?)
}
ProposalKind::ProposeTemplate => {
let kind = required_string(&data, "kind")?
.parse()
.map_err(EngineError::Validation)?;
let template = crate::engine::template::create_template(
conn,
&proposal.project_id,
required_string(&data, "title")?,
kind,
required_string(&data, "content")?,
)?;
let template = crate::engine::template::publish_template(conn, data_dir, &template.id)?;
Ok(serde_json::to_value(template)?)
}
ProposalKind::ProposeMediaTranslation => {
let translation = crate::engine::media::upsert_media_translation(
conn,
data_dir,
required_string(&data, "mediaId")?,
required_string(&data, "language")?,
optional_string(&data, "title"),
optional_string(&data, "alt"),
optional_string(&data, "caption"),
)?;
Ok(serde_json::to_value(translation)?)
}
ProposalKind::ProposeMediaMetadata => {
let media = crate::engine::media::update_media(
conn,
data_dir,
required_string(&data, "mediaId")?,
optional_optional_string(&data, "title"),
optional_optional_string(&data, "alt"),
optional_optional_string(&data, "caption"),
None,
None,
data.get("tags")
.is_some()
.then(|| string_array(&data, "tags")),
)?;
Ok(serde_json::to_value(media)?)
}
ProposalKind::ProposePostMetadata => {
let post = crate::engine::post::update_post(
conn,
data_dir,
required_string(&data, "postId")?,
optional_string(&data, "title"),
None,
optional_optional_string(&data, "excerpt"),
None,
data.get("tags")
.is_some()
.then(|| string_array(&data, "tags")),
data.get("categories")
.is_some()
.then(|| string_array(&data, "categories")),
None,
None,
None,
None,
)?;
Ok(serde_json::to_value(post)?)
}
}
}
pub(crate) fn active_project(
conn: &DbConnection,
) -> EngineResult<(crate::model::Project, PathBuf)> {
let project = project_q::get_active_project(conn)
.map_err(|_| EngineError::NotFound("active project".into()))?;
let data_dir = project
.data_path
.as_deref()
.map(PathBuf::from)
.ok_or_else(|| EngineError::Validation("active project has no data path".into()))?;
Ok((project, data_dir))
}
pub(crate) fn required_string<'a>(value: &'a Value, key: &str) -> EngineResult<&'a str> {
value
.get(key)
.and_then(Value::as_str)
.filter(|value| !value.trim().is_empty())
.ok_or_else(|| EngineError::Validation(format!("{key} is required")))
}
pub(crate) fn optional_string<'a>(value: &'a Value, key: &str) -> Option<&'a str> {
value.get(key).and_then(Value::as_str)
}
fn optional_optional_string<'a>(value: &'a Value, key: &str) -> Option<Option<&'a str>> {
value.get(key).map(|value| value.as_str())
}
pub(crate) fn string_array(value: &Value, key: &str) -> Vec<String> {
value
.get(key)
.and_then(Value::as_array)
.into_iter()
.flatten()
.filter_map(Value::as_str)
.map(str::trim)
.filter(|value| !value.is_empty())
.map(str::to_string)
.collect()
}
fn notify_proposals_changed() {
domain_events::settings_changed(None, PROPOSALS_EVENT_KEY);
}

View File

@@ -0,0 +1,130 @@
use serde_json::{Value, json};
use crate::engine::{EngineError, EngineResult};
use super::McpContext;
pub const MCP_PROTOCOL_VERSION: &str = "2025-06-18";
pub fn handle_rpc(context: &McpContext, request: &Value) -> Option<Value> {
let id = request.get("id").cloned();
let Some(method) = request.get("method").and_then(Value::as_str) else {
return Some(error(id.unwrap_or(Value::Null), -32600, "Invalid Request"));
};
if request.get("jsonrpc").and_then(Value::as_str) != Some("2.0") {
return Some(error(id.unwrap_or(Value::Null), -32600, "Invalid Request"));
}
let id = id?;
let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
let result = match method {
"initialize" => Ok(json!({
"protocolVersion": negotiated_version(&params),
"capabilities": {
"tools": {"listChanged": false},
"resources": {"subscribe": false, "listChanged": false}
},
"serverInfo": {
"name": "Blogging Desktop Server",
"version": env!("CARGO_PKG_VERSION")
}
})),
"ping" => Ok(json!({})),
"tools/list" => Ok(json!({"tools": context.list_tools()})),
"tools/call" => call_tool(context, &params),
"resources/list" => Ok(json!({"resources": context.list_resources()})),
"resources/templates/list" => Ok(json!({
"resourceTemplates": context.list_resource_templates()
})),
"resources/read" => read_resource(context, &params),
_ => return Some(error(id, -32601, "Method not found")),
};
Some(match result {
Ok(result) => success(id, result),
Err(EngineError::NotFound(message)) => error(id, -32004, &message),
Err(EngineError::Validation(message) | EngineError::Conflict(message)) => {
error(id, -32602, &message)
}
Err(error_value) => error(id, -32000, &error_value.to_string()),
})
}
fn negotiated_version(params: &Value) -> &str {
match params.get("protocolVersion").and_then(Value::as_str) {
Some("2025-03-26") => "2025-03-26",
Some("2025-06-18") => "2025-06-18",
_ => MCP_PROTOCOL_VERSION,
}
}
fn call_tool(context: &McpContext, params: &Value) -> EngineResult<Value> {
let name = params
.get("name")
.and_then(Value::as_str)
.filter(|name| !name.is_empty())
.ok_or_else(|| EngineError::Validation("tool name is required".into()))?;
let arguments = params
.get("arguments")
.cloned()
.unwrap_or_else(|| json!({}));
let result = context.call_tool(name, arguments)?;
Ok(json!({
"content": [{
"type": "text",
"text": serde_json::to_string(&result)?
}],
"structuredContent": result,
"isError": false
}))
}
fn read_resource(context: &McpContext, params: &Value) -> EngineResult<Value> {
let uri = params
.get("uri")
.and_then(Value::as_str)
.filter(|uri| !uri.is_empty())
.ok_or_else(|| EngineError::Validation("resource URI is required".into()))?;
let content = context.read_resource(uri)?;
let content = if let Some(blob) = content.blob {
json!({
"uri": content.uri,
"mimeType": content.mime_type,
"blob": blob
})
} else {
json!({
"uri": content.uri,
"mimeType": content.mime_type,
"text": content.text.unwrap_or_default()
})
};
Ok(json!({"contents": [content]}))
}
fn success(id: Value, result: Value) -> Value {
json!({"jsonrpc":"2.0","id":id,"result":result})
}
pub(crate) fn error(id: Value, code: i64, message: &str) -> Value {
json!({"jsonrpc":"2.0","id":id,"error":{"code":code,"message":message}})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn notifications_have_no_response_and_invalid_requests_are_rejected() {
let context = McpContext::new("missing.sqlite".into());
assert!(
handle_rpc(
&context,
&json!({"jsonrpc":"2.0","method":"notifications/initialized"})
)
.is_none()
);
assert_eq!(
handle_rpc(&context, &json!({"jsonrpc":"2.0","id":1})).unwrap()["error"]["code"],
-32600
);
}
}

View File

@@ -0,0 +1,382 @@
use std::path::Path;
use base64::Engine as _;
use serde_json::{Value, json};
use crate::db::DbConnection;
use crate::db::queries::{media as media_q, post as post_q, post_link, post_media, tag as tag_q};
use crate::engine::{EngineError, EngineResult};
use crate::model::{Media, Post};
use super::active_project;
const PAGE_SIZE: usize = 50;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ResourceContent {
pub uri: String,
pub mime_type: String,
pub text: Option<String>,
pub blob: Option<String>,
}
impl ResourceContent {
fn json(uri: &str, value: &Value) -> EngineResult<Self> {
Ok(Self {
uri: uri.to_string(),
mime_type: "application/json".into(),
text: Some(serde_json::to_string(value)?),
blob: None,
})
}
}
pub fn list() -> Vec<Value> {
[
("project", "Active project", "bds://project"),
("posts", "Blog posts", "bds://posts"),
("media", "Media", "bds://media"),
("tags", "Tags", "bds://tags"),
("categories", "Categories", "bds://categories"),
("stats", "Blog statistics", "bds://stats"),
]
.into_iter()
.map(|(name, title, uri)| {
json!({
"name": name,
"title": title,
"uri": uri,
"mimeType": "application/json"
})
})
.collect()
}
pub fn templates() -> Vec<Value> {
[
("posts", "Paginated blog posts", "bds://posts{?cursor}"),
("media", "Paginated media", "bds://media{?cursor}"),
(
"post media",
"Media linked to a post",
"bds://posts/{id}/media",
),
(
"media image",
"Original media bytes",
"bds://media/{id}/image",
),
]
.into_iter()
.map(|(name, title, uri_template)| {
json!({
"name": name,
"title": title,
"uriTemplate": uri_template
})
})
.collect()
}
pub fn read(conn: &DbConnection, uri: &str) -> EngineResult<ResourceContent> {
let url = url::Url::parse(uri)
.map_err(|_| EngineError::Validation("invalid MCP resource URI".into()))?;
if url.scheme() != "bds" {
return Err(EngineError::NotFound(uri.into()));
}
let host = url
.host_str()
.ok_or_else(|| EngineError::NotFound(uri.into()))?;
let path = url.path().trim_matches('/');
let (project, data_dir) = active_project(conn)?;
let value = match (host, path) {
("project", "") => project_resource(&project, &data_dir),
("posts", "") => posts_page(conn, &project.id, cursor_offset(&url)?),
("media", "") => media_page(conn, &project.id, cursor_offset(&url)?),
("tags", "") => tags(conn, &project.id),
("categories", "") => categories(conn, &project.id, &data_dir),
("stats", "") => stats(conn, &project.id, &data_dir),
("posts", path) => {
let parts = path.split('/').collect::<Vec<_>>();
match parts.as_slice() {
[post_id] => post_detail_by_id(conn, &project.id, &data_dir, post_id),
[post_id, "media"] => post_media_items(conn, &project.id, post_id),
_ => return Err(EngineError::NotFound(uri.into())),
}
}
("media", path) => {
let parts = path.split('/').collect::<Vec<_>>();
match parts.as_slice() {
[media_id] => media_detail_by_id(conn, &project.id, media_id),
[media_id, "image"] => {
return media_image(conn, &project.id, &data_dir, media_id, uri);
}
_ => return Err(EngineError::NotFound(uri.into())),
}
}
_ => return Err(EngineError::NotFound(uri.into())),
}?;
ResourceContent::json(uri, &value)
}
fn project_resource(project: &crate::model::Project, data_dir: &Path) -> EngineResult<Value> {
let metadata = crate::engine::meta::read_project_json(data_dir).ok();
Ok(json!({
"id": project.id,
"name": project.name,
"slug": project.slug,
"description": project.description,
"public_url": metadata.as_ref().and_then(|value| value.public_url.clone()),
"main_language": metadata.as_ref().and_then(|value| value.main_language.clone()),
"blog_languages": metadata.map(|value| value.blog_languages).unwrap_or_default()
}))
}
fn cursor_offset(url: &url::Url) -> EngineResult<usize> {
let Some(cursor) = url
.query_pairs()
.find_map(|(key, value)| (key == "cursor").then(|| value.into_owned()))
else {
return Ok(0);
};
if cursor.is_empty() {
return Err(EngineError::Validation("invalid cursor".into()));
}
let decoded = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(cursor)
.map_err(|_| EngineError::Validation("invalid cursor".into()))?;
let value: Value = serde_json::from_slice(&decoded)
.map_err(|_| EngineError::Validation("invalid cursor".into()))?;
value["offset"]
.as_u64()
.map(|offset| offset as usize)
.ok_or_else(|| EngineError::Validation("invalid cursor".into()))
}
fn encode_cursor(offset: usize) -> String {
base64::engine::general_purpose::URL_SAFE_NO_PAD
.encode(serde_json::to_vec(&json!({"offset": offset})).unwrap_or_default())
}
fn page(items: Vec<Value>, total: usize, offset: usize) -> Value {
let mut value = json!({
"items": items,
"total": total,
"offset": offset,
"limit": PAGE_SIZE
});
let next = offset.saturating_add(PAGE_SIZE);
if next < total {
value["nextCursor"] = Value::String(encode_cursor(next));
}
value
}
fn posts_page(conn: &DbConnection, project_id: &str, offset: usize) -> EngineResult<Value> {
let posts = post_q::list_posts_by_project(conn, project_id)?;
let total = posts.len();
let items = posts
.into_iter()
.skip(offset)
.take(PAGE_SIZE)
.map(|post| post_summary(conn, &post))
.collect::<EngineResult<Vec<_>>>()?;
Ok(page(items, total, offset))
}
fn media_page(conn: &DbConnection, project_id: &str, offset: usize) -> EngineResult<Value> {
let media = media_q::list_media_by_project(conn, project_id)?;
let total = media.len();
let items = media
.into_iter()
.skip(offset)
.take(PAGE_SIZE)
.map(|item| media_summary(&item))
.collect();
Ok(page(items, total, offset))
}
pub(crate) fn post_summary(conn: &DbConnection, post: &Post) -> EngineResult<Value> {
Ok(json!({
"id": post.id,
"title": post.title,
"slug": post.slug,
"status": post.status,
"tags": post.tags,
"categories": post.categories,
"created_at": post.created_at,
"backlinks": linked_posts(conn, &post.id, false)?,
"linksTo": linked_posts(conn, &post.id, true)?
}))
}
fn linked_posts(conn: &DbConnection, post_id: &str, outgoing: bool) -> EngineResult<Vec<Value>> {
let links = if outgoing {
post_link::list_links_by_source(conn, post_id)?
.into_iter()
.map(|link| link.target_post_id)
.collect::<Vec<_>>()
} else {
post_link::list_links_by_target(conn, post_id)?
.into_iter()
.map(|link| link.source_post_id)
.collect::<Vec<_>>()
};
Ok(links
.into_iter()
.filter_map(|id| post_q::get_post_by_id(conn, &id).ok())
.map(|post| json!({"id": post.id, "title": post.title, "slug": post.slug}))
.collect())
}
pub(crate) fn post_detail(
conn: &DbConnection,
data_dir: &Path,
post: &Post,
) -> EngineResult<Value> {
let mut value = serde_json::to_value(post)?;
value["content"] = Value::String(post_body(data_dir, post));
value["backlinks"] = Value::Array(linked_posts(conn, &post.id, false)?);
value["linksTo"] = Value::Array(linked_posts(conn, &post.id, true)?);
let translations =
crate::db::queries::post_translation::list_post_translations_by_post(conn, &post.id)?;
let mut languages = post.language.clone().into_iter().collect::<Vec<_>>();
languages.extend(
translations
.into_iter()
.map(|translation| translation.language),
);
languages.sort();
languages.dedup();
value["availableLanguages"] = serde_json::to_value(languages)?;
Ok(value)
}
fn post_body(data_dir: &Path, post: &Post) -> String {
if let Some(content) = &post.content {
return content.clone();
}
if post.file_path.is_empty() {
return String::new();
}
std::fs::read_to_string(data_dir.join(&post.file_path))
.ok()
.and_then(|source| {
crate::util::frontmatter::read_post_file(&source)
.ok()
.map(|(_, body)| body)
})
.unwrap_or_default()
}
fn post_detail_by_id(
conn: &DbConnection,
project_id: &str,
data_dir: &Path,
id: &str,
) -> EngineResult<Value> {
let post = post_q::get_post_by_id(conn, id)
.map_err(|_| EngineError::NotFound(format!("post {id}")))?;
ensure_project(project_id, &post.project_id, "post", id)?;
post_detail(conn, data_dir, &post)
}
pub(crate) fn media_summary(media: &Media) -> Value {
json!({
"id": media.id,
"filename": media.filename,
"title": media.title,
"alt": media.alt,
"caption": media.caption,
"tags": media.tags
})
}
fn media_detail_by_id(conn: &DbConnection, project_id: &str, id: &str) -> EngineResult<Value> {
let media = media_q::get_media_by_id(conn, id)
.map_err(|_| EngineError::NotFound(format!("media {id}")))?;
ensure_project(project_id, &media.project_id, "media", id)?;
Ok(serde_json::to_value(media)?)
}
fn post_media_items(conn: &DbConnection, project_id: &str, post_id: &str) -> EngineResult<Value> {
let post = post_q::get_post_by_id(conn, post_id)
.map_err(|_| EngineError::NotFound(format!("post {post_id}")))?;
ensure_project(project_id, &post.project_id, "post", post_id)?;
let items = post_media::list_post_media_by_post(conn, post_id)?
.into_iter()
.filter_map(|link| media_q::get_media_by_id(conn, &link.media_id).ok())
.map(|media| media_summary(&media))
.collect::<Vec<_>>();
Ok(json!({"items": items}))
}
fn media_image(
conn: &DbConnection,
project_id: &str,
data_dir: &Path,
media_id: &str,
uri: &str,
) -> EngineResult<ResourceContent> {
let media = media_q::get_media_by_id(conn, media_id)
.map_err(|_| EngineError::NotFound(format!("media {media_id}")))?;
ensure_project(project_id, &media.project_id, "media", media_id)?;
let bytes = std::fs::read(data_dir.join(&media.file_path))
.map_err(|_| EngineError::NotFound(format!("media file {media_id}")))?;
Ok(ResourceContent {
uri: uri.to_string(),
mime_type: media.mime_type,
text: None,
blob: Some(base64::engine::general_purpose::STANDARD.encode(bytes)),
})
}
fn tags(conn: &DbConnection, project_id: &str) -> EngineResult<Value> {
let posts = post_q::list_posts_by_project(conn, project_id)?;
let items = tag_q::list_tags_by_project(conn, project_id)?
.into_iter()
.map(|tag| {
let count = posts
.iter()
.filter(|post| post.tags.iter().any(|name| name == &tag.name))
.count();
json!({"name": tag.name, "color": tag.color, "post_count": count})
})
.collect::<Vec<_>>();
Ok(json!({"items": items}))
}
fn categories(conn: &DbConnection, project_id: &str, data_dir: &Path) -> EngineResult<Value> {
let posts = post_q::list_posts_by_project(conn, project_id)?;
let names = crate::engine::meta::read_categories_json(data_dir)
.unwrap_or_else(|_| post_q::distinct_post_categories(conn, project_id).unwrap_or_default());
let items = names
.into_iter()
.map(|name| {
let count = posts
.iter()
.filter(|post| post.categories.iter().any(|value| value == &name))
.count();
json!({"name": name, "post_count": count})
})
.collect::<Vec<_>>();
Ok(json!({"items": items}))
}
fn stats(conn: &DbConnection, project_id: &str, data_dir: &Path) -> EngineResult<Value> {
let categories = categories(conn, project_id, data_dir)?;
Ok(json!({
"post_count": post_q::count_posts_by_project(conn, project_id)?,
"media_count": media_q::count_media_by_project(conn, project_id)?,
"tag_count": tag_q::list_tags_by_project(conn, project_id)?.len(),
"category_count": categories["items"].as_array().map_or(0, Vec::len)
}))
}
fn ensure_project(expected: &str, actual: &str, entity: &str, id: &str) -> EngineResult<()> {
if expected == actual {
Ok(())
} else {
Err(EngineError::NotFound(format!("{entity} {id}")))
}
}

View File

@@ -0,0 +1,646 @@
use std::collections::BTreeMap;
use chrono::{Datelike, TimeZone as _, Utc};
use serde_json::{Map, Value, json};
use crate::db::DbConnection;
use crate::db::queries::{media as media_q, media_translation, post as post_q, post_translation};
use crate::engine::{EngineError, EngineResult};
use crate::model::{Post, ProposalKind};
use super::resources::{post_detail, post_summary};
use super::{active_project, create_proposal, optional_string, required_string, string_array};
const MAX_PAGE_SIZE: usize = 50;
pub fn list() -> Vec<Value> {
[
tool(
"check_term",
"Check Term",
"Check whether a term is a category, a tag, or both, with post counts.",
object_schema(json!({"term": string_schema("Term to check")}), &["term"]),
true,
),
tool(
"search_posts",
"Search Posts",
"Full-text and filtered post search with pagination, backlinks, and outgoing links.",
post_query_schema(false),
true,
),
tool(
"count_posts",
"Count Posts",
"Count filtered posts grouped by year, month, tag, category, or status.",
count_schema(),
true,
),
tool(
"read_post_by_slug",
"Read Post By Slug",
"Read full post content and metadata, optionally in a translated language.",
object_schema(
json!({
"slug": string_schema("Post slug"),
"language": string_schema("Optional translation language")
}),
&["slug"],
),
true,
),
tool(
"get_post_translations",
"Get Post Translations",
"List every translation for a post.",
id_schema("postId", "Post ID"),
true,
),
tool(
"get_media_translations",
"Get Media Translations",
"List every translated metadata record for a media item.",
id_schema("mediaId", "Media ID"),
true,
),
tool(
"upsert_media_translation",
"Propose Media Translation",
"Propose translated media metadata for explicit desktop approval.",
object_schema(
json!({
"mediaId": string_schema("Media ID"),
"language": string_schema("Language code"),
"title": string_schema("Translated title"),
"alt": string_schema("Translated alt text"),
"caption": string_schema("Translated caption")
}),
&["mediaId", "language"],
),
false,
),
tool(
"draft_post",
"Draft Post",
"Propose a post. No post is created until explicit desktop approval.",
object_schema(
json!({
"title": string_schema("Post title"),
"content": string_schema("Markdown body"),
"excerpt": string_schema("Excerpt"),
"tags": string_array_schema("Tags"),
"categories": string_array_schema("Categories"),
"author": string_schema("Author"),
"language": string_schema("Language")
}),
&["title", "content"],
),
false,
),
tool(
"propose_script",
"Propose Script",
"Validate and propose a Lua script for explicit desktop approval.",
object_schema(
json!({
"title": string_schema("Script title"),
"kind": {"type":"string","enum":["macro","utility","transform"]},
"content": string_schema("Lua source"),
"entrypoint": string_schema("Entrypoint function")
}),
&["title", "kind", "content"],
),
false,
),
tool(
"propose_template",
"Propose Template",
"Validate and propose a Liquid template for explicit desktop approval.",
object_schema(
json!({
"title": string_schema("Template title"),
"kind": {"type":"string","enum":["post","list","not-found","partial"]},
"content": string_schema("Liquid source")
}),
&["title", "kind", "content"],
),
false,
),
tool(
"propose_media_metadata",
"Propose Media Metadata",
"Propose media metadata changes for explicit desktop approval.",
object_schema(
json!({
"mediaId": string_schema("Media ID"),
"title": nullable_string_schema("Title"),
"alt": nullable_string_schema("Alt text"),
"caption": nullable_string_schema("Caption"),
"tags": string_array_schema("Tags")
}),
&["mediaId"],
),
false,
),
tool(
"propose_post_metadata",
"Propose Post Metadata",
"Propose post metadata changes for explicit desktop approval.",
object_schema(
json!({
"postId": string_schema("Post ID"),
"title": string_schema("Title"),
"excerpt": nullable_string_schema("Excerpt"),
"tags": string_array_schema("Tags"),
"categories": string_array_schema("Categories")
}),
&["postId"],
),
false,
),
]
.into_iter()
.collect()
}
pub fn call(conn: &DbConnection, name: &str, params: Value) -> EngineResult<Value> {
if !params.is_object() {
return Err(EngineError::Validation(
"tool arguments must be an object".into(),
));
}
match name {
"check_term" => check_term(conn, &params),
"search_posts" => search_posts(conn, &params),
"count_posts" => count_posts(conn, &params),
"read_post_by_slug" => read_post_by_slug(conn, &params),
"get_post_translations" => get_post_translations(conn, &params),
"get_media_translations" => get_media_translations(conn, &params),
"upsert_media_translation" => propose_media_translation(conn, &params),
"draft_post" => propose(conn, ProposalKind::DraftPost, None, &params),
"propose_script" => propose_script(conn, &params),
"propose_template" => propose_template(conn, &params),
"propose_media_metadata" => propose_media_metadata(conn, &params),
"propose_post_metadata" => propose_post_metadata(conn, &params),
_ => Err(EngineError::NotFound(format!("MCP tool {name}"))),
}
}
fn tool(name: &str, title: &str, description: &str, schema: Value, read_only: bool) -> Value {
json!({
"name": name,
"title": title,
"description": description,
"inputSchema": schema,
"annotations": {
"readOnlyHint": read_only,
"destructiveHint": false,
"openWorldHint": false
}
})
}
fn string_schema(description: &str) -> Value {
json!({"type":"string","description":description})
}
fn nullable_string_schema(description: &str) -> Value {
json!({"type":["string","null"],"description":description})
}
fn string_array_schema(description: &str) -> Value {
json!({"type":"array","items":{"type":"string"},"description":description})
}
fn object_schema(properties: Value, required: &[&str]) -> Value {
let mut schema = json!({
"type": "object",
"properties": properties,
"additionalProperties": false
});
if !required.is_empty() {
schema["required"] = serde_json::to_value(required).unwrap_or_default();
}
schema
}
fn id_schema(field: &str, description: &str) -> Value {
let mut properties = Map::new();
properties.insert(field.into(), string_schema(description));
object_schema(Value::Object(properties), &[field])
}
fn post_query_schema(query_required: bool) -> Value {
object_schema(
json!({
"query": string_schema("Full-text query"),
"category": string_schema("Category filter"),
"tags": string_array_schema("All required tags"),
"language": string_schema("Available language"),
"missingTranslationLanguage": string_schema("Missing translation language"),
"year": {"type":"integer"},
"month": {"type":"integer","minimum":1,"maximum":12},
"status": {"type":"string","enum":["draft","published","archived"]},
"offset": {"type":"integer","minimum":0},
"limit": {"type":"integer","minimum":1,"maximum":MAX_PAGE_SIZE}
}),
if query_required { &["query"] } else { &[] },
)
}
fn count_schema() -> Value {
let mut schema = post_query_schema(false);
schema["properties"]["groupBy"] = json!({"type":"array","items":{"type":"string","enum":["year","month","tag","category","status"]},"minItems":1});
schema["required"] = json!(["groupBy"]);
schema
}
fn check_term(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let term = required_string(params, "term")?;
let (project, _) = active_project(conn)?;
let posts = post_q::list_posts_by_project(conn, &project.id)?;
let normalized = term.to_lowercase();
let tag_count = posts
.iter()
.filter(|post| post.tags.iter().any(|tag| tag.to_lowercase() == normalized))
.count();
let category_count = posts
.iter()
.filter(|post| {
post.categories
.iter()
.any(|category| category.to_lowercase() == normalized)
})
.count();
Ok(json!({
"is_category": category_count > 0,
"category_post_count": category_count,
"is_tag": tag_count > 0,
"tag_post_count": tag_count
}))
}
fn filtered_posts(conn: &DbConnection, params: &Value) -> EngineResult<Vec<Post>> {
let (project, _) = active_project(conn)?;
let mut posts = post_q::list_posts_by_project(conn, &project.id)?;
let query = optional_string(params, "query")
.unwrap_or_default()
.trim()
.to_lowercase();
let fts_matches = if query.is_empty() {
None
} else {
let language = optional_string(params, "language").unwrap_or("en");
Some(
crate::db::fts::search_posts(conn, &query, language)?
.into_iter()
.collect::<std::collections::HashSet<_>>(),
)
};
let category = optional_string(params, "category").map(str::to_lowercase);
let tags = string_array(params, "tags")
.into_iter()
.map(|tag| tag.to_lowercase())
.collect::<Vec<_>>();
let language = optional_string(params, "language").map(str::to_lowercase);
let missing_language =
optional_string(params, "missingTranslationLanguage").map(str::to_lowercase);
let status = optional_string(params, "status");
let year = integer(params, "year")?.map(|value| value as i32);
let month = integer(params, "month")?.map(|value| value as u32);
if month.is_some() && year.is_none() {
return Err(EngineError::Validation("month requires year".into()));
}
posts.retain(|post| {
if fts_matches
.as_ref()
.is_some_and(|matches| !matches.contains(&post.id))
{
return false;
}
if category.as_ref().is_some_and(|wanted| {
!post
.categories
.iter()
.any(|value| value.to_lowercase() == *wanted)
}) {
return false;
}
if tags.iter().any(|wanted| {
!post
.tags
.iter()
.any(|value| value.to_lowercase() == *wanted)
}) {
return false;
}
if status.is_some_and(|wanted| post.status.as_str() != wanted) {
return false;
}
if let Some(wanted) = &language {
let canonical = post
.language
.as_deref()
.is_some_and(|value| value.eq_ignore_ascii_case(wanted));
let translated =
post_translation::get_post_translation_by_post_and_language(conn, &post.id, wanted)
.is_ok();
if !canonical && !translated {
return false;
}
}
if let Some(wanted) = &missing_language {
let canonical = post
.language
.as_deref()
.is_some_and(|value| value.eq_ignore_ascii_case(wanted));
let translated =
post_translation::get_post_translation_by_post_and_language(conn, &post.id, wanted)
.is_ok();
if canonical || translated {
return false;
}
}
if let Some(wanted_year) = year {
let Some(date) = Utc.timestamp_millis_opt(post.created_at).single() else {
return false;
};
if date.year() != wanted_year || month.is_some_and(|wanted| date.month() != wanted) {
return false;
}
}
true
});
Ok(posts)
}
fn search_posts(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let posts = filtered_posts(conn, params)?;
let total = posts.len();
let offset = unsigned(params, "offset")?.unwrap_or(0);
let limit = unsigned(params, "limit")?.unwrap_or(MAX_PAGE_SIZE);
if limit == 0 || limit > MAX_PAGE_SIZE {
return Err(EngineError::Validation(format!(
"limit must be between 1 and {MAX_PAGE_SIZE}"
)));
}
let posts = posts
.into_iter()
.skip(offset)
.take(limit)
.map(|post| post_summary(conn, &post))
.collect::<EngineResult<Vec<_>>>()?;
Ok(json!({
"total": total,
"offset": offset,
"limit": limit,
"hasMore": offset.saturating_add(limit) < total,
"posts": posts
}))
}
fn count_posts(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let group_by = string_array(params, "groupBy");
if group_by.is_empty()
|| group_by.iter().any(|dimension| {
!["year", "month", "tag", "category", "status"].contains(&dimension.as_str())
})
{
return Err(EngineError::Validation("invalid groupBy".into()));
}
let posts = filtered_posts(conn, params)?;
let total = posts.len();
let mut counts = BTreeMap::<String, (Map<String, Value>, usize)>::new();
for post in posts {
for row in group_rows(&post, &group_by) {
let key = serde_json::to_string(&row)?;
counts
.entry(key)
.and_modify(|(_, count)| *count += 1)
.or_insert((row, 1));
}
}
let groups = counts
.into_values()
.map(|(mut row, count)| {
row.insert("count".into(), json!(count));
Value::Object(row)
})
.collect::<Vec<_>>();
Ok(json!({"groups": groups, "totalPosts": total}))
}
fn group_rows(post: &Post, dimensions: &[String]) -> Vec<Map<String, Value>> {
let Some((dimension, rest)) = dimensions.split_first() else {
return vec![Map::new()];
};
let values = match dimension.as_str() {
"year" => Utc
.timestamp_millis_opt(post.created_at)
.single()
.map(|date| vec![json!(date.year())])
.unwrap_or_else(|| vec![Value::Null]),
"month" => Utc
.timestamp_millis_opt(post.created_at)
.single()
.map(|date| vec![json!(date.month())])
.unwrap_or_else(|| vec![Value::Null]),
"tag" => values_or_null(&post.tags),
"category" => values_or_null(&post.categories),
"status" => vec![json!(post.status.as_str())],
_ => vec![Value::Null],
};
values
.into_iter()
.flat_map(|value| {
group_rows(post, rest).into_iter().map({
let dimension = dimension.clone();
move |mut row| {
row.insert(dimension.clone(), value.clone());
row
}
})
})
.collect()
}
fn values_or_null(values: &[String]) -> Vec<Value> {
if values.is_empty() {
vec![Value::Null]
} else {
values.iter().map(|value| json!(value)).collect()
}
}
fn read_post_by_slug(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let slug = required_string(params, "slug")?;
let (project, data_dir) = active_project(conn)?;
let post = post_q::get_post_by_project_and_slug(conn, &project.id, slug)
.map_err(|_| EngineError::NotFound(format!("post slug {slug}")))?;
let Some(language) = optional_string(params, "language") else {
return Ok(json!({"post": post_detail(conn, &data_dir, &post)?}));
};
if post
.language
.as_deref()
.is_some_and(|canonical| canonical.eq_ignore_ascii_case(language))
{
return Ok(json!({"post": post_detail(conn, &data_dir, &post)?}));
}
let translation =
post_translation::get_post_translation_by_post_and_language(conn, &post.id, language)
.map_err(|_| EngineError::NotFound(format!("post translation {language}")))?;
let mut detail = post_detail(conn, &data_dir, &post)?;
detail["title"] = json!(translation.title);
detail["excerpt"] = json!(translation.excerpt);
detail["content"] = json!(translation_content(&data_dir, &translation));
detail["language"] = json!(translation.language);
detail["canonicalLanguage"] = json!(post.language);
Ok(json!({"post": detail}))
}
fn translation_content(
data_dir: &std::path::Path,
translation: &crate::model::PostTranslation,
) -> String {
if let Some(content) = &translation.content {
return content.clone();
}
std::fs::read_to_string(data_dir.join(&translation.file_path))
.ok()
.and_then(|source| {
crate::util::frontmatter::read_translation_file(&source)
.ok()
.map(|(_, body)| body)
})
.unwrap_or_default()
}
fn get_post_translations(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let id = required_string(params, "postId")?;
let (project, data_dir) = active_project(conn)?;
let post = post_q::get_post_by_id(conn, id)
.map_err(|_| EngineError::NotFound(format!("post {id}")))?;
if post.project_id != project.id {
return Err(EngineError::NotFound(format!("post {id}")));
}
let translations = post_translation::list_post_translations_by_post(conn, id)?
.into_iter()
.map(|translation| {
let mut value = serde_json::to_value(&translation)?;
value["content"] = json!(translation_content(&data_dir, &translation));
Ok(value)
})
.collect::<EngineResult<Vec<_>>>()?;
Ok(json!({"translations": translations}))
}
fn get_media_translations(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let id = required_string(params, "mediaId")?;
let (project, _) = active_project(conn)?;
let media = media_q::get_media_by_id(conn, id)
.map_err(|_| EngineError::NotFound(format!("media {id}")))?;
if media.project_id != project.id {
return Err(EngineError::NotFound(format!("media {id}")));
}
Ok(json!({
"translations": media_translation::list_media_translations_by_media(conn, id)?
}))
}
fn propose_media_translation(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let id = required_string(params, "mediaId")?;
required_string(params, "language")?;
let (project, _) = active_project(conn)?;
let media = media_q::get_media_by_id(conn, id)
.map_err(|_| EngineError::NotFound(format!("media {id}")))?;
if media.project_id != project.id {
return Err(EngineError::NotFound(format!("media {id}")));
}
propose(
conn,
ProposalKind::ProposeMediaTranslation,
Some(id),
params,
)
}
fn propose_script(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
required_string(params, "title")?;
let content = required_string(params, "content")?;
required_string(params, "kind")?
.parse::<crate::model::ScriptKind>()
.map_err(EngineError::Validation)?;
crate::engine::script::validate_script_syntax(content).map_err(EngineError::Validation)?;
propose(conn, ProposalKind::ProposeScript, None, params)
}
fn propose_template(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
required_string(params, "title")?;
let content = required_string(params, "content")?;
required_string(params, "kind")?
.parse::<crate::model::TemplateKind>()
.map_err(EngineError::Validation)?;
crate::engine::template::validate_template(content).map_err(EngineError::Validation)?;
propose(conn, ProposalKind::ProposeTemplate, None, params)
}
fn propose_media_metadata(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let id = required_string(params, "mediaId")?;
let (project, _) = active_project(conn)?;
let media = media_q::get_media_by_id(conn, id)
.map_err(|_| EngineError::NotFound(format!("media {id}")))?;
if media.project_id != project.id {
return Err(EngineError::NotFound(format!("media {id}")));
}
propose(conn, ProposalKind::ProposeMediaMetadata, Some(id), params)
}
fn propose_post_metadata(conn: &DbConnection, params: &Value) -> EngineResult<Value> {
let id = required_string(params, "postId")?;
let (project, _) = active_project(conn)?;
let post = post_q::get_post_by_id(conn, id)
.map_err(|_| EngineError::NotFound(format!("post {id}")))?;
if post.project_id != project.id {
return Err(EngineError::NotFound(format!("post {id}")));
}
propose(conn, ProposalKind::ProposePostMetadata, Some(id), params)
}
fn propose(
conn: &DbConnection,
kind: ProposalKind,
entity_id: Option<&str>,
params: &Value,
) -> EngineResult<Value> {
if kind == ProposalKind::DraftPost {
required_string(params, "title")?;
required_string(params, "content")?;
}
let (project, _) = active_project(conn)?;
let proposal = create_proposal(conn, kind, &project.id, entity_id, params)?;
Ok(json!({
"proposalId": proposal.id,
"status": proposal.status,
"expiresAt": proposal.expires_at,
"message": "Pending explicit approval in RuDS Settings"
}))
}
fn integer(value: &Value, key: &str) -> EngineResult<Option<i64>> {
match value.get(key) {
None => Ok(None),
Some(value) => value
.as_i64()
.map(Some)
.ok_or_else(|| EngineError::Validation(format!("{key} must be an integer"))),
}
}
fn unsigned(value: &Value, key: &str) -> EngineResult<Option<usize>> {
integer(value, key)?.map_or(Ok(None), |value| {
usize::try_from(value)
.map(Some)
.map_err(|_| EngineError::Validation(format!("{key} cannot be negative")))
})
}

View File

@@ -595,6 +595,7 @@ pub fn upsert_media_translation(
// Re-index FTS for parent media
fts_index_media(conn, &media)?;
emit_media(&media, NotificationAction::Updated);
Ok(translation)
}

View File

@@ -9,6 +9,7 @@ pub mod error;
pub mod gallery_import;
pub mod generation;
pub mod git;
pub mod mcp;
pub mod media;
pub mod menu;
pub mod meta;

View File

@@ -0,0 +1,128 @@
use serde::{Deserialize, Serialize};
#[derive(
Debug,
Clone,
Copy,
PartialEq,
Eq,
Serialize,
Deserialize,
diesel::AsExpression,
diesel::FromSqlRow,
)]
#[diesel(sql_type = diesel::sql_types::Text)]
#[serde(rename_all = "snake_case")]
pub enum ProposalKind {
DraftPost,
ProposeScript,
ProposeTemplate,
ProposeMediaTranslation,
ProposeMediaMetadata,
ProposePostMetadata,
}
impl ProposalKind {
pub const fn as_str(self) -> &'static str {
match self {
Self::DraftPost => "draft_post",
Self::ProposeScript => "propose_script",
Self::ProposeTemplate => "propose_template",
Self::ProposeMediaTranslation => "propose_media_translation",
Self::ProposeMediaMetadata => "propose_media_metadata",
Self::ProposePostMetadata => "propose_post_metadata",
}
}
}
impl std::str::FromStr for ProposalKind {
type Err = String;
fn from_str(value: &str) -> Result<Self, Self::Err> {
match value {
"draft_post" => Ok(Self::DraftPost),
"propose_script" => Ok(Self::ProposeScript),
"propose_template" => Ok(Self::ProposeTemplate),
"propose_media_translation" => Ok(Self::ProposeMediaTranslation),
"propose_media_metadata" => Ok(Self::ProposeMediaMetadata),
"propose_post_metadata" => Ok(Self::ProposePostMetadata),
_ => Err(format!("invalid proposal kind: {value}")),
}
}
}
#[derive(
Debug,
Clone,
Copy,
PartialEq,
Eq,
Serialize,
Deserialize,
diesel::AsExpression,
diesel::FromSqlRow,
)]
#[diesel(sql_type = diesel::sql_types::Text)]
#[serde(rename_all = "lowercase")]
pub enum ProposalStatus {
Pending,
Executing,
Accepted,
Rejected,
Expired,
}
impl ProposalStatus {
pub const fn as_str(self) -> &'static str {
match self {
Self::Pending => "pending",
Self::Executing => "executing",
Self::Accepted => "accepted",
Self::Rejected => "rejected",
Self::Expired => "expired",
}
}
}
impl std::str::FromStr for ProposalStatus {
type Err = String;
fn from_str(value: &str) -> Result<Self, Self::Err> {
match value {
"pending" => Ok(Self::Pending),
"executing" => Ok(Self::Executing),
"accepted" => Ok(Self::Accepted),
"rejected" => Ok(Self::Rejected),
"expired" => Ok(Self::Expired),
_ => Err(format!("invalid proposal status: {value}")),
}
}
}
#[derive(
Debug,
Clone,
PartialEq,
Eq,
Serialize,
Deserialize,
diesel::Queryable,
diesel::Selectable,
diesel::Insertable,
)]
#[diesel(
table_name = crate::db::schema::mcp_proposals,
check_for_backend(diesel::sqlite::Sqlite)
)]
pub struct McpProposal {
pub id: String,
pub project_id: String,
pub kind: ProposalKind,
pub status: ProposalStatus,
pub entity_id: Option<String>,
pub data: String,
pub result: Option<String>,
pub created_at: i64,
pub expires_at: i64,
pub resolved_at: Option<i64>,
}

View File

@@ -1,6 +1,7 @@
mod event;
mod generation;
mod import;
mod mcp;
mod media;
pub mod metadata;
mod post;
@@ -19,6 +20,7 @@ pub use import::{
ImportExecutionResult, ImportItemKind, ImportItemStatus, ImportMacroUsage, ImportPhase,
ImportProgress, ImportReport, ImportResolution, ImportedSite, TaxonomyCandidate, TaxonomyKind,
};
pub use mcp::{McpProposal, ProposalKind, ProposalStatus};
pub use media::{Media, MediaTranslation};
pub use metadata::{CategorySettings, ProjectMetadata, TagEntry};
pub use post::{Post, PostLink, PostMedia, PostStatus, PostTranslation};