diff --git a/multiapps-controller-persistence/src/main/java/org/cloudfoundry/multiapps/controller/persistence/model/AsyncUploadJobEntry.java b/multiapps-controller-persistence/src/main/java/org/cloudfoundry/multiapps/controller/persistence/model/AsyncUploadJobEntry.java index a74c152602..527d90441b 100644 --- a/multiapps-controller-persistence/src/main/java/org/cloudfoundry/multiapps/controller/persistence/model/AsyncUploadJobEntry.java +++ b/multiapps-controller-persistence/src/main/java/org/cloudfoundry/multiapps/controller/persistence/model/AsyncUploadJobEntry.java @@ -1,7 +1,9 @@ package org.cloudfoundry.multiapps.controller.persistence.model; import java.text.MessageFormat; +import java.time.Duration; import java.time.LocalDateTime; +import java.time.ZoneOffset; import org.cloudfoundry.multiapps.common.Nullable; import org.immutables.value.Value; @@ -11,6 +13,10 @@ public interface AsyncUploadJobEntry { String STALE_JOB_DETAILS_FORMAT = "Stale job details - id: {0}, state: {1}, updatedAt: {2}, addedAt: {3}, startedAt: {4}, bytesRead: {5}, url: {6}, space: {7}, namespace: {8}, user: {9}, instance: {10}"; + String ASYNC_UPLOAD_JOB_SUMMARY_FORMAT = "id: {0}, state: {1}, fileId: {2}, mtaId: {3}, schemaVersion: {4}, instanceIndex: {5}, bytesRead: {6}, addedAt: {7}, startedAt: {8}, finishedAt: {9}, updatedAt: {10}, queueWaitTime: {11}, uploadDuration: {12}, totalTime: {13}, error: {14}"; + + String NOT_AVAILABLE = "N/A"; + enum State { INITIAL, RUNNING, FINISHED, ERROR } @@ -61,4 +67,40 @@ default String buildStaleDetailsLogMessage() { return MessageFormat.format(STALE_JOB_DETAILS_FORMAT, getId(), getState(), getUpdatedAt(), getAddedAt(), getStartedAt(), getBytesRead(), getUrl(), getSpaceGuid(), getNamespace(), getUser(), getInstanceIndex()); } + + default String buildLogSummary() { + return MessageFormat.format(ASYNC_UPLOAD_JOB_SUMMARY_FORMAT, getId(), getState(), getFileId(), getMtaId(), getSchemaVersion(), + getInstanceIndex(), getBytesRead(), getAddedAt(), getStartedAt(), getFinishedAt(), getUpdatedAt(), + formatDuration(getQueueWaitTime()), formatDuration(getUploadDuration()), formatDuration(getTotalTime()), + getError()); + } + + @Nullable + default Duration getQueueWaitTime() { + return durationBetween(getAddedAt(), getStartedAt()); + } + + @Nullable + default Duration getUploadDuration() { + return durationBetween(getStartedAt(), getFinishedAt()); + } + + @Nullable + default Duration getTotalTime() { + return durationBetween(getAddedAt(), getFinishedAt()); + } + + private static Duration durationBetween(LocalDateTime start, LocalDateTime end) { + if (start == null || end == null) { + return null; + } + return Duration.between(start.toInstant(ZoneOffset.UTC), end.toInstant(ZoneOffset.UTC)); + } + + private static String formatDuration(Duration duration) { + if (duration == null) { + return NOT_AVAILABLE; + } + return duration.toMillis() + " ms"; + } } diff --git a/multiapps-controller-persistence/src/test/java/org/cloudfoundry/multiapps/controller/persistence/model/AsyncUploadJobEntryTest.java b/multiapps-controller-persistence/src/test/java/org/cloudfoundry/multiapps/controller/persistence/model/AsyncUploadJobEntryTest.java new file mode 100644 index 0000000000..b756636c1c --- /dev/null +++ b/multiapps-controller-persistence/src/test/java/org/cloudfoundry/multiapps/controller/persistence/model/AsyncUploadJobEntryTest.java @@ -0,0 +1,111 @@ +package org.cloudfoundry.multiapps.controller.persistence.model; + +import java.time.Duration; +import java.time.LocalDateTime; +import java.time.Month; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +class AsyncUploadJobEntryTest { + + private static final String JOB_ID = "job-id"; + private static final String FILE_ID = "file-id"; + private static final String SPACE_GUID = "space-guid"; + private static final String USER = "user"; + private static final String URL = "https://user:secret@example.com/my.mtar"; + private static final LocalDateTime ADDED_AT = LocalDateTime.of(2026, Month.AUGUST, 3, 10, 0, 0); + private static final LocalDateTime STARTED_AT = ADDED_AT.plusSeconds(30); + private static final LocalDateTime FINISHED_AT = STARTED_AT.plusSeconds(90); + + @Test + void testTimingsForFinishedJob() { + AsyncUploadJobEntry job = baseBuilder().state(AsyncUploadJobEntry.State.FINISHED) + .addedAt(ADDED_AT) + .startedAt(STARTED_AT) + .finishedAt(FINISHED_AT) + .build(); + + Assertions.assertEquals(Duration.ofSeconds(30), job.getQueueWaitTime()); + Assertions.assertEquals(Duration.ofSeconds(90), job.getUploadDuration()); + Assertions.assertEquals(Duration.ofSeconds(120), job.getTotalTime()); + } + + @Test + void testTimingsAreNullWhenTimestampsMissing() { + AsyncUploadJobEntry runningJob = baseBuilder().state(AsyncUploadJobEntry.State.RUNNING) + .addedAt(ADDED_AT) + .startedAt(STARTED_AT) + .build(); + + Assertions.assertEquals(Duration.ofSeconds(30), runningJob.getQueueWaitTime()); + Assertions.assertNull(runningJob.getUploadDuration()); + Assertions.assertNull(runningJob.getTotalTime()); + + AsyncUploadJobEntry initialJob = baseBuilder().state(AsyncUploadJobEntry.State.INITIAL) + .addedAt(ADDED_AT) + .build(); + + Assertions.assertNull(initialJob.getQueueWaitTime()); + Assertions.assertNull(initialJob.getUploadDuration()); + Assertions.assertNull(initialJob.getTotalTime()); + } + + @Test + void testLogSafeSummaryContainsTimingsAndFileId() { + AsyncUploadJobEntry job = baseBuilder().state(AsyncUploadJobEntry.State.FINISHED) + .mtaId("my-mta") + .bytesRead(2048L) + .addedAt(ADDED_AT) + .startedAt(STARTED_AT) + .finishedAt(FINISHED_AT) + .build(); + + String summary = job.buildLogSummary(); + + Assertions.assertTrue(summary.contains(JOB_ID), summary); + Assertions.assertTrue(summary.contains(FILE_ID), summary); + Assertions.assertTrue(summary.contains("my-mta"), summary); + Assertions.assertTrue(summary.contains("queueWaitTime: 30000 ms"), summary); + Assertions.assertTrue(summary.contains("uploadDuration: 90000 ms"), summary); + Assertions.assertTrue(summary.contains("totalTime: 120000 ms"), summary); + } + + @Test + void testLogSafeSummaryHidesSensitiveData() { + AsyncUploadJobEntry job = baseBuilder().state(AsyncUploadJobEntry.State.RUNNING) + .addedAt(ADDED_AT) + .startedAt(STARTED_AT) + .build(); + + String summary = job.buildLogSummary(); + + Assertions.assertFalse(summary.contains(URL), summary); + Assertions.assertFalse(summary.contains("secret"), summary); + Assertions.assertFalse(summary.contains(USER), summary); + Assertions.assertFalse(summary.contains(SPACE_GUID), summary); + } + + @Test + void testLogSafeSummaryRendersMissingTimingsAsNotAvailable() { + AsyncUploadJobEntry job = baseBuilder().state(AsyncUploadJobEntry.State.INITIAL) + .addedAt(ADDED_AT) + .build(); + + String summary = job.buildLogSummary(); + + Assertions.assertTrue(summary.contains("queueWaitTime: " + AsyncUploadJobEntry.NOT_AVAILABLE), summary); + Assertions.assertTrue(summary.contains("uploadDuration: " + AsyncUploadJobEntry.NOT_AVAILABLE), summary); + Assertions.assertTrue(summary.contains("totalTime: " + AsyncUploadJobEntry.NOT_AVAILABLE), summary); + } + + private ImmutableAsyncUploadJobEntry.Builder baseBuilder() { + return ImmutableAsyncUploadJobEntry.builder() + .id(JOB_ID) + .fileId(FILE_ID) + .user(USER) + .url(URL) + .spaceGuid(SPACE_GUID) + .instanceIndex(0); + } +} diff --git a/multiapps-controller-process/src/main/java/org/cloudfoundry/multiapps/controller/process/Messages.java b/multiapps-controller-process/src/main/java/org/cloudfoundry/multiapps/controller/process/Messages.java index c2589ad1f6..ac8279c2f8 100755 --- a/multiapps-controller-process/src/main/java/org/cloudfoundry/multiapps/controller/process/Messages.java +++ b/multiapps-controller-process/src/main/java/org/cloudfoundry/multiapps/controller/process/Messages.java @@ -822,6 +822,10 @@ public class Messages { public static final String IGNORING_NOT_FOUND_OPTIONAL_SERVICE = "Service {0} not found but is optional"; public static final String IGNORING_NOT_FOUND_INACTIVE_SERVICE = "Service {0} not found but is inactive"; public static final String MISSING_REQUIRED_0_CREDENTIAL_FROM_SCL_EXPORT = "Missing required {0} credential for SAP Cloud Logging export"; + + public static final String ASYNC_UPLOAD_JOB_FOR_OPERATION_0_IS_1 = "Async upload job for operation \"{0}\" - {1}"; + public static final String COULD_NOT_LOG_ASYNC_UPLOAD_JOBS_FOR_OPERATION_0 = "Could not log async upload jobs for operation \"{0}\""; + // Not log messages public static final String SERVICE_TYPE = "{0}/{1}"; public static final String PARSE_NULL_STRING_ERROR = "Cannot parse null string"; diff --git a/multiapps-controller-process/src/main/java/org/cloudfoundry/multiapps/controller/process/listeners/StartProcessListener.java b/multiapps-controller-process/src/main/java/org/cloudfoundry/multiapps/controller/process/listeners/StartProcessListener.java index 9d12282205..60b9c79814 100644 --- a/multiapps-controller-process/src/main/java/org/cloudfoundry/multiapps/controller/process/listeners/StartProcessListener.java +++ b/multiapps-controller-process/src/main/java/org/cloudfoundry/multiapps/controller/process/listeners/StartProcessListener.java @@ -18,8 +18,10 @@ import org.cloudfoundry.multiapps.controller.api.model.ProcessType; import org.cloudfoundry.multiapps.controller.core.util.ApplicationConfiguration; import org.cloudfoundry.multiapps.controller.core.util.LoggingUtil; +import org.cloudfoundry.multiapps.controller.persistence.model.AsyncUploadJobEntry; import org.cloudfoundry.multiapps.controller.persistence.model.HistoricOperationEvent; import org.cloudfoundry.multiapps.controller.persistence.model.ImmutableHistoricOperationEvent; +import org.cloudfoundry.multiapps.controller.persistence.services.AsyncUploadJobService; import org.cloudfoundry.multiapps.controller.persistence.services.FileService; import org.cloudfoundry.multiapps.controller.persistence.services.FileStorageException; import org.cloudfoundry.multiapps.controller.persistence.services.HistoricOperationEventService; @@ -56,6 +58,7 @@ public class StartProcessListener extends AbstractProcessExecutionListener { private final ProcessTypeToOperationMetadataMapper operationMetadataMapper; private final DynatracePublisher dynatracePublisher; private final FileService fileService; + private final AsyncUploadJobService asyncUploadJobService; @Inject public StartProcessListener(ProgressMessageService progressMessageService, StepLogger.Factory stepLoggerFactory, @@ -64,7 +67,7 @@ public StartProcessListener(ProgressMessageService progressMessageService, StepL ApplicationConfiguration configuration, ProcessTypeParser processTypeParser, OperationService operationService, ProcessTypeToOperationMetadataMapper operationMetadataMapper, DynatracePublisher dynatracePublisher, FileService fileService, - OperationLogsExporter operationLogsExporter) { + AsyncUploadJobService asyncUploadJobService, OperationLogsExporter operationLogsExporter) { super(progressMessageService, stepLoggerFactory, processLoggerProvider, @@ -78,6 +81,7 @@ public StartProcessListener(ProgressMessageService progressMessageService, StepL this.operationMetadataMapper = operationMetadataMapper; this.dynatracePublisher = dynatracePublisher; this.fileService = fileService; + this.asyncUploadJobService = asyncUploadJobService; } @Override @@ -94,6 +98,7 @@ protected void notifyInternal(DelegateExecution execution) { } updateOperationFiles(execution, correlationId); + logAsyncUploadJobs(execution, correlationId); getHistoricOperationEventService().add(ImmutableHistoricOperationEvent.of(correlationId, HistoricOperationEvent.EventType.STARTED)); logProcessEnvironment(); logProcessVariables(execution, processType, correlationId); @@ -160,6 +165,26 @@ private void updateOperationFiles(DelegateExecution execution, String correlatio } } + private void logAsyncUploadJobs(DelegateExecution execution, String correlationId) { + List operationFileIds = OperationFileIdsUtil.getOperationFileIds(execution); + if (operationFileIds.isEmpty()) { + return; + } + try { + List asyncUploadJobs = asyncUploadJobService.createQuery() + .withFileIds(operationFileIds) + .list(); + for (AsyncUploadJobEntry asyncUploadJob : asyncUploadJobs) { + if (LOGGER.isInfoEnabled()) { + LOGGER.info(MessageFormat.format(Messages.ASYNC_UPLOAD_JOB_FOR_OPERATION_0_IS_1, correlationId, + asyncUploadJob.buildLogSummary())); + } + } + } catch (Exception e) { + LOGGER.warn(MessageFormat.format(Messages.COULD_NOT_LOG_ASYNC_UPLOAD_JOBS_FOR_OPERATION_0, correlationId), e); + } + } + private void publishDynatraceEvent(DelegateExecution execution, ProcessType processType, String correlationId) { DynatraceProcessEvent startEvent = ImmutableDynatraceProcessEvent.builder() .processId(correlationId) diff --git a/multiapps-controller-process/src/test/java/org/cloudfoundry/multiapps/controller/process/listeners/StartProcessListenerTest.java b/multiapps-controller-process/src/test/java/org/cloudfoundry/multiapps/controller/process/listeners/StartProcessListenerTest.java index 213c66cf47..3421e98f7f 100644 --- a/multiapps-controller-process/src/test/java/org/cloudfoundry/multiapps/controller/process/listeners/StartProcessListenerTest.java +++ b/multiapps-controller-process/src/test/java/org/cloudfoundry/multiapps/controller/process/listeners/StartProcessListenerTest.java @@ -14,7 +14,11 @@ import org.cloudfoundry.multiapps.controller.api.model.Operation; import org.cloudfoundry.multiapps.controller.api.model.ProcessType; import org.cloudfoundry.multiapps.controller.core.util.ApplicationConfiguration; +import org.cloudfoundry.multiapps.controller.persistence.model.AsyncUploadJobEntry; +import org.cloudfoundry.multiapps.controller.persistence.model.ImmutableAsyncUploadJobEntry; import org.cloudfoundry.multiapps.controller.persistence.query.OperationQuery; +import org.cloudfoundry.multiapps.controller.persistence.query.AsyncUploadJobsQuery; +import org.cloudfoundry.multiapps.controller.persistence.services.AsyncUploadJobService; import org.cloudfoundry.multiapps.controller.persistence.services.FileService; import org.cloudfoundry.multiapps.controller.persistence.services.FileStorageException; import org.cloudfoundry.multiapps.controller.persistence.services.HistoricOperationEventService; @@ -35,6 +39,7 @@ import org.cloudfoundry.multiapps.controller.process.variables.Variables; import org.flowable.engine.delegate.DelegateExecution; import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; @@ -47,6 +52,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals; import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; class StartProcessListenerTest { @@ -79,6 +86,10 @@ class StartProcessListenerTest { private HistoricOperationEventService historicOperationEventService; @Mock private FileService fileService; + @Mock(answer = Answers.RETURNS_SELF) + private AsyncUploadJobService asyncUploadJobService; + @Mock(answer = Answers.RETURNS_SELF) + private AsyncUploadJobsQuery asyncUploadJobsQuery; @Spy private ProcessTypeToOperationMetadataMapper operationMetadataMapper; @Mock @@ -118,6 +129,7 @@ void setUp() throws Exception { operationMetadataMapper, dynatracePublisher, fileService, + asyncUploadJobService, operationLogsExporter); } @@ -134,25 +146,60 @@ void testVerify(String processInstanceId, ProcessType processType) throws Except verifyDynatracePublishEvent(); } + @Test + void testAssociatedAsyncUploadJobsAreLogged() { + this.processType = ProcessType.DEPLOY; + this.processInstanceId = "process-instance-id"; + prepare(); + AsyncUploadJobEntry asyncUploadJob = ImmutableAsyncUploadJobEntry.builder() + .id("async-job-id") + .state(AsyncUploadJobEntry.State.FINISHED) + .user(USER) + .url("https://example.com/my.mtar") + .spaceGuid(SPACE_ID) + .fileId(APP_ARCHIVE_IDS.split(",")[0]) + .instanceIndex(0) + .build(); + when(asyncUploadJobsQuery.list()).thenReturn(List.of(asyncUploadJob)); + + listener.notify(execution); + + List expectedFileIds = List.of(ArrayUtils.addAll(APP_ARCHIVE_IDS.split(","), EXT_DESCRIPTOR_IDS.split(","))); + verify(asyncUploadJobsQuery) + .withFileIds(expectedFileIds); + verify(asyncUploadJobsQuery) + .list(); + } + + @Test + void testFailureToLogAsyncUploadJobsDoesNotBreakOperation() throws Exception { + this.processType = ProcessType.DEPLOY; + this.processInstanceId = "process-instance-id"; + prepare(); + when(asyncUploadJobService.createQuery()).thenThrow(new IllegalStateException("Async upload job store is unavailable")); + + listener.notify(execution); + + verifyOperationFilesAreUpdated(); + verifyDynatracePublishEvent(); + } + private void prepare() { prepareContext(); - Mockito.when(stepLoggerFactory.create(any(), any(), any(), any(), any())) - .thenReturn(stepLogger); - Mockito.when(operationService.createQuery()) - .thenReturn(operationQuery); + when(stepLoggerFactory.create(any(), any(), any(), any(), any())).thenReturn(stepLogger); + when(operationService.createQuery()).thenReturn(operationQuery); Mockito.doReturn(null) .when(operationQuery) .singleResult(); + when(asyncUploadJobService.createQuery()).thenReturn(asyncUploadJobsQuery); + when(asyncUploadJobsQuery.list()).thenReturn(Collections.emptyList()); } private void prepareContext() { listener.currentTimeSupplier = currentTimeSupplier; - Mockito.when(execution.getProcessInstanceId()) - .thenReturn(processInstanceId); - Mockito.when(execution.getVariables()) - .thenReturn(Collections.emptyMap()); - Mockito.when(processTypeParser.getProcessType(execution)) - .thenReturn(processType); + when(execution.getProcessInstanceId()).thenReturn(processInstanceId); + when(execution.getVariables()).thenReturn(Collections.emptyMap()); + when(processTypeParser.getProcessType(execution)).thenReturn(processType); VariableHandling.set(execution, Variables.SPACE_GUID, SPACE_ID); VariableHandling.set(execution, Variables.MTA_ID, MTA_ID); VariableHandling.set(execution, Variables.USER, USER); @@ -178,13 +225,13 @@ private void verifyOperationInsertion() throws SLException { private void verifyOperationFilesAreUpdated() throws FileStorageException { List expectedFileIds = List.of(ArrayUtils.addAll(APP_ARCHIVE_IDS.split(","), EXT_DESCRIPTOR_IDS.split(","))); - Mockito.verify(fileService) + verify(fileService) .updateFilesOperationId(Mockito.eq(expectedFileIds), Mockito.anyString()); } private void verifyDynatracePublishEvent() { ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(DynatraceProcessEvent.class); - Mockito.verify(dynatracePublisher) + verify(dynatracePublisher) .publishProcessEvent(argumentCaptor.capture(), Mockito.any()); DynatraceProcessEvent actualDynatraceEvent = argumentCaptor.getValue(); assertEquals(MTA_ID, actualDynatraceEvent.getMtaId());