diff --git a/Makefile b/Makefile index 29f7c21e..91df7592 100644 --- a/Makefile +++ b/Makefile @@ -150,9 +150,9 @@ $(BIN_PATH)/wasm build/bin/wasm: FORCE # test targets # -test-gopkgs: go-generate ginkgo-tests test-ulimits test-rdt test-hook-injector test-writable-cgroups +test-gopkgs: go-generate ginkgo-tests test-ulimits test-rdt test-hook-injector test-writable-cgroups test-identity-injector -SKIPPED_PKGS="ulimit-adjuster,device-injector,rdt,hook-injector,writable-cgroups" +SKIPPED_PKGS="ulimit-adjuster,device-injector,rdt,hook-injector,writable-cgroups,identity-injector" ginkgo-tests: $(Q)$(GINKGO) run \ @@ -183,6 +183,9 @@ test-hook-injector: test-writable-cgroups: $(Q)cd ./plugins/writable-cgroups && $(GO_TEST) -v +test-identity-injector: + $(Q)cd ./plugins/identity-injector && $(GO_TEST) -v + codecov: SHELL := $(shell which bash) codecov: bash <(curl -s https://codecov.io/bash) -f $(COVERAGE_PATH)/coverprofile diff --git a/contrib/kustomize/identity-injector/base/daemonset.yaml b/contrib/kustomize/identity-injector/base/daemonset.yaml new file mode 100644 index 00000000..94cc62d4 --- /dev/null +++ b/contrib/kustomize/identity-injector/base/daemonset.yaml @@ -0,0 +1,48 @@ +apiVersion: apps/v1 +kind: DaemonSet +metadata: + name: nri-plugin-identity +spec: + template: + spec: + containers: + - name: nri-plugin-identity + image: plugin:latest + args: + - "--idx" + - "11" + - "--verbose" + - "true" + - "--spire-admin-socket" + - "unix:///run/spire/admin-socket/admin.sock" + - "--host-mount-path" + - "/var/run/spiffe/secrets/" + resources: + requests: + cpu: "2m" + memory: "5Mi" + securityContext: + allowPrivilegeEscalation: false + capabilities: + drop: + - ALL + volumeMounts: + - name: nri-socket + mountPath: /var/run/nri/nri.sock + - name: spire-admin-socket + mountPath: /run/spire/admin-socket/admin.sock + - name: spiffe-certs + mountPath: /var/run/spiffe/secrets/ + volumes: + - name: nri-socket + hostPath: + path: /var/run/nri/nri.sock + type: Socket + - name: spire-admin-socket + hostPath: + path: /run/spire/admin-socket/admin.sock + type: Socket + - name: spiffe-certs + hostPath: + path: /var/run/spiffe/secrets/ + type: DirectoryOrCreate diff --git a/contrib/kustomize/identity-injector/base/kustomization.yaml b/contrib/kustomize/identity-injector/base/kustomization.yaml new file mode 100644 index 00000000..dede43b3 --- /dev/null +++ b/contrib/kustomize/identity-injector/base/kustomization.yaml @@ -0,0 +1,13 @@ +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization +namespace: kube-system +resources: + - daemonset.yaml +images: + - name: plugin + newName: localhost:5000/nri-identity-injector + newTag: latest +labels: + - includeSelectors: true + pairs: + app.kubernetes.io/name: nri-plugin-identity diff --git a/contrib/kustomize/identity-injector/kustomization.yaml b/contrib/kustomize/identity-injector/kustomization.yaml new file mode 100644 index 00000000..c9646239 --- /dev/null +++ b/contrib/kustomize/identity-injector/kustomization.yaml @@ -0,0 +1,4 @@ +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization +resources: + - base/ diff --git a/plugins/identity-injector/go.mod b/plugins/identity-injector/go.mod new file mode 100644 index 00000000..8f78e8dc --- /dev/null +++ b/plugins/identity-injector/go.mod @@ -0,0 +1,37 @@ +module github.com/containerd/nri/plugins/identity-injector + +go 1.24.1 + +replace github.com/containerd/nri => ../.. + +require ( + github.com/containerd/nri v0.11.0 + github.com/sirupsen/logrus v1.9.4 + github.com/spiffe/go-spiffe/v2 v2.6.0 + github.com/stretchr/testify v1.11.1 +) + +require ( + github.com/davecgh/go-spew v1.1.1 // indirect + github.com/kr/text v0.2.0 // indirect + github.com/pmezard/go-difflib v1.0.0 // indirect + go.yaml.in/yaml/v2 v2.4.2 // indirect + golang.org/x/mod v0.32.0 // indirect + gopkg.in/yaml.v3 v3.0.1 // indirect +) + +require ( + github.com/containerd/log v0.1.0 // indirect + github.com/containerd/ttrpc v1.2.7 // indirect + github.com/knqyf263/go-plugin v0.9.0 // indirect + github.com/opencontainers/runtime-spec v1.3.0 // indirect + github.com/spiffe/spire-api-sdk v1.14.1 + github.com/tetratelabs/wazero v1.11.0 // indirect + golang.org/x/net v0.49.0 // indirect + golang.org/x/sys v0.40.0 // indirect + golang.org/x/text v0.33.0 // indirect + google.golang.org/genproto/googleapis/rpc v0.0.0-20250811230008-5f3141c8851a // indirect + google.golang.org/grpc v1.75.0 + google.golang.org/protobuf v1.36.7 // indirect + sigs.k8s.io/yaml v1.6.0 +) diff --git a/plugins/identity-injector/go.sum b/plugins/identity-injector/go.sum new file mode 100644 index 00000000..d1aee4f4 --- /dev/null +++ b/plugins/identity-injector/go.sum @@ -0,0 +1,94 @@ +github.com/Masterminds/semver/v3 v3.4.0 h1:Zog+i5UMtVoCU8oKka5P7i9q9HgrJeGzI9SA1Xbatp0= +github.com/Masterminds/semver/v3 v3.4.0/go.mod h1:4V+yj/TJE1HU9XfppCwVMZq3I84lprf4nC11bSS5beM= +github.com/brianvoe/gofakeit/v7 v7.12.1 h1:df1tiI4SL1dR5Ix4D/r6a3a+nXBJ/OBGU5jEKRBmmqg= +github.com/brianvoe/gofakeit/v7 v7.12.1/go.mod h1:QXuPeBw164PJCzCUZVmgpgHJ3Llj49jSLVkKPMtxtxA= +github.com/containerd/log v0.1.0 h1:TCJt7ioM2cr/tfR8GPbGf9/VRAX8D2B4PjzCpfX540I= +github.com/containerd/log v0.1.0/go.mod h1:VRRf09a7mHDIRezVKTRCrOq78v577GXq3bSa3EhrzVo= +github.com/containerd/ttrpc v1.2.7 h1:qIrroQvuOL9HQ1X6KHe2ohc7p+HP/0VE6XPU7elJRqQ= +github.com/containerd/ttrpc v1.2.7/go.mod h1:YCXHsb32f+Sq5/72xHubdiJRQY9inL4a4ZQrAbN1q9o= +github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= +github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= +github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= +github.com/go-logr/stdr v1.2.2/go.mod h1:mMo/vtBO5dYbehREoey6XUKy/eSumjCCveDpRre4VKE= +github.com/go-task/slim-sprig/v3 v3.0.0 h1:sUs3vkvUymDpBKi3qH1YSqBQk9+9D/8M2mN1vB6EwHI= +github.com/go-task/slim-sprig/v3 v3.0.0/go.mod h1:W848ghGpv3Qj3dhTPRyJypKRiqCdHZiAzKg9hl15HA8= +github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek= +github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= +github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= +github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/pprof v0.0.0-20260115054156-294ebfa9ad83 h1:z2ogiKUYzX5Is6zr/vP9vJGqPwcdqsWjOt+V8J7+bTc= +github.com/google/pprof v0.0.0-20260115054156-294ebfa9ad83/go.mod h1:MxpfABSjhmINe3F1It9d+8exIHFvUqtLIRCdOGNXqiI= +github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= +github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/knqyf263/go-plugin v0.9.0 h1:CQs2+lOPIlkZVtcb835ZYDEoyyWJWLbSTWeCs0EwTwI= +github.com/knqyf263/go-plugin v0.9.0/go.mod h1:2z5lCO1/pez6qGo8CvCxSlBFSEat4MEp1DrnA+f7w8Q= +github.com/kr/pretty v0.1.0 h1:L/CwN0zerZDmRFUapSPitk6f+Q3+0za1rQkzVuMiMFI= +github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo= +github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= +github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= +github.com/onsi/ginkgo/v2 v2.28.1 h1:S4hj+HbZp40fNKuLUQOYLDgZLwNUVn19N3Atb98NCyI= +github.com/onsi/ginkgo/v2 v2.28.1/go.mod h1:CLtbVInNckU3/+gC8LzkGUb9oF+e8W8TdUsxPwvdOgE= +github.com/onsi/gomega v1.39.1 h1:1IJLAad4zjPn2PsnhH70V4DKRFlrCzGBNrNaru+Vf28= +github.com/onsi/gomega v1.39.1/go.mod h1:hL6yVALoTOxeWudERyfppUcZXjMwIMLnuSfruD2lcfg= +github.com/opencontainers/runtime-spec v1.3.0 h1:YZupQUdctfhpZy3TM39nN9Ika5CBWT5diQ8ibYCRkxg= +github.com/opencontainers/runtime-spec v1.3.0/go.mod h1:jwyrGlmzljRJv/Fgzds9SsS/C5hL+LL3ko9hs6T5lQ0= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/prometheus/procfs v0.6.0 h1:mxy4L2jP6qMonqmq+aTtOx1ifVWUgG/TAmntgbh3xv4= +github.com/prometheus/procfs v0.6.0/go.mod h1:cz+aTbrPOrUb4q7XlbU9ygM+/jj0fzG6c1xBZuNvfVA= +github.com/sirupsen/logrus v1.9.4 h1:TsZE7l11zFCLZnZ+teH4Umoq5BhEIfIzfRDZ1Uzql2w= +github.com/sirupsen/logrus v1.9.4/go.mod h1:ftWc9WdOfJ0a92nsE2jF5u5ZwH8Bv2zdeOC42RjbV2g= +github.com/spiffe/go-spiffe/v2 v2.6.0 h1:l+DolpxNWYgruGQVV0xsfeya3CsC7m8iBzDnMpsbLuo= +github.com/spiffe/go-spiffe/v2 v2.6.0/go.mod h1:gm2SeUoMZEtpnzPNs2Csc0D/gX33k1xIx7lEzqblHEs= +github.com/spiffe/spire-api-sdk v1.14.1 h1:pAVoElOJseU+L4ecrPFv8krIDG8XYWBYKztGUG5qWfg= +github.com/spiffe/spire-api-sdk v1.14.1/go.mod h1:9hXJcMzatM1KwAtBDO3s6HccDCic++/5c2yOc5Iln8Y= +github.com/stretchr/testify v1.11.1 h1:7s2iGBzp5EwR7/aIZr8ao5+dra3wiQyKjjFuvgVKu7U= +github.com/stretchr/testify v1.11.1/go.mod h1:wZwfW3scLgRK+23gO65QZefKpKQRnfz6sD981Nm4B6U= +github.com/tetratelabs/wazero v1.11.0 h1:+gKemEuKCTevU4d7ZTzlsvgd1uaToIDtlQlmNbwqYhA= +github.com/tetratelabs/wazero v1.11.0/go.mod h1:eV28rsN8Q+xwjogd7f4/Pp4xFxO7uOGbLcD/LzB1wiU= +go.opentelemetry.io/auto/sdk v1.1.0 h1:cH53jehLUN6UFLY71z+NDOiNJqDdPRaXzTel0sJySYA= +go.opentelemetry.io/auto/sdk v1.1.0/go.mod h1:3wSPjt5PWp2RhlCcmmOial7AvC4DQqZb7a7wCow3W8A= +go.opentelemetry.io/otel v1.37.0 h1:9zhNfelUvx0KBfu/gb+ZgeAfAgtWrfHJZcAqFC228wQ= +go.opentelemetry.io/otel v1.37.0/go.mod h1:ehE/umFRLnuLa/vSccNq9oS1ErUlkkK71gMcN34UG8I= +go.opentelemetry.io/otel/metric v1.37.0 h1:mvwbQS5m0tbmqML4NqK+e3aDiO02vsf/WgbsdpcPoZE= +go.opentelemetry.io/otel/metric v1.37.0/go.mod h1:04wGrZurHYKOc+RKeye86GwKiTb9FKm1WHtO+4EVr2E= +go.opentelemetry.io/otel/sdk v1.37.0 h1:ItB0QUqnjesGRvNcmAcU0LyvkVyGJ2xftD29bWdDvKI= +go.opentelemetry.io/otel/sdk v1.37.0/go.mod h1:VredYzxUvuo2q3WRcDnKDjbdvmO0sCzOvVAiY+yUkAg= +go.opentelemetry.io/otel/sdk/metric v1.37.0 h1:90lI228XrB9jCMuSdA0673aubgRobVZFhbjxHHspCPc= +go.opentelemetry.io/otel/sdk/metric v1.37.0/go.mod h1:cNen4ZWfiD37l5NhS+Keb5RXVWZWpRE+9WyVCpbo5ps= +go.opentelemetry.io/otel/trace v1.37.0 h1:HLdcFNbRQBE2imdSEgm/kwqmQj1Or1l/7bW6mxVK7z4= +go.opentelemetry.io/otel/trace v1.37.0/go.mod h1:TlgrlQ+PtQO5XFerSPUYG0JSgGyryXewPGyayAWSBS0= +go.yaml.in/yaml/v2 v2.4.2 h1:DzmwEr2rDGHl7lsFgAHxmNz/1NlQ7xLIrlN2h5d1eGI= +go.yaml.in/yaml/v2 v2.4.2/go.mod h1:081UH+NErpNdqlCXm3TtEran0rJZGxAYx9hb/ELlsPU= +go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc= +go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg= +golang.org/x/mod v0.32.0 h1:9F4d3PHLljb6x//jOyokMv3eX+YDeepZSEo3mFJy93c= +golang.org/x/mod v0.32.0/go.mod h1:SgipZ/3h2Ci89DlEtEXWUk/HteuRin+HHhN+WbNhguU= +golang.org/x/net v0.49.0 h1:eeHFmOGUTtaaPSGNmjBKpbng9MulQsJURQUAfUwY++o= +golang.org/x/net v0.49.0/go.mod h1:/ysNB2EvaqvesRkuLAyjI1ycPZlQHM3q01F02UY/MV8= +golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4= +golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/sys v0.40.0 h1:DBZZqJ2Rkml6QMQsZywtnjnnGvHza6BTfYFWY9kjEWQ= +golang.org/x/sys v0.40.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/text v0.33.0 h1:B3njUFyqtHDUI5jMn1YIr5B0IE2U0qck04r6d4KPAxE= +golang.org/x/text v0.33.0/go.mod h1:LuMebE6+rBincTi9+xWTY8TztLzKHc/9C1uBCG27+q8= +golang.org/x/tools v0.41.0 h1:a9b8iMweWG+S0OBnlU36rzLp20z1Rp10w+IY2czHTQc= +golang.org/x/tools v0.41.0/go.mod h1:XSY6eDqxVNiYgezAVqqCeihT4j1U2CCsqvH3WhQpnlg= +gonum.org/v1/gonum v0.16.0 h1:5+ul4Swaf3ESvrOnidPp4GZbzf0mxVQpDCYUQE7OJfk= +gonum.org/v1/gonum v0.16.0/go.mod h1:fef3am4MQ93R2HHpKnLk4/Tbh/s0+wqD5nfa6Pnwy4E= +google.golang.org/genproto/googleapis/rpc v0.0.0-20250811230008-5f3141c8851a h1:tPE/Kp+x9dMSwUm/uM0JKK0IfdiJkwAbSMSeZBXXJXc= +google.golang.org/genproto/googleapis/rpc v0.0.0-20250811230008-5f3141c8851a/go.mod h1:gw1tLEfykwDz2ET4a12jcXt4couGAm7IwsVaTy0Sflo= +google.golang.org/grpc v1.75.0 h1:+TW+dqTd2Biwe6KKfhE5JpiYIBWq865PhKGSXiivqt4= +google.golang.org/grpc v1.75.0/go.mod h1:JtPAzKiq4v1xcAB2hydNlWI2RnF85XXcV0mhKXr2ecQ= +google.golang.org/protobuf v1.36.7 h1:IgrO7UwFQGJdRNXH/sQux4R1Dj1WAKcLElzeeRaXV2A= +google.golang.org/protobuf v1.36.7/go.mod h1:jduwjTPXsFjZGTmRluh+L6NjiWu7pchiJ2/5YcXBHnY= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c h1:Hei/4ADfdWqJk1ZMxUNpqntNwaWcugrBjAiHlqqRiVk= +gopkg.in/check.v1 v1.0.0-20201130134442-10cb98267c6c/go.mod h1:JHkPIbrfpd72SG/EVd6muEfDQjcINNoR0C8j2r3qZ4Q= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +sigs.k8s.io/yaml v1.6.0 h1:G8fkbMSAFqgEFgh4b1wmtzDnioxFCUgTZhlbj5P9QYs= +sigs.k8s.io/yaml v1.6.0/go.mod h1:796bPqUfzR/0jLAl6XjHl3Ck7MiyVv8dbTdyT3/pMf4= diff --git a/plugins/identity-injector/identity-injector.go b/plugins/identity-injector/identity-injector.go new file mode 100644 index 00000000..ed5b6e38 --- /dev/null +++ b/plugins/identity-injector/identity-injector.go @@ -0,0 +1,688 @@ +/* + Copyright The containerd Authors. + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +*/ + +package main + +import ( + "context" + "crypto/x509" + "encoding/pem" + "flag" + "fmt" + "os" + "path" + "path/filepath" + "sync" + "strings" + + "github.com/sirupsen/logrus" + "sigs.k8s.io/yaml" + + "github.com/spiffe/go-spiffe/v2/spiffeid" + delegatedidentityv1 "github.com/spiffe/spire-api-sdk/proto/spire/api/agent/delegatedidentity/v1" + + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + + "github.com/containerd/nri/pkg/api" + "github.com/containerd/nri/pkg/stub" +) + +const ( + // identityKey is the prefix of the key used for identity annotations in the podspec. + identityKey = "identity.noderesource.dev" + + // Default paths for certificate files in the container + defaultCertFileName = "svid.pem" + defaultKeyFileName = "key.pem" + defaultBundleFileName = "bundle.pem" + + // Path in the container where the identity artifacts will be stored + defaultMountPath = "/var/run/spiffe/secrets" +) + +var ( + log *logrus.Logger + verbose bool + spireAdminSocket string + hostMountPath string +) + +type plugin struct { + stub stub.Stub + delegatedIdentityConn *grpc.ClientConn + delegatedIdentityClient delegatedidentityv1.DelegatedIdentityClient + watchers map[string]*containerWatcher // key: pod-uid/container-name + watchersMu sync.RWMutex +} + +// identityConfig represents the configuration for identity injection +type identityConfig struct { + MountPath string `json:"mount_path,omitempty"` + CertFileName string `json:"cert_file_name,omitempty"` + KeyFileName string `json:"key_file_name,omitempty"` + BundleFileName string `json:"bundle_file_name,omitempty"` + SpiffeId string `json:"spiffe_id,omitempty"` +} + +// containerWatcher tracks a running certificate watcher for a container +type containerWatcher struct { + cancel context.CancelFunc + done chan struct{} +} + +func (p *plugin) CreateContainer(_ context.Context, pod *api.PodSandbox, container *api.Container) (*api.ContainerAdjustment, []*api.ContainerUpdate, error) { + if verbose { + dump("CreateContainer", "pod", pod, "container", container) + } + + adjust := &api.ContainerAdjustment{} + + config, err := parseIdentityConfig(container.Name, pod.Annotations) + if err != nil { + return nil, nil, err + } + + if config == nil { + return nil, nil, nil + } + + // Create host directory for certificates (will be populated in StartContainer) + hostDir := getHostDir(hostMountPath, pod.GetUid(), container.Name) + if err := os.MkdirAll(hostDir, 0755); err != nil { + return nil, nil, fmt.Errorf("failed to create host directory %s: %w", hostDir, err) + } + + // Add mount for the certificate directory + mount := &api.Mount{ + Source: hostDir, + Destination: config.MountPath, + Type: "bind", + Options: []string{"ro", "bind"}, + } + adjust.AddMount(mount) + + return adjust, nil, nil +} + +func (p *plugin) StartContainer(ctx context.Context, pod *api.PodSandbox, container *api.Container) error { + if verbose { + dump("StartContainer", "pod", pod, "container", container) + } + + return p.injectIdentity(ctx, pod, container) +} + +// Supporting UpdateContainer() is not required because the container pid does not change when a container is updated. +// Supporting StopContainer() addresses both cases - when the container is intentionally stopped under normal operations (graceful exit?) and when the container is stopped if a container crashes +// Remove the watcher and cleanup the certificates for that container on container stop. +func (p *plugin) StopContainer(ctx context.Context, pod *api.PodSandbox, container *api.Container) ([]*api.ContainerUpdate, error) { + if verbose { + dump("RemoveContainer", "pod", pod, "container", container) + } + + watcherKey := filepath.Join(pod.GetUid(), container.Name) + + // Stop the certificate watcher if it exists + p.watchersMu.Lock() + if p.watchers == nil { + p.watchersMu.Unlock() + log.Debugf("%s: no watchers map initialized", containerName(pod, container)) + } else if watcher, exists := p.watchers[watcherKey]; exists { + log.Infof("%s: stopping certificate watcher", containerName(pod, container)) + watcher.cancel() + delete(p.watchers, watcherKey) + p.watchersMu.Unlock() + + // Wait for watcher to finish (with timeout) + select { + case <-watcher.done: + log.Infof("%s: certificate watcher stopped gracefully", containerName(pod, container)) + case <-ctx.Done(): + log.Warnf("%s: timeout waiting for certificate watcher to stop", containerName(pod, container)) + } + } else { + p.watchersMu.Unlock() + } + + // Clean up certificate files + + hostDir := getHostDir(hostMountPath, pod.GetUid(), container.Name) + if err := os.RemoveAll(hostDir); err != nil { + log.Warnf("%s: failed to clean up certificate directory %s: %v", containerName(pod, container), hostDir, err) + } else { + log.Infof("%s: cleaned up certificate directory %s", containerName(pod, container), hostDir) + } + + return nil, nil +} + +func (p *plugin) Shutdown(ctx context.Context) { + log.Infof("Shutdown called, stopping all certificate watchers") + + // Stop all active watchers + p.watchersMu.Lock() + if p.watchers == nil { + p.watchersMu.Unlock() + log.Warnf("no watchers to stop") + return + } + + // Cancel all watchers + for _, watcher := range p.watchers { + watcher.cancel() + } + + // Wait for all watchers to finish (with timeout) + for _, watcher := range p.watchers { + select { + case <-watcher.done: + case <-ctx.Done(): + log.Warnf("timeout waiting for watcher to stop during shutdown") + } + } + + // Held the lock for the whole shutdown process. + p.watchersMu.Unlock() + + log.Infof("all certificate watchers stopped") +} + + +func (p *plugin) injectIdentity(ctx context.Context, pod *api.PodSandbox, container *api.Container) error { + // Check container PID + if container.Pid == 0 { + return fmt.Errorf("%s container PID not available", containerName(pod, container)) + } + + config, err := parseIdentityConfig(container.Name, pod.Annotations) + if err != nil { + return err + } + + if config == nil { + return nil + } + + if verbose { + dump(containerName(pod, container), "identity config", config) + } + + // Get host directory for certificates + hostDir := getHostDir(hostMountPath, pod.GetUid(), container.Name) + + // Start watching for certificate updates using streaming API + // This will automatically receive new certificates when they're rotated + if err := p.startCertificateWatcher(ctx, pod, container, int32(container.Pid), hostDir, config); err != nil { + return fmt.Errorf("failed to start certificate watcher: %w", err) + } + + return nil +} + +// startCertificateWatcher starts a goroutine that watches for certificate updates +// using the SPIRE Delegated Identity API streaming interface (SubscribeToX509SVIDs). +// This automatically receives new certificates when they're rotated by SPIRE. +func (p *plugin) startCertificateWatcher(ctx context.Context, pod *api.PodSandbox, ctr *api.Container, pid int32, hostDir string, config *identityConfig) error { + watcherKey := filepath.Join(pod.GetUid(), ctr.Name) + + // Check if watcher already exists + p.watchersMu.RLock() + if _, exists := p.watchers[watcherKey]; exists { + p.watchersMu.RUnlock() + log.Debugf("%s: certificate watcher already running", containerName(pod, ctr)) + return nil + } + p.watchersMu.RUnlock() + + // Create cancellable context for this watcher + watcherCtx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + + // Use WaitGroup to coordinate both goroutines + var wg sync.WaitGroup + wg.Add(2) // We have 2 goroutines: SVID watcher and bundle watcher + + // Store watcher info + p.watchersMu.Lock() + p.watchers[watcherKey] = &containerWatcher{ + cancel: cancel, + done: done, + } + p.watchersMu.Unlock() + + // Goroutine to wait for both watchers to finish and then close done channel + go func() { + wg.Wait() + close(done) + p.watchersMu.Lock() + delete(p.watchers, watcherKey) + p.watchersMu.Unlock() + log.Infof("%s: watcher stopped", containerName(pod, ctr)) + }() + + // Start the SVID watcher goroutine + go func() { + defer wg.Done() + defer cancel() // Cancel context on exit to stop the other goroutine + defer log.Infof("%s: SVID watcher stopped", containerName(pod, ctr)) + + req := &delegatedidentityv1.SubscribeToX509SVIDsRequest{ + Pid: pid, + } + + // Use SPIRE Delegated Identity API to subscribe to certificate updates + // This automatically handles certificate rotation - when SPIRE rotates certs, + // new certificates are pushed to the stream + stream, err := p.delegatedIdentityClient.SubscribeToX509SVIDs(watcherCtx, req) + + if err != nil { + // Error here means that there is no way to fetch certificates for the requested pid + // This can happen for example when there is no workload registered with the Spire Server + // associated with that pid. Therefore no retry logic needed here. + // Plugin will not stop the container. This is left to the application to deal with this error. + + // Just log the error. Returning the error not possible because the goroutine here does not return anything + log.Errorf("%s: failed to subscribe to X509 SVIDs: %v", containerName(pod, ctr), err) + return + } + + // Process streaming updates + // This loop also takes care of retrying to fetch certificates + for { + if err := watcherCtx.Err(); err != nil { + log.Errorf("%s: watcher context cancelled: %v", containerName(pod, ctr), err) + return + } + + resp, err := stream.Recv() + if err != nil { + log.Errorf("%s: bundle stream error: %v", containerName(pod, ctr), err) + return + } + + // Process the certificate update + if err := p.processSvidUpdate(containerName(pod, ctr), pid, hostDir, config, resp.X509Svids); err != nil { + log.Errorf("%s: failed to process certificate update: %v", containerName(pod, ctr), err) // Just log the error. Returning the error not possible because the goroutine here does not return anything + + // this error is a fatal error which should stop further execution/processing + return + } + + } + }() + + // Start the bundle watcher goroutine + go func() { + defer wg.Done() + defer cancel() // Cancel context on exit to stop the other goroutine + defer log.Infof("%s: bundle watcher stopped", containerName(pod, ctr)) + + log.Infof("%s: starting bundle stream", containerName(pod, ctr)) + + req := &delegatedidentityv1.SubscribeToX509BundlesRequest{} + + log.Infof("%s: requesting bundle stream using delegated identity API", containerName(pod, ctr)) + + // Use SPIRE Delegated Identity API to subscribe to bundle updates + // This automatically handles bundle rotation - when SPIRE rotates bundles, + // new bundles are pushed to the stream + stream, err := p.delegatedIdentityClient.SubscribeToX509Bundles(watcherCtx, req) + + if err != nil { + // Error here means that there is no way to fetch certificate bundles. + // Therefore no retry logic needed here. + // Plugin will not stop the container. This is left to the application to deal with this error. + + // Just log the error. Returning the error not possible because the goroutine here does not return anything + log.Errorf("%s: failed to subscribe to X509 Bundles: %v", containerName(pod, ctr), err) + return + } + + // Process streaming updates + for { + if err := watcherCtx.Err(); err != nil { + log.Errorf("%s: watcher context cancelled: %v", containerName(pod, ctr), err) + return + } + + resp, err := stream.Recv() + if err != nil { + log.Errorf("%s: bundle stream error: %v", containerName(pod, ctr), err) + return + } + + // Process the bundle update + if err := p.processBundleUpdate(containerName(pod, ctr), hostDir, config, resp.CaCertificates); err != nil { + log.Errorf("%s: failed to process bundle update: %v", containerName(pod, ctr), err) + return + } + + log.Infof("%s: bundle updated", containerName(pod, ctr)) + } + }() + + return nil +} + +// processSvidUpdate processes certificate updates from the delegated identity API +func (p *plugin) processSvidUpdate(containerName string, pid int32, hostDir string, config *identityConfig, x509Svids []*delegatedidentityv1.X509SVIDWithKey) error { + + if len(x509Svids) == 0 { + log.Warnf("%s: received empty SVID update for PID %d", containerName, pid) + + // It could take Spire some milliseconds to mint certificates. + // By not returning error ensures that the for loop in startCertificateWatcher() will also act as retry logic + return nil + } + + // TODO implement using hint to select relevant svid in case response has multiple svids + // Get the default SVID + svidWithKey := x509Svids[0] + + trustDomain, err := spiffeid.TrustDomainFromString(svidWithKey.X509Svid.Id.TrustDomain) + if err != nil { + log.Errorf("%s: failed to parse trust domain: %v", containerName, err) + return err + } + + // Compute SpiffeID and compare to the configured SpiffeID from the podspec + spiffeId, err := spiffeid.FromPath(trustDomain, svidWithKey.X509Svid.Id.Path) + if err != nil { + log.Errorf("%s: failed to parse spiffe id from path: %v", containerName, err) + return err + } + log.Infof("%s: parsed spiffe id %s for PID %d", containerName, spiffeId, pid) + + if config.SpiffeId != "" && config.SpiffeId != spiffeId.String() { + return fmt.Errorf("SpiffeId received from Spire Agent does not match the SpiffeId configured in the Podspec") + } + + + // Parse DER encoded certs + // We loop through each certificate because CertChain gives us certificates + // one at a time (like separate files). We can't use x509.ParseCertificates() + // because that function expects all certificates glued together into one big blob. + certs := make([]*x509.Certificate, 0, len(svidWithKey.X509Svid.CertChain)) + + for _, certDER := range svidWithKey.X509Svid.CertChain { + cert, err := x509.ParseCertificate(certDER) + if err != nil { + log.Errorf("%s: failed to parse certificate: %v", containerName, err) + return err + } + certs = append(certs, cert) + } + + + // Parse the DER-encoded private key + privateKey, err := x509.ParsePKCS8PrivateKey(svidWithKey.X509SvidKey) + if err != nil { + log.Errorf("%s: failed to parse private key for %s: %v", containerName, spiffeId, err) + return err + } + + // Re-marshal to PKCS#8 DER format + // (Even though we already have 'der', re-marshaling is good practice + // if you've modified the key or need to ensure standard formatting) + encodedPrivateKey, err := x509.MarshalPKCS8PrivateKey(privateKey) + if err != nil { + return fmt.Errorf("failed to marshal private key: %v", err) + } + + // Write updated certificates to host filesystem + if err := writeX509Content(hostDir, config, certs, encodedPrivateKey); err != nil { + return fmt.Errorf("failed to write certificates: %w", err) + } + + log.Infof("%s: certificates updated for PID %d (SPIFFE ID: %s, expires: %d)", + containerName, pid, svidWithKey.X509Svid.Id.String(), svidWithKey.X509Svid.ExpiresAt) + + return nil +} + +// processBundleUpdate processes bundle updates from the delegated identity API +func (p *plugin) processBundleUpdate(containerName string, hostDir string, config *identityConfig, caCertificates map[string][]byte) error { + if len(caCertificates) == 0 { + log.Warnf("%s: received empty bundle update", containerName) + + // By not returning error ensures that the for loop in startCertificateWatcher() will also act as retry logic + return nil + } + + log.Infof("%s: received bundle update with %d trust domains", containerName, len(caCertificates)) + + var bundleSet []*x509.Certificate + + // Among all the containers interacting with a Spire Agent through the NRI Identity Plugin, + // if a trust domain is not used by any container/pod interacting with this agent, + // it means that the agent is misconfigured and that trust domain should be removed from the agent. + // Parse all CA certificates from all trust domains + for trustDomain, certDERs := range caCertificates { + if verbose { + log.Infof("%s: processing bundle for trust domain %s", containerName, trustDomain) + } + + bundle, err := x509.ParseCertificates(certDERs) + if err != nil { + log.Errorf("%s: failed to parse CA certificate for trust domain %s: %v", containerName, trustDomain, err) + return err + } + + for _, cert := range bundle { + bundleSet = append(bundleSet, cert) + } + } + + // Write bundle to filesystem + bundleFile := path.Join(hostDir, config.BundleFileName) + if err := writeCerts(bundleFile, bundleSet); err != nil { + return fmt.Errorf("failed to write bundle: %w", err) + } + + log.Infof("%s: bundle written with %d CA certificates", containerName, len(bundleSet)) + return nil +} + +// writeCertificates writes certificates to the host filesystem +// Standalone utility function only performing file I/O operations and doesn't need access to plugin state. +// TODO implement SpiffeFS https://github.com/spiffe/spiffefs +func writeX509Content(hostDir string, config *identityConfig, certs []*x509.Certificate, privateKey []byte) error { + svidFile := path.Join(hostDir, config.CertFileName) + svidKeyFile := path.Join(hostDir, config.KeyFileName) + + if err := writeCerts(svidFile, certs); err != nil { + return err + } + + if err := writePrivateKey(svidKeyFile, privateKey); err != nil { + return err + } + + return nil +} + +func writeCerts(file string, certs []*x509.Certificate) error { + var pemData []byte + for _, cert := range certs { + + // TODO implement not writing expired certs + b := &pem.Block{ + Type: "CERTIFICATE", + Bytes: cert.Raw, + } + pemData = append(pemData, pem.EncodeToMemory(b)...) + } + return os.WriteFile(file, pemData, 0644) +} + +func writePrivateKey(file string, privateKey []byte) error { + b := &pem.Block{ + Type: "PRIVATE KEY", + Bytes: privateKey, + } + + return os.WriteFile(file, pem.EncodeToMemory(b), 0600) +} + +func parseIdentityConfig(ctr string, annotations map[string]string) (*identityConfig, error) { + var config identityConfig + + annotation := getAnnotation(annotations, identityKey, ctr) + if annotation == nil { + return nil, nil + } + + if err := yaml.Unmarshal(annotation, &config); err != nil { + return nil, fmt.Errorf("invalid identity annotation %q: %w", string(annotation), err) + } + + // Set default paths if not specified + if config.MountPath == "" { + config.MountPath = defaultMountPath + } + if config.CertFileName == "" { + config.CertFileName = defaultCertFileName + } + if config.KeyFileName == "" { + config.KeyFileName = defaultKeyFileName + } + if config.BundleFileName == "" { + config.BundleFileName = defaultBundleFileName + } + + return &config, nil +} + +func getAnnotation(annotations map[string]string, mainKey, ctr string) []byte { + for _, key := range []string { + mainKey + "/container." + ctr, + mainKey + "/pod", + mainKey, + } { + if key == "" || key[0] == '/' { + continue + } + if value, ok := annotations[key]; ok { + return []byte(value) + } + } + + return nil +} + +// Construct a container name for log messages. +func containerName(pod *api.PodSandbox, container *api.Container) string { + if pod != nil { + return pod.Name + "/" + container.Name + } + return container.Name +} + +func ensureUnixPrefix(path string) string { + if strings.HasPrefix(path, "unix://") { + return path + } + return "unix://" + path +} + +func getHostDir(hostMountPath, podUuid, containerName string ) string { + return filepath.Join(hostMountPath, podUuid, containerName) +} + +// Dump one or more objects, with an optional global prefix and per-object tags. +func dump(args ...interface{}) { + var ( + prefix string + idx int + ) + + if len(args)&0x1 == 1 { + prefix = args[0].(string) + idx++ + } + + for ; idx < len(args)-1; idx += 2 { + tag, obj := args[idx], args[idx+1] + msg, err := yaml.Marshal(obj) + if err != nil { + log.Infof("%s: %s: failed to dump object: %v", prefix, tag, err) + continue + } + + if prefix != "" { + log.Infof("%s: %s:", prefix, tag) + for _, line := range strings.Split(strings.TrimSpace(string(msg)), "\n") { + log.Infof("%s: %s", prefix, line) + } + } else { + log.Infof("%s:", tag) + for _, line := range strings.Split(strings.TrimSpace(string(msg)), "\n") { + log.Infof(" %s", line) + } + } + } +} + +func main() { + var ( + pluginIdx string + opts []stub.Option + err error + ) + + log = logrus.StandardLogger() + log.SetFormatter(&logrus.TextFormatter{ + PadLevelText: true, + }) + + flag.StringVar(&pluginIdx, "idx", "", "plugin index to register to NRI") + flag.BoolVar(&verbose, "verbose", false, "enable (more) verbose logging") + flag.StringVar(&spireAdminSocket, "spire-admin-socket", "/run/spire/admin-socket/admin.sock", "SPIRE Delegated Identity API socket path") + flag.StringVar(&hostMountPath, "host-mount-path", "/var/run/spiffe/secrets/", "Host Volume that will be used for writing identity artifacts of workloads") + flag.Parse() + + if pluginIdx != "" { + opts = append(opts, stub.WithPluginIdx(pluginIdx)) + } + + p := &plugin{ + watchers: make(map[string]*containerWatcher), + } + + p.delegatedIdentityConn, err = grpc.NewClient( + ensureUnixPrefix(spireAdminSocket), + grpc.WithTransportCredentials(insecure.NewCredentials()), + ) + + if err != nil { + log.Fatalf("failed to create SPIRE Delegated Identity API client: %v", err) + } + + defer p.delegatedIdentityConn.Close() + + p.delegatedIdentityClient = delegatedidentityv1.NewDelegatedIdentityClient(p.delegatedIdentityConn) + + if p.stub, err = stub.New(p, opts...); err != nil { + log.Fatalf("failed to create plugin stub: %v", err) + } + + err = p.stub.Run(context.Background()) + if err != nil { + log.Errorf("plugin exited with error %v", err) + os.Exit(1) + } +} diff --git a/plugins/identity-injector/identity-injector_test.go b/plugins/identity-injector/identity-injector_test.go new file mode 100644 index 00000000..f689a22b --- /dev/null +++ b/plugins/identity-injector/identity-injector_test.go @@ -0,0 +1,102 @@ +/* + Copyright The containerd Authors. + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. +*/ + +package main + +import ( + "testing" + + "github.com/stretchr/testify/require" +) + +func TestParseIdentityAnnotations(t *testing.T) { + type testCase struct { + name string + annotations map[string]string + result *identityConfig + } + + for _, tc := range []*testCase{ + { + name: "no identity annotations", + annotations: map[string]string{ + "foo": "bar", + }, + result: nil, + }, + { + name: "identity annotated", + annotations: map[string]string{ + "identity.noderesource.dev/container.c0": ` +mount_path: /var/run/secrets/spiffe +host_mount_path: /var/run/secrets/ +cert_file_name: svid.pem +key_file_name: svid_key.pem +bundle_file_name: svid_bundle.pem +spiffe_id: spiffe://example.org/p0/c0 +`, + }, + result: &identityConfig{ + MountPath: "/var/run/secrets/spiffe", + CertFileName: "svid.pem", + KeyFileName: "svid_key.pem", + BundleFileName: "svid_bundle.pem", + SpiffeId: "spiffe://example.org/p0/c0", + }, + }, + { + name: "container name mismatch", + annotations: map[string]string{ + "identity.noderesource.dev/container.c1": ` +mount_path: /var/run/secrets/spiffe +cert_file_name: svid.pem +key_file_name: svid_key.pem +bundle_file_name: svid_bundle.pem +spiffe_id: spiffe://example.org/p0/c0 +`, + }, + result: nil, + }, + { + name: "no container name", + annotations: map[string]string{ + "identity.noderesource.dev": ` +mount_path: /var/run/secrets/spiffe +cert_file_name: svid.pem +key_file_name: svid_key.pem +bundle_file_name: svid_bundle.pem +spiffe_id: spiffe://example.org/p0/c0 +`, + }, + result: &identityConfig{ + MountPath: "/var/run/secrets/spiffe", + CertFileName: "svid.pem", + KeyFileName: "svid_key.pem", + BundleFileName: "svid_bundle.pem", + SpiffeId: "spiffe://example.org/p0/c0", + }, + }, + } { + t.Run(tc.name, func(t *testing.T) { + config, err := parseIdentityConfig("c0", tc.annotations) + require.Nil(t, err, "config parsing error") + require.Equal(t, tc.result, config, "parsed config") + }) + } + +} + +// TODO create test cases for processDelegatedIdentityUpdate() diff --git a/plugins/identity-injector/setup.md b/plugins/identity-injector/setup.md new file mode 100644 index 00000000..69f0d0d6 --- /dev/null +++ b/plugins/identity-injector/setup.md @@ -0,0 +1,368 @@ +This document describes how to get a test setup up and running to test the NRI Identity Plugin + +Note 1: that this setup uses Containerd and Kubernetes both built and run from source. + +Note 2: This test setup is executed using `local-cluster-up.sh`. + +# Step 1: Spiffe/Spire + +The goal of this step is to have the Spire Server and Spire Agent up and running. https://spiffe.io/docs/latest/try/getting-started-k8s/ was used as the primary reference, with IBM Bob assisting with any gaps. + +## Step 1.1: Setup Spiffe/Spire + +**Why:** The SPIRE Server is the central authority that issues and manages X.509 SVIDs (SPIFFE Verifiable Identity Documents). The SPIRE Agent runs as a DaemonSet on every Kubernetes node and acts as the local broker between workloads and the SPIRE Server. Without these two components running, there is no identity infrastructure for the NRI Identity Plugin to delegate to. + +The SPIRE Server needs a `PersistentVolume` backed by a host directory (`/tmp/spire-data`) so that its SQLite database (used to store registration entries and keys) survives pod restarts. The `chmod 777` ensures the SPIRE Server container (which runs as a non-root user) can write to that directory. + +The manifests applied from the SPIRE tutorials repository install: +- A `spire` namespace to isolate all SPIRE resources. +- ServiceAccounts, ClusterRoles, and ClusterRoleBindings granting SPIRE the Kubernetes API access it needs for node attestation and workload attestation. +- A ConfigMap holding the SPIRE Server configuration (trust domain, port, etc.). +- A StatefulSet deploying a single SPIRE Server pod, and a Service exposing it to the SPIRE Agent. +- A ConfigMap holding the SPIRE Agent configuration (server address, trust bundle path, socket path, etc.). +- A DaemonSet deploying the SPIRE Agent on every node so that workloads on any node can obtain identities. + +**Result:** After running these commands, the `spire` namespace will contain a running `spire-server-0` StatefulSet pod and a `spire-agent` DaemonSet pod. The SPIRE Server will have written its initial trust bundle and be ready to accept registration entries. The SPIRE Agent will have completed node attestation (proving to the server which node it is running on) and will be ready to issue SVIDs to workloads. + +``` +sudo ./_output/bin/kubectl --kubeconfig=/var/run/kubernetes/admin.kubeconfig create namespace spire + +sudo mkdir /tmp/spire-data + +sudo chmod 777 /tmp/spire-data + +cat <