001/**
002 * Licensed to the Apache Software Foundation (ASF) under one
003 * or more contributor license agreements.  See the NOTICE file
004 * distributed with this work for additional information
005 * regarding copyright ownership.  The ASF licenses this file
006 * to you under the Apache License, Version 2.0 (the
007 * "License"); you may not use this file except in compliance
008 * with the License.  You may obtain a copy of the License at
009 *
010 *      http://www.apache.org/licenses/LICENSE-2.0
011 *
012 * Unless required by applicable law or agreed to in writing, software
013 * distributed under the License is distributed on an "AS IS" BASIS,
014 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
015 * See the License for the specific language governing permissions and
016 * limitations under the License.
017 */
018
019package org.apache.oozie.executor.jpa;
020
021import java.sql.Timestamp;
022import java.text.ParseException;
023import java.util.ArrayList;
024import java.util.Collections;
025import java.util.HashMap;
026import java.util.Iterator;
027import java.util.List;
028import java.util.Map;
029import java.util.Map.Entry;
030import java.util.regex.Matcher;
031import java.util.regex.Pattern;
032
033import javax.persistence.EntityManager;
034import javax.persistence.Query;
035
036import org.apache.oozie.BulkResponseInfo;
037import org.apache.oozie.BundleJobBean;
038import org.apache.oozie.CoordinatorActionBean;
039import org.apache.oozie.CoordinatorJobBean;
040import org.apache.oozie.ErrorCode;
041import org.apache.oozie.StringBlob;
042import org.apache.oozie.client.BundleJob;
043import org.apache.oozie.client.CoordinatorAction;
044import org.apache.oozie.client.CoordinatorJob;
045import org.apache.oozie.client.rest.BulkResponseImpl;
046import org.apache.oozie.service.Services;
047import org.apache.oozie.util.DateUtils;
048import org.apache.oozie.util.ParamChecker;
049
050/**
051 * The query executor class for bulk monitoring queries i.e. debugging bundle ->
052 * coord actions directly
053 */
054public class BulkJPAExecutor implements JPAExecutor<BulkResponseInfo> {
055    private Map<String, List<String>> bulkFilter;
056    // defaults
057    private int start = 1;
058    private int len = 50;
059    private enum PARAM_TYPE {
060        ID, NAME
061    }
062
063    public BulkJPAExecutor(Map<String, List<String>> bulkFilter, int start, int len) {
064        ParamChecker.notNull(bulkFilter, "bulkFilter");
065        this.bulkFilter = bulkFilter;
066        this.start = start;
067        this.len = len;
068    }
069
070    @Override
071    public String getName() {
072        return "BulkJPAExecutor";
073    }
074
075    @Override
076    public BulkResponseInfo execute(EntityManager em) throws JPAExecutorException {
077        List<BulkResponseImpl> responseList = new ArrayList<BulkResponseImpl>();
078        Map<String, Timestamp> actionTimes = new HashMap<String, Timestamp>();
079
080        try {
081            List<String> coords = bulkFilter.get(BulkResponseImpl.BULK_FILTER_COORD);
082            List<String> statuses = bulkFilter.get(BulkResponseImpl.BULK_FILTER_STATUS);
083            List<String> params = new ArrayList<String>();
084
085            // Lightweight Query 1 on Bundle level to fetch the bundle job(s)
086            // corresponding to names or ids
087            List<BundleJobBean> bundleBeans = bundleQuery(em);
088
089            // Join query between coordinator job and coordinator action tables
090            // to get entries for specific bundleId only
091            String conditions = actionQuery(coords, params, statuses, em, bundleBeans, actionTimes, responseList);
092
093            // Query to get the count of records
094            long total = countQuery(statuses, params, conditions, em, bundleBeans, actionTimes);
095
096            BulkResponseInfo bulk = new BulkResponseInfo(responseList, start, len, total);
097            return bulk;
098        }
099        catch (Exception e) {
100            throw new JPAExecutorException(ErrorCode.E0603, e.getMessage(), e);
101        }
102    }
103
104    /**
105     * build the bundle level query to get bundle beans for the specified ids or appnames
106     * @param em
107     * @return List BundleJobBeans
108     * @throws JPAExecutorException
109     */
110    @SuppressWarnings("unchecked")
111    private List<BundleJobBean> bundleQuery(EntityManager em) throws JPAExecutorException {
112        Query q = em.createNamedQuery("BULK_MONITOR_BUNDLE_QUERY");
113        StringBuilder bundleQuery = new StringBuilder(q.toString());
114
115        StringBuilder whereClause = null;
116        List<String> bundles = bulkFilter.get(BulkResponseImpl.BULK_FILTER_BUNDLE);
117        if (bundles != null) {
118            PARAM_TYPE type = getParamType(bundles.get(0), 'B');
119            if (type == PARAM_TYPE.NAME) {
120                whereClause = inClause(bundles.size(), "appName", 'b', "bundles");
121            }
122            else if (type == PARAM_TYPE.ID) {
123                whereClause = inClause(bundles.size(), "id", 'b', "bundles");
124            }
125
126            // Query: select <columns> from BundleJobBean b where b.id IN (...) _or_ b.appName IN (...)
127            bundleQuery.append(whereClause.replace(whereClause.indexOf("AND"), whereClause.indexOf("AND") + 3, "WHERE"));
128            Query tmp = em.createQuery(bundleQuery.toString());
129
130            fillParameters(tmp, "bundles", bundles);
131
132            List<Object[]> bundleObjs = (List<Object[]>) tmp.getResultList();
133            if (bundleObjs.isEmpty()) {
134                throw new JPAExecutorException(ErrorCode.E0603, "No entries found for given bundle(s)");
135            }
136
137            List<BundleJobBean> bundleBeans = new ArrayList<BundleJobBean>();
138            for (Object[] bundleElem : bundleObjs) {
139                bundleBeans.add(constructBundleBean(bundleElem));
140            }
141            return bundleBeans;
142        }
143        return null;
144    }
145
146    /**
147     * Validate and determine whether passed param is job-id or appname
148     * @param id
149     * @param job
150     * @return PARAM_TYPE
151     */
152    private PARAM_TYPE getParamType(String id, char job) {
153        Pattern p = Pattern.compile("\\d{7}-\\d{15}-" + Services.get().getSystemId() + "-" + job);
154        Matcher m = p.matcher(id);
155        if (m.matches()) {
156            return PARAM_TYPE.ID;
157        }
158        return PARAM_TYPE.NAME;
159    }
160
161    /**
162     * Compose the coord action level query comprising bundle id/appname filter and coord action
163     * status filter (if specified) and start-time or nominal-time filter (if specified)
164     * @param em
165     * @param bundles
166     * @param times
167     * @param responseList
168     * @return Query string
169     * @throws ParseException
170     */
171    @SuppressWarnings("unchecked")
172    private String actionQuery(final List<String> coords, final List<String> params, List<String> statuses, EntityManager em,
173                               List<BundleJobBean> bundles, Map<String, Timestamp> times, List<BulkResponseImpl> responseList)
174                               throws ParseException {
175        Query q = em.createNamedQuery("BULK_MONITOR_ACTIONS_QUERY");
176        StringBuilder getActions = new StringBuilder(q.toString());
177        int offset = getActions.indexOf("ORDER");
178        StringBuilder conditionClause = new StringBuilder();
179
180        // Query: Select <columns> from CoordinatorActionBean a, CoordinatorJobBean c WHERE a.jobId = c.id
181        // AND c.bundleId = :bundleId AND c.appName/id IN (...)
182
183        if (coords != null) {
184            PARAM_TYPE type = getParamType(coords.get(0), 'C');
185            if (type == PARAM_TYPE.NAME) {
186                conditionClause.append(inClause(coords.size(), "appName", 'c', "param"));
187                params.addAll(coords);
188            }
189            else if (type == PARAM_TYPE.ID) {
190                conditionClause.append(inClause(coords.size(), "id", 'c', "param"));
191                params.addAll(coords);
192            }
193        }
194        // Query: Select <columns> from CoordinatorActionBean a, CoordinatorJobBean c WHERE a.jobId = c.id
195        // AND c.bundleId = :bundleId AND c.appName/id IN (...) AND a.statusStr IN (...)
196        conditionClause.append(statusClause(statuses));
197
198        offset = getActions.indexOf("ORDER");
199        getActions.insert(offset - 1, conditionClause);
200
201        // Query: Select <columns> from CoordinatorActionBean a, CoordinatorJobBean c WHERE a.jobId = c.id
202        // AND c.bundleId = :bundleId AND c.appName/id IN (...) AND a.statusStr IN (...)
203        // AND a.createdTimestamp >= startCreated _or_ a.createdTimestamp <= endCreated
204        // AND a.nominalTimestamp >= startNominal _or_ a.nominalTimestamp <= endNominal
205        timesClause(getActions, offset, times);
206        q = em.createQuery(getActions.toString());
207        Iterator<Entry<String, Timestamp>> iter = times.entrySet().iterator();
208        while (iter.hasNext()) {
209            Entry<String, Timestamp> time = iter.next();
210            q.setParameter(time.getKey(), time.getValue());
211        }
212        // pagination
213        q.setFirstResult(start - 1);
214        q.setMaxResults(len);
215
216        if (coords != null) {
217            fillParameters(q, "param",  coords);
218        }
219
220        if (statuses != null) {
221            fillParameters(q, "status", statuses);
222        }
223
224        // repeatedly execute above query for each bundle
225        for (BundleJobBean bundle : bundles) {
226            q.setParameter("bundleId", bundle.getId());
227            List<Object[]> response = q.getResultList();
228            for (Object[] r : response) {
229                BulkResponseImpl br = getResponseFromObject(bundle, r);
230                responseList.add(br);
231            }
232        }
233        return q.toString();
234    }
235
236    /**
237     * Get total number of records for use with offset and len in API
238     * @param clause
239     * @param em
240     * @param bundles
241     * @return total count of coord actions
242     */
243    private long countQuery(List<String> statuses, List<String> params, String clause, EntityManager em,
244                            List<BundleJobBean> bundles, Map<String, Timestamp> times) {
245        Query q = em.createNamedQuery("BULK_MONITOR_COUNT_QUERY");
246        StringBuilder getTotal = new StringBuilder(q.toString() + " ");
247        // Query: select COUNT(a) from CoordinatorActionBean a, CoordinatorJobBean c
248        // get entire WHERE clause from above i.e. actionQuery() for all conditions on coordinator job
249        // and action status and times
250        getTotal.append(clause.substring(clause.indexOf("WHERE"), clause.indexOf("ORDER")));
251        int offset = getTotal.indexOf("bundleId");
252        List<String> bundleIds = new ArrayList<String>();
253        for (BundleJobBean bundle : bundles) {
254            bundleIds.add(bundle.getId());
255        }
256        // Query: select COUNT(a) from CoordinatorActionBean a, CoordinatorJobBean c WHERE ...
257        // AND c.bundleId IN (... list of bundle ids) i.e. replace single :bundleId with list
258        getTotal = getTotal.replace(offset - 6, offset + 20, inClause(bundleIds.size(), "bundleId", 'c', "count").toString());
259        q = em.createQuery(getTotal.toString());
260
261        fillParameters(q, "count", bundleIds);
262
263        if (statuses != null) {
264            fillParameters(q, "status", statuses);
265        }
266
267        if (params != null) {
268            fillParameters(q, "param",  params);
269        }
270
271        Iterator<Entry<String, Timestamp>> iter = times.entrySet().iterator();
272        while (iter.hasNext()) {
273            Entry<String, Timestamp> time = iter.next();
274            q.setParameter(time.getKey(), time.getValue());
275        }
276        long total = ((Long) q.getSingleResult()).longValue();
277        return total;
278    }
279
280    // Form the where clause to filter by coordinator appname/id
281    private StringBuilder inClause(int noOfValues, String col, char type, String paramPrefix) {
282        StringBuilder sb = new StringBuilder();
283        boolean firstVal = true;
284
285        for (int i = 0; i < noOfValues; i++) {
286            if (firstVal) {
287                sb.append(" AND " + type + "." + col + " IN (:" + paramPrefix + i);
288                firstVal = false;
289            }
290            else {
291                sb.append(", :" + paramPrefix + i);
292            }
293        }
294        if (!firstVal) {
295            sb.append(") ");
296        }
297
298        return sb;
299    }
300
301    // Form the where clause to filter by coord action status
302    private StringBuilder statusClause(List<String> statuses) {
303        StringBuilder sb = new StringBuilder();
304
305        if (statuses != null) {
306            sb = inClause(statuses.size(), "statusStr", 'a', "status");
307        }
308
309        if (sb.length() == 0) { // statuses was null. adding default
310            sb.append(" AND a.statusStr IN ('KILLED', 'FAILED') ");
311        }
312
313        return sb;
314    }
315
316    private void timesClause(StringBuilder sb, int offset, Map<String, Timestamp> eachTime) throws ParseException {
317        Timestamp ts = null;
318        List<String> times = bulkFilter.get(BulkResponseImpl.BULK_FILTER_START_CREATED_EPOCH);
319        if (times != null) {
320            ts = new Timestamp(DateUtils.parseDateUTC(times.get(0)).getTime());
321            sb.insert(offset - 1, " AND a.createdTimestamp >= :startCreated");
322            eachTime.put("startCreated", ts);
323        }
324        times = bulkFilter.get(BulkResponseImpl.BULK_FILTER_END_CREATED_EPOCH);
325        if (times != null) {
326            ts = new Timestamp(DateUtils.parseDateUTC(times.get(0)).getTime());
327            sb.insert(offset - 1, " AND a.createdTimestamp <= :endCreated");
328            eachTime.put("endCreated", ts);
329        }
330        times = bulkFilter.get(BulkResponseImpl.BULK_FILTER_START_NOMINAL_EPOCH);
331        if (times != null) {
332            ts = new Timestamp(DateUtils.parseDateUTC(times.get(0)).getTime());
333            sb.insert(offset - 1, " AND a.nominalTimestamp >= :startNominal");
334            eachTime.put("startNominal", ts);
335        }
336        times = bulkFilter.get(BulkResponseImpl.BULK_FILTER_END_NOMINAL_EPOCH);
337        if (times != null) {
338            ts = new Timestamp(DateUtils.parseDateUTC(times.get(0)).getTime());
339            sb.insert(offset - 1, " AND a.nominalTimestamp <= :endNominal");
340            eachTime.put("endNominal", ts);
341        }
342    }
343
344    private BulkResponseImpl getResponseFromObject(BundleJobBean bundleBean, Object arr[]) {
345        BulkResponseImpl bean = new BulkResponseImpl();
346        CoordinatorJobBean coordBean = new CoordinatorJobBean();
347        CoordinatorActionBean actionBean = new CoordinatorActionBean();
348        if (arr[0] != null) {
349            actionBean.setId((String) arr[0]);
350        }
351        if (arr[1] != null) {
352            actionBean.setActionNumber((Integer) arr[1]);
353        }
354        if (arr[2] != null) {
355            actionBean.setErrorCode((String) arr[2]);
356        }
357        if (arr[3] != null) {
358            actionBean.setErrorMessage((String) arr[3]);
359        }
360        if (arr[4] != null) {
361            actionBean.setExternalId((String) arr[4]);
362        }
363        if (arr[5] != null) {
364            actionBean.setExternalStatus((String) arr[5]);
365        }
366        if (arr[6] != null) {
367            actionBean.setStatus(CoordinatorAction.Status.valueOf((String) arr[6]));
368        }
369        if (arr[7] != null) {
370            actionBean.setCreatedTime(DateUtils.toDate((Timestamp) arr[7]));
371        }
372        if (arr[8] != null) {
373            actionBean.setNominalTime(DateUtils.toDate((Timestamp) arr[8]));
374        }
375        if (arr[9] != null) {
376            actionBean.setMissingDependenciesBlob((StringBlob) arr[9]);
377        }
378        if (arr[10] != null) {
379            coordBean.setId((String) arr[10]);
380            actionBean.setJobId((String) arr[10]);
381        }
382        if (arr[11] != null) {
383            coordBean.setAppName((String) arr[11]);
384        }
385        if (arr[12] != null) {
386            coordBean.setStatus(CoordinatorJob.Status.valueOf((String) arr[12]));
387        }
388        bean.setBundle(bundleBean);
389        bean.setCoordinator(coordBean);
390        bean.setAction(actionBean);
391        return bean;
392    }
393
394    private BundleJobBean constructBundleBean(Object[] barr) throws JPAExecutorException {
395        BundleJobBean bean = new BundleJobBean();
396        if (barr[0] != null) {
397            bean.setId((String) barr[0]);
398        }
399        else {
400            throw new JPAExecutorException(ErrorCode.E0603,
401                    "bundleId returned by query is null - cannot retrieve bulk results");
402        }
403        if (barr[1] != null) {
404            bean.setAppName((String) barr[1]);
405        }
406        if (barr[2] != null) {
407            bean.setStatus(BundleJob.Status.valueOf((String) barr[2]));
408        }
409        if (barr[3] != null) {
410            bean.setUser((String) barr[3]);
411        }
412        return bean;
413    }
414
415    private void fillParameters(Query query, String prefix, List<?> values) {
416        for (int i = 0; i < values.size(); i++) {
417            query.setParameter(prefix + i, values.get(i));
418        }
419    }
420
421    // null safeguard
422    public static List<String> nullToEmpty(List<String> input) {
423        return input == null ? Collections.<String> emptyList() : input;
424    }
425
426}