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
| Layer | Mechanism | Recovery Time |
|---|---|---|
| HTTP calls | Polly exponential backoff | Seconds (configurable retries) |
| SSE/WebSocket | Polly WaitAndRetryForever + state persistence | 5 seconds per retry |
| Pub/Sub messages | Managed retry (10s-600s) + DLQ after 5 attempts | Seconds to minutes |
| Cloud Tasks | 3 retries with exponential backoff (3s-1hr) | Seconds to minutes |
| Workflow batches | Per-batch retry (2 retries) + failure tracking | Seconds |
| Cache | Graceful degradation to database | Transparent |
| Token expiry | Auto-detection + re-auth + event ID preservation | Seconds |
| Container shutdown | Graceful state persistence to Redis | Sub-second |
| Bad deployment | Cloud Run revision rollback | Seconds |
Last updated: May 2026