Files
sless/terraform/provider/internal/resources/job_resource.go
T

384 lines
16 KiB
Go
Raw 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.
// deleteOnFail — откат: удаляем CR если upload провалился,
// иначе при следующем apply будет 409 (CR есть, state пустой).
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())
if delErr := r.client.DeleteJob(ctx, ns, plan.Name.ValueString()); delErr != nil {
resp.Diagnostics.AddWarning("rollback delete job", delErr.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())
if delErr := r.client.DeleteJob(ctx, ns, plan.Name.ValueString()); delErr != nil {
resp.Diagnostics.AddWarning("rollback delete job", delErr.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),
}
}