From c7e77a6f847a54ae6aacd6c0378f3abc744c783f Mon Sep 17 00:00:00 2001 From: youfak Date: Sun, 2 Aug 2026 07:59:39 +0800 Subject: [PATCH] feat: probe SOCKS5 upstreams --- README.md | 13 +- docs/api/control-plane.md | 7 +- docs/configuration/reference.md | 2 +- docs/design/project-structure.md | 2 +- docs/development/implementation-plan.md | 4 +- docs/operations/runbook.md | 2 +- docs/requirements/completion-audit.md | 6 +- docs/requirements/traceability.md | 2 +- internal/checker/probe/probe.go | 16 +- internal/checker/probe/probe_test.go | 205 +++++++++++++++++++++++- internal/checker/probe/socks5.go | 180 +++++++++++++++++++++ 11 files changed, 418 insertions(+), 21 deletions(-) create mode 100644 internal/checker/probe/socks5.go diff --git a/README.md b/README.md index 07e149c..ce3c9e3 100644 --- a/README.md +++ b/README.md @@ -80,7 +80,7 @@ flowchart LR ## 当前完成度 -截至 **2026-07-31**,实施计划中可直接勾选的检查项为 **53 / 74(71.6%)**。详情见 +截至 **2026-08-02**,实施计划中可直接勾选的检查项为 **53 / 74(71.6%)**。详情见 [实施计划](docs/development/implementation-plan.md)和 [交付完成度审计](docs/requirements/completion-audit.md)。 @@ -88,10 +88,11 @@ flowchart LR 提取与限流、Controller 的 Admin/Distribution/Metrics 监听,以及 PostgreSQL 管理状态;WorkerControlPlane 的 Register、Snapshot ACK、Runtime 心跳接收和 Redis 会话栅栏,以及 Gateway Outcome 上报的有界队列、序列确认与重试; - Controller 的 Redis 共享 BASIC 检查任务、按上游的有界调度、Checker 探测和 + Controller 的 Redis 共享 BASIC 检查任务、按上游的有界调度、HTTP/HTTPS/SOCKS5 + Checker 探测和 Observation 状态归并。 -- **部分完成**:EGRESS 与 TARGET 的任务编排和配置建模、SOCKS5 探测,Docker - Compose/Kubernetes 运行时 mTLS Overlay。 +- **部分完成**:EGRESS 与 TARGET 的任务编排和配置建模,Docker Compose/Kubernetes + 运行时 mTLS Overlay。 - **待完成**:loadgen、故障演练和代表性集群压测。 检查项数量不等于生产就绪度。静态部署清单与 protobuf descriptor 验证也不代表 @@ -168,7 +169,7 @@ Checker 的参数也可通过 `PROXY_POOL_CONTROL_PLANE_ADDRESS`、 `PROXY_POOL_CHECKER_ID`、`PROXY_POOL_CHECKER_INSTANCE_ID` 与 `PROXY_POOL_CHECKER_MAX_IN_FLIGHT` 提供。它不会访问 Redis/PostgreSQL;生产 Controller 在启用控制面时装配 Redis 共享任务 broker,并按启用的 Upstream 调度 -HTTP/HTTPS BASIC 检查。调度监督器每轮读取已发布配置,因此 reload 后的上游启停、 +HTTP/HTTPS/SOCKS5 BASIC 检查。调度监督器每轮读取已发布配置,因此 reload 后的上游启停、 检查间隔、抖动、超时、重试次数和 `maxInFlight` 会在下一轮生效;新启用的上游无需 重启 Controller。EGRESS 与 TARGET 尚未进入生产调度,loadgen 命令也尚未实现。 @@ -203,7 +204,7 @@ HTTP/HTTPS BASIC 检查。调度监督器每轮读取已发布配置,因此 re ownership 索引,以及 Gateway 快照客户端。 - **P0 - Checker 健康链**:BASIC 的共享调度、实际探测、Observation reducer 和 `FETCHED -> AVAILABLE / SUSPECT / UNHEALTHY` 状态链已完成;继续补齐 EGRESS、 - TARGET 和 SOCKS5。 + TARGET。 - **P1 - Gateway 与 Routing**:Gateway 进程、快照凭据分发、五种 Routing 策略与 `onUnavailable` 已接入;动态容量调整和 Drain 闭环待完成。 - **P1 - 可观测与部署**:低基数业务指标、完整 Compose/Kubernetes 进程拓扑, diff --git a/docs/api/control-plane.md b/docs/api/control-plane.md index ceb8285..4ddbcfc 100644 --- a/docs/api/control-plane.md +++ b/docs/api/control-plane.md @@ -43,7 +43,7 @@ Redis 只保留每个当前会话的一条序列和摘要;原始 Outcome、代 Redis 或 PostgreSQL。Gateway 只将结果写入本地有界队列,队列满时丢弃样本,不等待控制面 或存储。Checker 的有界任务领取、任务租约归属校验和 Observation 上报已经由 gRPC 契约测试覆盖;Controller 已装配 Redis 共享 due-index、生产任务 broker 与独立 -Checker 的 HTTP/HTTPS BASIC 探测进程。EGRESS 和 TARGET 尚未进入生产调度。 +Checker 的 HTTP/HTTPS/SOCKS5 BASIC 探测进程。EGRESS 和 TARGET 尚未进入生产调度。 `100,000 QPS` 仍是未验证的设计目标。 `WatchSnapshots` 建立时校验当前 session;每次签发快照引用时也把 `session_id` @@ -183,10 +183,9 @@ Routing 决定 AVAILABLE、SUSPECT 或 UNHEALTHY,并更新 Redis 活动池, `proxy-checker` 使用固定大小 worker-pool 执行每个 pull 批次,任务数不超过该请求的 `max_in_flight`;每次尝试都受 `deadline` 和 `timeout` 的较小值约束,失败可在同一 -deadline 内最多执行到 `max_attempts`。BASIC 针对 HTTP/HTTPS Proxy 验证到 Proxy 的 +deadline 内最多执行到 `max_attempts`。BASIC 针对 HTTP/HTTPS/SOCKS5 Proxy 验证到 Proxy 的 请求/认证握手;EGRESS 与 TARGET 通过 Proxy 请求任务指定的 HTTP/HTTPS 目标并将 -非成功状态作为事实。EGRESS 的出口身份响应解析、SOCKS5 与 EGRESS/TARGET 生产调度 -仍待后续实现。 +非成功状态作为事实。EGRESS 的出口身份响应解析与 EGRESS/TARGET 生产调度仍待后续实现。 ## 7. 兼容与演进 diff --git a/docs/configuration/reference.md b/docs/configuration/reference.md index 6968dfe..7091555 100644 --- a/docs/configuration/reference.md +++ b/docs/configuration/reference.md @@ -202,7 +202,7 @@ Checker 同样使用独立的可拨号地址:`proxy-checker` 的 `-control-pla `PROXY_POOL_*` 环境变量提供。mTLS 模式下该命令读取 `checkerTLS`,明文 fixture 模式只接受回环 Controller 地址。Checker 只从 gRPC 领取任务并批量上报事实,不读取 Redis/PostgreSQL;Controller 在生产启动拓扑中装配 Redis 共享任务队列,当前调度 -HTTP/HTTPS BASIC 检查。调度监督器在每轮从已发布配置读取启用的上游;Admin reload +HTTP/HTTPS/SOCKS5 BASIC 检查。调度监督器在每轮从已发布配置读取启用的上游;Admin reload 发布后,上游启停和有效 `check` 策略会在下一轮生效,新启用的上游无需重启 Controller。 EGRESS 和 TARGET 的生产调度仍在后续实施范围。 diff --git a/docs/design/project-structure.md b/docs/design/project-structure.md index 0bcd0a5..3f6e30d 100644 --- a/docs/design/project-structure.md +++ b/docs/design/project-structure.md @@ -47,7 +47,7 @@ Controller 是首版模块化单体。Provider、Pool、Routing 和 Extraction ### proxy-checker Checker 只产生 Observation。它从认证 gRPC 流领取有界任务,用固定 worker-pool 在任务 -deadline 内执行 HTTP/HTTPS BASIC、EGRESS 和 TARGET 探测并微批上报;EGRESS 的 +deadline 内执行 HTTP/HTTPS/SOCKS5 BASIC、EGRESS 和 TARGET 探测并微批上报;EGRESS 的 任务 URL 仅在执行期使用,回传全局事实不包含该 URL。最终状态迁移仍由 Controller 的 确定性 reducer 完成,避免多个检查实例同时写 Proxy 状态。Checker 不访问 Redis 或 PostgreSQL。 diff --git a/docs/development/implementation-plan.md b/docs/development/implementation-plan.md index 3df1ec8..da85c10 100644 --- a/docs/development/implementation-plan.md +++ b/docs/development/implementation-plan.md @@ -283,10 +283,10 @@ Checker Observation 上报 RPC 已复用既有控制面监听接入 Controller 有界 pull,并通过通用任务 broker 契约完成能力协商、同 Checker 并发窗口、租约到期回收、 领取者加不可预测 lease token 的栅栏和完成后重放;任务凭据仅由认证流在执行期下发。 Redis 共享 due-index/租约持久化、每 Upstream 的跨副本 in-flight 限制和生产 broker 已完成; -`proxy-checker` 独立进程、固定大小 worker-pool、任务期重试/微批上报和 HTTP/HTTPS +`proxy-checker` 独立进程、固定大小 worker-pool、任务期重试/微批上报和 HTTP/HTTPS/SOCKS5 BASIC 探测器已完成并有测试。TARGET 探测器具备任务执行能力,但尚无生产任务调度; EGRESS 已具备任务 URL 传输、HTTP/HTTPS 探测和全局事实回传契约,但出口身份响应解析、 -SOCKS5、EGRESS/TARGET 多维任务索引及部署运行态仍未实现, +EGRESS/TARGET 多维任务索引及部署运行态仍未实现, 因此本任务保持未完成。 补充进度(2026-08-02):BASIC 调度已改为配置驱动监督器。它每轮读取已发布快照并复用 diff --git a/docs/operations/runbook.md b/docs/operations/runbook.md index 26636f3..95cc5bd 100644 --- a/docs/operations/runbook.md +++ b/docs/operations/runbook.md @@ -24,7 +24,7 @@ `cmd/proxy-controller` 已完成配置单次加载、PostgreSQL 迁移、Redis 活动池、 Distribution/Admin/Metrics 独立监听和有界停机装配。Provider 自动补池、分布式 配额、动态重载和 Admin 低基数统计已装配;Controller 已装配 Redis BASIC 任务 broker, -`proxy-checker` 可执行 HTTP/HTTPS BASIC 探测。EGRESS/TARGET 调度、loadgen 与完整 +`proxy-checker` 可执行 HTTP/HTTPS/SOCKS5 BASIC 探测。EGRESS/TARGET 调度、loadgen 与完整 mTLS 环境 Overlay 仍属于 `implementation-plan.md` 后续任务。 因此 Compose/Kubernetes 资产当前仍用于评审网络、资源、探针和依赖关系,不能 视为完整可运行拓扑。 diff --git a/docs/requirements/completion-audit.md b/docs/requirements/completion-audit.md index c130400..bb0b202 100644 --- a/docs/requirements/completion-audit.md +++ b/docs/requirements/completion-audit.md @@ -62,7 +62,7 @@ Outcome 已实现为 Gateway 本地有界队列、微批确认重试和 Controll 记录归并,不改写 Proxy 全局状态。Controller 公用 Reducer 已作为 Observation 的唯一状态 归并边界,Checker Observation RPC 已在同一控制面监听以独立 SPIFFE 身份接入;有界任务领取、 租约归属与任务期凭据传输已由通用契约和 gRPC 往返测试覆盖。Redis 共享 due-index、 -任务/租约持久化、按上游的 in-flight 限制以及独立 Checker 的 HTTP/HTTPS BASIC 执行 +任务/租约持久化、按上游的 in-flight 限制以及独立 Checker 的 HTTP/HTTPS/SOCKS5 BASIC 执行 进程已经闭环;EGRESS、TARGET 的多维任务索引与生产调度尚未实现。Snapshot 签发在 Redis 中原子匹配当前 `session_id`,重注册会清除旧引用,迟到旧 Stream 不会覆盖新 session。Controller 在最近成功下发的 Snapshot `valid_until` 到达时关闭流;Gateway 的公用 @@ -114,9 +114,9 @@ Controller/Gateway 入口,完整 mTLS 运行时拓扑仍只有静态验证。 6. Worker 基础网络快照流、Proxy/Gateway Routing/凭据 Snapshot payload、Gateway Snapshot 客户端和进程装配、同版本 Routing 编译/动态匹配、五种策略上游选择与 reject/wait/direct 已完成; Outcome 上报已完成基础观测链;Redis ownership drain/ACK/过期回收及按 Worker 的可下发索引已完成。 -7. BASIC Checker 调度与 HTTP/HTTPS 探测器、全局与 TARGET Profile 的 Memory/Redis +7. BASIC Checker 调度与 HTTP/HTTPS/SOCKS5 探测器、全局与 TARGET Profile 的 Memory/Redis 原子归并、Controller Reducer 和 Observation 上报 RPC 已完成;EGRESS、TARGET - 生产任务调度、SOCKS5 探测与 REMOVE 编排仍待实现。 + 生产任务调度与 REMOVE 编排仍待实现。 8. Admin/Distribution 细粒度授权和审计查询;Distribution 分布式限流已完成。 9. 真实 Compose/Kubernetes 集成、故障演练和代表性集群负载测试。 10. 将 reject/wait/direct 接入 Distribution 运行链,补齐 Sequential 持久化恢复、跨实例 CAS diff --git a/docs/requirements/traceability.md b/docs/requirements/traceability.md index 50a6556..f319497 100644 --- a/docs/requirements/traceability.md +++ b/docs/requirements/traceability.md @@ -80,7 +80,7 @@ | ID | 最终需求 | 来源 | 验证证据 | |---|---|---|---| | HEALTH-001 | 全局健康与 Routing/目标健康分离 | 221-270, 8679-8708 | `domain/health` 已将 BASIC/EGRESS 全局 Reducer 与 TARGET Profile Reducer 分离;TARGET 在 Memory 和 Redis 独立、随代理 TTL 归并,不改写 Proxy 全局状态;Routing 消费待实现 | -| HEALTH-002 | 健康调度有 jitter、maxInFlight 和分级频率 | 8679-8736 | 配置有效合并、URL 校验、稳定抖动/优先级 Planner、有界 Scheduler tick、Redis due-index、跨副本 in-flight 原子限制,以及 HTTP/HTTPS BASIC 生产执行器已完成;EGRESS/TARGET 多维任务索引与调度待实现 | +| HEALTH-002 | 健康调度有 jitter、maxInFlight 和分级频率 | 8679-8736 | 配置有效合并、URL 校验、稳定抖动/优先级 Planner、有界 Scheduler tick、Redis due-index、跨副本 in-flight 原子限制,以及 HTTP/HTTPS/SOCKS5 BASIC 生产执行器已完成;EGRESS/TARGET 多维任务索引与调度待实现 | | HEALTH-003 | 失败分级 SUSPECT -> UNHEALTHY -> REMOVE | 8679-8736 | Controller 公用 Reducer 已通过 Memory/Redis 活动池原子提交全局连续失败、精确重放和成功恢复;BASIC 任务调度已完成,REMOVE 编排待实现 | | SEC-001 | API 认证与 Proxy 认证分离,Secret 统一脱敏 | 7528-8111, 8904-8945 | Config 脱敏、Provider Store -> SecretRef -> Gateway Resolver 跨包测试与格式化泄漏回归测试 | | SEC-002 | 非回环监听无保护时严格模式启动失败 | 8112-8441 | 配置校验测试 | diff --git a/internal/checker/probe/probe.go b/internal/checker/probe/probe.go index 0d4c254..1afefd7 100644 --- a/internal/checker/probe/probe.go +++ b/internal/checker/probe/probe.go @@ -59,12 +59,16 @@ func (executor *Executor) Execute(ctx context.Context, task *controlplanev1.Chec probeContext, cancel := context.WithDeadline(ctx, deadline) defer cancel() transport := &http.Transport{ - Proxy: http.ProxyURL(proxyURL), DialContext: (&net.Dialer{}).DialContext, ForceAttemptHTTP2: false, TLSHandshakeTimeout: minDuration(time.Until(deadline), 10*time.Second), ResponseHeaderTimeout: minDuration(time.Until(deadline), 10*time.Second), } + if proxyURL.Scheme == "socks5" { + transport.DialContext = newSOCKS5Dialer(proxyURL).DialContext + } else { + transport.Proxy = http.ProxyURL(proxyURL) + } defer transport.CloseIdleConnections() client := &http.Client{ Transport: transport, @@ -80,6 +84,9 @@ func (executor *Executor) Execute(ctx context.Context, task *controlplanev1.Chec if errors.Is(probeContext.Err(), context.DeadlineExceeded) || errors.Is(ctx.Err(), context.DeadlineExceeded) { return Result{FailureClass: FailureDeadlineExceeded, Latency: latency} } + if errors.Is(err, errSOCKS5Authentication) { + return Result{FailureClass: FailureProxyAuth, Latency: latency} + } return Result{FailureClass: FailureProxyRequest, Latency: latency} } _ = response.Body.Close() @@ -109,6 +116,13 @@ func prepare(task *controlplanev1.CheckTask, now time.Time) (time.Time, string, proxyScheme = "http" case controlplanev1.ProxyProtocol_PROXY_PROTOCOL_HTTPS: proxyScheme = "https" + case controlplanev1.ProxyProtocol_PROXY_PROTOCOL_SOCKS5: + if task.GetUsername() != "" || task.GetPassword() != "" { + if len(task.GetUsername()) == 0 || len(task.GetUsername()) > 255 || len(task.GetPassword()) == 0 || len(task.GetPassword()) > 255 { + return time.Time{}, "", nil, errors.New("invalid SOCKS5 credentials") + } + } + proxyScheme = "socks5" default: return time.Time{}, "", nil, errUnsupportedProtocol } diff --git a/internal/checker/probe/probe_test.go b/internal/checker/probe/probe_test.go index 0e8d216..bf00334 100644 --- a/internal/checker/probe/probe_test.go +++ b/internal/checker/probe/probe_test.go @@ -2,10 +2,13 @@ package probe import ( "context" + "encoding/binary" + "io" "net" "net/http" "net/http/httptest" "strconv" + "strings" "testing" "time" @@ -77,7 +80,7 @@ func TestExecutorRejectsEgressTaskWithRoutingProfile(t *testing.T) { func TestExecutorReportsUnsupportedProtocolAsFact(t *testing.T) { now := time.Now().UTC() result := NewExecutor().Execute(context.Background(), &controlplanev1.CheckTask{ - TaskId: "task-a", ProxyId: "proxy-a", Protocol: controlplanev1.ProxyProtocol_PROXY_PROTOCOL_SOCKS5, + TaskId: "task-a", ProxyId: "proxy-a", Protocol: controlplanev1.ProxyProtocol_PROXY_PROTOCOL_UNSPECIFIED, Host: "proxy.example", Port: 1080, Level: controlplanev1.CheckLevel_CHECK_LEVEL_BASIC, Timeout: durationpb.New(time.Second), Attempt: 1, MaxAttempts: 1, Deadline: timestamppb.New(now.Add(time.Second)), }) @@ -86,6 +89,41 @@ func TestExecutorReportsUnsupportedProtocolAsFact(t *testing.T) { } } +func TestExecutorBasicSupportsAuthenticatedSOCKS5Proxy(t *testing.T) { + target := httptest.NewServer(http.HandlerFunc(func(response http.ResponseWriter, request *http.Request) { + if request.Host != "example.invalid" { + t.Errorf("request host = %q", request.Host) + } + response.WriteHeader(http.StatusBadGateway) + })) + defer target.Close() + proxy := newSOCKS5TestProxy(t, target.URL, "user", "secret") + task := validSOCKS5Task(t, proxy.address, "user", "secret") + + result := NewExecutor().Execute(context.Background(), task) + if !result.Success || result.FailureClass != "" || result.Latency <= 0 { + t.Fatalf("Execute(SOCKS5 BASIC) = %+v", result) + } + select { + case address := <-proxy.requested: + if address != "example.invalid:80" { + t.Fatalf("SOCKS5 requested address = %q", address) + } + case <-time.After(time.Second): + t.Fatal("SOCKS5 proxy did not receive CONNECT request") + } +} + +func TestExecutorReportsSOCKS5AuthenticationFailure(t *testing.T) { + target := httptest.NewServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) {})) + defer target.Close() + proxy := newSOCKS5TestProxy(t, target.URL, "user", "secret") + result := NewExecutor().Execute(context.Background(), validSOCKS5Task(t, proxy.address, "user", "wrong")) + if result.Success || result.FailureClass != FailureProxyAuth || result.Latency <= 0 { + t.Fatalf("Execute(SOCKS5 auth failure) = %+v", result) + } +} + func validTask(t *testing.T, proxyAddress string, level controlplanev1.CheckLevel) *controlplanev1.CheckTask { t.Helper() host, port := splitProxyAddress(t, proxyAddress) @@ -97,3 +135,168 @@ func validTask(t *testing.T, proxyAddress string, level controlplanev1.CheckLeve Deadline: timestamppb.New(now.Add(time.Second)), } } + +func validSOCKS5Task(t *testing.T, address, username, password string) *controlplanev1.CheckTask { + t.Helper() + host, rawPort, err := net.SplitHostPort(address) + if err != nil { + t.Fatalf("net.SplitHostPort(%q): %v", address, err) + } + port, err := strconv.ParseUint(rawPort, 10, 16) + if err != nil { + t.Fatalf("strconv.ParseUint(%q): %v", rawPort, err) + } + now := time.Now().UTC() + return &controlplanev1.CheckTask{ + TaskId: "task-socks", ProxyId: "proxy-socks", Protocol: controlplanev1.ProxyProtocol_PROXY_PROTOCOL_SOCKS5, + Host: host, Port: uint32(port), Level: controlplanev1.CheckLevel_CHECK_LEVEL_BASIC, Username: username, Password: password, + Timeout: durationpb.New(time.Second), Attempt: 1, MaxAttempts: 1, Deadline: timestamppb.New(now.Add(time.Second)), + } +} + +type socks5TestProxy struct { + address string + target string + username string + password string + requested chan string + listener net.Listener +} + +func newSOCKS5TestProxy(t *testing.T, target, username, password string) *socks5TestProxy { + t.Helper() + listener, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("net.Listen(): %v", err) + } + proxy := &socks5TestProxy{ + address: listener.Addr().String(), target: target, username: username, password: password, + requested: make(chan string, 1), listener: listener, + } + go func() { + for { + connection, err := listener.Accept() + if err != nil { + return + } + go proxy.serve(connection) + } + }() + t.Cleanup(func() { _ = listener.Close() }) + return proxy +} + +func (proxy *socks5TestProxy) serve(connection net.Conn) { + defer connection.Close() + if _, err := readSOCKS5Greeting(connection); err != nil { + return + } + method := byte(0x00) + if proxy.username != "" || proxy.password != "" { + method = 0x02 + } + if _, err := connection.Write([]byte{0x05, method}); err != nil { + return + } + if method == 0x02 { + username, password, err := readSOCKS5Credentials(connection) + if err != nil { + return + } + if username != proxy.username || password != proxy.password { + _, _ = connection.Write([]byte{0x01, 0x01}) + return + } + if _, err := connection.Write([]byte{0x01, 0x00}); err != nil { + return + } + } + address, err := readSOCKS5ConnectRequest(connection) + if err != nil { + return + } + proxy.requested <- address + targetAddress := strings.TrimPrefix(proxy.target, "http://") + target, err := net.Dial("tcp", targetAddress) + if err != nil { + _, _ = connection.Write([]byte{0x05, 0x01, 0x00, 0x01, 0, 0, 0, 0, 0, 0}) + return + } + defer target.Close() + if _, err := connection.Write([]byte{0x05, 0x00, 0x00, 0x01, 0, 0, 0, 0, 0, 0}); err != nil { + return + } + done := make(chan struct{}, 2) + go func() { _, _ = io.Copy(target, connection); done <- struct{}{} }() + go func() { _, _ = io.Copy(connection, target); done <- struct{}{} }() + <-done +} + +func readSOCKS5Greeting(reader io.Reader) ([]byte, error) { + header := make([]byte, 2) + if _, err := io.ReadFull(reader, header); err != nil || header[0] != 0x05 || header[1] == 0 { + return nil, io.ErrUnexpectedEOF + } + methods := make([]byte, header[1]) + _, err := io.ReadFull(reader, methods) + return methods, err +} + +func readSOCKS5Credentials(reader io.Reader) (string, string, error) { + header := make([]byte, 2) + if _, err := io.ReadFull(reader, header); err != nil || header[0] != 0x01 || header[1] == 0 { + return "", "", io.ErrUnexpectedEOF + } + username := make([]byte, header[1]) + if _, err := io.ReadFull(reader, username); err != nil { + return "", "", err + } + passwordLength := make([]byte, 1) + if _, err := io.ReadFull(reader, passwordLength); err != nil || passwordLength[0] == 0 { + return "", "", io.ErrUnexpectedEOF + } + password := make([]byte, passwordLength[0]) + if _, err := io.ReadFull(reader, password); err != nil { + return "", "", err + } + return string(username), string(password), nil +} + +func readSOCKS5ConnectRequest(reader io.Reader) (string, error) { + header := make([]byte, 4) + if _, err := io.ReadFull(reader, header); err != nil || header[0] != 0x05 || header[1] != 0x01 || header[2] != 0x00 { + return "", io.ErrUnexpectedEOF + } + var host string + switch header[3] { + case 0x01: + address := make([]byte, net.IPv4len) + if _, err := io.ReadFull(reader, address); err != nil { + return "", err + } + host = net.IP(address).String() + case 0x03: + length := make([]byte, 1) + if _, err := io.ReadFull(reader, length); err != nil || length[0] == 0 { + return "", io.ErrUnexpectedEOF + } + address := make([]byte, length[0]) + if _, err := io.ReadFull(reader, address); err != nil { + return "", err + } + host = string(address) + case 0x04: + address := make([]byte, net.IPv6len) + if _, err := io.ReadFull(reader, address); err != nil { + return "", err + } + host = net.IP(address).String() + default: + return "", io.ErrUnexpectedEOF + } + port := make([]byte, 2) + if _, err := io.ReadFull(reader, port); err != nil { + return "", err + } + return net.JoinHostPort(host, strconv.Itoa(int(binary.BigEndian.Uint16(port)))), nil +} diff --git a/internal/checker/probe/socks5.go b/internal/checker/probe/socks5.go new file mode 100644 index 0000000..9ae72bb --- /dev/null +++ b/internal/checker/probe/socks5.go @@ -0,0 +1,180 @@ +package probe + +import ( + "context" + "encoding/binary" + "errors" + "fmt" + "io" + "net" + "net/url" + "strconv" + "time" +) + +const ( + socks5Version = 0x05 + socks5NoAuth = 0x00 + socks5UserPassword = 0x02 + socks5NoAcceptable = 0xff + socks5ConnectCommand = 0x01 + socks5AddressIPv4 = 0x01 + socks5AddressDomain = 0x03 + socks5AddressIPv6 = 0x04 +) + +var errSOCKS5Authentication = errors.New("SOCKS5 authentication failed") + +type socks5Dialer struct { + proxyAddress string + username string + password string + dialer net.Dialer +} + +func newSOCKS5Dialer(proxyURL *url.URL) socks5Dialer { + username := "" + password := "" + if proxyURL.User != nil { + username = proxyURL.User.Username() + password, _ = proxyURL.User.Password() + } + return socks5Dialer{proxyAddress: proxyURL.Host, username: username, password: password} +} + +func (dialer socks5Dialer) DialContext(ctx context.Context, _ string, address string) (net.Conn, error) { + connection, err := dialer.dialer.DialContext(ctx, "tcp", dialer.proxyAddress) + if err != nil { + return nil, err + } + if deadline, hasDeadline := ctx.Deadline(); hasDeadline { + if err := connection.SetDeadline(deadline); err != nil { + _ = connection.Close() + return nil, err + } + defer connection.SetDeadline(time.Time{}) + } + if err := dialer.connect(connection, address); err != nil { + _ = connection.Close() + return nil, err + } + return connection, nil +} + +func (dialer socks5Dialer) connect(connection net.Conn, address string) error { + method := byte(socks5NoAuth) + if dialer.username != "" || dialer.password != "" { + method = socks5UserPassword + } + if _, err := connection.Write([]byte{socks5Version, 0x01, method}); err != nil { + return err + } + selection := make([]byte, 2) + if _, err := io.ReadFull(connection, selection); err != nil { + return err + } + if selection[0] != socks5Version || selection[1] == socks5NoAcceptable { + return fmt.Errorf("invalid SOCKS5 method selection") + } + if selection[1] == socks5UserPassword { + if method != socks5UserPassword { + return errSOCKS5Authentication + } + if err := socks5Authenticate(connection, dialer.username, dialer.password); err != nil { + return err + } + } else if selection[1] != method { + return fmt.Errorf("unsupported SOCKS5 authentication method %d", selection[1]) + } + request, err := socks5ConnectRequest(address) + if err != nil { + return err + } + if _, err := connection.Write(request); err != nil { + return err + } + return readSOCKS5ConnectReply(connection) +} + +func socks5Authenticate(connection net.Conn, username, password string) error { + if len(username) == 0 || len(username) > 255 || len(password) == 0 || len(password) > 255 { + return errSOCKS5Authentication + } + request := make([]byte, 0, len(username)+len(password)+3) + request = append(request, 0x01, byte(len(username))) + request = append(request, username...) + request = append(request, byte(len(password))) + request = append(request, password...) + if _, err := connection.Write(request); err != nil { + return err + } + reply := make([]byte, 2) + if _, err := io.ReadFull(connection, reply); err != nil { + return err + } + if reply[0] != 0x01 || reply[1] != 0x00 { + return errSOCKS5Authentication + } + return nil +} + +func socks5ConnectRequest(address string) ([]byte, error) { + host, rawPort, err := net.SplitHostPort(address) + if err != nil || host == "" { + return nil, errors.New("invalid SOCKS5 target address") + } + port, err := strconv.ParseUint(rawPort, 10, 16) + if err != nil || port == 0 { + return nil, errors.New("invalid SOCKS5 target port") + } + request := []byte{socks5Version, socks5ConnectCommand, 0x00} + if parsed := net.ParseIP(host); parsed != nil { + if ipv4 := parsed.To4(); ipv4 != nil { + request = append(request, socks5AddressIPv4) + request = append(request, ipv4...) + } else { + request = append(request, socks5AddressIPv6) + request = append(request, parsed.To16()...) + } + } else { + if len(host) > 255 { + return nil, errors.New("SOCKS5 target hostname is too long") + } + request = append(request, socks5AddressDomain, byte(len(host))) + request = append(request, host...) + } + encodedPort := make([]byte, 2) + binary.BigEndian.PutUint16(encodedPort, uint16(port)) + return append(request, encodedPort...), nil +} + +func readSOCKS5ConnectReply(reader io.Reader) error { + header := make([]byte, 4) + if _, err := io.ReadFull(reader, header); err != nil { + return err + } + if header[0] != socks5Version || header[1] != 0x00 || header[2] != 0x00 { + return fmt.Errorf("SOCKS5 CONNECT failed with status %d", header[1]) + } + addressLength := 0 + switch header[3] { + case socks5AddressIPv4: + addressLength = net.IPv4len + case socks5AddressIPv6: + addressLength = net.IPv6len + case socks5AddressDomain: + length := make([]byte, 1) + if _, err := io.ReadFull(reader, length); err != nil { + return err + } + addressLength = int(length[0]) + if addressLength == 0 { + return errors.New("invalid SOCKS5 bound hostname") + } + default: + return errors.New("invalid SOCKS5 bound address type") + } + boundAddressAndPort := make([]byte, addressLength+2) + _, err := io.ReadFull(reader, boundAddressAndPort) + return err +}