From 8f9c4a684740ba65ac778361786178d6cd31ffea Mon Sep 17 00:00:00 2001 From: CodeMaster4711 Date: Thu, 23 Jul 2026 20:05:13 +0200 Subject: [PATCH 1/7] fix: point resource group Corefile at actual zone file directory --- agent/src/rg_dns.rs | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/agent/src/rg_dns.rs b/agent/src/rg_dns.rs index a9db6df8..a051744a 100644 --- a/agent/src/rg_dns.rs +++ b/agent/src/rg_dns.rs @@ -97,11 +97,12 @@ pub async fn write_corefile(resource_group_id: &str, listen_ip: &str) -> Result< .context("Failed to create dns zone directory")?; let zone = zone_name(resource_group_id); + let zone_path = zone_file_path(resource_group_id).display().to_string(); let contents = format!( - "{zone}:53 {{\n bind {listen_ip}\n file /zones/{rg}.zone\n reload 2s\n log\n errors\n}}\n", + "{zone}:53 {{\n bind {listen_ip}\n file {zone_path}\n reload 2s\n log\n errors\n}}\n", zone = zone, listen_ip = listen_ip, - rg = resource_group_id, + zone_path = zone_path, ); let path = corefile_path(resource_group_id); From c68a0eeb00ff924ac2b32f27becab88912dc6706 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" Date: Thu, 23 Jul 2026 18:08:32 +0000 Subject: [PATCH 2/7] style: apply cargo fmt and clippy fixes --- agent/src/agent_stream.rs | 9 +++- agent/src/docker.rs | 50 +++++++++---------- agent/src/main.rs | 34 +++++++++++-- agent/src/rg_dns.rs | 5 +- .../api-gateway/src/routes/agent_stream.rs | 9 ++-- control-plane/registry/src/services/pki.rs | 12 +++-- .../registry/src/services/registry.rs | 5 +- .../scheduler/src/handlers/workloads.rs | 8 +-- .../scheduler/src/services/gateway_notify.rs | 4 +- .../scheduler/src/services/scheduler.rs | 9 ++-- 10 files changed, 87 insertions(+), 58 deletions(-) diff --git a/agent/src/agent_stream.rs b/agent/src/agent_stream.rs index f425ac86..10cdb4e4 100644 --- a/agent/src/agent_stream.rs +++ b/agent/src/agent_stream.rs @@ -75,7 +75,14 @@ async fn run(gateway_url: String, api_key: String, agent_id: Uuid, tx: mpsc::Unb }, ); - match tokio_tungstenite::connect_async_tls_with_config(request, None, false, connector.clone()).await { + match tokio_tungstenite::connect_async_tls_with_config( + request, + None, + false, + connector.clone(), + ) + .await + { Ok((socket, _)) => { info!(agent_id = %agent_id, "agent stream connected"); delay = MIN_RECONNECT_DELAY; diff --git a/agent/src/docker.rs b/agent/src/docker.rs index 8c9fe3fc..fddf146d 100644 --- a/agent/src/docker.rs +++ b/agent/src/docker.rs @@ -251,30 +251,30 @@ impl DockerRuntime { let existing = self .docker - .inspect_container(&container_name, None::.map(|b| b.build())) + .inspect_container( + &container_name, + None::.map(|b| b.build()), + ) .await; - match existing { - Ok(inspect) => { - let running = inspect - .state - .and_then(|s| s.status) - .map(|status| status == ContainerStateStatusEnum::RUNNING) - .unwrap_or(false); - if running { - return Ok(()); - } - info!(resource_group_id = %resource_group_id, "Resource group dns container exists but not running, recreating"); - let remove_options = - bollard::query_parameters::RemoveContainerOptionsBuilder::default() - .force(true) - .build(); - self.docker - .remove_container(&container_name, Some(remove_options)) - .await - .context("Failed to remove stale resource group dns container")?; + if let Ok(inspect) = existing { + let running = inspect + .state + .and_then(|s| s.status) + .map(|status| status == ContainerStateStatusEnum::RUNNING) + .unwrap_or(false); + if running { + return Ok(()); } - Err(_) => {} + info!(resource_group_id = %resource_group_id, "Resource group dns container exists but not running, recreating"); + let remove_options = + bollard::query_parameters::RemoveContainerOptionsBuilder::default() + .force(true) + .build(); + self.docker + .remove_container(&container_name, Some(remove_options)) + .await + .context("Failed to remove stale resource group dns container")?; } crate::rg_dns::write_corefile(resource_group_id, &dns_ip) @@ -310,10 +310,7 @@ impl DockerRuntime { )])), }; - let corefile_container_path = format!( - "/zones/{}.Corefile", - resource_group_id - ); + let corefile_container_path = format!("/zones/{}.Corefile", resource_group_id); let config = ContainerCreateBody { image: Some(DNS_IMAGE.to_string()), @@ -822,8 +819,7 @@ impl crate::runtime::Runtime for DockerRuntime { workload_handle: &str, network_name: &str, ) -> Result> { - self.inspect_network_ip(workload_handle, network_name) - .await + self.inspect_network_ip(workload_handle, network_name).await } async fn list_managed_workloads(&self) -> Result> { diff --git a/agent/src/main.rs b/agent/src/main.rs index 9d7038f5..d9bcf6b4 100644 --- a/agent/src/main.rs +++ b/agent/src/main.rs @@ -273,7 +273,19 @@ async fn reconcile_tick( current_flake_rev: &mut String, ) { process_volumes(client, agent_id, api_key, mounted_volumes).await; - let resource_group_ids = process_workloads(client, api_key, docker, running_containers, workload_phases, mounted_volumes, restart_counts, service_dns_registry, rg_dns_registry, wg_private_key_b64).await; + let resource_group_ids = process_workloads( + client, + api_key, + docker, + running_containers, + workload_phases, + mounted_volumes, + restart_counts, + service_dns_registry, + rg_dns_registry, + wg_private_key_b64, + ) + .await; sync_wireguard_peers(client, api_key, agent_id, &resource_group_ids).await; sync_vpn_peers(client, api_key, &resource_group_ids).await; cleanup_stale_resource_groups(client, api_key, wg_private_key_b64).await; @@ -282,7 +294,10 @@ async fn reconcile_tick( push_workload_stats(client, api_key, docker, running_containers).await; let metrics = system::collect_metrics(); - match client.heartbeat(agent_id, api_key, Some(statuses), Some(metrics)).await { + match client + .heartbeat(agent_id, api_key, Some(statuses), Some(metrics)) + .await + { Ok(resp) => { if *failure_count > 0 { info!(agent_id = %agent_id, "Heartbeat recovered after {} failures", failure_count); @@ -465,7 +480,10 @@ async fn reap_stale_containers( let dns_entry = service_dns_registry.lock().await.remove(&workload_id); if let Some((resource_group_id, service_name)) = dns_entry { - if let Err(e) = rg_dns_registry.remove(&resource_group_id, &service_name).await { + if let Err(e) = rg_dns_registry + .remove(&resource_group_id, &service_name) + .await + { warn!(workload_id = %workload_id, error = %e, "Failed to remove service dns record"); } } @@ -572,7 +590,10 @@ async fn process_workloads( for (workload_id, container_id) in existing { containers.insert(workload_id, container_id); } - info!(count = containers.len(), "Reconciled running containers from docker state"); + info!( + count = containers.len(), + "Reconciled running containers from docker state" + ); } Err(e) => { warn!(error = %e, "Failed to reconcile running containers from docker state"); @@ -770,7 +791,10 @@ async fn register_service_dns( ) { let network_name = docker::rg_network_name(resource_group_id); - let ip_address = match runtime.service_network_ip(container_id, &network_name).await { + let ip_address = match runtime + .service_network_ip(container_id, &network_name) + .await + { Ok(Some(ip)) => ip, Ok(None) => { warn!(container_id = %container_id, network = %network_name, "No network ip found for service dns registration"); diff --git a/agent/src/rg_dns.rs b/agent/src/rg_dns.rs index a051744a..4ca50bec 100644 --- a/agent/src/rg_dns.rs +++ b/agent/src/rg_dns.rs @@ -62,10 +62,7 @@ impl RgDnsRegistry { } } -async fn write_zone_file( - resource_group_id: &str, - records: &HashMap, -) -> Result<()> { +async fn write_zone_file(resource_group_id: &str, records: &HashMap) -> Result<()> { fs::create_dir_all(DNS_DIR) .await .context("Failed to create dns zone directory")?; diff --git a/control-plane/api-gateway/src/routes/agent_stream.rs b/control-plane/api-gateway/src/routes/agent_stream.rs index 7171731b..25741655 100644 --- a/control-plane/api-gateway/src/routes/agent_stream.rs +++ b/control-plane/api-gateway/src/routes/agent_stream.rs @@ -56,9 +56,7 @@ impl AgentStreamRegistry { pub async fn notify_assignment(&self, agent_id: Uuid) -> bool { let senders = self.senders.lock().await; match senders.get(&agent_id) { - Some(tx) => tx - .send(Message::Text(ASSIGNMENT_SIGNAL.into())) - .is_ok(), + Some(tx) => tx.send(Message::Text(ASSIGNMENT_SIGNAL.into())).is_ok(), None => false, } } @@ -116,7 +114,10 @@ pub async fn notify_assignment_handler( return Err(StatusCode::FORBIDDEN); } - let delivered = state.agent_stream_registry.notify_assignment(agent_id).await; + let delivered = state + .agent_stream_registry + .notify_assignment(agent_id) + .await; info!(agent_id = %agent_id, delivered, "assignment notification processed"); Ok(StatusCode::NO_CONTENT) } diff --git a/control-plane/registry/src/services/pki.rs b/control-plane/registry/src/services/pki.rs index 5e756b67..a28875c7 100644 --- a/control-plane/registry/src/services/pki.rs +++ b/control-plane/registry/src/services/pki.rs @@ -62,9 +62,10 @@ impl PkiService { }); } - if let (Ok(cert_pem), Ok(key_pem)) = - (std::fs::read_to_string(CA_CERT_PATH), std::fs::read_to_string(CA_KEY_PATH)) - { + if let (Ok(cert_pem), Ok(key_pem)) = ( + std::fs::read_to_string(CA_CERT_PATH), + std::fs::read_to_string(CA_KEY_PATH), + ) { KeyPair::from_pem(&key_pem).map_err(|e| anyhow!("Failed to load CA key: {}", e))?; crate::log_info!("pki", "CA loaded from disk"); return Ok(CaState { @@ -74,7 +75,10 @@ impl PkiService { }); } - crate::log_warn!("pki", "no persisted CA found, generating and persisting new CA"); + crate::log_warn!( + "pki", + "no persisted CA found, generating and persisting new CA" + ); let ca = Self::generate_ca()?; std::fs::create_dir_all(CA_STATE_DIR) .map_err(|e| anyhow!("Failed to create CA state dir: {}", e))?; diff --git a/control-plane/registry/src/services/registry.rs b/control-plane/registry/src/services/registry.rs index 41afecf1..50769c48 100644 --- a/control-plane/registry/src/services/registry.rs +++ b/control-plane/registry/src/services/registry.rs @@ -271,7 +271,10 @@ impl AgentRegistry { crate::log_info!( "agent_registry", - &format!("Backfilled management tunnel IP agent={} ip={}", agent_id, ip) + &format!( + "Backfilled management tunnel IP agent={} ip={}", + agent_id, ip + ) ); } diff --git a/control-plane/scheduler/src/handlers/workloads.rs b/control-plane/scheduler/src/handlers/workloads.rs index 7e02f3fe..f23aa176 100644 --- a/control-plane/scheduler/src/handlers/workloads.rs +++ b/control-plane/scheduler/src/handlers/workloads.rs @@ -58,9 +58,7 @@ pub async fn stop_workload( match db::workloads::set_desired_state(&state.db, id, DesiredState::Stopped).await { Ok(model) => { if let Some(agent_id) = model.assigned_agent_id { - tokio::spawn(crate::services::gateway_notify::notify_assignment( - agent_id, - )); + tokio::spawn(crate::services::gateway_notify::notify_assignment(agent_id)); } (StatusCode::OK, Json(serde_json::json!(model))).into_response() } @@ -80,9 +78,7 @@ pub async fn restart_workload( match db::workloads::request_restart(&state.db, id).await { Ok(model) => { if let Some(agent_id) = model.assigned_agent_id { - tokio::spawn(crate::services::gateway_notify::notify_assignment( - agent_id, - )); + tokio::spawn(crate::services::gateway_notify::notify_assignment(agent_id)); } (StatusCode::OK, Json(serde_json::json!(model))).into_response() } diff --git a/control-plane/scheduler/src/services/gateway_notify.rs b/control-plane/scheduler/src/services/gateway_notify.rs index 1f779f0d..8111e063 100644 --- a/control-plane/scheduler/src/services/gateway_notify.rs +++ b/control-plane/scheduler/src/services/gateway_notify.rs @@ -4,8 +4,8 @@ const API_GATEWAY_URL_ENV: &str = "API_GATEWAY_URL"; const DEFAULT_API_GATEWAY_URL: &str = "http://localhost:8000"; pub async fn notify_assignment(agent_id: Uuid) { - let base_url = std::env::var(API_GATEWAY_URL_ENV) - .unwrap_or_else(|_| DEFAULT_API_GATEWAY_URL.to_string()); + let base_url = + std::env::var(API_GATEWAY_URL_ENV).unwrap_or_else(|_| DEFAULT_API_GATEWAY_URL.to_string()); let url = format!( "{}/api/internal/agents/{}/notify-assignment", base_url, agent_id diff --git a/control-plane/scheduler/src/services/scheduler.rs b/control-plane/scheduler/src/services/scheduler.rs index 33c06958..9562ef7c 100644 --- a/control-plane/scheduler/src/services/scheduler.rs +++ b/control-plane/scheduler/src/services/scheduler.rs @@ -68,7 +68,10 @@ impl SchedulerService { if let Some(mut ports) = req.ports.clone() { crate::services::port_allocator::allocate_node_ports( - &self.db, agent_id, workload.id, &mut ports, + &self.db, + agent_id, + workload.id, + &mut ports, ) .await?; let ports_json = crate::services::port_allocator::ports_to_json(&ports)?; @@ -90,9 +93,7 @@ impl SchedulerService { runtime_class: req.runtime_class.as_str().to_string(), }; - tokio::spawn(crate::services::gateway_notify::notify_assignment( - agent_id, - )); + tokio::spawn(crate::services::gateway_notify::notify_assignment(agent_id)); put_placement(&self.etcd, &record).await?; From 6da42500c582c412eda40ffe9a64dd715f2d2b1d Mon Sep 17 00:00:00 2001 From: CodeMaster4711 Date: Thu, 23 Jul 2026 20:33:22 +0200 Subject: [PATCH 3/7] feat: add container settings edit and fix compose port and status bugs --- agent/src/main.rs | 21 +++- app/src/lib/api/resource-groups.ts | 28 +++++ .../(app)/resource-groups/[id]/+page.svelte | 103 +++++++++++++++++- .../api-gateway/src/routes/workloads.rs | 22 +++- .../api-gateway/src/service_client.rs | 1 + control-plane/scheduler/src/db/workloads.rs | 42 +++++++ .../scheduler/src/handlers/workloads.rs | 25 ++++- .../scheduler/src/models/workload.rs | 10 ++ control-plane/scheduler/src/server.rs | 4 + .../scheduler/src/services/compose_parser.rs | 6 +- 10 files changed, 251 insertions(+), 11 deletions(-) diff --git a/agent/src/main.rs b/agent/src/main.rs index d9bcf6b4..5d2db51a 100644 --- a/agent/src/main.rs +++ b/agent/src/main.rs @@ -713,6 +713,17 @@ async fn start_or_restart_workload( } } + if let Some(max) = workload.max_restarts { + let count = *restart_counts.lock().await.get(&workload.id).unwrap_or(&0); + if count as i32 >= max { + workload_phases + .lock() + .await + .insert(workload.id.clone(), "failed".to_string()); + return false; + } + } + workload_phases .lock() .await @@ -775,8 +786,14 @@ async fn start_or_restart_workload( restart_requested } Err(e) => { - workload_phases.lock().await.remove(&workload.id); - warn!(workload_id = %workload.id, error = ?e, "Failed to start workload"); + let mut counts = restart_counts.lock().await; + let count = counts.entry(workload.id.clone()).or_insert(0); + *count += 1; + warn!(workload_id = %workload.id, error = ?e, attempt = *count, "Failed to start workload"); + workload_phases + .lock() + .await + .insert(workload.id.clone(), "failed".to_string()); false } } diff --git a/app/src/lib/api/resource-groups.ts b/app/src/lib/api/resource-groups.ts index 4d1b3939..ffc221a5 100644 --- a/app/src/lib/api/resource-groups.ts +++ b/app/src/lib/api/resource-groups.ts @@ -52,6 +52,7 @@ export interface Workload { status: string; assigned_agent_id: string | null; container_id: string | null; + env_vars: Record | null; ports: PortMapping[] | null; volume_mounts: VolumeMount[] | null; resource_group_id: string | null; @@ -84,6 +85,13 @@ export interface CreateWorkloadRequest { max_restarts?: number | null; } +export interface UpdateWorkloadRequest { + env_vars?: Record | null; + ports?: PortMapping[] | null; + restart_policy?: 'always' | 'on-failure' | 'never'; + max_restarts?: number | null; +} + export interface CreateVolumeRequest { name: string; size_gb: number; @@ -199,6 +207,26 @@ export async function deleteWorkload(token: string, id: string): Promise { if (!res.ok) throw new Error(`Failed to delete workload: ${res.status}`); } +export async function updateWorkload( + token: string, + id: string, + req: UpdateWorkloadRequest, +): Promise { + const res = await authedFetch(`${API_BASE}/workloads/${id}`, { + method: 'PATCH', + headers: { + Authorization: `Bearer ${token}`, + 'Content-Type': 'application/json', + }, + body: JSON.stringify(req), + }); + if (!res.ok) { + const err = await res.json().catch(() => ({ error: res.status })); + throw new Error(err.error ?? `Failed to update workload: ${res.status}`); + } + return res.json(); +} + export async function stopWorkload(token: string, id: string): Promise { const res = await authedFetch(`${API_BASE}/workloads/${id}/stop`, { method: 'POST', diff --git a/app/src/routes/(app)/resource-groups/[id]/+page.svelte b/app/src/routes/(app)/resource-groups/[id]/+page.svelte index 85a21d0b..73317c91 100644 --- a/app/src/routes/(app)/resource-groups/[id]/+page.svelte +++ b/app/src/routes/(app)/resource-groups/[id]/+page.svelte @@ -12,6 +12,7 @@ deleteWorkload, stopWorkload, restartWorkload, + updateWorkload, createVolume, deleteVolume, streamWorkloadLogs, @@ -124,11 +125,17 @@ let volFormSize = $state("10"); let containerDialog = $state(null); - let containerDialogTab = $state<"logs" | "shell" | "insights" | "network">("logs"); + let containerDialogTab = $state<"logs" | "shell" | "insights" | "network" | "settings">("logs"); let activeContainer = $state(null); let containerActionError = $state(null); let containerActionBusy = $state(false); + let settingsEnvText = $state(""); + let settingsRestartPolicy = $state<"always" | "on-failure" | "never">("always"); + let settingsMaxRestarts = $state(""); + let settingsError = $state(null); + let settingsSaving = $state(false); + let logsLines = $state([]); let logsError = $state(null); let logsAbort: AbortController | null = null; @@ -306,6 +313,15 @@ } } + function loadSettingsForm(workload: Workload) { + settingsEnvText = Object.entries(workload.env_vars ?? {}) + .map(([k, v]) => `${k}=${v}`) + .join("\n"); + settingsRestartPolicy = workload.restart_policy as "always" | "on-failure" | "never"; + settingsMaxRestarts = workload.max_restarts !== null ? String(workload.max_restarts) : ""; + settingsError = null; + } + function openContainer(workload: Workload) { activeContainer = workload; containerDialogTab = "logs"; @@ -313,6 +329,7 @@ containerDialog?.showModal(); startLogsStream(workload); resolveNodeIp(workload.assigned_agent_id); + loadSettingsForm(workload); } function closeContainer() { @@ -322,7 +339,7 @@ activeContainer = null; } - function switchTab(tab: "logs" | "shell" | "insights" | "network") { + function switchTab(tab: "logs" | "shell" | "insights" | "network" | "settings") { if (containerDialogTab === tab || !activeContainer) return; if (containerDialogTab === "logs") stopLogsStream(); @@ -449,6 +466,39 @@ } } + async function handleSaveSettings() { + if (!auth.token || !activeContainer) return; + settingsSaving = true; + settingsError = null; + try { + const env_vars: Record = {}; + for (const line of settingsEnvText.split("\n")) { + const trimmed = line.trim(); + if (!trimmed) continue; + const idx = trimmed.indexOf("="); + if (idx === -1) throw new Error(`Invalid env var line: "${trimmed}" (expected KEY=VALUE)`); + env_vars[trimmed.slice(0, idx)] = trimmed.slice(idx + 1); + } + const max_restarts = settingsMaxRestarts.trim() === "" ? null : Number(settingsMaxRestarts); + if (max_restarts !== null && (!Number.isInteger(max_restarts) || max_restarts < 0)) { + throw new Error("Max restarts must be a non-negative integer"); + } + + await updateWorkload(auth.token, activeContainer.id, { + env_vars, + restart_policy: settingsRestartPolicy, + max_restarts, + }); + workloads = await listResourceGroupWorkloads(auth.token, rgId); + activeContainer = workloads.find((w) => w.id === activeContainer?.id) ?? null; + if (activeContainer) loadSettingsForm(activeContainer); + } catch (e) { + settingsError = e instanceof Error ? e.message : "Failed to save settings"; + } finally { + settingsSaving = false; + } + } + async function handleDeleteContainer() { if (!auth.token || !activeContainer) return; containerActionBusy = true; @@ -1047,7 +1097,7 @@ {/if}
- {#each [["logs", "Logs"], ["shell", "Shell"], ["insights", "Performance"], ["network", "Network"]] as [tab, label]} + {#each [["logs", "Logs"], ["shell", "Shell"], ["insights", "Performance"], ["network", "Network"], ["settings", "Settings"]] as [tab, label]}
- {:else} + {:else if containerDialogTab === "network"}
{#if activeContainer.ports && activeContainer.ports.length > 0}
@@ -1187,6 +1237,51 @@

No ports configured for this container.

{/if}
+ {:else} +
+
+ + +

One KEY=VALUE per line.

+
+
+
+ + +
+
+ + +
+
+

+ Saving applies the new configuration by restarting the container. +

+ {#if settingsError} +

{settingsError}

+ {/if} +
+ +
+
{/if}
diff --git a/control-plane/api-gateway/src/routes/workloads.rs b/control-plane/api-gateway/src/routes/workloads.rs index 5dfd44a4..673f0d6b 100644 --- a/control-plane/api-gateway/src/routes/workloads.rs +++ b/control-plane/api-gateway/src/routes/workloads.rs @@ -3,7 +3,7 @@ use axum::{ extract::{Path, State}, http::{HeaderMap, StatusCode}, response::{IntoResponse, Json}, - routing::{delete, get, post}, + routing::{delete, get, patch, post}, Router, }; use serde_json::json; @@ -96,6 +96,25 @@ pub async fn delete_workload( .await } +pub async fn update_workload( + CanManageWorkloads(_claims): CanManageWorkloads, + State(state): State, + Path(id): Path, + headers: HeaderMap, + body: String, +) -> Result)> { + let body_json: Option = serde_json::from_str(&body).ok(); + let header_map = header_vec(&headers); + proxy_to_scheduler( + &state, + reqwest::Method::PATCH, + &format!("/workloads/{}", id), + body_json, + Some(header_map), + ) + .await +} + pub async fn stop_workload( CanManageWorkloads(_claims): CanManageWorkloads, State(state): State, @@ -160,6 +179,7 @@ pub fn workloads_routes() -> Router { .route("/workloads", post(create_workload)) .route("/workloads", get(list_workloads)) .route("/workloads/{id}", delete(delete_workload)) + .route("/workloads/{id}", patch(update_workload)) .route("/workloads/{id}/stop", post(stop_workload)) .route("/workloads/{id}/restart", post(restart_workload)) .route("/workload-stacks", post(create_stack)) diff --git a/control-plane/api-gateway/src/service_client.rs b/control-plane/api-gateway/src/service_client.rs index b2aa704a..1d096154 100644 --- a/control-plane/api-gateway/src/service_client.rs +++ b/control-plane/api-gateway/src/service_client.rs @@ -129,6 +129,7 @@ impl ServiceClient { reqwest::Method::GET => self.client.get(&url), reqwest::Method::POST => self.client.post(&url), reqwest::Method::DELETE => self.client.delete(&url), + reqwest::Method::PATCH => self.client.patch(&url), _ => return Err(anyhow::anyhow!("Unsupported HTTP method")), }; diff --git a/control-plane/scheduler/src/db/workloads.rs b/control-plane/scheduler/src/db/workloads.rs index fca28bd3..bf038795 100644 --- a/control-plane/scheduler/src/db/workloads.rs +++ b/control-plane/scheduler/src/db/workloads.rs @@ -100,6 +100,38 @@ pub async fn update_ports( Ok(()) } +pub async fn update_spec( + db: &DatabaseConnection, + workload_id: Uuid, + req: &crate::models::workload::UpdateWorkloadRequest, +) -> Result { + let workload = workloads::Entity::find_by_id(workload_id) + .one(db) + .await? + .ok_or(sea_orm::DbErr::RecordNotFound(workload_id.to_string()))?; + + let mut active: workloads::ActiveModel = workload.into(); + + if let Some(ref env_vars) = req.env_vars { + active.env_vars = Set(serde_json::to_value(env_vars).ok()); + } + if let Some(ref ports) = req.ports { + active.ports = Set(serde_json::to_value(ports).ok()); + } + if let Some(ref restart_policy) = req.restart_policy { + active.restart_policy = Set(restart_policy.as_str().to_string()); + } + if req.max_restarts.is_some() { + active.max_restarts = Set(req.max_restarts); + } + + active.restart_requested = Set(true); + active.status = Set(WorkloadStatus::Starting.as_str().to_string()); + active.updated_at = Set(Some(Utc::now().naive_utc())); + + active.update(db).await +} + pub async fn get_all(db: &DatabaseConnection) -> Result, sea_orm::DbErr> { let rows = workloads::Entity::find().all(db).await?; Ok(rows.into_iter().map(into_response).collect()) @@ -232,6 +264,14 @@ pub async fn delete(db: &DatabaseConnection, workload_id: Uuid) -> Result<(), se } fn into_response(m: workloads::Model) -> WorkloadResponse { + let env_vars = m + .env_vars + .as_ref() + .and_then(|e| serde_json::from_value(e.clone()).ok()); + let ports = m + .ports + .as_ref() + .and_then(|p| serde_json::from_value(p.clone()).ok()); let volume_mounts = m .volume_mounts .as_ref() @@ -247,6 +287,8 @@ fn into_response(m: workloads::Model) -> WorkloadResponse { status: WorkloadStatus::from_str(&m.status), assigned_agent_id: m.assigned_agent_id, container_id: m.container_id, + env_vars, + ports, volume_mounts, resource_group_id: m.resource_group_id, stack_id: m.stack_id, diff --git a/control-plane/scheduler/src/handlers/workloads.rs b/control-plane/scheduler/src/handlers/workloads.rs index f23aa176..ef9613c1 100644 --- a/control-plane/scheduler/src/handlers/workloads.rs +++ b/control-plane/scheduler/src/handlers/workloads.rs @@ -8,7 +8,7 @@ use uuid::Uuid; use crate::{ db, - models::workload::{CreateWorkloadRequest, DesiredState}, + models::workload::{CreateWorkloadRequest, DesiredState, UpdateWorkloadRequest}, server::AppState, }; @@ -71,6 +71,29 @@ pub async fn stop_workload( } } +pub async fn update_workload( + State(state): State, + Path(id): Path, + Json(req): Json, +) -> impl IntoResponse { + match db::workloads::update_spec(&state.db, id, &req).await { + Ok(model) => { + if let Some(agent_id) = model.assigned_agent_id { + tokio::spawn(crate::services::gateway_notify::notify_assignment( + agent_id, + )); + } + (StatusCode::OK, Json(serde_json::json!(model))).into_response() + } + Err(sea_orm::DbErr::RecordNotFound(_)) => StatusCode::NOT_FOUND.into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({ "error": e.to_string() })), + ) + .into_response(), + } +} + pub async fn restart_workload( State(state): State, Path(id): Path, diff --git a/control-plane/scheduler/src/models/workload.rs b/control-plane/scheduler/src/models/workload.rs index 3b3f3bf4..0c189f92 100644 --- a/control-plane/scheduler/src/models/workload.rs +++ b/control-plane/scheduler/src/models/workload.rs @@ -137,6 +137,14 @@ impl DesiredState { } } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct UpdateWorkloadRequest { + pub env_vars: Option>, + pub ports: Option>, + pub restart_policy: Option, + pub max_restarts: Option, +} + #[derive(Debug, Clone, Serialize, Deserialize)] pub struct CreateWorkloadRequest { pub name: String, @@ -187,6 +195,8 @@ pub struct WorkloadResponse { pub status: WorkloadStatus, pub assigned_agent_id: Option, pub container_id: Option, + pub env_vars: Option>, + pub ports: Option>, pub volume_mounts: Option>, pub resource_group_id: Option, pub stack_id: Option, diff --git a/control-plane/scheduler/src/server.rs b/control-plane/scheduler/src/server.rs index 58738673..01a3325a 100644 --- a/control-plane/scheduler/src/server.rs +++ b/control-plane/scheduler/src/server.rs @@ -34,6 +34,10 @@ pub fn create_router(state: AppState) -> Router { "/workloads/{id}", axum::routing::delete(workloads::delete_workload), ) + .route( + "/workloads/{id}", + axum::routing::patch(workloads::update_workload), + ) .route( "/workloads/{id}/stop", axum::routing::post(workloads::stop_workload), diff --git a/control-plane/scheduler/src/services/compose_parser.rs b/control-plane/scheduler/src/services/compose_parser.rs index 3aa7f7c1..d03f57b7 100644 --- a/control-plane/scheduler/src/services/compose_parser.rs +++ b/control-plane/scheduler/src/services/compose_parser.rs @@ -96,7 +96,7 @@ fn parse_compose_ports(raw: &[String]) -> Vec { } fn parse_compose_port(entry: &str) -> Option { - let (node_port, container_port) = match entry.split_once(':') { + let (rg_port, container_port) = match entry.split_once(':') { Some((host, container)) => (host.parse().ok(), container), None => (None, entry), }; @@ -104,7 +104,7 @@ fn parse_compose_port(entry: &str) -> Option { Some(PortMapping { container_port: container_port.parse().ok()?, protocol: None, - rg_port: None, - node_port, + rg_port, + node_port: None, }) } From ff3b1e72e43507e65e1e9f68e770a8a56e807fa8 Mon Sep 17 00:00:00 2001 From: CodeMaster4711 Date: Thu, 23 Jul 2026 20:45:43 +0200 Subject: [PATCH 4/7] feat: add redeploy and stack management for containers and compose stacks --- app/src/lib/api/resource-groups.ts | 50 +++++++ .../(app)/resource-groups/[id]/+page.svelte | 134 +++++++++++++++--- .../api-gateway/src/routes/workloads.rs | 74 ++++++++++ .../scheduler/src/db/workload_stacks.rs | 32 +++++ control-plane/scheduler/src/db/workloads.rs | 21 +++ .../scheduler/src/handlers/stacks.rs | 83 ++++++++++- control-plane/scheduler/src/models/compose.rs | 16 +++ .../scheduler/src/models/workload.rs | 1 + control-plane/scheduler/src/server.rs | 16 +++ .../scheduler/src/services/scheduler.rs | 83 ++++++++++- 10 files changed, 485 insertions(+), 25 deletions(-) diff --git a/app/src/lib/api/resource-groups.ts b/app/src/lib/api/resource-groups.ts index ffc221a5..4189a22b 100644 --- a/app/src/lib/api/resource-groups.ts +++ b/app/src/lib/api/resource-groups.ts @@ -86,6 +86,7 @@ export interface CreateWorkloadRequest { } export interface UpdateWorkloadRequest { + image?: string; env_vars?: Record | null; ports?: PortMapping[] | null; restart_policy?: 'always' | 'on-failure' | 'never'; @@ -199,6 +200,55 @@ export async function createWorkloadStack( return res.json(); } +export interface Stack { + id: string; + resource_group_id: string; + name: string; + compose_source: string | null; + status: string; + created_at: string; + updated_at: string | null; +} + +export async function getStack(token: string, id: string): Promise { + const res = await authedFetch(`${API_BASE}/workload-stacks/${id}`, { + headers: { Authorization: `Bearer ${token}` }, + }); + if (!res.ok) throw new Error(`Failed to get stack: ${res.status}`); + return res.json(); +} + +export async function deleteStack(token: string, id: string): Promise { + const res = await authedFetch(`${API_BASE}/workload-stacks/${id}`, { + method: 'DELETE', + headers: { Authorization: `Bearer ${token}` }, + }); + if (!res.ok) throw new Error(`Failed to delete stack: ${res.status}`); +} + +export async function restartStack(token: string, id: string): Promise { + const res = await authedFetch(`${API_BASE}/workload-stacks/${id}/restart`, { + method: 'POST', + headers: { Authorization: `Bearer ${token}` }, + }); + if (!res.ok) throw new Error(`Failed to restart stack: ${res.status}`); +} + +export async function redeployStack(token: string, id: string, compose_yaml: string): Promise { + const res = await authedFetch(`${API_BASE}/workload-stacks/${id}`, { + method: 'PATCH', + headers: { + Authorization: `Bearer ${token}`, + 'Content-Type': 'application/json', + }, + body: JSON.stringify({ compose_yaml }), + }); + if (!res.ok) { + const err = await res.json().catch(() => ({ error: res.status })); + throw new Error(err.error ?? `Failed to redeploy stack: ${res.status}`); + } +} + export async function deleteWorkload(token: string, id: string): Promise { const res = await authedFetch(`${API_BASE}/workloads/${id}`, { method: 'DELETE', diff --git a/app/src/routes/(app)/resource-groups/[id]/+page.svelte b/app/src/routes/(app)/resource-groups/[id]/+page.svelte index 73317c91..13593080 100644 --- a/app/src/routes/(app)/resource-groups/[id]/+page.svelte +++ b/app/src/routes/(app)/resource-groups/[id]/+page.svelte @@ -9,6 +9,10 @@ listResourceGroupVolumes, createWorkload, createWorkloadStack, + getStack, + deleteStack, + restartStack, + redeployStack, deleteWorkload, stopWorkload, restartWorkload, @@ -105,6 +109,7 @@ let composeStackName = $state(""); let composeYaml = $state(""); + let editingStackId = $state(null); let composePreview = $derived(parseComposePreview(composeYaml)); let composeLineCount = $derived(Math.max(composeYaml.split("\n").length, 1)); let composeGutter = $state(null); @@ -130,6 +135,7 @@ let containerActionError = $state(null); let containerActionBusy = $state(false); + let settingsImage = $state(""); let settingsEnvText = $state(""); let settingsRestartPolicy = $state<"always" | "on-failure" | "never">("always"); let settingsMaxRestarts = $state(""); @@ -239,15 +245,20 @@ } async function handleDeployStack() { - if (!auth.token || !composeYaml.trim() || !composeStackName) return; + if (!auth.token || !composeYaml.trim()) return; + if (!editingStackId && !composeStackName) return; deployingStack = true; composeError = null; try { - await createWorkloadStack(auth.token, { - name: composeStackName, - resource_group_id: rgId, - compose_yaml: composeYaml, - }); + if (editingStackId) { + await redeployStack(auth.token, editingStackId, composeYaml); + } else { + await createWorkloadStack(auth.token, { + name: composeStackName, + resource_group_id: rgId, + compose_yaml: composeYaml, + }); + } composeDialog?.close(); resetComposeForm(); workloads = await listResourceGroupWorkloads(auth.token, rgId); @@ -262,6 +273,43 @@ composeStackName = ""; composeYaml = ""; composeError = null; + editingStackId = null; + } + + async function openStackEditor(stackId: string) { + if (!auth.token) return; + composeError = null; + editingStackId = stackId; + composeStackName = ""; + composeYaml = ""; + composeDialog?.showModal(); + try { + const stack = await getStack(auth.token, stackId); + composeStackName = stack.name; + composeYaml = stack.compose_source ?? ""; + } catch (e) { + composeError = e instanceof Error ? e.message : "Failed to load stack"; + } + } + + async function handleRestartStack(stackId: string) { + if (!auth.token) return; + try { + await restartStack(auth.token, stackId); + workloads = await listResourceGroupWorkloads(auth.token, rgId); + } catch (e) { + error = e instanceof Error ? e.message : "Failed to restart stack"; + } + } + + async function handleDeleteStack(stackId: string) { + if (!auth.token) return; + try { + await deleteStack(auth.token, stackId); + workloads = workloads.filter((w) => w.stack_id !== stackId); + } catch (e) { + error = e instanceof Error ? e.message : "Failed to delete stack"; + } } function handleComposeFileUpload(event: Event) { @@ -314,6 +362,7 @@ } function loadSettingsForm(workload: Workload) { + settingsImage = workload.image; settingsEnvText = Object.entries(workload.env_vars ?? {}) .map(([k, v]) => `${k}=${v}`) .join("\n"); @@ -466,11 +515,12 @@ } } - async function handleSaveSettings() { + async function handleRedeployContainer() { if (!auth.token || !activeContainer) return; settingsSaving = true; settingsError = null; try { + if (!settingsImage.trim()) throw new Error("Image cannot be empty"); const env_vars: Record = {}; for (const line of settingsEnvText.split("\n")) { const trimmed = line.trim(); @@ -485,6 +535,7 @@ } await updateWorkload(auth.token, activeContainer.id, { + image: settingsImage.trim(), env_vars, restart_policy: settingsRestartPolicy, max_restarts, @@ -906,7 +957,7 @@ >
-

Deploy Docker Compose Stack

+

{editingStackId ? "Edit Compose Stack" : "Deploy Docker Compose Stack"}

- +
@@ -989,8 +1040,12 @@ {/if}
-
@@ -1239,6 +1294,15 @@
{:else}
+
+ + +