18 Commits

Author SHA1 Message Date
benvin b8ad59b053 Merge pull request 'Match zones defined by hosts entries' (#38) from benvin/hosts-zones into main
ci/woodpecker/tag/release Pipeline was successful
Reviewed-on: #38
2026-10-09 23:28:54 +11:00
unkin-agent 3174eabd94 Honour hosts routeback for same-interface intra-zone pairs
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
2026-10-09 23:15:50 +11:00
unkin-agent 0b68110220 Merge remote-tracking branch 'origin/main' into benvin/hosts-zones
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
# Conflicts:
#	internal/nftables/compiler.go
#	internal/nftables/compiler_test.go
2026-10-09 23:12:36 +11:00
benvin 70df237121 Merge pull request 'Accept intra-zone traffic between different interfaces' (#37) from benvin/intrazone-multi-iface into main
Reviewed-on: #37
2026-10-09 23:10:48 +11:00
unkin-agent 799c7f3524 Guard zone host exclusions with the address family
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-09 22:54:07 +11:00
unkin-agent ecc349cb6f Skip fw->fw policies and treat dest-side + as intra-zone override
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-09 22:48:33 +11:00
unkin-agent 695869c80b Match zones defined by hosts entries
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-09 22:48:12 +11:00
unkin-agent 96a1ba8351 Accept intra-zone traffic between different interfaces
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
2026-10-09 22:45:53 +11:00
benvin afa056b454 Merge pull request 'Add boot unit applying a local config' (#34) from benvin/boot-unit into main
ci/woodpecker/tag/release Pipeline was successful
Reviewed-on: #34
2026-10-05 21:46:58 +11:00
benvin f593f7625d Merge pull request 'Revert agent generations that cut off the control plane' (#35) from benvin/agent-safe-apply into main
Reviewed-on: #35
2026-10-05 21:46:08 +11:00
benvin 7c8bd87ec0 Merge pull request 'ci: use container-rpmbuilder image' (#36) from benvin/rpmbuilder-image into main
Reviewed-on: #36
2026-10-05 21:34:18 +11:00
unkin-agent 9854b0e7b6 ci: use container-rpmbuilder image
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-05 14:39:45 +11:00
unkin-agent c7e02c089a Apply the cached config without safe-apply and keep reverted generations in memory
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-05 13:56:11 +11:00
unkin-agent 9092b463a0 Apply agent generations as a pending try with a revert timer
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-05 13:52:32 +11:00
unkin-agent 502d06bdda Cap boot unit restarts so it fails open
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-05 13:46:50 +11:00
unkin-agent 4ad55fc65e Revert agent generations that cut off the control plane
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-05 13:46:20 +11:00
unkin-agent 7fcbb5fad8 Order boot unit before sysinit and retry on failure
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-05 13:44:58 +11:00
unkin-agent dc406c4f56 Add boot unit applying a local config
ci/woodpecker/pr/build Pipeline was successful
ci/woodpecker/pr/test Pipeline was successful
ci/woodpecker/pr/pre-commit Pipeline was successful
2026-10-05 13:41:57 +11:00
20 changed files with 1417 additions and 144 deletions
+1 -1
View File
@@ -59,7 +59,7 @@ steps:
cpu: 2
- name: package
image: git.unkin.net/unkin/almalinux9-rpmbuilder:latest
image: artifactapi.k8s.syd1.au.unkin.net/docker-internal/rpmbuilder:0.1.0-alma9
commands:
- ./scripts/build-rpm.sh ${CI_COMMIT_TAG}
depends_on: [build]
+7
View File
@@ -430,6 +430,13 @@ report the generation applied, giving a fleet-wide "converged / N behind" view.
source/dest disables its rule loudly, never opens it.
- **Adds fail closed, the control plane fails open.** Partial rollout blocks new
flows until every hop converges; a dead API leaves the last-good posture running.
- **A generation that severs the API is reverted.** The agent applies as a
`tomswall try` does (on-disk snapshot, systemd revert timer, shared lock), then
reports `applied` over a fresh connection. If that fails at the transport level,
or the apply errors, it records the generation in
`/var/lib/tomswall/reverted.json` (skipped until a newer one arrives), restores
the snapshot and reports `reverted`/`failed`. A failed restore leaves the timer
to retry it and is reported `failed`.
---
+2 -1
View File
@@ -29,7 +29,8 @@ func agentCmd() *cobra.Command {
Long: `Agent runs the control-plane pull loop: it fetches this device's compiled
config from tomswallapi, differentially applies it, and reports the applied
generation back. It caches the last known-good config and, if the control plane
is unreachable, keeps applying that cache — it never fails closed.
is unreachable, keeps applying that cache — it never fails closed. A new
generation that cuts the agent off from the API is reverted and reported as such.
The agent token defaults to the TOMSWALL_AGENT_TOKEN environment variable, and
the device name defaults to the system hostname.`,
+231 -29
View File
@@ -2,8 +2,13 @@ package agent
import (
"context"
"encoding/json"
"errors"
"fmt"
"log/slog"
"net/url"
"os"
"path/filepath"
"time"
"git.unkin.net/unkin/tomswall/internal/config"
@@ -12,11 +17,17 @@ import (
)
// Applier applies a translated config to the firewall. Abstracted so the run
// loop is testable without touching the kernel.
// loop is testable without touching the kernel. With safe, a change is applied
// as a pending try: revert restores the previous ruleset (a failed revert leaves
// the revert timer armed) and keep drops the snapshot. Both are nil when nothing
// changed or safe is false.
type Applier interface {
Apply(ctx context.Context, cfg *config.Config) error
Apply(ctx context.Context, cfg *config.Config, safe bool) (revert, keep func() error, err error)
}
// revertDelay is when the revert timer fires if the agent dies mid-apply.
const revertDelay = time.Minute
// Agent runs the pull-apply-report loop for one device.
type Agent struct {
Client *Client
@@ -25,6 +36,9 @@ type Agent struct {
Applier Applier
// Resolver overrides the DNS resolver (tests); nil derives it per-config.
Resolver *Resolver
// lastReverted covers a reverted generation whose persistence failed.
lastReverted *reverted
}
// Run loops until ctx is cancelled, applying one cycle per Interval (and once
@@ -62,16 +76,31 @@ func (a *Agent) RunOnce(ctx context.Context) error {
return fmt.Errorf("control plane unreachable and no cached config: %w", err)
}
// Re-apply last known-good; do not report a generation we didn't fetch.
return a.applyConfig(ctx, cached, false)
return a.applyConfig(ctx, cached, nil)
}
if err := a.Cache.Write(raw); err != nil {
slog.Warn("agent: caching config failed", "err", err)
rv, err := a.readReverted()
if err != nil {
return err
}
return a.applyConfig(ctx, rc, true)
if a.lastReverted != nil && (rv == nil || a.lastReverted.Generation > rv.Generation) {
rv = a.lastReverted
}
if rv != nil {
a.reportReverted(ctx, rv)
if rc.Generation <= rv.Generation {
slog.Info("agent: generation was reverted, waiting for a newer one", "generation", rc.Generation)
return nil
}
}
return a.applyConfig(ctx, rc, raw)
}
func (a *Agent) applyConfig(ctx context.Context, rc *RenderedConfig, report bool) error {
// applyConfig applies rc. A fetched config (raw != nil) is applied as a pending
// try, verified by reaching the API through the new ruleset and reverted if that
// fails; only then is it cached. The cached config is the last verified-good one,
// so it is applied plainly: there is nothing to verify it against.
func (a *Agent) applyConfig(ctx context.Context, rc *RenderedConfig, raw []byte) error {
resolver := a.Resolver
if resolver == nil {
resolver = NewResolver(rc.Resolver)
@@ -82,46 +111,219 @@ func (a *Agent) applyConfig(ctx context.Context, rc *RenderedConfig, report bool
if err != nil {
return fmt.Errorf("translate: %w", err)
}
if err := a.Applier.Apply(ctx, cfg); err != nil {
return fmt.Errorf("apply: %w", err)
unlock, err := tryapply.Acquire()
if errors.Is(err, tryapply.ErrPending) {
slog.Warn("agent: a 'tomswall try' is pending, skipping cycle")
return nil
}
if err != nil {
return err
}
defer unlock()
revert, keep, err := a.Applier.Apply(ctx, cfg, raw != nil)
if err != nil {
err = fmt.Errorf("apply: %w", err)
if raw != nil && revert != nil {
return a.revertGeneration(ctx, rc.Generation, StatusFailed, err, revert)
}
if revert != nil {
if rerr := revert(); rerr != nil {
err = fmt.Errorf("%w; restore: %v; revert timer pending", err, rerr)
}
}
if raw != nil {
if rerr := a.Client.ReportStatus(ctx, Status{Status: StatusFailed, Generation: rc.Generation, Error: err.Error()}); rerr != nil {
slog.Warn("agent: reporting status failed", "err", rerr)
}
}
return err
}
if raw == nil {
slog.Info("agent: applied cached config", "generation", rc.Generation, "rules", len(cfg.Rules))
return nil
}
if err := a.confirm(ctx, rc.Generation); err != nil {
if keep != nil && ctx.Err() == nil {
return a.revertGeneration(ctx, rc.Generation, StatusReverted, err, revert)
}
// Shutdown is not a verdict on the generation: keep it.
if keep != nil {
if kerr := keep(); kerr != nil {
slog.Warn("agent: dropping snapshot failed", "err", kerr)
}
}
return err
}
if keep != nil {
if err := keep(); err != nil {
return fmt.Errorf("dropping snapshot: %w", err)
}
}
slog.Info("agent: applied config", "generation", rc.Generation, "rules", len(cfg.Rules))
if report {
if err := a.Client.ReportStatus(ctx, rc.Generation); err != nil {
slog.Warn("agent: reporting status failed", "err", err)
}
// Report the FIB so the control plane can scope router enforcement.
if fib := CollectFIB(ctx); len(fib) > 0 {
if err := a.Client.ReportRoutes(ctx, fib); err != nil {
slog.Warn("agent: reporting routes failed", "err", err)
}
if err := a.Cache.Write(raw); err != nil {
slog.Warn("agent: caching config failed", "err", err)
}
a.lastReverted = nil
if err := os.Remove(a.revertedPath()); err != nil && !os.IsNotExist(err) {
slog.Warn("agent: clearing reverted generation failed", "err", err)
}
// Report the FIB so the control plane can scope router enforcement.
if fib := CollectFIB(ctx); len(fib) > 0 {
if err := a.Client.ReportRoutes(ctx, fib); err != nil {
slog.Warn("agent: reporting routes failed", "err", err)
}
}
return nil
}
// EngineApplier applies via the real nftables differential engine.
type EngineApplier struct{}
// revertGeneration marks generation as reverted before restoring, so a failed
// restore can never lead to re-applying it, then reports status. It is also kept
// in memory in case persisting fails. A failed
// restore leaves the snapshot and timer armed and is reported as failed.
func (a *Agent) revertGeneration(ctx context.Context, generation int64, status string, cause error, revert func() error) error {
rv := &reverted{Generation: generation, Status: status, Error: cause.Error()}
a.lastReverted = rv
if err := a.writeReverted(rv); err != nil {
slog.Error("agent: persisting reverted generation failed", "err", err)
}
suffix := ""
if rerr := revert(); rerr != nil {
suffix = fmt.Sprintf("; restore: %v; revert timer pending", rerr)
rv.Status = StatusFailed
rv.Error += suffix
if werr := a.writeReverted(rv); werr != nil {
slog.Error("agent: persisting reverted generation failed", "err", werr)
}
}
a.reportReverted(ctx, rv)
return fmt.Errorf("generation %d %s: %w%s", generation, rv.Status, cause, suffix)
}
// Apply computes and applies the differential change set for cfg. It refuses
// while a 'tomswall try' awaits confirmation.
func (EngineApplier) Apply(_ context.Context, cfg *config.Config) error {
unlock, err := tryapply.Acquire()
var (
verifyAttempts = 3
verifyDelay = 2 * time.Second
verifyTimeout = 5 * time.Second
)
// errUnreachable means the API could not be reached through the new ruleset.
var errUnreachable = errors.New("control plane unreachable after apply")
// confirm reports generation as applied over a fresh connection, which proves
// the API is reachable through the new ruleset. Any HTTP response counts as
// reachable; only repeated transport failures return errUnreachable.
func (a *Agent) confirm(ctx context.Context, generation int64) error {
var err error
for i := 0; i < verifyAttempts; i++ {
if i > 0 {
select {
case <-ctx.Done():
return ctx.Err()
case <-time.After(verifyDelay):
}
}
actx, cancel := context.WithTimeout(ctx, verifyTimeout)
err = a.Client.ReportStatus(actx, Status{Status: StatusApplied, Generation: generation})
cancel()
if ctx.Err() != nil {
return ctx.Err()
}
var uerr *url.Error
if !errors.As(err, &uerr) {
if err != nil {
slog.Warn("agent: reporting status failed", "err", err)
}
return nil
}
}
return fmt.Errorf("%w: %v", errUnreachable, err)
}
// reverted is a generation rolled back after a failed apply or for severing the
// API; persisted so it is not re-applied until a newer generation arrives.
type reverted struct {
Generation int64 `json:"generation"`
Status string `json:"status,omitempty"`
Error string `json:"error,omitempty"`
Reported bool `json:"reported"`
}
func (a *Agent) revertedPath() string {
return filepath.Join(filepath.Dir(a.Cache.Path), "reverted.json")
}
func (a *Agent) readReverted() (*reverted, error) {
b, err := os.ReadFile(a.revertedPath())
if os.IsNotExist(err) {
return nil, nil
}
if err != nil {
return nil, err
}
var rv reverted
if err := json.Unmarshal(b, &rv); err != nil {
return nil, fmt.Errorf("parsing %s: %w", a.revertedPath(), err)
}
return &rv, nil
}
func (a *Agent) writeReverted(rv *reverted) error {
b, err := json.Marshal(rv)
if err != nil {
return err
}
defer unlock()
return tryapply.WriteFile(a.revertedPath(), b)
}
// reportReverted reports rv until the control plane accepts it.
func (a *Agent) reportReverted(ctx context.Context, rv *reverted) {
if rv.Reported {
return
}
status := rv.Status
if status == "" {
status = StatusReverted
}
if err := a.Client.ReportStatus(ctx, Status{Status: status, Generation: rv.Generation, Error: rv.Error}); err != nil {
slog.Warn("agent: reporting reverted generation failed, retrying next cycle", "generation", rv.Generation, "err", err)
return
}
rv.Reported = true
if err := a.writeReverted(rv); err != nil {
slog.Warn("agent: persisting reverted generation failed", "err", err)
}
}
// EngineApplier applies via the real nftables differential engine.
type EngineApplier struct{}
// Apply computes and applies the differential change set for cfg, with safe
// under a pending try as 'tomswall try' does. The caller holds the try lock.
func (EngineApplier) Apply(_ context.Context, cfg *config.Config, safe bool) (revert, keep func() error, err error) {
engine, err := nftables.NewEngine(cfg)
if err != nil {
return fmt.Errorf("initializing nftables: %w", err)
return nil, nil, fmt.Errorf("initializing nftables: %w", err)
}
changes, err := engine.Plan()
if err != nil {
return fmt.Errorf("computing changes: %w", err)
return nil, nil, fmt.Errorf("computing changes: %w", err)
}
if changes.Empty() {
return nil
return nil, nil, nil
}
return engine.Apply(changes)
if !safe {
return nil, nil, engine.Apply(changes)
}
snap, err := engine.Snapshot()
if err != nil {
return nil, nil, fmt.Errorf("snapshotting ruleset: %w", err)
}
// PID 0: 'tomswall confirm' must not signal the agent.
if _, err := tryapply.Arm(snap, 0, revertDelay); err != nil {
return nil, nil, err
}
return tryapply.Abort, tryapply.Discard, engine.Apply(changes)
}
+2 -2
View File
@@ -132,10 +132,10 @@ type fakeApplier struct {
lastGen int
}
func (f *fakeApplier) Apply(_ context.Context, cfg *config.Config) error {
func (f *fakeApplier) Apply(_ context.Context, cfg *config.Config, _ bool) (func() error, func() error, error) {
atomic.AddInt32(&f.count, 1)
f.lastGen = len(cfg.Rules)
return nil
return nil, nil, nil
}
const renderedYAML = `generation: 7
+4 -10
View File
@@ -2,7 +2,8 @@ package agent
import (
"os"
"path/filepath"
"git.unkin.net/unkin/tomswall/internal/tryapply"
)
// Cache persists the last known-good rendered config to disk so the agent can
@@ -11,16 +12,9 @@ type Cache struct {
Path string
}
// Write atomically stores the raw config bytes.
// Write durably stores the raw config bytes.
func (c Cache) Write(raw []byte) error {
if err := os.MkdirAll(filepath.Dir(c.Path), 0o755); err != nil {
return err
}
tmp := c.Path + ".tmp"
if err := os.WriteFile(tmp, raw, 0o600); err != nil {
return err
}
return os.Rename(tmp, c.Path)
return tryapply.WriteFile(c.Path, raw)
}
// Read returns the cached config, or (nil, nil) when no cache exists yet.
+26 -4
View File
@@ -26,7 +26,9 @@ func NewClient(baseURL, device, token string) *Client {
BaseURL: baseURL,
Device: device,
Token: token,
HTTP: &http.Client{Timeout: 30 * time.Second},
// No keep-alives: every request, the post-apply check included, opens a
// fresh connection that must pass the current ruleset.
HTTP: &http.Client{Timeout: 30 * time.Second, Transport: noKeepAlive()},
}
}
@@ -94,10 +96,24 @@ func (c *Client) ReportRoutes(ctx context.Context, prefixes []string) error {
return nil
}
// ReportStatus tells the control plane which generation this device has applied.
func (c *Client) ReportStatus(ctx context.Context, generation int64) error {
// Status values reported to POST /api/v1/devices/{name}/status.
const (
StatusApplied = "applied"
StatusReverted = "reverted"
StatusFailed = "failed"
)
// Status is the outcome of applying one generation.
type Status struct {
Status string `json:"status"`
Generation int64 `json:"generation"`
Error string `json:"error,omitempty"`
}
// ReportStatus tells the control plane the outcome of applying a generation.
func (c *Client) ReportStatus(ctx context.Context, st Status) error {
url := fmt.Sprintf("%s/api/v1/devices/%s/status", c.BaseURL, c.Device)
payload, _ := json.Marshal(map[string]int64{"generation": generation})
payload, _ := json.Marshal(st)
req, err := http.NewRequestWithContext(ctx, http.MethodPost, url, bytes.NewReader(payload))
if err != nil {
return err
@@ -116,3 +132,9 @@ func (c *Client) ReportStatus(ctx context.Context, generation int64) error {
}
return nil
}
func noKeepAlive() http.RoundTripper {
t := http.DefaultTransport.(*http.Transport).Clone()
t.DisableKeepAlives = true
return t
}
+396
View File
@@ -0,0 +1,396 @@
package agent
import (
"context"
"encoding/json"
"errors"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"sync"
"sync/atomic"
"testing"
"time"
"git.unkin.net/unkin/tomswall/internal/config"
"git.unkin.net/unkin/tomswall/internal/nftables"
"git.unkin.net/unkin/tomswall/internal/tryapply"
)
func TestMain(m *testing.M) {
dir, err := os.MkdirTemp("", "tomswall-agent-test")
if err != nil {
panic(err)
}
tryapply.Dir = dir
tryapply.Run = func(name string, args ...string) error {
timerCmds = append(timerCmds, name)
return nil
}
verifyDelay = time.Millisecond
verifyTimeout = time.Second
code := m.Run()
os.RemoveAll(dir)
os.Exit(code)
}
// fakeAPI serves a config generation and records status reports; while cut it
// drops connections to the status endpoint, as a severing ruleset would.
type fakeAPI struct {
*httptest.Server
gen atomic.Int64
cut atomic.Bool
code atomic.Int32
mu sync.Mutex
reports []Status
}
func newFakeAPI(t *testing.T, gen int64) *fakeAPI {
f := &fakeAPI{}
f.gen.Store(gen)
f.code.Store(http.StatusNoContent)
f.Server = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch r.URL.Path {
case "/api/v1/devices/fw-a/config":
_, _ = w.Write([]byte(strings.Replace(renderedYAML, "generation: 7", "generation: "+itoa(f.gen.Load()), 1)))
case "/api/v1/devices/fw-a/status":
if f.cut.Load() {
conn, _, _ := w.(http.Hijacker).Hijack()
conn.Close()
return
}
var st Status
_ = json.NewDecoder(r.Body).Decode(&st)
f.mu.Lock()
f.reports = append(f.reports, st)
f.mu.Unlock()
w.WriteHeader(int(f.code.Load()))
default:
w.WriteHeader(http.StatusNotFound)
}
}))
t.Cleanup(f.Close)
return f
}
// timerCmds records the systemd commands tryapply runs.
var timerCmds []string
func itoa(n int64) string { b, _ := json.Marshal(n); return string(b) }
func (f *fakeAPI) last() Status {
f.mu.Lock()
defer f.mu.Unlock()
if len(f.reports) == 0 {
return Status{}
}
return f.reports[len(f.reports)-1]
}
// fakeEngine always changes the ruleset, when safe under a real tryapply pending
// try; onApply simulates its effect and restoreErr fails the restore.
type fakeEngine struct {
applies, plain, restores int
err, restoreErr error
onApply func()
onRestore func()
}
func (f *fakeEngine) Apply(_ context.Context, _ *config.Config, safe bool) (func() error, func() error, error) {
if !safe {
f.plain++
return nil, nil, f.err
}
if _, err := tryapply.Arm(&nftables.Snapshot{Table: "tomswall"}, 0, time.Minute); err != nil {
return nil, nil, err
}
tryapply.Restore = func(*nftables.Snapshot) error {
f.restores++
if f.onRestore != nil {
f.onRestore()
}
return f.restoreErr
}
f.applies++
if f.onApply != nil {
f.onApply()
}
return tryapply.Abort, tryapply.Discard, f.err
}
// pending reports whether a snapshot is still armed and its timer not stopped since.
func pending(t *testing.T) bool {
t.Helper()
_, err := os.Stat(filepath.Join(tryapply.Dir, "try-snapshot.json"))
armed := len(timerCmds) > 0 && timerCmds[len(timerCmds)-1] == "systemd-run"
if (err == nil) != armed {
t.Fatalf("snapshot present=%v but timer armed=%v", err == nil, armed)
}
return armed
}
func newAgent(t *testing.T, api *fakeAPI, eng *fakeEngine) *Agent {
return &Agent{
Client: NewClient(api.URL, "fw-a", "tok"),
Cache: Cache{Path: filepath.Join(t.TempDir(), "rendered.yaml")},
Applier: eng,
}
}
func cachedGen(t *testing.T, a *Agent) int64 {
rc, err := a.Cache.Read()
if err != nil {
t.Fatal(err)
}
if rc == nil {
return 0
}
return rc.Generation
}
func TestSafeApplyReachableApplies(t *testing.T) {
api := newFakeAPI(t, 7)
eng := &fakeEngine{}
a := newAgent(t, api, eng)
if err := a.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
if eng.restores != 0 || api.last() != (Status{Status: StatusApplied, Generation: 7}) || cachedGen(t, a) != 7 || pending(t) {
t.Fatalf("restores=%d last=%+v cache=%d", eng.restores, api.last(), cachedGen(t, a))
}
}
func TestSafeApplyUnreachableRevertsAndReports(t *testing.T) {
api := newFakeAPI(t, 7)
eng := &fakeEngine{onApply: func() { api.cut.Store(true) }, onRestore: func() { api.cut.Store(false) }}
a := newAgent(t, api, eng)
if err := a.RunOnce(context.Background()); !errors.Is(err, errUnreachable) {
t.Fatalf("want errUnreachable, got %v", err)
}
if eng.restores != 1 || cachedGen(t, a) != 0 || pending(t) {
t.Fatalf("restores=%d cache=%d", eng.restores, cachedGen(t, a))
}
if st := api.last(); st.Status != StatusReverted || st.Generation != 7 || st.Error == "" {
t.Fatalf("last report %+v", st)
}
rv, _ := a.readReverted()
if rv == nil || rv.Generation != 7 || !rv.Reported {
t.Fatalf("persisted %+v", rv)
}
}
func TestSafeApplyRevertReportedOnceReachable(t *testing.T) {
api := newFakeAPI(t, 7)
eng := &fakeEngine{onApply: func() { api.cut.Store(true) }}
a := newAgent(t, api, eng)
_ = a.RunOnce(context.Background())
if eng.restores != 1 || api.last().Status != "" {
t.Fatalf("restores=%d last=%+v", eng.restores, api.last())
}
api.cut.Store(false)
if err := a.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
if eng.applies != 1 || api.last() != (Status{Status: StatusReverted, Generation: 7, Error: api.last().Error}) {
t.Fatalf("applies=%d last=%+v", eng.applies, api.last())
}
}
func TestSafeApplyShutdownDoesNotRevert(t *testing.T) {
api := newFakeAPI(t, 7)
ctx, cancel := context.WithCancel(context.Background())
eng := &fakeEngine{onApply: func() { api.cut.Store(true); cancel() }}
a := newAgent(t, api, eng)
if err := a.RunOnce(ctx); !errors.Is(err, context.Canceled) {
t.Fatalf("want context.Canceled, got %v", err)
}
if rv, _ := a.readReverted(); eng.restores != 0 || rv != nil || cachedGen(t, a) != 0 {
t.Fatalf("restores=%d reverted=%+v cache=%d", eng.restores, rv, cachedGen(t, a))
}
}
func TestSafeApplyHTTPErrorDoesNotRevert(t *testing.T) {
api := newFakeAPI(t, 7)
api.code.Store(http.StatusInternalServerError)
eng := &fakeEngine{}
a := newAgent(t, api, eng)
if err := a.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
if eng.restores != 0 || cachedGen(t, a) != 7 {
t.Fatalf("restores=%d cache=%d", eng.restores, cachedGen(t, a))
}
}
func TestSafeApplyApplyErrorRestoresAndReportsFailed(t *testing.T) {
api := newFakeAPI(t, 7)
eng := &fakeEngine{err: errors.New("netlink: boom")}
a := newAgent(t, api, eng)
if err := a.RunOnce(context.Background()); err == nil {
t.Fatal("want error")
}
if st := api.last(); eng.restores != 1 || st.Status != StatusFailed || !strings.Contains(st.Error, "boom") || pending(t) {
t.Fatalf("restores=%d last=%+v", eng.restores, st)
}
if rv, _ := a.readReverted(); rv == nil || rv.Generation != 7 || !rv.Reported {
t.Fatalf("persisted %+v", rv)
}
}
func TestSafeApplyApplyErrorRestoreFailsKeepsTimer(t *testing.T) {
t.Cleanup(func() { _ = tryapply.Discard() })
api := newFakeAPI(t, 7)
eng := &fakeEngine{err: errors.New("netlink: boom"), restoreErr: errors.New("netlink: stuck")}
a := newAgent(t, api, eng)
if err := a.RunOnce(context.Background()); err == nil || !strings.Contains(err.Error(), "restore: restoring snapshot: netlink: stuck") {
t.Fatalf("got %v", err)
}
want := "apply: netlink: boom; restore: restoring snapshot: netlink: stuck; revert timer pending"
if st := api.last(); st != (Status{Status: StatusFailed, Generation: 7, Error: want}) || !pending(t) {
t.Fatalf("last=%+v", st)
}
if rv, _ := a.readReverted(); rv == nil || rv.Generation != 7 || rv.Status != StatusFailed {
t.Fatalf("persisted %+v", rv)
}
// The next cycle waits for the timer instead of re-applying.
if err := a.RunOnce(context.Background()); err != nil || eng.applies != 1 {
t.Fatalf("err=%v applies=%d", err, eng.applies)
}
}
func TestSafeApplyUnreachableRestoreFailsKeepsTimer(t *testing.T) {
t.Cleanup(func() { _ = tryapply.Discard() })
api := newFakeAPI(t, 7)
eng := &fakeEngine{onApply: func() { api.cut.Store(true) }, restoreErr: errors.New("netlink: stuck")}
eng.onRestore = func() { api.cut.Store(false) }
a := newAgent(t, api, eng)
if err := a.RunOnce(context.Background()); !errors.Is(err, errUnreachable) || !strings.Contains(err.Error(), "revert timer pending") {
t.Fatalf("got %v", err)
}
st := api.last()
if st.Status != StatusFailed || st.Generation != 7 || !strings.HasPrefix(st.Error, errUnreachable.Error()) ||
!strings.HasSuffix(st.Error, "; restore: restoring snapshot: netlink: stuck; revert timer pending") || !pending(t) {
t.Fatalf("last=%+v", st)
}
if rv, _ := a.readReverted(); rv == nil || rv.Generation != 7 || !rv.Reported || cachedGen(t, a) != 0 {
t.Fatalf("persisted %+v cache=%d", rv, cachedGen(t, a))
}
}
func TestSafeApplyRevertedGenerationSkippedAfterRestart(t *testing.T) {
api := newFakeAPI(t, 7)
eng := &fakeEngine{}
a := newAgent(t, api, eng)
if err := a.writeReverted(&reverted{Generation: 7, Reported: true}); err != nil {
t.Fatal(err)
}
if err := a.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
if eng.applies != 0 {
t.Fatalf("reverted generation re-applied")
}
api.gen.Store(8)
if err := a.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
if rv, _ := a.readReverted(); eng.applies != 1 || api.last().Generation != 8 || rv != nil {
t.Fatalf("applies=%d last=%+v reverted=%+v", eng.applies, api.last(), rv)
}
}
func TestSafeApplySkipsWhileTryPending(t *testing.T) {
marker := filepath.Join(tryapply.Dir, "try-snapshot.json")
if err := os.WriteFile(marker, []byte("{}"), 0o600); err != nil {
t.Fatal(err)
}
defer os.Remove(marker)
api := newFakeAPI(t, 7)
eng := &fakeEngine{}
a := newAgent(t, api, eng)
if err := a.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
if eng.applies != 0 || api.last().Status != "" {
t.Fatalf("applies=%d last=%+v", eng.applies, api.last())
}
}
// failArm makes arming the revert timer fail, as without systemd.
func failArm(t *testing.T) {
orig := tryapply.Run
tryapply.Run = func(name string, args ...string) error {
timerCmds = append(timerCmds, name)
if name == "systemd-run" {
return errors.New("no systemd")
}
return nil
}
t.Cleanup(func() { tryapply.Run = orig })
}
func TestSafeApplyArmFailureReportsFailedAndRetries(t *testing.T) {
failArm(t)
api := newFakeAPI(t, 7)
eng := &fakeEngine{}
a := newAgent(t, api, eng)
if err := a.RunOnce(context.Background()); err == nil || !strings.Contains(err.Error(), "no systemd") {
t.Fatalf("got %v", err)
}
if st := api.last(); eng.applies != 0 || st.Status != StatusFailed || st.Generation != 7 || pending(t) || cachedGen(t, a) != 0 {
t.Fatalf("applies=%d last=%+v", eng.applies, st)
}
if rv, _ := a.readReverted(); rv != nil || a.lastReverted != nil {
t.Fatalf("arm failure marked generation reverted: %+v", rv)
}
tryapply.Run = func(name string, args ...string) error {
timerCmds = append(timerCmds, name)
return nil
}
if err := a.RunOnce(context.Background()); err != nil || eng.applies != 1 || cachedGen(t, a) != 7 {
t.Fatalf("retry err=%v applies=%d", err, eng.applies)
}
}
func TestCachedConfigAppliesWithoutArm(t *testing.T) {
failArm(t)
api := newFakeAPI(t, 7)
eng := &fakeEngine{}
a := newAgent(t, api, eng)
if err := a.Cache.Write([]byte(renderedYAML)); err != nil {
t.Fatal(err)
}
api.Close()
if err := a.RunOnce(context.Background()); err != nil {
t.Fatal(err)
}
if eng.plain != 1 || eng.applies != 0 || pending(t) {
t.Fatalf("plain=%d safe=%d", eng.plain, eng.applies)
}
}
func TestSafeApplyRevertedKeptInMemoryWhenPersistFails(t *testing.T) {
api := newFakeAPI(t, 7)
eng := &fakeEngine{onRestore: func() { api.cut.Store(false) }}
a := newAgent(t, api, eng)
// A non-empty directory in its place makes persisting reverted.json fail.
eng.onApply = func() {
api.cut.Store(true)
_ = os.MkdirAll(filepath.Join(a.revertedPath(), "x"), 0o755)
}
if err := a.RunOnce(context.Background()); !errors.Is(err, errUnreachable) {
t.Fatalf("want errUnreachable, got %v", err)
}
if err := os.RemoveAll(a.revertedPath()); err != nil {
t.Fatal(err)
}
if eng.restores != 1 || api.last().Status != StatusReverted || pending(t) {
t.Fatalf("restores=%d last=%+v", eng.restores, api.last())
}
if err := a.RunOnce(context.Background()); err != nil || eng.applies != 1 {
t.Fatalf("reverted generation re-applied: err=%v applies=%d", err, eng.applies)
}
}
+7
View File
@@ -531,6 +531,13 @@ func TestValidateHosts(t *testing.T) {
},
wantErr: "interface required",
},
{
name: "invalid exclusion",
zones: map[string]Zone{"fw": {Type: ZoneFirewall}, "net": {Type: ZoneIP}, "loc": {Type: ZoneIP}},
interfaces: []Interface{{Zone: "net", Interface: "eth0"}},
hosts: []Host{{Zone: "loc", Interface: "eth0", Addresses: []string{"192.0.2.0/24"}, Exclusions: []string{"192.0.2.0/24!192.0.2.7"}}},
wantErr: "invalid address",
},
}
for _, tt := range tests {
+12 -1
View File
@@ -1,6 +1,10 @@
package config
import "fmt"
import (
"fmt"
"net/netip"
"slices"
)
type Host struct {
Zone string `yaml:"zone"`
@@ -52,6 +56,13 @@ func (c *Config) validateHosts() error {
if !h.Dynamic && len(h.Addresses) == 0 {
return fmt.Errorf("host[%d]: at least one address required (or set dynamic: true)", i)
}
for _, a := range slices.Concat(h.Addresses, h.Exclusions) {
if _, err := netip.ParsePrefix(a); err != nil {
if _, err := netip.ParseAddr(a); err != nil {
return fmt.Errorf("host[%d]: invalid address %q", i, a)
}
}
}
}
return nil
}
+234 -37
View File
@@ -5,6 +5,7 @@ import (
"fmt"
"log/slog"
"net"
"net/netip"
"slices"
"sort"
"strconv"
@@ -345,12 +346,12 @@ func (c *Compiler) compileConntrackPair(state *FirewallState, tag, chain string,
(dstAddr == "" || strings.HasPrefix(dstAddr, "!")) {
return fmt.Errorf("conntrack DEST zone %q needs an address in prerouting", dstZone)
}
srcIfaces, dstIfaces := c.resolveZoneInterfaces(srcZone, srcAddr), []string{""}
if chain == "raw_prerouting" && c.resolveZoneInterfaces(dstZone, dstAddr) == nil {
srcIfaces, dstIfaces := c.resolveZone(srcZone, srcAddr), []zoneMatch{{}}
if chain == "raw_prerouting" && c.resolveZone(dstZone, dstAddr) == nil {
return nil
}
if chain == "raw_output" {
srcIfaces, dstIfaces = []string{""}, c.resolveZoneInterfaces(dstZone, dstAddr)
srcIfaces, dstIfaces = []zoneMatch{{}}, c.resolveZone(dstZone, dstAddr)
}
out := chain
if ct.Action == config.ConntrackHelper {
@@ -674,8 +675,8 @@ func splitAddrs(addr string) []string {
func (c *Compiler) compileZonePair(state *FirewallState, tag, srcZone, srcAddr, dstZone, dstAddr, origDest, proto string,
dports, sports config.PortSpec, action config.RuleAction, logLevel string,
fwZone string, section config.RuleSection) error {
srcIfaces := c.resolveZoneInterfaces(srcZone, srcAddr)
dstIfaces := c.resolveZoneInterfaces(dstZone, dstAddr)
srcIfaces := c.resolveZone(srcZone, srcAddr)
dstIfaces := c.resolveZone(dstZone, dstAddr)
chain := c.selectChain(srcZone, dstZone, fwZone)
// ponytail: forward daddr is post-DNAT; lift with `ct original daddr` (expr.Ct Direction, google/nftables v0.3.0).
if origDest != "" && chain == "forward" {
@@ -744,7 +745,7 @@ func (c *Compiler) compileDNATRule(state *FirewallState, tag, srcZone, srcAddr,
dnatPort = uint16(p)
}
srcIfaces := c.resolveZoneInterfaces(srcZone, srcAddr)
srcIfaces := c.resolveZone(srcZone, srcAddr)
var odExprs []expr.Any
if origDest != "" {
@@ -765,12 +766,12 @@ func (c *Compiler) compileDNATRule(state *FirewallState, tag, srcZone, srcAddr,
}
for _, srcIface := range srcIfaces {
zm, err := zoneMatchExprs(srcIface, true)
if err != nil {
return err
}
for _, m := range matches {
var exprs []expr.Any
if srcIface != "" {
exprs = append(exprs, matchIfaceName(true, srcIface)...)
}
exprs := slices.Clone(zm)
if srcAddr != "" {
src, err := matchSourceCIDR(srcAddr)
@@ -857,31 +858,35 @@ func (c *Compiler) compileDNATRule(state *FirewallState, tag, srcZone, srcAddr,
func (c *Compiler) compilePolicies(state *FirewallState) error {
fwZone := c.cfg.FirewallZone()
overridden := map[string]bool{}
for i, pol := range c.cfg.Policy {
tag := fmt.Sprintf("policy:%d", i)
explicitIntra := pol.Source == pol.Dest && !isGlobalZone(pol.Source)
srcZones := c.expandZoneRef(pol.Source)
dstZones := c.expandZoneRef(pol.Dest)
for _, sz := range srcZones {
for _, dz := range dstZones {
if sz == dz && !strings.HasSuffix(pol.Source, "+") {
continue
if sz == dz {
if sz == fwZone || (!explicitIntra && !strings.HasSuffix(pol.Source, "+") && !strings.HasSuffix(pol.Dest, "+")) {
continue
}
overridden[sz] = true
}
chain := c.selectChain(sz, dz, fwZone)
srcIfaces := c.resolveZoneInterfaces(sz, "")
dstIfaces := c.resolveZoneInterfaces(dz, "")
srcIfaces := c.resolveZone(sz, "")
dstIfaces := c.resolveZone(dz, "")
for _, si := range srcIfaces {
for _, di := range dstIfaces {
var exprs []expr.Any
if si != "" {
exprs = append(exprs, matchIfaceName(true, si)...)
if sz == dz && intraZoneSkip(si, di) {
continue
}
if di != "" && chain != "input" {
exprs = append(exprs, matchIfaceName(false, di)...)
exprs, err := zonePairExprs(si, di, chain)
if err != nil {
return fmt.Errorf("policy[%d]: %w", i, err)
}
if pol.RateLimit != "" {
@@ -909,9 +914,66 @@ func (c *Compiler) compilePolicies(state *FirewallState) error {
}
}
return c.compileImplicitIntraZone(state, overridden)
}
// compileImplicitIntraZone accepts traffic between different interfaces of one zone, shorewall's implicit intra-zone ACCEPT policy.
func (c *Compiler) compileImplicitIntraZone(state *FirewallState, overridden map[string]bool) error {
fwZone := c.cfg.FirewallZone()
zones := make([]string, 0, len(c.cfg.Zones))
for z := range c.cfg.Zones {
zones = append(zones, z)
}
sort.Strings(zones)
for _, z := range zones {
if z == fwZone || overridden[z] {
continue
}
if len(c.cfg.ZoneInterfaces(z)) == 0 && !slices.ContainsFunc(c.cfg.Hosts, func(h config.Host) bool { return h.Zone == z }) {
continue
}
matches := c.resolveZone(z, "")
for _, si := range matches {
for _, di := range matches {
if intraZoneSkip(si, di) {
continue
}
exprs, err := zonePairExprs(si, di, "forward")
if err != nil {
return fmt.Errorf("zone %s: %w", z, err)
}
state.Rules["forward"] = append(state.Rules["forward"], ManagedRule{
Chain: "forward",
Exprs: append(exprs, &expr.Verdict{Kind: expr.VerdictAccept}),
Tag: "intra:" + z,
})
}
}
}
return nil
}
// intraZoneSkip drops intra-zone pairs on one interface unless both are routeback hosts entries
// (interface routeback is compileIntraZone's job), and pairs whose address families can never both match.
func intraZoneSkip(si, di zoneMatch) bool {
if si.iface != "" && si.iface == di.iface && !(si.routeback && di.routeback) {
return true
}
a, b := matchFamily(si), matchFamily(di)
return a != 0 && b != 0 && a != b
}
// matchFamily is the NFPROTO a zoneMatch is guarded by: its host address's family, else fam.
func matchFamily(m zoneMatch) byte {
if p, err := parsePrefix(m.addr); err == nil {
if p.Addr().Is4() {
return unix.NFPROTO_IPV4
}
return unix.NFPROTO_IPV6
}
return m.fam
}
func (c *Compiler) compileSNAT(state *FirewallState) error {
for i, snat := range c.cfg.SNAT {
tag := fmt.Sprintf("snat:%d", i)
@@ -1213,11 +1275,23 @@ func (c *Compiler) selectChain(srcZone, dstZone, fwZone string) string {
return "forward"
}
// resolveZoneInterfaces returns nil (fail closed) for an unknown zone, or one with no interfaces unless a non-negated address match narrows the rule.
func (c *Compiler) resolveZoneInterfaces(zone, addr string) []string {
// zoneMatch classifies a packet into a zone: an interface (empty: any) and, for a hosts entry, one host
// address; excl carves out hosts exclusions and the hosts of sub-zones, which shorewall matches first.
// fam (an NFPROTO; addr's family when set, else 0: any) guards addr and excl so IPv4 offsets are
// never compared against IPv6 bytes. routeback marks a hosts entry with the routeback option.
type zoneMatch struct {
iface, addr string
excl []string
fam byte
routeback bool
}
// resolveZone returns nil (fail closed) for an unknown zone, or one with neither interfaces nor hosts
// unless a non-negated address match narrows the rule.
func (c *Compiler) resolveZone(zone, addr string) []zoneMatch {
switch zone {
case "", "all", "all+", "any", "any+":
return []string{""}
return []zoneMatch{{}}
}
z, ok := c.cfg.Zones[zone]
if !ok {
@@ -1225,24 +1299,151 @@ func (c *Compiler) resolveZoneInterfaces(zone, addr string) []string {
return nil
}
if z.Type == config.ZoneFirewall {
return []string{""}
return []zoneMatch{{}}
}
if ifaces := c.cfg.ZoneInterfaces(zone); len(ifaces) > 0 {
return ifaces
var out []zoneMatch
hasHosts := false
for _, iface := range c.cfg.ZoneInterfaces(zone) {
sub, back := c.subZoneHosts(zone, iface)
if len(sub) == 0 {
out = append(out, zoneMatch{iface: iface})
} else {
v4, v6 := splitFamily(sub)
out = append(out, zoneMatch{iface: iface, excl: v4, fam: unix.NFPROTO_IPV4},
zoneMatch{iface: iface, excl: v6, fam: unix.NFPROTO_IPV6})
}
for _, b := range back {
if addrsOverlap(b, addr) {
out = append(out, zoneMatch{iface: iface, addr: b})
}
}
}
for _, h := range c.cfg.Hosts {
if h.Zone != zone {
continue
}
hasHosts = true
sub, _ := c.subZoneHosts(zone, h.Interface)
v4, v6 := splitFamily(slices.Concat(h.Exclusions, sub))
for _, a := range h.Addresses {
if addrsOverlap(a, addr) {
m := zoneMatch{iface: h.Interface, addr: a, excl: v6, routeback: h.Options.RouteBack}
if p, err := parsePrefix(a); err == nil && p.Addr().Is4() {
m.excl = v4
}
out = append(out, m)
}
}
}
if len(out) > 0 || hasHosts {
return out
}
if addr != "" && !strings.HasPrefix(addr, "!") {
return []string{""}
return []zoneMatch{{}}
}
if !c.warned[zone] {
if c.warned == nil {
c.warned = map[string]bool{}
}
c.warned[zone] = true
slog.Warn("compiler: zone has no interfaces, skipping its rules", "zone", zone)
slog.Warn("compiler: zone has no interfaces or hosts, skipping its rules", "zone", zone)
}
return nil
}
// subZoneHosts lists the host addresses on iface that belong to sub-zones of zone, and those
// sub-zone hosts' exclusions, which fall back to zone.
func (c *Compiler) subZoneHosts(zone, iface string) (sub, back []string) {
for _, h := range c.cfg.Hosts {
if h.Interface == iface && c.cfg.IsSubZone(h.Zone, zone) {
sub = append(sub, h.Addresses...)
back = append(back, h.Exclusions...)
}
}
return sub, back
}
// addrsOverlap reports whether a host address can match a rule address; unparsable or negated rule addresses keep the host.
func addrsOverlap(host, rule string) bool {
if rule == "" || strings.HasPrefix(rule, "!") {
return true
}
h, err := parsePrefix(host)
if err != nil {
return true
}
for _, r := range strings.Split(rule, ",") {
if p, err := parsePrefix(r); err != nil || p.Overlaps(h) {
return true
}
}
return false
}
// splitFamily partitions addresses by family; unparsable ones go to v6 so zoneMatchExprs still rejects them.
func splitFamily(addrs []string) (v4, v6 []string) {
for _, a := range addrs {
if p, err := parsePrefix(a); err == nil && p.Addr().Is4() {
v4 = append(v4, a)
} else {
v6 = append(v6, a)
}
}
return v4, v6
}
func parsePrefix(s string) (netip.Prefix, error) {
if a, err := netip.ParseAddr(s); err == nil {
return netip.PrefixFrom(a, a.BitLen()), nil
}
return netip.ParsePrefix(s)
}
// zoneMatchExprs matches a zone on the in (src) or out interface plus its host address and exclusions,
// all guarded by m.fam so an IPv4 address never matches IPv6 bytes in an inet table.
func zoneMatchExprs(m zoneMatch, src bool) ([]expr.Any, error) {
var out []expr.Any
if m.iface != "" {
out = matchIfaceName(src, m.iface)
}
var addr []expr.Any
if m.addr != "" {
p, err := parsePrefix(m.addr)
if err != nil {
return nil, fmt.Errorf("invalid host address %q", m.addr)
}
if addr, err = matchAddrCIDR(m.addr, src); err != nil {
return nil, err
}
m.fam = unix.NFPROTO_IPV6
if p.Addr().Is4() {
m.fam = unix.NFPROTO_IPV4
}
}
if m.fam != 0 {
out = append(out, matchNFProto(m.fam)...)
}
out = append(out, addr...)
if len(m.excl) > 0 {
e, err := matchAddrCIDR("!"+strings.Join(m.excl, ","), src)
if err != nil {
return nil, err
}
out = append(out, e...)
}
return out, nil
}
// zonePairExprs matches the source zone inbound and, outside input, the dest zone outbound.
func zonePairExprs(src, dst zoneMatch, chain string) ([]expr.Any, error) {
out, err := zoneMatchExprs(src, true)
if err != nil || chain == "input" {
return out, err
}
d, err := zoneMatchExprs(dst, false)
return append(out, d...), err
}
func (c *Compiler) expandZoneRef(ref string) []string {
base := ref
var excluded map[string]bool
@@ -1272,14 +1473,10 @@ func (c *Compiler) expandZoneRef(ref string) []string {
return []string{base}
}
func (c *Compiler) buildMatchExprs(srcIface, dstIface, chain, proto string, dports, sports config.PortSpec, srcAddr, dstAddr string) ([]l4Match, error) {
var exprs []expr.Any
if srcIface != "" {
exprs = append(exprs, matchIfaceName(true, srcIface)...)
}
if dstIface != "" && chain != "input" {
exprs = append(exprs, matchIfaceName(false, dstIface)...)
func (c *Compiler) buildMatchExprs(srcIface, dstIface zoneMatch, chain, proto string, dports, sports config.PortSpec, srcAddr, dstAddr string) ([]l4Match, error) {
exprs, err := zonePairExprs(srcIface, dstIface, chain)
if err != nil {
return nil, err
}
if srcAddr != "" {
+363 -20
View File
@@ -7,6 +7,7 @@ import (
"log/slog"
"net"
"reflect"
"slices"
"strings"
"testing"
@@ -144,19 +145,17 @@ func TestCompiler_ResolveZoneInterfaces(t *testing.T) {
}
c := NewCompiler(cfg)
ifaces := c.resolveZoneInterfaces("net", "")
if len(ifaces) != 1 || ifaces[0] != "eth0" {
t.Errorf("resolveZoneInterfaces(net) = %v, want [eth0]", ifaces)
}
ifaces = c.resolveZoneInterfaces("all", "")
if len(ifaces) != 1 || ifaces[0] != "" {
t.Errorf("resolveZoneInterfaces(all) = %v, want [\"\"]", ifaces)
}
ifaces = c.resolveZoneInterfaces("fw", "")
if len(ifaces) != 1 || ifaces[0] != "" {
t.Errorf("resolveZoneInterfaces(fw) = %v, want [\"\"]", ifaces)
for _, tt := range []struct {
zone string
want []zoneMatch
}{
{"net", []zoneMatch{{iface: "eth0"}}},
{"all", []zoneMatch{{}}},
{"fw", []zoneMatch{{}}},
} {
if got := c.resolveZone(tt.zone, ""); !reflect.DeepEqual(got, tt.want) {
t.Errorf("resolveZone(%s) = %v, want %v", tt.zone, got, tt.want)
}
}
}
@@ -1849,13 +1848,13 @@ func TestCompile_InterfacelessZonesFailClosed(t *testing.T) {
c.Zones["hst"] = config.Zone{Type: config.ZoneIP}
c.Hosts = []config.Host{{Zone: "hst", Interface: "eth0", Addresses: []string{"192.0.2.0/24"}}}
c.Policy = []config.Policy{{Source: "fw", Dest: "hst", Action: config.PolicyAccept}}
}, "policy:0", 0, []string{"hst"}},
}, "policy:0", 1, nil},
{"fw all expansion keeps zones with interfaces", func(c *config.Config) {
ipsec(c)
c.Zones["hst"] = config.Zone{Type: config.ZoneIP}
c.Hosts = []config.Host{{Zone: "hst", Interface: "eth0", Addresses: []string{"192.0.2.0/24"}}}
c.Policy = []config.Policy{{Source: "fw", Dest: "all", Action: config.PolicyDrop}}
}, "policy:0", 1, []string{"hst", "ips"}},
}, "policy:0", 2, []string{"ips"}},
{"negated address does not scope", rule("ips:!192.0.2.1"), "rule:0", 0, []string{"ips"}},
{"address scopes", rule("ips:192.0.2.1"), "rule:0", 1, []string{"ips"}},
}
@@ -1895,6 +1894,19 @@ func TestCompile_InterfacelessZonesFailClosed(t *testing.T) {
func describeRule(r ManagedRule) string {
var parts []string
for i, e := range r.Exprs {
if p, ok := e.(*expr.Payload); ok && p.Base == expr.PayloadBaseNetworkHeader && i+2 < len(r.Exprs) {
bw, okb := r.Exprs[i+1].(*expr.Bitwise)
cmp, okc := r.Exprs[i+2].(*expr.Cmp)
if okb && okc {
name := map[uint32]string{12: "saddr", 16: "daddr", 8: "saddr", 24: "daddr"}[p.Offset]
if cmp.Op == expr.CmpOpNeq {
name = "!" + name
}
ones, _ := net.IPMask(bw.Mask).Size()
parts = append(parts, fmt.Sprintf("%s=%s/%d", name, net.IP(cmp.Data), ones))
continue
}
}
cmp, ok := func() (*expr.Cmp, bool) {
if i+1 >= len(r.Exprs) {
return nil, false
@@ -1912,6 +1924,8 @@ func describeRule(r ManagedRule) string {
parts = append(parts, "iif="+strings.TrimRight(string(cmp.Data), "\x00"))
case expr.MetaKeyOIFNAME:
parts = append(parts, "oif="+strings.TrimRight(string(cmp.Data), "\x00"))
case expr.MetaKeyNFPROTO:
parts = append(parts, map[byte]string{unix.NFPROTO_IPV4: "ip4", unix.NFPROTO_IPV6: "ip6"}[cmp.Data[0]])
}
case *expr.Payload:
if m.Base == expr.PayloadBaseNetworkHeader && (m.Len == 4 || m.Len == 16) {
@@ -2050,25 +2064,25 @@ func TestCompile_CommaZoneLists(t *testing.T) {
{
name: "dnat origdest",
rule: config.Rule{Action: config.RuleDNAT, Source: "net", Dest: "svr:192.0.2.17", Proto: "tcp", DPort: config.PortSpec{"80"}, OrigDest: "203.0.113.5"},
want: map[string][]string{"prerouting": {"iif=eth0 daddr=203.0.113.5"},
want: map[string][]string{"prerouting": {"iif=eth0 ip4 daddr=203.0.113.5"},
"forward": {"iif=eth0 oif=eth2 daddr=192.0.2.17"}},
},
{
name: "dnat origdest list",
rule: config.Rule{Action: config.RuleDNAT, Source: "net", Dest: "svr:192.0.2.17", Proto: "tcp", DPort: config.PortSpec{"80"}, OrigDest: "203.0.113.5,203.0.113.6"},
want: map[string][]string{"prerouting": {"iif=eth0 daddr=203.0.113.5", "iif=eth0 daddr=203.0.113.6"},
want: map[string][]string{"prerouting": {"iif=eth0 ip4 daddr=203.0.113.5", "iif=eth0 ip4 daddr=203.0.113.6"},
"forward": {"iif=eth0 oif=eth2 daddr=192.0.2.17"}},
},
{
name: "dnat negated origdest list",
rule: config.Rule{Action: config.RuleDNAT, Source: "net", Dest: "svr:192.0.2.17", Proto: "tcp", DPort: config.PortSpec{"80"}, OrigDest: "!203.0.113.5,203.0.113.6"},
want: map[string][]string{"prerouting": {"iif=eth0 !daddr=203.0.113.5 !daddr=203.0.113.6"},
want: map[string][]string{"prerouting": {"iif=eth0 ip4 !daddr=203.0.113.5 !daddr=203.0.113.6"},
"forward": {"iif=eth0 oif=eth2 daddr=192.0.2.17"}},
},
{
name: "accept origdest",
rule: config.Rule{Action: config.RuleAccept, Source: "net", Dest: "fw", Proto: "tcp", DPort: config.PortSpec{"22"}, OrigDest: "203.0.113.5"},
want: map[string][]string{"input": {"iif=eth0 daddr=203.0.113.5"}},
want: map[string][]string{"input": {"iif=eth0 ip4 daddr=203.0.113.5"}},
},
{
name: "origdest does not scope interface-less zone",
@@ -2078,7 +2092,7 @@ func TestCompile_CommaZoneLists(t *testing.T) {
{
name: "accept ipv6 origdest",
rule: config.Rule{Action: config.RuleAccept, Source: "net", Dest: "fw", Proto: "tcp", DPort: config.PortSpec{"22"}, OrigDest: "2001:db8::5"},
want: map[string][]string{"input": {"iif=eth0 daddr=2001:db8::5"}},
want: map[string][]string{"input": {"iif=eth0 ip6 daddr=2001:db8::5"}},
},
{
name: "blrule zone list",
@@ -3320,3 +3334,332 @@ func TestCompile_LogLimitSplitsAroundExtrasAndNAT(t *testing.T) {
}
}
}
// hostsCfg models a shorewall setup where lan:net is defined by hosts on net's interfaces.
func hostsCfg(mod func(*config.Config)) *config.Config {
cfg := &config.Config{
Settings: config.Settings{TableName: "test", AddressFamily: config.FamilyINET},
Zones: map[string]config.Zone{
"fw": {Type: config.ZoneFirewall}, "net": {Type: config.ZoneIP},
"lan": {Type: config.ZoneIP, Parents: []string{"net"}}, "vpn": {Type: config.ZoneIP},
},
Interfaces: []config.Interface{
{Zone: "net", Interface: "wlo1"}, {Zone: "net", Interface: "enp2s0"}, {Zone: "vpn", Interface: "tun0"},
},
Hosts: []config.Host{
{Zone: "lan", Interface: "wlo1", Addresses: []string{"192.0.2.0/24"}},
{Zone: "lan", Interface: "enp2s0", Addresses: []string{"198.51.100.0/24"}},
},
Policy: []config.Policy{
{Source: "net", Dest: "all", Action: config.PolicyDrop},
{Source: "lan", Dest: "fw", Action: config.PolicyReject},
{Source: "all", Dest: "all", Action: config.PolicyReject},
},
Rules: []config.Rule{
{Action: config.RuleAccept, Source: "lan,vpn", Dest: "fw", Proto: "tcp", DPort: config.PortSpec{"6768"}},
},
PortGroups: make(map[string]config.PortGroup),
}
if mod != nil {
mod(cfg)
}
return cfg
}
func describeTagged(state *FirewallState, chain, tag string) []string {
var out []string
for _, r := range taggedRules(state, chain, tag) {
out = append(out, describeRule(r))
}
return out
}
func TestCompile_HostsZoneRules(t *testing.T) {
var logs bytes.Buffer
prev := slog.Default()
slog.SetDefault(slog.New(slog.NewTextHandler(&logs, nil)))
defer slog.SetDefault(prev)
state := mustCompile(t, hostsCfg(nil))
want := []string{"iif=wlo1 ip4 saddr=192.0.2.0/24", "iif=enp2s0 ip4 saddr=198.51.100.0/24", "iif=tun0"}
if got := describeTagged(state, "input", "rule:0"); !reflect.DeepEqual(got, want) {
t.Errorf("input rule:0 = %q, want %q", got, want)
}
if strings.Contains(logs.String(), "skipping") {
t.Errorf("unexpected warning:\n%s", logs.String())
}
}
func TestCompile_HostsSubZoneBeforeParent(t *testing.T) {
state := mustCompile(t, hostsCfg(nil))
want := []string{
"iif=wlo1 ip4 !saddr=192.0.2.0/24", "iif=wlo1 ip6",
"iif=enp2s0 ip4 !saddr=198.51.100.0/24", "iif=enp2s0 ip6",
}
if got := describeTagged(state, "input", "policy:0"); !reflect.DeepEqual(got, want) {
t.Errorf("net->fw policy = %q, want %q", got, want)
}
want = []string{"iif=wlo1 ip4 saddr=192.0.2.0/24", "iif=enp2s0 ip4 saddr=198.51.100.0/24"}
if got := describeTagged(state, "input", "policy:1"); !reflect.DeepEqual(got, want) {
t.Errorf("lan->fw policy = %q, want %q", got, want)
}
for _, r := range taggedRules(state, "input", "policy:1") {
if !slices.ContainsFunc(r.Exprs, func(e expr.Any) bool {
m, ok := e.(*expr.Meta)
return ok && m.Key == expr.MetaKeyNFPROTO
}) {
t.Errorf("host match lacks an nfproto guard: %v", describeRule(r))
}
}
}
func TestCompile_HostsZoneMatches(t *testing.T) {
tests := []struct {
name string
mod func(*config.Config)
chain string
tag string
want []string
}{
{"forward to hosts zone", func(c *config.Config) {
c.Rules = []config.Rule{{Action: config.RuleAccept, Source: "vpn", Dest: "lan", Proto: "tcp"}}
}, "forward", "rule:0", []string{"iif=tun0 oif=wlo1 ip4 daddr=192.0.2.0/24", "iif=tun0 oif=enp2s0 ip4 daddr=198.51.100.0/24"}},
{"rule address prunes non-overlapping hosts", func(c *config.Config) {
c.Rules = []config.Rule{{Action: config.RuleAccept, Source: "lan:192.0.2.5", Dest: "fw", Proto: "tcp"}}
}, "input", "rule:0", []string{"iif=wlo1 ip4 saddr=192.0.2.0/24 saddr=192.0.2.5"}},
{"host exclusions and sub-zone exclusions fall back to the parent", func(c *config.Config) {
c.Hosts[1].Exclusions = []string{"198.51.100.7"}
c.Rules = []config.Rule{{Action: config.RuleAccept, Source: "net", Dest: "fw", Proto: "tcp"}}
}, "input", "rule:0", []string{
"iif=wlo1 ip4 !saddr=192.0.2.0/24", "iif=wlo1 ip6",
"iif=enp2s0 ip4 !saddr=198.51.100.0/24", "iif=enp2s0 ip6",
"iif=enp2s0 ip4 saddr=198.51.100.7",
}},
{"host exclusion", func(c *config.Config) {
c.Hosts[1].Exclusions = []string{"198.51.100.7"}
}, "input", "policy:1", []string{"iif=wlo1 ip4 saddr=192.0.2.0/24", "iif=enp2s0 ip4 saddr=198.51.100.0/24 !saddr=198.51.100.7"}},
{"DNAT from hosts zone", func(c *config.Config) {
c.Rules = []config.Rule{{Action: config.RuleDNAT, Source: "lan", Dest: "vpn:203.0.113.10", Proto: "tcp", DPort: config.PortSpec{"80"}}}
}, "prerouting", "rule:0", []string{"iif=wlo1 ip4 saddr=192.0.2.0/24", "iif=enp2s0 ip4 saddr=198.51.100.0/24"}},
{"conntrack from hosts zone", func(c *config.Config) {
c.Conntrack = []config.ConntrackRule{{Action: config.ConntrackNoTrack, Source: "lan", Proto: "udp"}}
}, "raw_prerouting", "conntrack:0:raw_prerouting", []string{"iif=wlo1 ip4 saddr=192.0.2.0/24", "iif=enp2s0 ip4 saddr=198.51.100.0/24"}},
{"blrule from hosts zone", func(c *config.Config) {
c.Blrules = []config.BlruleRule{{Action: config.BlruleDrop, Source: "lan", Dest: "all"}}
}, "input", "blrule:0", []string{"iif=wlo1 ip4 saddr=192.0.2.0/24", "iif=enp2s0 ip4 saddr=198.51.100.0/24"}},
{"all expansion includes hosts zone", func(c *config.Config) {
c.Rules = []config.Rule{{Action: config.RuleAccept, Source: "all!net,vpn", Dest: "fw", Proto: "tcp"}}
}, "input", "rule:0", []string{"iif=wlo1 ip4 saddr=192.0.2.0/24", "iif=enp2s0 ip4 saddr=198.51.100.0/24"}},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
state := mustCompile(t, hostsCfg(tt.mod))
if got := describeTagged(state, tt.chain, tt.tag); !reflect.DeepEqual(got, tt.want) {
t.Errorf("%s %s = %q, want %q", tt.chain, tt.tag, got, tt.want)
}
})
}
}
func TestCompile_IntraZoneMultiInterface(t *testing.T) {
tests := []struct {
name string
policy []config.Policy
tag string
want []string
}{
{
name: "implicit accept between distinct interfaces",
policy: []config.Policy{{Source: "all", Dest: "all", Action: config.PolicyDrop}},
tag: "intra:lxd",
want: []string{"iif=lxdbr0 oif=docker0", "iif=lxdbr0 oif=br-", "iif=docker0 oif=lxdbr0", "iif=docker0 oif=br-", "iif=br- oif=lxdbr0", "iif=br- oif=docker0"},
},
{
name: "explicit zone policy overrides",
policy: []config.Policy{{Source: "lxd", Dest: "lxd", Action: config.PolicyDrop, Log: "info"}, {Source: "all", Dest: "all", Action: config.PolicyDrop}},
tag: "policy:0",
want: []string{"iif=lxdbr0 oif=docker0", "iif=lxdbr0 oif=br-", "iif=docker0 oif=lxdbr0", "iif=docker0 oif=br-", "iif=br- oif=lxdbr0", "iif=br- oif=docker0"},
},
}
for _, tc := range tests {
t.Run(tc.name, func(t *testing.T) {
cfg := &config.Config{
Settings: config.Settings{TableName: "test", AddressFamily: config.FamilyINET},
Zones: map[string]config.Zone{
"fw": {Type: config.ZoneFirewall},
"net": {Type: config.ZoneIP},
"lxd": {Type: config.ZoneIP},
},
Interfaces: []config.Interface{
{Zone: "net", Interface: "eth0"},
{Zone: "lxd", Interface: "lxdbr0"},
{Zone: "lxd", Interface: "docker0"},
{Zone: "lxd", Interface: "br-+"},
},
Policy: tc.policy,
PortGroups: map[string]config.PortGroup{},
}
state, err := NewCompiler(cfg).Compile()
if err != nil {
t.Fatalf("Compile() error: %v", err)
}
var got, all []string
for _, r := range state.Rules["forward"] {
if r.Tag == tc.tag {
got = append(got, describeRule(r))
}
if strings.HasPrefix(r.Tag, "intra:") {
all = append(all, r.Tag)
}
}
if !reflect.DeepEqual(got, tc.want) {
t.Errorf("%s rules = %q, want %q", tc.tag, got, tc.want)
}
if tc.tag != "intra:lxd" && len(all) != 0 {
t.Errorf("explicit policy must replace implicit accept, got %q", all)
}
last := taggedRules(state, "forward", tc.tag)
if len(last) == 0 {
return
}
if v, ok := last[0].Exprs[len(last[0].Exprs)-1].(*expr.Verdict); !ok || (tc.tag == "intra:lxd") != (v.Kind == expr.VerdictAccept) {
t.Errorf("%s verdict = %#v", tc.tag, last[0].Exprs[len(last[0].Exprs)-1])
}
})
}
}
func TestCompile_HostsAddressMatchesFamilyGuarded(t *testing.T) {
state := mustCompile(t, hostsCfg(func(c *config.Config) {
c.Hosts[0].Addresses = append(c.Hosts[0].Addresses, "2001:db8::/64")
c.Hosts[0].Exclusions = []string{"192.0.2.9"}
}))
want := []string{
"iif=wlo1 ip4 !saddr=192.0.2.0/24", "iif=wlo1 ip6 !saddr=2001:db8::/64",
"iif=wlo1 ip4 saddr=192.0.2.9",
"iif=enp2s0 ip4 !saddr=198.51.100.0/24", "iif=enp2s0 ip6",
}
if got := describeTagged(state, "input", "policy:0"); !reflect.DeepEqual(got, want) {
t.Errorf("net->fw DROP policy = %q, want %q", got, want)
}
want = []string{
"iif=wlo1 ip4 saddr=192.0.2.0/24 !saddr=192.0.2.9", "iif=wlo1 ip6 saddr=2001:db8::/64",
"iif=enp2s0 ip4 saddr=198.51.100.0/24",
}
if got := describeTagged(state, "input", "policy:1"); !reflect.DeepEqual(got, want) {
t.Errorf("lan->fw policy = %q, want %q", got, want)
}
for chain, rules := range state.Rules {
for _, r := range rules {
var fam byte
for i, e := range r.Exprs {
if m, ok := e.(*expr.Meta); ok && m.Key == expr.MetaKeyNFPROTO {
fam = r.Exprs[i+1].(*expr.Cmp).Data[0]
}
p, ok := e.(*expr.Payload)
if !ok || p.Base != expr.PayloadBaseNetworkHeader {
continue
}
if p.Len == 4 && fam != unix.NFPROTO_IPV4 || p.Len == 16 && fam != unix.NFPROTO_IPV6 {
t.Errorf("%s %s: %d-byte address compare without its family guard: %s", chain, r.Tag, p.Len, describeRule(r))
}
}
}
}
}
func TestCompile_FirewallSelfPolicySkipped(t *testing.T) {
for _, action := range []config.PolicyAction{config.PolicyAccept, config.PolicyDrop} {
t.Run(string(action), func(t *testing.T) {
state := mustCompile(t, listCfg(func(c *config.Config) {
c.Policy = []config.Policy{
{Source: "fw", Dest: "fw", Action: action},
{Source: "net", Dest: "fw", Action: config.PolicyDrop, Log: "info"},
}
}))
for _, chain := range []string{"input", "output", "forward"} {
if got := taggedRules(state, chain, "policy:0"); len(got) != 0 {
t.Errorf("fw->fw emitted %d rules in %s", len(got), chain)
}
}
got := taggedRules(state, "input", "policy:1")
if len(got) != 1 || describeRule(got[0]) != "iif=eth0" {
t.Errorf("net->fw input rules = %d, want one scoped to eth0", len(got))
}
})
}
}
func TestCompile_DestPlusOverridesIntraZone(t *testing.T) {
state := mustCompile(t, listCfg(func(c *config.Config) {
c.Zones["lxd"] = config.Zone{Type: config.ZoneIP}
c.Interfaces = append(c.Interfaces, config.Interface{Zone: "lxd", Interface: "lxdbr0"}, config.Interface{Zone: "lxd", Interface: "docker0"})
c.Policy = []config.Policy{{Source: "lxd", Dest: "all+", Action: config.PolicyDrop}}
}))
if got := taggedRules(state, "forward", "intra:lxd"); len(got) != 0 {
t.Errorf("lxd all+ must override implicit intra-zone accept, got %d rules", len(got))
}
if got := taggedRules(state, "forward", "policy:0"); len(got) == 0 {
t.Error("lxd all+ emitted no forward rules")
}
}
func TestCompile_HostsIntraZone(t *testing.T) {
state := mustCompile(t, hostsCfg(nil))
want := []string{
"iif=wlo1 ip4 saddr=192.0.2.0/24 oif=enp2s0 ip4 daddr=198.51.100.0/24",
"iif=enp2s0 ip4 saddr=198.51.100.0/24 oif=wlo1 ip4 daddr=192.0.2.0/24",
}
if got := describeTagged(state, "forward", "intra:lan"); !reflect.DeepEqual(got, want) {
t.Errorf("lan intra = %q, want %q", got, want)
}
want = []string{
"iif=wlo1 ip4 !saddr=192.0.2.0/24 oif=enp2s0 ip4 !daddr=198.51.100.0/24",
"iif=wlo1 ip6 oif=enp2s0 ip6",
"iif=enp2s0 ip4 !saddr=198.51.100.0/24 oif=wlo1 ip4 !daddr=192.0.2.0/24",
"iif=enp2s0 ip6 oif=wlo1 ip6",
}
if got := describeTagged(state, "forward", "intra:net"); !reflect.DeepEqual(got, want) {
t.Errorf("net intra = %q, want %q", got, want)
}
state = mustCompile(t, hostsCfg(func(c *config.Config) {
c.Policy = append([]config.Policy{{Source: "lan", Dest: "lan", Action: config.PolicyDrop}}, c.Policy...)
}))
if got := describeTagged(state, "forward", "intra:lan"); len(got) != 0 {
t.Errorf("explicit lan lan policy must replace implicit accept, got %q", got)
}
want = []string{
"iif=wlo1 ip4 saddr=192.0.2.0/24 oif=enp2s0 ip4 daddr=198.51.100.0/24",
"iif=enp2s0 ip4 saddr=198.51.100.0/24 oif=wlo1 ip4 daddr=192.0.2.0/24",
}
if got := describeTagged(state, "forward", "policy:0"); !reflect.DeepEqual(got, want) {
t.Errorf("lan lan policy = %q, want %q", got, want)
}
}
func TestCompile_HostsRouteBack(t *testing.T) {
sameIface := func(routeback bool) func(*config.Config) {
return func(c *config.Config) {
c.Hosts = []config.Host{
{Zone: "lan", Interface: "wlo1", Addresses: []string{"192.0.2.0/24"}, Options: config.HostOptions{RouteBack: routeback}},
{Zone: "lan", Interface: "wlo1", Addresses: []string{"198.51.100.0/24"}},
}
}
}
if got := describeTagged(mustCompile(t, hostsCfg(sameIface(false))), "forward", "intra:lan"); len(got) != 0 {
t.Errorf("same-interface hosts without routeback = %q, want none", got)
}
want := []string{"iif=wlo1 ip4 saddr=192.0.2.0/24 oif=wlo1 ip4 daddr=192.0.2.0/24"}
if got := describeTagged(mustCompile(t, hostsCfg(sameIface(true))), "forward", "intra:lan"); !reflect.DeepEqual(got, want) {
t.Errorf("routeback hosts intra = %q, want %q", got, want)
}
state := mustCompile(t, hostsCfg(func(c *config.Config) {
sameIface(true)(c)
c.Policy = append([]config.Policy{{Source: "lan", Dest: "lan", Action: config.PolicyDrop}}, c.Policy...)
}))
if got := describeTagged(state, "forward", "policy:0"); !reflect.DeepEqual(got, want) {
t.Errorf("routeback hosts lan lan policy = %q, want %q", got, want)
}
}
+16 -8
View File
@@ -248,6 +248,9 @@ func convertInterfaces(dir string, cfg *config.Config, params map[string]string)
for _, row := range rows {
zone := subst(field(row, 0), params)
iface := subst(field(row, 1), params)
if isDash(zone) {
zone = ""
}
intf := config.Interface{
Zone: zone,
@@ -392,12 +395,14 @@ func convertHosts(dir string, cfg *config.Config, params map[string]string) erro
zone := subst(field(row, 0), params)
hostDef := subst(field(row, 1), params)
hostDef, excl, _ := strings.Cut(hostDef, "!")
iface, addrs := splitHostDef(hostDef)
host := config.Host{
Zone: zone,
Interface: iface,
Addresses: addrs,
Zone: zone,
Interface: iface,
Addresses: addrs,
Exclusions: splitAddrList(excl),
}
optsStr := subst(field(row, 2), params)
@@ -415,16 +420,19 @@ func splitHostDef(s string) (string, []string) {
if idx < 0 {
return s, nil
}
iface := s[:idx]
addrPart := s[idx+1:]
return s[:idx], splitAddrList(s[idx+1:])
}
// splitAddrList splits a comma address list, unwrapping shorewall6 [addr]/len brackets.
func splitAddrList(s string) []string {
var addrs []string
for _, a := range strings.Split(addrPart, ",") {
a = strings.TrimSpace(a)
for _, a := range strings.Split(s, ",") {
a = strings.NewReplacer("[", "", "]", "").Replace(strings.TrimSpace(a))
if a != "" {
addrs = append(addrs, a)
}
}
return iface, addrs
return addrs
}
func parseHostOptions(s string) config.HostOptions {
+28
View File
@@ -3,6 +3,7 @@ package shorewall
import (
"os"
"path/filepath"
"reflect"
"testing"
"git.unkin.net/unkin/tomswall/internal/config"
@@ -776,3 +777,30 @@ func TestConvert_LogLimit(t *testing.T) {
}
}
}
func TestConvert_HostsExclusions(t *testing.T) {
dir := minimalShorewallDir(t)
writeFile(t, dir, "interfaces", `
net eth0
- eth1
`)
writeFile(t, dir, "hosts", `
loc eth0:192.0.2.0/24,198.51.100.0/24!192.0.2.7,192.0.2.8 routeback
loc eth1:[2001:db8::]/64
`)
cfg, err := Convert(dir)
if err != nil {
t.Fatalf("Convert: %v", err)
}
want := []config.Host{
{Zone: "loc", Interface: "eth0", Addresses: []string{"192.0.2.0/24", "198.51.100.0/24"},
Exclusions: []string{"192.0.2.7", "192.0.2.8"}, Options: config.HostOptions{RouteBack: true}},
{Zone: "loc", Interface: "eth1", Addresses: []string{"2001:db8::/64"}},
}
if !reflect.DeepEqual(cfg.Hosts, want) {
t.Errorf("hosts = %+v, want %+v", cfg.Hosts, want)
}
if err := cfg.Validate(); err != nil {
t.Errorf("Validate: %v", err)
}
}
+41 -24
View File
@@ -25,15 +25,15 @@ const Unit = "tomswall-try-revert"
var (
// Dir holds the lock and the pending snapshot.
Dir = "/var/lib/tomswall"
// run executes a systemd command; replaced in tests.
run = func(name string, args ...string) error {
// Run executes a systemd command. Test hook; production code must not reassign.
Run = func(name string, args ...string) error {
if out, err := exec.Command(name, args...).CombinedOutput(); err != nil {
return fmt.Errorf("%s %s: %w: %s", name, strings.Join(args, " "), err, strings.TrimSpace(string(out)))
}
return nil
}
// restore rolls the live table back to a snapshot; replaced in tests.
restore = func(s *nftables.Snapshot) error {
// Restore rolls the live table back to a snapshot. Test hook; production code must not reassign.
Restore = func(s *nftables.Snapshot) error {
engine, err := nftables.NewEngine(&config.Config{Settings: config.Settings{TableName: s.Table}})
if err != nil {
return err
@@ -94,23 +94,7 @@ func Arm(snap *nftables.Snapshot, pid int, delay time.Duration) (string, error)
if err != nil {
return "", err
}
f, err := os.CreateTemp(Dir, ".try-snapshot-*")
if err != nil {
return "", err
}
defer os.Remove(f.Name())
if _, err := f.Write(b); err != nil {
f.Close()
return "", err
}
if err := f.Sync(); err != nil {
f.Close()
return "", err
}
if err := f.Close(); err != nil {
return "", err
}
if err := os.Rename(f.Name(), snapshotPath()); err != nil {
if err := WriteFile(snapshotPath(), b); err != nil {
return "", err
}
@@ -119,13 +103,46 @@ func Arm(snap *nftables.Snapshot, pid int, delay time.Duration) (string, error)
return "", discardWith(err)
}
_ = disarm() // a leftover timer from an earlier try would block the unit name
if err := run("systemd-run", "--quiet", "--collect", "--unit", Unit,
if err := Run("systemd-run", "--quiet", "--collect", "--unit", Unit,
fmt.Sprintf("--on-active=%ds", int(delay.Round(time.Second).Seconds())), exe, "revert", "--id", id); err != nil {
return "", discardWith(fmt.Errorf("arming revert timer: %w", err))
}
return id, nil
}
// WriteFile durably replaces path with b: temp file, fsync, rename, fsync the directory.
func WriteFile(path string, b []byte) error {
dir := filepath.Dir(path)
if err := os.MkdirAll(dir, 0o755); err != nil {
return err
}
f, err := os.CreateTemp(dir, "."+filepath.Base(path)+"-*")
if err != nil {
return err
}
defer os.Remove(f.Name())
if _, err := f.Write(b); err != nil {
f.Close()
return err
}
if err := f.Sync(); err != nil {
f.Close()
return err
}
if err := f.Close(); err != nil {
return err
}
if err := os.Rename(f.Name(), path); err != nil {
return err
}
d, err := os.Open(dir)
if err != nil {
return err
}
defer d.Close()
return d.Sync()
}
// Discard drops the pending snapshot and timer without restoring. The caller must hold the lock.
func Discard() error {
_ = disarm()
@@ -143,7 +160,7 @@ func discardWith(err error) error {
}
func disarm() error {
return run("systemctl", "stop", Unit+".timer")
return Run("systemctl", "stop", Unit+".timer")
}
// Confirm keeps the tried ruleset. ok is false when no try was pending, i.e.
@@ -193,7 +210,7 @@ func Abort() error {
}
func restorePending(p *pending) error {
if err := restore(p.Snapshot); err != nil {
if err := Restore(p.Snapshot); err != nil {
return fmt.Errorf("restoring snapshot: %w", err)
}
return Discard()
+7 -7
View File
@@ -15,12 +15,12 @@ func setup(t *testing.T) *[]string {
t.Helper()
Dir = t.TempDir()
var cmds []string
orig := run
run = func(name string, args ...string) error {
orig := Run
Run = func(name string, args ...string) error {
cmds = append(cmds, name+" "+strings.Join(args, " "))
return nil
}
t.Cleanup(func() { run = orig })
t.Cleanup(func() { Run = orig })
return &cmds
}
@@ -91,7 +91,7 @@ func TestAcquireRefusesWhilePending(t *testing.T) {
func TestArmFailureDiscardsSnapshot(t *testing.T) {
setup(t)
run = func(name string, args ...string) error {
Run = func(name string, args ...string) error {
if name == "systemd-run" {
return errors.New("no systemd")
}
@@ -145,12 +145,12 @@ func TestConfirmAfterRevertFails(t *testing.T) {
func stubRestore(t *testing.T, err error) *[]*nftables.Snapshot {
t.Helper()
var got []*nftables.Snapshot
orig := restore
restore = func(s *nftables.Snapshot) error {
orig := Restore
Restore = func(s *nftables.Snapshot) error {
got = append(got, s)
return err
}
t.Cleanup(func() { restore = orig })
t.Cleanup(func() { Restore = orig })
return &got
}
+11
View File
@@ -47,6 +47,17 @@ contents:
file_info:
mode: 0640
# systemd unit + environment file for applying a local config at boot.
- src: packaging/tomswall.service
dst: /usr/lib/systemd/system/tomswall.service
file_info:
mode: 0644
- src: packaging/tomswall.env
dst: /etc/tomswall/tomswall.env
type: config|noreplace
file_info:
mode: 0644
# Shell completions (generated by scripts/build-rpm.sh before packaging).
- src: dist/completions/tomswall.bash
dst: /usr/share/bash-completion/completions/tomswall
+1
View File
@@ -3,6 +3,7 @@ Description=tomswall control-plane agent (pull and apply firewall config)
Documentation=https://git.unkin.net/unkin/tomswall
After=network-online.target
Wants=network-online.target
Conflicts=tomswall.service
[Service]
Type=simple
+3
View File
@@ -0,0 +1,3 @@
# Config applied by tomswall.service: a tomswall YAML file or a shorewall directory.
TOMSWALL_CONFIG=/etc/tomswall/tomswall.yaml
#TOMSWALL_CONFIG=/etc/shorewall
+25
View File
@@ -0,0 +1,25 @@
[Unit]
Description=tomswall firewall (apply local config at boot)
Documentation=https://git.unkin.net/unkin/tomswall
DefaultDependencies=no
Wants=network-pre.target
Before=network-pre.target shutdown.target
After=local-fs.target systemd-sysctl.service
Conflicts=shutdown.target tomswall-agent.service
StartLimitIntervalSec=60
StartLimitBurst=5
[Service]
Type=oneshot
RemainAfterExit=yes
Environment=TOMSWALL_CONFIG=/etc/tomswall/tomswall.yaml
EnvironmentFile=-/etc/tomswall/tomswall.env
ExecStart=/usr/sbin/tomswall apply -c ${TOMSWALL_CONFIG}
ExecReload=/usr/sbin/tomswall apply -c ${TOMSWALL_CONFIG}
# Fails open: after StartLimitBurst failures within StartLimitIntervalSec, boot continues without the ruleset.
Restart=on-failure
RestartSec=5
# No ExecStop: stopping the unit leaves the ruleset in place (flush would open the firewall).
[Install]
WantedBy=sysinit.target