From c234699ae75bdfdfbbff192e7419540c04d26476 Mon Sep 17 00:00:00 2001 From: Rogee Date: Mon, 31 Aug 2026 12:07:15 +0800 Subject: [PATCH] HH-834: tighten gateway lifecycle fences (#28) --- cmd/docker-gateway/main.go | 16 ++++--- cmd/docker-gateway/main_test.go | 56 +++++++++++++++++++++++- cmd/docker-gateway/proxy.go | 29 +++++++++++-- cmd/docker-gateway/proxy_test.go | 60 ++++++++++++++++++++++++++ docs/architecture/container-control.md | 3 +- 5 files changed, 153 insertions(+), 11 deletions(-) diff --git a/cmd/docker-gateway/main.go b/cmd/docker-gateway/main.go index 2e32dea..2e1da61 100644 --- a/cmd/docker-gateway/main.go +++ b/cmd/docker-gateway/main.go @@ -594,6 +594,9 @@ func (api gateway) changeState(c fiber.Ctx) error { if err != nil { return writeError(c, http.StatusBadRequest, err) } + if input.NetworkID == "" { + return writeError(c, http.StatusBadRequest, errors.New("network_id must identify the expected network generation")) + } _, exists, err := api.requireGeneration(id, input) if err != nil { return writeError(c, statusFor(err), err) @@ -798,7 +801,7 @@ func (api gateway) restoreProxy(c fiber.Ctx) error { removeStaleProxy() return writeError(c, http.StatusBadGateway, errors.New("restore isolated browser network")) } - if err = api.requireProxyNetworkGeneration(c.Params("id"), input, networkGeneration); err != nil { + if err = api.requireProxyNetworkGeneration(c.Params("id"), input, networkGeneration, bindHost); err != nil { removeStaleProxy() return writeError(c, statusFor(err), err) } @@ -807,7 +810,7 @@ func (api gateway) restoreProxy(c fiber.Ctx) error { removeStaleProxy() return writeError(c, statusFor(err), errors.Join(errors.New("restore in-memory proxy"), err)) } - if err = api.requireProxyNetworkGeneration(c.Params("id"), input, networkGeneration); err != nil { + if err = api.requireProxyNetworkGeneration(c.Params("id"), input, networkGeneration, bindHost); err != nil { undoProxy() return writeError(c, statusFor(err), err) } @@ -815,7 +818,7 @@ func (api gateway) restoreProxy(c fiber.Ctx) error { undoProxy() return writeError(c, http.StatusConflict, errGenerationConflict) } - if err = api.requireProxyNetworkGeneration(c.Params("id"), input, networkGeneration); err != nil { + if err = api.requireProxyNetworkGeneration(c.Params("id"), input, networkGeneration, bindHost); err != nil { undoProxy() return writeError(c, statusFor(err), err) } @@ -823,7 +826,7 @@ func (api gateway) restoreProxy(c fiber.Ctx) error { return c.SendStatus(http.StatusNoContent) } -func (api gateway) requireProxyNetworkGeneration(alias string, input proxyRestoreRequest, expected tenantNetworkGeneration) error { +func (api gateway) requireProxyNetworkGeneration(alias string, input proxyRestoreRequest, expected tenantNetworkGeneration, bindHost string) error { _, labels, err := api.requireProxyGeneration(alias, input) if err != nil { return err @@ -831,9 +834,10 @@ func (api gateway) requireProxyNetworkGeneration(alias string, input proxyRestor if networkID := labels[networkIDLabel]; networkID != "" && networkID != expected.ID { return errGenerationConflict } - current, _, exists, err := api.docker.inspectTenantNetwork(api.network, alias, input.BindingVersion, + current, addresses, exists, err := api.docker.inspectTenantNetwork(api.network, alias, input.BindingVersion, input.RuntimeID, api.self, expected.ID, false) - if err != nil || !exists || !sameTenantNetworkMembers(current, expected) { + host, _, _ := net.ParseCIDR(addresses[current.SelfMember]) + if err != nil || !exists || !sameTenantNetworkMembers(current, expected) || host == nil || host.String() != bindHost { return errGenerationConflict } return nil diff --git a/cmd/docker-gateway/main_test.go b/cmd/docker-gateway/main_test.go index dc8b703..5a137b4 100644 --- a/cmd/docker-gateway/main_test.go +++ b/cmd/docker-gateway/main_test.go @@ -1028,6 +1028,39 @@ func TestGatewayRestoreFinalFenceRemovesStaleProxy(t *testing.T) { } } +func TestGatewayProxyFenceRejectsGatewayAddressChange(t *testing.T) { + labels := map[string]string{managedLabel: "true", idLabel: "account-a", bindingVersionLabel: "1", + networkExitLabel: "exit-1", networkIDLabel: "network-n1"} + server := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + switch request.URL.Path { + case "/containers/" + namePrefix + "account-a/json": + _ = json.NewEncoder(response).Encode(map[string]any{"Id": "container-c1", "Config": map[string]any{"Labels": labels}}) + case "/containers/gateway-self/json": + _, _ = response.Write([]byte(`{"Id":"gateway-self","Config":{"Labels":{"` + gatewayMemberLabel + `":"true"}}}`)) + case "/networks/network-n1": + _ = json.NewEncoder(response).Encode(map[string]any{ + "Id": "network-n1", "Name": "creatorhub_browser-account-a", "Driver": "bridge", "Internal": false, "Attachable": false, "Ingress": false, + "Labels": map[string]string{managedLabel: "true", networkRoleLabel: browserNetworkRole, idLabel: "account-a", bindingVersionLabel: "1"}, + "Containers": map[string]any{ + "container-c1": map[string]string{"Name": namePrefix + "account-a", "IPv4Address": "127.0.0.2/8"}, + "gateway-self": map[string]string{"Name": "gateway-self", "IPv4Address": "127.0.0.4/8"}, + }, + }) + default: + t.Fatalf("unexpected Docker request %s %s", request.Method, request.URL.Path) + } + })) + defer server.Close() + docker := dockerClient{baseURL: server.URL, client: server.Client(), slow: server.Client()} + api := gateway{docker: docker, network: "creatorhub_browser", self: "gateway-self"} + input := proxyRestoreRequest{BindingVersion: 1, RuntimeID: "container-c1", NetworkID: "network-n1", NetworkExitID: "exit-1"} + expected := tenantNetworkGeneration{ID: "network-n1", Name: "creatorhub_browser-account-a", RuntimeAttached: true, + GatewayMembers: []string{"gateway-self"}, SelfMember: "gateway-self"} + if err := api.requireProxyNetworkGeneration("account-a", input, expected, "127.0.0.3"); err != errGenerationConflict { + t.Fatalf("gateway address change crossed proxy fence: %v", err) + } +} + func TestGatewayLifecycleUsesInspectedImmutableContainerID(t *testing.T) { tests := []struct { method string @@ -1072,6 +1105,23 @@ func TestGatewayLifecycleUsesInspectedImmutableContainerID(t *testing.T) { } } +func TestGatewayStartRejectsEmptyNetworkGeneration(t *testing.T) { + mutations := 0 + docker, server := testDocker(func(response http.ResponseWriter, request *http.Request) { + if request.Method != http.MethodGet { + mutations++ + } + _, _ = response.Write([]byte(`{"Id":"container-id","Config":{"Labels":{"` + managedLabel + `":"true","` + idLabel + `":"account-a","` + bindingVersionLabel + `":"1"}}}`)) + }) + defer server.Close() + response := httptest.NewRecorder() + adaptor.FiberApp(newGateway(docker, "creatorhub_browser", testToken)).ServeHTTP(response, + authed(http.MethodPost, "/v1/browsers/account-a/start", strings.NewReader(`{"binding_version":1,"runtime_id":"container-id"}`))) + if response.Code != http.StatusBadRequest || mutations != 0 { + t.Fatalf("empty network generation reached Docker: status=%d mutations=%d body=%s", response.Code, mutations, response.Body.String()) + } +} + func TestGatewayDeleteUsesImmutableNetworkIDAcrossCleanupRetry(t *testing.T) { containerExists, cleanupFails, containerDeletes := true, true, 0 networkExists := true @@ -2212,7 +2262,11 @@ func TestGatewayRejectsStaleGenerationBeforeDockerMutation(t *testing.T) { defer server.Close() handler := newGateway(docker, "creatorhub_browser", testToken) response := httptest.NewRecorder() - adaptor.FiberApp(handler).ServeHTTP(response, authed(request.method, request.path, strings.NewReader(testGenerationBody))) + body := testGenerationBody + if strings.HasSuffix(request.path, "/start") { + body = `{"binding_version":1,"runtime_id":"container-id","network_id":"network-id"}` + } + adaptor.FiberApp(handler).ServeHTTP(response, authed(request.method, request.path, strings.NewReader(body))) if response.Code != http.StatusConflict || mutations != 0 { t.Fatalf("stale generation reached Docker mutation: status=%d mutations=%d body=%s", response.Code, mutations, response.Body.String()) } diff --git a/cmd/docker-gateway/proxy.go b/cmd/docker-gateway/proxy.go index bc7f649..ae92734 100644 --- a/cmd/docker-gateway/proxy.go +++ b/cmd/docker-gateway/proxy.go @@ -33,6 +33,7 @@ type memoryProxy struct { bindHost string listener net.Listener server *http.Server + tunnels map[net.Conn]net.Conn url string } @@ -65,7 +66,7 @@ func (registry *memoryProxyRegistry) configure(alias string, bindingVersion int6 } actualPort := listener.Addr().(*net.TCPAddr).Port proxy := &memoryProxy{exit: exit, bindingVersion: bindingVersion, networkID: networkID, bindHost: bindHost, listener: listener, - url: "http://" + net.JoinHostPort(browserProxyHost, strconv.Itoa(actualPort))} + tunnels: make(map[net.Conn]net.Conn), url: "http://" + net.JoinHostPort(browserProxyHost, strconv.Itoa(actualPort))} proxy.server = &http.Server{Handler: proxy, ReadHeaderTimeout: 10 * time.Second, IdleTimeout: 60 * time.Second} registry.proxies[alias] = proxy go func() { _ = proxy.server.Serve(listener) }() @@ -89,8 +90,16 @@ func (registry *memoryProxyRegistry) removeObject(alias string, proxy *memoryPro } func closeMemoryProxy(proxy *memoryProxy) { + proxy.mu.Lock() + tunnels := proxy.tunnels + proxy.tunnels = nil + proxy.mu.Unlock() _ = proxy.listener.Close() _ = proxy.server.Close() + for client, upstream := range tunnels { + _ = client.Close() + _ = upstream.Close() + } } func (registry *memoryProxyRegistry) ready(alias string, port int, runtimeID string, networkIDs ...string) bool { @@ -188,12 +197,26 @@ func (proxy *memoryProxy) tunnel(response http.ResponseWriter, request *http.Req _ = upstream.Close() return } + proxy.mu.Lock() + if proxy.tunnels == nil { + proxy.mu.Unlock() + _ = client.Close() + _ = upstream.Close() + return + } + proxy.tunnels[client] = upstream + proxy.mu.Unlock() + defer func() { + proxy.mu.Lock() + delete(proxy.tunnels, client) + proxy.mu.Unlock() + _ = client.Close() + _ = upstream.Close() + }() done := make(chan struct{}, 2) go func() { _, _ = io.Copy(upstream, client); done <- struct{}{} }() go func() { _, _ = io.Copy(client, upstream); done <- struct{}{} }() <-done - _ = client.Close() - _ = upstream.Close() } func (proxy *memoryProxy) dialContext(ctx context.Context, _, target string) (net.Conn, error) { diff --git a/cmd/docker-gateway/proxy_test.go b/cmd/docker-gateway/proxy_test.go index 163296a..5467e07 100644 --- a/cmd/docker-gateway/proxy_test.go +++ b/cmd/docker-gateway/proxy_test.go @@ -139,6 +139,66 @@ func TestMemoryProxyUsesAbsoluteFormForHTTPUpstream(t *testing.T) { } } +func TestMemoryProxyCleanupClosesHijackedTunnel(t *testing.T) { + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatal(err) + } + defer listener.Close() + upstreamClosed := make(chan error, 1) + go func() { + connection, err := listener.Accept() + if err != nil { + upstreamClosed <- err + return + } + defer connection.Close() + request, err := http.ReadRequest(bufio.NewReader(connection)) + if err == nil && request.Method == http.MethodConnect { + _, err = fmt.Fprint(connection, "HTTP/1.1 200 Connection Established\r\n\r\n") + } + if err == nil { + var data [1]byte + _, err = connection.Read(data[:]) + } + upstreamClosed <- err + }() + host, portText, _ := net.SplitHostPort(listener.Addr().String()) + port, _ := strconv.Atoi(portText) + registry := newMemoryProxyRegistry() + proxyURL, cleanup, err := registry.configure("account-a", 1, "127.0.0.1", 0, + gatewayProxyExit{Protocol: "http", Host: host, Port: port}) + if err != nil { + t.Fatal(err) + } + defer cleanup() + proxyAddress := strings.Replace(strings.TrimPrefix(proxyURL, "http://"), browserProxyHost, "127.0.0.1", 1) + client, err := net.Dial("tcp", proxyAddress) + if err != nil { + t.Fatal(err) + } + defer client.Close() + if _, err = fmt.Fprint(client, "CONNECT example.com:443 HTTP/1.1\r\nHost: example.com:443\r\n\r\n"); err != nil { + t.Fatal(err) + } + if response, err := http.ReadResponse(bufio.NewReader(client), &http.Request{Method: http.MethodConnect}); err != nil || response.StatusCode != http.StatusOK { + t.Fatalf("open CONNECT tunnel: response=%v err=%v", response, err) + } + cleanup() + _ = client.SetReadDeadline(time.Now().Add(time.Second)) + if _, err = client.Read(make([]byte, 1)); err == nil { + t.Fatal("proxy cleanup left the client tunnel open") + } + select { + case err = <-upstreamClosed: + if err == nil { + t.Fatal("proxy cleanup left the upstream tunnel open") + } + case <-time.After(time.Second): + t.Fatal("proxy cleanup did not close the upstream tunnel") + } +} + func TestMemoryProxyRejectsCrossAliasAddress(t *testing.T) { registry := newMemoryProxyRegistry() proxyURL, cleanup, err := registry.configure("account-a", 1, "127.0.0.1", 0, gatewayProxyExit{Protocol: "http", Host: "127.0.0.1", Port: 1}) diff --git a/docs/architecture/container-control.md b/docs/architecture/container-control.md index 5039e65..5a269e0 100644 --- a/docs/architecture/container-control.md +++ b/docs/architecture/container-control.md @@ -32,7 +32,8 @@ React ──> control-plane ── /api/browsers ──(Bearer token)──> doc - 只有 `docker-gateway` 挂载 socket,控制面和浏览器容器均不可见;网关加入 control 与 browser 网络,浏览器只拿到无凭据的内存转发代理地址,`/v1` 仍必须通过容器内不可见的网关令牌; - 网关只暴露面向领域的路由,不提供通用 Docker 代理;`/v1` 全部接口校验 `Authorization: Bearer `(常数时间比较),令牌由部署者在网关环境变量与平台注册表中保持一致; -- 网关直连的 `POST /v1/browsers/{alias}/start|stop` 仅供内部维护使用,必须提交并精确匹配容器标签中的 `{binding_version,runtime_id,network_id}`;generation 不匹配返回 `409`,控制面生命周期编排不依赖无 fence 的直连 start; +- 网关直连的 `POST /v1/browsers/{alias}/start|stop` 仅供内部维护使用,必须提交并精确匹配容器标签中的 `{binding_version,runtime_id,network_id}`;直连 start 拒绝空 `network_id`,generation 不匹配返回 `409`,控制面生命周期编排不依赖无 fence 的直连 start; +- 网关恢复内存代理时会同时 fence 隔离网络成员及网关成员 IPv4,地址变化返回 `409` 并关闭刚恢复的监听;代理移除或代际替换会立即关闭所有已 hijack 的 CONNECT 双向连接,任一端先关闭也会关闭隧道两端,不等待优雅 drain; - 网关固定命令、网络、挂载和资源限制;外部输入是受校验的别名,以及平台下发的镜像引用、启动参数和卷名——镜像引用来自平台维护的版本表,新增/变更由人工在页面审核启用,不再写死在代码中; - 启停和删除前必须同时匹配固定名称前缀及 `io.creatorhub.managed`、`io.creatorhub.runtime-id` 标签; - 动态容器使用只读根文件系统、非 root `1000:1000` 与固定镜像入口、全部 capability drop、`no-new-privileges`、CPU/内存/PID 限制,且无宿主机端口和目录挂载;