what ee5e5b96af feat: 添加 WithStreamErrors 查询流式调用执行过程中的异常
流式输出/双向流的 handler(Python 生成器)如果执行过程中抛异常,
Invoke[chan T] 本身的 err 只描述"调用有没有发起成功",跟这个异常
无关(永远是 nil),channel 只会静默提前关闭,调用方原本完全无法
感知。新增 WithStreamErrors(ctx) 返回一个包过的 ctx 和一个查询函数
streamErr,opt-in 之后可以查到具体错误。

错误记录挂在 WithStreamErrors 返回的 ctx 的对象图里(context.WithValue),
不是全局表——调用方不再引用 ctx/channel 时会被 GC 自然回收,不需要
任何显式清理逻辑,也不依赖 ctx.Done(),即使用 context.Background()
也能正常释放;ctx 之后被别的 context.With*(包括 StickyCtx)再包一层
也不影响查询。

同时补充完整的自动化测试覆盖 example/main.go 里演示过的所有功能:
四种调用模式 × int/struct/slice/[]byte 的组合(client_test.go)、
WithHandlers/call_go 全双工(handlers_test.go)、NewSession 隔离性
与 StickyCtx 路由(session_test.go),之前这些只能靠人肉跑 go run
看输出,现在都有真实断言。
2026-07-23 16:24:05 +08:00
2026-04-10 09:24:28 +08:00
2026-04-10 09:24:28 +08:00

gobridge

在 Go 与 Python 之间建立双向通信桥接,将 Go 的 channel 与 Python 的 yield 原生对接,支持普通调用、单向流、双向流。

底层通过 Unix Domain Socket (UDS) 通信,Go 侧维护 Worker 进程池,Python 侧以多线程方式并发处理请求。

特性

  • 零配置序列化Go struct/slice ↔ Python dict/list 通过 JSON 自动互转
  • 原生流语义Go chan T 对应 Python Iterator[T],无需额外 API
  • 进程池:Go 自动启动并管理多个 Python 子进程,崩溃后自动重启
  • ctx 取消Go context 取消时自动中断 Python 计算,无需函数内检查
  • 四种调用模式:普通、流式输出、流式输入、双向流,同一个 Invoke 函数自动推断

安装

go get git.fsdpf.net/go/gobridge

Python 端直接复制 python/gobridge/ 目录到项目中,无需安装依赖(仅用标准库)。

快速开始

Python 端(worker.py):

from gobridge import gobridge, run
from typing import Iterator

@expose
def add(a: int, b: int) -> int:
    return a + b

run()

Go 端:

pool, _ := gobridge.NewPool("worker.py")
defer pool.Close()

ctx := context.Background()
sum, _ := gobridge.Invoke[int](ctx, pool, "add", 3, 4)
fmt.Println(sum) // 7

四种调用模式

1. 普通调用

// Go
sum, err := gobridge.Invoke[int](ctx, pool, "add", 3, 4)

user, err := gobridge.Invoke[User](ctx, pool, "get_user", 42)

result, err := gobridge.Invoke[[]User](ctx, pool, "enrich_users", users)
# Python
@expose
def add(a: int, b: int) -> int:
    return a + b

@expose
def get_user(uid: int) -> dict:
    return {"id": uid, "name": f"user_{uid}", "score": uid * 1.5}

@expose
def enrich_users(users: list) -> list:
    for u in users:
        u["level"] = "gold" if u["score"] >= 10 else "silver"
    return users

2. 流式输出(Python yield → Go channel

返回类型为 chan T 时自动进入流式输出模式,Python 函数使用 yieldGo 侧通过 range 消费。

// Go
ch, err := gobridge.Invoke[chan int](ctx, pool, "range_gen", 1, 6)
for v := range ch {
    fmt.Println(v) // 1 2 3 4 5
}

userCh, err := gobridge.Invoke[chan User](ctx, pool, "gen_users", 3)
for u := range userCh {
    fmt.Println(u)
}
# Python
@expose
def range_gen(start: int, stop: int) -> Iterator[int]:
    for i in range(start, stop):
        yield i

@expose
def gen_users(count: int) -> Iterator[dict]:
    for i in range(1, count + 1):
        yield {"id": i, "name": f"user_{i}", "score": float(i * 3)}

3. 流式输入(Go channel → Python Iterator

参数中含 chan T 且返回非 chan 时自动进入流式输入模式。

// Go
inputCh := make(chan int, 10)
go func() {
    for i := 1; i <= 5; i++ {
        inputCh <- i
    }
    close(inputCh)
}()
total, err := gobridge.Invoke[int](ctx, pool, "sum_stream", inputCh)
fmt.Println(total) // 15
# Python
@expose
def sum_stream(numbers: Iterator[int]) -> int:
    return sum(numbers)

4. 双向流(Go channel 输入 + Go channel 输出)

参数含 chan T 且返回类型也为 chan R 时自动进入双向流模式。

// Go
inCh := make(chan User, 5)
go func() {
    for _, u := range users {
        inCh <- u
    }
    close(inCh)
}()
outCh, err := gobridge.Invoke[chan User](ctx, pool, "process_users", inCh)
for u := range outCh {
    fmt.Println(u)
}
# Python
@expose
def process_users(users: Iterator[dict]) -> Iterator[dict]:
    for u in users:
        yield {"id": u["id"], "name": u["name"].upper(), "score": u["score"] * 2}

ctx 取消

Go 的 context 取消会自动中断 Python 侧的执行,无需在 Python 函数中做任何检查:

ctx, cancel := context.WithTimeout(context.Background(), 2*time.Second)
defer cancel()

// 超时或手动 cancel() 后,Python 计算立即中断,返回 context.DeadlineExceeded
result, err := gobridge.Invoke[int](ctx, pool, "slow_compute", 1000000)
@expose
def slow_compute(n: int) -> int:
    total = 0
    for i in range(n):
        total += i   # ctx 取消时此处自动抛出 InterruptedError,无需手动检查
    return total

实现机制:

  1. Go ctx 取消 → 发送 cancel 消息给 Python
  2. Python 单 reader 线程收到 cancel → 向执行线程注入 InterruptedErrorPyThreadState_SetAsyncExc
  3. Python 函数在下一条字节码指令处中断
  4. Go 同时关闭连接,解除阻塞的读写操作

限制:长时间不释放 GIL 的 C 扩展(如大规模 numpy 矩阵运算)无法被中断,需等其释放 GIL 后才触发。

默认超时(WithDefaultTimeout

Invoke 本身不带任何默认超时——如果传入的 ctx 没有 deadline(比如直接用 context.Background()),且池子被占满(所有 worker 的连接都在处理别的请求),调用会永久阻塞,不会自动放弃。

WithDefaultTimeout 用于兜底这种情况:仅当调用方传入的 ctx 未设置 deadline 时才生效,调用方显式设置的 context.WithTimeout 优先级更高,不会被覆盖。

池子被占满时新调用是排队阻塞等待,不是立刻失败或被跳过——example/main.go 中的 demoBlockingworkers=2, maxConns=2(总容量 4)故意占满连接池, 再发起一个不设超时的调用,实测会阻塞约 1.8s(等到某个占位任务释放连接)才返回,而不是瞬间失败。

pool, _ := gobridge.NewPool("worker.py",
    gobridge.WithDefaultTimeout(500 * time.Millisecond),
)

// 未设置 deadline,池的默认超时自动生效,500ms 后返回 context.DeadlineExceeded
_, err := gobridge.Invoke[string](context.Background(), pool, "sleep_seconds", 2.0)

// 显式传入的 deadline 优先,不受默认超时影响
ctx, cancel := context.WithTimeout(context.Background(), 3*time.Second)
defer cancel()
result, err := gobridge.Invoke[string](ctx, pool, "sleep_seconds", 1.0) // 正常返回

WithDefaultTimeout 只对普通调用和流式输入生效,对流式输出/双向流完全不生效(即使不设置也不会被套用)。 原因:Invoke[chan T] 返回的 channel 由调用方通过 range 自行决定消费多久——典型场景是流式聊天回复, 可能正常持续几十秒甚至更久,如果套一个全局默认超时,会在输出还在正常进行时把它腰斩。 流式输出/双向流的调用只认调用方显式传入的 deadline;不想让它无限期阻塞(包括等待连接、等待数据), 必须自己 context.WithTimeout 包一层传进去,WithDefaultTimeout 在这里不提供任何兜底。

超时是否会通过 error 返回,取决于调用模式:

调用模式 WithDefaultTimeout 是否生效 超时表现
普通调用 / 流式输入 生效 Invoke 返回非 nil 的 errorcontext.DeadlineExceeded / context.Canceled
流式输出 / 双向流 不生效,只认显式 deadline Invoke 建立阶段失败会返回 error建立成功后若显式 ctx 超时,只会静默关闭已返回的 channel,不会有第二次 error

流式模式下,for v := range ch 结束后无法区分"正常读完"还是"被超时打断",需要调用方自行检查传入的 ctx.Err()

// 流式输出必须自己设置超时,WithDefaultTimeout 不会兜底
ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
defer cancel()
ch, err := gobridge.Invoke[chan int](ctx, pool, "slow_range_gen", 1, 10, 100)
if err != nil {
    // 建立阶段失败/超时
}
for v := range ch {
    fmt.Println(v)
}
if ctx.Err() != nil {
    // channel 是因为超时/取消提前关闭的,不是正常 yield 完
}

完整可运行示例见 example/main.go 中的 demoTimeout(默认超时 / 显式 deadline 优先级 / 流式超时静默关闭 channel)和 demoBlocking(池子占满后阻塞排队)。

查询流式调用执行过程中的错误(WithStreamErrors

流式输出/双向流的 handler(Python 生成器)如果执行过程中抛异常(不是等全部 yield 完才失败,是还没跑完就失败),默认情况下调用方完全无法感知——Invoke[chan T] 返回的那个 err 只描述"调用有没有成功发起"(参数序列化、抢连接、写 call 消息),跟 Python 函数体最终是正常结束还是执行过程中抛异常没有任何关系,因为 handler 真正执行是在 Invoke 返回之后的后台 goroutine 里才发生的。Python 侧异常只会让 ch 静默提前关闭,现象上跟"正常读完"一模一样,ctx.Err() 也查不出来(不是 ctx 取消导致的)。

想知道是不是这种情况,用 WithStreamErrors 包一层 ctx(opt-in),它会额外返回一个查询函数,之后不需要再传 ctx:

ctx, streamErr := gobridge.WithStreamErrors(ctx) // opt-in,不包这一层就是上面说的默认行为
ch, err := gobridge.Invoke[chan int](ctx, pool, "range_gen", 1, 10)
if err != nil {
    // 建立阶段失败,比如池子占满/参数错误
}
for v := range ch {
    fmt.Println(v)
}
if err := streamErr(ch); err != nil {
    // channel 是因为 Python 侧执行过程中抛异常提前关闭的,不是正常 yield 完
}

错误记录挂在 WithStreamErrors 返回的那个 ctx 的对象图里(通过 context.WithValue 携带一个可变的存储,streamErr 直接持有这块存储,不需要再查 ctx),不是全局表——调用方不再引用这个 ctx(和对应的 ch)时,整条链会被 GC 自然回收,不需要任何显式清理逻辑,也不依赖 ctx.Done(),即使用 context.Background() 也能正常释放。WithStreamErrors 之后即使 ctx 又被别的 context.With*(包括库里自己的 StickyCtx)再包一层,也不影响——context.Value() 的查找是逐层往上找的,不会因为后续再包装而丢失(streamErr 是从最初那次 WithStreamErrors 调用里直接拿到的闭包,跟后续怎么包 ctx 完全无关)。

不调用 WithStreamErrors 是完全零成本的默认行为,现有的流式调用代码不需要做任何改动。完整可运行示例见 example/main.godemoTimeout 的示例5。

查询流式调用中途的错误(StreamError

流式输出/双向流的 handler(Python 生成器)如果执行过程中抛异常(不是等全部 yield 完才失败,是还没跑完就失败),默认情况下调用方完全无法感知——Invoke[chan T] 返回的那个 err 只描述"调用有没有成功发起"(参数序列化、抢连接、写 call 消息),跟 Python 函数体最终是正常结束还是中途抛异常没有任何关系,因为 handler 真正执行是在 Invoke 返回之后的后台 goroutine 里才发生的。Python 侧异常只会让 ch 静默提前关闭,现象上跟"正常读完"一模一样,ctx.Err() 也查不出来(不是 ctx 取消导致的)。

想知道是不是这种情况,用 StreamError 查——不需要任何额外配置,直接传入 Invoke[chan T] 返回的那个 channel

ch, err := gobridge.Invoke[chan int](ctx, pool, "range_gen", 1, 10)
if err != nil {
    // 建立阶段失败,比如池子占满/参数错误
}
for v := range ch {
    fmt.Println(v)
}
if streamErr := gobridge.StreamError(ch); streamErr != nil {
    // channel 是因为 Python 侧执行过程中抛异常提前关闭的,不是正常 yield 完
}

StreamError 内部用调用方传给 Invoke 的那个 ctx 挂了 context.AfterFuncctx 结束(取消/超时/调用方自己 cancel())时会自动清理对应记录,不需要调用方主动查询才释放——前提是 ctx 本身会结束(context.Background() 永不 Done(),如果还从来不查 StreamError,记录会一直留着,这也是本文档反复强调"别用没有 deadline 的 ctx"的场景之一)。

Session 亲和路由

默认情况下,每次 Invoke 通过轮询分配 worker 进程。当多次调用需要共享同一 Python 进程的状态时,可以使用 Session 或 StickyCtx 将调用固定到同一进程。

NewSession

NewSession 返回一个固定到某个 worker 的 Pool 视图,通过它发起的所有 Invoke 始终路由到同一 Python 进程:

pool, _ := gobridge.NewPool("worker.py", gobridge.WithWorkers(2))

sessA := gobridge.NewSession(pool) // 固定到 worker 1
sessB := gobridge.NewSession(pool) // 固定到 worker 0(轮询下一个)
sessC := gobridge.NewSession(pool) // 固定到 worker 1(与 sessA 同进程)

gobridge.Invoke(ctx, sessA, "init", "A", 100)
gobridge.Invoke(ctx, sessA, "step", "A", 10)  // 始终走 worker 1
gobridge.Invoke(ctx, sessC, "step", "C", 99)  // 也走 worker 1,但 session_id 不同

Python 侧用模块级 dict 以 session_id 为 key 隔离状态:

_sessions = {}

@expose
def init(session_id: str, value: int):
    _sessions[session_id] = {"value": value}

@expose
def step(session_id: str, delta: int) -> int:
    _sessions[session_id]["value"] += delta
    return _sessions[session_id]["value"]

NewSession 不拥有底层 pool 的生命周期,调用 session.Close() 是空操作,只需关闭原始 pool。

StickyCtx

StickyCtx 将亲和键写入 ctx,相同 key 通过哈希稳定路由到同一 worker,无需持有 session 对象:

// 每次调用前附加 key,相同 key 始终走同一 worker
ctx = gobridge.StickyCtx(ctx, "user-42")
gobridge.Invoke(ctx, pool, "method_a", ...)
gobridge.Invoke(ctx, pool, "method_b", ...)

适合按用户 ID、租户 ID 等自然键做路由,不需要显式创建 session 对象。

进程级全局变量

同一 worker 进程内的所有 session 共享进程级变量(如模块级 dict、计数器等),不同 worker 进程之间完全隔离:

_counter = 0  # 进程级全局变量

@expose
def increment(delta: int) -> int:
    global _counter
    _counter += delta
    return _counter
sessA := gobridge.NewSession(pool) // worker 1
sessC := gobridge.NewSession(pool) // worker 1(与 A 同进程)
sessB := gobridge.NewSession(pool) // worker 0(独立进程)

gobridge.Invoke(ctx, sessA, "increment", 10) // worker 1: counter = 10
gobridge.Invoke(ctx, sessC, "increment", 5)  // worker 1: counter = 15(共享)
gobridge.Invoke(ctx, sessB, "increment", 99) // worker 0: counter = 99(独立)

两种方式对比

NewSession StickyCtx
路由方式 创建时轮询确定 worker 按 key 哈希确定 worker
适用场景 显式会话管理 按自然键(用户 ID 等)路由
携带方式 替换 pool 参数 写入 ctx
超时支持 context.WithTimeout(ctx, d) 正常使用 同左

[]byte ↔ bytes 支持

Go 的 []byte 通过 base64 编码在 JSON 帧中传输,框架在 Python 侧自动完成编解码,用户侧完全透明。

Python 侧

函数参数注解为 bytes 时,框架自动将 Go 传入的 base64 字符串解码为 bytes
返回值为 bytes 时,框架自动将其编码为 base64 字符串再发送给 Go。

from gobridge import expose
from typing import Iterator

@expose
def bytes_reverse(data: bytes) -> bytes:
    return data[::-1]

@expose
def bytes_concat(a: bytes, b: bytes) -> bytes:
    return a + b

# 流式输出 bytes:对应 Go Invoke[chan []byte]
@expose
def bytes_chunks(data: bytes, size: int) -> Iterator[bytes]:
    for i in range(0, len(data), size):
        yield data[i:i + size]

Go 侧

Go 直接使用 []byteencoding/json 自动处理 base64 编解码:

// 普通 []byte 参数与返回值
rev, err := gobridge.Invoke[[]byte](ctx, pool, "bytes_reverse", []byte("hello"))
fmt.Printf("%s\n", rev) // olleh

cat, err := gobridge.Invoke[[]byte](ctx, pool, "bytes_concat", []byte("foo"), []byte("bar"))
fmt.Printf("%s\n", cat) // foobar

// 流式输出 []byte
ch, err := gobridge.Invoke[chan []byte](ctx, pool, "bytes_chunks", []byte("abcdefgh"), 3)
for chunk := range ch {
    fmt.Printf("%s ", chunk) // abc def gh
}

效率说明:base64 编码约使数据体积增大 33%,并有少量 CPU 开销。
对于小块二进制数据(< 1 MB)的 RPC 调用,这通常可以忽略不计;
若需传输大量原始二进制流,建议改用独立的 socket/文件通道。

注意事项

在 handler 中使用 threading.Thread

handler 函数内部可以启动后台线程,但行为取决于是否等待:

# ✅ fire-and-forget:立即返回,连接立即释放,后台线程独立运行
@expose
def do_something():
    threading.Thread(target=long_task, daemon=True).start()
    return "ok"

# ⚠️ 等待线程:连接被占用直到线程结束,等同于直接在 handler 里执行
@expose
def do_something():
    t = threading.Thread(target=long_task)
    t.start()
    t.join()  # 连接在此阻塞
    return "ok"

call_go() 只能在 handler 的原始线程中调用。 call_go() 依赖线程局部变量 _local.mux 获取当前连接,后台线程中该变量不存在,调用会抛出 RuntimeError

@expose
def do_something():
    def bg():
        call_go("Method")  # ❌ RuntimeError: call_go() must be called within a gobridge handler
    threading.Thread(target=bg, daemon=True).start()
    return "ok"

如果后台线程的结果需要回调 Go,应在 handler 线程中通过 queue.Queue 等待后台线程结果后再调用:

@expose
def do_something():
    result_q = queue.Queue()
    threading.Thread(target=lambda: result_q.put(compute()), daemon=True).start()
    result = result_q.get()          # 等待后台线程
    return call_go[str]("Process", result)  # ✅ 在 handler 线程中调用

进程自动重启

Python worker 进程崩溃时自动重启,调用方无感知:

Python 进程崩溃
  → monitor goroutine 检测到退出
  → 排空失效连接
  → 指数退避重启(100ms → 200ms → ... → 30s)
  → 新进程就绪后恢复连接池

配置

NewPool 使用函数选项模式,第一个参数为脚本路径:

// 最简调用
pool, err := gobridge.NewPool("worker.py")

// 完整配置
pool, err := gobridge.NewPool("worker.py",
    gobridge.WithWorkers(4),                         // Python 进程数量,默认 2
    gobridge.WithMaxConns(8),                        // 每进程最大连接数,默认 4
    gobridge.WithPythonExe("python3"),               // 可执行文件,默认 "python3"
    gobridge.WithWorkDir("/path/to/workdir"),        // 工作目录,默认继承当前进程
    gobridge.WithEnv("PYTHONUNBUFFERED=1", "K=V"),  // 附加环境变量,与当前进程环境合并
    gobridge.WithSocketDir("/var/run/myapp"),        // socket 文件目录,默认 /tmp
    gobridge.WithStdout(os.Stdout),                 // 子进程 stdout,默认 os.Stdout
    gobridge.WithStderr(os.Stderr),                 // 子进程 stderr,默认 os.Stderr
    gobridge.WithDefaultTimeout(10*time.Second),    // Invoke 默认超时,默认不启用
)

// 静默模式:丢弃子进程输出
pool, err := gobridge.NewPool("worker.py",
    gobridge.WithStdout(io.Discard),
    gobridge.WithStderr(io.Discard),
)
Option 说明 默认值
WithWorkers(n) Python 进程数量 2
WithMaxConns(n) 每进程最大连接数 4
WithPythonExe(exe) 可执行文件 "python3"
WithScriptArgs(args...) 脚本路径之后的附加参数
WithWorkDir(dir) 子进程工作目录 继承当前进程
WithEnv(kv...) 附加环境变量 "K=V"
WithSocketDir(dir) UDS socket 文件目录 "/tmp"
WithStdout(w) 子进程标准输出 os.Stdout
WithStderr(w) 子进程标准错误 os.Stderr
WithDefaultTimeout(d) Invoke 默认超时,仅在传入的 ctx 未设置 deadline 时生效 不启用

使用 uv 管理 Python 环境

推荐使用 uv 管理 Python 版本和虚拟环境。

方式一:uv run(推荐,无需手动激活环境)

// 等价于执行:uv run worker.py
pool, err := gobridge.NewPool("run",
    gobridge.WithPythonExe("uv"),
    gobridge.WithScriptArgs("worker.py"),
    gobridge.WithWorkDir("./worker"), // uv 项目目录(含 pyproject.toml
)

方式二:直接使用虚拟环境的 python

cd worker && uv sync
venvPython, _ := exec.LookPath("worker/.venv/bin/python")
pool, err := gobridge.NewPool("worker/worker.py",
    gobridge.WithPythonExe(venvPython),
)

方式三:shell 脚本封装(适合 CI/部署)

#!/bin/sh
# run_worker.sh
cd "$(dirname "$0")"
exec uv run python worker.py
pool, err := gobridge.NewPool("./run_worker.sh",
    gobridge.WithPythonExe("/bin/sh"),
)

典型项目结构:

myproject/
├── main.go
├── go.mod
└── worker/
    ├── pyproject.toml
    ├── uv.lock
    ├── .venv/
    └── worker.py

pyproject.toml

[project]
name = "worker"
version = "0.1.0"
requires-python = ">=3.11"
dependencies = []

[tool.uv.sources]
# 从 git 仓库安装(推荐)
gobridge = { git = "https://git.fsdpf.net/go/gobridge.git", subdirectory = "python" }

# 本地开发时改用本地路径
# gobridge = { path = "../../python", editable = true }

通信协议

整体架构

  Go 进程
  ┌─────────────────────────────────────────────────────────┐
  │  Invoke[R](ctx, pool, method, args...)                  │
  │       │                                                 │
  │  ┌────▼─────────────────────────────────────────────┐  │
  │  │  Pool                                            │  │
  │  │  workers[0]  workers[1]  ...  workers[N-1]       │  │
  │  │  (轮询选择)                                     │  │
  │  └────┬─────────────────────────────────────────────┘  │
  │       │  每 worker 维护 M 个可复用连接                  │
  └───────┼─────────────────────────────────────────────────┘
          │ Unix Domain Socket(每 worker 独立 .sock 文件)
  ┌───────▼──────────────┐  ┌──────────────────────────┐
  │  Python 进程 0        │  │  Python 进程 1            │
  │  worker.py           │  │  worker.py               │
  └──────────────────────┘  └──────────────────────────┘

Python Worker 内部结构

  Python 进程
  ┌──────────────────────────────────────────────────────────────┐
  │  run()  ──  UDS server.accept() 循环                         │
  │                │                                             │
  │         每个连接 → 独立线程 _handle_conn()                    │
  │                                                              │
  │  ┌─────────────────────────────────────────────────────┐     │
  │  │  _handle_conn(连接线程)                            │     │
  │  │                                                     │     │
  │  │  ┌──────────────────────────────────────────────┐   │     │
  │  │  │  _ConnMux(单 reader 线程)                   │   │     │
  │  │  │                                              │   │     │
  │  │  │  socket ──► 读消息                           │   │     │
  │  │  │                  │                           │   │     │
  │  │  │         ┌────────┼──────────┐                │   │     │
  │  │  │         ▼        ▼          ▼                │   │     │
  │  │  │      call_q   chunk_q    cancel              │   │     │
  │  │  └─────────┬────────┬──────────┼───────────────┘   │     │
  │  │            │        │          │                    │     │
  │  │            ▼        │          ▼                    │     │
  │  │      主循环读取      │    PyThreadState_SetAsyncExc  │     │
  │  │      _dispatch()    │    → 执行线程抛 InterruptedError│    │
  │  │            │        │                              │     │
  │  │     ┌──────▼──────┐ │                              │     │
  │  │     │ @expose fn│ │                              │     │
  │  │     │             │ │                              │     │
  │  │     │  普通函数    │ │                              │     │
  │  │     │  return val ──────────────► result/error     │     │
  │  │     │             │ │                              │     │
  │  │     │  生成器函数  │ │                              │     │
  │  │     │  yield val ───────────────► chunk × N        │     │
  │  │     │             │ │              + end           │     │
  │  │     │  流式输入    │ │                              │     │
  │  │     │  Iterator ◄─┘ │  ← chunk_q                  │     │
  │  │     └─────────────┘                               │     │
  │  └─────────────────────────────────────────────────────┘     │
  └──────────────────────────────────────────────────────────────┘

消息帧: [4字节大端长度][JSON载荷]

消息类型:

type 方向 含义
call Go → Python 调用请求
result Python → Go 普通返回值
chunk 双向 流数据块
end 双向 流结束标记
error 双向 错误响应
cancel Go → Python 取消请求,触发 InterruptedError

项目结构

gobridge/
├── protocol.go          # Message 结构与类型常量
├── framing.go           # 帧读写(4字节长度前缀 + JSON
├── worker.go            # Python 子进程管理 + UDS 连接池 + 自动重启
├── pool.go              # 多进程池(轮询负载均衡)+ Option 函数
├── client.go            # Invoke[R] 泛型函数(四种模式自动推断)
├── example/
│   ├── main.go          # 完整调用示例
│   └── worker.py        # Python 函数示例
└── python/
    └── gobridge/
        └── __init__.py  # Python 库(expose、run、_ConnMux

类型对应关系

Go 类型 Python 类型
int int
float64 float
string str
bool bool
[]byte bytes
struct dict
[]T list
chan T Iterator[T]

[]byte 经 base64 编码在 JSON 帧中传输,框架自动完成编解码,用户侧透明。chan T 中的 T 同样支持 []byte,即 chan []byteIterator[bytes]

参考

本项目的进程池、UDS 通信、帧协议设计参考自 pyproc,在此基础上增加了 Go channel 与 Python yield 的流式对接及 ctx 取消支持。

S
Description
No description provided
Readme
297 KiB
Languages
Go 76.2%
Python 23.8%