Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion deploy/stacks/self-managed/global.yaml.gotmpl
Original file line number Diff line number Diff line change
Expand Up @@ -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. */}}
Expand Down Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions deploy/stacks/self-managed/tests/ha-value-wiring.sh
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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" ||
Expand Down
13 changes: 8 additions & 5 deletions docs/self-managed/high-availability.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,17 +16,21 @@
*/
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;
import lombok.extern.slf4j.Slf4j;
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 {
Expand All @@ -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;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
private Duration connectionTimeout = Duration.ZERO;
private Duration pingInterval = Duration.ZERO;
private Duration reconnectWait = Duration.ZERO;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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())
Expand Down Expand Up @@ -166,6 +180,7 @@ private StreamConfiguration streamConfiguration(String name, String subject) {
.retentionPolicy(RetentionPolicy.WorkQueue)
.maxMessages(MAX_MESSAGES)
.maxAge(natsConfigurationProperties.getMessageTtl())
.replicas(natsConfigurationProperties.getReplicas())
.build();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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));

Expand All @@ -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
Expand All @@ -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));

Expand All @@ -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
Expand Down Expand Up @@ -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) {
Expand Down
Loading