blob: 284e3935bd56c485f08c32ceba6d0651cd32605f [file] [log] [blame]
giof6ad2982024-08-23 17:42:49 +04001package installer
2
3import (
4 "bytes"
giof6ad2982024-08-23 17:42:49 +04005 "encoding/json"
giob7df27f2026-07-28 10:36:17 +04006 "errors"
giof6ad2982024-08-23 17:42:49 +04007 "fmt"
8 "io"
9 "net"
10 "net/http"
11 "path/filepath"
12 "strconv"
13 "strings"
14 "text/template"
15
16 "github.com/giolekva/pcloud/core/installer/soft"
17
gio721c0042025-04-03 11:56:36 +040018 corev1 "k8s.io/api/core/v1"
19 "k8s.io/apimachinery/pkg/util/intstr"
giof6ad2982024-08-23 17:42:49 +040020 "sigs.k8s.io/yaml"
21)
22
23type ClusterNetworkConfigurator interface {
24 AddCluster(name string, ingressIP net.IP) error
25 RemoveCluster(name string, ingressIP net.IP) error
gio721c0042025-04-03 11:56:36 +040026 AddProxy(src int, dst string, protocol Protocol) (string, error)
27 AddIngressProxy(src, dst string) error
28 RemoveProxy(src int, dst string, protocol Protocol) error
29 RemoveIngressProxy(src, dst string) error
giof6ad2982024-08-23 17:42:49 +040030}
31
32type NginxProxyConfigurator struct {
33 PrivateSubdomain string
34 DNSAPIAddr string
35 Repo soft.RepoIO
gio721c0042025-04-03 11:56:36 +040036 ConfigPath string
37 ServicePath string
giof6ad2982024-08-23 17:42:49 +040038}
39
40type createARecordReq struct {
41 Entry string `json:"entry"`
42 IP net.IP `json:"text"`
43}
44
45func (c *NginxProxyConfigurator) AddCluster(name string, ingressIP net.IP) error {
46 req := createARecordReq{
47 Entry: fmt.Sprintf("*.%s.cluster.%s", name, c.PrivateSubdomain),
48 IP: ingressIP,
49 }
50 var buf bytes.Buffer
51 if err := json.NewEncoder(&buf).Encode(req); err != nil {
52 return err
53 }
54 resp, err := http.Post(fmt.Sprintf("%s/create-a-record", c.DNSAPIAddr), "application/json", &buf)
55 if err != nil {
56 return err
57 }
58 if resp.StatusCode != http.StatusOK {
59 var buf bytes.Buffer
60 io.Copy(&buf, resp.Body)
giob7df27f2026-07-28 10:36:17 +040061 return errors.New(buf.String())
giof6ad2982024-08-23 17:42:49 +040062 }
63 return nil
64}
65
66func (c *NginxProxyConfigurator) RemoveCluster(name string, ingressIP net.IP) error {
67 req := createARecordReq{
68 Entry: fmt.Sprintf("*.%s.cluster.%s", name, c.PrivateSubdomain),
69 IP: ingressIP,
70 }
71 var buf bytes.Buffer
72 if err := json.NewEncoder(&buf).Encode(req); err != nil {
73 return err
74 }
75 resp, err := http.Post(fmt.Sprintf("%s/delete-a-record", c.DNSAPIAddr), "application/json", &buf)
76 if err != nil {
77 return err
78 }
79 if resp.StatusCode != http.StatusOK {
80 var buf bytes.Buffer
81 io.Copy(&buf, resp.Body)
giob7df27f2026-07-28 10:36:17 +040082 return errors.New(buf.String())
giof6ad2982024-08-23 17:42:49 +040083 }
84 return nil
85}
86
gio721c0042025-04-03 11:56:36 +040087func (c *NginxProxyConfigurator) AddProxy(src int, dst string, protocol Protocol) (string, error) {
88 var namespace string
giof6ad2982024-08-23 17:42:49 +040089 _, err := c.Repo.Do(func(fs soft.RepoFS) (string, error) {
gio721c0042025-04-03 11:56:36 +040090 if err := func() error {
91 r, err := fs.Reader(c.ServicePath)
92 if err != nil {
93 return err
94 }
95 defer r.Close()
96 var buf bytes.Buffer
97 if _, err := io.Copy(&buf, r); err != nil {
98 return err
99 }
100 var svc corev1.Service
101 if err := yaml.Unmarshal(buf.Bytes(), &svc); err != nil {
102 return err
103 }
104 svc.Spec.Ports = append(svc.Spec.Ports, corev1.ServicePort{
105 Name: fmt.Sprintf("p%d", src),
106 Protocol: corev1.Protocol(ProtocolToString(protocol)),
107 Port: int32(src),
108 TargetPort: intstr.FromInt(src),
109 })
110 w, err := fs.Writer(c.ServicePath)
111 if err != nil {
112 return err
113 }
114 defer w.Close()
115 tmp, err := yaml.Marshal(svc)
116 if err != nil {
117 return err
118 }
119 if _, err := io.Copy(w, bytes.NewReader(tmp)); err != nil {
120 return err
121 }
122 return nil
123 }(); err != nil {
124 return "", err
125 }
giof6ad2982024-08-23 17:42:49 +0400126 cfg, err := func() (NginxProxyConfig, error) {
gio721c0042025-04-03 11:56:36 +0400127 r, err := fs.Reader(c.ConfigPath)
giof6ad2982024-08-23 17:42:49 +0400128 if err != nil {
129 return NginxProxyConfig{}, err
130 }
131 defer r.Close()
132 return ParseNginxProxyConfig(r)
133 }()
134 if err != nil {
135 return "", err
136 }
gio721c0042025-04-03 11:56:36 +0400137 namespace = cfg.Namespace
138 var proxyMap map[int]string
139 switch protocol {
140 case ProtocolTCP:
141 proxyMap = cfg.TCP
142 case ProtocolUDP:
143 proxyMap = cfg.UDP
144 default:
145 return "", fmt.Errorf("invalid protocol: %v", protocol)
giof6ad2982024-08-23 17:42:49 +0400146 }
gio721c0042025-04-03 11:56:36 +0400147 // TODO(gio): check for already existing mapping
148 proxyMap[src] = dst
149 w, err := fs.Writer(c.ConfigPath)
giof6ad2982024-08-23 17:42:49 +0400150 if err != nil {
151 return "", err
152 }
153 defer w.Close()
giof55ab362025-04-11 17:48:17 +0400154 if err := cfg.Render(w); err != nil {
giof6ad2982024-08-23 17:42:49 +0400155 return "", err
156 }
gio721c0042025-04-03 11:56:36 +0400157 nginxPath := filepath.Join(filepath.Dir(c.ConfigPath), "ingress-nginx.yaml")
giof6ad2982024-08-23 17:42:49 +0400158 nginx, err := func() (map[string]any, error) {
159 r, err := fs.Reader(nginxPath)
160 if err != nil {
161 return nil, err
162 }
163 defer r.Close()
164 var buf bytes.Buffer
165 if _, err := io.Copy(&buf, r); err != nil {
166 return nil, err
167 }
168 ret := map[string]any{}
169 if err := yaml.Unmarshal(buf.Bytes(), &ret); err != nil {
170 return nil, err
171 }
172 return ret, nil
173 }()
174 if err != nil {
175 return "", err
176 }
giof6ad2982024-08-23 17:42:49 +0400177 buf, err := yaml.Marshal(nginx)
178 if err != nil {
179 return "", err
180 }
181 w, err = fs.Writer(nginxPath)
182 if err != nil {
183 return "", err
184 }
185 defer w.Close()
186 if _, err := io.Copy(w, bytes.NewReader(buf)); err != nil {
187 return "", err
188 }
gio721c0042025-04-03 11:56:36 +0400189 return fmt.Sprintf("add proxy mapping: %d %s", src, dst), nil
giof6ad2982024-08-23 17:42:49 +0400190 })
gio721c0042025-04-03 11:56:36 +0400191 if err != nil {
192 return "", err
193 }
194 return namespace, nil
giof6ad2982024-08-23 17:42:49 +0400195}
196
gio721c0042025-04-03 11:56:36 +0400197func (c *NginxProxyConfigurator) AddIngressProxy(src, dst string) error {
giof6ad2982024-08-23 17:42:49 +0400198 _, err := c.Repo.Do(func(fs soft.RepoFS) (string, error) {
199 cfg, err := func() (NginxProxyConfig, error) {
gio721c0042025-04-03 11:56:36 +0400200 r, err := fs.Reader(c.ConfigPath)
giof6ad2982024-08-23 17:42:49 +0400201 if err != nil {
202 return NginxProxyConfig{}, err
203 }
204 defer r.Close()
205 return ParseNginxProxyConfig(r)
206 }()
207 if err != nil {
208 return "", err
209 }
gio721c0042025-04-03 11:56:36 +0400210 if v, ok := cfg.Ingress[src]; ok && v != dst {
211 return "", fmt.Errorf("wrong mapping %s already exists (%s)", src, v)
giof6ad2982024-08-23 17:42:49 +0400212 }
gio721c0042025-04-03 11:56:36 +0400213 cfg.Ingress[src] = dst
214 w, err := fs.Writer(c.ConfigPath)
giof6ad2982024-08-23 17:42:49 +0400215 if err != nil {
216 return "", err
217 }
218 defer w.Close()
giof55ab362025-04-11 17:48:17 +0400219 if err := cfg.Render(w); err != nil {
giof6ad2982024-08-23 17:42:49 +0400220 return "", err
221 }
gio721c0042025-04-03 11:56:36 +0400222 nginxPath := filepath.Join(filepath.Dir(c.ConfigPath), "ingress-nginx.yaml")
223 nginx, err := func() (map[string]any, error) {
224 r, err := fs.Reader(nginxPath)
225 if err != nil {
226 return nil, err
227 }
228 defer r.Close()
229 var buf bytes.Buffer
230 if _, err := io.Copy(&buf, r); err != nil {
231 return nil, err
232 }
233 ret := map[string]any{}
234 if err := yaml.Unmarshal(buf.Bytes(), &ret); err != nil {
235 return nil, err
236 }
237 return ret, nil
238 }()
239 if err != nil {
240 return "", err
241 }
gio721c0042025-04-03 11:56:36 +0400242 buf, err := yaml.Marshal(nginx)
243 if err != nil {
244 return "", err
245 }
246 w, err = fs.Writer(nginxPath)
247 if err != nil {
248 return "", err
249 }
250 defer w.Close()
251 if _, err := io.Copy(w, bytes.NewReader(buf)); err != nil {
252 return "", err
253 }
254 return fmt.Sprintf("add ingress proxy mapping: %s %s", src, dst), nil
255 })
256 return err
257}
258
259func (c *NginxProxyConfigurator) RemoveProxy(src int, dst string, protocol Protocol) error {
260 _, err := c.Repo.Do(func(fs soft.RepoFS) (string, error) {
261 if err := func() error {
262 r, err := fs.Reader(c.ServicePath)
263 if err != nil {
264 return err
265 }
266 defer r.Close()
267 var buf bytes.Buffer
268 if _, err := io.Copy(&buf, r); err != nil {
269 return err
270 }
271 var svc corev1.Service
272 if err := yaml.Unmarshal(buf.Bytes(), &svc); err != nil {
273 return err
274 }
275 for i, p := range svc.Spec.Ports {
276 if p.Port == int32(src) {
277 svc.Spec.Ports = append(svc.Spec.Ports[:i], svc.Spec.Ports[i+1:]...)
278 break
279 }
280 }
281 w, err := fs.Writer(c.ServicePath)
282 if err != nil {
283 return err
284 }
285 defer w.Close()
286 tmp, err := yaml.Marshal(svc)
287 if err != nil {
288 return err
289 }
290 if _, err := io.Copy(w, bytes.NewReader(tmp)); err != nil {
291 return err
292 }
293 return nil
294 }(); err != nil {
295 return "", err
296 }
297 cfg, err := func() (NginxProxyConfig, error) {
298 r, err := fs.Reader(c.ConfigPath)
299 if err != nil {
300 return NginxProxyConfig{}, err
301 }
302 defer r.Close()
303 return ParseNginxProxyConfig(r)
304 }()
305 if err != nil {
306 return "", err
307 }
308 var proxyMap map[int]string
309 switch protocol {
310 case ProtocolTCP:
311 proxyMap = cfg.TCP
312 case ProtocolUDP:
313 proxyMap = cfg.UDP
314 default:
315 return "", fmt.Errorf("invalid protocol: %v", protocol)
316 }
317 // TODO(gio): check for already existing mapping
318 delete(proxyMap, src)
319 w, err := fs.Writer(c.ConfigPath)
320 if err != nil {
321 return "", err
322 }
323 defer w.Close()
giof55ab362025-04-11 17:48:17 +0400324 if err := cfg.Render(w); err != nil {
gio721c0042025-04-03 11:56:36 +0400325 return "", err
326 }
gio721c0042025-04-03 11:56:36 +0400327 nginxPath := filepath.Join(filepath.Dir(c.ConfigPath), "ingress-nginx.yaml")
328 nginx, err := func() (map[string]any, error) {
329 r, err := fs.Reader(nginxPath)
330 if err != nil {
331 return nil, err
332 }
333 defer r.Close()
334 var buf bytes.Buffer
335 if _, err := io.Copy(&buf, r); err != nil {
336 return nil, err
337 }
338 ret := map[string]any{}
339 if err := yaml.Unmarshal(buf.Bytes(), &ret); err != nil {
340 return nil, err
341 }
342 return ret, nil
343 }()
344 if err != nil {
345 return "", err
346 }
gio721c0042025-04-03 11:56:36 +0400347 buf, err := yaml.Marshal(nginx)
348 if err != nil {
349 return "", err
350 }
351 w, err = fs.Writer(nginxPath)
352 if err != nil {
353 return "", err
354 }
355 defer w.Close()
356 if _, err := io.Copy(w, bytes.NewReader(buf)); err != nil {
357 return "", err
358 }
359 return fmt.Sprintf("remove proxy mapping: %d %s", src, dst), nil
360 })
361 return err
362}
363
364func (c *NginxProxyConfigurator) RemoveIngressProxy(src, dst string) error {
365 _, err := c.Repo.Do(func(fs soft.RepoFS) (string, error) {
366 cfg, err := func() (NginxProxyConfig, error) {
367 r, err := fs.Reader(c.ConfigPath)
368 if err != nil {
369 return NginxProxyConfig{}, err
370 }
371 defer r.Close()
372 return ParseNginxProxyConfig(r)
373 }()
374 if err != nil {
375 return "", err
376 }
377 if v, ok := cfg.Ingress[src]; !ok || v != dst {
378 return "", fmt.Errorf("wrong mapping from source: %s actual: %s expected: %s", src, v, dst)
379 }
380 delete(cfg.Ingress, src)
381 w, err := fs.Writer(c.ConfigPath)
382 if err != nil {
383 return "", err
384 }
385 defer w.Close()
giof55ab362025-04-11 17:48:17 +0400386 if err := cfg.Render(w); err != nil {
gio721c0042025-04-03 11:56:36 +0400387 return "", err
388 }
gio721c0042025-04-03 11:56:36 +0400389 nginxPath := filepath.Join(filepath.Dir(c.ConfigPath), "ingress-nginx.yaml")
giof6ad2982024-08-23 17:42:49 +0400390 nginx, err := func() (map[string]any, error) {
391 r, err := fs.Reader(nginxPath)
392 if err != nil {
393 return nil, err
394 }
395 defer r.Close()
396 var buf bytes.Buffer
397 if _, err := io.Copy(&buf, r); err != nil {
398 return nil, err
399 }
400 ret := map[string]any{}
401 if err := yaml.Unmarshal(buf.Bytes(), &ret); err != nil {
402 return nil, err
403 }
404 return ret, nil
405 }()
406 if err != nil {
407 return "", err
408 }
giof6ad2982024-08-23 17:42:49 +0400409 buf, err := yaml.Marshal(nginx)
410 if err != nil {
411 return "", err
412 }
413 w, err = fs.Writer(nginxPath)
414 if err != nil {
415 return "", err
416 }
417 defer w.Close()
418 if _, err := io.Copy(w, bytes.NewReader(buf)); err != nil {
419 return "", err
420 }
421 return fmt.Sprintf("remove proxy mapping: %s %s", src, dst), nil
422 })
423 return err
424}
425
gio721c0042025-04-03 11:56:36 +0400426type Protocol int
427
428const (
429 ProtocolTCP Protocol = iota
430 ProtocolUDP
431)
432
433func ProtocolToString(p Protocol) string {
434 if p == ProtocolTCP {
435 return "TCP"
436 } else {
437 return "UDP"
438 }
439}
440
giof6ad2982024-08-23 17:42:49 +0400441type NginxProxyConfig struct {
gio721c0042025-04-03 11:56:36 +0400442 Namespace string
giof55ab362025-04-11 17:48:17 +0400443 PID string
gio721c0042025-04-03 11:56:36 +0400444 IngressPort int
445 Resolvers []net.IP
446 Ingress map[string]string
447 TCP map[int]string
448 UDP map[int]string
449 PreConf []string
450}
451
452func parseProtocol(s string) (Protocol, error) {
453 switch strings.ToLower(s) {
454 case "tcp":
455 return ProtocolTCP, nil
456 case "udp":
457 return ProtocolUDP, nil
458 default:
459 return ProtocolUDP, fmt.Errorf("invalid protocol: %s", s)
460 }
giof6ad2982024-08-23 17:42:49 +0400461}
462
463func ParseNginxProxyConfig(r io.Reader) (NginxProxyConfig, error) {
464 var buf strings.Builder
465 if _, err := io.Copy(&buf, r); err != nil {
466 return NginxProxyConfig{}, err
467 }
468 ret := NginxProxyConfig{
gio721c0042025-04-03 11:56:36 +0400469 IngressPort: -1,
470 Resolvers: nil,
471 Ingress: make(map[string]string),
472 TCP: make(map[int]string),
473 UDP: make(map[int]string),
giof6ad2982024-08-23 17:42:49 +0400474 }
475 lines := strings.Split(buf.String(), "\n")
gio721c0042025-04-03 11:56:36 +0400476 insidePreConf := true
477 insideHttp := false
giof6ad2982024-08-23 17:42:49 +0400478 insideMap := false
gio721c0042025-04-03 11:56:36 +0400479 insideStream := false
480 streamPort := -1
481 streamPortProtocol := ProtocolTCP
giof6ad2982024-08-23 17:42:49 +0400482 for _, l := range lines {
483 items := strings.Fields(strings.TrimSuffix(l, ";"))
484 if len(items) == 0 {
485 continue
486 }
487 if strings.Contains(l, "nginx.conf") {
488 ret.PreConf = append(ret.PreConf, l)
gio721c0042025-04-03 11:56:36 +0400489 insidePreConf = false
490 } else if insidePreConf {
giof6ad2982024-08-23 17:42:49 +0400491 ret.PreConf = append(ret.PreConf, l)
gio721c0042025-04-03 11:56:36 +0400492 items := strings.Fields(l)
493 if items[0] == "namespace:" {
494 ret.Namespace = items[1]
495 }
giof55ab362025-04-11 17:48:17 +0400496 } else if items[0] == "pid" {
497
498 ret.PID = items[1]
gio721c0042025-04-03 11:56:36 +0400499 } else if items[0] == "http" {
500 insideHttp = true
501 } else if insideHttp && items[0] == "map" {
502 insideMap = true
503 } else if items[0] == "stream" {
504 insideHttp = false
505 insideMap = false
506 insideStream = true
giof6ad2982024-08-23 17:42:49 +0400507 } else if strings.Contains(l, "listen") {
508 if len(items) < 2 {
gio721c0042025-04-03 11:56:36 +0400509 return NginxProxyConfig{}, fmt.Errorf("invalid listen: %s", l)
giof6ad2982024-08-23 17:42:49 +0400510 }
511 port, err := strconv.Atoi(items[1])
512 if err != nil {
513 return NginxProxyConfig{}, err
514 }
gio721c0042025-04-03 11:56:36 +0400515 if insideHttp {
516 if len(items) > 2 {
517 return NginxProxyConfig{}, fmt.Errorf("invalid http listen: %s", l)
518 }
519 ret.IngressPort = port
520 } else {
521 if !insideStream {
522 return NginxProxyConfig{}, fmt.Errorf("invalid state, expected to be inside stream section")
523 }
524 streamPort = port
525 if len(items) == 3 {
526 streamPortProtocol, err = parseProtocol(items[2])
527 if err != nil {
528 return NginxProxyConfig{}, err
529 }
530 } else {
531 streamPortProtocol = ProtocolTCP
532 }
533 }
534 } else if insideHttp && strings.Contains(l, "resolver") {
giof6ad2982024-08-23 17:42:49 +0400535 if len(items) < 2 {
536 return NginxProxyConfig{}, fmt.Errorf("invalid resolver: %s", l)
537 }
538 ip := net.ParseIP(items[1])
539 if ip == nil {
540 return NginxProxyConfig{}, fmt.Errorf("invalid resolver ip: %s", l)
541 }
542 ret.Resolvers = append(ret.Resolvers, ip)
gio721c0042025-04-03 11:56:36 +0400543 } else if insideHttp && insideMap {
giof6ad2982024-08-23 17:42:49 +0400544 if items[0] == "}" {
545 insideMap = false
546 continue
547 }
548 if len(items) < 2 {
549 return NginxProxyConfig{}, fmt.Errorf("invalid map: %s", l)
550 }
gio721c0042025-04-03 11:56:36 +0400551 ret.Ingress[items[0]] = items[1]
552 } else if insideStream && strings.Contains(l, "proxy_pass") {
553 if streamPort == -1 {
554 return NginxProxyConfig{}, fmt.Errorf("invalid state, expected server port to be defined")
555 }
556 if len(items) < 2 {
557 return NginxProxyConfig{}, fmt.Errorf("invalid proxy_pass: %s", l)
558 }
559 if streamPortProtocol == ProtocolTCP {
560 ret.TCP[streamPort] = items[1]
561 } else {
562 ret.UDP[streamPort] = items[1]
563 }
giof6ad2982024-08-23 17:42:49 +0400564 }
565 }
566 return ret, nil
567}
568
569func (c NginxProxyConfig) Render(w io.Writer) error {
570 for _, l := range c.PreConf {
571 fmt.Fprintln(w, l)
572 }
573 tmpl, err := template.New("nginx.conf").Parse(nginxConfigTmpl)
574 if err != nil {
575 return err
576 }
577 return tmpl.Execute(w, c)
578}
579
580const nginxConfigTmpl = ` worker_processes 1;
581 worker_rlimit_nofile 8192;
giof55ab362025-04-11 17:48:17 +0400582 {{- if .PID }}
583 pid {{ .PID }};
584 {{- end }}
giof6ad2982024-08-23 17:42:49 +0400585 events {
586 worker_connections 1024;
587 }
588 http {
589 map $http_host $backend {
gio721c0042025-04-03 11:56:36 +0400590 {{- range $from, $to := .Ingress }}
giof6ad2982024-08-23 17:42:49 +0400591 {{ $from }} {{ $to }};
592 {{- end }}
593 }
594 server {
gio721c0042025-04-03 11:56:36 +0400595 listen {{ .IngressPort }};
giof6ad2982024-08-23 17:42:49 +0400596 location / {
597 {{- range .Resolvers }}
598 resolver {{ . }};
599 {{- end }}
600 proxy_pass http://$backend;
601 }
602 }
gio721c0042025-04-03 11:56:36 +0400603 }
604 {{- if or .TCP .UDP }}
605 stream {
606 {{- range $port, $upstream := .TCP }}
607 server {
608 listen {{ $port }};
609 resolver 100.100.100.100;
610 proxy_pass {{ $upstream }};
611 }
612 {{- end }}
613 {{- range $port, $upstream := .UDP }}
614 server {
615 listen {{ $port }} udp;
gio721c0042025-04-03 11:56:36 +0400616 proxy_pass {{ $upstream }};
617 }
618 {{- end }}
619 }
620 {{- end }}`