1
2
3
4
5
6
7
8
9
10
11
12
13
14
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;
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 }