Files

868 lines
33 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-17 12:20
// client.go — HTTP-клиент для REST API sless оператора.
// Изолирован от terraform-plugin-framework — зависит только от stdlib и net/http.
// Все методы принимают ctx для правильной работы с таймаутами terraform.
//
// Архитектура namespace:
// - Namespace вычисляется из JWT-токена провайдером ОДИН РАЗ при Configure().
// - Алгоритм: JWT.sub → SHA256 → hex первые 16 байт → "sless-{hex}"
// - Client хранит уже вычисленный Namespace — ресурсы просто читают его.
// - SubFromJWT и NamespaceFromSub — package-level функции (не методы),
// вызываются из provider.Configure() до создания Client.
//
// Валидация токена:
// - PingNubesAPI делает GET запрос к nubes API с Bearer токеном.
// - 401/403 → токен невалиден → ошибка инициализации провайдера.
// - Любой другой ответ → токен принят сервером.
package client
import (
"bytes"
"context"
"crypto/sha256"
"encoding/base64"
"encoding/json"
"errors"
"fmt"
"io"
"mime/multipart"
"net/http"
"os"
"path/filepath"
"strings"
"time"
)
// Client — HTTP-клиент к sless operator REST API.
type Client struct {
httpClient *http.Client
endpoint string
token string
// Namespace — k8s namespace пользователя, вычисленный из JWT-токена.
// Устанавливается один раз при создании Client в provider.Configure().
// Алгоритм вычисления: SubFromJWT → NamespaceFromSub.
// Все ресурсы (FunctionResource, TriggerResource, JobResource) читают это поле.
Namespace string
}
// ErrJobAlreadyExists возвращается при попытке создать FunctionJob с уже существующим именем.
var ErrJobAlreadyExists = errors.New("job already exists")
// New создаёт клиент.
// - endpoint — базовый URL оператора (без trailing slash), например "https://sless-api.kube5s.ru"
// - token — Bearer JWT-токен облака
// - namespace — k8s namespace пользователя (вычислен через NamespaceFromSub)
func New(endpoint, token, namespace string) *Client {
return &Client{
httpClient: &http.Client{Timeout: 30 * time.Second},
endpoint: endpoint,
token: token,
Namespace: namespace,
}
}
// SubFromJWT декодирует JWT payload (base64url) и возвращает claim "sub".
// Не проверяет подпись — только структуру и наличие sub.
// Проверка подписи не нужна: токен val идирован через PingNubesAPI запросом к реальному API.
func SubFromJWT(token string) (string, error) {
parts := strings.Split(token, ".")
if len(parts) != 3 {
return "", fmt.Errorf("invalid JWT: expected 3 parts, got %d", len(parts))
}
// JWT использует base64url без padding — добавляем padding
payload := parts[1]
switch len(payload) % 4 {
case 2:
payload += "=="
case 3:
payload += "="
}
decoded, err := base64.URLEncoding.DecodeString(payload)
if err != nil {
// Пробуем StdEncoding на случай нестандартного токена
decoded, err = base64.StdEncoding.DecodeString(payload)
if err != nil {
return "", fmt.Errorf("decode JWT payload: %w", err)
}
}
var claims struct {
Sub string `json:"sub"`
Exp int64 `json:"exp"`
}
if err := json.Unmarshal(decoded, &claims); err != nil {
return "", fmt.Errorf("parse JWT claims: %w", err)
}
if claims.Sub == "" {
return "", fmt.Errorf("JWT missing 'sub' claim")
}
if claims.Exp > 0 && claims.Exp < time.Now().Unix() {
return "", fmt.Errorf("JWT token expired")
}
return claims.Sub, nil
}
// NamespaceFromSub вычисляет имя k8s namespace из JWT subject (sub claim).
// Алгоритм: SHA256(sub) → берём первые 8 байт → hex → "sless-{16 hex символов}".
// Итоговая длина: 6 + 16 = 22 символа — укладывается в лимит k8s (63 символа).
// SHA256 необратим — sub пользователя не раскрывается через имя namespace.
// Детерминирован: один и тот же sub всегда даёт один и тот же namespace.
func NamespaceFromSub(sub string) string {
hash := sha256.Sum256([]byte(sub))
return fmt.Sprintf("sless-%x", hash[:8])
}
// PingNubesAPI делает GET запрос к nubes API для проверки валидности токена.
// endpoint — базовый URL nubes API (например "https://deck-api.ngcloud.ru/api/v1").
// Логика проверки:
// - 401 или 403 → токен невалиден или истёк → возвращаем ошибку
// - ошибка соединения → API недоступен → возвращаем ошибку
// - любой другой HTTP статус → API ответил, токен не отклонён → OK
func PingNubesAPI(ctx context.Context, endpoint, token string) error {
req, err := http.NewRequestWithContext(ctx, http.MethodGet, endpoint, nil)
if err != nil {
return fmt.Errorf("build nubes ping request: %w", err)
}
req.Header.Set("Authorization", "Bearer "+token)
c := &http.Client{Timeout: 10 * time.Second}
resp, err := c.Do(req)
if err != nil {
return fmt.Errorf("nubes API unreachable at %s: %w", endpoint, err)
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusUnauthorized || resp.StatusCode == http.StatusForbidden {
return fmt.Errorf("nubes API rejected token (HTTP %d) — check api_token", resp.StatusCode)
}
return nil
}
// --- JSON-структуры (зеркало handler/functions.go и handler/triggers.go) ---
// FunctionRequest — тело POST/PUT /v1/namespaces/{ns}/functions[/{name}]
type FunctionRequest struct {
Name string `json:"name"`
Runtime string `json:"runtime"`
Entrypoint string `json:"entrypoint,omitempty"`
MemoryMB int32 `json:"memory_mb,omitempty"`
TimeoutSec int32 `json:"timeout_sec,omitempty"`
Env map[string]string `json:"env_vars,omitempty"`
}
// FunctionResponse — ответ GET /v1/namespaces/{ns}/functions/{name}
type FunctionResponse struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
Runtime string `json:"runtime"`
Entrypoint string `json:"entrypoint"`
MemoryMB int32 `json:"memory_mb"`
TimeoutSec int32 `json:"timeout_sec"`
Env map[string]string `json:"env_vars"`
S3Bucket string `json:"s3_bucket"`
S3Key string `json:"s3_key"`
Phase string `json:"phase"`
ImageRef string `json:"image_ref"`
Message string `json:"message"`
}
// TriggerRequest — тело POST /v1/namespaces/{ns}/triggers
type TriggerRequest struct {
Name string `json:"name"`
Type string `json:"type"`
FunctionRef string `json:"function"`
Schedule string `json:"schedule,omitempty"`
// Enabled: nil = не передавать (по умолчанию true)
Enabled *bool `json:"enabled"`
}
// TriggerUpdateRequest — тело PATCH /v1/namespaces/{ns}/triggers/{name}
// Используем *bool чтобы различать nil (не передано) от false (явно выключен).
type TriggerUpdateRequest struct {
Enabled *bool `json:"enabled"`
}
// TriggerResponse — ответ GET /v1/namespaces/{ns}/triggers/{name}
type TriggerResponse struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
Type string `json:"type"`
FunctionRef string `json:"function"`
Schedule string `json:"schedule"`
Enabled bool `json:"enabled"`
Active bool `json:"active"`
URL string `json:"url"`
Message string `json:"message"`
}
// --- Внутренний хелпер: выполнить JSON-запрос с Bearer-токеном ---
func (c *Client) doJSON(ctx context.Context, method, url string, body interface{}) (*http.Response, error) {
var bodyReader io.Reader
if body != nil {
data, err := json.Marshal(body)
if err != nil {
return nil, fmt.Errorf("marshal request: %w", err)
}
bodyReader = bytes.NewReader(data)
}
req, err := http.NewRequestWithContext(ctx, method, url, bodyReader)
if err != nil {
return nil, err
}
if body != nil {
req.Header.Set("Content-Type", "application/json")
}
req.Header.Set("Authorization", "Bearer "+c.token)
return c.httpClient.Do(req)
}
// --- Function CRUD ---
// CreateFunction — POST /v1/namespaces/{ns}/functions → 201
func (c *Client) CreateFunction(ctx context.Context, ns string, req FunctionRequest) (*FunctionResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/functions", c.endpoint, ns)
resp, err := c.doJSON(ctx, http.MethodPost, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusCreated {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("create function: status %d: %s", resp.StatusCode, body)
}
var fn FunctionResponse
return &fn, json.NewDecoder(resp.Body).Decode(&fn)
}
// GetFunction — GET /v1/namespaces/{ns}/functions/{name} → nil если 404
func (c *Client) GetFunction(ctx context.Context, ns, name string) (*FunctionResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/functions/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodGet, url, nil)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusNotFound {
return nil, nil
}
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("get function: status %d: %s", resp.StatusCode, body)
}
var fn FunctionResponse
return &fn, json.NewDecoder(resp.Body).Decode(&fn)
}
// UpdateFunction — PUT /v1/namespaces/{ns}/functions/{name} → 200
func (c *Client) UpdateFunction(ctx context.Context, ns, name string, req FunctionRequest) (*FunctionResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/functions/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodPut, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("update function: status %d: %s", resp.StatusCode, body)
}
var fn FunctionResponse
return &fn, json.NewDecoder(resp.Body).Decode(&fn)
}
// DeleteFunction — DELETE /v1/namespaces/{ns}/functions/{name} → 204
func (c *Client) DeleteFunction(ctx context.Context, ns, name string) error {
url := fmt.Sprintf("%s/v1/namespaces/%s/functions/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodDelete, url, nil)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusNoContent && resp.StatusCode != http.StatusNotFound {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("delete function: status %d: %s", resp.StatusCode, body)
}
return nil
}
// UploadCode — POST /v1/namespaces/{ns}/functions/{name}/upload (multipart, field=code)
// После вызова оператор начинает kaniko-сборку образа.
func (c *Client) UploadCode(ctx context.Context, ns, name, zipPath string) error {
f, err := os.Open(zipPath)
if err != nil {
return fmt.Errorf("open zip %q: %w", zipPath, err)
}
defer f.Close()
return c.UploadCodeReader(ctx, ns, name, filepath.Base(zipPath), f)
}
// UploadCodeReader — загружает код из произвольного io.Reader (например in-memory zip).
// filename используется только как имя файла в multipart-форме.
func (c *Client) UploadCodeReader(ctx context.Context, ns, name, filename string, r io.Reader) error {
return c.uploadCodeToURL(ctx, fmt.Sprintf("%s/v1/namespaces/%s/functions/%s/upload", c.endpoint, ns, name), filename, r)
}
// UploadJobCode — POST /v1/namespaces/{ns}/jobs/{name}/upload
// Аналогично UploadCodeReader но для FunctionJob — запускает kaniko через FunctionJob CRD.
func (c *Client) UploadJobCode(ctx context.Context, ns, name, filename string, r io.Reader) error {
return c.uploadCodeToURL(ctx, fmt.Sprintf("%s/v1/namespaces/%s/jobs/%s/upload", c.endpoint, ns, name), filename, r)
}
// uploadCodeToURL — внутренний хелпер: пакует io.Reader в multipart и POST-ит по указанному URL.
func (c *Client) uploadCodeToURL(ctx context.Context, uploadURL, filename string, r io.Reader) error {
var buf bytes.Buffer
mw := multipart.NewWriter(&buf)
fw, err := mw.CreateFormFile("code", filename)
if err != nil {
return fmt.Errorf("create form file: %w", err)
}
if _, err := io.Copy(fw, r); err != nil {
return fmt.Errorf("copy zip: %w", err)
}
mw.Close()
req, err := http.NewRequestWithContext(ctx, http.MethodPost, uploadURL, &buf)
if err != nil {
return err
}
req.Header.Set("Content-Type", mw.FormDataContentType())
req.Header.Set("Authorization", "Bearer "+c.token)
resp, err := c.httpClient.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("upload code: status %d: %s", resp.StatusCode, body)
}
return nil
}
// WaitReady опрашивает функцию каждые 5 секунд пока phase != Ready/Failed.
// Нужен после UploadCode — kaniko-сборка занимает ~1 минуту.
func (c *Client) WaitReady(ctx context.Context, ns, name string, timeout time.Duration) (*FunctionResponse, error) {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
fn, err := c.GetFunction(ctx, ns, name)
if err != nil {
return nil, err
}
if fn == nil {
return nil, fmt.Errorf("function %s/%s not found while waiting", ns, name)
}
switch fn.Phase {
case "Ready":
return fn, nil
case "Failed":
return nil, fmt.Errorf("function build failed: %s", fn.Message)
}
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(5 * time.Second):
}
}
return nil, fmt.Errorf("timeout waiting for function %s/%s to become Ready", ns, name)
}
// --- Trigger CRUD ---
// CreateTrigger — POST /v1/namespaces/{ns}/triggers → 201
func (c *Client) CreateTrigger(ctx context.Context, ns string, req TriggerRequest) (*TriggerResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/triggers", c.endpoint, ns)
resp, err := c.doJSON(ctx, http.MethodPost, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusCreated {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("create trigger: status %d: %s", resp.StatusCode, body)
}
var tr TriggerResponse
return &tr, json.NewDecoder(resp.Body).Decode(&tr)
}
// GetTrigger — GET /v1/namespaces/{ns}/triggers/{name} → nil если 404
func (c *Client) GetTrigger(ctx context.Context, ns, name string) (*TriggerResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/triggers/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodGet, url, nil)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusNotFound {
return nil, nil
}
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("get trigger: status %d: %s", resp.StatusCode, body)
}
var tr TriggerResponse
return &tr, json.NewDecoder(resp.Body).Decode(&tr)
}
// DeleteTrigger — DELETE /v1/namespaces/{ns}/triggers/{name} → 204
func (c *Client) DeleteTrigger(ctx context.Context, ns, name string) error {
url := fmt.Sprintf("%s/v1/namespaces/%s/triggers/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodDelete, url, nil)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusNoContent && resp.StatusCode != http.StatusNotFound {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("delete trigger: status %d: %s", resp.StatusCode, body)
}
return nil
}
// UpdateTrigger — PATCH /v1/namespaces/{ns}/triggers/{name} → 200
// Позволяет изменить enabled без пересоздания триггера.
func (c *Client) UpdateTrigger(ctx context.Context, ns, name string, req TriggerUpdateRequest) (*TriggerResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/triggers/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodPatch, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("update trigger: status %d: %s", resp.StatusCode, body)
}
var tr TriggerResponse
return &tr, json.NewDecoder(resp.Body).Decode(&tr)
}
// --- Job CRUD ---
// JobRequest — тело POST /v1/namespaces/{ns}/jobs
// Изменено: 2026-03-20 (merge: убран FunctionRef, добавлены Runtime/Entrypoint/Env)
type JobRequest struct {
Name string `json:"name"`
Runtime string `json:"runtime"`
Entrypoint string `json:"entrypoint"`
MemoryMB int32 `json:"memory_mb,omitempty"`
TimeoutSec int32 `json:"timeout_sec,omitempty"`
Env map[string]string `json:"env_vars,omitempty"`
S3Bucket string `json:"s3_bucket,omitempty"`
S3Key string `json:"s3_key,omitempty"`
EventJSON string `json:"event_json,omitempty"`
// RunID: 0 = создать без запуска, >0 = запустить
RunID int64 `json:"run_id"`
}
// JobResponse — ответ GET /v1/namespaces/{ns}/jobs/{name}
type JobResponse struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
Runtime string `json:"runtime"`
Entrypoint string `json:"entrypoint"`
EventJSON string `json:"event_json"`
RunID int64 `json:"run_id"`
Phase string `json:"phase"`
ImageRef string `json:"image_ref"`
JobName string `json:"job_name"`
StartTime string `json:"start_time"`
CompletionTime string `json:"completion_time"`
Message string `json:"message"`
}
// CreateJob — POST /v1/namespaces/{ns}/jobs → 201
func (c *Client) CreateJob(ctx context.Context, ns string, req JobRequest) (*JobResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/jobs", c.endpoint, ns)
resp, err := c.doJSON(ctx, http.MethodPost, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusCreated {
body, _ := io.ReadAll(resp.Body)
if resp.StatusCode == http.StatusConflict {
return nil, fmt.Errorf("%w: %s", ErrJobAlreadyExists, strings.TrimSpace(string(body)))
}
return nil, fmt.Errorf("create job: status %d: %s", resp.StatusCode, body)
}
var j JobResponse
return &j, json.NewDecoder(resp.Body).Decode(&j)
}
// GetJob — GET /v1/namespaces/{ns}/jobs/{name} → nil если 404
func (c *Client) GetJob(ctx context.Context, ns, name string) (*JobResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/jobs/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodGet, url, nil)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusNotFound {
return nil, nil
}
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("get job: status %d: %s", resp.StatusCode, body)
}
var j JobResponse
return &j, json.NewDecoder(resp.Body).Decode(&j)
}
// DeleteJob — DELETE /v1/namespaces/{ns}/jobs/{name} → 204
func (c *Client) DeleteJob(ctx context.Context, ns, name string) error {
url := fmt.Sprintf("%s/v1/namespaces/%s/jobs/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodDelete, url, nil)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusNoContent && resp.StatusCode != http.StatusNotFound {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("delete job: status %d: %s", resp.StatusCode, body)
}
return nil
}
// --- IoTDevice CRUD ---
// Изменено: 2026-04-04
// IoTDeviceRequest — тело POST /v1/namespaces/{ns}/iot/devices
type IoTDeviceRequest struct {
Name string `json:"name"`
DeviceID string `json:"device_id"`
Enabled *bool `json:"enabled"`
Metadata map[string]string `json:"metadata,omitempty"`
}
// IoTDeviceUpdateRequest — тело PATCH /v1/namespaces/{ns}/iot/devices/{name}
type IoTDeviceUpdateRequest struct {
Enabled *bool `json:"enabled"`
}
// IoTDeviceResponse — ответ GET /v1/namespaces/{ns}/iot/devices/{name}
// mqtt_password заполняется только если устройство в фазе Active.
type IoTDeviceResponse struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
DeviceID string `json:"device_id"`
Enabled bool `json:"enabled"`
Phase string `json:"phase"`
MQTTUsername string `json:"mqtt_username"`
MQTTPassword string `json:"mqtt_password"`
SecretName string `json:"secret_name"`
TopicPrefix string `json:"topic_prefix"`
LastConnected string `json:"last_connected,omitempty"`
Message string `json:"message,omitempty"`
Metadata map[string]string `json:"metadata,omitempty"`
CreatedAt string `json:"created_at,omitempty"`
}
// CreateIoTDevice — POST /v1/namespaces/{ns}/iot/devices → 201
func (c *Client) CreateIoTDevice(ctx context.Context, ns string, req IoTDeviceRequest) (*IoTDeviceResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/iot/devices", c.endpoint, ns)
resp, err := c.doJSON(ctx, http.MethodPost, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusCreated {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("create iot device: status %d: %s", resp.StatusCode, body)
}
var d IoTDeviceResponse
return &d, json.NewDecoder(resp.Body).Decode(&d)
}
// GetIoTDevice — GET /v1/namespaces/{ns}/iot/devices/{name} → nil если 404
// Возвращает mqtt_password из Secret. После создания ждать phase=Active.
func (c *Client) GetIoTDevice(ctx context.Context, ns, name string) (*IoTDeviceResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/iot/devices/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodGet, url, nil)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusNotFound {
return nil, nil
}
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("get iot device: status %d: %s", resp.StatusCode, body)
}
var d IoTDeviceResponse
return &d, json.NewDecoder(resp.Body).Decode(&d)
}
// DeleteIoTDevice — DELETE /v1/namespaces/{ns}/iot/devices/{name} → 204
func (c *Client) DeleteIoTDevice(ctx context.Context, ns, name string) error {
url := fmt.Sprintf("%s/v1/namespaces/%s/iot/devices/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodDelete, url, nil)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusNoContent && resp.StatusCode != http.StatusNotFound {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("delete iot device: status %d: %s", resp.StatusCode, body)
}
return nil
}
// UpdateIoTDevice — PATCH /v1/namespaces/{ns}/iot/devices/{name} → 200
// Позволяет изменить enabled без пересоздания устройства.
func (c *Client) UpdateIoTDevice(ctx context.Context, ns, name string, req IoTDeviceUpdateRequest) (*IoTDeviceResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/iot/devices/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodPatch, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("update iot device: status %d: %s", resp.StatusCode, body)
}
var d IoTDeviceResponse
return &d, json.NewDecoder(resp.Body).Decode(&d)
}
// WaitIoTDeviceActive опрашивает устройство каждые 2 секунды пока phase != Active/Error/Disabled.
// Нужен после CreateIoTDevice — контроллер асинхронно создаёт Secret с credentials.
func (c *Client) WaitIoTDeviceActive(ctx context.Context, ns, name string, timeout time.Duration) (*IoTDeviceResponse, error) {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
d, err := c.GetIoTDevice(ctx, ns, name)
if err != nil {
return nil, err
}
if d == nil {
return nil, fmt.Errorf("iot device %s/%s not found while waiting", ns, name)
}
switch d.Phase {
case "Active":
return d, nil
case "Error":
return nil, fmt.Errorf("iot device provisioning failed: %s", d.Message)
case "Disabled":
// enabled=false при создании — credentials не создаются, это ожидаемо
return d, nil
}
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(2 * time.Second):
}
}
return nil, fmt.Errorf("timeout waiting for iot device %s/%s to become Active", ns, name)
}
// EnsureNamespace — POST /v1/namespaces/{ns}/ensure
// Создаёт k8s namespace пользователя если не существует. Идемпотентен.
// Вызывается ОДИН РАЗ из provider.Configure() до любых ресурсных операций.
// 200 OK = namespace уже был, 201 Created = создан сейчас, остальное = ошибка.
func (c *Client) EnsureNamespace(ctx context.Context, ns string) error {
url := fmt.Sprintf("%s/v1/namespaces/%s/ensure", c.endpoint, ns)
resp, err := c.doJSON(ctx, http.MethodPost, url, nil)
if err != nil {
return fmt.Errorf("ensure namespace: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK && resp.StatusCode != http.StatusCreated {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("ensure namespace: status %d: %s", resp.StatusCode, body)
}
return nil
}
// WaitJobDone опрашивает job каждые 5 секунд пока phase не Succeeded или Failed.
// Блокирует terraform apply до завершения джоба.
func (c *Client) WaitJobDone(ctx context.Context, ns, name string, timeout time.Duration) (*JobResponse, error) {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
j, err := c.GetJob(ctx, ns, name)
if err != nil {
return nil, err
}
if j == nil {
return nil, fmt.Errorf("job %s/%s not found while waiting", ns, name)
}
switch j.Phase {
case "Succeeded":
return j, nil
case "Failed":
return nil, fmt.Errorf("job failed: %s", j.Message)
}
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(5 * time.Second):
}
}
return nil, fmt.Errorf("timeout waiting for job %s/%s to complete", ns, name)
}
// --- Service CRUD (sless_service — long-running Deployment + URL) ---
// ServiceRequest — тело POST/PUT /v1/namespaces/{ns}/services[/{name}]
type ServiceRequest struct {
Name string `json:"name"`
Runtime string `json:"runtime"`
Entrypoint string `json:"entrypoint,omitempty"`
MemoryMB int32 `json:"memory_mb,omitempty"`
TimeoutSec int32 `json:"timeout_sec,omitempty"`
Env map[string]string `json:"env_vars,omitempty"`
}
// ServiceResponse — ответ GET /v1/namespaces/{ns}/services/{name}
// URL — ключевое поле: заполняется оператором после деплоя Deployment+Ingress.
type ServiceResponse struct {
Name string `json:"name"`
Namespace string `json:"namespace"`
Runtime string `json:"runtime"`
Entrypoint string `json:"entrypoint"`
MemoryMB int32 `json:"memory_mb"`
TimeoutSec int32 `json:"timeout_sec"`
Env map[string]string `json:"env_vars"`
S3Bucket string `json:"s3_bucket"`
S3Key string `json:"s3_key"`
Phase string `json:"phase"`
ImageRef string `json:"image_ref"`
URL string `json:"url"`
Message string `json:"message"`
}
// CreateService — POST /v1/namespaces/{ns}/services → 201
func (c *Client) CreateService(ctx context.Context, ns string, req ServiceRequest) (*ServiceResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/services", c.endpoint, ns)
resp, err := c.doJSON(ctx, http.MethodPost, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusCreated {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("create service: status %d: %s", resp.StatusCode, body)
}
var svc ServiceResponse
return &svc, json.NewDecoder(resp.Body).Decode(&svc)
}
// GetService — GET /v1/namespaces/{ns}/services/{name} → nil если 404
func (c *Client) GetService(ctx context.Context, ns, name string) (*ServiceResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/services/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodGet, url, nil)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode == http.StatusNotFound {
return nil, nil
}
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("get service: status %d: %s", resp.StatusCode, body)
}
var svc ServiceResponse
return &svc, json.NewDecoder(resp.Body).Decode(&svc)
}
// UpdateService — PUT /v1/namespaces/{ns}/services/{name} → 200
func (c *Client) UpdateService(ctx context.Context, ns, name string, req ServiceRequest) (*ServiceResponse, error) {
url := fmt.Sprintf("%s/v1/namespaces/%s/services/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodPut, url, req)
if err != nil {
return nil, err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return nil, fmt.Errorf("update service: status %d: %s", resp.StatusCode, body)
}
var svc ServiceResponse
return &svc, json.NewDecoder(resp.Body).Decode(&svc)
}
// DeleteService — DELETE /v1/namespaces/{ns}/services/{name} → 204
func (c *Client) DeleteService(ctx context.Context, ns, name string) error {
url := fmt.Sprintf("%s/v1/namespaces/%s/services/%s", c.endpoint, ns, name)
resp, err := c.doJSON(ctx, http.MethodDelete, url, nil)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusNoContent && resp.StatusCode != http.StatusNotFound {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("delete service: status %d: %s", resp.StatusCode, body)
}
return nil
}
// UploadServiceCode — POST /v1/namespaces/{ns}/services/{name}/upload (multipart, field=code)
// После вызова оператор начинает kaniko-сборку и затем деплоит Deployment.
func (c *Client) UploadServiceCode(ctx context.Context, ns, name, zipPath string) error {
f, err := os.Open(zipPath)
if err != nil {
return fmt.Errorf("open zip %q: %w", zipPath, err)
}
defer f.Close()
return c.UploadServiceCodeReader(ctx, ns, name, filepath.Base(zipPath), f)
}
// UploadServiceCodeReader — загружает код сервиса из произвольного io.Reader.
func (c *Client) UploadServiceCodeReader(ctx context.Context, ns, name, filename string, r io.Reader) error {
var buf bytes.Buffer
mw := multipart.NewWriter(&buf)
fw, err := mw.CreateFormFile("code", filename)
if err != nil {
return fmt.Errorf("create form file: %w", err)
}
if _, err := io.Copy(fw, r); err != nil {
return fmt.Errorf("copy zip: %w", err)
}
mw.Close()
url := fmt.Sprintf("%s/v1/namespaces/%s/services/%s/upload", c.endpoint, ns, name)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, &buf)
if err != nil {
return err
}
req.Header.Set("Content-Type", mw.FormDataContentType())
req.Header.Set("Authorization", "Bearer "+c.token)
resp, err := c.httpClient.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
body, _ := io.ReadAll(resp.Body)
return fmt.Errorf("upload service code: status %d: %s", resp.StatusCode, body)
}
return nil
}
// WaitServiceReady опрашивает сервис каждые 5 секунд пока phase != Ready/Failed.
// Нужен после UploadServiceCode — kaniko-сборка + деплой занимают ~1-2 минуты.
func (c *Client) WaitServiceReady(ctx context.Context, ns, name string, timeout time.Duration) (*ServiceResponse, error) {
deadline := time.Now().Add(timeout)
for time.Now().Before(deadline) {
svc, err := c.GetService(ctx, ns, name)
if err != nil {
return nil, err
}
if svc == nil {
return nil, fmt.Errorf("service %s/%s not found while waiting", ns, name)
}
switch svc.Phase {
case "Ready":
return svc, nil
case "Failed":
return nil, fmt.Errorf("service build failed: %s", svc.Message)
}
select {
case <-ctx.Done():
return nil, ctx.Err()
case <-time.After(5 * time.Second):
}
}
return nil, fmt.Errorf("timeout waiting for service %s/%s to become Ready", ns, name)
}