Last active
July 13, 2026 01:20
-
-
Save s1037989/c54562692d82b03a2219b4c581410a90 to your computer and use it in GitHub Desktop.
s3 resumable multipart uploads; extracts tars on the fly. All object uploads are broken into multiparts and resumable.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| #!/usr/bin/env perl | |
| use v5.22; | |
| use strict; | |
| use warnings; | |
| use feature 'signatures'; | |
| use MIME::Base64 qw(encode_base64); | |
| use constant { | |
| POLY => 0x9A6C9329AC4BC9B5, | |
| MASK64 => 0xFFFFFFFFFFFFFFFF, | |
| }; | |
| sub make_crc64nvme_table { | |
| my @table; | |
| for my $byte (0 .. 255) { | |
| my $crc = $byte; | |
| for (1 .. 8) { | |
| $crc = $crc & 1 | |
| ? ($crc >> 1) ^ POLY | |
| : $crc >> 1; | |
| } | |
| push @table, $crc; | |
| } | |
| return \@table; | |
| } | |
| my $table = make_crc64nvme_table(); | |
| sub crc64nvme_file ($filename) { | |
| open my $fh, '<:raw', $filename | |
| or die "Cannot open '$filename': $!\n"; | |
| my $crc = MASK64; | |
| my $buffer; | |
| while (1) { | |
| my $length = read $fh, $buffer, 1024 * 1024; | |
| die "Cannot read '$filename': $!\n" | |
| unless defined $length; | |
| last if $length == 0; | |
| for my $byte (unpack 'C*', $buffer) { | |
| $crc = $table->[($crc ^ $byte) & 0xff] ^ ($crc >> 8); | |
| } | |
| } | |
| close $fh | |
| or die "Cannot close '$filename': $!\n"; | |
| return $crc ^ MASK64; | |
| } | |
| sub usage { | |
| die <<"USAGE"; | |
| Usage: $0 [--hex | --base64 | --both] FILE | |
| --hex Print the hexadecimal CRC64NVME value | |
| --base64 Print the Base64 value used by AWS S3 | |
| --both Print both values | |
| USAGE | |
| } | |
| my $format = 'base64'; | |
| if (@ARGV && $ARGV[0] =~ /^--(hex|base64|both)$/) { | |
| $format = $1; | |
| shift @ARGV; | |
| } | |
| usage() unless @ARGV == 1; | |
| my $filename = shift; | |
| my $crc = crc64nvme_file($filename); | |
| my $binary = pack 'Q>', $crc; | |
| my $hex = unpack 'H*', $binary; | |
| my $base64 = encode_base64($binary, ''); | |
| if ($format eq 'hex') { | |
| say $hex; | |
| } | |
| elsif ($format eq 'base64') { | |
| say $base64; | |
| } | |
| else { | |
| say "hex: $hex"; | |
| say "base64: $base64"; | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| #!/usr/bin/env bash | |
| set -Eeuo pipefail | |
| # Upload one large file to S3 using multipart upload only. | |
| # | |
| # SOURCE may be: | |
| # * a generic blob -> one S3 object | |
| # * a compressed or uncompressed tar -> one S3 prefix containing its members | |
| # | |
| # Persistent state is one SQLite file dedicated to this source upload. SQLite | |
| # remembers only source/object identity and multipart UploadIds. S3 list-parts is | |
| # authoritative for which parts have succeeded. | |
| # | |
| # Usage: | |
| # s3-upload SOURCE BUCKET [DESTINATION] | |
| # | |
| # Optional environment: | |
| # PART_SIZE_MIB=64 | |
| # STATE_DIR=/persistent/state | |
| # AWS_ENDPOINT_URL=http://127.0.0.1:9000 | |
| # KEEP_STATE_ON_SUCCESS=1 | |
| # FAIL_AFTER_PARTS=N | |
| PROGRAM=${0##*/} | |
| MIN_PART_SIZE=$((5 * 1024 * 1024)) | |
| MAX_PART_SIZE=$((5 * 1024 * 1024 * 1024)) | |
| MAX_PARTS=10000 | |
| log() { printf '%s: %s\n' "$PROGRAM" "$*" >&2; } | |
| die() { log "ERROR: $*"; exit 1; } | |
| require_command() { command -v "$1" >/dev/null 2>&1 || die "missing command: $1"; } | |
| sql_quote() { | |
| local value=${1//\'/\'\'} | |
| printf "'%s'" "$value" | |
| } | |
| aws_s3api() { | |
| local options=() | |
| [[ -n ${AWS_ENDPOINT_URL:-} ]] && options+=(--endpoint-url "$AWS_ENDPOINT_URL") | |
| aws "${options[@]}" s3api "$@" | |
| } | |
| sha256_file() { | |
| perl -MDigest::SHA -e ' | |
| my $sha = Digest::SHA->new(256); | |
| $sha->addfile($ARGV[0]); | |
| print $sha->hexdigest, "\n"; | |
| ' "$1" | |
| } | |
| sha256_text() { | |
| perl -MDigest::SHA=sha256_hex -e 'local $/; print sha256_hex(<STDIN>), "\n"' | |
| } | |
| normalize_key() { | |
| local key=$1 | |
| while [[ $key == /* ]]; do key=${key#/}; done | |
| while [[ $key == *//* ]]; do key=${key//\/\//\/}; done | |
| printf '%s\n' "$key" | |
| } | |
| archive_base_name() { | |
| local name=${1##*/} | |
| case "$name" in | |
| *.tar.gz) printf '%s\n' "${name%.tar.gz}" ;; | |
| *.tar.bz2) printf '%s\n' "${name%.tar.bz2}" ;; | |
| *.tar.xz) printf '%s\n' "${name%.tar.xz}" ;; | |
| *.tar.zst) printf '%s\n' "${name%.tar.zst}" ;; | |
| *.tar.lzma) printf '%s\n' "${name%.tar.lzma}" ;; | |
| *.tgz) printf '%s\n' "${name%.tgz}" ;; | |
| *.tbz|*.tbz2) printf '%s\n' "${name%.*}" ;; | |
| *.txz) printf '%s\n' "${name%.txz}" ;; | |
| *.tzst) printf '%s\n' "${name%.tzst}" ;; | |
| *.tar) printf '%s\n' "${name%.tar}" ;; | |
| *) printf '%s\n' "$name" ;; | |
| esac | |
| } | |
| safe_member_name() { | |
| local name=${1#./} component | |
| [[ -n $name && $name != /* && $name != *$'\n'* && $name != *$'\r'* ]] || return 1 | |
| IFS='/' read -r -a components <<<"$name" | |
| for component in "${components[@]}"; do | |
| [[ $component != '..' ]] || return 1 | |
| done | |
| printf '%s\n' "$name" | |
| } | |
| choose_part_size() { | |
| local size=$1 | |
| local mib=$((1024 * 1024)) | |
| local selected=$(( ${PART_SIZE_MIB:-64} * mib )) | |
| local required=$(( (size + MAX_PARTS - 1) / MAX_PARTS )) | |
| (( selected < MIN_PART_SIZE )) && selected=$MIN_PART_SIZE | |
| (( required > selected )) && selected=$required | |
| selected=$(( ((selected + mib - 1) / mib) * mib )) | |
| (( selected <= MAX_PART_SIZE )) || die "object exceeds S3 multipart limits" | |
| printf '%s\n' "$selected" | |
| } | |
| init_database() { | |
| sqlite3 "$DB" <<'SQL' | |
| PRAGMA journal_mode=WAL; | |
| PRAGMA synchronous=FULL; | |
| PRAGMA busy_timeout=30000; | |
| CREATE TABLE IF NOT EXISTS source ( | |
| singleton INTEGER PRIMARY KEY CHECK (singleton=1), | |
| source_id TEXT NOT NULL, | |
| source_path TEXT NOT NULL, | |
| source_size INTEGER NOT NULL, | |
| source_mtime INTEGER NOT NULL, | |
| source_sha256 TEXT NOT NULL, | |
| source_kind TEXT NOT NULL, | |
| bucket TEXT NOT NULL, | |
| destination TEXT NOT NULL | |
| ); | |
| CREATE TABLE IF NOT EXISTS object_upload ( | |
| object_key TEXT PRIMARY KEY, | |
| member_name TEXT NOT NULL, | |
| object_size INTEGER NOT NULL, | |
| member_mtime INTEGER NOT NULL, | |
| part_size INTEGER NOT NULL, | |
| upload_id TEXT, | |
| complete INTEGER NOT NULL DEFAULT 0 CHECK (complete IN (0,1)), | |
| updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP | |
| ); | |
| SQL | |
| } | |
| write_source_record() { | |
| local existing expected | |
| existing=$(sqlite3 "$DB" \ | |
| "SELECT source_id||'|'||bucket||'|'||destination||'|'||source_kind | |
| FROM source WHERE singleton=1;") | |
| expected="$SOURCE_ID|$BUCKET|$DESTINATION|$SOURCE_KIND" | |
| [[ -z $existing || $existing == "$expected" ]] || | |
| die "state database belongs to another upload: $DB" | |
| sqlite3 "$DB" " | |
| INSERT OR IGNORE INTO source VALUES( | |
| 1, $(sql_quote "$SOURCE_ID"), $(sql_quote "$SOURCE"), | |
| $SOURCE_SIZE, $SOURCE_MTIME, $(sql_quote "$SOURCE_SHA256"), | |
| $(sql_quote "$SOURCE_KIND"), $(sql_quote "$BUCKET"), | |
| $(sql_quote "$DESTINATION") | |
| ); | |
| " | |
| } | |
| member_id() { printf '%s' "$1" | sha256_text; } | |
| head_matches() { | |
| local key=$1 member=$2 size=$3 response expected_member | |
| expected_member=$(member_id "$member") | |
| response=$(aws_s3api head-object \ | |
| --bucket "$BUCKET" --key "$key" --output json 2>/dev/null) || return 1 | |
| EXPECT_SOURCE=$SOURCE_ID EXPECT_MEMBER=$expected_member EXPECT_SIZE=$size \ | |
| perl -MJSON::PP -e ' | |
| local $/; | |
| my $j = decode_json(<STDIN>); | |
| my $m = $j->{Metadata} || {}; | |
| exit !(($j->{ContentLength}//-1) == $ENV{EXPECT_SIZE} | |
| && ($m->{"upload-source-id"}//"") eq $ENV{EXPECT_SOURCE} | |
| && ($m->{"upload-member-id"}//"") eq $ENV{EXPECT_MEMBER} | |
| && ($m->{"upload-complete"}//"") eq "1"); | |
| ' <<<"$response" | |
| } | |
| ensure_object_row() { | |
| local key=$1 member=$2 size=$3 mtime=$4 part_size=$5 | |
| sqlite3 "$DB" " | |
| INSERT INTO object_upload(object_key,member_name,object_size,member_mtime,part_size) | |
| VALUES($(sql_quote "$key"),$(sql_quote "$member"),$size,$mtime,$part_size) | |
| ON CONFLICT(object_key) DO UPDATE SET | |
| member_name=excluded.member_name, | |
| object_size=excluded.object_size, | |
| member_mtime=excluded.member_mtime, | |
| part_size=excluded.part_size, | |
| updated_at=CURRENT_TIMESTAMP | |
| WHERE object_upload.complete=0; | |
| " | |
| } | |
| saved_upload_id() { | |
| sqlite3 "$DB" "SELECT COALESCE(upload_id,'') FROM object_upload | |
| WHERE object_key=$(sql_quote "$1");" | |
| } | |
| mark_complete() { | |
| sqlite3 "$DB" "UPDATE object_upload SET complete=1,upload_id=NULL, | |
| updated_at=CURRENT_TIMESTAMP | |
| WHERE object_key=$(sql_quote "$1");" | |
| } | |
| clear_upload_id() { | |
| sqlite3 "$DB" "UPDATE object_upload SET upload_id=NULL, | |
| updated_at=CURRENT_TIMESTAMP | |
| WHERE object_key=$(sql_quote "$1");" | |
| } | |
| create_upload() { | |
| local key=$1 member=$2 size=$3 upload_id metadata | |
| metadata="upload-source-id=$SOURCE_ID,upload-member-id=$(member_id "$member"),upload-complete=1,upload-size=$size" | |
| upload_id=$(aws_s3api create-multipart-upload \ | |
| --bucket "$BUCKET" --key "$key" --metadata "$metadata" \ | |
| --query UploadId --output text) | |
| sqlite3 "$DB" "UPDATE object_upload SET upload_id=$(sql_quote "$upload_id"), | |
| updated_at=CURRENT_TIMESTAMP | |
| WHERE object_key=$(sql_quote "$key");" | |
| printf '%s\n' "$upload_id" | |
| } | |
| # Populate these associative arrays from one authoritative S3 list-parts call. | |
| declare -A S3_ETAG=() | |
| declare -A S3_SIZE=() | |
| load_s3_parts() { | |
| local key=$1 upload_id=$2 response line number size etag | |
| S3_ETAG=() | |
| S3_SIZE=() | |
| response=$(aws_s3api list-parts \ | |
| --bucket "$BUCKET" --key "$key" --upload-id "$upload_id" \ | |
| --output json 2>/dev/null) || return 1 | |
| while IFS=$'\t' read -r number size etag; do | |
| [[ -n $number ]] || continue | |
| S3_SIZE[$number]=$size | |
| S3_ETAG[$number]=$etag | |
| done < <(perl -MJSON::PP -e ' | |
| local $/; | |
| my $j=decode_json(<STDIN>); | |
| for my $p (@{$j->{Parts}||[]}) { | |
| print join("\t",$p->{PartNumber},$p->{Size},$p->{ETag}),"\n"; | |
| } | |
| ' <<<"$response") | |
| } | |
| PREPARED_UPLOAD_ID= | |
| prepare_object() { | |
| local key=$1 member=$2 size=$3 mtime=$4 part_size=$5 upload_id | |
| PREPARED_UPLOAD_ID= | |
| ensure_object_row "$key" "$member" "$size" "$mtime" "$part_size" | |
| if head_matches "$key" "$member" "$size"; then | |
| mark_complete "$key" | |
| PREPARED_UPLOAD_ID= | |
| return | |
| fi | |
| upload_id=$(saved_upload_id "$key") | |
| if [[ -n $upload_id ]] && ! load_s3_parts "$key" "$upload_id"; then | |
| log "$key: saved upload no longer exists; starting a new one" | |
| clear_upload_id "$key" | |
| upload_id= | |
| fi | |
| if [[ -z $upload_id ]]; then | |
| upload_id=$(create_upload "$key" "$member" "$size") | |
| load_s3_parts "$key" "$upload_id" || die "$key: cannot list new upload" | |
| fi | |
| PREPARED_UPLOAD_ID=$upload_id | |
| } | |
| completion_json() { | |
| local expected_parts=$1 output=$2 number | |
| { | |
| printf '{"Parts":[' | |
| local separator= | |
| for ((number=1; number<=expected_parts; number++)); do | |
| [[ -n ${S3_ETAG[$number]+x} ]] || die "missing S3 part $number" | |
| printf '%s{"PartNumber":%d,"ETag":%s}' \ | |
| "$separator" "$number" \ | |
| "$(perl -MJSON::PP -e 'print encode_json($ARGV[0])' "${S3_ETAG[$number]}")" | |
| separator=, | |
| done | |
| printf ']}\n' | |
| } >"$output" | |
| } | |
| maybe_fail() { | |
| [[ -n ${FAIL_AFTER_PARTS:-} ]] || return | |
| local file="$WORK/fault-count" count=0 | |
| [[ -f $file ]] && count=$(<"$file") | |
| count=$((count+1)); printf '%s\n' "$count" >"$file" | |
| (( count < FAIL_AFTER_PARTS )) || die "intentional failure after $count uploaded part(s)" | |
| } | |
| upload_part_file() { | |
| local key=$1 upload_id=$2 number=$3 file=$4 size=$5 etag | |
| etag=$(aws_s3api upload-part \ | |
| --bucket "$BUCKET" --key "$key" --upload-id "$upload_id" \ | |
| --part-number "$number" --body "$file" --query ETag --output text) | |
| S3_SIZE[$number]=$size | |
| S3_ETAG[$number]=$etag | |
| maybe_fail | |
| } | |
| complete_object() { | |
| local key=$1 member=$2 size=$3 upload_id=$4 part_count=$5 manifest response | |
| manifest="$WORK/complete.$$.json" | |
| completion_json "$part_count" "$manifest" | |
| response=$(aws_s3api complete-multipart-upload \ | |
| --bucket "$BUCKET" --key "$key" --upload-id "$upload_id" \ | |
| --multipart-upload "file://$manifest" --output json) | |
| rm -f "$manifest" | |
| head_matches "$key" "$member" "$size" || | |
| die "$key: completed object failed remote verification" | |
| mark_complete "$key" | |
| log "$key: complete" | |
| } | |
| # Upload one object from a seekable file. | |
| upload_file_object() { | |
| local key=$1 file=$2 member=$3 size=$4 mtime=$5 | |
| local part_size upload_id part_count number offset bytes staged actual | |
| part_size=$(choose_part_size "$size") | |
| prepare_object "$key" "$member" "$size" "$mtime" "$part_size" | |
| upload_id=$PREPARED_UPLOAD_ID | |
| [[ -n $upload_id ]] || { log "$key: already complete"; return; } | |
| part_count=$(( size == 0 ? 1 : (size + part_size - 1) / part_size )) | |
| offset=0 | |
| for ((number=1; number<=part_count; number++)); do | |
| bytes=$part_size | |
| (( size == 0 )) && bytes=0 | |
| (( offset + bytes > size )) && bytes=$((size-offset)) | |
| if [[ ${S3_SIZE[$number]:--1} == "$bytes" ]]; then | |
| log "$key: part $number/$part_count already present" | |
| offset=$((offset+bytes)) | |
| continue | |
| fi | |
| staged="$WORK/part.bin" | |
| if (( bytes == 0 )); then | |
| : >"$staged" | |
| else | |
| dd if="$file" of="$staged" bs=1M skip="$offset" count="$bytes" \ | |
| iflag=skip_bytes,count_bytes status=none | |
| fi | |
| actual=$(wc -c <"$staged" | tr -d '[:space:]') | |
| (( actual == bytes )) || die "$key: failed to stage part $number" | |
| log "$key: uploading part $number/$part_count ($bytes bytes)" | |
| upload_part_file "$key" "$upload_id" "$number" "$staged" "$bytes" | |
| rm -f "$staged" | |
| offset=$((offset+bytes)) | |
| done | |
| complete_object "$key" "$member" "$size" "$upload_id" "$part_count" | |
| } | |
| # Upload one object from tar-provided stdin. Existing parts are consumed and | |
| # discarded so stdin advances, but they are not retransmitted. | |
| upload_stream_object() { | |
| local key=$1 member=$2 size=$3 mtime=$4 | |
| local part_size upload_id part_count number bytes staged actual | |
| part_size=$(choose_part_size "$size") | |
| prepare_object "$key" "$member" "$size" "$mtime" "$part_size" | |
| upload_id=$PREPARED_UPLOAD_ID | |
| [[ -n $upload_id ]] || { log "$key: already complete"; return; } | |
| part_count=$(( size == 0 ? 1 : (size + part_size - 1) / part_size )) | |
| for ((number=1; number<=part_count; number++)); do | |
| if (( size == 0 )); then | |
| bytes=0 | |
| elif (( number < part_count )); then | |
| bytes=$part_size | |
| else | |
| bytes=$((size-(number-1)*part_size)) | |
| fi | |
| staged="$WORK/part.bin" | |
| if (( bytes == 0 )); then | |
| : >"$staged" | |
| else | |
| dd of="$staged" bs="$bytes" count=1 iflag=fullblock status=none | |
| fi | |
| actual=$(wc -c <"$staged" | tr -d '[:space:]') | |
| (( actual == bytes )) || die "$key: short tar member part $number" | |
| if [[ ${S3_SIZE[$number]:--1} == "$bytes" ]]; then | |
| log "$key: part $number/$part_count already present; discarded input" | |
| rm -f "$staged" | |
| continue | |
| fi | |
| log "$key: uploading part $number/$part_count ($bytes bytes)" | |
| upload_part_file "$key" "$upload_id" "$number" "$staged" "$bytes" | |
| rm -f "$staged" | |
| done | |
| complete_object "$key" "$member" "$size" "$upload_id" "$part_count" | |
| } | |
| completion_marker_content() { | |
| printf '{"source_id":"%s","source_sha256":"%s","kind":"tar"}\n' \ | |
| "$SOURCE_ID" "$SOURCE_SHA256" | |
| } | |
| upload_tar_marker() { | |
| local prefix=$1 file="$WORK/tar-complete.json" | |
| completion_marker_content >"$file" | |
| upload_file_object "${prefix%/}/.s3-upload-complete" "$file" \ | |
| '__tar_completion_marker__' "$(wc -c <"$file" | tr -d '[:space:]')" "$SOURCE_MTIME" | |
| } | |
| tar_marker_matches() { | |
| local prefix=$1 size | |
| size=$(completion_marker_content | wc -c | tr -d '[:space:]') | |
| head_matches "${prefix%/}/.s3-upload-complete" '__tar_completion_marker__' "$size" | |
| } | |
| tar_member_mode() { | |
| DB=${IDU_DB:?}; WORK=${IDU_WORK:?}; BUCKET=${IDU_BUCKET:?} | |
| SOURCE_ID=${IDU_SOURCE_ID:?}; AWS_ENDPOINT_URL=${IDU_AWS_ENDPOINT_URL:-} | |
| FAIL_AFTER_PARTS=${IDU_FAIL_AFTER_PARTS:-}; PROGRAM=${IDU_PROGRAM:-$PROGRAM} | |
| local raw=${TAR_FILENAME:?} size=${TAR_SIZE:?} mtime=${TAR_MTIME:-0} | |
| local member key | |
| member=$(safe_member_name "$raw") || die "unsafe tar member: $raw" | |
| key=$(normalize_key "${IDU_PREFIX%/}/$member") | |
| upload_stream_object "$key" "$member" "$size" "$mtime" | |
| exit 0 | |
| } | |
| [[ ${1:-} == --tar-member ]] && { shift; tar_member_mode; } | |
| main() { | |
| require_command aws; require_command sqlite3; require_command perl | |
| require_command tar; require_command dd | |
| [[ $# -ge 2 && $# -le 3 ]] || die "usage: $PROGRAM SOURCE BUCKET [DESTINATION]" | |
| SOURCE=$1; BUCKET=$2; DESTINATION=${3:-${SOURCE##*/}} | |
| [[ -f $SOURCE ]] || die "not a regular file: $SOURCE" | |
| SOURCE=$(cd "$(dirname "$SOURCE")" && pwd -P)/$(basename "$SOURCE") | |
| DESTINATION=$(normalize_key "$DESTINATION") | |
| [[ -n $DESTINATION ]] || die "empty destination" | |
| SOURCE_SIZE=$(wc -c <"$SOURCE" | tr -d '[:space:]') | |
| SOURCE_MTIME=$(perl -e 'print((stat($ARGV[0]))[9],"\n")' "$SOURCE") | |
| SOURCE_SHA256=$(sha256_file "$SOURCE") | |
| SOURCE_ID=$(printf '%s\0%s\0%s' "$SOURCE_SIZE" "$SOURCE_MTIME" "$SOURCE_SHA256" | sha256_text) | |
| if tar -tf "$SOURCE" >/dev/null 2>&1; then SOURCE_KIND=tar; else SOURCE_KIND=blob; fi | |
| local state_root=${STATE_DIR:-"$(dirname "$SOURCE")/.s3-upload-state"} | |
| local state_name | |
| mkdir -p "$state_root" | |
| state_name=$(printf '%s\0%s\0%s' "$SOURCE_ID" "$BUCKET" "$DESTINATION" | sha256_text) | |
| DB="$state_root/$state_name.sqlite" | |
| WORK="$state_root/$state_name.work" | |
| mkdir -p "$WORK" | |
| init_database | |
| write_source_record | |
| if [[ $SOURCE_KIND == blob ]]; then | |
| upload_file_object "$DESTINATION" "$SOURCE" '__blob__' "$SOURCE_SIZE" "$SOURCE_MTIME" | |
| else | |
| local directory filename prefix command | |
| directory=${DESTINATION%/*}; filename=${DESTINATION##*/} | |
| [[ $directory == "$DESTINATION" ]] && directory= | |
| prefix=$(archive_base_name "$filename") | |
| [[ -n $directory ]] && prefix="$directory/$prefix" | |
| prefix=$(normalize_key "$prefix") | |
| if tar_marker_matches "$prefix"; then | |
| log "tar already complete: s3://$BUCKET/$prefix/" | |
| else | |
| export IDU_DB=$DB IDU_WORK=$WORK IDU_BUCKET=$BUCKET IDU_PREFIX=$prefix | |
| export IDU_SOURCE_ID=$SOURCE_ID IDU_AWS_ENDPOINT_URL=${AWS_ENDPOINT_URL:-} | |
| export IDU_FAIL_AFTER_PARTS=${FAIL_AFTER_PARTS:-} IDU_PROGRAM=$PROGRAM | |
| printf -v command '%q --tar-member' "$0" | |
| # Do not use -O: --to-command supplies each member on stdin. | |
| tar -xf "$SOURCE" --to-command="$command" | |
| upload_tar_marker "$prefix" | |
| tar_marker_matches "$prefix" || die "tar marker verification failed" | |
| log "tar complete: s3://$BUCKET/$prefix/" | |
| fi | |
| fi | |
| if [[ ${KEEP_STATE_ON_SUCCESS:-0} != 1 ]]; then | |
| rm -rf "$WORK" | |
| rm -f "$DB" "$DB-wal" "$DB-shm" | |
| rmdir "$state_root" 2>/dev/null || true | |
| fi | |
| log "all work complete" | |
| } | |
| main "$@" |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| #!/usr/bin/env perl | |
| use v5.22; | |
| use strict; | |
| use warnings; | |
| use feature 'signatures'; | |
| use Digest::SHA qw(sha256_hex); | |
| use File::Basename qw(basename dirname); | |
| use File::Path qw(make_path remove_tree); | |
| use File::Spec; | |
| use File::Temp qw(tempfile tempdir); | |
| use Fcntl qw(SEEK_SET); | |
| use Getopt::Long qw(GetOptions); | |
| use IO::Uncompress::Gunzip qw($GunzipError); | |
| use IO::Uncompress::Bunzip2 qw($Bunzip2Error); | |
| use IPC::Open3; | |
| use JSON::PP; | |
| use Symbol qw(gensym); | |
| # --------------------------------------------------------------------------- | |
| # s3-upload.pl | |
| # | |
| # Upload one local blob or TAR archive to S3. | |
| # | |
| # * Uses only modules distributed with core Perl. | |
| # * Uses the AWS CLI for S3 API calls. | |
| # * Parses TAR archives directly; GNU tar is not required. | |
| # * Uncompressed TAR members are seekable and may be skipped without reading. | |
| # * gzip/bzip2 TAR members are sequential because ordinary compressed streams | |
| # cannot seek to arbitrary uncompressed offsets. | |
| # * Every non-empty S3 object is uploaded through multipart upload. | |
| # * Multipart state is resumed from one durable JSON file per source upload. | |
| # * S3 ListParts is authoritative; successful parts are not retransmitted. | |
| # * Logical files larger than the configured object limit are split into | |
| # numbered S3 objects plus a manifest uploaded last. | |
| # | |
| # Usage: | |
| # s3-upload.pl [options] SOURCE BUCKET [DESTINATION] | |
| # | |
| # Examples: | |
| # ./s3-upload.pl disk.img my-bucket backups/disk.img | |
| # ./s3-upload.pl archive.tar.gz my-bucket imports/archive.tar.gz | |
| # | |
| # Options: | |
| # --state-dir DIR Persistent state directory | |
| # --work-dir DIR Directory for one-part staging files | |
| # --part-size SIZE Preferred multipart part size (default 64MiB) | |
| # --object-limit SIZE Split logical files at this size (default 5TiB) | |
| # --endpoint-url URL AWS CLI endpoint override for integration testing | |
| # --keep-state Keep state after successful completion | |
| # --help | |
| # | |
| # SIZE accepts plain bytes or KiB/MiB/GiB/TiB suffixes. | |
| # --------------------------------------------------------------------------- | |
| use constant { | |
| DEFAULT_LIMIT => 5 * 1024 * 1024 * 1024 * 1024, # 5 TiB | |
| DEFAULT_PART => 64 * 1024 * 1024, # 64 MiB | |
| MAX_META_ENTRY => 16 * 1024 * 1024, # 16 MiB | |
| MAX_PARTS => 10_000, # 5 TiB | |
| MAX_PART_SIZE => 5 * 1024 * 1024 * 1024, # 5 GiB | |
| MIN_PART_SIZE => 5 * 1024 * 1024, # 5 MiB | |
| STATE_VERSION => 1, | |
| TAR_BLOCK => 512, | |
| }; | |
| my %opt = ( | |
| part_size => DEFAULT_PART, | |
| object_limit => DEFAULT_LIMIT, | |
| keep_state => 0, | |
| ); | |
| GetOptions( | |
| 'state-dir=s' => \$opt{state_dir}, | |
| 'work-dir=s' => \$opt{work_dir}, | |
| 'part-size=s' => \$opt{part_size_text}, | |
| 'object-limit=s' => \$opt{object_limit_text}, | |
| 'endpoint-url=s' => \$opt{endpoint_url}, | |
| 'keep-state!' => \$opt{keep_state}, | |
| 'help' => \$opt{help}, | |
| ) or usage(2); | |
| usage(0) if $opt{help}; | |
| @ARGV >= 2 && @ARGV <= 3 or usage(2); | |
| my ($source, $bucket, $destination) = @ARGV; | |
| $opt{part_size} = parse_size($opt{part_size_text}) if defined $opt{part_size_text}; | |
| $opt{object_limit} = parse_size($opt{object_limit_text}) if defined $opt{object_limit_text}; | |
| $opt{part_size} >= MIN_PART_SIZE or fatal("part size must be at least 5 MiB"); | |
| $opt{part_size} <= MAX_PART_SIZE or fatal("part size cannot exceed 5 GiB"); | |
| $opt{object_limit} > 0 or fatal("object limit must be positive"); | |
| -f $source or fatal("source is not a regular file: $source"); | |
| $source = File::Spec->rel2abs($source); | |
| $destination //= basename($source); | |
| $destination = normalize_key($destination); | |
| length $destination or fatal("destination cannot be empty"); | |
| command_exists('aws') or fatal("aws CLI was not found in PATH"); | |
| my @source_stat = stat($source); | |
| my $source_size = $source_stat[7]; | |
| my $source_mtime = $source_stat[9]; | |
| my $state_root = $opt{state_dir} // File::Spec->catdir(dirname($source), '.s3-upload-state'); | |
| my $work_root = $opt{work_dir} // File::Spec->catdir($state_root, 'work'); | |
| make_path($state_root, $work_root); | |
| my $state_identity = sha256_hex(join "\0", $source, $bucket, $destination); | |
| my $state_file = File::Spec->catfile($state_root, "$state_identity.json"); | |
| my $work_dir = File::Spec->catdir($work_root, $state_identity); | |
| make_path($work_dir); | |
| my $state = load_state($state_file); | |
| if ($state) { | |
| validate_existing_state( | |
| $state, | |
| source => $source, | |
| source_size => $source_size, | |
| source_mtime => $source_mtime, | |
| bucket => $bucket, | |
| destination => $destination, | |
| ); | |
| } | |
| else { | |
| $state = { | |
| version => STATE_VERSION, | |
| source => { | |
| path => $source, | |
| size => 0 + $source_size, | |
| mtime => 0 + $source_mtime, | |
| bucket => $bucket, | |
| destination => $destination, | |
| }, | |
| objects => {}, | |
| }; | |
| save_state($state_file, $state); | |
| } | |
| my $source_kind = $state->{source}{kind}; | |
| if (!$source_kind) { | |
| $source_kind = classify_source($source); | |
| $state->{source}{kind} = $source_kind; | |
| save_state($state_file, $state); | |
| } | |
| if (!$state->{source}{sha256}) { | |
| log_message("calculating SHA-256 of source"); | |
| $state->{source}{sha256} = sha256_range($source, 0, $source_size); | |
| save_state($state_file, $state); | |
| } | |
| if ($source_kind eq 'blob') { | |
| upload_logical_range( | |
| state => $state, | |
| state_file => $state_file, | |
| bucket => $bucket, | |
| logical_key => $destination, | |
| logical_name => basename($destination), | |
| sha256 => $state->{source}{sha256}, | |
| source_file => $source, | |
| source_offset => 0, | |
| size => $source_size, | |
| ); | |
| } | |
| elsif ($source_kind eq 'tar') { | |
| process_uncompressed_tar( | |
| state => $state, | |
| state_file => $state_file, | |
| source => $source, | |
| bucket => $bucket, | |
| destination => $destination, | |
| ); | |
| } | |
| elsif ($source_kind eq 'tar.gz' || $source_kind eq 'tar.bz2') { | |
| process_compressed_tar( | |
| state => $state, | |
| state_file => $state_file, | |
| source => $source, | |
| kind => $source_kind, | |
| bucket => $bucket, | |
| destination => $destination, | |
| ); | |
| } | |
| else { | |
| fatal("unsupported source kind: $source_kind"); | |
| } | |
| unless ($opt{keep_state}) { | |
| unlink $state_file | |
| or fatal("cannot remove completed state file $state_file: $!"); | |
| remove_tree($work_dir) if -d $work_dir; | |
| rmdir $work_root; | |
| rmdir $state_root; | |
| } | |
| log_message("upload complete"); | |
| exit 0; | |
| # =========================================================================== | |
| # High-level source processing | |
| # =========================================================================== | |
| sub process_uncompressed_tar (%arg) { | |
| my $state = $arg{state}; | |
| my $state_file = $arg{state_file}; | |
| my $source = $arg{source}; | |
| my $bucket = $arg{bucket}; | |
| my $destination = $arg{destination}; | |
| my $prefix = tar_destination_prefix($destination); | |
| my $members = $state->{source}{members}; | |
| if (!$members) { | |
| log_message("indexing uncompressed TAR"); | |
| my $reader = FileReader->new($source); | |
| $members = scan_tar( | |
| reader => $reader, | |
| record_offsets => 1, | |
| calculate_hashes => 1, | |
| ); | |
| $state->{source}{members} = $members; | |
| save_state($state_file, $state); | |
| } | |
| for my $member (@$members) { | |
| my $key = normalize_key("$prefix/$member->{path}"); | |
| upload_logical_range( | |
| state => $state, | |
| state_file => $state_file, | |
| bucket => $bucket, | |
| logical_key => $key, | |
| logical_name => basename($member->{path}), | |
| sha256 => $member->{sha256}, | |
| source_file => $source, | |
| source_offset => $member->{data_offset}, | |
| size => $member->{size}, | |
| ); | |
| } | |
| upload_tar_completion_marker( | |
| state => $state, | |
| state_file => $state_file, | |
| bucket => $bucket, | |
| prefix => $prefix, | |
| source_sha256 => $state->{source}{sha256}, | |
| member_count => scalar @$members, | |
| ); | |
| } | |
| sub process_compressed_tar (%arg) { | |
| my $state = $arg{state}; | |
| my $state_file = $arg{state_file}; | |
| my $source = $arg{source}; | |
| my $kind = $arg{kind}; | |
| my $bucket = $arg{bucket}; | |
| my $destination = $arg{destination}; | |
| my $prefix = tar_destination_prefix($destination); | |
| my $members = $state->{source}{members}; | |
| if (!$members) { | |
| log_message("indexing compressed TAR; this requires one sequential pass"); | |
| my $reader = compressed_reader($source, $kind); | |
| $members = scan_tar( | |
| reader => $reader, | |
| record_offsets => 0, | |
| calculate_hashes => 1, | |
| ); | |
| $state->{source}{members} = $members; | |
| save_state($state_file, $state); | |
| } | |
| my %expected = map { $_->{path} => $_ } @$members; | |
| my %seen; | |
| log_message("streaming compressed TAR for upload"); | |
| my $reader = compressed_reader($source, $kind); | |
| iterate_tar( | |
| reader => $reader, | |
| on_regular_file => sub ($entry, $data_reader) { | |
| my $path = $entry->{path}; | |
| my $known = $expected{$path} | |
| or fatal("TAR changed since indexing: unexpected member $path"); | |
| $entry->{size} == $known->{size} | |
| or fatal("TAR changed since indexing: size mismatch for $path"); | |
| $seen{$path}++; | |
| my $key = normalize_key("$prefix/$path"); | |
| upload_logical_stream( | |
| state => $state, | |
| state_file => $state_file, | |
| bucket => $bucket, | |
| logical_key => $key, | |
| logical_name => basename($path), | |
| sha256 => $known->{sha256}, | |
| size => $known->{size}, | |
| reader => $data_reader, | |
| ); | |
| }, | |
| ); | |
| for my $path (keys %expected) { | |
| $seen{$path} | |
| or fatal("TAR changed since indexing: missing member $path"); | |
| } | |
| upload_tar_completion_marker( | |
| state => $state, | |
| state_file => $state_file, | |
| bucket => $bucket, | |
| prefix => $prefix, | |
| source_sha256 => $state->{source}{sha256}, | |
| member_count => scalar @$members, | |
| ); | |
| } | |
| # =========================================================================== | |
| # Logical object splitting and manifests | |
| # =========================================================================== | |
| sub upload_logical_range (%arg) { | |
| my $plan = build_segment_plan( | |
| logical_key => $arg{logical_key}, | |
| logical_name => $arg{logical_name}, | |
| sha256 => $arg{sha256}, | |
| size => $arg{size}, | |
| ); | |
| for my $segment (@{$plan->{segments}}) { | |
| upload_s3_object_range( | |
| state => $arg{state}, | |
| state_file => $arg{state_file}, | |
| bucket => $arg{bucket}, | |
| key => $segment->{key}, | |
| source_file => $arg{source_file}, | |
| source_offset => $arg{source_offset} + $segment->{offset}, | |
| size => $segment->{size}, | |
| logical_sha256 => $arg{sha256}, | |
| segment_number => $segment->{number}, | |
| segment_count => $plan->{segment_count}, | |
| ); | |
| } | |
| upload_split_manifest( | |
| %arg, | |
| plan => $plan, | |
| ) if $plan->{segment_count} > 1; | |
| } | |
| sub upload_logical_stream (%arg) { | |
| my $plan = build_segment_plan( | |
| logical_key => $arg{logical_key}, | |
| logical_name => $arg{logical_name}, | |
| sha256 => $arg{sha256}, | |
| size => $arg{size}, | |
| ); | |
| for my $segment (@{$plan->{segments}}) { | |
| upload_s3_object_stream( | |
| state => $arg{state}, | |
| state_file => $arg{state_file}, | |
| bucket => $arg{bucket}, | |
| key => $segment->{key}, | |
| reader => $arg{reader}, | |
| size => $segment->{size}, | |
| logical_sha256 => $arg{sha256}, | |
| segment_number => $segment->{number}, | |
| segment_count => $plan->{segment_count}, | |
| ); | |
| } | |
| upload_split_manifest( | |
| %arg, | |
| plan => $plan, | |
| ) if $plan->{segment_count} > 1; | |
| } | |
| sub build_segment_plan (%arg) { | |
| my $size = $arg{size}; | |
| my $limit = $opt{object_limit}; | |
| my $count = $size == 0 ? 1 : int(($size + $limit - 1) / $limit); | |
| if ($count == 1) { | |
| return { | |
| segment_count => 1, | |
| segments => [{ | |
| number => 1, | |
| offset => 0, | |
| size => 0 + $size, | |
| key => $arg{logical_key}, | |
| }], | |
| }; | |
| } | |
| my $width = length($count); | |
| $width = 4 if $width < 4; | |
| my $parent = dirname($arg{logical_key}); | |
| $parent = '' if $parent eq '.'; | |
| my $base = basename($arg{logical_key}); | |
| my $prefix = "$arg{logical_key}.parts"; | |
| my $short_sha = substr($arg{sha256}, 0, 12); | |
| my @segments; | |
| for my $number (1 .. $count) { | |
| my $offset = ($number - 1) * $limit; | |
| my $length = $limit; | |
| $length = $size - $offset if $offset + $length > $size; | |
| my $filename = sprintf( | |
| "%s.%s.%0*d-of-%0*d", | |
| $base, | |
| $short_sha, | |
| $width, | |
| $number, | |
| $width, | |
| $count, | |
| ); | |
| push @segments, { | |
| number => $number, | |
| offset => 0 + $offset, | |
| size => 0 + $length, | |
| key => normalize_key("$prefix/$filename"), | |
| }; | |
| } | |
| return { | |
| segment_count => $count, | |
| manifest_key => normalize_key("$prefix/manifest.json"), | |
| segments => \@segments, | |
| }; | |
| } | |
| sub upload_split_manifest (%arg) { | |
| my $plan = $arg{plan}; | |
| my $manifest = { | |
| format => 's3-upload-split-object-v1', | |
| logical_key => $arg{logical_key}, | |
| logical_name => $arg{logical_name}, | |
| logical_size => 0 + $arg{size}, | |
| sha256 => $arg{sha256}, | |
| segment_count => 0 + $plan->{segment_count}, | |
| segments => [ | |
| map { | |
| +{ | |
| number => 0 + $_->{number}, | |
| key => $_->{key}, | |
| offset => 0 + $_->{offset}, | |
| size => 0 + $_->{size}, | |
| } | |
| } @{$plan->{segments}} | |
| ], | |
| }; | |
| my ($fh, $filename) = tempfile( | |
| 'manifest-XXXXXX', | |
| DIR => $work_dir, | |
| UNLINK => 0, | |
| ); | |
| binmode $fh; | |
| print {$fh} canonical_json($manifest); | |
| close $fh or fatal("cannot close $filename: $!"); | |
| my $size = -s $filename; | |
| my $sha = sha256_range($filename, 0, $size); | |
| upload_s3_object_range( | |
| state => $arg{state}, | |
| state_file => $arg{state_file}, | |
| bucket => $arg{bucket}, | |
| key => $plan->{manifest_key}, | |
| source_file => $filename, | |
| source_offset => 0, | |
| size => $size, | |
| logical_sha256 => $sha, | |
| segment_number => 1, | |
| segment_count => 1, | |
| ); | |
| unlink $filename; | |
| } | |
| sub upload_tar_completion_marker (%arg) { | |
| my $key = normalize_key("$arg{prefix}/.s3-upload-complete.json"); | |
| my $marker = { | |
| format => 's3-upload-tar-v1', | |
| source_sha256 => $arg{source_sha256}, | |
| member_count => 0 + $arg{member_count}, | |
| }; | |
| my ($fh, $filename) = tempfile( | |
| 'tar-marker-XXXXXX', | |
| DIR => $work_dir, | |
| UNLINK => 0, | |
| ); | |
| binmode $fh; | |
| print {$fh} canonical_json($marker); | |
| close $fh or fatal("cannot close $filename: $!"); | |
| my $size = -s $filename; | |
| my $sha = sha256_range($filename, 0, $size); | |
| upload_s3_object_range( | |
| state => $arg{state}, | |
| state_file => $arg{state_file}, | |
| bucket => $arg{bucket}, | |
| key => $key, | |
| source_file => $filename, | |
| source_offset => 0, | |
| size => $size, | |
| logical_sha256 => $sha, | |
| segment_number => 1, | |
| segment_count => 1, | |
| ); | |
| unlink $filename; | |
| } | |
| # =========================================================================== | |
| # Multipart upload engine | |
| # =========================================================================== | |
| sub upload_s3_object_range (%arg) { | |
| my $source_file = $arg{source_file}; | |
| my $source_offset = $arg{source_offset}; | |
| my $size = $arg{size}; | |
| my $read_part = sub ($relative_offset, $length, $output_file) { | |
| open my $input, '<:raw', $source_file | |
| or fatal("cannot open $source_file: $!"); | |
| sysseek($input, $source_offset + $relative_offset, SEEK_SET) | |
| // fatal("cannot seek $source_file: $!"); | |
| copy_exact($input, $output_file, $length); | |
| close $input; | |
| }; | |
| upload_s3_object( | |
| %arg, | |
| read_part => $read_part, | |
| skip_part => undef, | |
| ); | |
| } | |
| sub upload_s3_object_stream (%arg) { | |
| my $reader = $arg{reader}; | |
| my $read_part = sub ($relative_offset, $length, $output_file) { | |
| $reader->copy_exact_to_file($output_file, $length); | |
| }; | |
| my $skip_part = sub ($length) { | |
| $reader->skip_exact($length); | |
| }; | |
| upload_s3_object( | |
| %arg, | |
| read_part => $read_part, | |
| skip_part => $skip_part, | |
| ); | |
| } | |
| sub upload_s3_object (%arg) { | |
| my $state = $arg{state}; | |
| my $state_file = $arg{state_file}; | |
| my $bucket = $arg{bucket}; | |
| my $key = $arg{key}; | |
| my $size = $arg{size}; | |
| my $object_state = $state->{objects}{$key} //= { | |
| size => 0 + $size, | |
| logical_sha256 => $arg{logical_sha256}, | |
| segment_number => 0 + $arg{segment_number}, | |
| segment_count => 0 + $arg{segment_count}, | |
| status => 'pending', | |
| }; | |
| validate_object_state($object_state, %arg); | |
| if (remote_object_matches(%arg)) { | |
| $object_state->{status} = 'complete'; | |
| delete $object_state->{upload_id}; | |
| save_state($state_file, $state); | |
| if ($arg{skip_part}) { | |
| $arg{skip_part}->($size); | |
| } | |
| log_message("$key: already complete"); | |
| return; | |
| } | |
| # S3 multipart upload cannot represent an object with no uploaded bytes in | |
| # a portable way. Zero-byte logical files use PutObject as the sole | |
| # exception to the multipart-only rule. | |
| if ($size == 0) { | |
| my ($fh, $empty) = tempfile( | |
| 'empty-XXXXXX', | |
| DIR => $work_dir, | |
| UNLINK => 0, | |
| ); | |
| close $fh; | |
| aws_json( | |
| 's3api', 'put-object', | |
| '--bucket', $bucket, | |
| '--key', $key, | |
| '--body', $empty, | |
| '--metadata', metadata_string(%arg), | |
| ); | |
| unlink $empty; | |
| remote_object_matches(%arg) | |
| or fatal("$key: zero-byte object failed verification"); | |
| $object_state->{status} = 'complete'; | |
| save_state($state_file, $state); | |
| return; | |
| } | |
| my $part_size = choose_part_size($size); | |
| $object_state->{part_size} = 0 + $part_size; | |
| my ($upload_id, $parts) = prepare_multipart_upload( | |
| state => $state, | |
| state_file => $state_file, | |
| object_state => $object_state, | |
| bucket => $bucket, | |
| key => $key, | |
| metadata => metadata_string(%arg), | |
| ); | |
| my $part_count = int(($size + $part_size - 1) / $part_size); | |
| my $offset = 0; | |
| for my $part_number (1 .. $part_count) { | |
| my $length = $part_size; | |
| $length = $size - $offset if $offset + $length > $size; | |
| if ( | |
| $parts->{$part_number} | |
| && $parts->{$part_number}{Size} == $length | |
| ) { | |
| if ($arg{skip_part}) { | |
| $arg{skip_part}->($length); | |
| } | |
| log_message( | |
| "$key: part $part_number/$part_count already present" | |
| ); | |
| $offset += $length; | |
| next; | |
| } | |
| my ($fh, $part_file) = tempfile( | |
| 'part-XXXXXX', | |
| DIR => $work_dir, | |
| UNLINK => 0, | |
| ); | |
| close $fh; | |
| $arg{read_part}->($offset, $length, $part_file); | |
| my $actual = -s $part_file; | |
| $actual == $length | |
| or fatal( | |
| "$key: staged part $part_number has $actual bytes; " | |
| . "expected $length" | |
| ); | |
| log_message( | |
| "$key: uploading part $part_number/$part_count ($length bytes)" | |
| ); | |
| my $response = aws_json( | |
| 's3api', 'upload-part', | |
| '--bucket', $bucket, | |
| '--key', $key, | |
| '--upload-id', $upload_id, | |
| '--part-number', $part_number, | |
| '--body', $part_file, | |
| ); | |
| unlink $part_file; | |
| my $etag = $response->{ETag} | |
| // fatal("$key: UploadPart returned no ETag"); | |
| $parts->{$part_number} = { | |
| PartNumber => 0 + $part_number, | |
| ETag => $etag, | |
| Size => 0 + $length, | |
| }; | |
| $offset += $length; | |
| } | |
| my @completion_parts = map { | |
| +{ | |
| PartNumber => 0 + $_, | |
| ETag => $parts->{$_}{ETag}, | |
| } | |
| } 1 .. $part_count; | |
| my ($manifest_fh, $manifest_file) = tempfile( | |
| 'complete-XXXXXX', | |
| DIR => $work_dir, | |
| UNLINK => 0, | |
| ); | |
| binmode $manifest_fh; | |
| print {$manifest_fh} canonical_json({Parts => \@completion_parts}); | |
| close $manifest_fh | |
| or fatal("cannot close $manifest_file: $!"); | |
| aws_json( | |
| 's3api', 'complete-multipart-upload', | |
| '--bucket', $bucket, | |
| '--key', $key, | |
| '--upload-id', $upload_id, | |
| '--multipart-upload', "file://$manifest_file", | |
| ); | |
| unlink $manifest_file; | |
| remote_object_matches(%arg) | |
| or fatal("$key: completed object failed metadata/size verification"); | |
| $object_state->{status} = 'complete'; | |
| delete $object_state->{upload_id}; | |
| save_state($state_file, $state); | |
| log_message("$key: complete"); | |
| } | |
| sub prepare_multipart_upload (%arg) { | |
| my $state = $arg{state}; | |
| my $state_file = $arg{state_file}; | |
| my $object_state = $arg{object_state}; | |
| my $bucket = $arg{bucket}; | |
| my $key = $arg{key}; | |
| my $upload_id = $object_state->{upload_id}; | |
| if ($upload_id) { | |
| my ($ok, $response) = try_aws_json( | |
| 's3api', 'list-parts', | |
| '--bucket', $bucket, | |
| '--key', $key, | |
| '--upload-id', $upload_id, | |
| ); | |
| if ($ok) { | |
| return ( | |
| $upload_id, | |
| parts_hash($response->{Parts} // []), | |
| ); | |
| } | |
| log_message("$key: saved multipart upload no longer exists"); | |
| delete $object_state->{upload_id}; | |
| $object_state->{status} = 'pending'; | |
| save_state($state_file, $state); | |
| } | |
| my $response = aws_json( | |
| 's3api', 'create-multipart-upload', | |
| '--bucket', $bucket, | |
| '--key', $key, | |
| '--metadata', $arg{metadata}, | |
| ); | |
| $upload_id = $response->{UploadId} | |
| // fatal("$key: CreateMultipartUpload returned no UploadId"); | |
| $object_state->{upload_id} = $upload_id; | |
| $object_state->{status} = 'uploading'; | |
| save_state($state_file, $state); | |
| return ($upload_id, {}); | |
| } | |
| sub parts_hash ($parts) { | |
| my %result; | |
| for my $part (@$parts) { | |
| my $number = $part->{PartNumber}; | |
| next unless defined $number; | |
| $result{$number} = { | |
| PartNumber => 0 + $number, | |
| ETag => $part->{ETag}, | |
| Size => 0 + ($part->{Size} // 0), | |
| }; | |
| } | |
| return \%result; | |
| } | |
| sub choose_part_size ($object_size) { | |
| my $required = int(($object_size + MAX_PARTS - 1) / MAX_PARTS); | |
| my $selected = $opt{part_size}; | |
| $selected = $required if $required > $selected; | |
| $selected = MIN_PART_SIZE if $selected < MIN_PART_SIZE; | |
| my $mib = 1024 * 1024; | |
| $selected = int(($selected + $mib - 1) / $mib) * $mib; | |
| $selected <= MAX_PART_SIZE | |
| or fatal( | |
| "object segment of $object_size bytes cannot fit within " | |
| . "10,000 parts of at most 5 GiB" | |
| ); | |
| return $selected; | |
| } | |
| sub metadata_string (%arg) { | |
| return join ',', | |
| 'idu-complete=1', | |
| 'idu-sha256=' . $arg{logical_sha256}, | |
| 'idu-size=' . $arg{size}, | |
| 'idu-segment=' . $arg{segment_number}, | |
| 'idu-segments=' . $arg{segment_count}; | |
| } | |
| sub remote_object_matches (%arg) { | |
| my ($ok, $response) = try_aws_json( | |
| 's3api', 'head-object', | |
| '--bucket', $arg{bucket}, | |
| '--key', $arg{key}, | |
| ); | |
| return 0 unless $ok; | |
| my $metadata = $response->{Metadata} // {}; | |
| return | |
| ($response->{ContentLength} // -1) == $arg{size} | |
| && ($metadata->{'idu-complete'} // '') eq '1' | |
| && ($metadata->{'idu-sha256'} // '') eq $arg{logical_sha256} | |
| && ($metadata->{'idu-size'} // '') eq "$arg{size}" | |
| && ($metadata->{'idu-segment'} // '') eq "$arg{segment_number}" | |
| && ($metadata->{'idu-segments'} // '') eq "$arg{segment_count}"; | |
| } | |
| sub validate_object_state ($object_state, %arg) { | |
| my %expected = ( | |
| size => 0 + $arg{size}, | |
| logical_sha256 => $arg{logical_sha256}, | |
| segment_number => 0 + $arg{segment_number}, | |
| segment_count => 0 + $arg{segment_count}, | |
| ); | |
| for my $field (keys %expected) { | |
| next unless exists $object_state->{$field}; | |
| "$object_state->{$field}" eq "$expected{$field}" | |
| or fatal( | |
| "$arg{key}: state mismatch for $field; " | |
| . "remove the state file only after investigating" | |
| ); | |
| } | |
| } | |
| # =========================================================================== | |
| # TAR parsing | |
| # =========================================================================== | |
| sub classify_source ($filename) { | |
| open my $fh, '<:raw', $filename | |
| or fatal("cannot open $filename: $!"); | |
| my $prefix = ''; | |
| read($fh, $prefix, TAR_BLOCK); | |
| close $fh; | |
| if (length($prefix) == TAR_BLOCK && tar_header_is_valid($prefix)) { | |
| return 'tar'; | |
| } | |
| if (substr($prefix, 0, 2) eq "\x1f\x8b") { | |
| my $reader = eval { compressed_reader($filename, 'tar.gz') }; | |
| return 'tar.gz' | |
| if $reader && first_tar_header_is_valid($reader); | |
| return 'blob'; | |
| } | |
| if (substr($prefix, 0, 3) eq 'BZh') { | |
| my $reader = eval { compressed_reader($filename, 'tar.bz2') }; | |
| return 'tar.bz2' | |
| if $reader && first_tar_header_is_valid($reader); | |
| return 'blob'; | |
| } | |
| return 'blob'; | |
| } | |
| sub first_tar_header_is_valid ($reader) { | |
| my $block = $reader->read_exact_or_eof(TAR_BLOCK); | |
| return defined($block) | |
| && length($block) == TAR_BLOCK | |
| && tar_header_is_valid($block); | |
| } | |
| sub scan_tar (%arg) { | |
| my @members; | |
| iterate_tar( | |
| reader => $arg{reader}, | |
| on_regular_file => sub ($entry, $data_reader) { | |
| my $sha = Digest::SHA->new(256); | |
| my $remaining = $entry->{size}; | |
| while ($remaining > 0) { | |
| my $want = $remaining > 1024 * 1024 | |
| ? 1024 * 1024 | |
| : $remaining; | |
| my $chunk = $data_reader->read_exact($want); | |
| $sha->add($chunk); | |
| $remaining -= length $chunk; | |
| } | |
| push @members, { | |
| path => $entry->{path}, | |
| size => 0 + $entry->{size}, | |
| mtime => 0 + $entry->{mtime}, | |
| sha256 => $sha->hexdigest, | |
| ( | |
| $arg{record_offsets} | |
| ? (data_offset => 0 + $entry->{data_offset}) | |
| : () | |
| ), | |
| }; | |
| }, | |
| ); | |
| return \@members; | |
| } | |
| sub iterate_tar (%arg) { | |
| my $reader = $arg{reader}; | |
| my $callback = $arg{on_regular_file}; | |
| my %global_pax; | |
| my $next_pax = {}; | |
| my $long_name; | |
| my %seen_paths; | |
| my $zero_blocks = 0; | |
| while (1) { | |
| my $header_offset = $reader->tell_position; | |
| my $block = $reader->read_exact_or_eof(TAR_BLOCK); | |
| last unless defined $block; | |
| length($block) == TAR_BLOCK | |
| or fatal("truncated TAR header at offset $header_offset"); | |
| if ($block eq "\0" x TAR_BLOCK) { | |
| $zero_blocks++; | |
| last if $zero_blocks >= 2; | |
| next; | |
| } | |
| $zero_blocks = 0; | |
| my $entry = parse_tar_header($block, $header_offset); | |
| my $type = $entry->{type}; | |
| if ($type eq 'L') { | |
| $entry->{size} <= MAX_META_ENTRY | |
| or fatal("GNU long-name record is unreasonably large"); | |
| my $name = $reader->read_exact($entry->{size}); | |
| $name =~ s/\0.*\z//s; | |
| $name =~ s/\n\z//; | |
| $long_name = $name; | |
| skip_tar_padding($reader, $entry->{size}); | |
| next; | |
| } | |
| if ($type eq 'x' || $type eq 'g') { | |
| $entry->{size} <= MAX_META_ENTRY | |
| or fatal("PAX record is unreasonably large"); | |
| my $payload = $reader->read_exact($entry->{size}); | |
| my $pax = parse_pax_records($payload); | |
| if ($type eq 'g') { | |
| %global_pax = (%global_pax, %$pax); | |
| } | |
| else { | |
| $next_pax = $pax; | |
| } | |
| skip_tar_padding($reader, $entry->{size}); | |
| next; | |
| } | |
| my %pax = (%global_pax, %$next_pax); | |
| $next_pax = {}; | |
| $entry->{path} = $long_name if defined $long_name; | |
| $long_name = undef; | |
| $entry->{path} = $pax{path} if exists $pax{path}; | |
| $entry->{size} = 0 + $pax{size} if exists $pax{size}; | |
| $entry->{mtime} = int($pax{mtime}) if exists $pax{mtime}; | |
| for my $key (keys %pax) { | |
| if ($key =~ /\AGNU\.sparse/ || $key =~ /\ASCHILY\.realsize\z/) { | |
| fatal("sparse TAR member is not supported: $entry->{path}"); | |
| } | |
| } | |
| $entry->{path} = safe_tar_path($entry->{path}); | |
| $entry->{data_offset} = $reader->tell_position; | |
| if (is_regular_type($type)) { | |
| $seen_paths{$entry->{path}}++ | |
| and fatal("duplicate regular TAR member: $entry->{path}"); | |
| my $limited = LimitedReader->new($reader, $entry->{size}); | |
| $callback->($entry, $limited); | |
| $limited->remaining == 0 | |
| or fatal( | |
| "internal error: callback did not consume " | |
| . "$entry->{path}" | |
| ); | |
| skip_tar_padding($reader, $entry->{size}); | |
| } | |
| elsif ($type eq 'S') { | |
| fatal("GNU sparse TAR member is not supported: $entry->{path}"); | |
| } | |
| else { | |
| $reader->skip_exact($entry->{size}); | |
| skip_tar_padding($reader, $entry->{size}); | |
| } | |
| } | |
| } | |
| sub parse_tar_header ($block, $offset) { | |
| tar_header_is_valid($block) | |
| or fatal("invalid TAR header checksum at offset $offset"); | |
| my $name = tar_string(substr($block, 0, 100)); | |
| my $size = tar_number(substr($block, 124, 12)); | |
| my $mtime = tar_number(substr($block, 136, 12)); | |
| my $type = substr($block, 156, 1); | |
| my $prefix = tar_string(substr($block, 345, 155)); | |
| $type = "\0" if $type eq ''; | |
| my $path = length($prefix) ? "$prefix/$name" : $name; | |
| return { | |
| path => $path, | |
| size => 0 + $size, | |
| mtime => 0 + $mtime, | |
| type => $type, | |
| offset => 0 + $offset, | |
| }; | |
| } | |
| sub tar_header_is_valid ($block) { | |
| return 0 unless length($block) == TAR_BLOCK; | |
| return 1 if $block eq "\0" x TAR_BLOCK; | |
| my $stored = eval { tar_number(substr($block, 148, 8)) }; | |
| return 0 if $@; | |
| my $copy = $block; | |
| substr($copy, 148, 8, ' ' x 8); | |
| my $unsigned = 0; | |
| $unsigned += ord($_) for split //, $copy; | |
| return $stored == $unsigned; | |
| } | |
| sub tar_number ($field) { | |
| my @bytes = unpack('C*', $field); | |
| if ($bytes[0] & 0x80) { | |
| # POSIX/GNU base-256 encoding. Clear the signal bit and interpret the | |
| # remaining bytes as an unsigned big-endian integer. | |
| $bytes[0] &= 0x7f; | |
| my $value = 0; | |
| $value = ($value << 8) | $_ for @bytes; | |
| return $value; | |
| } | |
| $field =~ s/\0.*\z//s; | |
| $field =~ s/^\s+|\s+$//g; | |
| return 0 if $field eq ''; | |
| $field =~ /\A[0-7]+\z/ | |
| or fatal("invalid octal number in TAR header"); | |
| return oct($field); | |
| } | |
| sub tar_string ($field) { | |
| $field =~ s/\0.*\z//s; | |
| return $field; | |
| } | |
| sub parse_pax_records ($payload) { | |
| my %result; | |
| my $offset = 0; | |
| my $length = length $payload; | |
| while ($offset < $length) { | |
| my $space = index($payload, ' ', $offset); | |
| $space >= 0 or fatal("invalid PAX record"); | |
| my $record_length = substr($payload, $offset, $space - $offset); | |
| $record_length =~ /\A\d+\z/ | |
| or fatal("invalid PAX record length"); | |
| my $record = substr($payload, $offset, $record_length); | |
| length($record) == $record_length | |
| or fatal("truncated PAX record"); | |
| $record =~ s/\n\z//; | |
| $record =~ s/\A\d+ //; | |
| my ($key, $value) = split /=/, $record, 2; | |
| defined $value or fatal("invalid PAX key/value record"); | |
| $result{$key} = $value; | |
| $offset += $record_length; | |
| } | |
| return \%result; | |
| } | |
| sub skip_tar_padding ($reader, $size) { | |
| my $padding = (TAR_BLOCK - ($size % TAR_BLOCK)) % TAR_BLOCK; | |
| $reader->skip_exact($padding) if $padding; | |
| } | |
| sub is_regular_type ($type) { | |
| return $type eq "\0" || $type eq '0' || $type eq '7'; | |
| } | |
| sub safe_tar_path ($path) { | |
| defined $path && length $path | |
| or fatal("empty TAR member path"); | |
| $path =~ s{^\./}{}; | |
| $path !~ m{\A/} | |
| or fatal("absolute TAR member path rejected: $path"); | |
| $path !~ /[\r\n]/ | |
| or fatal("TAR member path contains a newline"); | |
| for my $part (split m{/+}, $path) { | |
| $part ne '..' | |
| or fatal("parent traversal in TAR member path: $path"); | |
| } | |
| return normalize_key($path); | |
| } | |
| # =========================================================================== | |
| # Reader classes | |
| # =========================================================================== | |
| package FileReader { | |
| use Fcntl qw(SEEK_CUR); | |
| sub new ($class, $filename) { | |
| open my $fh, '<:raw', $filename | |
| or main::fatal("cannot open $filename: $!"); | |
| return bless { | |
| fh => $fh, | |
| position => 0, | |
| }, $class; | |
| } | |
| sub read_exact_or_eof ($self, $length) { | |
| return undef if $length && eof($self->{fh}); | |
| my $data = ''; | |
| while (length($data) < $length) { | |
| my $count = read( | |
| $self->{fh}, | |
| my $chunk, | |
| $length - length($data), | |
| ); | |
| defined $count | |
| or main::fatal("read failed: $!"); | |
| last if $count == 0; | |
| $data .= $chunk; | |
| } | |
| $self->{position} += length $data; | |
| return $data; | |
| } | |
| sub read_exact ($self, $length) { | |
| my $data = $self->read_exact_or_eof($length); | |
| defined($data) && length($data) == $length | |
| or main::fatal("unexpected end of input"); | |
| return $data; | |
| } | |
| sub skip_exact ($self, $length) { | |
| return if $length == 0; | |
| my $new = sysseek($self->{fh}, $length, SEEK_CUR); | |
| defined $new or main::fatal("seek failed: $!"); | |
| $self->{position} += $length; | |
| } | |
| sub copy_exact_to_file ($self, $filename, $length) { | |
| open my $output, '>:raw', $filename | |
| or main::fatal("cannot create $filename: $!"); | |
| my $remaining = $length; | |
| while ($remaining > 0) { | |
| my $want = $remaining > 1024 * 1024 | |
| ? 1024 * 1024 | |
| : $remaining; | |
| my $data = $self->read_exact($want); | |
| print {$output} $data | |
| or main::fatal("cannot write $filename: $!"); | |
| $remaining -= length $data; | |
| } | |
| close $output | |
| or main::fatal("cannot close $filename: $!"); | |
| } | |
| sub tell_position ($self) { | |
| return $self->{position}; | |
| } | |
| } | |
| package StreamReader { | |
| sub new ($class, $fh) { | |
| return bless { | |
| fh => $fh, | |
| position => 0, | |
| }, $class; | |
| } | |
| sub read_exact_or_eof ($self, $length) { | |
| my $data = ''; | |
| while (length($data) < $length) { | |
| my $count = $self->{fh}->read( | |
| my $chunk, | |
| $length - length($data), | |
| ); | |
| defined $count | |
| or main::fatal("compressed read failed"); | |
| last if $count == 0; | |
| $data .= $chunk; | |
| } | |
| $self->{position} += length $data; | |
| return undef if $data eq '' && $length; | |
| return $data; | |
| } | |
| sub read_exact ($self, $length) { | |
| my $data = $self->read_exact_or_eof($length); | |
| defined($data) && length($data) == $length | |
| or main::fatal("unexpected end of compressed input"); | |
| return $data; | |
| } | |
| sub skip_exact ($self, $length) { | |
| my $remaining = $length; | |
| while ($remaining > 0) { | |
| my $want = $remaining > 1024 * 1024 | |
| ? 1024 * 1024 | |
| : $remaining; | |
| my $data = $self->read_exact($want); | |
| $remaining -= length $data; | |
| } | |
| } | |
| sub copy_exact_to_file ($self, $filename, $length) { | |
| open my $output, '>:raw', $filename | |
| or main::fatal("cannot create $filename: $!"); | |
| my $remaining = $length; | |
| while ($remaining > 0) { | |
| my $want = $remaining > 1024 * 1024 | |
| ? 1024 * 1024 | |
| : $remaining; | |
| my $data = $self->read_exact($want); | |
| print {$output} $data | |
| or main::fatal("cannot write $filename: $!"); | |
| $remaining -= length $data; | |
| } | |
| close $output | |
| or main::fatal("cannot close $filename: $!"); | |
| } | |
| sub tell_position ($self) { | |
| return $self->{position}; | |
| } | |
| } | |
| package LimitedReader { | |
| sub new ($class, $reader, $remaining) { | |
| return bless { | |
| reader => $reader, | |
| remaining => $remaining, | |
| }, $class; | |
| } | |
| sub remaining ($self) { | |
| return $self->{remaining}; | |
| } | |
| sub read_exact ($self, $length) { | |
| $length <= $self->{remaining} | |
| or main::fatal("attempt to read beyond TAR member"); | |
| my $data = $self->{reader}->read_exact($length); | |
| $self->{remaining} -= length $data; | |
| return $data; | |
| } | |
| sub skip_exact ($self, $length) { | |
| $length <= $self->{remaining} | |
| or main::fatal("attempt to skip beyond TAR member"); | |
| $self->{reader}->skip_exact($length); | |
| $self->{remaining} -= $length; | |
| } | |
| sub copy_exact_to_file ($self, $filename, $length) { | |
| $length <= $self->{remaining} | |
| or main::fatal("attempt to copy beyond TAR member"); | |
| $self->{reader}->copy_exact_to_file($filename, $length); | |
| $self->{remaining} -= $length; | |
| } | |
| sub tell_position ($self) { | |
| return $self->{reader}->tell_position; | |
| } | |
| } | |
| package main; | |
| sub compressed_reader ($filename, $kind) { | |
| if ($kind eq 'tar.gz') { | |
| my $fh = IO::Uncompress::Gunzip->new( | |
| $filename, | |
| Transparent => 0, | |
| MultiStream => 1, | |
| ) or fatal("cannot open gzip stream $filename: $GunzipError"); | |
| return StreamReader->new($fh); | |
| } | |
| if ($kind eq 'tar.bz2') { | |
| my $fh = IO::Uncompress::Bunzip2->new( | |
| $filename, | |
| Transparent => 0, | |
| MultiStream => 1, | |
| ) or fatal("cannot open bzip2 stream $filename: $Bunzip2Error"); | |
| return StreamReader->new($fh); | |
| } | |
| fatal("unsupported compressed TAR kind: $kind"); | |
| } | |
| # =========================================================================== | |
| # AWS CLI | |
| # =========================================================================== | |
| sub aws_json (@arguments) { | |
| my ($ok, $value, $error) = run_aws_json(@arguments); | |
| $ok or fatal("AWS CLI failed: $error"); | |
| return $value; | |
| } | |
| sub try_aws_json (@arguments) { | |
| my ($ok, $value) = run_aws_json(@arguments); | |
| return ($ok, $value); | |
| } | |
| sub run_aws_json (@arguments) { | |
| my @command = ('aws'); | |
| if (defined $opt{endpoint_url}) { | |
| push @command, '--endpoint-url', $opt{endpoint_url}; | |
| } | |
| push @command, @arguments, '--output', 'json'; | |
| my $stderr = gensym; | |
| my $pid = open3(my $stdin, my $stdout, $stderr, @command); | |
| close $stdin; | |
| local $/; | |
| my $out = <$stdout> // ''; | |
| my $err = <$stderr> // ''; | |
| waitpid($pid, 0); | |
| my $status = $? >> 8; | |
| if ($status != 0) { | |
| $err =~ s/\s+\z//; | |
| return (0, undef, $err || "exit status $status"); | |
| } | |
| return (1, {}) if $out =~ /\A\s*\z/; | |
| my $decoded = eval { JSON::PP->new->decode($out) }; | |
| if (!$decoded && $@) { | |
| return (0, undef, "invalid JSON from AWS CLI: $@"); | |
| } | |
| return (1, $decoded); | |
| } | |
| # =========================================================================== | |
| # State and utility functions | |
| # =========================================================================== | |
| sub load_state ($filename) { | |
| return undef unless -f $filename; | |
| open my $fh, '<:raw', $filename | |
| or fatal("cannot open state file $filename: $!"); | |
| local $/; | |
| my $json = <$fh>; | |
| close $fh; | |
| my $state = eval { JSON::PP->new->decode($json) }; | |
| $state && ref($state) eq 'HASH' | |
| or fatal("invalid state file $filename"); | |
| return $state; | |
| } | |
| sub save_state ($filename, $state) { | |
| my $temporary = "$filename.tmp.$$"; | |
| open my $fh, '>:raw', $temporary | |
| or fatal("cannot create $temporary: $!"); | |
| print {$fh} canonical_json($state) | |
| or fatal("cannot write $temporary: $!"); | |
| close $fh | |
| or fatal("cannot close $temporary: $!"); | |
| rename $temporary, $filename | |
| or fatal("cannot rename $temporary to $filename: $!"); | |
| } | |
| sub validate_existing_state ($state, %expected) { | |
| $state->{version} == STATE_VERSION | |
| or fatal("unsupported state file version"); | |
| my $source = $state->{source} // {}; | |
| for my $field (qw(path size mtime bucket destination)) { | |
| my $wanted = $expected{ | |
| $field eq 'path' ? 'source' | |
| : $field eq 'size' ? 'source_size' | |
| : $field eq 'mtime' ? 'source_mtime' | |
| : $field | |
| }; | |
| "$source->{$field}" eq "$wanted" | |
| or fatal( | |
| "source or destination changed since the state file " | |
| . "was created ($field mismatch)" | |
| ); | |
| } | |
| } | |
| sub canonical_json ($value) { | |
| return JSON::PP->new->canonical(1)->pretty(1)->encode($value); | |
| } | |
| sub sha256_range ($filename, $offset, $length) { | |
| open my $fh, '<:raw', $filename | |
| or fatal("cannot open $filename: $!"); | |
| sysseek($fh, $offset, SEEK_SET) | |
| // fatal("cannot seek $filename: $!"); | |
| my $sha = Digest::SHA->new(256); | |
| my $remaining = $length; | |
| while ($remaining > 0) { | |
| my $want = $remaining > 1024 * 1024 | |
| ? 1024 * 1024 | |
| : $remaining; | |
| my $count = read($fh, my $buffer, $want); | |
| defined $count or fatal("cannot read $filename: $!"); | |
| $count > 0 or fatal("unexpected EOF in $filename"); | |
| $sha->add(substr($buffer, 0, $count)); | |
| $remaining -= $count; | |
| } | |
| close $fh; | |
| return $sha->hexdigest; | |
| } | |
| sub copy_exact ($input, $output_file, $length) { | |
| open my $output, '>:raw', $output_file | |
| or fatal("cannot create $output_file: $!"); | |
| my $remaining = $length; | |
| while ($remaining > 0) { | |
| my $want = $remaining > 1024 * 1024 | |
| ? 1024 * 1024 | |
| : $remaining; | |
| my $count = read($input, my $buffer, $want); | |
| defined $count or fatal("read failed: $!"); | |
| $count > 0 or fatal("unexpected EOF while staging part"); | |
| print {$output} substr($buffer, 0, $count) | |
| or fatal("cannot write $output_file: $!"); | |
| $remaining -= $count; | |
| } | |
| close $output | |
| or fatal("cannot close $output_file: $!"); | |
| } | |
| sub parse_size ($text) { | |
| defined $text or fatal("missing size"); | |
| $text =~ /\A(\d+)([KMGT]iB|B)?\z/i | |
| or fatal("invalid size: $text"); | |
| my ($number, $suffix) = ($1, uc($2 // 'B')); | |
| my %multiplier = ( | |
| B => 1, | |
| KIB => 1024, | |
| MIB => 1024 ** 2, | |
| GIB => 1024 ** 3, | |
| TIB => 1024 ** 4, | |
| ); | |
| return $number * $multiplier{$suffix}; | |
| } | |
| sub tar_destination_prefix ($destination) { | |
| my $directory = dirname($destination); | |
| $directory = '' if $directory eq '.'; | |
| my $name = basename($destination); | |
| $name =~ s/\.(?:tar\.(?:gz|bz2)|tgz|tbz2?|tar)\z//i; | |
| return normalize_key( | |
| length($directory) ? "$directory/$name" : $name | |
| ); | |
| } | |
| sub normalize_key ($key) { | |
| $key =~ s{\A/+}{}; | |
| $key =~ s{//+}{/}g; | |
| return $key; | |
| } | |
| sub command_exists ($command) { | |
| for my $directory (File::Spec->path) { | |
| my $candidate = File::Spec->catfile($directory, $command); | |
| return 1 if -x $candidate; | |
| } | |
| return 0; | |
| } | |
| sub log_message ($message) { | |
| print STDERR "s3-upload: $message\n"; | |
| } | |
| sub fatal ($message) { | |
| die "s3-upload: ERROR: $message\n"; | |
| } | |
| sub usage ($status) { | |
| print STDERR <<'USAGE'; | |
| Usage: | |
| s3-upload.pl [options] SOURCE BUCKET [DESTINATION] | |
| Options: | |
| --state-dir DIR | |
| --work-dir DIR | |
| --part-size SIZE Default: 64MiB | |
| --object-limit SIZE Default: 5TiB | |
| --endpoint-url URL | |
| --keep-state | |
| --help | |
| Supported TAR inputs: | |
| uncompressed TAR | |
| gzip-compressed TAR | |
| bzip2-compressed TAR | |
| Examples: | |
| s3-upload.pl disk.img my-bucket backups/disk.img | |
| s3-upload.pl archive.tar.gz my-bucket imports/archive.tar.gz | |
| USAGE | |
| exit $status; | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| #!/usr/bin/env bash | |
| set -Eeuo pipefail | |
| # s3-upload | |
| # | |
| # Upload one source file to S3 using multipart upload only. | |
| # | |
| # * A generic blob becomes one S3 object. | |
| # * A tar archive becomes a prefix named after the archive; each regular member | |
| # becomes an object below that prefix. | |
| # * One SQLite database is dedicated to this source upload. | |
| # * Incomplete multipart uploads and completed parts are resumed. | |
| # * The database and staging directory are removed only after the entire source | |
| # is complete. | |
| # | |
| # Requirements: | |
| # bash, GNU tar, aws CLI v2, sqlite3, perl with JSON::PP and Digest::SHA | |
| # | |
| # Usage: | |
| # s3-upload SOURCE BUCKET [DESTINATION] | |
| # | |
| # Examples: | |
| # ./s3-upload disk.img my-bucket backups/disk.img | |
| # ./s3-upload logs.tar.gz my-bucket imports/logs.tar.gz | |
| # | |
| # For a tar source whose destination is imports/logs.tar.gz, member a/b.txt is | |
| # stored as: | |
| # s3://my-bucket/imports/logs/a/b.txt | |
| # | |
| # Optional environment: | |
| # PART_SIZE_MIB=64 | |
| # STATE_DIR=/persistent/state | |
| # AWS_ENDPOINT_URL=http://127.0.0.1:9000 | |
| # KEEP_STATE_ON_SUCCESS=1 | |
| # FAIL_AFTER_PARTS=N # integration-test fault injection | |
| PROGRAM=${0##*/} | |
| MIN_PART_SIZE=$((5 * 1024 * 1024)) | |
| MAX_PART_SIZE=$((5 * 1024 * 1024 * 1024)) | |
| MAX_PARTS=10000 | |
| DEFAULT_PART_SIZE=$((64 * 1024 * 1024)) | |
| log() { | |
| printf '%s: %s\n' "$PROGRAM" "$*" >&2 | |
| } | |
| die() { | |
| log "ERROR: $*" | |
| exit 1 | |
| } | |
| require_command() { | |
| command -v "$1" >/dev/null 2>&1 || die "required command not found: $1" | |
| } | |
| sql_quote() { | |
| # SQLite string literal, including surrounding quotes. | |
| local s=${1//\'/\'\'} | |
| printf "'%s'" "$s" | |
| } | |
| aws_s3api() { | |
| local args=() | |
| if [[ -n ${AWS_ENDPOINT_URL:-} ]]; then | |
| args+=(--endpoint-url "$AWS_ENDPOINT_URL") | |
| fi | |
| aws "${args[@]}" s3api "$@" | |
| } | |
| sha256_file() { | |
| perl -MDigest::SHA -e ' | |
| my $sha = Digest::SHA->new(256); | |
| $sha->addfile($ARGV[0]); | |
| print $sha->hexdigest, "\n"; | |
| ' "$1" | |
| } | |
| sha256_text() { | |
| perl -MDigest::SHA=sha256_hex -e ' | |
| local $/; | |
| print sha256_hex(<STDIN>), "\n"; | |
| ' | |
| } | |
| json_string() { | |
| perl -MJSON::PP -e 'print encode_json($ARGV[0])' "$1" | |
| } | |
| archive_base_name() { | |
| local name=${1##*/} | |
| case "$name" in | |
| *.tar.gz) printf '%s\n' "${name%.tar.gz}" ;; | |
| *.tar.bz2) printf '%s\n' "${name%.tar.bz2}" ;; | |
| *.tar.xz) printf '%s\n' "${name%.tar.xz}" ;; | |
| *.tar.zst) printf '%s\n' "${name%.tar.zst}" ;; | |
| *.tar.lz) printf '%s\n' "${name%.tar.lz}" ;; | |
| *.tar.lzma) printf '%s\n' "${name%.tar.lzma}" ;; | |
| *.tgz) printf '%s\n' "${name%.tgz}" ;; | |
| *.tbz) printf '%s\n' "${name%.tbz}" ;; | |
| *.tbz2) printf '%s\n' "${name%.tbz2}" ;; | |
| *.txz) printf '%s\n' "${name%.txz}" ;; | |
| *.tzst) printf '%s\n' "${name%.tzst}" ;; | |
| *.tar) printf '%s\n' "${name%.tar}" ;; | |
| *) printf '%s\n' "$name" ;; | |
| esac | |
| } | |
| normalize_key() { | |
| local key=$1 | |
| while [[ $key == /* ]]; do key=${key#/}; done | |
| while [[ $key == *//* ]]; do key=${key//\/\//\/}; done | |
| printf '%s\n' "$key" | |
| } | |
| safe_member_name() { | |
| local name=$1 | |
| name=${name#./} | |
| [[ -n $name ]] || return 1 | |
| [[ $name != /* ]] || return 1 | |
| [[ $name != *$'\n'* ]] || return 1 | |
| [[ $name != *$'\r'* ]] || return 1 | |
| local component | |
| IFS='/' read -r -a components <<<"$name" | |
| for component in "${components[@]}"; do | |
| [[ $component != ".." ]] || return 1 | |
| done | |
| printf '%s\n' "$name" | |
| } | |
| choose_part_size() { | |
| local size=$1 | |
| local configured=$(( ${PART_SIZE_MIB:-64} * 1024 * 1024 )) | |
| local required=$(( (size + MAX_PARTS - 1) / MAX_PARTS )) | |
| local mib=$((1024 * 1024)) | |
| local selected=$configured | |
| (( selected < MIN_PART_SIZE )) && selected=$MIN_PART_SIZE | |
| (( required > selected )) && selected=$required | |
| # Round upward to a whole MiB. | |
| selected=$(( ((selected + mib - 1) / mib) * mib )) | |
| (( selected <= MAX_PART_SIZE )) || | |
| die "object is too large for S3 multipart limits: $size bytes" | |
| printf '%s\n' "$selected" | |
| } | |
| init_database() { | |
| sqlite3 "$DB" <<'SQL' | |
| PRAGMA journal_mode=WAL; | |
| PRAGMA synchronous=FULL; | |
| PRAGMA foreign_keys=ON; | |
| PRAGMA busy_timeout=30000; | |
| CREATE TABLE IF NOT EXISTS source ( | |
| singleton INTEGER PRIMARY KEY CHECK (singleton = 1), | |
| source_path TEXT NOT NULL, | |
| source_size INTEGER NOT NULL, | |
| source_mtime INTEGER NOT NULL, | |
| source_sha256 TEXT NOT NULL, | |
| source_id TEXT NOT NULL, | |
| source_kind TEXT NOT NULL, | |
| bucket TEXT NOT NULL, | |
| destination TEXT NOT NULL, | |
| created_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP | |
| ); | |
| CREATE TABLE IF NOT EXISTS object_upload ( | |
| object_key TEXT PRIMARY KEY, | |
| member_name TEXT, | |
| object_size INTEGER NOT NULL, | |
| member_mtime INTEGER, | |
| part_size INTEGER NOT NULL, | |
| upload_id TEXT, | |
| status TEXT NOT NULL DEFAULT 'pending' | |
| CHECK (status IN ('pending','uploading','complete')), | |
| etag TEXT, | |
| updated_at TEXT NOT NULL DEFAULT CURRENT_TIMESTAMP | |
| ); | |
| CREATE TABLE IF NOT EXISTS part ( | |
| object_key TEXT NOT NULL REFERENCES object_upload(object_key) | |
| ON DELETE CASCADE, | |
| part_number INTEGER NOT NULL, | |
| part_size INTEGER NOT NULL, | |
| etag TEXT NOT NULL, | |
| PRIMARY KEY (object_key, part_number) | |
| ); | |
| SQL | |
| } | |
| write_source_record() { | |
| local existing | |
| existing=$(sqlite3 "$DB" \ | |
| "SELECT source_id || '|' || bucket || '|' || destination || | |
| '|' || source_kind FROM source WHERE singleton=1;") | |
| local expected="$SOURCE_ID|$BUCKET|$DESTINATION|$SOURCE_KIND" | |
| if [[ -n $existing && $existing != "$expected" ]]; then | |
| die "state database belongs to a different source or destination: $DB" | |
| fi | |
| sqlite3 "$DB" " | |
| INSERT OR IGNORE INTO source( | |
| singleton, source_path, source_size, source_mtime, | |
| source_sha256, source_id, source_kind, bucket, destination | |
| ) VALUES ( | |
| 1, | |
| $(sql_quote "$SOURCE"), | |
| $SOURCE_SIZE, | |
| $SOURCE_MTIME, | |
| $(sql_quote "$SOURCE_SHA256"), | |
| $(sql_quote "$SOURCE_ID"), | |
| $(sql_quote "$SOURCE_KIND"), | |
| $(sql_quote "$BUCKET"), | |
| $(sql_quote "$DESTINATION") | |
| ); | |
| " | |
| } | |
| head_matches() { | |
| local key=$1 | |
| local member=$2 | |
| local size=$3 | |
| local response | |
| response=$(aws_s3api head-object \ | |
| --bucket "$BUCKET" \ | |
| --key "$key" \ | |
| --output json 2>/dev/null) || return 1 | |
| EXPECT_SOURCE_ID=$SOURCE_ID \ | |
| EXPECT_MEMBER=$member \ | |
| EXPECT_SIZE=$size \ | |
| perl -MJSON::PP -e ' | |
| local $/; | |
| my $j = decode_json(<STDIN>); | |
| my $m = $j->{Metadata} || {}; | |
| exit !( | |
| ($j->{ContentLength} // -1) == $ENV{EXPECT_SIZE} | |
| && ($m->{"upload-source-id"} // "") eq $ENV{EXPECT_SOURCE_ID} | |
| && ($m->{"upload-member"} // "") eq $ENV{EXPECT_MEMBER} | |
| && ($m->{"upload-complete"} // "") eq "1" | |
| ); | |
| ' <<<"$response" | |
| } | |
| ensure_object_row() { | |
| local key=$1 member=$2 size=$3 mtime=$4 part_size=$5 | |
| sqlite3 "$DB" " | |
| INSERT INTO object_upload( | |
| object_key, member_name, object_size, member_mtime, part_size | |
| ) VALUES ( | |
| $(sql_quote "$key"), | |
| $(sql_quote "$member"), | |
| $size, | |
| $mtime, | |
| $part_size | |
| ) | |
| ON CONFLICT(object_key) DO UPDATE SET | |
| member_name = excluded.member_name, | |
| object_size = excluded.object_size, | |
| member_mtime = excluded.member_mtime, | |
| part_size = excluded.part_size, | |
| updated_at = CURRENT_TIMESTAMP | |
| WHERE object_upload.status != 'complete'; | |
| " | |
| } | |
| create_upload() { | |
| local key=$1 member=$2 size=$3 | |
| local metadata upload_id | |
| # Values are deliberately limited to hex, decimal, and URI-safe text. | |
| # Member identity itself is kept in SQLite; the metadata member value is a | |
| # SHA-256 so arbitrary tar names never have to fit an HTTP metadata header. | |
| local member_id | |
| member_id=$(printf '%s' "$member" | sha256_text) | |
| metadata="upload-source-id=$SOURCE_ID,upload-member=$member_id,upload-complete=1,upload-size=$size" | |
| upload_id=$(aws_s3api create-multipart-upload \ | |
| --bucket "$BUCKET" \ | |
| --key "$key" \ | |
| --metadata "$metadata" \ | |
| --query UploadId \ | |
| --output text) | |
| sqlite3 "$DB" " | |
| UPDATE object_upload | |
| SET upload_id=$(sql_quote "$upload_id"), | |
| status='uploading', | |
| updated_at=CURRENT_TIMESTAMP | |
| WHERE object_key=$(sql_quote "$key"); | |
| " | |
| printf '%s\n' "$upload_id" | |
| } | |
| member_identity() { | |
| printf '%s' "$1" | sha256_text | |
| } | |
| head_matches_member() { | |
| local key=$1 member=$2 size=$3 | |
| head_matches "$key" "$(member_identity "$member")" "$size" | |
| } | |
| load_upload_id() { | |
| local key=$1 | |
| sqlite3 "$DB" \ | |
| "SELECT COALESCE(upload_id,'') FROM object_upload | |
| WHERE object_key=$(sql_quote "$key");" | |
| } | |
| clear_invalid_upload() { | |
| local key=$1 | |
| sqlite3 "$DB" " | |
| BEGIN IMMEDIATE; | |
| DELETE FROM part WHERE object_key=$(sql_quote "$key"); | |
| UPDATE object_upload | |
| SET upload_id=NULL, status='pending', updated_at=CURRENT_TIMESTAMP | |
| WHERE object_key=$(sql_quote "$key"); | |
| COMMIT; | |
| " | |
| } | |
| reconcile_parts_from_s3() { | |
| local key=$1 upload_id=$2 response | |
| # AWS CLI paginates list-parts automatically unless --no-paginate is used. | |
| if ! response=$(aws_s3api list-parts \ | |
| --bucket "$BUCKET" \ | |
| --key "$key" \ | |
| --upload-id "$upload_id" \ | |
| --output json 2>/dev/null) | |
| then | |
| return 1 | |
| fi | |
| local sql_file="$WORK/reconcile.$$.sql" | |
| { | |
| printf 'BEGIN IMMEDIATE;\n' | |
| printf 'DELETE FROM part WHERE object_key=%s;\n' "$(sql_quote "$key")" | |
| OBJECT_KEY=$key perl -MJSON::PP -e ' | |
| local $/; | |
| my $j = decode_json(<STDIN>); | |
| my $key = $ENV{OBJECT_KEY}; | |
| my $q = chr 39; | |
| $key =~ s/$q/$q$q/g; | |
| for my $p (@{$j->{Parts} || []}) { | |
| my $etag = $p->{ETag} // q{}; | |
| $etag =~ s/$q/$q$q/g; | |
| printf "INSERT INTO part(object_key,part_number,part_size,etag)" | |
| . " VALUES(%s%s%s,%d,%d,%s%s%s);\n", | |
| $q, $key, $q, $p->{PartNumber}, $p->{Size}, | |
| $q, $etag, $q; | |
| } | |
| ' <<<"$response" | |
| printf 'COMMIT;\n' | |
| } >"$sql_file" | |
| sqlite3 "$DB" <"$sql_file" | |
| rm -f "$sql_file" | |
| } | |
| record_part() { | |
| local key=$1 number=$2 size=$3 etag=$4 | |
| sqlite3 "$DB" " | |
| BEGIN IMMEDIATE; | |
| INSERT INTO part(object_key,part_number,part_size,etag) | |
| VALUES ( | |
| $(sql_quote "$key"), $number, $size, $(sql_quote "$etag") | |
| ) | |
| ON CONFLICT(object_key,part_number) DO UPDATE SET | |
| part_size=excluded.part_size, | |
| etag=excluded.etag; | |
| UPDATE object_upload SET updated_at=CURRENT_TIMESTAMP | |
| WHERE object_key=$(sql_quote "$key"); | |
| COMMIT; | |
| " | |
| } | |
| part_etag() { | |
| local key=$1 number=$2 | |
| sqlite3 "$DB" \ | |
| "SELECT COALESCE(etag,'') FROM part | |
| WHERE object_key=$(sql_quote "$key") AND part_number=$number;" | |
| } | |
| part_recorded_size() { | |
| local key=$1 number=$2 | |
| sqlite3 "$DB" \ | |
| "SELECT COALESCE(part_size,-1) FROM part | |
| WHERE object_key=$(sql_quote "$key") AND part_number=$number;" | |
| } | |
| # Avoid requiring DBI, which is not core, by replacing build_completion_json | |
| # with a sqlite3 TSV stream parsed by core JSON::PP. | |
| build_completion_json_core() { | |
| local key=$1 output=$2 | |
| sqlite3 -separator $'\t' "$DB" \ | |
| "SELECT part_number, etag FROM part | |
| WHERE object_key=$(sql_quote "$key") ORDER BY part_number;" | | |
| perl -MJSON::PP -e ' | |
| my @parts; | |
| while (<STDIN>) { | |
| chomp; | |
| my ($n, $etag) = split /\t/, $_, 2; | |
| push @parts, {PartNumber => 0 + $n, ETag => $etag}; | |
| } | |
| print encode_json({Parts => \@parts}), "\n"; | |
| ' >"$output" | |
| } | |
| complete_object() { | |
| local key=$1 upload_id=$2 expected_parts=$3 | |
| local actual_parts completion response etag | |
| actual_parts=$(sqlite3 "$DB" \ | |
| "SELECT COUNT(*) FROM part WHERE object_key=$(sql_quote "$key");") | |
| (( actual_parts == expected_parts )) || | |
| die "$key: expected $expected_parts parts, state has $actual_parts" | |
| completion="$WORK/complete.$$.json" | |
| build_completion_json_core "$key" "$completion" | |
| response=$(aws_s3api complete-multipart-upload \ | |
| --bucket "$BUCKET" \ | |
| --key "$key" \ | |
| --upload-id "$upload_id" \ | |
| --multipart-upload "file://$completion" \ | |
| --output json) | |
| rm -f "$completion" | |
| etag=$(perl -MJSON::PP -e ' | |
| local $/; | |
| my $j=decode_json(<STDIN>); | |
| print $j->{ETag} // q{}; | |
| ' <<<"$response") | |
| sqlite3 "$DB" " | |
| UPDATE object_upload | |
| SET status='complete', | |
| etag=$(sql_quote "$etag"), | |
| updated_at=CURRENT_TIMESTAMP | |
| WHERE object_key=$(sql_quote "$key"); | |
| " | |
| } | |
| maybe_fault_inject() { | |
| [[ -n ${FAIL_AFTER_PARTS:-} ]] || return 0 | |
| local count_file="$WORK/fault-count" | |
| local count=0 | |
| [[ -f $count_file ]] && count=$(<"$count_file") | |
| count=$((count + 1)) | |
| printf '%s\n' "$count" >"$count_file" | |
| if (( count >= FAIL_AFTER_PARTS )); then | |
| die "intentional test failure after $count uploaded part(s)" | |
| fi | |
| } | |
| upload_staged_part() { | |
| local key=$1 upload_id=$2 number=$3 staged=$4 size=$5 | |
| local etag | |
| etag=$(aws_s3api upload-part \ | |
| --bucket "$BUCKET" \ | |
| --key "$key" \ | |
| --upload-id "$upload_id" \ | |
| --part-number "$number" \ | |
| --body "$staged" \ | |
| --query ETag \ | |
| --output text) | |
| # The part is durable in S3 before this transaction is committed. If the | |
| # process dies in this tiny interval, the next run's list-parts reconciliation | |
| # rediscovers it, so it still need not be retransmitted. | |
| record_part "$key" "$number" "$size" "$etag" | |
| maybe_fault_inject | |
| } | |
| prepare_object_upload() { | |
| local key=$1 member=$2 size=$3 mtime=$4 part_size=$5 | |
| local upload_id | |
| ensure_object_row "$key" "$member" "$size" "$mtime" "$part_size" | |
| if head_matches_member "$key" "$member" "$size"; then | |
| sqlite3 "$DB" " | |
| UPDATE object_upload SET status='complete', updated_at=CURRENT_TIMESTAMP | |
| WHERE object_key=$(sql_quote "$key"); | |
| " | |
| printf '\n' | |
| return 0 | |
| fi | |
| upload_id=$(load_upload_id "$key") | |
| if [[ -n $upload_id ]]; then | |
| if ! reconcile_parts_from_s3 "$key" "$upload_id"; then | |
| log "$key: saved multipart upload no longer exists; creating another" | |
| clear_invalid_upload "$key" | |
| upload_id= | |
| fi | |
| fi | |
| if [[ -z $upload_id ]]; then | |
| upload_id=$(create_upload "$key" "$member" "$size") | |
| reconcile_parts_from_s3 "$key" "$upload_id" || | |
| die "$key: cannot inspect newly created multipart upload" | |
| fi | |
| printf '%s\n' "$upload_id" | |
| } | |
| upload_blob_object() { | |
| local key=$1 file=$2 member=$3 size=$4 mtime=$5 | |
| local part_size upload_id part_count part_number offset bytes staged | |
| local recorded_size | |
| part_size=$(choose_part_size "$size") | |
| upload_id=$(prepare_object_upload "$key" "$member" "$size" "$mtime" "$part_size") | |
| if [[ -z $upload_id ]]; then | |
| log "$key: already complete" | |
| return | |
| fi | |
| part_count=$(( size == 0 ? 1 : (size + part_size - 1) / part_size )) | |
| offset=0 | |
| for ((part_number=1; part_number<=part_count; part_number++)); do | |
| if (( size == 0 )); then | |
| bytes=0 | |
| else | |
| bytes=$part_size | |
| (( offset + bytes > size )) && bytes=$((size - offset)) | |
| fi | |
| recorded_size=$(part_recorded_size "$key" "$part_number") | |
| if (( recorded_size == bytes )); then | |
| log "$key: part $part_number/$part_count already present" | |
| offset=$((offset + bytes)) | |
| continue | |
| fi | |
| staged="$WORK/part.bin" | |
| rm -f "$staged" | |
| if (( bytes == 0 )); then | |
| : >"$staged" | |
| else | |
| # bs=1M keeps skip/count efficient and byte-exact via skip_bytes/count_bytes. | |
| dd if="$file" of="$staged" \ | |
| bs=1M skip="$offset" count="$bytes" \ | |
| iflag=skip_bytes,count_bytes status=none | |
| fi | |
| [[ $(wc -c <"$staged" | tr -d '[:space:]') == "$bytes" ]] || | |
| die "$key: failed to stage part $part_number" | |
| log "$key: uploading part $part_number/$part_count ($bytes bytes)" | |
| upload_staged_part "$key" "$upload_id" "$part_number" "$staged" "$bytes" | |
| rm -f "$staged" | |
| offset=$((offset + bytes)) | |
| done | |
| complete_object "$key" "$upload_id" "$part_count" | |
| head_matches_member "$key" "$member" "$size" || | |
| die "$key: completed object failed metadata/size verification" | |
| log "$key: complete" | |
| } | |
| upload_stream_object() { | |
| local key=$1 member=$2 size=$3 mtime=$4 | |
| local part_size upload_id part_count part_number bytes staged | |
| local recorded_size actual | |
| part_size=$(choose_part_size "$size") | |
| upload_id=$(prepare_object_upload "$key" "$member" "$size" "$mtime" "$part_size") | |
| if [[ -z $upload_id ]]; then | |
| # tar is still writing this member to our stdin. Drain it so tar does not | |
| # receive SIGPIPE and can continue to the next member. | |
| cat >/dev/null | |
| log "$key: already complete" | |
| return | |
| fi | |
| part_count=$(( size == 0 ? 1 : (size + part_size - 1) / part_size )) | |
| for ((part_number=1; part_number<=part_count; part_number++)); do | |
| if (( size == 0 )); then | |
| bytes=0 | |
| elif (( part_number < part_count )); then | |
| bytes=$part_size | |
| else | |
| bytes=$((size - (part_number - 1) * part_size)) | |
| fi | |
| staged="$WORK/part.bin" | |
| rm -f "$staged" | |
| if (( bytes == 0 )); then | |
| : >"$staged" | |
| else | |
| # fullblock is necessary because stdin is a pipe from GNU tar. | |
| dd of="$staged" bs="$bytes" count=1 iflag=fullblock status=none | |
| fi | |
| actual=$(wc -c <"$staged" | tr -d '[:space:]') | |
| (( actual == bytes )) || | |
| die "$key: unexpected EOF in member $member part $part_number; expected $bytes, got $actual" | |
| recorded_size=$(part_recorded_size "$key" "$part_number") | |
| if (( recorded_size == bytes )); then | |
| log "$key: part $part_number/$part_count already present; consumed and discarded" | |
| rm -f "$staged" | |
| continue | |
| fi | |
| log "$key: uploading part $part_number/$part_count ($bytes bytes)" | |
| upload_staged_part "$key" "$upload_id" "$part_number" "$staged" "$bytes" | |
| rm -f "$staged" | |
| done | |
| complete_object "$key" "$upload_id" "$part_count" | |
| head_matches_member "$key" "$member" "$size" || | |
| die "$key: completed object failed metadata/size verification" | |
| log "$key: complete" | |
| } | |
| upload_completion_marker() { | |
| local prefix=$1 | |
| local key="${prefix%/}/.s3-upload-complete" | |
| local marker="$WORK/completion-marker.json" | |
| printf '{"source_id":%s,"source_sha256":%s,"kind":"tar"}\n' \ | |
| "$(json_string "$SOURCE_ID")" \ | |
| "$(json_string "$SOURCE_SHA256")" >"$marker" | |
| upload_blob_object \ | |
| "$key" "$marker" "__completion_marker__" \ | |
| "$(wc -c <"$marker" | tr -d '[:space:]')" \ | |
| "$SOURCE_MTIME" | |
| } | |
| tar_marker_matches() { | |
| local prefix=$1 | |
| local key="${prefix%/}/.s3-upload-complete" | |
| head_matches_member "$key" "__completion_marker__" \ | |
| "$(printf '{"source_id":%s,"source_sha256":%s,"kind":"tar"}\n' \ | |
| "$(json_string "$SOURCE_ID")" \ | |
| "$(json_string "$SOURCE_SHA256")" | wc -c | tr -d '[:space:]')" | |
| } | |
| tar_member_mode() { | |
| : "${IDU_DB:?}" | |
| : "${IDU_WORK:?}" | |
| : "${IDU_BUCKET:?}" | |
| : "${IDU_PREFIX:?}" | |
| : "${IDU_SOURCE_ID:?}" | |
| DB=$IDU_DB | |
| WORK=$IDU_WORK | |
| BUCKET=$IDU_BUCKET | |
| SOURCE_ID=$IDU_SOURCE_ID | |
| AWS_ENDPOINT_URL=${IDU_AWS_ENDPOINT_URL:-} | |
| FAIL_AFTER_PARTS=${IDU_FAIL_AFTER_PARTS:-} | |
| PROGRAM=${IDU_PROGRAM:-$PROGRAM} | |
| local raw=${TAR_FILENAME:?GNU tar did not provide TAR_FILENAME} | |
| local size=${TAR_SIZE:?GNU tar did not provide TAR_SIZE} | |
| local mtime=${TAR_MTIME:-0} | |
| local member key | |
| member=$(safe_member_name "$raw") || | |
| die "unsafe tar member name rejected: $raw" | |
| key=$(normalize_key "${IDU_PREFIX%/}/$member") | |
| upload_stream_object "$key" "$member" "$size" "$mtime" | |
| exit 0 | |
| } | |
| if [[ ${1:-} == --tar-member ]]; then | |
| shift | |
| tar_member_mode | |
| fi | |
| main() { | |
| require_command aws | |
| require_command sqlite3 | |
| require_command perl | |
| require_command tar | |
| require_command dd | |
| [[ $# -ge 2 && $# -le 3 ]] || | |
| die "usage: $PROGRAM SOURCE BUCKET [DESTINATION]" | |
| local self | |
| self=$(cd "$(dirname "$0")" && pwd -P)/$(basename "$0") | |
| SOURCE=$1 | |
| BUCKET=$2 | |
| DESTINATION=${3:-${SOURCE##*/}} | |
| [[ -f $SOURCE ]] || die "source is not a regular file: $SOURCE" | |
| SOURCE=$(cd "$(dirname "$SOURCE")" && pwd -P)/$(basename "$SOURCE") | |
| DESTINATION=$(normalize_key "$DESTINATION") | |
| [[ -n $DESTINATION ]] || die "destination key may not be empty" | |
| SOURCE_SIZE=$(wc -c <"$SOURCE" | tr -d '[:space:]') | |
| SOURCE_MTIME=$(perl -e 'print((stat($ARGV[0]))[9], "\n")' "$SOURCE") | |
| SOURCE_SHA256=$(sha256_file "$SOURCE") | |
| SOURCE_ID=$(printf '%s\0%s\0%s' \ | |
| "$SOURCE_SIZE" "$SOURCE_MTIME" "$SOURCE_SHA256" | sha256_text) | |
| if tar -tf "$SOURCE" >/dev/null 2>&1; then | |
| SOURCE_KIND=tar | |
| else | |
| SOURCE_KIND=blob | |
| fi | |
| local state_root=${STATE_DIR:-"$(dirname "$SOURCE")/.s3-upload-state"} | |
| mkdir -p "$state_root" | |
| local state_name | |
| state_name=$(printf '%s\0%s\0%s' "$SOURCE_ID" "$BUCKET" "$DESTINATION" | | |
| sha256_text) | |
| DB="$state_root/$state_name.sqlite" | |
| WORK="$state_root/$state_name.work" | |
| mkdir -p "$WORK" | |
| init_database | |
| write_source_record | |
| if [[ $SOURCE_KIND == blob ]]; then | |
| upload_blob_object \ | |
| "$DESTINATION" "$SOURCE" "__blob__" \ | |
| "$SOURCE_SIZE" "$SOURCE_MTIME" | |
| else | |
| local dest_dir dest_name prefix | |
| dest_dir=${DESTINATION%/*} | |
| dest_name=${DESTINATION##*/} | |
| if [[ $dest_dir == "$DESTINATION" ]]; then | |
| dest_dir= | |
| fi | |
| prefix=$(archive_base_name "$dest_name") | |
| [[ -n $dest_dir ]] && prefix="$dest_dir/$prefix" | |
| prefix=$(normalize_key "$prefix") | |
| if tar_marker_matches "$prefix"; then | |
| log "tar upload already complete: s3://$BUCKET/$prefix/" | |
| else | |
| export IDU_DB=$DB | |
| export IDU_WORK=$WORK | |
| export IDU_BUCKET=$BUCKET | |
| export IDU_PREFIX=$prefix | |
| export IDU_SOURCE_ID=$SOURCE_ID | |
| export IDU_AWS_ENDPOINT_URL=${AWS_ENDPOINT_URL:-} | |
| export IDU_FAIL_AFTER_PARTS=${FAIL_AFTER_PARTS:-} | |
| export IDU_PROGRAM=$PROGRAM | |
| # GNU tar invokes this program once per regular member and provides | |
| # that member's bytes on stdin plus TAR_FILENAME/TAR_SIZE/TAR_MTIME. | |
| tar -xf "$SOURCE" \ | |
| --to-command="$(printf '%q ' "$self")--tar-member" | |
| upload_completion_marker "$prefix" | |
| tar_marker_matches "$prefix" || | |
| die "tar completion marker failed verification" | |
| log "tar upload complete: s3://$BUCKET/$prefix/" | |
| fi | |
| fi | |
| if [[ ${KEEP_STATE_ON_SUCCESS:-0} != 1 ]]; then | |
| rm -rf "$WORK" | |
| rm -f "$DB" "$DB-wal" "$DB-shm" | |
| rmdir "$state_root" 2>/dev/null || true | |
| fi | |
| log "all work complete" | |
| } | |
| main "$@" |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment