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.service;
020
021import java.io.IOException;
022import java.text.MessageFormat;
023import java.util.Collection;
024import java.util.List;
025import java.util.Properties;
026import java.util.concurrent.Callable;
027
028import javax.persistence.EntityManager;
029import javax.persistence.EntityManagerFactory;
030import javax.persistence.EntityTransaction;
031import javax.persistence.NoResultException;
032import javax.persistence.Persistence;
033import javax.persistence.Query;
034
035import org.apache.commons.collections.CollectionUtils;
036import org.apache.commons.dbcp.BasicDataSource;
037import org.apache.commons.lang.StringUtils;
038import org.apache.hadoop.conf.Configuration;
039import org.apache.oozie.BundleActionBean;
040import org.apache.oozie.BundleJobBean;
041import org.apache.oozie.CoordinatorActionBean;
042import org.apache.oozie.CoordinatorJobBean;
043import org.apache.oozie.ErrorCode;
044import org.apache.oozie.FaultInjection;
045import org.apache.oozie.SLAEventBean;
046import org.apache.oozie.WorkflowActionBean;
047import org.apache.oozie.WorkflowJobBean;
048import org.apache.oozie.client.rest.JsonBean;
049import org.apache.oozie.client.rest.JsonSLAEvent;
050import org.apache.oozie.command.SkipCommitFaultInjection;
051import org.apache.oozie.compression.CodecFactory;
052import org.apache.oozie.executor.jpa.JPAExecutor;
053import org.apache.oozie.executor.jpa.JPAExecutorException;
054import org.apache.oozie.sla.SLARegistrationBean;
055import org.apache.oozie.sla.SLASummaryBean;
056import org.apache.oozie.util.IOUtils;
057import org.apache.oozie.util.Instrumentable;
058import org.apache.oozie.util.Instrumentation;
059import org.apache.oozie.util.XLog;
060import org.apache.oozie.util.db.OperationRetryHandler;
061import org.apache.oozie.util.db.PersistenceExceptionSubclassFilterRetryPredicate;
062import org.apache.openjpa.lib.jdbc.DecoratingDataSource;
063import org.apache.openjpa.persistence.InvalidStateException;
064import org.apache.openjpa.persistence.OpenJPAEntityManagerFactorySPI;
065
066/**
067 * Service that manages JPA and executes {@link JPAExecutor}.
068 */
069@SuppressWarnings("deprecation")
070public class JPAService implements Service, Instrumentable {
071    private static final String INSTRUMENTATION_GROUP_JPA = "jpa";
072
073    public static final long DEFAULT_INITIAL_WAIT_TIME = 100;
074    public static final long DEFAULT_MAX_WAIT_TIME = 30_000;
075    public static final int DEFAULT_MAX_RETRY_COUNT = 1;
076
077    public static final String CONF_DB_SCHEMA = "oozie.db.schema.name";
078
079    public static final String CONF_PREFIX = Service.CONF_PREFIX + "JPAService.";
080    public static final String CONF_URL = CONF_PREFIX + "jdbc.url";
081    public static final String CONF_DRIVER = CONF_PREFIX + "jdbc.driver";
082    public static final String CONF_USERNAME = CONF_PREFIX + "jdbc.username";
083    public static final String CONF_PASSWORD = CONF_PREFIX + "jdbc.password";
084    public static final String CONF_CONN_DATA_SOURCE = CONF_PREFIX + "connection.data.source";
085    public static final String CONF_CONN_PROPERTIES = CONF_PREFIX + "connection.properties";
086    public static final String CONF_MAX_ACTIVE_CONN = CONF_PREFIX + "pool.max.active.conn";
087    public static final String CONF_CREATE_DB_SCHEMA = CONF_PREFIX + "create.db.schema";
088    public static final String CONF_VALIDATE_DB_CONN = CONF_PREFIX + "validate.db.connection";
089    public static final String CONF_VALIDATE_DB_CONN_EVICTION_INTERVAL = CONF_PREFIX + "validate.db.connection.eviction.interval";
090    public static final String CONF_VALIDATE_DB_CONN_EVICTION_NUM = CONF_PREFIX + "validate.db.connection.eviction.num";
091    public static final String CONF_OPENJPA_BROKER_IMPL = CONF_PREFIX + "openjpa.BrokerImpl";
092    public static final String INITIAL_WAIT_TIME = CONF_PREFIX + "retry.initial-wait-time.ms";
093    public static final String MAX_WAIT_TIME = CONF_PREFIX + "maximum-wait-time.ms";
094    public static final String MAX_RETRY_COUNT = CONF_PREFIX + "retry.max-retries";
095    public static final String SKIP_COMMIT_FAULT_INJECTION_CLASS = SkipCommitFaultInjection.class.getName();
096
097    private EntityManagerFactory factory;
098    private Instrumentation instr;
099
100    private static XLog LOG;
101    private OperationRetryHandler retryHandler;
102
103    /**
104     * Return the public interface of the service.
105     *
106     * @return {@link JPAService}.
107     */
108    public Class<? extends Service> getInterface() {
109        return JPAService.class;
110    }
111
112    @Override
113    public void instrument(final Instrumentation instr) {
114        this.instr = instr;
115
116        final BasicDataSource dataSource = getBasicDataSource();
117        if (dataSource != null) {
118            instr.addSampler("jdbc", "connections.active", 60, 1, new Instrumentation.Variable<Long>() {
119                @Override
120                public Long getValue() {
121                    return (long) dataSource.getNumActive();
122                }
123            });
124            instr.addSampler("jdbc", "connections.idle", 60, 1, new Instrumentation.Variable<Long>() {
125                @Override
126                public Long getValue() {
127                    return (long) dataSource.getNumIdle();
128                }
129            });
130        }
131    }
132
133    private BasicDataSource getBasicDataSource() {
134        // Get the BasicDataSource object; it could be wrapped in a DecoratingDataSource
135        // It might also not be a BasicDataSource if the user configured something different
136        BasicDataSource basicDataSource = null;
137        final OpenJPAEntityManagerFactorySPI spi = (OpenJPAEntityManagerFactorySPI) factory;
138        final Object connectionFactory = spi.getConfiguration().getConnectionFactory();
139        if (connectionFactory instanceof DecoratingDataSource) {
140            final DecoratingDataSource decoratingDataSource = (DecoratingDataSource) connectionFactory;
141            basicDataSource = (BasicDataSource) decoratingDataSource.getInnermostDelegate();
142        } else if (connectionFactory instanceof BasicDataSource) {
143            basicDataSource = (BasicDataSource) connectionFactory;
144        }
145        return basicDataSource;
146    }
147
148    /**
149     * Initializes the {@link JPAService}.
150     *
151     * @param services services instance.
152     */
153    public void init(final Services services) throws ServiceException {
154        LOG = XLog.getLog(JPAService.class);
155        final Configuration conf = services.getConf();
156        final String dbSchema = ConfigurationService.get(conf, CONF_DB_SCHEMA);
157        String url = ConfigurationService.get(conf, CONF_URL);
158        final String driver = ConfigurationService.get(conf, CONF_DRIVER);
159        final String user = ConfigurationService.get(conf, CONF_USERNAME);
160        final String password = ConfigurationService.getPassword(conf, CONF_PASSWORD).trim();
161        final String maxConn = ConfigurationService.get(conf, CONF_MAX_ACTIVE_CONN).trim();
162        final String dataSource = ConfigurationService.get(conf, CONF_CONN_DATA_SOURCE);
163        final String connPropsConfig = ConfigurationService.get(conf, CONF_CONN_PROPERTIES);
164        final String brokerImplConfig = ConfigurationService.get(conf, CONF_OPENJPA_BROKER_IMPL);
165        final boolean autoSchemaCreation = ConfigurationService.getBoolean(conf, CONF_CREATE_DB_SCHEMA);
166        final boolean validateDbConn = ConfigurationService.getBoolean(conf, CONF_VALIDATE_DB_CONN);
167        final String evictionInterval = ConfigurationService.get(conf, CONF_VALIDATE_DB_CONN_EVICTION_INTERVAL).trim();
168        final String evictionNum = ConfigurationService.get(conf, CONF_VALIDATE_DB_CONN_EVICTION_NUM).trim();
169
170        if (!url.startsWith("jdbc:")) {
171            throw new ServiceException(ErrorCode.E0608, url, "invalid JDBC URL, must start with 'jdbc:'");
172        }
173        String dbType = url.substring("jdbc:".length());
174        if (dbType.indexOf(":") <= 0) {
175            throw new ServiceException(ErrorCode.E0608, url, "invalid JDBC URL, missing vendor 'jdbc:[VENDOR]:...'");
176        }
177        dbType = dbType.substring(0, dbType.indexOf(":"));
178
179        final String persistentUnit = "oozie-" + dbType;
180
181        // Checking existince of ORM file for DB type
182        final String ormFile = "META-INF/" + persistentUnit + "-orm.xml";
183        try {
184            IOUtils.getResourceAsStream(ormFile, -1);
185        }
186        catch (final IOException ex) {
187            throw new ServiceException(ErrorCode.E0609, dbType, ormFile);
188        }
189
190        // support for mysql replication urls "jdbc:mysql:replication://master:port,slave:port[,slave:port]/db"
191        if (url.startsWith("jdbc:mysql:replication")) {
192            url = "\"".concat(url).concat("\"");
193            LOG.info("A jdbc replication url is provided. Url: [{0}]", url);
194        }
195
196        String connProps = "DriverClassName={0},Url={1},Username={2},Password={3},MaxActive={4}";
197        connProps = MessageFormat.format(connProps, driver, url, user, password, maxConn);
198        final Properties props = new Properties();
199        if (autoSchemaCreation) {
200            connProps += ",TestOnBorrow=false,TestOnReturn=false,TestWhileIdle=false";
201            props.setProperty("openjpa.jdbc.SynchronizeMappings", "buildSchema(ForeignKeys=true)");
202        }
203        else if (validateDbConn) {
204            // validation can be done only if the schema already exist, else a
205            // connection cannot be obtained to create the schema.
206            final String interval = "timeBetweenEvictionRunsMillis=" + evictionInterval;
207            final String num = "numTestsPerEvictionRun=" + evictionNum;
208            connProps += ",TestOnBorrow=true,TestOnReturn=true,TestWhileIdle=true," + interval + "," + num;
209            connProps += ",ValidationQuery=select count(*) from VALIDATE_CONN";
210            connProps = MessageFormat.format(connProps, dbSchema);
211        }
212        else {
213            connProps += ",TestOnBorrow=false,TestOnReturn=false,TestWhileIdle=false";
214        }
215        if (connPropsConfig != null) {
216            connProps += "," + connPropsConfig;
217        }
218        props.setProperty("openjpa.ConnectionProperties", connProps);
219
220        props.setProperty("openjpa.ConnectionDriverName", dataSource);
221        if (!StringUtils.isEmpty(brokerImplConfig)) {
222            props.setProperty("openjpa.BrokerImpl", brokerImplConfig);
223            LOG.info("Setting openjpa.BrokerImpl to {0}", brokerImplConfig);
224        }
225
226        initRetryHandler();
227
228        factory = Persistence.createEntityManagerFactory(persistentUnit, props);
229
230        final EntityManager entityManager = getEntityManager();
231        findRetrying(entityManager, WorkflowActionBean.class, 1);
232        findRetrying(entityManager, WorkflowJobBean.class, 1);
233        findRetrying(entityManager, CoordinatorActionBean.class, 1);
234        findRetrying(entityManager, CoordinatorJobBean.class, 1);
235        findRetrying(entityManager, SLAEventBean.class, 1);
236        findRetrying(entityManager, JsonSLAEvent.class, 1);
237        findRetrying(entityManager, BundleActionBean.class, 1);
238        findRetrying(entityManager, BundleJobBean.class, 1);
239        findRetrying(entityManager, SLARegistrationBean.class, 1);
240        findRetrying(entityManager, SLASummaryBean.class, 1);
241
242        LOG.info(XLog.STD, "All entities initialized");
243        // need to use a pseudo no-op transaction so all entities, datasource
244        // and connection pool are initialized one time only
245        entityManager.getTransaction().begin();
246        final OpenJPAEntityManagerFactorySPI spi = (OpenJPAEntityManagerFactorySPI) factory;
247        // Mask the password with '***'
248        final String logMsg = spi.getConfiguration().getConnectionProperties().replaceAll("Password=.*?,", "Password=***,");
249        LOG.info("JPA configuration: {0}", logMsg);
250        entityManager.getTransaction().commit();
251        entityManager.close();
252        try {
253            CodecFactory.initialize(conf);
254        }
255        catch (final Exception ex) {
256            throw new ServiceException(ErrorCode.E0100, getClass().getName(), ex);
257        }
258
259    }
260
261    private void initRetryHandler() {
262        final long initialWaitTime = ConfigurationService.getInt(INITIAL_WAIT_TIME, (int) DEFAULT_INITIAL_WAIT_TIME);
263        final long maxWaitTime = ConfigurationService.getInt(MAX_WAIT_TIME, (int) DEFAULT_MAX_WAIT_TIME);
264        final int maxRetryCount = ConfigurationService.getInt(MAX_RETRY_COUNT, DEFAULT_MAX_RETRY_COUNT);
265
266        LOG.info(XLog.STD, "Failing database operations will be retried {0} times, with an initial sleep time of {1} ms,"
267                + "max sleep time {2} ms", maxRetryCount, initialWaitTime, maxWaitTime);
268        retryHandler = new OperationRetryHandler(maxRetryCount,
269                initialWaitTime,
270                maxWaitTime,
271                new PersistenceExceptionSubclassFilterRetryPredicate());
272    }
273
274    private void findRetrying(final EntityManager entityManager, final Class entityClass, final int primaryKey)
275            throws ServiceException {
276        try {
277            retryHandler.executeWithRetry(new Callable<Void>() {
278                @Override
279                public Void call() throws Exception {
280                    if (!entityManager.getTransaction().isActive()) {
281                        entityManager.getTransaction().begin();
282                    }
283
284                    entityManager.find(entityClass, primaryKey);
285
286                    if (entityManager.getTransaction().isActive()) {
287                        entityManager.getTransaction().commit();
288                    }
289                    return null;
290                }
291            });
292        }
293        catch (final Exception e) {
294            throw new ServiceException(ErrorCode.E0603, e);
295        }
296    }
297
298    /**
299     * Destroy the JPAService
300     */
301    public void destroy() {
302        if (factory != null && factory.isOpen()) {
303            try {
304                factory.close();
305            }
306            catch (final InvalidStateException ise) {
307                LOG.warn("Cannot close EntityManagerFactory. [ise.message={0}]", ise.getMessage());
308            }
309        }
310    }
311
312    /**
313     * Execute a {@link JPAExecutor}.
314     *
315     * @param executor JPAExecutor to execute.
316     * @return return value of the JPAExecutor.
317     * @throws JPAExecutorException thrown if an jpa executor failed
318     */
319    public <T> T execute(final JPAExecutor<T> executor) throws JPAExecutorException {
320        final EntityManager em = getEntityManager();
321        final Instrumentation.Cron cron = new Instrumentation.Cron();
322        try {
323            LOG.trace("Executing JPAExecutor [{0}]", executor.getName());
324            if (instr != null) {
325                instr.incr(INSTRUMENTATION_GROUP_JPA, executor.getName(), 1);
326            }
327            cron.start();
328
329            return retryHandler.executeWithRetry(new Callable<T>() {
330                @Override
331                public T call() throws Exception {
332                    if (!em.getTransaction().isActive()) {
333                        em.getTransaction().begin();
334                    }
335
336                    final T t = executor.execute(em);
337
338                    checkAndCommit(em.getTransaction());
339
340                    return t;
341                }
342            });
343        }
344        catch (final Exception e) {
345            throw getTargetException(e);
346        }
347        finally {
348            cron.stop();
349            if (instr != null) {
350                instr.addCron(INSTRUMENTATION_GROUP_JPA, executor.getName(), cron);
351            }
352            try {
353                if (em.getTransaction().isActive()) {
354                    LOG.warn("JPAExecutor [{0}] ended with an active transaction, rolling back", executor.getName());
355                    em.getTransaction().rollback();
356                }
357            }
358            catch (final Exception ex) {
359                LOG.warn("Could not check/rollback transaction after JPAExecutor [{0}], {1}", executor.getName(), ex
360                        .getMessage(), ex);
361            }
362            try {
363                if (em.isOpen()) {
364                    em.close();
365                }
366                else {
367                    LOG.warn("JPAExecutor [{0}] closed the EntityManager, it should not!", executor.getName());
368                }
369            }
370            catch (final Exception ex) {
371                LOG.warn("Could not close EntityManager after JPAExecutor [{0}], {1}", executor.getName(), ex
372                        .getMessage(), ex);
373            }
374        }
375    }
376
377    private void checkAndCommit(final EntityTransaction tx) throws JPAExecutorException {
378        if (tx.isActive()) {
379            if (FaultInjection.isActive(SKIP_COMMIT_FAULT_INJECTION_CLASS)) {
380                throw new JPAExecutorException(ErrorCode.E0603, "Skipping Commit for Failover Testing");
381            }
382
383            tx.commit();
384        }
385    }
386
387    /**
388     * Execute an UPDATE query
389     * @param namedQueryName the name of query to be executed
390     * @param query query instance to be executed
391     * @param em Entity Manager
392     * @return Integer that query returns, which corresponds to the number of rows updated
393     * @throws JPAExecutorException
394     */
395    public int executeUpdate(final String namedQueryName, final Query query, final EntityManager em) throws JPAExecutorException {
396        final Instrumentation.Cron cron = new Instrumentation.Cron();
397        try {
398
399            LOG.trace("Executing Update/Delete Query [{0}]", namedQueryName);
400            if (instr != null) {
401                instr.incr(INSTRUMENTATION_GROUP_JPA, namedQueryName, 1);
402            }
403            cron.start();
404
405            return retryHandler.executeWithRetry(new Callable<Integer>() {
406                @Override
407                public Integer call() throws Exception {
408                    if (!em.getTransaction().isActive()) {
409                        em.getTransaction().begin();
410                    }
411                    final int ret = query.executeUpdate();
412
413                    checkAndCommit(em.getTransaction());
414
415                    return ret;
416                }
417            });
418        }
419        catch (final Exception e) {
420            throw getTargetException(e);
421        }
422        finally {
423            processFinally(em, cron, namedQueryName, true);
424        }
425    }
426
427    public static class QueryEntry<E extends Enum<E>> {
428        E namedQuery;
429        Query query;
430
431        public QueryEntry(final E namedQuery, final Query query) {
432            this.namedQuery = namedQuery;
433            this.query = query;
434        }
435
436        public Query getQuery() {
437            return this.query;
438        }
439
440        public E getQueryName() {
441            return this.namedQuery;
442        }
443    }
444
445    private void processFinally(final EntityManager em,
446                                final Instrumentation.Cron cron,
447                                final String name,
448                                final boolean checkActive) {
449        cron.stop();
450        if (instr != null) {
451            instr.addCron(INSTRUMENTATION_GROUP_JPA, name, cron);
452        }
453        if (checkActive) {
454            try {
455                if (em.getTransaction().isActive()) {
456                    LOG.warn("[{0}] ended with an active transaction, rolling back", name);
457                    em.getTransaction().rollback();
458                }
459            }
460            catch (final Exception ex) {
461                LOG.warn("Could not check/rollback transaction after [{0}], {1}", name,
462                        ex.getMessage(), ex);
463            }
464        }
465        try {
466            if (em.isOpen()) {
467                em.close();
468            }
469            else {
470                LOG.warn("[{0}] closed the EntityManager, it should not!", name);
471            }
472        }
473        catch (final Exception ex) {
474            LOG.warn("Could not close EntityManager after [{0}], {1}", name, ex.getMessage(), ex);
475        }
476    }
477
478    /**
479     * Execute multiple update/insert queries in one transaction
480     * @param insertBeans list of beans to be inserted
481     * @param updateQueryList list of update queries
482     * @param deleteBeans list of beans to be deleted
483     * @param em Entity Manager
484     * @throws JPAExecutorException
485     */
486    public void executeBatchInsertUpdateDelete(final Collection<JsonBean> insertBeans, final List<QueryEntry> updateQueryList,
487            final Collection<JsonBean> deleteBeans, final EntityManager em) throws JPAExecutorException {
488        final Instrumentation.Cron cron = new Instrumentation.Cron();
489        try {
490
491            LOG.trace("Executing Queries in Batch");
492            cron.start();
493
494            retryHandler.executeWithRetry(new Callable<Void>() {
495                @Override
496                public Void call() throws Exception {
497                    if (em.getTransaction().isActive()) {
498                        try {
499                            em.getTransaction().rollback();
500                        }
501                        catch (final Exception e) {
502                            LOG.warn("Rollback failed - ignoring");
503                        }
504                    }
505
506                    em.getTransaction().begin();
507
508                    if (CollectionUtils.isNotEmpty(updateQueryList)) {
509                        for (final QueryEntry q : updateQueryList) {
510                            if (instr != null) {
511                                instr.incr(INSTRUMENTATION_GROUP_JPA, q.getQueryName().name(), 1);
512                            }
513                            q.getQuery().executeUpdate();
514                        }
515                    }
516
517                    if (CollectionUtils.isNotEmpty(insertBeans)) {
518                        for (final JsonBean bean : insertBeans) {
519                            em.persist(bean);
520                        }
521                    }
522
523                    if (CollectionUtils.isNotEmpty(deleteBeans)) {
524                        for (final JsonBean bean : deleteBeans) {
525                            em.remove(em.merge(bean));
526                        }
527                    }
528
529                    checkAndCommit(em.getTransaction());
530
531                    return null;
532                }
533            });
534        }
535        catch (final Exception e) {
536            throw getTargetException(e);
537        }
538        finally {
539            processFinally(em, cron, "batchqueryexecutor", true);
540        }
541    }
542
543    /**
544     * Execute a SELECT query
545     * @param namedQueryName the name of query to be executed
546     * @param query query instance to be executed
547     * @param em Entity Manager
548     * @return object that matches the query
549     */
550    public Object executeGet(final String namedQueryName, final Query query, final EntityManager em) throws JPAExecutorException {
551        final Instrumentation.Cron cron = new Instrumentation.Cron();
552        try {
553            LOG.trace("Executing Select Query to Get a Single row  [{0}]", namedQueryName);
554            if (instr != null) {
555                instr.incr(INSTRUMENTATION_GROUP_JPA, namedQueryName, 1);
556            }
557
558            cron.start();
559
560            return retryHandler.executeWithRetry(new Callable<Object>() {
561                @Override
562                public Object call() throws Exception {
563                    Object obj = null;
564                    try {
565                        obj = query.getSingleResult();
566                    }
567                    catch (final NoResultException e) {
568                        LOG.info("No results found");
569                        // return null when no matched result
570                    }
571                    return obj;
572                }
573            });
574        }
575        catch (final Exception e) {
576            throw getTargetException(e);
577        }
578        finally {
579            processFinally(em, cron, namedQueryName, false);
580        }
581    }
582
583    /**
584     * Execute a SELECT query to get list of results
585     * @param namedQueryName the name of query to be executed
586     * @param query query instance to be executed
587     * @param em Entity Manager
588     * @return list containing results that match the query
589     */
590    public List<?> executeGetList(final String namedQueryName, final Query query, final EntityManager em)
591            throws JPAExecutorException {
592        final Instrumentation.Cron cron = new Instrumentation.Cron();
593        try {
594
595            LOG.trace("Executing Select Query to Get Multiple Rows [{0}]", namedQueryName);
596            if (instr != null) {
597                instr.incr(INSTRUMENTATION_GROUP_JPA, namedQueryName, 1);
598            }
599
600            cron.start();
601
602            return retryHandler.executeWithRetry(new Callable<List<?>>() {
603                @Override
604                public List<?> call() throws Exception {
605                    List<?> resultList = null;
606                    try {
607                        resultList = query.getResultList();
608                    }
609                    catch (final NoResultException e) {
610                        LOG.info("No results found");
611                        // return null when no matched result
612                    }
613                    return resultList;
614                }
615            });
616        }
617        catch (final Exception e) {
618            throw getTargetException(e);
619        }
620        finally {
621            processFinally(em, cron, namedQueryName, false);
622        }
623    }
624
625    /**
626     * Return an EntityManager. Used by the StoreService. Once the StoreService is removed this method must be removed.
627     *
628     * @return an entity manager
629     */
630    public EntityManager getEntityManager() {
631        return factory.createEntityManager();
632    }
633
634    private JPAExecutorException getTargetException(final Exception e) {
635        if (e instanceof JPAExecutorException) {
636            return (JPAExecutorException) e;
637        }
638        else {
639            return new JPAExecutorException(ErrorCode.E0603, e.getMessage());
640        }
641    }
642}