Compare commits

...
Author SHA1 Message Date
panxiao81 1f9577c816 fix: 使用合法的 JetStream consumer 名称
test / shell (pull_request) Successful in 52s
test / python (pull_request) Successful in 1m24s
test / go (pull_request) Successful in 3m57s
2026-09-25 18:41:47 +00:00
panxiao81 6ea78d4704 Merge pull request '重构 workload class 与 placement driver' (#54) from refactor/workload-placement into main
test / shell (push) Successful in 1m7s
test / python (push) Successful in 3m3s
publish controller image / publish-controller (push) Successful in 5m44s
test / go (push) Successful in 9m7s
Reviewed-on: #54
2026-09-25 17:11:01 +00:00
2 changed files with 7 additions and 3 deletions
+6 -2
View File
@@ -186,9 +186,13 @@ func OpenConsumer(ctx context.Context, manager consumerManager, stream, subjectB
return nil, err
}
key := placement.Key()
// JetStream consumer names may not contain dots even though subjects do.
// Keep the placement subject hierarchical while using a stable, legal name
// for the durable cursor.
consumerName := strings.ReplaceAll(key, ".", "-")
consumer, err := manager.CreateOrUpdateConsumer(ctx, stream, jetstream.ConsumerConfig{
Name: key,
Durable: key,
Name: consumerName,
Durable: consumerName,
FilterSubject: strings.TrimSuffix(subjectBase, ".") + "." + key,
AckPolicy: jetstream.AckExplicitPolicy,
AckWait: 5 * time.Minute,
+1 -1
View File
@@ -215,7 +215,7 @@ func TestOpenConsumerUsesIndependentDurablePerPlacement(t *testing.T) {
if _, err := OpenConsumer(context.Background(), manager, "CI_RUNNER", "ci.assignment", taskassignment.KubernetesContainer, 4); err != nil {
t.Fatal(err)
}
if manager.config.Durable != "container.kubernetes" || manager.config.FilterSubject != "ci.assignment.container.kubernetes" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 || manager.config.MaxDeliver != 1000 {
if manager.config.Name != "container-kubernetes" || manager.config.Durable != "container-kubernetes" || manager.config.FilterSubject != "ci.assignment.container.kubernetes" || manager.config.AckPolicy != jetstream.AckExplicitPolicy || manager.config.MaxAckPending != 4 || manager.config.MaxDeliver != 1000 {
t.Fatalf("config = %#v", manager.config)
}
}