diff --git a/universal_rebuild/internal/core/client.go b/universal_rebuild/internal/core/client.go new file mode 100644 index 0000000..71288bf --- /dev/null +++ b/universal_rebuild/internal/core/client.go @@ -0,0 +1,1334 @@ +// ВНИМАНИЕ: НЕ ИЗМЕНЯТЬ НИЧЕГО В ЯДРЕ БЕЗ ПРЯМОГО РАЗРЕШЕНИЯ ОПЕРАТОРА. +package core + +import ( + "bytes" + "context" + "encoding/json" + "errors" + "fmt" + "io" + "net/http" + "net/url" + "os" + "regexp" + "strings" + "time" +) + +// UniversalClient handles Nubes API logic. +type UniversalClient struct { + HttpClient *http.Client + ApiEndpoint string + ApiToken string + ProviderVersion string + OperationTimeouts *OperationTimeouts + // LogLevel задаёт уровень вывода этапов операции: ""/"none" | "info" | "debug". + // Может быть переопределён на уровне ресурса через context (ключ ctxKeyLogLevel). + LogLevel string +} + +// ctxKeyLogLevel — ключ для переопределения LogLevel на уровне ресурса через context.WithValue. +type ctxKeyLogLevelType struct{} + +var ctxKeyLogLevel = ctxKeyLogLevelType{} + +// CtxWithLogLevel возвращает ctx с переопределённым уровнем логирования. +func CtxWithLogLevel(ctx context.Context, level string) context.Context { + return context.WithValue(ctx, ctxKeyLogLevel, level) +} + +// reTimestamp удаляет временны́е метки вида [2026-03-24T05:07:47.716Z +0.860s] из строк. +var reTimestamp = regexp.MustCompile(`\[[0-9]{4}-[0-9]{2}-[0-9]{2}T[^\]]+\]\s*`) + +// formatStageMsg парсит stageMsg (JSON-массив пар [имя, текст]) и возвращает +// отформатированные строки для вывода. Временны́е метки убираются. +func formatStageMsg(raw string) []string { + raw = strings.TrimSpace(raw) + if raw == "" { + return nil + } + // Пытаемся распарсить как [["name","text"],...] + var pairs [][]string + if err := json.Unmarshal([]byte(raw), &pairs); err != nil || len(pairs) == 0 { + return nil + } + // Вычисляем максимальную ширину имени для выравнивания + maxLen := 0 + for _, p := range pairs { + if len(p) >= 1 && len(p[0]) > maxLen { + maxLen = len(p[0]) + } + } + var result []string + for _, p := range pairs { + if len(p) < 2 { + continue + } + name := p[0] + text := reTimestamp.ReplaceAllString(p[1], "") + lines := strings.Split(strings.TrimSpace(text), "\n") + padding := strings.Repeat(" ", maxLen-len(name)) + for i, line := range lines { + line = strings.TrimSpace(line) + if line == "" { + continue + } + if i == 0 { + result = append(result, fmt.Sprintf(" %s:%s %s", name, padding, line)) + } else { + result = append(result, fmt.Sprintf(" %s %s", strings.Repeat(" ", maxLen), line)) + } + } + } + return result +} + +type genericInstanceReq struct { + ServiceId int `json:"serviceId"` + DisplayName string `json:"displayName"` + Descr string `json:"descr"` +} + +type genericOpReq struct { + InstanceUid string `json:"instanceUid"` + Operation string `json:"operation"` +} + +type genericParamReq struct { + InstanceOperationUid string `json:"instanceOperationUid"` + SvcOperationCfsParamId int `json:"svcOperationCfsParamId"` + ParamValue string `json:"paramValue"` +} + +// ===== UNIVERSAL FLOW V6 (baseline) ===== + +type universalOpResponse struct { + InstanceOperation universalOperation `json:"instanceOperation"` +} + +type universalOperation struct { + CfsParams []universalCfsParam `json:"cfsParams"` +} + +type universalCfsParam struct { + SvcOperationCfsParamId int `json:"svcOperationCfsParamId"` + ParamValue *string `json:"paramValue"` + DefaultValue *string `json:"defaultValue"` + IsRequired bool `json:"isRequired"` + DataType string `json:"dataType"` + Name string `json:"name"` + Code string `json:"code"` + Label string `json:"label"` + SvcOperationCfsParam string `json:"svcOperationCfsParam"` + RefSvcId *int `json:"refSvcId"` +} + +// CreateGenericInstanceUniversalV6 implements the universal flow: +// instances -> instanceOperations -> get cfsParams -> submit params -> validate -> run +func (c *UniversalClient) CreateGenericInstanceUniversalV6(ctx context.Context, serviceId int, displayName string, params map[int]string) (string, error) { + instPayload := genericInstanceReq{ + ServiceId: serviceId, + DisplayName: displayName, + Descr: c.instanceDescr(), + } + + instResp, instHeaders, err := c.doRequest(ctx, "POST", "/instances", instPayload) + if err != nil { + return "", err + } + + instanceUid := extractUIDFromLocation(instHeaders.Get("Location")) + if instanceUid == "" { + var instResult struct { + InstanceUid string `json:"instanceUid"` + } + if err := json.Unmarshal(instResp, &instResult); err == nil && instResult.InstanceUid != "" { + instanceUid = instResult.InstanceUid + } else { + var justId string + if err2 := json.Unmarshal(instResp, &justId); err2 == nil && justId != "" { + instanceUid = justId + } + } + } + if instanceUid == "" { + return "", fmt.Errorf("не удалось извлечь instanceUid из ответа (Header: %s)", instHeaders.Get("Location")) + } + + opPayload := genericOpReq{ + InstanceUid: instanceUid, + Operation: "create", + } + + opResp, opHeaders, err := c.doRequest(ctx, "POST", "/instanceOperations", opPayload) + if err != nil { + return "", err + } + + opUid := extractUIDFromLocation(opHeaders.Get("Location")) + if opUid == "" { + var opResult struct { + InstanceOperationUid string `json:"instanceOperationUid"` + } + if err := json.Unmarshal(opResp, &opResult); err == nil && opResult.InstanceOperationUid != "" { + opUid = opResult.InstanceOperationUid + } else { + var justId string + if err2 := json.Unmarshal(opResp, &justId); err2 == nil && justId != "" { + opUid = justId + } + } + } + if opUid == "" { + return "", fmt.Errorf("не удалось извлечь instanceOperationUid (Header: %s)", opHeaders.Get("Location")) + } + + opDetailsResp, _, err := c.doRequest(ctx, "GET", fmt.Sprintf("/instanceOperations/%s?fields=cfsParams", opUid), nil) + if err != nil { + return "", fmt.Errorf("не удалось получить детали операции: %w", err) + } + var opDetails universalOpResponse + if err := json.Unmarshal(opDetailsResp, &opDetails); err != nil { + return "", fmt.Errorf("не удалось разобрать детали операции: %w", err) + } + + params, err = c.resolveRefSvcParamValues(ctx, opDetails.InstanceOperation.CfsParams, params) + if err != nil { + return "", err + } + + sent := make(map[int]bool) + for paramId, value := range params { + pPayload := genericParamReq{ + InstanceOperationUid: opUid, + SvcOperationCfsParamId: paramId, + ParamValue: value, + } + _, _, err := c.doRequest(ctx, "POST", "/instanceOperationCfsParams", pPayload) + if err != nil { + return "", fmt.Errorf("не удалось установить параметр %d: %w", paramId, err) + } + sent[paramId] = true + } + + for _, param := range opDetails.InstanceOperation.CfsParams { + if sent[param.SvcOperationCfsParamId] { + continue + } + if !param.IsRequired { + continue + } + + val := "" + if param.ParamValue != nil { + val = *param.ParamValue + } else if param.DefaultValue != nil { + val = *param.DefaultValue + } + val = normalizeUniversalValueV6(val, param) + + pPayload := genericParamReq{ + InstanceOperationUid: opUid, + SvcOperationCfsParamId: param.SvcOperationCfsParamId, + ParamValue: val, + } + _, _, err := c.doRequest(ctx, "POST", "/instanceOperationCfsParams", pPayload) + if err != nil { + return "", fmt.Errorf("не удалось отправить параметр по умолчанию %d: %w", param.SvcOperationCfsParamId, err) + } + } + + _, _, err = c.doRequest(ctx, "GET", fmt.Sprintf("/instanceOperations/%s/validate-cfs", opUid), nil) + if err != nil { + return "", fmt.Errorf("валидация не пройдена: %w", err) + } + + _, _, err = c.doRequest(ctx, "POST", fmt.Sprintf("/instanceOperations/%s/run", opUid), map[string]interface{}{}) + if err != nil { + return "", fmt.Errorf("выполнение не удалось: %w", err) + } + + // НЕ МЕНЯТЬ: завершение операции определяется по dtFinish + if err := c.waitForOperationFinish(ctx, opUid, c.operationTimeoutForContext(ctx, serviceId, "create")); err != nil { + return "", err + } + if err := c.ensureInstanceCreated(ctx, instanceUid); err != nil { + return "", err + } + + return instanceUid, nil +} + +func normalizeUniversalValueV6(val string, param universalCfsParam) string { + trimmed := strings.TrimSpace(val) + if strings.EqualFold(trimmed, "null") { + trimmed = "" + } + if trimmed == "\"\"" { + trimmed = "" + } + if trimmed != "" { + return trimmed + } + + dataType := strings.ToLower(param.DataType) + nameHint := strings.ToLower(param.Name + " " + param.Code + " " + param.Label + " " + param.SvcOperationCfsParam) + + if strings.Contains(dataType, "array") || strings.Contains(nameHint, "array") || strings.Contains(nameHint, "list") { + return "[]" + } + if strings.Contains(dataType, "map") || strings.Contains(dataType, "json") || strings.Contains(nameHint, "map") || strings.Contains(nameHint, "json") { + return "{}" + } + + return trimmed +} + +// RunInstanceOperationUniversal runs an available operation (modify/suspend/delete/resume) if possible. +func (c *UniversalClient) RunInstanceOperationUniversal(ctx context.Context, instanceUid string, action string, params map[int]string) error { + state, err := c.GetInstanceState(ctx, instanceUid) + if err != nil { + return err + } + if state.OperationIsPending || state.OperationIsInProgress { + if err := c.waitForInstanceIdle(ctx, instanceUid, c.idleTimeoutFor(state.ServiceId)); err != nil { + return err + } + } + + var opId int + for _, op := range state.AvailableOperations { + if strings.EqualFold(op.Operation, action) { + opId = op.SvcOperationId + break + } + } + if opId == 0 { + return fmt.Errorf("операция %s недоступна для экземпляра %s", action, instanceUid) + } + + payload := map[string]interface{}{ + "instanceUid": instanceUid, + "svcOperationId": opId, + "operation": action, + } + + opUid, err := c.postIgnoreResponse(ctx, "/instanceOperations", payload, true) + if err != nil { + return fmt.Errorf("не удалось создать операцию %s: %w", action, err) + } + if opUid == "" { + return fmt.Errorf("не удалось получить UID операции для %s", action) + } + + for paramId, value := range params { + pPayload := genericParamReq{ + InstanceOperationUid: opUid, + SvcOperationCfsParamId: paramId, + ParamValue: value, + } + _, _, err := c.doRequest(ctx, "POST", "/instanceOperationCfsParams", pPayload) + if err != nil { + return fmt.Errorf("не удалось установить параметр %d: %w", paramId, err) + } + } + + _, _, err = c.doRequest(ctx, "POST", fmt.Sprintf("/instanceOperations/%s/run", opUid), map[string]interface{}{}) + if err != nil { + return err + } + + // НЕ МЕНЯТЬ: завершение операции определяется по dtFinish + return c.waitForOperationFinish(ctx, opUid, c.operationTimeoutForContext(ctx, state.ServiceId, action)) +} + +// RunInstanceOperationUniversalWithDefaults runs an operation and submits required params (including defaults). +func (c *UniversalClient) RunInstanceOperationUniversalWithDefaults(ctx context.Context, instanceUid string, action string, params map[int]string) error { + state, err := c.GetInstanceState(ctx, instanceUid) + if err != nil { + return err + } + if state.OperationIsPending || state.OperationIsInProgress { + if err := c.waitForInstanceIdle(ctx, instanceUid, c.idleTimeoutFor(state.ServiceId)); err != nil { + return err + } + } + + var opId int + for _, op := range state.AvailableOperations { + if strings.EqualFold(op.Operation, action) { + opId = op.SvcOperationId + break + } + } + if opId == 0 { + return fmt.Errorf("операция %s недоступна для экземпляра %s", action, instanceUid) + } + + payload := map[string]interface{}{ + "instanceUid": instanceUid, + "svcOperationId": opId, + "operation": action, + } + + opUid, err := c.postIgnoreResponse(ctx, "/instanceOperations", payload, true) + if err != nil { + return fmt.Errorf("не удалось создать операцию %s: %w", action, err) + } + if opUid == "" { + return fmt.Errorf("не удалось получить UID операции для %s", action) + } + + opDetailsResp, _, err := c.doRequest(ctx, "GET", fmt.Sprintf("/instanceOperations/%s?fields=cfsParams", opUid), nil) + if err != nil { + return fmt.Errorf("не удалось получить детали операции: %w", err) + } + var opDetails universalOpResponse + if err := json.Unmarshal(opDetailsResp, &opDetails); err != nil { + return fmt.Errorf("не удалось разобрать детали операции: %w", err) + } + + params, err = c.resolveRefSvcParamValues(ctx, opDetails.InstanceOperation.CfsParams, params) + if err != nil { + return err + } + + sent := make(map[int]bool) + for paramId, value := range params { + pPayload := genericParamReq{ + InstanceOperationUid: opUid, + SvcOperationCfsParamId: paramId, + ParamValue: value, + } + _, _, err := c.doRequest(ctx, "POST", "/instanceOperationCfsParams", pPayload) + if err != nil { + return fmt.Errorf("не удалось установить параметр %d: %w", paramId, err) + } + sent[paramId] = true + } + + for _, param := range opDetails.InstanceOperation.CfsParams { + if sent[param.SvcOperationCfsParamId] { + continue + } + + val := "" + if param.ParamValue != nil { + val = *param.ParamValue + } else if param.DefaultValue != nil { + val = *param.DefaultValue + } + val = normalizeUniversalValueV6(val, param) + + pPayload := genericParamReq{ + InstanceOperationUid: opUid, + SvcOperationCfsParamId: param.SvcOperationCfsParamId, + ParamValue: val, + } + _, _, err := c.doRequest(ctx, "POST", "/instanceOperationCfsParams", pPayload) + if err != nil { + return fmt.Errorf("не удалось отправить параметр по умолчанию %d: %w", param.SvcOperationCfsParamId, err) + } + } + + _, _, err = c.doRequest(ctx, "GET", fmt.Sprintf("/instanceOperations/%s/validate-cfs", opUid), nil) + if err != nil { + return fmt.Errorf("валидация не пройдена: %w", err) + } + + _, _, err = c.doRequest(ctx, "POST", fmt.Sprintf("/instanceOperations/%s/run", opUid), map[string]interface{}{}) + if err != nil { + return err + } + + // НЕ МЕНЯТЬ: завершение операции определяется по dtFinish + return c.waitForOperationFinish(ctx, opUid, c.operationTimeoutForContext(ctx, state.ServiceId, action)) +} + +// Instance state structures + +type ApiOperation struct { + SvcOperationId int `json:"svcOperationId"` + Operation string `json:"operation"` +} + +type InstanceStateResponse struct { + InstanceUid string `json:"instanceUid"` + ServiceId int `json:"serviceId"` + ExplainedStatus string `json:"explainedStatus"` + IsDeleted bool `json:"isDeleted"` + OperationIsInProgress bool `json:"operationIsInProgress"` + OperationIsPending bool `json:"operationIsPending"` + AvailableOperations []ApiOperation `json:"availableOperations"` +} + +// FindInstanceByDisplayName finds an instance by display_name for a given serviceId. +// Если найдено больше одного non-deleted инстанса — возвращает ошибку. +func (c *UniversalClient) FindInstanceByDisplayName(ctx context.Context, serviceId int, displayName string) (*InstanceStateResponse, error) { + var found []InstanceStateResponse + + // Быстрый поиск через search API + searchQuery := url.Values{} + searchQuery.Set("fields", "instanceUid,displayName,serviceId") + searchQuery.Set("page", "1") + searchQuery.Set("pageSize", "100") + searchQuery.Set("search", displayName) + searchQuery.Set("isAuxiliary", "false") + searchQuery.Set("isDeleted", "false") + if serviceId > 0 { + searchQuery.Set("serviceId", fmt.Sprintf("%d", serviceId)) + } + respBody, _, err := c.doRequest(ctx, "GET", "/instances?"+searchQuery.Encode(), nil) + if err == nil { + var res struct { + Results []struct { + InstanceUid string `json:"instanceUid"` + DisplayName string `json:"displayName"` + ServiceId int `json:"serviceId"` + } `json:"results"` + } + if err := json.Unmarshal(respBody, &res); err == nil { + for _, item := range res.Results { + if item.ServiceId == serviceId && strings.EqualFold(item.DisplayName, displayName) { + state, err := c.GetInstanceStateRaw(ctx, item.InstanceUid) + if err != nil { + continue + } + if state != nil && !isInstanceDeleted(state) { + found = append(found, *state) + } + } + } + } + } + + // Если быстрый поиск не дал результатов — пагинированный fallback + if len(found) == 0 { + page := 1 + for { + reqURL := fmt.Sprintf("%s/instances?page=%d&size=100", c.ApiEndpoint, page) + req, err := http.NewRequestWithContext(ctx, "GET", reqURL, nil) + if err != nil { + return nil, err + } + if c.ApiToken != "" { + req.Header.Set("Authorization", "Bearer "+c.ApiToken) + } + + resp, err := c.HttpClient.Do(req) + if err != nil { + return nil, err + } + if resp.Body != nil { + defer resp.Body.Close() + } + + if resp.StatusCode != 200 { + return nil, fmt.Errorf("HTTP статус %d", resp.StatusCode) + } + + var res struct { + Results []struct { + InstanceUid string `json:"instanceUid"` + DisplayName string `json:"displayName"` + ServiceId int `json:"serviceId"` + } `json:"results"` + } + if err := json.NewDecoder(resp.Body).Decode(&res); err != nil { + return nil, err + } + + if len(res.Results) == 0 { + break + } + + for _, item := range res.Results { + if item.ServiceId == serviceId && strings.EqualFold(item.DisplayName, displayName) { + state, err := c.GetInstanceStateRaw(ctx, item.InstanceUid) + if err != nil { + continue + } + if state != nil && !isInstanceDeleted(state) { + found = append(found, *state) + } + } + } + + page++ + if page > 100 { + break + } + } + } + + if len(found) == 0 { + return nil, nil + } + if len(found) == 1 { + return &found[0], nil + } + + // Множественные non-deleted инстансы — ошибка + details := make([]string, 0, len(found)) + for _, inst := range found { + details = append(details, fmt.Sprintf(" - %s (статус: %s)", inst.InstanceUid, inst.ExplainedStatus)) + } + return nil, fmt.Errorf( + "обнаружено %d инстансов с именем '%s' (serviceId=%d):\n%s\nНевозможно определить какой adopt-ить. Удалите лишние через ЛК (Управление облаком) или укажите конкретный UUID через 'terraform import'", + len(found), displayName, serviceId, strings.Join(details, "\n"), + ) +} + +func (c *UniversalClient) GetInstanceState(ctx context.Context, instanceUid string) (*InstanceStateResponse, error) { + url := fmt.Sprintf("%s/instances/%s", c.ApiEndpoint, instanceUid) + req, err := http.NewRequestWithContext(ctx, "GET", url, nil) + if err != nil { + return nil, err + } + if c.ApiToken != "" { + req.Header.Set("Authorization", "Bearer "+c.ApiToken) + } + + resp, err := c.HttpClient.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != 200 { + return nil, fmt.Errorf("HTTP статус %d", resp.StatusCode) + } + + var res struct { + Instance InstanceStateResponse `json:"instance"` + } + if err := json.NewDecoder(resp.Body).Decode(&res); err != nil { + return nil, err + } + + if err := validateInstanceStatus(&res.Instance); err != nil { + return nil, err + } + + return &res.Instance, nil +} + +func isInstanceDeleted(state *InstanceStateResponse) bool { + if state == nil { + return false + } + if state.IsDeleted { + return true + } + return strings.EqualFold(strings.TrimSpace(state.ExplainedStatus), "deleted") +} + +// GetInstanceStateRaw получает состояние инстанса БЕЗ валидации статуса. +// Используется для проверки ref-параметров: нужно читать даже deleted/suspended инстансы. +func (c *UniversalClient) GetInstanceStateRaw(ctx context.Context, instanceUid string) (*InstanceStateResponse, error) { + reqURL := fmt.Sprintf("%s/instances/%s", c.ApiEndpoint, instanceUid) + req, err := http.NewRequestWithContext(ctx, "GET", reqURL, nil) + if err != nil { + return nil, err + } + if c.ApiToken != "" { + req.Header.Set("Authorization", "Bearer "+c.ApiToken) + } + + resp, err := c.HttpClient.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + + if resp.StatusCode != 200 { + return nil, fmt.Errorf("HTTP статус %d", resp.StatusCode) + } + + var res struct { + Instance InstanceStateResponse `json:"instance"` + } + if err := json.NewDecoder(resp.Body).Decode(&res); err != nil { + return nil, err + } + + return &res.Instance, nil +} + +// ===== UNIVERSAL OPERATION WAIT (APPEND-ONLY) ===== +// +// КРИТИЧЕСКИ ВАЖНО: +// 1. Критерий завершения операции — наличие dtFinish. +// 2. Эта логика едина для всех сервисов/операций и является контрактом поведения. +// 3. ЗАПРЕЩЕНО ПРАВИТЬ ЭТОТ КОД БЕЗ ЯВНОГО СОГЛАСОВАНИЯ С ОПЕРАТОРОМ. +// Любые изменения (таймауты, критерии завершения, частота опроса, обработка ошибок) +// должны быть согласованы заранее. +const defaultOperationTimeout = 30 * time.Minute + +// ttyOut открывает /dev/tty для прямого вывода в терминал, минуя перехват terraform. +func ttyOut() *os.File { + if f, err := os.OpenFile("/dev/tty", os.O_WRONLY, 0); err == nil { + return f + } + return os.Stderr +} + +type opStage struct { + InstanceOperationStageUid string `json:"instanceOperationStageUid"` + Stage string `json:"stage"` + IsSuccessful bool `json:"isSuccessful"` + DtFinish *string `json:"dtFinish"` + Duration float64 `json:"duration"` + StageMsg *string `json:"stageMsg"` +} + +type operationStatusResponse struct { + InstanceOperation struct { + DtFinish *string `json:"dtFinish"` + IsSuccessful *bool `json:"isSuccessful"` + ErrorLog *string `json:"errorLog"` + IsInProgress bool `json:"isInProgress"` + IsPending bool `json:"isPending"` + Duration *float64 `json:"duration"` + Stages []opStage `json:"stages"` + } `json:"instanceOperation"` +} + +func (c *UniversalClient) waitForOperationFinish(ctx context.Context, opUid string, timeout time.Duration) error { + deadline := time.Now().Add(timeout) + ticker := time.NewTicker(5 * time.Second) + defer ticker.Stop() + + printedStages := make(map[string]bool) + + for { + select { + case <-ctx.Done(): + return fmt.Errorf("операция %s отменена", opUid) + case <-ticker.C: + if time.Now().After(deadline) { + return fmt.Errorf("операция %s не завершилась за установленный таймаут в %.0f секунд", opUid, timeout.Seconds()) + } + + respBody, _, err := c.doRequest(ctx, "GET", + fmt.Sprintf("/instanceOperations/%s?fields=dtFinish,isSuccessful,errorLog,isInProgress,isPending,duration,stages", opUid), nil) + if err != nil { + return fmt.Errorf("не удалось проверить статус операции %s: %w", opUid, err) + } + + var status operationStatusResponse + if err := json.Unmarshal(respBody, &status); err != nil { + return fmt.Errorf("не удалось разобрать статус операции %s: %w", opUid, err) + } + + // Определяем уровень логирования: ctx перекрывает c.LogLevel + logLevel := c.LogLevel + if v, ok := ctx.Value(ctxKeyLogLevel).(string); ok && v != "" { + logLevel = v + } + showStages := logLevel == "info" || logLevel == "debug" + showDetails := logLevel == "debug" + + // Печатаем завершённые этапы по мере появления + tty := ttyOut() + defer tty.Close() + for _, stage := range status.InstanceOperation.Stages { + if printedStages[stage.InstanceOperationStageUid] { + continue + } + if stage.DtFinish == nil || *stage.DtFinish == "" { + continue + } + printedStages[stage.InstanceOperationStageUid] = true + if !showStages { + continue + } + status2 := "OK " + if !stage.IsSuccessful { + status2 = "FAIL" + } + fmt.Fprintf(tty, " [%s] %s — %.1f sec\n", status2, stage.Stage, stage.Duration) + if showDetails && stage.StageMsg != nil { + for _, line := range formatStageMsg(*stage.StageMsg) { + fmt.Fprintf(tty, "%s\n", line) + } + } + } + + // Критерий завершения — dtFinish (НЕ МЕНЯТЬ) + if status.InstanceOperation.DtFinish != nil && strings.TrimSpace(*status.InstanceOperation.DtFinish) != "" { + if status.InstanceOperation.IsSuccessful != nil && !*status.InstanceOperation.IsSuccessful { + if status.InstanceOperation.ErrorLog != nil && strings.TrimSpace(*status.InstanceOperation.ErrorLog) != "" { + return fmt.Errorf("операция %s завершилась с ошибкой: %s", opUid, *status.InstanceOperation.ErrorLog) + } + return fmt.Errorf("операция %s завершилась с ошибкой", opUid) + } + if showStages && status.InstanceOperation.Duration != nil { + fmt.Fprintf(tty, " [DONE] %.1f sec\n", *status.InstanceOperation.Duration) + } + return nil + } + } + } +} + +// waitForInstanceIdle waits until no operation is pending/in-progress for the instance. +func (c *UniversalClient) waitForInstanceIdle(ctx context.Context, instanceUid string, timeout time.Duration) error { + deadline := time.Now().Add(timeout) + ticker := time.NewTicker(5 * time.Second) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return fmt.Errorf("ожидание операции отменено для экземпляра %s", instanceUid) + case <-ticker.C: + if time.Now().After(deadline) { + return fmt.Errorf("экземпляр %s не перешёл в состояние ожидания за установленный таймаут в %.0f секунд", instanceUid, timeout.Seconds()) + } + + state, err := c.GetInstanceState(ctx, instanceUid) + if err != nil { + return fmt.Errorf("не удалось проверить состояние экземпляра %s: %w", instanceUid, err) + } + if !state.OperationIsPending && !state.OperationIsInProgress { + return nil + } + } + } +} + +// Internal HTTP helpers + +func (c *UniversalClient) doRequest(ctx context.Context, method, path string, payload interface{}) ([]byte, http.Header, error) { + var body io.Reader + if payload != nil { + b, err := json.Marshal(payload) + if err != nil { + return nil, nil, err + } + body = bytes.NewBuffer(b) + } + + req, err := http.NewRequestWithContext(ctx, method, c.ApiEndpoint+path, body) + if err != nil { + return nil, nil, err + } + + req.Close = true + req.Header.Set("Content-Type", "application/json") + if c.ApiToken != "" { + req.Header.Set("Authorization", "Bearer "+c.ApiToken) + } + + resp, err := c.HttpClient.Do(req) + if err != nil { + return nil, nil, err + } + defer resp.Body.Close() + + respBody, err := io.ReadAll(resp.Body) + if err != nil { + return nil, nil, err + } + + if resp.StatusCode >= 400 { + return nil, nil, formatAPIError(resp.StatusCode, respBody) + } + + return respBody, resp.Header, nil +} + +func (c *UniversalClient) postIgnoreResponse(ctx context.Context, path string, payload interface{}, returnLocation bool) (string, error) { + respBody, headers, err := c.doRequest(ctx, "POST", path, payload) + if err != nil { + return "", err + } + + if returnLocation { + if loc := headers.Get("Location"); loc != "" { + return extractUIDFromLocation(loc), nil + } + } + + var justId string + if err := json.Unmarshal(respBody, &justId); err == nil && justId != "" { + return justId, nil + } + + return "", nil +} + +func extractUIDFromLocation(loc string) string { + if loc == "" { + return "" + } + return strings.TrimPrefix(loc, "./") +} + +// NOTE(resourceRealm): доступные resourceRealm для сервиса получают через +// GET /resourceRealms/available?svcId= (через index.cfm proxy). +// Пример: svcId=1 (dummy) возвращает results="dummy". + +// NOTE(refSvcId lookup): если параметр операции create ссылается на другой сервис, +// то в cfsParams будет проставлен refSvcId. Это означает, что значение параметра +// должно быть UUID инстанса указанного сервиса. +// Пример: s3UserUid (service_id=13, S3 бакет) имеет refSvcId=12 (S3 Object Storage), +// значит поле s3_user_uid должно быть UUID S3-инстанса, а не имя вида "s3-111805". +// Для S3 (refSvcId=12) запрещен резолв из displayName — принимается только UUID. +// +// ВАЖНО (uuid-case, 2026-03): API облака возвращает UUID всегда в нижнем регистре. +// Пользователь может написать UUID в ВЕРХНЕМ регистре — оба варианта принимаются API. +// Здесь мы нормализуем UUID к lowercase ДЛЯ ОТПРАВКИ В API (это нужно API). +// Однако в state UUID должен сохраняться в том регистре, который написал пользователь +// (иначе plan != state → "Provider produced inconsistent result after apply"). +// Восстановление регистра делается в Create/Read/Update generated resource, см. шаблон +// instanceTemplate в tools/gen_v2/generate_resources_v2.go (блоки "Restore user-provided casing"). + +func (c *UniversalClient) resolveRefSvcParamValues(ctx context.Context, opParams []universalCfsParam, params map[int]string) (map[int]string, error) { + if len(params) == 0 || len(opParams) == 0 { + return params, nil + } + + refById := make(map[int]int) + for _, p := range opParams { + if p.RefSvcId != nil && *p.RefSvcId > 0 { + refById[p.SvcOperationCfsParamId] = *p.RefSvcId + } + } + if len(refById) == 0 { + return params, nil + } + + resolved := make(map[int]string, len(params)) + for paramId, value := range params { + newValue := value + if refSvcId, ok := refById[paramId]; ok { + if refSvcId == 12 { + if !isUUIDLike(value) { + return nil, fmt.Errorf("параметр %d должен быть UUID (S3 Object Storage, serviceId=12): %s", paramId, value) + } + // FIX(uuid-case): нормализуем UUID к lowercase для отправки в API. + // Пользователь мог написать "6214BA32-...", API принимает оба варианта, + // но сам всегда возвращает нижний регистр. Регистр пользователя будет + // восстановлен в state после RefreshResourceState (generated resource). + newValue = strings.ToLower(value) + } else if !isUUIDLike(value) { + uid, err := c.findInstanceUidByDisplayNameRefSvc(ctx, refSvcId, value) + if err != nil { + return nil, err + } + if uid == "" { + return nil, fmt.Errorf("не удалось разрешить параметр %d: не найден экземпляр для serviceId=%d displayName=%s", paramId, refSvcId, value) + } + newValue = uid + } else { + // FIX(uuid-case): UUID-like value passed directly — normalise to lowercase. + // Аналогично refSvcId==12: нормализуем для API, регистр восстановится в state. + newValue = strings.ToLower(value) + } + } + resolved[paramId] = newValue + } + + return resolved, nil +} + +func (c *UniversalClient) findInstanceUidByDisplayNameRefSvc(ctx context.Context, serviceId int, displayName string) (string, error) { + // Собираем ВСЕ совпадения, затем выбираем лучшее (running > suspended > остальные). + // Deleted инстансы пропускаются — они не должны участвовать в resolve. + type candidate struct { + uid string + status string + } + var candidates []candidate + + page := 1 + for { + path := fmt.Sprintf("/instances?page=%d&size=100", page) + respBody, _, err := c.doRequest(ctx, "GET", path, nil) + if err != nil { + return "", err + } + + var res struct { + Results []struct { + InstanceUid string `json:"instanceUid"` + DisplayName string `json:"displayName"` + ServiceId int `json:"serviceId"` + IsDeleted bool `json:"isDeleted"` + ExplainedStatus string `json:"explainedStatus"` + } `json:"results"` + } + if err := json.Unmarshal(respBody, &res); err != nil { + return "", err + } + + if len(res.Results) == 0 { + break + } + + for _, item := range res.Results { + if item.ServiceId != serviceId || !strings.EqualFold(item.DisplayName, displayName) { + continue + } + // Пропускаем deleted инстансы + if item.IsDeleted { + continue + } + status := strings.ToLower(strings.TrimSpace(item.ExplainedStatus)) + if status == "deleted" { + continue + } + candidates = append(candidates, candidate{uid: item.InstanceUid, status: status}) + } + + page++ + if page > 100 { + break + } + } + + if len(candidates) == 0 { + return "", nil + } + if len(candidates) == 1 { + return candidates[0].uid, nil + } + // Несколько non-deleted совпадений — предпочитаем running + for _, c := range candidates { + if strings.Contains(c.status, "running") || strings.Contains(c.status, "active") { + return c.uid, nil + } + } + // Нет running — предпочитаем suspended + for _, c := range candidates { + if strings.Contains(c.status, "suspend") { + return c.uid, nil + } + } + // Вернуть первый (лучше чем ничего) + return candidates[0].uid, nil +} + +func isUUIDLike(value string) bool { + trimmed := strings.TrimSpace(value) + if len(trimmed) != 36 { + return false + } + for i, r := range trimmed { + switch i { + case 8, 13, 18, 23: + if r != '-' { + return false + } + default: + if !isHexDigit(r) { + return false + } + } + } + return true +} + +func isHexDigit(r rune) bool { + return (r >= '0' && r <= '9') || (r >= 'a' && r <= 'f') || (r >= 'A' && r <= 'F') +} + +func (c *UniversalClient) ensureInstanceCreated(ctx context.Context, instanceUid string) error { + state, err := c.GetInstanceState(ctx, instanceUid) + if err != nil { + return fmt.Errorf("не удалось получить состояние экземпляра после создания: %w", err) + } + if state == nil { + return fmt.Errorf("отсутствует состояние экземпляра после создания для %s", instanceUid) + } + if isInstanceDeleted(state) { + return fmt.Errorf("экземпляр %s удален после создания", instanceUid) + } + status := strings.ToLower(strings.TrimSpace(state.ExplainedStatus)) + if status == "not created" { + return fmt.Errorf("экземпляр %s не создан после завершения операции. Видимо, такой инстанс уже существует в статусе Not Created. %s", instanceUid, formatInstanceStateDetails(state)) + } + return nil +} + +// findInstanceDisplayNameByUidRefSvc resolves display_name by instance UID and service ID. +func (c *UniversalClient) findInstanceDisplayNameByUidRefSvc(ctx context.Context, serviceId int, instanceUid string) (string, error) { + page := 1 + for { + path := fmt.Sprintf("/instances?page=%d&size=100", page) + respBody, _, err := c.doRequest(ctx, "GET", path, nil) + if err != nil { + return "", err + } + + var res struct { + Results []struct { + InstanceUid string `json:"instanceUid"` + DisplayName string `json:"displayName"` + ServiceId int `json:"serviceId"` + } `json:"results"` + } + if err := json.Unmarshal(respBody, &res); err != nil { + return "", err + } + + if len(res.Results) == 0 { + break + } + + for _, item := range res.Results { + if item.ServiceId == serviceId && strings.EqualFold(item.InstanceUid, instanceUid) { + return item.DisplayName, nil + } + } + + page++ + if page > 100 { + break + } + } + + return "", nil +} + +// getInstanceDisplayNameByUidRefSvc fetches display_name by instance UID and checks service ID. +func (c *UniversalClient) getInstanceDisplayNameByUidRefSvc(ctx context.Context, serviceId int, instanceUid string) (string, error) { + if strings.TrimSpace(instanceUid) == "" { + return "", nil + } + + respBody, _, err := c.doRequest(ctx, "GET", fmt.Sprintf("/instances/%s", instanceUid), nil) + if err != nil { + return "", err + } + + var res struct { + Instance struct { + InstanceUid string `json:"instanceUid"` + DisplayName string `json:"displayName"` + ServiceId int `json:"serviceId"` + } `json:"instance"` + } + if err := json.Unmarshal(respBody, &res); err != nil { + return "", err + } + if strings.TrimSpace(res.Instance.InstanceUid) == "" { + return "", nil + } + if serviceId > 0 && res.Instance.ServiceId != serviceId { + return "", nil + } + + return res.Instance.DisplayName, nil +} + +func validateInstanceStatus(state *InstanceStateResponse) error { + if state == nil { + return fmt.Errorf("отсутствует состояние экземпляра") + } + if isInstanceDeleted(state) { + return &instanceDeletedError{InstanceUID: state.InstanceUid} + } + if state.OperationIsPending || state.OperationIsInProgress { + return fmt.Errorf("экземпляр %s не готов: операция в ожидании", state.InstanceUid) + } + status := strings.ToLower(strings.TrimSpace(state.ExplainedStatus)) + if strings.Contains(status, "not created") { + return fmt.Errorf("экземпляр %s не создан. Видимо, такой инстанс уже существует в статусе Not Created. %s", state.InstanceUid, formatInstanceStateDetails(state)) + } + if strings.Contains(status, "pending") { + return fmt.Errorf("экземпляр %s в ожидании. %s", state.InstanceUid, formatInstanceStateDetails(state)) + } + if strings.Contains(status, "failed") || strings.Contains(status, "error") { + return fmt.Errorf("экземпляр %s завершился с ошибкой: %s. %s", state.InstanceUid, state.ExplainedStatus, formatInstanceStateDetails(state)) + } + return nil +} + +func formatInstanceStateDetails(state *InstanceStateResponse) string { + if state == nil { + return "details: state=nil" + } + statusRaw := strings.TrimSpace(state.ExplainedStatus) + statusNorm := strings.ToLower(statusRaw) + return "details: instance_uid=" + state.InstanceUid + ", status=\"" + statusNorm + "\", status_raw=\"" + statusRaw + "\", operation_pending=" + formatBool(state.OperationIsPending) + ", operation_in_progress=" + formatBool(state.OperationIsInProgress) +} + +func formatBool(value bool) string { + if value { + return "true" + } + return "false" +} + +// RunInstanceOperationUniversalByCode runs an operation using params keyed by code. +// It resolves param codes to IDs via operation manifest, applies defaults, validates, and runs. +func (c *UniversalClient) RunInstanceOperationUniversalByCode(ctx context.Context, instanceUid string, action string, params map[string]string) error { + state, err := c.GetInstanceState(ctx, instanceUid) + if err != nil { + return err + } + if state.OperationIsPending || state.OperationIsInProgress { + if err := c.waitForInstanceIdle(ctx, instanceUid, c.idleTimeoutFor(state.ServiceId)); err != nil { + return err + } + } + + var opId int + for _, op := range state.AvailableOperations { + if strings.EqualFold(op.Operation, action) { + opId = op.SvcOperationId + break + } + } + if opId == 0 { + return fmt.Errorf("операция %s недоступна для экземпляра %s", action, instanceUid) + } + + payload := map[string]interface{}{ + "instanceUid": instanceUid, + "svcOperationId": opId, + "operation": action, + } + + opUid, err := c.postIgnoreResponse(ctx, "/instanceOperations", payload, true) + if err != nil { + return fmt.Errorf("не удалось создать операцию %s: %w", action, err) + } + if opUid == "" { + return fmt.Errorf("не удалось получить UID операции для %s", action) + } + + opDetailsResp, _, err := c.doRequest(ctx, "GET", fmt.Sprintf("/instanceOperations/%s?fields=cfsParams", opUid), nil) + if err != nil { + return fmt.Errorf("не удалось получить детали операции: %w", err) + } + var opDetails universalOpResponse + if err := json.Unmarshal(opDetailsResp, &opDetails); err != nil { + return fmt.Errorf("не удалось разобрать детали операции: %w", err) + } + + codeToParam := make(map[string]universalCfsParam) + for _, p := range opDetails.InstanceOperation.CfsParams { + if key := strings.ToLower(strings.TrimSpace(p.Code)); key != "" { + codeToParam[key] = p + } + if key := strings.ToLower(strings.TrimSpace(p.SvcOperationCfsParam)); key != "" { + codeToParam[key] = p + } + } + + paramsByID := map[int]string{} + for code, value := range params { + key := strings.ToLower(strings.TrimSpace(code)) + p, ok := codeToParam[key] + if !ok { + return fmt.Errorf("код параметра %s не найден для операции %s", code, action) + } + paramsByID[p.SvcOperationCfsParamId] = value + } + + paramsByID, err = c.resolveRefSvcParamValues(ctx, opDetails.InstanceOperation.CfsParams, paramsByID) + if err != nil { + return err + } + + sent := make(map[int]bool) + for paramId, value := range paramsByID { + pPayload := genericParamReq{ + InstanceOperationUid: opUid, + SvcOperationCfsParamId: paramId, + ParamValue: value, + } + _, _, err := c.doRequest(ctx, "POST", "/instanceOperationCfsParams", pPayload) + if err != nil { + return fmt.Errorf("не удалось установить параметр %d: %w", paramId, err) + } + sent[paramId] = true + } + + for _, param := range opDetails.InstanceOperation.CfsParams { + if sent[param.SvcOperationCfsParamId] { + continue + } + if !param.IsRequired { + continue + } + val := "" + if param.ParamValue != nil { + val = *param.ParamValue + } else if param.DefaultValue != nil { + val = *param.DefaultValue + } + val = normalizeUniversalValueV6(val, param) + if !param.IsRequired && strings.TrimSpace(val) == "" { + continue + } + + pPayload := genericParamReq{ + InstanceOperationUid: opUid, + SvcOperationCfsParamId: param.SvcOperationCfsParamId, + ParamValue: val, + } + _, _, err := c.doRequest(ctx, "POST", "/instanceOperationCfsParams", pPayload) + if err != nil { + return fmt.Errorf("не удалось отправить параметр по умолчанию %d: %w", param.SvcOperationCfsParamId, err) + } + } + + _, _, err = c.doRequest(ctx, "GET", fmt.Sprintf("/instanceOperations/%s/validate-cfs", opUid), nil) + if err != nil { + return fmt.Errorf("валидация не пройдена: %w", err) + } + + _, _, err = c.doRequest(ctx, "POST", fmt.Sprintf("/instanceOperations/%s/run", opUid), map[string]interface{}{}) + if err != nil { + return err + } + + // НЕ МЕНЯТЬ: завершение операции определяется по dtFinish + return c.waitForOperationFinish(ctx, opUid, c.operationTimeoutForContext(ctx, state.ServiceId, action)) +} + +// instanceDescr returns a stable description with provider version when available. +func (c *UniversalClient) instanceDescr() string { + version := strings.TrimSpace(c.ProviderVersion) + if version == "" { + return "Создано через Nubes Terraform Universal Provider" + } + return fmt.Sprintf("Создано через Nubes Terraform Universal Provider %s", version) +} + +// instanceDeletedError marks deleted instances to be treated as not found. +type instanceDeletedError struct { + InstanceUID string +} + +func (e *instanceDeletedError) Error() string { + return fmt.Sprintf("экземпляр %s удален", e.InstanceUID) +} + +func isInstanceDeletedError(err error) bool { + var e *instanceDeletedError + return errors.As(err, &e) +} + +// formatAPIError парсит JSON-ответ API и возвращает читаемое сообщение. +// Если тело не является JSON с полем ERROR — возвращает сырой текст. +func formatAPIError(statusCode int, body []byte) error { + var parsed struct { + Error string `json:"ERROR"` + Detail string `json:"DETAIL"` + } + if json.Unmarshal(body, &parsed) == nil && parsed.Error != "" { + msg := parsed.Error + if d := strings.TrimSpace(parsed.Detail); d != "" { + msg += ": " + d + } + return fmt.Errorf("ошибка API %d: %s", statusCode, msg) + } + return fmt.Errorf("ошибка API %d: %s", statusCode, strings.TrimSpace(string(body))) +} diff --git a/universal_rebuild/internal/core/client_test.go b/universal_rebuild/internal/core/client_test.go new file mode 100644 index 0000000..935dbf0 --- /dev/null +++ b/universal_rebuild/internal/core/client_test.go @@ -0,0 +1,45 @@ +package core + +import ( + "testing" +) + +func TestNormalizeValue_EmptyString(t *testing.T) { + param := universalCfsParam{DataType: "string"} + result := normalizeUniversalValueV6("", param) + if result != "" { + t.Errorf("expected empty string, got: %q", result) + } +} + +func TestNormalizeValue_NullString(t *testing.T) { + param := universalCfsParam{DataType: "string"} + result := normalizeUniversalValueV6("null", param) + if result != "" { + t.Errorf("expected empty string for 'null' input, got: %q", result) + } +} + +func TestNormalizeValue_MapType(t *testing.T) { + param := universalCfsParam{DataType: "map"} + result := normalizeUniversalValueV6("", param) + if result != "{}" { + t.Errorf("expected '{}' for map dataType, got: %q", result) + } +} + +func TestNormalizeValue_ArrayType(t *testing.T) { + param := universalCfsParam{DataType: "array"} + result := normalizeUniversalValueV6("", param) + if result != "[]" { + t.Errorf("expected '[]' for array dataType, got: %q", result) + } +} + +func TestNormalizeValue_TrimSpace(t *testing.T) { + param := universalCfsParam{DataType: "string"} + result := normalizeUniversalValueV6(" hello ", param) + if result != "hello" { + t.Errorf("expected 'hello' after trim, got: %q", result) + } +} diff --git a/universal_rebuild/internal/core/instance_outputs.go b/universal_rebuild/internal/core/instance_outputs.go new file mode 100644 index 0000000..8bfbeb8 --- /dev/null +++ b/universal_rebuild/internal/core/instance_outputs.go @@ -0,0 +1,213 @@ +package core + +import ( + "context" + "encoding/json" + "fmt" + "net/url" + "strings" +) + +type InstanceStateDetails struct { + Params map[string]string + Out map[string]string + RawParams map[string]interface{} + RawOut map[string]interface{} + Vault InstanceVaultMeta +} + +type InstanceVaultMeta struct { + Url string + UserPath string + Fields []string +} + +// GetInstanceStateDetails returns instance.state params/out and vault metadata. +func (c *UniversalClient) GetInstanceStateDetails(ctx context.Context, instanceUid string) (*InstanceStateDetails, error) { + if strings.TrimSpace(instanceUid) == "" { + return nil, fmt.Errorf("missing instance uid") + } + + respBody, _, err := c.doRequest(ctx, "GET", fmt.Sprintf("/instances/%s", instanceUid), nil) + if err != nil { + return nil, err + } + + var res map[string]interface{} + if err := json.Unmarshal(respBody, &res); err != nil { + return nil, err + } + + inst, ok := res["instance"].(map[string]interface{}) + if !ok { + return &InstanceStateDetails{Params: map[string]string{}, Out: map[string]string{}}, nil + } + + state, _ := inst["state"].(map[string]interface{}) + if state == nil { + return &InstanceStateDetails{Params: map[string]string{}, Out: map[string]string{}}, nil + } + + paramsRaw := extractStateRawMap(state["params"]) + outRaw := extractStateRawMap(state["out"]) + params := extractStateStringMap(paramsRaw) + out := extractStateStringMap(outRaw) + vault := extractVaultMeta(state["vault"]) + + return &InstanceStateDetails{ + Params: params, + Out: out, + RawParams: paramsRaw, + RawOut: outRaw, + Vault: vault, + }, nil +} + +// GetInstanceVaultSecrets fetches all named secrets from the vault endpoint. +func (c *UniversalClient) GetInstanceVaultSecrets(ctx context.Context, instanceUid string, fields []string) (map[string]string, error) { + secrets := make(map[string]string) + if strings.TrimSpace(instanceUid) == "" || len(fields) == 0 { + return secrets, nil + } + + var errs []string + for _, field := range fields { + name := strings.TrimSpace(field) + if name == "" { + continue + } + value, err := c.GetInstanceVaultSecret(ctx, instanceUid, name) + if err != nil { + errs = append(errs, fmt.Sprintf("%s: %v", name, err)) + continue + } + secrets[name] = value + } + + if len(errs) > 0 { + return secrets, fmt.Errorf(strings.Join(errs, "; ")) + } + return secrets, nil +} + +// GetInstanceVaultSecret fetches a single secret value by name. +func (c *UniversalClient) GetInstanceVaultSecret(ctx context.Context, instanceUid string, secretName string) (string, error) { + if strings.TrimSpace(instanceUid) == "" { + return "", fmt.Errorf("missing instance uid") + } + if strings.TrimSpace(secretName) == "" { + return "", fmt.Errorf("missing secret name") + } + + path := fmt.Sprintf("/instances/%s/vault/%s", instanceUid, url.PathEscape(secretName)) + respBody, _, err := c.doRequest(ctx, "GET", path, nil) + if err != nil { + return "", err + } + + return extractVaultSecretValue(respBody) +} + +func extractStateRawMap(value interface{}) map[string]interface{} { + raw, _ := value.(map[string]interface{}) + if raw == nil { + return map[string]interface{}{} + } + return raw +} + +func extractStateStringMap(raw map[string]interface{}) map[string]string { + if raw == nil { + return map[string]string{} + } + + out := make(map[string]string, len(raw)) + for key, val := range raw { + out[key] = normalizeInstanceParamValue(val) + } + return out +} + +func extractVaultMeta(value interface{}) InstanceVaultMeta { + raw, _ := value.(map[string]interface{}) + if raw == nil { + return InstanceVaultMeta{} + } + + return InstanceVaultMeta{ + Url: normalizeInstanceParamValue(raw["url"]), + UserPath: normalizeInstanceParamValue(raw["userPath"]), + Fields: extractStringSlice(raw["fields"]), + } +} + +func extractStringSlice(value interface{}) []string { + list, _ := value.([]interface{}) + if list == nil { + return []string{} + } + + out := make([]string, 0, len(list)) + for _, item := range list { + val := strings.TrimSpace(normalizeInstanceParamValue(item)) + if val == "" { + continue + } + out = append(out, val) + } + return out +} + +func extractVaultSecretValue(body []byte) (string, error) { + var decoded interface{} + if err := json.Unmarshal(body, &decoded); err != nil { + return strings.TrimSpace(string(body)), nil + } + + switch v := decoded.(type) { + case string: + return strings.TrimSpace(v), nil + case map[string]interface{}: + if value, ok := extractSecretFromMap(v); ok { + return value, nil + } + b, _ := json.Marshal(v) + return strings.TrimSpace(string(b)), nil + default: + b, _ := json.Marshal(v) + return strings.TrimSpace(string(b)), nil + } +} + +func extractSecretFromMap(raw map[string]interface{}) (string, bool) { + for _, key := range []string{"value", "secret", "password", "token"} { + if value, ok := raw[key]; ok { + out := strings.TrimSpace(normalizeInstanceParamValue(value)) + if out != "" { + return out, true + } + } + } + + if data, ok := raw["data"].(map[string]interface{}); ok { + if value, ok := extractSecretFromMap(data); ok { + return value, true + } + } + if record, ok := raw["record"].(map[string]interface{}); ok { + if value, ok := extractSecretFromMap(record); ok { + return value, true + } + } + + if len(raw) == 1 { + for _, value := range raw { + out := strings.TrimSpace(normalizeInstanceParamValue(value)) + if out != "" { + return out, true + } + } + } + + return "", false +} diff --git a/universal_rebuild/internal/core/instance_params.go b/universal_rebuild/internal/core/instance_params.go new file mode 100644 index 0000000..5b338fa --- /dev/null +++ b/universal_rebuild/internal/core/instance_params.go @@ -0,0 +1,76 @@ +package core + +import ( + "context" + "encoding/json" + "fmt" + "math" + "strconv" + "strings" +) + +// GetInstanceStateParams returns instance.state.params as a map of string values. +func (c *UniversalClient) GetInstanceStateParams(ctx context.Context, instanceUid string) (map[string]string, error) { + if strings.TrimSpace(instanceUid) == "" { + return nil, fmt.Errorf("missing instance uid") + } + + respBody, _, err := c.doRequest(ctx, "GET", fmt.Sprintf("/instances/%s", instanceUid), nil) + if err != nil { + return nil, err + } + + var res map[string]interface{} + if err := json.Unmarshal(respBody, &res); err != nil { + return nil, err + } + + inst, ok := res["instance"].(map[string]interface{}) + if !ok { + return nil, fmt.Errorf("instance field missing in response") + } + + state, _ := inst["state"].(map[string]interface{}) + if state == nil { + return map[string]string{}, nil + } + + paramsRaw, _ := state["params"].(map[string]interface{}) + if paramsRaw == nil { + return map[string]string{}, nil + } + + params := make(map[string]string, len(paramsRaw)) + for key, val := range paramsRaw { + params[key] = normalizeInstanceParamValue(val) + } + + return params, nil +} + +func normalizeInstanceParamValue(value interface{}) string { + switch v := value.(type) { + case nil: + return "" + case string: + return strings.TrimSpace(v) + case bool: + if v { + return "true" + } + return "false" + case float64: + if math.Mod(v, 1) == 0 { + return strconv.FormatInt(int64(v), 10) + } + return strconv.FormatFloat(v, 'f', -1, 64) + case json.Number: + return v.String() + default: + b, err := json.Marshal(v) + if err == nil { + return strings.TrimSpace(string(b)) + } + return fmt.Sprintf("%v", v) + } +} diff --git a/universal_rebuild/internal/core/operation_timeout_override.go b/universal_rebuild/internal/core/operation_timeout_override.go new file mode 100644 index 0000000..9d0a43d --- /dev/null +++ b/universal_rebuild/internal/core/operation_timeout_override.go @@ -0,0 +1,46 @@ +package core + +import ( + "context" + "fmt" + "strings" + "time" +) + +type operationTimeoutOverrideKey struct{} + +// WithOperationTimeout attaches optional per-request operation timeout override to context. +// Empty value means no override and falls back to configured defaults. +func WithOperationTimeout(ctx context.Context, raw string) (context.Context, error) { + trimmed := strings.TrimSpace(raw) + if trimmed == "" { + return ctx, nil + } + d, err := time.ParseDuration(trimmed) + if err != nil { + return nil, fmt.Errorf("некорректный operation_timeout %q: %w", raw, err) + } + if d <= 0 { + return nil, fmt.Errorf("некорректный operation_timeout %q: значение должно быть больше 0", raw) + } + return context.WithValue(ctx, operationTimeoutOverrideKey{}, d), nil +} + +func operationTimeoutOverrideFromContext(ctx context.Context) (time.Duration, bool) { + if ctx == nil { + return 0, false + } + v := ctx.Value(operationTimeoutOverrideKey{}) + d, ok := v.(time.Duration) + if !ok || d <= 0 { + return 0, false + } + return d, true +} + +func (c *UniversalClient) operationTimeoutForContext(ctx context.Context, serviceID int, action string) time.Duration { + if d, ok := operationTimeoutOverrideFromContext(ctx); ok { + return d + } + return c.operationTimeoutFor(serviceID, action) +} diff --git a/universal_rebuild/internal/core/operation_timeouts.go b/universal_rebuild/internal/core/operation_timeouts.go new file mode 100644 index 0000000..58dfc74 --- /dev/null +++ b/universal_rebuild/internal/core/operation_timeouts.go @@ -0,0 +1,168 @@ +package core + +import ( + "encoding/json" + "fmt" + "strconv" + "strings" + "time" +) + +const ( + defaultOperationTimeoutValue = 30 * time.Minute + defaultIdleTimeoutValue = 30 * time.Minute +) + +type OperationTimeouts struct { + DefaultTimeout time.Duration + IdleTimeout time.Duration + ByOperationName map[string]time.Duration + ByServiceOperation map[int]map[string]time.Duration +} + +type operationTimeoutsConfig struct { + DefaultTimeout string `json:"default_timeout"` + IdleTimeout string `json:"idle_timeout"` + Operations map[string]string `json:"operations"` + Overrides map[string]string `json:"overrides"` +} + +func DefaultOperationTimeouts() *OperationTimeouts { + return &OperationTimeouts{ + DefaultTimeout: defaultOperationTimeoutValue, + IdleTimeout: defaultIdleTimeoutValue, + ByOperationName: map[string]time.Duration{ + "create": 30 * time.Minute, + "modify": 30 * time.Minute, + "resume": 30 * time.Minute, + "suspend": 30 * time.Minute, + "delete_user": 10 * time.Minute, + "create_user": 10 * time.Minute, + "delete_database": 10 * time.Minute, + "create_database": 10 * time.Minute, + }, + ByServiceOperation: map[int]map[string]time.Duration{}, + } +} + +func LoadOperationTimeoutsFromBytes(raw []byte) (*OperationTimeouts, error) { + var cfg operationTimeoutsConfig + if err := json.Unmarshal(raw, &cfg); err != nil { + return nil, fmt.Errorf("failed to parse operation timeouts config: %w", err) + } + + result := DefaultOperationTimeouts() + + if strings.TrimSpace(cfg.DefaultTimeout) != "" { + d, err := time.ParseDuration(strings.TrimSpace(cfg.DefaultTimeout)) + if err != nil { + return nil, fmt.Errorf("invalid default_timeout value %q: %w", cfg.DefaultTimeout, err) + } + if d <= 0 { + return nil, fmt.Errorf("invalid default_timeout value %q: must be > 0", cfg.DefaultTimeout) + } + result.DefaultTimeout = d + } + + if strings.TrimSpace(cfg.IdleTimeout) != "" { + d, err := time.ParseDuration(strings.TrimSpace(cfg.IdleTimeout)) + if err != nil { + return nil, fmt.Errorf("invalid idle_timeout value %q: %w", cfg.IdleTimeout, err) + } + if d <= 0 { + return nil, fmt.Errorf("invalid idle_timeout value %q: must be > 0", cfg.IdleTimeout) + } + result.IdleTimeout = d + } + + for operation, timeoutRaw := range cfg.Operations { + operationKey := strings.ToLower(strings.TrimSpace(operation)) + if operationKey == "" { + continue + } + d, err := time.ParseDuration(strings.TrimSpace(timeoutRaw)) + if err != nil { + return nil, fmt.Errorf("invalid timeout for operation %q: %w", operation, err) + } + if d <= 0 { + return nil, fmt.Errorf("invalid timeout for operation %q: must be > 0", operation) + } + result.ByOperationName[operationKey] = d + } + + for overrideKey, timeoutRaw := range cfg.Overrides { + serviceID, operationKey, err := parseServiceOperationOverrideKey(overrideKey) + if err != nil { + return nil, err + } + d, err := time.ParseDuration(strings.TrimSpace(timeoutRaw)) + if err != nil { + return nil, fmt.Errorf("invalid timeout for override %q: %w", overrideKey, err) + } + if d <= 0 { + return nil, fmt.Errorf("invalid timeout for override %q: must be > 0", overrideKey) + } + if _, ok := result.ByServiceOperation[serviceID]; !ok { + result.ByServiceOperation[serviceID] = map[string]time.Duration{} + } + result.ByServiceOperation[serviceID][operationKey] = d + } + + return result, nil +} + +func parseServiceOperationOverrideKey(raw string) (int, string, error) { + key := strings.ToLower(strings.TrimSpace(raw)) + parts := strings.Split(key, ".") + if len(parts) != 3 || parts[0] != "services" || strings.TrimSpace(parts[1]) == "" || strings.TrimSpace(parts[2]) == "" { + return 0, "", fmt.Errorf("invalid override key %q: expected services..", raw) + } + serviceID, err := strconv.Atoi(parts[1]) + if err != nil || serviceID <= 0 { + return 0, "", fmt.Errorf("invalid override key %q: serviceId must be positive integer", raw) + } + return serviceID, parts[2], nil +} + +func (c *UniversalClient) operationTimeoutFor(serviceID int, action string) time.Duration { + if c == nil || c.OperationTimeouts == nil { + return defaultOperationTimeoutValue + } + key := strings.ToLower(strings.TrimSpace(action)) + if serviceID > 0 && key != "" { + if serviceOverrides, ok := c.OperationTimeouts.ByServiceOperation[serviceID]; ok { + if d, ok := serviceOverrides[key]; ok && d > 0 { + return d + } + } + } + if key != "" { + if d, ok := c.OperationTimeouts.ByOperationName[key]; ok && d > 0 { + return d + } + } + if c.OperationTimeouts.DefaultTimeout > 0 { + return c.OperationTimeouts.DefaultTimeout + } + return defaultOperationTimeoutValue +} + +func (c *UniversalClient) idleTimeoutFor(serviceID int) time.Duration { + if c == nil || c.OperationTimeouts == nil { + return defaultIdleTimeoutValue + } + if serviceID > 0 { + if serviceOverrides, ok := c.OperationTimeouts.ByServiceOperation[serviceID]; ok { + if d, ok := serviceOverrides["idle"]; ok && d > 0 { + return d + } + } + } + if c.OperationTimeouts.IdleTimeout > 0 { + return c.OperationTimeouts.IdleTimeout + } + if c.OperationTimeouts.DefaultTimeout > 0 { + return c.OperationTimeouts.DefaultTimeout + } + return defaultIdleTimeoutValue +} diff --git a/universal_rebuild/internal/core/refsvc_resolve.go b/universal_rebuild/internal/core/refsvc_resolve.go new file mode 100644 index 0000000..61bdf6d --- /dev/null +++ b/universal_rebuild/internal/core/refsvc_resolve.go @@ -0,0 +1,199 @@ +package core + +import ( + "context" + "encoding/json" + "fmt" + "strings" +) + +type RefServiceInstanceOption struct { + InstanceUID string + DisplayName string +} + +// ResolveRefSvcParamValue normalizes a refSvcId parameter to a UUID if possible. +// If value already looks like a UUID, it is returned as-is. +func (c *UniversalClient) ResolveRefSvcParamValue(ctx context.Context, refSvcId int, value string) (string, error) { + if c == nil { + return "", nil + } + if refSvcId <= 0 { + return value, nil + } + trimmed := strings.TrimSpace(value) + if trimmed == "" { + return "", nil + } + // FIX(uuid-case): UUID-like values are normalised to lowercase so that + // plan (user input may be uppercase) and API response (always lowercase) + // produce identical strings and do not cause Terraform"s plan/state + // inconsistency error. + if isUUIDLike(trimmed) { + return strings.ToLower(trimmed), nil + } + uid, err := c.findInstanceUidByDisplayNameRefSvc(ctx, refSvcId, trimmed) + if err != nil { + return "", err + } + return uid, nil +} + +// ResolveRefSvcParamDisplayName maps UUID values back to display_name when possible. +func (c *UniversalClient) ResolveRefSvcParamDisplayName(ctx context.Context, refSvcId int, value string) (string, error) { + if c == nil { + return "", nil + } + if refSvcId <= 0 { + return value, nil + } + trimmed := strings.TrimSpace(value) + if trimmed == "" { + return "", nil + } + if !isUUIDLike(trimmed) { + return trimmed, nil + } + // Prefer direct lookup by instance UID to avoid paging the full instances list. + displayName, err := c.getInstanceDisplayNameByUidRefSvc(ctx, refSvcId, trimmed) + if err != nil { + return "", err + } + if strings.TrimSpace(displayName) == "" { + // Fallback to list scan if the direct lookup did not match by serviceId. + displayName, err = c.findInstanceDisplayNameByUidRefSvc(ctx, refSvcId, trimmed) + if err != nil { + return "", err + } + if strings.TrimSpace(displayName) == "" { + return trimmed, nil + } + } + return displayName, nil +} + +// ListRefServiceInstances returns active instance options for a referenced service. +func (c *UniversalClient) ListRefServiceInstances(ctx context.Context, serviceId int) ([]RefServiceInstanceOption, error) { + if c == nil || serviceId <= 0 { + return nil, nil + } + + var out []RefServiceInstanceOption + page := 1 + for { + // isDeleted=false — API-уровень фильтрации, не тратим трафик на удалённые инстансы + path := fmt.Sprintf("/instances?page=%d&size=100&isDeleted=false", page) + respBody, _, err := c.doRequest(ctx, "GET", path, nil) + if err != nil { + return nil, err + } + + var res struct { + Results []struct { + InstanceUid string `json:"instanceUid"` + DisplayName string `json:"displayName"` + ServiceId int `json:"serviceId"` + IsDeleted bool `json:"isDeleted"` + ExplainedStatus string `json:"explainedStatus"` + } `json:"results"` + } + if err := json.Unmarshal(respBody, &res); err != nil { + return nil, err + } + + if len(res.Results) == 0 { + break + } + + for _, item := range res.Results { + if item.ServiceId != serviceId { + continue + } + // показываем только running инстансы — deleted/not_created не пригодны к использованию + if item.IsDeleted || !strings.EqualFold(strings.TrimSpace(item.ExplainedStatus), "running") { + continue + } + uid := strings.TrimSpace(item.InstanceUid) + name := strings.TrimSpace(item.DisplayName) + if uid == "" { + continue + } + out = append(out, RefServiceInstanceOption{InstanceUID: uid, DisplayName: name}) + } + + page++ + if page > 100 { + break + } + } + + return out, nil +} + +// ListAvailableResourceRealms returns realm names available for the target service. +func (c *UniversalClient) ListAvailableResourceRealms(ctx context.Context, serviceId int) ([]string, error) { + if c == nil || serviceId <= 0 { + return nil, nil + } + + respBody, _, err := c.doRequest(ctx, "GET", fmt.Sprintf("/resourceRealms/available?svcId=%d", serviceId), nil) + if err != nil { + return nil, err + } + + var payload struct { + Results json.RawMessage `json:"results"` + } + if err := json.Unmarshal(respBody, &payload); err != nil { + return nil, err + } + + parseArray := func(raw json.RawMessage) []string { + var values []string + _ = json.Unmarshal(raw, &values) + return values + } + + parseString := func(raw json.RawMessage) []string { + var single string + if err := json.Unmarshal(raw, &single); err != nil { + return nil + } + single = strings.TrimSpace(single) + if single == "" { + return nil + } + parts := strings.Split(single, ",") + result := make([]string, 0, len(parts)) + for _, part := range parts { + trimmed := strings.TrimSpace(part) + if trimmed != "" { + result = append(result, trimmed) + } + } + if len(result) == 0 { + result = append(result, single) + } + return result + } + + values := parseArray(payload.Results) + if len(values) == 0 { + values = parseString(payload.Results) + } + + out := make([]string, 0, len(values)) + seen := map[string]struct{}{} + for _, v := range values { + trimmed := strings.TrimSpace(v) + if trimmed == "" { + continue + } + if _, ok := seen[trimmed]; ok { + continue + } + seen[trimmed] = struct{}{} + out = append(out, trimmed) + } + return out, nil +}