HH-834: tighten gateway lifecycle fences (#28)
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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())
|
||||
}
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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})
|
||||
|
||||
@@ -32,7 +32,8 @@ React ──> control-plane ── /api/browsers ──(Bearer token)──> doc
|
||||
|
||||
- 只有 `docker-gateway` 挂载 socket,控制面和浏览器容器均不可见;网关加入 control 与 browser 网络,浏览器只拿到无凭据的内存转发代理地址,`/v1` 仍必须通过容器内不可见的网关令牌;
|
||||
- 网关只暴露面向领域的路由,不提供通用 Docker 代理;`/v1` 全部接口校验 `Authorization: Bearer <GATEWAY_TOKEN>`(常数时间比较),令牌由部署者在网关环境变量与平台注册表中保持一致;
|
||||
- 网关直连的 `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 限制,且无宿主机端口和目录挂载;
|
||||
|
||||
Reference in New Issue
Block a user