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..5d2db51a 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"); @@ -692,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 @@ -754,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 } } @@ -770,7 +808,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 a9db6df8..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")?; @@ -97,11 +94,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); diff --git a/app/src/lib/api/resource-groups.ts b/app/src/lib/api/resource-groups.ts index 4d1b3939..fca551ff 100644 --- a/app/src/lib/api/resource-groups.ts +++ b/app/src/lib/api/resource-groups.ts @@ -9,6 +9,9 @@ export interface ResourceGroup { description: string | null; internal_cidr: string; status: string; + icon: string; + color: string; + pinned: boolean; created_at: string; updated_at: string | null; } @@ -17,6 +20,16 @@ export interface CreateResourceGroupRequest { name: string; description?: string; internal_cidr: string; + icon?: string; + color?: string; +} + +export interface UpdateResourceGroupRequest { + name?: string; + description?: string; + icon?: string; + color?: string; + pinned?: boolean; } export interface PortMapping { @@ -52,6 +65,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 +98,14 @@ export interface CreateWorkloadRequest { max_restarts?: number | null; } +export interface UpdateWorkloadRequest { + image?: string; + 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; @@ -116,6 +138,15 @@ export async function listResourceGroups(token: string): Promise { + const res = await authedFetch(`${API_BASE}/resource-groups/suggest-cidr`, { + headers: { Authorization: `Bearer ${token}` }, + }); + if (!res.ok) throw new Error(`Failed to suggest cidr: ${res.status}`); + const data = await res.json(); + return data.internal_cidr; +} + export async function getResourceGroup(token: string, id: string): Promise { const res = await authedFetch(`${API_BASE}/resource-groups/${id}`, { headers: { Authorization: `Bearer ${token}` }, @@ -140,6 +171,26 @@ export async function createResourceGroup( return res.json(); } +export async function updateResourceGroup( + token: string, + id: string, + req: UpdateResourceGroupRequest, +): Promise { + const res = await authedFetch(`${API_BASE}/resource-groups/${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 resource group: ${res.status}`); + } + return res.json(); +} + export async function deleteResourceGroup(token: string, id: string): Promise { const res = await authedFetch(`${API_BASE}/resource-groups/${id}`, { method: 'DELETE', @@ -191,6 +242,63 @@ 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 stopStack(token: string, id: string): Promise { + const res = await authedFetch(`${API_BASE}/workload-stacks/${id}/stop`, { + method: 'POST', + headers: { Authorization: `Bearer ${token}` }, + }); + if (!res.ok) throw new Error(`Failed to stop 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', @@ -199,6 +307,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/lib/components/icon-picker.svelte b/app/src/lib/components/icon-picker.svelte new file mode 100644 index 00000000..168549c0 --- /dev/null +++ b/app/src/lib/components/icon-picker.svelte @@ -0,0 +1,80 @@ + + +
+
+ +
+
+ + +
+
+
+ {#each SUGGESTED_ICONS as suggestion (suggestion)} + + {/each} +
+
+ +
+ +
+ {#each SUGGESTED_COLORS as suggestion (suggestion)} + + {/each} +
+
+
diff --git a/app/src/lib/components/ui/sidebar/sidebar-inset.svelte b/app/src/lib/components/ui/sidebar/sidebar-inset.svelte index 19a787cb..524a1980 100644 --- a/app/src/lib/components/ui/sidebar/sidebar-inset.svelte +++ b/app/src/lib/components/ui/sidebar/sidebar-inset.svelte @@ -13,7 +13,7 @@
{@render children?.()} diff --git a/app/src/routes/(app)/resource-groups/+page.svelte b/app/src/routes/(app)/resource-groups/+page.svelte index 42d37fcc..da5ea67f 100644 --- a/app/src/routes/(app)/resource-groups/+page.svelte +++ b/app/src/routes/(app)/resource-groups/+page.svelte @@ -4,21 +4,32 @@ import { listResourceGroups, createResourceGroup, - deleteResourceGroup, + updateResourceGroup, + suggestCidr, type ResourceGroup, } from "$lib/api/resource-groups"; import * as Sidebar from "$lib/components/ui/sidebar/index.js"; import { Button } from "$lib/components/ui/button/index.js"; + import Icon from "@iconify/svelte"; + import IconPicker from "$lib/components/icon-picker.svelte"; let groups = $state([]); let loading = $state(true); let error = $state(null); let creating = $state(false); let createDialog = $state(null); + let searchText = $state(""); + let pinnedScroller = $state(null); + + function scrollPinned(direction: -1 | 1) { + pinnedScroller?.scrollBy({ left: direction * 280, behavior: "smooth" }); + } let newName = $state(""); let newCidr = $state("10.100.0.0/24"); let newDescription = $state(""); + let newIcon = $state("mdi:cube-outline"); + let newColor = $state("#6366f1"); let createError = $state(null); async function load() { @@ -32,6 +43,17 @@ } } + async function openCreateDialog() { + if (auth.token) { + try { + newCidr = await suggestCidr(auth.token); + } catch { + // keep default fallback below + } + } + createDialog?.showModal(); + } + async function handleCreate() { if (!auth.token || !newName || !newCidr) return; creating = true; @@ -41,12 +63,16 @@ name: newName, description: newDescription || undefined, internal_cidr: newCidr, + icon: newIcon, + color: newColor, }); groups = [...groups, created]; createDialog?.close(); newName = ""; newCidr = "10.100.0.0/24"; newDescription = ""; + newIcon = "mdi:cube-outline"; + newColor = "#6366f1"; } catch (e) { createError = e instanceof Error ? e.message : "Failed to create resource group"; } finally { @@ -54,13 +80,13 @@ } } - async function handleDelete(id: string) { + async function handleTogglePin(group: ResourceGroup) { if (!auth.token) return; try { - await deleteResourceGroup(auth.token, id); - groups = groups.filter((g) => g.id !== id); + const updated = await updateResourceGroup(auth.token, group.id, { pinned: !group.pinned }); + groups = groups.map((g) => (g.id === updated.id ? updated : g)); } catch (e) { - error = e instanceof Error ? e.message : "Failed to delete resource group"; + error = e instanceof Error ? e.message : "Failed to update resource group"; } } @@ -73,6 +99,11 @@ } } + let pinnedGroups = $derived(groups.filter((g) => g.pinned)); + let filteredGroups = $derived( + groups.filter((g) => !searchText || g.name.toLowerCase().includes(searchText.toLowerCase())), + ); + let loadStarted = false; $effect(() => { @@ -127,6 +158,9 @@ bind:value={newDescription} /> +
+ +
{#if createError}

{createError}

@@ -146,7 +180,7 @@ Resource Groups -
+

Resource Groups

@@ -154,7 +188,7 @@ Isolated namespaces with dedicated internal networks

- + {/each} +
+
+ + +
+
+ {/if} + +
+ + +
+
- - {#if loading} - + - {:else if groups.length === 0} + {:else if filteredGroups.length === 0} - {:else} - {#each groups as group (group.id)} + {#each filteredGroups as group (group.id)} goto(`/resource-groups/${group.id}`)} > - - + - {/each} {/if} diff --git a/app/src/routes/(app)/resource-groups/[id]/+page.svelte b/app/src/routes/(app)/resource-groups/[id]/+page.svelte index 85a21d0b..5c107e7e 100644 --- a/app/src/routes/(app)/resource-groups/[id]/+page.svelte +++ b/app/src/routes/(app)/resource-groups/[id]/+page.svelte @@ -5,13 +5,21 @@ import { auth } from "$lib/auth/store.svelte"; import { getResourceGroup, + updateResourceGroup, + deleteResourceGroup, listResourceGroupWorkloads, listResourceGroupVolumes, createWorkload, createWorkloadStack, + getStack, + deleteStack, + stopStack, + restartStack, + redeployStack, deleteWorkload, stopWorkload, restartWorkload, + updateWorkload, createVolume, deleteVolume, streamWorkloadLogs, @@ -32,6 +40,7 @@ import Icon from "@iconify/svelte"; import * as Sidebar from "$lib/components/ui/sidebar/index.js"; import { Button } from "$lib/components/ui/button/index.js"; + import IconPicker from "$lib/components/icon-picker.svelte"; const rgId: string = $page.params.id; @@ -41,6 +50,18 @@ let loading = $state(true); let error = $state(null); + let appearanceDialog = $state(null); + let editIcon = $state("mdi:cube-outline"); + let editColor = $state("#6366f1"); + let savingAppearance = $state(false); + let appearanceError = $state(null); + + let rgSettingsDialog = $state(null); + let editName = $state(""); + let editDescription = $state(""); + let savingRgSettings = $state(false); + let rgSettingsError = $state(null); + let activeTab = $state<"all" | "container" | "volume">("all"); let filterText = $state(""); @@ -104,6 +125,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); @@ -124,11 +146,18 @@ 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 settingsImage = $state(""); + 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; @@ -232,15 +261,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); @@ -255,6 +289,53 @@ 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 handleStopStack(stackId: string) { + if (!auth.token) return; + try { + await stopStack(auth.token, stackId); + workloads = await listResourceGroupWorkloads(auth.token, rgId); + } catch (e) { + error = e instanceof Error ? e.message : "Failed to stop 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) { @@ -306,6 +387,16 @@ } } + function loadSettingsForm(workload: Workload) { + settingsImage = workload.image; + 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 +404,7 @@ containerDialog?.showModal(); startLogsStream(workload); resolveNodeIp(workload.assigned_agent_id); + loadSettingsForm(workload); } function closeContainer() { @@ -322,7 +414,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 +541,41 @@ } } + 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(); + 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, { + image: settingsImage.trim(), + 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; @@ -642,8 +769,70 @@ return `${bytes} B`; } - function initials(name: string): string { - return name.slice(0, 2).toUpperCase(); + function openAppearanceDialog() { + if (!group) return; + editIcon = group.icon; + editColor = group.color; + appearanceError = null; + appearanceDialog?.showModal(); + } + + async function handleSaveAppearance() { + if (!auth.token || !group) return; + savingAppearance = true; + appearanceError = null; + try { + group = await updateResourceGroup(auth.token, group.id, { icon: editIcon, color: editColor }); + appearanceDialog?.close(); + } catch (e) { + appearanceError = e instanceof Error ? e.message : "Failed to update appearance"; + } finally { + savingAppearance = false; + } + } + + async function handleTogglePin() { + if (!auth.token || !group) return; + try { + group = await updateResourceGroup(auth.token, group.id, { pinned: !group.pinned }); + } catch (e) { + error = e instanceof Error ? e.message : "Failed to update resource group"; + } + } + + function openSettingsDialog() { + if (!group) return; + editName = group.name; + editDescription = group.description ?? ""; + rgSettingsError = null; + rgSettingsDialog?.showModal(); + } + + async function handleSaveSettings() { + if (!auth.token || !group || !editName) return; + savingRgSettings = true; + rgSettingsError = null; + try { + group = await updateResourceGroup(auth.token, group.id, { + name: editName, + description: editDescription, + }); + rgSettingsDialog?.close(); + } catch (e) { + rgSettingsError = e instanceof Error ? e.message : "Failed to update resource group"; + } finally { + savingRgSettings = false; + } + } + + async function handleDeleteGroup() { + if (!auth.token || !group) return; + try { + await deleteResourceGroup(auth.token, group.id); + goto("/resource-groups"); + } catch (e) { + error = e instanceof Error ? e.message : "Failed to delete resource group"; + } } let totalCpu = $derived(workloads.reduce((s, w) => s + w.cpu_millicores, 0)); @@ -856,7 +1045,7 @@ >
-

Deploy Docker Compose Stack

+

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

- +
@@ -939,8 +1128,81 @@ {/if}
- +
+
+ + + { appearanceError = null; }} +> +
+
+

Edit Appearance

+ +
+ + {#if appearanceError} +

{appearanceError}

+ {/if} +
+ + +
+
+
+ + { rgSettingsError = null; }} +> +
+
+

Resource Group Settings

+ +
+
+
+ + +
+
+ + +
+
+ {#if rgSettingsError} +

{rgSettingsError}

+ {/if} +
+ +
@@ -1047,7 +1309,7 @@ {/if}
- {#each [["logs", "Logs"], ["shell", "Shell"], ["insights", "Performance"], ["network", "Network"]] as [tab, label]} + {#each [["logs", "Logs"], ["shell", "Shell"], ["insights", "Performance"], ["network", "Network"], ...(activeContainer.stack_id ? [] : [["settings", "Settings"]])] as [tab, label]}
- {:else} + {:else if containerDialogTab === "network"}
{#if activeContainer.ports && activeContainer.ports.length > 0}
@@ -1187,6 +1449,60 @@

No ports configured for this container.

{/if}
+ {:else if !activeContainer.stack_id} +
+
+ + +
+
+ + +

One KEY=VALUE per line.

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

+ Redeploying applies the new configuration by recreating the container, pulling the image again if changed. +

+ {#if settingsError} +

{settingsError}

+ {/if} +
+ +
+
{/if}
@@ -1213,9 +1529,15 @@

{error}

{:else if group}
-
- {initials(group.name)} -
+

{group.name}

@@ -1227,6 +1549,35 @@

{group.internal_cidr}

+ + +
+ openStackEditor(stack.stack_id)} + > - + {#if expanded} {#each stack.children as child (child.id)} 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/api-gateway/src/routes/resource_groups.rs b/control-plane/api-gateway/src/routes/resource_groups.rs index 1bc55955..002658af 100644 --- a/control-plane/api-gateway/src/routes/resource_groups.rs +++ b/control-plane/api-gateway/src/routes/resource_groups.rs @@ -25,11 +25,25 @@ use crate::{ AppState, }; +const DEFAULT_ICON: &str = "mdi:cube-outline"; +const DEFAULT_COLOR: &str = "#6366f1"; + #[derive(Debug, Deserialize)] pub struct CreateResourceGroupRequest { pub name: String, pub description: Option, pub internal_cidr: String, + pub icon: Option, + pub color: Option, +} + +#[derive(Debug, Deserialize)] +pub struct UpdateResourceGroupRequest { + pub name: Option, + pub description: Option, + pub icon: Option, + pub color: Option, + pub pinned: Option, } #[derive(Debug, Serialize)] @@ -40,6 +54,9 @@ pub struct ResourceGroupResponse { pub description: Option, pub internal_cidr: String, pub status: String, + pub icon: String, + pub color: String, + pub pinned: bool, pub created_at: chrono::NaiveDateTime, pub updated_at: Option, } @@ -53,6 +70,9 @@ impl From for ResourceGroupResponse { description: m.description, internal_cidr: m.internal_cidr, status: m.status, + icon: m.icon, + color: m.color, + pinned: m.pinned, created_at: m.created_at, updated_at: m.updated_at, } @@ -73,6 +93,8 @@ pub async fn list_resource_groups( let groups = ResourceGroups::find() .filter(resource_groups::Column::OrganizationId.eq(org_id)) + .order_by_desc(resource_groups::Column::Pinned) + .order_by_asc(resource_groups::Column::Name) .all(&state.db_conn) .await .map_err(|e| { @@ -136,6 +158,9 @@ pub async fn create_resource_group( description: Set(req.description), internal_cidr: Set(req.internal_cidr), status: Set("active".to_string()), + icon: Set(req.icon.unwrap_or_else(|| DEFAULT_ICON.to_string())), + color: Set(req.color.unwrap_or_else(|| DEFAULT_COLOR.to_string())), + pinned: Set(false), created_at: Set(now), updated_at: Set(None), }; @@ -154,6 +179,55 @@ pub async fn create_resource_group( )) } +const SUGGESTED_CIDR_BASE: u32 = 0x0A640000; +const SUGGESTED_CIDR_PREFIX: u8 = 24; +const SUGGESTED_CIDR_MAX_SUBNETS: u32 = 256; + +pub async fn suggest_cidr( + CanViewResourceGroups(_claims): CanViewResourceGroups, + State(state): State, +) -> Result)> { + let org_id = get_org_id(&state); + + let existing_groups = ResourceGroups::find() + .filter(resource_groups::Column::OrganizationId.eq(org_id)) + .all(&state.db_conn) + .await + .map_err(|e| { + tracing::error!(error = %e, "failed to list resource groups for cidr suggestion"); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({ "error": "database error" })), + ) + })?; + + let existing_cidrs: Vec = existing_groups + .iter() + .filter_map(|g| parse_cidr(&g.internal_cidr)) + .collect(); + + for subnet_index in 0..SUGGESTED_CIDR_MAX_SUBNETS { + let network = SUGGESTED_CIDR_BASE + (subnet_index << (32 - SUGGESTED_CIDR_PREFIX)); + let candidate = Cidr { + network, + prefix_len: SUGGESTED_CIDR_PREFIX, + }; + + if !existing_cidrs.iter().any(|c| c.overlaps(&candidate)) { + let [a, b, c, d] = network.to_be_bytes(); + return Ok(( + StatusCode::OK, + Json(json!({ "internal_cidr": format!("{}.{}.{}.{}/{}", a, b, c, d, SUGGESTED_CIDR_PREFIX) })), + )); + } + } + + Err(( + StatusCode::CONFLICT, + Json(json!({ "error": "no free cidr subnet available" })), + )) +} + pub async fn get_resource_group( CanViewResourceGroups(_claims): CanViewResourceGroups, State(state): State, @@ -185,6 +259,64 @@ pub async fn get_resource_group( )) } +pub async fn update_resource_group( + CanManageResourceGroups(_claims): CanManageResourceGroups, + State(state): State, + Path(id): Path, + Json(req): Json, +) -> Result)> { + let org_id = get_org_id(&state); + + let group = ResourceGroups::find_by_id(id) + .filter(resource_groups::Column::OrganizationId.eq(org_id)) + .one(&state.db_conn) + .await + .map_err(|e| { + tracing::error!(error = %e, "failed to find resource group"); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({ "error": "database error" })), + ) + })? + .ok_or_else(|| { + ( + StatusCode::NOT_FOUND, + Json(json!({ "error": "resource group not found" })), + ) + })?; + + let mut active: resource_groups::ActiveModel = group.into(); + if let Some(name) = req.name { + active.name = Set(name); + } + if let Some(description) = req.description { + active.description = Set(Some(description)); + } + if let Some(icon) = req.icon { + active.icon = Set(icon); + } + if let Some(color) = req.color { + active.color = Set(color); + } + if let Some(pinned) = req.pinned { + active.pinned = Set(pinned); + } + active.updated_at = Set(Some(Utc::now().naive_utc())); + + let updated = active.update(&state.db_conn).await.map_err(|e| { + tracing::error!(error = %e, "failed to update resource group"); + ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(json!({ "error": "database error" })), + ) + })?; + + Ok(( + StatusCode::OK, + Json(json!(ResourceGroupResponse::from(updated))), + )) +} + pub async fn delete_resource_group( CanManageResourceGroups(_claims): CanManageResourceGroups, State(state): State, @@ -787,9 +919,12 @@ pub fn resource_groups_routes() -> Router { "/resource-groups", get(list_resource_groups).post(create_resource_group), ) + .route("/resource-groups/suggest-cidr", get(suggest_cidr)) .route( "/resource-groups/{id}", - get(get_resource_group).delete(delete_resource_group), + get(get_resource_group) + .patch(update_resource_group) + .delete(delete_resource_group), ) .route( "/resource-groups/{id}/workloads", diff --git a/control-plane/api-gateway/src/routes/workloads.rs b/control-plane/api-gateway/src/routes/workloads.rs index 5dfd44a4..dbdc66e5 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, @@ -148,6 +167,93 @@ pub async fn create_stack( .await } +pub async fn get_stack( + CanViewWorkloads(_claims): CanViewWorkloads, + State(state): State, + Path(id): Path, + headers: HeaderMap, +) -> Result)> { + let header_map = header_vec(&headers); + proxy_to_scheduler( + &state, + reqwest::Method::GET, + &format!("/workload-stacks/{}", id), + None, + Some(header_map), + ) + .await +} + +pub async fn delete_stack( + CanManageWorkloads(_claims): CanManageWorkloads, + State(state): State, + Path(id): Path, + headers: HeaderMap, +) -> Result)> { + let header_map = header_vec(&headers); + proxy_to_scheduler( + &state, + reqwest::Method::DELETE, + &format!("/workload-stacks/{}", id), + None, + Some(header_map), + ) + .await +} + +pub async fn stop_stack( + CanManageWorkloads(_claims): CanManageWorkloads, + State(state): State, + Path(id): Path, + headers: HeaderMap, +) -> Result)> { + let header_map = header_vec(&headers); + proxy_to_scheduler( + &state, + reqwest::Method::POST, + &format!("/workload-stacks/{}/stop", id), + None, + Some(header_map), + ) + .await +} + +pub async fn restart_stack( + CanManageWorkloads(_claims): CanManageWorkloads, + State(state): State, + Path(id): Path, + headers: HeaderMap, +) -> Result)> { + let header_map = header_vec(&headers); + proxy_to_scheduler( + &state, + reqwest::Method::POST, + &format!("/workload-stacks/{}/restart", id), + None, + Some(header_map), + ) + .await +} + +pub async fn redeploy_stack( + 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!("/workload-stacks/{}", id), + body_json, + Some(header_map), + ) + .await +} + fn header_vec(headers: &HeaderMap) -> Vec<(String, String)> { headers .iter() @@ -160,7 +266,13 @@ 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)) + .route("/workload-stacks/{id}", get(get_stack)) + .route("/workload-stacks/{id}", delete(delete_stack)) + .route("/workload-stacks/{id}/stop", post(stop_stack)) + .route("/workload-stacks/{id}/restart", post(restart_stack)) + .route("/workload-stacks/{id}", patch(redeploy_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/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/db/workload_stacks.rs b/control-plane/scheduler/src/db/workload_stacks.rs index d5f56c84..99a42ddf 100644 --- a/control-plane/scheduler/src/db/workload_stacks.rs +++ b/control-plane/scheduler/src/db/workload_stacks.rs @@ -39,3 +39,35 @@ pub async fn update_status( active.update(db).await?; Ok(()) } + +pub async fn update_compose_source( + db: &DatabaseConnection, + stack_id: Uuid, + compose_source: &str, +) -> Result<(), sea_orm::DbErr> { + let stack = workload_stacks::Entity::find_by_id(stack_id) + .one(db) + .await? + .ok_or(sea_orm::DbErr::RecordNotFound(stack_id.to_string()))?; + + let mut active: workload_stacks::ActiveModel = stack.into(); + active.compose_source = Set(Some(compose_source.to_string())); + active.updated_at = Set(Some(Utc::now().naive_utc())); + + active.update(db).await?; + Ok(()) +} + +pub async fn get_by_id( + db: &DatabaseConnection, + stack_id: Uuid, +) -> Result, sea_orm::DbErr> { + workload_stacks::Entity::find_by_id(stack_id).one(db).await +} + +pub async fn delete(db: &DatabaseConnection, stack_id: Uuid) -> Result<(), sea_orm::DbErr> { + workload_stacks::Entity::delete_by_id(stack_id) + .exec(db) + .await?; + Ok(()) +} diff --git a/control-plane/scheduler/src/db/workloads.rs b/control-plane/scheduler/src/db/workloads.rs index fca28bd3..565749be 100644 --- a/control-plane/scheduler/src/db/workloads.rs +++ b/control-plane/scheduler/src/db/workloads.rs @@ -100,11 +100,64 @@ 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 image) = req.image { + active.image = Set(image.clone()); + } + 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()) } +pub async fn get_by_stack_id( + db: &DatabaseConnection, + stack_id: Uuid, +) -> Result, sea_orm::DbErr> { + workloads::Entity::find() + .filter(workloads::Column::StackId.eq(stack_id)) + .all(db) + .await +} + +pub async fn delete_by_stack_id(db: &DatabaseConnection, stack_id: Uuid) -> Result<(), sea_orm::DbErr> { + workloads::Entity::delete_many() + .filter(workloads::Column::StackId.eq(stack_id)) + .exec(db) + .await?; + Ok(()) +} + pub async fn get_pending(db: &DatabaseConnection) -> Result, sea_orm::DbErr> { workloads::Entity::find() .filter(workloads::Column::Status.eq(WorkloadStatus::Pending.as_str())) @@ -232,6 +285,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 +308,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/stacks.rs b/control-plane/scheduler/src/handlers/stacks.rs index 7fefacfe..031f409d 100644 --- a/control-plane/scheduler/src/handlers/stacks.rs +++ b/control-plane/scheduler/src/handlers/stacks.rs @@ -1,6 +1,15 @@ -use axum::{extract::State, http::StatusCode, response::IntoResponse, Json}; +use axum::{ + extract::{Path, State}, + http::StatusCode, + response::IntoResponse, + Json, +}; +use uuid::Uuid; -use crate::{models::compose::CreateStackRequest, server::AppState}; +use crate::{ + models::compose::{CreateStackRequest, RedeployStackRequest, StackResponse}, + server::AppState, +}; pub async fn create_stack( State(state): State, @@ -15,3 +24,87 @@ pub async fn create_stack( .into_response(), } } + +pub async fn get_stack( + State(state): State, + Path(id): Path, +) -> impl IntoResponse { + match crate::db::workload_stacks::get_by_id(&state.db, id).await { + Ok(Some(model)) => ( + StatusCode::OK, + Json(serde_json::json!(StackResponse { + id: model.id, + resource_group_id: model.resource_group_id, + name: model.name, + compose_source: model.compose_source, + status: model.status, + created_at: model.created_at.and_utc(), + updated_at: model.updated_at.map(|dt| dt.and_utc()), + })), + ) + .into_response(), + Ok(None) => StatusCode::NOT_FOUND.into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({ "error": e.to_string() })), + ) + .into_response(), + } +} + +pub async fn delete_stack( + State(state): State, + Path(id): Path, +) -> impl IntoResponse { + match state.scheduler.delete_stack(id).await { + Ok(()) => StatusCode::NO_CONTENT.into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({ "error": e })), + ) + .into_response(), + } +} + +pub async fn stop_stack( + State(state): State, + Path(id): Path, +) -> impl IntoResponse { + match state.scheduler.stop_stack(id).await { + Ok(()) => StatusCode::OK.into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({ "error": e })), + ) + .into_response(), + } +} + +pub async fn restart_stack( + State(state): State, + Path(id): Path, +) -> impl IntoResponse { + match state.scheduler.restart_stack(id).await { + Ok(()) => StatusCode::OK.into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({ "error": e })), + ) + .into_response(), + } +} + +pub async fn redeploy_stack( + State(state): State, + Path(id): Path, + Json(req): Json, +) -> impl IntoResponse { + match state.scheduler.redeploy_stack(id, &req.compose_yaml).await { + Ok(()) => StatusCode::OK.into_response(), + Err(e) => ( + StatusCode::INTERNAL_SERVER_ERROR, + Json(serde_json::json!({ "error": e })), + ) + .into_response(), + } +} diff --git a/control-plane/scheduler/src/handlers/workloads.rs b/control-plane/scheduler/src/handlers/workloads.rs index 7e02f3fe..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, }; @@ -56,6 +56,27 @@ pub async fn stop_workload( Path(id): Path, ) -> impl IntoResponse { 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)); + } + (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 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( @@ -80,9 +101,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/models/compose.rs b/control-plane/scheduler/src/models/compose.rs index abd47fe3..e3569c07 100644 --- a/control-plane/scheduler/src/models/compose.rs +++ b/control-plane/scheduler/src/models/compose.rs @@ -26,3 +26,19 @@ pub struct CreateStackResponse { pub stack_id: Uuid, pub workloads: Vec, } + +#[derive(Debug, Clone, serde::Deserialize)] +pub struct RedeployStackRequest { + pub compose_yaml: String, +} + +#[derive(Debug, Clone, serde::Serialize)] +pub struct StackResponse { + pub id: Uuid, + pub resource_group_id: Uuid, + pub name: String, + pub compose_source: Option, + pub status: String, + pub created_at: chrono::DateTime, + pub updated_at: Option>, +} diff --git a/control-plane/scheduler/src/models/workload.rs b/control-plane/scheduler/src/models/workload.rs index 3b3f3bf4..dcc90f46 100644 --- a/control-plane/scheduler/src/models/workload.rs +++ b/control-plane/scheduler/src/models/workload.rs @@ -137,6 +137,15 @@ impl DesiredState { } } +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct UpdateWorkloadRequest { + pub image: Option, + 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 +196,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..9dda7f07 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), @@ -46,6 +50,26 @@ pub fn create_router(state: AppState) -> Router { "/workload-stacks", axum::routing::post(stacks::create_stack), ) + .route( + "/workload-stacks/{id}", + axum::routing::get(stacks::get_stack), + ) + .route( + "/workload-stacks/{id}", + axum::routing::delete(stacks::delete_stack), + ) + .route( + "/workload-stacks/{id}/stop", + axum::routing::post(stacks::stop_stack), + ) + .route( + "/workload-stacks/{id}/restart", + axum::routing::post(stacks::restart_stack), + ) + .route( + "/workload-stacks/{id}", + axum::routing::patch(stacks::redeploy_stack), + ) .route( "/internal/workloads/status", axum::routing::post(internal::update_container_statuses), 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, }) } 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..4c252d94 100644 --- a/control-plane/scheduler/src/services/scheduler.rs +++ b/control-plane/scheduler/src/services/scheduler.rs @@ -7,7 +7,8 @@ use uuid::Uuid; use crate::models::compose::{ComposeServiceSpec, CreateStackRequest, CreateStackResponse}; use crate::models::workload::{ - AgentResources, CreateWorkloadRequest, CreateWorkloadResponse, WorkloadStatus, + AgentResources, CreateWorkloadRequest, CreateWorkloadResponse, UpdateWorkloadRequest, + WorkloadStatus, }; use crate::services::compose_parser::{self, ComposeParseError}; use crate::services::etcd::{delete_placement, put_placement, PlacementRecord}; @@ -68,7 +69,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 +94,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?; @@ -165,6 +167,107 @@ impl SchedulerService { }) } + pub async fn delete_stack(&self, stack_id: Uuid) -> Result<(), String> { + let workloads = crate::db::workloads::get_by_stack_id(&self.db, stack_id) + .await + .map_err(|e| format!("Failed to fetch stack workloads: {}", e))?; + + for workload in &workloads { + delete_placement(&self.etcd, workload.id).await?; + } + + crate::db::workloads::delete_by_stack_id(&self.db, stack_id) + .await + .map_err(|e| format!("Failed to delete stack workloads: {}", e))?; + + crate::db::workload_stacks::delete(&self.db, stack_id) + .await + .map_err(|e| format!("Failed to delete stack: {}", e))?; + + Ok(()) + } + + pub async fn stop_stack(&self, stack_id: Uuid) -> Result<(), String> { + let workloads = crate::db::workloads::get_by_stack_id(&self.db, stack_id) + .await + .map_err(|e| format!("Failed to fetch stack workloads: {}", e))?; + + for workload in &workloads { + let model = crate::db::workloads::set_desired_state( + &self.db, + workload.id, + crate::models::workload::DesiredState::Stopped, + ) + .await + .map_err(|e| format!("Failed to stop workload: {}", e))?; + if let Some(agent_id) = model.assigned_agent_id { + tokio::spawn(crate::services::gateway_notify::notify_assignment(agent_id)); + } + } + + Ok(()) + } + + pub async fn restart_stack(&self, stack_id: Uuid) -> Result<(), String> { + let workloads = crate::db::workloads::get_by_stack_id(&self.db, stack_id) + .await + .map_err(|e| format!("Failed to fetch stack workloads: {}", e))?; + + for workload in &workloads { + let model = crate::db::workloads::request_restart(&self.db, workload.id) + .await + .map_err(|e| format!("Failed to restart workload: {}", e))?; + if let Some(agent_id) = model.assigned_agent_id { + tokio::spawn(crate::services::gateway_notify::notify_assignment(agent_id)); + } + } + + Ok(()) + } + + pub async fn redeploy_stack( + &self, + stack_id: Uuid, + compose_yaml: &str, + ) -> Result<(), String> { + let services = compose_parser::parse_compose(compose_yaml) + .map_err(|e| format_compose_error(&e))?; + + let existing = crate::db::workloads::get_by_stack_id(&self.db, stack_id) + .await + .map_err(|e| format!("Failed to fetch stack workloads: {}", e))?; + + for service in &services { + let Some(workload) = existing + .iter() + .find(|w| w.service_name.as_deref() == Some(service.service_name.as_str())) + else { + continue; + }; + + let update = UpdateWorkloadRequest { + image: Some(service.image.clone()), + env_vars: Some(service.env_vars.clone().unwrap_or_default()), + ports: service.ports.clone(), + restart_policy: None, + max_restarts: None, + }; + + let model = crate::db::workloads::update_spec(&self.db, workload.id, &update) + .await + .map_err(|e| format!("Failed to update workload: {}", e))?; + if let Some(agent_id) = model.assigned_agent_id { + tokio::spawn(crate::services::gateway_notify::notify_assignment(agent_id)); + } + } + + crate::db::workload_stacks::update_compose_source(&self.db, stack_id, compose_yaml) + .await + .map_err(|e| format!("Failed to update stack source: {}", e))?; + + Ok(()) + } + async fn place_stack_workloads( &self, stack_id: Uuid, diff --git a/control-plane/shared/entity/src/entities/resource_groups.rs b/control-plane/shared/entity/src/entities/resource_groups.rs index 0a153d50..3994ea42 100644 --- a/control-plane/shared/entity/src/entities/resource_groups.rs +++ b/control-plane/shared/entity/src/entities/resource_groups.rs @@ -11,6 +11,9 @@ pub struct Model { pub description: Option, pub internal_cidr: String, pub status: String, + pub icon: String, + pub color: String, + pub pinned: bool, pub created_at: chrono::NaiveDateTime, pub updated_at: Option, } diff --git a/control-plane/shared/migration/src/lib.rs b/control-plane/shared/migration/src/lib.rs index a23d2b9c..07ceeee3 100644 --- a/control-plane/shared/migration/src/lib.rs +++ b/control-plane/shared/migration/src/lib.rs @@ -29,6 +29,7 @@ mod m20260711_000000_add_runtime_class; mod m20260711_010000_add_agent_cordoned; mod m20260712_000000_add_workload_lifecycle; mod m20260712_010000_add_resource_group_vpn_peers; +mod m20260723_000000_add_resource_group_appearance; pub struct Migrator; @@ -65,6 +66,7 @@ impl MigratorTrait for Migrator { Box::new(m20260711_010000_add_agent_cordoned::Migration), Box::new(m20260712_000000_add_workload_lifecycle::Migration), Box::new(m20260712_010000_add_resource_group_vpn_peers::Migration), + Box::new(m20260723_000000_add_resource_group_appearance::Migration), ] } } diff --git a/control-plane/shared/migration/src/m20260723_000000_add_resource_group_appearance.rs b/control-plane/shared/migration/src/m20260723_000000_add_resource_group_appearance.rs new file mode 100644 index 00000000..67267ce9 --- /dev/null +++ b/control-plane/shared/migration/src/m20260723_000000_add_resource_group_appearance.rs @@ -0,0 +1,48 @@ +use sea_orm_migration::prelude::*; + +#[derive(DeriveMigrationName)] +pub struct Migration; + +#[async_trait::async_trait] +impl MigrationTrait for Migration { + async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> { + manager + .alter_table( + Table::alter() + .table(Alias::new("resource_groups")) + .add_column_if_not_exists( + ColumnDef::new(Alias::new("icon")) + .string() + .not_null() + .default("mdi:cube-outline"), + ) + .add_column_if_not_exists( + ColumnDef::new(Alias::new("color")) + .string() + .not_null() + .default("#6366f1"), + ) + .add_column_if_not_exists( + ColumnDef::new(Alias::new("pinned")) + .boolean() + .not_null() + .default(false), + ) + .to_owned(), + ) + .await + } + + async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> { + manager + .alter_table( + Table::alter() + .table(Alias::new("resource_groups")) + .drop_column(Alias::new("icon")) + .drop_column(Alias::new("color")) + .drop_column(Alias::new("pinned")) + .to_owned(), + ) + .await + } +}
ID Name CIDR Description Status Created
Loading...Loading...
- No resource groups. Create one to get started. + + {groups.length === 0 ? "No resource groups. Create one to get started." : "No resource groups match search."}
{group.id.slice(0, 8)}{group.name} +
+
+ +
+ {group.name} +
+
{group.internal_cidr} {group.description ?? "-"} {group.status} {group.created_at.slice(0, 10)} - -
-
- +

{stack.stack_name}

{stack.children.length} services

- - +
Compose Stack @@ -1412,7 +1770,34 @@ {stack.children.filter((c) => c.status === 'running').length}/{stack.children.length} running +
+ + + +
+