Skip to main content

Fault Tolerance

Fault Tolerance

This page describes how the CM Marketplace platform achieves fault tolerance — the ability to continue operating correctly when individual components fail. For infrastructure-level HA (database failover, backup config, monitoring alerts), see Redundancy & High Availability.


1. Retry Patterns​

Polly — Exponential Backoff (Livechat)​

HTTP calls to Salesforce Partner Messaging use Polly with configurable retry and exponential backoff:

// src/app-livechat/.../Services/AsyncRetryService.cs

public AsyncRetryPolicy<Response<object>> GetRetryPolicyAsync()
{
TryParse(_pollyMaxRetries, out var maxRetries);

return Policy<Response<object>>
.Handle<HttpRequestException>(
ex => ex.StatusCode != HttpStatusCode.OK)
.WaitAndRetryAsync(
maxRetries,
attemptNumber => TimeSpan.FromSeconds(attemptNumber * 2),
(exception, sleepDuration, attemptNumber, context) =>
{
logger.LogInformation(
"Retrying in {SleepDuration}. {AttemptNumber} / {MaxRetries}",
sleepDuration, attemptNumber, maxRetries);
});
}

Retry schedule: 2s, 4s, 6s, 8s... (linear * 2)

Polly — Infinite Retry (SSE Connections)​

Background SSE connections to Salesforce retry forever with fixed 5-second intervals:

// src/app-bgservice/.../Services/SalesforceSseBackgroundService.cs

var retryPolicy = Policy
.Handle<Exception>()
.WaitAndRetryForeverAsync(
_ => TimeSpan.FromSeconds(5),
(exception, retryCount, timeSpan) =>
{
logger.LogInformation(
"SSE retry {RetryCount} for adapter {AdapterId}: {Message}",
retryCount, tenantAdapterId, exception.Message);
});

Pub/Sub — Managed Retry with Dead Letter​

Every subscription has exponential backoff retry and a dead-letter queue:

# infra/environments/production/pubsub.tf

retry_policy = {
minimum_backoff = "10s"
maximum_backoff = "600s" # 10 minutes max
}
dead_letter_policy = {
dead_letter_topic = "projects/.../topics/Events-error"
max_delivery_attempts = 5 # DLQ after 5 failures
}

Cloud Tasks — Exponential Backoff with Doublings​

# infra/environments/production/cloudtask.tf

retry_config {
max_attempts = 3
min_backoff = "3s"
max_backoff = "3600s" # 1 hour max
max_doublings = 16 # 3s -> 6s -> 12s -> ... -> 3600s
}

Workflow Batch Retry​

Data sync workflows retry individual batch failures without blocking the pipeline:

# infra/modules/workflow/templates/productfeed.yaml

retry:
predicate: ${custom_predicate} # Retry on HTTP 4xx/5xx
max_retries: 2
backoff:
initial_delay: 2
max_delay: 10
multiplier: 2
except:
steps:
- track_failure:
assign:
- failedBatches: ${list.concat(failedBatches, [v])}

Failed batches are tracked and reported in the finish-sync step — remaining batches continue processing.


2. Graceful Shutdown​

SSE State Persistence on Shutdown​

When the background service container shuts down (deployment, scaling, restart), all active SSE connections persist their state to Redis:

// src/app-bgservice/.../Services/SalesforceSseBackgroundService.cs

public override async Task StopAsync(CancellationToken cancellationToken)
{
logger.LogInformation(
"StopAsync called - persisting all active SSE connection states");
await PersistAllConnectionStatesAsync();
await base.StopAsync(cancellationToken);
}

private async Task PersistAllConnectionStatesAsync()
{
var tasks = _activeConnections.Values
.Where(state => !string.IsNullOrEmpty(state.CurrentEventId))
.Select(async state =>
{
await PersistLastEventIdAsync(
state.CurrentEventId,
state.EsDeveloperName,
state.OrganizationId,
state.TenantAdapterSession);
});
await Task.WhenAll(tasks);
}

On restart, each SSE connection reads the last event ID from cache and sends it as Last-Event-Id header — the Salesforce stream resumes from exactly where it left off.

SSE Error Recovery​

If an SSE connection drops due to an error (not shutdown), state is persisted before the retry loop reconnects:

catch (Exception ex)
{
logger.LogError(ex,
"SSE connection error for adapter {AdapterId}. " +
"Persisting lastEventId for recovery.", tenantAdapterId);
await PersistLastEventIdAsync(
connectionState.CurrentEventId, esDeveloperName,
organizationId, tenantAdapterSession);
throw; // Polly catches this and retries
}

3. Dead-Letter Queue (DLQ) Architecture​

Service publishes to "Events" topic
|
v
Pub/Sub routes to subscription (via attribute filter)
|
v
Push delivers to Cloud Run endpoint
|
+-- Success (2xx) --> Message acknowledged
|
+-- Failure --> Retry with exponential backoff (10s - 600s)
|
+-- After 5 failures --> Route to "Events-error" topic
|
v
"dead-letter" subscription
(for manual investigation)

All 11 push subscriptions share this pattern. The dead-letter subscription allows the team to inspect, replay, or discard failed messages.


4. Cache Fault Tolerance​

The caching layer is designed so that cache failures never cause application failures:

// src/lib-shared/Cache/CacheProvider.cs

public async Task<T?> GetCacheAsync<T>(string key) where T : class
{
var result = await _cache.GetStringAsync(searchKey);
if (string.IsNullOrEmpty(result))
{
_logger.LogInformation("Cache miss for Key: {@Key}", searchKey);
return null; // Caller falls back to database
}
return JsonSerializer.Deserialize<T>(result);
}
  • Cache miss returns null — callers read from the database instead
  • Cache unavailability degrades performance, not correctness
  • All cache operations are wrapped with logging for observability

5. Token Expiry Handling​

SSE connections automatically detect expired tokens and re-authenticate without losing event position:

// src/app-bgservice/.../Services/SalesforceSseBackgroundService.cs

if (IsTokenExpired(tenantAdapterSession.TokenExpiry))
{
logger.LogInformation("Token expired for adapter {Id}. Re-authenticating...",
tenantAdapterId);
await PersistLastEventIdAndUpdateAccessTokenAsync(
connectionState.CurrentEventId, esDeveloperName,
organizationId, scrtUrl);
isTokenExpired = true;
break; // Exits SSE loop, retry policy reconnects with new token
}

The new access token is cached alongside the last event ID, so the reconnection resumes seamlessly.


6. Thread-Safe Connection Management​

Multiple SSE connections run concurrently using a ConcurrentDictionary:

private readonly ConcurrentDictionary<string, SalesforceSseConnectionState>
_activeConnections = new();

// Key format: "{organizationId}_{esDeveloperName}"
_activeConnections[connectionKey] = connectionState;

// Cleanup in finally block
_activeConnections.TryRemove(connectionKey, out _);

This ensures that connection add/remove/lookup operations are safe even when multiple tenants connect and disconnect simultaneously.


7. Deployment Fault Tolerance​

Cloud Run provides deployment-level fault tolerance:

  • Revision-based: Each deploy creates a new revision. 100% traffic shifts only after health checks pass
  • Instant rollback: Previous revisions remain available — traffic can be rerouted in seconds
  • Concurrency control: CI/CD pipeline prevents concurrent deploys (cancel-in-progress: false)
# .github/workflows/deploy-prod.yaml

concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
cancel-in-progress: false # Never cancel a running deploy

Summary — Fault Tolerance by Layer​

LayerMechanismRecovery Time
HTTP callsPolly exponential backoffSeconds (configurable retries)
SSE/WebSocketPolly WaitAndRetryForever + state persistence5 seconds per retry
Pub/Sub messagesManaged retry (10s-600s) + DLQ after 5 attemptsSeconds to minutes
Cloud Tasks3 retries with exponential backoff (3s-1hr)Seconds to minutes
Workflow batchesPer-batch retry (2 retries) + failure trackingSeconds
CacheGraceful degradation to databaseTransparent
Token expiryAuto-detection + re-auth + event ID preservationSeconds
Container shutdownGraceful state persistence to RedisSub-second
Bad deploymentCloud Run revision rollbackSeconds

Last updated: May 2026