Skip to content

Commit f106f31

Browse files
committed
fix(sdk): 审查发现修复——schemadiff 枚举方向按数据流判定 + manifest 上传 fire-and-forget
审查发现 #1(设计缺陷):schemadiff 枚举方向标反 - input_schema:收窄 = breaking(旧调用发被删值会被拒)、扩张 = compatible - output_schema:扩张 = breaking(消费方见到新值)、收窄 = compatible - 原实现只按「新增 = breaking」单向判定,测试固化了错误语义;按 source 分方向修正,测试改为四象限断言(input/output × 收窄/扩张) 审查发现 #2(健壮性):manifest 上传阻塞注册主路径 - Go:Register 拆 maybeRegisterCapabilitiesAsync(快照 functions/ serviceID 后台 goroutine,避免与重注册竞态),注册主路径不再被 控制面连接/调用阻塞最长 10s - JS:fire-and-forget + .catch 兜底;Python:守护线程;Java:守护线程; C#:discard await,并把 transport.Connect 移入 fail-open try(顺带 修复 CroupierClientLifecycleTests 文档化的死控制面中止连接 bug, 测试改为断言 fail-open 新语义)
1 parent ca324c5 commit f106f31

13 files changed

Lines changed: 625 additions & 54 deletions

File tree

internal/function/schemadiff/diff.go

Lines changed: 40 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -187,31 +187,61 @@ func diffRequired(source, path string, oldMap, newMap map[string]interface{}, fi
187187
}
188188
}
189189

190+
// diffEnum 枚举方向性判定(审查修正):破坏方向取决于数据流向——
191+
// - input_schema:收窄 = breaking(旧调用方发被删的值会被服务端拒绝),
192+
// 扩张 = compatible(旧调用方不受影响);
193+
// - output_schema:扩张 = breaking(消费方会见到新值),
194+
// 收窄 = compatible(消费方只会见到更少的值)。
190195
func diffEnum(source, path string, oldMap, newMap map[string]interface{}, findings *[]Finding) {
191196
oldEnum, okOld := oldMap["enum"].([]interface{})
192197
newEnum, okNew := newMap["enum"].([]interface{})
193198
if !okOld || !okNew {
194199
return
195200
}
196-
// 枚举收窄 = breaking:新取值域必须是旧取值域的子集
201+
isInput := source == "input_schema"
197202
oldValues := make(map[string]bool, len(oldEnum))
198203
for _, item := range oldEnum {
199204
oldValues[fmt.Sprint(item)] = true
200205
}
201-
newValues := make([]string, 0, len(newEnum))
206+
newValues := make(map[string]bool, len(newEnum))
207+
changed := make([]string, 0)
202208
for _, item := range newEnum {
203209
key := fmt.Sprint(item)
204-
newValues = append(newValues, key)
210+
newValues[key] = true
205211
if !oldValues[key] {
206-
*findings = append(*findings, Finding{
207-
Severity: SeverityBreaking,
208-
Source: source,
209-
Path: path + "/enum/" + key,
210-
Reason: fmt.Sprintf("枚举新增取值 %q", item),
211-
})
212+
changed = append(changed, key)
213+
}
214+
}
215+
for _, item := range oldEnum {
216+
key := fmt.Sprint(item)
217+
if !newValues[key] {
218+
changed = append(changed, key)
219+
}
220+
}
221+
sort.Strings(changed)
222+
for _, key := range changed {
223+
added := newValues[key]
224+
breaking := added != isInput // input:删旧值破坏;output:增新值破坏
225+
severity := SeverityCompatible
226+
reason := fmt.Sprintf("枚举新增取值 %q", key)
227+
if !added {
228+
reason = fmt.Sprintf("枚举删除取值 %q", key)
229+
}
230+
if breaking {
231+
severity = SeverityBreaking
232+
if added {
233+
reason = fmt.Sprintf("枚举新增取值 %q(消费方会见到新值)", key)
234+
} else {
235+
reason = fmt.Sprintf("枚举删除取值 %q(旧调用方发被删值将被拒绝)", key)
236+
}
212237
}
238+
*findings = append(*findings, Finding{
239+
Severity: severity,
240+
Source: source,
241+
Path: path + "/enum/" + key,
242+
Reason: reason,
243+
})
213244
}
214-
sort.Strings(newValues)
215245
}
216246

217247
func requiredSet(node map[string]interface{}) map[string]bool {

internal/function/schemadiff/diff_test.go

Lines changed: 19 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -58,21 +58,28 @@ func TestDiffTypeChanged(t *testing.T) {
5858
}
5959
}
6060

61-
// 4. enum 收窄/扩张 = breaking(新取值不在旧取值域)
62-
func TestDiffEnumNarrowed(t *testing.T) {
61+
// 4. enum 方向性(审查修正):input 收窄 = breaking、扩张 = compatible;
62+
// output 恰好相反。
63+
func TestDiffEnumDirectional(t *testing.T) {
6364
oldRaw := mustRaw(t, `{"type":"object","properties":{"level":{"type":"string","enum":["low","high"]}}}`)
64-
// 收窄(去掉 high)兼容;扩张(加 critical)breaking
65-
compatible := mustRaw(t, `{"type":"object","properties":{"level":{"type":"string","enum":["low"]}}}`)
66-
if findings := DiffSchemas("input_schema", oldRaw, compatible); HasBreaking(findings) {
67-
t.Fatalf("enum narrowing (subset) should be compatible, got %+v", findings)
68-
}
65+
narrowed := mustRaw(t, `{"type":"object","properties":{"level":{"type":"string","enum":["low"]}}}`)
6966
expanded := mustRaw(t, `{"type":"object","properties":{"level":{"type":"string","enum":["low","high","critical"]}}}`)
70-
findings := DiffSchemas("input_schema", oldRaw, expanded)
71-
if !HasBreaking(findings) {
72-
t.Fatalf("expected breaking finding for enum expansion, got %+v", findings)
67+
68+
// input_schema:收窄 = breaking(旧调用发被删值会被拒)
69+
if findings := DiffSchemas("input_schema", oldRaw, narrowed); !HasBreaking(findings) {
70+
t.Fatalf("input enum narrowing should be breaking, got %+v", findings)
71+
}
72+
// input_schema:扩张 = compatible(旧调用方不受影响)
73+
if findings := DiffSchemas("input_schema", oldRaw, expanded); HasBreaking(findings) {
74+
t.Fatalf("input enum expansion should be compatible, got %+v", findings)
75+
}
76+
// output_schema:扩张 = breaking(消费方会见到新值)
77+
if findings := DiffSchemas("output_schema", oldRaw, expanded); !HasBreaking(findings) {
78+
t.Fatalf("output enum expansion should be breaking, got %+v", findings)
7379
}
74-
if _, ok := findByPath(findings, "$/level/enum/critical"); !ok {
75-
t.Fatalf("expected enum/critical finding, got %+v", findings)
80+
// output_schema:收窄 = compatible
81+
if findings := DiffSchemas("output_schema", oldRaw, narrowed); HasBreaking(findings) {
82+
t.Fatalf("output enum narrowing should be compatible, got %+v", findings)
7683
}
7784
}
7885

sdks/csharp/src/Croupier.Sdk.Tests/CroupierClientLifecycleTests.cs

Lines changed: 11 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -223,12 +223,10 @@ public async Task ConnectAsync_WithControlAddr_RegistersCapabilitiesWithGzippedM
223223
}
224224

225225
[Fact]
226-
public async Task ConnectAsync_WithDeadControlAddr_FailsProviderConnect_DocumentsBug()
226+
public async Task ConnectAsync_WithDeadControlAddr_IsFailOpen()
227227
{
228-
// Documents a bug: in RegisterCapabilitiesAsync, transport.Connect()
229-
// runs OUTSIDE the try/catch that guards CallAsync. A dead ControlAddr
230-
// therefore aborts the whole provider connect, even though capability
231-
// registration is supposed to be best-effort (warning + continue).
228+
// 审查发现 #2 修复后:manifest 上传 fire-and-forget 且整体 fail-open
229+
// (原实现在 try 外 Connect,死控制面会中止整个 provider connect)。
232230
var probe = new System.Net.Sockets.TcpListener(System.Net.IPAddress.Loopback, 0);
233231
probe.Start();
234232
var deadPort = ((System.Net.IPEndPoint)probe.LocalEndpoint).Port;
@@ -241,10 +239,8 @@ public async Task ConnectAsync_WithDeadControlAddr_FailsProviderConnect_Document
241239
}));
242240
client.RegisterFunction(Descriptor("fn.cap-dead"), (ctx, payload) => Task.FromResult("{}"));
243241

244-
var action = () => client.ConnectAsync();
245-
246-
await action.Should().ThrowAsync<Exception>();
247-
client.IsConnected.Should().BeFalse();
242+
await client.ConnectAsync();
243+
client.IsConnected.Should().BeTrue();
248244
}
249245

250246
[Fact]
@@ -259,6 +255,12 @@ public async Task ConnectAsync_WithAliveControlAddr_ConnectSucceeds()
259255
await client.ConnectAsync();
260256

261257
client.IsConnected.Should().BeTrue();
258+
// fire-and-forget:能力帧异步到达,轮询等待
259+
var deadline = DateTime.UtcNow.AddSeconds(5);
260+
while (control.CapabilityRequests.Count == 0 && DateTime.UtcNow < deadline)
261+
{
262+
await Task.Delay(50);
263+
}
262264
control.CapabilityRequests.Should().ContainSingle();
263265
}
264266

sdks/csharp/src/Croupier.Sdk/CroupierClient.cs

Lines changed: 6 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -677,7 +677,9 @@ private async Task ConnectAndRegisterAsync(CancellationToken cancellationToken)
677677
// Register inbound request handler for InvokeRequest from Agent
678678
transport.SetInboundRequestHandler(HandleInboundRequestAsync);
679679

680-
await RegisterCapabilitiesAsync(cancellationToken);
680+
// 审查发现 #2:manifest 上传 fire-and-forget——控制面慢/不可达
681+
// 不得拖慢注册主路径(方法内部已 fail-open)。
682+
_ = RegisterCapabilitiesAsync(cancellationToken);
681683
}
682684
catch
683685
{
@@ -739,10 +741,12 @@ private async Task RegisterCapabilitiesAsync(CancellationToken cancellationToken
739741
? _config.ControlAddr["tcp://".Length..]
740742
: _config.ControlAddr;
741743
using var transport = _transportFactory(address, _config.TimeoutSeconds * 1000, _config.ConnectTimeoutSeconds * 1000, _logger);
742-
transport.Connect();
743744

744745
try
745746
{
747+
// 审查发现 #2:连接也在 fail-open 范围内(原实现在 try 外,
748+
// 死控制面会中止整个 provider connect)
749+
transport.Connect();
746750
await transport.CallAsync(
747751
Protocol.MsgRegisterCapabilitiesReq,
748752
BuildRegisterCapabilitiesRequestData(),

sdks/go/pkg/croupier/manifest_upload_test.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@ func TestTCPManager_BuildManifest(t *testing.T) {
4444
{Id: " "}, // 空 id 跳过
4545
}
4646

47-
manifest := m.buildManifest()
47+
manifest := m.buildManifest("go-1", "2.0.0", m.functions)
4848
provider, ok := manifest["provider"].(map[string]interface{})
4949
if !ok {
5050
t.Fatalf("provider missing: %+v", manifest)
@@ -135,7 +135,7 @@ func TestTCPManager_MaybeRegisterCapabilities(t *testing.T) {
135135
m.serviceID = "go-1"
136136
m.functions = []*sdkv1.ProviderFunctionDescriptor{{Id: "f1"}}
137137

138-
m.maybeRegisterCapabilities()
138+
m.maybeRegisterCapabilities("go-1", "2.0.0", m.functions)
139139
if !gotCapabilities {
140140
t.Fatal("control plane did not receive a valid manifest")
141141
}
@@ -148,7 +148,7 @@ func TestTCPManager_MaybeRegisterCapabilities_NoControlAddr(t *testing.T) {
148148
}
149149
m := manager.(*TCPManager)
150150
// 未配置 ControlAddr:应直接返回(无连接、无 panic)
151-
m.maybeRegisterCapabilities()
151+
m.maybeRegisterCapabilities("go-1", "2.0.0", m.functions)
152152
if strings.TrimSpace(m.config.ControlAddr) != "" {
153153
t.Fatal("expected empty ControlAddr")
154154
}

sdks/go/pkg/croupier/tcp_manager.go

Lines changed: 31 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -294,23 +294,26 @@ func (m *TCPManager) RegisterWithAgent(ctx context.Context, serviceID, serviceVe
294294
m.heartbeatStop = cancel
295295
go m.heartbeatLoop(ctx)
296296

297-
// 控制面 manifest 上传(best-effort,不阻断注册结果)
298-
m.maybeRegisterCapabilities()
297+
// 控制面 manifest 上传(审查发现 #2):fire-and-forget——控制面慢/
298+
// 不可达时不得拖慢注册主路径(同步执行最长 10s)。参数全部取
299+
// 快照传值,避免与后续重注册竞态。
300+
m.maybeRegisterCapabilitiesAsync()
299301

300302
return m.sessionID, nil
301303
}
302304

303305
// buildManifest 构建供控制面注册的能力清单(provider 元数据 + 函数摘要),
304306
// 与 JS/C# 参考实现同构。
305-
func (m *TCPManager) buildManifest() map[string]interface{} {
307+
func (m *TCPManager) buildManifest(serviceID, serviceVersion string,
308+
functions []*sdkv1.ProviderFunctionDescriptor) map[string]interface{} {
306309
provider := map[string]interface{}{
307-
"id": orDefault(m.serviceID, "go-service"),
308-
"version": orDefault(m.serviceVersion, "1.0.0"),
310+
"id": orDefault(serviceID, "go-service"),
311+
"version": orDefault(serviceVersion, "1.0.0"),
309312
"lang": orDefault(m.config.ProviderLang, "go"),
310313
"sdk": orDefault(m.config.ProviderSDK, "croupier-go-sdk"),
311314
}
312-
functions := make([]map[string]interface{}, 0, len(m.functions))
313-
for _, descriptor := range m.functions {
315+
entries := make([]map[string]interface{}, 0, len(functions))
316+
for _, descriptor := range functions {
314317
if descriptor == nil || strings.TrimSpace(descriptor.Id) == "" {
315318
continue
316319
}
@@ -339,9 +342,9 @@ func (m *TCPManager) buildManifest() map[string]interface{} {
339342
if descriptor.OutputSchema != "" {
340343
entry["outputSchema"] = json.RawMessage(descriptor.OutputSchema)
341344
}
342-
functions = append(functions, entry)
345+
entries = append(entries, entry)
343346
}
344-
return map[string]interface{}{"provider": provider, "functions": functions}
347+
return map[string]interface{}{"provider": provider, "functions": entries}
345348
}
346349

347350
func orDefault(value, fallback string) string {
@@ -364,14 +367,31 @@ func gzipBytes(data []byte) ([]byte, error) {
364367
return buf.Bytes(), nil
365368
}
366369

370+
// maybeRegisterCapabilitiesAsync 异步 fire-and-forget:先快照上传所需
371+
// 状态(functions 会在重注册时整体替换),再交后台执行——注册主路径
372+
// 不再被控制面连接/调用阻塞(审查发现 #2)。
373+
func (m *TCPManager) maybeRegisterCapabilitiesAsync() {
374+
if strings.TrimSpace(m.config.ControlAddr) == "" {
375+
return
376+
}
377+
m.mu.RLock()
378+
serviceID, serviceVersion := m.serviceID, m.serviceVersion
379+
functions := make([]*sdkv1.ProviderFunctionDescriptor, len(m.functions))
380+
copy(functions, m.functions)
381+
m.mu.RUnlock()
382+
383+
go m.maybeRegisterCapabilities(serviceID, serviceVersion, functions)
384+
}
385+
367386
// maybeRegisterCapabilities 向控制面(config.ControlAddr)上传能力清单。
368387
// 独立短连接 + 5s 超时;任何失败仅告警,不影响已完成的函数注册。
369-
func (m *TCPManager) maybeRegisterCapabilities() {
388+
func (m *TCPManager) maybeRegisterCapabilities(serviceID, serviceVersion string,
389+
functions []*sdkv1.ProviderFunctionDescriptor) {
370390
if strings.TrimSpace(m.config.ControlAddr) == "" {
371391
return
372392
}
373393

374-
manifestJSON, err := json.Marshal(m.buildManifest())
394+
manifestJSON, err := json.Marshal(m.buildManifest(serviceID, serviceVersion, functions))
375395
if err != nil {
376396
logWarnf("register capabilities: marshal manifest: %v", err)
377397
return

sdks/java/src/main/java/io/github/cuihairu/croupier/sdk/CroupierClientImpl.java

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -146,8 +146,11 @@ public CompletableFuture<Void> connect() {
146146

147147
logger.info("Successfully connected");
148148

149-
// F:控制面 manifest 上传(best-effort,不阻断连接结果)
150-
maybeRegisterCapabilities();
149+
// F/审查发现 #2:控制面 manifest 上传 fire-and-forget——
150+
// 控制面慢/不可达不得阻塞注册主路径(方法内部已 fail-open)
151+
Thread capabilitiesThread = new Thread(this::maybeRegisterCapabilities);
152+
capabilitiesThread.setDaemon(true);
153+
capabilitiesThread.start();
151154
} catch (Exception e) {
152155
connected.set(false);
153156
sessionId = "";

0 commit comments

Comments
 (0)