Overview
The problem and what I owned
Agencies run sales across WhatsApp threads, supplier PDFs, phone notes, and spreadsheets. Leads go cold while agents assemble itineraries by hand, and adding LLM features without metering makes costs unpredictable. The goal was one system for the sales and fulfilment loop, with AI built into the steps agents already take and billed like any other commercial SaaS resource.
Role & ownership
- Designed the domain model, multi-tenancy, and auth.
- Wrote the Spring Boot backend (3 Maven modules), the React web app, and the AWS setup, including two supporting Lambdas.
- Designed and implemented every AI workflow: prompts, JSON contracts, grounding, metering.
- Integrated Meta (WhatsApp, Lead Ads), Gmail/IMAP, Sarvam, Anthropic, and FCM.
- Ran customer demos with agencies and turned their feedback into product constraints.
[TODO: note any contributors (e.g. design, Flutter app) if applicable]
Core requirements
- Agency-scoped tenancy for auth, roles, and data.
- Lead → trip → package → payment as the core loop.
- Public surfaces (webhooks, share pages) kept apart from authenticated APIs.
- AI that is grounded in agency data, reviewable, metered, and async where needed.
- Integrations that tolerate retries and partial failure.
Major workflows
- 01
Capture
Leads arrive from Meta Lead Ads webhooks, WhatsApp, Gmail/IMAP, or manual entry; lead-form submissions are scored by Claude.
- 02
Qualify & plan
A lead becomes a trip. Agents build day-by-day packages by hand, by voice, from a supplier PDF, or through the WhatsApp assistant.
- 03
Propose
Packages have versions and options, policies, media, and a public share page for the customer.
- 04
Close & collect
Cost sheets, payment links with UPI QR codes, and AI extraction of payment screenshots into receipts.
- 05
Follow up
Automations send templates, assign owners, create tasks, and stop when the customer replies.
Architecture
One backend, one domain, AI as a layer inside it
AI is part of the CRM, not a separate service. Voice, WhatsApp, and PDF workflows read and write the same agency-scoped entities as the manual UI, and all of them go through one metering path.
Technology choices
| Choice | How it's used | Why |
|---|---|---|
| Spring Boot modular monolith | Maven modules: common (auth, entities, AWS clients), crm, mobile. Each app deploys as a WAR. | AI workflows, billing, and CRM writes share one database transaction model. For one engineer, module boundaries give separation without distributed-system overhead. |
| MySQL + JPA/Hibernate | ddl-auto=none; schema changes are hand-written SQL migrations. | Relational domain (agency → users → leads → trips → packages → payments) with strong consistency needs. |
| SQS FIFO | WhatsApp inbound and PDF-scan queues; scheduled pollers with semaphore-bounded workers. | Managed queue with per-group ordering, which matches per-conversation processing. Kafka-level throughput isn't needed. |
| Claude over plain HTTP | Spring WebClient to /v1/messages; Haiku 4.5 by default, Sonnet 4.6 for PDF scans and hotel onboarding. | Full control over request shape, retries, and usage parsing for metering. Model choice is per workflow, by cost. |
| Sarvam speech-to-text | saaras:v3, ≤30s per sync clip, parallel multi-clip. | Agents dictate in Indian languages and mixed language; Sarvam's models target that. |
| Elastic Beanstalk | WAR bundles with .platform / .ebextensions; Tomcat; JSON logs to CloudWatch. | Managed deploys and scaling with little ops work for a solo engineer. |
| React 18 + Vite + TypeScript | TanStack Query for server state, Redux Toolkit for client state; separate Flutter app. | Typed contracts against a large API surface; query caching suits list-heavy CRM screens. |
Multi-tenancy
- The tenant is the agency. RSA-signed access tokens carry agencyId and role.
- On every request, a servlet filter checks that the session is active, the user exists within that agency and is active, and the agency itself is ACTIVE. Otherwise it returns 401 and revokes refresh and device tokens.
- Repositories are queried by ID and agencyId (findByIdAndAgencyId), so a guessed ID from another tenant returns nothing.
- Unauthenticated surfaces (Meta webhooks, share pages, demo booking) live under /public/; webhooks verify HMAC signatures.
- Subscription expiry suspends the agency, and the same filter then blocks access, so billing and tenancy share one enforcement point.
Trade-off: scoping is by convention in each query, not a framework filter. It's explicit and easy to read, but one missed query leaks data. Hardening this is first on the list below.
AWS infrastructure & operations
- Compute: Elastic Beanstalk (Tomcat) for the crm and mobile WARs; two Java 21 Lambdas.
- Data: RDS MySQL; S3 for PDFs, media, and galleries.
- Messaging: SQS FIFO for WhatsApp inbound and PDF scans; SES for email; SNS for alerts.
- Secrets: Secrets Manager for DB and provider keys; KMS for field-level encryption of stored credentials.
- Observability: JSON logs (Logstash encoder) with request, user, and agency IDs in MDC → CloudWatch. Tagged alert lines → Lambda → SNS.
- Scheduled work: Spring schedulers (subscription lifecycle, automation resume, Gmail sync) and a scheduled Lambda for booking expiry and token cleanup.
[TODO: instance count / sizing and environments (dev, prod) if you want to show them]
Deep dives
Four systems in detail
Each section follows the same order: problem, requirements, workflow, decisions, edge cases, code, trade-offs, and known limitations. Code excerpts are from the production repository, trimmed for length.
Deep dive 01
Voice → itinerary drafts
Applied AI inside a production workflow: speech becomes structured intent, then a database-grounded draft the user approves.
Problem
Agents capture trip requirements on calls and on the move. Spoken requests are ambiguous (“the Munnar one, but add a houseboat night”), and writing model output straight into packages, payments, or tasks is unsafe. In demos, agencies consistently preferred reviewing a draft over automation that changes records silently.
Requirements
- Accept one or more audio clips, or a typed chat message, through the same pipeline
- Classify intent: save note, add task, record payment, create package/template, or ask a question about CRM data
- Resolve spoken references to real leads, trips, destinations, hotels, and sightseeing in the agency's catalog
- Ask targeted clarification questions instead of guessing, with a hard cap
- Never write data. Return a draft that the client confirms through existing APIs
- Meter every STT and LLM call for billing
- 1Validate & scope
JWT gives userId and agencyId. Clip count/size validated (max 5 clips; 1 on follow-ups). User loaded with findByIdAndAgencyId.
- 2Open usage operation
ClaudeUsageService.beginSoft binds an operation to the thread; every provider call becomes a usage hit.
- 3Speech-to-textExternal
Sarvam saaras:v3. Clips are transcribed in parallel on a dedicated 5-thread pool and merged.
- 4Intent + extractionLLM
Two-stage Claude pipeline: classify the intent, then extract fields into an intent-specific JSON shape.
- 5Entity resolutionMySQL
Backend matchers (stay, sightseeing, country/destination) map names to catalog rows. The model never supplies IDs.
- 6Clarify or draft
Missing or unresolved fields → NEEDS_CLARIFICATION with conversation state. Otherwise a preview-ready draft.
- 7Finalize billing
In finally: timing log, endSoft() freezes credits_used and debits the wallet, MDC cleared.
- 8User confirms in UI
The frontend calls the existing create-note / task / receipt / package endpoints. Same validation and permissions as manual entry.
Key technical decisions
The AI path has no write access
VoiceProcessService never calls a write API. The response is always a draft, so AI-created records pass the same validation and permission checks as manual ones and need no separate audit path.Stateless clarification
ChatConversationState object is returned to the client and echoed back. That avoids session storage and expiry logic, and any backend instance can serve the next turn. Rounds are capped at MAX_CLARIFICATION_ROUNDS = 3, shared by spoken and typed replies.Backend owns entity resolution
Contextual vs universal entry
CONTEXTUAL with a known intent and entity, which skips classification. The global mic is UNIVERSAL and classifies first. Same pipeline, less model work when context exists.Edge cases handled
| Scenario | What the system does |
|---|---|
| Hotel, day, or destination missing from speech | Returns NEEDS_CLARIFICATION with the specific question instead of filling a default. |
| Spoken name doesn't match the catalog | Matcher fails to resolve → backend adds a mandatory clarification; the model's spelling is never saved as a new entity. |
| User keeps answering vaguely | After 3 rounds the flow stops and sends the user to the manual edit screen. |
| Clarification answered by typing instead of speaking | Structured answers skip STT entirely; the same state object continues the flow. |
| Long dictation | Sarvam's sync API accepts ≤30s per clip, so the client sends up to 5 clips, transcribed in parallel and merged into one run. |
Code
/**
* Sarvam STT -> Claude intent classification/extraction -> entity resolution ->
* precondition checks -> a preview-ready draft. No write API is ever called from here —
* the response is always a draft the frontend shows the user for confirmation ...
*/
static final int MAX_CLARIFICATION_ROUNDS = 3;
static final int MAX_VOICE_CLIPS_PER_REQUEST = 5;
public ApiResponse<VoiceProcessResponse> process(List<MultipartFile> audioFiles,
VoiceProcessRequest payload, HttpServletRequest httpRequest, String requestId) {
long overallStartMs = System.currentTimeMillis();
try {
Long userId = JwtAuthUtil.getTokenUserId(httpRequest);
Long agencyId = JwtAuthUtil.getAgencyId(httpRequest);
MDC.put(LoggingFields.REQUEST_ID.getField(), requestId);
MDC.put(LoggingFields.AGENCY_ID.getField(), agencyId.toString());
// … validate clips, load user scoped to agency …
claudeUsageService.beginSoft(ClaudeUsageLabels.OPERATION_VOICE_PROCESS,
agencyId, userId, claudeLlmProperties.resolveModel(), null, null);
// … CONTEXTUAL / UNIVERSAL / follow-up flows, each returning a draft …
} finally {
log.info("[VOICE_PROCESS_TIMING] stage=TOTAL durationMs={}",
System.currentTimeMillis() - overallStartMs);
claudeUsageService.endSoft();
MDC.clear();
}
}Trade-offs
- Preview, then confirm over auto-save. One extra tap, but agencies trust the output and no write path bypasses validation.
- Several specialised prompts over one large prompt. Each stage has a narrow JSON contract that is easier to debug and cheaper; most stages run on Haiku.
- Client-held conversation state over server sessions. No session store or TTL to manage; the cost is a larger request payload.
Validation & evidence
- Per-stage timing logs ([VOICE_PROCESS_TIMING] stage=…) with the request ID in MDC, queryable in CloudWatch.
- Each run links request → usage operation → STT and Claude hits → credits, so a bad draft can be traced back to the exact model calls and their cost.
- Agency demos mainly exposed wrong entity matches and over-eager actions, which led to backend-owned resolution and the draft-only rule.
- Measured latency per stage: [TODO: p50/p95 from CloudWatch, if you want to publish them]
Known limitations & what I'd build next
- Pre-call credit check: today usage is billed after the run, so an empty wallet doesn't stop STT/LLM spend.
- Client-supplied idempotency key so a mobile retry can't create a second billed run.
- Golden-transcript regression tests for extraction and entity resolution.
In the demo, use the mic to dictate a trip and watch the draft and any clarification questions before anything is saved.
Deep dive 02
Event-driven WhatsApp sales assistant
Webhook ingestion, queue-based processing, and an LLM orchestrator grounded in each agency's catalog, with escalation to a human.
Problem
Agency deals happen on WhatsApp, but the context (packages, hotels, prices) lives in the CRM. Replying from inside the webhook request ties the response to model latency and Meta's retry behaviour, and duplicate deliveries become duplicate customer messages. Agencies were clear that the assistant should hand off to a person rather than bluff.
Requirements
- Verify Meta webhook signatures; acknowledge fast
- Process asynchronously with ordering per customer conversation and bounded concurrency
- Ignore duplicate deliveries of the same message
- Route by intent: package questions, new package, modify destinations/stays/transport/sightseeing, escalate
- Ground replies in the agency's own packages, hotels, and transport
- Respect WhatsApp's 24-hour customer-service window and template rules
- Notify a human on low confidence; stop drip automations when the customer replies
- 1WebhookExternal
POST /public/webhooks/whatsapp. HMAC-SHA256 of the raw body checked against X-Hub-Signature-256 (constant-time compare).
- 2EnqueueSQS
Raw payload sent to an SQS FIFO queue. MessageGroupId = our number + customer number; dedup ID = SHA-256 of the payload.
- 3Poll with backpressureSQS
Scheduled long-poll sizes each receive batch to free slots in a semaphore (default 5) and dispatches to a thread pool.
- 4Persist inboundMySQL
Skip if provider_message_id (wamid) already exists; otherwise save the message and cancel the lead's waiting automations.
- 5Classify intentLLM
Claude returns intent_type, needs_clarification, needs_human, and a package query.
- 6Route to a flowLLM
Package create/modify reuses the voice package pipeline through a bridge service. Hotel, transport, and Q&A stages each have their own prompt and context JSON.
- 7Persist, publish, reply
Updated packages are saved and a share link is added to the reply. Messaging-window rules decide between a session message and a template.
- 8Bill or hand off
Usage closes against the conversation. needs_human or confidence < 0.4 → handoff notification (in-app + FCM push).
Key technical decisions
FIFO groups per conversation
At-least-once, with duplicate suppression
Orchestration in backend code, not tool-calling
Reuse the voice package pipeline
WhatsAppVoicePackageBridgeService sends package create/modify requests through the same matchers and draft builder as voice. One grounding implementation serves both channels.Edge cases handled
| Scenario | What the system does |
|---|---|
| Meta retries the same webhook | Identical payload → same SQS dedup ID within 5 minutes; after that, the wamid existence check drops it. |
| Burst of messages from one customer | Same FIFO group → processed strictly in order; context from earlier messages is in the DB before later ones run. |
| Workers saturated | Receive batch shrinks to free semaphore slots; if dispatch is refused, the message stays in flight and is redelivered after the visibility timeout. |
| Processing throws | Message is not deleted, so SQS redelivers it after the visibility timeout. |
| Low-confidence or unsupported request | needs_human or confidence < 0.4 → customer told a person will follow up; agent notified in-app and via push. |
| Outside the 24h service window | WhatsAppMessagingWindowService switches to SESSION, TEMPLATE_ONLY_FREE (72h free-entry window), or TEMPLATE_ONLY_PAID. |
| Customer replies mid-drip | Waiting automation enrollments for that lead are cancelled with reason CUSTOMER_REPLIED. |
Code
if (whatsAppInboundSqsQueueUrlProvider.isFifoQueue()) {
requestBuilder
.messageGroupId(resolveFifoMessageGroupId(payloadJson))
.messageDeduplicationId(sha256Hex(payloadJson));
}
// resolveFifoMessageGroupId:
// - messages webhooks → our line + customer WhatsApp id (ordering per thread)
// - statuses-only → provider message id, so one message's lifecycle stays ordered
JsonNode messages = value.path("messages");
if (messages.isArray() && !messages.isEmpty()) {
String from = messages.get(0).path("from").asText("");
String ourLine = !displayPhoneNumber.isBlank()
? displayPhoneNumber
: (!phoneNumberId.isBlank() ? "pid_" + phoneNumberId : "");
if (!ourLine.isBlank() && !from.isBlank()) {
return sanitizeFifoToken("wa_msg_" + ourLine + "_" + from);
}
}@Scheduled(fixedDelayString = "${hj.sqs.whatsapp-inbound.poll-interval-ms:1000}")
public void poll() {
int batchSize = SqsWorkerSupport.receiveBatchSize(
whatsAppInboundSqsInFlight, properties.getMaxNumberOfMessages());
if (batchSize <= 0) {
return;
}
// … long-poll receive (≤20s) …
for (Message message : messages) {
boolean accepted = SqsWorkerSupport.dispatch(
whatsAppInboundSqsInFlight, whatsAppInboundSqsExecutor,
() -> processMessage(message));
if (!accepted) {
log.warn("[BUSINESS_ALERT] SQS whatsapp worker at capacity; leaving message in-flight ...");
}
}
}
private void processMessage(Message message) {
try {
// … deserialize …
whatsAppInboundSqsMessageHandler.handleInboundWebhookPayload(payload.getPayloadJson());
deleteMessage(message); // only on success
} catch (Exception e) {
log.error("[CRITICAL_ALERT] Failed processing whatsapp inbound SQS message ...", e);
}
}Trade-offs
- Async replies via SQS over replying inside the webhook. Adds a few seconds of latency, but model slowness can't cause Meta timeouts and retries.
- Multi-stage orchestrator over a thin chat bot. Much more code (the orchestrator alone is ~3,000 lines), but every customer-visible action is deterministic and logged.
- Escalate on low confidence over always answering. Fewer automated replies, more agency trust.
Validation & evidence
- Failures are logged with [CRITICAL_ALERT] / [BUSINESS_ALERT] tags. A CloudWatch Logs subscription feeds a separate Java 21 Lambda that publishes them to SNS, so alerts reach me without watching dashboards.
- Every LLM stage logs its input context and output per conversation ID, so any reply can be reconstructed.
- Agency demos and replays focused on duplicate suppression and handoff behaviour.
- Production volume and reply latency: [TODO: add only from real CloudWatch data]
Known limitations & what I'd build next
- Add a database UNIQUE constraint on whatsapp_messages.provider_message_id. Today the existence check is application-level, safe mainly because FIFO serialises each conversation.
- Configure a dead-letter queue with a redrive policy. No DLQ is defined in the codebase, and a poison message in a FIFO group blocks that conversation until it expires.
- Confirm the AWS queue config: [TODO: is a redrive policy set on the queue outside the codebase?]
- Make signature verification fail closed in production. It currently logs a warning and skips the check when the app secret is unset, which is meant for local testing.
- Record outbound sends keyed to the inbound wamid so that a redelivery after a send can't send the reply twice.
Open a conversation in the demo to see how WhatsApp threads, leads, and packages connect.
Deep dive 03
AI usage metering, credits & subscriptions
SaaS economics in backend code: per-call metering, pricing, wallets, and tenant lifecycle.
Problem
LLM and speech costs scale with usage, not seats. Without per-tenant metering, one busy agency can erase the margin, and nobody can explain a bill. Usage also comes from places with no logged-in user, such as WhatsApp webhooks and background jobs, and still needs an owner.
Requirements
- Record every successful provider call (Claude tokens, Sarvam audio seconds) under one parent operation per workflow run
- Convert usage to credits using a configurable tariff (per-model token rates, STT rate, FX, GST)
- Per-user wallets with a monthly expirable allotment and non-expiring top-ups
- Billing failures must never break the user's workflow
- Subscription lifecycle: active → past due with grace → expired → tenant suspended
- Idempotent finalize and idempotent admin grants and transfers
- 1beginSoft(operation)MySQL
Creates an ai_usage_operations row (REQUIRES_NEW tx) and binds its ID to the worker thread.
- 2Provider callsLLM
Each successful Claude/Sarvam response is recorded as a hit: input/output tokens or billable duration. Cache hits record nothing.
- 3endSoft()
Unbinds the thread and finalizes credits in a separate transaction; errors are logged, never thrown.
- 4Price
Σ hits × tariff (Haiku vs Sonnet rates; STT hours × INR/hr ÷ FX) × GST → USD → ceil(USD ÷ credit_usd).
- 5Debit walletMySQL
Expirable first, then non-expiring, then overdraw on expirable. Before/after deltas stored on the operation.
- 6Daily lifecycle job
00:15 IST: period end → PAST_DUE with 3 grace days → EXPIRED → agency SUSPENDED; monthly credit drip per seat; history rows.
- 7Every request
JWT filter rejects users of non-ACTIVE agencies and revokes their refresh and device tokens.
Key technical decisions
Meter at the call site, bill at the workflow boundary
Usage writes in REQUIRES_NEW transactions
Credits instead of raw cost pass-through
Attribute system usage to the agency owner
Edge cases handled
| Scenario | What the system does |
|---|---|
| finalize called twice for one operation | No-op if credits_used is already set. |
| Wallet runs out mid-workflow | Run completes; the overdraw lands on the expirable pool, which can go negative. This is a deliberate policy: never cut off a reply mid-conversation. |
| Usage from a webhook with no user | Attributed to the agency owner; logged as a business alert if no owner exists. |
| Tariff row missing | Critical alert, no debit. Usage rows remain, so the run can be re-billed later. |
| Subscription period ends | PAST_DUE with a 3-day grace countdown written to history each day; then EXPIRED and the agency is SUSPENDED. |
| Suspended tenant still has a valid token | JWT filter checks agency status on every request; returns 401 and revokes refresh/device tokens. |
| Retried credit grant or transfer | Unique idempotency_key on credit transactions; transfers check source balance first. |
Code
/**
* Freezes ai_usage_operations.credits_used from hits × credit_tariff,
* then debits the attributed user's wallet (expirable → non-expiring → negative expirable).
*/
@Transactional
public void finalizeOperationCredits(Long operationId) {
// …
ClaudeUsageOperation op = opOpt.get();
if (op.getCreditsUsed() != null) {
return; // idempotent
}
List<ClaudeUsageHit> hits = hitRepository.findByOperationIdOrderByIdAsc(operationId);
int credits = computeCredits(op.getModel(), hits, tariff);
op.setCreditsUsed(credits);
DebitSnapshot snapshot = debitWallet(op.getCrmUserId(), op.getAgencyId(), credits);
op.setExpirableDelta(snapshot.expirableDelta());
op.setNonExpiringDelta(snapshot.nonExpiringDelta());
// …
}
private DebitSnapshot debitWallet(Long crmUserId, Long agencyId, int credits) {
// …
if (remaining > 0 && exp > 0) { takenExp = Math.min(exp, remaining); exp -= takenExp; remaining -= takenExp; }
if (remaining > 0 && non > 0) { takenNon = Math.min(non, remaining); non -= takenNon; remaining -= takenNon; }
if (remaining > 0) {
// Overdraw lands on expirable (may go negative).
takenExp += remaining;
exp -= remaining;
}
// …
}Long agencyId = jwtTokenService.getAgencyIdFromToken(claims);
CrmUser user = crmUserRepository.findByIdAndAgencyId(userId, agencyId).orElse(null);
// … reject inactive users …
AgencyStatus agencyStatus = user.getAgency() != null ? user.getAgency().getStatus() : null;
if (agencyStatus == null || agencyStatus != AgencyStatus.ACTIVE) {
log.warn("[CRITICAL_ALERT] CRM access blocked for agency {} (status={}) user {}", ...);
response.setStatus(HttpServletResponse.SC_UNAUTHORIZED);
deviceTokenService.deleteAllTokensForUser(UserType.CRM, userId);
refreshTokenService.deleteAllTokensForUser(UserType.CRM, userId);
refreshTokenService.clearRefreshTokenCookie(response);
return;
}
request.setAttribute(JwtClaim.AGENCY_ID.getClaimName(), agencyId);Trade-offs
- Bill after the run, allow overdraw over blocking before each call. A customer conversation is never cut off mid-reply. The cost is bounded by subscription status, not by balance.
- Store pricing inputs in the database over hard-coded rates. Rates, FX, and GST change; more rows to manage, but no deploy needed.
- One wallet per user over one pool per agency. Supports per-agent allocation and reporting; transfers add complexity.
Validation & evidence
- Each operation stores credits_used plus wallet deltas and post-debit balances, so a balance can be reconciled from usage history.
- Subscription transitions write history events (PAST_DUE, GRACE_TICK, EXPIRED) for support and audit.
- Automated tests for billing paths: [TODO: add once written (none exist today)]
Known limitations & what I'd build next
- Concurrency hardening: wallet debits are read-modify-write with no row lock or @Version, and the credits_used check is check-then-set. Two concurrent finalizes on one wallet could lose an update. Fix with a pessimistic lock on credit_balance, or a conditional UPDATE … WHERE credits_used IS NULL.
- Add a pre-call balance gate (or reservation) for expensive workflows such as PDF scans, while keeping overdraw for live conversations.
- Add concurrency tests that run parallel finalize calls against one wallet.
Credit balances and AI usage history are shown in the web app. Run a voice command in the demo, then check usage.
Deep dive 04
Automations runtime
Workflow engineering: a state machine, delayed resumption, audit records, and exits driven by humans.
Problem
Agencies want consistent follow-up when leads arrive or change status: a WhatsApp template, an email, an owner assignment, a task. A visual builder is the easy part. The runtime has to know when to wait, when to stop, and what happened. In discovery calls the clearest requirement was “don't keep messaging after the customer has replied.”
Requirements
- Triggers: lead imported, lead status changed; conditional branches
- At most one open (active or waiting) journey per lead
- ACTION and WAIT steps; resume after the wait without holding threads
- Actions: WhatsApp template, email template, assign owner (round-robin), change status, create task, notify agent
- Per-step audit trail
- Exit when the customer replies; skip automated sends when an agent has taken over
- 1CRM event
Lead import or status change calls AutomationEnrollmentService.
- 2EnrollMySQL
Skip if the lead already has an ACTIVE/WAITING enrollment. Match enabled workflows for the agency and trigger; choose a branch.
- 3Execute
Run up to 20 steps per pass. ACTION → executor → AutomationStepRun (DONE / SKIPPED / FAILED) → advance.
- 4WAIT
Check for a customer reply first; otherwise status = WAITING, next_run_at = now + duration (default 24h).
- 5Resume
Scheduler every 30s loads due WAITING enrollments, re-checks for replies, advances past the WAIT, continues.
- 6Exit
COMPLETED at end of branch; CANCELLED with a reason (CUSTOMER_REPLIED, MANUAL, …).
Key technical decisions
The enrollment row is the state machine
Human signals come from the messaging data
customerReplied looks for inbound WhatsApp after the latest outbound. agentTakeover looks for outbound WhatsApp or email after enrollment that this enrollment's own step runs didn't produce, by comparing message IDs stored in step results.Step runs as the audit log
Round-robin under a row lock
Edge cases handled
| Scenario | What the system does |
|---|---|
| Lead already in a journey | New enrollment skipped and logged. |
| Customer replies while waiting | Cancelled with CUSTOMER_REPLIED, both on the inbound path and when the wait expires. |
| Agent messages the customer manually | Later WhatsApp/email steps are recorded as SKIPPED with AGENT_TAKEOVER. |
| Lead deleted mid-journey | Enrollment cancelled (MANUAL). |
| Misconfigured branch loops | MAX_STEPS_PER_PASS = 20 bounds a single pass. |
| Two leads assigned at the same instant | Locked cursor gives each one the next agent in turn; no duplicate picks. |
Code
@Transactional
public void processEnrollment(Long enrollmentId) {
AutomationEnrollment enrollment = /* … */;
if (enrollment.getStatus().isTerminal()) {
return;
}
if (enrollment.getStatus() == AutomationEnrollmentStatus.WAITING) {
if (conversationSignals.customerReplied(enrollment.getLeadId(), enrollment.getStartedAt())) {
cancel(enrollment, AutomationCancelReason.CUSTOMER_REPLIED);
return;
}
if (enrollment.getNextRunAt() != null && enrollment.getNextRunAt().isAfter(LocalDateTime.now())) {
return;
}
advancePastCurrent(enrollment); // wait completed
enrollment.setStatus(AutomationEnrollmentStatus.ACTIVE);
enrollment.setNextRunAt(null);
enrollmentRepository.save(enrollment);
}
int guard = 0;
while (enrollment.getStatus() == AutomationEnrollmentStatus.ACTIVE && guard++ < MAX_STEPS_PER_PASS) {
// … load step; WAIT → enterWait(...) and return …
AutomationActionExecutor.ActionResult result = actionExecutor.execute(enrollment, step, lead);
recordRun(enrollment, step, result);
advancePastCurrent(enrollment);
enrollmentRepository.save(enrollment);
}
}AutomationRoundRobinCursor cursor = cursorRepository
.findForUpdate(agencyId, workflowId, stepId) // @Lock(PESSIMISTIC_WRITE)
.orElse(null);
if (cursor == null) {
try {
cursorRepository.saveAndFlush(newCursor(agencyId, workflowId, stepId));
} catch (DataIntegrityViolationException ignored) {
// Concurrent create — fall through to locked read.
}
cursor = cursorRepository.findForUpdate(agencyId, workflowId, stepId).orElseThrow(...);
}
int index = Math.floorMod(cursor.getNextIndex(), userIds.size());
cursor.setNextIndex((index + 1) % userIds.size());Trade-offs
- Polling scheduler on a DB column over a delay queue per wait. Simple, durable, inspectable with SQL; resolution is ~30s, which is fine for hour-scale waits.
- Conversation-aware exits over pure timer drips. Fewer messages sent, but the engine now depends on the messaging tables.
Validation & evidence
- An enrollment's history (trigger → step runs → skips → cancel reason) can be reconstructed from the database alone.
- Agency demos and discovery calls drove the reply-cancel and agent-takeover rules.
- Automated tests for WAIT/cancel paths: [TODO: add once written (none exist today)]
Known limitations & what I'd build next
- Add a FAILED-step policy: today a failed action is recorded and the journey still moves on. Options are retry with backoff, pause, or notify.
- Claim due enrollments with SELECT … FOR UPDATE SKIP LOCKED (or a lease column) before running more than one backend instance. The scheduler currently has no claim step.
- Confirm the deployment: [TODO: number of Elastic Beanstalk instances running the scheduler]
- Store a per (enrollment, step) idempotency key on outbound sends so a retried step can't send twice.
- Also cancel ACTIVE enrollments on reply. The inbound hook only cancels WAITING ones; active ones are checked when they reach their next WAIT.
In the demo, open Automations to see workflows, branches, and per-lead step history.
Supporting capabilities
The rest of the product
Not covered in depth, but part of the same system and data model.
PDF → package
Supplier brochure → S3 → SQS FIFO → staged Claude (Sonnet) pipeline: extract itinerary, enrich media, create package, add days, finalize. Each stage reports progress and records failures.
Lead capture
Meta Lead Ads webhook (signature-verified, de-duplicated per agency), Gmail API OAuth sync every 2 minutes, IMAP/SMTP for other providers.
Lead scoring
Claude Haiku scores incoming lead-form submissions with a dedicated prompt template.
Payments
Cost sheets, payment-request links, UPI QR generation, and Claude extraction of payment screenshots.
Package authoring
Day-by-day itineraries, versions/options, policies, media galleries on S3, public share pages.
Catalogs
Destinations, sightseeing, properties, vendors, transport. These ground every AI flow.
Notifications
Firebase Cloud Messaging for handoffs, bookings, and automation alerts.
Security
RSA-signed JWTs with refresh sessions, TOTP MFA, KMS-encrypted mail passwords, banking details, and TOTP secrets, Secrets Manager for all credentials.
Product development
How customer feedback shaped the engineering
I ran product demos and ongoing conversations with travel agencies throughout development. The feedback below led directly to design constraints. Agency counts and commercial figures are left out until I can verify and publish them.
| What agencies told me | What it changed in the system |
|---|---|
| “We sell on WhatsApp, not through web forms.” | WhatsApp Cloud API as a first-class channel on the agency's own number, with an AI assistant and full CRM context. |
| “Don't let AI change our data without asking.” | Voice produces drafts only; existing APIs write after the agent confirms. |
| “Answers have to come from our packages and our hotels.” | Backend entity resolution against each agency's catalog; the model never supplies IDs. |
| “A person must be able to step in.” | Low-confidence handoff with push notifications; automations skip sends after an agent takes over. |
| “Stop follow-ups once the customer replies.” | Reply-driven cancellation of automation enrollments. |
| “Supplier inventory arrives as PDFs.” | The staged PDF → package pipeline. |
Verified outcomes
- PeakInsight is live in production at peakinsight.in, with a public demo that needs no signup.
- The four systems above are shipped and running on the same multi-tenant backend.
- [TODO: agencies onboarded / paying. Include only if verified and OK to publish]
- [TODO: one measured result, e.g. time to build an itinerary by voice vs manually, with how it was measured]
A specific example
[TODO: one short story from a demo: what an agency tried, what broke or confused them, and what you changed in response]
Evolving the system
What I'd improve next
Each deep dive lists its own limitations. These are the cross-cutting ones, roughly in priority order. Most come from the system being built quickly by one person and are cheap to fix now.
- 01Adopt Flyway so schema migrations are versioned and applied automatically instead of run by hand.
- 02Enforce tenant scoping at the framework level (Hibernate filter or a scoped repository base) and add tests that fail on unscoped queries. Today isolation relies on every query including agencyId.
- 03Build a test suite around the riskiest paths: concurrent billing, webhook redelivery, and automation resume. Only a small set of unit tests exists today.
- 04Add a DLQ and idempotent-send records for every external side effect (WhatsApp, email), tracked as one reliability workstream.
- 05Publish latency and reliability numbers once there's enough production volume to make them meaningful.