Subversion Repositories SmartDukaan

Rev

Details | Last modification | View Log | RSS feed

Rev Author Line No. Line
37651 vikas 1
package com.spice.profitmandi.service.lms;
2
 
3
import com.spice.profitmandi.dao.entity.user.LeadCall;
4
import com.spice.profitmandi.dao.repository.dtr.LeadCallRepository;
5
import com.spice.profitmandi.service.storage.ObjectStoreService;
6
import org.apache.logging.log4j.LogManager;
7
import org.apache.logging.log4j.Logger;
8
import org.springframework.beans.factory.annotation.Autowired;
9
import org.springframework.scheduling.annotation.Async;
10
import org.springframework.stereotype.Service;
11
import org.springframework.transaction.annotation.Propagation;
12
import org.springframework.transaction.annotation.Transactional;
13
 
14
import java.io.ByteArrayInputStream;
15
import java.io.ByteArrayOutputStream;
16
import java.io.InputStream;
17
import java.time.LocalDateTime;
18
import java.time.format.DateTimeFormatter;
19
 
20
/**
21
 * Archives a single call recording into the object store.
22
 *
23
 * <p><b>Why this is its own bean.</b> Spring's {@code @Async} and {@code @Transactional} are applied
24
 * by a proxy around the bean, so a self-invocation from a sweep method in the same class would run
25
 * inline on the scheduler thread, in the scheduler's transaction — silently losing both annotations.
26
 * Keeping the per-item work in a separate injected bean is what makes them actually take effect.
27
 */
28
@Service
29
public class RecordingArchiveWorker {
30
 
31
    private static final Logger LOGGER = LogManager.getLogger(RecordingArchiveWorker.class);
32
 
33
    private static final DateTimeFormatter KEY_DATE = DateTimeFormatter.ofPattern("yyyy/MM/dd");
34
 
35
    /**
36
     * Guard against a runaway download. Vonage mp3 runs roughly 1 MB/minute, so this allows a very
37
     * long call while still refusing to buffer something pathological into heap.
38
     */
39
    private static final long MAX_RECORDING_BYTES = 200L * 1024 * 1024;
40
 
41
    @Autowired
42
    private LeadCallRepository leadCallRepository;
43
 
44
    @Autowired
45
    private ObjectStoreService objectStoreService;
46
 
47
    @Autowired
48
    private LmsDialerProvider dialerProvider;
49
 
50
    /** Separate bean on purpose — see its javadoc; a self-call here would lose REQUIRES_NEW. */
51
    @Autowired
52
    private RecordingArchiveFailureRecorder failureRecorder;
53
 
54
    /**
55
     * Fetch one recording from the provider and store it.
56
     *
57
     * <p>Runs in its own transaction so one bad recording cannot roll back a whole sweep. Failures
58
     * are logged and swallowed: the row keeps its null {@code recording_document_id}, stays in the
59
     * pending set, and the next sweep retries it. That is the retry mechanism — there is no queue to
60
     * lose across a restart.
61
     */
62
    @Async("recordingArchiveExecutor")
63
    @Transactional(propagation = Propagation.REQUIRES_NEW, rollbackFor = Throwable.class)
64
    public void archiveOne(int leadCallId) {
65
        LeadCall call = leadCallRepository.selectById(leadCallId);
66
        if (call == null || call.getRecordingUuid() == null || call.getRecordingObjectKey() != null) {
67
            return;
68
        }
69
 
70
        InputStream stream = null;
71
        try {
72
            stream = dialerProvider.fetchRecording(call);
73
            if (stream == null) {
74
                // Counted as an attempt: at the provider's retention boundary this becomes permanent,
75
                // and without the counter the row would be retried every ten minutes forever.
76
                LOGGER.warn("Provider has no recording for call {}", leadCallId);
77
                failureRecorder.recordFailure(leadCallId);
78
                return;
79
            }
80
            byte[] audio = readCapped(stream);
81
            if (audio.length == 0) {
82
                LOGGER.warn("Empty recording body for call {} — leaving it pending", leadCallId);
83
                failureRecorder.recordFailure(leadCallId);
84
                return;
85
            }
86
 
87
            String key = objectKey(call);
88
            objectStoreService.put(key, new ByteArrayInputStream(audio), audio.length, "audio/mpeg");
89
 
90
            // The key goes on the call row itself. It used to be written to dtr.document.path, which
91
            // is VARCHAR(30) against a ~58-character key — under STRICT_TRANS_TABLES that threw on
92
            // every single upload, and because failures here are swallowed to drive the retry, the
93
            // symptom was an endless loop that stored nothing while the provider window expired.
94
            call.setRecordingObjectKey(key);
95
            leadCallRepository.persist(call);
96
            LOGGER.info("Archived recording for call {} to {} ({} bytes)", leadCallId, key, audio.length);
97
        } catch (Exception e) {
98
            LOGGER.error("Could not archive the recording for call {} — will retry", leadCallId, e);
99
            failureRecorder.recordFailure(leadCallId);
100
        } finally {
101
            closeQuietly(stream);
102
        }
103
    }
104
 
105
    /**
106
     * {@code yyyy/MM/dd/lead-<leadId>/call-<callId>-<recordingUuid>.mp3}
107
     *
108
     * <p>Date-partitioned so a lifecycle rule can expire a day's prefix without scanning, and the ids
109
     * are in the key so an object found on its own is still traceable back to a lead.
110
     *
111
     * <p>No {@code lms-recordings/} prefix: the bucket is already named that, and prefixing would
112
     * give every object a path of {@code lms-recordings/lms-recordings/...}.
113
     */
114
    private String objectKey(LeadCall call) {
115
        LocalDateTime when = call.getStartedAt() != null ? call.getStartedAt() : call.getCreatedTimestamp();
116
        if (when == null) {
117
            when = LocalDateTime.now();
118
        }
119
        return KEY_DATE.format(when)
120
                + "/lead-" + call.getLeadId()
121
                + "/call-" + call.getId() + "-" + call.getRecordingUuid() + ".mp3";
122
    }
123
 
124
    /** Buffered because S3 needs an accurate content length up front, but bounded so it cannot OOM. */
125
    private byte[] readCapped(InputStream in) throws Exception {
126
        ByteArrayOutputStream out = new ByteArrayOutputStream();
127
        byte[] buffer = new byte[8192];
128
        long total = 0;
129
        int read;
130
        while ((read = in.read(buffer)) != -1) {
131
            total += read;
132
            if (total > MAX_RECORDING_BYTES) {
133
                throw new IllegalStateException("Recording exceeds " + MAX_RECORDING_BYTES + " bytes");
134
            }
135
            out.write(buffer, 0, read);
136
        }
137
        return out.toByteArray();
138
    }
139
 
140
    private void closeQuietly(InputStream stream) {
141
        if (stream == null) {
142
            return;
143
        }
144
        try {
145
            stream.close();
146
        } catch (Exception e) {
147
            LOGGER.debug("Ignoring a failure closing the recording stream", e);
148
        }
149
    }
150
}