Skip to content

Instantly share code, notes, and snippets.

@s1037989
Last active July 13, 2026 01:20
Show Gist options
  • Select an option

  • Save s1037989/c54562692d82b03a2219b4c581410a90 to your computer and use it in GitHub Desktop.

Select an option

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.
#!/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";
}
#!/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 "$@"
#!/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;
}
#!/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