| 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 |
}
|