Files
sless/terraform/provider/internal/resources/job_resource.go
Naeel 8ca8faedd1 feat(job): merge sless_function into sless_job — self-contained build+run
- FunctionJobSpec: убран FunctionRef, добавлены Runtime/Entrypoint/Env/S3Key/MemoryMB/TimeoutSec
- FunctionJobStatus: новый ImageRef, новая фаза Building
- FunctionJobReconciler: Building фаза (kaniko), убрана зависимость от Function CRD
  Builder+OperatorNamespace как поля struct; аннотация sless.kube5s.ru/build-job guard
- main.go: Builder+OperatorNamespace переданы в FunctionJobReconciler
- jobs.go handler: jobRequest/jobResponse без FunctionRef; новый UploadJobCode handler
- router.go: /jobs/{name}/upload маршрут
- client.go: JobRequest/JobResponse обновлены; UploadJobCode; uploadCodeToURL общий хелпер
- job_resource.go: полная переработка — источник/среда встроены в JobModel, ModifyPlan,
  Create с upload, wait_timeout_sec=900 по умолчанию (kaniko + выполнение)
- examples/POSTGRES/functions.tf: раскомментирован, sless_function удалён,
  sless_job самодостаточен (inline source_dir/runtime/entrypoint/env_vars)
2026-03-20 21:27:35 +03:00

376 lines
15 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 2026-03-20 (merge: убран FunctionRef, добавлены Runtime/Entrypoint/SourceDir — sless_job теперь самодостаточен)
// job_resource.go — Terraform ресурс sless_job.
//
// Lifecycle:
//
// Create: POST /v1/namespaces/{ns}/jobs
// Если source_dir задан → zipDir → UploadJobCode → контроллер запускает kaniko
// WaitJobDone покрывает: Building (kaniko) + Running (функция) + Succeeded/Failed
// Если run_id=0 — джоб создаётся в k8s, но не запускается (phase=Skipped).
// Read: GET /v1/namespaces/{ns}/jobs/{name} → sync phase/timing в state
// Delete: DELETE /v1/namespaces/{ns}/jobs/{name}
//
// run_id: 0 = создать без запуска; >0 = запустить.
// RequiresReplace на run_id + code_hash: изменение кода или run_id всегда пересоздаёт джоб.
// wait_timeout_sec должен покрывать оба этапа: kaniko (~5 мин) + выполнение функции.
// По умолчанию 900 сек (15 мин).
package resources
import (
"bytes"
"context"
"errors"
"fmt"
"time"
"terraform-provider-sless/internal/client"
"github.com/hashicorp/terraform-plugin-framework-validators/int64validator"
"github.com/hashicorp/terraform-plugin-framework/resource"
"github.com/hashicorp/terraform-plugin-framework/resource/schema"
"github.com/hashicorp/terraform-plugin-framework/resource/schema/int64default"
"github.com/hashicorp/terraform-plugin-framework/resource/schema/int64planmodifier"
"github.com/hashicorp/terraform-plugin-framework/resource/schema/planmodifier"
"github.com/hashicorp/terraform-plugin-framework/resource/schema/stringplanmodifier"
"github.com/hashicorp/terraform-plugin-framework/schema/validator"
"github.com/hashicorp/terraform-plugin-framework/types"
)
// defaultJobWaitTimeoutSec — дефолтный таймаут ожидания завершения sless_job.
// Покрывает kaniko сборку (~5 мин) + выполнение функции (~остаток).
// Переопределяется через wait_timeout_sec.
const defaultJobWaitTimeoutSec = 900
var _ resource.Resource = &JobResource{}
var _ resource.ResourceWithModifyPlan = &JobResource{}
type JobResource struct {
client *client.Client
}
func NewJobResource() resource.Resource {
return &JobResource{}
}
// JobModel — модель состояния terraform для sless_job.
type JobModel struct {
Name types.String `tfsdk:"name"`
Runtime types.String `tfsdk:"runtime"`
Entrypoint types.String `tfsdk:"entrypoint"`
MemoryMB types.Int64 `tfsdk:"memory_mb"`
TimeoutSec types.Int64 `tfsdk:"timeout_sec"`
EnvVars types.Map `tfsdk:"env_vars"`
// source_dir — директория с исходниками. Провайдер сам упакует в zip и загрузит.
SourceDir types.String `tfsdk:"source_dir"`
// code_hash — SHA256 содержимого source_dir. Computed; при изменении RequiresReplace.
CodeHash types.String `tfsdk:"code_hash"`
EventJSON types.String `tfsdk:"event_json"`
// RunID: 0 = не запускать (Skipped), 1+ = запустить/перезапустить.
// RequiresReplace: изменение = пересоздание FunctionJob → новый запуск.
RunID types.Int64 `tfsdk:"run_id"`
// wait_timeout_sec — максимальное ожидание (kaniko + выполнение). Дефолт 900 сек.
WaitTimeoutSec types.Int64 `tfsdk:"wait_timeout_sec"`
Phase types.String `tfsdk:"phase"`
ImageRef types.String `tfsdk:"image_ref"`
StartTime types.String `tfsdk:"start_time"`
CompletionTime types.String `tfsdk:"completion_time"`
Message types.String `tfsdk:"message"`
}
func (r *JobResource) Metadata(_ context.Context, req resource.MetadataRequest, resp *resource.MetadataResponse) {
resp.TypeName = req.ProviderTypeName + "_job"
}
func (r *JobResource) Schema(_ context.Context, _ resource.SchemaRequest, resp *resource.SchemaResponse) {
resp.Schema = schema.Schema{
MarkdownDescription: "Одноразовый запуск serverless функции. terraform apply блокируется до завершения. Включает kaniko сборку образа.",
Attributes: map[string]schema.Attribute{
"name": schema.StringAttribute{
Required: true,
PlanModifiers: []planmodifier.String{
stringplanmodifier.RequiresReplace(),
},
},
"runtime": schema.StringAttribute{
Required: true,
MarkdownDescription: "Среда выполнения: python3.11, nodejs20, go1.23.",
PlanModifiers: []planmodifier.String{
stringplanmodifier.RequiresReplace(),
},
},
"entrypoint": schema.StringAttribute{
Required: true,
MarkdownDescription: "Точка входа: module.function (например sql_runner.run_sql).",
PlanModifiers: []planmodifier.String{
stringplanmodifier.RequiresReplace(),
},
},
"memory_mb": schema.Int64Attribute{
Optional: true,
Computed: true,
Default: int64default.StaticInt64(128),
MarkdownDescription: "Лимит оперативной памяти в MB. По умолчанию 128.",
},
"timeout_sec": schema.Int64Attribute{
Optional: true,
Computed: true,
Default: int64default.StaticInt64(30),
MarkdownDescription: "Таймаут выполнения функции в секундах. По умолчанию 30.",
},
"env_vars": schema.MapAttribute{
ElementType: types.StringType,
Optional: true,
MarkdownDescription: "Переменные окружения для функции.",
},
"source_dir": schema.StringAttribute{
Optional: true,
MarkdownDescription: "Директория с исходниками. Провайдер сам запакует zip и загрузит.",
PlanModifiers: []planmodifier.String{
stringplanmodifier.RequiresReplace(),
},
},
// code_hash — вычисляется автоматически из source_dir (ModifyPlan).
// RequiresReplace: если код изменился (hash != state) → пересоздать джоб.
"code_hash": schema.StringAttribute{
Computed: true,
PlanModifiers: []planmodifier.String{
stringplanmodifier.RequiresReplace(),
},
},
"event_json": schema.StringAttribute{
Optional: true,
MarkdownDescription: `JSON-объект передаваемый в handle(event). По умолчанию "{}".`,
PlanModifiers: []planmodifier.String{
stringplanmodifier.RequiresReplace(),
},
},
// run_id: 0 = создать без запуска (Skipped), >0 = запустить.
// Изменение run_id (1→2→3...) триггерирует пересоздание = новый запуск.
"run_id": schema.Int64Attribute{
Optional: true,
Computed: true,
Default: int64default.StaticInt64(0),
MarkdownDescription: "0 = не запускать; >0 = запустить. Увеличьте run_id для повторного запуска.",
PlanModifiers: []planmodifier.Int64{
int64planmodifier.RequiresReplace(),
},
Validators: []validator.Int64{
int64validator.AtLeast(0),
},
},
// wait_timeout_sec покрывает kaniko сборку + выполнение функции.
"wait_timeout_sec": schema.Int64Attribute{
Optional: true,
Computed: true,
Default: int64default.StaticInt64(defaultJobWaitTimeoutSec),
MarkdownDescription: "Таймаут ожидания (kaniko + выполнение) в секундах. По умолчанию 900.",
},
// Computed — заполняются контроллером и финализируются после завершения
"phase": schema.StringAttribute{
Computed: true,
MarkdownDescription: "Фаза: Pending, Building, Running, Succeeded, Failed.",
},
"image_ref": schema.StringAttribute{
Computed: true,
MarkdownDescription: "Docker image собранный kaniko.",
},
"start_time": schema.StringAttribute{
Computed: true,
MarkdownDescription: "Время запуска k8s Job (RFC3339).",
},
"completion_time": schema.StringAttribute{
Computed: true,
MarkdownDescription: "Время завершения k8s Job (RFC3339).",
},
"message": schema.StringAttribute{
Computed: true,
MarkdownDescription: "Результат выполнения или сообщение об ошибке.",
},
},
}
}
func (r *JobResource) Configure(_ context.Context, req resource.ConfigureRequest, resp *resource.ConfigureResponse) {
if req.ProviderData == nil {
return
}
c, ok := req.ProviderData.(*client.Client)
if !ok {
resp.Diagnostics.AddError(
"unexpected provider data",
fmt.Sprintf("expected *client.Client, got: %T", req.ProviderData),
)
return
}
r.client = c
}
func (r *JobResource) Create(ctx context.Context, req resource.CreateRequest, resp *resource.CreateResponse) {
var plan JobModel
resp.Diagnostics.Append(req.Plan.Get(ctx, &plan)...)
if resp.Diagnostics.HasError() {
return
}
ns := r.client.Namespace
eventJSON := plan.EventJSON.ValueString()
if eventJSON == "" {
eventJSON = "{}"
}
envMap, diags := mapToStringMap(ctx, plan.EnvVars)
resp.Diagnostics.Append(diags...)
if resp.Diagnostics.HasError() {
return
}
runID := plan.RunID.ValueInt64()
createdJob, err := r.client.CreateJob(ctx, ns, client.JobRequest{
Name: plan.Name.ValueString(),
Runtime: plan.Runtime.ValueString(),
Entrypoint: plan.Entrypoint.ValueString(),
MemoryMB: int32(plan.MemoryMB.ValueInt64()),
TimeoutSec: int32(plan.TimeoutSec.ValueInt64()),
Env: envMap,
EventJSON: eventJSON,
RunID: runID,
})
if err != nil {
if errors.Is(err, client.ErrJobAlreadyExists) {
existingJob, getErr := r.client.GetJob(ctx, ns, plan.Name.ValueString())
if getErr != nil {
resp.Diagnostics.AddError("create job", fmt.Sprintf("job already exists and get failed: %s", getErr.Error()))
return
}
if existingJob == nil {
resp.Diagnostics.AddError("create job", "job already exists but cannot be read after conflict")
return
}
createdJob = existingJob
} else {
resp.Diagnostics.AddError("create job", err.Error())
return
}
}
// Загружаем код если задан source_dir — контроллер начнёт kaniko сборку после upload
if !plan.SourceDir.IsNull() && plan.SourceDir.ValueString() != "" {
zipData, hash, err := zipDir(plan.SourceDir.ValueString())
if err != nil {
resp.Diagnostics.AddError("zip source_dir", err.Error())
return
}
if err := r.client.UploadJobCode(ctx, ns, plan.Name.ValueString(), "function.zip", bytes.NewReader(zipData)); err != nil {
resp.Diagnostics.AddError("upload job code", err.Error())
return
}
plan.CodeHash = types.StringValue(hash)
}
// run_id=0: не ждём, пишем state сразу (job в состоянии Skipped)
if runID == 0 {
resp.Diagnostics.Append(resp.State.Set(ctx, jobToModel(plan, createdJob))...)
return
}
// Блокируем apply до завершения джоба (охватывает Building + Running → Succeeded/Failed)
waitSec := plan.WaitTimeoutSec.ValueInt64()
if waitSec <= 0 {
waitSec = defaultJobWaitTimeoutSec
}
j, err := r.client.WaitJobDone(ctx, ns, plan.Name.ValueString(), time.Duration(waitSec)*time.Second)
if err != nil {
resp.Diagnostics.AddError("waiting for job to complete", err.Error())
return
}
resp.Diagnostics.Append(resp.State.Set(ctx, jobToModel(plan, j))...)
}
func (r *JobResource) Read(ctx context.Context, req resource.ReadRequest, resp *resource.ReadResponse) {
var state JobModel
resp.Diagnostics.Append(req.State.Get(ctx, &state)...)
if resp.Diagnostics.HasError() {
return
}
j, err := r.client.GetJob(ctx, r.client.Namespace, state.Name.ValueString())
if err != nil {
resp.Diagnostics.AddError("read job", err.Error())
return
}
if j == nil {
resp.State.RemoveResource(ctx)
return
}
resp.Diagnostics.Append(resp.State.Set(ctx, jobToModel(state, j))...)
}
// Update не реализован — все input-поля имеют RequiresReplace.
// terraform-plugin-framework никогда не вызовет Update для этого ресурса.
func (r *JobResource) Update(_ context.Context, _ resource.UpdateRequest, resp *resource.UpdateResponse) {
resp.Diagnostics.AddError("update not supported", "sless_job does not support in-place updates; increment run_id to re-run")
}
func (r *JobResource) Delete(ctx context.Context, req resource.DeleteRequest, resp *resource.DeleteResponse) {
var state JobModel
resp.Diagnostics.Append(req.State.Get(ctx, &state)...)
if resp.Diagnostics.HasError() {
return
}
if err := r.client.DeleteJob(ctx, r.client.Namespace, state.Name.ValueString()); err != nil {
resp.Diagnostics.AddError("delete job", err.Error())
}
}
// ModifyPlan вычисляет hash директории source_dir на фазе terraform plan.
// Если hash изменился относительно state — code_hash в plan будет другим,
// что вместе с RequiresReplace на code_hash триггерирует пересоздание джоба.
func (r *JobResource) ModifyPlan(ctx context.Context, req resource.ModifyPlanRequest, resp *resource.ModifyPlanResponse) {
if req.Plan.Raw.IsNull() {
return
}
var plan JobModel
resp.Diagnostics.Append(req.Plan.Get(ctx, &plan)...)
if resp.Diagnostics.HasError() {
return
}
if plan.SourceDir.IsNull() || plan.SourceDir.ValueString() == "" {
return
}
_, hash, err := zipDir(plan.SourceDir.ValueString())
if err != nil {
return // директория может не существовать при первом init — не фейлим plan
}
plan.CodeHash = types.StringValue(hash)
resp.Diagnostics.Append(resp.Plan.Set(ctx, plan)...)
}
// jobToModel конвертирует API-ответ + plan (для локальных полей) → state модель.
func jobToModel(plan JobModel, j *client.JobResponse) JobModel {
waitSec := plan.WaitTimeoutSec
if waitSec.IsNull() || waitSec.IsUnknown() || waitSec.ValueInt64() <= 0 {
waitSec = types.Int64Value(defaultJobWaitTimeoutSec)
}
return JobModel{
Name: types.StringValue(j.Name),
Runtime: types.StringValue(j.Runtime),
Entrypoint: types.StringValue(j.Entrypoint),
MemoryMB: plan.MemoryMB,
TimeoutSec: plan.TimeoutSec,
EnvVars: plan.EnvVars,
SourceDir: plan.SourceDir,
CodeHash: plan.CodeHash,
EventJSON: plan.EventJSON,
RunID: types.Int64Value(j.RunID),
WaitTimeoutSec: waitSec,
Phase: types.StringValue(j.Phase),
ImageRef: types.StringValue(j.ImageRef),
StartTime: types.StringValue(j.StartTime),
CompletionTime: types.StringValue(j.CompletionTime),
Message: types.StringValue(j.Message),
}
}