From ef646452bc89cbc6f2576c52891544610c6aea45 Mon Sep 17 00:00:00 2001 From: shobham Date: Wed, 30 Sep 2026 10:21:18 +0530 Subject: [PATCH 1/2] fix(icms): create NVCA NATS streams with the HA replica count icms created CreateNvcaFunctionTaskStream and TerminateNvcaStream without a replica count, so they stayed at one copy even when the self-managed stack runs a 3-server NATS cluster under highAvailability.mode. Add icms.nats.replicas (default 1), apply it when creating both streams, and raise the replica count of an existing stream that is below it. The stack now sets ICMS_NATS_REPLICAS from the same helper as nvcf-api and invocation-service. Signed-off-by: shobham --- deploy/stacks/self-managed/global.yaml.gotmpl | 5 +- .../self-managed/tests/ha-value-wiring.sh | 5 ++ docs/self-managed/high-availability.md | 13 ++++-- .../nats/NatsConfigurationProperties.java | 1 + .../icms/outbound/nats/NatsStreamManager.java | 19 +++++++- ...atsMessageSenderClientIntegrationTest.java | 1 + .../NatsStreamManagerIntegrationTest.java | 1 + .../outbound/nats/NatsStreamManagerTest.java | 46 ++++++++++++++++++- 8 files changed, 82 insertions(+), 9 deletions(-) diff --git a/deploy/stacks/self-managed/global.yaml.gotmpl b/deploy/stacks/self-managed/global.yaml.gotmpl index 6e5759b951..0aac320f52 100644 --- a/deploy/stacks/self-managed/global.yaml.gotmpl +++ b/deploy/stacks/self-managed/global.yaml.gotmpl @@ -100,7 +100,7 @@ merge: {{- end -}} {{- end -}} -{{/* JetStream stream replicas for the stream creators (nvcf-api, invocation), +{{/* JetStream stream replicas for the stream creators (nvcf-api, invocation, icms), derived from the NATS server count and capped at 3 as NATS recommends. Empty for a single server so the services keep their own default. Override per service with api.env.NVCF_NATS_REPLICAS or invocation.env.NATS_PROPERTIES__REPLICAS. */}} @@ -1336,6 +1336,9 @@ sis: {{- if .Values.global.observability.tracing.enabled }} MANAGEMENT_OTLP_TRACING_ENDPOINT: "{{ .Values.global.observability.tracing.collectorProtocol }}://{{ .Values.global.observability.tracing.collectorEndpoint }}:{{ .Values.global.observability.tracing.collectorPort }}/v1/traces" {{- end }} + {{- with include "nvcf.natsStreamReplicas" .Values }} + ICMS_NATS_REPLICAS: {{ . | quote }} + {{- end }} adminIssuerProxy: fullnameOverride: admin-token-issuer-proxy diff --git a/deploy/stacks/self-managed/tests/ha-value-wiring.sh b/deploy/stacks/self-managed/tests/ha-value-wiring.sh index 4681b8558a..23fa227d97 100755 --- a/deploy/stacks/self-managed/tests/ha-value-wiring.sh +++ b/deploy/stacks/self-managed/tests/ha-value-wiring.sh @@ -100,6 +100,9 @@ fi if awk '/^api:/{p=1;next} /^[a-zA-Z]/{p=0} p' "$work_dir/api-off.yaml" | grep -q "NVCF_NATS_REPLICAS:"; then fail "api: JetStream RF env leaked while highAvailability.mode=none" fi +if awk '/^sis:/{p=1;next} /^[a-zA-Z]/{p=0} p' "$work_dir/api-off.yaml" | grep -q "ICMS_NATS_REPLICAS:"; then + fail "icms (sis release): JetStream RF env leaked while highAvailability.mode=none" +fi render_chart_values invocation-service "$work_dir/invocation-off.yaml" "$core" || fail "render invocation (ha none)" if grep -q "NATS_PROPERTIES__REPLICAS:" "$work_dir/invocation-off.yaml"; then fail "invocation: JetStream RF env leaked while highAvailability.mode=none" @@ -139,6 +142,8 @@ awk '/^api:/{p=1;next} /^[a-zA-Z]/{p=0} p' "$work_dir/api-on.yaml" | grep -q "po fail "api: expected podDisruptionBudget when highAvailability.mode=preferred" grep -q 'NVCF_NATS_REPLICAS: "3"' "$work_dir/api-on.yaml" || fail "api: expected JetStream RF NVCF_NATS_REPLICAS=3 when highAvailability.mode=preferred" +awk '/^sis:/{p=1;next} /^[a-zA-Z]/{p=0} p' "$work_dir/api-on.yaml" | grep -q 'ICMS_NATS_REPLICAS: "3"' || + fail "icms (sis release): expected JetStream RF ICMS_NATS_REPLICAS=3 when highAvailability.mode=preferred" render_chart_values invocation-service "$work_dir/invocation-on.yaml" "$core" || fail "render invocation (preferred)" grep -q 'NATS_PROPERTIES__REPLICAS: "3"' "$work_dir/invocation-on.yaml" || diff --git a/docs/self-managed/high-availability.md b/docs/self-managed/high-availability.md index 587d7535d3..0796297361 100644 --- a/docs/self-managed/high-availability.md +++ b/docs/self-managed/high-availability.md @@ -281,17 +281,20 @@ Beyond placement, HA also raises the data-durability settings: Streams default to a single replica. When NATS runs more than one server, the stack derives the JetStream replica factor (RF) from the server count, capped at **3** as the NATS documentation recommends, so the HA cluster of 3 servers -gets RF=3. It is set on the two services that create streams: `nvcf-api` (via -`NVCF_NATS_REPLICAS`) and `invocation-service` (via `NATS_PROPERTIES__REPLICAS`). -To override it for one service, set that variable in `api.env` or -`invocation.env`. JetStream streams use Raft quorum: RF=3 tolerates the loss of +gets RF=3. It is set on the three services that create streams: `nvcf-api` (via +`NVCF_NATS_REPLICAS`), `invocation-service` (via `NATS_PROPERTIES__REPLICAS`), and +`icms` (the `sis` release, via `ICMS_NATS_REPLICAS`, for `CreateNvcaFunctionTaskStream` and +`TerminateNvcaStream`). To override it for nvcf-api or invocation-service, set +that variable in `api.env` or `invocation.env`. JetStream streams use Raft quorum: RF=3 tolerates the loss of one replica. **RF=2 is not sufficient** — a 2-member Raft group loses quorum the moment either replica is unavailable, so it provides no resilience benefit over RF=1. RF applies when a stream is **created**. Streams already created at RF=1 are not rewritten by changing this value — after enabling HA, recreate or edit those -streams (for example `nats stream edit`) to raise their replica factor. +streams (for example `nats stream edit`) to raise their replica factor. The +exception is `icms`, which raises its two streams to the configured RF when it +starts and on its periodic stream validation. #### Cassandra replication and consistency diff --git a/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/configuration/nats/NatsConfigurationProperties.java b/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/configuration/nats/NatsConfigurationProperties.java index 708e783480..87be74c220 100644 --- a/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/configuration/nats/NatsConfigurationProperties.java +++ b/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/configuration/nats/NatsConfigurationProperties.java @@ -35,6 +35,7 @@ public class NatsConfigurationProperties { private String natsUrl; private int maxPoolSize = 8; private boolean createNatsStreams; + private int replicas = 1; private Duration connectionTimeout = Duration.ZERO; private Duration pingInterval = Duration.ZERO; private Duration reconnectWait = Duration.ZERO; diff --git a/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/outbound/nats/NatsStreamManager.java b/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/outbound/nats/NatsStreamManager.java index e377ad3598..32b5d64e0e 100644 --- a/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/outbound/nats/NatsStreamManager.java +++ b/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/outbound/nats/NatsStreamManager.java @@ -95,7 +95,7 @@ public void createStream(StreamConfiguration streamConfig) throws IOException, JetStreamApiException { try { var streamInfo = jetStreamManagement.getStreamInfo(streamConfig.getName()); - validateStreamConfiguration(streamConfig, streamInfo.getConfiguration()); + reconcileExistingStream(streamConfig, streamInfo.getConfiguration()); return; } catch (JetStreamApiException e) { // non-404 related error gets passed back up @@ -111,13 +111,27 @@ public void createStream(StreamConfiguration streamConfig) } catch (JetStreamApiException e) { if (e.getApiErrorCode() == 10058) { var streamInfo = jetStreamManagement.getStreamInfo(streamConfig.getName()); - validateStreamConfiguration(streamConfig, streamInfo.getConfiguration()); + reconcileExistingStream(streamConfig, streamInfo.getConfiguration()); return; } throw e; } } + private void reconcileExistingStream(StreamConfiguration expected, StreamConfiguration actual) + throws IOException, JetStreamApiException { + validateStreamConfiguration(expected, actual); + // Replicas are only raised: a stream created under HA keeps its copies if the + // configured count later drops. + if (actual.getReplicas() < expected.getReplicas()) { + log.info("Raising NATS stream {} replicas from {} to {}", expected.getName(), + actual.getReplicas(), expected.getReplicas()); + jetStreamManagement.updateStream(StreamConfiguration.builder(actual) + .replicas(expected.getReplicas()) + .build()); + } + } + private static void validateStreamConfiguration( StreamConfiguration expected, StreamConfiguration actual) { if (!expected.getSubjects().equals(actual.getSubjects()) @@ -166,6 +180,7 @@ private StreamConfiguration streamConfiguration(String name, String subject) { .retentionPolicy(RetentionPolicy.WorkQueue) .maxMessages(MAX_MESSAGES) .maxAge(natsConfigurationProperties.getMessageTtl()) + .replicas(natsConfigurationProperties.getReplicas()) .build(); } diff --git a/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsMessageSenderClientIntegrationTest.java b/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsMessageSenderClientIntegrationTest.java index bf8db08300..d84c6bef18 100644 --- a/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsMessageSenderClientIntegrationTest.java +++ b/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsMessageSenderClientIntegrationTest.java @@ -70,6 +70,7 @@ void setUp() { when(natsConfigurationProperties.getNkeySeed()).thenReturn(Optional.empty()); when(natsConfigurationProperties.getDelayBetweenMessages()).thenReturn(Duration.ZERO); when(natsConfigurationProperties.isCreateNatsStreams()).thenReturn(true); + when(natsConfigurationProperties.getReplicas()).thenReturn(1); when(natsConfigurationProperties.isEnabled()).thenReturn(true); when(natsConfigurationProperties.getMaxPoolSize()).thenReturn(1); diff --git a/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsStreamManagerIntegrationTest.java b/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsStreamManagerIntegrationTest.java index 8235c09f7f..ec58c52a81 100644 --- a/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsStreamManagerIntegrationTest.java +++ b/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsStreamManagerIntegrationTest.java @@ -58,6 +58,7 @@ void setUp() { when(natsConfigurationProperties.getReconnectJitter()).thenReturn(Duration.ZERO); when(natsConfigurationProperties.getNkeySeed()).thenReturn(Optional.empty()); when(natsConfigurationProperties.getMessageTtl()).thenReturn(Duration.ofHours(24)); + when(natsConfigurationProperties.getReplicas()).thenReturn(1); when(natsConfigurationProperties.isEnabled()).thenReturn(true); when(natsConfigurationProperties.getMaxPoolSize()).thenReturn(1); diff --git a/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsStreamManagerTest.java b/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsStreamManagerTest.java index 5089ae4dbc..6e81c38e28 100644 --- a/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsStreamManagerTest.java +++ b/src/control-plane-services/instance-cluster-management/icms-core/src/test/java/com/nvidia/icms/outbound/nats/NatsStreamManagerTest.java @@ -63,6 +63,7 @@ void setUp() throws Exception { @Test void validateNatsStreams_createsNvcaStreamsWithExistingConfiguration() throws Exception { when(natsConfigurationProperties.getMessageTtl()).thenReturn(Duration.ofHours(24)); + when(natsConfigurationProperties.getReplicas()).thenReturn(3); when(management.getStreamInfo(any())) .thenThrow(apiException(Status.NOT_FOUND_CODE)); @@ -76,6 +77,8 @@ void validateNatsStreams_createsNvcaStreamsWithExistingConfiguration() throws Ex "Create.NVCA.>"); assertStream(streams.get(1), NatsStreamManager.TERMINATE_NVCA_STREAM_NAME, "Terminate.NVCA.>"); + assertEquals(3, streams.get(0).getReplicas()); + assertEquals(3, streams.get(1).getReplicas()); } @Test @@ -102,6 +105,7 @@ void init_createsBothStreams() throws Exception { when(natsConfigurationProperties.isEnabled()).thenReturn(true); when(natsConfigurationProperties.isCreateNatsStreams()).thenReturn(true); when(natsConfigurationProperties.getMessageTtl()).thenReturn(Duration.ofHours(24)); + when(natsConfigurationProperties.getReplicas()).thenReturn(1); when(management.getStreamInfo(any())) .thenThrow(apiException(Status.NOT_FOUND_CODE)); @@ -120,6 +124,38 @@ void createStream_doesNotModifyExistingStream() throws Exception { natsStreamManager.createStream(configuration); verify(management, never()).addStream(configuration); + verify(management, never()).updateStream(any()); + } + + @Test + void createStream_raisesReplicasOfExistingStream() throws Exception { + var configuration = streamConfiguration(3); + var existing = StreamConfiguration.builder(streamConfiguration(1)) + .maxBytes(1024) + .build(); + var streamInfo = org.mockito.Mockito.mock(StreamInfo.class); + when(streamInfo.getConfiguration()).thenReturn(existing); + when(management.getStreamInfo(configuration.getName())).thenReturn(streamInfo); + + natsStreamManager.createStream(configuration); + + var captor = ArgumentCaptor.forClass(StreamConfiguration.class); + verify(management).updateStream(captor.capture()); + assertEquals(3, captor.getValue().getReplicas()); + assertEquals(1024, captor.getValue().getMaxBytes()); + verify(management, never()).addStream(any()); + } + + @Test + void createStream_doesNotLowerReplicasOfExistingStream() throws Exception { + var configuration = streamConfiguration(1); + var streamInfo = org.mockito.Mockito.mock(StreamInfo.class); + when(streamInfo.getConfiguration()).thenReturn(streamConfiguration(3)); + when(management.getStreamInfo(configuration.getName())).thenReturn(streamInfo); + + natsStreamManager.createStream(configuration); + + verify(management, never()).updateStream(any()); } @Test @@ -171,7 +207,15 @@ private static void assertStream( } private static StreamConfiguration streamConfiguration() { - return StreamConfiguration.builder().name("stream").subjects("subject.>").build(); + return streamConfiguration(1); + } + + private static StreamConfiguration streamConfiguration(int replicas) { + return StreamConfiguration.builder() + .name("stream") + .subjects("subject.>") + .replicas(replicas) + .build(); } private static JetStreamApiException apiException(int statusCode) { From 3cf6c0f5be4556ab5ed42079d35c3af6a5a51fc1 Mon Sep 17 00:00:00 2001 From: shobham Date: Wed, 30 Sep 2026 11:00:50 +0530 Subject: [PATCH 2/2] fix(icms): reject an out-of-range NATS stream replica count at startup jnats only accepts 1 to 5 replicas. Validate icms.nats.replicas when it is bound so a bad ICMS_NATS_REPLICAS fails startup immediately instead of after the stream init retries. Signed-off-by: shobham --- .../configuration/nats/NatsConfigurationProperties.java | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/configuration/nats/NatsConfigurationProperties.java b/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/configuration/nats/NatsConfigurationProperties.java index 87be74c220..cbebf2a588 100644 --- a/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/configuration/nats/NatsConfigurationProperties.java +++ b/src/control-plane-services/instance-cluster-management/icms-core/src/main/java/com/nvidia/icms/configuration/nats/NatsConfigurationProperties.java @@ -16,6 +16,8 @@ */ package com.nvidia.icms.configuration.nats; +import jakarta.validation.constraints.Max; +import jakarta.validation.constraints.Min; import java.time.Duration; import java.util.Optional; import lombok.Data; @@ -23,10 +25,12 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.cloud.context.config.annotation.RefreshScope; import org.springframework.context.annotation.Configuration; +import org.springframework.validation.annotation.Validated; @RefreshScope @Configuration @ConfigurationProperties(prefix = "icms.nats") +@Validated @Data @Slf4j public class NatsConfigurationProperties { @@ -35,6 +39,8 @@ public class NatsConfigurationProperties { private String natsUrl; private int maxPoolSize = 8; private boolean createNatsStreams; + @Min(1) + @Max(5) private int replicas = 1; private Duration connectionTimeout = Duration.ZERO; private Duration pingInterval = Duration.ZERO;