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..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,9 @@ 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; 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) {