// 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), } }