feat: upload endpoint POST /v1/namespaces/{ns}/functions/{name}/upload

- s3/client.go: Bucket() accessor + UploadContext() для tar.gz build context
- handler/upload.go: принимает zip, генерирует Dockerfile (FROM naeel/sless-runtime-{runtime}),
  перепаковывает в tar.gz, загружает в S3, обновляет Function CRD → kaniko запускается
- router.go: маршрут POST .../upload зарегистрирован
This commit is contained in:
“Naeel”
2026-03-07 09:36:32 +04:00
parent c61e822308
commit 74458f848d
3 changed files with 213 additions and 0 deletions
+191
View File
@@ -0,0 +1,191 @@
// Изменено: 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
}
// Сбрасываем phase в Pending чтобы контроллер начал сборку
fn.Status.Phase = slessv1alpha1.FunctionPhasePending
fn.Status.Message = "code uploaded, build queued"
if err := h.K8s.Status().Update(r.Context(), fn); err != nil {
// Не фатально — контроллер подберёт через reconcile
h.Log.Warn("failed to reset function phase", "name", name, "ns", ns, "err", err)
}
writeJSON(w, http.StatusOK, map[string]string{
"s3_key": s3Key,
"phase": string(slessv1alpha1.FunctionPhasePending),
"message": "build queued",
})
}
+3
View File
@@ -33,6 +33,9 @@ func NewRouter(h *handler.Handler, apiToken string, log *slog.Logger) http.Handl
// Invocation logs
v1.HandleFunc("/namespaces/{namespace}/functions/{name}/invocations", h.ListInvocations).Methods(http.MethodGet)
// Upload code — принимает zip, генерирует Dockerfile, кладёт tar.gz в S3, запускает сборку
v1.HandleFunc("/namespaces/{namespace}/functions/{name}/upload", h.UploadCode).Methods(http.MethodPost)
// Triggers CRUD
v1.HandleFunc("/namespaces/{namespace}/triggers", h.ListTriggers).Methods(http.MethodGet)
v1.HandleFunc("/namespaces/{namespace}/triggers", h.CreateTrigger).Methods(http.MethodPost)
+19
View File
@@ -80,3 +80,22 @@ func (c *Client) Delete(ctx context.Context, key string) error {
}
return nil
}
// Bucket возвращает имя бакета, с которым работает клиент.
func (c *Client) Bucket() string {
return c.bucket
}
// UploadContext загружает tar.gz с контекстом сборки (Dockerfile + код) в S3.
// Ключ: contexts/{namespace}/{name}/{version}.tar.gz
// Kaniko читает этот tar.gz как build context (--context=s3://bucket/key).
func (c *Client) UploadContext(ctx context.Context, namespace, funcName, version string, r io.Reader, size int64) (string, error) {
key := fmt.Sprintf("contexts/%s/%s/%s.tar.gz", namespace, funcName, version)
_, err := c.mc.PutObject(ctx, c.bucket, key, r, size, minio.PutObjectOptions{
ContentType: "application/gzip",
})
if err != nil {
return "", fmt.Errorf("upload build context: %w", err)
}
return key, nil
}