diff --git a/internal/assignmentqueue/jetstream.go b/internal/assignmentqueue/jetstream.go index 80de9d5..67ec0d2 100644 --- a/internal/assignmentqueue/jetstream.go +++ b/internal/assignmentqueue/jetstream.go @@ -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, diff --git a/internal/assignmentqueue/jetstream_test.go b/internal/assignmentqueue/jetstream_test.go index 9c1490f..74a6857 100644 --- a/internal/assignmentqueue/jetstream_test.go +++ b/internal/assignmentqueue/jetstream_test.go @@ -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) } }