From 3441745ff4fe0af5a383a8be563a44d09c991cd7 Mon Sep 17 00:00:00 2001 From: dan-s1 Date: Fri, 4 Sep 2026 18:10:27 +0000 Subject: [PATCH] NIFI-16195 Wrote part headers into attributes --- .../email/ExtractEmailAttachments.java | 39 ++++++++++++++++--- .../email/TestExtractEmailAttachments.java | 16 +++++++- 2 files changed, 49 insertions(+), 6 deletions(-) diff --git a/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/main/java/org/apache/nifi/processors/email/ExtractEmailAttachments.java b/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/main/java/org/apache/nifi/processors/email/ExtractEmailAttachments.java index c2763df5c0cb..e8ca2940d438 100644 --- a/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/main/java/org/apache/nifi/processors/email/ExtractEmailAttachments.java +++ b/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/main/java/org/apache/nifi/processors/email/ExtractEmailAttachments.java @@ -19,6 +19,7 @@ import jakarta.activation.DataSource; import jakarta.mail.Address; import jakarta.mail.BodyPart; +import jakarta.mail.Header; import jakarta.mail.MessagingException; import jakarta.mail.Multipart; import jakarta.mail.Session; @@ -46,6 +47,7 @@ import java.io.IOException; import java.io.InputStream; import java.util.ArrayList; +import java.util.Enumeration; import java.util.HashMap; import java.util.List; import java.util.Map; @@ -61,7 +63,8 @@ @WritesAttribute(attribute = "filename ", description = "The filename of the attachment"), @WritesAttribute(attribute = "email.attachment.parent.filename ", description = "The filename of the parent FlowFile"), @WritesAttribute(attribute = "email.attachment.parent.uuid", description = "The UUID of the original FlowFile."), - @WritesAttribute(attribute = "mime.type", description = "The mime type of the attachment.")}) + @WritesAttribute(attribute = "mime.type", description = "The mime type of the attachment."), + @WritesAttribute(attribute = ExtractEmailAttachments.ATTACHMENT_HEADER_ATTRIBUTE_PREFIX + "", description = "Attachment header.")}) public class ExtractEmailAttachments extends AbstractProcessor { public static final String ATTACHMENT_ORIGINAL_FILENAME = "email.attachment.parent.filename"; @@ -80,6 +83,7 @@ public class ExtractEmailAttachments extends AbstractProcessor { .description("FlowFiles that could not be parsed") .build(); + static final String ATTACHMENT_HEADER_ATTRIBUTE_PREFIX = "email.attachment.header."; private static final String ATTACHMENT_DISPOSITION = "attachment"; private static final Set RELATIONSHIPS = Set.of( @@ -117,12 +121,13 @@ public void onTrigger(final ProcessContext context, final ProcessSession session final String originalFlowFileName = originalFlowFile.getAttribute(CoreAttributes.FILENAME.key()); try { - final List attachments = new ArrayList<>(); + final List attachments = new ArrayList<>(); parseAttachments(attachments, originalMessage, 0); - for (final DataSource data : attachments) { + for (final Attachment attachment : attachments) { FlowFile split = session.create(originalFlowFile); final Map attributes = new HashMap<>(); + final DataSource data = attachment.dataSource(); final String name = data.getName(); if (name != null && !name.isBlank()) { attributes.put(CoreAttributes.FILENAME.key(), name); @@ -131,6 +136,13 @@ public void onTrigger(final ProcessContext context, final ProcessSession session if (contentType != null && !contentType.isBlank()) { attributes.put(CoreAttributes.MIME_TYPE.key(), contentType); } + + for (Map.Entry entry : attachment.headers().entrySet()) { + final String headerAttributeName = ATTACHMENT_HEADER_ATTRIBUTE_PREFIX + entry.getKey(); + final String headerAttributeValue = entry.getValue(); + attributes.put(headerAttributeName, headerAttributeValue); + } + String parentUuid = originalFlowFile.getAttribute(CoreAttributes.UUID.key()); attributes.put(ATTACHMENT_ORIGINAL_UUID, parentUuid); attributes.put(ATTACHMENT_ORIGINAL_FILENAME, originalFlowFileName); @@ -176,7 +188,7 @@ public Set getRelationships() { return RELATIONSHIPS; } - private void parseAttachments(final List attachments, final MimePart parentPart, final int depth) throws MessagingException, IOException { + private void parseAttachments(final List attachments, final MimePart parentPart, final int depth) throws MessagingException, IOException { final String disposition = parentPart.getDisposition(); final Object parentContent = parentPart.getContent(); @@ -191,7 +203,24 @@ private void parseAttachments(final List attachments, final MimePart } } else if (ATTACHMENT_DISPOSITION.equalsIgnoreCase(disposition) || depth > 0) { final DataSource dataSource = parentPart.getDataHandler().getDataSource(); - attachments.add(dataSource); + final Map extractedHeaders = new HashMap<>(); + + if (parentPart instanceof final MimeBodyPart mimeBodyPart) { + final Enumeration
headers = mimeBodyPart.getAllHeaders(); + while (headers.hasMoreElements()) { + final Header header = headers.nextElement(); + final String name = header.getName(); + if (name != null && !name.isBlank()) { + final String value = header.getValue(); + extractedHeaders.put(name, value); + } + } + } + + final Attachment attachment = new Attachment(dataSource, extractedHeaders); + attachments.add(attachment); } } } + +record Attachment(DataSource dataSource, Map headers) { } diff --git a/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/test/java/org/apache/nifi/processors/email/TestExtractEmailAttachments.java b/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/test/java/org/apache/nifi/processors/email/TestExtractEmailAttachments.java index 6707290353d2..727cc0af2c55 100644 --- a/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/test/java/org/apache/nifi/processors/email/TestExtractEmailAttachments.java +++ b/nifi-extension-bundles/nifi-email-bundle/nifi-email-processors/src/test/java/org/apache/nifi/processors/email/TestExtractEmailAttachments.java @@ -51,7 +51,9 @@ public void testValidEmailWithAttachments() { runner.assertTransferCount(ExtractEmailAttachments.REL_ATTACHMENTS, 1); // Have a look at the attachments... final List splits = runner.getFlowFilesForRelationship(ExtractEmailAttachments.REL_ATTACHMENTS); - splits.get(0).assertAttributeEquals("filename", "pom.xml-0"); + final MockFlowFile split = splits.getFirst(); + split.assertAttributeEquals("filename", "pom.xml-0"); + assertAttachmentHeaderAttributes(split); } @Test @@ -76,6 +78,10 @@ public void testValidEmailWithMultipleAttachments() { } assertTrue(filenames.containsAll(Arrays.asList("pom.xml-0", "pom.xml-1", "pom.xml-2"))); + + for (MockFlowFile split : splits) { + assertAttachmentHeaderAttributes(split); + } } @Test @@ -102,4 +108,12 @@ public void testInvalidEmail() { runner.assertTransferCount(ExtractEmailAttachments.REL_FAILURE, 1); runner.assertTransferCount(ExtractEmailAttachments.REL_ATTACHMENTS, 0); } + + private void assertAttachmentHeaderAttributes(MockFlowFile split) { + final String regex = "^" + ExtractEmailAttachments.ATTACHMENT_HEADER_ATTRIBUTE_PREFIX + ".+"; + final boolean match = split.getAttributes().keySet().stream() + .anyMatch(attribute -> attribute.matches(regex)); + + assertTrue(match, "FlowFile did not have any attributes which began with " + ExtractEmailAttachments.ATTACHMENT_HEADER_ATTRIBUTE_PREFIX); + } }