diff --git a/.github/buildkitd.toml b/.github/buildkitd.toml deleted file mode 100644 index 1ca32ce..0000000 --- a/.github/buildkitd.toml +++ /dev/null @@ -1,2 +0,0 @@ -[worker.oci] -max-parallelism = 1 diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 3e2fece..5d77a93 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -46,62 +46,134 @@ jobs: run: | cargo test --release + # One image build per architecture, on a native runner. We push by digest; + # a final `merge` job stitches the two digests into a multi-arch manifest list. + # Native arm64 runners avoid QEMU emulation, which dominated build time. image: - name: ๐Ÿณ Publish Image - runs-on: ubuntu-latest + name: ๐Ÿณ Build (${{ matrix.platform }}) needs: - test - outputs: - tag: ${{ steps.tag.outputs.tag }} + runs-on: ${{ matrix.runner }} + strategy: + fail-fast: false + matrix: + include: + - platform: linux/amd64 + runner: ubuntu-latest + - platform: linux/arm64 + runner: ubuntu-24.04-arm steps: + - name: ๐Ÿงท Platform pair + run: | + platform="${{ matrix.platform }}" + echo "PLATFORM_PAIR=${platform//\//-}" >> "${GITHUB_ENV}" + - name: ๐Ÿ›Ž๏ธ Checkout uses: actions/checkout@v5 - name: ๐Ÿ—ฝ Free disk space uses: ShubhamTatvamasi/free-disk-space-action@master - - name: ๐Ÿท Get tag - id: tag - run: | - if [ "${{ github.event_name }}" -eq "push" ] - then - echo "tag=${GITHUB_REF#refs/tags/}" >> "$GITHUB_OUTPUT" - else - echo "tag=ci-${{ github.run_number }}" >> "$GITHUB_OUTPUT" - fi - - name: ๐Ÿณ Login to DockerHub + if: github.event_name == 'push' uses: docker/login-action@v3 with: username: ${{ secrets.DOCKERHUB_USERNAME }} password: ${{ secrets.DOCKERHUB_TOKEN }} - - name: ๐Ÿ› ๏ธ Set up QEMU - uses: docker/setup-qemu-action@v3 - - name: ๐Ÿ› ๏ธ Set up Docker Buildx uses: docker/setup-buildx-action@v3 - with: - buildkitd-config: .github/buildkitd.toml - platforms: linux/amd64,linux/arm64 - - name: ๐Ÿณ Build and push + # PR validation: build the image but don't push. Cache scoped per arch so + # the two matrix entries don't clobber each other's GHA cache. + - name: ๐Ÿณ Build (PR validation) + if: github.event_name != 'push' uses: docker/build-push-action@v6 with: context: . build-args: BASE_IMAGE=platzio/base:v8 - push: ${{ github.event_name == 'push' }} - platforms: linux/amd64,linux/arm64 - tags: ${{ env.DOCKER_REPO }}:${{ steps.tag.outputs.tag }} - cache-from: type=gha - cache-to: type=gha,mode=max + platforms: ${{ matrix.platform }} + push: false + cache-from: type=gha,scope=build-${{ env.PLATFORM_PAIR }} + cache-to: type=gha,mode=max,scope=build-${{ env.PLATFORM_PAIR }} + + # Tag push: build and push by digest. The merge job below assembles the + # final tagged manifest. + - name: ๐Ÿณ Build & push by digest + if: github.event_name == 'push' + id: build + uses: docker/build-push-action@v6 + with: + context: . + build-args: BASE_IMAGE=platzio/base:v8 + platforms: ${{ matrix.platform }} + outputs: type=image,name=${{ env.DOCKER_REPO }},push-by-digest=true,name-canonical=true,push=true + cache-from: type=gha,scope=build-${{ env.PLATFORM_PAIR }} + cache-to: type=gha,mode=max,scope=build-${{ env.PLATFORM_PAIR }} + + - name: ๐Ÿ“ค Export digest + if: github.event_name == 'push' + run: | + mkdir -p /tmp/digests + digest="${{ steps.build.outputs.digest }}" + touch "/tmp/digests/${digest#sha256:}" + + - name: ๐Ÿ“ฆ Upload digest + if: github.event_name == 'push' + uses: actions/upload-artifact@v4 + with: + name: digests-${{ env.PLATFORM_PAIR }} + path: /tmp/digests/* + if-no-files-found: error + retention-days: 1 + + merge: + name: ๐Ÿงฌ Assemble manifest list + if: github.event_name == 'push' + needs: + - image + runs-on: ubuntu-latest + outputs: + tag: ${{ steps.tag.outputs.tag }} + steps: + - name: ๐Ÿท Get tag + id: tag + run: | + echo "tag=${GITHUB_REF#refs/tags/}" >> "${GITHUB_OUTPUT}" + + - name: ๐Ÿ“ฅ Download digests + uses: actions/download-artifact@v4 + with: + path: /tmp/digests + pattern: digests-* + merge-multiple: true + + - name: ๐Ÿณ Login to DockerHub + uses: docker/login-action@v3 + with: + username: ${{ secrets.DOCKERHUB_USERNAME }} + password: ${{ secrets.DOCKERHUB_TOKEN }} + + - name: ๐Ÿ› ๏ธ Set up Docker Buildx + uses: docker/setup-buildx-action@v3 + + - name: ๐Ÿงฌ Create manifest list and push + working-directory: /tmp/digests + run: | + docker buildx imagetools create \ + -t "${DOCKER_REPO}:${{ steps.tag.outputs.tag }}" \ + $(printf "${DOCKER_REPO}@sha256:%s " *) + + - name: ๐Ÿ” Inspect image + run: | + docker buildx imagetools inspect "${DOCKER_REPO}:${{ steps.tag.outputs.tag }}" release: name: ๐Ÿš€ Create Release if: github.event_name == 'push' runs-on: ubuntu-latest needs: - - image + - merge permissions: contents: write packages: write @@ -109,7 +181,7 @@ jobs: - name: ๐Ÿ—๏ธ Generate OpenAPI schema run: | docker run --rm \ - ${{ env.DOCKER_REPO }}:${{ needs.image.outputs.tag }} \ + ${{ env.DOCKER_REPO }}:${{ needs.merge.outputs.tag }} \ /root/platz-api openapi schema \ > openapi.yaml @@ -120,8 +192,8 @@ jobs: GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }} with: draft: true - tag_name: ${{ needs.image.outputs.tag }} - release_name: ${{ needs.image.outputs.tag }} + tag_name: ${{ needs.merge.outputs.tag }} + release_name: ${{ needs.merge.outputs.tag }} - name: ๐Ÿ“ฆ Upload OpenAPI schema uses: actions/upload-release-asset@v1 diff --git a/Cargo.lock b/Cargo.lock index 2028ad1..e588fa3 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1587,8 +1587,9 @@ dependencies = [ [[package]] name = "diesel_json" -version = "0.2.1" -source = "git+https://github.com/popen2/diesel_json#e376b6e3d30daa0c50bd09f77f40ecbba5c4cfed" +version = "0.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3b4e35cdca3d007ccca4cac79b2dd52469b74f7baf0cbdeac47bdcc7a31662e5" dependencies = [ "diesel", "serde", @@ -3249,16 +3250,19 @@ dependencies = [ "chrono", "clap", "futures", + "humantime", "itertools", "platz-chart-ext", "platz-db", "platz-otel", "regex", + "reqwest", "serde", "serde_json", "titlecase", "tokio", "tracing", + "url", "uuid", ] diff --git a/Dockerfile b/Dockerfile index 27c423f..17f2a19 100644 --- a/Dockerfile +++ b/Dockerfile @@ -1,22 +1,112 @@ +# Multi-stage build using cargo-chef for dep-only layer caching and per-arch +# cache mounts so amd64 and arm64 don't fight over the same target dir. +# +# Drives the workspace into a static musl binary so the runtime image (alpine- +# based platzio/base) doesn't need a libc. Architecture is selected via Docker +# Buildx' automatic TARGETARCH build arg โ€” set platforms in the build invocation, +# not here. + ARG BASE_IMAGE +ARG RUST_IMAGE=rust:1-trixie -FROM rust:1-trixie AS builder -RUN apt-get update && \ - apt-get install -y \ +# --------------------------------------------------------------------------- +# 1. chef โ€” Rust toolchain + musl tooling + cargo-chef. Shared by planner and +# builder so both layers reuse the same toolchain image. +# --------------------------------------------------------------------------- +FROM ${RUST_IMAGE} AS chef +RUN apt-get update && apt-get install -y --no-install-recommends \ musl \ musl-dev \ - musl-tools + musl-tools \ + && rm -rf /var/lib/apt/lists/* +RUN cargo install cargo-chef --locked --version ^0.1 +WORKDIR /build -FROM builder AS build +# --------------------------------------------------------------------------- +# 2. planner โ€” strip the workspace down to a "recipe" describing the dep graph. +# This stage is invalidated by *any* source change, but it's cheap (no compile). +# --------------------------------------------------------------------------- +FROM chef AS planner +COPY . . +RUN cargo chef prepare --recipe-path recipe.json + +# --------------------------------------------------------------------------- +# 3. builder โ€” cook deps from the recipe, then build the workspace. The cooked +# deps live in a buildkit cache mount keyed by TARGETARCH, so each architecture +# keeps its own warm target dir across CI runs. +# --------------------------------------------------------------------------- +FROM chef AS builder ARG RELEASE_BUILD=1 +ARG TARGETARCH + +RUN set -eux; \ + case "${TARGETARCH}" in \ + amd64) target=x86_64-unknown-linux-musl ;; \ + arm64) target=aarch64-unknown-linux-musl ;; \ + *) echo "Unsupported TARGETARCH: ${TARGETARCH}" >&2; exit 1 ;; \ + esac; \ + echo "${target}" > /target.txt; \ + rustup target add "${target}" + +COPY --from=planner /build/recipe.json recipe.json +RUN --mount=type=cache,id=platz-cargo-target-${TARGETARCH},target=/build/target,sharing=locked \ + --mount=type=cache,id=platz-cargo-git,target=/usr/local/cargo/git,sharing=locked \ + --mount=type=cache,id=platz-cargo-registry,target=/usr/local/cargo/registry,sharing=locked \ + set -eux; \ + target="$(cat /target.txt)"; \ + if [ "${RELEASE_BUILD}" = "1" ]; then \ + cargo chef cook --release --target "${target}" --recipe-path recipe.json; \ + else \ + cargo chef cook --target "${target}" --recipe-path recipe.json; \ + fi + +COPY . . +RUN --mount=type=cache,id=platz-cargo-target-${TARGETARCH},target=/build/target,sharing=locked \ + --mount=type=cache,id=platz-cargo-git,target=/usr/local/cargo/git,sharing=locked \ + --mount=type=cache,id=platz-cargo-registry,target=/usr/local/cargo/registry,sharing=locked \ + set -eux; \ + target="$(cat /target.txt)"; \ + if [ "${RELEASE_BUILD}" = "1" ]; then \ + cargo build --release --target "${target}"; \ + out_dir="target/${target}/release"; \ + else \ + cargo build --target "${target}"; \ + out_dir="target/${target}/debug"; \ + fi; \ + mkdir -p /out; \ + find "${out_dir}" -maxdepth 1 -type f -executable -exec cp -v {} /out/ \; + +# --------------------------------------------------------------------------- +# 4a. dev โ€” debug-mode build kept in the Rust toolchain image so Tilt's +# live_update can run `cargo build` *inside* the container after syncing +# source. No cache mount on the build step: the warm /build/target dir +# survives into the resulting image and incremental rebuilds in the running +# container reuse it. Dynamic-linked debian runtime (libpq5) โ€” the musl- +# static release path isn't useful when we're recompiling at runtime. +# +# Selected by `--target=dev` (the Tiltfile in platzio/dev sets this). Not +# referenced by any other stage, so default builds skip it. +# --------------------------------------------------------------------------- +FROM ${RUST_IMAGE} AS dev +RUN apt-get update && apt-get install -y --no-install-recommends \ + ca-certificates \ + libpq5 \ + libpq-dev \ + && rm -rf /var/lib/apt/lists/* WORKDIR /build -RUN mkdir -p /build/outputs -COPY . /build -RUN --mount=type=cache,id=platz-backend-cargo-target,target=/build/target,sharing=locked \ - --mount=type=cache,id=platz-backend-cargo-git,target=/usr/local/cargo/git,sharing=locked \ - --mount=type=cache,id=platz-backend-cargo-registry,target=/usr/local/cargo/registry,sharing=locked \ - ./scripts/container-build.sh "${RELEASE_BUILD}" "/build/outputs" - -FROM $BASE_IMAGE +COPY . . +RUN cargo build --workspace --bins \ + && mkdir -p /root \ + && cp target/debug/platz-api /root/platz-api \ + && cp target/debug/platz-k8s-agent /root/platz-k8s-agent \ + && cp target/debug/platz-chart-discovery /root/platz-chart-discovery \ + && cp target/debug/platz-status-updates /root/platz-status-updates \ + && cp target/debug/platz-resource-sync /root/platz-resource-sync +WORKDIR /root + +# --------------------------------------------------------------------------- +# 4. runtime โ€” small base image carrying just the static musl binaries. +# --------------------------------------------------------------------------- +FROM ${BASE_IMAGE} WORKDIR /root/ -COPY --from=build /build/outputs/* /root/ +COPY --from=builder /out/* /root/ diff --git a/README.md b/README.md index b54df34..4f9a347 100644 --- a/README.md +++ b/README.md @@ -7,45 +7,28 @@ This repo contains Platz's backend. It's written in Rust ๐Ÿฆ€ and is broken down * `platz-k8s-agent` * `platz-chart-discovery` * `platz-status-updates` +* `platz-resource-sync` ## Running Locally -```bash -./scripts/run-api.sh -``` +Local development uses Tilt + kind, set up from the sibling +[`platzio/dev`](https://github.com/platzio/dev) repo. Clone it next to this +repo and follow its README. -This script starts a local Postgres and Dex containers. +### Config knobs -Dex starts up with one user `admin@example.com` with a password of `password`. +The workers honor a few environment variables that the local setup wires up +for you and that production deployments set via the platzio helm chart: -Forwarded ports: - -* Port `3000`: API server -* Port `9000`: OIDC provider -* Port `15432`: Postgres database - -To connect to the development database: - -```bash -PGHOST=127.0.0.1 PGPORT=15432 PGUSER=postgres PGPASSWORD=postgres psql -``` - -## Developing Against Production Database - -> โš ๏ธ This is potentially dangerous, don't use this method after large or dangerous changes. - -Run the following command in a separate tab: - -```bash -kubectl --context=control -n platz port-forward platz-postgresql-0 5432:5432 -``` - -Then run any of the workers with the `DATABASE_URL` environment variable set from the directory of the crate you'd like to execute: - -```bash -cd api -DATABASE_URL=postgres://postgres:postgres@localhost:5432/platz cargo run -- api -``` +* `platz-k8s-agent` reads `PLATZ_CLUSTER_PROVIDER` (default `eks`): + * `eks` โ€” discovers EKS clusters across all AWS regions in the running account. + * `local` โ€” registers a single cluster from a kubeconfig context. +* `platz-chart-discovery` reads `PLATZ_REGISTRY_PROVIDER` (default `ecr`): + * `ecr` โ€” watches an SQS queue fed by ECR push/delete events. + * `oci` โ€” periodically polls a generic OCI registry (set by `PLATZ_OCI_REGISTRY_URL`) + for new chart artifacts. +* `platz-api` reads `OIDC_*` environment variables for OIDC config and + `ADMIN_EMAILS` as a space-delimited allow-list. ## Crates Overview @@ -58,14 +41,17 @@ In addition, this crate is responsible for distributing database notifications. ### `platz-api` The API is a worker that serves the API and handles user authentication. - -For authenticating users, it gets the OIDC information from AWS SSM parameters which should be passed as command line arguments. +OIDC parameters are passed via the `OIDC_*` environment variables. +`ADMIN_EMAILS` is a space-delimited allowlist. ### `platz-k8s-agent` This worker tracks Kubernetes clusters, updates their status in the database, and keeps a fresh copy of credentials allowing other parts in the worker to communicate with Kubernetes clusters. -Discovering and communicating with Kubernetes clusters is based on permissions to AWS EKS. This worker automatically discovers EKS clusters from all regions in the same AWS account it's running in. +In `eks` mode, the worker discovers EKS clusters across all regions in the +same AWS account it's running in. In `local` mode it registers a single +cluster from the configured kubeconfig context (`PLATZ_LOCAL_CONTEXT`, +defaulting to the kubeconfig's `current-context`). The first part that needs access to Kubernetes clusters is the `deploy` module. This module watches for pending deployment tasks and runs them one by one. @@ -88,7 +74,11 @@ The second part is the `k8s/tracker` module. It watches Kubernetes resources and ### `platz-chart-discovery` -This worker discovers charts uploaded to ECR: Terraform defines the appropriate AWS resources to receive these notifications, and this worker watches for events in an SQS queue. +This worker discovers Helm charts pushed to a registry. + +* In `ecr` mode it watches an SQS queue fed by ECR push/delete events. +* In `oci` mode it polls a generic OCI registry (e.g. Docker Distribution + `registry:2` or zot) for new helm-typed artifacts. See *Helm Chart Extensions* below for more information. @@ -98,22 +88,10 @@ This worker is responsible for watching Platz deployments that have enabled the For each deployment, Platz queries the status endpoint and updates the deployment's status in the database. The frontend can then display this information. -## Terraform - -Platz has a Terraform part that creates the AWS resources necessary for running some of the workers above. - -After making changes, you can apply them by running: - -``` -./run-terraform.sh -``` - -with the appropriate AWS credentials (`aws sso login`, define `AWS_PROFILE`, etc.) - -The main Terraform files are: +### `platz-resource-sync` -* `api-role.tf` allows `platz-api` to access the OIDC SSM parameters. -* `ecr-events.tf` creates the SQS queue for receiving events on ECR activity. +Watches Kubernetes Namespaces, Deployments, StatefulSets, and Jobs that +belong to Platz deployments and reflects their state into the database. ## Helm Chart Extensions diff --git a/chart-discovery/Cargo.toml b/chart-discovery/Cargo.toml index 4589b05..d66b6e1 100644 --- a/chart-discovery/Cargo.toml +++ b/chart-discovery/Cargo.toml @@ -16,14 +16,17 @@ chrono = { version = "0.4.41", default-features = false, features = [ ] } clap = { version = "4.5.41", features = ["derive", "env"] } futures = "0.3.31" +humantime = "2.2.0" itertools = "0.14.0" platz-chart-ext = { workspace = true } regex = "1.11.1" +reqwest = { version = "0.12.23", default-features = false, features = ["json", "rustls-tls"] } serde = { version = "1.0.219", features = ["derive"] } serde_json = "1.0.141" titlecase = "3.6.0" tokio = { version = "1.47.0", features = ["rt-multi-thread", "signal"] } tracing = "0.1.41" +url = { version = "2.5.4", features = ["serde"] } uuid = { version = "1.17.0", features = ["serde", "v4"] } [dependencies.platz-db] diff --git a/chart-discovery/src/charts.rs b/chart-discovery/src/charts.rs index f42fd41..0fdac05 100644 --- a/chart-discovery/src/charts.rs +++ b/chart-discovery/src/charts.rs @@ -2,16 +2,18 @@ use crate::ecr_events::{EcrEvent, EcrEventDetail}; use crate::tag_parser::parse_image_tag; use anyhow::{Result, anyhow}; use aws_smithy_types_convert::date_time::DateTimeExt; +use chrono::prelude::*; use platz_chart_ext::ChartExt; use platz_db::{ Json, schema::helm_chart::{HelmChart, HelmChartTagInfo, NewHelmChart, UpdateHelmChart}, }; -use std::path::PathBuf; +use std::path::{Path, PathBuf}; use tokio::process::Command; use tracing::{debug, info, warn}; +use uuid::Uuid; -const HELM_ARTIFACT_MEDIA_TYPE: &str = "application/vnd.cncf.helm.config.v1+json"; +pub const HELM_ARTIFACT_MEDIA_TYPE: &str = "application/vnd.cncf.helm.config.v1+json"; const TEMP_DOWNLOAD_PATH: &str = "/tmp/platz-chart-download"; pub async fn add_helm_chart(ecr: &aws_sdk_ecr::Client, event: EcrEvent) -> Result<()> { @@ -45,14 +47,38 @@ pub async fn add_helm_chart(ecr: &aws_sdk_ecr::Client, event: EcrEvent) -> Resul let helm_registry = event.find_or_create_ecr_repo().await?; - let chart_ext = match download_chart(&event).await { + let image_digest = event.detail.image_digest.clone(); + let image_tag = event.detail.image_tag.clone(); + + let chart_ext = match download_chart_via_ecr(&event).await { Ok(path) => ChartExt::from_path(&path).await?, Err(err) => ChartExt::new_with_error(err.to_string()), }; - let tag_info = match chart_ext.metadata { + record_helm_chart( + helm_registry.id, + image_digest, + image_tag, + created_at, + chart_ext, + ) + .await?; + + Ok(()) +} + +/// Provider-agnostic insertion of a `HelmChart` row given a prepared `ChartExt` +/// and the chart's identity in its registry. +pub async fn record_helm_chart( + helm_registry_id: Uuid, + image_digest: String, + image_tag: String, + created_at: DateTime, + chart_ext: ChartExt, +) -> Result<()> { + let tag_info = match chart_ext.metadata.clone() { Some(metadata) if metadata.git_commit.is_none() && metadata.git_branch.is_none() => { - let parsed = parse_image_tag(&event.detail.image_tag).await?; + let parsed = parse_image_tag(&image_tag).await?; HelmChartTagInfo { tag_format_id: parsed.tag_format_id, parsed_commit: parsed.parsed_commit, @@ -68,14 +94,14 @@ pub async fn add_helm_chart(ecr: &aws_sdk_ecr::Client, event: EcrEvent) -> Resul parsed_version: Some(metadata.version), parsed_revision: None, }, - None => parse_image_tag(&event.detail.image_tag).await?, + None => parse_image_tag(&image_tag).await?, }; let chart = NewHelmChart { created_at, - helm_registry_id: helm_registry.id, - image_digest: event.detail.image_digest, - image_tag: event.detail.image_tag, + helm_registry_id, + image_digest, + image_tag, values_ui: chart_ext.ui_schema.map(Json), actions_schema: chart_ext.actions.map(Json), features: chart_ext.features.map(Json), @@ -94,9 +120,7 @@ pub async fn add_helm_chart(ecr: &aws_sdk_ecr::Client, event: EcrEvent) -> Resul Ok(()) } -async fn download_chart(event: &EcrEvent) -> Result { - info!("Downloading chart"); - let path = PathBuf::from(TEMP_DOWNLOAD_PATH); +async fn download_chart_via_ecr(event: &EcrEvent) -> Result { let script = [ "rm -rf $TEMP_DOWNLOAD_PATH", "mkdir -p $TEMP_DOWNLOAD_PATH", @@ -105,22 +129,62 @@ async fn download_chart(event: &EcrEvent) -> Result { "helm pull oci://$HELM_REGISTRY/$HELM_REPO --version $HELM_CHART_TAG -d ./ --untar", ].join(" && "); - let output = Command::new("/bin/bash") - .arg("-euxc") - .arg(&script) - .env("TEMP_DOWNLOAD_PATH", TEMP_DOWNLOAD_PATH) - .env("HELM_REGISTRY_REGION", &event.region) - .env("HELM_REGISTRY", event.helm_registry_domain_name()) - .env("HELM_REPO", &event.detail.repository_name) - .env("HELM_CHART_TAG", &event.detail.image_tag) - .spawn()? - .wait_with_output() - .await?; + let extra_env = [ + ("HELM_REGISTRY_REGION", event.region.clone()), + ("HELM_REGISTRY", event.helm_registry_domain_name()), + ("HELM_REPO", event.detail.repository_name.clone()), + ("HELM_CHART_TAG", event.detail.image_tag.clone()), + ]; + + download_via_helm_pull(&script, &extra_env).await +} + +/// Downloads and untars a chart from a generic OCI registry. The registry must be +/// reachable anonymously; no `helm registry login` is performed. +pub async fn download_chart_via_oci( + registry_domain: &str, + repo_name: &str, + image_tag: &str, +) -> Result { + let script = [ + "rm -rf $TEMP_DOWNLOAD_PATH", + "mkdir -p $TEMP_DOWNLOAD_PATH", + "cd $TEMP_DOWNLOAD_PATH", + "helm pull oci://$HELM_REGISTRY/$HELM_REPO --version $HELM_CHART_TAG -d ./ --untar", + ] + .join(" && "); + + let extra_env = [ + ("HELM_REGISTRY", registry_domain.to_owned()), + ("HELM_REPO", repo_name.to_owned()), + ("HELM_CHART_TAG", image_tag.to_owned()), + ]; + + download_via_helm_pull(&script, &extra_env).await +} + +async fn download_via_helm_pull(script: &str, extra_env: &[(&str, String)]) -> Result { + info!("Downloading chart"); + let path = PathBuf::from(TEMP_DOWNLOAD_PATH); + + let mut cmd = Command::new("/bin/bash"); + cmd.arg("-euxc") + .arg(script) + .env("TEMP_DOWNLOAD_PATH", TEMP_DOWNLOAD_PATH); + for (key, value) in extra_env { + cmd.env(key, value); + } + + let output = cmd.spawn()?.wait_with_output().await?; info!("Finished downloading chart ({:?})", output.status); info!("Stdout: {}", String::from_utf8_lossy(&output.stdout)); info!("Stderr: {}", String::from_utf8_lossy(&output.stderr)); + pick_single_entry(&path).await +} + +async fn pick_single_entry(path: &Path) -> Result { let mut dir = tokio::fs::read_dir(path).await?; let mut files = Vec::new(); while let Some(entry) = dir.next_entry().await? { diff --git a/chart-discovery/src/ecr_events.rs b/chart-discovery/src/ecr_events.rs index 9b1f7de..f86de57 100644 --- a/chart-discovery/src/ecr_events.rs +++ b/chart-discovery/src/ecr_events.rs @@ -13,12 +13,29 @@ use uuid::Uuid; #[group(skip)] pub struct Config { #[clap(long, env = "PLATZ_ECR_EVENTS_QUEUE")] - ecr_events_queue: String, + ecr_events_queue: Option, #[clap(long, env = "PLATZ_ECR_EVENTS_REGION")] + ecr_events_region: Option, +} + +#[derive(Debug)] +pub struct ResolvedConfig { + ecr_events_queue: String, ecr_events_region: String, } +impl Config { + /// Returns a `ResolvedConfig` only when both fields are present, since they + /// are now optional at the CLI level (they're only required in `provider=ecr` mode). + pub fn resolved(&self) -> Option { + Some(ResolvedConfig { + ecr_events_queue: self.ecr_events_queue.clone()?, + ecr_events_region: self.ecr_events_region.clone()?, + }) + } +} + #[derive(Debug, Deserialize)] #[allow(dead_code)] pub struct EcrEvent { @@ -134,9 +151,9 @@ async fn handle_ecr_event(ecr: &aws_sdk_ecr::Client, event: EcrEvent) -> Result< } } -pub async fn run(config: &Config) -> Result<()> { +pub async fn run(config: &ResolvedConfig) -> Result<()> { info!("Starting to watch for ECR events"); - let region = Region::new(config.ecr_events_region.to_owned()); + let region = Region::new(config.ecr_events_region.clone()); let shared_config = aws_config::load_defaults(aws_config::BehaviorVersion::latest()).await; let ecr_config = aws_sdk_ecr::config::Builder::from(&shared_config) diff --git a/chart-discovery/src/main.rs b/chart-discovery/src/main.rs index 6f840bb..9284983 100644 --- a/chart-discovery/src/main.rs +++ b/chart-discovery/src/main.rs @@ -1,5 +1,5 @@ -use anyhow::Result; -use clap::Parser; +use anyhow::{Result, anyhow}; +use clap::{Parser, ValueEnum}; use platz_db::{DbTable, NotificationListeningOpts, init_db}; use tokio::{ select, @@ -10,15 +10,32 @@ use tracing::{info, warn}; mod charts; mod ecr_events; mod kind; +mod oci_poll; mod registries; mod sqs; mod tag_parser; +#[derive(Copy, Clone, Debug, PartialEq, Eq, ValueEnum)] +#[clap(rename_all = "lowercase")] +enum RegistryProvider { + /// Watch an SQS queue fed by ECR push/delete events. Requires AWS credentials. + Ecr, + /// Periodically poll a generic OCI registry (e.g. Docker Distribution / zot) + /// for new chart artifacts. Used by local dev and air-gapped setups. + Oci, +} + #[derive(Debug, Parser)] pub struct Config { + #[clap(long, env = "PLATZ_REGISTRY_PROVIDER", value_enum, default_value = "ecr")] + provider: RegistryProvider, + #[clap(flatten)] ecr_events: ecr_events::Config, + #[clap(flatten)] + oci_poll: oci_poll::Config, + #[clap(long, default_value_t = false)] enable_tag_parser: bool, } @@ -40,6 +57,21 @@ async fn main() -> Result<()> { } }; + let provider_fut = async { + match config.provider { + RegistryProvider::Ecr => { + let ecr_config = config + .ecr_events + .resolved() + .ok_or_else(|| anyhow!( + "PLATZ_ECR_EVENTS_QUEUE and PLATZ_ECR_EVENTS_REGION must be set when provider=ecr" + ))?; + ecr_events::run(&ecr_config).await + } + RegistryProvider::Oci => oci_poll::run(&config.oci_poll).await, + } + }; + select! { _ = sigterm.recv() => { info!("SIGTERM received, exiting"); @@ -58,7 +90,7 @@ async fn main() -> Result<()> { result.map_err(Into::into) } - result = ecr_events::run(&config.ecr_events) => { + result = provider_fut => { result } diff --git a/chart-discovery/src/oci_poll.rs b/chart-discovery/src/oci_poll.rs new file mode 100644 index 0000000..81106c3 --- /dev/null +++ b/chart-discovery/src/oci_poll.rs @@ -0,0 +1,274 @@ +use crate::charts::{ + HELM_ARTIFACT_MEDIA_TYPE, download_chart_via_oci, record_helm_chart, +}; +use crate::kind::get_or_create_kind; +use anyhow::{Result, anyhow}; +use chrono::prelude::*; +use clap::Parser; +use platz_chart_ext::ChartExt; +use platz_db::schema::helm_chart::HelmChart; +use platz_db::schema::helm_registry::{ + HelmRegistry, HelmRegistryProvider, NewHelmRegistry, +}; +use serde::Deserialize; +use std::collections::HashSet; +use tokio::time; +use tracing::{debug, info, warn}; +use url::Url; +use uuid::Uuid; + +#[derive(Debug, Parser)] +#[group(skip)] +pub struct Config { + /// Base URL of the OCI registry to poll, e.g. `http://localhost:5001`. + /// Required when `PLATZ_REGISTRY_PROVIDER=oci`. + #[clap(long, env = "PLATZ_OCI_REGISTRY_URL")] + pub oci_registry_url: Option, + + /// How often to poll the registry's `_catalog` and tag lists for new charts. + #[clap(long, env = "PLATZ_OCI_POLL_INTERVAL", default_value = "5s")] + pub oci_poll_interval: humantime::Duration, +} + +pub async fn run(config: &Config) -> Result<()> { + let url = config + .oci_registry_url + .clone() + .ok_or_else(|| anyhow!("PLATZ_OCI_REGISTRY_URL is required when provider=oci"))?; + let domain = registry_domain(&url)?; + let client = reqwest::Client::builder().build()?; + + info!("Polling OCI registry at {url}"); + + let mut interval = time::interval(*config.oci_poll_interval); + let mut seen_pairs: HashSet<(String, String)> = HashSet::new(); + + loop { + interval.tick().await; + if let Err(err) = poll_once(&client, &url, &domain, &mut seen_pairs).await { + warn!("OCI poll iteration failed: {err:?}"); + } + } +} + +async fn poll_once( + client: &reqwest::Client, + url: &Url, + domain: &str, + seen_pairs: &mut HashSet<(String, String)>, +) -> Result<()> { + let catalog = list_repos(client, url).await?; + debug!("Catalog: {} repos", catalog.len()); + + for repo in catalog { + let tags = match list_tags(client, url, &repo).await { + Ok(tags) => tags, + Err(err) => { + warn!("Failed listing tags for {repo}: {err:?}"); + continue; + } + }; + + for tag in tags { + let key = (repo.clone(), tag.clone()); + if seen_pairs.contains(&key) { + continue; + } + match handle_tag(client, url, domain, &repo, &tag).await { + Ok(()) => { + seen_pairs.insert(key); + } + Err(err) => { + warn!("Failed handling {repo}:{tag}: {err:?}"); + } + } + } + } + + Ok(()) +} + +async fn handle_tag( + client: &reqwest::Client, + url: &Url, + domain: &str, + repo: &str, + tag: &str, +) -> Result<()> { + let manifest = fetch_manifest(client, url, repo, tag).await?; + + if manifest.body.config.media_type != HELM_ARTIFACT_MEDIA_TYPE { + debug!( + "Skipping {repo}:{tag} (media type {})", + manifest.body.config.media_type + ); + return Ok(()); + } + + // The Distribution API doesn't expose pushed-at, so we synthesise it from + // "now" the first time we observe the tag. Chart.yaml metadata still wins + // for any version/branch/commit info. + let created_at = Utc::now(); + + let registry = ensure_registry(domain, repo).await?; + + if HelmChart::find_by_registry_and_digest(registry.id, manifest.digest.clone()) + .await? + .is_some() + { + debug!("Chart {repo}:{tag} already in DB, skipping"); + return Ok(()); + } + + let chart_ext = match download_chart_via_oci(domain, repo, tag).await { + Ok(path) => ChartExt::from_path(&path).await?, + Err(err) => ChartExt::new_with_error(err.to_string()), + }; + + record_helm_chart( + registry.id, + manifest.digest, + tag.to_owned(), + created_at, + chart_ext, + ) + .await?; + + Ok(()) +} + +async fn ensure_registry(domain: &str, repo: &str) -> Result { + if let Some(reg) = + HelmRegistry::find_by_domain_and_repo(domain.to_owned(), repo.to_owned()).await? + { + return Ok(reg); + } + let new = NewHelmRegistry { + created_at: Utc::now(), + domain_name: domain.to_owned(), + repo_name: repo.to_owned(), + kind_id: get_or_create_kind(repo).await?.id, + provider: HelmRegistryProvider::Oci, + }; + info!("Saving new OCI helm registry {:?}", new); + Ok(new.insert().await?) +} + +#[derive(Debug, Deserialize)] +struct CatalogResponse { + repositories: Vec, +} + +#[derive(Debug, Deserialize)] +struct TagsResponse { + tags: Option>, +} + +#[derive(Debug)] +struct ManifestFetchResult { + digest: String, + body: ManifestBody, +} + +#[derive(Debug, Deserialize)] +struct ManifestBody { + config: ManifestConfig, +} + +#[derive(Debug, Deserialize)] +struct ManifestConfig { + #[serde(rename = "mediaType")] + media_type: String, +} + +async fn list_repos(client: &reqwest::Client, base: &Url) -> Result> { + let url = base.join("v2/_catalog")?; + let resp: CatalogResponse = client + .get(url) + .send() + .await? + .error_for_status()? + .json() + .await?; + Ok(resp.repositories) +} + +async fn list_tags(client: &reqwest::Client, base: &Url, repo: &str) -> Result> { + let url = base.join(&format!("v2/{repo}/tags/list"))?; + let resp: TagsResponse = client + .get(url) + .send() + .await? + .error_for_status()? + .json() + .await?; + Ok(resp.tags.unwrap_or_default()) +} + +async fn fetch_manifest( + client: &reqwest::Client, + base: &Url, + repo: &str, + tag: &str, +) -> Result { + let url = base.join(&format!("v2/{repo}/manifests/{tag}"))?; + let resp = client + .get(url) + .header( + "Accept", + "application/vnd.oci.image.manifest.v1+json, \ + application/vnd.docker.distribution.manifest.v2+json", + ) + .send() + .await? + .error_for_status()?; + + let digest = resp + .headers() + .get("Docker-Content-Digest") + .and_then(|v| v.to_str().ok()) + .map(|s| s.to_owned()) + .unwrap_or_else(|| { + // Fallback: synthesise a non-colliding key from repo+tag if the registry + // doesn't return the digest header (some bare-bones implementations). + format!("oci-pseudo:{repo}:{tag}:{}", Uuid::new_v4()) + }); + + let body: ManifestBody = resp.json().await?; + Ok(ManifestFetchResult { digest, body }) +} + +fn registry_domain(url: &Url) -> Result { + let host = url + .host_str() + .ok_or_else(|| anyhow!("OCI registry URL has no host: {url}"))?; + let domain = match url.port() { + Some(port) => format!("{host}:{port}"), + None => host.to_owned(), + }; + Ok(domain) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn registry_domain_includes_port() { + let url = Url::parse("http://localhost:5001").unwrap(); + assert_eq!(registry_domain(&url).unwrap(), "localhost:5001"); + } + + #[test] + fn registry_domain_omits_default_port() { + let url = Url::parse("https://example.com").unwrap(); + assert_eq!(registry_domain(&url).unwrap(), "example.com"); + } + + #[test] + fn registry_domain_with_explicit_port() { + let url = Url::parse("https://oci.example.com:443/").unwrap(); + // url crate normalises the default port for the scheme away. + assert_eq!(registry_domain(&url).unwrap(), "oci.example.com"); + } +} diff --git a/chart-discovery/src/registries.rs b/chart-discovery/src/registries.rs index 2b8f328..217c6cb 100644 --- a/chart-discovery/src/registries.rs +++ b/chart-discovery/src/registries.rs @@ -4,7 +4,7 @@ use aws_sdk_ecr::types::Repository; use aws_smithy_types_convert::date_time::DateTimeExt; use aws_types::region::Region; use itertools::Itertools; -use platz_db::schema::helm_registry::{HelmRegistry, NewHelmRegistry}; +use platz_db::schema::helm_registry::{HelmRegistry, HelmRegistryProvider, NewHelmRegistry}; use tracing::info; pub async fn find_and_save_ecr_repo(region: Region, repo_name: &str) -> Result { @@ -51,6 +51,7 @@ async fn save_ecr_repo_in_db(repo: Repository) -> Result { domain_name: domain_name.to_owned(), repo_name: repo_name.to_owned(), kind_id: get_or_create_kind(repo_name).await?.id, + provider: HelmRegistryProvider::Ecr, }; info!("Saving {:?}", new_registry); Ok(new_registry.insert().await?) diff --git a/db/Cargo.toml b/db/Cargo.toml index 1649480..e92cbbe 100644 --- a/db/Cargo.toml +++ b/db/Cargo.toml @@ -24,7 +24,7 @@ diesel-async = { version = "0.6.1", features = [ diesel_enum_derive = { git = "https://github.com/popen2/diesel-enum-derive", default-features = false, features = [ "plain", ] } -diesel_json = { git = "https://github.com/popen2/diesel_json" } +diesel_json = "0.3.0" diesel_migrations = "2.2.0" itertools = "0.14.0" lazy_static = "1.5.0" diff --git a/db/migrations/2026-05-09-120000_helm_registry_provider/down.sql b/db/migrations/2026-05-09-120000_helm_registry_provider/down.sql new file mode 100644 index 0000000..3f84519 --- /dev/null +++ b/db/migrations/2026-05-09-120000_helm_registry_provider/down.sql @@ -0,0 +1,2 @@ +ALTER TABLE helm_registries + DROP COLUMN provider; diff --git a/db/migrations/2026-05-09-120000_helm_registry_provider/up.sql b/db/migrations/2026-05-09-120000_helm_registry_provider/up.sql new file mode 100644 index 0000000..1e008b9 --- /dev/null +++ b/db/migrations/2026-05-09-120000_helm_registry_provider/up.sql @@ -0,0 +1,2 @@ +ALTER TABLE helm_registries + ADD COLUMN provider VARCHAR NOT NULL DEFAULT 'Ecr'; diff --git a/db/src/schema/helm_registry.rs b/db/src/schema/helm_registry.rs index 31880d4..00621fc 100644 --- a/db/src/schema/helm_registry.rs +++ b/db/src/schema/helm_registry.rs @@ -2,11 +2,13 @@ use crate::{DbError, DbResult, db_conn}; use chrono::prelude::*; use diesel::prelude::*; use diesel_async::RunQueryDsl; +use diesel_enum_derive::DieselEnum; use diesel_filter::DieselFilter; use diesel_pagination::{Paginate, Paginated, PaginationParams}; use itertools::Itertools; use serde::{Deserialize, Serialize}; use std::ops::DerefMut; +use strum::{Display, EnumString}; use utoipa::ToSchema; use uuid::Uuid; @@ -19,9 +21,34 @@ table! { kind_id -> Uuid, available -> Bool, fa_icon -> Varchar, + provider -> Varchar, } } +#[derive( + Debug, + Clone, + Copy, + Default, + PartialEq, + Eq, + Serialize, + Deserialize, + EnumString, + Display, + DieselEnum, + ToSchema, +)] +pub enum HelmRegistryProvider { + /// AWS Elastic Container Registry. Authenticated with `aws ecr get-login-password`. + #[default] + Ecr, + /// Generic OCI registry (e.g. Docker Distribution `registry:2`, zot, ghcr.io). + /// No authentication is performed by Platz; the registry is expected to be + /// anonymous-readable from the cluster running the Helm pod. + Oci, +} + #[derive(Debug, Identifiable, Queryable, Serialize, DieselFilter, ToSchema)] #[diesel(table_name = helm_registries)] pub struct HelmRegistry { @@ -34,6 +61,7 @@ pub struct HelmRegistry { pub kind_id: Uuid, pub available: bool, pub fa_icon: String, + pub provider: HelmRegistryProvider, } impl HelmRegistry { @@ -72,13 +100,20 @@ impl HelmRegistry { .optional()?) } - pub fn region_name(&self) -> DbResult { - let (_aws_account_id, _dkr, _ecr, region_name, _amazonaws, _com) = self - .domain_name - .split('.') - .collect_tuple() - .ok_or(DbError::RegionNameParseError)?; - Ok(region_name.into()) + /// Returns the AWS region embedded in the registry's domain name, when applicable. + /// Returns `Ok(None)` for non-ECR providers, since the concept doesn't apply. + pub fn region_name(&self) -> DbResult> { + match self.provider { + HelmRegistryProvider::Oci => Ok(None), + HelmRegistryProvider::Ecr => { + let (_aws_account_id, _dkr, _ecr, region_name, _amazonaws, _com) = self + .domain_name + .split('.') + .collect_tuple() + .ok_or(DbError::RegionNameParseError)?; + Ok(Some(region_name.into())) + } + } } } @@ -89,6 +124,8 @@ pub struct NewHelmRegistry { pub domain_name: String, pub repo_name: String, pub kind_id: Uuid, + #[serde(default)] + pub provider: HelmRegistryProvider, } impl NewHelmRegistry { diff --git a/docker-compose.yaml b/docker-compose.yaml deleted file mode 100644 index f528d4d..0000000 --- a/docker-compose.yaml +++ /dev/null @@ -1,32 +0,0 @@ -version: "3.9" -services: - db: - image: postgres:17-alpine - environment: - POSTGRES_PASSWORD: postgres - ports: - - 15432:5432 - idp: - image: dexidp/dex - volumes: - - ./scripts/dex.config.yaml:/etc/dex/config.docker.yaml:ro - ports: - - 5556:5556 - platz-api: - build: - context: . - args: - BASE_IMAGE: platzio/base:v8 - RELEASE_BUILD: 0 - command: - - /root/platz-api - - run - environment: - ADMIN_EMAILS: admin@example.com - DATABASE_URL: postgres://postgres:postgres@db:5432 - OIDC_SERVER_URL: http://idp:5556 - OIDC_CLIENT_ID: foo - OIDC_CLIENT_SECRET: bar - PLATZ_OWN_URL: https://localhost:8080 - ports: - - 3000:3000 diff --git a/k8s-agent/src/config.rs b/k8s-agent/src/config.rs index 0cb92a8..05293d3 100644 --- a/k8s-agent/src/config.rs +++ b/k8s-agent/src/config.rs @@ -15,7 +15,7 @@ pub struct Config { #[arg(long, env = "PLATZ_HELM_IMAGE")] pub helm_image: String, - #[arg(long, default_value = "false")] + #[arg(long, env = "PLATZ_DISABLE_DEPLOYMENT_CREDENTIALS", default_value = "false")] pub disable_deployment_credentials: bool, #[arg(long, env = "PLATZ_OWN_URL")] diff --git a/k8s-agent/src/k8s/cluster_discovery.rs b/k8s-agent/src/k8s/cluster_discovery.rs index e110c33..be010dc 100644 --- a/k8s-agent/src/k8s/cluster_discovery.rs +++ b/k8s-agent/src/k8s/cluster_discovery.rs @@ -1,16 +1,45 @@ -use super::{cluster_type::K8s, tracker::K8S_TRACKER}; +use super::{ + cluster_type::{K8s, LocalCluster}, + tracker::K8S_TRACKER, +}; use anyhow::{Result, anyhow}; use aws_types::region::Region; +use clap::ValueEnum; use futures::future::try_join_all; +use std::path::PathBuf; use std::sync::Arc; use tokio::time; -use tracing::{debug, error}; +use tracing::{debug, error, info}; + +#[derive(Copy, Clone, Debug, PartialEq, Eq, ValueEnum)] +#[clap(rename_all = "lowercase")] +pub enum ClusterProvider { + /// Discover clusters by scanning AWS EKS in every region of the running account. + Eks, + /// Register a single cluster from a kubeconfig file. + Local, +} #[derive(clap::Args)] #[group(skip)] pub struct Config { #[arg(long, env = "K8S_REFRESH_INTERVAL", default_value = "1h")] pub k8s_refresh_interval: humantime::Duration, + + /// Selects how clusters are discovered. Defaults to `eks` (production behaviour); + /// set to `local` for laptop/dev workflows that target a kubeconfig context. + #[arg(long, env = "PLATZ_CLUSTER_PROVIDER", value_enum, default_value = "eks")] + pub provider: ClusterProvider, + + /// Path to the kubeconfig file used in `local` mode. + /// Falls back to `$KUBECONFIG`, then `~/.kube/config`. + #[arg(long, env = "PLATZ_LOCAL_KUBECONFIG")] + pub local_kubeconfig: Option, + + /// Name of the kubeconfig context to register in `local` mode. + /// Defaults to the kubeconfig's `current-context`. + #[arg(long, env = "PLATZ_LOCAL_CONTEXT")] + pub local_context: Option, } pub async fn run_cluster_discovery(config: &Config) -> Result<()> { @@ -18,16 +47,16 @@ pub async fn run_cluster_discovery(config: &Config) -> Result<()> { loop { interval.tick().await; - if let Err(err) = load_clusters().await { + if let Err(err) = load_clusters(config).await { error!("Error scanning for clusters: {:?}", err); } } } -async fn load_clusters() -> Result<()> { +async fn load_clusters(config: &Config) -> Result<()> { let tracker_tx = K8S_TRACKER.inbound_requests_tx().await; - for cluster in discover_clusters().await?.into_iter() { + for cluster in discover_clusters(config).await?.into_iter() { tracing::debug!(%cluster); tracker_tx.send(Arc::new(cluster))?; } @@ -35,9 +64,16 @@ async fn load_clusters() -> Result<()> { Ok(()) } -#[tracing::instrument(err, ret)] -async fn discover_clusters() -> Result> { - debug!("starting..."); +#[tracing::instrument(skip_all, err, ret)] +async fn discover_clusters(config: &Config) -> Result> { + match config.provider { + ClusterProvider::Eks => discover_eks_clusters().await, + ClusterProvider::Local => discover_local_cluster(config).await.map(|c| vec![c]), + } +} + +async fn discover_eks_clusters() -> Result> { + debug!("starting EKS discovery..."); let shared_config = aws_config::load_defaults(aws_config::BehaviorVersion::latest()).await; let ec2 = aws_sdk_ec2::Client::new(&shared_config); debug!("discovering regions..."); @@ -86,3 +122,73 @@ async fn get_clusters(region: Region) -> Result> { Ok(eks_clusters.into_iter().map(K8s::from).collect()) } + +#[tracing::instrument(skip_all, err)] +async fn discover_local_cluster(config: &Config) -> Result { + let kubeconfig = load_local_kubeconfig(config.local_kubeconfig.as_deref()).await?; + + let context_name = match &config.local_context { + Some(name) => name.clone(), + None => kubeconfig + .current_context + .clone() + .ok_or_else(|| anyhow!("Kubeconfig has no current-context and PLATZ_LOCAL_CONTEXT is not set"))?, + }; + + let context = kubeconfig + .contexts + .iter() + .find(|c| c.name == context_name) + .ok_or_else(|| anyhow!("Context {context_name:?} not found in kubeconfig"))?; + let context_inner = context + .context + .as_ref() + .ok_or_else(|| anyhow!("Context {context_name:?} has no body"))?; + + let cluster_name = context_inner.cluster.clone(); + let user_name = context_inner + .user + .clone() + .ok_or_else(|| anyhow!("Context {context_name:?} has no user"))?; + + let cluster = kubeconfig + .clusters + .iter() + .find(|c| c.name == cluster_name) + .ok_or_else(|| anyhow!("Cluster {cluster_name:?} not found in kubeconfig"))? + .clone(); + let auth_info = kubeconfig + .auth_infos + .iter() + .find(|a| a.name == user_name) + .ok_or_else(|| anyhow!("User {user_name:?} not found in kubeconfig"))? + .clone(); + + let scoped_kubeconfig = kube::config::Kubeconfig { + api_version: kubeconfig.api_version.clone(), + kind: kubeconfig.kind.clone(), + preferences: kubeconfig.preferences.clone(), + current_context: Some(context_name.clone()), + contexts: vec![context.clone()], + clusters: vec![cluster], + auth_infos: vec![auth_info], + extensions: kubeconfig.extensions.clone(), + }; + + info!("Registering local cluster from context {context_name:?}"); + + Ok(K8s::Local(Box::new(LocalCluster { + name: context_name.clone(), + provider_id: format!("local:{context_name}"), + kubeconfig: scoped_kubeconfig, + }))) +} + +async fn load_local_kubeconfig(explicit_path: Option<&std::path::Path>) -> Result { + if let Some(path) = explicit_path { + debug!("Loading kubeconfig from {}", path.display()); + return Ok(kube::config::Kubeconfig::read_from(path)?); + } + debug!("Loading kubeconfig from KUBECONFIG / default location"); + Ok(kube::config::Kubeconfig::read()?) +} diff --git a/k8s-agent/src/k8s/cluster_type.rs b/k8s-agent/src/k8s/cluster_type.rs index cf9694e..204830d 100644 --- a/k8s-agent/src/k8s/cluster_type.rs +++ b/k8s-agent/src/k8s/cluster_type.rs @@ -6,9 +6,24 @@ use std::convert::TryFrom; use std::fmt; use tracing::debug; +const LOCAL_REGION: &str = "local"; + #[derive(Debug)] pub enum K8s { - Eks(aws_sdk_eks::types::Cluster), + Eks(Box), + Local(Box), +} + +#[derive(Debug, Clone)] +pub struct LocalCluster { + /// Display name for the cluster, derived from the kubeconfig context name. + pub name: String, + /// Synthetic provider id (e.g. `local:platz-local`) used as a stable key in + /// the `k8s_clusters.provider_id` column. Avoids colliding with EKS ARNs. + pub provider_id: String, + /// The kubeconfig that drives both the in-process kube client and the + /// helm pod's `KUBECONFIG_BASE64`. + pub kubeconfig: kube::config::Kubeconfig, } impl fmt::Display for K8s { @@ -22,13 +37,20 @@ impl fmt::Display for K8s { .as_ref() .unwrap_or(&String::from("unknown")) ), + Self::Local(cluster) => write!(f, "Local({})", cluster.name), } } } impl From for K8s { fn from(cluster: aws_sdk_eks::types::Cluster) -> Self { - Self::Eks(cluster) + Self::Eks(Box::new(cluster)) + } +} + +impl From for K8s { + fn from(cluster: LocalCluster) -> Self { + Self::Local(Box::new(cluster)) } } @@ -43,32 +65,30 @@ impl K8s { .name .as_ref() .ok_or_else(|| anyhow!("Cluster has empty name"))?, + K8s::Local(cluster) => &cluster.name, }) } - fn server_url(&self) -> Result<&str> { + fn provider_id(&self) -> Result { Ok(match self { K8s::Eks(cluster) => cluster - .endpoint + .arn .as_ref() - .ok_or_else(|| anyhow!("Got empty endpoint"))?, + .ok_or_else(|| anyhow!("Cluster has no ARN"))? + .clone(), + K8s::Local(cluster) => cluster.provider_id.clone(), }) } - fn ca_data(&self) -> Result<&str> { + fn region_name(&self) -> Result { Ok(match self { - K8s::Eks(cluster) => cluster - .certificate_authority - .as_ref() - .ok_or_else(|| anyhow!("No certificate_authority for cluster"))? - .data - .as_ref() - .ok_or_else(|| anyhow!("certificate_authority didn't contain any data"))?, + K8s::Eks(_) => self.eks_region()?.into(), + K8s::Local(_) => LOCAL_REGION.to_owned(), }) } - fn region(&self) -> Result { - Ok(match self { + fn eks_region(&self) -> Result { + match self { K8s::Eks(cluster) => { let resource_name: aws_arn::ResourceName = cluster .arn @@ -78,17 +98,35 @@ impl K8s { .map_err(|err| anyhow!("Failed parsing region from ARN: {}", err))?; resource_name .region - .ok_or_else(|| anyhow!("Cluster ARN has no region"))? + .ok_or_else(|| anyhow!("Cluster ARN has no region")) } - }) + K8s::Local(_) => Err(anyhow!("Local clusters have no AWS region")), + } } pub async fn kube_config(&self) -> Result { let kubeconfig = kube::config::Kubeconfig::try_from(self)?; + let context = kubeconfig + .current_context + .clone() + .or_else(|| kubeconfig.contexts.first().map(|c| c.name.clone())) + .ok_or_else(|| anyhow!("Kubeconfig has no contexts"))?; + let cluster = kubeconfig + .clusters + .first() + .ok_or_else(|| anyhow!("Kubeconfig has no clusters"))? + .name + .clone(); + let user = kubeconfig + .auth_infos + .first() + .ok_or_else(|| anyhow!("Kubeconfig has no auth infos"))? + .name + .clone(); let kubeconfig_options = kube::config::KubeConfigOptions { - context: Some(kubeconfig.contexts.first().unwrap().name.clone()), - cluster: Some(kubeconfig.clusters.first().unwrap().name.clone()), - user: Some(kubeconfig.auth_infos.first().unwrap().name.clone()), + context: Some(context), + cluster: Some(cluster), + user: Some(user), }; Ok(kube::Config::from_custom_kubeconfig(kubeconfig, &kubeconfig_options).await?) } @@ -103,14 +141,11 @@ impl K8s { impl From<&K8s> for NewK8sCluster { fn from(cluster: &K8s) -> Self { - let region_name = cluster.region().unwrap().into(); - match cluster { - K8s::Eks(cluster) => Self { - provider_id: cluster.arn.as_ref().unwrap().clone(), - name: cluster.name.as_ref().unwrap().clone(), - env_id: None, - region_name, - }, + Self { + provider_id: cluster.provider_id().unwrap(), + name: cluster.name().unwrap().to_owned(), + env_id: None, + region_name: cluster.region_name().unwrap(), } } } @@ -119,55 +154,87 @@ impl TryFrom<&K8s> for kube::config::Kubeconfig { type Error = anyhow::Error; fn try_from(k8s: &K8s) -> Result { - let cluster = k8s.name()?; - let user = "user"; - let server_url = k8s.server_url()?; - Ok(Self { - api_version: Some("v1".to_owned()), - kind: Some("Config".to_owned()), - clusters: vec![kube::config::NamedCluster { - name: cluster.into(), - cluster: Some(kube::config::Cluster { - server: Some(server_url.into()), - insecure_skip_tls_verify: Some(false), - certificate_authority_data: Some(k8s.ca_data()?.into()), - ..Default::default() - }), - }], - auth_infos: vec![kube::config::NamedAuthInfo { - name: user.to_owned(), - auth_info: Some(kube::config::AuthInfo { - exec: Some(kube::config::ExecConfig { - command: Some("aws".into()), - args: Some(vec![ - "eks".into(), - "get-token".into(), - "--region".into(), - k8s.region()?.into(), - "--cluster-name".into(), - cluster.into(), - ]), - api_version: Some("client.authentication.k8s.io/v1".to_owned()), - interactive_mode: Some(ExecInteractiveMode::Never), - env: None, - drop_env: None, - provide_cluster_info: false, - cluster: None, - }), - ..Default::default() - }), - }], - contexts: vec![kube::config::NamedContext { - name: "default".to_owned(), - context: Some(kube::config::Context { - cluster: cluster.into(), - user: Some(user.to_owned()), - namespace: None, - extensions: None, - }), - }], - current_context: Some("default".to_owned()), - ..Default::default() - }) + match k8s { + K8s::Eks(cluster) => eks_kubeconfig(cluster), + K8s::Local(cluster) => Ok(cluster.kubeconfig.clone()), + } } } + +fn eks_kubeconfig(cluster: &aws_sdk_eks::types::Cluster) -> Result { + let cluster_name = cluster + .name + .as_ref() + .ok_or_else(|| anyhow!("Cluster has empty name"))?; + let server_url = cluster + .endpoint + .as_ref() + .ok_or_else(|| anyhow!("Got empty endpoint"))?; + let ca_data = cluster + .certificate_authority + .as_ref() + .ok_or_else(|| anyhow!("No certificate_authority for cluster"))? + .data + .as_ref() + .ok_or_else(|| anyhow!("certificate_authority didn't contain any data"))?; + let region: String = { + let resource_name: aws_arn::ResourceName = cluster + .arn + .as_ref() + .ok_or_else(|| anyhow!("Cluster has no ARN"))? + .parse() + .map_err(|err| anyhow!("Failed parsing region from ARN: {}", err))?; + resource_name + .region + .ok_or_else(|| anyhow!("Cluster ARN has no region"))? + .into() + }; + let user = "user"; + Ok(kube::config::Kubeconfig { + api_version: Some("v1".to_owned()), + kind: Some("Config".to_owned()), + clusters: vec![kube::config::NamedCluster { + name: cluster_name.clone(), + cluster: Some(kube::config::Cluster { + server: Some(server_url.clone()), + insecure_skip_tls_verify: Some(false), + certificate_authority_data: Some(ca_data.clone()), + ..Default::default() + }), + }], + auth_infos: vec![kube::config::NamedAuthInfo { + name: user.to_owned(), + auth_info: Some(kube::config::AuthInfo { + exec: Some(kube::config::ExecConfig { + command: Some("aws".into()), + args: Some(vec![ + "eks".into(), + "get-token".into(), + "--region".into(), + region, + "--cluster-name".into(), + cluster_name.clone(), + ]), + api_version: Some("client.authentication.k8s.io/v1".to_owned()), + interactive_mode: Some(ExecInteractiveMode::Never), + env: None, + drop_env: None, + provide_cluster_info: false, + cluster: None, + }), + ..Default::default() + }), + }], + contexts: vec![kube::config::NamedContext { + name: "default".to_owned(), + context: Some(kube::config::Context { + cluster: cluster_name.clone(), + user: Some(user.to_owned()), + namespace: None, + extensions: None, + }), + }], + current_context: Some("default".to_owned()), + ..Default::default() + }) +} diff --git a/k8s-agent/src/main.rs b/k8s-agent/src/main.rs index 9f741b7..ca141fe 100644 --- a/k8s-agent/src/main.rs +++ b/k8s-agent/src/main.rs @@ -49,7 +49,7 @@ async fn main() -> Result<()> { } result = run_cluster_discovery(&config.cluster_discovery) => { - warn!("EKS discovery task finished"); + warn!("Cluster discovery task finished"); result } diff --git a/k8s-agent/src/task_runner/helm.rs b/k8s-agent/src/task_runner/helm.rs index ff5f35e..837d909 100644 --- a/k8s-agent/src/task_runner/helm.rs +++ b/k8s-agent/src/task_runner/helm.rs @@ -10,7 +10,9 @@ use k8s_openapi::{ apimachinery::pkg::apis::meta::v1::ObjectMeta, }; use platz_db::schema::{ - deployment::Deployment, deployment_task::DeploymentTask, helm_registry::HelmRegistry, + deployment::Deployment, + deployment_task::DeploymentTask, + helm_registry::{HelmRegistry, HelmRegistryProvider}, }; use tracing::debug; @@ -54,18 +56,69 @@ async fn helm_pod( let chart = task.helm_chart().await?; let registry = HelmRegistry::find(chart.helm_registry_id).await?; - let script = [ - "mkdir -p /root/.kube", - "echo $KUBECONFIG_BASE64 | base64 -d > /root/.kube/config", - "chmod 400 /root/.kube/config", - "aws ecr get-login-password --region $HELM_REGISTRY_REGION | helm registry login --username AWS --password-stdin $HELM_REGISTRY", - "helm pull oci://$HELM_REGISTRY/$HELM_REPO --version $HELM_CHART_TAG", - "echo $VALUES_BASE64 | base64 -d > values.yaml", - "echo $VALUES_OVERRIDE_BASE64 | base64 -d > values-override.yaml", - &format!( + let mut script_lines: Vec = vec![ + "mkdir -p /root/.kube".into(), + "echo $KUBECONFIG_BASE64 | base64 -d > /root/.kube/config".into(), + "chmod 400 /root/.kube/config".into(), + ]; + if registry.provider == HelmRegistryProvider::Ecr { + script_lines.push( + "aws ecr get-login-password --region $HELM_REGISTRY_REGION | helm registry login --username AWS --password-stdin $HELM_REGISTRY".into(), + ); + } + script_lines.extend([ + "helm pull oci://$HELM_REGISTRY/$HELM_REPO --version $HELM_CHART_TAG".into(), + "echo $VALUES_BASE64 | base64 -d > values.yaml".into(), + "echo $VALUES_OVERRIDE_BASE64 | base64 -d > values-override.yaml".into(), + format!( "helm --debug --kubeconfig=/root/.kube/config {command} {namespace_name} oci://$HELM_REGISTRY/$HELM_REPO --version $HELM_CHART_TAG --namespace={namespace_name} -f values.yaml -f values-override.yaml", ), - ].join(" && "); + ]); + let script = script_lines.join(" && "); + + let mut env_vars = vec![ + EnvVar { + name: "KUBECONFIG_BASE64".into(), + value: Some(kubeconfig), + ..Default::default() + }, + EnvVar { + name: "HELM_REGISTRY".into(), + value: Some(registry.domain_name.clone()), + ..Default::default() + }, + EnvVar { + name: "HELM_REPO".into(), + value: Some(registry.repo_name.clone()), + ..Default::default() + }, + EnvVar { + name: "HELM_CHART_TAG".into(), + value: Some(chart.image_tag), + ..Default::default() + }, + EnvVar { + name: "VALUES_BASE64".into(), + value: Some(BASE64_STANDARD.encode(serde_yaml::to_string(&values)?)), + ..Default::default() + }, + EnvVar { + name: "VALUES_OVERRIDE_BASE64".into(), + value: if let Some(values_override) = &deployment.values_override { + Some(BASE64_STANDARD.encode(serde_yaml::to_string(values_override)?)) + } else { + None + }, + ..Default::default() + }, + ]; + if let Some(region_name) = registry.region_name()? { + env_vars.push(EnvVar { + name: "HELM_REGISTRY_REGION".into(), + value: Some(region_name), + ..Default::default() + }); + } Ok(Pod { metadata: ObjectMeta { @@ -80,47 +133,7 @@ async fn helm_pod( image: Some(config.helm_image.to_owned()), image_pull_policy: Some("Always".into()), command: Some(vec!["/bin/bash".into(), "-cex".into(), script]), - env: Some(vec![ - EnvVar { - name: "KUBECONFIG_BASE64".into(), - value: Some(kubeconfig), - ..Default::default() - }, - EnvVar { - name: "HELM_REGISTRY_REGION".into(), - value: Some(registry.region_name()?), - ..Default::default() - }, - EnvVar { - name: "HELM_REGISTRY".into(), - value: Some(registry.domain_name), - ..Default::default() - }, - EnvVar { - name: "HELM_REPO".into(), - value: Some(registry.repo_name), - ..Default::default() - }, - EnvVar { - name: "HELM_CHART_TAG".into(), - value: Some(chart.image_tag), - ..Default::default() - }, - EnvVar { - name: "VALUES_BASE64".into(), - value: Some(BASE64_STANDARD.encode(serde_yaml::to_string(&values)?)), - ..Default::default() - }, - EnvVar { - name: "VALUES_OVERRIDE_BASE64".into(), - value: if let Some(values_override) = &deployment.values_override { - Some(BASE64_STANDARD.encode(serde_yaml::to_string(values_override)?)) - } else { - None - }, - ..Default::default() - }, - ]), + env: Some(env_vars), ..Default::default() }], restart_policy: Some("Never".into()), diff --git a/scripts/container-build.sh b/scripts/container-build.sh deleted file mode 100755 index 1a0e6f7..0000000 --- a/scripts/container-build.sh +++ /dev/null @@ -1,21 +0,0 @@ -#!/usr/bin/env bash - -set -euo pipefail - -RELEASE_BUILD="$1" -BUILD_DEST="$2" - -export CARGO_BUILD_TARGET="$(arch)-unknown-linux-musl" -rustup target add "${CARGO_BUILD_TARGET}" - -if [ "${RELEASE_BUILD}" = "1" ] -then - CARGO_FLAGS="--release" - CARGO_TARGET_DIR="target/${CARGO_BUILD_TARGET}/release/" -else - CARGO_FLAGS="" - CARGO_TARGET_DIR="target/${CARGO_BUILD_TARGET}/debug/" -fi - -cargo build ${CARGO_FLAGS} -find "${CARGO_TARGET_DIR}" -maxdepth 1 -type f -executable -exec cp -v {} "${BUILD_DEST}/" \; diff --git a/scripts/dex.config.yaml b/scripts/dex.config.yaml deleted file mode 100644 index 6359960..0000000 --- a/scripts/dex.config.yaml +++ /dev/null @@ -1,37 +0,0 @@ -issuer: http://127.0.0.1:5556 - -storage: - type: sqlite3 - config: - file: /tmp/dex.db - -web: - http: 0.0.0.0:5556 - -telemetry: - http: 0.0.0.0:5558 - -grpc: - addr: 0.0.0.0:5557 - -staticClients: - - id: platz - redirectURIs: - - "http://127.0.0.1:5173/auth/google/callback" - - "https://127.0.0.1:5173/auth/google/callback" - name: "Platz" - secret: platz - -connectors: - - type: mockCallback - id: mock - name: Example - -enablePasswordDB: true - -staticPasswords: - - email: "admin@example.com" - # bcrypt hash of the string "password": $(echo password | htpasswd -BinC 10 admin | cut -d: -f2) - hash: "$2a$10$2b2cU8CPhOTaGrs1HRQuAueS7JTT5ZHsHSzYiFPm1leZck7Mc8T4W" - username: "admin" - userID: "08a8684b-db88-4b73-90a9-3cd1661f5466" diff --git a/scripts/run-api.sh b/scripts/run-api.sh deleted file mode 100755 index d4b6991..0000000 --- a/scripts/run-api.sh +++ /dev/null @@ -1,34 +0,0 @@ -#!/usr/bin/env bash - -set -eu - -HERE="$(dirname "$0")" -SCRIPT="$(basename "$0")" -export API_HOST="127.0.0.1" -export PLATZ_FRONTEND_PORT="5173" - -if [ -z "${DATABASE_URL:-}" ] -then - source "${HERE}/run-db.sh" -fi - -if [ -z "${OIDC_SERVER_URL:-}" ] -then - source "${HERE}/run-oidc.sh" -fi - -if ! which cargo-watch &>/dev/null -then - echo "[${SCRIPT}] ๐Ÿฆ€ Installing cargo watch" - cargo install cargo-watch -fi - -export RUST_LOG="debug" -export RUST_BACKTRACE="1" -export PLATZ_OWN_URL="http://127.0.0.1:${PLATZ_FRONTEND_PORT}" -# From oidc-users.json -export ADMIN_EMAILS="admin@example.com" - -echo "[${SCRIPT}] ๐Ÿš€ Running API server" -args=("$@") -exec cargo watch -x "run --bin=platz-api -- run ${args[*]}" diff --git a/scripts/run-db.sh b/scripts/run-db.sh deleted file mode 100755 index 109e129..0000000 --- a/scripts/run-db.sh +++ /dev/null @@ -1,48 +0,0 @@ -#!/bin/bash - -set -eu - -LOG_PREFIX="[$(basename $0)]" - -export PGHOST="localhost" -export PGUSER="postgres" -export PGPORT="5432" -export PGDATABASE="test" -export PGPASSWORD="postgres" - -DB_CONTAINER="${PGDATABASE}-postgres" - -if ps -ef | grep kubectl | grep port-forward | grep "${PGPORT}" &> /dev/null -then - echo "${LOG_PREFIX} โ›”๏ธ kubectl port-forward is running" >&2 - echo "${LOG_PREFIX}" >&2 - echo "${LOG_PREFIX} An active kubectl port-forward may potentially point to" >&2 - echo "${LOG_PREFIX} a production database." >&2 - echo "${LOG_PREFIX} Please turn it off first, then re-run this script." >&2 - exit 2 -fi - -if [ ! -z "${RM_DB:-}" ] -then - echo "${LOG_PREFIX} ๐Ÿช“ Deleting existing database container" >&2 - docker stop "${DB_CONTAINER}" >&2 || true - docker rm "${DB_CONTAINER}" >&2 || true -fi - -existing=$(docker ps -aq --no-trunc --filter name="${DB_CONTAINER}") -if [ -z "${existing}" ] -then - echo "${LOG_PREFIX} โœจ Creating database container" >&2 - docker run -d \ - --name "${DB_CONTAINER}" \ - -p "${PGPORT}":"${PGPORT}" \ - -e POSTGRES_PASSWORD="${PGPASSWORD}" \ - -e POSTGRES_DB="${PGDATABASE}" \ - postgres >&2 -fi - -if [ "$( docker container inspect -f '{{.State.Status}}' ${DB_CONTAINER} )" != "running" ] -then - echo "${LOG_PREFIX} ๐Ÿ Starting database container" >&2 - docker start "${DB_CONTAINER}" >&2 -fi diff --git a/scripts/run-oidc.sh b/scripts/run-oidc.sh deleted file mode 100755 index 0f249d2..0000000 --- a/scripts/run-oidc.sh +++ /dev/null @@ -1,27 +0,0 @@ -#!/usr/bin/env bash - -set -eu - -HERE="$(realpath $(dirname "$0"))" -LOG_PREFIX="[$(basename $0)]" - -OIDC_CONFIG_FILE="${HERE}/dex.config.yaml" -export OIDC_PORT="5556" -export OIDC_SERVER_URL="http://127.0.0.1:${OIDC_PORT}" -export OIDC_CLIENT_ID="$(yq '.staticClients[0].id' "${OIDC_CONFIG_FILE}")" -export OIDC_CLIENT_SECRET="$(yq '.staticClients[0].secret' "${OIDC_CONFIG_FILE}")" -OIDC_CONTAINER="oidc-provider" - -existing=$(docker ps -aq --no-trunc --filter name="${OIDC_CONTAINER}") -if [ ! -z "${existing}" ] -then - docker stop "${OIDC_CONTAINER}" || true - docker rm "${OIDC_CONTAINER}" -fi - -echo "${LOG_PREFIX} โœจ Creating OIDC server" >&2 -docker run -d \ - --name "${OIDC_CONTAINER}" \ - -p "${OIDC_PORT}":"${OIDC_PORT}" \ - -v "${OIDC_CONFIG_FILE}:/etc/dex/config.docker.yaml:ro" \ - dexidp/dex >&2