11package ru .javaops .masterjava .service ;
22
3- import java .util .Collections ;
3+ import java .util .ArrayList ;
44import java .util .List ;
55import java .util .Set ;
6+ import java .util .concurrent .*;
7+ import java .util .stream .Collectors ;
68
79public 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