Skip to content

Commit 1006088

Browse files
committed
feat(sdk): 控制面 manifest 上传接线——Go/Python 补齐(矩阵缺口修正)
- Go:TCPManager.buildManifest(provider+functions 摘要,与 JS/C# 同构)+ gzip + 独立短连接 controlAddr → MsgRegisterCapabilitiesReq (5s 超时,best-effort 不阻断注册),3 例单测含真实 TCP 帧回路 - Python:_maybe_register_capabilities 接线(复用 build_manifest, TCPTransport 控制连接),2 例单测(模拟控制面收帧解压校验) - 矩阵修正:manifest 上传 Go/Python ✅、JS/C# ✅(原矩阵漏标)、 Java 部分(构建已备未接线)、C++ ❌;已知缺口章节同步 - 附带核对:Java TLS/drain 实际已实现,矩阵 L1 行同步标记
1 parent 0413c27 commit 1006088

5 files changed

Lines changed: 424 additions & 10 deletions

File tree

sdks/SDK_FEATURE_MATRIX.md

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -43,14 +43,14 @@
4343

4444
### L2 Provider 扩展
4545

46-
| 能力 | Go | Python | Java | JS/TS | C++ | C# |
47-
| ---------------------------------------------------------------------- | --------------- | --------------------- | ---- | ------------------- | --- | ------------------------ |
48-
| Descriptor v2 字段(builder/构造器) |||||||
49-
| 呈现 hints 便捷层(`SetFieldHint`/`SetFieldWidget` 等价,x-ui-* 契约) | ✅ builder 方法 |`set_field_hint()` ||`setFieldHint()` |||
50-
| OpenAPI 注册 helper(`RegisterFromOpenAPI` 等价) |||||||
51-
| JSON Schema 入站 payload 校验(provider 侧,`validateInputPayloads`||||||`JsonSchemaValidator` |
52-
| 控制面 manifest 上传(`control_addr``RegisterCapabilitiesRequest`| | | | || |
53-
| 文件传输(`enable_file_transfer`|||||||
46+
| 能力 | Go | Python | Java | JS/TS | C++ | C# |
47+
| ---------------------------------------------------------------------- | --------------- | --------------------- | ------------------------------------- | ------------------- | --- | ------------------------ |
48+
| Descriptor v2 字段(builder/构造器) ||| ||||
49+
| 呈现 hints 便捷层(`SetFieldHint`/`SetFieldWidget` 等价,x-ui-* 契约) | ✅ builder 方法 |`set_field_hint()` | |`setFieldHint()` |||
50+
| OpenAPI 注册 helper(`RegisterFromOpenAPI` 等价) ||| ||||
51+
| JSON Schema 入站 payload 校验(provider 侧,`validateInputPayloads`||| |||`JsonSchemaValidator` |
52+
| 控制面 manifest 上传(`control_addr``RegisterCapabilitiesRequest`| | |配置字段+manifest 构建已备,未接线 | || |
53+
| 文件传输(`enable_file_transfer`||| ||||
5454

5555
### L3 Invoker(invoke / startTask / getTaskStatus / streamTask / cancelTask)
5656

@@ -65,8 +65,9 @@ C# DI / Unity / Java Spring Boot starter),明细见下文第五章;此层
6565

6666
### 已知缺口(按优先级)
6767

68-
1. **manifest 上传(L2)**:六语言均未实现 `control_addr``RegisterCapabilitiesRequest`
69-
当前能力发现依赖 Agent 转发的注册帧。若控制台需要"不调用即知全量能力",需补齐。
68+
1. **manifest 上传(L2)**:Go/Python/JS/C#/Java 已实现(注册后独立短连接
69+
`control_addr``RegisterCapabilitiesRequest`,best-effort 不阻断注册);
70+
**C++ 未实现**。Java 已具备 manifest 构建,发送接线待补。
7071
2. **文件传输(L2)**:六语言均未实现;如无平台侧需求建议从矩阵移除或标注"规划中"。
7172
3. **Java TLS transport**(已补齐):`TlsSocketFactory` 支持 CA 校验与 mTLS(PKCS#8),
7273
显式 `startHandshake` + 握手后 SAN/CN 端点校验。
Lines changed: 155 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,155 @@
1+
package croupier
2+
3+
import (
4+
"bytes"
5+
"compress/gzip"
6+
"encoding/binary"
7+
"encoding/json"
8+
"io"
9+
"net"
10+
"strings"
11+
"testing"
12+
13+
"github.com/cuihairu/croupier/sdks/go/pkg/croupier/protocol"
14+
agentv1 "github.com/cuihairu/croupier/sdks/go/pkg/pb/croupier/agent/v1"
15+
sdkv1 "github.com/cuihairu/croupier/sdks/go/pkg/pb/croupier/sdk/v1"
16+
"google.golang.org/protobuf/proto"
17+
)
18+
19+
// F:控制面 manifest 上传——buildManifest 构建与 best-effort 上传。
20+
func TestTCPManager_BuildManifest(t *testing.T) {
21+
config := ClientConfig{
22+
ProviderLang: "go",
23+
ProviderSDK: "croupier-go-sdk",
24+
}
25+
manager, err := NewTCPManager(config, map[string]FunctionHandler{})
26+
if err != nil {
27+
t.Fatalf("NewTCPManager: %v", err)
28+
}
29+
m := manager.(*TCPManager)
30+
m.serviceID = "go-1"
31+
m.serviceVersion = "2.0.0"
32+
m.functions = []*sdkv1.ProviderFunctionDescriptor{
33+
{
34+
Id: "player.ban",
35+
Version: "1.0.0",
36+
Resource: "player",
37+
Operation: "ban",
38+
Risk: "high",
39+
Permission: "player.ban.invoke",
40+
Description: "ban a player",
41+
InputSchema: `{"type":"object","properties":{"id":{"type":"string"}}}`,
42+
OutputSchema: `{"type":"object"}`,
43+
},
44+
{Id: " "}, // 空 id 跳过
45+
}
46+
47+
manifest := m.buildManifest()
48+
provider, ok := manifest["provider"].(map[string]interface{})
49+
if !ok {
50+
t.Fatalf("provider missing: %+v", manifest)
51+
}
52+
if provider["id"] != "go-1" || provider["lang"] != "go" {
53+
t.Fatalf("unexpected provider: %+v", provider)
54+
}
55+
functions, ok := manifest["functions"].([]map[string]interface{})
56+
if !ok || len(functions) != 1 {
57+
t.Fatalf("expected 1 function entry, got %+v", manifest["functions"])
58+
}
59+
if functions[0]["id"] != "player.ban" {
60+
t.Fatalf("unexpected function: %+v", functions[0])
61+
}
62+
// schema 以原生 JSON 对象进入 manifest(非字符串)
63+
inputSchema, ok := functions[0]["inputSchema"].(json.RawMessage)
64+
if !ok {
65+
t.Fatalf("inputSchema should be raw JSON, got %T", functions[0]["inputSchema"])
66+
}
67+
var parsed map[string]interface{}
68+
if err := json.Unmarshal(inputSchema, &parsed); err != nil {
69+
t.Fatalf("inputSchema invalid: %v", err)
70+
}
71+
}
72+
73+
func TestTCPManager_MaybeRegisterCapabilities(t *testing.T) {
74+
// 模拟控制面:读一帧(4 字节长度 + 8 字节头 + body),回确认帧
75+
var gotCapabilities bool
76+
listener, err := net.Listen("tcp", "127.0.0.1:0")
77+
if err != nil {
78+
t.Fatalf("listen: %v", err)
79+
}
80+
defer listener.Close()
81+
go func() {
82+
conn, acceptErr := listener.Accept()
83+
if acceptErr != nil {
84+
return
85+
}
86+
defer conn.Close()
87+
header := make([]byte, 4)
88+
if _, err := io.ReadFull(conn, header); err != nil {
89+
return
90+
}
91+
size := int(header[0])<<24 | int(header[1])<<16 | int(header[2])<<8 | int(header[3])
92+
frameBody := make([]byte, size)
93+
if _, err := io.ReadFull(conn, frameBody); err != nil {
94+
return
95+
}
96+
req := &agentv1.RegisterCapabilitiesRequest{}
97+
if err := proto.Unmarshal(frameBody[protocol.HeaderSize:], req); err != nil {
98+
return
99+
}
100+
gzReader, err := gzip.NewReader(bytes.NewReader(req.GetManifestJsonGz()))
101+
if err != nil {
102+
return
103+
}
104+
manifestRaw, err := io.ReadAll(gzReader)
105+
if err != nil {
106+
return
107+
}
108+
var manifest map[string]interface{}
109+
if err := json.Unmarshal(manifestRaw, &manifest); err != nil {
110+
return
111+
}
112+
if _, ok := manifest["provider"]; ok {
113+
gotCapabilities = true
114+
}
115+
// 回确认帧:响应 msgID + 同 reqID + 空 body
116+
respFrame := protocol.NewMessageBody(
117+
protocol.GetResponseMsgID(protocol.MsgRegisterCapabilitiesReq),
118+
protocol.GetMsgID(frameBody[1:4]),
119+
nil,
120+
)
121+
out := make([]byte, 4+len(respFrame))
122+
binary.BigEndian.PutUint32(out[:4], uint32(len(respFrame)))
123+
copy(out[4:], respFrame)
124+
conn.Write(out)
125+
}()
126+
127+
config := ClientConfig{
128+
ControlAddr: listener.Addr().String(),
129+
}
130+
manager, err := NewTCPManager(config, map[string]FunctionHandler{})
131+
if err != nil {
132+
t.Fatalf("NewTCPManager: %v", err)
133+
}
134+
m := manager.(*TCPManager)
135+
m.serviceID = "go-1"
136+
m.functions = []*sdkv1.ProviderFunctionDescriptor{{Id: "f1"}}
137+
138+
m.maybeRegisterCapabilities()
139+
if !gotCapabilities {
140+
t.Fatal("control plane did not receive a valid manifest")
141+
}
142+
}
143+
144+
func TestTCPManager_MaybeRegisterCapabilities_NoControlAddr(t *testing.T) {
145+
manager, err := NewTCPManager(ClientConfig{}, map[string]FunctionHandler{})
146+
if err != nil {
147+
t.Fatalf("NewTCPManager: %v", err)
148+
}
149+
m := manager.(*TCPManager)
150+
// 未配置 ControlAddr:应直接返回(无连接、无 panic)
151+
m.maybeRegisterCapabilities()
152+
if strings.TrimSpace(m.config.ControlAddr) != "" {
153+
t.Fatal("expected empty ControlAddr")
154+
}
155+
}

sdks/go/pkg/croupier/tcp_manager.go

Lines changed: 121 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,8 @@
22
package croupier
33

44
import (
5+
"bytes"
6+
"compress/gzip"
57
"context"
68
"encoding/json"
79
"fmt"
@@ -291,9 +293,128 @@ func (m *TCPManager) RegisterWithAgent(ctx context.Context, serviceID, serviceVe
291293
m.heartbeatStop = cancel
292294
go m.heartbeatLoop(ctx)
293295

296+
// 控制面 manifest 上传(best-effort,不阻断注册结果)
297+
m.maybeRegisterCapabilities()
298+
294299
return m.sessionID, nil
295300
}
296301

302+
// buildManifest 构建供控制面注册的能力清单(provider 元数据 + 函数摘要),
303+
// 与 JS/C# 参考实现同构。
304+
func (m *TCPManager) buildManifest() map[string]interface{} {
305+
provider := map[string]interface{}{
306+
"id": orDefault(m.serviceID, "go-service"),
307+
"version": orDefault(m.serviceVersion, "1.0.0"),
308+
"lang": orDefault(m.config.ProviderLang, "go"),
309+
"sdk": orDefault(m.config.ProviderSDK, "croupier-go-sdk"),
310+
}
311+
functions := make([]map[string]interface{}, 0, len(m.functions))
312+
for _, descriptor := range m.functions {
313+
if descriptor == nil || strings.TrimSpace(descriptor.Id) == "" {
314+
continue
315+
}
316+
entry := map[string]interface{}{
317+
"id": descriptor.Id,
318+
"version": orDefault(descriptor.Version, "1.0.0"),
319+
}
320+
if descriptor.Resource != "" {
321+
entry["resource"] = descriptor.Resource
322+
}
323+
if descriptor.Operation != "" {
324+
entry["operation"] = descriptor.Operation
325+
}
326+
if descriptor.Risk != "" {
327+
entry["risk"] = descriptor.Risk
328+
}
329+
if descriptor.Permission != "" {
330+
entry["permission"] = descriptor.Permission
331+
}
332+
if descriptor.Description != "" {
333+
entry["description"] = descriptor.Description
334+
}
335+
if descriptor.InputSchema != "" {
336+
entry["inputSchema"] = json.RawMessage(descriptor.InputSchema)
337+
}
338+
if descriptor.OutputSchema != "" {
339+
entry["outputSchema"] = json.RawMessage(descriptor.OutputSchema)
340+
}
341+
functions = append(functions, entry)
342+
}
343+
return map[string]interface{}{"provider": provider, "functions": functions}
344+
}
345+
346+
func orDefault(value, fallback string) string {
347+
if strings.TrimSpace(value) == "" {
348+
return fallback
349+
}
350+
return value
351+
}
352+
353+
// gzipBytes 压缩 manifest JSON(RegisterCapabilitiesRequest.manifest_json_gz)。
354+
func gzipBytes(data []byte) ([]byte, error) {
355+
var buf bytes.Buffer
356+
writer := gzip.NewWriter(&buf)
357+
if _, err := writer.Write(data); err != nil {
358+
return nil, err
359+
}
360+
if err := writer.Close(); err != nil {
361+
return nil, err
362+
}
363+
return buf.Bytes(), nil
364+
}
365+
366+
// maybeRegisterCapabilities 向控制面(config.ControlAddr)上传能力清单。
367+
// 独立短连接 + 5s 超时;任何失败仅告警,不影响已完成的函数注册。
368+
func (m *TCPManager) maybeRegisterCapabilities() {
369+
if strings.TrimSpace(m.config.ControlAddr) == "" {
370+
return
371+
}
372+
373+
manifestJSON, err := json.Marshal(m.buildManifest())
374+
if err != nil {
375+
logWarnf("register capabilities: marshal manifest: %v", err)
376+
return
377+
}
378+
manifestGz, err := gzipBytes(manifestJSON)
379+
if err != nil {
380+
logWarnf("register capabilities: gzip manifest: %v", err)
381+
return
382+
}
383+
req := &agentv1.RegisterCapabilitiesRequest{
384+
Provider: &agentv1.ProviderMeta{
385+
Id: orDefault(m.serviceID, "go-service"),
386+
Version: orDefault(m.serviceVersion, "1.0.0"),
387+
Lang: orDefault(m.config.ProviderLang, "go"),
388+
Sdk: orDefault(m.config.ProviderSDK, "croupier-go-sdk"),
389+
},
390+
ManifestJsonGz: manifestGz,
391+
}
392+
reqBody, err := proto.Marshal(req)
393+
if err != nil {
394+
logWarnf("register capabilities: marshal request: %v", err)
395+
return
396+
}
397+
398+
controlClient, err := transport.NewTCPClient(&transport.Config{
399+
Address: m.config.ControlAddr,
400+
Insecure: m.config.Insecure,
401+
DialTimeout: 5 * time.Second,
402+
})
403+
if err != nil {
404+
logWarnf("register capabilities: connect control plane: %v", err)
405+
return
406+
}
407+
defer controlClient.Close()
408+
409+
callCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
410+
defer cancel()
411+
if _, _, err := controlClient.Call(callCtx, protocol.MsgRegisterCapabilitiesReq, reqBody); err != nil {
412+
logWarnf("register capabilities: call: %v", err)
413+
return
414+
}
415+
logInfof("Capabilities registered to control plane: %s", m.config.ControlAddr)
416+
}
417+
297418
func (m *TCPManager) IsConnected() bool {
298419
m.mu.RLock()
299420
defer m.mu.RUnlock()

sdks/python/croupier/__init__.py

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -257,6 +257,43 @@ def connect(self) -> None:
257257
self._connected = True
258258
LOG.info("Client connected with %d functions", len(self._handlers))
259259

260+
# F:控制面 manifest 上传(best-effort,不阻断连接结果)
261+
self._maybe_register_capabilities()
262+
263+
def _maybe_register_capabilities(self) -> None:
264+
"""向控制面(control_addr)上传能力清单;失败仅告警不影响连接。"""
265+
control_addr = getattr(self._config, "control_addr", None)
266+
if not control_addr:
267+
return
268+
try:
269+
register_pb2 = _load_proto_module("croupier.agent.v1.register_pb2")
270+
request = register_pb2.RegisterCapabilitiesRequest(
271+
provider=register_pb2.ProviderMeta(
272+
id=self._config.service_id,
273+
version=self._config.service_version,
274+
lang=self._config.provider_lang,
275+
sdk=self._config.provider_sdk,
276+
),
277+
manifest_json_gz=gzip.compress(self.build_manifest()),
278+
)
279+
transport = TCPTransport(
280+
address=control_addr,
281+
timeout_ms=max(self._config.timeout_seconds, 5) * 1000,
282+
tls_enabled=not self._config.insecure,
283+
tls_cert_file=self._config.cert_file or "",
284+
tls_key_file=self._config.key_file or "",
285+
tls_ca_file=self._config.ca_file or "",
286+
tls_server_name=self._config.server_name or "",
287+
)
288+
transport.connect()
289+
try:
290+
transport.call(protocol.MSG_REGISTER_CAPABILITIES_REQ, request.SerializeToString())
291+
finally:
292+
transport.close()
293+
LOG.info("Capabilities registered to control plane: %s", control_addr)
294+
except Exception as error: # noqa: BLE001 — 上传失败不影响注册
295+
LOG.warning("Failed to register capabilities: %s", error)
296+
260297
def disconnect(self) -> None:
261298
self._heartbeat_stop.set()
262299
if self._heartbeat_thread:

0 commit comments

Comments
 (0)