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.helper;
17  
18  import java.util.HashMap;
19  import java.util.Map;
20  
21  import org.codelibs.core.lang.StringUtil;
22  import org.codelibs.core.timer.TimeoutManager;
23  import org.codelibs.core.timer.TimeoutTarget;
24  import org.codelibs.core.timer.TimeoutTask;
25  import org.codelibs.fess.Constants;
26  import org.codelibs.fess.es.config.exbhv.JobLogBhv;
27  import org.codelibs.fess.es.config.exbhv.ScheduledJobBhv;
28  import org.codelibs.fess.es.config.exentity.JobLog;
29  import org.codelibs.fess.es.config.exentity.ScheduledJob;
30  import org.codelibs.fess.exception.JobNotFoundException;
31  import org.codelibs.fess.job.ScheduledJobException;
32  import org.codelibs.fess.mylasta.direction.FessConfig;
33  import org.codelibs.fess.util.ComponentUtil;
34  import org.dbflute.optional.OptionalThing;
35  import org.lastaflute.job.JobManager;
36  import org.lastaflute.job.LaCron;
37  import org.lastaflute.job.LaScheduledJob;
38  import org.lastaflute.job.key.LaJobUnique;
39  import org.lastaflute.job.subsidiary.CronParamsSupplier;
40  import org.slf4j.Logger;
41  import org.slf4j.LoggerFactory;
42  
43  public class JobHelper {
44      private static final Logger logger = LoggerFactory.getLogger(JobHelper.class);
45      private int monitorInterval = 60 * 60;// 1hour
46  
47      public void register(final ScheduledJob scheduledJob) {
48          final JobManager jobManager = ComponentUtil.getJobManager();
49          jobManager.schedule(cron -> register(cron, scheduledJob));
50      }
51  
52      public void register(final LaCron cron, final ScheduledJob scheduledJob) {
53          if (scheduledJob == null) {
54              throw new ScheduledJobException("No job.");
55          }
56  
57          final String id = scheduledJob.getId();
58          if (!Constants.T.equals(scheduledJob.getAvailable())) {
59              logger.info("Inactive Job " + id + ":" + scheduledJob.getName());
60              try {
61                  unregister(scheduledJob);
62              } catch (final Exception e) {
63                  if (logger.isDebugEnabled()) {
64                      logger.debug("Failed to delete Job " + scheduledJob, e);
65                  }
66              }
67              return;
68          }
69  
70          final FessConfig fessConfig = ComponentUtil.getFessConfig();
71          final CronParamsSupplier paramsOp =
72                  () -> {
73                      final Map<String, Object> params = new HashMap<>();
74                      params.put(Constants.SCHEDULED_JOB, ComponentUtil.getComponent(ScheduledJobBhv.class).selectByPK(scheduledJob.getId())
75                              .orElseThrow(() -> new JobNotFoundException(scheduledJob)));
76                      return params;
77                  };
78          findJobByUniqueOf(LaJobUnique.of(id)).ifPresent(job -> {
79              if (!job.isUnscheduled()) {
80                  if (StringUtil.isNotBlank(scheduledJob.getCronExpression())) {
81                      logger.info("Starting Job " + id + ":" + scheduledJob.getName());
82                      final String cronExpression = scheduledJob.getCronExpression();
83                      job.reschedule(cronExpression, op -> op.changeNoticeLogToDebug().params(paramsOp));
84                  } else {
85                      logger.info("Inactive Job " + id + ":" + scheduledJob.getName());
86                      job.becomeNonCron();
87                  }
88              } else if (StringUtil.isNotBlank(scheduledJob.getCronExpression())) {
89                  logger.info("Starting Job " + id + ":" + scheduledJob.getName());
90                  final String cronExpression = scheduledJob.getCronExpression();
91                  job.reschedule(cronExpression, op -> op.changeNoticeLogToDebug().params(paramsOp));
92              }
93          }).orElse(
94                  () -> {
95                      if (StringUtil.isNotBlank(scheduledJob.getCronExpression())) {
96                          logger.info("Starting Job " + id + ":" + scheduledJob.getName());
97                          final String cronExpression = scheduledJob.getCronExpression();
98                          cron.register(cronExpression, fessConfig.getSchedulerJobClassAsClass(),
99                                  fessConfig.getSchedulerConcurrentExecModeAsEnum(),
100                                 op -> op.uniqueBy(id).changeNoticeLogToDebug().params(paramsOp));
101                     } else {
102                         logger.info("Inactive Job " + id + ":" + scheduledJob.getName());
103                         cron.registerNonCron(fessConfig.getSchedulerJobClassAsClass(), fessConfig.getSchedulerConcurrentExecModeAsEnum(),
104                                 op -> op.uniqueBy(id).changeNoticeLogToDebug().params(paramsOp));
105                     }
106                 });
107     }
108 
109     private OptionalThing<LaScheduledJob> findJobByUniqueOf(final LaJobUnique jobUnique) {
110         final JobManager jobManager = ComponentUtil.getJobManager();
111         try {
112             return jobManager.findJobByUniqueOf(jobUnique);
113         } catch (final Exception e) {
114             return OptionalThing.empty();
115         }
116     }
117 
118     public void unregister(final ScheduledJob scheduledJob) {
119         try {
120             final JobManager jobManager = ComponentUtil.getJobManager();
121             if (jobManager.isSchedulingDone()) {
122                 jobManager.findJobByUniqueOf(LaJobUnique.of(scheduledJob.getId())).ifPresent(job -> {
123                     job.unschedule();
124                 }).orElse(() -> logger.debug("Job {} is not scheduled.", scheduledJob.getId()));
125             }
126         } catch (final Exception e) {
127             throw new ScheduledJobException("Failed to delete Job: " + scheduledJob, e);
128         }
129     }
130 
131     public void remove(final ScheduledJob scheduledJob) {
132         try {
133             final JobManager jobManager = ComponentUtil.getJobManager();
134             if (jobManager.isSchedulingDone()) {
135                 jobManager.findJobByUniqueOf(LaJobUnique.of(scheduledJob.getId())).ifPresent(job -> {
136                     job.disappear();
137                 }).orElse(() -> logger.debug("Job {} is not scheduled.", scheduledJob.getId()));
138             }
139         } catch (final Exception e) {
140             throw new ScheduledJobException("Failed to delete Job: " + scheduledJob, e);
141         }
142     }
143 
144     public boolean isAvailable(final String id) {
145         return ComponentUtil.getComponent(ScheduledJobBhv.class).selectByPK(id).filter(e -> Boolean.TRUE.equals(e.getAvailable()))
146                 .isPresent();
147     }
148 
149     public void store(final JobLog jobLog) {
150         ComponentUtil.getComponent(JobLogBhv.class).insertOrUpdate(jobLog, op -> {
151             op.setRefreshPolicy(Constants.TRUE);
152         });
153     }
154 
155     public TimeoutTask startMonitorTask(final JobLog jobLog) {
156         final TimeoutTarget target = new MonitorTarget(jobLog);
157         return TimeoutManager.getInstance().addTimeoutTarget(target, monitorInterval, true);
158     }
159 
160     public void setMonitorInterval(final int monitorInterval) {
161         this.monitorInterval = monitorInterval;
162     }
163 
164     static class MonitorTarget implements TimeoutTarget {
165 
166         private final JobLog jobLog;
167 
168         public MonitorTarget(final JobLog jobLog) {
169             this.jobLog = jobLog;
170         }
171 
172         @Override
173         public void expired() {
174             if (jobLog.getEndTime() == null) {
175                 jobLog.setLastUpdated(ComponentUtil.getSystemHelper().getCurrentTimeAsLong());
176                 if (logger.isDebugEnabled()) {
177                     logger.debug("Update " + jobLog);
178                 }
179                 ComponentUtil.getComponent(JobLogBhv.class).insertOrUpdate(jobLog, op -> {
180                     op.setRefreshPolicy(Constants.TRUE);
181                 });
182             }
183         }
184 
185     }
186 
187 }