Blame | Last modification | View Log | RSS feed
package com.spice.profitmandi.service.lms;import com.spice.profitmandi.dao.entity.user.LeadCall;import com.spice.profitmandi.dao.repository.dtr.LeadCallRepository;import com.spice.profitmandi.service.storage.ObjectStoreService;import org.apache.logging.log4j.LogManager;import org.apache.logging.log4j.Logger;import org.springframework.beans.factory.annotation.Autowired;import org.springframework.scheduling.annotation.Async;import org.springframework.stereotype.Service;import org.springframework.transaction.annotation.Propagation;import org.springframework.transaction.annotation.Transactional;import java.io.ByteArrayInputStream;import java.io.ByteArrayOutputStream;import java.io.InputStream;import java.time.LocalDateTime;import java.time.format.DateTimeFormatter;/*** Archives a single call recording into the object store.** <p><b>Why this is its own bean.</b> Spring's {@code @Async} and {@code @Transactional} are applied* by a proxy around the bean, so a self-invocation from a sweep method in the same class would run* inline on the scheduler thread, in the scheduler's transaction — silently losing both annotations.* Keeping the per-item work in a separate injected bean is what makes them actually take effect.*/@Servicepublic class RecordingArchiveWorker {private static final Logger LOGGER = LogManager.getLogger(RecordingArchiveWorker.class);private static final DateTimeFormatter KEY_DATE = DateTimeFormatter.ofPattern("yyyy/MM/dd");/*** Guard against a runaway download. Vonage mp3 runs roughly 1 MB/minute, so this allows a very* long call while still refusing to buffer something pathological into heap.*/private static final long MAX_RECORDING_BYTES = 200L * 1024 * 1024;@Autowiredprivate LeadCallRepository leadCallRepository;@Autowiredprivate ObjectStoreService objectStoreService;@Autowiredprivate LmsDialerProvider dialerProvider;/** Separate bean on purpose — see its javadoc; a self-call here would lose REQUIRES_NEW. */@Autowiredprivate RecordingArchiveFailureRecorder failureRecorder;/*** Fetch one recording from the provider and store it.** <p>Runs in its own transaction so one bad recording cannot roll back a whole sweep. Failures* are logged and swallowed: the row keeps its null {@code recording_document_id}, stays in the* pending set, and the next sweep retries it. That is the retry mechanism — there is no queue to* lose across a restart.*/@Async("recordingArchiveExecutor")@Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Throwable.class)public void archiveOne(int leadCallId) {LeadCall call = leadCallRepository.selectById(leadCallId);if (call == null || call.getRecordingUuid() == null || call.getRecordingObjectKey() != null) {return;}InputStream stream = null;try {stream = dialerProvider.fetchRecording(call);if (stream == null) {// Counted as an attempt: at the provider's retention boundary this becomes permanent,// and without the counter the row would be retried every ten minutes forever.LOGGER.warn("Provider has no recording for call {}", leadCallId);failureRecorder.recordFailure(leadCallId);return;}byte[] audio = readCapped(stream);if (audio.length == 0) {LOGGER.warn("Empty recording body for call {} — leaving it pending", leadCallId);failureRecorder.recordFailure(leadCallId);return;}String key = objectKey(call);objectStoreService.put(key, new ByteArrayInputStream(audio), audio.length, "audio/mpeg");// The key goes on the call row itself. It used to be written to dtr.document.path, which// is VARCHAR(30) against a ~58-character key — under STRICT_TRANS_TABLES that threw on// every single upload, and because failures here are swallowed to drive the retry, the// symptom was an endless loop that stored nothing while the provider window expired.call.setRecordingObjectKey(key);leadCallRepository.persist(call);LOGGER.info("Archived recording for call {} to {} ({} bytes)", leadCallId, key, audio.length);} catch (Exception e) {LOGGER.error("Could not archive the recording for call {} — will retry", leadCallId, e);failureRecorder.recordFailure(leadCallId);} finally {closeQuietly(stream);}}/*** {@code yyyy/MM/dd/lead-<leadId>/call-<callId>-<recordingUuid>.mp3}** <p>Date-partitioned so a lifecycle rule can expire a day's prefix without scanning, and the ids* are in the key so an object found on its own is still traceable back to a lead.** <p>No {@code lms-recordings/} prefix: the bucket is already named that, and prefixing would* give every object a path of {@code lms-recordings/lms-recordings/...}.*/private String objectKey(LeadCall call) {LocalDateTime when = call.getStartedAt() != null ? call.getStartedAt() : call.getCreatedTimestamp();if (when == null) {when = LocalDateTime.now();}return KEY_DATE.format(when)+ "/lead-" + call.getLeadId()+ "/call-" + call.getId() + "-" + call.getRecordingUuid() + ".mp3";}/** Buffered because S3 needs an accurate content length up front, but bounded so it cannot OOM. */private byte[] readCapped(InputStream in) throws Exception {ByteArrayOutputStream out = new ByteArrayOutputStream();byte[] buffer = new byte[8192];long total = 0;int read;while ((read = in.read(buffer)) != -1) {total += read;if (total > MAX_RECORDING_BYTES) {throw new IllegalStateException("Recording exceeds " + MAX_RECORDING_BYTES + " bytes");}out.write(buffer, 0, read);}return out.toByteArray();}private void closeQuietly(InputStream stream) {if (stream == null) {return;}try {stream.close();} catch (Exception e) {LOGGER.debug("Ignoring a failure closing the recording stream", e);}}}