From 607f8e361c9ff17adafb529959b957d341cfe206 Mon Sep 17 00:00:00 2001 From: phil Date: Wed, 12 Aug 2026 23:30:44 +0900 Subject: [PATCH] =?UTF-8?q?feat(scrape,discovery,promtext):=20=EC=88=98?= =?UTF-8?q?=EC=A7=91=20=EA=B3=84=EC=B8=B5=20=EC=9D=B4=EA=B4=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Prometheus 텍스트 파서·타겟 발견(static/kubernetes)·스크레이프 루프를 옮겼다. promtext 와 discovery 는 의존이 없어 그대로 왔고, scrape 만 개조했다. scrape 는 *tsdb.DB 대신 Appender 인터페이스를 받는다 — promql 이 저장 계층 인터페이스를 스스로 정의한 것과 같은 이유다. tsdb.DB 가 그 시그니처를 이미 만족하므로 어댑터가 필요 없고, 테스트는 저장 엔진 없이 가짜 하나로 끝난다. --- internal/discovery/discovery.go | 237 ++++++++++++++++++++ internal/discovery/discovery_test.go | 319 +++++++++++++++++++++++++++ internal/promtext/parse.go | 260 ++++++++++++++++++++++ internal/promtext/parse_test.go | 302 +++++++++++++++++++++++++ internal/scrape/scrape.go | 201 +++++++++++++++++ internal/scrape/scrape_test.go | 286 ++++++++++++++++++++++++ 6 files changed, 1605 insertions(+) create mode 100644 internal/discovery/discovery.go create mode 100644 internal/discovery/discovery_test.go create mode 100644 internal/promtext/parse.go create mode 100644 internal/promtext/parse_test.go create mode 100644 internal/scrape/scrape.go create mode 100644 internal/scrape/scrape_test.go diff --git a/internal/discovery/discovery.go b/internal/discovery/discovery.go new file mode 100644 index 0000000..3e5e2d8 --- /dev/null +++ b/internal/discovery/discovery.go @@ -0,0 +1,237 @@ +// Package discovery 는 스크레이프 대상(nodevitals 파드)을 찾는다. 두 가지 방법을 +// 제공한다 — 고정 목록(Static)과 in-cluster Kubernetes API 폴링(Kubernetes). 둘 다 +// client-go 없이 stdlib 만으로 동작한다(외부 의존 0 계약). +package discovery + +import ( + "context" + "crypto/tls" + "crypto/x509" + "encoding/json" + "fmt" + "io" + "net" + "net/http" + "net/url" + "os" + "strings" +) + +// Target 은 스크레이프 대상 하나다. URL 은 완전한 metrics URL +// (예: http://10.0.7.101:9847/metrics), Name 은 instance 라벨 값. +type Target struct { + Name string + URL string +} + +// Discoverer 는 스크레이프 대상 목록을 낸다. +type Discoverer interface { + Discover(ctx context.Context) ([]Target, error) +} + +// ---- (a) static ---- + +// Static 은 고정된 타겟 목록을 그대로 내는 Discoverer 다. +type Static struct{ targets []Target } + +var _ Discoverer = (*Static)(nil) + +// NewStatic 은 고정 목록을 그대로 내는 Static 을 만든다. 인자 슬라이스를 복사해 +// 보관한다 — 호출자가 이후 원본 슬라이스를 변조해도 내부 상태에 영향 없다. +func NewStatic(targets []Target) *Static { + cp := make([]Target, len(targets)) + copy(cp, targets) + return &Static{targets: cp} +} + +// Discover 는 생성 시 받은 목록의 복사본을 반환한다(호출자 변조 격리) — 반환된 +// 슬라이스를 호출자가 수정해도 다음 Discover 호출 결과에는 영향이 없다. +func (s *Static) Discover(_ context.Context) ([]Target, error) { + out := make([]Target, len(s.targets)) + copy(out, s.targets) + return out, nil +} + +// ParseStaticTargets 는 콤마 구분 URL 목록("http://a:9847/metrics,http://b:9847/metrics")을 +// Target 으로 변환한다. Name = url.Host. 빈 문자열은 빈 목록(에러 아님). 파싱 불가 +// URL 이나 scheme/host 가 없는 URL 은 에러. +func ParseStaticTargets(s string) ([]Target, error) { + if strings.TrimSpace(s) == "" { + return nil, nil + } + + parts := strings.Split(s, ",") + out := make([]Target, 0, len(parts)) + for _, p := range parts { + p = strings.TrimSpace(p) + if p == "" { + continue + } + u, err := url.Parse(p) + if err != nil { + return nil, fmt.Errorf("discovery: 정적 타겟 URL 파싱 실패 %q: %w", p, err) + } + if u.Scheme == "" || u.Host == "" { + return nil, fmt.Errorf("discovery: 정적 타겟 URL 에 scheme 또는 host 가 없다 %q", p) + } + out = append(out, Target{Name: u.Host, URL: p}) + } + return out, nil +} + +// ---- (b) kubernetes (in-cluster, stdlib 직접 REST) ---- + +const ( + defaultBaseURL = "https://kubernetes.default.svc" + defaultTokenFile = "/var/run/secrets/kubernetes.io/serviceaccount/token" + defaultCAFile = "/var/run/secrets/kubernetes.io/serviceaccount/ca.crt" + defaultMetricsPath = "/metrics" +) + +// KubeConfig 는 Kubernetes 파드 발견 설정이다. +type KubeConfig struct { + // BaseURL 기본 "https://kubernetes.default.svc". 테스트는 httptest.Server.URL 주입. + BaseURL string + // TokenFile 기본 "/var/run/secrets/kubernetes.io/serviceaccount/token". + TokenFile string + // CAFile 기본 "/var/run/secrets/kubernetes.io/serviceaccount/ca.crt". + // Client 가 nil 일 때만 사용해 TLS RootCAs 를 구성한다. + CAFile string + + Namespace string // 예: platform-system + LabelSelector string // 예: app.kubernetes.io/name=nodevitals + Port int // 예: 9847 + MetricsPath string // 기본 "/metrics" + + // Client 주입 시 CAFile 무시(테스트: httptest 용 평문/자체서명 클라이언트). + Client *http.Client +} + +// Kubernetes 는 in-cluster kube-apiserver REST 호출로 파드를 발견하는 Discoverer 다. +type Kubernetes struct { + cfg KubeConfig // 기본값 채운 정규화본(Client 필드는 참고용, 실사용은 client) + client *http.Client +} + +var _ Discoverer = (*Kubernetes)(nil) + +// NewKubernetes 는 기본값을 채우고 cfg.Client 가 nil 이면 CAFile 로 TLS pool 을 +// 구성한 클라이언트를 만든다. CAFile 읽기·파싱 실패는 여기서 즉시 에러다 — +// Discover 호출 시점까지 미루지 않는다(fail fast). +func NewKubernetes(cfg KubeConfig) (*Kubernetes, error) { + if cfg.BaseURL == "" { + cfg.BaseURL = defaultBaseURL + } + if cfg.TokenFile == "" { + cfg.TokenFile = defaultTokenFile + } + if cfg.CAFile == "" { + cfg.CAFile = defaultCAFile + } + if cfg.MetricsPath == "" { + cfg.MetricsPath = defaultMetricsPath + } + + client := cfg.Client + if client == nil { + caPEM, err := os.ReadFile(cfg.CAFile) + if err != nil { + return nil, fmt.Errorf("discovery: CA 파일 읽기 실패 %q: %w", cfg.CAFile, err) + } + pool := x509.NewCertPool() + if !pool.AppendCertsFromPEM(caPEM) { + return nil, fmt.Errorf("discovery: CA 파일에 유효한 인증서가 없다 %q", cfg.CAFile) + } + client = &http.Client{ + Transport: &http.Transport{ + TLSClientConfig: &tls.Config{RootCAs: pool}, + }, + } + } + + return &Kubernetes{cfg: cfg, client: client}, nil +} + +// podList 는 kube-apiserver PodList 응답에서 필요한 필드만 담는 비공개 구조체다 +// (client-go/k8s.io/api 금지 계약 — 외부 의존 0). +type podList struct { + Items []struct { + Metadata struct { + Name string `json:"name"` + } `json:"metadata"` + Spec struct { + NodeName string `json:"nodeName"` + } `json:"spec"` + Status struct { + Phase string `json:"phase"` + PodIP string `json:"podIP"` + } `json:"status"` + } `json:"items"` +} + +// Discover 는 GET {BaseURL}/api/v1/namespaces/{ns}/pods?labelSelector={sel} 을 +// 호출한다. Authorization: Bearer . 응답 JSON 에서 status.phase=="Running" && status.podIP!="" 인 파드만 +// Target 으로 변환한다: Name = spec.nodeName(hostNetwork 라 노드 식별이 유의미; +// 비면 metadata.name), URL = "http://" + net.JoinHostPort(podIP, port) + MetricsPath. +// 비-200 응답은 상태코드 + 본문 앞 256바이트를 담은 에러다. +func (k *Kubernetes) Discover(ctx context.Context) ([]Target, error) { + token, err := os.ReadFile(k.cfg.TokenFile) + if err != nil { + return nil, fmt.Errorf("discovery: 토큰 파일 읽기 실패 %q: %w", k.cfg.TokenFile, err) + } + + q := url.Values{} + q.Set("labelSelector", k.cfg.LabelSelector) + reqURL := fmt.Sprintf("%s/api/v1/namespaces/%s/pods?%s", + strings.TrimRight(k.cfg.BaseURL, "/"), + url.PathEscape(k.cfg.Namespace), + q.Encode()) + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, reqURL, nil) + if err != nil { + return nil, fmt.Errorf("discovery: 요청 생성 실패: %w", err) + } + req.Header.Set("Authorization", "Bearer "+strings.TrimSpace(string(token))) + req.Header.Set("Accept", "application/json") + + resp, err := k.client.Do(req) + if err != nil { + return nil, fmt.Errorf("discovery: apiserver 요청 실패: %w", err) + } + defer resp.Body.Close() + + body, err := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) + if err != nil { + return nil, fmt.Errorf("discovery: 응답 본문 읽기 실패: %w", err) + } + + if resp.StatusCode != http.StatusOK { + snippet := body + if len(snippet) > 256 { + snippet = snippet[:256] + } + return nil, fmt.Errorf("discovery: apiserver 비정상 응답 %d: %s", resp.StatusCode, snippet) + } + + var list podList + if err := json.Unmarshal(body, &list); err != nil { + return nil, fmt.Errorf("discovery: 파드 목록 JSON 파싱 실패: %w", err) + } + + targets := make([]Target, 0, len(list.Items)) + for _, item := range list.Items { + if item.Status.Phase != "Running" || item.Status.PodIP == "" { + continue + } + name := item.Spec.NodeName + if name == "" { + name = item.Metadata.Name + } + targets = append(targets, Target{ + Name: name, + URL: "http://" + net.JoinHostPort(item.Status.PodIP, fmt.Sprintf("%d", k.cfg.Port)) + k.cfg.MetricsPath, + }) + } + return targets, nil +} diff --git a/internal/discovery/discovery_test.go b/internal/discovery/discovery_test.go new file mode 100644 index 0000000..ac11d72 --- /dev/null +++ b/internal/discovery/discovery_test.go @@ -0,0 +1,319 @@ +package discovery + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" +) + +// ---- Static ---- + +func TestStatic_Discover_고정목록을_그대로_낸다(t *testing.T) { + want := []Target{ + {Name: "e101", URL: "http://10.0.7.101:9847/metrics"}, + {Name: "e102", URL: "http://10.0.7.102:9847/metrics"}, + } + s := NewStatic(want) + + got, err := s.Discover(context.Background()) + if err != nil { + t.Fatalf("Discover: %v", err) + } + if len(got) != len(want) { + t.Fatalf("타겟 개수: got %d, want %d", len(got), len(want)) + } + for i := range want { + if got[i] != want[i] { + t.Errorf("타겟[%d]: got %+v, want %+v", i, got[i], want[i]) + } + } +} + +// NewStatic/Discover 가 슬라이스를 공유하면(복사 생략) 한쪽 변조가 다른 쪽에 +// 새어나간다 — 계약("Discover 는 복사본 반환")을 실제로 어기는 구현이면 실패한다. +func TestStatic_Discover는_복사본을_반환한다(t *testing.T) { + orig := []Target{{Name: "a", URL: "http://a:9847/metrics"}} + s := NewStatic(orig) + orig[0].Name = "mutated-after-new-static" // NewStatic 이 원본 슬라이스를 그대로 잡고 있으면 오염 + + got1, err := s.Discover(context.Background()) + if err != nil { + t.Fatalf("Discover: %v", err) + } + if got1[0].Name != "a" { + t.Fatalf("NewStatic 이 입력 슬라이스를 복사하지 않았다: got %q, want %q", got1[0].Name, "a") + } + + got1[0].Name = "mutated-after-discover" // Discover 반환값이 내부 상태와 공유되면 다음 호출이 오염 + + got2, err := s.Discover(context.Background()) + if err != nil { + t.Fatalf("Discover: %v", err) + } + if got2[0].Name != "a" { + t.Fatalf("Discover 가 내부 상태를 공유 반환했다: got %q, want %q", got2[0].Name, "a") + } +} + +// ---- ParseStaticTargets ---- + +func TestParseStaticTargets(t *testing.T) { + got, err := ParseStaticTargets("http://a:9847/metrics,http://b:9847/metrics") + if err != nil { + t.Fatalf("정상 입력인데 에러: %v", err) + } + want := []Target{ + {Name: "a:9847", URL: "http://a:9847/metrics"}, + {Name: "b:9847", URL: "http://b:9847/metrics"}, + } + if len(got) != len(want) { + t.Fatalf("타겟 개수: got %d, want %d (%+v)", len(got), len(want), got) + } + for i := range want { + if got[i] != want[i] { + t.Errorf("타겟[%d]: got %+v, want %+v", i, got[i], want[i]) + } + } + + if _, err := ParseStaticTargets("://x"); err == nil { + t.Fatal("깨진 URL(\"://x\")인데 에러가 없다") + } + + // scheme 은 있으나 host 가 없는 경우 — url.Parse 자체는 에러를 내지 않으므로 + // 별도 Host 검증 분기가 실제로 도달·동작하는지 확인한다. + if _, err := ParseStaticTargets("http://"); err == nil { + t.Fatal("host 없는 URL(\"http://\")인데 에러가 없다") + } + + got, err = ParseStaticTargets("") + if err != nil { + t.Fatalf("빈 문자열인데 에러: %v", err) + } + if len(got) != 0 { + t.Fatalf("빈 문자열 결과: got %d 개, want 0", len(got)) + } +} + +// ---- Kubernetes ---- + +func writeTokenFile(t *testing.T, content string) string { + t.Helper() + p := filepath.Join(t.TempDir(), "token") + if err := os.WriteFile(p, []byte(content), 0o600); err != nil { + t.Fatal(err) + } + return p +} + +func TestKubernetes_fake_apiserver_URL과_토큰(t *testing.T) { + tokenPath := writeTokenFile(t, "test-token-abc") + + var gotAuth, gotPath, gotSelector string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotAuth = r.Header.Get("Authorization") + gotPath = r.URL.Path + gotSelector = r.URL.Query().Get("labelSelector") + + w.Header().Set("Content-Type", "application/json") + // 5 파드로 필터 조건(phase==Running && podIP!="") 의 양쪽을 독립적으로 검증한다: + // pod-a: Running+IP → 포함, Name=nodeName + // pod-b: Running+IP, nodeName 없음 → 포함, Name=metadata.name(대체) + // pod-c: Pending, IP 없음 → 제외(둘 다 거짓) + // pod-d: Running 이지만 IP 없음 → 제외(podIP 조건 단독 검증) + // pod-e: Succeeded 인데 IP 있음 → 제외(phase 조건 단독 검증) + fmt.Fprint(w, `{"items":[ + {"metadata":{"name":"pod-a"},"spec":{"nodeName":"e101"},"status":{"phase":"Running","podIP":"10.0.7.101"}}, + {"metadata":{"name":"pod-b"},"spec":{},"status":{"phase":"Running","podIP":"10.0.7.102"}}, + {"metadata":{"name":"pod-c"},"spec":{"nodeName":"e103"},"status":{"phase":"Pending","podIP":""}}, + {"metadata":{"name":"pod-d"},"spec":{"nodeName":"e104"},"status":{"phase":"Running","podIP":""}}, + {"metadata":{"name":"pod-e"},"spec":{"nodeName":"e105"},"status":{"phase":"Succeeded","podIP":"10.0.7.105"}} + ]}`) + })) + defer srv.Close() + + k, err := NewKubernetes(KubeConfig{ + BaseURL: srv.URL, + TokenFile: tokenPath, + Namespace: "platform-system", + LabelSelector: "app.kubernetes.io/name=nodevitals", + Port: 9847, + Client: srv.Client(), + }) + if err != nil { + t.Fatalf("NewKubernetes: %v", err) + } + + targets, err := k.Discover(context.Background()) + if err != nil { + t.Fatalf("Discover: %v", err) + } + + if gotAuth != "Bearer test-token-abc" { + t.Errorf("Authorization 헤더: got %q, want %q", gotAuth, "Bearer test-token-abc") + } + if gotPath != "/api/v1/namespaces/platform-system/pods" { + t.Errorf("요청 경로: got %q", gotPath) + } + if gotSelector != "app.kubernetes.io/name=nodevitals" { + t.Errorf("labelSelector 쿼리: got %q", gotSelector) + } + + if len(targets) != 2 { + t.Fatalf("타겟 개수: got %d, want 2 (Running && podIP!=\"\" 인 파드만): %+v", len(targets), targets) + } + + byName := map[string]string{} + for _, tg := range targets { + byName[tg.Name] = tg.URL + } + if got, want := byName["e101"], "http://10.0.7.101:9847/metrics"; got != want { + t.Errorf("e101(nodeName) 타겟 URL: got %q, want %q", got, want) + } + if got, want := byName["pod-b"], "http://10.0.7.102:9847/metrics"; got != want { + t.Errorf("pod-b(nodeName 없음 → metadata.name 대체) 타겟 URL: got %q, want %q", got, want) + } + for _, excluded := range []string{"e103", "e104", "e105", "pod-d", "pod-e"} { + if _, ok := byName[excluded]; ok { + t.Errorf("제외되어야 할 파드 %q 가 결과에 포함되었다: %+v", excluded, targets) + } + } +} + +func TestKubernetes_비200은_에러(t *testing.T) { + tokenPath := writeTokenFile(t, "tok") + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusForbidden) + fmt.Fprint(w, "forbidden: no rbac") + })) + defer srv.Close() + + k, err := NewKubernetes(KubeConfig{ + BaseURL: srv.URL, + TokenFile: tokenPath, + Namespace: "platform-system", + Port: 9847, + Client: srv.Client(), + }) + if err != nil { + t.Fatalf("NewKubernetes: %v", err) + } + + targets, err := k.Discover(context.Background()) + if err == nil { + t.Fatalf("Discover 가 nil 에러를 반환했다(403 응답인데): targets=%+v", targets) + } + if !strings.Contains(err.Error(), "403") { + t.Errorf("에러 메시지에 상태코드가 없다: %v", err) + } + if !strings.Contains(err.Error(), "forbidden") { + t.Errorf("에러 메시지에 응답 본문이 없다: %v", err) + } +} + +// 본문이 256바이트를 넘으면 잘려야 한다는 계약(§3) — 자르기 로직이 없는 +// 구현이면(snippet := body 로 캡을 생략) 마커가 그대로 노출되어 실패한다. +func TestKubernetes_비200_본문은_256바이트로_잘린다(t *testing.T) { + tokenPath := writeTokenFile(t, "tok") + + long := strings.Repeat("x", 300) + "TAIL-MARKER" + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + fmt.Fprint(w, long) + })) + defer srv.Close() + + k, err := NewKubernetes(KubeConfig{BaseURL: srv.URL, TokenFile: tokenPath, Client: srv.Client()}) + if err != nil { + t.Fatalf("NewKubernetes: %v", err) + } + + _, err = k.Discover(context.Background()) + if err == nil { + t.Fatal("500 응답인데 에러가 없다") + } + if strings.Contains(err.Error(), "TAIL-MARKER") { + t.Fatalf("본문이 256바이트로 잘리지 않았다(300바이트 뒤 마커가 노출됨): %v", err) + } + if !strings.Contains(err.Error(), "500") { + t.Errorf("상태코드가 메시지에 없다: %v", err) + } +} + +// Discover 는 매 호출마다 토큰 파일을 재독해야 한다(로테이션 대응) — 캐시로 +// 퇴화한 구현이면 2차 요청도 구 토큰을 보내 이 테스트가 실패한다. +func TestKubernetes_토큰_재독(t *testing.T) { + tokenPath := writeTokenFile(t, "token-v1") + + var gotAuths []string + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + gotAuths = append(gotAuths, r.Header.Get("Authorization")) + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"items":[]}`) + })) + defer srv.Close() + + k, err := NewKubernetes(KubeConfig{ + BaseURL: srv.URL, + TokenFile: tokenPath, + Namespace: "platform-system", + Port: 9847, + Client: srv.Client(), + }) + if err != nil { + t.Fatalf("NewKubernetes: %v", err) + } + + if _, err := k.Discover(context.Background()); err != nil { + t.Fatalf("1차 Discover: %v", err) + } + + if err := os.WriteFile(tokenPath, []byte("token-v2"), 0o600); err != nil { + t.Fatal(err) + } + + if _, err := k.Discover(context.Background()); err != nil { + t.Fatalf("2차 Discover: %v", err) + } + + if len(gotAuths) != 2 { + t.Fatalf("요청 횟수: got %d, want 2 (%v)", len(gotAuths), gotAuths) + } + if gotAuths[0] != "Bearer token-v1" { + t.Errorf("1차 요청 토큰: got %q, want %q", gotAuths[0], "Bearer token-v1") + } + if gotAuths[1] != "Bearer token-v2" { + t.Errorf("2차 요청 토큰(재독 되었어야 함): got %q, want %q", gotAuths[1], "Bearer token-v2") + } +} + +// NewKubernetes 는 CAFile 읽기 실패를 Discover 시점까지 미루지 않고 즉시 반환해야 +// 한다(fail fast 계약, §3). +func TestNewKubernetes_CA파일_읽기실패는_즉시_에러(t *testing.T) { + _, err := NewKubernetes(KubeConfig{ + CAFile: filepath.Join(t.TempDir(), "no-such-ca.crt"), + }) + if err == nil { + t.Fatal("존재하지 않는 CAFile 인데 NewKubernetes 가 에러를 내지 않았다") + } +} + +// Client 가 주입되면 CAFile 은 무시되어야 한다(§3) — 무시하지 않는 구현이면 +// 존재하지 않는 CAFile 경로 때문에 이 테스트가 실패한다. +func TestNewKubernetes_Client_주입시_CAFile_무시(t *testing.T) { + k, err := NewKubernetes(KubeConfig{ + CAFile: filepath.Join(t.TempDir(), "no-such-ca.crt"), + Client: &http.Client{}, + }) + if err != nil { + t.Fatalf("Client 가 주입되었는데도 CAFile 읽기를 시도해 에러가 났다: %v", err) + } + if k == nil { + t.Fatal("반환된 *Kubernetes 가 nil") + } +} diff --git a/internal/promtext/parse.go b/internal/promtext/parse.go new file mode 100644 index 0000000..bda3b03 --- /dev/null +++ b/internal/promtext/parse.go @@ -0,0 +1,260 @@ +// Package promtext 는 Prometheus text exposition format(v0.0.4)을 파싱한다. +// +// stdlib 만 사용하고 정규식은 쓰지 않는다 — 3800 시리즈 × 10 타겟 × 15s 스크레이프 +// 주기를 감당해야 하므로, bufio.Scanner 로 줄을 읽은 뒤 수동 문자 스캔으로 +// 이름 / 라벨블록 / 값 / 옵션 타임스탬프를 읽는다. 라벨 블록 안은 콤마로 +// split 하지 않는다 — 라벨 값 자체에 콤마가 올 수 있어 split 이 값을 깨뜨린다. +package promtext + +import ( + "bufio" + "fmt" + "io" + "strconv" + "strings" +) + +// maxLineBytes 는 bufio.Scanner 가 허용하는 한 줄의 최대 바이트 수다. +// 비정상적으로 긴 줄이 무제한 메모리를 잡아먹지 않도록 상한을 둔다. +const maxLineBytes = 1 << 20 // 1MiB + +// Series 는 exposition 한 줄이다. TimestampMS==0 은 "타임스탬프 없음"을 뜻한다. +type Series struct { + Name string + Labels map[string]string // 라벨 없으면 nil 또는 빈 맵 — 호출자는 len 으로만 판단 + Value float64 + TimestampMS int64 +} + +// Parse 는 r 이 담은 Prometheus text exposition format 을 파싱해 Series 목록을 +// 낸다. +// - `# HELP` / `# TYPE` / 기타 `#` 로 시작하는 주석 줄, 빈 줄은 무시한다. +// - 라벨 값 이스케이프 3종(\\, \", \n)을 해석한다. +// - 값은 strconv.ParseFloat 이 받아들이는 전부(NaN/±Inf 대소문자 무관 포함)를 받는다. +// - 옵션 타임스탬프(밀리초 정수)가 없으면 TimestampMS=0. +// - 형식이 어긋난 줄을 만나면 즉시 줄번호를 포함한 에러를 반환한다 — 조용히 +// 건너뛰지 않는다. +func Parse(r io.Reader) ([]Series, error) { + sc := bufio.NewScanner(r) + sc.Buffer(make([]byte, 0, 64*1024), maxLineBytes) + + var out []Series + lineNo := 0 + for sc.Scan() { + lineNo++ + line := sc.Text() + trimmed := strings.TrimSpace(line) + if trimmed == "" || trimmed[0] == '#' { + continue + } + s, err := parseLine(line, lineNo) + if err != nil { + return nil, err + } + out = append(out, s) + } + if err := sc.Err(); err != nil { + return nil, fmt.Errorf("promtext: 스캔 실패: %w", err) + } + return out, nil +} + +// parseLine 은 데이터 한 줄(이미 주석·빈줄이 아님이 확인됨)을 +// 이름 → [라벨블록] → 값 → [타임스탬프] 순서로 스캔한다. +func parseLine(line string, lineNo int) (Series, error) { + p := &lineScanner{s: line, lineNo: lineNo} + p.skipSpaces() // 관례상 선행 공백은 없지만 방어적으로 허용 + + name := p.scanName() + if name == "" { + return Series{}, p.errorf("메트릭 이름이 비어 있다") + } + p.skipSpaces() + + var labels map[string]string + if p.peek() == '{' { + var err error + labels, err = p.scanLabelBlock() + if err != nil { + return Series{}, err + } + p.skipSpaces() + } + + valTok := p.scanToken() + if valTok == "" { + return Series{}, p.errorf("값이 없다") + } + val, err := strconv.ParseFloat(valTok, 64) + if err != nil { + return Series{}, p.errorf("값 파싱 실패 %q: %v", valTok, err) + } + p.skipSpaces() + + var tsMS int64 + if !p.atEnd() { + tsTok := p.scanToken() + tsMS, err = strconv.ParseInt(tsTok, 10, 64) + if err != nil { + return Series{}, p.errorf("타임스탬프 파싱 실패 %q: %v", tsTok, err) + } + p.skipSpaces() + if !p.atEnd() { + return Series{}, p.errorf("줄 끝에 처리되지 않은 내용이 남아 있다: %q", p.s[p.pos:]) + } + } + + return Series{Name: name, Labels: labels, Value: val, TimestampMS: tsMS}, nil +} + +// lineScanner 는 한 줄을 바이트 위치 기준으로 훑는 최소 상태 스캐너다. +// UTF-8 멀티바이트 문자의 continuation byte 는 항상 0x80 이상이라, 아래에서 +// 비교하는 ASCII 구두점(`{` `}` `"` `\` `,` `=` 공백)과 절대 충돌하지 않는다 — +// 바이트 단위로 스캔해도 UTF-8 라벨 값을 안전하게 그대로 복사할 수 있다. +type lineScanner struct { + s string + pos int + lineNo int +} + +func (p *lineScanner) errorf(format string, args ...any) error { + return fmt.Errorf("promtext: %d번째 줄: %s", p.lineNo, fmt.Sprintf(format, args...)) +} + +func (p *lineScanner) atEnd() bool { return p.pos >= len(p.s) } + +func (p *lineScanner) peek() byte { + if p.atEnd() { + return 0 + } + return p.s[p.pos] +} + +func isSpace(b byte) bool { return b == ' ' || b == '\t' } + +func (p *lineScanner) skipSpaces() { + for !p.atEnd() && isSpace(p.s[p.pos]) { + p.pos++ + } +} + +// scanName 은 공백 또는 '{' 전까지를 메트릭 이름으로 읽는다. +func (p *lineScanner) scanName() string { + start := p.pos + for !p.atEnd() && !isSpace(p.s[p.pos]) && p.s[p.pos] != '{' { + p.pos++ + } + return p.s[start:p.pos] +} + +// scanToken 은 공백으로 구분되는 다음 토큰(값 또는 타임스탬프)을 읽는다. +func (p *lineScanner) scanToken() string { + start := p.pos + for !p.atEnd() && !isSpace(p.s[p.pos]) { + p.pos++ + } + return p.s[start:p.pos] +} + +// scanLabelBlock 은 이미 '{' 를 확인한 상태에서 호출되어 라벨 블록 전체를 +// 문자 단위로 스캔한다. 콤마로 split 하지 않는다 — 라벨 값 안에 콤마가 올 수 +// 있어(이스케이프 대상이 아님) split 은 값을 깨뜨린다. +func (p *lineScanner) scanLabelBlock() (map[string]string, error) { + p.pos++ // consume '{' + labels := map[string]string{} + + p.skipSpaces() + if p.peek() == '}' { + p.pos++ + return labels, nil + } + + for { + p.skipSpaces() + nameStart := p.pos + for !p.atEnd() && p.s[p.pos] != '=' && p.s[p.pos] != '}' && p.s[p.pos] != ',' && !isSpace(p.s[p.pos]) { + p.pos++ + } + labelName := p.s[nameStart:p.pos] + if labelName == "" { + return nil, p.errorf("라벨 이름이 비어 있다") + } + p.skipSpaces() + if p.atEnd() || p.s[p.pos] != '=' { + return nil, p.errorf("라벨 %q 뒤에 '=' 이 없다", labelName) + } + p.pos++ // consume '=' + p.skipSpaces() + if p.atEnd() || p.s[p.pos] != '"' { + return nil, p.errorf("라벨 %q 값이 큰따옴표로 시작하지 않는다", labelName) + } + p.pos++ // consume opening '"' + + val, err := p.scanQuotedValue(labelName) + if err != nil { + return nil, err + } + labels[labelName] = val + + p.skipSpaces() + if p.atEnd() { + return nil, p.errorf("라벨 블록이 '}' 로 닫히지 않았다") + } + switch p.s[p.pos] { + case ',': + p.pos++ + // Prometheus 명세상 라벨 리스트의 trailing comma(`m{a="1",}`)는 + // 유효하다 — 콤마 소비 직후 곧바로 '}' 면 새 라벨을 기대하지 않고 + // 블록을 닫는다. 이 검사가 없으면 다음 루프의 라벨이름 스캔이 + // '}' 를 종료문자로만 취급해 0바이트를 읽고 "라벨 이름이 비어 + // 있다" 로 오분류한다. + p.skipSpaces() + if p.peek() == '}' { + p.pos++ + return labels, nil + } + continue + case '}': + p.pos++ + return labels, nil + default: + return nil, p.errorf("라벨 %q 값 뒤에 예기치 않은 문자 %q", labelName, p.s[p.pos]) + } + } +} + +// scanQuotedValue 는 여는 '"' 를 이미 소비한 상태에서 시작해 닫는 '"' 까지 +// Prometheus 명세의 이스케이프 3종(\\, \", \n)을 해석하며 문자 단위로 스캔한다. +func (p *lineScanner) scanQuotedValue(labelName string) (string, error) { + var b strings.Builder + for { + if p.atEnd() { + return "", p.errorf("라벨 %q 값의 큰따옴표가 닫히지 않았다", labelName) + } + c := p.s[p.pos] + switch c { + case '"': + p.pos++ + return b.String(), nil + case '\\': + p.pos++ + if p.atEnd() { + return "", p.errorf("라벨 %q 값이 이스케이프 문자로 끝났다", labelName) + } + switch p.s[p.pos] { + case '\\': + b.WriteByte('\\') + case '"': + b.WriteByte('"') + case 'n': + b.WriteByte('\n') + default: + return "", p.errorf("라벨 %q 값에 알 수 없는 이스케이프 시퀀스 \\%c", labelName, p.s[p.pos]) + } + p.pos++ + default: + b.WriteByte(c) + p.pos++ + } + } +} diff --git a/internal/promtext/parse_test.go b/internal/promtext/parse_test.go new file mode 100644 index 0000000..d375cda --- /dev/null +++ b/internal/promtext/parse_test.go @@ -0,0 +1,302 @@ +package promtext + +import ( + "bufio" + "errors" + "math" + "reflect" + "strings" + "testing" +) + +// TestParse_라벨_이스케이프_3종 은 백슬래시·큰따옴표·개행 이스케이프가 정확히 +// 복원되는지 확인한다. 원본은 raw 문자열(백틱)이라 Go 자체의 이스케이프 +// 해석을 거치지 않는다 — 실제 바이트가 곧 Prometheus exposition 규격의 +// 이스케이프 시퀀스다: +// +// a="x\\y" -> 바이트: x \ \ y -> 디코드: x \ y (이스케이프된 백슬래시 1개) +// b="q\"z" -> 바이트: q \ " z -> 디코드: q " z (이스케이프된 큰따옴표) +// c="l1\nl2"-> 바이트: l1 \ n l2 -> 디코드: l1 l2 (이스케이프된 개행) +func TestParse_라벨_이스케이프_3종(t *testing.T) { + input := `m{a="x\\y",b="q\"z",c="l1\nl2"} 1` + "\n" + + series, err := Parse(strings.NewReader(input)) + if err != nil { + t.Fatalf("Parse: 예기치 않은 에러: %v", err) + } + if len(series) != 1 { + t.Fatalf("시리즈 개수: got %d, want 1", len(series)) + } + s := series[0] + + if got, want := s.Labels["a"], `x\y`; got != want { + t.Errorf(`라벨 a (\\ 이스케이프): got %q, want %q`, got, want) + } + if got, want := s.Labels["b"], `q"z`; got != want { + t.Errorf(`라벨 b (\" 이스케이프): got %q, want %q`, got, want) + } + if got, want := s.Labels["c"], "l1\nl2"; got != want { + t.Errorf(`라벨 c (\n 이스케이프): got %q, want %q`, got, want) + } + if !strings.Contains(s.Labels["c"], "\n") { + t.Error("라벨 c 에 실제 개행 문자가 없다 — 리터럴 \\n 두 글자로 남아있으면 안 된다") + } +} + +// TestParse_NaN_Inf_와_타임스탬프 는 특수값 3종(NaN/+Inf/-Inf)과 타임스탬프 +// 유무 분기(있음/없음/있음)를 모두 실행한다. +func TestParse_NaN_Inf_와_타임스탬프(t *testing.T) { + input := "m1 NaN 1690000000000\nm2 +Inf\nm3 -Inf 5\n" + + series, err := Parse(strings.NewReader(input)) + if err != nil { + t.Fatalf("Parse: 예기치 않은 에러: %v", err) + } + if len(series) != 3 { + t.Fatalf("시리즈 개수: got %d, want 3", len(series)) + } + m1, m2, m3 := series[0], series[1], series[2] + + if !math.IsNaN(m1.Value) { + t.Errorf("m1.Value: got %v, want NaN", m1.Value) + } + if m1.TimestampMS != 1690000000000 { + t.Errorf("m1.TimestampMS: got %d, want 1690000000000", m1.TimestampMS) + } + + if !math.IsInf(m2.Value, 1) { + t.Errorf("m2.Value: got %v, want +Inf", m2.Value) + } + if m2.TimestampMS != 0 { + t.Errorf("m2.TimestampMS(타임스탬프 없음): got %d, want 0", m2.TimestampMS) + } + + if !math.IsInf(m3.Value, -1) { + t.Errorf("m3.Value: got %v, want -Inf", m3.Value) + } + if m3.TimestampMS != 5 { + t.Errorf("m3.TimestampMS: got %d, want 5", m3.TimestampMS) + } +} + +// TestParse_형식오류는_에러 는 constraints §테스트 함정("에러 경로 실행 의무")에 +// 따라, parse.go 가 반환하는 모든 에러 분기를 개별 입력으로 도달시키고 +// 각 분기 고유의 메시지 부분 문자열로 "맞는 분기가 실행됐는지"까지 단정한다 +// (단순 err!=nil 만 보면 엉뚱한 분기가 에러를 내도 통과하는 동어반복이 된다). +func TestParse_형식오류는_에러(t *testing.T) { + cases := []struct { + name string + input string + wantSubstr string + }{ + // 설계 계약이 지정한 2개 필수 픽스처. + {"라벨_블록_미닫힘", `m{a="1" 2`, "예기치 않은 문자"}, + {"값_없는_줄", `m{}`, "값이 없다"}, + // 아래는 구현이 실제로 만드는 나머지 에러 분기 — 각각 미도달 상태로 + // 남지 않도록 자체 보강. + {"메트릭_이름_없음", `{a="1"} 2`, "메트릭 이름이 비어 있다"}, + {"라벨_이름_없음", `m{="1"} 2`, "라벨 이름이 비어 있다"}, + {"라벨_등호_없음", `m{a} 2`, `'=' 이 없다`}, + {"라벨값_따옴표로_시작안함", `m{a=1} 2`, "큰따옴표로 시작하지 않는다"}, + {"라벨값_따옴표_안닫힘_EOF", `m{a="unterminated} 2`, "값의 큰따옴표가 닫히지 않았다"}, + {"라벨_블록이_EOF까지_안닫힘", `m{a="1"`, "라벨 블록이 '}'"}, + {"이스케이프_문자로_끝남", `m{a="x\`, "이스케이프 문자로 끝났다"}, + {"알수없는_이스케이프", `m{a="x\tb"} 1`, "알 수 없는 이스케이프"}, + {"값_파싱_실패", `m abc`, "값 파싱 실패"}, + {"타임스탬프_파싱_실패", `m 1 abc`, "타임스탬프 파싱 실패"}, + {"타임스탬프_뒤_잉여문자", `m 1 5 extra`, "처리되지 않은 내용"}, + } + + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + series, err := Parse(strings.NewReader(c.input + "\n")) + if err == nil { + t.Fatalf("입력 %q: 에러가 나야 하는데 성공했다 (series=%v)", c.input, series) + } + if !strings.Contains(err.Error(), c.wantSubstr) { + t.Errorf("에러 메시지에 %q 가 없다: %v", c.wantSubstr, err) + } + if !strings.Contains(err.Error(), "1번째 줄") { + t.Errorf("에러 메시지에 줄번호(1번째 줄)가 없다: %v", err) + } + }) + } +} + +// TestParse_여러줄_에러_줄번호가_정확하다 는 lineNo 추적이 하드코딩(예: 항상 1)이 +// 아니라 실제 스캔 진행에 연동됨을 별도로 확인한다 — 위 표의 모든 케이스가 +// 단일 줄이라 이 분기(2줄 이상 진행 후 에러)는 그것만으로 검증되지 않는다. +func TestParse_여러줄_에러_줄번호가_정확하다(t *testing.T) { + input := "a 1\nb 2\nc{x=1} 3\n" // 1·2번째 줄은 정상, 3번째 줄만 라벨값 따옴표 누락 + + _, err := Parse(strings.NewReader(input)) + if err == nil { + t.Fatal("3번째 줄에서 에러가 나야 한다") + } + if !strings.Contains(err.Error(), "3번째 줄") { + t.Fatalf("에러 메시지에 3번째 줄 표기가 없다: %v", err) + } + if strings.Contains(err.Error(), "1번째 줄") || strings.Contains(err.Error(), "2번째 줄") { + t.Fatalf("엉뚱한 줄번호가 찍혔다(하드코딩 의심): %v", err) + } +} + +// TestParse_주석과_라벨없는_평문 은 `# HELP`/`# TYPE` 사이에 낀 라벨 없는 평문 +// 한 줄이 주석은 결과에서 빠진 채 정확히 파싱되는지 확인한다. +func TestParse_주석과_라벨없는_평문(t *testing.T) { + input := "# HELP up 스크레이프 대상 상태\n# TYPE up gauge\nup 1\n" + + series, err := Parse(strings.NewReader(input)) + if err != nil { + t.Fatalf("Parse: 예기치 않은 에러: %v", err) + } + if len(series) != 1 { + t.Fatalf("주석 줄이 결과에 섞여 들어갔다: got %d 개, want 1", len(series)) + } + s := series[0] + if s.Name != "up" { + t.Errorf("Name: got %q, want %q", s.Name, "up") + } + if len(s.Labels) != 0 { + t.Errorf("Labels: got %v, want 빈 라벨", s.Labels) + } + if s.Value != 1 { + t.Errorf("Value: got %v, want 1", s.Value) + } +} + +// TestParse_공백만있는_줄도_빈줄로_취급 은 "빈 줄" 판정이 원본 문자열(line=="") +// 이 아니라 TrimSpace 결과에 근거함을 확인한다 — 공백/탭만 있는 줄을 빈 +// 문자열과 다르게 취급하는 구현이면 이 테스트가 실패한다. +func TestParse_공백만있는_줄도_빈줄로_취급(t *testing.T) { + input := "# c\n \nup 1\n\t\nup 2\n" + + series, err := Parse(strings.NewReader(input)) + if err != nil { + t.Fatalf("Parse: 예기치 않은 에러: %v", err) + } + if len(series) != 2 { + t.Fatalf("공백만 있는 줄이 파싱을 방해했다: got %d, want 2", len(series)) + } +} + +// TestParse_실제_nodevitals_픽스처 는 프롬프트가 지정한 실제 /metrics 형식 +// (nodevitals_hw_* 자체 메트릭 + node_exporter 계열 node_* 메트릭)을 HELP/TYPE +// 주석·빈 줄과 함께 섞어 라벨 맵·값까지 개별 단정한다. +func TestParse_실제_nodevitals_픽스처(t *testing.T) { + input := `# HELP nodevitals_hw_temp_celsius 하드웨어 온도 센서 값(섭씨) +# TYPE nodevitals_hw_temp_celsius gauge +nodevitals_hw_temp_celsius{node="e101",tier="core",device="hwmon0"} 42.5 + +# HELP node_cpu_seconds_total CPU 가 각 모드에서 소비한 누적 초 +# TYPE node_cpu_seconds_total counter +node_cpu_seconds_total{cpu="0",mode="idle"} 12345.6 +node_cpu_seconds_total{cpu="0",mode="user"} 678.9 +` + series, err := Parse(strings.NewReader(input)) + if err != nil { + t.Fatalf("Parse: 예기치 않은 에러: %v", err) + } + if len(series) != 3 { + t.Fatalf("시리즈 개수: got %d, want 3 (주석/빈줄 제외)", len(series)) + } + + temp := series[0] + if temp.Name != "nodevitals_hw_temp_celsius" { + t.Errorf("Name: got %q", temp.Name) + } + if want := (map[string]string{"node": "e101", "tier": "core", "device": "hwmon0"}); !reflect.DeepEqual(temp.Labels, want) { + t.Errorf("Labels: got %v, want %v", temp.Labels, want) + } + if temp.Value != 42.5 { + t.Errorf("Value: got %v, want 42.5", temp.Value) + } + if temp.TimestampMS != 0 { + t.Errorf("TimestampMS: got %d, want 0 (타임스탬프 없음)", temp.TimestampMS) + } + + idle := series[1] + if idle.Name != "node_cpu_seconds_total" { + t.Errorf("Name: got %q", idle.Name) + } + if want := (map[string]string{"cpu": "0", "mode": "idle"}); !reflect.DeepEqual(idle.Labels, want) { + t.Errorf("Labels: got %v, want %v", idle.Labels, want) + } + if idle.Value != 12345.6 { + t.Errorf("Value: got %v, want 12345.6", idle.Value) + } + + user := series[2] + if want := (map[string]string{"cpu": "0", "mode": "user"}); !reflect.DeepEqual(user.Labels, want) { + t.Errorf("Labels: got %v, want %v", user.Labels, want) + } + if user.Value != 678.9 { + t.Errorf("Value: got %v, want 678.9", user.Value) + } +} + +// TestParse_라벨블록_trailing_comma_허용 은 Prometheus 명세상 유효한 라벨 리스트 +// trailing comma(`m{a="1",} 1`)가 에러 없이 파싱되는지 확인한다 — 적대검증 +// 지적사항 회귀 테스트: scanLabelBlock 이 콤마 소비 후 재루프에서 '}' 를 +// 라벨 이름 시작 문자로 오인해 "라벨 이름이 비어 있다" 에러를 내던 결함. +// 콤마 뒤 공백이 낀 경우(`m{a="1", } 1`)까지 별도 입력으로 함께 확인해 +// 새로 추가한 skipSpaces 분기까지 실행시킨다(동어반복 방지). +func TestParse_라벨블록_trailing_comma_허용(t *testing.T) { + cases := []struct { + name string + input string + }{ + {"콤마_직후_닫힘", `m{a="1",} 1`}, + {"콤마_공백_닫힘", `m{a="1", } 1`}, + } + + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + series, err := Parse(strings.NewReader(c.input + "\n")) + if err != nil { + t.Fatalf("Parse(%q): trailing comma 인데 에러가 났다: %v", c.input, err) + } + if len(series) != 1 { + t.Fatalf("시리즈 개수: got %d, want 1", len(series)) + } + s := series[0] + if s.Name != "m" { + t.Errorf("Name: got %q, want %q", s.Name, "m") + } + if want := (map[string]string{"a": "1"}); !reflect.DeepEqual(s.Labels, want) { + t.Errorf("Labels: got %v, want %v", s.Labels, want) + } + if s.Value != 1 { + t.Errorf("Value: got %v, want 1", s.Value) + } + }) + } +} + +// TestParse_빈_입력은_빈_결과 는 0바이트 입력이 에러 없이 빈 결과를 내는지 +// 확인한다 (스캔 루프가 한 번도 안 돈 경계 케이스). +func TestParse_빈_입력은_빈_결과(t *testing.T) { + series, err := Parse(strings.NewReader("")) + if err != nil { + t.Fatalf("Parse(빈 입력): 예기치 않은 에러: %v", err) + } + if len(series) != 0 { + t.Fatalf("빈 입력인데 시리즈가 나왔다: %v", series) + } +} + +// TestParse_줄이_1MiB_초과면_스캔에러 는 Parse 가 sc.Err() 를 확인하는 분기 +// (bufio.Scanner 버퍼 상한 초과)가 실제로 도달 가능하고 bufio.ErrTooLong 을 +// %w 로 래핑함을 errors.Is 로 단정한다 — constraints §테스트 함정의 +// "특정 에러는 errors.Is 로 단정" 지침 준수. +func TestParse_줄이_1MiB_초과면_스캔에러(t *testing.T) { + huge := strings.Repeat("x", maxLineBytes+1024) + + _, err := Parse(strings.NewReader(huge + "\n")) + if err == nil { + t.Fatal("1MiB 초과 줄인데 에러가 안 났다") + } + if !errors.Is(err, bufio.ErrTooLong) { + t.Fatalf("bufio.ErrTooLong 이 래핑되어 있지 않다: %v", err) + } +} diff --git a/internal/scrape/scrape.go b/internal/scrape/scrape.go new file mode 100644 index 0000000..a1f6da7 --- /dev/null +++ b/internal/scrape/scrape.go @@ -0,0 +1,201 @@ +// Package scrape 는 discovery.Discoverer 가 낸 타겟들의 /metrics 를 주기적으로 +// GET 해 promtext 로 파싱하고 tsdb.DB 에 적재한다. 타겟 하나의 실패가 다른 +// 타겟의 적재를 막지 않도록 격리한다(m2-design.md §4 계약). +package scrape + +import ( + "context" + "fmt" + "log/slog" + "net/http" + "sync/atomic" + "time" + + "github.com/KeiaiLab/nodevitals-observatory/internal/discovery" + "github.com/KeiaiLab/nodevitals-observatory/internal/labels" + "github.com/KeiaiLab/nodevitals-observatory/internal/promtext" +) + +const ( + defaultInterval = 15 * time.Second + defaultTimeout = 10 * time.Second + + // 자체 관측 메트릭 이름 + 두 라벨 매핑에 공통으로 쓰는 instance 라벨 이름. + metricUp = "observatory_up" + metricScrapeSamples = "observatory_scrape_samples" + instanceLabel = "instance" +) + +// Options 는 Scraper 동작을 결정한다. +type Options struct { + Interval time.Duration // Run 루프 주기. 0 이면 15s. + Timeout time.Duration // 타겟당 HTTP 타임아웃. 0 이면 10s. + Client *http.Client // nil 이면 &http.Client{Timeout: Timeout} + // Now 는 밀리초 epoch 을 낸다. nil 이면 time.Now 기반(cmd 전용) — 테스트는 + // 반드시 고정 클록을 주입해야 한다(결정론 계약). + Now func() int64 +} + +// Appender 는 스크레이퍼가 저장 계층에 요구하는 전부다. tsdb.DB 가 이 +// 시그니처를 만족하므로 별도 어댑터가 필요 없고, 테스트는 저장 엔진 없이 +// 가짜 하나로 끝난다. +type Appender interface { + Append(lset labels.Labels, t int64, v float64) error +} + +// Scraper 는 하나의 discovery.Discoverer 로 얻은 타겟들을 주기적으로 긁어 +// Appender 에 적재한다. +type Scraper struct { + db Appender + d discovery.Discoverer + opts Options + + ready atomic.Bool +} + +// NewScraper 는 opts 의 미설정 필드를 기본값으로 채운 뒤 Scraper 를 만든다. +func NewScraper(db Appender, d discovery.Discoverer, opts Options) *Scraper { + if opts.Interval <= 0 { + opts.Interval = defaultInterval + } + if opts.Timeout <= 0 { + opts.Timeout = defaultTimeout + } + if opts.Client == nil { + opts.Client = &http.Client{Timeout: opts.Timeout} + } + if opts.Now == nil { + opts.Now = func() int64 { return time.Now().UnixMilli() } + } + return &Scraper{db: db, d: d, opts: opts} +} + +// ScrapeOnce 는 Discover → 타겟별 GET → promtext.Parse → Append 1 사이클을 +// 수행한다. 타겟 간 실패는 격리한다 — 실패한 타겟은 observatory_up=0 만 +// 남기고 다음 타겟으로 진행한다. 반환 에러는 Discover 실패처럼 사이클 +// 전체가 불가능한 경우만이다. +func (s *Scraper) ScrapeOnce(ctx context.Context) error { + // Ready 는 "첫 사이클이 완주(에러 무관)했는가" 계약이다 — Discover 가 + // 실패해도 시도 자체는 끝난 것으로 보고 readyz 를 열어야 dead-target + // 환경에서도 API 조회가 가능하다. defer 로 모든 반환 경로를 덮는다. + defer s.ready.Store(true) + + targets, err := s.d.Discover(ctx) + if err != nil { + return fmt.Errorf("scrape: 타겟 발견 실패: %w", err) + } + + now := s.opts.Now() + for _, target := range targets { + s.scrapeTarget(ctx, target, now) + } + return nil +} + +// scrapeTarget 은 타겟 하나를 긁어 저장하고, 성공/실패와 무관하게 자체 관측 +// 메트릭(observatory_up, observatory_scrape_samples)을 남긴다. 타겟 실패를 +// 이 함수 밖으로 전파하지 않는다 — 실패 격리는 ScrapeOnce 의 계약이다. +func (s *Scraper) scrapeTarget(ctx context.Context, target discovery.Target, now int64) { + n, err := s.scrapeAndStore(ctx, target, now) + + up := 1.0 + if err != nil { + slog.Warn("scrape: 타겟 실패", "target", target.Name, "url", target.URL, "error", err) + up = 0 + n = 0 + } + s.appendSelfMetric(metricUp, target.Name, now, up) + s.appendSelfMetric(metricScrapeSamples, target.Name, now, float64(n)) +} + +// scrapeAndStore 는 타겟의 /metrics 를 GET 해 파싱하고 DB 에 적재한다. +// 반환값은 성공적으로 Append 된 샘플 수다. GET·비200·파싱 실패만 에러로 +// 반환한다 — 개별 샘플의 Append 실패는 그 샘플만 건너뛰고 계속한다(부분 +// 성공이 전체 실패보다 낫다). +func (s *Scraper) scrapeAndStore(ctx context.Context, target discovery.Target, now int64) (int, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, target.URL, nil) + if err != nil { + return 0, fmt.Errorf("요청 생성 실패: %w", err) + } + + resp, err := s.opts.Client.Do(req) + if err != nil { + return 0, fmt.Errorf("GET 실패: %w", err) + } + defer resp.Body.Close() + + if resp.StatusCode != http.StatusOK { + return 0, fmt.Errorf("비정상 응답 %d", resp.StatusCode) + } + + series, err := promtext.Parse(resp.Body) + if err != nil { + return 0, fmt.Errorf("파싱 실패: %w", err) + } + + n := 0 + for _, sr := range series { + lset := buildTargetLabels(sr, target.Name) + if err := s.db.Append(lset, now, sr.Value); err != nil { + slog.Warn("scrape: append 실패", "target", target.Name, "metric", sr.Name, "error", err) + continue + } + n++ + } + return n, nil +} + +// buildTargetLabels 는 파싱된 시리즈 하나를 TSDB 라벨셋으로 변환한다. +// exposition 라벨을 베이스로 하되 __name__ 과 instance 는 스크레이퍼가 +// 확정한 값으로 덮는다 — instance 는 타겟 exposition 자체의 동명 라벨보다 +// 우선한다(M2 계약, m2-design.md §4 "스크레이퍼 관례"). +func buildTargetLabels(sr promtext.Series, instance string) labels.Labels { + m := make(map[string]string, len(sr.Labels)+2) + for k, v := range sr.Labels { + m[k] = v + } + m[labels.MetricName] = sr.Name + m[instanceLabel] = instance + return labels.LabelsFromMap(m) +} + +// appendSelfMetric 은 observatory_up/observatory_scrape_samples 한 점을 +// 남긴다. Append 자체가 실패해도(디스크 문제 등) 사이클을 막지 않고 경고만 +// 남긴다 — 자체 관측 메트릭의 실패로 다른 타겟 처리를 중단할 이유가 없다. +func (s *Scraper) appendSelfMetric(name, instance string, t int64, v float64) { + lset := labels.LabelsFromMap(map[string]string{ + labels.MetricName: name, + instanceLabel: instance, + }) + if err := s.db.Append(lset, t, v); err != nil { + slog.Warn("scrape: 자체 메트릭 append 실패", "metric", name, "instance", instance, "error", err) + } +} + +// Run 은 즉시 1회 ScrapeOnce 를 실행한 뒤 Interval 마다 반복한다. ctx 취소로 +// 종료한다(에러 반환 없음) — 사이클 에러는 로그만 남기고 다음 tick 에 +// 재시도한다. +func (s *Scraper) Run(ctx context.Context) { + s.runOnce(ctx) + + ticker := time.NewTicker(s.opts.Interval) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + s.runOnce(ctx) + } + } +} + +func (s *Scraper) runOnce(ctx context.Context) { + if err := s.ScrapeOnce(ctx); err != nil { + slog.Error("scrape: 사이클 실패", "error", err) + } +} + +// Ready 는 첫 ScrapeOnce 가 완주(에러 무관)한 뒤 true 다 — apiserver +// /readyz 배선용. +func (s *Scraper) Ready() bool { return s.ready.Load() } diff --git a/internal/scrape/scrape_test.go b/internal/scrape/scrape_test.go new file mode 100644 index 0000000..3f7893f --- /dev/null +++ b/internal/scrape/scrape_test.go @@ -0,0 +1,286 @@ +package scrape + +import ( + "context" + "errors" + "fmt" + "github.com/KeiaiLab/nodevitals-observatory/internal/labels" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/KeiaiLab/nodevitals-observatory/internal/discovery" + "github.com/KeiaiLab/nodevitals-observatory/internal/tsdb" +) + +// discoverFunc 는 함수 하나를 discovery.Discoverer 로 어댑트한다(테스트 전용 +// stub — 실 네트워크·k8s 의존 없이 Discover 의 반환값을 자유롭게 주입한다). +type discoverFunc func(ctx context.Context) ([]discovery.Target, error) + +func (f discoverFunc) Discover(ctx context.Context) ([]discovery.Target, error) { return f(ctx) } + +// openTestDB 는 임시 디렉터리에 실 tsdb.DB 를 연다(왕복 검증 — mock 아님). +func openTestDB(t *testing.T) *tsdb.DB { + t.Helper() + db, err := tsdb.Open(tsdb.DefaultOptions(t.TempDir())) + if err != nil { + t.Fatalf("tsdb.Open: %v", err) + } + t.Cleanup(func() { _ = db.Close() }) + return db +} + +// latestSample 은 name{instance=instance} 시리즈의 최신(유일) 샘플 값을 낸다. +// 시리즈가 정확히 1개가 아니면 즉시 실패한다 — self-metric 은 사이클마다 +// 타겟당 한 점만 남아야 하므로 개수 자체가 계약이다. +func latestSample(t *testing.T, q *tsdb.Querier, name, instance string) float64 { + t.Helper() + mName, err := labels.NewMatcher(labels.MatchEqual, labels.MetricName, name) + if err != nil { + t.Fatalf("NewMatcher(%q): %v", name, err) + } + mInst, err := labels.NewMatcher(labels.MatchEqual, instanceLabel, instance) + if err != nil { + t.Fatalf("NewMatcher(instance=%q): %v", instance, err) + } + series := q.Select(mName, mInst) + if len(series) != 1 { + t.Fatalf("%s{instance=%s} 시리즈 개수 = %d, want 1", name, instance, len(series)) + } + it := series[0].Iterator() + if !it.Next() { + t.Fatalf("%s{instance=%s} 샘플이 없다", name, instance) + } + _, v := it.At() + return v +} + +// ---- ScrapeOnce: 저장 + 라벨 매핑 (왕복) ---- + +func TestScrapeOnce_저장과_라벨매핑(t *testing.T) { + const t0 = int64(1_700_000_000_000) // 임의 고정 ms epoch — 결정론 계약 + + mux := http.NewServeMux() + mux.HandleFunc("/metrics", func(w http.ResponseWriter, r *http.Request) { + fmt.Fprintln(w, `node_load1{node="e101"} 1.5`) + }) + srv := httptest.NewServer(mux) + defer srv.Close() + + db := openTestDB(t) + disc := discovery.NewStatic([]discovery.Target{{Name: "e101-9847", URL: srv.URL + "/metrics"}}) + scr := NewScraper(db, disc, Options{Now: func() int64 { return t0 }}) + + if err := scr.ScrapeOnce(context.Background()); err != nil { + t.Fatalf("ScrapeOnce: %v", err) + } + + q, closeQ, err := db.Querier(t0-1, t0+1) + if err != nil { + t.Fatalf("Querier: %v", err) + } + defer closeQ() + + m, err := labels.NewMatcher(labels.MatchEqual, labels.MetricName, "node_load1") + if err != nil { + t.Fatalf("NewMatcher: %v", err) + } + series := q.Select(m) + if len(series) != 1 { + t.Fatalf("node_load1 시리즈 개수 = %d, want 1", len(series)) + } + + lset := series[0].Labels() + // __name__ 변환이 틀렸으면 애초에 위 Select 가 0개를 냈을 것이므로, 여기서는 + // exposition 라벨(node)과 스크레이퍼가 덮은 instance 라벨을 개별 단정한다. + if got := lset.Get("node"); got != "e101" { + t.Errorf("node 라벨 = %q, want e101", got) + } + if got := lset.Get(instanceLabel); got != "e101-9847" { + t.Errorf("instance 라벨 = %q, want e101-9847 (target.Name 오버라이드 계약)", got) + } + + it := series[0].Iterator() + if !it.Next() { + t.Fatalf("Iterator.Next() = false, want true") + } + gotT, gotV := it.At() + // gotT 가 t0 와 다르면 Now() 주입이 무시되고 실제 wall clock 이 쓰였다는 + // 뜻이다 — 그 경우 애초에 [t0-1,t0+1] 창 밖이라 위 Select 에서 이미 0개로 + // 걸렸겠지만, 값 자체도 한 번 더 명시적으로 단정해 둔다. + if gotT != t0 { + t.Errorf("샘플 시각 = %d, want %d (nowMS 주입 계약)", gotT, t0) + } + if gotV != 1.5 { + t.Errorf("샘플 값 = %v, want 1.5", gotV) + } + if it.Next() { + t.Errorf("샘플이 1개여야 하는데 2개 이상 존재한다") + } + if err := it.Err(); err != nil { + t.Errorf("Iterator.Err() = %v, want nil", err) + } +} + +// ---- ScrapeOnce: 실패 타겟 격리 + up/samples 메트릭 ---- + +func TestScrapeOnce_실패타겟_격리와_up메트릭(t *testing.T) { + const t0 = int64(1_700_000_000_000) + + okMux := http.NewServeMux() + okMux.HandleFunc("/metrics", func(w http.ResponseWriter, r *http.Request) { + fmt.Fprintln(w, `up 1`) + fmt.Fprintln(w, `node_load1{node="e102"} 2.5`) + }) + okSrv := httptest.NewServer(okMux) + defer okSrv.Close() + + failMux := http.NewServeMux() + failMux.HandleFunc("/metrics", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + fmt.Fprintln(w, "boom") + }) + failSrv := httptest.NewServer(failMux) + defer failSrv.Close() + + db := openTestDB(t) + disc := discovery.NewStatic([]discovery.Target{ + {Name: "ok-target", URL: okSrv.URL + "/metrics"}, + {Name: "fail-target", URL: failSrv.URL + "/metrics"}, + }) + scr := NewScraper(db, disc, Options{Now: func() int64 { return t0 }}) + + if err := scr.ScrapeOnce(context.Background()); err != nil { + t.Fatalf("ScrapeOnce 반환 에러 = %v, want nil (타겟 실패는 격리되어야 한다)", err) + } + + q, closeQ, err := db.Querier(t0-1, t0+1) + if err != nil { + t.Fatalf("Querier: %v", err) + } + defer closeQ() + + // 정상 타겟의 실 샘플(node_load1)이 착지했는지 — 격리가 안 되고 사이클 + // 전체가 죽는 구현이면 여기서 0개가 나와 실패한다. + mLoad, err := labels.NewMatcher(labels.MatchEqual, labels.MetricName, "node_load1") + if err != nil { + t.Fatalf("NewMatcher: %v", err) + } + if got := q.Select(mLoad); len(got) != 1 { + t.Fatalf("정상 타겟 node_load1 시리즈 개수 = %d, want 1", len(got)) + } + + if v := latestSample(t, q, metricUp, "fail-target"); v != 0 { + t.Errorf("observatory_up{instance=fail-target} = %v, want 0", v) + } + if v := latestSample(t, q, metricUp, "ok-target"); v != 1 { + t.Errorf("observatory_up{instance=ok-target} = %v, want 1", v) + } + // ok 타겟은 "up 1" + "node_load1{...} 2.5" 두 줄을 노출한다 → 실 샘플 수 2. + // 이 값을 하드코딩하지 않고 타겟이 실제로 낸 라인 수와 맞춰, 구현이 + // 상수를 그냥 반환해도 통과하는 동어반복이 되지 않게 한다. + if v := latestSample(t, q, metricScrapeSamples, "ok-target"); v != 2 { + t.Errorf("observatory_scrape_samples{instance=ok-target} = %v, want 2", v) + } + // 실패 타겟도 매 사이클·타겟마다 samples=0 을 남긴다는 계약(design §4)까지 + // 확인한다 — "실패 타겟은 관측 메트릭 자체를 안 남긴다"로 잘못 구현하면 + // latestSample 이 시리즈 0개를 만나 Fatalf 로 잡아낸다. + if v := latestSample(t, q, metricScrapeSamples, "fail-target"); v != 0 { + t.Errorf("observatory_scrape_samples{instance=fail-target} = %v, want 0", v) + } +} + +// ---- ScrapeOnce: Discover 자체 실패 ---- + +func TestScrapeOnce_Discover실패시_에러반환_그러나_Ready는_true(t *testing.T) { + db := openTestDB(t) + wantErr := errors.New("발견 실패 주입") + disc := discoverFunc(func(_ context.Context) ([]discovery.Target, error) { return nil, wantErr }) + scr := NewScraper(db, disc, Options{Now: func() int64 { return 1_700_000_000_000 }}) + + err := scr.ScrapeOnce(context.Background()) + if err == nil { + t.Fatal("ScrapeOnce() 에러 = nil, want non-nil (Discover 실패가 전파돼야 한다)") + } + if !errors.Is(err, wantErr) { + t.Errorf("ScrapeOnce() 에러 = %v, want %v 를 wrap", err, wantErr) + } + if !scr.Ready() { + t.Errorf("Discover 실패 후 Ready() = false, want true (완주=에러 무관 계약)") + } +} + +// ---- Ready ---- + +func TestReady_첫사이클_후_true(t *testing.T) { + db := openTestDB(t) + + failMux := http.NewServeMux() + failMux.HandleFunc("/metrics", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusInternalServerError) + }) + failSrv := httptest.NewServer(failMux) + defer failSrv.Close() + + disc := discovery.NewStatic([]discovery.Target{{Name: "dead", URL: failSrv.URL + "/metrics"}}) + scr := NewScraper(db, disc, Options{Now: func() int64 { return 1_700_000_000_000 }}) + + if scr.Ready() { + t.Fatal("생성 직후 Ready() = true, want false") + } + + // "실패-only" 사이클 — 유일한 타겟이 500 을 반환해 샘플이 하나도 안 붙는다. + // "성공 시에만 ready" 로 잘못 구현하면(예: samples>0 조건부 Store) 아래 + // 단정이 실패한다 — dead-target 환경에서도 파드는 Ready 여야 한다. + if err := scr.ScrapeOnce(context.Background()); err != nil { + t.Fatalf("ScrapeOnce(실패-only 사이클) 에러 = %v, want nil (타겟 실패는 사이클 에러가 아니다)", err) + } + + if !scr.Ready() { + t.Error("실패-only 사이클 후 Ready() = false, want true") + } +} + +// ---- Run: 즉시 1회 + ctx 취소 종료 ---- + +func TestRun_즉시_1회_실행하고_ctx취소로_종료한다(t *testing.T) { + db := openTestDB(t) + + called := make(chan struct{}, 1) + disc := discoverFunc(func(_ context.Context) ([]discovery.Target, error) { + select { + case called <- struct{}{}: + default: + } + return nil, nil + }) + + scr := NewScraper(db, disc, Options{ + // tick 이 테스트 창 안에서 절대 일어나지 않을 만큼 크게 잡는다 — 검증 + // 대상은 "즉시 1회 실행" 과 "ctx 취소 종료" 뿐, 재-tick 타이밍이 아니다. + Interval: time.Hour, + Now: func() int64 { return 1_700_000_000_000 }, + }) + + ctx, cancel := context.WithCancel(context.Background()) + done := make(chan struct{}) + go func() { + scr.Run(ctx) + close(done) + }() + + select { + case <-called: + case <-time.After(2 * time.Second): + t.Fatal("Run() 이 즉시 1회 Discover(ScrapeOnce) 를 호출하지 않았다") + } + + cancel() + + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("ctx 취소 후에도 Run() 이 반환하지 않았다(ctx.Done() 미선택 의심)") + } +}