Проблема: upload handler сбрасывал phase в Pending ПОСЛЕ того как контроллер уже выставил Building → бесконечный цикл (94 job'а). Решение: - controllers/function_controller.go: startBuild сначала ставит аннотацию last-built-s3key = spec.S3Key (idempotency guard), потом status. Reconcile стартует сборку только если spec.S3Key != last-built-s3key. - handler/upload.go: убран Status().Update() — контроллер сам управляет фазой. - builder.go: IsAlreadyExists при создании job не ошибка (parallel reconcile).
184 lines
6.6 KiB
Go
184 lines
6.6 KiB
Go
// Изменено: 2026-03-07
|
||
// upload.go — обработчик загрузки кода функции.
|
||
// Принимает zip от пользователя, генерирует Dockerfile, упаковывает в tar.gz,
|
||
// кладёт в S3 и обновляет Function CRD чтобы контроллер запустил kaniko.
|
||
//
|
||
// Почему tar.gz а не zip: kaniko читает build context только в формате tar (или OCI layout).
|
||
// Почему генерируем Dockerfile здесь: пользователь не должен думать про образы —
|
||
// это детали платформы, скрытые от него.
|
||
|
||
package handler
|
||
|
||
import (
|
||
"archive/tar"
|
||
"archive/zip"
|
||
"bytes"
|
||
"compress/gzip"
|
||
"fmt"
|
||
"io"
|
||
"net/http"
|
||
"time"
|
||
|
||
"k8s.io/apimachinery/pkg/api/errors"
|
||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||
|
||
slessv1alpha1 "gitea-naeel.giteak8s.services.ngcloud.ru/naeel/sless/api/v1alpha1"
|
||
)
|
||
|
||
// runtimeBaseImage возвращает Docker образ базового runtime для данного runtime-идентификатора.
|
||
// Соглашение: образы лежат на DockerHub под аккаунтом naeel.
|
||
// Возвращает ошибку если runtime не поддерживается — это граница валидации.
|
||
func runtimeBaseImage(runtime string) (string, error) {
|
||
switch runtime {
|
||
case "python3.11":
|
||
return "naeel/sless-runtime-python3.11:latest", nil
|
||
default:
|
||
return "", fmt.Errorf("unsupported runtime: %q (supported: python3.11)", runtime)
|
||
}
|
||
}
|
||
|
||
// generateDockerfile генерирует Dockerfile для kaniko.
|
||
// Базовый образ содержит HTTP-обёртку (server.py).
|
||
// Пользовательский код копируется в /app/function/ поверх базового образа.
|
||
func generateDockerfile(runtime string) ([]byte, error) {
|
||
baseImage, err := runtimeBaseImage(runtime)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
content := fmt.Sprintf("FROM %s\nCOPY . /app/function/\n", baseImage)
|
||
return []byte(content), nil
|
||
}
|
||
|
||
// zipToTarGz распаковывает zip и упаковывает содержимое + Dockerfile в tar.gz.
|
||
// Результат кладётся в переданный buf.
|
||
// Почему распаковываем zip и перепаковываем: kaniko не умеет читать zip-контекст,
|
||
// только tar(.gz) или OCI.
|
||
func zipToTarGz(zipData []byte, dockerfileContent []byte, buf *bytes.Buffer) error {
|
||
zr, err := zip.NewReader(bytes.NewReader(zipData), int64(len(zipData)))
|
||
if err != nil {
|
||
return fmt.Errorf("parse zip: %w", err)
|
||
}
|
||
|
||
gw := gzip.NewWriter(buf)
|
||
tw := tar.NewWriter(gw)
|
||
|
||
// Первым файлом пишем Dockerfile — kaniko ищет его в корне контекста
|
||
if err := tw.WriteHeader(&tar.Header{
|
||
Name: "Dockerfile",
|
||
Mode: 0644,
|
||
Size: int64(len(dockerfileContent)),
|
||
ModTime: time.Now(),
|
||
}); err != nil {
|
||
return fmt.Errorf("write Dockerfile header: %w", err)
|
||
}
|
||
if _, err := tw.Write(dockerfileContent); err != nil {
|
||
return fmt.Errorf("write Dockerfile: %w", err)
|
||
}
|
||
|
||
// Копируем файлы из zip в tar
|
||
for _, f := range zr.File {
|
||
if f.FileInfo().IsDir() {
|
||
continue // пустые директории не нужны
|
||
}
|
||
rc, err := f.Open()
|
||
if err != nil {
|
||
return fmt.Errorf("open zip entry %s: %w", f.Name, err)
|
||
}
|
||
data, err := io.ReadAll(rc)
|
||
rc.Close()
|
||
if err != nil {
|
||
return fmt.Errorf("read zip entry %s: %w", f.Name, err)
|
||
}
|
||
|
||
if err := tw.WriteHeader(&tar.Header{
|
||
Name: f.Name,
|
||
Mode: 0644,
|
||
Size: int64(len(data)),
|
||
ModTime: f.Modified,
|
||
}); err != nil {
|
||
return fmt.Errorf("write tar header %s: %w", f.Name, err)
|
||
}
|
||
if _, err := tw.Write(data); err != nil {
|
||
return fmt.Errorf("write tar entry %s: %w", f.Name, err)
|
||
}
|
||
}
|
||
|
||
if err := tw.Close(); err != nil {
|
||
return fmt.Errorf("close tar: %w", err)
|
||
}
|
||
return gw.Close()
|
||
}
|
||
|
||
// UploadCode — POST /v1/namespaces/{namespace}/functions/{name}/upload
|
||
// Принимает multipart/form-data с полем "code" (zip архив с кодом функции).
|
||
// Генерирует Dockerfile, пакует tar.gz, загружает в S3, обновляет Function CRD.
|
||
func (h *Handler) UploadCode(w http.ResponseWriter, r *http.Request) {
|
||
ns := namespace(r)
|
||
name := pathVar(r, "name")
|
||
|
||
// Проверяем что Function CRD существует и получаем runtime
|
||
fn := &slessv1alpha1.Function{}
|
||
if err := h.K8s.Get(r.Context(), client.ObjectKey{Name: name, Namespace: ns}, fn); err != nil {
|
||
if errors.IsNotFound(err) {
|
||
writeJSON(w, http.StatusNotFound, errResp("function not found"))
|
||
return
|
||
}
|
||
writeJSON(w, http.StatusInternalServerError, errResp(err.Error()))
|
||
return
|
||
}
|
||
|
||
// Лимит 32MB на загрузку кода функции
|
||
if err := r.ParseMultipartForm(32 << 20); err != nil {
|
||
writeJSON(w, http.StatusBadRequest, errResp("invalid multipart form: "+err.Error()))
|
||
return
|
||
}
|
||
file, _, err := r.FormFile("code")
|
||
if err != nil {
|
||
writeJSON(w, http.StatusBadRequest, errResp(`field "code" is required (zip file)`))
|
||
return
|
||
}
|
||
defer file.Close()
|
||
|
||
zipData, err := io.ReadAll(file)
|
||
if err != nil {
|
||
writeJSON(w, http.StatusInternalServerError, errResp("read upload: "+err.Error()))
|
||
return
|
||
}
|
||
|
||
// Генерируем Dockerfile под runtime функции
|
||
dockerfileContent, err := generateDockerfile(fn.Spec.Runtime)
|
||
if err != nil {
|
||
writeJSON(w, http.StatusBadRequest, errResp(err.Error()))
|
||
return
|
||
}
|
||
|
||
// Упаковываем Dockerfile + код пользователя в tar.gz для kaniko
|
||
var buf bytes.Buffer
|
||
if err := zipToTarGz(zipData, dockerfileContent, &buf); err != nil {
|
||
writeJSON(w, http.StatusInternalServerError, errResp("pack build context: "+err.Error()))
|
||
return
|
||
}
|
||
|
||
// Версия на основе timestamp — каждый upload → новый уникальный ключ в S3
|
||
version := time.Now().Format("20060102150405")
|
||
s3Key, err := h.S3.UploadContext(r.Context(), ns, name, version, &buf, int64(buf.Len()))
|
||
if err != nil {
|
||
writeJSON(w, http.StatusInternalServerError, errResp("upload to S3: "+err.Error()))
|
||
return
|
||
}
|
||
|
||
// Обновляем Function CRD: новый s3Key и bucket → контроллер запустит kaniko
|
||
fn.Spec.S3Key = s3Key
|
||
fn.Spec.S3Bucket = h.S3.Bucket()
|
||
if err := h.K8s.Update(r.Context(), fn); err != nil {
|
||
writeJSON(w, http.StatusInternalServerError, errResp("update function: "+err.Error()))
|
||
return
|
||
}
|
||
|
||
writeJSON(w, http.StatusOK, map[string]string{
|
||
"s3_key": s3Key,
|
||
"phase": string(slessv1alpha1.FunctionPhasePending),
|
||
"message": "build queued",
|
||
})
|
||
}
|