1
2
3
4
5
6
7
8
9
10
11
12
13
14
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
92 addDocument(paramMap, dataMap);
93 } else if (getParamValue(paramMap, "event.delete", "delete").equals(eventType)) {
94
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
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
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
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 }