Skip to content

Commit 051b854

Browse files
committed
fix(clone): Phase 1 — clone-time source rewrite and auth header seeding
Autonomous decisions: - Used existing with_inputs() for input_values — no with_extra_headers builder exists on InstalledServer - Added PartialEq to ServerSource — needed for assert_eq! in the new unit test Signed-off-by: crimsonsunset <jsangio1@gmail.com>
1 parent aa76566 commit 051b854

2 files changed

Lines changed: 323 additions & 3 deletions

File tree

crates/mcpmux-core/src/application/server.rs

Lines changed: 320 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,7 @@ use uuid::Uuid;
1010

1111
use crate::domain::{
1212
config::UserServerEntry, DomainEvent, InstallationSource, InstalledServer, ServerDefinition,
13-
UpdatePolicy,
13+
ServerSource, UpdatePolicy,
1414
};
1515
use crate::event_bus::EventSender;
1616
use crate::repository::{CredentialRepository, InstalledServerRepository, ServerFeatureRepository};
@@ -409,15 +409,18 @@ impl ServerAppService {
409409
definition.id = new_server_id.clone();
410410
definition.name = format!("{} ({})", source.display_name(), normalized_suffix);
411411
definition.alias = Some(alias);
412+
definition.source = ServerSource::ManualEntry;
412413

413-
let server = InstalledServer::new(&space_id_str, &new_server_id)
414+
let mut server = InstalledServer::new(&space_id_str, &new_server_id)
414415
.with_definition(&definition)
415416
.with_source(InstallationSource::ManualEntry)
416417
.with_cloned_from(source_server_id)
417418
.with_display_name_override(display_name_override)
418419
.with_update_policy(source.update_policy)
419420
.with_pinned_version(source.pinned_version.clone())
421+
.with_inputs(source.input_values.clone())
420422
.with_enabled(false);
423+
server.extra_headers = source.extra_headers.clone();
421424

422425
self.server_repo.install(&server).await?;
423426

@@ -545,3 +548,318 @@ impl ServerAppService {
545548
Ok(())
546549
}
547550
}
551+
552+
#[cfg(test)]
553+
mod tests {
554+
use super::*;
555+
use async_trait::async_trait;
556+
use chrono::Utc;
557+
use std::collections::HashMap;
558+
use std::path::PathBuf;
559+
use std::sync::Arc;
560+
use tokio::sync::RwLock;
561+
562+
use crate::domain::{
563+
AuthConfig, InputDefinition, ServerSource, TransportConfig, TransportMetadata,
564+
};
565+
use crate::event_bus::EventBus;
566+
use crate::repository::InstalledServerRepository;
567+
568+
/// Minimal in-memory repo for clone_server unit tests.
569+
struct InMemoryInstalledServerRepository {
570+
servers: RwLock<Vec<InstalledServer>>,
571+
}
572+
573+
impl InMemoryInstalledServerRepository {
574+
fn new() -> Self {
575+
Self {
576+
servers: RwLock::new(Vec::new()),
577+
}
578+
}
579+
580+
async fn seed(&self, server: InstalledServer) {
581+
self.servers.write().await.push(server);
582+
}
583+
584+
async fn get_by_server_id(
585+
&self,
586+
space_id: &str,
587+
server_id: &str,
588+
) -> Option<InstalledServer> {
589+
self.servers
590+
.read()
591+
.await
592+
.iter()
593+
.find(|s| s.space_id == space_id && s.server_id == server_id)
594+
.cloned()
595+
}
596+
}
597+
598+
#[async_trait]
599+
impl InstalledServerRepository for InMemoryInstalledServerRepository {
600+
async fn list(&self) -> crate::repository::RepoResult<Vec<InstalledServer>> {
601+
Ok(self.servers.read().await.clone())
602+
}
603+
604+
async fn list_for_space(
605+
&self,
606+
space_id: &str,
607+
) -> crate::repository::RepoResult<Vec<InstalledServer>> {
608+
Ok(self
609+
.servers
610+
.read()
611+
.await
612+
.iter()
613+
.filter(|s| s.space_id == space_id)
614+
.cloned()
615+
.collect())
616+
}
617+
618+
async fn list_by_source_file(
619+
&self,
620+
_file_path: &std::path::Path,
621+
) -> crate::repository::RepoResult<Vec<InstalledServer>> {
622+
Ok(Vec::new())
623+
}
624+
625+
async fn get(&self, id: &Uuid) -> crate::repository::RepoResult<Option<InstalledServer>> {
626+
Ok(self
627+
.servers
628+
.read()
629+
.await
630+
.iter()
631+
.find(|s| &s.id == id)
632+
.cloned())
633+
}
634+
635+
async fn get_by_server_id(
636+
&self,
637+
space_id: &str,
638+
server_id: &str,
639+
) -> crate::repository::RepoResult<Option<InstalledServer>> {
640+
Ok(self.get_by_server_id(space_id, server_id).await)
641+
}
642+
643+
async fn install(&self, server: &InstalledServer) -> crate::repository::RepoResult<()> {
644+
self.servers.write().await.push(server.clone());
645+
Ok(())
646+
}
647+
648+
async fn update(&self, server: &InstalledServer) -> crate::repository::RepoResult<()> {
649+
let mut servers = self.servers.write().await;
650+
if let Some(existing) = servers.iter_mut().find(|s| s.id == server.id) {
651+
*existing = server.clone();
652+
Ok(())
653+
} else {
654+
anyhow::bail!("server not found for update: {}", server.server_id)
655+
}
656+
}
657+
658+
async fn uninstall(&self, id: &Uuid) -> crate::repository::RepoResult<()> {
659+
self.servers.write().await.retain(|s| &s.id != id);
660+
Ok(())
661+
}
662+
663+
async fn list_enabled(
664+
&self,
665+
space_id: &str,
666+
) -> crate::repository::RepoResult<Vec<InstalledServer>> {
667+
Ok(self
668+
.servers
669+
.read()
670+
.await
671+
.iter()
672+
.filter(|s| s.space_id == space_id && s.enabled)
673+
.cloned()
674+
.collect())
675+
}
676+
677+
async fn list_enabled_all(&self) -> crate::repository::RepoResult<Vec<InstalledServer>> {
678+
Ok(self
679+
.servers
680+
.read()
681+
.await
682+
.iter()
683+
.filter(|s| s.enabled)
684+
.cloned()
685+
.collect())
686+
}
687+
688+
async fn set_enabled(&self, id: &Uuid, enabled: bool) -> crate::repository::RepoResult<()> {
689+
if let Some(server) = self.servers.write().await.iter_mut().find(|s| &s.id == id) {
690+
server.enabled = enabled;
691+
}
692+
Ok(())
693+
}
694+
695+
async fn set_oauth_connected(
696+
&self,
697+
id: &Uuid,
698+
connected: bool,
699+
) -> crate::repository::RepoResult<()> {
700+
if let Some(server) = self.servers.write().await.iter_mut().find(|s| &s.id == id) {
701+
server.oauth_connected = connected;
702+
}
703+
Ok(())
704+
}
705+
706+
async fn update_inputs(
707+
&self,
708+
id: &Uuid,
709+
input_values: HashMap<String, String>,
710+
) -> crate::repository::RepoResult<()> {
711+
if let Some(server) = self.servers.write().await.iter_mut().find(|s| &s.id == id) {
712+
server.input_values = input_values;
713+
}
714+
Ok(())
715+
}
716+
717+
async fn update_cached_definition(
718+
&self,
719+
id: &Uuid,
720+
server_name: Option<String>,
721+
cached_definition: Option<String>,
722+
) -> crate::repository::RepoResult<()> {
723+
if let Some(server) = self.servers.write().await.iter_mut().find(|s| &s.id == id) {
724+
server.server_name = server_name;
725+
server.cached_definition = cached_definition;
726+
}
727+
Ok(())
728+
}
729+
730+
async fn set_display_name_override(
731+
&self,
732+
id: &Uuid,
733+
value: Option<String>,
734+
) -> crate::repository::RepoResult<()> {
735+
if let Some(server) = self.servers.write().await.iter_mut().find(|s| &s.id == id) {
736+
server.display_name_override = value;
737+
}
738+
Ok(())
739+
}
740+
741+
async fn update_version_cache(
742+
&self,
743+
id: &Uuid,
744+
latest_available_version: Option<String>,
745+
current_version: Option<String>,
746+
version_checked_at: chrono::DateTime<Utc>,
747+
) -> crate::repository::RepoResult<()> {
748+
if let Some(server) = self.servers.write().await.iter_mut().find(|s| &s.id == id) {
749+
server.latest_available_version = latest_available_version;
750+
server.current_version = current_version;
751+
server.version_checked_at = Some(version_checked_at);
752+
}
753+
Ok(())
754+
}
755+
}
756+
757+
fn user_space_http_definition(server_id: &str) -> ServerDefinition {
758+
ServerDefinition {
759+
id: server_id.to_string(),
760+
name: "PostHog Personal".to_string(),
761+
description: None,
762+
alias: Some("posthog".to_string()),
763+
auth: Some(AuthConfig::ApiKey {
764+
instructions: None,
765+
}),
766+
icon: None,
767+
transport: TransportConfig::Http {
768+
url: "https://mcp.posthog.com/mcp".to_string(),
769+
headers: HashMap::new(),
770+
metadata: TransportMetadata {
771+
inputs: vec![InputDefinition {
772+
id: "POSTHOG_API_KEY".to_string(),
773+
label: "API Key".to_string(),
774+
r#type: "password".to_string(),
775+
required: true,
776+
secret: true,
777+
description: None,
778+
default: None,
779+
placeholder: None,
780+
obtain_url: None,
781+
obtain_instructions: None,
782+
}],
783+
},
784+
},
785+
categories: vec![],
786+
publisher: None,
787+
source: ServerSource::UserSpace {
788+
space_id: "space-1".to_string(),
789+
file_path: PathBuf::from("/tmp/posthog.json"),
790+
},
791+
badges: vec![],
792+
hosting_type: Default::default(),
793+
license: None,
794+
license_url: None,
795+
installation: None,
796+
capabilities: None,
797+
sponsored: None,
798+
media: None,
799+
changelog_url: None,
800+
}
801+
}
802+
803+
#[tokio::test]
804+
async fn clone_server_rewrites_source_and_seeds_auth_headers() {
805+
let space_id = Uuid::new_v4();
806+
let repo = Arc::new(InMemoryInstalledServerRepository::new());
807+
let event_bus = EventBus::new();
808+
809+
let parent_headers = HashMap::from([
810+
(
811+
"Authorization".to_string(),
812+
"Bearer phx_parent_token".to_string(),
813+
),
814+
(
815+
"x-posthog-project-id".to_string(),
816+
"345911".to_string(),
817+
),
818+
]);
819+
let parent_inputs = HashMap::from([(
820+
"POSTHOG_API_KEY".to_string(),
821+
"phc_parent_key".to_string(),
822+
)]);
823+
824+
let definition = user_space_http_definition("posthog-personal");
825+
let mut source = InstalledServer::new(space_id.to_string(), "posthog-personal")
826+
.with_definition(&definition)
827+
.with_source(InstallationSource::UserConfig {
828+
file_path: PathBuf::from("/tmp/posthog.json"),
829+
})
830+
.with_inputs(parent_inputs.clone());
831+
source.extra_headers = parent_headers.clone();
832+
repo.seed(source).await;
833+
834+
let service = ServerAppService::new(
835+
repo.clone(),
836+
None,
837+
None,
838+
event_bus.sender(),
839+
);
840+
841+
let cloned = service
842+
.clone_server(space_id, "posthog-personal", "mesh", None, None)
843+
.await
844+
.expect("clone should succeed");
845+
846+
assert_eq!(cloned.source, InstallationSource::ManualEntry);
847+
assert_eq!(cloned.extra_headers, parent_headers);
848+
assert_eq!(cloned.input_values, parent_inputs);
849+
850+
let definition = cloned
851+
.get_definition()
852+
.expect("clone should cache definition");
853+
assert!(
854+
!matches!(definition.source, ServerSource::UserSpace { .. }),
855+
"cloned definition source must not remain UserSpace"
856+
);
857+
assert_eq!(definition.source, ServerSource::ManualEntry);
858+
859+
let persisted = repo
860+
.get_by_server_id(&space_id.to_string(), "posthog-personal-mesh")
861+
.await
862+
.expect("cloned row should be persisted");
863+
assert_eq!(persisted.extra_headers, parent_headers);
864+
}
865+
}

crates/mcpmux-core/src/domain/server.rs

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -75,7 +75,7 @@ impl ServerDefinition {
7575
}
7676
}
7777

78-
#[derive(Debug, Clone, Serialize, Deserialize, Default)]
78+
#[derive(Debug, Clone, Serialize, Deserialize, Default, PartialEq)]
7979
#[serde(tag = "type")]
8080
pub enum ServerSource {
8181
/// Loaded from a user-defined JSON file in the spaces directory
@@ -88,6 +88,8 @@ pub enum ServerSource {
8888
Bundled,
8989
/// Loaded from a remote or custom registry (API, NPM, etc.)
9090
Registry { url: String, name: String },
91+
/// Manually installed clone stored only in SQLite (not a user-space JSON file)
92+
ManualEntry,
9193
}
9294

9395
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]

0 commit comments

Comments
 (0)