chore: refactor read scheme to openapi from kubernetes api definition

Previously, only a single embedded Kubernetes API version could be specified, resulting in a lack of support for certain GVKs.
To address this, the goal is to create and utilize a unified scheme that consolidates Kubernetes API definitions.
As a preliminary step, the current method of loading API definitions will be improved.
This commit is contained in:
koba1t
2026-07-17 17:16:44 +09:00
committed by Yugo Kobayashi
parent f7d33c2610
commit 9216c6f573
29 changed files with 1791 additions and 45006 deletions

View File

@@ -0,0 +1,387 @@
// Copyright 2026 The Kubernetes Authors.
// SPDX-License-Identifier: Apache-2.0
// openapi-bundle compiles a Kubernetes OpenAPI v2 protobuf document into the
// compact, deterministic bundle embedded by kyaml.
package main
import (
"bufio"
"compress/gzip"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"flag"
"fmt"
"io"
"os"
"path/filepath"
"strings"
"time"
openapi_v2 "github.com/google/gnostic-models/openapiv2"
"google.golang.org/protobuf/proto"
"k8s.io/kube-openapi/pkg/validation/spec"
"sigs.k8s.io/kustomize/kyaml/openapi/internal/builtinopenapi"
)
const gvkExtension = "x-kubernetes-group-version-kind"
const maxInputSize = 64 << 20
type options struct {
input string
output string
legacyProtoOutput string
kubernetesVersion string
}
func main() {
var opts options
flag.StringVar(&opts.input, "input", "", "path to a Kubernetes OpenAPI v2 protobuf document (optionally gzip-compressed)")
flag.StringVar(&opts.output, "output", "", "path to the generated .json.gz bundle")
flag.StringVar(&opts.legacyProtoOutput, "legacy-proto-output", "", "optional path to a deterministic gzip archive of the input protobuf")
flag.StringVar(&opts.kubernetesVersion, "kubernetes-version", "", "Kubernetes version represented by the input")
flag.Parse()
if err := run(opts); err != nil {
fmt.Fprintln(os.Stderr, err)
os.Exit(1)
}
}
func run(opts options) error {
if opts.input == "" || opts.output == "" || opts.kubernetesVersion == "" {
return errors.New("-input, -output, and -kubernetes-version are required")
}
input, err := readInput(opts.input)
if err != nil {
return fmt.Errorf("read input: %w", err)
}
bundle, err := compile(input, opts.kubernetesVersion)
if err != nil {
return fmt.Errorf("compile bundle: %w", err)
}
if err := writeBundle(opts.output, bundle); err != nil {
return fmt.Errorf("write bundle: %w", err)
}
if opts.legacyProtoOutput != "" {
if err := writeGzip(opts.legacyProtoOutput, input); err != nil {
return fmt.Errorf("write legacy protobuf archive: %w", err)
}
}
return nil
}
func readInput(path string) ([]byte, error) {
return readInputWithLimit(path, maxInputSize)
}
func readInputWithLimit(path string, limit int64) (result []byte, resultErr error) {
file, err := os.Open(path)
if err != nil {
return nil, fmt.Errorf("open %q: %w", path, err)
}
defer func() {
if err := file.Close(); err != nil {
resultErr = errors.Join(resultErr, fmt.Errorf("close input %q: %w", path, err))
}
}()
reader := bufio.NewReader(file)
magic, err := reader.Peek(2)
if err != nil && !errors.Is(err, io.EOF) {
return nil, fmt.Errorf("inspect input %q: %w", path, err)
}
if len(magic) < 2 || magic[0] != 0x1f || magic[1] != 0x8b {
return readLimitedInput(reader, limit, "input", path)
}
gzipReader, err := gzip.NewReader(reader)
if err != nil {
return nil, fmt.Errorf("open gzip input %q: %w", path, err)
}
uncompressed, readErr := readLimitedInput(gzipReader, limit, "decompressed input", path)
closeErr := gzipReader.Close()
if readErr != nil {
return nil, readErr
}
if closeErr != nil {
return nil, fmt.Errorf("close gzip input %q: %w", path, closeErr)
}
return uncompressed, nil
}
func readLimitedInput(reader io.Reader, limit int64, description, path string) ([]byte, error) {
contents, err := io.ReadAll(io.LimitReader(reader, limit+1))
if err != nil {
return nil, fmt.Errorf("read %s %q: %w", description, path, err)
}
if int64(len(contents)) > limit {
return nil, fmt.Errorf("%s %q exceeds %d bytes", description, path, limit)
}
return contents, nil
}
func compile(input []byte, kubernetesVersion string) (*builtinopenapi.Bundle, error) {
document := &openapi_v2.Document{}
if err := proto.Unmarshal(input, document); err != nil {
return nil, fmt.Errorf("unmarshal OpenAPI protobuf: %w", err)
}
var swagger spec.Swagger
ok, err := swagger.FromGnostic(document)
if err != nil {
return nil, fmt.Errorf("convert gnostic document: %w", err)
}
if !ok {
return nil, errors.New("gnostic document cannot be converted without data loss")
}
resources, err := collectResources(&swagger)
if err != nil {
return nil, err
}
if err := validateDefinitionReferences(swagger.Definitions); err != nil {
return nil, err
}
digest := sha256.Sum256(input)
bundle := &builtinopenapi.Bundle{
FormatVersion: builtinopenapi.FormatVersion,
Coverage: builtinopenapi.Coverage{
Floor: kubernetesVersion,
Ceiling: kubernetesVersion,
},
SelectionPolicy: builtinopenapi.SelectionPolicy,
Sources: []builtinopenapi.Source{{
KubernetesVersion: kubernetesVersion,
SHA256: hex.EncodeToString(digest[:]),
}},
Definitions: swagger.Definitions,
Resources: resources,
}
if err := bundle.Validate(); err != nil {
return nil, fmt.Errorf("validate compiled bundle: %w", err)
}
return bundle, nil
}
func collectResources(swagger *spec.Swagger) ([]builtinopenapi.Resource, error) {
resources := map[string]builtinopenapi.Resource{}
for definitionName, definition := range swagger.Definitions {
extension, found := definition.Extensions[gvkExtension]
if !found {
continue
}
entries, ok := extension.([]interface{})
if !ok {
return nil, fmt.Errorf("definition %q has a malformed %s extension", definitionName, gvkExtension)
}
for _, entry := range entries {
apiVersion, kind, err := parseGVK(entry)
if err != nil {
return nil, fmt.Errorf("definition %q: %w", definitionName, err)
}
key := resourceKey(apiVersion, kind)
resource := resources[key]
if resource.Definition != "" && resource.Definition != definitionName {
return nil, fmt.Errorf("GVK %s/%s is advertised by definitions %q and %q",
apiVersion, kind, resource.Definition, definitionName)
}
resource.APIVersion = apiVersion
resource.Kind = kind
resource.Definition = definitionName
resources[key] = resource
}
}
if err := collectPathResources(swagger.Paths, resources); err != nil {
return nil, err
}
result := make([]builtinopenapi.Resource, 0, len(resources))
for _, resource := range resources {
result = append(result, resource)
}
builtinopenapi.SortResources(result)
return result, nil
}
func collectPathResources(paths *spec.Paths, resources map[string]builtinopenapi.Resource) error {
if paths == nil {
return nil
}
for path, pathInfo := range paths.Paths {
if pathInfo.Get == nil {
continue
}
extension, found := pathInfo.Get.Extensions[gvkExtension]
if !found {
continue
}
apiVersion, kind, err := parseGVK(extension)
if err != nil {
return fmt.Errorf("path %q: %w", path, err)
}
key := resourceKey(apiVersion, kind)
resource := resources[key]
resource.APIVersion = apiVersion
resource.Kind = kind
if strings.Contains(path, "namespaces/{namespace}") {
resource.Scope = builtinopenapi.ScopeNamespaced
} else if resource.Scope == builtinopenapi.ScopeUnknown {
resource.Scope = builtinopenapi.ScopeCluster
}
resources[key] = resource
}
return nil
}
func parseGVK(value interface{}) (string, string, error) {
entry, ok := value.(map[string]interface{})
if !ok {
return "", "", fmt.Errorf("malformed %s extension entry", gvkExtension)
}
version, versionOK := entry["version"].(string)
kind, kindOK := entry["kind"].(string)
if !versionOK || version == "" || !kindOK || kind == "" {
return "", "", fmt.Errorf("incomplete %s extension entry", gvkExtension)
}
group, groupOK := entry["group"].(string)
if groupOK && group != "" {
return group + "/" + version, kind, nil
}
return version, kind, nil
}
func resourceKey(apiVersion, kind string) string {
return apiVersion + "\x00" + kind
}
func validateDefinitionReferences(definitions spec.Definitions) error {
b, err := json.Marshal(definitions)
if err != nil {
return fmt.Errorf("marshal definitions for reference validation: %w", err)
}
var value interface{}
if err := json.Unmarshal(b, &value); err != nil {
return fmt.Errorf("unmarshal definitions for reference validation: %w", err)
}
return walkReferences(value, definitions)
}
func walkReferences(value interface{}, definitions spec.Definitions) error {
switch typed := value.(type) {
case []interface{}:
for _, item := range typed {
if err := walkReferences(item, definitions); err != nil {
return err
}
}
case map[string]interface{}:
for key, item := range typed {
if key == "$ref" {
if err := validateReference(item, definitions); err != nil {
return err
}
}
if err := walkReferences(item, definitions); err != nil {
return err
}
}
}
return nil
}
func validateReference(value interface{}, definitions spec.Definitions) error {
ref, ok := value.(string)
// A schema may itself describe an object with a property named "$ref".
// Such a property has a schema object as its value and is not an OpenAPI
// reference keyword.
if !ok {
return nil
}
const prefix = "#/definitions/"
if !strings.HasPrefix(ref, prefix) {
return fmt.Errorf("OpenAPI definition contains unsupported reference %q", ref)
}
name := strings.TrimPrefix(ref, prefix)
name = strings.ReplaceAll(strings.ReplaceAll(name, "~1", "/"), "~0", "~")
if _, found := definitions[name]; !found {
return fmt.Errorf("OpenAPI definition references missing definition %q", name)
}
return nil
}
func writeBundle(path string, bundle *builtinopenapi.Bundle) error {
jsonBytes, err := json.Marshal(bundle)
if err != nil {
return fmt.Errorf("marshal bundle: %w", err)
}
return writeGzip(path, jsonBytes)
}
func writeGzip(path string, contents []byte) (resultErr error) {
dir := filepath.Dir(path)
if err := os.MkdirAll(dir, 0o755); err != nil {
return fmt.Errorf("create output directory %q: %w", dir, err)
}
tmp, err := os.CreateTemp(dir, ".openapi-bundle-*")
if err != nil {
return fmt.Errorf("create temporary output: %w", err)
}
tmpName := tmp.Name()
tmpClosed := false
defer func() {
if !tmpClosed {
if err := tmp.Close(); err != nil {
resultErr = errors.Join(resultErr, fmt.Errorf("close temporary output: %w", err))
}
}
if err := os.Remove(tmpName); err != nil && !errors.Is(err, os.ErrNotExist) {
resultErr = errors.Join(resultErr, fmt.Errorf("remove temporary output: %w", err))
}
}()
writer, err := gzip.NewWriterLevel(tmp, gzip.BestCompression)
if err != nil {
return fmt.Errorf("create gzip writer: %w", err)
}
writerClosed := false
defer func() {
if !writerClosed {
if err := writer.Close(); err != nil {
resultErr = errors.Join(resultErr, fmt.Errorf("close gzip writer: %w", err))
}
}
}()
writer.Header.ModTime = time.Time{}
writer.Header.Name = ""
writer.Header.Comment = ""
writer.Header.Extra = nil
writer.Header.OS = 255
if _, err := writer.Write(contents); err != nil {
return fmt.Errorf("write compressed output: %w", err)
}
closeWriterErr := writer.Close()
writerClosed = true
if closeWriterErr != nil {
return fmt.Errorf("close gzip writer: %w", closeWriterErr)
}
if err := tmp.Chmod(0o644); err != nil {
return fmt.Errorf("set output permissions: %w", err)
}
closeTempErr := tmp.Close()
tmpClosed = true
if closeTempErr != nil {
return fmt.Errorf("close temporary output: %w", closeTempErr)
}
if err := os.Rename(tmpName, path); err != nil {
return fmt.Errorf("replace output %q: %w", path, err)
}
return nil
}

View File

@@ -0,0 +1,149 @@
// Copyright 2026 The Kubernetes Authors.
// SPDX-License-Identifier: Apache-2.0
package main
import (
"bytes"
"compress/gzip"
"encoding/json"
"io"
"os"
"path/filepath"
"testing"
"time"
"github.com/stretchr/testify/require"
"sigs.k8s.io/kustomize/kyaml/openapi/internal/builtinopenapi"
)
func TestGeneratedBundleIsCurrentAndDeterministic(t *testing.T) {
source := filepath.Join("..", "..", "kubernetesapi", "v1_21_2", "swagger.pb.gz")
checkedIn := filepath.Join("..", "..", "kubernetesapi", "data",
"kubernetes-openapi-union-v1.21.2.bundle-v1.json.gz")
tempDir := t.TempDir()
first := filepath.Join(tempDir, "first.json.gz")
second := filepath.Join(tempDir, "second.json.gz")
legacy := filepath.Join(tempDir, "swagger.pb.gz")
for i, output := range []string{first, second} {
legacyOutput := ""
if i == 0 {
legacyOutput = legacy
}
require.NoError(t, run(options{
input: source,
output: output,
legacyProtoOutput: legacyOutput,
kubernetesVersion: "v1.21.2",
}))
}
want, err := os.ReadFile(checkedIn)
require.NoError(t, err)
got, err := os.ReadFile(first)
require.NoError(t, err)
again, err := os.ReadFile(second)
require.NoError(t, err)
require.Equal(t, want, got, "checked-in bundle is stale")
require.Equal(t, got, again, "bundle generation is not deterministic")
sourceArchive, err := os.ReadFile(source)
require.NoError(t, err)
legacyArchive, err := os.ReadFile(legacy)
require.NoError(t, err)
require.Equal(t, sourceArchive, legacyArchive, "compiler input archive is not deterministic")
reader, err := gzip.NewReader(bytes.NewReader(got))
require.NoError(t, err)
require.True(t, reader.ModTime.IsZero())
require.Empty(t, reader.Name)
require.Empty(t, reader.Comment)
decoder := json.NewDecoder(reader)
var bundle builtinopenapi.Bundle
require.NoError(t, decoder.Decode(&bundle))
var trailing interface{}
require.ErrorIs(t, decoder.Decode(&trailing), io.EOF)
require.NoError(t, reader.Close())
require.NoError(t, bundle.Validate())
require.Len(t, bundle.Definitions, 618)
require.Len(t, bundle.Resources, 275)
require.Equal(t, "5d171b55e9601912807a870d73ffe70bb306f5889a00e76986042a0f2d7b6bc2",
bundle.Sources[0].SHA256)
}
func TestWriteGzipUsesStableHeader(t *testing.T) {
path := filepath.Join(t.TempDir(), "data.gz")
require.NoError(t, writeGzip(path, []byte("data")))
b, err := os.ReadFile(path)
require.NoError(t, err)
reader, err := gzip.NewReader(bytes.NewReader(b))
require.NoError(t, err)
require.Equal(t, time.Time{}, reader.ModTime)
require.Empty(t, reader.Name)
require.Empty(t, reader.Comment)
require.Equal(t, byte(255), reader.OS)
require.NoError(t, reader.Close())
}
func TestReadInputWithLimit(t *testing.T) {
const limit = int64(4)
testDir := t.TempDir()
write := func(t *testing.T, name string, contents []byte) string {
t.Helper()
path := filepath.Join(testDir, name)
require.NoError(t, os.WriteFile(path, contents, 0o600))
return path
}
gzipContents := func(t *testing.T, contents []byte) []byte {
t.Helper()
var buffer bytes.Buffer
writer := gzip.NewWriter(&buffer)
_, err := writer.Write(contents)
require.NoError(t, err)
require.NoError(t, writer.Close())
return buffer.Bytes()
}
t.Run("raw at limit", func(t *testing.T) {
path := write(t, "raw", []byte("data"))
got, err := readInputWithLimit(path, limit)
require.NoError(t, err)
require.Equal(t, []byte("data"), got)
})
t.Run("raw exceeds limit", func(t *testing.T) {
path := write(t, "raw-large", []byte("large"))
_, err := readInputWithLimit(path, limit)
require.ErrorContains(t, err, "input")
require.ErrorContains(t, err, "exceeds 4 bytes")
})
t.Run("gzip at limit", func(t *testing.T) {
path := write(t, "gzip", gzipContents(t, []byte("data")))
got, err := readInputWithLimit(path, limit)
require.NoError(t, err)
require.Equal(t, []byte("data"), got)
})
t.Run("gzip exceeds limit", func(t *testing.T) {
path := write(t, "gzip-large", gzipContents(t, []byte("large")))
_, err := readInputWithLimit(path, limit)
require.ErrorContains(t, err, "decompressed input")
require.ErrorContains(t, err, "exceeds 4 bytes")
})
t.Run("invalid gzip header", func(t *testing.T) {
path := write(t, "gzip-invalid-header", []byte{0x1f, 0x8b})
_, err := readInputWithLimit(path, limit)
require.ErrorContains(t, err, "open gzip input")
})
t.Run("invalid gzip body", func(t *testing.T) {
contents := gzipContents(t, []byte("data"))
contents[len(contents)-1]++
path := write(t, "gzip-invalid-body", contents)
_, err := readInputWithLimit(path, limit)
require.ErrorContains(t, err, "read decompressed input")
})
}