Subversion Repositories SmartDukaan

Rev

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.
 */
@Service
public 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;

    @Autowired
    private LeadCallRepository leadCallRepository;

    @Autowired
    private ObjectStoreService objectStoreService;

    @Autowired
    private LmsDialerProvider dialerProvider;

    /** Separate bean on purpose — see its javadoc; a self-call here would lose REQUIRES_NEW. */
    @Autowired
    private 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);
        }
    }
}