diff --git a/k8s/backend-data-pvc.yaml b/k8s/backend-data-pvc.yaml deleted file mode 100644 index c588af96..00000000 --- a/k8s/backend-data-pvc.yaml +++ /dev/null @@ -1,14 +0,0 @@ -apiVersion: v1 -kind: PersistentVolumeClaim -metadata: - name: ushadow-data - namespace: ushadow - labels: - app: ushadow - component: data -spec: - accessModes: - - ReadWriteOnce - resources: - requests: - storage: 1Gi diff --git a/k8s/backend-deployment.yaml b/k8s/backend-deployment.yaml deleted file mode 100644 index 8a1b5d2d..00000000 --- a/k8s/backend-deployment.yaml +++ /dev/null @@ -1,49 +0,0 @@ -apiVersion: apps/v1 -kind: Deployment -metadata: - labels: - io.kompose.service: backend - name: backend -spec: - replicas: 1 - selector: - matchLabels: - io.kompose.service: backend - template: - metadata: - labels: - io.kompose.service: backend - spec: - containers: - - envFrom: - # K8s env vars sourced from .env.k8s (KC_HOSTNAME_URL, CORS_ORIGINS, etc.) - - configMapRef: - name: backend-k8s-env - # Secrets (MONGODB_URI, KEYCLOAK_ADMIN_PASSWORD) — apply once manually: - # kubectl apply -f k8s/secret.yaml - - secretRef: - name: ushadow-secret - optional: true - image: registry.temp.skaffold/ushadow-backend - name: ushadow-backend - ports: - - containerPort: 8000 - protocol: TCP - volumeMounts: - - name: config - mountPath: /config - - name: data - mountPath: /app/data - tolerations: - - key: "gpu" - operator: "Equal" - value: "amd" - effect: "NoSchedule" - volumes: - - name: config - persistentVolumeClaim: - claimName: ushadow-config - - name: data - persistentVolumeClaim: - claimName: ushadow-data - restartPolicy: Always diff --git a/k8s/backend-pvc.yaml b/k8s/backend-pvc.yaml deleted file mode 100644 index 0f7fcc44..00000000 --- a/k8s/backend-pvc.yaml +++ /dev/null @@ -1,14 +0,0 @@ -apiVersion: v1 -kind: PersistentVolumeClaim -metadata: - name: ushadow-config - namespace: ushadow - labels: - app: ushadow - component: config -spec: - accessModes: - - ReadWriteOnce - resources: - requests: - storage: 1Gi diff --git a/k8s/base/backend-deployment.yaml b/k8s/base/backend-deployment.yaml new file mode 100644 index 00000000..8ae06581 --- /dev/null +++ b/k8s/base/backend-deployment.yaml @@ -0,0 +1,42 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + labels: + io.kompose.service: backend + name: backend +spec: + replicas: 1 + selector: + matchLabels: + io.kompose.service: backend + template: + metadata: + labels: + io.kompose.service: backend + spec: + serviceAccountName: ushadow-backend + containers: + - envFrom: + - configMapRef: + name: ushadow-config + - secretRef: + name: ushadow-secret + env: + - name: BACKEND_PORT + value: "8000" + - name: CORS_ORIGINS + value: http://localhost:3000 + - name: HOST + value: 0.0.0.0 + - name: MONGODB_DATABASE + value: ushadow + - name: PORT + value: "8000" + - name: REDIS_URL + value: redis://redis:6379/0 + image: registry.temp.skaffold/ushadow-backend + name: ushadow-backend + ports: + - containerPort: 8000 + protocol: TCP + restartPolicy: Always diff --git a/k8s/backend-service.yaml b/k8s/base/backend-service.yaml similarity index 89% rename from k8s/backend-service.yaml rename to k8s/base/backend-service.yaml index 9085adfa..a976dfe5 100644 --- a/k8s/backend-service.yaml +++ b/k8s/base/backend-service.yaml @@ -3,7 +3,7 @@ kind: Service metadata: labels: io.kompose.service: backend - name: ushadow-backend + name: backend spec: ports: - name: "8000" diff --git a/k8s/base/kustomization.yaml b/k8s/base/kustomization.yaml new file mode 100644 index 00000000..3eb8f24b --- /dev/null +++ b/k8s/base/kustomization.yaml @@ -0,0 +1,9 @@ +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization + +resources: + - rbac.yaml + - backend-deployment.yaml + - backend-service.yaml + - webui-deployment.yaml + - webui-service.yaml diff --git a/k8s/base/rbac.yaml b/k8s/base/rbac.yaml new file mode 100644 index 00000000..196559bd --- /dev/null +++ b/k8s/base/rbac.yaml @@ -0,0 +1,25 @@ +apiVersion: v1 +kind: ServiceAccount +metadata: + name: ushadow-backend +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: Role +metadata: + name: ushadow-backend-reader +rules: + - apiGroups: [""] + resources: ["pods", "services"] + verbs: ["get", "list"] +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: RoleBinding +metadata: + name: ushadow-backend-reader +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: Role + name: ushadow-backend-reader +subjects: + - kind: ServiceAccount + name: ushadow-backend diff --git a/k8s/webui-deployment.yaml b/k8s/base/webui-deployment.yaml similarity index 82% rename from k8s/webui-deployment.yaml rename to k8s/base/webui-deployment.yaml index 8d64e613..1f16587f 100644 --- a/k8s/webui-deployment.yaml +++ b/k8s/base/webui-deployment.yaml @@ -18,6 +18,9 @@ spec: - env: - name: NODE_ENV value: production + - name: VITE_BACKEND_URL + value: http://localhost:8000 + - name: VITE_ENV_NAME image: registry.temp.skaffold/ushadow-frontend name: ushadow-webui ports: diff --git a/k8s/webui-service.yaml b/k8s/base/webui-service.yaml similarity index 73% rename from k8s/webui-service.yaml rename to k8s/base/webui-service.yaml index f7c69523..19e632db 100644 --- a/k8s/webui-service.yaml +++ b/k8s/base/webui-service.yaml @@ -3,11 +3,11 @@ kind: Service metadata: labels: io.kompose.service: webui - name: ushadow-frontend + name: webui spec: ports: - - name: "80" - port: 80 + - name: "3000" + port: 3000 targetPort: 80 selector: io.kompose.service: webui diff --git a/k8s/certificate.yaml b/k8s/certificate.yaml deleted file mode 100644 index ea847a88..00000000 --- a/k8s/certificate.yaml +++ /dev/null @@ -1,16 +0,0 @@ -apiVersion: cert-manager.io/v1 -kind: Certificate -metadata: - name: ushadow-tls - namespace: ushadow -spec: - dnsNames: - - ushadow.chakra - issuerRef: - group: cert-manager.io - kind: ClusterIssuer - name: local-ca-issuer - secretName: ushadow-tls - usages: - - digital signature - - key encipherment diff --git a/k8s/infra/kustomization.yaml b/k8s/infra/kustomization.yaml index 0e6d43a7..71dfed81 100644 --- a/k8s/infra/kustomization.yaml +++ b/k8s/infra/kustomization.yaml @@ -2,9 +2,9 @@ apiVersion: kustomize.config.k8s.io/v1beta1 kind: Kustomization resources: - - mongo-pvc.yaml - - mongo-deployment.yaml - - mongo-service.yaml + # mongo is managed separately — do not apply via kustomize to avoid selector conflicts + # - mongo-deployment.yaml + # - mongo-service.yaml - postgres-deployment.yaml - postgres-service.yaml - qdrant-deployment.yaml diff --git a/k8s/infra/mongo-deployment.yaml b/k8s/infra/mongo-deployment.yaml index bfda567a..85bfd611 100644 --- a/k8s/infra/mongo-deployment.yaml +++ b/k8s/infra/mongo-deployment.yaml @@ -17,15 +17,15 @@ spec: containers: - image: mongo:8.0 name: mongo - command: ["mongod", "--noauth"] ports: - containerPort: 27017 protocol: TCP - volumeMounts: - - name: mongo-data - mountPath: /data/db - volumes: - - name: mongo-data - persistentVolumeClaim: - claimName: mongo-data + env: + - name: MONGO_INITDB_ROOT_USERNAME + value: "root" + - name: MONGO_INITDB_ROOT_PASSWORD + valueFrom: + secretKeyRef: + name: ushadow-secret + key: MONGO_ROOT_PASSWORD restartPolicy: Always diff --git a/k8s/ingress.yaml b/k8s/ingress.yaml deleted file mode 100644 index 9cee6ab0..00000000 --- a/k8s/ingress.yaml +++ /dev/null @@ -1,57 +0,0 @@ -apiVersion: networking.k8s.io/v1 -kind: Ingress -metadata: - name: ushadow-ingress - namespace: ushadow -spec: - ingressClassName: nginx - tls: - - hosts: - - ushadow.chakra - secretName: ushadow-tls - rules: - - host: ushadow.chakra - http: - paths: - - path: /ws - pathType: Prefix - backend: - service: - name: ushadow-backend - port: - number: 8000 - - path: /api - pathType: Prefix - backend: - service: - name: ushadow-backend - port: - number: 8000 - - path: /openapi.json - pathType: Exact - backend: - service: - name: ushadow-backend - port: - number: 8000 - - path: /health - pathType: Exact - backend: - service: - name: ushadow-backend - port: - number: 8000 - - path: /docs - pathType: Exact - backend: - service: - name: ushadow-backend - port: - number: 8000 - - path: / - pathType: Prefix - backend: - service: - name: ushadow-frontend - port: - number: 80 diff --git a/k8s/kustomization.yaml b/k8s/kustomization.yaml index 67255274..ad559a9d 100644 --- a/k8s/kustomization.yaml +++ b/k8s/kustomization.yaml @@ -4,39 +4,22 @@ kind: Kustomization namespace: ushadow resources: - # - namespace.yaml # Namespace managed separately - DO NOT let Skaffold delete it! - # secret.yaml is gitignored (contains plaintext creds) - apply once manually: - # kubectl apply -f k8s/secret.yaml (copy from secret.yaml.example and fill in values) - # - # PVCs are NOT listed here — managed by seed-pvcs.sh so Skaffold never deletes them. - # Skaffold runs `skaffold delete` on `dev` teardown which would wipe PVC-stored overrides. - - backend-deployment.yaml - - backend-service.yaml - - webui-deployment.yaml - - webui-service.yaml - - ingress.yaml - - certificate.yaml - # - infra/ # Infrastructure managed separately - -# K8s-specific env vars loaded from .env.k8s (gitignored, copy from .env and update for K8s) -# Provides: KC_HOSTNAME_URL, KC_URL, CORS_ORIGINS, USHADOW_PUBLIC_URL, etc. -configMapGenerator: - - name: backend-k8s-env - envs: - - .env.k8s -generatorOptions: - disableNameSuffixHash: true + - namespace.yaml + - configmap.yaml + # secret.yaml intentionally NOT included. + # ushadow-secret is the single responsibility of `just k8s-apply-secret`, + # which sources it from config/SECRETS/secrets_k8s.yaml (gitignored). + # skaffold's `deploy.hooks.before` runs that command pre-deploy. + - infra/ + - base/ + # - tweaks/ingress-example.yaml # Uncomment when ready # Add common labels to all resources -labels: - - pairs: - app: ushadow - env: purple +commonLabels: + app: ushadow + env: ushadow # Add annotations commonAnnotations: managed-by: kustomize deployed-from: docker-compose - -# Note: Image replacement handled by Skaffold's --default-repo flag -# No Kustomize image transformers needed when using Skaffold diff --git a/k8s/scripts/seed-pvcs.sh b/k8s/scripts/seed-pvcs.sh deleted file mode 100644 index d7adda6a..00000000 --- a/k8s/scripts/seed-pvcs.sh +++ /dev/null @@ -1,132 +0,0 @@ -#!/usr/bin/env bash -# seed-pvcs.sh — copy local directories into Kubernetes PVCs. -# -# Used by Skaffold as a post-deploy hook and can also be run manually: -# ./k8s/scripts/seed-pvcs.sh [namespace] [kubeconfig] -# -# Seeding rules: -# - Only seeds if the PVC is missing a specific sentinel file -# - On check failure (PVC busy / pod error), SKIPS seeding (safe default) -# - Never overwrites user-data files: config.overrides.yaml, SECRETS/, instance-overrides.yaml -# -# Usage: -# ./k8s/scripts/seed-pvcs.sh ushadow # uses current kubeconfig -# ./k8s/scripts/seed-pvcs.sh ushadow ~/.kube/config - -set -euo pipefail - -NAMESPACE="${1:-ushadow}" -KUBECONFIG_PATH="${2:-}" -REPO_ROOT="$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)" - -KUBECTL="kubectl" -if [[ -n "$KUBECONFIG_PATH" ]]; then - KUBECTL="kubectl --kubeconfig=$KUBECONFIG_PATH" -fi - -KUBECTL="$KUBECTL -n $NAMESPACE" - -# ─── helpers ──────────────────────────────────────────────────────────────── - -log() { echo "[seed-pvcs] $*"; } -warn() { echo "[seed-pvcs] WARNING: $*" >&2; } - -# Seed a single PVC from a local directory. -# -# Checks for a sentinel file (config.defaults.yaml or compose/) to decide -# whether the PVC is already seeded. Defaults to SKIP if the check fails -# (e.g. PVC busy with ReadWriteOnce storage — better to skip than overwrite). -# -# Never overwrites user-data files even when seeding fresh: -# config.overrides.yaml, SECRETS/, instance-overrides.yaml -# -seed_pvc() { - local pvc_name="$1" - local source_dir="$2" - local sentinel_file="${3:-}" # relative path inside PVC to check for (e.g. "config.defaults.yaml") - - if [[ ! -d "$source_dir" ]]; then - warn "Source dir $source_dir does not exist, skipping $pvc_name" - return 0 - fi - - # Check for the sentinel file inside the PVC - local check_pod="check-${pvc_name:0:16}-$(head -c4 /dev/urandom | xxd -p)" - local sentinel_check_cmd="ls /data/${sentinel_file} > /dev/null 2>&1 && echo FOUND || echo MISSING" - - log "Checking PVC $pvc_name for sentinel: ${sentinel_file:-}..." - - local check_result="SKIP" # Safe default: skip if we can't determine state - if $KUBECTL run "$check_pod" \ - --image=busybox:1.36 \ - --restart=Never \ - --overrides="{\"spec\":{\"volumes\":[{\"name\":\"pvc\",\"persistentVolumeClaim\":{\"claimName\":\"$pvc_name\"}}],\"containers\":[{\"name\":\"check\",\"image\":\"busybox:1.36\",\"command\":[\"sh\",\"-c\",\"${sentinel_check_cmd}\"],\"volumeMounts\":[{\"name\":\"pvc\",\"mountPath\":\"/data\"}]}]}}" \ - --wait --timeout=60s 2>/dev/null; then - check_result=$($KUBECTL logs "$check_pod" 2>/dev/null | tr -d '[:space:]') || check_result="SKIP" - else - warn "Check pod failed to run for $pvc_name — skipping seed (safe default)" - fi - $KUBECTL delete pod "$check_pod" --ignore-not-found --wait=false 2>/dev/null || true - - if [[ "$check_result" == "FOUND" ]]; then - log "PVC $pvc_name already seeded (sentinel found), skipping" - return 0 - elif [[ "$check_result" != "MISSING" ]]; then - # SKIP or unexpected output — don't risk overwriting user data - warn "PVC $pvc_name check inconclusive (result='$check_result'), skipping seed" - return 0 - fi - - log "Seeding PVC $pvc_name from $source_dir ..." - - # Create a pod that mounts the PVC and sleeps so we can cp into it - local pod_name="seed-${pvc_name:0:18}-$(head -c4 /dev/urandom | xxd -p)" - $KUBECTL run "$pod_name" \ - --image=busybox:1.36 \ - --restart=Never \ - --overrides="{\"spec\":{\"volumes\":[{\"name\":\"pvc\",\"persistentVolumeClaim\":{\"claimName\":\"$pvc_name\"}}],\"containers\":[{\"name\":\"seeder\",\"image\":\"busybox:1.36\",\"command\":[\"sh\",\"-c\",\"sleep 300\"],\"volumeMounts\":[{\"name\":\"pvc\",\"mountPath\":\"/seed-data\"}]}]}}" - - log "Waiting for seeder pod $pod_name to be Running..." - $KUBECTL wait pod "$pod_name" --for=condition=Ready --timeout=120s - - log "Copying files from $source_dir into PVC $pvc_name (preserving user-data files)..." - - # Copy all files EXCEPT user-data files that should never be overwritten. - # We use a temporary staging approach: copy everything then delete protected files. - # Protected files: config.overrides.yaml, instance-overrides.yaml, SECRETS/ - $KUBECTL cp "$source_dir/." "$pod_name:/seed-data/" - - # Remove any user-data files we just copied — they should only exist if the user created them. - # This prevents factory-default versions from polluting the PVC. - $KUBECTL exec "$pod_name" -- sh -c " - rm -f /seed-data/config.overrides.yaml - rm -f /seed-data/instance-overrides.yaml - rm -f /seed-data/secrets.yml - rm -rf /seed-data/SECRETS - " 2>/dev/null || true - - # Clean up - $KUBECTL delete pod "$pod_name" --grace-period=0 --force 2>/dev/null || true - - log "PVC $pvc_name seeded successfully" -} - -# ─── main ─────────────────────────────────────────────────────────────────── - -log "Seeding PVCs in namespace: $NAMESPACE" -log "Repository root: $REPO_ROOT" - -# Apply PVC manifests first (idempotent - kubectl apply is a no-op if they exist). -# PVCs are NOT in kustomization.yaml to prevent Skaffold from deleting them on teardown -# (skaffold dev runs `skaffold delete` on Ctrl+C, which would wipe stored overrides/secrets). -log "Ensuring PVCs exist..." -$KUBECTL apply -f "$REPO_ROOT/k8s/backend-pvc.yaml" -$KUBECTL apply -f "$REPO_ROOT/k8s/backend-data-pvc.yaml" - -# Seed the shared config PVC — sentinel: config.defaults.yaml (always present after seeding) -seed_pvc "ushadow-config" "$REPO_ROOT/config" "config.defaults.yaml" - -# Seed the shared compose PVC — sentinel: the directory itself having content -seed_pvc "ushadow-compose" "$REPO_ROOT/compose" "docker-compose.yml" - -log "PVC seeding complete" diff --git a/k8s/scripts/setup-tailscale-operator.sh b/k8s/scripts/setup-tailscale-operator.sh deleted file mode 100644 index 4536cb15..00000000 --- a/k8s/scripts/setup-tailscale-operator.sh +++ /dev/null @@ -1,207 +0,0 @@ -#!/usr/bin/env bash -# setup-tailscale-operator.sh — Install the Tailscale Kubernetes operator and -# expose the ushadow ingress via a Tailscale-issued *.ts.net hostname. -# -# The operator installs once per cluster. After installation, annotating any -# Service with tailscale.com/expose=true gives it a MagicDNS hostname and a -# Tailscale-backed TLS cert (trusted by all devices, including mobile). -# -# Prerequisites: -# - kubectl connected to the target cluster -# - helm 3.x installed -# - Tailscale OAuth client (NOT an auth key) — create one at: -# https://login.tailscale.com/admin/settings/oauth -# Scopes required: Core (devices read/write) + Keys (auth keys write) -# Tag required: ensure tag:k8s-operator exists in your tailnet ACL policy -# -# Usage: -# ./k8s/scripts/setup-tailscale-operator.sh \ -# --client-id tskey-client-xxx \ -# --client-secret tskey-secret-xxx \ -# --hostname ushadow-chakra -# -# Idempotent: safe to re-run. Upgrades the operator if already installed. -# The annotated service is left as-is if already annotated. - -set -euo pipefail - -# ─── defaults ──────────────────────────────────────────────────────────────── - -OPERATOR_NAMESPACE="tailscale-system" -INGRESS_NAMESPACE="ingress-nginx" -INGRESS_SERVICE="ingress-nginx-controller" -TS_HOSTNAME="ushadow-chakra" -CLIENT_ID="" -CLIENT_SECRET="" -KUBECONFIG_PATH="" -WAIT_TIMEOUT=120 # seconds to wait for Tailscale IP - -# ─── helpers ───────────────────────────────────────────────────────────────── - -log() { echo "[ts-operator] $*"; } -warn() { echo "[ts-operator] WARNING: $*" >&2; } -err() { echo "[ts-operator] ERROR: $*" >&2; exit 1; } - -usage() { - grep '^#' "$0" | sed 's/^# \{0,1\}//' | head -30 - exit 0 -} - -require_cmd() { - command -v "$1" >/dev/null 2>&1 || err "'$1' is required but not found in PATH" -} - -kubectl_cmd() { - if [[ -n "$KUBECONFIG_PATH" ]]; then - kubectl --kubeconfig="$KUBECONFIG_PATH" "$@" - else - kubectl "$@" - fi -} - -# ─── arg parsing ───────────────────────────────────────────────────────────── - -while [[ $# -gt 0 ]]; do - case "$1" in - --client-id) CLIENT_ID="$2"; shift 2 ;; - --client-secret) CLIENT_SECRET="$2"; shift 2 ;; - --hostname) TS_HOSTNAME="$2"; shift 2 ;; - --operator-ns) OPERATOR_NAMESPACE="$2"; shift 2 ;; - --ingress-ns) INGRESS_NAMESPACE="$2"; shift 2 ;; - --ingress-svc) INGRESS_SERVICE="$2"; shift 2 ;; - --kubeconfig) KUBECONFIG_PATH="$2"; shift 2 ;; - --wait-timeout) WAIT_TIMEOUT="$2"; shift 2 ;; - -h|--help) usage ;; - *) err "Unknown argument: $1" ;; - esac -done - -[[ -n "$CLIENT_ID" ]] || err "--client-id is required (create at https://login.tailscale.com/admin/settings/oauth)" -[[ -n "$CLIENT_SECRET" ]] || err "--client-secret is required" - -# ─── preflight ─────────────────────────────────────────────────────────────── - -require_cmd kubectl -require_cmd helm - -log "Checking cluster connectivity..." -kubectl_cmd cluster-info --request-timeout=5s >/dev/null \ - || err "Cannot reach cluster. Check your kubeconfig / VPN." - -log "Target cluster: $(kubectl_cmd config current-context 2>/dev/null || echo '(unknown)')" -log "Ingress service: $INGRESS_NAMESPACE/$INGRESS_SERVICE" -log "Tailscale hostname: $TS_HOSTNAME" - -# ─── step 1: Helm repo ─────────────────────────────────────────────────────── - -log "Adding Tailscale Helm repo..." -helm repo add tailscale https://pkgs.tailscale.com/helmcharts 2>/dev/null || true -helm repo update tailscale - -# ─── step 2: Install / upgrade operator ────────────────────────────────────── - -log "Installing/upgrading Tailscale operator in namespace '$OPERATOR_NAMESPACE'..." -kubectl_cmd create namespace "$OPERATOR_NAMESPACE" --dry-run=client -o yaml \ - | kubectl_cmd apply -f - - -helm upgrade --install tailscale-operator tailscale/tailscale-operator \ - --namespace "$OPERATOR_NAMESPACE" \ - --set-string oauth.clientId="$CLIENT_ID" \ - --set-string oauth.clientSecret="$CLIENT_SECRET" \ - --set operatorConfig.defaultTags=tag:k8s-operator \ - --wait --timeout 120s - -log "Operator installed. Waiting for it to become ready..." -kubectl_cmd rollout status deployment/operator \ - -n "$OPERATOR_NAMESPACE" --timeout=60s - -# ─── step 3: Annotate ingress service ──────────────────────────────────────── - -log "Checking ingress service $INGRESS_NAMESPACE/$INGRESS_SERVICE..." -kubectl_cmd get service "$INGRESS_SERVICE" -n "$INGRESS_NAMESPACE" >/dev/null \ - || err "Service $INGRESS_NAMESPACE/$INGRESS_SERVICE not found. Check --ingress-ns / --ingress-svc." - -# Check if already annotated -EXISTING_HOSTNAME=$(kubectl_cmd get service "$INGRESS_SERVICE" \ - -n "$INGRESS_NAMESPACE" \ - -o jsonpath='{.metadata.annotations.tailscale\.com/hostname}' 2>/dev/null || true) - -if [[ "$EXISTING_HOSTNAME" == "$TS_HOSTNAME" ]]; then - log "Service already annotated with hostname '$TS_HOSTNAME', skipping annotation." -else - log "Annotating $INGRESS_SERVICE with Tailscale hostname '$TS_HOSTNAME'..." - kubectl_cmd annotate service "$INGRESS_SERVICE" \ - -n "$INGRESS_NAMESPACE" \ - "tailscale.com/expose=true" \ - "tailscale.com/hostname=$TS_HOSTNAME" \ - --overwrite -fi - -# ─── step 4: Wait for Tailscale IP ─────────────────────────────────────────── - -log "Waiting up to ${WAIT_TIMEOUT}s for Tailscale to provision the proxy..." - -TS_IP="" -ELAPSED=0 -INTERVAL=5 -while [[ $ELAPSED -lt $WAIT_TIMEOUT ]]; do - TS_IP=$(kubectl_cmd get service "$INGRESS_SERVICE" \ - -n "$INGRESS_NAMESPACE" \ - -o jsonpath='{.status.loadBalancer.ingress[?(@.hostname)].hostname}' 2>/dev/null || true) - - # Also check for IP (some setups return IP not hostname at first) - if [[ -z "$TS_IP" ]]; then - TS_IP=$(kubectl_cmd get service "$INGRESS_SERVICE" \ - -n "$INGRESS_NAMESPACE" \ - -o jsonpath='{.status.loadBalancer.ingress[*].ip}' 2>/dev/null || true) - fi - - if [[ -n "$TS_IP" ]]; then - break - fi - - printf "." - sleep $INTERVAL - ELAPSED=$((ELAPSED + INTERVAL)) -done -echo "" - -# ─── step 5: Summary ───────────────────────────────────────────────────────── - -FULL_HOSTNAME="${TS_HOSTNAME}.$(kubectl_cmd get secret \ - -n "$OPERATOR_NAMESPACE" \ - -l "tailscale.com/managed=true" \ - -o jsonpath='{.items[0].data.config}' 2>/dev/null \ - | base64 -d 2>/dev/null \ - | grep -o '"ServerName":"[^"]*"' \ - | head -1 \ - | sed 's/"ServerName":"//;s/"//' \ - | sed "s/^${TS_HOSTNAME}\.//" \ - || echo "spangled-kettle.ts.net")" - -log "" -log "━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━" -log " Tailscale operator setup complete" -log "━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━" -log "" -if [[ -n "$TS_IP" ]]; then - log " Tailscale proxy: $TS_IP" -fi -log " MagicDNS hostname: https://${TS_HOSTNAME}..ts.net" -log "" -log " Next steps:" -log " 1. Check your Tailscale admin panel for the actual hostname:" -log " https://login.tailscale.com/admin/machines" -log "" -log " 2. Update k8s/backend-deployment.yaml:" -log " USHADOW_PUBLIC_URL: https://${TS_HOSTNAME}..ts.net" -log " CORS_ORIGINS: https://${TS_HOSTNAME}..ts.net" -log "" -log " 3. Update k8s/ingress.yaml to add the new hostname:" -log " - host: ${TS_HOSTNAME}..ts.net" -log "" -log " 4. Delete the old local-ca certificate (no longer needed):" -log " kubectl delete certificate ushadow-tls -n ushadow" -log " kubectl delete clusterissuer local-ca-issuer" -log "" -log "━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━" diff --git a/k8s/secret.yaml.example b/k8s/secret.yaml.example deleted file mode 100644 index 6489de89..00000000 --- a/k8s/secret.yaml.example +++ /dev/null @@ -1,17 +0,0 @@ -apiVersion: v1 -kind: Secret -metadata: - name: ushadow-secret - namespace: ushadow -type: Opaque -stringData: - # MongoDB connection URI with credentials - # Format: mongodb://user:password@host:port - MONGODB_URI: "mongodb://root:CHANGE_ME@mongo:27017" - - # Keycloak admin password - KEYCLOAK_ADMIN_PASSWORD: "CHANGE_ME" - - # TODO: Add other secrets as needed - # REDIS_PASSWORD: "your-redis-password" - # API_KEYS: "your-api-keys" diff --git a/k8s/tweaks/README.md b/k8s/tweaks/README.md index dce1a0d3..689d3a59 100644 --- a/k8s/tweaks/README.md +++ b/k8s/tweaks/README.md @@ -60,4 +60,4 @@ Add NetworkPolicy for security (optional but recommended) - `base/` - Application services (backend, frontend) - `namespace.yaml` - Namespace definition - `configmap.yaml` - Configuration data -- `secret.yaml` - Secrets (template only) +- `ushadow-secret` - Applied by `just k8s-apply-secret` from `config/SECRETS/secrets_k8s.yaml`. Not managed by kustomize (skaffold runs it as a pre-deploy hook). diff --git a/scripts/just/k8s.just b/scripts/just/k8s.just new file mode 100644 index 00000000..acb084e3 --- /dev/null +++ b/scripts/just/k8s.just @@ -0,0 +1,44 @@ +# Kubernetes recipes + +_secrets_k8s := "./config/SECRETS/secrets_k8s.yaml" +_k8s_namespace := "ushadow" + +# Apply ushadow-secret to k8s from config/SECRETS/secrets_k8s.yaml (k8s_secret section) +k8s-apply-secret: + #!/usr/bin/env python3 + import subprocess, sys + try: + import yaml + except ImportError: + subprocess.run([sys.executable, "-m", "pip", "install", "-q", "pyyaml"], check=True) + import yaml + + secrets_file = "{{_secrets_k8s}}" + namespace = "{{_k8s_namespace}}" + + with open(secrets_file) as f: + data = yaml.safe_load(f) + + k8s_secret = data.get("k8s_secret") + if not k8s_secret: + print("ERROR: no k8s_secret section found in", secrets_file, file=sys.stderr) + sys.exit(1) + + manifest = { + "apiVersion": "v1", + "kind": "Secret", + "metadata": {"name": "ushadow-secret", "namespace": namespace}, + "type": "Opaque", + "stringData": {k: str(v) for k, v in k8s_secret.items()}, + } + + proc = subprocess.run( + ["kubectl", "apply", "-f", "-"], + input=yaml.dump(manifest), + text=True, + capture_output=True, + ) + print(proc.stdout.strip() or proc.stderr.strip()) + if proc.returncode != 0: + sys.exit(proc.returncode) + print("✓ ushadow-secret applied from", secrets_file) diff --git a/skaffold.yaml b/skaffold.yaml index 365141b0..5556b246 100644 --- a/skaffold.yaml +++ b/skaffold.yaml @@ -23,13 +23,16 @@ build: useDockerCLI: true deploy: - kubectl: - hooks: - after: - - host: - command: ["bash", "k8s/scripts/seed-pvcs.sh", "ushadow"] - dir: . - os: [linux, darwin] + kubectl: {} + # ushadow-secret is intentionally NOT in the kustomize bundle. + # It's sourced from config/SECRETS/secrets_k8s.yaml (gitignored) and applied + # by `just k8s-apply-secret`. Running it here makes `skaffold dev` idempotent + # and guarantees the secret exists before any Deployment references it. + hooks: + before: + - host: + command: ["just", "k8s-apply-secret"] + os: [darwin, linux] manifests: kustomize: diff --git a/ushadow/backend/pyproject.toml b/ushadow/backend/pyproject.toml index 3015a75b..c424179e 100644 --- a/ushadow/backend/pyproject.toml +++ b/ushadow/backend/pyproject.toml @@ -132,6 +132,8 @@ max-returns = 6 # Simplify return logic "scripts/*" = ["T201"] # Allow assert in tests "tests/*" = ["S101", "PLR2004"] +# FastAPI routers: Depends() in default args (B008) and unused current_user auth deps (ARG001) are idiomatic +"src/routers/*.py" = ["B008", "ARG001"] [tool.ruff.lint.isort] # Organize imports for readability diff --git a/ushadow/backend/src/routers/connect.py b/ushadow/backend/src/routers/connect.py index 8024c3a0..9f04dec9 100644 --- a/ushadow/backend/src/routers/connect.py +++ b/ushadow/backend/src/routers/connect.py @@ -13,7 +13,6 @@ import logging import os -from typing import Optional from fastapi import APIRouter, HTTPException, Query from pydantic import BaseModel @@ -35,7 +34,7 @@ class ConnectionInfo(BaseModel): @router.get("", response_model=ConnectionInfo) async def get_connection_info( - target_id: Optional[str] = Query( + target_id: str | None = Query( None, description="Deploy target ID (e.g. 'my-cluster.k8s.prod' or 'orange-public.unode.orange'). " "Auto-detected when omitted.", @@ -48,8 +47,8 @@ async def get_connection_info( After connecting, clients use the returned api_url and keycloak_mobile_url to authenticate and access the API. """ - from src.utils.environment import is_kubernetes from src.config.casdoor_settings import get_casdoor_config + from src.utils.environment import is_kubernetes casdoor_config = get_casdoor_config() realm = casdoor_config.get("organization", "ushadow") @@ -67,8 +66,8 @@ async def get_connection_info( async def _connection_info_docker(realm: str, auth_url: str) -> ConnectionInfo: """Connection info for a private Docker/unode deployment.""" - from src.services.unode_manager import get_unode_manager from src.models.unode import UNodeRole + from src.services.unode_manager import get_unode_manager from src.utils.tailscale_serve import get_tailscale_status unode_manager = await get_unode_manager() @@ -101,23 +100,25 @@ async def _connection_info_docker(realm: str, auth_url: str) -> ConnectionInfo: def _connection_info_k8s(realm: str) -> ConnectionInfo: """Connection info for a Kubernetes deployment.""" + from src.config.casdoor_settings import get_casdoor_config + api_url = os.getenv("USHADOW_PUBLIC_URL", "").rstrip("/") - kc_mobile_url = os.getenv("KC_HOSTNAME_URL", "").rstrip("/") + casdoor_url = get_casdoor_config().get("public_url", "").rstrip("/") if not api_url: raise HTTPException( status_code=503, detail="USHADOW_PUBLIC_URL is not set. Configure the K8s deployment.", ) - if not kc_mobile_url: + if not casdoor_url: raise HTTPException( status_code=503, - detail="KC_HOSTNAME_URL is not set. Configure the K8s deployment.", + detail="CASDOOR_EXTERNAL_URL is not set. Configure the K8s deployment.", ) return ConnectionInfo( api_url=api_url, - keycloak_mobile_url=kc_mobile_url, + keycloak_mobile_url=casdoor_url, realm=realm, mobile_client_id="ushadow-mobile", platform="kubernetes", @@ -132,7 +133,7 @@ async def _connection_info_for_target(target_id: str, realm: str) -> ConnectionI try: target = await DeployTarget.from_id(target_id) except ValueError as e: - raise HTTPException(status_code=404, detail=str(e)) + raise HTTPException(status_code=404, detail=str(e)) from e is_public = (target.raw_metadata.get("labels") or {}).get("zone") == "public" diff --git a/ushadow/backend/src/routers/deployments.py b/ushadow/backend/src/routers/deployments.py index 5bb70a8b..dbe6b1de 100644 --- a/ushadow/backend/src/routers/deployments.py +++ b/ushadow/backend/src/routers/deployments.py @@ -1,34 +1,33 @@ """API routes for service deployments.""" import logging -from typing import List, Optional, Dict, Any +from typing import Any -from fastapi import APIRouter, HTTPException, Depends, Query +from fastapi import APIRouter, Depends, HTTPException, Query from pydantic import BaseModel +from src.models.deploy_target import DeployTarget from src.models.deployment import ( - ServiceDefinition, - ServiceDefinitionCreate, - ServiceDefinitionUpdate, + AdoptRequest, Deployment, DeployRequest, DiscoveredWorkload, - AdoptRequest, + ServiceDefinition, + ServiceDefinitionCreate, + ServiceDefinitionUpdate, ) +from src.services.auth import get_current_user +from src.services.deployment_manager import get_deployment_manager +from src.services.kubernetes import get_kubernetes_manager +from src.services.unode_manager import get_unode_manager + # Slim view of a deployment — excludes deployed_config (which contains the full env map). -# Used by the instance page to avoid transmitting large environment variable payloads. _SLIM_FIELDS = { 'id', 'config_id', 'service_id', 'unode_hostname', 'status', 'container_id', 'container_name', 'created_at', 'deployed_at', 'exposed_port', 'access_url', 'public_url', 'metadata', 'backend_type', 'healthy', 'health_message', 'error', } -from src.services.deployment_manager import get_deployment_manager -from src.services.auth import get_current_user -from src.services.unode_manager import get_unode_manager -from src.services.kubernetes import get_kubernetes_manager -from src.models.deploy_target import DeployTarget -from src.models.unode import UNodeType logger = logging.getLogger(__name__) @@ -39,7 +38,7 @@ # Deployment Targets Endpoint # ============================================================================= -@router.get("/targets", response_model=List[Dict[str, Any]]) +@router.get("/targets", response_model=list[dict[str, Any]]) async def list_deployment_targets( current_user: dict = Depends(get_current_user) ): @@ -118,9 +117,9 @@ async def list_deployment_targets( if infra: logger.info(f" Infrastructure services: {list(infra.keys())}") else: - logger.info(f" ⚠️ No infrastructure found in any namespace") + logger.info(" ⚠️ No infrastructure found in any namespace") else: - logger.info(f" ⚠️ No infra_scans available for cluster") + logger.info(" ⚠️ No infra_scans available for cluster") # Try to infer provider from labels or use default provider = cluster.labels.get("provider", "kubernetes") @@ -159,14 +158,14 @@ async def get_target_infrastructure_env_vars( - source='infrastructure' : from cluster infra scan (locked=True) - source='override' : user-saved manual override (locked=False) """ - from src.models.deploy_target import DeployTarget from src.config import get_settings_store + from src.models.deploy_target import DeployTarget from src.services.infrastructure_env_service import resolve_infrastructure_env_vars try: target = await DeployTarget.from_id(target_id) except ValueError as e: - raise HTTPException(status_code=404, detail=str(e)) + raise HTTPException(status_code=404, detail=str(e)) from e store = get_settings_store() env_vars = await resolve_infrastructure_env_vars(target, store) @@ -176,17 +175,17 @@ async def get_target_infrastructure_env_vars( @router.put("/targets/{target_id}/infrastructure-env-vars") async def save_target_infrastructure_overrides( target_id: str, - overrides: Dict[str, str], + overrides: dict[str, str], current_user: dict = Depends(get_current_user), ): """Save manual infrastructure overrides for a deploy target.""" - from src.models.deploy_target import DeployTarget from src.config import get_settings_store + from src.models.deploy_target import DeployTarget try: target = await DeployTarget.from_id(target_id) except ValueError as e: - raise HTTPException(status_code=404, detail=str(e)) + raise HTTPException(status_code=404, detail=str(e)) from e if target.type != "k8s": raise HTTPException(status_code=400, detail="Infrastructure overrides only supported for K8s targets") @@ -209,17 +208,16 @@ async def create_service_definition( """Create a new service definition.""" manager = get_deployment_manager() try: - service = await manager.create_service( + return await manager.create_service( data, created_by=current_user.get("email") ) - return service except Exception as e: logger.error(f"Failed to create service: {e}") - raise HTTPException(status_code=400, detail=str(e)) + raise HTTPException(status_code=400, detail=str(e)) from e -@router.get("/services", response_model=List[ServiceDefinition]) +@router.get("/services", response_model=list[ServiceDefinition]) async def list_service_definitions( current_user: dict = Depends(get_current_user) ): @@ -268,7 +266,7 @@ async def delete_service_definition( raise HTTPException(status_code=404, detail="Service not found") return {"success": True, "message": f"Service {service_id} deleted"} except ValueError as e: - raise HTTPException(status_code=400, detail=str(e)) + raise HTTPException(status_code=400, detail=str(e)) from e # ============================================================================= @@ -283,27 +281,24 @@ async def deploy_service( """Deploy a service to a u-node.""" manager = get_deployment_manager() try: - # If config_id not provided, use service_id (template as config) config_id = data.config_id or data.service_id - - deployment = await manager.deploy_service( + return await manager.deploy_service( data.service_id, data.unode_hostname, config_id=config_id, force_rebuild=data.force_rebuild ) - return deployment except ValueError as e: - raise HTTPException(status_code=400, detail=str(e)) + raise HTTPException(status_code=400, detail=str(e)) from e except Exception as e: logger.exception(f"Deployment failed ({type(e).__name__}): {e}") - raise HTTPException(status_code=500, detail=str(e)) + raise HTTPException(status_code=500, detail=str(e)) from e @router.get("") async def list_deployments( - service_id: Optional[str] = None, - unode_hostname: Optional[str] = None, + service_id: str | None = None, + unode_hostname: str | None = None, slim: bool = False, current_user: dict = Depends(get_current_user) ): @@ -331,11 +326,11 @@ async def list_deployments( # ============================================================================= @router.get("/exposed-urls") -async def get_exposed_urls( - url_type: Optional[str] = Query(None, alias="type", description="Filter by URL type (e.g., 'audio', 'http')"), - url_name: Optional[str] = Query(None, alias="name", description="Filter by URL name (e.g., 'audio_intake')"), - format: Optional[str] = Query(None, description="Filter by audio format (e.g., 'opus', 'pcm')"), - status: Optional[str] = Query(None, description="Filter by instance status (e.g., 'running')"), +async def get_exposed_urls( # noqa: C901 + url_type: str | None = Query(None, alias="type", description="Filter by URL type (e.g., 'audio', 'http')"), + url_name: str | None = Query(None, alias="name", description="Filter by URL name (e.g., 'audio_intake')"), + audio_format: str | None = Query(None, alias="format", description="Filter by audio format (e.g., 'opus', 'pcm')"), + status: str | None = Query(None, description="Filter by instance status (e.g., 'running')"), current_user: dict = Depends(get_current_user) ): """ @@ -347,9 +342,8 @@ async def get_exposed_urls( Example: GET /api/deployments/exposed-urls?type=audio&name=audio_intake&format=opus&status=running Returns: List of audio intake endpoints that support Opus format from Chronicle, Mycelia, etc. """ - from src.services.service_config_manager import get_service_config_manager - logger.info(f"[exposed-urls] Filtering by status={status}, url_type={url_type}, url_name={url_name}, format={format}") + logger.info(f"[exposed-urls] Filtering by status={status}, url_type={url_type}, url_name={url_name}, format={audio_format}") result = [] seen_containers = set() # Track container names to avoid duplicates @@ -359,17 +353,17 @@ async def get_exposed_urls( # Also check running docker containers from MANAGEABLE_SERVICES # This handles services started via docker compose that don't have service_config entries - from src.services.docker_manager import get_docker_manager - from src.services.compose_registry import get_compose_registry import os + from src.services.compose_registry import get_compose_registry + from src.services.docker_manager import get_docker_manager + docker_mgr = get_docker_manager() compose_registry = get_compose_registry() - project_name = os.getenv("COMPOSE_PROJECT_NAME", "ushadow") - logger.info(f"[exposed-urls] Checking MANAGEABLE_SERVICES for additional exposed URLs") + logger.info("[exposed-urls] Checking MANAGEABLE_SERVICES for additional exposed URLs") - for service_name in docker_mgr.MANAGEABLE_SERVICES.keys(): + for service_name in docker_mgr.MANAGEABLE_SERVICES: # Get service info to check if it's running service_info = docker_mgr.get_service_info(service_name) @@ -410,10 +404,8 @@ async def get_exposed_urls( continue # Filter by format if requested (check metadata.formats array) - if format: - supported_formats = exp_metadata.get('formats', []) - if format not in supported_formats: - continue + if audio_format and audio_format not in exp_metadata.get('formats', []): + continue # Build internal URL (for relay to connect to) # Audio endpoints use WebSocket protocol @@ -440,7 +432,7 @@ async def get_exposed_urls( # Also check running deployments on the local leader # This handles services started via deployment manager (not docker compose) - logger.info(f"[exposed-urls] Checking deployments for additional exposed URLs") + logger.info("[exposed-urls] Checking deployments for additional exposed URLs") # Get local hostname to filter for only local deployments # Use COMPOSE_PROJECT_NAME first as that's what unodes are registered with @@ -502,10 +494,8 @@ async def get_exposed_urls( continue # Filter by format if requested (check metadata.formats array) - if format: - supported_formats = exp_metadata.get('formats', []) - if format not in supported_formats: - continue + if audio_format and audio_format not in exp_metadata.get('formats', []): + continue # Build internal URL (for relay to connect to) # Audio endpoints use WebSocket protocol @@ -564,9 +554,8 @@ async def get_exposed_urls( continue if url_name and exp_name != url_name: continue - if format: - if format not in exp_metadata.get("formats", []): - continue + if audio_format and audio_format not in exp_metadata.get("formats", []): + continue protocol = "ws" if exp_type == "audio" else "http" exp_url = f"{protocol}://{k8s_host}:{port}{path}" result.append({ @@ -583,9 +572,66 @@ async def get_exposed_urls( except Exception as e: logger.warning(f"[exposed-urls] K8s discovery failed: {e}") + # ---- In-cluster discovery (when running inside a k8s pod, no MongoDB needed) ---- + from src.utils.environment import is_kubernetes + if is_kubernetes(): + try: + from pathlib import Path as _Path + + from kubernetes import client as _k8s_client + from kubernetes import config as _k8s_config + + _k8s_config.load_incluster_config() + _core_api = _k8s_client.CoreV1Api() + + _ns_file = _Path("/var/run/secrets/kubernetes.io/serviceaccount/namespace") + _namespace = _ns_file.read_text().strip() if _ns_file.exists() else "ushadow" + + _pods = _core_api.list_namespaced_pod(namespace=_namespace) + for _pod in _pods.items: + if _pod.status.phase != "Running": + continue + _labels = _pod.metadata.labels or {} + _svc_name = _labels.get("app.kubernetes.io/name") + if not _svc_name: + continue + _compose_svc = compose_registry.get_service_by_name(_svc_name) + if not _compose_svc or not getattr(_compose_svc, "exposes", None): + continue + _k8s_host = f"{_svc_name}.{_namespace}.svc.cluster.local" + if _k8s_host in seen_containers: + continue + for _expose in _compose_svc.exposes: + _exp_type = _expose.get("type") + _exp_name = _expose.get("name") + _path = _expose.get("path", "") + _port = _expose.get("port") + _exp_meta = _expose.get("metadata", {}) + if url_type and _exp_type != url_type: + continue + if url_name and _exp_name != url_name: + continue + if audio_format and audio_format not in _exp_meta.get("formats", []): + continue + _proto = "ws" if _exp_type == "audio" else "http" + _exp_url = f"{_proto}://{_k8s_host}:{_port}{_path}" + result.append({ + "instance_id": f"incluster:{_namespace}:{_pod.metadata.name}", + "instance_name": _compose_svc.display_name or _svc_name, + "url": _exp_url, + "type": _exp_type, + "name": _exp_name, + "metadata": _exp_meta, + "status": "running", + }) + seen_containers.add(_k8s_host) + logger.info(f"[exposed-urls] Added in-cluster URL: {_exp_name} -> {_exp_url} (pod {_pod.metadata.name})") + except Exception as _e: + logger.warning(f"[exposed-urls] In-cluster discovery failed: {_e}") + logger.info("=" * 80) logger.info(f"[exposed-urls] RETURNING {len(result)} TOTAL EXPOSED URLs") - logger.info(f"[exposed-urls] PARAMS: status={status}, url_type={url_type}, url_name={url_name}, format={format}") + logger.info(f"[exposed-urls] PARAMS: status={status}, url_type={url_type}, url_name={url_name}, format={audio_format}") logger.info(f"[exposed-urls] Unique containers collected: {len(seen_containers)}") logger.info("=" * 80) @@ -602,7 +648,7 @@ async def get_exposed_urls( return result -@router.get("/find/{service_id}", response_model=List[DiscoveredWorkload]) +@router.get("/find/{service_id}", response_model=list[DiscoveredWorkload]) async def find_workloads( service_id: str, current_user: dict = Depends(get_current_user), @@ -635,7 +681,7 @@ async def adopt_workload( return await manager.adopt_workload(service_id, req) except Exception as e: logger.error(f"Failed to adopt workload {service_id}: {e}") - raise HTTPException(status_code=500, detail=str(e)) + raise HTTPException(status_code=500, detail=str(e)) from e @router.get("/{deployment_id}", response_model=Deployment) @@ -659,13 +705,12 @@ async def stop_deployment( """Stop a deployment.""" manager = get_deployment_manager() try: - deployment = await manager.stop_deployment(deployment_id) - return deployment + return await manager.stop_deployment(deployment_id) except ValueError as e: - raise HTTPException(status_code=404, detail=str(e)) + raise HTTPException(status_code=404, detail=str(e)) from e except Exception as e: logger.error(f"Stop failed: {e}") - raise HTTPException(status_code=500, detail=str(e)) + raise HTTPException(status_code=500, detail=str(e)) from e @router.post("/{deployment_id}/restart", response_model=Deployment) @@ -676,18 +721,17 @@ async def restart_deployment( """Restart a deployment.""" manager = get_deployment_manager() try: - deployment = await manager.restart_deployment(deployment_id) - return deployment + return await manager.restart_deployment(deployment_id) except ValueError as e: - raise HTTPException(status_code=404, detail=str(e)) + raise HTTPException(status_code=404, detail=str(e)) from e except Exception as e: logger.error(f"Restart failed: {e}") - raise HTTPException(status_code=500, detail=str(e)) + raise HTTPException(status_code=500, detail=str(e)) from e class UpdateDeploymentRequest(BaseModel): """Request to update a deployment's environment variables.""" - env_vars: Dict[str, str] + env_vars: dict[str, str] @router.put("/{deployment_id}", response_model=Deployment) @@ -722,10 +766,10 @@ async def update_deployment( return updated_deployment except ValueError as e: - raise HTTPException(status_code=404, detail=str(e)) + raise HTTPException(status_code=404, detail=str(e)) from e except Exception as e: logger.error(f"Update deployment failed: {e}") - raise HTTPException(status_code=500, detail=str(e)) + raise HTTPException(status_code=500, detail=str(e)) from e @router.delete("/{deployment_id}") @@ -742,7 +786,7 @@ async def remove_deployment( raise except Exception as e: logger.error(f"Remove failed: {e}") - raise HTTPException(status_code=500, detail=str(e)) + raise HTTPException(status_code=500, detail=str(e)) from e @router.get("/{deployment_id}/logs") @@ -767,7 +811,7 @@ async def get_deployment_logs( async def get_funnel_configuration( deployment_id: str, current_user: dict = Depends(get_current_user) -) -> Dict[str, Any]: +) -> dict[str, Any]: """Get funnel configuration for a deployment. Returns funnel status, route, and public URL if configured. @@ -816,15 +860,15 @@ async def get_funnel_configuration( @router.patch("/{deployment_id}/funnel") async def configure_funnel_route( deployment_id: str, - request: Dict[str, Any], + request: dict[str, Any], current_user: dict = Depends(get_current_user) -) -> Dict[str, Any]: +) -> dict[str, Any]: """Configure funnel route for a deployment. Enables public internet access via Tailscale Funnel. """ - from src.services.tailscale_manager import get_tailscale_manager from src.services.service_config_manager import get_service_config_manager + from src.services.tailscale_manager import get_tailscale_manager route = request.get("route") save_to_config = request.get("save_to_config", False) @@ -920,13 +964,13 @@ async def remove_funnel_route( deployment_id: str, save_to_config: bool = Query(False), current_user: dict = Depends(get_current_user) -) -> Dict[str, Any]: +) -> dict[str, Any]: """Remove funnel route for a deployment. Disables public internet access for this deployment. """ - from src.services.tailscale_manager import get_tailscale_manager from src.services.service_config_manager import get_service_config_manager + from src.services.tailscale_manager import get_tailscale_manager manager = get_deployment_manager() diff --git a/ushadow/backend/src/services/deployment_manager.py b/ushadow/backend/src/services/deployment_manager.py index e0bba05c..8398e036 100644 --- a/ushadow/backend/src/services/deployment_manager.py +++ b/ushadow/backend/src/services/deployment_manager.py @@ -1,27 +1,26 @@ """Deployment manager for orchestrating services across u-nodes.""" -import asyncio import logging import os import uuid -from datetime import datetime, timezone -from typing import Any, Dict, List, Optional +from datetime import UTC, datetime +from typing import Any import aiohttp from motor.motor_asyncio import AsyncIOMotorDatabase +from src.models.deploy_target import DeployTarget from src.models.deployment import ( - ServiceDefinition, - ServiceDefinitionCreate, - ServiceDefinitionUpdate, + AdoptRequest, Deployment, DeploymentStatus, - ResolvedServiceDefinition, DiscoveredWorkload, - AdoptRequest, + ResolvedServiceDefinition, + ServiceDefinition, + ServiceDefinitionCreate, + ServiceDefinitionUpdate, ) from src.models.unode import UNode -from src.models.deploy_target import DeployTarget from src.services.compose_registry import get_compose_registry from src.services.deployment_platforms import get_deploy_platform from src.utils.environment import is_local_deployment as env_is_local_deployment @@ -54,8 +53,7 @@ def _update_tailscale_serve_route(service_id: str, container_name: str, port: in if add: return add_service_route(service_id, container_name, port) - else: - return remove_service_route(service_id) + return remove_service_route(service_id) except Exception as e: logger.warning(f"Failed to update tailscale serve route: {e}") return False @@ -82,7 +80,7 @@ def __init__(self, db: AsyncIOMotorDatabase): # self.deployments_collection = db.deployments self.unodes_collection = db.unodes self.adopted_workloads_collection = db.adopted_workloads - self._http_session: Optional[aiohttp.ClientSession] = None + self._http_session: aiohttp.ClientSession | None = None async def initialize(self): """Initialize indexes.""" @@ -150,6 +148,7 @@ async def _ensure_mycelia_tokens(self, force_regenerate: bool = False) -> dict: """ import subprocess from pathlib import Path + from src.config import get_settings settings = get_settings() @@ -206,11 +205,11 @@ async def _ensure_mycelia_tokens(self, force_regenerate: bool = False) -> dict: # Centralized Service Resolution # ========================================================================= - async def resolve_service_for_deployment( + async def resolve_service_for_deployment( # noqa: C901 self, service_id: str, - deploy_target: Optional[str] = None, - config_id: Optional[str] = None + deploy_target: str | None = None, + config_id: str | None = None ) -> "ResolvedServiceDefinition": """ Resolve all variables for a service using the new Settings API. @@ -246,8 +245,10 @@ async def resolve_service_for_deployment( ValueError: If service not found or resolution fails """ import subprocess - import yaml from pathlib import Path + + import yaml + from src.models.deployment import ResolvedServiceDefinition compose_registry = get_compose_registry() @@ -289,13 +290,12 @@ async def resolve_service_for_deployment( } # Build subprocess environment for docker-compose config (needs all vars for ${VAR} substitution) - import os # Strip K8s-injected service discovery vars (e.g. MONGODB_PORT=tcp://10.x.x.x:27017). # Kubernetes auto-injects these for every service in the namespace; Docker Compose # misinterprets the tcp:// prefix as a port bind address and raises "invalid IP address". subprocess_env = { k: v for k, v in os.environ.items() - if not (isinstance(v, str) and (v.startswith("tcp://") or v.startswith("udp://"))) + if not (isinstance(v, str) and v.startswith(("tcp://", "udp://"))) } subprocess_env.update(container_env) @@ -416,7 +416,6 @@ async def resolve_service_for_deployment( elif isinstance(vol, dict): # Long format: {"type": "volume", "source": "name", "target": "/path"} # or {"type": "bind", "source": "/host/path", "target": "/container/path"} - vol_type = vol.get("type", "volume") source = vol.get("source", "") target = vol.get("target", "") read_only = vol.get("read_only", False) @@ -489,14 +488,14 @@ async def resolve_service_for_deployment( return resolved - except subprocess.TimeoutExpired: - raise ValueError("docker-compose config timed out") + except subprocess.TimeoutExpired as e: + raise ValueError("docker-compose config timed out") from e except Exception as e: import traceback logger.error(f"Failed to resolve service {service_id}: {e}") logger.error(f"Exception type: {type(e).__name__}") logger.error(f"Traceback: {traceback.format_exc()}") - raise ValueError(f"Service resolution failed: {e}") + raise ValueError(f"Service resolution failed: {e}") from e # ========================================================================= # Service Definition CRUD @@ -505,10 +504,10 @@ async def resolve_service_for_deployment( async def create_service( self, data: ServiceDefinitionCreate, - created_by: Optional[str] = None + created_by: str | None = None ) -> ServiceDefinition: """Create a new service definition.""" - now = datetime.now(timezone.utc) + now = datetime.now(UTC) service = ServiceDefinition( service_id=data.service_id, @@ -534,7 +533,7 @@ async def create_service( logger.info(f"Created service definition: {service.service_id}") return service - async def list_services(self) -> List[ServiceDefinition]: + async def list_services(self) -> list[ServiceDefinition]: """List all service definitions.""" cursor = self.services_collection.find({}) services = [] @@ -542,7 +541,7 @@ async def list_services(self) -> List[ServiceDefinition]: services.append(ServiceDefinition(**doc)) return services - async def get_service(self, service_id: str) -> Optional[ServiceDefinition]: + async def get_service(self, service_id: str) -> ServiceDefinition | None: """Get a service definition by ID.""" doc = await self.services_collection.find_one({"service_id": service_id}) if doc: @@ -553,13 +552,13 @@ async def update_service( self, service_id: str, data: ServiceDefinitionUpdate - ) -> Optional[ServiceDefinition]: + ) -> ServiceDefinition | None: """Update a service definition.""" update_data = data.model_dump(exclude_unset=True) if not update_data: return await self.get_service(service_id) - update_data["updated_at"] = datetime.now(timezone.utc) + update_data["updated_at"] = datetime.now(UTC) result = await self.services_collection.find_one_and_update( {"service_id": service_id}, @@ -596,12 +595,12 @@ async def delete_service(self, service_id: str) -> bool: # Deployment Operations # ========================================================================= - async def deploy_service( + async def deploy_service( # noqa: C901 self, service_id: str, unode_hostname: str, config_id: str, - namespace: Optional[str] = None, + namespace: str | None = None, force_rebuild: bool = False, ) -> Deployment: """ @@ -628,7 +627,7 @@ async def deploy_service( logger.info(f"[DEBUG deploy_service] Resolved service has service_id={resolved_service.service_id}, name={resolved_service.name}") except ValueError as e: logger.error(f"Failed to resolve service {service_id}: {e}") - raise ValueError(f"Service resolution failed: {e}") + raise ValueError(f"Service resolution failed: {e}") from e # For mycelia services: generate tokens in MongoDB before deploying. # The script writes directly via pymongo — Mycelia does not need to be running. @@ -652,7 +651,7 @@ async def deploy_service( deployment_id = str(uuid.uuid4())[:8] # Create deployment target from unode with standardized fields - from src.models.unode import UNodeType, UNodeRole + from src.models.unode import UNodeRole, UNodeType from src.utils.deployment_targets import parse_deployment_target_id parsed = parse_deployment_target_id(unode.deployment_target_id) @@ -777,7 +776,7 @@ async def deploy_service( logger.warning(f"Could not configure Tailscale access URL: {e}") logger.debug("Deployment will continue without Tailscale URL") - deployment.deployed_at = datetime.now(timezone.utc) + deployment.deployed_at = datetime.now(UTC) except Exception as e: logger.error(f"Deploy failed for {service_id} on {unode_hostname}: {e}") @@ -814,7 +813,7 @@ async def stop_deployment(self, deployment_id: str) -> Deployment: # Refresh container status container.reload() deployment.status = DeploymentStatus.STOPPED - deployment.stopped_at = datetime.now(timezone.utc) + deployment.stopped_at = datetime.now(UTC) except Exception as e: logger.error(f"Failed to stop local deployment {deployment_id}: {e}") @@ -831,7 +830,7 @@ async def stop_deployment(self, deployment_id: str) -> Deployment: unode = UNode(**unode_dict) # Create deployment target from unode with standardized fields - from src.models.unode import UNodeType, UNodeRole + from src.models.unode import UNodeRole, UNodeType from src.utils.deployment_targets import parse_deployment_target_id parsed = parse_deployment_target_id(unode.deployment_target_id) @@ -860,7 +859,7 @@ async def stop_deployment(self, deployment_id: str) -> Deployment: if success: deployment.status = DeploymentStatus.STOPPED - deployment.stopped_at = datetime.now(timezone.utc) + deployment.stopped_at = datetime.now(UTC) except Exception as e: logger.error(f"Failed to stop remote deployment {deployment_id}: {e}") @@ -909,7 +908,7 @@ async def restart_deployment(self, deployment_id: str) -> Deployment: async def update_deployment( self, deployment_id: str, - env_vars: Dict[str, str] + env_vars: dict[str, str] ) -> Deployment: """ Update a deployment's environment variables and redeploy. @@ -917,9 +916,9 @@ async def update_deployment( Compares provided env_vars against what Settings would normally resolve (layers 1-5) and only saves actual overrides to ServiceConfig. """ - from src.services.service_config_manager import get_service_config_manager - from src.models.service_config import ServiceConfigCreate, ServiceConfigUpdate from src.config import get_settings + from src.models.service_config import ServiceConfigCreate, ServiceConfigUpdate + from src.services.service_config_manager import get_service_config_manager # Get existing deployment deployment = await self.get_deployment(deployment_id) @@ -978,7 +977,7 @@ async def update_deployment( id=config_id, template_id=deployment.service_id, name=f"{deployment.service_id} ({deployment.unode_hostname})", - description=f"Deployment configuration", + description="Deployment configuration", config=overrides_only, ) ) @@ -1024,8 +1023,8 @@ async def _remove_orphaned_container(self, deployment_id: str) -> bool: async def _stop_k8s_deployment(self, deployment: Deployment) -> Deployment: """Scale a Kubernetes deployment to 0 replicas.""" - from src.services.kubernetes import get_kubernetes_manager from src.services.deployment_platforms import KubernetesDeployPlatform + from src.services.kubernetes import get_kubernetes_manager from src.utils.environment import get_env_name cluster_id = deployment.backend_metadata.get("cluster_id") @@ -1056,7 +1055,7 @@ async def _stop_k8s_deployment(self, deployment: Deployment) -> Deployment: success = await platform.stop(target, deployment) if success: deployment.status = DeploymentStatus.STOPPED - deployment.stopped_at = datetime.now(timezone.utc) + deployment.stopped_at = datetime.now(UTC) except RuntimeError: pass # KubernetesManager not initialized @@ -1067,10 +1066,29 @@ async def _stop_k8s_deployment(self, deployment: Deployment) -> Deployment: return deployment + async def _unadopt_k8s_workload(self, deployment: Deployment) -> bool: + """Remove the adoption record for a K8s workload without touching the running workload.""" + cluster_id = deployment.backend_metadata.get("cluster_id") + container_name = deployment.container_name + query: dict = {"container_name": container_name, "backend_type": "kubernetes"} + if cluster_id: + query["cluster_id"] = cluster_id + result = await self.adopted_workloads_collection.delete_one(query) + if result.deleted_count == 0 and cluster_id: + # Retry without cluster_id in case the record predates that field + result = await self.adopted_workloads_collection.delete_one( + {"container_name": container_name, "backend_type": "kubernetes"} + ) + if result.deleted_count > 0: + logger.info(f"[unadopt] Removed adoption record for K8s workload {container_name}") + return True + logger.warning(f"[unadopt] No adoption record found for {container_name} (cluster={cluster_id})") + return False + async def _remove_k8s_deployment(self, deployment: Deployment) -> bool: """Remove a Kubernetes deployment via KubernetesDeployPlatform.""" - from src.services.kubernetes import get_kubernetes_manager from src.services.deployment_platforms import KubernetesDeployPlatform + from src.services.kubernetes import get_kubernetes_manager from src.utils.environment import get_env_name cluster_id = deployment.backend_metadata.get("cluster_id") @@ -1114,6 +1132,17 @@ async def remove_deployment(self, deployment_id: str) -> bool: # K8s deployments: route directly to KubernetesDeployPlatform if deployment.backend_type == "kubernetes": + # Check MongoDB for an adoption record — this covers both the case where + # get_deployment returned an adopted record (metadata.adopted=True) and the + # case where it returned a live-scan record that happens to have an adoption. + is_adopted = deployment.metadata.get("adopted") or bool( + await self.adopted_workloads_collection.find_one({ + "container_name": deployment.container_name, + "backend_type": "kubernetes", + }) + ) + if is_adopted: + return await self._unadopt_k8s_workload(deployment) return await self._remove_k8s_deployment(deployment) unode_dict = await self.unodes_collection.find_one({ @@ -1126,7 +1155,7 @@ async def remove_deployment(self, deployment_id: str) -> bool: unode = UNode(**unode_dict) - from src.models.unode import UNodeType, UNodeRole + from src.models.unode import UNodeRole, UNodeType from src.utils.deployment_targets import parse_deployment_target_id parsed = parse_deployment_target_id(unode.deployment_target_id) @@ -1183,13 +1212,13 @@ async def remove_deployment(self, deployment_id: str) -> bool: logger.info(f"Removed deployment: {deployment_id}") return True - async def get_deployment(self, deployment_id: str) -> Optional[Deployment]: + async def get_deployment(self, deployment_id: str) -> Deployment | None: # noqa: C901 """ Get a deployment by ID by querying runtime. Queries all online unodes and K8s clusters until deployment is found. """ - from src.models.unode import UNodeType, UNodeRole + from src.models.unode import UNodeRole, UNodeType from src.utils.deployment_targets import parse_deployment_target_id # Query all online unodes (Docker deployments) @@ -1225,8 +1254,8 @@ async def get_deployment(self, deployment_id: str) -> Optional[Deployment]: # Also search K8s clusters directly (mirrors list_deployments) try: - from src.services.kubernetes import get_kubernetes_manager from src.services.deployment_platforms import KubernetesDeployPlatform + from src.services.kubernetes import get_kubernetes_manager from src.utils.environment import get_env_name k8s_mgr = await get_kubernetes_manager() @@ -1258,14 +1287,63 @@ async def get_deployment(self, deployment_id: str) -> Optional[Deployment]: except Exception as e: logger.error(f"Failed to search K8s clusters in get_deployment: {e}") + # Also check adopted workloads in MongoDB (mirrors list_deployments adopted section) + try: + async for doc in self.adopted_workloads_collection.find({}): + backend_type = doc.get("backend_type", "docker") + container_name = doc.get("container_name") + cluster_id = doc.get("cluster_id", "unknown") + dep_id = ( + f"adopted-k8s-{cluster_id}-{container_name}" + if backend_type == "kubernetes" + else f"adopted-{container_name}" + ) + if dep_id != deployment_id: + continue + ports = doc.get("ports", []) + exposed_port = None + if ports: + import contextlib + with contextlib.suppress(ValueError, IndexError): + exposed_port = int(ports[0].split(":")[0]) + dep_backend_meta = ( + { + "cluster_id": cluster_id, + "namespace": doc.get("namespace"), + "k8s_deployment_name": doc.get("k8s_deployment_name") or container_name, + } + if backend_type == "kubernetes" + else {"compose_project": doc.get("compose_project")} + ) + return Deployment( + id=dep_id, + service_id=doc["service_id"], + config_id=doc.get("config_id"), + unode_hostname=( + f"{cluster_id}.k8s" if backend_type == "kubernetes" + else doc.get("node_hostname", "local") + ), + status=DeploymentStatus.RUNNING if doc.get("status") == "running" else DeploymentStatus.STOPPED, + container_name=container_name, + container_id=doc.get("container_id"), + deployed_config={"image": doc.get("image", ""), "ports": ports}, + backend_type=backend_type, + backend_metadata=dep_backend_meta, + metadata={"adopted": True}, + exposed_port=exposed_port, + access_url=doc.get("access_url"), + ) + except Exception as e: + logger.error(f"Failed to search adopted workloads in get_deployment: {e}") + return None - async def list_deployments( + async def list_deployments( # noqa: C901 self, - service_id: Optional[str] = None, - unode_hostname: Optional[str] = None, + service_id: str | None = None, + unode_hostname: str | None = None, local_only: bool = False, - ) -> List[Deployment]: + ) -> list[Deployment]: """ List deployments by querying runtime (Docker/K8s). @@ -1276,7 +1354,7 @@ async def list_deployments( Use this for proxy routing to prevent forwarding to remote environments (e.g. another ushadow instance's services). """ - from src.models.unode import UNodeType, UNodeRole + from src.models.unode import UNodeRole, UNodeType from src.utils.deployment_targets import parse_deployment_target_id all_deployments = [] @@ -1336,10 +1414,14 @@ async def list_deployments( logger.debug(f"[list_deployments] Checked {unode_count} unodes, found {len(all_deployments)} deployments so far") # Also query registered K8s clusters directly (stateless — K8s is source of truth) - if not unode_hostname: # Skip K8s scan when filtering by specific unode hostname + # In Docker mode: skip when local_only=True (K8s services are remote, different AUTH_SECRET_KEY). + # In K8s mode: always scan — K8s services are the local services, routable via cluster DNS. + from src.utils.environment import is_kubernetes as _is_k8s_env + _k8s_is_remote = local_only and not _is_k8s_env() + if not unode_hostname and not _k8s_is_remote: try: - from src.services.kubernetes import get_kubernetes_manager from src.services.deployment_platforms import KubernetesDeployPlatform + from src.services.kubernetes import get_kubernetes_manager from src.utils.environment import get_env_name k8s_mgr = await get_kubernetes_manager() @@ -1383,13 +1465,16 @@ async def list_deployments( adopted_query["service_id"] = service_id async for doc in self.adopted_workloads_collection.find(adopted_query): backend_type = doc.get("backend_type", "docker") + # Exclude K8s adopted workloads only in Docker mode with local_only=True. + # In K8s mode they are same-cluster services and should be reachable. + if local_only and backend_type == "kubernetes" and not _is_k8s_env(): + continue ports = doc.get("ports", []) exposed_port = None if ports: - try: + import contextlib + with contextlib.suppress(ValueError, IndexError): exposed_port = int(ports[0].split(":")[0]) - except (ValueError, IndexError): - pass if backend_type == "kubernetes": dep_id = f"adopted-k8s-{doc.get('cluster_id', 'unknown')}-{doc['container_name']}" @@ -1422,6 +1507,31 @@ async def list_deployments( except Exception as e: logger.warning(f"[list_deployments] Failed to query adopted workloads: {e}") + # Merge adopted K8s records into their live-scan twins to eliminate duplicates. + # Strategy: keep the live-scan record (natural ID, fresh state) and annotate it + # with adopted=True so that remove_deployment knows to un-adopt rather than delete. + # Adopted records with no live-scan twin (e.g., cross-namespace) are kept as-is. + adopted_by_container: dict[str, Deployment] = { + dep.container_name: dep + for dep in all_deployments + if dep.metadata.get("adopted") and dep.backend_type == "kubernetes" + } + if adopted_by_container: + merged_containers: set[str] = set() + result: list[Deployment] = [] + for dep in all_deployments: + if dep.metadata.get("adopted") and dep.backend_type == "kubernetes": + continue # Handled below after live-scan pass + if dep.backend_type == "kubernetes" and dep.container_name in adopted_by_container: + dep.metadata["adopted"] = True # Annotate: prefer un-adopt on remove + merged_containers.add(dep.container_name) + result.append(dep) + # Keep adopted records that have no live-scan twin (e.g., different namespace) + for dep in adopted_by_container.values(): + if dep.container_name not in merged_containers: + result.append(dep) + all_deployments = result + logger.debug(f"[list_deployments] Total deployments: {len(all_deployments)}") return all_deployments @@ -1429,7 +1539,7 @@ async def get_deployment_logs( self, deployment_id: str, tail: int = 100 - ) -> Optional[str]: + ) -> str | None: """Get logs for a deployment.""" deployment = await self.get_deployment(deployment_id) if not deployment: @@ -1444,7 +1554,7 @@ async def get_deployment_logs( unode = UNode(**unode_dict) # Create deployment target from unode with standardized fields - from src.models.unode import UNodeType, UNodeRole + from src.models.unode import UNodeRole, UNodeType from src.utils.deployment_targets import parse_deployment_target_id parsed = parse_deployment_target_id(unode.deployment_target_id) @@ -1479,13 +1589,13 @@ async def get_deployment_logs( # Node Communication # ========================================================================= - async def _get_node_url(self, unode: Dict[str, Any]) -> str: + async def _get_node_url(self, unode: dict[str, Any]) -> str: """Get the manager API URL for a u-node.""" # Prefer Tailscale IP for cross-node communication ip = unode.get("tailscale_ip") or unode.get("hostname") return f"http://{ip}:{MANAGER_PORT}" - async def _get_node_secret(self, unode: Dict[str, Any]) -> str: + async def _get_node_secret(self, unode: dict[str, Any]) -> str: """Get the secret for authenticating with a u-node.""" # Secret is stored encrypted - need to decrypt it encrypted_secret = unode.get("unode_secret_encrypted", "") @@ -1508,10 +1618,10 @@ async def _get_node_secret(self, unode: Dict[str, Any]) -> str: async def _send_deploy_command( self, - unode: Dict[str, Any], + unode: dict[str, Any], resolved_service: ResolvedServiceDefinition, container_name: str - ) -> Dict[str, Any]: + ) -> dict[str, Any]: """ Send deploy command to a u-node. @@ -1563,9 +1673,9 @@ async def _send_deploy_command( async def _send_stop_command( self, - unode: Dict[str, Any], + unode: dict[str, Any], container_name: str - ) -> Dict[str, Any]: + ) -> dict[str, Any]: """Send stop command to a u-node.""" session = await self._get_session() url = await self._get_node_url(unode) @@ -1582,9 +1692,9 @@ async def _send_stop_command( async def _send_restart_command( self, - unode: Dict[str, Any], + unode: dict[str, Any], container_name: str - ) -> Dict[str, Any]: + ) -> dict[str, Any]: """Send restart command to a u-node.""" session = await self._get_session() url = await self._get_node_url(unode) @@ -1601,9 +1711,9 @@ async def _send_restart_command( async def _send_remove_command( self, - unode: Dict[str, Any], + unode: dict[str, Any], container_name: str - ) -> Dict[str, Any]: + ) -> dict[str, Any]: """Send remove command to a u-node.""" session = await self._get_session() url = await self._get_node_url(unode) @@ -1620,10 +1730,10 @@ async def _send_remove_command( async def _send_logs_command( self, - unode: Dict[str, Any], + unode: dict[str, Any], container_name: str, tail: int = 100 - ) -> Dict[str, Any]: + ) -> dict[str, Any]: """Get logs from a container on a u-node.""" session = await self._get_session() url = await self._get_node_url(unode) @@ -1643,7 +1753,7 @@ async def _send_logs_command( # Find & Adopt # ========================================================================= - async def find_workloads(self, service_name: str) -> List[DiscoveredWorkload]: + async def find_workloads(self, service_name: str) -> list[DiscoveredWorkload]: # noqa: C901 """ Search Docker and all K8s clusters for workloads matching service_name. @@ -1652,7 +1762,7 @@ async def find_workloads(self, service_name: str) -> List[DiscoveredWorkload]: Returns both already-adopted and unadopted results. """ - results: List[DiscoveredWorkload] = [] + results: list[DiscoveredWorkload] = [] name_lower = service_name.lower() # Pre-load adopted container names from MongoDB for both Docker and K8s. @@ -1713,7 +1823,6 @@ async def find_workloads(self, service_name: str) -> List[DiscoveredWorkload]: if name_lower not in dep.metadata.name.lower(): continue ns = dep.metadata.namespace - dep_labels = dep.metadata.labels or {} containers = dep.spec.template.spec.containers or [] image = containers[0].image if containers else "unknown" ready = dep.status.ready_replicas or 0 @@ -1770,7 +1879,7 @@ async def adopt_workload(self, service_id: str, req: AdoptRequest) -> Deployment The workload is NOT restarted or otherwise modified. """ import re - now = datetime.now(timezone.utc) + now = datetime.now(UTC) # Derive a valid ServiceConfig ID from the container name safe_name = re.sub(r'[^a-z0-9-]', '-', req.container_name.lower()) @@ -1780,8 +1889,8 @@ async def adopt_workload(self, service_id: str, req: AdoptRequest) -> Deployment # Create a ServiceConfig so the adopted service appears in the wiring UI # and is configurable. Skip silently if it already exists. try: - from src.services.service_config_manager import get_service_config_manager from src.models.service_config import ServiceConfigCreate + from src.services.service_config_manager import get_service_config_manager config_manager = get_service_config_manager() if not config_manager.get_service_config(config_id): backend_label = "Kubernetes" if req.backend_type == "kubernetes" else "Docker" @@ -1823,10 +1932,9 @@ async def adopt_workload(self, service_id: str, req: AdoptRequest) -> Deployment namespace_val = req.namespace or "default" port = 8000 if req.ports: - try: + import contextlib + with contextlib.suppress(ValueError, IndexError): port = int(str(req.ports[0]).split(":")[-1]) - except (ValueError, IndexError): - pass svc_access_url = await k8s_mgr.get_service_access_url( req.cluster_id, req.container_name, namespace_val, port ) @@ -1903,7 +2011,7 @@ async def adopt_workload(self, service_id: str, req: AdoptRequest) -> Deployment ) - async def resolve_service_url(self, name: str) -> str: + async def resolve_service_url(self, name: str) -> str: # noqa: C901 """ Resolve the internal URL for a named service. @@ -1921,10 +2029,8 @@ async def resolve_service_url(self, name: str) -> str: Raises: ValueError: service not found or not reachable """ - import os - from src.services.docker_manager import get_docker_manager from src.services.compose_registry import get_compose_registry - from src.utils.environment import is_kubernetes + from src.services.docker_manager import get_docker_manager compose_registry = get_compose_registry() docker_mgr = get_docker_manager() @@ -1969,11 +2075,10 @@ async def resolve_service_url(self, name: str) -> str: ports = docker_mgr.get_service_ports(name) port = 8000 if ports: + import contextlib raw = ports[0].get("container_port", 8000) - try: + with contextlib.suppress(ValueError, TypeError): port = int(raw) - except (ValueError, TypeError): - pass try: container = docker_mgr._client.containers.get(info.container_id) return f"http://{container.name}:{port}" @@ -2030,7 +2135,7 @@ async def _resolve_k8s_url(self, dep, port: int) -> str: # Global instance -_deployment_manager: Optional[DeploymentManager] = None +_deployment_manager: DeploymentManager | None = None def get_deployment_manager() -> DeploymentManager: diff --git a/ushadow/frontend/src/components/conversations/ConversationCard.tsx b/ushadow/frontend/src/components/conversations/ConversationCard.tsx index b4aaae4a..b0488fb4 100644 --- a/ushadow/frontend/src/components/conversations/ConversationCard.tsx +++ b/ushadow/frontend/src/components/conversations/ConversationCard.tsx @@ -19,12 +19,10 @@ export default function ConversationCard({ conversation, source, onClick }: Conv has_memory, } = conversation - // Mycelia stores data differently than Chronicle - const myceliaConv = conversation as any - - // Extract start time - Mycelia uses timeRanges[0].start for actual conversation time + // Extract start time - Mycelia uses started_at (timeRanges[0].start) for actual conversation time // created_at in Mycelia is the processing timestamp, not the conversation time - const conversationDate = myceliaConv?.timeRanges?.[0]?.start || created_at + const myceliaConv = conversation as any + const conversationDate = myceliaConv?.started_at || created_at // Format date const date = conversationDate ? new Date(conversationDate) : null diff --git a/ushadow/frontend/src/pages/ConversationDetailPage.tsx b/ushadow/frontend/src/pages/ConversationDetailPage.tsx index ebb2e3de..28dbd84e 100644 --- a/ushadow/frontend/src/pages/ConversationDetailPage.tsx +++ b/ushadow/frontend/src/pages/ConversationDetailPage.tsx @@ -292,8 +292,9 @@ export default function ConversationDetailPage() { const hasValidSegments = conversation.segments && conversation.segments.length > 0 // Extract start/end times - const startTime = myceliaConv?.timeRanges?.[0]?.start || conversation.created_at - const endTime = myceliaConv?.timeRanges?.[0]?.end || conversation.completed_at + // For Mycelia: use started_at (actual recording start) over created_at (processing time) + const startTime = myceliaConv?.started_at || myceliaConv?.timeRanges?.[0]?.start || conversation.created_at + const endTime = myceliaConv?.completed_at || myceliaConv?.timeRanges?.[0]?.end // Format duration const formatDuration = (seconds?: number) => { @@ -318,23 +319,6 @@ export default function ConversationDetailPage() { return `${mins}m ${secs}s` } - // Format date - const formatDate = (dateString?: string) => { - if (!dateString) { - // Mycelia uses createdAt - if (myceliaConv?.createdAt) { - dateString = myceliaConv.createdAt - } else { - return 'Unknown' - } - } - try { - return new Date(dateString).toLocaleString() - } catch { - return dateString - } - } - const sourceColor = source === 'chronicle' ? 'blue' : 'purple' const sourceLabel = source === 'chronicle' ? 'Chronicle' : 'Mycelia' @@ -553,7 +537,7 @@ export default function ConversationDetailPage() { {/* Segmented transcript (only if segments have actual text) */} {hasValidSegments ? (
- {conversation.segments.map((segment, idx) => { + {conversation.segments!.map((segment, idx) => { const segmentId = `segment-${idx}` const isPlaying = playingSegment === segmentId diff --git a/ushadow/frontend/src/pages/LoginPage.tsx b/ushadow/frontend/src/pages/LoginPage.tsx index 40d65d0f..caf7462b 100644 --- a/ushadow/frontend/src/pages/LoginPage.tsx +++ b/ushadow/frontend/src/pages/LoginPage.tsx @@ -1,6 +1,7 @@ import React from 'react' import { useNavigate, useLocation } from 'react-router-dom' import { useCasdoorAuth } from '../contexts/CasdoorAuthContext' +import { useSettings } from '../contexts/SettingsContext' import AuthHeader from '../components/auth/AuthHeader' import { LogIn, ExternalLink, UserPlus } from 'lucide-react' import { setupApi } from '../services/api' @@ -9,6 +10,7 @@ export default function LoginPage() { const navigate = useNavigate() const location = useLocation() const { isAuthenticated, isLoading, login, register } = useCasdoorAuth() + const { isLoading: settingsLoading } = useSettings() const [hasUsers, setHasUsers] = React.useState(null) // Parse query parameters once @@ -178,7 +180,7 @@ export default function LoginPage() { )}