View Javadoc
1   /*
2    * Copyright 2012-2017 CodeLibs Project and the Others.
3    *
4    * Licensed under the Apache License, Version 2.0 (the "License");
5    * you may not use this file except in compliance with the License.
6    * You may obtain a copy of the License at
7    *
8    *     http://www.apache.org/licenses/LICENSE-2.0
9    *
10   * Unless required by applicable law or agreed to in writing, software
11   * distributed under the License is distributed on an "AS IS" BASIS,
12   * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND,
13   * either express or implied. See the License for the specific language
14   * governing permissions and limitations under the License.
15   */
16  package org.codelibs.fess.ds.impl;
17  
18  import static org.codelibs.core.stream.StreamUtil.stream;
19  
20  import java.util.ArrayList;
21  import java.util.List;
22  import java.util.Map;
23  import java.util.concurrent.ExecutorService;
24  import java.util.concurrent.LinkedBlockingQueue;
25  import java.util.concurrent.ThreadPoolExecutor;
26  import java.util.concurrent.TimeUnit;
27  import java.util.stream.Collectors;
28  
29  import org.codelibs.core.io.SerializeUtil;
30  import org.codelibs.fess.Constants;
31  import org.codelibs.fess.crawler.builder.RequestDataBuilder;
32  import org.codelibs.fess.crawler.client.CrawlerClient;
33  import org.codelibs.fess.crawler.client.CrawlerClientFactory;
34  import org.codelibs.fess.crawler.entity.ResponseData;
35  import org.codelibs.fess.crawler.entity.ResultData;
36  import org.codelibs.fess.crawler.exception.ChildUrlsException;
37  import org.codelibs.fess.crawler.exception.CrawlerSystemException;
38  import org.codelibs.fess.crawler.processor.ResponseProcessor;
39  import org.codelibs.fess.crawler.processor.impl.DefaultResponseProcessor;
40  import org.codelibs.fess.crawler.rule.Rule;
41  import org.codelibs.fess.crawler.rule.RuleManager;
42  import org.codelibs.fess.crawler.transformer.Transformer;
43  import org.codelibs.fess.ds.IndexUpdateCallback;
44  import org.codelibs.fess.es.client.FessEsClient;
45  import org.codelibs.fess.exception.DataStoreCrawlingException;
46  import org.codelibs.fess.helper.IndexingHelper;
47  import org.codelibs.fess.mylasta.direction.FessConfig;
48  import org.codelibs.fess.util.ComponentUtil;
49  import org.lastaflute.di.core.SingletonLaContainer;
50  import org.slf4j.Logger;
51  import org.slf4j.LoggerFactory;
52  
53  public class FileListIndexUpdateCallbackImpl implements IndexUpdateCallback {
54      private static final Logger logger = LoggerFactory.getLogger(FileListIndexUpdateCallbackImpl.class);
55  
56      protected IndexUpdateCallback indexUpdateCallback;
57  
58      protected CrawlerClientFactory crawlerClientFactory;
59  
60      protected List<String> deleteUrlList = new ArrayList<>(100);
61  
62      protected int maxDeleteDocumentCacheSize = 100;
63  
64      protected int maxRedirectCount = 10;
65  
66      private final ExecutorService executor;
67  
68      private int executorTerminationTimeout = 300;
69  
70      protected FileListIndexUpdateCallbackImpl(final IndexUpdateCallback indexUpdateCallback,
71              final CrawlerClientFactory crawlerClientFactory, final int nThreads) {
72          this.indexUpdateCallback = indexUpdateCallback;
73          this.crawlerClientFactory = crawlerClientFactory;
74          executor = newFixedThreadPool(nThreads < 1 ? 1 : nThreads);
75      }
76  
77      protected ExecutorService newFixedThreadPool(final int nThreads) {
78          if (logger.isDebugEnabled()) {
79              logger.debug("Executor Thread Pool: " + nThreads);
80          }
81          return new ThreadPoolExecutor(nThreads, nThreads, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>(nThreads),
82                  new ThreadPoolExecutor.CallerRunsPolicy());
83      }
84  
85      @Override
86      public void store(final Map<String, String> paramMap, final Map<String, Object> dataMap) {
87          executor.execute(() -> {
88              final Object eventType = dataMap.remove(getParamValue(paramMap, "field.event_type", "event_type"));
89              if (getParamValue(paramMap, "event.create", "create").equals(eventType)
90                      || getParamValue(paramMap, "event.modify", "modify").equals(eventType)) {
91                  // updated file
92                  addDocument(paramMap, dataMap);
93              } else if (getParamValue(paramMap, "event.delete", "delete").equals(eventType)) {
94                  // deleted file
95                  deleteDocument(paramMap, dataMap);
96              } else {
97                  logger.warn("unknown event: " + eventType + ", data: " + dataMap);
98              }
99          });
100     }
101 
102     protected String getParamValue(final Map<String, String> paramMap, final String key, final String defaultValue) {
103         return paramMap.getOrDefault(key, defaultValue);
104     }
105 
106     protected void addDocument(final Map<String, String> paramMap, final Map<String, Object> dataMap) {
107         final FessConfig fessConfig = ComponentUtil.getFessConfig();
108         synchronized (indexUpdateCallback) {
109             // required check
110             if (!dataMap.containsKey(fessConfig.getIndexFieldUrl()) || dataMap.get(fessConfig.getIndexFieldUrl()) == null) {
111                 logger.warn("Could not add a doc. Invalid data: " + dataMap);
112                 return;
113             }
114 
115             final String url = dataMap.get(fessConfig.getIndexFieldUrl()).toString();
116             final CrawlerClient client = crawlerClientFactory.getClient(url);
117             if (client == null) {
118                 logger.warn("CrawlerClient is null. Data: " + dataMap);
119                 return;
120             }
121 
122             String processingUrl = url;
123             for (int i = 0; i < maxRedirectCount; i++) {
124                 processingUrl = processRequest(paramMap, dataMap, processingUrl, client);
125                 if (processingUrl == null) {
126                     break;
127                 }
128                 dataMap.put(fessConfig.getIndexFieldUrl(), processingUrl);
129             }
130         }
131     }
132 
133     protected String processRequest(final Map<String, String> paramMap, final Map<String, Object> dataMap, final String url,
134             final CrawlerClient client) {
135         final long startTime = System.currentTimeMillis();
136         try (final ResponseData responseData = client.execute(RequestDataBuilder.newRequestData().get().url(url).build())) {
137             if (responseData.getRedirectLocation() != null) {
138                 return responseData.getRedirectLocation();
139             }
140             responseData.setExecutionTime(System.currentTimeMillis() - startTime);
141             if (dataMap.containsKey(Constants.SESSION_ID)) {
142                 responseData.setSessionId((String) dataMap.get(Constants.SESSION_ID));
143             } else {
144                 responseData.setSessionId(paramMap.get(Constants.CRAWLING_INFO_ID));
145             }
146 
147             final RuleManager ruleManager = SingletonLaContainer.getComponent(RuleManager.class);
148             final Rule rule = ruleManager.getRule(responseData);
149             if (rule == null) {
150                 logger.warn("No url rule. Data: " + dataMap);
151             } else {
152                 responseData.setRuleId(rule.getRuleId());
153                 final ResponseProcessor responseProcessor = rule.getResponseProcessor();
154                 if (responseProcessor instanceof DefaultResponseProcessor) {
155                     final Transformer transformer = ((DefaultResponseProcessor) responseProcessor).getTransformer();
156                     final ResultData resultData = transformer.transform(responseData);
157                     final byte[] data = resultData.getData();
158                     if (data != null) {
159                         try {
160                             @SuppressWarnings("unchecked")
161                             final Map<String, Object> responseDataMap = (Map<String, Object>) SerializeUtil.fromBinaryToObject(data);
162                             dataMap.putAll(responseDataMap);
163                         } catch (final Exception e) {
164                             throw new CrawlerSystemException("Could not create an instance from bytes.", e);
165                         }
166                     }
167 
168                     // remove
169                     String[] ignoreFields;
170                     if (paramMap.containsKey("ignore.field.names")) {
171                         ignoreFields = paramMap.get("ignore.field.names").split(",");
172                     } else {
173                         ignoreFields = new String[] { Constants.INDEXING_TARGET, Constants.SESSION_ID };
174                     }
175                     stream(ignoreFields).of(stream -> stream.map(s -> s.trim()).forEach(s -> dataMap.remove(s)));
176 
177                     indexUpdateCallback.store(paramMap, dataMap);
178                 } else {
179                     logger.warn("The response processor is not DefaultResponseProcessor. responseProcessor: " + responseProcessor
180                             + ", Data: " + dataMap);
181                 }
182             }
183             return null;
184         } catch (final ChildUrlsException e) {
185             throw new DataStoreCrawlingException(url, "Redirected to "
186                     + e.getChildUrlList().stream().map(r -> r.getUrl()).collect(Collectors.joining(", ")), e);
187         } catch (final Exception e) {
188             throw new DataStoreCrawlingException(url, "Failed to add: " + dataMap, e);
189         }
190     }
191 
192     protected boolean deleteDocument(final Map<String, String> paramMap, final Map<String, Object> dataMap) {
193 
194         if (logger.isDebugEnabled()) {
195             logger.debug("Deleting " + dataMap);
196         }
197 
198         final FessConfig fessConfig = ComponentUtil.getFessConfig();
199 
200         // required check
201         if (!dataMap.containsKey(fessConfig.getIndexFieldUrl()) || dataMap.get(fessConfig.getIndexFieldUrl()) == null) {
202             logger.warn("Could not delete a doc. Invalid data: " + dataMap);
203             return false;
204         }
205 
206         synchronized (indexUpdateCallback) {
207             deleteUrlList.add(dataMap.get(fessConfig.getIndexFieldUrl()).toString());
208 
209             if (deleteUrlList.size() >= maxDeleteDocumentCacheSize) {
210                 deleteDocuments();
211             }
212 
213         }
214         return true;
215     }
216 
217     @Override
218     public void commit() {
219         try {
220             if (logger.isDebugEnabled()) {
221                 logger.debug("Shutting down thread executor.");
222             }
223             executor.shutdown();
224             executor.awaitTermination(executorTerminationTimeout, TimeUnit.SECONDS);
225         } catch (final InterruptedException e) {
226             if (logger.isDebugEnabled()) {
227                 logger.debug("Failed to interrupt executor.", e);
228             }
229         } finally {
230             executor.shutdownNow();
231         }
232 
233         if (!deleteUrlList.isEmpty()) {
234             deleteDocuments();
235         }
236         indexUpdateCallback.commit();
237     }
238 
239     protected void deleteDocuments() {
240         final FessEsClient fessEsClient = ComponentUtil.getFessEsClient();
241         final IndexingHelper indexingHelper = ComponentUtil.getIndexingHelper();
242         for (final String url : deleteUrlList) {
243             indexingHelper.deleteDocumentByUrl(fessEsClient, url);
244         }
245         if (logger.isDebugEnabled()) {
246             logger.debug("Deleted " + deleteUrlList);
247         }
248         deleteUrlList.clear();
249     }
250 
251     @Override
252     public long getDocumentSize() {
253         return indexUpdateCallback.getDocumentSize();
254     }
255 
256     @Override
257     public long getExecuteTime() {
258         return indexUpdateCallback.getExecuteTime();
259     }
260 
261     public void setMaxDeleteDocumentCacheSize(final int maxDeleteDocumentCacheSize) {
262         this.maxDeleteDocumentCacheSize = maxDeleteDocumentCacheSize;
263     }
264 
265     public void setMaxRedirectCount(final int maxRedirectCount) {
266         this.maxRedirectCount = maxRedirectCount;
267     }
268 
269     public void setExecutorTerminationTimeout(final int executorTerminationTimeout) {
270         this.executorTerminationTimeout = executorTerminationTimeout;
271     }
272 
273 }