View Javadoc
1   /*
2    * Copyright 2012-2021 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.callback;
17  
18  import static org.codelibs.core.stream.StreamUtil.stream;
19  
20  import java.util.ArrayList;
21  import java.util.Deque;
22  import java.util.LinkedList;
23  import java.util.List;
24  import java.util.Map;
25  import java.util.concurrent.ExecutorService;
26  import java.util.concurrent.LinkedBlockingQueue;
27  import java.util.concurrent.ThreadPoolExecutor;
28  import java.util.concurrent.TimeUnit;
29  import java.util.stream.Collectors;
30  
31  import org.apache.logging.log4j.LogManager;
32  import org.apache.logging.log4j.Logger;
33  import org.codelibs.core.io.SerializeUtil;
34  import org.codelibs.fesen.index.query.QueryBuilders;
35  import org.codelibs.fess.Constants;
36  import org.codelibs.fess.crawler.builder.RequestDataBuilder;
37  import org.codelibs.fess.crawler.client.CrawlerClient;
38  import org.codelibs.fess.crawler.client.CrawlerClientFactory;
39  import org.codelibs.fess.crawler.entity.RequestData;
40  import org.codelibs.fess.crawler.entity.ResponseData;
41  import org.codelibs.fess.crawler.entity.ResultData;
42  import org.codelibs.fess.crawler.exception.ChildUrlsException;
43  import org.codelibs.fess.crawler.exception.CrawlerSystemException;
44  import org.codelibs.fess.crawler.processor.ResponseProcessor;
45  import org.codelibs.fess.crawler.processor.impl.DefaultResponseProcessor;
46  import org.codelibs.fess.crawler.rule.Rule;
47  import org.codelibs.fess.crawler.rule.RuleManager;
48  import org.codelibs.fess.crawler.transformer.Transformer;
49  import org.codelibs.fess.es.client.SearchEngineClient;
50  import org.codelibs.fess.exception.DataStoreCrawlingException;
51  import org.codelibs.fess.helper.IndexingHelper;
52  import org.codelibs.fess.mylasta.direction.FessConfig;
53  import org.codelibs.fess.util.ComponentUtil;
54  import org.lastaflute.di.core.SingletonLaContainer;
55  
56  public class FileListIndexUpdateCallbackImpl implements IndexUpdateCallback {
57      private static final Logger logger = LogManager.getLogger(FileListIndexUpdateCallbackImpl.class);
58  
59      protected IndexUpdateCallback indexUpdateCallback;
60  
61      protected CrawlerClientFactory crawlerClientFactory;
62  
63      protected List<String> deleteUrlList = new ArrayList<>(100);
64  
65      protected int maxDeleteDocumentCacheSize;
66  
67      protected int maxRedirectCount;
68  
69      private final ExecutorService executor;
70  
71      private int executorTerminationTimeout = 300;
72  
73      public FileListIndexUpdateCallbackImpl(final IndexUpdateCallback indexUpdateCallback, final CrawlerClientFactory crawlerClientFactory,
74              final int nThreads) {
75          this.indexUpdateCallback = indexUpdateCallback;
76          this.crawlerClientFactory = crawlerClientFactory;
77          executor = newFixedThreadPool(nThreads < 1 ? 1 : nThreads);
78          final FessConfig fessConfig = ComponentUtil.getFessConfig();
79          maxDeleteDocumentCacheSize = fessConfig.getIndexerDataMaxDeleteCacheSizeAsInteger();
80          maxRedirectCount = fessConfig.getIndexerDataMaxRedirectCountAsInteger();
81      }
82  
83      protected ExecutorService newFixedThreadPool(final int nThreads) {
84          if (logger.isDebugEnabled()) {
85              logger.debug("Executor Thread Pool: {}", nThreads);
86          }
87          return new ThreadPoolExecutor(nThreads, nThreads, 0L, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<Runnable>(nThreads),
88                  new ThreadPoolExecutor.CallerRunsPolicy());
89      }
90  
91      @Override
92      public void store(final Map<String, String> paramMap, final Map<String, Object> dataMap) {
93          executor.execute(() -> {
94              final Object eventType = dataMap.remove(getParamValue(paramMap, "field.event_type", "event_type"));
95              if (getParamValue(paramMap, "event.create", "create").equals(eventType)
96                      || getParamValue(paramMap, "event.modify", "modify").equals(eventType)) {
97                  // updated file
98                  addDocument(paramMap, dataMap);
99              } else if (getParamValue(paramMap, "event.delete", "delete").equals(eventType)) {
100                 // deleted file
101                 deleteDocument(paramMap, dataMap);
102             } else {
103                 logger.warn("unknown event: {}, data: {}", eventType, dataMap);
104             }
105         });
106     }
107 
108     protected String getParamValue(final Map<String, String> paramMap, final String key, final String defaultValue) {
109         return paramMap.getOrDefault(key, defaultValue);
110     }
111 
112     protected void addDocument(final Map<String, String> paramMap, final Map<String, Object> dataMap) {
113         final FessConfig fessConfig = ComponentUtil.getFessConfig();
114         synchronized (indexUpdateCallback) {
115             // required check
116             if (!dataMap.containsKey(fessConfig.getIndexFieldUrl()) || dataMap.get(fessConfig.getIndexFieldUrl()) == null) {
117                 logger.warn("Could not add a doc. Invalid data: {}", dataMap);
118                 return;
119             }
120 
121             final String url = dataMap.get(fessConfig.getIndexFieldUrl()).toString();
122             final CrawlerClient client = crawlerClientFactory.getClient(url);
123             if (client == null) {
124                 logger.warn("CrawlerClient is null. Data: {}", dataMap);
125                 return;
126             }
127 
128             final long maxAccessCount = getMaxAccessCount(paramMap, dataMap);
129             long counter = 0;
130             final Deque<String> urlQueue = new LinkedList<>();
131             urlQueue.offer(url);
132             while (!urlQueue.isEmpty() && (maxAccessCount < 0 || counter < maxAccessCount)) {
133                 final Map<String, Object> localDataMap =
134                         dataMap.entrySet().stream().collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue));
135                 String processingUrl = urlQueue.poll();
136                 if (deleteUrlList.contains(processingUrl)) {
137                     deleteDocuments(); // delete before indexing
138                 }
139                 try {
140                     for (int i = 0; i < maxRedirectCount; i++) {
141                         processingUrl = processRequest(paramMap, localDataMap, processingUrl, client);
142                         if (processingUrl == null) {
143                             break;
144                         }
145                         counter++;
146                         localDataMap.put(fessConfig.getIndexFieldUrl(), processingUrl);
147                     }
148                 } catch (final ChildUrlsException e) {
149                     e.getChildUrlList().stream().map(RequestData::getUrl).forEach(urlQueue::offer);
150                 } catch (final DataStoreCrawlingException e) {
151                     final Throwable cause = e.getCause();
152                     if (cause instanceof ChildUrlsException) {
153                         ((ChildUrlsException) cause).getChildUrlList().stream().map(RequestData::getUrl).forEach(urlQueue::offer);
154                     } else if (maxAccessCount != 1L) {
155                         throw e;
156                     } else {
157                         logger.warn("Failed to access {}.", processingUrl, e);
158                     }
159                 }
160             }
161         }
162     }
163 
164     protected long getMaxAccessCount(final Map<String, String> paramMap, final Map<String, Object> dataMap) {
165         final Object recursive = dataMap.remove(getParamValue(paramMap, "field.recursive", "recursive"));
166         if (recursive == null || Constants.FALSE.equalsIgnoreCase(recursive.toString())) {
167             return 1L;
168         }
169         if (Constants.TRUE.equalsIgnoreCase(recursive.toString())) {
170             return -1L;
171         }
172         try {
173             return Long.parseLong(recursive.toString());
174         } catch (final NumberFormatException e) {
175             return 1L;
176         }
177     }
178 
179     protected String processRequest(final Map<String, String> paramMap, final Map<String, Object> dataMap, final String url,
180             final CrawlerClient client) {
181         final long startTime = System.currentTimeMillis();
182         try (final ResponseData responseData = client.execute(RequestDataBuilder.newRequestData().get().url(url).build())) {
183             if (responseData.getRedirectLocation() != null) {
184                 return responseData.getRedirectLocation();
185             }
186             responseData.setExecutionTime(System.currentTimeMillis() - startTime);
187             if (dataMap.containsKey(Constants.SESSION_ID)) {
188                 responseData.setSessionId((String) dataMap.get(Constants.SESSION_ID));
189             } else {
190                 responseData.setSessionId(paramMap.get(Constants.CRAWLING_INFO_ID));
191             }
192 
193             final RuleManager ruleManager = SingletonLaContainer.getComponent(RuleManager.class);
194             final Rule rule = ruleManager.getRule(responseData);
195             if (rule == null) {
196                 logger.warn("No url rule. Data: {}", dataMap);
197             } else {
198                 responseData.setRuleId(rule.getRuleId());
199                 final ResponseProcessor responseProcessor = rule.getResponseProcessor();
200                 if (responseProcessor instanceof DefaultResponseProcessor) {
201                     final Transformer transformer = ((DefaultResponseProcessor) responseProcessor).getTransformer();
202                     final ResultData resultData = transformer.transform(responseData);
203                     final byte[] data = resultData.getData();
204                     if (data != null) {
205                         try {
206                             @SuppressWarnings("unchecked")
207                             final Map<String, Object> responseDataMap = (Map<String, Object>) SerializeUtil.fromBinaryToObject(data);
208                             dataMap.putAll(responseDataMap);
209                         } catch (final Exception e) {
210                             throw new CrawlerSystemException("Could not create an instance from bytes.", e);
211                         }
212                     }
213 
214                     // remove
215                     String[] ignoreFields;
216                     if (paramMap.containsKey("ignore.field.names")) {
217                         ignoreFields = paramMap.get("ignore.field.names").split(",");
218                     } else {
219                         ignoreFields = new String[] { Constants.INDEXING_TARGET, Constants.SESSION_ID };
220                     }
221                     stream(ignoreFields).of(stream -> stream.map(String::trim).forEach(s -> dataMap.remove(s)));
222 
223                     indexUpdateCallback.store(paramMap, dataMap);
224                 } else {
225                     logger.warn("The response processor is not DefaultResponseProcessor. responseProcessor: {}, Data: {}",
226                             responseProcessor, dataMap);
227                 }
228             }
229             return null;
230         } catch (final ChildUrlsException e) {
231             throw new DataStoreCrawlingException(url,
232                     "Redirected to " + e.getChildUrlList().stream().map(RequestData::getUrl).collect(Collectors.joining(", ")), e);
233         } catch (final Exception e) {
234             throw new DataStoreCrawlingException(url, "Failed to add: " + dataMap, e);
235         }
236     }
237 
238     protected boolean deleteDocument(final Map<String, String> paramMap, final Map<String, Object> dataMap) {
239 
240         if (logger.isDebugEnabled()) {
241             logger.debug("Deleting {}", dataMap);
242         }
243 
244         final FessConfig fessConfig = ComponentUtil.getFessConfig();
245 
246         // required check
247         if (!dataMap.containsKey(fessConfig.getIndexFieldUrl()) || dataMap.get(fessConfig.getIndexFieldUrl()) == null) {
248             logger.warn("Could not delete a doc. Invalid data: {}", dataMap);
249             return false;
250         }
251 
252         synchronized (indexUpdateCallback) {
253             final long maxAccessCount = getMaxAccessCount(paramMap, dataMap);
254             final String url = dataMap.get(fessConfig.getIndexFieldUrl()).toString();
255             if (maxAccessCount != 1L) {
256                 final SearchEngineClient searchEngineClient = ComponentUtil.getSearchEngineClient();
257                 final IndexingHelper indexingHelper = ComponentUtil.getIndexingHelper();
258                 final long count = indexingHelper.deleteDocumentByQuery(searchEngineClient,
259                         QueryBuilders.prefixQuery(fessConfig.getIndexFieldUrl(), url));
260                 if (logger.isDebugEnabled()) {
261                     logger.debug("Deleted {} docs for {}*", count, url);
262                 }
263             } else {
264                 deleteUrlList.add(url);
265 
266                 if (deleteUrlList.size() >= maxDeleteDocumentCacheSize) {
267                     deleteDocuments();
268                 }
269             }
270         }
271         return true;
272     }
273 
274     @Override
275     public void commit() {
276         try {
277             if (logger.isDebugEnabled()) {
278                 logger.debug("Shutting down thread executor.");
279             }
280             executor.shutdown();
281             executor.awaitTermination(executorTerminationTimeout, TimeUnit.SECONDS);
282         } catch (final InterruptedException e) {
283             if (logger.isDebugEnabled()) {
284                 logger.debug("Failed to interrupt executor.", e);
285             }
286         } finally {
287             executor.shutdownNow();
288         }
289 
290         if (!deleteUrlList.isEmpty()) {
291             deleteDocuments();
292         }
293         indexUpdateCallback.commit();
294     }
295 
296     protected void deleteDocuments() {
297         final SearchEngineClient searchEngineClient = ComponentUtil.getSearchEngineClient();
298         final IndexingHelper indexingHelper = ComponentUtil.getIndexingHelper();
299         for (final String url : deleteUrlList) {
300             indexingHelper.deleteDocumentByUrl(searchEngineClient, url);
301         }
302         if (logger.isDebugEnabled()) {
303             logger.debug("Deleted {}", deleteUrlList);
304         }
305         deleteUrlList.clear();
306     }
307 
308     @Override
309     public long getDocumentSize() {
310         return indexUpdateCallback.getDocumentSize();
311     }
312 
313     @Override
314     public long getExecuteTime() {
315         return indexUpdateCallback.getExecuteTime();
316     }
317 
318     public void setMaxDeleteDocumentCacheSize(final int maxDeleteDocumentCacheSize) {
319         this.maxDeleteDocumentCacheSize = maxDeleteDocumentCacheSize;
320     }
321 
322     public void setMaxRedirectCount(final int maxRedirectCount) {
323         this.maxRedirectCount = maxRedirectCount;
324     }
325 
326     public void setExecutorTerminationTimeout(final int executorTerminationTimeout) {
327         this.executorTerminationTimeout = executorTerminationTimeout;
328     }
329 
330 }