| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 1 | package installer |
| 2 | |
| 3 | import ( |
| 4 | "bytes" |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 5 | "encoding/json" |
| gio | b7df27f | 2026-07-28 10:36:17 +0400 | [diff] [blame] | 6 | "errors" |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 7 | "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 | |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 18 | corev1 "k8s.io/api/core/v1" |
| 19 | "k8s.io/apimachinery/pkg/util/intstr" |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 20 | "sigs.k8s.io/yaml" |
| 21 | ) |
| 22 | |
| 23 | type ClusterNetworkConfigurator interface { |
| 24 | AddCluster(name string, ingressIP net.IP) error |
| 25 | RemoveCluster(name string, ingressIP net.IP) error |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 26 | 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 |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 30 | } |
| 31 | |
| 32 | type NginxProxyConfigurator struct { |
| 33 | PrivateSubdomain string |
| 34 | DNSAPIAddr string |
| 35 | Repo soft.RepoIO |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 36 | ConfigPath string |
| 37 | ServicePath string |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 38 | } |
| 39 | |
| 40 | type createARecordReq struct { |
| 41 | Entry string `json:"entry"` |
| 42 | IP net.IP `json:"text"` |
| 43 | } |
| 44 | |
| 45 | func (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) |
| gio | b7df27f | 2026-07-28 10:36:17 +0400 | [diff] [blame] | 61 | return errors.New(buf.String()) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 62 | } |
| 63 | return nil |
| 64 | } |
| 65 | |
| 66 | func (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) |
| gio | b7df27f | 2026-07-28 10:36:17 +0400 | [diff] [blame] | 82 | return errors.New(buf.String()) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 83 | } |
| 84 | return nil |
| 85 | } |
| 86 | |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 87 | func (c *NginxProxyConfigurator) AddProxy(src int, dst string, protocol Protocol) (string, error) { |
| 88 | var namespace string |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 89 | _, err := c.Repo.Do(func(fs soft.RepoFS) (string, error) { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 90 | 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 | } |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 126 | cfg, err := func() (NginxProxyConfig, error) { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 127 | r, err := fs.Reader(c.ConfigPath) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 128 | 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 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 137 | 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) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 146 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 147 | // TODO(gio): check for already existing mapping |
| 148 | proxyMap[src] = dst |
| 149 | w, err := fs.Writer(c.ConfigPath) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 150 | if err != nil { |
| 151 | return "", err |
| 152 | } |
| 153 | defer w.Close() |
| gio | f55ab36 | 2025-04-11 17:48:17 +0400 | [diff] [blame] | 154 | if err := cfg.Render(w); err != nil { |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 155 | return "", err |
| 156 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 157 | nginxPath := filepath.Join(filepath.Dir(c.ConfigPath), "ingress-nginx.yaml") |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 158 | 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 | } |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 177 | 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 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 189 | return fmt.Sprintf("add proxy mapping: %d %s", src, dst), nil |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 190 | }) |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 191 | if err != nil { |
| 192 | return "", err |
| 193 | } |
| 194 | return namespace, nil |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 195 | } |
| 196 | |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 197 | func (c *NginxProxyConfigurator) AddIngressProxy(src, dst string) error { |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 198 | _, err := c.Repo.Do(func(fs soft.RepoFS) (string, error) { |
| 199 | cfg, err := func() (NginxProxyConfig, error) { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 200 | r, err := fs.Reader(c.ConfigPath) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 201 | 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 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 210 | if v, ok := cfg.Ingress[src]; ok && v != dst { |
| 211 | return "", fmt.Errorf("wrong mapping %s already exists (%s)", src, v) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 212 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 213 | cfg.Ingress[src] = dst |
| 214 | w, err := fs.Writer(c.ConfigPath) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 215 | if err != nil { |
| 216 | return "", err |
| 217 | } |
| 218 | defer w.Close() |
| gio | f55ab36 | 2025-04-11 17:48:17 +0400 | [diff] [blame] | 219 | if err := cfg.Render(w); err != nil { |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 220 | return "", err |
| 221 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 222 | 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 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 242 | 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 | |
| 259 | func (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() |
| gio | f55ab36 | 2025-04-11 17:48:17 +0400 | [diff] [blame] | 324 | if err := cfg.Render(w); err != nil { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 325 | return "", err |
| 326 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 327 | 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 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 347 | 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 | |
| 364 | func (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() |
| gio | f55ab36 | 2025-04-11 17:48:17 +0400 | [diff] [blame] | 386 | if err := cfg.Render(w); err != nil { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 387 | return "", err |
| 388 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 389 | nginxPath := filepath.Join(filepath.Dir(c.ConfigPath), "ingress-nginx.yaml") |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 390 | 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 | } |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 409 | 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 | |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 426 | type Protocol int |
| 427 | |
| 428 | const ( |
| 429 | ProtocolTCP Protocol = iota |
| 430 | ProtocolUDP |
| 431 | ) |
| 432 | |
| 433 | func ProtocolToString(p Protocol) string { |
| 434 | if p == ProtocolTCP { |
| 435 | return "TCP" |
| 436 | } else { |
| 437 | return "UDP" |
| 438 | } |
| 439 | } |
| 440 | |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 441 | type NginxProxyConfig struct { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 442 | Namespace string |
| gio | f55ab36 | 2025-04-11 17:48:17 +0400 | [diff] [blame] | 443 | PID string |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 444 | 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 | |
| 452 | func 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 | } |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 461 | } |
| 462 | |
| 463 | func 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{ |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 469 | IngressPort: -1, |
| 470 | Resolvers: nil, |
| 471 | Ingress: make(map[string]string), |
| 472 | TCP: make(map[int]string), |
| 473 | UDP: make(map[int]string), |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 474 | } |
| 475 | lines := strings.Split(buf.String(), "\n") |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 476 | insidePreConf := true |
| 477 | insideHttp := false |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 478 | insideMap := false |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 479 | insideStream := false |
| 480 | streamPort := -1 |
| 481 | streamPortProtocol := ProtocolTCP |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 482 | 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) |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 489 | insidePreConf = false |
| 490 | } else if insidePreConf { |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 491 | ret.PreConf = append(ret.PreConf, l) |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 492 | items := strings.Fields(l) |
| 493 | if items[0] == "namespace:" { |
| 494 | ret.Namespace = items[1] |
| 495 | } |
| gio | f55ab36 | 2025-04-11 17:48:17 +0400 | [diff] [blame] | 496 | } else if items[0] == "pid" { |
| 497 | |
| 498 | ret.PID = items[1] |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 499 | } 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 |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 507 | } else if strings.Contains(l, "listen") { |
| 508 | if len(items) < 2 { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 509 | return NginxProxyConfig{}, fmt.Errorf("invalid listen: %s", l) |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 510 | } |
| 511 | port, err := strconv.Atoi(items[1]) |
| 512 | if err != nil { |
| 513 | return NginxProxyConfig{}, err |
| 514 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 515 | 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") { |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 535 | 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) |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 543 | } else if insideHttp && insideMap { |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 544 | if items[0] == "}" { |
| 545 | insideMap = false |
| 546 | continue |
| 547 | } |
| 548 | if len(items) < 2 { |
| 549 | return NginxProxyConfig{}, fmt.Errorf("invalid map: %s", l) |
| 550 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 551 | 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 | } |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 564 | } |
| 565 | } |
| 566 | return ret, nil |
| 567 | } |
| 568 | |
| 569 | func (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 | |
| 580 | const nginxConfigTmpl = ` worker_processes 1; |
| 581 | worker_rlimit_nofile 8192; |
| gio | f55ab36 | 2025-04-11 17:48:17 +0400 | [diff] [blame] | 582 | {{- if .PID }} |
| 583 | pid {{ .PID }}; |
| 584 | {{- end }} |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 585 | events { |
| 586 | worker_connections 1024; |
| 587 | } |
| 588 | http { |
| 589 | map $http_host $backend { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 590 | {{- range $from, $to := .Ingress }} |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 591 | {{ $from }} {{ $to }}; |
| 592 | {{- end }} |
| 593 | } |
| 594 | server { |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 595 | listen {{ .IngressPort }}; |
| gio | f6ad298 | 2024-08-23 17:42:49 +0400 | [diff] [blame] | 596 | location / { |
| 597 | {{- range .Resolvers }} |
| 598 | resolver {{ . }}; |
| 599 | {{- end }} |
| 600 | proxy_pass http://$backend; |
| 601 | } |
| 602 | } |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 603 | } |
| 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; |
| gio | 721c004 | 2025-04-03 11:56:36 +0400 | [diff] [blame] | 616 | proxy_pass {{ $upstream }}; |
| 617 | } |
| 618 | {{- end }} |
| 619 | } |
| 620 | {{- end }}` |