commit c953c94dcf60553a8ce0eeb2501c5cd6a3bffcf6 Author: tiennm99 Date: Sun Oct 11 03:01:45 2026 +0700 feat: sync Cloudflare DNS records with Traefik container hosts Create a record for each Host(...) in running containers' Traefik router labels, tagged with a configurable comment and optional tags, and delete owned records once their host has been down for DELETE_AFTER. diff --git a/.dockerignore b/.dockerignore new file mode 100644 index 0000000..1a9a46d --- /dev/null +++ b/.dockerignore @@ -0,0 +1,4 @@ +.git +.github +.env +*.md diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml new file mode 100644 index 0000000..e060e32 --- /dev/null +++ b/.github/workflows/ci.yml @@ -0,0 +1,54 @@ +name: ci + +on: + push: + branches: [main] + tags: ["v*.*.*"] + pull_request: + +permissions: + contents: read + +jobs: + test: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v5 + - uses: actions/setup-go@v6 + with: + go-version-file: go.mod + - run: test -z "$(gofmt -l .)" + - run: go vet ./... + - run: go test -race ./... + + image: + needs: test + if: github.event_name == 'push' + runs-on: ubuntu-latest + permissions: + contents: read + packages: write + steps: + - uses: actions/checkout@v5 + - uses: docker/setup-buildx-action@v3 + - uses: docker/login-action@v3 + with: + registry: ghcr.io + username: ${{ github.actor }} + password: ${{ secrets.GITHUB_TOKEN }} + - id: meta + uses: docker/metadata-action@v5 + with: + images: ghcr.io/${{ github.repository }} + tags: | + type=semver,pattern={{version}} + type=semver,pattern={{major}}.{{minor}} + type=semver,pattern={{major}} + type=edge,branch=main + - uses: docker/build-push-action@v6 + with: + context: . + platforms: linux/amd64,linux/arm64 + push: true + tags: ${{ steps.meta.outputs.tags }} + labels: ${{ steps.meta.outputs.labels }} diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..5014fd6 --- /dev/null +++ b/.gitignore @@ -0,0 +1,4 @@ +.env +/traefik-cloudflare-dns +*.test +*.out diff --git a/Dockerfile b/Dockerfile new file mode 100644 index 0000000..13c5b5d --- /dev/null +++ b/Dockerfile @@ -0,0 +1,11 @@ +FROM --platform=$BUILDPLATFORM golang:1-alpine AS build +ARG TARGETOS TARGETARCH +WORKDIR /src +COPY go.mod ./ +COPY *.go ./ +RUN CGO_ENABLED=0 GOOS=$TARGETOS GOARCH=$TARGETARCH go build -trimpath -ldflags="-s -w" -o /out/traefik-cloudflare-dns . + +FROM gcr.io/distroless/static-debian12:nonroot +COPY --from=build /out/traefik-cloudflare-dns /traefik-cloudflare-dns +USER nonroot:nonroot +ENTRYPOINT ["/traefik-cloudflare-dns"] diff --git a/LICENSE b/LICENSE new file mode 100644 index 0000000..261eeb9 --- /dev/null +++ b/LICENSE @@ -0,0 +1,201 @@ + Apache License + Version 2.0, January 2004 + http://www.apache.org/licenses/ + + TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION + + 1. Definitions. + + "License" shall mean the terms and conditions for use, reproduction, + and distribution as defined by Sections 1 through 9 of this document. + + "Licensor" shall mean the copyright owner or entity authorized by + the copyright owner that is granting the License. + + "Legal Entity" shall mean the union of the acting entity and all + other entities that control, are controlled by, or are under common + control with that entity. For the purposes of this definition, + "control" means (i) the power, direct or indirect, to cause the + direction or management of such entity, whether by contract or + otherwise, or (ii) ownership of fifty percent (50%) or more of the + outstanding shares, or (iii) beneficial ownership of such entity. + + "You" (or "Your") shall mean an individual or Legal Entity + exercising permissions granted by this License. + + "Source" form shall mean the preferred form for making modifications, + including but not limited to software source code, documentation + source, and configuration files. + + "Object" form shall mean any form resulting from mechanical + transformation or translation of a Source form, including but + not limited to compiled object code, generated documentation, + and conversions to other media types. + + "Work" shall mean the work of authorship, whether in Source or + Object form, made available under the License, as indicated by a + copyright notice that is included in or attached to the work + (an example is provided in the Appendix below). + + "Derivative Works" shall mean any work, whether in Source or Object + form, that is based on (or derived from) the Work and for which the + editorial revisions, annotations, elaborations, or other modifications + represent, as a whole, an original work of authorship. For the purposes + of this License, Derivative Works shall not include works that remain + separable from, or merely link (or bind by name) to the interfaces of, + the Work and Derivative Works thereof. + + "Contribution" shall mean any work of authorship, including + the original version of the Work and any modifications or additions + to that Work or Derivative Works thereof, that is intentionally + submitted to Licensor for inclusion in the Work by the copyright owner + or by an individual or Legal Entity authorized to submit on behalf of + the copyright owner. For the purposes of this definition, "submitted" + means any form of electronic, verbal, or written communication sent + to the Licensor or its representatives, including but not limited to + communication on electronic mailing lists, source code control systems, + and issue tracking systems that are managed by, or on behalf of, the + Licensor for the purpose of discussing and improving the Work, but + excluding communication that is conspicuously marked or otherwise + designated in writing by the copyright owner as "Not a Contribution." + + "Contributor" shall mean Licensor and any individual or Legal Entity + on behalf of whom a Contribution has been received by Licensor and + subsequently incorporated within the Work. + + 2. Grant of Copyright License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + copyright license to reproduce, prepare Derivative Works of, + publicly display, publicly perform, sublicense, and distribute the + Work and such Derivative Works in Source or Object form. + + 3. Grant of Patent License. Subject to the terms and conditions of + this License, each Contributor hereby grants to You a perpetual, + worldwide, non-exclusive, no-charge, royalty-free, irrevocable + (except as stated in this section) patent license to make, have made, + use, offer to sell, sell, import, and otherwise transfer the Work, + where such license applies only to those patent claims licensable + by such Contributor that are necessarily infringed by their + Contribution(s) alone or by combination of their Contribution(s) + with the Work to which such Contribution(s) was submitted. If You + institute patent litigation against any entity (including a + cross-claim or counterclaim in a lawsuit) alleging that the Work + or a Contribution incorporated within the Work constitutes direct + or contributory patent infringement, then any patent licenses + granted to You under this License for that Work shall terminate + as of the date such litigation is filed. + + 4. Redistribution. You may reproduce and distribute copies of the + Work or Derivative Works thereof in any medium, with or without + modifications, and in Source or Object form, provided that You + meet the following conditions: + + (a) You must give any other recipients of the Work or + Derivative Works a copy of this License; and + + (b) You must cause any modified files to carry prominent notices + stating that You changed the files; and + + (c) You must retain, in the Source form of any Derivative Works + that You distribute, all copyright, patent, trademark, and + attribution notices from the Source form of the Work, + excluding those notices that do not pertain to any part of + the Derivative Works; and + + (d) If the Work includes a "NOTICE" text file as part of its + distribution, then any Derivative Works that You distribute must + include a readable copy of the attribution notices contained + within such NOTICE file, excluding those notices that do not + pertain to any part of the Derivative Works, in at least one + of the following places: within a NOTICE text file distributed + as part of the Derivative Works; within the Source form or + documentation, if provided along with the Derivative Works; or, + within a display generated by the Derivative Works, if and + wherever such third-party notices normally appear. The contents + of the NOTICE file are for informational purposes only and + do not modify the License. You may add Your own attribution + notices within Derivative Works that You distribute, alongside + or as an addendum to the NOTICE text from the Work, provided + that such additional attribution notices cannot be construed + as modifying the License. + + You may add Your own copyright statement to Your modifications and + may provide additional or different license terms and conditions + for use, reproduction, or distribution of Your modifications, or + for any such Derivative Works as a whole, provided Your use, + reproduction, and distribution of the Work otherwise complies with + the conditions stated in this License. + + 5. Submission of Contributions. Unless You explicitly state otherwise, + any Contribution intentionally submitted for inclusion in the Work + by You to the Licensor shall be under the terms and conditions of + this License, without any additional terms or conditions. + Notwithstanding the above, nothing herein shall supersede or modify + the terms of any separate license agreement you may have executed + with Licensor regarding such Contributions. + + 6. Trademarks. This License does not grant permission to use the trade + names, trademarks, service marks, or product names of the Licensor, + except as required for reasonable and customary use in describing the + origin of the Work and reproducing the content of the NOTICE file. + + 7. Disclaimer of Warranty. Unless required by applicable law or + agreed to in writing, Licensor provides the Work (and each + Contributor provides its Contributions) on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or + implied, including, without limitation, any warranties or conditions + of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A + PARTICULAR PURPOSE. You are solely responsible for determining the + appropriateness of using or redistributing the Work and assume any + risks associated with Your exercise of permissions under this License. + + 8. Limitation of Liability. In no event and under no legal theory, + whether in tort (including negligence), contract, or otherwise, + unless required by applicable law (such as deliberate and grossly + negligent acts) or agreed to in writing, shall any Contributor be + liable to You for damages, including any direct, indirect, special, + incidental, or consequential damages of any character arising as a + result of this License or out of the use or inability to use the + Work (including but not limited to damages for loss of goodwill, + work stoppage, computer failure or malfunction, or any and all + other commercial damages or losses), even if such Contributor + has been advised of the possibility of such damages. + + 9. Accepting Warranty or Additional Liability. While redistributing + the Work or Derivative Works thereof, You may choose to offer, + and charge a fee for, acceptance of support, warranty, indemnity, + or other liability obligations and/or rights consistent with this + License. However, in accepting such obligations, You may act only + on Your own behalf and on Your sole responsibility, not on behalf + of any other Contributor, and only if You agree to indemnify, + defend, and hold each Contributor harmless for any liability + incurred by, or claims asserted against, such Contributor by reason + of your accepting any such warranty or additional liability. + + END OF TERMS AND CONDITIONS + + APPENDIX: How to apply the Apache License to your work. + + To apply the Apache License to your work, attach the following + boilerplate notice, with the fields enclosed by brackets "[]" + replaced with your own identifying information. (Don't include + the brackets!) The text should be enclosed in the appropriate + comment syntax for the file format. We also recommend that a + file or class name and description of purpose be included on the + same "printed page" as the copyright notice for easier + identification within third-party archives. + + Copyright [yyyy] [name of copyright owner] + + Licensed under the Apache License, Version 2.0 (the "License"); + you may not use this file except in compliance with the License. + You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + + Unless required by applicable law or agreed to in writing, software + distributed under the License is distributed on an "AS IS" BASIS, + WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + See the License for the specific language governing permissions and + limitations under the License. diff --git a/README.md b/README.md new file mode 100644 index 0000000..58abd0f --- /dev/null +++ b/README.md @@ -0,0 +1,95 @@ +# traefik-cloudflare-dns + +**Cloudflare DNS records that follow your Traefik containers.** traefik-cloudflare-dns watches Docker for containers whose Traefik router rules name a host, creates a Cloudflare record for each new host, and deletes the records it created once their host has been gone long enough. + +- **Follows Docker events.** It reacts when a container starts, stops or is removed, and resyncs every minute in case an event was missed. +- **Never touches records it didn't make.** It only treats a record as its own when the record points at this instance's `TARGET` and carries its comment and tags. Records made by hand, or by another instance pointing at another server, stay as they are. +- **Waits before deleting.** A host must be down for `DELETE_AFTER` (1 hour by default) before its record is removed, so redeploys and restarts never cause a DNS gap. +- **Small and safe.** A static, distroless, non-root image with no ports, built to sit behind a read-only Docker socket proxy. + +## How it decides what to create + +Every running container's `traefik.http.routers..rule` labels are read. Each hostname in a `Host(...)` matcher that is `DOMAIN` or a name under it gets a record: + +| `TARGET` | Record | +| --- | --- | +| IPv4 address | `A` | +| IPv6 address | `AAAA` | +| hostname | `CNAME` | + +If a record with that name already exists and isn't owned by this instance, it is left alone, and a warning is logged once. + +A typical setup is one wildcard record (`*.example.com`) for your main server, plus one instance of this tool on every other server, with `TARGET` set to that server's IP. Explicit records take priority over the wildcard, so each app's name resolves to the server it runs on. + +## Quick start + +```yaml +services: + traefik-cloudflare-dns: + image: ghcr.io/tiennm99/traefik-cloudflare-dns:1 + restart: unless-stopped + environment: + CF_API_TOKEN: ${CF_API_TOKEN:?required} + CF_ZONE_ID: ${CF_ZONE_ID:?required} + DOMAIN: example.com + TARGET: 192.0.2.10 + DOCKER_HOST: tcp://dockerproxy:2375 + depends_on: + - dockerproxy + + dockerproxy: + image: tecnativa/docker-socket-proxy:latest + restart: unless-stopped + environment: + CONTAINERS: 1 + volumes: + - /var/run/docker.sock:/var/run/docker.sock:ro +``` + +The proxy only needs `CONTAINERS`. Events and ping are allowed by the proxy's defaults, and every `POST` is refused, so the tool can read Docker but never change it. + +Create the Cloudflare API token with **Zone → DNS → Edit**, limited to the one zone. + +## Configuration + +| Variable | Default | Purpose | +| --- | --- | --- | +| `CF_API_TOKEN` | required | Cloudflare API token with DNS edit on the zone | +| `CF_ZONE_ID` | required | Zone ID, from the zone's overview page | +| `DOMAIN` | required | Only hosts equal to or under this name are managed | +| `TARGET` | required | Record content: this server's IP, or a hostname for a CNAME | +| `PROXIED` | `false` | Create orange-cloud (proxied) records | +| `TTL` | `1` | Record TTL in seconds; `1` means automatic | +| `RECORD_COMMENT` | `managed by traefik-cloudflare-dns` | Comment set on every record, and part of ownership | +| `RECORD_TAGS` | empty | Comma-separated `name:value` tags, set on every record and part of ownership | +| `DELETE_AFTER` | `1h` | How long a host must be down before its record is deleted | +| `RESYNC_INTERVAL` | `1m` | Full resync period, on top of Docker events | +| `DRY_RUN` | `false` | Log creates and deletes without making them | +| `DOCKER_HOST` | `unix:///var/run/docker.sock` | Docker API endpoint, `unix://` or `tcp://` | +| `LOG_LEVEL` | `info` | `debug`, `info`, `warn` or `error` | + +### Comments and tags + +Cloudflare allows record comments on every plan. They are limited to 100 characters on Free and 500 on paid plans. **Tags are only available on paid plans**, and on a Free zone Cloudflare rejects any record that has them, so leave `RECORD_TAGS` empty there. + +Ownership is matched exactly. If you change `RECORD_COMMENT` or `RECORD_TAGS`, records created under the old values are no longer this instance's, so they won't be deleted automatically. Update or remove them by hand. + +### Deletion timing + +The clock for `DELETE_AFTER` starts when a resync finds an owned record whose host no longer belongs to any running container. It resets if the host comes back. The timers live in memory, so a restart starts them over. That can only delay a deletion, never bring one forward. + +Give each server's instance its own `TARGET`. Two instances pointing at the same target with the same comment would each treat the other's records as their own. + +## Development + +```bash +go vet ./... +go test -race ./... +docker build -t traefik-cloudflare-dns:dev . +``` + +CI runs the same checks on every push. Pushes to `main` publish `ghcr.io/tiennm99/traefik-cloudflare-dns:edge`. A `vX.Y.Z` tag publishes `X.Y.Z`, `X.Y` and `X`. + +## License + +[Apache 2.0](LICENSE) diff --git a/client_test.go b/client_test.go new file mode 100644 index 0000000..e1b4e0a --- /dev/null +++ b/client_test.go @@ -0,0 +1,122 @@ +package main + +import ( + "context" + "encoding/json" + "fmt" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" +) + +func TestCloudflareClientPaginatesAndFilters(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Header.Get("Authorization") != "Bearer tok" { + t.Errorf("missing token") + } + page := r.URL.Query().Get("page") + results := map[string]string{ + "1": `[{"id":"1","type":"A","name":"a.example.com","content":"192.0.2.1"},{"id":"2","type":"TXT","name":"t.example.com","content":"x"}]`, + "2": `[{"id":"3","type":"CNAME","name":"c.example.com","content":"h.example.net","comment":"managed"}]`, + }[page] + fmt.Fprintf(w, `{"success":true,"errors":[],"result":%s,"result_info":{"page":%s,"total_pages":2}}`, results, page) + })) + defer srv.Close() + + c := NewCloudflareClient("tok", "zone") + c.base = srv.URL + records, err := c.ListRecords(context.Background()) + if err != nil { + t.Fatal(err) + } + if len(records) != 2 || records[1].Comment != "managed" { + t.Fatalf("records = %+v", records) + } +} + +func TestCloudflareClientCreateSendsCommentAndTags(t *testing.T) { + var got Record + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodPost || r.URL.Path != "/zones/zone/dns_records" { + t.Errorf("unexpected %s %s", r.Method, r.URL.Path) + } + _ = json.NewDecoder(r.Body).Decode(&got) + got.ID = "new" + b, _ := json.Marshal(got) + fmt.Fprintf(w, `{"success":true,"errors":[],"result":%s}`, b) + })) + defer srv.Close() + + c := NewCloudflareClient("tok", "zone") + c.base = srv.URL + rec, err := c.CreateRecord(context.Background(), Record{ + Type: "A", Name: "a.example.com", Content: "192.0.2.1", TTL: 1, + Comment: "managed by traefik-cloudflare-dns", Tags: []string{"owner:x"}, + }) + if err != nil { + t.Fatal(err) + } + if rec.ID != "new" || got.Comment != "managed by traefik-cloudflare-dns" || len(got.Tags) != 1 { + t.Fatalf("sent %+v, got %+v", got, rec) + } +} + +func TestCloudflareClientReportsAPIErrors(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusForbidden) + fmt.Fprint(w, `{"success":false,"errors":[{"code":10000,"message":"Authentication error"}],"result":null}`) + })) + defer srv.Close() + + c := NewCloudflareClient("tok", "zone") + c.base = srv.URL + _, err := c.ListRecords(context.Background()) + if err == nil || !strings.Contains(err.Error(), "10000 Authentication error") { + t.Fatalf("err = %v", err) + } +} + +func TestDockerClientListsAndStreamsEvents(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch r.URL.Path { + case "/containers/json": + fmt.Fprint(w, `[{"Id":"abc","Labels":{"traefik.http.routers.r.rule":"Host(`+"`a.example.com`"+`)"}}]`) + case "/events": + if !strings.Contains(r.URL.Query().Get("filters"), `"start"`) { + t.Errorf("filters = %s", r.URL.Query().Get("filters")) + } + fmt.Fprint(w, `{"Type":"container","Action":"start"}`+"\n"+`{"Type":"container","Action":"die"}`+"\n") + default: + http.NotFound(w, r) + } + })) + defer srv.Close() + + d, err := NewDockerClient("tcp://" + strings.TrimPrefix(srv.URL, "http://")) + if err != nil { + t.Fatal(err) + } + containers, err := d.RunningContainers(context.Background()) + if err != nil { + t.Fatal(err) + } + if len(containers) != 1 || HostsFromLabels(containers[0].Labels)[0] != "a.example.com" { + t.Fatalf("containers = %+v", containers) + } + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + var actions []string + _ = d.WatchContainerEvents(ctx, func(a string) { actions = append(actions, a) }) + if strings.Join(actions, ",") != "start,die" { + t.Fatalf("actions = %v", actions) + } +} + +func TestNewDockerClientRejectsUnknownScheme(t *testing.T) { + if _, err := NewDockerClient("ssh://host"); err == nil { + t.Fatal("expected an error") + } +} diff --git a/cloudflare.go b/cloudflare.go new file mode 100644 index 0000000..d355f27 --- /dev/null +++ b/cloudflare.go @@ -0,0 +1,137 @@ +package main + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "strings" + "time" +) + +const cloudflareAPI = "https://api.cloudflare.com/client/v4" + +// Record is a Cloudflare DNS record, limited to the fields this tool uses. +type Record struct { + ID string `json:"id,omitempty"` + Type string `json:"type"` + Name string `json:"name"` + Content string `json:"content"` + TTL int `json:"ttl,omitempty"` + Proxied bool `json:"proxied"` + Comment string `json:"comment,omitempty"` + Tags []string `json:"tags,omitempty"` +} + +// CloudflareClient manages DNS records in one zone with a scoped API token. +type CloudflareClient struct { + base string + token string + zoneID string + http *http.Client +} + +// NewCloudflareClient returns a client for zoneID. +func NewCloudflareClient(token, zoneID string) *CloudflareClient { + return &CloudflareClient{ + base: cloudflareAPI, + token: token, + zoneID: zoneID, + http: &http.Client{Timeout: 30 * time.Second}, + } +} + +type cfResponse struct { + Success bool `json:"success"` + Errors []struct { + Code int `json:"code"` + Message string `json:"message"` + } `json:"errors"` + Result json.RawMessage `json:"result"` + ResultInfo struct { + Page int `json:"page"` + TotalPages int `json:"total_pages"` + } `json:"result_info"` +} + +func (c *CloudflareClient) do(ctx context.Context, method, path string, body any) (*cfResponse, error) { + var reader io.Reader + if body != nil { + b, err := json.Marshal(body) + if err != nil { + return nil, err + } + reader = bytes.NewReader(b) + } + req, err := http.NewRequestWithContext(ctx, method, c.base+path, reader) + if err != nil { + return nil, err + } + req.Header.Set("Authorization", "Bearer "+c.token) + req.Header.Set("Content-Type", "application/json") + resp, err := c.http.Do(req) + if err != nil { + return nil, err + } + defer resp.Body.Close() + var out cfResponse + if err := json.NewDecoder(io.LimitReader(resp.Body, 32<<20)).Decode(&out); err != nil { + return nil, fmt.Errorf("HTTP %d: undecodable response: %w", resp.StatusCode, err) + } + if !out.Success { + var msgs []string + for _, e := range out.Errors { + msgs = append(msgs, fmt.Sprintf("%d %s", e.Code, e.Message)) + } + return nil, fmt.Errorf("HTTP %d: %s", resp.StatusCode, strings.Join(msgs, "; ")) + } + return &out, nil +} + +// ListRecords returns every A, AAAA and CNAME record in the zone. +func (c *CloudflareClient) ListRecords(ctx context.Context) ([]Record, error) { + var all []Record + for page := 1; ; page++ { + q := url.Values{"per_page": {"100"}, "page": {fmt.Sprint(page)}} + out, err := c.do(ctx, http.MethodGet, "/zones/"+c.zoneID+"/dns_records?"+q.Encode(), nil) + if err != nil { + return nil, fmt.Errorf("list records: %w", err) + } + var records []Record + if err := json.Unmarshal(out.Result, &records); err != nil { + return nil, fmt.Errorf("list records: %w", err) + } + for _, r := range records { + if r.Type == "A" || r.Type == "AAAA" || r.Type == "CNAME" { + all = append(all, r) + } + } + if out.ResultInfo.TotalPages <= page { + return all, nil + } + } +} + +// CreateRecord creates r and returns the stored record. +func (c *CloudflareClient) CreateRecord(ctx context.Context, r Record) (Record, error) { + out, err := c.do(ctx, http.MethodPost, "/zones/"+c.zoneID+"/dns_records", r) + if err != nil { + return Record{}, fmt.Errorf("create %s: %w", r.Name, err) + } + var created Record + if err := json.Unmarshal(out.Result, &created); err != nil { + return Record{}, fmt.Errorf("create %s: %w", r.Name, err) + } + return created, nil +} + +// DeleteRecord deletes the record with id. +func (c *CloudflareClient) DeleteRecord(ctx context.Context, id string) error { + if _, err := c.do(ctx, http.MethodDelete, "/zones/"+c.zoneID+"/dns_records/"+id, nil); err != nil { + return fmt.Errorf("delete %s: %w", id, err) + } + return nil +} diff --git a/config.go b/config.go new file mode 100644 index 0000000..1dc22bc --- /dev/null +++ b/config.go @@ -0,0 +1,163 @@ +package main + +import ( + "fmt" + "net" + "os" + "slices" + "strconv" + "strings" + "time" +) + +const defaultComment = "managed by traefik-cloudflare-dns" + +// Config holds every setting, read from environment variables. +type Config struct { + APIToken string + ZoneID string + Domain string + Target string + RecordType string + Proxied bool + TTL int + Comment string + Tags []string + DeleteAfter time.Duration + ResyncInterval time.Duration + DryRun bool + DockerHost string + LogLevel string +} + +// LoadConfig reads the configuration through getenv and validates it. +func LoadConfig(getenv func(string) string) (Config, error) { + c := Config{ + APIToken: strings.TrimSpace(getenv("CF_API_TOKEN")), + ZoneID: strings.TrimSpace(getenv("CF_ZONE_ID")), + Domain: strings.ToLower(strings.Trim(strings.TrimSpace(getenv("DOMAIN")), ".")), + Target: strings.TrimSpace(getenv("TARGET")), + Comment: getenv("RECORD_COMMENT"), + DockerHost: strings.TrimSpace(getenv("DOCKER_HOST")), + LogLevel: strings.ToLower(strings.TrimSpace(getenv("LOG_LEVEL"))), + } + + var missing []string + for name, value := range map[string]string{ + "CF_API_TOKEN": c.APIToken, + "CF_ZONE_ID": c.ZoneID, + "DOMAIN": c.Domain, + "TARGET": c.Target, + } { + if value == "" { + missing = append(missing, name) + } + } + if len(missing) > 0 { + slices.Sort(missing) + return Config{}, fmt.Errorf("missing required variables: %s", strings.Join(missing, ", ")) + } + + c.RecordType = recordTypeFor(c.Target) + if c.RecordType == "CNAME" { + c.Target = strings.ToLower(strings.TrimSuffix(c.Target, ".")) + } + + if c.Comment == "" { + c.Comment = defaultComment + } + if strings.ContainsAny(c.Comment, "\r\n") { + return Config{}, fmt.Errorf("RECORD_COMMENT must not contain line breaks") + } + + for _, tag := range strings.Split(getenv("RECORD_TAGS"), ",") { + tag = strings.TrimSpace(tag) + if tag == "" { + continue + } + if !strings.Contains(tag, ":") { + return Config{}, fmt.Errorf("RECORD_TAGS entry %q is not in name:value form", tag) + } + c.Tags = append(c.Tags, tag) + } + + var err error + if c.Proxied, err = parseBool(getenv, "PROXIED", false); err != nil { + return Config{}, err + } + if c.DryRun, err = parseBool(getenv, "DRY_RUN", false); err != nil { + return Config{}, err + } + if c.TTL, err = parseInt(getenv, "TTL", 1); err != nil { + return Config{}, err + } + if c.DeleteAfter, err = parseDuration(getenv, "DELETE_AFTER", time.Hour); err != nil { + return Config{}, err + } + if c.ResyncInterval, err = parseDuration(getenv, "RESYNC_INTERVAL", time.Minute); err != nil { + return Config{}, err + } + if c.ResyncInterval <= 0 { + return Config{}, fmt.Errorf("RESYNC_INTERVAL must be positive") + } + + if c.DockerHost == "" { + c.DockerHost = "unix:///var/run/docker.sock" + } + if c.LogLevel == "" { + c.LogLevel = "info" + } + return c, nil +} + +// recordTypeFor returns A or AAAA for an IP target and CNAME for a hostname. +func recordTypeFor(target string) string { + ip := net.ParseIP(target) + switch { + case ip == nil: + return "CNAME" + case ip.To4() != nil: + return "A" + default: + return "AAAA" + } +} + +func parseBool(getenv func(string) string, name string, def bool) (bool, error) { + v := strings.TrimSpace(getenv(name)) + if v == "" { + return def, nil + } + b, err := strconv.ParseBool(v) + if err != nil { + return false, fmt.Errorf("%s: %q is not a boolean", name, v) + } + return b, nil +} + +func parseInt(getenv func(string) string, name string, def int) (int, error) { + v := strings.TrimSpace(getenv(name)) + if v == "" { + return def, nil + } + n, err := strconv.Atoi(v) + if err != nil { + return 0, fmt.Errorf("%s: %q is not an integer", name, v) + } + return n, nil +} + +func parseDuration(getenv func(string) string, name string, def time.Duration) (time.Duration, error) { + v := strings.TrimSpace(getenv(name)) + if v == "" { + return def, nil + } + d, err := time.ParseDuration(v) + if err != nil || d < 0 { + return 0, fmt.Errorf("%s: %q is not a non-negative duration such as 1h or 30m", name, v) + } + return d, nil +} + +// osGetenv is the production getenv. +var osGetenv = os.Getenv diff --git a/config_test.go b/config_test.go new file mode 100644 index 0000000..cde9f15 --- /dev/null +++ b/config_test.go @@ -0,0 +1,99 @@ +package main + +import ( + "strings" + "testing" + "time" +) + +func envOf(m map[string]string) func(string) string { + return func(k string) string { return m[k] } +} + +func requiredEnv() map[string]string { + return map[string]string{ + "CF_API_TOKEN": "token", + "CF_ZONE_ID": "zone", + "DOMAIN": "Example.com.", + "TARGET": "192.0.2.10", + } +} + +func TestLoadConfigDefaults(t *testing.T) { + c, err := LoadConfig(envOf(requiredEnv())) + if err != nil { + t.Fatal(err) + } + if c.Domain != "example.com" || c.RecordType != "A" || c.TTL != 1 || c.Proxied || c.DryRun { + t.Fatalf("unexpected config: %+v", c) + } + if c.Comment != "managed by traefik-cloudflare-dns" { + t.Fatalf("comment = %q", c.Comment) + } + if c.DeleteAfter != time.Hour || c.ResyncInterval != time.Minute { + t.Fatalf("durations = %v, %v", c.DeleteAfter, c.ResyncInterval) + } + if c.DockerHost != "unix:///var/run/docker.sock" { + t.Fatalf("docker host = %q", c.DockerHost) + } +} + +func TestLoadConfigMissing(t *testing.T) { + _, err := LoadConfig(envOf(map[string]string{"DOMAIN": "example.com"})) + if err == nil || !strings.Contains(err.Error(), "CF_API_TOKEN, CF_ZONE_ID, TARGET") { + t.Fatalf("err = %v", err) + } +} + +func TestLoadConfigRecordType(t *testing.T) { + for target, want := range map[string]string{ + "192.0.2.10": "A", + "2001:db8::1": "AAAA", + "Host.Example.net.": "CNAME", + } { + env := requiredEnv() + env["TARGET"] = target + c, err := LoadConfig(envOf(env)) + if err != nil { + t.Fatal(err) + } + if c.RecordType != want { + t.Errorf("TARGET %q: type %s, want %s", target, c.RecordType, want) + } + } +} + +func TestLoadConfigOptional(t *testing.T) { + env := requiredEnv() + env["RECORD_COMMENT"] = "created by the jp companion" + env["RECORD_TAGS"] = "managed-by:traefik-cloudflare-dns, server:jp" + env["PROXIED"] = "true" + env["DELETE_AFTER"] = "30m" + env["DRY_RUN"] = "1" + c, err := LoadConfig(envOf(env)) + if err != nil { + t.Fatal(err) + } + if c.Comment != "created by the jp companion" || !c.Proxied || !c.DryRun || c.DeleteAfter != 30*time.Minute { + t.Fatalf("unexpected config: %+v", c) + } + if len(c.Tags) != 2 || c.Tags[1] != "server:jp" { + t.Fatalf("tags = %v", c.Tags) + } +} + +func TestLoadConfigInvalid(t *testing.T) { + for name, kv := range map[string][2]string{ + "bad tag": {"RECORD_TAGS", "nocolon"}, + "bad bool": {"PROXIED", "maybe"}, + "bad duration": {"DELETE_AFTER", "soon"}, + "line break": {"RECORD_COMMENT", "a\nb"}, + "zero resync": {"RESYNC_INTERVAL", "0s"}, + } { + env := requiredEnv() + env[kv[0]] = kv[1] + if _, err := LoadConfig(envOf(env)); err == nil { + t.Errorf("%s: expected an error", name) + } + } +} diff --git a/docker.go b/docker.go new file mode 100644 index 0000000..e064996 --- /dev/null +++ b/docker.go @@ -0,0 +1,115 @@ +package main + +import ( + "bufio" + "context" + "encoding/json" + "fmt" + "io" + "net" + "net/http" + "net/url" + "strings" + "time" +) + +// Container is the part of a running container this tool reads. +type Container struct { + ID string `json:"Id"` + Labels map[string]string `json:"Labels"` +} + +// DockerClient talks to the Docker Engine API over a unix socket or TCP, +// such as a read-only docker-socket-proxy. +type DockerClient struct { + base string + http *http.Client + stream *http.Client +} + +// NewDockerClient builds a client from a DOCKER_HOST value +// (unix:///path or tcp://host:port). +func NewDockerClient(dockerHost string) (*DockerClient, error) { + u, err := url.Parse(dockerHost) + if err != nil { + return nil, fmt.Errorf("DOCKER_HOST %q: %w", dockerHost, err) + } + transport := &http.Transport{} + base := "" + switch u.Scheme { + case "unix": + socket := u.Path + transport.DialContext = func(ctx context.Context, _, _ string) (net.Conn, error) { + var d net.Dialer + return d.DialContext(ctx, "unix", socket) + } + base = "http://docker" + case "tcp", "http": + base = "http://" + u.Host + default: + return nil, fmt.Errorf("DOCKER_HOST %q: scheme must be unix or tcp", dockerHost) + } + return &DockerClient{ + base: base, + http: &http.Client{Transport: transport, Timeout: 30 * time.Second}, + stream: &http.Client{Transport: transport}, + }, nil +} + +// RunningContainers lists running containers with their labels. +func (d *DockerClient) RunningContainers(ctx context.Context) ([]Container, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, d.base+"/containers/json", nil) + if err != nil { + return nil, err + } + resp, err := d.http.Do(req) + if err != nil { + return nil, fmt.Errorf("list containers: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + return nil, fmt.Errorf("list containers: HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) + } + var containers []Container + if err := json.NewDecoder(resp.Body).Decode(&containers); err != nil { + return nil, fmt.Errorf("list containers: %w", err) + } + return containers, nil +} + +// WatchContainerEvents streams container start, die and destroy events and +// calls onEvent for each one. It returns when the stream ends or ctx is done. +func (d *DockerClient) WatchContainerEvents(ctx context.Context, onEvent func(action string)) error { + filters := `{"type":["container"],"event":["start","die","destroy"]}` + req, err := http.NewRequestWithContext(ctx, http.MethodGet, + d.base+"/events?filters="+url.QueryEscape(filters), nil) + if err != nil { + return err + } + resp, err := d.stream.Do(req) + if err != nil { + return fmt.Errorf("watch events: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + body, _ := io.ReadAll(io.LimitReader(resp.Body, 512)) + return fmt.Errorf("watch events: HTTP %d: %s", resp.StatusCode, strings.TrimSpace(string(body))) + } + dec := json.NewDecoder(bufio.NewReader(resp.Body)) + for { + var ev struct { + Type string `json:"Type"` + Action string `json:"Action"` + } + if err := dec.Decode(&ev); err != nil { + if ctx.Err() != nil { + return ctx.Err() + } + return fmt.Errorf("watch events: %w", err) + } + if ev.Type == "container" { + onEvent(ev.Action) + } + } +} diff --git a/go.mod b/go.mod new file mode 100644 index 0000000..2e782e0 --- /dev/null +++ b/go.mod @@ -0,0 +1,3 @@ +module github.com/tiennm99/traefik-cloudflare-dns + +go 1.27.1 diff --git a/hosts.go b/hosts.go new file mode 100644 index 0000000..3ebbea9 --- /dev/null +++ b/hosts.go @@ -0,0 +1,44 @@ +package main + +import ( + "regexp" + "strings" +) + +var ( + routerRuleLabel = regexp.MustCompile(`^traefik\.http\.routers\.[^.]+\.rule$`) + hostMatcher = regexp.MustCompile(`Host\(([^)]*)\)`) + backtickValue = regexp.MustCompile("`([^`]+)`") +) + +// HostsFromLabels returns the hostnames that a container's Traefik HTTP +// router rules match with Host(...), lowercased and deduplicated. +// HostRegexp and other matchers are ignored. A container labelled +// traefik.enable=false yields nothing. +func HostsFromLabels(labels map[string]string) []string { + if strings.EqualFold(strings.TrimSpace(labels["traefik.enable"]), "false") { + return nil + } + seen := map[string]bool{} + var hosts []string + for key, rule := range labels { + if !routerRuleLabel.MatchString(key) { + continue + } + for _, m := range hostMatcher.FindAllStringSubmatch(rule, -1) { + for _, v := range backtickValue.FindAllStringSubmatch(m[1], -1) { + host := strings.ToLower(strings.TrimSuffix(strings.TrimSpace(v[1]), ".")) + if host != "" && !seen[host] { + seen[host] = true + hosts = append(hosts, host) + } + } + } + } + return hosts +} + +// InDomain reports whether host is domain itself or a name under it. +func InDomain(host, domain string) bool { + return host == domain || strings.HasSuffix(host, "."+domain) +} diff --git a/hosts_test.go b/hosts_test.go new file mode 100644 index 0000000..09fda24 --- /dev/null +++ b/hosts_test.go @@ -0,0 +1,71 @@ +package main + +import ( + "slices" + "testing" +) + +func TestHostsFromLabels(t *testing.T) { + tests := []struct { + name string + labels map[string]string + want []string + }{ + { + name: "coolify rule with path prefix", + labels: map[string]string{ + "traefik.enable": "true", + "traefik.http.routers.https-0-x.rule": "Host(`App.Example.com`) && PathPrefix(`/`)", + "traefik.http.routers.http-0-x.rule": "Host(`app.example.com`) && PathPrefix(`/`)", + "traefik.http.services.x.loadbalancer": "ignored", + }, + want: []string{"app.example.com"}, + }, + { + name: "several hosts in one matcher and several matchers", + labels: map[string]string{ + "traefik.http.routers.r.rule": "Host(`a.example.com`, `b.example.com`) || Host(`c.example.com.`)", + }, + want: []string{"a.example.com", "b.example.com", "c.example.com"}, + }, + { + name: "host regexp and tcp routers are ignored", + labels: map[string]string{ + "traefik.http.routers.r.rule": "HostRegexp(`{sub:[a-z]+}.example.com`)", + "traefik.tcp.routers.t.rule": "HostSNI(`db.example.com`)", + }, + want: nil, + }, + { + name: "traefik disabled", + labels: map[string]string{ + "traefik.enable": "false", + "traefik.http.routers.r.rule": "Host(`a.example.com`)", + }, + want: nil, + }, + } + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := HostsFromLabels(tt.labels) + slices.Sort(got) + if !slices.Equal(got, tt.want) { + t.Fatalf("got %v, want %v", got, tt.want) + } + }) + } +} + +func TestInDomain(t *testing.T) { + for host, want := range map[string]bool{ + "example.com": true, + "app.example.com": true, + "a.b.example.com": true, + "notexample.com": false, + "example.com.evil.net": false, + } { + if got := InDomain(host, "example.com"); got != want { + t.Errorf("InDomain(%q) = %v, want %v", host, got, want) + } + } +} diff --git a/main.go b/main.go new file mode 100644 index 0000000..6ec94e9 --- /dev/null +++ b/main.go @@ -0,0 +1,123 @@ +// Command traefik-cloudflare-dns creates Cloudflare DNS records for the +// hostnames in Traefik router labels of running Docker containers, and +// deletes the records it created once their host has been down long enough. +package main + +import ( + "context" + "log/slog" + "os" + "os/signal" + "syscall" + "time" +) + +const ( + debounce = 2 * time.Second + retryAfter = 10 * time.Second +) + +func main() { + cfg, err := LoadConfig(osGetenv) + if err != nil { + slog.Error("invalid configuration", "err", err) + os.Exit(2) + } + log := slog.New(slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{Level: parseLevel(cfg.LogLevel)})) + + docker, err := NewDockerClient(cfg.DockerHost) + if err != nil { + log.Error("invalid configuration", "err", err) + os.Exit(2) + } + rec := NewReconciler(cfg, NewCloudflareClient(cfg.APIToken, cfg.ZoneID), docker, log) + + ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer stop() + + log.Info("starting", + "domain", cfg.Domain, "target", cfg.Target, "type", cfg.RecordType, + "proxied", cfg.Proxied, "comment", cfg.Comment, "tags", cfg.Tags, + "delete_after", cfg.DeleteAfter, "resync_interval", cfg.ResyncInterval, + "dry_run", cfg.DryRun, "docker_host", cfg.DockerHost) + + trigger := make(chan struct{}, 1) + go watchEvents(ctx, docker, trigger, log) + run(ctx, rec, trigger, cfg.ResyncInterval, log) + log.Info("stopped") +} + +// run reconciles once at start, then after each burst of container events, +// on every resync tick, and shortly after a failed pass, until ctx is done. +func run(ctx context.Context, rec *Reconciler, trigger <-chan struct{}, every time.Duration, log *slog.Logger) { + var retry <-chan time.Time + pass := func() { + retry = nil + if err := rec.Reconcile(ctx); err != nil && ctx.Err() == nil { + log.Error("reconcile failed; retrying", "err", err, "in", retryAfter) + retry = time.After(retryAfter) + } + } + pass() + ticker := time.NewTicker(every) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + pass() + case <-retry: + pass() + case <-trigger: + select { + case <-ctx.Done(): + return + case <-time.After(debounce): + } + pass() + } + } +} + +// watchEvents keeps a Docker event stream open and signals trigger on every +// container start, die or destroy, reconnecting with backoff when it drops. +func watchEvents(ctx context.Context, docker *DockerClient, trigger chan<- struct{}, log *slog.Logger) { + backoff := time.Second + for ctx.Err() == nil { + started := time.Now() + err := docker.WatchContainerEvents(ctx, func(action string) { + log.Debug("container event", "action", action) + select { + case trigger <- struct{}{}: + default: + } + }) + if ctx.Err() != nil { + return + } + if time.Since(started) > time.Minute { + backoff = time.Second + } + log.Warn("docker event stream ended; reconnecting", "err", err, "in", backoff) + select { + case <-ctx.Done(): + return + case <-time.After(backoff): + } + backoff = min(backoff*2, time.Minute) + } +} + +func parseLevel(s string) slog.Level { + switch s { + case "debug": + return slog.LevelDebug + case "warn": + return slog.LevelWarn + case "error": + return slog.LevelError + default: + return slog.LevelInfo + } +} diff --git a/reconcile.go b/reconcile.go new file mode 100644 index 0000000..ef8076e --- /dev/null +++ b/reconcile.go @@ -0,0 +1,173 @@ +package main + +import ( + "context" + "log/slog" + "slices" + "time" +) + +// DNS is the record store the reconciler writes to. +type DNS interface { + ListRecords(ctx context.Context) ([]Record, error) + CreateRecord(ctx context.Context, r Record) (Record, error) + DeleteRecord(ctx context.Context, id string) error +} + +// ContainerSource lists the running containers. +type ContainerSource interface { + RunningContainers(ctx context.Context) ([]Container, error) +} + +// Reconciler keeps the zone in line with the hosts of running containers. +// +// A record is owned when its content is this instance's target and it +// carries the configured comment and every configured tag. Only owned +// records are deleted, so records made by hand, or by another instance +// pointing at another server, are never touched. +type Reconciler struct { + cfg Config + dns DNS + containers ContainerSource + log *slog.Logger + now func() time.Time + + // missingSince tracks when each owned record's host was last served. + missingSince map[string]time.Time + // warned avoids repeating the same conflict warning every pass. + warned map[string]bool +} + +// NewReconciler returns a reconciler for cfg. +func NewReconciler(cfg Config, dns DNS, containers ContainerSource, log *slog.Logger) *Reconciler { + return &Reconciler{ + cfg: cfg, + dns: dns, + containers: containers, + log: log, + now: time.Now, + missingSince: map[string]time.Time{}, + warned: map[string]bool{}, + } +} + +func (r *Reconciler) owns(rec Record) bool { + if rec.Type != r.cfg.RecordType || rec.Content != r.cfg.Target || rec.Comment != r.cfg.Comment { + return false + } + for _, tag := range r.cfg.Tags { + if !slices.Contains(rec.Tags, tag) { + return false + } + } + return true +} + +// Reconcile runs one pass: create records for new hosts, and delete owned +// records whose host has been absent for at least DeleteAfter. +func (r *Reconciler) Reconcile(ctx context.Context) error { + containers, err := r.containers.RunningContainers(ctx) + if err != nil { + return err + } + desired := map[string]bool{} + for _, c := range containers { + for _, h := range HostsFromLabels(c.Labels) { + if InDomain(h, r.cfg.Domain) { + desired[h] = true + } + } + } + + records, err := r.dns.ListRecords(ctx) + if err != nil { + return err + } + byName := map[string][]Record{} + for _, rec := range records { + byName[rec.Name] = append(byName[rec.Name], rec) + } + + for _, host := range sortedKeys(desired) { + existing := byName[host] + if slices.ContainsFunc(existing, r.owns) { + delete(r.missingSince, host) + continue + } + if len(existing) > 0 { + if !r.warned[host] { + r.log.Warn("record exists and is not managed by this instance; leaving it alone", + "host", host, "type", existing[0].Type, "content", existing[0].Content) + r.warned[host] = true + } + continue + } + rec := Record{ + Type: r.cfg.RecordType, + Name: host, + Content: r.cfg.Target, + TTL: r.cfg.TTL, + Proxied: r.cfg.Proxied, + Comment: r.cfg.Comment, + Tags: r.cfg.Tags, + } + if r.cfg.DryRun { + r.log.Info("dry run: would create record", "host", host, "type", rec.Type, "content", rec.Content) + continue + } + if _, err := r.dns.CreateRecord(ctx, rec); err != nil { + r.log.Error("create record failed", "host", host, "err", err) + continue + } + r.log.Info("created record", "host", host, "type", rec.Type, "content", rec.Content) + } + for host := range r.warned { + if !desired[host] { + delete(r.warned, host) + } + } + + now := r.now() + owned := map[string]bool{} + for _, rec := range records { + if !r.owns(rec) || desired[rec.Name] { + continue + } + owned[rec.Name] = true + since, tracked := r.missingSince[rec.Name] + if !tracked { + r.missingSince[rec.Name] = now + r.log.Info("host no longer served; record will be deleted if it stays down", + "host", rec.Name, "delete_after", r.cfg.DeleteAfter) + since = now + } + if now.Sub(since) < r.cfg.DeleteAfter { + continue + } + if r.cfg.DryRun { + r.log.Info("dry run: would delete record", "host", rec.Name, "down_for", now.Sub(since).Round(time.Second)) + continue + } + if err := r.dns.DeleteRecord(ctx, rec.ID); err != nil { + r.log.Error("delete record failed", "host", rec.Name, "err", err) + continue + } + r.log.Info("deleted record", "host", rec.Name, "down_for", now.Sub(since).Round(time.Second)) + delete(r.missingSince, rec.Name) + } + for host := range r.missingSince { + if !owned[host] { + delete(r.missingSince, host) + } + } + return nil +} + +func sortedKeys(m map[string]bool) []string { + keys := make([]string, 0, len(m)) + for k := range m { + keys = append(keys, k) + } + slices.Sort(keys) + return keys +} diff --git a/reconcile_test.go b/reconcile_test.go new file mode 100644 index 0000000..52b71a1 --- /dev/null +++ b/reconcile_test.go @@ -0,0 +1,201 @@ +package main + +import ( + "context" + "fmt" + "io" + "log/slog" + "testing" + "time" +) + +type fakeDNS struct { + records []Record + nextID int + created []string + deleted []string +} + +func (f *fakeDNS) ListRecords(context.Context) ([]Record, error) { + return append([]Record(nil), f.records...), nil +} + +func (f *fakeDNS) CreateRecord(_ context.Context, r Record) (Record, error) { + f.nextID++ + r.ID = fmt.Sprint(f.nextID) + f.records = append(f.records, r) + f.created = append(f.created, r.Name) + return r, nil +} + +func (f *fakeDNS) DeleteRecord(_ context.Context, id string) error { + for i, r := range f.records { + if r.ID == id { + f.records = append(f.records[:i], f.records[i+1:]...) + f.deleted = append(f.deleted, r.Name) + return nil + } + } + return fmt.Errorf("no record %s", id) +} + +type fakeContainers struct{ hosts []string } + +func (f *fakeContainers) RunningContainers(context.Context) ([]Container, error) { + var out []Container + for i, h := range f.hosts { + out = append(out, Container{ + ID: fmt.Sprint(i), + Labels: map[string]string{"traefik.http.routers.r.rule": "Host(`" + h + "`)"}, + }) + } + return out, nil +} + +func newTestReconciler(t *testing.T, cfg Config, dns *fakeDNS, c *fakeContainers) (*Reconciler, *time.Time) { + t.Helper() + now := time.Date(2026, 10, 11, 0, 0, 0, 0, time.UTC) + r := NewReconciler(cfg, dns, c, slog.New(slog.NewTextHandler(io.Discard, nil))) + r.now = func() time.Time { return now } + return r, &now +} + +func testConfig() Config { + return Config{ + Domain: "example.com", + Target: "192.0.2.10", + RecordType: "A", + TTL: 1, + Comment: defaultComment, + DeleteAfter: time.Hour, + } +} + +func TestReconcileCreatesOnlyHostsInDomain(t *testing.T) { + dns := &fakeDNS{} + r, _ := newTestReconciler(t, testConfig(), dns, &fakeContainers{hosts: []string{"app.example.com", "other.example.net"}}) + if err := r.Reconcile(context.Background()); err != nil { + t.Fatal(err) + } + if len(dns.created) != 1 || dns.created[0] != "app.example.com" { + t.Fatalf("created = %v", dns.created) + } + got := dns.records[0] + if got.Type != "A" || got.Content != "192.0.2.10" || got.Comment != defaultComment { + t.Fatalf("record = %+v", got) + } + if err := r.Reconcile(context.Background()); err != nil { + t.Fatal(err) + } + if len(dns.created) != 1 { + t.Fatalf("second pass created again: %v", dns.created) + } +} + +func TestReconcileLeavesForeignRecords(t *testing.T) { + dns := &fakeDNS{records: []Record{ + {ID: "a", Type: "A", Name: "manual.example.com", Content: "192.0.2.99"}, + {ID: "b", Type: "A", Name: "other-server.example.com", Content: "192.0.2.20", Comment: defaultComment}, + {ID: "c", Type: "A", Name: "same-ip-no-comment.example.com", Content: "192.0.2.10"}, + }} + r, now := newTestReconciler(t, testConfig(), dns, &fakeContainers{hosts: []string{"manual.example.com"}}) + for range 3 { + if err := r.Reconcile(context.Background()); err != nil { + t.Fatal(err) + } + *now = now.Add(2 * time.Hour) + } + if len(dns.created) != 0 || len(dns.deleted) != 0 { + t.Fatalf("created %v, deleted %v", dns.created, dns.deleted) + } +} + +func TestReconcileDeletesAfterGracePeriod(t *testing.T) { + dns := &fakeDNS{} + containers := &fakeContainers{hosts: []string{"app.example.com"}} + r, now := newTestReconciler(t, testConfig(), dns, containers) + ctx := context.Background() + if err := r.Reconcile(ctx); err != nil { + t.Fatal(err) + } + + containers.hosts = nil + for _, step := range []time.Duration{0, 30 * time.Minute, 29 * time.Minute} { + *now = now.Add(step) + if err := r.Reconcile(ctx); err != nil { + t.Fatal(err) + } + } + if len(dns.deleted) != 0 { + t.Fatalf("deleted before the grace period: %v", dns.deleted) + } + + *now = now.Add(time.Minute) + if err := r.Reconcile(ctx); err != nil { + t.Fatal(err) + } + if len(dns.deleted) != 1 || dns.deleted[0] != "app.example.com" { + t.Fatalf("deleted = %v", dns.deleted) + } +} + +func TestReconcileComebackResetsGracePeriod(t *testing.T) { + dns := &fakeDNS{} + containers := &fakeContainers{hosts: []string{"app.example.com"}} + r, now := newTestReconciler(t, testConfig(), dns, containers) + ctx := context.Background() + steps := []struct { + hosts []string + after time.Duration + }{ + {[]string{"app.example.com"}, 0}, + {nil, 0}, + {[]string{"app.example.com"}, 50 * time.Minute}, + {nil, time.Minute}, + {nil, 50 * time.Minute}, + } + for _, s := range steps { + containers.hosts = s.hosts + *now = now.Add(s.after) + if err := r.Reconcile(ctx); err != nil { + t.Fatal(err) + } + } + if len(dns.deleted) != 0 { + t.Fatalf("grace period was not reset by the host coming back: %v", dns.deleted) + } +} + +func TestReconcileRequiresTagsForOwnership(t *testing.T) { + cfg := testConfig() + cfg.Tags = []string{"managed-by:traefik-cloudflare-dns"} + dns := &fakeDNS{records: []Record{ + {ID: "a", Type: "A", Name: "untagged.example.com", Content: "192.0.2.10", Comment: defaultComment}, + }} + r, now := newTestReconciler(t, cfg, dns, &fakeContainers{}) + for range 2 { + if err := r.Reconcile(context.Background()); err != nil { + t.Fatal(err) + } + *now = now.Add(2 * time.Hour) + } + if len(dns.deleted) != 0 { + t.Fatalf("deleted a record without the configured tag: %v", dns.deleted) + } +} + +func TestReconcileDryRunWritesNothing(t *testing.T) { + cfg := testConfig() + cfg.DryRun = true + cfg.DeleteAfter = 0 + dns := &fakeDNS{records: []Record{ + {ID: "a", Type: "A", Name: "gone.example.com", Content: "192.0.2.10", Comment: defaultComment}, + }} + r, _ := newTestReconciler(t, cfg, dns, &fakeContainers{hosts: []string{"app.example.com"}}) + if err := r.Reconcile(context.Background()); err != nil { + t.Fatal(err) + } + if len(dns.created) != 0 || len(dns.deleted) != 0 { + t.Fatalf("dry run wrote: created %v, deleted %v", dns.created, dns.deleted) + } +}