Skip to content

Commit decf68e

Browse files
authored
Merge branch 'main' into feat/execd-websocket-pty
2 parents 9ce847a + ebc764c commit decf68e

24 files changed

Lines changed: 657 additions & 118 deletions

File tree

‎components/egress/pkg/nftables/dynamic.go‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -24,9 +24,12 @@ import (
2424
const (
2525
dynAllowV4Set = "dyn_allow_v4"
2626
dynAllowV6Set = "dyn_allow_v6"
27-
dynSetTimeoutS = 300
27+
dynSetTimeoutS = 360
28+
// nftTTLSlackSec is added to the DNS TTL before clamping, so allow entries
29+
// slightly outlive the resolver cache and reduce races with short TTLs.
30+
nftTTLSlackSec = 60
2831
minTTLSec = 60
29-
maxTTLSec = 300
32+
maxTTLSec = 360 // max DNS TTL (300) + nftTTLSlackSec
3033
)
3134

3235
// ResolvedIP is a single IP learned from DNS with TTL for dynamic nft set.
@@ -58,7 +61,7 @@ func buildAddResolvedIPsScript(table string, ips []ResolvedIP) string {
5861
}
5962

6063
func clampTTL(d time.Duration) int {
61-
sec := int(d.Seconds())
64+
sec := int(d.Seconds()) + nftTTLSlackSec
6265
if sec < minTTLSec {
6366
return minTTLSec
6467
}
Lines changed: 112 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,112 @@
1+
// Copyright 2026 Alibaba Group Holding Ltd.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package nftables
16+
17+
import (
18+
"fmt"
19+
"net/netip"
20+
)
21+
22+
// normalizeNFTIntervalSet drops entries that are redundant for nftables sets with
23+
// `flags interval`: a host or smaller CIDR that falls strictly inside another
24+
// listed CIDR would make `nft add element` fail with "conflicting intervals specified".
25+
// Used on every ApplyStatic (startup and /policy updates).
26+
func normalizeNFTIntervalSet(elems []string) ([]string, error) {
27+
if len(elems) == 0 {
28+
return nil, nil
29+
}
30+
prefs := make([]netip.Prefix, 0, len(elems))
31+
for _, s := range elems {
32+
p, err := parseAsPrefix(s)
33+
if err != nil {
34+
return nil, err
35+
}
36+
prefs = append(prefs, p.Masked())
37+
}
38+
prefs = uniquePrefixes(prefs)
39+
prefs = removeStrictSubnets(prefs)
40+
out := make([]string, 0, len(prefs))
41+
for _, p := range prefs {
42+
out = append(out, formatPrefixForNFT(p))
43+
}
44+
return out, nil
45+
}
46+
47+
func parseAsPrefix(s string) (netip.Prefix, error) {
48+
if p, err := netip.ParsePrefix(s); err == nil {
49+
return p, nil
50+
}
51+
if a, err := netip.ParseAddr(s); err == nil {
52+
if a.Is4() {
53+
return a.Prefix(32)
54+
}
55+
return a.Prefix(128)
56+
}
57+
return netip.Prefix{}, fmt.Errorf("nftables: invalid IP or CIDR %q", s)
58+
}
59+
60+
func uniquePrefixes(prefs []netip.Prefix) []netip.Prefix {
61+
seen := make(map[string]struct{}, len(prefs))
62+
out := make([]netip.Prefix, 0, len(prefs))
63+
for _, p := range prefs {
64+
k := p.String()
65+
if _, ok := seen[k]; ok {
66+
continue
67+
}
68+
seen[k] = struct{}{}
69+
out = append(out, p)
70+
}
71+
return out
72+
}
73+
74+
// strictSupernet reports whether super covers all addresses of sub, with super
75+
// strictly larger (smaller prefix length) than sub.
76+
func strictSupernet(super, sub netip.Prefix) bool {
77+
if super.Addr().Is4() != sub.Addr().Is4() {
78+
return false
79+
}
80+
if super.Bits() >= sub.Bits() {
81+
return false
82+
}
83+
return super.Contains(sub.Addr())
84+
}
85+
86+
func removeStrictSubnets(prefs []netip.Prefix) []netip.Prefix {
87+
out := make([]netip.Prefix, 0, len(prefs))
88+
for _, p := range prefs {
89+
redundant := false
90+
for _, q := range prefs {
91+
if strictSupernet(q, p) {
92+
redundant = true
93+
break
94+
}
95+
}
96+
if !redundant {
97+
out = append(out, p)
98+
}
99+
}
100+
return out
101+
}
102+
103+
func formatPrefixForNFT(p netip.Prefix) string {
104+
p = p.Masked()
105+
if p.Addr().Is4() && p.Bits() == 32 {
106+
return p.Addr().String()
107+
}
108+
if p.Addr().Is6() && p.Bits() == 128 {
109+
return p.Addr().String()
110+
}
111+
return p.String()
112+
}
Lines changed: 58 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,58 @@
1+
// Copyright 2026 Alibaba Group Holding Ltd.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package nftables
16+
17+
import (
18+
"testing"
19+
20+
"github.com/stretchr/testify/require"
21+
)
22+
23+
func TestNormalizeNFTIntervalSet_cgnatContainsNameservers(t *testing.T) {
24+
got, err := normalizeNFTIntervalSet([]string{
25+
"100.64.0.0/10",
26+
"127.0.0.1",
27+
"100.100.2.136",
28+
"100.100.2.138",
29+
})
30+
require.NoError(t, err)
31+
require.Equal(t, []string{"100.64.0.0/10", "127.0.0.1"}, got)
32+
}
33+
34+
func TestNormalizeNFTIntervalSet_noChange(t *testing.T) {
35+
got, err := normalizeNFTIntervalSet([]string{"1.1.1.1", "2.2.0.0/16"})
36+
require.NoError(t, err)
37+
require.Equal(t, []string{"1.1.1.1", "2.2.0.0/16"}, got)
38+
}
39+
40+
func TestNormalizeNFTIntervalSet_dedupe(t *testing.T) {
41+
got, err := normalizeNFTIntervalSet([]string{"8.8.8.8", "8.8.8.8"})
42+
require.NoError(t, err)
43+
require.Equal(t, []string{"8.8.8.8"}, got)
44+
}
45+
46+
func TestNormalizeNFTIntervalSet_ipv6(t *testing.T) {
47+
got, err := normalizeNFTIntervalSet([]string{
48+
"2001:db8::/32",
49+
"2001:db8::1",
50+
})
51+
require.NoError(t, err)
52+
require.Equal(t, []string{"2001:db8::/32"}, got)
53+
}
54+
55+
func TestNormalizeNFTIntervalSet_invalid(t *testing.T) {
56+
_, err := normalizeNFTIntervalSet([]string{"not-an-ip"})
57+
require.Error(t, err)
58+
}

‎components/egress/pkg/nftables/manager.go‎

Lines changed: 43 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -82,6 +82,9 @@ func NewManagerWithOptions(opts Options) *Manager {
8282
// Uses the same mutex as AddResolvedIPs so a /policy update never overlaps a DNS
8383
// callback: without this, add-element could run while the table is being deleted/recreated
8484
// and fail, causing a transient deny for a client that already got an allowed DNS answer.
85+
//
86+
// On every call (startup and /policy updates), static allow/deny and DoH blocklist
87+
// interval sets are normalized so overlapping CIDR/host pairs do not make nft fail.
8588
func (m *Manager) ApplyStatic(ctx context.Context, p *policy.NetworkPolicy) error {
8689
if p == nil {
8790
p = policy.DefaultDenyPolicy()
@@ -91,7 +94,10 @@ func (m *Manager) ApplyStatic(ctx context.Context, p *policy.NetworkPolicy) erro
9194
p.DefaultAction, len(allowV4), len(allowV6), len(denyV4), len(denyV6))
9295
m.mu.Lock()
9396
defer m.mu.Unlock()
94-
script := buildRuleset(p, m.opts)
97+
script, err := buildRuleset(p, m.opts)
98+
if err != nil {
99+
return err
100+
}
95101
if _, err := m.run(ctx, script); err != nil {
96102
// On a fresh host the delete-table may fail; retry once without the delete line.
97103
if isMissingTableError(err) {
@@ -109,7 +115,8 @@ func (m *Manager) ApplyStatic(ctx context.Context, p *policy.NetworkPolicy) erro
109115
}
110116

111117
// AddResolvedIPs adds DNS-learned IPs to dynamic allow sets with TTL-based timeout.
112-
// TTL is clamped to minTTLSec–maxTTLSec. Call only when table exists (dns+nft mode).
118+
// Each element timeout is DNS TTL + 60s, then clamped to minTTLSec–maxTTLSec.
119+
// Call only when table exists (dns+nft mode).
113120
func (m *Manager) AddResolvedIPs(ctx context.Context, ips []ResolvedIP) error {
114121
if len(ips) == 0 {
115122
return nil
@@ -126,8 +133,33 @@ func (m *Manager) AddResolvedIPs(ctx context.Context, ips []ResolvedIP) error {
126133
return err
127134
}
128135

129-
func buildRuleset(p *policy.NetworkPolicy, opts Options) string {
136+
func buildRuleset(p *policy.NetworkPolicy, opts Options) (string, error) {
130137
allowV4, allowV6, denyV4, denyV6 := p.StaticIPSets()
138+
var err error
139+
if allowV4, err = normalizeNFTIntervalSet(allowV4); err != nil {
140+
return "", err
141+
}
142+
if allowV6, err = normalizeNFTIntervalSet(allowV6); err != nil {
143+
return "", err
144+
}
145+
if denyV4, err = normalizeNFTIntervalSet(denyV4); err != nil {
146+
return "", err
147+
}
148+
if denyV6, err = normalizeNFTIntervalSet(denyV6); err != nil {
149+
return "", err
150+
}
151+
dohBlockV4 := opts.DoHBlocklistV4
152+
dohBlockV6 := opts.DoHBlocklistV6
153+
if len(dohBlockV4) > 0 {
154+
if dohBlockV4, err = normalizeNFTIntervalSet(dohBlockV4); err != nil {
155+
return "", err
156+
}
157+
}
158+
if len(dohBlockV6) > 0 {
159+
if dohBlockV6, err = normalizeNFTIntervalSet(dohBlockV6); err != nil {
160+
return "", err
161+
}
162+
}
131163

132164
var b strings.Builder
133165
// Reset and re-create table, sets, and chain.
@@ -141,19 +173,19 @@ func buildRuleset(p *policy.NetworkPolicy, opts Options) string {
141173
fmt.Fprintf(&b, "add set inet %s %s { type ipv4_addr; timeout %ds; }\n", tableName, dynAllowV4Set, dynSetTimeoutS)
142174
fmt.Fprintf(&b, "add set inet %s %s { type ipv6_addr; timeout %ds; }\n", tableName, dynAllowV6Set, dynSetTimeoutS)
143175

144-
if len(opts.DoHBlocklistV4) > 0 {
176+
if len(dohBlockV4) > 0 {
145177
fmt.Fprintf(&b, "add set inet %s %s { type ipv4_addr; flags interval; }\n", tableName, dohBlockV4Set)
146178
}
147-
if len(opts.DoHBlocklistV6) > 0 {
179+
if len(dohBlockV6) > 0 {
148180
fmt.Fprintf(&b, "add set inet %s %s { type ipv6_addr; flags interval; }\n", tableName, dohBlockV6Set)
149181
}
150182

151183
writeElements(&b, allowV4Set, allowV4)
152184
writeElements(&b, denyV4Set, denyV4)
153185
writeElements(&b, allowV6Set, allowV6)
154186
writeElements(&b, denyV6Set, denyV6)
155-
writeElements(&b, dohBlockV4Set, opts.DoHBlocklistV4)
156-
writeElements(&b, dohBlockV6Set, opts.DoHBlocklistV6)
187+
writeElements(&b, dohBlockV4Set, dohBlockV4)
188+
writeElements(&b, dohBlockV6Set, dohBlockV6)
157189

158190
chainPolicy := "drop"
159191
if p.DefaultAction == policy.ActionAllow {
@@ -168,14 +200,14 @@ func buildRuleset(p *policy.NetworkPolicy, opts Options) string {
168200
fmt.Fprintf(&b, "add rule inet %s %s udp dport 853 drop\n", tableName, chainName)
169201
}
170202
if opts.BlockDoH443 {
171-
if len(opts.DoHBlocklistV4) == 0 && len(opts.DoHBlocklistV6) == 0 {
203+
if len(dohBlockV4) == 0 && len(dohBlockV6) == 0 {
172204
// strict: drop all 443 when enabled but no blocklist provided
173205
fmt.Fprintf(&b, "add rule inet %s %s tcp dport 443 drop\n", tableName, chainName)
174206
} else {
175-
if len(opts.DoHBlocklistV4) > 0 {
207+
if len(dohBlockV4) > 0 {
176208
fmt.Fprintf(&b, "add rule inet %s %s ip daddr @%s tcp dport 443 drop\n", tableName, chainName, dohBlockV4Set)
177209
}
178-
if len(opts.DoHBlocklistV6) > 0 {
210+
if len(dohBlockV6) > 0 {
179211
fmt.Fprintf(&b, "add rule inet %s %s ip6 daddr @%s tcp dport 443 drop\n", tableName, chainName, dohBlockV6Set)
180212
}
181213
}
@@ -190,7 +222,7 @@ func buildRuleset(p *policy.NetworkPolicy, opts Options) string {
190222
fmt.Fprintf(&b, "add rule inet %s %s counter drop\n", tableName, chainName)
191223
}
192224

193-
return b.String()
225+
return b.String(), nil
194226
}
195227

196228
func writeElements(b *strings.Builder, setName string, elems []string) {

‎components/egress/pkg/nftables/manager_test.go‎

Lines changed: 24 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -50,8 +50,8 @@ func TestApplyStatic_BuildsRuleset_DefaultDeny(t *testing.T) {
5050
expectContains(t, rendered, "add rule inet opensandbox egress oifname \"lo\" accept")
5151
expectContains(t, rendered, "add rule inet opensandbox egress tcp dport 853 drop")
5252
expectContains(t, rendered, "add rule inet opensandbox egress udp dport 853 drop")
53-
expectContains(t, rendered, "add set inet opensandbox dyn_allow_v4 { type ipv4_addr; timeout 300s; }")
54-
expectContains(t, rendered, "add set inet opensandbox dyn_allow_v6 { type ipv6_addr; timeout 300s; }")
53+
expectContains(t, rendered, "add set inet opensandbox dyn_allow_v4 { type ipv4_addr; timeout 360s; }")
54+
expectContains(t, rendered, "add set inet opensandbox dyn_allow_v6 { type ipv6_addr; timeout 360s; }")
5555
expectContains(t, rendered, "add element inet opensandbox allow_v4 { 1.1.1.1, 2.2.0.0/16 }")
5656
expectContains(t, rendered, "add element inet opensandbox deny_v6 { 2001:db8::/32 }")
5757
expectContains(t, rendered, "add rule inet opensandbox egress ip daddr @dyn_allow_v4 accept")
@@ -137,8 +137,8 @@ func TestAddResolvedIPs_BuildsDynamicElements(t *testing.T) {
137137
{Addr: netip.MustParseAddr("2001:db8::1"), TTL: 60 * time.Second},
138138
}
139139
require.NoError(t, m.AddResolvedIPs(context.Background(), ips), "AddResolvedIPs returned error")
140-
expectContains(t, rendered, "add element inet opensandbox dyn_allow_v4 { 1.1.1.1 timeout 120s }")
141-
expectContains(t, rendered, "add element inet opensandbox dyn_allow_v6 { 2001:db8::1 timeout 60s }")
140+
expectContains(t, rendered, "add element inet opensandbox dyn_allow_v4 { 1.1.1.1 timeout 180s }")
141+
expectContains(t, rendered, "add element inet opensandbox dyn_allow_v6 { 2001:db8::1 timeout 120s }")
142142
}
143143

144144
func TestAddResolvedIPs_ClampsTTL(t *testing.T) {
@@ -152,8 +152,8 @@ func TestAddResolvedIPs_ClampsTTL(t *testing.T) {
152152
{Addr: netip.MustParseAddr("10.0.0.2"), TTL: 9999 * time.Second},
153153
}
154154
require.NoError(t, m.AddResolvedIPs(context.Background(), ips), "AddResolvedIPs returned error")
155-
expectContains(t, rendered, "10.0.0.1 timeout 60s")
156-
expectContains(t, rendered, "10.0.0.2 timeout 300s")
155+
expectContains(t, rendered, "10.0.0.1 timeout 70s")
156+
expectContains(t, rendered, "10.0.0.2 timeout 360s")
157157
}
158158

159159
func TestAddResolvedIPs_EmptyNoOp(t *testing.T) {
@@ -164,3 +164,21 @@ func TestAddResolvedIPs_EmptyNoOp(t *testing.T) {
164164
require.NoError(t, m.AddResolvedIPs(context.Background(), nil), "AddResolvedIPs returned error")
165165
require.NoError(t, m.AddResolvedIPs(context.Background(), []ResolvedIP{}), "AddResolvedIPs returned error")
166166
}
167+
168+
func TestApplyStatic_NormalizesOverlappingAllow(t *testing.T) {
169+
var rendered string
170+
m := NewManagerWithRunner(func(_ context.Context, script string) ([]byte, error) {
171+
rendered = script
172+
return nil, nil
173+
})
174+
p, err := policy.ParsePolicy(`{
175+
"defaultAction":"deny",
176+
"egress":[
177+
{"action":"allow","target":"100.64.0.0/10"},
178+
{"action":"allow","target":"100.100.2.136"}
179+
]
180+
}`)
181+
require.NoError(t, err)
182+
require.NoError(t, m.ApplyStatic(context.Background(), p))
183+
expectContains(t, rendered, "add element inet opensandbox allow_v4 { 100.64.0.0/10 }")
184+
}

0 commit comments

Comments
 (0)