From 9f59117ec1403f4790863fc43aebb5c626554227 Mon Sep 17 00:00:00 2001 From: Jan Bauer Nielsen Date: Fri, 4 Sep 2026 07:54:13 +0200 Subject: [PATCH] sink/periodic-jobs: deliver and report one item at a time (DI-3015) Migrate the sink from MessageConsumerAdapter to SinkMessageConsumerAdapter. job-store now dispatches one JMS message per item, and the sink implements deliverItem(ConsumedMessage, ChunkItem) instead of assembling a DELIVERED chunk. Header reading, result reporting and the tracking-id scope move to the framework. The sink aggregates a whole job before delivering anything, so it overrides usesDeliveryWatermark() to false. It still reports every item individually, since the DELIVERING counters and the per-job gate are driven by those reports. The job termination item arrives as an ordinary item message carrying ChunkItem.Type.JOB_END. Drop the five second sleep that preceded finalization. It guessed at a window the chunk protocol opened by reporting the result before committing the sink's own transaction. deliverItem commits first and the framework reports afterwards, so a reported item is an item whose datablocks are durable and the termination chunk is released only against complete data. InvalidMessageException from finalization is now reported as a FAILED termination item rather than discarding the message, which completes the job and sets its fatal error flag instead of leaving it forever incomplete. --- docs/chunk-scheduling-redesign.md | 20 ++ sink/CLAUDE.md | 39 ++- sink/periodic-jobs/pom.xml | 11 + .../PeriodicJobsConfigurationBean.java | 25 +- .../PeriodicJobsFinalizerBean.java | 39 +-- .../PeriodicJobsMessageConsumer.java | 171 ++++++++---- .../pickup/PeriodicJobsFtpFinalizerBean.java | 32 ++- .../pickup/PeriodicJobsHttpFinalizerBean.java | 16 +- .../pickup/PeriodicJobsMailFinalizerBean.java | 30 +-- .../pickup/PeriodicJobsPickupFinalizer.java | 22 +- .../pickup/PeriodicJobsSFtpFinalizerBean.java | 30 +-- .../PeriodicJobsConfigurationBeanIT.java | 63 ++--- .../PeriodicJobsMessageConsumerIT.java | 57 ++-- .../PeriodicJobsMessageConsumerTest.java | 249 ++++++++++++++++++ .../pickup/PeriodicJobsFinalizerBeanIT.java | 30 ++- .../PeriodicJobsFtpFinalizerBeanIT.java | 10 +- .../PeriodicJobsHttpFinalizerBeanIT.java | 13 +- .../PeriodicJobsMailFinalizerBeanIT.java | 38 +-- .../PeriodicJobsSFtpFinalizerBeanIT.java | 13 +- 19 files changed, 627 insertions(+), 281 deletions(-) create mode 100644 sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumerTest.java diff --git a/docs/chunk-scheduling-redesign.md b/docs/chunk-scheduling-redesign.md index 5ebff895f0..8a4ee5e135 100644 --- a/docs/chunk-scheduling-redesign.md +++ b/docs/chunk-scheduling-redesign.md @@ -1723,6 +1723,26 @@ What it does **not** skip is result reporting: every item is still reported indi because the DELIVERING phase counters and the per-job gate are driven by those reports, and a job whose items are never reported never completes. +##### Job-end work runs against complete data by construction + +An aggregating sink's job-end work — the `PeriodicJobs*FinalizerBean`s, marcconv's +`ConversionFinalizer` — reads what every preceding item of the job persisted, so it is +only correct if all of those writes are committed before it starts. The per-item protocol +gives that for free, from the order in which one item is handled: `deliverItem` commits +its own transaction and returns, and `SinkMessageConsumerAdapter` reports the result only +afterwards. A reported item is therefore an item whose writes are durable, and the +termination chunk is released only once every data chunk of the job has reported (by +`waitingOn` today, by `gate_open` once the dispatch filter lands — both driven by +`chunkDeliveringDone`, which fires when a chunk's last item result is committed). + +The chunk protocol had this the other way round: each sink called +`sendResultToJobStore(result)` *before* committing its own transaction, so job-store could +see a chunk as delivered while that chunk's data was still uncommitted, and the +termination chunk could be released against an incomplete set. `periodic-jobs` covered +that window with a fixed five second sleep before finalizing, removed in DI-3015 along +with the ordering problem it guessed at. `marcconv` has the same shape and the same +argument applies to it. + ### Watermark calls (`job-store-service-connector`) An earlier version of this document specified a separate `WatermarkServiceConnector`. diff --git a/sink/CLAUDE.md b/sink/CLAUDE.md index 5ee94f2a4b..4438357313 100644 --- a/sink/CLAUDE.md +++ b/sink/CLAUDE.md @@ -26,15 +26,35 @@ mvn package -pl marc-client -P generate-source ## Architecture -All sinks follow the same structural pattern built on the `jse-artemis` framework: +All sinks share the same structure, built on the `jse-artemis` framework: 1. **`*SinkApp`** — entry point, extends `MessageConsumerApp`. Creates a `ServiceHub` and a `Supplier`, then calls `go(serviceHub, messageConsumer)`. Database-backed sinks also call `JPAHelper.migrate()` and `JPAHelper.makeEntityManagerFactory()` here. -2. **`*MessageConsumer`** — extends `MessageConsumerAdapter`. Implements `handleConsumedMessage(ConsumedMessage)`, which is the main processing loop. The standard flow is: - - `unmarshallPayload(consumedMessage)` → `Chunk` of type PROCESSED - - Iterate over `ChunkItem`s; handle SUCCESS / FAILURE / IGNORE status - - Build a new `Chunk` of type DELIVERED - - `sendResultToJobStore(deliveredChunk)` +2. **`*MessageConsumer`** — the consumer, extending `SinkMessageConsumerAdapter` (epic DI-2946, see `docs/chunk-scheduling-redesign.md`). job-store dispatches one JMS message per item with `payload = ITEM_PAYLOAD_TYPE` and a single `ChunkItem` body, and the sink implements one method: + + ```java + protected ItemDeliveryResult deliverItem(ConsumedMessage message, ChunkItem item) + ``` + + returning `ItemDeliveryResult.of(status, outcomeItem)`. `handleConsumedMessage` is `final` in the base class. The framework — not the sink — owns the header reads, the delivery watermark check, reporting the result to job-store, the `DBCTrackedLogContext` tracking-id scope, and the `dataio_item_delivery_count` metric. In particular, never read `JMSHeader.recordKey` or re-derive a watermark key from record content. + + Four verdicts, three of them the sink's to return: + + | Verdict | Meaning | DELIVERING counter | Returned by | + |---|---|---|---| + | `DELIVERED` | sent to the target | succeeded | the sink | + | `IGNORED` | not sent, nothing to send | ignored | the sink | + | `FAILED` | attempted, rejected in a way retrying will not fix | failed | the sink | + | `SUPERSEDED` | a newer version of the record was already delivered | ignored | the framework only | + + Rules that are easy to get wrong: + - **Throwing means "retry"** — it rolls the JMS session back and the item is redelivered until the broker gives up. A terminal failure must be returned as `FAILED`, not thrown. + - **A processing outcome passed through without being sent is `IGNORED`, not `DELIVERED` with an `IGNORE` item.** `DELIVERED` is the only verdict that advances the watermark, so using it for an unsent item both overstates the succeeded count and makes a false claim about what is at the target. + - **Commit your own writes before returning.** The framework reports the result after `deliverItem` returns, so a reported item is an item whose writes are durable — which is what lets an aggregating sink's job-end work run against complete data. + - Sinks that aggregate a whole job before delivering anything (`periodic-jobs`, `marcconv`) override `usesDeliveryWatermark()` to `false`. They still report every item individually: the phase counters and the per-job gate are driven by those reports. + - The job termination item arrives as an ordinary item message carrying `ChunkItem.Type.JOB_END`; sinks needing job-end work branch on that. + + `dlq-errorhandler` and `job-processor2` are not sinks in this sense and stay on the chunk-level `MessageConsumerAdapter` by design: they implement `handleConsumedMessage(ConsumedMessage)` themselves and report whole `Chunk`s via `sendResultToJobStore`. 3. **`SinkConfig`** — enum implementing `EnvConfig`. Each constant maps to an environment variable. Values are read at startup; default values can be provided in the constructor. @@ -44,7 +64,12 @@ Configuration that varies per job (e.g. endpoint, credentials) comes from the ** Files named `*IT.java` are integration tests run by `maven-failsafe-plugin` during `verify`. They use **Testcontainers** (PostgreSQL) via `PostgresContainerJPAUtils`. Unit tests (`*Test.java`) use Mockito and run with `maven-surefire-plugin` during `test`. -The `testutil` module provides `ObjectFactory.createConsumedMessage(Chunk)` — the standard way to build a `ConsumedMessage` in tests. +There is no shared helper for building a per-item `ConsumedMessage`: build it from a header map (`JMSHeader.payload` = `ITEM_PAYLOAD_TYPE`, plus `jobId`, `chunkId`, `itemId`, `sinkId` and, unless the sink opted out, `recordKey`) and a `JSONBContext`-marshalled `ChunkItem` body. `DummyMessageConsumerTest` and `PeriodicJobsMessageConsumerTest` are the models. Construct the consumer with `new ServiceHub.Builder().withJobStoreServiceConnector(mock).test()` — `test()` rather than `build()`, so no HTTP service is started. + +Two things worth knowing before writing such a test: + +- **A test that drives `handleConsumedMessage` and then asserts the reported result is re-testing the framework.** Header reading, the watermark comparison and result reporting all live in `SinkMessageConsumerAdapter` and are covered by `SinkMessageConsumerAdapterTest`. A sink's own surface is `deliverItem` plus its `usesDeliveryWatermark()` choice — call `deliverItem` directly. The one exception is asserting the watermark opt-out, which is observable only as the lookup being (or not being) made. +- **A unit test that constructs a consumer needs `APP_NAME`** — `UserAgent.forInternalRequests()` reads it while the consumer builds its connectors, and without it the test fails with "APP_NAME environment variable has not been set". Several modules set it only for failsafe; add the same `` block to `maven-surefire-plugin` (see `dpf` and `periodic-jobs`). ## Notable modules diff --git a/sink/periodic-jobs/pom.xml b/sink/periodic-jobs/pom.xml index e0dc207fe5..ec7e7d6106 100644 --- a/sink/periodic-jobs/pom.xml +++ b/sink/periodic-jobs/pom.xml @@ -146,6 +146,17 @@ org.codehaus.mojo exec-maven-plugin + + org.apache.maven.plugins + maven-surefire-plugin + + + + ${app-name} + + + org.apache.maven.plugins maven-failsafe-plugin diff --git a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsConfigurationBean.java b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsConfigurationBean.java index a1295d5f0f..48249b5e17 100644 --- a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsConfigurationBean.java +++ b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsConfigurationBean.java @@ -2,7 +2,6 @@ import dk.dbc.dataio.common.utils.flowstore.FlowStoreServiceConnector; import dk.dbc.dataio.common.utils.flowstore.FlowStoreServiceConnectorException; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.HarvesterToken; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnector; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnectorException; @@ -25,13 +24,13 @@ public class PeriodicJobsConfigurationBean { JobStoreServiceConnector jobStoreServiceConnector; /** - * Returns delivery configuration for given chunk + * Returns delivery configuration for given job * - * @param chunk {@link Chunk} for which get delivery configuration + * @param chunkId id of the chunk the lookup is made for, which gates whether the + * delivery entity is persisted * @return delivery configuration as {@link PeriodicJobsDelivery} */ - public PeriodicJobsDelivery getDelivery(Chunk chunk, EntityManager entityManager) { - Integer jobId = Math.toIntExact(chunk.getJobId()); + public PeriodicJobsDelivery getDelivery(int jobId, int chunkId, EntityManager entityManager) { PeriodicJobsDelivery periodicJobsDelivery = deliveryCache.getIfPresent(jobId); if (periodicJobsDelivery != null) { // Return delivery entity from local bean cache. @@ -41,11 +40,11 @@ public PeriodicJobsDelivery getDelivery(Chunk chunk, EntityManager entityManager if (periodicJobsDelivery == null) { // Retrieve harvester config from flow-store and create new // delivery entity. - PeriodicJobsHarvesterConfig periodicJobsHarvesterConfig = getHarvesterConfig(chunk); - periodicJobsDelivery = new PeriodicJobsDelivery(Math.toIntExact(chunk.getJobId())); + PeriodicJobsHarvesterConfig periodicJobsHarvesterConfig = getHarvesterConfig(jobId); + periodicJobsDelivery = new PeriodicJobsDelivery(jobId); periodicJobsDelivery.setConfig(periodicJobsHarvesterConfig); } - if (chunk.getChunkId() == 0) { + if (chunkId == 0) { // Only allow the first chunk to persist the delivery entity entityManager.persist(periodicJobsDelivery); } @@ -54,8 +53,8 @@ public PeriodicJobsDelivery getDelivery(Chunk chunk, EntityManager entityManager return periodicJobsDelivery; } - private PeriodicJobsHarvesterConfig getHarvesterConfig(Chunk chunk) { - HarvesterToken harvesterToken = getHarvesterToken(chunk); + private PeriodicJobsHarvesterConfig getHarvesterConfig(int jobId) { + HarvesterToken harvesterToken = getHarvesterToken(jobId); try { return flowStoreServiceConnector .getHarvesterConfig(harvesterToken.getId(), PeriodicJobsHarvesterConfig.class); @@ -65,16 +64,16 @@ private PeriodicJobsHarvesterConfig getHarvesterConfig(Chunk chunk) { } } - private HarvesterToken getHarvesterToken(Chunk chunk) { + private HarvesterToken getHarvesterToken(int jobId) { try { JobListCriteria findJobCriteria = new JobListCriteria() .where(new ListFilter<>(JobListCriteria.Field.JOB_ID, - ListFilter.Op.EQUAL, chunk.getJobId())); + ListFilter.Op.EQUAL, jobId)); JobInfoSnapshot jobInfoSnapshot = jobStoreServiceConnector.listJobs(findJobCriteria).get(0); return HarvesterToken.of(jobInfoSnapshot.getSpecification().getAncestry().getHarvesterToken()); } catch (RuntimeException | JobStoreServiceConnectorException e) { throw new RuntimeException( - String.format("Failed to find job %d", chunk.getJobId()), e); + String.format("Failed to find job %d", jobId), e); } } diff --git a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsFinalizerBean.java b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsFinalizerBean.java index e611d3fd97..7dc77333c4 100644 --- a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsFinalizerBean.java +++ b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsFinalizerBean.java @@ -1,6 +1,5 @@ package dk.dbc.dataio.sink.periodicjobs; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; import dk.dbc.dataio.harvester.types.FtpPickup; @@ -28,23 +27,32 @@ public class PeriodicJobsFinalizerBean { PeriodicJobsFtpFinalizerBean periodicJobsFtpFinalizerBean; PeriodicJobsSFtpFinalizerBean periodicJobsSFtpFinalizerBean; - public Chunk handleTerminationChunk(Chunk chunk, EntityManager entityManager) throws InvalidMessageException { - LOGGER.info("Finalizing periodic job {}", chunk.getJobId()); + /** + * Delivers the job's accumulated datablocks to its pickup destination + * + * @param chunkId id of the job's termination chunk, needed only to tell an empty job + * from one that has data, see + * {@link dk.dbc.dataio.sink.periodicjobs.pickup.PeriodicJobsPickupFinalizer#isEmptyJob(int, int)} + * @return the job's delivering outcome, as the JOB_END item reported for its + * termination item + */ + public ChunkItem finalizeJob(int jobId, int chunkId, EntityManager entityManager) throws InvalidMessageException { + LOGGER.info("Finalizing periodic job {}", jobId); - PeriodicJobsDelivery delivery = periodicJobsConfigurationBean.getDelivery(chunk, entityManager); + PeriodicJobsDelivery delivery = periodicJobsConfigurationBean.getDelivery(jobId, chunkId, entityManager); Pickup pickup = delivery.getConfig().getContent().getPickup(); - Chunk result; + ChunkItem result; if (pickup instanceof HttpPickup) { - result = periodicJobsHttpFinalizerBean.deliver(chunk, delivery, entityManager); + result = periodicJobsHttpFinalizerBean.deliver(jobId, chunkId, delivery, entityManager); } else if (pickup instanceof MailPickup) { - result = periodicJobsMailFinalizerBean.deliver(chunk, delivery, entityManager); + result = periodicJobsMailFinalizerBean.deliver(jobId, chunkId, delivery, entityManager); } else if (pickup instanceof FtpPickup) { - result = periodicJobsFtpFinalizerBean.deliver(chunk, delivery, entityManager); + result = periodicJobsFtpFinalizerBean.deliver(jobId, chunkId, delivery, entityManager); } else if (pickup instanceof SFtpPickup) { - result = periodicJobsSFtpFinalizerBean.deliver(chunk, delivery, entityManager); + result = periodicJobsSFtpFinalizerBean.deliver(jobId, chunkId, delivery, entityManager); } else { - result = getUnhandledPickupTypeResult(chunk, pickup); + result = unhandledPickupTypeResult(pickup); } LOGGER.info("Deleted {} data blocks for job {}", @@ -68,13 +76,10 @@ public int deleteDelivery(Integer jobId, EntityManager entityManager) { .executeUpdate(); } - private Chunk getUnhandledPickupTypeResult(Chunk chunk, Pickup pickupType) { - final Chunk result = new Chunk(chunk.getJobId(), chunk.getChunkId(), Chunk.Type.DELIVERED); - result.insertItem( - ChunkItem.failedChunkItem() - .withType(ChunkItem.Type.JOB_END) - .withData("Unhandled pickup type: " + pickupType)); - return result; + private ChunkItem unhandledPickupTypeResult(Pickup pickupType) { + return ChunkItem.failedChunkItem() + .withType(ChunkItem.Type.JOB_END) + .withData("Unhandled pickup type: " + pickupType); } public PeriodicJobsFinalizerBean withPeriodicJobsConfigurationBean(PeriodicJobsConfigurationBean periodicJobsConfigurationBean) { this.periodicJobsConfigurationBean = periodicJobsConfigurationBean; diff --git a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumer.java b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumer.java index beafcb5452..b9880e1b1d 100644 --- a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumer.java +++ b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumer.java @@ -9,15 +9,15 @@ import dk.dbc.dataio.commons.conversion.Conversion; import dk.dbc.dataio.commons.conversion.ConversionException; import dk.dbc.dataio.commons.conversion.ConversionFactory; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.types.ConsumedMessage; import dk.dbc.dataio.commons.types.Diagnostic; -import dk.dbc.dataio.commons.types.Tools; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; +import dk.dbc.dataio.commons.types.jms.JMSHeader; import dk.dbc.dataio.commons.utils.lang.StringUtil; import dk.dbc.dataio.filestore.service.connector.FileStoreServiceConnector; -import dk.dbc.dataio.jse.artemis.common.jms.MessageConsumerAdapter; +import dk.dbc.dataio.jobstore.types.ItemDeliveryResult; +import dk.dbc.dataio.jse.artemis.common.jms.SinkMessageConsumerAdapter; import dk.dbc.dataio.jse.artemis.common.service.ServiceHub; import dk.dbc.dataio.sink.periodicjobs.mail.MailSession; import dk.dbc.dataio.sink.periodicjobs.pickup.PeriodicJobsFtpFinalizerBean; @@ -25,7 +25,6 @@ import dk.dbc.dataio.sink.periodicjobs.pickup.PeriodicJobsMailFinalizerBean; import dk.dbc.dataio.sink.periodicjobs.pickup.PeriodicJobsSFtpFinalizerBean; import dk.dbc.httpclient.FailSafeHttpClient; -import dk.dbc.log.DBCTrackedLogContext; import dk.dbc.proxy.ProxyBean; import dk.dbc.weekresolver.connector.WeekResolverConnector; import jakarta.persistence.EntityManager; @@ -48,7 +47,7 @@ import java.util.List; import java.util.Set; -public class PeriodicJobsMessageConsumer extends MessageConsumerAdapter { +public class PeriodicJobsMessageConsumer extends SinkMessageConsumerAdapter { private static final RetryPolicy RETRY_POLICY = new RetryPolicy() .handle(ProcessingException.class) @@ -119,33 +118,58 @@ void initializeFinalizers(ServiceHub serviceHub) { .withWeekResolverConnector(weekResolverConnector)); } + /** + * An item here is converted and persisted as a datablock rather than sent to a target + * system, and nothing leaves this sink until the job's termination item triggers the + * pickup, so there is no "a newer version of this record was already delivered" + * question to ask about a single item + *

+ * See docs/chunk-scheduling-redesign.md, Watermark opt-out. + */ + @Override + protected boolean usesDeliveryWatermark() { + return false; + } + /** + * Converts one item into datablocks, or, for the job's termination item, delivers the + * datablocks accumulated by every preceding item to the job's pickup destination + *

+ * The transaction is committed before this method returns, and only then does + * {@link SinkMessageConsumerAdapter} report the result. That order is what lets the + * job-end finalization run against complete data: a reported item is an item whose + * datablocks are durable, and the termination chunk is released only once every data + * item of the job has reported. + */ @Override - public void handleConsumedMessage(ConsumedMessage consumedMessage) - throws InvalidMessageException, NullPointerException { - Chunk chunk = unmarshallPayload(consumedMessage); + protected ItemDeliveryResult deliverItem(ConsumedMessage message, ChunkItem item) { + // Not null-checked: SinkMessageConsumerAdapter has already rejected the message as + // invalid if any of the three is missing. + int jobId = JMSHeader.jobId.getHeader(message, Integer.class); + int chunkId = JMSHeader.chunkId.getHeader(message, Long.class).intValue(); + short itemId = JMSHeader.itemId.getHeader(message, Short.class); + EntityManager entityManager = entityManagerFactory.createEntityManager(); EntityTransaction transaction = entityManager.getTransaction(); try { - Chunk result; transaction.begin(); - if (chunk.isTerminationChunk()) { - // Give the before-last message enough time to commit - // its datablocks to the database before initiating - // the finalization process. - // (The result is uploaded to the job-store before the - // implicit commit, so without the sleep pause, there was a - // small risk that the end-chunk would reach this bean - // before all data was available.) - Tools.sleep(5000); - result = periodicJobsFinalizerBean.handleTerminationChunk(chunk, entityManager); - } else { - result = handleChunk(chunk, entityManager); - } - sendResultToJobStore(result); + ChunkItem outcome = isTerminationItem(item) + ? periodicJobsFinalizerBean.finalizeJob(jobId, chunkId, entityManager) + : convertItem(item, jobId, chunkId, itemId, entityManager); transaction.commit(); + return ItemDeliveryResult.of(verdictOf(outcome), outcome); + } catch (InvalidMessageException e) { + // Thrown by the job-end finalization alone, and reported as failed rather than + // rethrown. A failed termination item completes the job and sets its fatal error + // flag. Nothing is committed, since the delivery it was rejected by did not happen. + LOGGER.error("Finalization of periodic job {} was rejected", jobId, e); + transaction.rollback(); + return ItemDeliveryResult.of(ItemDeliveryResult.Status.FAILED, jobEndFailure(item, e)); } finally { - if(transaction.isActive()) transaction.rollback(); + if (transaction.isActive()) { + transaction.rollback(); + } + entityManager.close(); } } @@ -160,6 +184,7 @@ public void abortJob(int jobId) { LOGGER.info("Aborted job {}", jobId); } finally { if(transaction.isActive()) transaction.commit(); + entityManager.close(); } } @@ -173,57 +198,87 @@ public String getAddress() { return ADDRESS; } - Chunk handleChunk(Chunk chunk, EntityManager entityManager) { - Chunk result = new Chunk(chunk.getJobId(), chunk.getChunkId(), Chunk.Type.DELIVERED); - try { - for (ChunkItem chunkItem : chunk.getItems()) { - DBCTrackedLogContext.setTrackingId(chunkItem.getTrackingId()); - result.insertItem(handleChunkItem(chunkItem, chunk, entityManager)); - } - } finally { - DBCTrackedLogContext.remove(); - } - return result; + /** + * Recognizes the job termination item the same way job-store does on its own side of + * the protocol ({@code PgJobStore.isTerminationItem}) + */ + private static boolean isTerminationItem(ChunkItem item) { + return item.isTyped() && item.getType().getFirst() == ChunkItem.Type.JOB_END; + } + + /** + * Maps a delivering outcome onto the verdict job-store counts the item by + *

+ * The mapping belongs here rather than in job-store, which reads the verdict alone: + * this sink owns both the outcome item and the verdict and is free to derive one from + * the other. + */ + private static ItemDeliveryResult.Status verdictOf(ChunkItem outcome) { + return switch (outcome.getStatus()) { + case SUCCESS -> ItemDeliveryResult.Status.DELIVERED; + case IGNORE -> ItemDeliveryResult.Status.IGNORED; + case FAILURE -> ItemDeliveryResult.Status.FAILED; + }; } - private ChunkItem handleChunkItem(ChunkItem chunkItem, Chunk chunk, EntityManager entityManager) { - ChunkItem result = new ChunkItem() - .withId(chunkItem.getId()) - .withTrackingId(chunkItem.getTrackingId()) + /** + * The delivering outcome recorded for a job whose finalization was rejected, keeping + * the item's JOB_END type so the job view still shows it for what it is + */ + private static ChunkItem jobEndFailure(ChunkItem item, Exception cause) { + return new ChunkItem() + .withId(item.getId()) + .withStatus(ChunkItem.Status.FAILURE) + .withType(ChunkItem.Type.JOB_END) + .withTrackingId(item.getTrackingId()) + .withDiagnostics(new Diagnostic(Diagnostic.Level.FATAL, cause.getMessage(), cause)) + .withData(cause.getMessage()); + } + + /** + * Converts one processed item into datablocks, to be delivered when the job ends + * + * @return delivering outcome for the item, which {@link #verdictOf(ChunkItem)} turns + * into the verdict reported for it + */ + ChunkItem convertItem(ChunkItem item, int jobId, int chunkId, short itemId, EntityManager entityManager) { + ChunkItem outcome = new ChunkItem() + .withId(itemId) + .withTrackingId(item.getTrackingId()) .withType(ChunkItem.Type.STRING) .withEncoding(StandardCharsets.UTF_8); try { - switch (chunkItem.getStatus()) { - case FAILURE: - return result - .withStatus(ChunkItem.Status.IGNORE) - .withData("Failed by processor"); - case IGNORE: - return result - .withStatus(ChunkItem.Status.IGNORE) - .withData("Ignored by processor"); - default: - convertChunkItem(chunkItem, chunk, entityManager); - return result + return switch (item.getStatus()) { + case FAILURE -> outcome + .withStatus(ChunkItem.Status.IGNORE) + .withData("Failed by processor"); + case IGNORE -> outcome + .withStatus(ChunkItem.Status.IGNORE) + .withData("Ignored by processor"); + case SUCCESS -> { + convertToDataBlocks(item, jobId, chunkId, itemId, entityManager); + yield outcome .withStatus(ChunkItem.Status.SUCCESS) .withData("Converted"); - } + } + }; } catch (RuntimeException e) { - return result + return outcome .withStatus(ChunkItem.Status.FAILURE) .withDiagnostics(new Diagnostic(Diagnostic.Level.FATAL, e.getMessage(), e)) .withData(e.getMessage()); } } - private void convertChunkItem(ChunkItem chunkItem, Chunk chunk, EntityManager entityManager) { + private void convertToDataBlocks(ChunkItem item, int jobId, int chunkId, short itemId, + EntityManager entityManager) { + int recordNumber = getRecordNumber(chunkId, itemId); try { - AddiReader addiReader = new AddiReader(new ByteArrayInputStream(chunkItem.getData())); + AddiReader addiReader = new AddiReader(new ByteArrayInputStream(item.getData())); byte[] data; int recordPart = 0; while (addiReader != null && addiReader.hasNext()) { - PeriodicJobsDataBlock.Key key = new PeriodicJobsDataBlock.Key(chunk.getJobId(), - getRecordNumber((int) chunk.getChunkId(), (int) chunkItem.getId()), recordPart); + PeriodicJobsDataBlock.Key key = new PeriodicJobsDataBlock.Key(jobId, recordNumber, recordPart); AddiRecord addiRecord; PeriodicJobsConversionParam conversionParam; String sortkey; @@ -241,7 +296,7 @@ private void convertChunkItem(ChunkItem chunkItem, Chunk chunk, EntityManager en } catch (IOException e) { // We assume that the IOException was caused by non-addi chunk item content addiReader = null; - data = chunkItem.getData(); + data = item.getData(); if (data == null || data.length == 0) { throw new IOException("Chunk item has empty data"); } diff --git a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFtpFinalizerBean.java b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFtpFinalizerBean.java index cd30233c3c..4a7834b2e9 100644 --- a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFtpFinalizerBean.java +++ b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFtpFinalizerBean.java @@ -1,6 +1,5 @@ package dk.dbc.dataio.sink.periodicjobs.pickup; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; import dk.dbc.dataio.harvester.types.FtpPickup; @@ -31,14 +30,15 @@ public PeriodicJobsFtpFinalizerBean() { @Timed @Override - public Chunk deliver(Chunk chunk, PeriodicJobsDelivery delivery, EntityManager entityManager) throws InvalidMessageException { - if (isEmptyJob(chunk)) { - return deliverEmptyFile(chunk, delivery); + public ChunkItem deliver(int jobId, int chunkId, PeriodicJobsDelivery delivery, + EntityManager entityManager) throws InvalidMessageException { + if (isEmptyJob(jobId, chunkId)) { + return deliverEmptyFile(delivery); } - return deliverDatablocks(chunk, delivery, entityManager); + return deliverDatablocks(jobId, chunkId, delivery, entityManager); } - private Chunk deliverEmptyFile(Chunk chunk, PeriodicJobsDelivery delivery) { + private ChunkItem deliverEmptyFile(PeriodicJobsDelivery delivery) { String remoteFile = getRemoteFilename(delivery) + ".EMPTY"; FtpPickup ftpPickup = (FtpPickup) delivery.getConfig().getContent().getPickup(); FtpClient ftpClient = null; @@ -50,11 +50,12 @@ private Chunk deliverEmptyFile(Chunk chunk, PeriodicJobsDelivery delivery) { ftpClient.close(); } } - return newResultChunk(chunk, + return newResultItem( String.format("Empty file %s uploaded to ftp host '%s'", remoteFile, ftpPickup.getFtpHost())); } - private Chunk deliverDatablocks(Chunk chunk, PeriodicJobsDelivery delivery, EntityManager entityManager) throws InvalidMessageException { + private ChunkItem deliverDatablocks(int jobId, int chunkId, PeriodicJobsDelivery delivery, + EntityManager entityManager) throws InvalidMessageException { String remoteFile = getRemoteFilename(delivery); FtpPickup ftpPickup = (FtpPickup) delivery.getConfig().getContent().getPickup(); File localFile = null; @@ -68,20 +69,20 @@ private Chunk deliverDatablocks(Chunk chunk, PeriodicJobsDelivery delivery, Enti .createLocalFile(); if (localFile.length() > 0) { uploadLocalFileToFtp(ftpPickup, localFile, remoteFile); - LOGGER.info("jobId '{}' uploaded to ftp host '{}'.", chunk.getJobId(), ftpPickup.getFtpHost()); + LOGGER.info("jobId '{}' uploaded to ftp host '{}'.", jobId, ftpPickup.getFtpHost()); } else { LOGGER.warn("jobId '{}' NOT uploaded to ftp host '{}' - no datablocks", - chunk.getJobId(), ftpPickup.getFtpHost()); + jobId, ftpPickup.getFtpHost()); } } catch (IOException e) { throw new InvalidMessageException(String.format("Unable to deliver datablocks for chuk: %d/%d", - chunk.getJobId(), chunk.getChunkId()),e); + jobId, chunkId),e); } finally { if (localFile != null) { if(!localFile.delete()) LOGGER.warn("Unable to delete file " + localFile); } } - return newResultChunk(chunk, + return newResultItem( String.format("File %s uploaded to ftp host '%s'", remoteFile, ftpPickup.getFtpHost())); } @@ -100,15 +101,12 @@ private void uploadLocalFileToFtp(FtpPickup ftpPickup, File local, String remote } } - private Chunk newResultChunk(Chunk chunk, String data) { - Chunk result = new Chunk(chunk.getJobId(), chunk.getChunkId(), Chunk.Type.DELIVERED); - ChunkItem chunkItem = ChunkItem.successfulChunkItem() + private ChunkItem newResultItem(String data) { + return ChunkItem.successfulChunkItem() .withId(0) .withType(ChunkItem.Type.JOB_END) .withData(data) .withEncoding(StandardCharsets.UTF_8); - result.insertItem(chunkItem); - return result; } FtpClient open(FtpPickup ftpPickup) { diff --git a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsHttpFinalizerBean.java b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsHttpFinalizerBean.java index 81f2a9574c..77dcafb512 100644 --- a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsHttpFinalizerBean.java +++ b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsHttpFinalizerBean.java @@ -6,7 +6,6 @@ import dk.dbc.commons.jpa.ResultSet; import dk.dbc.dataio.commons.conversion.ConversionMetadata; import dk.dbc.dataio.commons.macroexpansion.MacroSubstitutor; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; import dk.dbc.dataio.filestore.service.connector.FileStoreServiceConnector; @@ -34,8 +33,9 @@ public class PeriodicJobsHttpFinalizerBean extends PeriodicJobsPickupFinalizer { public FileStoreServiceConnector fileStoreServiceConnector; @Override - public Chunk deliver(Chunk chunk, PeriodicJobsDelivery delivery, EntityManager entityManager) throws InvalidMessageException { - boolean isEmptyJob = isEmptyJob(chunk); + public ChunkItem deliver(int jobId, int chunkId, PeriodicJobsDelivery delivery, + EntityManager entityManager) throws InvalidMessageException { + boolean isEmptyJob = isEmptyJob(jobId, chunkId); HttpPickup httpPickup = (HttpPickup) delivery.getConfig().getContent().getPickup(); ConversionMetadata fileMetadata = new ConversionMetadata(ORIGIN) .withJobId(delivery.getJobId()) @@ -63,7 +63,7 @@ public Chunk deliver(Chunk chunk, PeriodicJobsDelivery delivery, EntityManager e uploadMetadata(fileStoreServiceConnector, fileId, fileMetadata, delivery); } } - return newResultChunk(fileStoreServiceConnector, chunk, fileId, fileMetadata); + return newResultItem(fileStoreServiceConnector, fileId, fileMetadata); } private Optional fileAlreadyExists(FileStoreServiceConnector fileStoreServiceConnector, @@ -166,9 +166,8 @@ public PeriodicJobsHttpFinalizerBean withFileStoreServiceConnector(FileStoreServ return this; } - private Chunk newResultChunk(FileStoreServiceConnector fileStoreServiceConnector, Chunk chunk, - String fileId, ConversionMetadata fileMetadata) { - Chunk result = new Chunk(chunk.getJobId(), chunk.getChunkId(), Chunk.Type.DELIVERED); + private ChunkItem newResultItem(FileStoreServiceConnector fileStoreServiceConnector, + String fileId, ConversionMetadata fileMetadata) { ChunkItem chunkItem = ChunkItem.successfulChunkItem() .withId(0) .withType(ChunkItem.Type.JOB_END) @@ -180,8 +179,7 @@ private Chunk newResultChunk(FileStoreServiceConnector fileStoreServiceConnector } else { chunkItem.withData("No file uploaded"); } - result.insertItem(chunkItem); - return result; + return chunkItem; } @JsonIgnoreProperties(ignoreUnknown = true) diff --git a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsMailFinalizerBean.java b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsMailFinalizerBean.java index 93d561e09a..9c1f90b107 100644 --- a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsMailFinalizerBean.java +++ b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsMailFinalizerBean.java @@ -3,7 +3,6 @@ import dk.dbc.commons.jpa.ResultSet; import dk.dbc.dataio.common.utils.io.UncheckedByteArrayOutputStream; import dk.dbc.dataio.commons.macroexpansion.MacroSubstitutor; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; import dk.dbc.dataio.commons.utils.lang.StringUtil; @@ -42,33 +41,34 @@ public class PeriodicJobsMailFinalizerBean extends PeriodicJobsPickupFinalizer { @Timed @Override - public Chunk deliver(Chunk chunk, PeriodicJobsDelivery delivery, EntityManager entityManager) throws InvalidMessageException { + public ChunkItem deliver(int jobId, int chunkId, PeriodicJobsDelivery delivery, + EntityManager entityManager) throws InvalidMessageException { final MacroSubstitutor macroSubstitutor = getMacroSubstitutor(delivery); final MailPickup mailPickup = (MailPickup) delivery.getConfig().getContent().getPickup(); try { InternetAddress.parse(mailPickup.getRecipients()); } catch (AddressException e) { - return newFailedResultChunk(chunk, "Invalid mail recipient: " + e.getMessage()); + return newFailedResultItem("Invalid mail recipient: " + e.getMessage()); } final String content; - if (isEmptyJob(chunk)) { + if (isEmptyJob(jobId, chunkId)) { content = I18n.get("mail.empty_job.body"); } else { try { content = datablocksMailBody(delivery, macroSubstitutor, entityManager); } catch (IllegalStateException e) { - return newFailedResultChunk(chunk, "IllegalStateException: " + e.getMessage()); + return newFailedResultItem("IllegalStateException: " + e.getMessage()); } } if (!content.trim().isEmpty()) { sendMail(mailPickup, content, macroSubstitutor); - LOGGER.info("Job {}: mail sent to {}", chunk.getJobId(), mailPickup.getRecipients()); + LOGGER.info("Job {}: mail sent to {}", jobId, mailPickup.getRecipients()); } else { - LOGGER.warn("Job {}: no mail sent", chunk.getJobId()); + LOGGER.warn("Job {}: no mail sent", jobId); } - return newResultChunk(chunk, mailPickup); + return newResultItem(mailPickup); } private String datablocksMailBody(PeriodicJobsDelivery delivery, MacroSubstitutor macroSubstitutor, EntityManager entityManager) throws InvalidMessageException { @@ -171,26 +171,20 @@ public PeriodicJobsMailFinalizerBean withSession(Session session) { return this; } - private Chunk newResultChunk(Chunk chunk, MailPickup mailPickup) { - final Chunk result = new Chunk(chunk.getJobId(), chunk.getChunkId(), Chunk.Type.DELIVERED); - final ChunkItem chunkItem = ChunkItem.successfulChunkItem() + private ChunkItem newResultItem(MailPickup mailPickup) { + return ChunkItem.successfulChunkItem() .withId(0) .withType(ChunkItem.Type.JOB_END) .withEncoding(StandardCharsets.UTF_8) .withData(String.format("Mail sent to '%s' with subject '%s'", mailPickup.getRecipients(), mailPickup.getSubject())); - result.insertItem(chunkItem); - return result; } - private Chunk newFailedResultChunk(Chunk chunk, String cause) { - final Chunk result = new Chunk(chunk.getJobId(), chunk.getChunkId(), Chunk.Type.DELIVERED); - final ChunkItem chunkItem = ChunkItem.failedChunkItem() + private ChunkItem newFailedResultItem(String cause) { + return ChunkItem.failedChunkItem() .withId(0) .withType(ChunkItem.Type.JOB_END) .withEncoding(StandardCharsets.UTF_8) .withData(cause); - result.insertItem(chunkItem); - return result; } } diff --git a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsPickupFinalizer.java b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsPickupFinalizer.java index b9f6e49f73..9d4391db23 100644 --- a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsPickupFinalizer.java +++ b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsPickupFinalizer.java @@ -1,7 +1,7 @@ package dk.dbc.dataio.sink.periodicjobs.pickup; import dk.dbc.dataio.commons.macroexpansion.MacroSubstitutor; -import dk.dbc.dataio.commons.types.Chunk; +import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnector; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnectorException; @@ -26,12 +26,15 @@ public abstract class PeriodicJobsPickupFinalizer { JobStoreServiceConnector jobStoreServiceConnector; - public boolean isEmptyJob(Chunk endChunk) throws InvalidMessageException { - if (endChunk.getChunkId() == 0) { - // End chunk having ID 0 means job is empty + /** + * @param chunkId id of the job's termination chunk, which job-store sets to the job's + * data-chunk count, so zero means the job has no data + */ + public boolean isEmptyJob(int jobId, int chunkId) throws InvalidMessageException { + if (chunkId == 0) { return true; } - return isIgnoredJob(endChunk.getJobId()); + return isIgnoredJob(jobId); } public MacroSubstitutor getMacroSubstitutor(PeriodicJobsDelivery delivery) { @@ -108,7 +111,14 @@ private String getWeekCode(String catalogueCode, LocalDate localDate) { } } - public abstract Chunk deliver(Chunk chunk, PeriodicJobsDelivery delivery, EntityManager entityManager) throws InvalidMessageException; + /** + * Delivers the job's accumulated datablocks to its pickup destination + * + * @return the job's delivering outcome, as the JOB_END item reported for its + * termination item + */ + public abstract ChunkItem deliver(int jobId, int chunkId, PeriodicJobsDelivery delivery, + EntityManager entityManager) throws InvalidMessageException; public PeriodicJobsPickupFinalizer withJobStoreServiceConnector(JobStoreServiceConnector jobStoreServiceConnector) { this.jobStoreServiceConnector = jobStoreServiceConnector; diff --git a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsSFtpFinalizerBean.java b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsSFtpFinalizerBean.java index 6a98a71d41..8c6437a241 100644 --- a/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsSFtpFinalizerBean.java +++ b/sink/periodic-jobs/src/main/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsSFtpFinalizerBean.java @@ -3,7 +3,6 @@ import dk.dbc.commons.sftpclient.SFTPConfig; import dk.dbc.commons.sftpclient.SFtpClient; import dk.dbc.commons.sftpclient.SFtpClientException; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; import dk.dbc.dataio.harvester.types.SFtpPickup; @@ -30,24 +29,26 @@ public class PeriodicJobsSFtpFinalizerBean extends PeriodicJobsPickupFinalizer { @Timed @Override - public Chunk deliver(Chunk chunk, PeriodicJobsDelivery delivery, EntityManager entityManager) throws InvalidMessageException { - if (isEmptyJob(chunk)) { - return deliverEmptyFile(chunk, delivery); + public ChunkItem deliver(int jobId, int chunkId, PeriodicJobsDelivery delivery, + EntityManager entityManager) throws InvalidMessageException { + if (isEmptyJob(jobId, chunkId)) { + return deliverEmptyFile(delivery); } - return deliverDatablocks(chunk, delivery, entityManager); + return deliverDatablocks(jobId, delivery, entityManager); } - private Chunk deliverEmptyFile(Chunk chunk, PeriodicJobsDelivery delivery) { + private ChunkItem deliverEmptyFile(PeriodicJobsDelivery delivery) { final String remoteFile = getRemoteFilename(delivery) + ".EMPTY"; final SFtpPickup sFtpPickup = (SFtpPickup) delivery.getConfig().getContent().getPickup(); try (SFtpClient sFtpClient = open(sFtpPickup)) { sFtpClient.putContent(remoteFile, new ByteArrayInputStream("".getBytes())); } - return newResultChunk(chunk, + return newResultItem( String.format("Empty file %s uploaded to sftp host '%s'", remoteFile, sFtpPickup.getsFtpHost())); } - private Chunk deliverDatablocks(Chunk chunk, PeriodicJobsDelivery delivery, EntityManager entityManager) throws InvalidMessageException { + private ChunkItem deliverDatablocks(int jobId, PeriodicJobsDelivery delivery, + EntityManager entityManager) throws InvalidMessageException { final String remoteFile = getRemoteFilename(delivery); final SFtpPickup sftpPickup = (SFtpPickup) delivery.getConfig().getContent().getPickup(); File localFile = null; @@ -61,17 +62,17 @@ private Chunk deliverDatablocks(Chunk chunk, PeriodicJobsDelivery delivery, Enti .createLocalFile(); if (localFile.length() > 0) { uploadLocalFileToSFtp(sftpPickup, localFile, remoteFile); - LOGGER.info("jobId '{}' uploaded to sftp host '{}'.", chunk.getJobId(), sftpPickup.getsFtpHost()); + LOGGER.info("jobId '{}' uploaded to sftp host '{}'.", jobId, sftpPickup.getsFtpHost()); } else { LOGGER.warn("jobId '{}' NOT uploaded to sftp host '{}' - no datablocks", - chunk.getJobId(), sftpPickup.getsFtpHost()); + jobId, sftpPickup.getsFtpHost()); } } catch (IOException e) { throw new InvalidMessageException(String.format("Unable to deliver datablocks for:%d", delivery.getJobId()),e); } finally { if (localFile != null && !localFile.delete()) LOGGER.warn("Unable to delete file " + localFile); } - return newResultChunk(chunk, + return newResultItem( String.format("File %s uploaded to sftp host '%s'", remoteFile, sftpPickup.getsFtpHost())); } @@ -97,15 +98,12 @@ private void uploadLocalFileToSFtp(SFtpPickup sFtpPickup, File local, String rem } } - private Chunk newResultChunk(Chunk chunk, String data) { - final Chunk result = new Chunk(chunk.getJobId(), chunk.getChunkId(), Chunk.Type.DELIVERED); - final ChunkItem chunkItem = ChunkItem.successfulChunkItem() + private ChunkItem newResultItem(String data) { + return ChunkItem.successfulChunkItem() .withId(0) .withType(ChunkItem.Type.JOB_END) .withData(data) .withEncoding(StandardCharsets.UTF_8); - result.insertItem(chunkItem); - return result; } public PeriodicJobsSFtpFinalizerBean withProxyBean(ProxyBean proxyBean) { diff --git a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsConfigurationBeanIT.java b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsConfigurationBeanIT.java index eb14d7d776..b1f5198f09 100644 --- a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsConfigurationBeanIT.java +++ b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsConfigurationBeanIT.java @@ -4,13 +4,11 @@ import dk.dbc.dataio.common.utils.flowstore.FlowStoreServiceConnector; import dk.dbc.dataio.common.utils.flowstore.FlowStoreServiceConnectorException; import dk.dbc.dataio.common.utils.flowstore.ejb.FlowStoreServiceConnectorBean; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.JobSpecification; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnector; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnectorException; import dk.dbc.dataio.commons.utils.jobstore.ejb.JobStoreServiceConnectorBean; -import dk.dbc.dataio.commons.utils.test.model.ChunkBuilder; import dk.dbc.dataio.harvester.types.PeriodicJobsHarvesterConfig; import dk.dbc.dataio.jobstore.types.JobInfoSnapshot; import dk.dbc.dataio.jobstore.types.criteria.JobListCriteria; @@ -50,20 +48,19 @@ public void setupMocks() { @Test public void getDelivery_throwsOnFailureToResolveJob() throws JobStoreServiceConnectorException { - final Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED).setJobId(0).build(); when(jobStoreServiceConnector.listJobs(any(JobListCriteria.class))) .thenReturn(Collections.emptyList()); final PeriodicJobsConfigurationBean periodicJobsConfigurationBean = newPeriodicJobsConfigurationBean(); - assertThat(() -> periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager()), isThrowing(RuntimeException.class)); + assertThat(() -> periodicJobsConfigurationBean.getDelivery(0, 1, env().getEntityManager()), isThrowing(RuntimeException.class)); } @Test public void getDelivery_throwsOnFailureToResolveHarvesterConfig() throws JobStoreServiceConnectorException, FlowStoreServiceConnectorException { - Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED).build(); + final int jobId = 3; JobInfoSnapshot jobInfoSnapshot = new JobInfoSnapshot() - .withJobId(chunk.getJobId()) + .withJobId(jobId) .withSpecification( new JobSpecification() .withAncestry(new JobSpecification.Ancestry() @@ -75,18 +72,16 @@ public void getDelivery_throwsOnFailureToResolveHarvesterConfig() .thenThrow(new FlowStoreServiceConnectorException("DIED")); PeriodicJobsConfigurationBean periodicJobsConfigurationBean = newPeriodicJobsConfigurationBean(); - assertThat(() -> periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager()), isThrowing(RuntimeException.class)); + assertThat(() -> periodicJobsConfigurationBean.getDelivery(jobId, 1, env().getEntityManager()), isThrowing(RuntimeException.class)); } @Test public void getDelivery_onlyFirstChunkPersists() throws JobStoreServiceConnectorException, FlowStoreServiceConnectorException, SQLException { - final Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED) - .setJobId(1) - .setChunkId(1) - .build(); + final int jobId = 1; + final int chunkId = 1; final JobInfoSnapshot jobInfoSnapshot = new JobInfoSnapshot() - .withJobId(chunk.getJobId()) + .withJobId(jobId) .withSpecification( new JobSpecification() .withAncestry(new JobSpecification.Ancestry() @@ -101,12 +96,12 @@ public void getDelivery_onlyFirstChunkPersists() final PeriodicJobsConfigurationBean periodicJobsConfigurationBean = newPeriodicJobsConfigurationBean(); PeriodicJobsDelivery delivery = env().getPersistenceContext().run(() -> - periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager())); + periodicJobsConfigurationBean.getDelivery(jobId, chunkId, env().getEntityManager())); - assertThat("delivery.jobId", delivery.getJobId(), is(chunk.getJobId())); + assertThat("delivery.jobId", delivery.getJobId(), is(jobId)); assertThat("delivery.config", delivery.getConfig(), is(periodicJobsHarvesterConfig)); assertThat("delivery is cached", - periodicJobsConfigurationBean.deliveryCache.getIfPresent(chunk.getJobId()), is(notNullValue())); + periodicJobsConfigurationBean.deliveryCache.getIfPresent(jobId), is(notNullValue())); try (Connection conn = connectToPeriodicJobsDB()) { assertThat("number of persisted deliveries", @@ -117,12 +112,10 @@ public void getDelivery_onlyFirstChunkPersists() @Test public void getDelivery_firstChunkPersists() throws JobStoreServiceConnectorException, FlowStoreServiceConnectorException, SQLException { - final Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED) - .setJobId(1) - .setChunkId(0) - .build(); + final int jobId = 1; + final int chunkId = 0; final JobInfoSnapshot jobInfoSnapshot = new JobInfoSnapshot() - .withJobId(chunk.getJobId()) + .withJobId(jobId) .withSpecification( new JobSpecification() .withAncestry(new JobSpecification.Ancestry() @@ -137,12 +130,12 @@ public void getDelivery_firstChunkPersists() final PeriodicJobsConfigurationBean periodicJobsConfigurationBean = newPeriodicJobsConfigurationBean(); PeriodicJobsDelivery delivery = env().getPersistenceContext().run(() -> - periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager())); + periodicJobsConfigurationBean.getDelivery(jobId, chunkId, env().getEntityManager())); - assertThat("delivery.jobId", delivery.getJobId(), is(chunk.getJobId())); + assertThat("delivery.jobId", delivery.getJobId(), is(jobId)); assertThat("delivery.config", delivery.getConfig(), is(periodicJobsHarvesterConfig)); assertThat("delivery is cached", - periodicJobsConfigurationBean.deliveryCache.getIfPresent(chunk.getJobId()), is(notNullValue())); + periodicJobsConfigurationBean.deliveryCache.getIfPresent(jobId), is(notNullValue())); try (Connection conn = connectToPeriodicJobsDB()) { assertThat("number of persisted deliveries", @@ -152,38 +145,34 @@ public void getDelivery_firstChunkPersists() @Test public void getDelivery_servesFromCache() throws InvalidMessageException { - final Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED) - .setJobId(42) - .setChunkId(5) - .build(); + final int jobId = 42; + final int chunkId = 5; final PeriodicJobsHarvesterConfig periodicJobsHarvesterConfig = new PeriodicJobsHarvesterConfig(1, 1, new PeriodicJobsHarvesterConfig.Content()); - final PeriodicJobsDelivery expectedDelivery = new PeriodicJobsDelivery(chunk.getJobId()); + final PeriodicJobsDelivery expectedDelivery = new PeriodicJobsDelivery(jobId); expectedDelivery.setConfig(periodicJobsHarvesterConfig); final PeriodicJobsConfigurationBean periodicJobsConfigurationBean = newPeriodicJobsConfigurationBean(); - periodicJobsConfigurationBean.deliveryCache.put(chunk.getJobId(), expectedDelivery); - assertThat(periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager()), is(expectedDelivery)); + periodicJobsConfigurationBean.deliveryCache.put(jobId, expectedDelivery); + assertThat(periodicJobsConfigurationBean.getDelivery(jobId, chunkId, env().getEntityManager()), is(expectedDelivery)); } @Test public void getDelivery_servesFromDatabase() throws InvalidMessageException { - final Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED) - .setJobId(42) - .setChunkId(5) - .build(); + final int jobId = 42; + final int chunkId = 5; final PeriodicJobsHarvesterConfig periodicJobsHarvesterConfig = new PeriodicJobsHarvesterConfig(1, 1, new PeriodicJobsHarvesterConfig.Content()); - final PeriodicJobsDelivery expectedDelivery = new PeriodicJobsDelivery(chunk.getJobId()); + final PeriodicJobsDelivery expectedDelivery = new PeriodicJobsDelivery(jobId); expectedDelivery.setConfig(periodicJobsHarvesterConfig); env().getPersistenceContext().run(() -> env().getEntityManager().persist(expectedDelivery)); final PeriodicJobsConfigurationBean periodicJobsConfigurationBean = newPeriodicJobsConfigurationBean(); - assertThat(periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager()), is(expectedDelivery)); + assertThat(periodicJobsConfigurationBean.getDelivery(jobId, chunkId, env().getEntityManager()), is(expectedDelivery)); assertThat("delivery is cached", - periodicJobsConfigurationBean.deliveryCache.getIfPresent(chunk.getJobId()), is(notNullValue())); + periodicJobsConfigurationBean.deliveryCache.getIfPresent(jobId), is(notNullValue())); } private PeriodicJobsConfigurationBean newPeriodicJobsConfigurationBean() { diff --git a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumerIT.java b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumerIT.java index d1904f1c8d..12b99fcc08 100644 --- a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumerIT.java +++ b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumerIT.java @@ -2,9 +2,7 @@ import dk.dbc.commons.addi.AddiRecord; import dk.dbc.dataio.commons.conversion.ConversionParam; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.ChunkItem; -import dk.dbc.dataio.commons.utils.test.model.ChunkBuilder; import dk.dbc.dataio.commons.utils.test.model.ChunkItemBuilder; import dk.dbc.dataio.jse.artemis.common.service.ServiceHub; import org.junit.Test; @@ -12,8 +10,8 @@ import org.testcontainers.shaded.com.fasterxml.jackson.databind.ObjectMapper; import java.nio.charset.StandardCharsets; +import java.util.ArrayList; import java.util.Arrays; -import java.util.Collections; import java.util.List; import static org.hamcrest.CoreMatchers.is; @@ -24,8 +22,13 @@ public class PeriodicJobsMessageConsumerIT extends IntegrationTest { private final AddiRecord addiRecord1 = newAddiRecord(new ConversionParam(), "record-1"); private final AddiRecord addiRecord2 = newAddiRecord(new PeriodicJobsConversionParam().withSortkey("custom-sortkey").withRecordHeader("custom-header\n"), "record-2"); + /** + * The datablock keys asserted here are what pins the record number against the item + * ids the message headers carry, since they decide where each converted record lands + * in the sort order of the file delivered when the job ends. + */ @Test - public void handleChunk() { + public void convertItems() { PeriodicJobsMessageConsumer periodicJobsMessageConsumer = newMessageConsumerBean(); List chunkItems = Arrays.asList( @@ -35,15 +38,23 @@ public void handleChunk() { new ChunkItemBuilder().setId(3L).setStatus(ChunkItem.Status.SUCCESS).setData(addiRecord1.getBytes()).build(), new ChunkItemBuilder().setId(4L).setStatus(ChunkItem.Status.SUCCESS).setData(addiRecord2.getBytes()).build()); final int jobId = 42; - Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED).setJobId(jobId).setChunkId(0L).setItems(chunkItems).build(); + final int chunkId = 0; + + List outcomes = env().getPersistenceContext().run(() -> { + List results = new ArrayList<>(); + for (ChunkItem chunkItem : chunkItems) { + results.add(periodicJobsMessageConsumer.convertItem(chunkItem, jobId, chunkId, + (short) chunkItem.getId(), env().getEntityManager())); + } + return results; + }); - Chunk result = env().getPersistenceContext().run(() -> periodicJobsMessageConsumer.handleChunk(chunk, env().getEntityManager())); - assertThat("number of chunk items", result.size(), is(5)); - assertThat("1st chunk item", result.getItems().get(0).getStatus(), is(ChunkItem.Status.IGNORE)); - assertThat("2nd chunk item", result.getItems().get(1).getStatus(), is(ChunkItem.Status.SUCCESS)); - assertThat("3rd chunk item", result.getItems().get(2).getStatus(), is(ChunkItem.Status.IGNORE)); - assertThat("4th chunk item", result.getItems().get(3).getStatus(), is(ChunkItem.Status.SUCCESS)); - assertThat("5th chunk item", result.getItems().get(4).getStatus(), is(ChunkItem.Status.SUCCESS)); + assertThat("number of outcomes", outcomes.size(), is(5)); + assertThat("1st outcome", outcomes.get(0).getStatus(), is(ChunkItem.Status.IGNORE)); + assertThat("2nd outcome", outcomes.get(1).getStatus(), is(ChunkItem.Status.SUCCESS)); + assertThat("3rd outcome", outcomes.get(2).getStatus(), is(ChunkItem.Status.IGNORE)); + assertThat("4th outcome", outcomes.get(3).getStatus(), is(ChunkItem.Status.SUCCESS)); + assertThat("5th outcome", outcomes.get(4).getStatus(), is(ChunkItem.Status.SUCCESS)); PeriodicJobsDataBlock.Key key1 = new PeriodicJobsDataBlock.Key(jobId, 1, 0); PeriodicJobsDataBlock datablock1 = env().getPersistenceContext().run(() -> env().getEntityManager().find(PeriodicJobsDataBlock.class, key1)); @@ -73,8 +84,9 @@ public void handleChunk() { @Test public void overwriteExistingDataBlock() { final int jobId = 42; - Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED).setJobId(jobId).setChunkId(7L) - .setItems(Collections.singletonList(new ChunkItemBuilder().setId(0L).setStatus(ChunkItem.Status.SUCCESS).setData(addiRecord1.getBytes()).build())).build(); + final int chunkId = 7; + ChunkItem chunkItem = new ChunkItemBuilder().setId(0L).setStatus(ChunkItem.Status.SUCCESS) + .setData(addiRecord1.getBytes()).build(); PeriodicJobsDataBlock.Key key = new PeriodicJobsDataBlock.Key(jobId, 70, 0); PeriodicJobsDataBlock existingDatablock = new PeriodicJobsDataBlock(); @@ -88,20 +100,25 @@ public void overwriteExistingDataBlock() { PeriodicJobsMessageConsumer periodicJobsMessageConsumer = newMessageConsumerBean(); - Chunk result = env().getPersistenceContext().run(() -> periodicJobsMessageConsumer.handleChunk(chunk, env().getEntityManager())); + ChunkItem outcome = env().getPersistenceContext().run(() -> + periodicJobsMessageConsumer.convertItem(chunkItem, jobId, chunkId, (short) 0, + env().getEntityManager())); - assertThat("1st chunk item", result.getItems().get(0).getStatus(), is(ChunkItem.Status.SUCCESS)); + assertThat("outcome", outcome.getStatus(), is(ChunkItem.Status.SUCCESS)); } @Test public void emptyConversionResultsFails() { final int jobId = 42; + final int chunkId = 0; PeriodicJobsMessageConsumer periodicJobsMessageConsumer = newMessageConsumerBean(); - List chunkItems = Collections.singletonList(new ChunkItemBuilder().setId(0L).setStatus(ChunkItem.Status.SUCCESS).setData(newAddiRecord(new ConversionParam(), "").getBytes()).build()); - Chunk chunk = new ChunkBuilder(Chunk.Type.PROCESSED).setJobId(jobId).setChunkId(0L).setItems(chunkItems).build(); + ChunkItem chunkItem = new ChunkItemBuilder().setId(0L).setStatus(ChunkItem.Status.SUCCESS) + .setData(newAddiRecord(new ConversionParam(), "").getBytes()).build(); - Chunk result = env().getPersistenceContext().run(() -> periodicJobsMessageConsumer.handleChunk(chunk, env().getEntityManager())); - assertThat("1st chunk item", result.getItems().get(0).getStatus(), is(ChunkItem.Status.FAILURE)); + ChunkItem outcome = env().getPersistenceContext().run(() -> + periodicJobsMessageConsumer.convertItem(chunkItem, jobId, chunkId, (short) 0, + env().getEntityManager())); + assertThat("outcome", outcome.getStatus(), is(ChunkItem.Status.FAILURE)); PeriodicJobsDataBlock.Key key = new PeriodicJobsDataBlock.Key(jobId, 0, 0); PeriodicJobsDataBlock datablock = env().getPersistenceContext().run(() -> env().getEntityManager().find(PeriodicJobsDataBlock.class, key)); diff --git a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumerTest.java b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumerTest.java new file mode 100644 index 0000000000..22beda2772 --- /dev/null +++ b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/PeriodicJobsMessageConsumerTest.java @@ -0,0 +1,249 @@ +package dk.dbc.dataio.sink.periodicjobs; + +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import dk.dbc.commons.addi.AddiRecord; +import dk.dbc.commons.jsonb.JSONBContext; +import dk.dbc.commons.jsonb.JSONBException; +import dk.dbc.dataio.commons.conversion.ConversionParam; +import dk.dbc.dataio.commons.types.ChunkItem; +import dk.dbc.dataio.commons.types.ConsumedMessage; +import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; +import dk.dbc.dataio.commons.types.jms.JMSHeader; +import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnector; +import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnectorException; +import dk.dbc.dataio.commons.utils.lang.StringUtil; +import dk.dbc.dataio.jobstore.types.ItemDeliveryResult; +import dk.dbc.dataio.jse.artemis.common.service.ServiceHub; +import jakarta.persistence.EntityManager; +import jakarta.persistence.EntityManagerFactory; +import jakarta.persistence.EntityTransaction; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; +import org.mockito.InOrder; + +import java.nio.charset.StandardCharsets; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.hamcrest.CoreMatchers.is; +import static org.hamcrest.CoreMatchers.notNullValue; +import static org.hamcrest.MatcherAssert.assertThat; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.anyInt; +import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.inOrder; +import static org.mockito.Mockito.mock; +import static org.mockito.ArgumentMatchers.anyShort; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +/** + * Covers this sink's own delivery surface: the verdict each processing outcome maps to, + * the job-end branch, and the watermark opt-out. Header reading, the watermark decision + * itself and result reporting belong to {@code SinkMessageConsumerAdapter} and are covered + * by its own test. + */ +class PeriodicJobsMessageConsumerTest { + private static final int JOB_ID = 42; + private static final int CHUNK_ID = 7; + private static final short ITEM_ID = 3; + private static final long SINK_ID = 15; + private static final String TRACKING_ID = "rr:1223io:12534"; + + private final JSONBContext jsonbContext = new JSONBContext(); + private final JobStoreServiceConnector jobStoreServiceConnector = mock(JobStoreServiceConnector.class); + private final EntityManagerFactory entityManagerFactory = mock(EntityManagerFactory.class); + private final EntityManager entityManager = mock(EntityManager.class); + private final EntityTransaction transaction = mock(EntityTransaction.class); + private final PeriodicJobsFinalizerBean periodicJobsFinalizerBean = mock(PeriodicJobsFinalizerBean.class); + + private PeriodicJobsMessageConsumer consumer; + + @BeforeEach + void setupConsumer() { + when(entityManagerFactory.createEntityManager()).thenReturn(entityManager); + when(entityManager.getTransaction()).thenReturn(transaction); + consumer = new PeriodicJobsMessageConsumer( + new ServiceHub.Builder().withJobStoreServiceConnector(jobStoreServiceConnector).test(), + entityManagerFactory); + consumer.periodicJobsFinalizerBean = periodicJobsFinalizerBean; + } + + @Test + void successfulConversion_isDelivered() { + ItemDeliveryResult result = consumer.deliverItem(itemMessage(), addiItem("record-1")); + + assertThat("verdict", result.status(), is(ItemDeliveryResult.Status.DELIVERED)); + assertThat("outcome status", result.chunkItem().getStatus(), is(ChunkItem.Status.SUCCESS)); + assertThat("outcome data", StringUtil.asString(result.chunkItem().getData()), is("Converted")); + verify(entityManager).persist(any(PeriodicJobsDataBlock.class)); + } + + /** + * IGNORED rather than DELIVERED is what keeps an item this sink converted nothing for + * counted as ignored in the delivering phase, as it was when whole chunks were + * delivered. + */ + @Test + void processorFailure_isIgnored() { + ItemDeliveryResult result = consumer.deliverItem(itemMessage(), item(ChunkItem.Status.FAILURE, "irrelevant")); + + assertThat("verdict", result.status(), is(ItemDeliveryResult.Status.IGNORED)); + assertThat("outcome status", result.chunkItem().getStatus(), is(ChunkItem.Status.IGNORE)); + assertThat("outcome data", StringUtil.asString(result.chunkItem().getData()), is("Failed by processor")); + verify(entityManager, never()).persist(any()); + } + + @Test + void processorIgnore_isIgnored() { + ItemDeliveryResult result = consumer.deliverItem(itemMessage(), item(ChunkItem.Status.IGNORE, "irrelevant")); + + assertThat("verdict", result.status(), is(ItemDeliveryResult.Status.IGNORED)); + assertThat("outcome status", result.chunkItem().getStatus(), is(ChunkItem.Status.IGNORE)); + assertThat("outcome data", StringUtil.asString(result.chunkItem().getData()), is("Ignored by processor")); + verify(entityManager, never()).persist(any()); + } + + @Test + void conversionFailure_isFailed() { + ItemDeliveryResult result = consumer.deliverItem(itemMessage(), addiItem("")); + + assertThat("verdict", result.status(), is(ItemDeliveryResult.Status.FAILED)); + assertThat("outcome status", result.chunkItem().getStatus(), is(ChunkItem.Status.FAILURE)); + assertThat("outcome diagnostics", result.chunkItem().getDiagnostics(), is(notNullValue())); + verify(entityManager, never()).persist(any()); + } + + @Test + void outcomeCarriesItemIdAndTrackingId() { + ChunkItem outcome = consumer.deliverItem(itemMessage(), addiItem("record-1")).chunkItem(); + + assertThat("id", outcome.getId(), is((long) ITEM_ID)); + assertThat("tracking id", outcome.getTrackingId(), is(TRACKING_ID)); + } + + /** + * The job and chunk ids the finalizer chain is handed must be the ones the item + * message carries, since the chunk id is what tells an empty job from one with data. + */ + @Test + void terminationItem_isFinalized() throws InvalidMessageException { + when(periodicJobsFinalizerBean.finalizeJob(anyInt(), anyInt(), any(EntityManager.class))) + .thenReturn(deliveredJobEndItem()); + + ItemDeliveryResult result = consumer.deliverItem(itemMessage(), terminationItem(ChunkItem.Status.SUCCESS)); + + verify(periodicJobsFinalizerBean).finalizeJob(eq(JOB_ID), eq(CHUNK_ID), any(EntityManager.class)); + + assertThat("verdict", result.status(), is(ItemDeliveryResult.Status.DELIVERED)); + assertThat("outcome type", result.chunkItem().getType(), is(List.of(ChunkItem.Type.JOB_END))); + } + + /** + * On the chunk path this exception had the message discarded without a result, leaving + * the job forever incomplete. A FAILED verdict completes it and sets its fatal error + * flag instead. + */ + @Test + void terminationItem_finalizationRejected_isFailed() throws InvalidMessageException { + when(periodicJobsFinalizerBean.finalizeJob(anyInt(), anyInt(), any(EntityManager.class))) + .thenThrow(new InvalidMessageException("no pickup for you")); + + ItemDeliveryResult result = consumer.deliverItem(itemMessage(), terminationItem(ChunkItem.Status.SUCCESS)); + + assertThat("verdict", result.status(), is(ItemDeliveryResult.Status.FAILED)); + assertThat("outcome status", result.chunkItem().getStatus(), is(ChunkItem.Status.FAILURE)); + assertThat("outcome diagnostics", result.chunkItem().getDiagnostics(), is(notNullValue())); + } + + /** + * The datablock write must be committed before the framework reports the item, since + * that report is what releases the job's termination chunk, and the finalization it + * triggers reads the datablocks with a different EntityManager in a different + * transaction. Reporting first is the ordering the chunk protocol had, and the reason + * it needed a sleep before finalizing. + */ + @Test + void commitPrecedesResultReport() throws InvalidMessageException, JobStoreServiceConnectorException { + consumer.handleConsumedMessage(itemMessage()); + + InOrder ordered = inOrder(transaction, jobStoreServiceConnector); + ordered.verify(transaction).commit(); + ordered.verify(jobStoreServiceConnector).addItemDelivered(any(ItemDeliveryResult.class), anyInt(), + anyInt(), anyShort()); + } + + /** + * This sink opts out of the delivery watermark, which is observable only as the lookup + * not being made. + */ + @Test + void watermarkIsNotConsulted() throws InvalidMessageException, JobStoreServiceConnectorException { + consumer.handleConsumedMessage(itemMessage()); + + verify(jobStoreServiceConnector, never()).getWatermark(anyInt(), anyString()); + verify(jobStoreServiceConnector).addItemDelivered(any(ItemDeliveryResult.class), anyInt(), anyInt(), + anyShort()); + } + + private ChunkItem item(ChunkItem.Status status, String data) { + return new ChunkItem() + .withId(ITEM_ID) + .withStatus(status) + .withType(ChunkItem.Type.STRING) + .withTrackingId(TRACKING_ID) + .withData(data); + } + + private ChunkItem addiItem(String record) { + return new ChunkItem() + .withId(ITEM_ID) + .withStatus(ChunkItem.Status.SUCCESS) + .withType(ChunkItem.Type.ADDI) + .withTrackingId(TRACKING_ID) + .withData(newAddiRecord(record).getBytes()); + } + + private ChunkItem terminationItem(ChunkItem.Status status) { + return new ChunkItem() + .withId(0) + .withStatus(status) + .withType(ChunkItem.Type.JOB_END) + .withTrackingId(JOB_ID + ".JOB_END") + .withData("Job termination item"); + } + + private ChunkItem deliveredJobEndItem() { + return ChunkItem.successfulChunkItem() + .withId(0) + .withType(ChunkItem.Type.JOB_END) + .withData("delivered to pickup"); + } + + private AddiRecord newAddiRecord(String record) { + try { + byte[] metadata = new ObjectMapper().writeValueAsBytes(new ConversionParam()); + return new AddiRecord(metadata, record.getBytes(StandardCharsets.UTF_8)); + } catch (JsonProcessingException e) { + throw new IllegalStateException(e); + } + } + + private ConsumedMessage itemMessage() { + Map headers = new HashMap<>(); + headers.put(JMSHeader.payload.name, JMSHeader.ITEM_PAYLOAD_TYPE); + headers.put(JMSHeader.jobId.name, JOB_ID); + headers.put(JMSHeader.chunkId.name, (long) CHUNK_ID); + headers.put(JMSHeader.itemId.name, ITEM_ID); + headers.put(JMSHeader.sinkId.name, SINK_ID); + try { + return new ConsumedMessage("id", headers, jsonbContext.marshall(addiItem("record-1"))); + } catch (JSONBException e) { + throw new IllegalStateException(e); + } + } +} diff --git a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFinalizerBeanIT.java b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFinalizerBeanIT.java index fc6c2a6a8f..cbf63b54be 100644 --- a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFinalizerBeanIT.java +++ b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFinalizerBeanIT.java @@ -1,7 +1,7 @@ package dk.dbc.dataio.sink.periodicjobs.pickup; import dk.dbc.commons.jdbc.util.JDBCUtil; -import dk.dbc.dataio.commons.types.Chunk; +import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.types.exceptions.InvalidMessageException; import dk.dbc.dataio.commons.utils.lang.StringUtil; import dk.dbc.dataio.harvester.types.HttpPickup; @@ -23,6 +23,9 @@ import static org.mockito.Mockito.when; public class PeriodicJobsFinalizerBeanIT extends IntegrationTest { + /* Id of the job's termination chunk, non-zero so the job counts as having data */ + private static final int CHUNK_ID = 3; + private final PeriodicJobsConfigurationBean periodicJobsConfigurationBean = mock(PeriodicJobsConfigurationBean.class); private final PeriodicJobsHttpFinalizerBean periodicJobsHttpFinalizerBean = @@ -60,13 +63,12 @@ public void deletesDataBlocks() throws SQLException { delivery.setConfig(new PeriodicJobsHarvesterConfig(1, 1, new PeriodicJobsHarvesterConfig.Content() .withPickup(new HttpPickup()))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); - when(periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager())) + when(periodicJobsConfigurationBean.getDelivery(jobId, CHUNK_ID, env().getEntityManager())) .thenReturn(delivery); PeriodicJobsFinalizerBean periodicJobsFinalizerBean = newPeriodicJobsFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsFinalizerBean.handleTerminationChunk(chunk, env().getEntityManager())); + periodicJobsFinalizerBean.finalizeJob(jobId, CHUNK_ID, env().getEntityManager())); try (Connection conn = connectToPeriodicJobsDB()) { assertThat("number of remaining persisted data blocks", @@ -97,13 +99,12 @@ public void deletesDelivery() throws SQLException { env().getEntityManager().persist(delivery1); }); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); - when(periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager())) + when(periodicJobsConfigurationBean.getDelivery(jobId, CHUNK_ID, env().getEntityManager())) .thenReturn(delivery1); PeriodicJobsFinalizerBean periodicJobsFinalizerBean = newPeriodicJobsFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsFinalizerBean.handleTerminationChunk(chunk, env().getEntityManager())); + periodicJobsFinalizerBean.finalizeJob(jobId, CHUNK_ID, env().getEntityManager())); try (Connection conn = connectToPeriodicJobsDB()) { assertThat("number of remaining persisted deliveries", @@ -123,19 +124,20 @@ public void returnsResultOfDelivery() throws InvalidMessageException { new PeriodicJobsHarvesterConfig.Content() .withPickup(new HttpPickup()))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); - when(periodicJobsConfigurationBean.getDelivery(chunk, env().getEntityManager())) + when(periodicJobsConfigurationBean.getDelivery(jobId, CHUNK_ID, env().getEntityManager())) .thenReturn(delivery); - Chunk expectedResult = new Chunk(jobId, 3, Chunk.Type.DELIVERED); - when(periodicJobsHttpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())) + ChunkItem expectedResult = ChunkItem.successfulChunkItem() + .withId(0) + .withType(ChunkItem.Type.JOB_END); + when(periodicJobsHttpFinalizerBean.deliver(jobId, CHUNK_ID, delivery, env().getEntityManager())) .thenReturn(expectedResult); PeriodicJobsFinalizerBean periodicJobsFinalizerBean = newPeriodicJobsFinalizerBean(); - Chunk result = env().getPersistenceContext().run(() -> - periodicJobsFinalizerBean.handleTerminationChunk(chunk, env().getEntityManager())); + ChunkItem result = env().getPersistenceContext().run(() -> + periodicJobsFinalizerBean.finalizeJob(jobId, CHUNK_ID, env().getEntityManager())); - assertThat("result chunk", result, is(sameInstance(expectedResult))); + assertThat("result item", result, is(sameInstance(expectedResult))); } private PeriodicJobsFinalizerBean newPeriodicJobsFinalizerBean() { diff --git a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFtpFinalizerBeanIT.java b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFtpFinalizerBeanIT.java index 6aab8ef053..4003295942 100644 --- a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFtpFinalizerBeanIT.java +++ b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsFtpFinalizerBeanIT.java @@ -1,6 +1,5 @@ package dk.dbc.dataio.sink.periodicjobs.pickup; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnector; import dk.dbc.dataio.commons.utils.lang.StringUtil; import dk.dbc.dataio.harvester.types.FtpPickup; @@ -66,10 +65,9 @@ public void deliver_onNonEmptyJobNoDataBlocks() { .withFtpUser(USERNAME) .withFtpPassword(PASSWORD) .withFtpSubdirectory(PUT_DIR)))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsFtpFinalizerBean periodicJobsFtpFinalizerBean = newPeriodicJobsFtpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsFtpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsFtpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); FtpClient ftpClient = new FtpClient() .withHost("localhost") .withPort(fakeFtpServer.getServerControlPort()) @@ -117,10 +115,9 @@ public void deliver_onNonEmptyJob() throws IOException { .withFtpUser(USERNAME) .withFtpPassword(PASSWORD) .withFtpSubdirectory(PUT_DIR)))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsFtpFinalizerBean periodicJobsFtpFinalizerBean = newPeriodicJobsFtpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsFtpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsFtpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); FtpClient ftpClient = new FtpClient() .withHost("localhost") .withPort(fakeFtpServer.getServerControlPort()) @@ -146,11 +143,10 @@ public void deliver_onEmptyJob() throws IOException { .withFtpUser(USERNAME) .withFtpPassword(PASSWORD) .withFtpSubdirectory(PUT_DIR)))); - Chunk chunk = new Chunk(jobId, 0, Chunk.Type.PROCESSED); PeriodicJobsFtpFinalizerBean periodicJobsFtpFinalizerBean = newPeriodicJobsFtpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsFtpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsFtpFinalizerBean.deliver(jobId, 0, delivery, env().getEntityManager())); FtpClient ftpClient = new FtpClient() .withHost("localhost") diff --git a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsHttpFinalizerBeanIT.java b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsHttpFinalizerBeanIT.java index bf9c96e766..57b2660aac 100644 --- a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsHttpFinalizerBeanIT.java +++ b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsHttpFinalizerBeanIT.java @@ -1,7 +1,6 @@ package dk.dbc.dataio.sink.periodicjobs.pickup; import dk.dbc.dataio.commons.conversion.ConversionMetadata; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnector; import dk.dbc.dataio.commons.utils.lang.StringUtil; import dk.dbc.dataio.filestore.service.connector.FileStoreServiceConnector; @@ -90,11 +89,10 @@ public void deliver_onNonEmptyJob() throws FileStoreServiceConnectorUnexpectedSt .withSubmitterNumber("111111") .withPickup(new HttpPickup() .withReceivingAgency(String.valueOf(receivingAgency))))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsHttpFinalizerBean periodicJobsHttpFinalizerBean = newPeriodicJobsHttpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsHttpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsHttpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); InOrder orderVerifier = Mockito.inOrder(fileStoreServiceConnector); orderVerifier.verify(fileStoreServiceConnector).appendToFile(FILE_ID, block1.getBytes()); @@ -127,11 +125,10 @@ public void deliver_autoprintJob() throws FileStoreServiceConnectorUnexpectedSta .withPickup(new HttpPickup() .withReceivingAgency(String.valueOf(receivingAgency)) .withOverrideFilename("autoprint")))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsHttpFinalizerBean periodicJobsHttpFinalizerBean = newPeriodicJobsHttpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsHttpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsHttpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); ConversionMetadata expectedMetadata = new ConversionMetadata(PeriodicJobsHttpFinalizerBean.ORIGIN) .withJobId(delivery.getJobId()) @@ -153,11 +150,10 @@ public void deliver_onEmptyJob() throws FileStoreServiceConnectorUnexpectedStatu .withSubmitterNumber("111111") .withPickup(new HttpPickup() .withReceivingAgency(String.valueOf(receivingAgency))))); - Chunk chunk = new Chunk(jobId, 0, Chunk.Type.PROCESSED); PeriodicJobsHttpFinalizerBean periodicJobsHttpFinalizerBean = newPeriodicJobsHttpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsHttpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsHttpFinalizerBean.deliver(jobId, 0, delivery, env().getEntityManager())); ConversionMetadata expectedMetadata = new ConversionMetadata(PeriodicJobsHttpFinalizerBean.ORIGIN) .withJobId(delivery.getJobId()) @@ -210,11 +206,10 @@ public void deliver_file_with_header_and_footer() throws FileStoreServiceConnect .withReceivingAgency(String.valueOf(receivingAgency)) .withContentHeader("Ugekorrektur uge ${__WEEKCODE_EMO__}\n") .withContentFooter("\nslut uge ${__WEEKCODE_EMO__}")))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsHttpFinalizerBean periodicJobsHttpFinalizerBean = newPeriodicJobsHttpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsHttpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsHttpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); InOrder orderVerifier = Mockito.inOrder(fileStoreServiceConnector); orderVerifier.verify(fileStoreServiceConnector).appendToFile(FILE_ID, block1.getBytes()); diff --git a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsMailFinalizerBeanIT.java b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsMailFinalizerBeanIT.java index 62cbb1781e..ee20837878 100644 --- a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsMailFinalizerBeanIT.java +++ b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsMailFinalizerBeanIT.java @@ -1,6 +1,5 @@ package dk.dbc.dataio.sink.periodicjobs.pickup; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.types.ChunkItem; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnector; import dk.dbc.dataio.commons.utils.jobstore.ejb.JobStoreServiceConnectorBean; @@ -69,10 +68,9 @@ public void deliver_onNonEmptyJobNoDatablocks() throws AddressException { .withPickup(new MailPickup() .withRecipients(recipients) .withSubject(subject)))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsMailFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); Mailbox inbox = Mailbox.get("someone_out_there@outthere.dk"); assertThat("Inbox size", inbox.size(), is(0)); } @@ -111,10 +109,9 @@ public void deliver_onNonEmptyJob() throws IOException, MessagingException { .withRecipients(recipients) .withSubject(subject) .withRecordLimit(3)))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsMailFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); Mailbox inbox = Mailbox.get("someone_out_there@outthere.dk"); assertThat("Inbox size", inbox.size(), is(1)); Message receivedMail = inbox.get(0); @@ -136,10 +133,9 @@ public void deliver_onEmptyJob() throws IOException, MessagingException { .withPickup(new MailPickup() .withRecipients(recipients) .withSubject(subject)))); - Chunk chunk = new Chunk(jobId, 0, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsMailFinalizerBean.deliver(jobId, 0, delivery, env().getEntityManager())); Mailbox inbox = Mailbox.get("someone_out_there@outthere.dk"); assertThat("Inbox size", inbox.size(), is(1)); Message receivedMail = inbox.get(0); @@ -162,11 +158,10 @@ public void onInvalidMailRecipients() { .withPickup(new MailPickup() .withRecipients("not a valid email address") .withSubject(subject)))); - Chunk chunk = new Chunk(jobId, 0, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); - Chunk result = env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); - assertThat(result.getItems().get(0).getStatus(), is(ChunkItem.Status.FAILURE)); + ChunkItem result = env().getPersistenceContext().run(() -> + periodicJobsMailFinalizerBean.deliver(jobId, 0, delivery, env().getEntityManager())); + assertThat(result.getStatus(), is(ChunkItem.Status.FAILURE)); } @Test @@ -186,10 +181,9 @@ public void weekcodeInMailSubject() throws WeekResolverConnectorException, Messa .withPickup(new MailPickup() .withRecipients(recipients) .withSubject("mail for week ${__WEEKCODE_EMO__}")))); - Chunk chunk = new Chunk(jobId, 0, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsMailFinalizerBean.deliver(jobId, 0, delivery, env().getEntityManager())); List inbox = Mailbox.get("someone_out_there@outthere.dk"); assertThat("Inbox size", inbox.size(), is(1)); Message receivedMail = inbox.get(0); @@ -235,10 +229,9 @@ public void deliver_file_with_header_and_footer() throws MessagingException, IOE .withSubject(subject) .withContentHeader("Ugekorrektur uge ${__WEEKCODE_EMO__}\n") .withContentFooter("\nslut")))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsMailFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); List inbox = Mailbox.get("someone_out_there@outthere.dk"); assertThat("Inbox size", inbox.size(), is(1)); Message receivedMail = inbox.get(0); @@ -287,10 +280,9 @@ public void deliver_mail_as_attachment() throws MessagingException, IOException, .withMimetype("text/html") .withContentHeader("Ugekorrektur uge ${__WEEKCODE_EMO__}\n") .withContentFooter("\nslut")))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsMailFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); List inbox = Mailbox.get("someone_out_there@outthere.dk"); assertThat("Inbox size", inbox.size(), is(1)); Message receivedMail = inbox.get(0); @@ -333,10 +325,9 @@ public void deliver_mail_as_configured_body_with_attachment() .withSubject(subject) .withMimetype("text/html") .withBody("Ugekorrektur uge ${__WEEKCODE_EMO__}")))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsMailFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); List inbox = Mailbox.get("someone_out_there@outthere.dk"); assertThat("Inbox size", inbox.size(), is(1)); Message receivedMail = inbox.get(0); @@ -393,13 +384,12 @@ public void deliver_recordCountExceedsLimit() throws AddressException { .withRecipients(recipients) .withSubject(subject) .withRecordLimit(1)))); - Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); PeriodicJobsMailFinalizerBean periodicJobsMailFinalizerBean = newPeriodicJobsMailFinalizerBean(); - Chunk result = env().getPersistenceContext().run(() -> - periodicJobsMailFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + ChunkItem result = env().getPersistenceContext().run(() -> + periodicJobsMailFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); List inbox = Mailbox.get("someone_out_there@outthere.dk"); - assertThat(result.getItems().get(0).getStatus(), is(ChunkItem.Status.FAILURE)); - assertThat(result.getItems().get(0).getData(), + assertThat(result.getStatus(), is(ChunkItem.Status.FAILURE)); + assertThat(result.getData(), is("IllegalStateException: Record count exceeded record limit of 1".getBytes(StandardCharsets.UTF_8))); assertThat("Inbox size", inbox.size(), is(0)); } diff --git a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsSFtpFinalizerBeanIT.java b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsSFtpFinalizerBeanIT.java index 653f13b393..c2a7c2c5e5 100644 --- a/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsSFtpFinalizerBeanIT.java +++ b/sink/periodic-jobs/src/test/java/dk/dbc/dataio/sink/periodicjobs/pickup/PeriodicJobsSFtpFinalizerBeanIT.java @@ -1,6 +1,5 @@ package dk.dbc.dataio.sink.periodicjobs.pickup; -import dk.dbc.dataio.commons.types.Chunk; import dk.dbc.dataio.commons.utils.jobstore.JobStoreServiceConnector; import dk.dbc.dataio.commons.utils.lang.StringUtil; import dk.dbc.dataio.harvester.types.PeriodicJobsHarvesterConfig; @@ -74,10 +73,9 @@ public void deliver_onNonEmptyJobNoDataBlocks() { .withSubmitterNumber("22222222") .withTimeOfLastHarvest(new Date()) .withPickup(getPickup()))); - final Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); final PeriodicJobsSFtpFinalizerBean periodicJobsSFtpFinalizerBean = newPeriodicJobsSFtpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsSFtpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsSFtpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); final String fileName = String.format("%s/periodisk-job-%d.data", testDir, jobId); assertThat("File is NOT present", fakeSFtpServer.existsFile(fileName), is(false)); } @@ -94,10 +92,9 @@ public void deliver_onNonEmptyJob() throws IOException { .withSubmitterNumber("111111") .withTimeOfLastHarvest(new Date()) .withPickup(getPickup()))); - final Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); final PeriodicJobsSFtpFinalizerBean periodicJobsSFtpFinalizerBean = newPeriodicJobsSFtpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsSFtpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsSFtpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); String dataSentUsingSFtp = fakeSFtpServer.getFileContent(String .format("%s/%s", testDir, String. @@ -118,10 +115,9 @@ public void deliver_fileWithOverrideFilename() throws IOException { .withTimeOfLastHarvest(new Date()) .withPickup(getPickup() .withOverrideFilename("testMyNewFileName.data")))); - final Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); final PeriodicJobsSFtpFinalizerBean periodicJobsSFtpFinalizerBean = newPeriodicJobsSFtpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsSFtpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsSFtpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); String dataSentUsingSFtp = fakeSFtpServer.getFileContent( String.format("%s/%s", testDir, "testMyNewFileName.data"), StandardCharsets.UTF_8); @@ -147,10 +143,9 @@ public void deliver_fileWithHeaderAndFooter() throws IOException, WeekResolverCo .withOverrideFilename("testMyNewFileName${__WEEKCODE_EMO__}.data") .withContentHeader("Ugekorrektur uge ${__WEEKCODE_EMO__}\n") .withContentFooter("\nslut uge ${__WEEKCODE_EMO__}")))); - final Chunk chunk = new Chunk(jobId, 3, Chunk.Type.PROCESSED); final PeriodicJobsSFtpFinalizerBean periodicJobsSFtpFinalizerBean = newPeriodicJobsSFtpFinalizerBean(); env().getPersistenceContext().run(() -> - periodicJobsSFtpFinalizerBean.deliver(chunk, delivery, env().getEntityManager())); + periodicJobsSFtpFinalizerBean.deliver(jobId, 3, delivery, env().getEntityManager())); String dataSentUsingSFtp = fakeSFtpServer.getFileContent( String.format("%s/%s", testDir, "testMyNewFileName202041.data"), StandardCharsets.UTF_8);