Skip to content

Commit d8dc608

Browse files
committed
1_1 MailService
1 parent 66dbb57 commit d8dc608

1 file changed

Lines changed: 69 additions & 3 deletions

File tree

src/main/java/ru/javaops/masterjava/service/MailService.java

Lines changed: 69 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,10 @@
11
package ru.javaops.masterjava.service;
22

3-
import java.util.Collections;
3+
import java.util.ArrayList;
44
import java.util.List;
55
import java.util.Set;
6+
import java.util.concurrent.*;
7+
import java.util.stream.Collectors;
68

79
public class MailService {
810
private static final String OK = "OK";
@@ -11,10 +13,74 @@ public class MailService {
1113
private static final String INTERRUPTED_BY_TIMEOUT = "+++ Interrupted by timeout";
1214
private static final String INTERRUPTED_EXCEPTION = "+++ InterruptedException";
1315

16+
private final ExecutorService mailExecutor = Executors.newFixedThreadPool(8);
17+
1418
public GroupResult sendToList(final String template, final Set<String> emails) throws Exception {
15-
return new GroupResult(0, Collections.emptyList(), null);
16-
}
19+
final CompletionService<MailResult> completionService = new ExecutorCompletionService<>(mailExecutor);
20+
21+
List<Future<MailResult>> futures = emails.stream()
22+
.map(email -> completionService.submit(() -> sendToUser(template, email)))
23+
.collect(Collectors.toList());
1724

25+
return new Callable<GroupResult>() {
26+
private int success = 0;
27+
private List<MailResult> failed = new ArrayList<>();
28+
29+
@Override
30+
public GroupResult call() {
31+
while (!futures.isEmpty()) {
32+
try {
33+
Future<MailResult> future = completionService.poll(10, TimeUnit.SECONDS);
34+
if (future == null) {
35+
return cancelWithFail(INTERRUPTED_BY_TIMEOUT);
36+
}
37+
futures.remove(future);
38+
MailResult mailResult = future.get();
39+
if (mailResult.isOk()) {
40+
success++;
41+
} else {
42+
failed.add(mailResult);
43+
if (failed.size() >= 5) {
44+
return cancelWithFail(INTERRUPTED_BY_FAULTS_NUMBER);
45+
}
46+
}
47+
} catch (ExecutionException e) {
48+
return cancelWithFail(e.getCause().toString());
49+
} catch (InterruptedException e) {
50+
return cancelWithFail(INTERRUPTED_EXCEPTION);
51+
}
52+
}
53+
/*
54+
for (Future<MailResult> future : futures) {
55+
MailResult mailResult;
56+
try {
57+
mailResult = future.get(10, TimeUnit.SECONDS);
58+
} catch (InterruptedException e) {
59+
return cancelWithFail(INTERRUPTED_EXCEPTION);
60+
} catch (ExecutionException e) {
61+
return cancelWithFail(e.getCause().toString());
62+
} catch (TimeoutException e) {
63+
return cancelWithFail(INTERRUPTED_BY_TIMEOUT);
64+
}
65+
if (mailResult.isOk()) {
66+
success++;
67+
} else {
68+
failed.add(mailResult);
69+
if (failed.size() >= 5) {
70+
return cancelWithFail(INTERRUPTED_BY_FAULTS_NUMBER);
71+
}
72+
}
73+
}
74+
*/
75+
return new GroupResult(success, failed, null);
76+
}
77+
78+
private GroupResult cancelWithFail(String cause) {
79+
futures.forEach(f -> f.cancel(true));
80+
return new GroupResult(success, failed, cause);
81+
}
82+
}.call();
83+
}
1884

1985
// dummy realization
2086
public MailResult sendToUser(String template, String email) throws Exception {

0 commit comments

Comments
 (0)