Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
54 changes: 46 additions & 8 deletions funcake_dags/scripts/transform.sh
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ rm -f /tmp/all-identifiers-$DAG_ID.*

SAXON_VERSION=9.9.1-5
SAXON_DOWNLOAD_SHA1=c1f413a1b810dbf0d673ffd3b27c8829a82ac31c
SAXON_CP=/tmp/saxon/saxon-$SAXON_VERSION.jar
SAXON_CP=${SAXON_CP:-/tmp/saxon/saxon-$SAXON_VERSION.jar}

if [ ! -f $SAXON_CP ]; then
mkdir -p /tmp/saxon && \
Expand All @@ -41,8 +41,35 @@ fi

TOTAL_TRANSFORMED=0
RESP=`aws s3api list-objects --bucket $BUCKET --prefix ${DAG_ID}/${DAG_TS}/${SOURCE}`
for SOURCE_XML in `echo $RESP | jq -r '.Contents[].Key'`
OBJECT_COUNT=$(
printf "%s\n" "$RESP" |
jq -r '(.Contents // []) | map(select(.Key | endswith("/") | not)) | length'
)

if [ "$OBJECT_COUNT" -eq 0 ]; then
echo "No source files found at s3://$BUCKET/${DAG_ID}/${DAG_TS}/${SOURCE}" >&2
exit 1
fi

OBJECTS=$(
printf "%s\n" "$RESP" |
jq -r '(.Contents // [])[] | select(.Key | endswith("/") | not) | [.Key, (.Size | tostring)] | @tsv'
)

SKIPPED_EMPTY_FILES=0
PROCESSED_FILES=0

while IFS=$'\t' read -r SOURCE_XML SOURCE_SIZE
do
[ -n "${SOURCE_XML:-}" ] || continue

if [ "${SOURCE_SIZE:-0}" -eq 0 ]; then
echo "Skipping empty source file: $SOURCE_XML"
SKIPPED_EMPTY_FILES=$((SKIPPED_EMPTY_FILES + 1))
continue
fi

PROCESSED_FILES=$((PROCESSED_FILES + 1))
SOURCE_URL=$(aws s3 presign s3://$BUCKET/$SOURCE_XML)
echo Reading from $SOURCE_URL

Expand All @@ -56,22 +83,33 @@ do
echo "</collection>" >> $SOURCE_XML-2.xml

java -jar $SAXON_CP -xsl:$SCRIPTS_PATH/batch-transform.xsl -s:$SOURCE_XML-2.xml -o:$SOURCE_XML-transformed.xml -t
COUNT=$(grep -o "<oai_dc:dc" "$SOURCE_XML-transformed.xml" | wc -l || echo 0)
COUNT=$(grep -o "<oai_dc:dc" "$SOURCE_XML-transformed.xml" | wc -l || :)
TOTAL_TRANSFORMED=$((TOTAL_TRANSFORMED + COUNT))
aws s3 cp $SOURCE_XML-transformed.xml s3://$BUCKET/$TRANSFORM_XML

TEMPFILE=$(mktemp /tmp/identifier-output-$DAG_ID.XXXXXX)
grep "^<dcterms:identifier>\|</dcterms:identifier>$" "$SOURCE_XML-transformed.xml" >> "$TEMPFILE" || true
done
done <<< "$OBJECTS"

IDENTIFIER_FILE=$(mktemp /tmp/all-identifiers-$DAG_ID.XXXXXX)
for file in /tmp/identifier-output-$DAG_ID.*;
do
sort --u $file
done | sort -u > $IDENTIFIER_FILE
shopt -s nullglob
IDENTIFIER_FILES=(/tmp/identifier-output-$DAG_ID.*)

if [ ${#IDENTIFIER_FILES[@]} -gt 0 ]; then
for file in "${IDENTIFIER_FILES[@]}"
do
sort --u "$file"
done | sort -u > "$IDENTIFIER_FILE"
else
: > "$IDENTIFIER_FILE"
fi

shopt -u nullglob

UNIQUE_RECORD_COUNT=$(wc -l < "$IDENTIFIER_FILE")


echo "Total Records transformed: $TOTAL_TRANSFORMED"
echo "Unique Record Count: $UNIQUE_RECORD_COUNT"
echo "Files transformed: $PROCESSED_FILES"
echo "Empty files skipped: $SKIPPED_EMPTY_FILES"
229 changes: 229 additions & 0 deletions tests/transform_script_test.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,229 @@
import os
import shutil
import stat
import subprocess
import tempfile
import textwrap
import unittest
from pathlib import Path


REPO_ROOT = Path(__file__).resolve().parents[1]
SCRIPT_PATH = REPO_ROOT / "funcake_dags" / "scripts" / "transform.sh"


class TransformScriptTest(unittest.TestCase):
def test_skips_zero_byte_s3_objects(self):
result, uploads, transformed_files = self._run_script(
'{"Contents":[{"Key":"funcake_test/2021-03-23_17-25-10/new-updated-filtered/empty.xml","Size":0}]}'
)

self.assertEqual(result.returncode, 0, msg=result.stderr)
self.assertIn("Skipping empty source file", result.stdout)
self.assertIn("Files transformed: 0", result.stdout)
self.assertIn("Empty files skipped: 1", result.stdout)
self.assertEqual(uploads, "")
self.assertEqual(transformed_files, [])

def test_fails_when_no_s3_objects_are_found(self):
result, uploads, transformed_files = self._run_script('{"Contents":[]}')

self.assertNotEqual(result.returncode, 0)
self.assertIn("No source files found", result.stderr)
self.assertEqual(uploads, "")
self.assertEqual(transformed_files, [])

def test_fails_when_listing_only_contains_directory_markers(self):
result, uploads, transformed_files = self._run_script(
'{"Contents":[{"Key":"funcake_test/2021-03-23_17-25-10/new-updated-filtered/","Size":0}]}'
)

self.assertNotEqual(result.returncode, 0)
self.assertIn("No source files found", result.stderr)
self.assertEqual(uploads, "")
self.assertEqual(transformed_files, [])

def test_processes_non_empty_objects_and_skips_empty_ones(self):
listing = """{
"Contents": [
{"Key":"funcake_test/2021-03-23_17-25-10/new-updated-filtered/file1.xml","Size":100},
{"Key":"funcake_test/2021-03-23_17-25-10/new-updated-filtered/empty.xml","Size":0},
{"Key":"funcake_test/2021-03-23_17-25-10/new-updated-filtered/file2.xml","Size":200}
]
}"""
result, uploads, transformed_files = self._run_script(listing, java_mode="transform")

self.assertEqual(result.returncode, 0, msg=result.stderr)
self.assertIn("Files transformed: 2", result.stdout)
self.assertIn("Empty files skipped: 1", result.stdout)
self.assertEqual(
uploads.strip().splitlines(),
[
"s3://test-bucket/funcake_test/2021-03-23_17-25-10/transformed/file1.xml",
"s3://test-bucket/funcake_test/2021-03-23_17-25-10/transformed/file2.xml",
],
)
self.assertEqual(len(transformed_files), 2)

def test_handles_transformed_files_with_zero_records(self):
listing = """{
"Contents": [
{"Key":"funcake_test/2021-03-23_17-25-10/new-updated-filtered/file1.xml","Size":100}
]
}"""
result, uploads, transformed_files = self._run_script(listing, java_mode="transform_zero_records")

self.assertEqual(result.returncode, 0, msg=result.stderr)
self.assertIn("Total Records transformed: 0", result.stdout)
self.assertIn("Files transformed: 1", result.stdout)
self.assertIn("Empty files skipped: 0", result.stdout)
self.assertEqual(
uploads.strip().splitlines(),
["s3://test-bucket/funcake_test/2021-03-23_17-25-10/transformed/file1.xml"],
)
self.assertEqual(len(transformed_files), 1)

def _run_script(self, listing_json: str, java_mode: str = "fail"):
tempdir = tempfile.mkdtemp()
self.addCleanup(shutil.rmtree, tempdir, ignore_errors=True)
tempdir_path = Path(tempdir)
bin_dir = tempdir_path / "bin"
bin_dir.mkdir()
uploads_log = tempdir_path / "uploads.log"
uploads_log.write_text("", encoding="utf-8")
saxon_jar = tempdir_path / "saxon.jar"
saxon_jar.write_text("fake-jar", encoding="utf-8")

self._write_executable(
bin_dir / "aws",
f"""#!/usr/bin/env bash
set -euo pipefail
if [ "$1" = "s3api" ] && [ "$2" = "list-objects" ]; then
cat <<'EOF'
{listing_json}
EOF
elif [ "$1" = "s3" ] && [ "$2" = "presign" ]; then
printf '%s\\n' "file:///$3"
elif [ "$1" = "s3" ] && [ "$2" = "cp" ]; then
printf '%s\\n' "$4" >> "{uploads_log}"
else
printf 'unexpected aws invocation: %s\\n' "$*" >&2
exit 1
fi
""",
)
self._write_executable(
bin_dir / "jq",
"""#!/usr/bin/env python3
import json
import sys

query = sys.argv[-1]
payload = json.load(sys.stdin)
contents = payload.get("Contents") or []
if query == '(.Contents // []) | map(select(.Key | endswith("/") | not)) | length':
print(len([item for item in contents if not item["Key"].endswith("/")]))
elif query == '(.Contents // [])[] | select(.Key | endswith("/") | not) | [.Key, (.Size | tostring)] | @tsv':
for item in contents:
if not item["Key"].endswith("/"):
print(f"{item['Key']}\\t{item['Size']}")
else:
raise SystemExit(f"unexpected jq query: {query}")
""",
)
self._write_executable(bin_dir / "java", self._java_script(java_mode))
self._write_executable(
bin_dir / "curl",
"""#!/usr/bin/env bash
printf 'curl should not run when Saxon jar already exists\\n' >&2
exit 1
""",
)
self._write_executable(
bin_dir / "sha1sum",
"""#!/usr/bin/env bash
printf 'sha1sum should not run when Saxon jar already exists\\n' >&2
exit 1
""",
)

env = os.environ.copy()
env.update(
{
"BUCKET": "test-bucket",
"DAG_ID": "funcake_test",
"DAG_TS": "2021-03-23_17-25-10",
"DEST": "transformed",
"HOME": tempdir,
"PATH": f"{bin_dir}:{env['PATH']}",
"SAXON_CP": str(saxon_jar),
"SCRIPTS_PATH": str(REPO_ROOT / "funcake_dags" / "scripts"),
"SOURCE": "new-updated-filtered",
"TMPDIR": tempdir,
"XSL_BRANCH": "main",
"XSL_FILENAME": "transforms/test.xsl",
"XSL_REPO": "tulibraries/aggregator_mdx",
}
)

result = subprocess.run(
[str(SCRIPT_PATH)],
check=False,
capture_output=True,
text=True,
cwd=tempdir,
env=env,
)

uploads = uploads_log.read_text(encoding="utf-8")
transformed_files = sorted(path.name for path in tempdir_path.rglob("*.xml-transformed.xml"))
return result, uploads, transformed_files

def _java_script(self, java_mode: str):
if java_mode == "fail":
return """#!/usr/bin/env bash
printf 'java should not run for zero-byte inputs\\n' >&2
exit 1
"""

transformed_payload = """<root>
<oai_dc:dc xmlns:oai_dc="urn:oai_dc"/>
<dcterms:identifier>id</dcterms:identifier>
</root>"""
if java_mode == "transform_zero_records":
transformed_payload = """<root>
<dcterms:identifier>id</dcterms:identifier>
</root>"""

return f"""#!/usr/bin/env bash
set -euo pipefail
output=""
for arg in "$@"; do
case "$arg" in
-o:*)
output="${{arg#-o:}}"
;;
esac
done

if [ -z "$output" ]; then
printf 'missing output arg\\n' >&2
exit 1
fi

mkdir -p "$(dirname "$output")"
if [[ "$output" == *"-transformed.xml" ]]; then
cat <<'EOF' > "$output"
{transformed_payload}
EOF
else
cat <<'EOF' > "$output"
<?xml version="1.0"?>
<record/>
EOF
fi
"""

def _write_executable(self, path: Path, content: str):
path.write_text(textwrap.dedent(content), encoding="utf-8")
path.chmod(path.stat().st_mode | stat.S_IEXEC)