fix: fix circuit breaker failure recording

fix: fix circuit breaker failure recording
chore: provide advanced configuration option in UI for campaigns
This commit is contained in:
Abhishek Kumar 2026-03-05 13:43:13 +05:30
parent 628132f29b
commit 3ea235a666
17 changed files with 448 additions and 58 deletions

View file

@ -22,6 +22,7 @@ class CampaignClient(BaseDBClient):
retry_config: Optional[dict] = None, retry_config: Optional[dict] = None,
max_concurrency: Optional[int] = None, max_concurrency: Optional[int] = None,
schedule_config: Optional[dict] = None, schedule_config: Optional[dict] = None,
circuit_breaker: Optional[dict] = None,
) -> CampaignModel: ) -> CampaignModel:
"""Create a new campaign""" """Create a new campaign"""
async with self.async_session() as session: async with self.async_session() as session:
@ -31,6 +32,8 @@ class CampaignClient(BaseDBClient):
orchestrator_metadata["max_concurrency"] = max_concurrency orchestrator_metadata["max_concurrency"] = max_concurrency
if schedule_config is not None: if schedule_config is not None:
orchestrator_metadata["schedule_config"] = schedule_config orchestrator_metadata["schedule_config"] = schedule_config
if circuit_breaker is not None:
orchestrator_metadata["circuit_breaker"] = circuit_breaker
campaign = CampaignModel( campaign = CampaignModel(
name=name, name=name,
@ -68,6 +71,21 @@ class CampaignClient(BaseDBClient):
result = await session.execute(query) result = await session.execute(query)
return list(result.scalars().all()) return list(result.scalars().all())
async def get_latest_campaign(
self,
organization_id: int,
) -> Optional[CampaignModel]:
"""Get the most recently created campaign for an organization"""
async with self.async_session() as session:
query = (
select(CampaignModel)
.where(CampaignModel.organization_id == organization_id)
.order_by(CampaignModel.created_at.desc())
.limit(1)
)
result = await session.execute(query)
return result.scalars().first()
async def get_campaign( async def get_campaign(
self, self,
campaign_id: int, campaign_id: int,

View file

@ -6,7 +6,10 @@ from zoneinfo import ZoneInfo
from fastapi import APIRouter, Depends, HTTPException, Query from fastapi import APIRouter, Depends, HTTPException, Query
from pydantic import BaseModel, Field, field_validator, model_validator from pydantic import BaseModel, Field, field_validator, model_validator
from api.constants import DEFAULT_CAMPAIGN_RETRY_CONFIG, DEFAULT_ORG_CONCURRENCY_LIMIT from api.constants import (
DEFAULT_CAMPAIGN_RETRY_CONFIG,
DEFAULT_ORG_CONCURRENCY_LIMIT,
)
from api.db import db_client from api.db import db_client
from api.db.models import UserModel from api.db.models import UserModel
from api.enums import OrganizationConfigurationKey from api.enums import OrganizationConfigurationKey
@ -126,6 +129,20 @@ class ScheduleConfigResponse(BaseModel):
slots: List[TimeSlotResponse] slots: List[TimeSlotResponse]
class CircuitBreakerConfigRequest(BaseModel):
enabled: bool = True
failure_threshold: float = Field(default=0.5, ge=0.0, le=1.0)
window_seconds: int = Field(default=120, ge=30, le=600)
min_calls_in_window: int = Field(default=5, ge=1, le=100)
class CircuitBreakerConfigResponse(BaseModel):
enabled: bool
failure_threshold: float
window_seconds: int
min_calls_in_window: int
class CreateCampaignRequest(BaseModel): class CreateCampaignRequest(BaseModel):
name: str = Field(..., min_length=1, max_length=255) name: str = Field(..., min_length=1, max_length=255)
workflow_id: int workflow_id: int
@ -134,6 +151,7 @@ class CreateCampaignRequest(BaseModel):
retry_config: Optional[RetryConfigRequest] = None retry_config: Optional[RetryConfigRequest] = None
max_concurrency: Optional[int] = Field(default=None, ge=1, le=100) max_concurrency: Optional[int] = Field(default=None, ge=1, le=100)
schedule_config: Optional[ScheduleConfigRequest] = None schedule_config: Optional[ScheduleConfigRequest] = None
circuit_breaker: Optional[CircuitBreakerConfigRequest] = None
class UpdateCampaignRequest(BaseModel): class UpdateCampaignRequest(BaseModel):
@ -141,6 +159,7 @@ class UpdateCampaignRequest(BaseModel):
retry_config: Optional[RetryConfigRequest] = None retry_config: Optional[RetryConfigRequest] = None
max_concurrency: Optional[int] = Field(default=None, ge=1, le=100) max_concurrency: Optional[int] = Field(default=None, ge=1, le=100)
schedule_config: Optional[ScheduleConfigRequest] = None schedule_config: Optional[ScheduleConfigRequest] = None
circuit_breaker: Optional[CircuitBreakerConfigRequest] = None
class CampaignResponse(BaseModel): class CampaignResponse(BaseModel):
@ -160,6 +179,7 @@ class CampaignResponse(BaseModel):
retry_config: RetryConfigResponse retry_config: RetryConfigResponse
max_concurrency: Optional[int] = None max_concurrency: Optional[int] = None
schedule_config: Optional[ScheduleConfigResponse] = None schedule_config: Optional[ScheduleConfigResponse] = None
circuit_breaker: Optional[CircuitBreakerConfigResponse] = None
class CampaignsResponse(BaseModel): class CampaignsResponse(BaseModel):
@ -209,9 +229,10 @@ def _build_campaign_response(campaign, workflow_name: str) -> CampaignResponse:
else DEFAULT_CAMPAIGN_RETRY_CONFIG else DEFAULT_CAMPAIGN_RETRY_CONFIG
) )
# Get max_concurrency and schedule_config from orchestrator_metadata # Get max_concurrency, schedule_config, circuit_breaker from orchestrator_metadata
max_concurrency = None max_concurrency = None
schedule_config = None schedule_config = None
circuit_breaker_config = None
if campaign.orchestrator_metadata: if campaign.orchestrator_metadata:
max_concurrency = campaign.orchestrator_metadata.get("max_concurrency") max_concurrency = campaign.orchestrator_metadata.get("max_concurrency")
sc = campaign.orchestrator_metadata.get("schedule_config") sc = campaign.orchestrator_metadata.get("schedule_config")
@ -221,6 +242,9 @@ def _build_campaign_response(campaign, workflow_name: str) -> CampaignResponse:
timezone=sc.get("timezone", "UTC"), timezone=sc.get("timezone", "UTC"),
slots=[TimeSlotResponse(**slot) for slot in sc.get("slots", [])], slots=[TimeSlotResponse(**slot) for slot in sc.get("slots", [])],
) )
cb = campaign.orchestrator_metadata.get("circuit_breaker")
if cb:
circuit_breaker_config = CircuitBreakerConfigResponse(**cb)
return CampaignResponse( return CampaignResponse(
id=campaign.id, id=campaign.id,
@ -239,6 +263,7 @@ def _build_campaign_response(campaign, workflow_name: str) -> CampaignResponse:
retry_config=RetryConfigResponse(**retry_config), retry_config=RetryConfigResponse(**retry_config),
max_concurrency=max_concurrency, max_concurrency=max_concurrency,
schedule_config=schedule_config, schedule_config=schedule_config,
circuit_breaker=circuit_breaker_config,
) )
@ -276,6 +301,11 @@ async def create_campaign(
if request.schedule_config: if request.schedule_config:
schedule_config = request.schedule_config.model_dump() schedule_config = request.schedule_config.model_dump()
# Build circuit_breaker dict if provided
circuit_breaker_config = None
if request.circuit_breaker:
circuit_breaker_config = request.circuit_breaker.model_dump()
campaign = await db_client.create_campaign( campaign = await db_client.create_campaign(
name=request.name, name=request.name,
workflow_id=request.workflow_id, workflow_id=request.workflow_id,
@ -286,6 +316,7 @@ async def create_campaign(
retry_config=retry_config, retry_config=retry_config,
max_concurrency=request.max_concurrency, max_concurrency=request.max_concurrency,
schedule_config=schedule_config, schedule_config=schedule_config,
circuit_breaker=circuit_breaker_config,
) )
return _build_campaign_response(campaign, workflow_name) return _build_campaign_response(campaign, workflow_name)
@ -436,6 +467,10 @@ async def update_campaign(
metadata["schedule_config"] = request.schedule_config.model_dump() metadata["schedule_config"] = request.schedule_config.model_dump()
metadata_changed = True metadata_changed = True
if request.circuit_breaker is not None:
metadata["circuit_breaker"] = request.circuit_breaker.model_dump()
metadata_changed = True
if metadata_changed: if metadata_changed:
update_kwargs["orchestrator_metadata"] = metadata update_kwargs["orchestrator_metadata"] = metadata

View file

@ -1,4 +1,4 @@
from typing import Union from typing import List, Optional, Union
from fastapi import APIRouter, Depends, HTTPException from fastapi import APIRouter, Depends, HTTPException
from pydantic import BaseModel from pydantic import BaseModel
@ -257,14 +257,41 @@ class RetryConfigResponse(BaseModel):
retry_on_voicemail: bool retry_on_voicemail: bool
class CampaignLimitsResponse(BaseModel): class TimeSlotResponse(BaseModel):
day_of_week: int
start_time: str
end_time: str
class ScheduleConfigResponse(BaseModel):
enabled: bool
timezone: str
slots: List[TimeSlotResponse]
class CircuitBreakerConfigResponse(BaseModel):
enabled: bool
failure_threshold: float
window_seconds: int
min_calls_in_window: int
class LastCampaignSettingsResponse(BaseModel):
retry_config: Optional[RetryConfigResponse] = None
max_concurrency: Optional[int] = None
schedule_config: Optional[ScheduleConfigResponse] = None
circuit_breaker: Optional[CircuitBreakerConfigResponse] = None
class CampaignDefaultsResponse(BaseModel):
concurrent_call_limit: int concurrent_call_limit: int
from_numbers_count: int from_numbers_count: int
default_retry_config: RetryConfigResponse default_retry_config: RetryConfigResponse
last_campaign_settings: Optional[LastCampaignSettingsResponse] = None
@router.get("/campaign-limits", response_model=CampaignLimitsResponse) @router.get("/campaign-defaults", response_model=CampaignDefaultsResponse)
async def get_campaign_limits(user: UserModel = Depends(get_user)): async def get_campaign_defaults(user: UserModel = Depends(get_user)):
"""Get campaign limits for the user's organization. """Get campaign limits for the user's organization.
Returns the organization's concurrent call limit and default retry configuration. Returns the organization's concurrent call limit and default retry configuration.
@ -299,8 +326,47 @@ async def get_campaign_limits(user: UserModel = Depends(get_user)):
except Exception: except Exception:
pass pass
return CampaignLimitsResponse( # Get last campaign settings for pre-population
last_campaign_settings = None
try:
last_campaign = await db_client.get_latest_campaign(
user.selected_organization_id
)
if last_campaign:
retry = None
if last_campaign.retry_config:
retry = RetryConfigResponse(**last_campaign.retry_config)
max_conc = None
sched = None
cb = None
if last_campaign.orchestrator_metadata:
max_conc = last_campaign.orchestrator_metadata.get("max_concurrency")
sc = last_campaign.orchestrator_metadata.get("schedule_config")
if sc:
sched = ScheduleConfigResponse(
enabled=sc.get("enabled", False),
timezone=sc.get("timezone", "UTC"),
slots=[
TimeSlotResponse(**slot) for slot in sc.get("slots", [])
],
)
cb_data = last_campaign.orchestrator_metadata.get("circuit_breaker")
if cb_data:
cb = CircuitBreakerConfigResponse(**cb_data)
last_campaign_settings = LastCampaignSettingsResponse(
retry_config=retry,
max_concurrency=max_conc,
schedule_config=sched,
circuit_breaker=cb,
)
except Exception:
pass
return CampaignDefaultsResponse(
concurrent_call_limit=concurrent_limit, concurrent_call_limit=concurrent_limit,
from_numbers_count=from_numbers_count, from_numbers_count=from_numbers_count,
default_retry_config=RetryConfigResponse(**DEFAULT_CAMPAIGN_RETRY_CONFIG), default_retry_config=RetryConfigResponse(**DEFAULT_CAMPAIGN_RETRY_CONFIG),
last_campaign_settings=last_campaign_settings,
) )

View file

@ -783,7 +783,8 @@ async def _process_status_update(workflow_run_id: int, status: StatusCallbackReq
if workflow_run.campaign_id: if workflow_run.campaign_id:
await campaign_call_dispatcher.release_call_slot(workflow_run_id) await campaign_call_dispatcher.release_call_slot(workflow_run_id)
await circuit_breaker.record_and_evaluate( await circuit_breaker.record_and_evaluate(
workflow_run.campaign_id, is_failure=True workflow_run.campaign_id,
is_failure=status.status == "error",
) )
# Check if retry is needed for campaign calls (busy/no-answer) # Check if retry is needed for campaign calls (busy/no-answer)
@ -1209,6 +1210,7 @@ async def handle_cloudonix_status_callback(
return {"status": "success"} return {"status": "success"}
@router.post("/cloudonix/amd-callback/{workflow_run_id}") @router.post("/cloudonix/amd-callback/{workflow_run_id}")
async def handle_cloudonix_amd_callback( async def handle_cloudonix_amd_callback(
workflow_run_id: int, workflow_run_id: int,

View file

@ -9,6 +9,7 @@ from api.constants import DEFAULT_ORG_CONCURRENCY_LIMIT
from api.db import db_client from api.db import db_client
from api.db.models import QueuedRunModel, WorkflowRunModel from api.db.models import QueuedRunModel, WorkflowRunModel
from api.enums import OrganizationConfigurationKey, WorkflowRunState from api.enums import OrganizationConfigurationKey, WorkflowRunState
from api.services.campaign.circuit_breaker import circuit_breaker
from api.services.campaign.errors import ( from api.services.campaign.errors import (
ConcurrentSlotAcquisitionError, ConcurrentSlotAcquisitionError,
PhoneNumberPoolExhaustedError, PhoneNumberPoolExhaustedError,
@ -315,6 +316,9 @@ class CampaignCallDispatcher:
}, },
) )
# Record call initiation failure in circuit breaker
await circuit_breaker.record_and_evaluate(campaign.id, is_failure=True)
# Release concurrent slot on failure # Release concurrent slot on failure
mapping = await rate_limiter.get_workflow_slot_mapping(workflow_run.id) mapping = await rate_limiter.get_workflow_slot_mapping(workflow_run.id)
if mapping: if mapping:

View file

@ -14,7 +14,6 @@ setup_logging()
import asyncio import asyncio
import json import json
import signal import signal
import time
from typing import Dict, Optional, Set from typing import Dict, Optional, Set
from urllib.parse import urlparse from urllib.parse import urlparse

View file

@ -45,13 +45,16 @@ class ARIBridgeSwapStrategy(TransferStrategy):
from api.services.telephony.call_transfer_manager import ( from api.services.telephony.call_transfer_manager import (
get_call_transfer_manager, get_call_transfer_manager,
) )
auth = BasicAuth(app_name, app_password) auth = BasicAuth(app_name, app_password)
# Get call transfer manager instance # Get call transfer manager instance
call_transfer_manager = await get_call_transfer_manager() call_transfer_manager = await get_call_transfer_manager()
# 1. Find active transfer context for this caller channel # 1. Find active transfer context for this caller channel
transfer_context = await call_transfer_manager.find_transfer_context_for_call(channel_id) transfer_context = (
await call_transfer_manager.find_transfer_context_for_call(channel_id)
)
if not transfer_context: if not transfer_context:
logger.error( logger.error(
f"[ARI Transfer] No active transfer context found for caller {channel_id}" f"[ARI Transfer] No active transfer context found for caller {channel_id}"
@ -178,6 +181,7 @@ class ARIBridgeSwapStrategy(TransferStrategy):
logger.exception(f"Failed to execute ARI transfer: {e}") logger.exception(f"Failed to execute ARI transfer: {e}")
return False return False
class ARIHangupStrategy(HangupStrategy): class ARIHangupStrategy(HangupStrategy):
"""Implements hangup for Asterisk ARI channels.""" """Implements hangup for Asterisk ARI channels."""

View file

@ -455,7 +455,9 @@ class ARIProvider(TelephonyProvider):
} }
except Exception as e: except Exception as e:
logger.error(f"[ARI Transfer] Failed to originate call transfer destination channel: {e}") logger.error(
f"[ARI Transfer] Failed to originate call transfer destination channel: {e}"
)
await call_transfer_manager.remove_transfer_context(transfer_id) await call_transfer_manager.remove_transfer_context(transfer_id)
raise raise

View file

@ -107,9 +107,10 @@ class CloudonixProvider(TelephonyProvider):
} }
data["machineDetection"] = "DetectMessageEnd" data["machineDetection"] = "DetectMessageEnd"
data["asyncAmd"] = True data["asyncAmd"] = True
data["asyncAmdStatusCallback"] = f"{backend_endpoint}/api/v1/telephony/cloudonix/amd-callback/{workflow_run_id}" data["asyncAmdStatusCallback"] = (
data["asyncAmdStatusCallbackMethod"]= "POST" f"{backend_endpoint}/api/v1/telephony/cloudonix/amd-callback/{workflow_run_id}"
)
data["asyncAmdStatusCallbackMethod"] = "POST"
# TODO: Cloudonix status callbacks are spammy, so commenting it out. Can send it to # TODO: Cloudonix status callbacks are spammy, so commenting it out. Can send it to
# some persistent logging system instead of transcational database. # some persistent logging system instead of transcational database.

View file

@ -76,20 +76,26 @@ class TwilioConferenceStrategy(TransferStrategy):
) )
# 3. Clean up transfer context after successful transfer # 3. Clean up transfer context after successful transfer
await self._cleanup_transfer_context(transfer_context.transfer_id) await self._cleanup_transfer_context(
transfer_context.transfer_id
)
return True return True
elif response.status == 404: elif response.status == 404:
logger.error( logger.error(
f"Failed to transfer Twilio call {call_sid}: Call not found (404)" f"Failed to transfer Twilio call {call_sid}: Call not found (404)"
) )
await self._cleanup_transfer_context(transfer_context.transfer_id) await self._cleanup_transfer_context(
transfer_context.transfer_id
)
return False return False
else: else:
logger.error( logger.error(
f"Failed to transfer Twilio call {call_sid} to conference {conference_name}: " f"Failed to transfer Twilio call {call_sid} to conference {conference_name}: "
f"Status {response.status}, Response: {response_text}" f"Status {response.status}, Response: {response_text}"
) )
await self._cleanup_transfer_context(transfer_context.transfer_id) await self._cleanup_transfer_context(
transfer_context.transfer_id
)
return False return False
except Exception as e: except Exception as e:

View file

@ -69,6 +69,16 @@ async def process_knowledge_base_document(
file_size = os.path.getsize(temp_file_path) file_size = os.path.getsize(temp_file_path)
logger.info(f"Downloaded file size: {file_size} bytes") logger.info(f"Downloaded file size: {file_size} bytes")
# Validate file size (max 5MB)
max_file_size = 5 * 1024 * 1024
if file_size > max_file_size:
error_message = f"File size ({file_size / (1024 * 1024):.1f}MB) exceeds the maximum allowed size of 5MB."
logger.warning(f"Document {document_id}: {error_message}")
await db_client.update_document_status(
document_id, "failed", error_message=error_message
)
return
# Compute file hash and get mime type # Compute file hash and get mime type
file_hash = db_client.compute_file_hash(temp_file_path) file_hash = db_client.compute_file_hash(temp_file_path)
mime_type = db_client.get_mime_type(temp_file_path) mime_type = db_client.get_mime_type(temp_file_path)

View file

@ -46,6 +46,15 @@ export interface CampaignAdvancedSettingsProps {
onScheduleTimezoneChange: (value: ITimezoneOption | string) => void; onScheduleTimezoneChange: (value: ITimezoneOption | string) => void;
timeSlots: TimeSlot[]; timeSlots: TimeSlot[];
onTimeSlotsChange: (value: TimeSlot[]) => void; onTimeSlotsChange: (value: TimeSlot[]) => void;
// Circuit breaker config
circuitBreakerEnabled: boolean;
onCircuitBreakerEnabledChange: (value: boolean) => void;
circuitBreakerFailureThreshold: string;
onCircuitBreakerFailureThresholdChange: (value: string) => void;
circuitBreakerWindowSeconds: string;
onCircuitBreakerWindowSecondsChange: (value: string) => void;
circuitBreakerMinCalls: string;
onCircuitBreakerMinCallsChange: (value: string) => void;
} }
/** Extract the string timezone value from ITimezoneOption | string */ /** Extract the string timezone value from ITimezoneOption | string */
@ -101,6 +110,10 @@ export default function CampaignAdvancedSettings({
retryOnVoicemail, onRetryOnVoicemailChange, retryOnVoicemail, onRetryOnVoicemailChange,
scheduleEnabled, onScheduleEnabledChange, scheduleTimezone, onScheduleTimezoneChange, scheduleEnabled, onScheduleEnabledChange, scheduleTimezone, onScheduleTimezoneChange,
timeSlots, onTimeSlotsChange, timeSlots, onTimeSlotsChange,
circuitBreakerEnabled, onCircuitBreakerEnabledChange,
circuitBreakerFailureThreshold, onCircuitBreakerFailureThresholdChange,
circuitBreakerWindowSeconds, onCircuitBreakerWindowSecondsChange,
circuitBreakerMinCalls, onCircuitBreakerMinCallsChange,
}: CampaignAdvancedSettingsProps) { }: CampaignAdvancedSettingsProps) {
const timezoneSelectId = useId(); const timezoneSelectId = useId();
@ -295,6 +308,68 @@ export default function CampaignAdvancedSettings({
</div> </div>
)} )}
</div> </div>
<Separator />
{/* Circuit Breaker */}
<div className="space-y-4">
<div className="flex items-center justify-between">
<div>
<Label htmlFor="circuit-breaker-enabled">Circuit Breaker</Label>
<p className="text-sm text-muted-foreground">
Auto-pause campaign on high failure rates
</p>
</div>
<Switch
id="circuit-breaker-enabled"
checked={circuitBreakerEnabled}
onCheckedChange={onCircuitBreakerEnabledChange}
/>
</div>
{circuitBreakerEnabled && (
<div className="space-y-4 pl-4 border-l-2 border-muted">
<div className="space-y-2">
<Label htmlFor="cb-failure-threshold">Failure Threshold (%)</Label>
<Input
id="cb-failure-threshold"
type="number"
value={circuitBreakerFailureThreshold}
onChange={(e) => onCircuitBreakerFailureThresholdChange(e.target.value)}
min={1}
max={100}
/>
<p className="text-sm text-muted-foreground">
Pause when failure rate exceeds this percentage
</p>
</div>
<div className="grid grid-cols-2 gap-4">
<div className="space-y-2">
<Label htmlFor="cb-window">Window (seconds)</Label>
<Input
id="cb-window"
type="number"
value={circuitBreakerWindowSeconds}
onChange={(e) => onCircuitBreakerWindowSecondsChange(e.target.value)}
min={30}
max={600}
/>
</div>
<div className="space-y-2">
<Label htmlFor="cb-min-calls">Min Calls in Window</Label>
<Input
id="cb-min-calls"
type="number"
value={circuitBreakerMinCalls}
onChange={(e) => onCircuitBreakerMinCallsChange(e.target.value)}
min={1}
max={100}
/>
</div>
</div>
</div>
)}
</div>
</div> </div>
); );
} }

View file

@ -8,8 +8,8 @@ import { toast } from 'sonner';
import { import {
getCampaignApiV1CampaignCampaignIdGet, getCampaignApiV1CampaignCampaignIdGet,
getCampaignLimitsApiV1OrganizationsCampaignLimitsGet, getCampaignDefaultsApiV1OrganizationsCampaignDefaultsGet,
updateCampaignApiV1CampaignCampaignIdPatch, updateCampaignApiV1CampaignCampaignIdPatch
} from '@/client/sdk.gen'; } from '@/client/sdk.gen';
import type { CampaignResponse } from '@/client/types.gen'; import type { CampaignResponse } from '@/client/types.gen';
import { Button } from '@/components/ui/button'; import { Button } from '@/components/ui/button';
@ -55,6 +55,11 @@ export default function EditCampaignPage() {
const [timeSlots, setTimeSlots] = useState<TimeSlot[]>([ const [timeSlots, setTimeSlots] = useState<TimeSlot[]>([
{ day_of_week: 0, start_time: '09:00', end_time: '17:00' }, { day_of_week: 0, start_time: '09:00', end_time: '17:00' },
]); ]);
// Circuit breaker config state
const [circuitBreakerEnabled, setCircuitBreakerEnabled] = useState(true);
const [circuitBreakerFailureThreshold, setCircuitBreakerFailureThreshold] = useState<string>('50');
const [circuitBreakerWindowSeconds, setCircuitBreakerWindowSeconds] = useState<string>('120');
const [circuitBreakerMinCalls, setCircuitBreakerMinCalls] = useState<string>('5');
// Redirect if not authenticated // Redirect if not authenticated
useEffect(() => { useEffect(() => {
@ -104,6 +109,15 @@ export default function EditCampaignPage() {
setTimeSlots(c.schedule_config.slots.map(s => ({ ...s }))); setTimeSlots(c.schedule_config.slots.map(s => ({ ...s })));
} }
} }
// Circuit breaker config
const cb = (c as unknown as { circuit_breaker?: { enabled: boolean; failure_threshold: number; window_seconds: number; min_calls_in_window: number } }).circuit_breaker;
if (cb) {
setCircuitBreakerEnabled(cb.enabled);
setCircuitBreakerFailureThreshold(String(Math.round(cb.failure_threshold * 100)));
setCircuitBreakerWindowSeconds(String(cb.window_seconds));
setCircuitBreakerMinCalls(String(cb.min_calls_in_window));
}
} }
} catch (error) { } catch (error) {
console.error('Failed to fetch campaign:', error); console.error('Failed to fetch campaign:', error);
@ -115,11 +129,11 @@ export default function EditCampaignPage() {
}, [user, getAccessToken, campaignId, router]); }, [user, getAccessToken, campaignId, router]);
// Fetch campaign limits // Fetch campaign limits
const fetchCampaignLimits = useCallback(async () => { const fetchCampaignDefaults = useCallback(async () => {
if (!user) return; if (!user) return;
try { try {
const accessToken = await getAccessToken(); const accessToken = await getAccessToken();
const response = await getCampaignLimitsApiV1OrganizationsCampaignLimitsGet({ const response = await getCampaignDefaultsApiV1OrganizationsCampaignDefaultsGet({
headers: { 'Authorization': `Bearer ${accessToken}` }, headers: { 'Authorization': `Bearer ${accessToken}` },
}); });
@ -136,9 +150,9 @@ export default function EditCampaignPage() {
useEffect(() => { useEffect(() => {
if (user) { if (user) {
fetchCampaign(); fetchCampaign();
fetchCampaignLimits(); fetchCampaignDefaults();
} }
}, [fetchCampaign, fetchCampaignLimits, user]); }, [fetchCampaign, fetchCampaignDefaults, user]);
// Effective concurrency limit // Effective concurrency limit
const effectiveLimit = fromNumbersCount > 0 const effectiveLimit = fromNumbersCount > 0
@ -213,6 +227,14 @@ export default function EditCampaignPage() {
slots: [{ day_of_week: 0, start_time: '09:00', end_time: '17:00' }], slots: [{ day_of_week: 0, start_time: '09:00', end_time: '17:00' }],
}; };
const circuitBreakerConfig = {
enabled: circuitBreakerEnabled,
failure_threshold: (parseInt(circuitBreakerFailureThreshold) || 50) / 100,
window_seconds: parseInt(circuitBreakerWindowSeconds) || 120,
min_calls_in_window: parseInt(circuitBreakerMinCalls) || 5,
};
const response = await updateCampaignApiV1CampaignCampaignIdPatch({ const response = await updateCampaignApiV1CampaignCampaignIdPatch({
path: { campaign_id: campaignId }, path: { campaign_id: campaignId },
body: { body: {
@ -220,6 +242,7 @@ export default function EditCampaignPage() {
retry_config: retryConfig, retry_config: retryConfig,
max_concurrency: maxConcurrencyValue, max_concurrency: maxConcurrencyValue,
schedule_config: scheduleConfig, schedule_config: scheduleConfig,
circuit_breaker: circuitBreakerConfig,
}, },
headers: { 'Authorization': `Bearer ${accessToken}` }, headers: { 'Authorization': `Bearer ${accessToken}` },
}); });
@ -332,6 +355,14 @@ export default function EditCampaignPage() {
onScheduleTimezoneChange={setScheduleTimezone} onScheduleTimezoneChange={setScheduleTimezone}
timeSlots={timeSlots} timeSlots={timeSlots}
onTimeSlotsChange={setTimeSlots} onTimeSlotsChange={setTimeSlots}
circuitBreakerEnabled={circuitBreakerEnabled}
onCircuitBreakerEnabledChange={setCircuitBreakerEnabled}
circuitBreakerFailureThreshold={circuitBreakerFailureThreshold}
onCircuitBreakerFailureThresholdChange={setCircuitBreakerFailureThreshold}
circuitBreakerWindowSeconds={circuitBreakerWindowSeconds}
onCircuitBreakerWindowSecondsChange={setCircuitBreakerWindowSeconds}
circuitBreakerMinCalls={circuitBreakerMinCalls}
onCircuitBreakerMinCallsChange={setCircuitBreakerMinCalls}
/> />
{submitError && ( {submitError && (

View file

@ -8,7 +8,7 @@ import { toast } from 'sonner';
import { import {
createCampaignApiV1CampaignCreatePost, createCampaignApiV1CampaignCreatePost,
getCampaignLimitsApiV1OrganizationsCampaignLimitsGet, getCampaignDefaultsApiV1OrganizationsCampaignDefaultsGet,
getWorkflowsSummaryApiV1WorkflowSummaryGet getWorkflowsSummaryApiV1WorkflowSummaryGet
} from '@/client/sdk.gen'; } from '@/client/sdk.gen';
import type { WorkflowSummaryResponse } from '@/client/types.gen'; import type { WorkflowSummaryResponse } from '@/client/types.gen';
@ -72,6 +72,11 @@ export default function NewCampaignPage() {
const [timeSlots, setTimeSlots] = useState<TimeSlot[]>([ const [timeSlots, setTimeSlots] = useState<TimeSlot[]>([
{ day_of_week: 0, start_time: '09:00', end_time: '17:00' }, { day_of_week: 0, start_time: '09:00', end_time: '17:00' },
]); ]);
// Circuit breaker config state
const [circuitBreakerEnabled, setCircuitBreakerEnabled] = useState(true);
const [circuitBreakerFailureThreshold, setCircuitBreakerFailureThreshold] = useState<string>('50');
const [circuitBreakerWindowSeconds, setCircuitBreakerWindowSeconds] = useState<string>('120');
const [circuitBreakerMinCalls, setCircuitBreakerMinCalls] = useState<string>('5');
// Redirect if not authenticated // Redirect if not authenticated
useEffect(() => { useEffect(() => {
@ -104,11 +109,11 @@ export default function NewCampaignPage() {
}, [user, getAccessToken]); }, [user, getAccessToken]);
// Fetch campaign limits // Fetch campaign limits
const fetchCampaignLimits = useCallback(async () => { const fetchCampaignDefaults = useCallback(async () => {
if (!user) return; if (!user) return;
try { try {
const accessToken = await getAccessToken(); const accessToken = await getAccessToken();
const response = await getCampaignLimitsApiV1OrganizationsCampaignLimitsGet({ const response = await getCampaignDefaultsApiV1OrganizationsCampaignDefaultsGet({
headers: { headers: {
'Authorization': `Bearer ${accessToken}`, 'Authorization': `Bearer ${accessToken}`,
} }
@ -117,14 +122,56 @@ export default function NewCampaignPage() {
if (response.data) { if (response.data) {
setOrgConcurrentLimit(response.data.concurrent_call_limit); setOrgConcurrentLimit(response.data.concurrent_call_limit);
setFromNumbersCount(response.data.from_numbers_count); setFromNumbersCount(response.data.from_numbers_count);
// Initialize retry config from defaults
const retryConfig = response.data.default_retry_config; const last = (response.data as { last_campaign_settings?: {
setRetryEnabled(retryConfig.enabled); retry_config?: { enabled: boolean; max_retries: number; retry_delay_seconds: number; retry_on_busy: boolean; retry_on_no_answer: boolean; retry_on_voicemail: boolean };
setMaxRetries(String(retryConfig.max_retries)); max_concurrency?: number | null;
setRetryDelaySeconds(String(retryConfig.retry_delay_seconds)); schedule_config?: { enabled: boolean; timezone: string; slots: TimeSlot[] } | null;
setRetryOnBusy(retryConfig.retry_on_busy); circuit_breaker?: { enabled: boolean; failure_threshold: number; window_seconds: number; min_calls_in_window: number } | null;
setRetryOnNoAnswer(retryConfig.retry_on_no_answer); } | null }).last_campaign_settings;
setRetryOnVoicemail(retryConfig.retry_on_voicemail);
if (last) {
// Pre-populate from last campaign
if (last.retry_config) {
setRetryEnabled(last.retry_config.enabled);
setMaxRetries(String(last.retry_config.max_retries));
setRetryDelaySeconds(String(last.retry_config.retry_delay_seconds));
setRetryOnBusy(last.retry_config.retry_on_busy);
setRetryOnNoAnswer(last.retry_config.retry_on_no_answer);
setRetryOnVoicemail(last.retry_config.retry_on_voicemail);
} else {
const retryConfig = response.data.default_retry_config;
setRetryEnabled(retryConfig.enabled);
setMaxRetries(String(retryConfig.max_retries));
setRetryDelaySeconds(String(retryConfig.retry_delay_seconds));
setRetryOnBusy(retryConfig.retry_on_busy);
setRetryOnNoAnswer(retryConfig.retry_on_no_answer);
setRetryOnVoicemail(retryConfig.retry_on_voicemail);
}
if (last.max_concurrency) {
setMaxConcurrency(String(last.max_concurrency));
}
if (last.schedule_config) {
setScheduleEnabled(last.schedule_config.enabled);
setScheduleTimezone(last.schedule_config.timezone);
setTimeSlots(last.schedule_config.slots);
}
if (last.circuit_breaker) {
setCircuitBreakerEnabled(last.circuit_breaker.enabled);
setCircuitBreakerFailureThreshold(String(Math.round(last.circuit_breaker.failure_threshold * 100)));
setCircuitBreakerWindowSeconds(String(last.circuit_breaker.window_seconds));
setCircuitBreakerMinCalls(String(last.circuit_breaker.min_calls_in_window));
}
} else {
// No previous campaign — use defaults
const retryConfig = response.data.default_retry_config;
setRetryEnabled(retryConfig.enabled);
setMaxRetries(String(retryConfig.max_retries));
setRetryDelaySeconds(String(retryConfig.retry_delay_seconds));
setRetryOnBusy(retryConfig.retry_on_busy);
setRetryOnNoAnswer(retryConfig.retry_on_no_answer);
setRetryOnVoicemail(retryConfig.retry_on_voicemail);
}
} }
} catch (error) { } catch (error) {
console.error('Failed to fetch campaign limits:', error); console.error('Failed to fetch campaign limits:', error);
@ -135,9 +182,9 @@ export default function NewCampaignPage() {
useEffect(() => { useEffect(() => {
if (user) { if (user) {
fetchWorkflows(); fetchWorkflows();
fetchCampaignLimits(); fetchCampaignDefaults();
} }
}, [fetchWorkflows, fetchCampaignLimits, user]); }, [fetchWorkflows, fetchCampaignDefaults, user]);
// Effective concurrency limit considering both org limit and available CLIs // Effective concurrency limit considering both org limit and available CLIs
const effectiveLimit = fromNumbersCount > 0 const effectiveLimit = fromNumbersCount > 0
@ -195,6 +242,15 @@ export default function NewCampaignPage() {
} }
: undefined; : undefined;
// Build circuit_breaker config
const circuitBreakerConfig = {
enabled: circuitBreakerEnabled,
failure_threshold: (parseInt(circuitBreakerFailureThreshold) || 50) / 100,
window_seconds: parseInt(circuitBreakerWindowSeconds) || 120,
min_calls_in_window: parseInt(circuitBreakerMinCalls) || 5,
};
const response = await createCampaignApiV1CampaignCreatePost({ const response = await createCampaignApiV1CampaignCreatePost({
body: { body: {
name: campaignName, name: campaignName,
@ -204,6 +260,7 @@ export default function NewCampaignPage() {
retry_config: retryConfig, retry_config: retryConfig,
max_concurrency: maxConcurrencyValue, max_concurrency: maxConcurrencyValue,
schedule_config: scheduleConfig, schedule_config: scheduleConfig,
circuit_breaker: circuitBreakerConfig,
}, },
headers: { headers: {
'Authorization': `Bearer ${accessToken}`, 'Authorization': `Bearer ${accessToken}`,
@ -401,6 +458,14 @@ export default function NewCampaignPage() {
onScheduleTimezoneChange={setScheduleTimezone} onScheduleTimezoneChange={setScheduleTimezone}
timeSlots={timeSlots} timeSlots={timeSlots}
onTimeSlotsChange={setTimeSlots} onTimeSlotsChange={setTimeSlots}
circuitBreakerEnabled={circuitBreakerEnabled}
onCircuitBreakerEnabledChange={setCircuitBreakerEnabled}
circuitBreakerFailureThreshold={circuitBreakerFailureThreshold}
onCircuitBreakerFailureThresholdChange={setCircuitBreakerFailureThreshold}
circuitBreakerWindowSeconds={circuitBreakerWindowSeconds}
onCircuitBreakerWindowSecondsChange={setCircuitBreakerWindowSeconds}
circuitBreakerMinCalls={circuitBreakerMinCalls}
onCircuitBreakerMinCallsChange={setCircuitBreakerMinCalls}
/> />
</CollapsibleContent> </CollapsibleContent>
</Collapsible> </Collapsible>

View file

@ -17,7 +17,7 @@ interface DocumentUploadProps {
onUploadSuccess: () => void; onUploadSuccess: () => void;
} }
const MAX_FILE_SIZE = 100 * 1024 * 1024; // 100MB const MAX_FILE_SIZE = 5 * 1024 * 1024; // 5MB
const ACCEPTED_FILE_TYPES = ['.pdf', '.docx', '.doc', '.txt']; const ACCEPTED_FILE_TYPES = ['.pdf', '.docx', '.doc', '.txt'];
export default function DocumentUpload({ onUploadSuccess }: DocumentUploadProps) { export default function DocumentUpload({ onUploadSuccess }: DocumentUploadProps) {
@ -36,7 +36,7 @@ export default function DocumentUpload({ onUploadSuccess }: DocumentUploadProps)
// Validate file size // Validate file size
if (file.size > MAX_FILE_SIZE) { if (file.size > MAX_FILE_SIZE) {
toast.error('File size must be less than 100MB'); toast.error('File size must be less than 5MB');
return false; return false;
} }
@ -44,7 +44,13 @@ export default function DocumentUpload({ onUploadSuccess }: DocumentUploadProps)
}; };
const uploadFile = async (file: File) => { const uploadFile = async (file: File) => {
if (!validateFile(file)) return; if (!validateFile(file)) {
// Reset file input so the same file can be re-selected
if (fileInputRef.current) {
fileInputRef.current.value = '';
}
return;
}
setUploading(true); setUploading(true);
setUploadProgress(0); setUploadProgress(0);
@ -182,7 +188,7 @@ export default function DocumentUpload({ onUploadSuccess }: DocumentUploadProps)
or click to browse or click to browse
</p> </p>
<p className="text-xs text-muted-foreground"> <p className="text-xs text-muted-foreground">
Supported formats: {ACCEPTED_FILE_TYPES.join(', ')} (Max 100MB) Supported formats: {ACCEPTED_FILE_TYPES.join(', ')} (Max 5MB)
</p> </p>
</div> </div>

File diff suppressed because one or more lines are too long

View file

@ -82,10 +82,11 @@ export type AuthUserResponse = {
export type CallType = 'inbound' | 'outbound'; export type CallType = 'inbound' | 'outbound';
export type CampaignLimitsResponse = { export type CampaignDefaultsResponse = {
concurrent_call_limit: number; concurrent_call_limit: number;
from_numbers_count: number; from_numbers_count: number;
default_retry_config: RetryConfigResponse; default_retry_config: RetryConfigResponse;
last_campaign_settings?: LastCampaignSettingsResponse | null;
}; };
export type CampaignProgressResponse = { export type CampaignProgressResponse = {
@ -120,6 +121,7 @@ export type CampaignResponse = {
retry_config: RetryConfigResponse; retry_config: RetryConfigResponse;
max_concurrency?: number | null; max_concurrency?: number | null;
schedule_config?: ScheduleConfigResponse | null; schedule_config?: ScheduleConfigResponse | null;
circuit_breaker?: CircuitBreakerConfigResponse | null;
}; };
/** /**
@ -192,6 +194,20 @@ export type ChunkSearchResponseSchema = {
total_results: number; total_results: number;
}; };
export type CircuitBreakerConfigRequest = {
enabled?: boolean;
failure_threshold?: number;
window_seconds?: number;
min_calls_in_window?: number;
};
export type CircuitBreakerConfigResponse = {
enabled: boolean;
failure_threshold: number;
window_seconds: number;
min_calls_in_window: number;
};
/** /**
* Request schema for Cloudonix configuration. * Request schema for Cloudonix configuration.
*/ */
@ -241,6 +257,7 @@ export type CreateCampaignRequest = {
retry_config?: RetryConfigRequest | null; retry_config?: RetryConfigRequest | null;
max_concurrency?: number | null; max_concurrency?: number | null;
schedule_config?: ScheduleConfigRequest | null; schedule_config?: ScheduleConfigRequest | null;
circuit_breaker?: CircuitBreakerConfigRequest | null;
}; };
/** /**
@ -715,6 +732,13 @@ export type IntegrationResponse = {
export type ItemKind = 'node' | 'edge' | 'workflow'; export type ItemKind = 'node' | 'edge' | 'workflow';
export type LastCampaignSettingsResponse = {
retry_config?: RetryConfigResponse | null;
max_concurrency?: number | null;
schedule_config?: ScheduleConfigResponse | null;
circuit_breaker?: CircuitBreakerConfigResponse | null;
};
export type LoadTestStatsResponse = { export type LoadTestStatsResponse = {
total: number; total: number;
pending: number; pending: number;
@ -949,7 +973,7 @@ export type ToolResponse = {
*/ */
export type TransferCallConfig = { export type TransferCallConfig = {
/** /**
* Phone number to transfer the call to (E.164 format, e.g., +1234567890) * Phone number or SIP endpoint to transfer the call to (E.164 format e.g., +1234567890, or SIP endpoint e.g., PJSIP/1234)
*/ */
destination: string; destination: string;
/** /**
@ -1058,6 +1082,7 @@ export type UpdateCampaignRequest = {
retry_config?: RetryConfigRequest | null; retry_config?: RetryConfigRequest | null;
max_concurrency?: number | null; max_concurrency?: number | null;
schedule_config?: ScheduleConfigRequest | null; schedule_config?: ScheduleConfigRequest | null;
circuit_breaker?: CircuitBreakerConfigRequest | null;
}; };
/** /**
@ -1574,6 +1599,35 @@ export type HandleCloudonixStatusCallbackApiV1TelephonyCloudonixStatusCallbackWo
200: unknown; 200: unknown;
}; };
export type HandleCloudonixAmdCallbackApiV1TelephonyCloudonixAmdCallbackWorkflowRunIdPostData = {
body?: never;
path: {
workflow_run_id: number;
};
query?: never;
url: '/api/v1/telephony/cloudonix/amd-callback/{workflow_run_id}';
};
export type HandleCloudonixAmdCallbackApiV1TelephonyCloudonixAmdCallbackWorkflowRunIdPostErrors = {
/**
* Not found
*/
404: unknown;
/**
* Validation Error
*/
422: HttpValidationError;
};
export type HandleCloudonixAmdCallbackApiV1TelephonyCloudonixAmdCallbackWorkflowRunIdPostError = HandleCloudonixAmdCallbackApiV1TelephonyCloudonixAmdCallbackWorkflowRunIdPostErrors[keyof HandleCloudonixAmdCallbackApiV1TelephonyCloudonixAmdCallbackWorkflowRunIdPostErrors];
export type HandleCloudonixAmdCallbackApiV1TelephonyCloudonixAmdCallbackWorkflowRunIdPostResponses = {
/**
* Successful Response
*/
200: unknown;
};
export type HandleVobizHangupCallbackByWorkflowApiV1TelephonyVobizHangupCallbackWorkflowWorkflowIdPostData = { export type HandleVobizHangupCallbackByWorkflowApiV1TelephonyVobizHangupCallbackWorkflowWorkflowIdPostData = {
body?: never; body?: never;
headers?: { headers?: {
@ -3593,7 +3647,7 @@ export type SaveTelephonyConfigurationApiV1OrganizationsTelephonyConfigPostRespo
200: unknown; 200: unknown;
}; };
export type GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetData = { export type GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetData = {
body?: never; body?: never;
headers?: { headers?: {
authorization?: string | null; authorization?: string | null;
@ -3601,10 +3655,10 @@ export type GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetData = {
}; };
path?: never; path?: never;
query?: never; query?: never;
url: '/api/v1/organizations/campaign-limits'; url: '/api/v1/organizations/campaign-defaults';
}; };
export type GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetErrors = { export type GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetErrors = {
/** /**
* Not found * Not found
*/ */
@ -3615,16 +3669,16 @@ export type GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetErrors = {
422: HttpValidationError; 422: HttpValidationError;
}; };
export type GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetError = GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetErrors[keyof GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetErrors]; export type GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetError = GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetErrors[keyof GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetErrors];
export type GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetResponses = { export type GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetResponses = {
/** /**
* Successful Response * Successful Response
*/ */
200: CampaignLimitsResponse; 200: CampaignDefaultsResponse;
}; };
export type GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetResponse = GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetResponses[keyof GetCampaignLimitsApiV1OrganizationsCampaignLimitsGetResponses]; export type GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetResponse = GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetResponses[keyof GetCampaignDefaultsApiV1OrganizationsCampaignDefaultsGetResponses];
export type GetSignedUrlApiV1S3SignedUrlGetData = { export type GetSignedUrlApiV1S3SignedUrlGetData = {
body?: never; body?: never;