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.coord;
020
021import com.google.common.collect.Lists;
022import org.apache.commons.lang.StringUtils;
023import org.apache.hadoop.conf.Configuration;
024import org.apache.oozie.ErrorCode;
025import org.apache.oozie.client.OozieClient;
026import org.apache.oozie.command.CommandException;
027import org.apache.oozie.dependency.URIHandler;
028import org.apache.oozie.dependency.URIHandler.Context;
029import org.apache.oozie.service.Services;
030import org.apache.oozie.service.URIHandlerService;
031import org.apache.oozie.util.DateUtils;
032import org.apache.oozie.util.ELEvaluator;
033import org.apache.oozie.util.ParamChecker;
034import org.apache.oozie.util.XLog;
035
036import java.net.URI;
037import java.util.ArrayList;
038import java.util.Calendar;
039import java.util.Date;
040import java.util.GregorianCalendar;
041import java.util.List;
042import java.util.TimeZone;
043
044/**
045 * This class implements the EL function related to coordinator
046 */
047
048public class CoordELFunctions {
049    final public static String DATASET = "oozie.coord.el.dataset.bean";
050    final public static String COORD_ACTION = "oozie.coord.el.app.bean";
051    final public static String CONFIGURATION = "oozie.coord.el.conf";
052    final public static String LATEST_EL_USE_CURRENT_TIME = "oozie.service.ELService.latest-el.use-current-time";
053    // INSTANCE_SEPARATOR is used to separate multiple directories into one tag.
054    final public static String INSTANCE_SEPARATOR = "#";
055    final public static String DIR_SEPARATOR = ",";
056    // TODO: in next release, support flexibility
057    private static String END_OF_OPERATION_INDICATOR_FILE = "_SUCCESS";
058
059    public static final long MINUTE_MSEC = 60 * 1000L;
060    public static final long HOUR_MSEC = 60 * MINUTE_MSEC;
061    public static final long DAY_MSEC = 24 * HOUR_MSEC;
062    public static final long MONTH_MSEC = 30 * DAY_MSEC;
063    public static final long YEAR_MSEC = 365 * DAY_MSEC;
064
065    /**
066     * Used in defining the frequency in 'day' unit. <p/> domain: <code> val &gt; 0</code> and should be integer.
067     *
068     * @param val frequency in number of days.
069     * @return number of days and also set the frequency timeunit to "day"
070     */
071    public static int ph1_coord_days(int val) {
072        val = ParamChecker.checkGTZero(val, "n");
073        ELEvaluator eval = ELEvaluator.getCurrent();
074        eval.setVariable("timeunit", TimeUnit.DAY);
075        eval.setVariable("endOfDuration", TimeUnit.NONE);
076        return val;
077    }
078
079    /**
080     * Used in defining the frequency in 'month' unit. <p/> domain: <code> val &gt; 0</code> and should be integer.
081     *
082     * @param val frequency in number of months.
083     * @return number of months and also set the frequency timeunit to "month"
084     */
085    public static int ph1_coord_months(int val) {
086        val = ParamChecker.checkGTZero(val, "n");
087        ELEvaluator eval = ELEvaluator.getCurrent();
088        eval.setVariable("timeunit", TimeUnit.MONTH);
089        eval.setVariable("endOfDuration", TimeUnit.NONE);
090        return val;
091    }
092
093    /**
094     * Used in defining the frequency in 'hour' unit. <p/> parameter value domain: <code> val &gt; 0</code> and should
095     * be integer.
096     *
097     * @param val frequency in number of hours.
098     * @return number of minutes and also set the frequency timeunit to "minute"
099     */
100    public static int ph1_coord_hours(int val) {
101        val = ParamChecker.checkGTZero(val, "n");
102        ELEvaluator eval = ELEvaluator.getCurrent();
103        eval.setVariable("timeunit", TimeUnit.MINUTE);
104        eval.setVariable("endOfDuration", TimeUnit.NONE);
105        return val * 60;
106    }
107
108    /**
109     * Used in defining the frequency in 'minute' unit. <p/> domain: <code> val &gt; 0</code> and should be integer.
110     *
111     * @param val frequency in number of minutes.
112     * @return number of minutes and also set the frequency timeunit to "minute"
113     */
114    public static int ph1_coord_minutes(int val) {
115        val = ParamChecker.checkGTZero(val, "n");
116        ELEvaluator eval = ELEvaluator.getCurrent();
117        eval.setVariable("timeunit", TimeUnit.MINUTE);
118        eval.setVariable("endOfDuration", TimeUnit.NONE);
119        return val;
120    }
121
122    /**
123     * Used in defining the frequency in 'day' unit and specify the "end of day" property. <p/> Every instance will
124     * start at 00:00 hour of each day. <p/> domain: <code> val &gt; 0</code> and should be integer.
125     *
126     * @param val frequency in number of days.
127     * @return number of days and also set the frequency timeunit to "day" and end_of_duration flag to "day"
128     */
129    public static int ph1_coord_endOfDays(int val) {
130        val = ParamChecker.checkGTZero(val, "n");
131        ELEvaluator eval = ELEvaluator.getCurrent();
132        eval.setVariable("timeunit", TimeUnit.DAY);
133        eval.setVariable("endOfDuration", TimeUnit.END_OF_DAY);
134        return val;
135    }
136
137    /**
138     * Used in defining the frequency in 'month' unit and specify the "end of month" property. <p/> Every instance will
139     * start at first day of each month at 00:00 hour. <p/> domain: <code> val &gt; 0</code> and should be integer.
140     *
141     * @param val: frequency in number of months.
142     * @return number of months and also set the frequency timeunit to "month" and end_of_duration flag to "month"
143     */
144    public static int ph1_coord_endOfMonths(int val) {
145        val = ParamChecker.checkGTZero(val, "n");
146        ELEvaluator eval = ELEvaluator.getCurrent();
147        eval.setVariable("timeunit", TimeUnit.MONTH);
148        eval.setVariable("endOfDuration", TimeUnit.END_OF_MONTH);
149        return val;
150    }
151
152    /**
153     * Calculate the difference of timezone offset in minutes between dataset and coordinator job. <p/> Depends on: <p/>
154     * 1. Timezone of both dataset and job <p/> 2. Action creation Time
155     *
156     * @return difference in minutes (DataSet TZ Offset - Application TZ offset)
157     */
158    public static int ph2_coord_tzOffset() {
159        long actionCreationTime = getActionCreationtime().getTime();
160        TimeZone dsTZ = ParamChecker.notNull(getDatasetTZ(), "DatasetTZ");
161        TimeZone jobTZ = ParamChecker.notNull(getJobTZ(), "JobTZ");
162        return (dsTZ.getOffset(actionCreationTime) - jobTZ.getOffset(actionCreationTime)) / (1000 * 60);
163    }
164
165    public static int ph3_coord_tzOffset() {
166        return ph2_coord_tzOffset();
167    }
168
169    /**
170     * Returns a date string that is offset from 'strBaseDate' by the amount specified.  The unit can be one of
171     * DAY, MONTH, HOUR, MINUTE, MONTH.
172     *
173     * @param strBaseDate The base date
174     * @param offset any number
175     * @param unit one of DAY, MONTH, HOUR, MINUTE, MONTH
176     * @return the offset date string
177     * @throws Exception
178     */
179    public static String ph2_coord_dateOffset(String strBaseDate, int offset, String unit) throws Exception {
180        Calendar baseCalDate = DateUtils.getCalendar(strBaseDate);
181        StringBuilder buffer = new StringBuilder();
182        baseCalDate.add(TimeUnit.valueOf(unit).getCalendarUnit(), offset);
183        buffer.append(DateUtils.formatDateOozieTZ(baseCalDate));
184        return buffer.toString();
185    }
186
187    public static String ph3_coord_dateOffset(String strBaseDate, int offset, String unit) throws Exception {
188        return ph2_coord_dateOffset(strBaseDate, offset, unit);
189    }
190
191    /**
192     * Returns a date string that is offset from 'strBaseDate' by the difference from Oozie processing timezone to the given
193     * timezone. It will account for daylight saving time based on the given 'strBaseDate' and 'timezone'.
194     *
195     * @param strBaseDate The base date
196     * @param timezone
197     * @return the offset date string
198     * @throws Exception
199     */
200    public static String ph2_coord_dateTzOffset(String strBaseDate, String timezone) throws Exception {
201        Calendar baseCalDate = DateUtils.getCalendar(strBaseDate);
202        StringBuilder buffer = new StringBuilder();
203        baseCalDate.setTimeZone(DateUtils.getTimeZone(timezone));
204        buffer.append(DateUtils.formatDate(baseCalDate));
205        return buffer.toString();
206    }
207
208    public static String ph3_coord_dateTzOffset(String strBaseDate, String timezone) throws Exception{
209        return ph2_coord_dateTzOffset(strBaseDate, timezone);
210    }
211
212    /**
213     * Determine the date-time in Oozie processing timezone of n-th future available dataset instance
214     * from nominal Time but not beyond the instance specified as 'instance.
215     * <p/>
216     * It depends on:
217     * <p/>
218     * 1. Data set frequency
219     * <p/>
220     * 2. Data set Time unit (day, month, minute)
221     * <p/>
222     * 3. Data set Time zone/DST
223     * <p/>
224     * 4. End Day/Month flag
225     * <p/>
226     * 5. Data set initial instance
227     * <p/>
228     * 6. Action Creation Time
229     * <p/>
230     * 7. Existence of dataset's directory
231     *
232     * @param n :instance count
233     *        <p/>
234     *        domain: n >= 0, n is integer
235     * @param instance: How many future instance it should check? value should
236     *        be >=0
237     * @return date-time in Oozie processing timezone of the n-th instance
238     *         <p/>
239     * @throws Exception
240     */
241    public static String ph3_coord_future(int n, int instance) throws Exception {
242        ParamChecker.checkGEZero(n, "future:n");
243        ParamChecker.checkGTZero(instance, "future:instance");
244        if (isSyncDataSet()) {// For Sync Dataset
245            return coord_future_sync(n, instance);
246        }
247        else {
248            throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet");
249        }
250    }
251
252    /**
253     * Determine the date-time in Oozie processing timezone of the future available dataset instances
254     * from start to end offsets from nominal Time but not beyond the instance specified as 'instance'.
255     * <p/>
256     * It depends on:
257     * <p/>
258     * 1. Data set frequency
259     * <p/>
260     * 2. Data set Time unit (day, month, minute)
261     * <p/>
262     * 3. Data set Time zone/DST
263     * <p/>
264     * 4. End Day/Month flag
265     * <p/>
266     * 5. Data set initial instance
267     * <p/>
268     * 6. Action Creation Time
269     * <p/>
270     * 7. Existence of dataset's directory
271     *
272     * @param start : start instance offset
273     *        <p/>
274     *        domain: start >= 0, start is integer
275     * @param end : end instance offset
276     *        <p/>
277     *        domain: end >= 0, end is integer
278     * @param instance: How many future instance it should check? value should
279     *        be >=0
280     * @return date-time in Oozie processing timezone of the instances from start to end offsets
281     *        delimited by comma.
282     *         <p/>
283     * @throws Exception
284     */
285    public static String ph3_coord_futureRange(int start, int end, int instance) throws Exception {
286        ParamChecker.checkGEZero(start, "future:n");
287        ParamChecker.checkGEZero(end, "future:n");
288        ParamChecker.checkGTZero(instance, "future:instance");
289        if (isSyncDataSet()) {// For Sync Dataset
290            return coord_futureRange_sync(start, end, instance);
291        }
292        else {
293            throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet");
294        }
295    }
296
297    private static String coord_future_sync(int n, int instance) throws Exception {
298        return coord_futureRange_sync(n, n, instance);
299    }
300
301    private static String coord_futureRange_sync(int startOffset, int endOffset, int instance) throws Exception {
302        final XLog LOG = XLog.getLog(CoordELFunctions.class);
303        final Thread currentThread = Thread.currentThread();
304        ELEvaluator eval = ELEvaluator.getCurrent();
305        String retVal = "";
306        int datasetFrequency = (int) getDSFrequency();// in minutes
307        TimeUnit dsTimeUnit = getDSTimeUnit();
308        int[] instCount = new int[1];
309        Calendar nominalInstanceCal = getCurrentInstance(getActionCreationtime(), instCount);
310        StringBuilder resolvedInstances = new StringBuilder();
311        StringBuilder resolvedURIPaths = new StringBuilder();
312        if (nominalInstanceCal != null) {
313            Calendar initInstance = getInitialInstanceCal();
314            nominalInstanceCal = (Calendar) initInstance.clone();
315            nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), instCount[0] * datasetFrequency);
316
317            SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET);
318            if (ds == null) {
319                throw new RuntimeException("Associated Dataset should be defined with key " + DATASET);
320            }
321            String uriTemplate = ds.getUriTemplate();
322            Configuration conf = (Configuration) eval.getVariable(CONFIGURATION);
323            if (conf == null) {
324                throw new RuntimeException("Associated Configuration should be defined with key " + CONFIGURATION);
325            }
326            int available = 0, checkedInstance = 0;
327            boolean resolved = false;
328            String user = ParamChecker
329                    .notEmpty((String) eval.getVariable(OozieClient.USER_NAME), OozieClient.USER_NAME);
330            String doneFlag = ds.getDoneFlag();
331            URIHandlerService uriService = Services.get().get(URIHandlerService.class);
332            URIHandler uriHandler = null;
333            Context uriContext = null;
334            try {
335                while (instance >= checkedInstance && !currentThread.isInterrupted()) {
336                    ELEvaluator uriEval = getUriEvaluator(nominalInstanceCal);
337                    String uriPath = uriEval.evaluate(uriTemplate, String.class);
338                    if (uriHandler == null) {
339                        URI uri = new URI(uriPath);
340                        uriHandler = uriService.getURIHandler(uri);
341                        uriContext = uriHandler.getContext(uri, conf, user);
342                    }
343                    String uriWithDoneFlag = uriHandler.getURIWithDoneFlag(uriPath, doneFlag);
344                    if (uriHandler.exists(new URI(uriWithDoneFlag), uriContext)) {
345                        if (available == endOffset) {
346                            LOG.debug("Matched future(" + available + "): " + uriWithDoneFlag);
347                            resolved = true;
348                            resolvedInstances.append(DateUtils.formatDateOozieTZ(nominalInstanceCal));
349                            resolvedURIPaths.append(uriPath);
350                            retVal = resolvedInstances.toString();
351                            eval.setVariable("resolved_path", resolvedURIPaths.toString());
352                            break;
353                        }
354                        else if (available >= startOffset) {
355                            LOG.debug("Matched future(" + available + "): " + uriWithDoneFlag);
356                            resolvedInstances.append(DateUtils.formatDateOozieTZ(nominalInstanceCal)).append(
357                                    INSTANCE_SEPARATOR);
358                            resolvedURIPaths.append(uriPath).append(INSTANCE_SEPARATOR);
359                        }
360                        available++;
361                    }
362                    // nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), datasetFrequency);
363                    nominalInstanceCal = (Calendar) initInstance.clone();
364                    instCount[0]++;
365                    nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), instCount[0] * datasetFrequency);
366                    checkedInstance++;
367                    // DateUtils.moveToEnd(nominalInstanceCal, getDSEndOfFlag());
368                }
369            }
370            finally {
371                if (uriContext != null) {
372                    uriContext.destroy();
373                }
374            }
375            if (!resolved) {
376                // return unchanged future function with variable 'is_resolved'
377                // to 'false'
378                eval.setVariable("is_resolved", Boolean.FALSE);
379                if (startOffset == endOffset) {
380                    retVal = "${coord:future(" + startOffset + ", " + instance + ")}";
381                }
382                else {
383                    retVal = "${coord:futureRange(" + startOffset + ", " + endOffset + ", " + instance + ")}";
384                }
385            }
386            else {
387                eval.setVariable("is_resolved", Boolean.TRUE);
388            }
389        }
390        else {// No feasible nominal time
391            eval.setVariable("is_resolved", Boolean.TRUE);
392            retVal = "";
393        }
394        return retVal;
395    }
396
397    /**
398     * Return nominal time or Action Creation Time.
399     * <p/>
400     *
401     * @return coordinator action creation or materialization date time
402     * @throws Exception if unable to format the Date object to String
403     */
404    public static String ph2_coord_nominalTime() throws Exception {
405        ELEvaluator eval = ELEvaluator.getCurrent();
406        SyncCoordAction action = ParamChecker.notNull((SyncCoordAction) eval.getVariable(COORD_ACTION),
407                "Coordinator Action");
408        return DateUtils.formatDateOozieTZ(action.getNominalTime());
409    }
410
411    public static String ph3_coord_nominalTime() throws Exception {
412        return ph2_coord_nominalTime();
413    }
414
415    /**
416     * Convert from standard date-time formatting to a desired format.
417     * <p/>
418     * @param dateTimeStr - A timestamp in standard (ISO8601) format.
419     * @param format - A string representing the desired format.
420     * @return coordinator action creation or materialization date time
421     * @throws Exception if unable to format the Date object to String
422     */
423    public static String ph2_coord_formatTime(String dateTimeStr, String format)
424            throws Exception {
425        Date dateTime = DateUtils.parseDateOozieTZ(dateTimeStr);
426        return DateUtils.formatDateCustom(dateTime, format);
427    }
428
429    public static String ph3_coord_formatTime(String dateTimeStr, String format)
430            throws Exception {
431        return ph2_coord_formatTime(dateTimeStr, format);
432    }
433
434    /**
435     * Return Action Id. <p/>
436     * Convert from standard date-time formatting to a Unix epoch time.
437     * <p/>
438     * @param dateTimeStr - A timestamp in standard (ISO8601) format.
439     * @param millis - "true" to include millis; otherwise will only include seconds
440     * @return coordinator action creation or materialization date time
441     * @throws Exception if unable to format the Date object to String
442     */
443    public static String ph2_coord_epochTime(String dateTimeStr, String millis)
444            throws Exception {
445        Date dateTime = DateUtils.parseDateOozieTZ(dateTimeStr);
446        return DateUtils.formatDateEpoch(dateTime, Boolean.valueOf(millis));
447    }
448
449    public static String ph3_coord_epochTime(String dateTimeStr, String millis)
450            throws Exception {
451        return  ph2_coord_epochTime(dateTimeStr, millis);
452    }
453
454    /**
455     * Return Action Id. <p>
456     *
457     * @return coordinator action Id
458     */
459    public static String ph2_coord_actionId() throws Exception {
460        ELEvaluator eval = ELEvaluator.getCurrent();
461        SyncCoordAction action = ParamChecker.notNull((SyncCoordAction) eval.getVariable(COORD_ACTION),
462                "Coordinator Action");
463        return action.getActionId();
464    }
465
466    public static String ph3_coord_actionId() throws Exception {
467        return ph2_coord_actionId();
468    }
469
470    /**
471     * Return Job Name. <p/>
472     *
473     * @return coordinator name
474     */
475    public static String ph2_coord_name() throws Exception {
476        ELEvaluator eval = ELEvaluator.getCurrent();
477        SyncCoordAction action = ParamChecker.notNull((SyncCoordAction) eval.getVariable(COORD_ACTION),
478                "Coordinator Action");
479        return action.getName();
480    }
481
482    public static String ph3_coord_name() throws Exception {
483        return ph2_coord_name();
484    }
485
486    /**
487     * Return Action Start time. <p/>
488     *
489     * @return coordinator action start time
490     * @throws Exception if unable to format the Date object to String
491     */
492    public static String ph2_coord_actualTime() throws Exception {
493        ELEvaluator eval = ELEvaluator.getCurrent();
494        SyncCoordAction coordAction = (SyncCoordAction) eval.getVariable(COORD_ACTION);
495        if (coordAction == null) {
496            throw new RuntimeException("Associated Application instance should be defined with key " + COORD_ACTION);
497        }
498        return DateUtils.formatDateOozieTZ(coordAction.getActualTime());
499    }
500
501    public static String ph3_coord_actualTime() throws Exception {
502        return ph2_coord_actualTime();
503    }
504
505    /**
506     * Used to specify a list of URI's that are used as input dir to the workflow job. <p/> Look for two evaluator-level
507     * variables <p/> A) .datain.<DATAIN_NAME> B) .datain.<DATAIN_NAME>.unresolved <p/> A defines the current list of
508     * URI. <p/> B defines whether there are any unresolved EL-function (i.e latest) <p/> If there are something
509     * unresolved, this function will echo back the original function <p/> otherwise it sends the uris.
510     *
511     * @param dataInName : Datain name
512     * @return the list of URI's separated by INSTANCE_SEPARATOR <p/> if there are unresolved EL function (i.e. latest)
513     *         , echo back <p/> the function without resolving the function.
514     */
515    public static String ph3_coord_dataIn(String dataInName) {
516        String uris = "";
517        ELEvaluator eval = ELEvaluator.getCurrent();
518        uris = (String) eval.getVariable(".datain." + dataInName);
519        Boolean unresolved = (Boolean) eval.getVariable(".datain." + dataInName + ".unresolved");
520        if (unresolved != null && unresolved.booleanValue() == true) {
521            return "${coord:dataIn('" + dataInName + "')}";
522        }
523        return uris;
524    }
525
526    /**
527     * Used to specify a list of URI's that are output dir of the workflow job. <p/> Look for one evaluator-level
528     * variable <p/> dataout.<DATAOUT_NAME> <p/> It defines the current list of URI. <p/> otherwise it sends the uris.
529     *
530     * @param dataOutName : Dataout name
531     * @return the list of URI's separated by INSTANCE_SEPARATOR
532     */
533    public static String ph3_coord_dataOut(String dataOutName) {
534        String uris = "";
535        ELEvaluator eval = ELEvaluator.getCurrent();
536        uris = (String) eval.getVariable(".dataout." + dataOutName);
537        return uris;
538    }
539
540    /**
541     * Determine the date-time in Oozie processing timezone of n-th dataset instance. <p/> It depends on: <p/> 1.
542     * Data set frequency <p/> 2.
543     * Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST <p/> 4. End Day/Month flag <p/> 5. Data
544     * set initial instance <p/> 6. Action Creation Time
545     *
546     * @param n instance count domain: n is integer
547     * @return date-time in Oozie processing timezone of the n-th instance returns 'null' means n-th instance is
548     * earlier than Initial-Instance of DS
549     * @throws Exception
550     */
551    public static String ph2_coord_current(int n) throws Exception {
552        if (isSyncDataSet()) { // For Sync Dataset
553            return coord_current_sync(n);
554        }
555        else {
556            throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet");
557        }
558    }
559
560    /**
561     * Determine the date-time in Oozie processing timezone of current dataset instances
562     * from start to end offsets from the nominal time. <p/> It depends
563     * on: <p/> 1. Data set frequency <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST
564     * <p/> 4. End Day/Month flag <p/> 5. Data set initial instance <p/> 6. Action Creation Time
565     *
566     * @param start :start instance offset <p/> domain: start <= 0, start is integer
567     * @param end :end instance offset <p/> domain: end <= 0, end is integer
568     * @return date-time in Oozie processing timezone of the instances from start to end offsets
569     *        delimited by comma. <p/> If the current instance time of the dataset based on the Action Creation Time
570     *        is earlier than the Initial-Instance of DS an empty string is returned.
571     *        If an instance within the range is earlier than Initial-Instance of DS that instance is ignored
572     * @throws Exception
573     */
574    public static String ph2_coord_currentRange(int start, int end) throws Exception {
575        if (isSyncDataSet()) { // For Sync Dataset
576            return coord_currentRange_sync(start, end);
577        }
578        else {
579            throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet");
580        }
581    }
582    /**
583     * Determine the date-time in Oozie processing timezone of the given offset from the dataset effective nominal time. <p/> It
584     * depends on: <p> 1. Data set frequency <p/> 2. Data set Time Unit <p/> 3. Data set Time zone/DST
585     * <p/> 4. Data set initial instance <p/> 5. Action Creation Time
586     *
587     * @param n offset amount (integer)
588     * @param timeUnit TimeUnit for offset n ("MINUTE", "HOUR", "DAY", "MONTH", "YEAR")
589     * @return date-time in Oozie processing timezone of the given offset from the dataset effective nominal time
590     * @throws Exception if there was a problem formatting
591     */
592    public static String ph2_coord_offset(int n, String timeUnit) throws Exception {
593        if (isSyncDataSet()) { // For Sync Dataset
594            return coord_offset_sync(n, timeUnit);
595        }
596        else {
597            throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet");
598        }
599    }
600
601    /**
602     * Determine how many hours is on the date of n-th dataset instance. <p/> It depends on: <p/> 1. Data set frequency
603     * <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST <p/> 4. End Day/Month flag <p/> 5.
604     * Data set initial instance <p/> 6. Action Creation Time
605     *
606     * @param n instance count <p/> domain: n is integer
607     * @return number of hours on that day <p/> returns -1 means n-th instance is earlier than Initial-Instance of DS
608     * @throws Exception
609     */
610    public static int ph2_coord_hoursInDay(int n) throws Exception {
611        int datasetFrequency = (int) getDSFrequency();
612        // /Calendar nominalInstanceCal =
613        // getCurrentInstance(getActionCreationtime());
614        Calendar nominalInstanceCal = getEffectiveNominalTime();
615        if (nominalInstanceCal == null) {
616            return -1;
617        }
618        nominalInstanceCal.add(getDSTimeUnit().getCalendarUnit(), datasetFrequency * n);
619        /*
620         * if (nominalInstanceCal.getTime().compareTo(getInitialInstance()) < 0)
621         * { return -1; }
622         */
623        nominalInstanceCal.setTimeZone(getDatasetTZ());// Use Dataset TZ
624        // DateUtils.moveToEnd(nominalInstanceCal, getDSEndOfFlag());
625        return DateUtils.hoursInDay(nominalInstanceCal);
626    }
627
628    public static int ph3_coord_hoursInDay(int n) throws Exception {
629        return ph2_coord_hoursInDay(n);
630    }
631
632    /**
633     * Calculate number of days in one month for n-th dataset instance. <p/> It depends on: <p/> 1. Data set frequency .
634     * <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST <p/> 4. End Day/Month flag <p/> 5.
635     * Data set initial instance <p/> 6. Action Creation Time
636     *
637     * @param n instance count. domain: n is integer
638     * @return number of days in that month <p/> returns -1 means n-th instance is earlier than Initial-Instance of DS
639     * @throws Exception
640     */
641    public static int ph2_coord_daysInMonth(int n) throws Exception {
642        int datasetFrequency = (int) getDSFrequency();// in minutes
643        // Calendar nominalInstanceCal =
644        // getCurrentInstance(getActionCreationtime());
645        Calendar nominalInstanceCal = getEffectiveNominalTime();
646        if (nominalInstanceCal == null) {
647            return -1;
648        }
649        nominalInstanceCal.add(getDSTimeUnit().getCalendarUnit(), datasetFrequency * n);
650        /*
651         * if (nominalInstanceCal.getTime().compareTo(getInitialInstance()) < 0)
652         * { return -1; }
653         */
654        nominalInstanceCal.setTimeZone(getDatasetTZ());// Use Dataset TZ
655        // DateUtils.moveToEnd(nominalInstanceCal, getDSEndOfFlag());
656        return nominalInstanceCal.getActualMaximum(Calendar.DAY_OF_MONTH);
657    }
658
659    public static int ph3_coord_daysInMonth(int n) throws Exception {
660        return ph2_coord_daysInMonth(n);
661    }
662
663    /**
664     * Determine the date-time in Oozie processing timezone of n-th latest available dataset instance. <p/> It depends
665     * on: <p/> 1. Data set frequency <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST
666     * <p/> 4. End Day/Month flag <p/> 5. Data set initial instance <p/> 6. Action Creation Time <p/> 7. Existence of
667     * dataset's directory
668     *
669     * @param n :instance count <p/> domain: n <= 0, n is integer
670     * @return date-time in Oozie processing timezone of the n-th instance <p/> returns 'null' means n-th instance is
671     * earlier than Initial-Instance of DS
672     * @throws Exception
673     */
674    public static String ph3_coord_latest(int n) throws Exception {
675        ParamChecker.checkLEZero(n, "latest:n");
676        if (isSyncDataSet()) {// For Sync Dataset
677            return coord_latest_sync(n);
678        }
679        else {
680            throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet");
681        }
682    }
683
684    /**
685     * Determine the date-time in Oozie processing timezone of latest available dataset instances
686     * from start to end offsets from the nominal time. <p/> It depends
687     * on: <p/> 1. Data set frequency <p/> 2. Data set Time unit (day, month, minute) <p/> 3. Data set Time zone/DST
688     * <p/> 4. End Day/Month flag <p/> 5. Data set initial instance <p/> 6. Action Creation Time <p/> 7. Existence of
689     * dataset's directory
690     *
691     * @param start :start instance offset <p/> domain: start <= 0, start is integer
692     * @param end :end instance offset <p/> domain: end <= 0, end is integer
693     * @return date-time in Oozie processing timezone of the instances from start to end offsets
694     *        delimited by comma. <p/> returns 'null' means start offset instance is
695     *        earlier than Initial-Instance of DS
696     * @throws Exception
697     */
698    public static String ph3_coord_latestRange(int start, int end) throws Exception {
699        ParamChecker.checkLEZero(start, "latest:n");
700        ParamChecker.checkLEZero(end, "latest:n");
701        if (isSyncDataSet()) {// For Sync Dataset
702            return coord_latestRange_sync(start, end);
703        }
704        else {
705            throw new UnsupportedOperationException("Asynchronous Dataset is not supported yet");
706        }
707    }
708
709    /**
710     * Configure an evaluator with data set and application specific information. <p/> Helper method of associating
711     * dataset and application object
712     *
713     * @param evaluator : to set variables
714     * @param ds : Data Set object
715     * @param coordAction : Application instance
716     */
717    public static void configureEvaluator(ELEvaluator evaluator, SyncCoordDataset ds, SyncCoordAction coordAction) {
718        evaluator.setVariable(COORD_ACTION, coordAction);
719        evaluator.setVariable(DATASET, ds);
720    }
721
722    /**
723     * Helper method to wrap around with "${..}". <p/>
724     *
725     *
726     * @param eval :EL evaluator
727     * @param expr : expression to evaluate
728     * @return Resolved expression or echo back the same expression
729     * @throws Exception
730     */
731    public static String evalAndWrap(ELEvaluator eval, String expr) throws Exception {
732        try {
733            eval.setVariable(".wrap", null);
734            String result = eval.evaluate(expr, String.class);
735            if (eval.getVariable(".wrap") != null) {
736                return "${" + result + "}";
737            }
738            else {
739                return result;
740            }
741        }
742        catch (Exception e) {
743            throw new Exception("Unable to evaluate :" + expr + ":\n", e);
744        }
745    }
746
747    // Set of echo functions
748
749    public static String ph1_coord_current_echo(String n) {
750        return echoUnResolved("current", n);
751    }
752
753    public static String ph1_coord_absolute_echo(String date) {
754        return echoUnResolved("absolute", date);
755    }
756
757    public static String ph1_coord_currentRange_echo(String start, String end) {
758        return echoUnResolved("currentRange", start + ", " + end);
759    }
760
761    public static String ph1_coord_offset_echo(String n, String timeUnit) {
762        return echoUnResolved("offset", n + " , " + timeUnit);
763    }
764
765    public static String ph2_coord_current_echo(String n) {
766        return echoUnResolved("current", n);
767    }
768
769    public static String ph2_coord_currentRange_echo(String start, String end) {
770        return echoUnResolved("currentRange", start + ", " + end);
771    }
772
773    public static String ph2_coord_offset_echo(String n, String timeUnit) {
774        return echoUnResolved("offset", n + " , " + timeUnit);
775    }
776
777    public static String ph2_coord_absolute_echo(String date) {
778        return echoUnResolved("absolute", date);
779    }
780
781    public static String ph2_coord_absolute_range(String startInstance, int end) throws Exception {
782        int[] instanceCount = new int[1];
783
784        // getCurrentInstance() returns null, which means startInstance is less
785        // than initial instance
786        if (getCurrentInstance(DateUtils.getCalendar(startInstance).getTime(), instanceCount) == null) {
787            throw new CommandException(ErrorCode.E1010,
788                    "intial-instance should be equal or earlier than the start-instance. intial-instance is "
789                            + getInitialInstance() + " and start-instance is " + startInstance);
790        }
791        int[] nominalCount = new int[1];
792        if (getCurrentInstance(getActionCreationtime(), nominalCount) == null) {
793            throw new CommandException(ErrorCode.E1010,
794                    "intial-instance should be equal or earlier than the nominal time. intial-instance is "
795                            + getInitialInstance() + " and nominal time is " + getActionCreationtime());
796        }
797        // getCurrentInstance return offset relative to initial instance.
798        // start instance offset - nominal offset = start offset relative to
799        // nominal time-stamp.
800        int start = instanceCount[0] - nominalCount[0];
801        if (start > end) {
802            throw new CommandException(ErrorCode.E1010,
803                    "start-instance should be equal or earlier than the end-instance. startInstance is "
804                            + startInstance + " which is equivalent to current (" + instanceCount[0]
805                            + ") but end is specified as current (" + end + ")");
806        }
807        return ph2_coord_currentRange(start, end);
808    }
809
810    public static String ph1_coord_dateOffset_echo(String n, String offset, String unit) {
811        return echoUnResolved("dateOffset", n + " , " + offset + " , " + unit);
812    }
813
814    public static String ph1_coord_dateTzOffset_echo(String n, String timezone) {
815        return echoUnResolved("dateTzOffset", n + " , " + timezone);
816    }
817
818    public static String ph1_coord_epochTime_echo(String dateTime, String millis) {
819        // Quote the dateTime value since it would contain a ':'.
820        return echoUnResolved("epochTime", "'"+dateTime+"'" + " , " + millis);
821    }
822
823    public static String ph1_coord_formatTime_echo(String dateTime, String format) {
824        // Quote the dateTime value since it would contain a ':'.
825        return echoUnResolved("formatTime", "'"+dateTime+"'" + " , " + format);
826    }
827
828    public static String ph1_coord_latest_echo(String n) {
829        return echoUnResolved("latest", n);
830    }
831
832    public static String ph2_coord_latest_echo(String n) {
833        return ph1_coord_latest_echo(n);
834    }
835
836    public static String ph1_coord_future_echo(String n, String instance) {
837        return echoUnResolved("future", n + ", " + instance + "");
838    }
839
840    public static String ph2_coord_future_echo(String n, String instance) {
841        return ph1_coord_future_echo(n, instance);
842    }
843
844    public static String ph1_coord_latestRange_echo(String start, String end) {
845        return echoUnResolved("latestRange", start + ", " + end);
846    }
847
848    public static String ph2_coord_latestRange_echo(String start, String end) {
849        return ph1_coord_latestRange_echo(start, end);
850    }
851
852    public static String ph1_coord_futureRange_echo(String start, String end, String instance) {
853        return echoUnResolved("futureRange", start + ", " + end + ", " + instance);
854    }
855
856    public static String ph2_coord_futureRange_echo(String start, String end, String instance) {
857        return ph1_coord_futureRange_echo(start, end, instance);
858    }
859
860    public static String ph1_coord_dataIn_echo(String n) {
861        ELEvaluator eval = ELEvaluator.getCurrent();
862        String val = (String) eval.getVariable("oozie.dataname." + n);
863        if (val == null || val.equals("data-in") == false) {
864            XLog.getLog(CoordELFunctions.class).error("data_in_name " + n + " is not valid");
865            throw new RuntimeException("data_in_name " + n + " is not valid");
866        }
867        return echoUnResolved("dataIn", "'" + n + "'");
868    }
869
870    public static String ph1_coord_dataOut_echo(String n) {
871        ELEvaluator eval = ELEvaluator.getCurrent();
872        String val = (String) eval.getVariable("oozie.dataname." + n);
873        if (val == null || val.equals("data-out") == false) {
874            XLog.getLog(CoordELFunctions.class).error("data_out_name " + n + " is not valid");
875            throw new RuntimeException("data_out_name " + n + " is not valid");
876        }
877        return echoUnResolved("dataOut", "'" + n + "'");
878    }
879
880    public static String ph1_coord_nominalTime_echo() {
881        return echoUnResolved("nominalTime", "");
882    }
883
884    public static String ph1_coord_nominalTime_echo_wrap() {
885        // return "${coord:nominalTime()}"; // no resolution
886        return echoUnResolved("nominalTime", "");
887    }
888
889    public static String ph1_coord_nominalTime_echo_fixed() {
890        return "2009-03-06T010:00"; // Dummy resolution
891    }
892
893    public static String ph1_coord_actualTime_echo_wrap() {
894        // return "${coord:actualTime()}"; // no resolution
895        return echoUnResolved("actualTime", "");
896    }
897
898    public static String ph1_coord_actionId_echo() {
899        return echoUnResolved("actionId", "");
900    }
901
902    public static String ph1_coord_name_echo() {
903        return echoUnResolved("name", "");
904    }
905
906    // The following echo functions are not used in any phases yet
907    // They are here for future purpose.
908    public static String coord_minutes_echo(String n) {
909        return echoUnResolved("minutes", n);
910    }
911
912    public static String coord_hours_echo(String n) {
913        return echoUnResolved("hours", n);
914    }
915
916    public static String coord_days_echo(String n) {
917        return echoUnResolved("days", n);
918    }
919
920    public static String coord_endOfDay_echo(String n) {
921        return echoUnResolved("endOfDay", n);
922    }
923
924    public static String coord_months_echo(String n) {
925        return echoUnResolved("months", n);
926    }
927
928    public static String coord_endOfMonth_echo(String n) {
929        return echoUnResolved("endOfMonth", n);
930    }
931
932    public static String coord_actualTime_echo() {
933        return echoUnResolved("actualTime", "");
934    }
935
936    // This echo function will always return "24" for validation only.
937    // This evaluation ****should not**** replace the original XML
938    // Create a temporary string and validate the function
939    // This is **required** for evaluating an expression like
940    // coord:HoursInDay(0) + 3
941    // actual evaluation will happen in phase 2 or phase 3.
942    public static String ph1_coord_hoursInDay_echo(String n) {
943        return "24";
944        // return echoUnResolved("hoursInDay", n);
945    }
946
947    // This echo function will always return "30" for validation only.
948    // This evaluation ****should not**** replace the original XML
949    // Create a temporary string and validate the function
950    // This is **required** for evaluating an expression like
951    // coord:daysInMonth(0) + 3
952    // actual evaluation will happen in phase 2 or phase 3.
953    public static String ph1_coord_daysInMonth_echo(String n) {
954        // return echoUnResolved("daysInMonth", n);
955        return "30";
956    }
957
958    // This echo function will always return "3" for validation only.
959    // This evaluation ****should not**** replace the original XML
960    // Create a temporary string and validate the function
961    // This is **required** for evaluating an expression like coord:tzOffset + 2
962    // actual evaluation will happen in phase 2 or phase 3.
963    public static String ph1_coord_tzOffset_echo() {
964        // return echoUnResolved("tzOffset", "");
965        return "3";
966    }
967
968    // Local methods
969    /**
970     * @param n
971     * @return n-th instance Date-Time from current instance for data-set <p/> return empty string ("") if the
972     *         Action_Creation_time or the n-th instance <p/> is earlier than the Initial_Instance of dataset.
973     * @throws Exception
974     */
975    private static String coord_current_sync(int n) throws Exception {
976        return coord_currentRange_sync(n, n);
977    }
978
979    private static String coord_currentRange_sync(int start, int end) throws Exception {
980        final XLog LOG = XLog.getLog(CoordELFunctions.class);
981        int datasetFrequency = getDSFrequency();// in minutes
982        TimeUnit dsTimeUnit = getDSTimeUnit();
983        int[] instCount = new int[1];// used as pass by ref
984        Calendar nominalInstanceCal = getCurrentInstance(getActionCreationtime(), instCount);
985        if (nominalInstanceCal == null) {
986            LOG.warn("If the initial instance of the dataset is later than the nominal time, an empty string is"
987                    + " returned. This means that no data is available at the current-instance specified by the user"
988                    + " and the user could try modifying his initial-instance to an earlier time.");
989            return "";
990        } else {
991            Calendar initInstance = getInitialInstanceCal();
992            // Add in the reverse order - newest instance first.
993            nominalInstanceCal = (Calendar) initInstance.clone();
994            nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), (instCount[0] + start) * datasetFrequency);
995            List<String> instances = new ArrayList<String>();
996            for (int i = start; i <= end; i++) {
997                if (nominalInstanceCal.compareTo(initInstance) < 0) {
998                    LOG.warn("If the initial instance of the dataset is later than the current-instance specified,"
999                            + " such as coord:current({0}) in this case, an empty string is returned. This means that"
1000                            + " no data is available at the current-instance specified by the user and the user could"
1001                            + " try modifying his initial-instance to an earlier time.", start);
1002                }
1003                else {
1004                    instances.add(DateUtils.formatDateOozieTZ(nominalInstanceCal));
1005                }
1006                nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), datasetFrequency);
1007            }
1008            instances = Lists.reverse(instances);
1009            return StringUtils.join(instances, CoordELFunctions.INSTANCE_SEPARATOR);
1010        }
1011    }
1012
1013    /**
1014     *
1015     * @param n offset amount (integer)
1016     * @param timeUnit TimeUnit for offset n ("MINUTE", "HOUR", "DAY", "MONTH", "YEAR")
1017     * @return the offset time from the effective nominal time <p/> return empty string ("") if the Action_Creation_time or the
1018     *         offset instance <p/> is earlier than the Initial_Instance of dataset.
1019     * @throws Exception
1020     */
1021    private static String coord_offset_sync(int n, String timeUnit) throws Exception {
1022        Calendar rawCal = resolveOffsetRawTime(n, TimeUnit.valueOf(timeUnit), null);
1023        if (rawCal == null) {
1024            // warning already logged by resolveOffsetRawTime()
1025            return "";
1026        }
1027
1028        int freq = getDSFrequency();
1029        TimeUnit freqUnit = getDSTimeUnit();
1030        int freqCount = 0;
1031        // We're going to manually turn back/forward cal by decrements/increments of freq and then check that it gives the same
1032        // time as rawCal; this is to check that the offset time resolves to a frequency offset of the effective nominal time
1033        // In other words, that there exists an integer x, such that coord:offset(n, timeUnit) == coord:current(x) is true
1034        // If not, then we'll "rewind" rawCal to the latest instance earlier than rawCal and use that.
1035        Calendar cal = getInitialInstanceCal();
1036        if (rawCal.before(cal)) {
1037            while (cal.after(rawCal)) {
1038                cal.add(freqUnit.getCalendarUnit(), -freq);
1039                freqCount--;
1040            }
1041        }
1042        else if (rawCal.after(cal)) {
1043            while (cal.before(rawCal)) {
1044                cal.add(freqUnit.getCalendarUnit(), freq);
1045                freqCount++;
1046            }
1047        }
1048        if (cal.before(rawCal)) {
1049            rawCal = cal;
1050        }
1051        else if (cal.after(rawCal)) {
1052            cal.add(freqUnit.getCalendarUnit(), -freq);
1053            rawCal = cal;
1054            freqCount--;
1055        }
1056        String rawCalStr = DateUtils.formatDateOozieTZ(rawCal);
1057
1058        Calendar nominalInstanceCal = getInitialInstanceCal();
1059        nominalInstanceCal.add(freqUnit.getCalendarUnit(), freq * freqCount);
1060        if (nominalInstanceCal.getTime().compareTo(getInitialInstance()) < 0) {
1061            XLog.getLog(CoordELFunctions.class).warn("If the initial instance of the dataset is later than the offset instance"
1062                    + " specified, such as coord:offset({0}, {1}) in this case, an empty string is returned. This means that no"
1063                    + " data is available at the offset instance specified by the user and the user could try modifying his"
1064                    + " initial-instance to an earlier time.", n, timeUnit);
1065            return "";
1066        }
1067        String nominalCalStr = DateUtils.formatDateOozieTZ(nominalInstanceCal);
1068
1069        if (!rawCalStr.equals(nominalCalStr)) {
1070            throw new RuntimeException("Shouldn't happen");
1071        }
1072        return rawCalStr;
1073    }
1074
1075    /**
1076     * @param offset
1077     * @return n-th available latest instance Date-Time for SYNC data-set
1078     * @throws Exception
1079     */
1080    private static String coord_latest_sync(int offset) throws Exception {
1081        return coord_latestRange_sync(offset, offset);
1082    }
1083
1084    private static String coord_latestRange_sync(int startOffset, int endOffset) throws Exception {
1085        final XLog LOG = XLog.getLog(CoordELFunctions.class);
1086        final Thread currentThread = Thread.currentThread();
1087        ELEvaluator eval = ELEvaluator.getCurrent();
1088        String retVal = "";
1089        int datasetFrequency = (int) getDSFrequency();// in minutes
1090        TimeUnit dsTimeUnit = getDSTimeUnit();
1091        int[] instCount = new int[1];
1092        boolean useCurrentTime = Services.get().getConf().getBoolean(LATEST_EL_USE_CURRENT_TIME, false);
1093        Calendar nominalInstanceCal;
1094        if (useCurrentTime) {
1095            nominalInstanceCal = getCurrentInstance(new Date(), instCount);
1096        }
1097        else {
1098            nominalInstanceCal = getCurrentInstance(getActualTime(), instCount);
1099        }
1100        StringBuilder resolvedInstances = new StringBuilder();
1101        StringBuilder resolvedURIPaths = new StringBuilder();
1102        if (nominalInstanceCal != null) {
1103            Calendar initInstance = getInitialInstanceCal();
1104            SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET);
1105            if (ds == null) {
1106                throw new RuntimeException("Associated Dataset should be defined with key " + DATASET);
1107            }
1108            String uriTemplate = ds.getUriTemplate();
1109            Configuration conf = (Configuration) eval.getVariable(CONFIGURATION);
1110            if (conf == null) {
1111                throw new RuntimeException("Associated Configuration should be defined with key " + CONFIGURATION);
1112            }
1113            int available = 0;
1114            boolean resolved = false;
1115            String user = ParamChecker
1116                    .notEmpty((String) eval.getVariable(OozieClient.USER_NAME), OozieClient.USER_NAME);
1117            String doneFlag = ds.getDoneFlag();
1118            URIHandlerService uriService = Services.get().get(URIHandlerService.class);
1119            URIHandler uriHandler = null;
1120            Context uriContext = null;
1121            try {
1122                while (nominalInstanceCal.compareTo(initInstance) >= 0 && !currentThread.isInterrupted()) {
1123                    ELEvaluator uriEval = getUriEvaluator(nominalInstanceCal);
1124                    String uriPath = uriEval.evaluate(uriTemplate, String.class);
1125                    if (uriHandler == null) {
1126                        URI uri = new URI(uriPath);
1127                        uriHandler = uriService.getURIHandler(uri);
1128                        uriContext = uriHandler.getContext(uri, conf, user);
1129                    }
1130                    String uriWithDoneFlag = uriHandler.getURIWithDoneFlag(uriPath, doneFlag);
1131                    if (uriHandler.exists(new URI(uriWithDoneFlag), uriContext)) {
1132                        XLog.getLog(CoordELFunctions.class)
1133                        .debug("Found latest(" + available + "): " + uriWithDoneFlag);
1134                        if (available == startOffset) {
1135                            LOG.debug("Matched latest(" + available + "): " + uriWithDoneFlag);
1136                            resolved = true;
1137                            resolvedInstances.append(DateUtils.formatDateOozieTZ(nominalInstanceCal));
1138                            resolvedURIPaths.append(uriPath);
1139                            retVal = resolvedInstances.toString();
1140                            eval.setVariable("resolved_path", resolvedURIPaths.toString());
1141                            break;
1142                        }
1143                        else if (available <= endOffset) {
1144                            LOG.debug("Matched latest(" + available + "): " + uriWithDoneFlag);
1145                            resolvedInstances.append(DateUtils.formatDateOozieTZ(nominalInstanceCal)).append(
1146                                    INSTANCE_SEPARATOR);
1147                            resolvedURIPaths.append(uriPath).append(INSTANCE_SEPARATOR);
1148                        }
1149
1150                        available--;
1151                    }
1152                    // nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), -datasetFrequency);
1153                    nominalInstanceCal = (Calendar) initInstance.clone();
1154                    instCount[0]--;
1155                    nominalInstanceCal.add(dsTimeUnit.getCalendarUnit(), instCount[0] * datasetFrequency);
1156                    // DateUtils.moveToEnd(nominalInstanceCal, getDSEndOfFlag());
1157                }
1158            }
1159            finally {
1160                if (uriContext != null) {
1161                    uriContext.destroy();
1162                }
1163            }
1164            if (!resolved) {
1165                // return unchanged latest function with variable 'is_resolved'
1166                // to 'false'
1167                eval.setVariable("is_resolved", Boolean.FALSE);
1168                if (startOffset == endOffset) {
1169                    retVal = "${coord:latest(" + startOffset + ")}";
1170                }
1171                else {
1172                    retVal = "${coord:latestRange(" + startOffset + "," + endOffset + ")}";
1173                }
1174            }
1175            else {
1176                eval.setVariable("is_resolved", Boolean.TRUE);
1177            }
1178        }
1179        else {// No feasible nominal time
1180            eval.setVariable("is_resolved", Boolean.FALSE);
1181        }
1182        return retVal;
1183    }
1184
1185    /**
1186     * @param tm
1187     * @return a new Evaluator to be used for URI-template evaluation
1188     */
1189    private static ELEvaluator getUriEvaluator(Calendar tm) {
1190        tm.setTimeZone(DateUtils.getOozieProcessingTimeZone());
1191        ELEvaluator retEval = new ELEvaluator();
1192        retEval.setVariable("YEAR", tm.get(Calendar.YEAR));
1193        retEval.setVariable("MONTH", (tm.get(Calendar.MONTH) + 1) < 10 ? "0" + (tm.get(Calendar.MONTH) + 1) : (tm
1194                .get(Calendar.MONTH) + 1));
1195        retEval.setVariable("DAY", tm.get(Calendar.DAY_OF_MONTH) < 10 ? "0" + tm.get(Calendar.DAY_OF_MONTH) : tm
1196                .get(Calendar.DAY_OF_MONTH));
1197        retEval.setVariable("HOUR", tm.get(Calendar.HOUR_OF_DAY) < 10 ? "0" + tm.get(Calendar.HOUR_OF_DAY) : tm
1198                .get(Calendar.HOUR_OF_DAY));
1199        retEval.setVariable("MINUTE", tm.get(Calendar.MINUTE) < 10 ? "0" + tm.get(Calendar.MINUTE) : tm
1200                .get(Calendar.MINUTE));
1201        return retEval;
1202    }
1203
1204    /**
1205     * @return whether a data set is SYNCH or ASYNC
1206     */
1207    private static boolean isSyncDataSet() {
1208        ELEvaluator eval = ELEvaluator.getCurrent();
1209        SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET);
1210        if (ds == null) {
1211            throw new RuntimeException("Associated Dataset should be defined with key " + DATASET);
1212        }
1213        return ds.getType().equalsIgnoreCase("SYNC");
1214    }
1215
1216    /**
1217     * Check whether a function should be resolved.
1218     *
1219     * @param functionName
1220     * @param n
1221     * @return null if the functionName needs to be resolved otherwise return the calling function unresolved.
1222     */
1223    private static String checkIfResolved(String functionName, String n) {
1224        ELEvaluator eval = ELEvaluator.getCurrent();
1225        String replace = (String) eval.getVariable("resolve_" + functionName);
1226        if (replace == null || (replace != null && replace.equalsIgnoreCase("false"))) { // Don't
1227            // resolve
1228            // return "${coord:" + functionName + "(" + n +")}"; //Unresolved
1229            eval.setVariable(".wrap", "true");
1230            return "coord:" + functionName + "(" + n + ")"; // Unresolved
1231        }
1232        return null; // Resolved it
1233    }
1234
1235    private static String echoUnResolved(String functionName, String n) {
1236        return echoUnResolvedPre(functionName, n, "coord:");
1237    }
1238
1239    private static String echoUnResolvedPre(String functionName, String n, String prefix) {
1240        ELEvaluator eval = ELEvaluator.getCurrent();
1241        eval.setVariable(".wrap", "true");
1242        return prefix + functionName + "(" + n + ")"; // Unresolved
1243    }
1244
1245    /**
1246     * @return the initial instance of a DataSet in DATE
1247     */
1248    private static Date getInitialInstance() {
1249        ELEvaluator eval = ELEvaluator.getCurrent();
1250        return getInitialInstance(eval);
1251    }
1252
1253    /**
1254     * @return the initial instance of a DataSet in DATE
1255     */
1256    private static Date getInitialInstance(ELEvaluator eval) {
1257        return getInitialInstanceCal(eval).getTime();
1258        // return ds.getInitInstance();
1259    }
1260
1261    /**
1262     * @return the initial instance of a DataSet in Calendar
1263     */
1264    private static Calendar getInitialInstanceCal() {
1265        ELEvaluator eval = ELEvaluator.getCurrent();
1266        return getInitialInstanceCal(eval);
1267    }
1268
1269    /**
1270     * @return the initial instance of a DataSet in Calendar
1271     */
1272    private static Calendar getInitialInstanceCal(ELEvaluator eval) {
1273        SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET);
1274        if (ds == null) {
1275            throw new RuntimeException("Associated Dataset should be defined with key " + DATASET);
1276        }
1277        Calendar effInitTS = new GregorianCalendar(ds.getTimeZone());
1278        effInitTS.setTime(ds.getInitInstance());
1279        // To adjust EOD/EOM
1280        DateUtils.moveToEnd(effInitTS, getDSEndOfFlag(eval));
1281        return effInitTS;
1282        // return ds.getInitInstance();
1283    }
1284
1285    /**
1286     * @return Nominal or action creation Time when all the dependencies of an application instance are met.
1287     */
1288    private static Date getActionCreationtime() {
1289        ELEvaluator eval = ELEvaluator.getCurrent();
1290        return getActionCreationtime(eval);
1291    }
1292
1293    /**
1294     * @return Nominal or action creation Time when all the dependencies of an application instance are met.
1295     */
1296    private static Date getActionCreationtime(ELEvaluator eval) {
1297        SyncCoordAction coordAction = (SyncCoordAction) eval.getVariable(COORD_ACTION);
1298        if (coordAction == null) {
1299            throw new RuntimeException("Associated Application instance should be defined with key " + COORD_ACTION);
1300        }
1301        return coordAction.getNominalTime();
1302    }
1303
1304    /**
1305     * @return Actual Time when all the dependencies of an application instance are met.
1306     */
1307    private static Date getActualTime() {
1308        ELEvaluator eval = ELEvaluator.getCurrent();
1309        SyncCoordAction coordAction = (SyncCoordAction) eval.getVariable(COORD_ACTION);
1310        if (coordAction == null) {
1311            throw new RuntimeException("Associated Application instance should be defined with key " + COORD_ACTION);
1312        }
1313        return coordAction.getActualTime();
1314    }
1315
1316    /**
1317     * @return TimeZone for the application or job.
1318     */
1319    private static TimeZone getJobTZ() {
1320        ELEvaluator eval = ELEvaluator.getCurrent();
1321        SyncCoordAction coordAction = (SyncCoordAction) eval.getVariable(COORD_ACTION);
1322        if (coordAction == null) {
1323            throw new RuntimeException("Associated Application instance should be defined with key " + COORD_ACTION);
1324        }
1325        return coordAction.getTimeZone();
1326    }
1327
1328    /**
1329     * Find the current instance based on effectiveTime (i.e Action_Creation_Time or Action_Start_Time)
1330     *
1331     * @return current instance i.e. current(0) returns null if effectiveTime is earlier than Initial Instance time of
1332     *         the dataset.
1333     */
1334    public static Calendar getCurrentInstance(Date effectiveTime, int instanceCount[]) {
1335        ELEvaluator eval = ELEvaluator.getCurrent();
1336        return getCurrentInstance(effectiveTime, instanceCount, eval);
1337    }
1338
1339    /**
1340     * Find the current instance based on effectiveTime (i.e Action_Creation_Time or Action_Start_Time)
1341     *
1342     * @return current instance i.e. current(0) returns null if effectiveTime is earlier than Initial Instance time of
1343     *         the dataset.
1344     */
1345    private static Calendar getCurrentInstance(Date effectiveTime, int instanceCount[], ELEvaluator eval) {
1346        Date datasetInitialInstance = getInitialInstance(eval);
1347        TimeUnit dsTimeUnit = getDSTimeUnit(eval);
1348        TimeZone dsTZ = getDatasetTZ(eval);
1349        int dsFreq = getDSFrequency(eval);
1350        // Convert Date to Calendar for corresponding TZ
1351        Calendar current = Calendar.getInstance(dsTZ);
1352        current.setTime(datasetInitialInstance);
1353
1354        Calendar calEffectiveTime = new GregorianCalendar(dsTZ);
1355        calEffectiveTime.setTime(effectiveTime);
1356        if (instanceCount == null) {    // caller doesn't care about this value
1357            instanceCount = new int[1];
1358        }
1359        instanceCount[0] = 0;
1360        if (current.compareTo(calEffectiveTime) > 0) {
1361            return null;
1362        }
1363
1364        switch(dsTimeUnit) {
1365            case MINUTE:
1366                instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / MINUTE_MSEC);
1367                break;
1368            case HOUR:
1369                instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / HOUR_MSEC);
1370                break;
1371            case DAY:
1372            case END_OF_DAY:
1373                instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / DAY_MSEC);
1374                break;
1375            case MONTH:
1376            case END_OF_MONTH:
1377                instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / MONTH_MSEC);
1378                break;
1379            case YEAR:
1380                instanceCount[0] = (int) ((effectiveTime.getTime() - datasetInitialInstance.getTime()) / YEAR_MSEC);
1381                break;
1382            default:
1383                throw new IllegalArgumentException("Unhandled dataset time unit " + dsTimeUnit);
1384        }
1385
1386        if (instanceCount[0] > 2) {
1387            instanceCount[0] = (instanceCount[0] / dsFreq);
1388            current.add(dsTimeUnit.getCalendarUnit(), instanceCount[0] * dsFreq);
1389        } else {
1390            instanceCount[0] = 0;
1391        }
1392        while (!current.getTime().after(effectiveTime)) {
1393            current.add(dsTimeUnit.getCalendarUnit(), dsFreq);
1394            instanceCount[0]++;
1395        }
1396        current.add(dsTimeUnit.getCalendarUnit(), -dsFreq);
1397        instanceCount[0]--;
1398        return current;
1399    }
1400
1401    /**
1402     * Find the current instance based on effectiveTime (i.e Action_Creation_Time or Action_Start_Time)
1403     *
1404     * @return current instance i.e. current(0) returns null if effectiveTime is earlier than Initial Instance time of
1405     *         the dataset.
1406     */
1407    private static Calendar getCurrentInstance_old(Date effectiveTime, int instanceCount[], ELEvaluator eval) {
1408        Date datasetInitialInstance = getInitialInstance(eval);
1409        TimeUnit dsTimeUnit = getDSTimeUnit(eval);
1410        TimeZone dsTZ = getDatasetTZ(eval);
1411        int dsFreq = getDSFrequency(eval);
1412        // Convert Date to Calendar for corresponding TZ
1413        Calendar current = Calendar.getInstance();
1414        current.setTime(datasetInitialInstance);
1415        current.setTimeZone(dsTZ);
1416
1417        Calendar calEffectiveTime = Calendar.getInstance();
1418        calEffectiveTime.setTime(effectiveTime);
1419        calEffectiveTime.setTimeZone(dsTZ);
1420        if (instanceCount == null) {    // caller doesn't care about this value
1421            instanceCount = new int[1];
1422        }
1423        instanceCount[0] = 0;
1424        if (current.compareTo(calEffectiveTime) > 0) {
1425            return null;
1426        }
1427        Calendar origCurrent = (Calendar) current.clone();
1428        while (current.compareTo(calEffectiveTime) <= 0) {
1429            current = (Calendar) origCurrent.clone();
1430            instanceCount[0]++;
1431            current.add(dsTimeUnit.getCalendarUnit(), instanceCount[0] * dsFreq);
1432        }
1433        instanceCount[0]--;
1434
1435        current = (Calendar) origCurrent.clone();
1436        current.add(dsTimeUnit.getCalendarUnit(), instanceCount[0] * dsFreq);
1437        return current;
1438    }
1439
1440    public static Calendar getEffectiveNominalTime() {
1441        Date datasetInitialInstance = getInitialInstance();
1442        TimeZone dsTZ = getDatasetTZ();
1443        // Convert Date to Calendar for corresponding TZ
1444        Calendar current = Calendar.getInstance();
1445        current.setTime(datasetInitialInstance);
1446        current.setTimeZone(dsTZ);
1447
1448        Calendar calEffectiveTime = Calendar.getInstance();
1449        calEffectiveTime.setTime(getActionCreationtime());
1450        calEffectiveTime.setTimeZone(dsTZ);
1451        if (current.compareTo(calEffectiveTime) > 0) {
1452            // Nominal Time < initial Instance
1453            // TODO: getClass() call doesn't work from static method.
1454            // XLog.getLog("CoordELFunction.class").warn("ACTION CREATED BEFORE INITIAL INSTACE "+
1455            // current.getTime());
1456            return null;
1457        }
1458        return calEffectiveTime;
1459    }
1460
1461    /**
1462     * @return dataset frequency in minutes
1463     */
1464    private static int getDSFrequency() {
1465        ELEvaluator eval = ELEvaluator.getCurrent();
1466        return getDSFrequency(eval);
1467    }
1468
1469    /**
1470     * @return dataset frequency in minutes
1471     */
1472    private static int getDSFrequency(ELEvaluator eval) {
1473        SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET);
1474        if (ds == null) {
1475            throw new RuntimeException("Associated Dataset should be defined with key " + DATASET);
1476        }
1477        return ds.getFrequency();
1478    }
1479
1480    /**
1481     * @return dataset TimeUnit
1482     */
1483    private static TimeUnit getDSTimeUnit() {
1484        ELEvaluator eval = ELEvaluator.getCurrent();
1485        return getDSTimeUnit(eval);
1486    }
1487
1488    /**
1489     * @return dataset TimeUnit
1490     */
1491    public static TimeUnit getDSTimeUnit(ELEvaluator eval) {
1492        SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET);
1493        if (ds == null) {
1494            throw new RuntimeException("Associated Dataset should be defined with key " + DATASET);
1495        }
1496        return ds.getTimeUnit();
1497    }
1498
1499    /**
1500     * @return dataset TimeZone
1501     */
1502    public static TimeZone getDatasetTZ() {
1503        ELEvaluator eval = ELEvaluator.getCurrent();
1504        return getDatasetTZ(eval);
1505    }
1506
1507    /**
1508     * @return dataset TimeZone
1509     */
1510    private static TimeZone getDatasetTZ(ELEvaluator eval) {
1511        SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET);
1512        if (ds == null) {
1513            throw new RuntimeException("Associated Dataset should be defined with key " + DATASET);
1514        }
1515        return ds.getTimeZone();
1516    }
1517
1518    /**
1519     * @return dataset TimeUnit
1520     */
1521    private static TimeUnit getDSEndOfFlag() {
1522        ELEvaluator eval = ELEvaluator.getCurrent();
1523        return getDSEndOfFlag(eval);
1524    }
1525
1526    /**
1527     * @return dataset TimeUnit
1528     */
1529    private static TimeUnit getDSEndOfFlag(ELEvaluator eval) {
1530        SyncCoordDataset ds = (SyncCoordDataset) eval.getVariable(DATASET);
1531        if (ds == null) {
1532            throw new RuntimeException("Associated Dataset should be defined with key " + DATASET);
1533        }
1534        return ds.getEndOfDuration();// == null ? "": ds.getEndOfDuration();
1535    }
1536
1537    /**
1538     * Return a job configuration property for the coordinator.
1539     *
1540     * @param property property name.
1541     * @return the value of the property, <code>null</code> if the property is undefined.
1542     */
1543    public static String coord_conf(String property) {
1544        ELEvaluator eval = ELEvaluator.getCurrent();
1545        return (String) eval.getVariable(property);
1546    }
1547
1548    /**
1549     * Return the user that submitted the coordinator job.
1550     *
1551     * @return the user that submitted the coordinator job.
1552     */
1553    public static String coord_user() {
1554        ELEvaluator eval = ELEvaluator.getCurrent();
1555        return (String) eval.getVariable(OozieClient.USER_NAME);
1556    }
1557
1558    /**
1559     * Takes two offset times and returns a list of multiples of the frequency offset from the effective nominal time that occur
1560     * between them.  The caller should make sure that startCal is earlier than endCal.
1561     * <p>
1562     * As a simple example, assume its the same day: startCal is 1:00, endCal is 2:00, frequency is 20min, and effective nominal
1563     * time is 1:20 -- then this method would return a list containing: -20, 0, 20, 40, 60
1564     *
1565     * @param startCal The earlier offset time
1566     * @param endCal The later offset time
1567     * @param eval The ELEvaluator to use; cannot be null
1568     * @return A list of multiple of the frequency offset from the effective nominal time that occur between the startCal and endCal
1569     */
1570    public static List<Integer> expandOffsetTimes(Calendar startCal, Calendar endCal, ELEvaluator eval) {
1571        List<Integer> expandedFreqs = new ArrayList<Integer>();
1572        // Use eval because the "current" eval isn't set
1573        int freq = getDSFrequency(eval);
1574        TimeUnit freqUnit = getDSTimeUnit(eval);
1575        Calendar cal = getCurrentInstance(getActionCreationtime(eval), null, eval);
1576        int totalFreq = 0;
1577        if (startCal.before(cal)) {
1578            while (cal.after(startCal)) {
1579                cal.add(freqUnit.getCalendarUnit(), -freq);
1580                totalFreq += -freq;
1581            }
1582            if (cal.before(startCal)) {
1583                cal.add(freqUnit.getCalendarUnit(), freq);
1584                totalFreq += freq;
1585            }
1586        }
1587        else if (startCal.after(cal)) {
1588            while (cal.before(startCal)) {
1589                cal.add(freqUnit.getCalendarUnit(), freq);
1590                totalFreq += freq;
1591            }
1592        }
1593        // At this point, cal is the smallest multiple of the dataset frequency that is >= to the startCal and offset from the
1594        // effective nominal time.  Now we can find all of the instances that occur between startCal and endCal, inclusive.
1595        while (cal.before(endCal) || cal.equals(endCal)) {
1596            expandedFreqs.add(totalFreq);
1597            cal.add(freqUnit.getCalendarUnit(), freq);
1598            totalFreq += freq;
1599        }
1600        return expandedFreqs;
1601    }
1602
1603    /**
1604     * Resolve the offset time from the effective nominal time
1605     *
1606     * @param n offset amount (integer)
1607     * @param timeUnit TimeUnit for offset n ("MINUTE", "HOUR", "DAY", "MONTH", "YEAR")
1608     * @param eval The ELEvaluator to use; or null to use the "current" eval
1609     * @return A Calendar of the offset time
1610     */
1611    public static Calendar resolveOffsetRawTime(int n, TimeUnit timeUnit, ELEvaluator eval) {
1612        // Use eval if given (for when the "current" eval isn't set)
1613        Calendar cal;
1614        if (eval == null) {
1615            cal = getCurrentInstance(getActionCreationtime(), null);
1616        }
1617        else {
1618            cal = getCurrentInstance(getActionCreationtime(eval), null, eval);
1619        }
1620        if (cal == null) {
1621            XLog.getLog(CoordELFunctions.class).warn("If the initial instance of the dataset is later than the nominal time, an"
1622                    + " empty string is returned. This means that no data is available at the offset instance specified by the user"
1623                    + " and the user could try modifying his or her initial-instance to an earlier time.");
1624            return null;
1625        }
1626        cal.add(timeUnit.getCalendarUnit(), n);
1627        return cal;
1628    }
1629}