Replacing NFS With S3 Without Changing Your Code: Automating Amazon S3 Files
You can replace NFS with Amazon S3 Files, a managed NFS v4.2 file system backed by an S3 bucket, and mount it with the mount.s3files helper. It works for workloads that only use the mount, but direct S3 API writes, sync lag and small file overhead need testing first.
For most of S3’s life, the honest answer to “can I just mount a bucket and point my application at it?” was some version of “sort of, if you are willing to live with FUSE and its compromises.” Tools like s3fs and Mountpoint for Amazon S3 gave you a path that looked like a directory, but the moment your application did something ordinary for a file system, such as rewriting part of a file, renaming it or taking a lock, you discovered that you were talking to an object store wearing a costume. That changed in April 2026 when AWS made Amazon S3 Files generally available, and it is worth looking at carefully because it is the first option where swapping an NFS mount for S3 can genuinely be a configuration change rather than a refactor.
This post walks through what S3 Files actually is, where the “no code changes” promise holds and where it quietly does not, and then gives you a set of scripts for the whole path: finding every file mount you already have, provisioning the bucket, role and file system, preparing each client, checking the existing data for features S3 Files cannot represent, migrating it, cutting over the mount, rolling back if you need to, and running smoke tests afterwards. The scripts are a starting point that you should read and run in a non production account first, not something to point at a production estate on day one. I have exercised their logic against simulated environments and reviewed them against the AWS documentation, but every estate has its own surprises.
1. What S3 Files Actually Is
S3 Files is a managed file system whose authoritative store is an S3 bucket, or a prefix within one. It is built on EFS infrastructure, clients mount it over NFS 4.1 or 4.2 using the mount.s3files helper that ships in amazon-efs-utils 3.0.0 and later, and according to the AWS mounting documentation it always uses TLS in transit and IAM authentication, neither of which can be turned off. The file system keeps a high performance storage layer for your active working set and synchronises in both directions with the bucket, so anything written through the mount eventually becomes an ordinary object version, and anything written to the bucket through the S3 API eventually shows up in the mount (AWS: S3 Files overview).
The details of that synchronisation explain most of the caveats later in this post, so it is worth being precise about them. When you first list a directory or open a file in it, S3 Files imports the metadata for every file in that directory, along with the data for files smaller than an import threshold that defaults to 128 KiB; larger files get metadata only, and their data is read from the bucket when you access it. Reads are served from the high performance layer when the data is there and the read is small, but the documentation is explicit that reads of 1 MiB or more are streamed directly from S3 even when the data is also on the file system, as are reads of data that has not been imported. Recently modified data that has not yet been exported is always served from the file system. Data that has not been read for a configurable period (30 days by default) and has already been synchronised is expired from the high performance layer, with the metadata kept so it can be fetched again on demand (AWS: synchronization).
In the other direction, S3 Files waits for 60 seconds of write inactivity on a file before exporting it to the bucket, and rapid successive writes are captured in a single PUT rather than one object version per write. That is an inactivity timer, not an upper bound: a file that is appended to every 30 seconds will not be exported until the writes stop for a full minute. Changes made directly to the bucket are picked up through S3 Event Notifications, which is why the service needs permission to manage EventBridge rules for your bucket. The pricing follows the same split: you pay a storage rate for the fraction of data resident on the high performance layer, access charges for reads and writes against that layer, and synchronisation charges for imports and exports, while reads streamed directly from S3 carry no file system data charge.
From the application’s point of view, none of this is visible. It sees a POSIX path, opens files, writes to them, renames them and closes them, and that is the whole point.
2. Where “No Code Changes” Holds, and Where It Does Not
The claim that you can replace an NFS mount without touching code is true for a large class of workloads, but I would treat it as a hypothesis to test rather than a guarantee, because the failure modes are subtle and tend to show up under load or months later rather than in a smoke test.
The good news first. AWS documents read after write consistency, advisory file locking and POSIX permissions for clients of the file system, and stores each file’s owner, group and mode as user defined object metadata on export. Classmethod’s hands on test after launch showed flock exclusive and non blocking locks behaving as expected between two processes on one instance, and observed modification times recorded in object metadata alongside the documented fields (Classmethod: S3 Files GA test). Renames, directory creation and in place writes all work as normal file operations, which is the real difference from FUSE based approaches.
Now the caveats, which are the things the preflight and validation scripts later in this post are designed to surface:
- The S3 API ignores file system locks. In the same Classmethod test, an object overwritten through the S3 API while an NFS client held an exclusive lock on it simply replaced the content, and the lock holder saw the new data. NFS locks coordinate NFS clients and nothing else, so if part of your migration plan is to move some consumers onto the SDK while others keep using the mount against the same keys, you need application level coordination, and that is a code change.
- Conflicts resolve in favour of the bucket. If a file is changed through the mount and the corresponding object also changes before the local change is exported, S3 Files moves the local version into a
.s3files-lost+found-<file-system-id>directory at the file system root and imports the bucket’s version. Files in that directory are not exported to S3, stay there until you delete them, and count towards your file system storage costs. Nothing is lost, but nothing is merged either, and somebody has to look. - Export timing depends on write patterns. Because export waits for 60 seconds without writes, the delay before a downstream consumer can see a file in the bucket is at least a minute and can be much longer for files written continuously. Classmethod measured a median of roughly 63 to 66 seconds from a single write through the mount to the object being visible with
head-object, and roughly 30 seconds in the other direction, on one instance in one region; your numbers will depend on your write patterns and on event notification delivery. If a consumer reads from S3 the moment a producer closes a file on the mount, it will see stale or missing data. - Some file system features are simply not supported. The AWS limitations page lists hard links, NFSv4 ACLs, custom user extended attributes, block and character devices, setuid on directories, mandatory locking, pNFS, delegations, Kerberos security and the
nconnectmount option as unsupported, along with limits of 255 bytes per path component and 1,024 bytes for the full object key, which includes your prefix. Objects already in Glacier storage classes cannot be read through the mount until they are restored (AWS: limits and unsupported features). If your application depends on any of these, the migration is not a configuration change, and copying with flags that try to preserve them will either fail or silently lose them. - Small file workloads pay a tax. Classmethod compared S3 Files with EFS in Elastic throughput mode on a single r7gd.medium instance, taking the median of three runs with simple
ddandcptests. Large sequential writes were within a few percent, a 1 GB read was about 9% slower, writing 1,000 files of 1 KB took about 24% longer (11.3 s against 9.0 s) and reading them back about 34% longer. Those are rough indicators from one small instance rather than a benchmark, but the direction is consistent with the architecture, and every exported file also becomes S3 request work. - It is NFS on Linux, in one VPC, against a bucket in the same Region. There is no SMB support, so Windows shares still belong on FSx for Windows File Server, a file system can only have mount targets in one VPC, and the bucket must be in the same Region as the file system. Clients in other VPCs or Regions can mount with the
mounttargetipoption and some extra configuration, according to the mounting documentation.
If your current NFS usage is “one application reads and writes files on a shared mount, and nothing else touches the underlying storage,” you are very likely in the safe zone. If your share is also being scraped by a batch job that you were planning to move onto the S3 API on day one, or it is full of hard links from a backup tool, be more careful.
3. Finding Every File Mount Before You Move Any of Them
Every storage migration I have seen go badly started with an inventory that was wrong. The NFS server everyone knew about was migrated cleanly, and then a month later somebody discovered that a reporting box had a hand typed mount in /etc/fstab pointing at the old server, or that a Lambda function was reading from an EFS access point nobody had written down, or that a Windows host had a persistent drive letter mapped to a share that was quietly being retired. So before provisioning anything, it is worth spending an hour finding out what you actually have, and that is what file_mount_recon.py is for.
The script looks at the problem from both ends. From the AWS APIs it builds the server side picture in every enabled region: EFS file systems with their mount targets, access points, backup and lifecycle policies and every CloudWatch metric they publish; S3 Files file systems and their prefixes; every FSx flavour, including ONTAP SVM endpoints and volumes, OpenZFS exports, Windows File Server aliases and Lustre mount names; Storage Gateway NFS and SMB shares; and the consumers you cannot log into, which are Lambda functions with file system configs, ECS task definitions with EFS volumes and EKS clusters running the EFS CSI driver. It also pulls every instance, network interface and security group so that it can check, rule by rule, whether a client can actually reach port 2049 for NFS, 445 for SMB or 988 for Lustre.
With --ssm, it then runs a read only collector on every SSM managed host, which is where most of the useful detail lives. I deliberately did not limit this to NFS, because the same exercise that finds your NFS mounts should also find the SMB shares, the Lustre clients and the FUSE mounts such as Mountpoint for S3, s3fs, goofys and rclone that people set up years ago and forgot about. On Linux it reads /proc/mounts, /etc/fstab (with passwords and credential paths redacted), autofs maps and systemd mount units, and it fully parses /proc/self/mountstats for every NFS mount, which gives you the negotiated options, the mount age, the byte counters, the RPC transport counters, and per operation counts with round trip time, retransmissions, timeouts and errors. It records which processes hold files open on each mount, which bucket each FUSE daemon is serving, the CIFS and Lustre client state, the TLS proxy state for EFS and S3 Files, and whether the host is itself quietly acting as a file server. On Windows it collects SMB mappings and live sessions with their dialect, encryption and signing, mapped and persistent per user drive letters from the registry, NFS client mounts, and any shares the host serves.
Every remote mount is also tested for responsiveness, without writing anything to it. On Linux the collector runs four probes, each under its own timeout and each timed: a statfs (what df does, and the one call that always goes to the server), a getattr on the mount root, a read of the first few directory entries, and a fresh TCP connection to the server’s port on the address the client resolves for it. The getattr and directory read can be answered from the client’s caches, so they tell you whether an application could use the mount right now rather than proving the server is alive, which is why the statfs and TCP probes are there as well. Put together, they let the report separate a server that is unreachable (probes time out and the port is closed) from one that is reachable but whose file system is not answering (the port accepts connections but calls time out) and from one that is merely slow, which are three very different conversations with three different teams. On Windows the collector does a timed listing of each mapped share and a TCP check to the server, with one caveat worth knowing: SSM runs as SYSTEM, so a share that only grants access to a particular user may show up as an access error rather than a timeout, and the report labels it accordingly.
Beyond whether a mount works at all, the script checks the things that tend to cause the next incident. For capacity, it looks at space and inode usage on every mount with a real size (EFS, S3 Files and S3 FUSE mounts report a virtual exabyte figure, so they are skipped), FSx storage utilisation from CloudWatch, Storage Gateway cache usage and the share of the cache not yet uploaded to S3, and EFS burst credits running down in bursting mode. For latency, it uses the kernel’s own RPC statistics rather than a synthetic benchmark: average round trip times for metadata operations and for reads and writes since the mount was made, and the gap between execute time and round trip time, which shows requests queueing on the client before they are even sent. It also flags FSx utilisation metrics that have reached 90%, where latency usually starts to climb, and mounts that cross Availability Zones, especially when a mount target exists in the client’s own zone. For cross region use, it flags clients and file servers in different regions, targets reached over inter region peering, and FUSE mounts of buckets in another region, all of which add latency and data transfer charges. For permissions, it reports share roots that are world writable without the sticky bit, root being denied access to a share root (normal with root squashing, but worth confirming the application user can get in), CIFS credentials files readable by anyone other than root, passwords written into /etc/fstab, EFS file systems without a file system policy (EFS’s default then lets any client that can reach a mount target mount, write and act as root) or without a TLS requirement, Storage Gateway shares open to any address or configured with NoSquash or guest access, and OpenZFS exports with no_root_squash. ONTAP export policies live inside ONTAP rather than the AWS API, so the report says so instead of guessing.
Each client mount is then resolved back to the resource behind it, whether that is EFS, S3 Files, an FSx file system, a Storage Gateway, an S3 bucket mounted through FUSE, a self managed NFS server running on EC2, or something the account cannot see at all. A mount through 127.0.0.1 is traced back to its real file system through the fstab entry or the TLS proxy state file, so EFS mounts using TLS do not show up as mysteries. The script then raises findings, which range from the operational (hung mounts, security groups blocking the port, soft NFS mounts, mounts that will not survive a reboot, file systems with no backups or no clients) to the ones that matter for this migration specifically. If mountstats shows a workload taking POSIX locks, renaming files heavily or writing in small chunks, or the mount uses nconnect, those map directly onto the caveats in section 2, and it is much better to see them flagged against a specific mount on a specific host now than to discover them after cutover.
Cross account mounts deserve their own mention, because they are the ones most likely to be missing from everybody’s diagram. A host in one account mounting an EFS file system that lives in another, over a VPC peering connection or a Transit Gateway, looks perfectly ordinary from inside the host, and the owning team often has no idea the dependency exists. The script looks for these from both directions. From the client side, it takes the address each mount actually connects to and walks the client subnet’s route table to see how traffic gets there: a local route means the same VPC (and, in a shared VPC, the subnet owner tells you whose network it is), a peering route names the peer account and VPC directly, and a Transit Gateway, Cloud WAN, VPN or Direct Connect route tells you the target is outside the VPC even when it cannot tell you exactly whose it is. EFS and S3 Files IDs, S3 buckets and Lambda access points that belong to no scanned account are flagged as well. From the server side, it reports file system policies that grant other accounts, security group rules that admit another account’s security group or address ranges outside the file system’s own VPC, and, with --flow-logs, the source addresses VPC Flow Logs have actually seen connecting to each mount target, attributed to an account where the network data allows. A single account run can only say “possibly another account” for anything behind a Transit Gateway; running with --org (or a list of --accounts) and a read only role in each member account resolves both ends exactly, so the report can say “this mount on this host in account A is that file system in account B”.
# API only, every enabled region
python3 file_mount_recon.py
# Include client side collection on SSM managed hosts, 30 days of metrics
python3 file_mount_recon.py --ssm --days 30
# Busy hosts with many mounts can exceed the 24,000 character inline SSM output limit
python3 file_mount_recon.py --ssm --ssm-bucket my-ssm-output-bucket --regions af-south-1
# Whole organisation, with flow log evidence of who connects to each file server
python3 file_mount_recon.py --ssm --flow-logs --org --role-name FileMountReconReadOnlyThe output is a directory containing report.md for humans, mounts.csv with roughly 170 columns per client mount (including the probe results, capacity, latency, permissions and the network path and account attribution), servers.csv, findings.csv, flow_sources.csv when flow logs are queried, a fully correlated inventory.json, and the raw API responses and host bundles under raw/ in case you need to check how a conclusion was reached. The report also lists its own blind spots, such as running instances that are not managed by SSM, accounts it could not assume a role into, Transit Gateway paths it could not attribute, and Kubernetes persistent volumes, because an inventory that does not tell you what it could not see is the kind that causes the surprises described above.
The script is strictly read only: the probes only read, and the flow log option runs CloudWatch Logs Insights queries, which are billed by the volume of logs scanned, so it is off by default. I chose not to have it copy any data, even though that would have been convenient, because something that runs across every host in an account should be reviewable as an audit rather than a change; copying belongs in the migration script further down, run deliberately one share at a time. The S3 Files calls follow the published CLI reference, which documents status values in both upper and lower case, so the script normalises them; even so, run it against a single region first and check that the S3 Files section of the output looks sensible before you rely on it across the estate. The full script is in the appendix at the end of this post.
4. The Automation Plan
Once you know what you have, the scripts below break the migration of one share into five pieces, all driven from a single configuration file:
s3files-provision.shcreates or reconciles the bucket settings, the IAM role that S3 Files assumes, the file system, the mount target security group and one mount target per Availability Zone.s3files-client-prep.shruns on each client host, installs a recent enoughamazon-efs-utils, and checks that TCP 2049 to the mount target is actually reachable before anything tries to mount.s3files-migrate.shhas four modes:preflightscans the existing share for features S3 Files cannot represent,bulkcopies the data while the application is still running,finaldoes the last pass through the mount during the cutover window and verifies content and metadata, andverify-syncwaits until every file is visible in the bucket.s3files-cutover.shrepoints the application’s mount path at S3 Files, and rolls it back safely, which in this context means refusing to roll back over data the old share does not have.s3files-validate.shruns smoke tests for the behaviours applications usually depend on, including two host lock contention and rename visibility.
The configuration file exists to fix a subtle problem: if provisioning, migration and cutover each have their own idea of where data lives, a file that was /data/report.csv on the old share can quietly end up at /data/data/report.csv on the new one, and the “no code changes” promise is broken by a path rather than by a semantic. Here the file system is created with --prefix, so its root is that prefix in the bucket; the bulk copy writes to s3://BUCKET/PREFIX, the final pass and the cutover both use the file system root, and so /data/report.csv on the mount is always s3://BUCKET/PREFIX/report.csv in the bucket. Every script reads the same s3files.env, and provisioning writes the file system ID back into it.
cat > s3files.env << 'EOF'
# s3files.env: one config file shared by every script, so the bucket prefix,
# file system and mount path cannot drift between provisioning, migration,
# cutover and validation.
BUCKET=my-company-shared-data
REGION=af-south-1
# S3 key prefix the file system is scoped to. The file system root maps to this
# prefix, so /data/report.csv on the mount becomes s3://BUCKET/apps/billing/report.csv
PREFIX=apps/billing/
# Path the application uses today. The old NFS share is mounted here until
# cutover, and S3 Files is mounted here afterwards.
MOUNT_PATH=/data
VPC_SUBNETS="subnet-aaaa1111 subnet-bbbb2222"
CLIENT_SG=sg-0123456789abcdef0
# Written by s3files-provision.sh
FS_ID=
EOF
chmod 644 s3files.envI have tried to make the scripts safe to rerun after a partial failure rather than claim they are idempotent in the strict sense. Provisioning reapplies the role’s trust and permission policies and the security group rule on every run, reuses a file system it finds for the same bucket and prefix, and skips Availability Zones that already have a mount target; cutover recognises a completed cutover and does nothing; the migration modes can be repeated. Rerun behaviour is still worth testing in a non production account before you rely on it. Each script below is wrapped in a cat > ... << 'EOF' block followed by chmod +x, so you can paste the whole block into a shell and end up with an executable file, and any heredocs inside the scripts use their own delimiters so they do not end the outer block early. All of them need AWS CLI 2.34 or later, since older versions do not have the aws s3files commands at all.
5. Provisioning the Bucket, Role and File System
S3 Files requires versioning on the bucket, and its default encryption must be SSE-S3 or SSE-KMS (AWS: prerequisites). The script enables versioning only if it is not already on, and it is worth knowing that versioning, once enabled, can be suspended but never removed. It never replaces an existing default encryption setting: a bucket already using SSE-KMS keeps its key, and the role gets the KMS permissions from the AWS prerequisites policy instead, while a bucket with no default encryption gets SSE-S3. If the bucket uses a customer managed key whose key policy does not delegate to IAM, you also need to add the role to the key policy, which the script reminds you about but cannot do for you.
The IAM role is assumed by the elasticfilesystem.amazonaws.com service principal, and the permission policy follows the scoped version in the AWS prerequisites documentation rather than the events:* shortcut you will find in a few tutorials, because in a regulated environment you do not want a storage service role that can rewrite every EventBridge rule in the account. The polling loops have deadlines and treat error, deleting and deleted as terminal failures rather than waiting forever, and the mount target loop counts targets that are not yet available across the whole list, so one target being available can never end the wait while another is still being created.
cat > s3files-provision.sh << 'EOF'
#!/usr/bin/env bash
# s3files-provision.sh
# Creates or reconciles the bucket settings, the S3 Files service role, the file
# system (scoped to PREFIX), the mount target security group and one mount target
# per subnet. Reads and updates s3files.env (override with S3FILES_ENV).
set -euo pipefail
ENV_FILE="${S3FILES_ENV:-./s3files.env}"
[[ -f "$ENV_FILE" ]] && source "$ENV_FILE"
: "${BUCKET:?set BUCKET}"
: "${REGION:?set REGION}"
: "${VPC_SUBNETS:?set VPC_SUBNETS (space separated, one subnet per AZ)}"
: "${CLIENT_SG:?set CLIENT_SG (security group of the NFS clients)}"
norm_prefix() { local p="${1#/}"; [[ -n "$p" && "$p" != */ ]] && p="$p/"; printf '%s' "$p"; }
PREFIX=$(norm_prefix "${PREFIX:-}")
ROLE_NAME="${ROLE_NAME:-S3FilesAccessRole-${BUCKET}}"
MT_SG_NAME="${MT_SG_NAME:-s3files-mt-${BUCKET}}"
WAIT_SECS="${WAIT_SECS:-1200}"
ACCOUNT_ID=$(aws sts get-caller-identity --query Account --output text)
BUCKET_ARN="arn:aws:s3:::${BUCKET}"
log() { echo "[$(date +%H:%M:%S)] $*"; }
die() { echo "ERROR: $*" >&2; exit 1; }
lc() { tr '[:upper:]' '[:lower:]'; }
TMP=$(mktemp -d); trap 'rm -rf "$TMP"' EXIT
# 1. Bucket: create if missing; enable versioning only if needed; never replace
# an existing default encryption setting.
if ! aws s3api head-bucket --bucket "$BUCKET" 2>/dev/null; then
log "Creating bucket $BUCKET"
if [[ "$REGION" == "us-east-1" ]]; then
aws s3api create-bucket --bucket "$BUCKET" --region "$REGION" >/dev/null
else
aws s3api create-bucket --bucket "$BUCKET" --region "$REGION" \
--create-bucket-configuration LocationConstraint="$REGION" >/dev/null
fi
fi
VER=$(aws s3api get-bucket-versioning --bucket "$BUCKET" --query Status --output text)
if [[ "$VER" != "Enabled" ]]; then
log "Enabling versioning (required by S3 Files; once enabled it can be suspended but not removed)"
aws s3api put-bucket-versioning --bucket "$BUCKET" --versioning-configuration Status=Enabled
fi
SSE_ALG=$(aws s3api get-bucket-encryption --bucket "$BUCKET" \
--query 'ServerSideEncryptionConfiguration.Rules[0].ApplyServerSideEncryptionByDefault.SSEAlgorithm' \
--output text 2>/dev/null || echo NONE)
KMS_KEY=$(aws s3api get-bucket-encryption --bucket "$BUCKET" \
--query 'ServerSideEncryptionConfiguration.Rules[0].ApplyServerSideEncryptionByDefault.KMSMasterKeyID' \
--output text 2>/dev/null || echo None)
case "$SSE_ALG" in
AES256) log "Bucket default encryption is SSE-S3; leaving it unchanged" ;;
aws:kms) log "Bucket default encryption is SSE-KMS (${KMS_KEY}); leaving it unchanged and granting the role KMS use" ;;
NONE|None|"")
log "No default encryption configured; setting SSE-S3"
aws s3api put-bucket-encryption --bucket "$BUCKET" --server-side-encryption-configuration \
'{"Rules":[{"ApplyServerSideEncryptionByDefault":{"SSEAlgorithm":"AES256"},"BucketKeyEnabled":true}]}'
SSE_ALG=AES256 ;;
*) die "bucket default encryption ${SSE_ALG} is not supported by S3 Files (SSE-S3 or SSE-KMS only)" ;;
esac
# 2. Service role assumed by S3 Files. Trust and permission policies are applied on
# every run, so an existing role is reconciled rather than trusted as is.
cat > "$TMP/trust.json" <<JSON
{
"Version": "2012-10-17",
"Statement": [{
"Sid": "AllowS3FilesAssumeRole",
"Effect": "Allow",
"Principal": { "Service": "elasticfilesystem.amazonaws.com" },
"Action": "sts:AssumeRole",
"Condition": {
"StringEquals": { "aws:SourceAccount": "${ACCOUNT_ID}" },
"ArnLike": { "aws:SourceArn": "arn:aws:s3files:${REGION}:${ACCOUNT_ID}:file-system/*" }
}
}]
}
JSON
KMS_STMT=""
if [[ "$SSE_ALG" == "aws:kms" ]]; then
KMS_RESOURCE="arn:aws:kms:${REGION}:${ACCOUNT_ID}:*"
[[ "$KMS_KEY" == arn:aws*:kms:*:key/* ]] && KMS_RESOURCE="$KMS_KEY"
KMS_STMT=$(cat <<JSON
,{ "Sid": "UseKmsKeyWithS3Files", "Effect": "Allow",
"Action": ["kms:GenerateDataKey","kms:Encrypt","kms:Decrypt","kms:ReEncryptFrom","kms:ReEncryptTo"],
"Condition": { "StringLike": {
"kms:ViaService": "s3.${REGION}.amazonaws.com",
"kms:EncryptionContext:aws:s3:arn": ["${BUCKET_ARN}", "${BUCKET_ARN}/*"] } },
"Resource": "${KMS_RESOURCE}" }
JSON
)
fi
cat > "$TMP/policy.json" <<JSON
{
"Version": "2012-10-17",
"Statement": [
{ "Sid": "S3BucketPermissions", "Effect": "Allow",
"Action": ["s3:ListBucket","s3:ListBucketVersions"],
"Resource": "${BUCKET_ARN}",
"Condition": { "StringEquals": { "aws:ResourceAccount": "${ACCOUNT_ID}" } } },
{ "Sid": "S3ObjectPermissions", "Effect": "Allow",
"Action": ["s3:AbortMultipartUpload","s3:DeleteObject*","s3:GetObject*","s3:List*","s3:PutObject*"],
"Resource": "${BUCKET_ARN}/*",
"Condition": { "StringEquals": { "aws:ResourceAccount": "${ACCOUNT_ID}" } } }
${KMS_STMT}
,{ "Sid": "EventBridgeManage", "Effect": "Allow",
"Action": ["events:DeleteRule","events:DisableRule","events:EnableRule",
"events:PutRule","events:PutTargets","events:RemoveTargets"],
"Resource": ["arn:aws:events:*:*:rule/DO-NOT-DELETE-S3-Files*"],
"Condition": { "StringEquals": { "events:ManagedBy": "elasticfilesystem.amazonaws.com" } } },
{ "Sid": "EventBridgeRead", "Effect": "Allow",
"Action": ["events:DescribeRule","events:ListRuleNamesByTarget","events:ListRules","events:ListTargetsByRule"],
"Resource": ["arn:aws:events:*:*:rule/*"] }
]
}
JSON
NEW_ROLE=0
if aws iam get-role --role-name "$ROLE_NAME" >/dev/null 2>&1; then
log "Reconciling trust policy on existing role $ROLE_NAME"
aws iam update-assume-role-policy --role-name "$ROLE_NAME" --policy-document "file://$TMP/trust.json"
else
log "Creating role $ROLE_NAME"
aws iam create-role --role-name "$ROLE_NAME" --assume-role-policy-document "file://$TMP/trust.json" >/dev/null
NEW_ROLE=1
fi
aws iam put-role-policy --role-name "$ROLE_NAME" --policy-name S3FilesBucketAccess \
--policy-document "file://$TMP/policy.json"
ROLE_ARN=$(aws iam get-role --role-name "$ROLE_NAME" --query Role.Arn --output text)
if [[ "$NEW_ROLE" == 1 ]]; then log "Waiting for the new role to propagate"; sleep 15; fi
[[ "$SSE_ALG" == "aws:kms" ]] && log "NOTE: the KMS key policy must also allow ${ROLE_ARN} unless it delegates to IAM"
# 3. File system for exactly this bucket and prefix: reuse if present, else create.
FS_ID="${FS_ID:-}"
if [[ -z "$FS_ID" ]]; then
for id in $(aws s3files list-file-systems --region "$REGION" --bucket "$BUCKET_ARN" \
--query 'fileSystems[].fileSystemId' --output text); do
[[ "$id" == "None" ]] && continue
read -r st p < <(aws s3files get-file-system --region "$REGION" --file-system-id "$id" \
--query '[status, prefix]' --output text)
[[ "$p" == "None" ]] && p=""
st=$(lc <<<"$st")
if [[ "$(norm_prefix "$p")" == "$PREFIX" && "$st" != "deleting" && "$st" != "deleted" ]]; then
FS_ID="$id"; break
fi
done
fi
if [[ -z "$FS_ID" ]]; then
TOKEN="prov-$(printf '%s|%s' "$BUCKET" "$PREFIX" | sha1sum | cut -c1-40)"
log "Creating S3 Files file system for s3://${BUCKET}/${PREFIX}"
ARGS=(--region "$REGION" --bucket "$BUCKET_ARN" --role-arn "$ROLE_ARN" --client-token "$TOKEN"
--tags "key=Name,value=s3files-${BUCKET}")
[[ -n "$PREFIX" ]] && ARGS+=(--prefix "$PREFIX")
FS_ID=$(aws s3files create-file-system "${ARGS[@]}" --query fileSystemId --output text)
fi
log "File system: $FS_ID"
deadline=$((SECONDS + WAIT_SECS))
while :; do
read -r st msg < <(aws s3files get-file-system --region "$REGION" --file-system-id "$FS_ID" \
--query '[status, statusMessage]' --output text)
st=$(lc <<<"$st")
case "$st" in
available) break ;;
error|deleting|deleted) die "file system $FS_ID is ${st}: ${msg}" ;;
esac
(( SECONDS > deadline )) && die "timed out waiting for $FS_ID (last status: $st)"
log "File system status: $st"; sleep 15
done
# 4. Subnets must share one VPC and sit in distinct AZs (one mount target per AZ).
aws ec2 describe-subnets --region "$REGION" --subnet-ids $VPC_SUBNETS \
--query 'Subnets[].[SubnetId,VpcId,AvailabilityZoneId]' --output text > "$TMP/subnets"
[[ $(cut -f2 "$TMP/subnets" | sort -u | wc -l) -eq 1 ]] || die "subnets span more than one VPC"
[[ -z $(cut -f3 "$TMP/subnets" | sort | uniq -d) ]] || die "two subnets share an AZ; give one subnet per AZ"
VPC_ID=$(head -1 "$TMP/subnets" | cut -f2)
# 5. Mount target security group, with the NFS rule ensured on every run.
MT_SG=$(aws ec2 describe-security-groups --region "$REGION" \
--filters "Name=group-name,Values=${MT_SG_NAME}" "Name=vpc-id,Values=${VPC_ID}" \
--query 'SecurityGroups[0].GroupId' --output text)
if [[ "$MT_SG" == "None" ]]; then
MT_SG=$(aws ec2 create-security-group --region "$REGION" --vpc-id "$VPC_ID" \
--group-name "$MT_SG_NAME" --description "S3 Files mount targets for ${BUCKET}" \
--query GroupId --output text)
fi
if ! out=$(aws ec2 authorize-security-group-ingress --region "$REGION" --group-id "$MT_SG" \
--ip-permissions "[{\"IpProtocol\":\"tcp\",\"FromPort\":2049,\"ToPort\":2049,\"UserIdGroupPairs\":[{\"GroupId\":\"${CLIENT_SG}\",\"Description\":\"NFS from clients\"}]}]" 2>&1); then
grep -q 'InvalidPermission.Duplicate' <<<"$out" || die "$out"
fi
log "Mount target security group: $MT_SG (TCP 2049 from $CLIENT_SG ensured)"
# 6. One mount target per subnet; skip AZs that already have one.
COVERED=$(aws s3files list-mount-targets --region "$REGION" --file-system-id "$FS_ID" --no-paginate \
--query 'mountTargets[].availabilityZoneId' --output text)
while read -r SUBNET _ AZ_ID; do
if grep -qw -- "$AZ_ID" <<<"$COVERED"; then
log "Mount target already exists in $AZ_ID"; continue
fi
log "Creating mount target in $SUBNET ($AZ_ID)"
aws s3files create-mount-target --region "$REGION" --file-system-id "$FS_ID" \
--subnet-id "$SUBNET" --security-groups "$MT_SG" >/dev/null
done < "$TMP/subnets"
# 7. Wait until every mount target is available. Counts come from JMESPath over the
# full list, so one available target cannot mask another that is still creating.
EXPECTED=$(wc -l < "$TMP/subnets")
deadline=$((SECONDS + WAIT_SECS))
while :; do
read -r total notready failed < <(aws s3files list-mount-targets --region "$REGION" \
--file-system-id "$FS_ID" --no-paginate --output text --query \
"[length(mountTargets), length(mountTargets[?status!='available' && status!='AVAILABLE']), length(mountTargets[?contains(['error','ERROR','deleting','DELETING','deleted','DELETED'], status)])]")
if (( failed > 0 )); then
aws s3files list-mount-targets --region "$REGION" --file-system-id "$FS_ID" --no-paginate \
--query 'mountTargets[].[mountTargetId,subnetId,status,statusMessage]' --output table >&2
die "one or more mount targets failed"
fi
(( total >= EXPECTED && notready == 0 )) && break
(( SECONDS > deadline )) && die "timed out waiting for mount targets ($notready of $total not available)"
log "Mount targets: $((total - notready))/$EXPECTED available"; sleep 20
done
# 8. Record the file system ID for the other scripts.
touch "$ENV_FILE"
if grep -q '^FS_ID=' "$ENV_FILE"; then
sed -i "s/^FS_ID=.*/FS_ID=${FS_ID}/" "$ENV_FILE"
else
echo "FS_ID=${FS_ID}" >> "$ENV_FILE"
fi
log "Done. FS_ID=${FS_ID} written to ${ENV_FILE}"
EOF
chmod +x s3files-provision.shThe file system typically takes a few minutes to become available, and each mount target takes several minutes more because it is creating an elastic network interface in your subnet, so the deadlines are generous. Separately, the client hosts need their own permissions: attach the managed AmazonS3FilesClientFullAccess policy (or the read only variant for consumers that never write) to the instance role, and give the instance role direct read access to the bucket (s3:GetObject, s3:GetObjectVersion and s3:ListBucket), because large reads are streamed straight from S3 using the client’s own credentials and are governed by IAM and bucket policies rather than POSIX permissions.
6. Preparing the Client Hosts
The most common way this goes wrong is not IAM, it is networking. If TCP 2049 is blocked between the client and the mount target, the mount does not fail cleanly; it retries and eventually times out, and you spend twenty minutes wondering whether the problem is your role. The client preparation script therefore finds the mount target in the host’s own Availability Zone and checks reachability explicitly before it attempts a trial mount.
On Amazon Linux 2023 the default repositories carried only the 2.x series of amazon-efs-utils for a while after launch, and that series does not include mount.s3files, so the script adds the official repository when the installed version is too old.
cat > s3files-client-prep.sh << 'EOF'
#!/usr/bin/env bash
# s3files-client-prep.sh
# Run on each client host: installs amazon-efs-utils 3.0.0+ (which provides
# mount.s3files), checks TCP 2049 to the mount target in this host's AZ, and does
# a trial mount on a scratch path. Usage: sudo ./s3files-client-prep.sh
set -euo pipefail
ENV_FILE="${S3FILES_ENV:-./s3files.env}"
[[ -f "$ENV_FILE" ]] && source "$ENV_FILE"
: "${FS_ID:?set FS_ID (s3files-provision.sh writes it to s3files.env)}"
: "${REGION:?set REGION}"
log() { echo "[$(date +%H:%M:%S)] $*"; }
die() { echo "ERROR: $*" >&2; exit 1; }
ver_ge() { [[ "$(printf '%s\n%s\n' "$2" "$1" | sort -V | head -1)" == "$2" ]]; }
# 1. AWS CLI must know about s3files
CLI_VER=$(aws --version 2>&1 | sed -E 's#aws-cli/([0-9.]+).*#\1#')
ver_ge "$CLI_VER" "2.34.0" || die "AWS CLI $CLI_VER is too old, need 2.34+"
# 2. amazon-efs-utils 3.0.0 or later
current_efs_utils() {
rpm -q --qf '%{VERSION}' amazon-efs-utils 2>/dev/null \
|| dpkg-query -W -f='${Version}' amazon-efs-utils 2>/dev/null || echo 0
}
if ! ver_ge "$(current_efs_utils)" "3.0.0"; then
if command -v dnf >/dev/null; then
log "Adding the official efs-utils repository"
cat > /etc/yum.repos.d/efs-utils.repo <<'REPO'
[efs-utils]
name=efs-utils repository
baseurl=https://amazon-efs-utils.aws.com/repo/rpm/amazon/2023
priority=1
enabled=1
repo_gpgcheck=1
type=rpm
gpgcheck=1
gpgkey=file:///etc/pki/rpm-gpg/RPM-GPG-KEY-efs-utils.gpg
REPO
curl -fsSL https://amazon-efs-utils.aws.com/efs-utils-armored.gpg \
-o /etc/pki/rpm-gpg/RPM-GPG-KEY-efs-utils.gpg
rpm --import /etc/pki/rpm-gpg/RPM-GPG-KEY-efs-utils.gpg
dnf install -y amazon-efs-utils
else
log "Installing efs-utils via the AWS installer"
curl -fsSL https://amazon-efs-utils.aws.com/efs-utils-installer.sh | sh -s -- --install
fi
python3 -m pip install --quiet botocore || true
fi
command -v mount.s3files >/dev/null || die "mount.s3files not found after install"
log "mount.s3files $(mount.s3files --version 2>&1 | head -1)"
# 3. Mount target in this host's AZ, and TCP 2049 reachability
TOKEN=$(curl -sX PUT http://169.254.169.254/latest/api/token -H 'X-aws-ec2-metadata-token-ttl-seconds: 60')
AZ_ID=$(curl -s -H "X-aws-ec2-metadata-token: $TOKEN" \
http://169.254.169.254/latest/meta-data/placement/availability-zone-id)
MT_IP=$(aws s3files list-mount-targets --region "$REGION" --file-system-id "$FS_ID" --no-paginate \
--query "mountTargets[?availabilityZoneId=='${AZ_ID}'].ipv4Address | [0]" --output text)
[[ -n "$MT_IP" && "$MT_IP" != "None" ]] || die "no mount target in AZ $AZ_ID"
if timeout 5 bash -c "echo > /dev/tcp/${MT_IP}/2049" 2>/dev/null; then
log "TCP 2049 to $MT_IP is reachable"
else
die "TCP 2049 to $MT_IP is BLOCKED: check the mount target security group"
fi
# 4. Trial mount on a scratch path, then unmount
PROBE=/mnt/.s3files-probe
mkdir -p "$PROBE"
mount -t s3files "${FS_ID}:/" "$PROBE"
df -hT "$PROBE"
umount "$PROBE" && rmdir "$PROBE"
log "Client is ready"
EOF
chmod +x s3files-client-prep.shIf you see 127.0.0.1:/ as the source in df and an 8.0E size, that is expected. The NFS client connects through a local TLS proxy, and the size is a virtual figure reflecting the bucket rather than anything you are paying for.
7. Migrating the Existing Data
The migration script starts with a preflight scan, because the cheapest time to discover that your share is full of hard links is before you have copied any of it. The scan walks the existing share read only and treats every unsupported feature from section 2 as a compatibility gate: hard links, block and character devices, setuid directories, path components over 255 bytes, paths that would exceed 1,024 bytes once the prefix is added, symlinks with empty or overlong targets, user extended attributes, POSIX ACLs, a sample of NFSv4 ACLs when the share is mounted over NFSv4, and objects already sitting under the prefix in Glacier storage classes. FIFOs and sockets are reported too, since the copy does not transfer them. Any gate that finds something fails the run, and the only way past it is to fix the source or to list the gate in ACCEPT after you have confirmed the application does not depend on that feature. Simply dropping the flags that would have tried to preserve these things is not the same as knowing you do not need them.
The copy itself uses two paths. The bulk pass uses aws s3 sync straight into the bucket under the prefix, which is much faster for large volumes and can run as often as you like while the application is still live; objects written this way carry no POSIX owner or mode, and symlinks are skipped. The final pass runs in the cutover window with the application stopped: it mounts the file system on a staging path and runs rsync -a with deletion through the mount, which sets owner, group, mode and times as file system metadata, creates the symlinks, and, because both paths end in a slash, applies the old share’s root ownership and mode to the file system root. It deliberately does not use -H, -A or -X, since hard links, ACLs and extended attributes are exactly what the preflight gates are about.
Verification is not a file count. The final mode compares a metadata manifest of both trees (type, mode, owner, group, size, modification time and symlink target for every file; type, mode, owner and group for directories) and then runs a checksum comparison that reads every file on both sides, so a file with the right size and timestamp but the wrong content is still caught. Finally, because export waits for 60 seconds without writes rather than happening within a fixed time, verify-sync does not sleep for a minute and hope; it lists the bucket under the prefix until every regular file is present with the expected size, and checks that the conflict directory is empty, before it tells you to cut over.
cat > s3files-migrate.sh << 'EOF'
#!/usr/bin/env bash
# s3files-migrate.sh
# Copies an existing NFS share into S3 Files in four steps:
# preflight (app running) scan the share for features S3 Files cannot represent
# bulk (app running) aws s3 sync the share into s3://BUCKET/PREFIX
# final (app stopped) rsync through an S3 Files mount, then verify content and metadata
# verify-sync (app stopped) wait until every file is visible in S3 with the right size
# Usage: sudo ./s3files-migrate.sh <mode>
set -euo pipefail
ENV_FILE="${S3FILES_ENV:-./s3files.env}"
[[ -f "$ENV_FILE" ]] && source "$ENV_FILE"
MODE="${1:-${MODE:-}}"
: "${MODE:?usage: s3files-migrate.sh preflight|bulk|final|verify-sync}"
: "${MOUNT_PATH:?set MOUNT_PATH}"
: "${BUCKET:?set BUCKET}"
norm_prefix() { local p="${1#/}"; [[ -n "$p" && "$p" != */ ]] && p="$p/"; printf '%s' "$p"; }
PREFIX=$(norm_prefix "${PREFIX:-}")
SRC="${SRC:-$MOUNT_PATH}" # the old NFS share, still mounted at the application path
STAGE="${STAGE:-/mnt/.s3files-stage}"
REPORTS="${REPORTS:-./s3files-reports}"
EXCL='.s3files-lost+found-*'
RSYNC_FLAGS=(-a --numeric-ids --no-devices --no-specials --exclude="$EXCL")
log() { echo "[$(date +%H:%M:%S)] $*"; }
die() { echo "ERROR: $*" >&2; exit 1; }
mkdir -p "$REPORTS"
export LC_ALL=C
mount_stage() {
: "${FS_ID:?set FS_ID}"
mkdir -p "$STAGE"
mountpoint -q "$STAGE" || mount -t s3files "${FS_ID}:/" "$STAGE"
}
# Metadata manifest: type, mode, owner, group, size and mtime for files; type, mode,
# owner and group for directories (directory sizes and times are not comparable
# across file systems); symlinks are compared by target.
manifest() {
(cd "$1" && find . -path "./${EXCL}" -prune -o -printf '%y\t%m\t%U\t%G\t%s\t%Ts\t%l\t%p\n') \
| awk -F'\t' 'BEGIN { OFS = "\t" }
$1 == "d" { $5 = ""; $6 = "" }
$1 == "l" { $2 = ""; $5 = ""; $6 = "" }
$1 != "f" && $1 != "d" && $1 != "l" { next }
{ print }' | sort
}
gate() { # gate NAME COUNT DESCRIPTION
local name="$1" count="$2" desc="$3" status=PASS
if (( count > 0 )); then
if [[ ",${ACCEPT:-}," == *",${name},"* ]]; then status=ACCEPTED; else status=FAIL; FAILED=1; fi
fi
printf '%-14s %-9s %8s %s\n' "$name" "$status" "$count" "$desc" | tee -a "$REPORTS/preflight.txt"
}
case "$MODE" in
preflight)
[[ "$(findmnt -n -o FSTYPE --mountpoint "$SRC" 2>/dev/null)" == nfs* ]] \
|| log "WARNING: $SRC is not an NFS mount; scanning anyway"
: > "$REPORTS/preflight.txt"; FAILED=0
log "Scanning $SRC (read only); this walks the whole tree"
find "$SRC" -xdev -type f -links +1 -printf '%i\t%p\n' > "$REPORTS/hardlinks.txt"
find "$SRC" -xdev \( -type b -o -type c \) > "$REPORTS/devices.txt"
find "$SRC" -xdev \( -type p -o -type s \) > "$REPORTS/fifos_sockets.txt"
find "$SRC" -xdev -type d -perm -4000 > "$REPORTS/setuid_dirs.txt"
find "$SRC" -xdev -printf '%f\n' | awk 'length($0) > 255' > "$REPORTS/long_names.txt"
(cd "$SRC" && find . -xdev -printf '%P\n') | awk -v p="$PREFIX" 'length(p $0) > 1024' > "$REPORTS/long_keys.txt"
find "$SRC" -xdev -type l -printf '%l\t%p\n' | awk -F'\t' 'length($1) > 4080 || length($1) == 0' > "$REPORTS/bad_symlinks.txt"
if command -v getfattr >/dev/null; then
getfattr -R -d -m '^user\.' --absolute-names "$SRC" 2>/dev/null | grep '^# file:' > "$REPORTS/xattrs.txt" || true
else
log "getfattr not installed: user xattrs NOT checked (install attr)"; : > "$REPORTS/xattrs.txt"
fi
if command -v getfacl >/dev/null; then
getfacl -R -s -p "$SRC" 2>/dev/null | grep '^# file:' > "$REPORTS/posix_acls.txt" || true
else
log "getfacl not installed: POSIX ACLs NOT checked (install acl)"; : > "$REPORTS/posix_acls.txt"
fi
: > "$REPORTS/nfs4_acls.txt"
if [[ "$(findmnt -n -o FSTYPE --mountpoint "$SRC" 2>/dev/null)" == nfs4 ]]; then
if command -v nfs4_getfacl >/dev/null; then
find "$SRC" -xdev | head -n "${ACL_SAMPLE:-5000}" | while IFS= read -r f; do
nfs4_getfacl "$f" 2>/dev/null | grep -v '^#' | grep -Eqv ':(OWNER|GROUP|EVERYONE)@:|^$' \
&& echo "$f"
done > "$REPORTS/nfs4_acls.txt" || true
else
log "nfs4_getfacl not installed: NFSv4 ACLs NOT checked (install nfs4-acl-tools)"
fi
fi
ARCHIVED=0
if [[ -n "$(aws s3api list-objects-v2 --bucket "$BUCKET" --prefix "$PREFIX" --max-items 1 --query 'Contents[0].Key' --output text 2>/dev/null | grep -v None)" ]]; then
ARCHIVED=$(aws s3api list-objects-v2 --bucket "$BUCKET" --prefix "$PREFIX" --output text \
--query "Contents[?StorageClass=='GLACIER' || StorageClass=='DEEP_ARCHIVE'].[Key]" | grep -vc '^None$' || true)
fi
HL=$(cut -f1 "$REPORTS/hardlinks.txt" | sort -u | wc -l)
{
echo "Preflight for $SRC -> s3://${BUCKET}/${PREFIX} ($(date -u +%FT%TZ))"
printf '%-14s %-9s %8s %s\n' GATE STATUS COUNT DETAIL
} >> "$REPORTS/preflight.txt"
gate hardlinks "$HL" "inodes with more than one name (not supported; each name would become a separate copy)"
gate devices "$(wc -l < "$REPORTS/devices.txt")" "block or character devices (not supported)"
gate specials "$(wc -l < "$REPORTS/fifos_sockets.txt")" "FIFOs and sockets (not copied by this script)"
gate setuid_dirs "$(wc -l < "$REPORTS/setuid_dirs.txt")" "setuid directories (not supported)"
gate long_names "$(wc -l < "$REPORTS/long_names.txt")" "path components over 255 bytes"
gate long_keys "$(wc -l < "$REPORTS/long_keys.txt")" "prefix plus path over 1,024 bytes (cannot be exported to S3)"
gate symlinks "$(wc -l < "$REPORTS/bad_symlinks.txt")" "empty symlink targets or targets over 4,080 bytes"
gate xattrs "$(wc -l < "$REPORTS/xattrs.txt")" "files with user extended attributes (not supported)"
gate posix_acls "$(wc -l < "$REPORTS/posix_acls.txt")" "files with extended POSIX ACLs (NFS ACLs are not supported)"
gate nfs4_acls "$(wc -l < "$REPORTS/nfs4_acls.txt")" "files with named NFSv4 ACEs (sampled, first ${ACL_SAMPLE:-5000} entries)"
gate archived "$ARCHIVED" "objects already under the prefix in Glacier storage classes (unreadable through the mount)"
if (( FAILED )); then
echo "FAIL" > "$REPORTS/preflight.status"
die "compatibility gates failed; see $REPORTS. Fix the source, or set ACCEPT=gate1,gate2 once you have confirmed the application does not depend on that feature"
fi
echo "PASS" > "$REPORTS/preflight.status"
log "Preflight passed; details in $REPORTS"
;;
bulk)
[[ "$(cat "$REPORTS/preflight.status" 2>/dev/null)" == PASS ]] || die "run preflight first"
log "Bulk copy $SRC -> s3://${BUCKET}/${PREFIX}"
# Symlinks are skipped here and created by the final rsync pass. Objects written
# this way carry no POSIX owner or mode, which the final pass sets.
aws s3 sync "$SRC" "s3://${BUCKET}/${PREFIX}" --only-show-errors --exact-timestamps --no-follow-symlinks
log "Bulk pass complete; rerun as often as you like before the cutover window"
;;
final)
[[ "$(cat "$REPORTS/preflight.status" 2>/dev/null)" == PASS ]] || die "run preflight first"
if fuser -m "$SRC" >/dev/null 2>&1; then die "processes still use $SRC; stop the application first"; fi
mount_stage
log "Final pass with metadata: $SRC -> S3 Files root (s3://${BUCKET}/${PREFIX})"
# The trailing slash on both sides makes rsync apply the source root's owner,
# group and mode to the file system root as well as copying its contents.
rsync "${RSYNC_FLAGS[@]}" --delete --info=stats2 "${SRC%/}/" "$STAGE/"
log "Verifying metadata (type, mode, owner, group, size, mtime, symlink target)"
manifest "$SRC" > "$REPORTS/src.manifest"
manifest "$STAGE" > "$REPORTS/dst.manifest"
if ! diff -q "$REPORTS/src.manifest" "$REPORTS/dst.manifest" >/dev/null; then
diff "$REPORTS/src.manifest" "$REPORTS/dst.manifest" | head -20 >&2
die "metadata differs between source and S3 Files; full manifests in $REPORTS"
fi
if [[ "${VERIFY:-checksum}" == checksum ]]; then
log "Verifying content with checksums (reads every file on both sides)"
rsync "${RSYNC_FLAGS[@]}" --delete -n -i --checksum "${SRC%/}/" "$STAGE/" \
| grep -v '^\.d' > "$REPORTS/content_diff.txt" || true
[[ -s "$REPORTS/content_diff.txt" ]] && { head -20 "$REPORTS/content_diff.txt" >&2; die "content differs; see $REPORTS/content_diff.txt"; }
fi
umount "$STAGE" && rmdir "$STAGE"
log "Final pass verified. Next: verify-sync, then cutover"
;;
verify-sync)
# S3 Files exports a file after 60 seconds without writes to it, so there is no
# fixed wait that is always enough. Instead, wait until every regular file is
# present in the bucket with the expected size.
(cd "$SRC" && find . -type f -printf '%P\t%s\n') | awk -v p="$PREFIX" '{ print p $0 }' | sort > "$REPORTS/want.txt"
deadline=$((SECONDS + ${SYNC_WAIT:-3600}))
while :; do
aws s3api list-objects-v2 --bucket "$BUCKET" --prefix "$PREFIX" \
--query 'Contents[].[Key,Size]' --output text | grep -v '^None$' | sort > "$REPORTS/have.txt" || true
missing=$(comm -23 "$REPORTS/want.txt" "$REPORTS/have.txt" | wc -l)
(( missing == 0 )) && break
(( SECONDS > deadline )) && { comm -23 "$REPORTS/want.txt" "$REPORTS/have.txt" | head -20 >&2; die "$missing files not exported"; }
log "$missing of $(wc -l < "$REPORTS/want.txt") files not yet in S3 with the expected size; waiting"
sleep 30
done
mount_stage
LF=$(find "$STAGE" -maxdepth 1 -name "$EXCL" -type d | head -1)
if [[ -n "$LF" && -n "$(ls -A "$LF" 2>/dev/null)" ]]; then
die "conflict files present in $LF; resolve them before cutover"
fi
umount "$STAGE" && rmdir "$STAGE"
log "All $(wc -l < "$REPORTS/want.txt") files are in s3://${BUCKET}/${PREFIX} and no conflicts were recorded"
;;
*) die "unknown mode $MODE (preflight|bulk|final|verify-sync)" ;;
esac
EOF
chmod +x s3files-migrate.shRun the bulk pass before any client has the new file system mounted for production use, because once both the mount and the S3 API are being written you are in the territory where conflicts resolve in favour of the bucket. The checksum verification reads the whole share twice, so on a very large share you may want VERIFY=metadata for a first rehearsal, but I would not skip it in the real window.
8. Cutting Over the Mount Path
This is the step that delivers the “no code changes” outcome. Your application is configured to read and write /data (or whatever your path is), and that path is currently an NFS mount in /etc/fstab. If the new entry mounts the S3 Files root at exactly the same path, and the root holds exactly what the old share’s root held, the application has no way of knowing anything changed. The script refuses to start while anything still has files open on the mount, records the old fstab line and the old root’s owner and mode, swaps the entry, and if the S3 Files mount fails it restores the old entry and remounts the old share. After a successful mount it makes sure the root’s owner and mode match what the old share had.
Rollback deserves more care than a fstab backup, because the moment the application writes to S3 Files the old share becomes stale, and remounting it would quietly discard everything written since cutover. So the script records a cutover marker after its own changes and, on rollback, looks for any file or directory modified or changed after that marker. If it finds none, rollback is a simple swap back. If it finds some, it refuses unless you choose: REVERSE_COPY=1 copies the S3 Files tree back onto the old share with the application stopped and verifies it with checksums before switching, while FORCE=1 discards the new writes deliberately. The marker is set a few seconds ahead of the client clock (SKEW_MARGIN) so that small clock differences between the client and the file system cannot make the script’s own changes look like application writes; the application is stopped during cutover, so nothing legitimate happens in that margin. The check only sees changes made through the file system, so if anything writes to the prefix directly through the S3 API after cutover you need to account for that separately. And if umount fails at any point, the script stops rather than carrying on with a half switched host.
cat > s3files-cutover.sh << 'EOF'
#!/usr/bin/env bash
# s3files-cutover.sh
# Repoints the application's mount path from the old NFS share to S3 Files, and
# rolls it back safely.
# sudo ./s3files-cutover.sh cut over (rerunning after success is a no op)
# sudo ./s3files-cutover.sh rollback only if nothing was written since cutover
# sudo REVERSE_COPY=1 ./s3files-cutover.sh rollback
# copy S3 Files back to the old share, verify, then roll back
set -euo pipefail
ENV_FILE="${S3FILES_ENV:-./s3files.env}"
[[ -f "$ENV_FILE" ]] && source "$ENV_FILE"
ACTION="${1:-cutover}"
: "${MOUNT_PATH:?set MOUNT_PATH to the path your application uses}"
FSTAB="${FSTAB:-/etc/fstab}"
STATE_DIR="${STATE_DIR:-/var/lib/s3files-cutover}"
KEY=$(printf '%s' "$MOUNT_PATH" | tr '/' '_')
STATE="$STATE_DIR/${KEY}.state"
MARK="$STATE_DIR/${KEY}.cutover-time"
EXCL='.s3files-lost+found-*'
RSYNC_FLAGS=(-a --numeric-ids --no-devices --no-specials --exclude="$EXCL")
log() { echo "[$(date +%H:%M:%S)] $*"; }
die() { echo "ERROR: $*" >&2; exit 1; }
entry_field() { awk -v p="$MOUNT_PATH" -v f="$1" '$1 !~ /^#/ && $2 == p { print $f; exit }' "$FSTAB"; }
busy() { fuser -m "$MOUNT_PATH" >/dev/null 2>&1; }
# S3 Files always mounts through the local TLS proxy, so its source is 127.0.0.1
is_s3files_mounted() { [[ "$(findmnt -n -o SOURCE --mountpoint "$MOUNT_PATH" 2>/dev/null)" == 127.0.0.1:* ]]; }
write_fstab() { # write_fstab NEW_LINE COMMENT_OLD(0|1)
local tmp; tmp=$(mktemp)
awk -v p="$MOUNT_PATH" -v nl="$1" -v keep="$2" '
$1 == "#" && $2 == "pre-s3files:" && $4 == p { if (keep == 0) next }
$1 !~ /^#/ && $2 == p { if (keep == 1) print "# pre-s3files: " $0; print nl; next }
{ print }' "$FSTAB" > "$tmp"
cat "$tmp" > "$FSTAB"; rm -f "$tmp" # cat keeps the original file's owner and mode
}
cutover() {
: "${FS_ID:?set FS_ID}"
local type dev; type=$(entry_field 3); dev=$(entry_field 1)
if [[ "$type" == "s3files" ]]; then
[[ "$dev" == "${FS_ID}:/" ]] || die "fstab already points $MOUNT_PATH at $dev, not ${FS_ID}:/"
if is_s3files_mounted; then log "Already cut over and mounted; nothing to do"; exit 0; fi
log "fstab already cut over; mounting"; mount "$MOUNT_PATH"; findmnt -T "$MOUNT_PATH"; exit 0
fi
[[ "$type" =~ ^nfs ]] || die "expected an nfs entry for $MOUNT_PATH in $FSTAB, found '${type:-none}'"
busy && die "processes still use $MOUNT_PATH; stop the application first"
mkdir -p "$STATE_DIR"
local old_line root_stat opts="_netdev"
old_line=$(awk -v p="$MOUNT_PATH" '$1 !~ /^#/ && $2 == p { print; exit }' "$FSTAB")
root_stat=$(stat -c '%u:%g %a' "$MOUNT_PATH")
cp -p "$FSTAB" "$STATE_DIR/${KEY}.fstab.$(date +%Y%m%d%H%M%S)"
printf 'OLD_LINE=%q\nROOT_STAT=%q\nCUT_FS_ID=%q\n' "$old_line" "$root_stat" "$FS_ID" > "$STATE"
log "Unmounting the old NFS share at $MOUNT_PATH"
umount "$MOUNT_PATH" || die "umount failed; nothing was changed"
[[ "${NOFAIL:-0}" == 1 ]] && opts+=",nofail"
write_fstab "${FS_ID}:/ ${MOUNT_PATH} s3files ${opts} 0 0" 1
if ! mount "$MOUNT_PATH"; then
log "S3 Files mount failed; restoring the old entry"
write_fstab "$old_line" 0
mount "$MOUNT_PATH" || true
die "cutover aborted; old NFS entry restored"
fi
# Match the mount root's owner, group and mode to the old share's root
local now; now=$(stat -c '%u:%g %a' "$MOUNT_PATH")
if [[ "$now" != "$root_stat" ]]; then
chown "${root_stat% *}" "$MOUNT_PATH"; chmod "${root_stat#* }" "$MOUNT_PATH"
log "Mount root set to ${root_stat} (was ${now})"
fi
# Cutover marker, taken after the root fix and set slightly ahead so that client and
# server clock skew (up to SKEW_MARGIN seconds) cannot make our own changes look like
# application writes. The application is stopped, so nothing writes in that margin.
touch -d "@$(( $(date +%s) + ${SKEW_MARGIN:-5} ))" "$MARK"
findmnt -T "$MOUNT_PATH"
log "Cutover complete. Start the application and run s3files-validate.sh"
}
reverse_copy() {
local dev mp type opts rest stage=/mnt/.s3files-rollback
read -r dev mp type opts rest <<<"$OLD_LINE"
mkdir -p "$stage"
mountpoint -q "$stage" || mount -t "$type" -o "$opts" "$dev" "$stage"
log "Copying S3 Files back to the old share ($dev)"
rsync "${RSYNC_FLAGS[@]}" --delete --info=stats2 "${MOUNT_PATH%/}/" "$stage/"
log "Verifying the old share matches S3 Files (checksums)"
local diffs; diffs=$(rsync "${RSYNC_FLAGS[@]}" --delete -n -i --checksum "${MOUNT_PATH%/}/" "$stage/" | grep -v '^\.d' || true)
[[ -z "$diffs" ]] || { echo "$diffs" | head -20 >&2; die "reverse copy did not verify; still on S3 Files"; }
umount "$stage" && rmdir "$stage"
}
rollback() {
[[ -f "$STATE" ]] || die "no cutover state for $MOUNT_PATH in $STATE_DIR"
source "$STATE"
[[ "$(entry_field 3)" == "s3files" ]] || die "fstab entry for $MOUNT_PATH is not s3files; nothing to roll back"
busy && die "processes still use $MOUNT_PATH; stop the application first"
local lf; lf=$(find "$MOUNT_PATH" -maxdepth 1 -name "$EXCL" -type d 2>/dev/null | head -1)
[[ -n "$lf" && -n "$(ls -A "$lf" 2>/dev/null)" ]] && log "WARNING: conflict files exist in $lf; review them first"
# Any file or directory modified (mtime) or changed (ctime) after cutover means the
# application wrote to S3 Files, and the old share no longer has the latest data.
local changed
changed=$(find "$MOUNT_PATH" -path "${MOUNT_PATH%/}/${EXCL}" -prune -o \
\( -newer "$MARK" -o -cnewer "$MARK" \) -print -quit 2>/dev/null || true)
if [[ -n "$changed" ]]; then
if [[ "${REVERSE_COPY:-0}" == 1 ]]; then
reverse_copy
elif [[ "${FORCE:-0}" == 1 ]]; then
log "WARNING: discarding writes made after cutover (first seen: $changed)"
else
die "writes after cutover detected (first seen: $changed). The old share does not have them.
Rerun with REVERSE_COPY=1 (application stopped) to copy them back first, or FORCE=1 to discard them."
fi
fi
log "Unmounting S3 Files at $MOUNT_PATH"
umount "$MOUNT_PATH" || die "umount failed; rollback stopped with S3 Files still in fstab"
write_fstab "$OLD_LINE" 0
mount "$MOUNT_PATH" || die "old NFS entry restored in $FSTAB but mount failed"
findmnt -T "$MOUNT_PATH"
log "Rolled back to $(awk '{print $1}' <<<"$OLD_LINE")"
}
case "$ACTION" in
cutover) cutover ;;
rollback) rollback ;;
*) die "usage: s3files-cutover.sh [cutover|rollback]" ;;
esac
EOF
chmod +x s3files-cutover.shThe _netdev option makes the system wait for networking before mounting, and the AWS mounting documentation warns that leaving it out can make an instance unresponsive at boot. nofail is off by default here and can be enabled with NOFAIL=1; whether you want it depends on the application, since some services are better off failing to start than starting against an empty directory, so make that a deliberate choice rather than a default you copied.
For Kubernetes workloads the equivalent move is to repoint the PersistentVolume behind your existing claim, so that the pod spec and its mountPath stay the same. Check the current EKS documentation for the supported driver and volume definition for S3 Files before you script that part.
9. Smoke Testing That Nothing Changed
I want to be precise about what the validation script can and cannot tell you, because a green run is easy to over read. These are smoke tests. The local mode checks on one host that renaming over an existing file leaves the new content, that a byte range can be rewritten in place, that appends work, that two processes contend correctly for a lock, how long a thousand small writes take, how long a new file takes to appear in the bucket, and that the conflict directory is empty.
The behaviours that actually break shared file system applications involve more than one client, so there are two host modes as well. Both hosts use the same RUN_ID and therefore the same lock file and directory: lock-holder takes an exclusive lock on one host and holds it, and lock-contender on the other host must fail to take the lock while it is held and then succeed after it is released. For renames, rename-writer on one host repeatedly replaces a file by writing a temporary file and renaming it over the target, while rename-reader on the other host reads the target in a loop and checks an embedded checksum, counting any read that sees a partial file. Zero partial reads over a few minutes is useful evidence that readers on another client see either the old file or the new one, but it is evidence over a window, not a proof of atomicity under every timing. Data integrity is established by the migration’s checksum verification, not by these tests.
cat > s3files-validate.sh << 'EOF'
#!/usr/bin/env bash
# s3files-validate.sh: smoke tests for the behaviours applications usually depend on.
# Passing tests are evidence, not proof; run your own integration tests as well.
#
# ./s3files-validate.sh local one host: rename, in place write, append,
# lock contention, small writes, export lag
# RUN_ID=x ./s3files-validate.sh lock-holder host A: holds a shared lock
# RUN_ID=x ./s3files-validate.sh lock-contender host B: must be blocked, then succeed
# RUN_ID=x ./s3files-validate.sh rename-writer host A: replaces a file by rename in a loop
# RUN_ID=x ./s3files-validate.sh rename-reader host B: checks it never sees a partial file
#
# Two host tests use one shared directory and lock file, so both hosts contend for
# the same lock. Test files are exported to S3 like any other file and removed at the end.
set -euo pipefail
ENV_FILE="${S3FILES_ENV:-./s3files.env}"
[[ -f "$ENV_FILE" ]] && source "$ENV_FILE"
MODE="${1:-local}"
: "${MOUNT_PATH:?set MOUNT_PATH}"
RUN_ID="${RUN_ID:-default}"
D="${MOUNT_PATH%/}/.s3files-validate/${RUN_ID}"
HOLD="${HOLD:-180}" # seconds the holder keeps the lock
DURATION="${DURATION:-120}" # seconds the rename writer and reader run
mkdir -p "$D"
FAILED=0
pass() { echo "PASS $*"; }
fail() { echo "FAIL $*"; FAILED=1; }
info() { echo "INFO $*"; }
case "$MODE" in
local)
norm_prefix() { local p="${1#/}"; [[ -n "$p" && "$p" != */ ]] && p="$p/"; printf '%s' "$p"; }
PREFIX=$(norm_prefix "${PREFIX:-}")
T="$D/local-$(hostname)-$$"; mkdir -p "$T"
echo v1 > "$T/target"; echo v2 > "$T/target.tmp"; mv "$T/target.tmp" "$T/target"
[[ "$(cat "$T/target")" == "v2" ]] && pass "rename over an existing file (final content)" || fail "rename"
printf 'AAAAAAAAAA' > "$T/inplace"
printf 'BB' | dd of="$T/inplace" bs=1 seek=4 conv=notrunc status=none
[[ "$(cat "$T/inplace")" == "AAAABBAAAA" ]] && pass "in place write of a byte range" || fail "in place write"
echo one > "$T/append"; echo two >> "$T/append"
[[ "$(wc -l < "$T/append")" -eq 2 ]] && pass "append" || fail "append"
: > "$T/lock"
( flock -x 9; sleep 3 ) 9>"$T/lock" &
sleep 1
if flock -n -x "$T/lock" true; then fail "lock exclusion between processes on one host"
else pass "lock exclusion between processes on one host"; fi
wait
start=$(date +%s%N)
for i in $(seq 1 1000); do echo "$i" > "$T/small_$i"; done
info "1000 small file writes took $(( ($(date +%s%N) - start) / 1000000 )) ms"
if [[ -n "${BUCKET:-}" ]]; then
key="${PREFIX}${T#${MOUNT_PATH%/}/}/exportcheck"
date +%s > "$T/exportcheck"
secs=0
until aws s3api head-object --bucket "$BUCKET" --key "$key" >/dev/null 2>&1; do
sleep 5; secs=$((secs + 5))
(( secs >= 600 )) && { fail "file not exported to S3 within 600s"; break; }
done
(( secs < 600 )) && info "export to S3 visible after ~${secs}s (S3 Files exports after 60s without writes)"
else
info "BUCKET not set; export lag not measured"
fi
LF=$(find "$MOUNT_PATH" -maxdepth 1 -name '.s3files-lost+found-*' -type d | head -1)
if [[ -n "$LF" && -n "$(ls -A "$LF" 2>/dev/null)" ]]; then fail "conflict files present in $LF"
else pass "no conflict files"; fi
rm -rf "$T"
;;
lock-holder)
: > "$D/shared.lock"
echo "holder $(hostname) waiting for the lock"
(
flock -x 9
echo "$(hostname) $(date +%s)" > "$D/held"
echo "holding $D/shared.lock for ${HOLD}s; start lock-contender on the other host now"
sleep "$HOLD"
rm -f "$D/held"
) 9>"$D/shared.lock"
echo "released"
;;
lock-contender)
deadline=$((SECONDS + 120))
until [[ -f "$D/held" ]]; do
(( SECONDS > deadline )) && { fail "never saw the holder's marker in $D"; exit 1; }
ls "$D" >/dev/null; sleep 2
done
info "holder is $(cat "$D/held")"
if flock -n -x "$D/shared.lock" true; then fail "acquired the lock while another host held it"
else pass "blocked while another host held the lock"; fi
if flock -w $((HOLD + 60)) -x "$D/shared.lock" true; then pass "acquired the lock after the holder released it"
else fail "lock was never released to this host"; fi
;;
rename-writer)
end=$((SECONDS + DURATION)); n=0
while (( SECONDS < end )); do
body=$(head -c 49152 /dev/urandom | base64 -w0)
tmp="$D/.target.tmp.$(hostname).$$"
printf '%s\n%s\n' "$(printf '%s' "$body" | md5sum | cut -d' ' -f1)" "$body" > "$tmp"
mv -f "$tmp" "$D/target"; n=$((n + 1))
done
info "writer replaced $D/target $n times by rename"
;;
rename-reader)
end=$((SECONDS + DURATION)); ok=0; torn=0; absent=0
while (( SECONDS < end )); do
if ! content=$(cat "$D/target" 2>/dev/null); then absent=$((absent + 1)); sleep 0.2; continue; fi
sum=$(head -1 <<<"$content"); body=$(sed -n 2p <<<"$content")
if [[ -n "$sum" && "$(printf '%s' "$body" | md5sum | cut -d' ' -f1)" == "$sum" ]]; then ok=$((ok + 1))
else torn=$((torn + 1)); fi
done
info "reads: complete=$ok partial=$torn absent=$absent"
(( ok > 0 )) || fail "no complete reads; start rename-writer first"
(( torn == 0 )) && pass "no partial files observed during ${DURATION}s of renames" || fail "$torn partial reads"
;;
cleanup) rm -rf "$D"; info "removed $D" ;;
*) echo "unknown mode $MODE"; exit 2 ;;
esac
exit $FAILED
EOF
chmod +x s3files-validate.shA clean run tells you the basic semantics match. It does not tell you that performance is acceptable for your workload, so the last step should always be your real integration or load tests against the new mount, compared with a baseline taken on the old share. If small file performance or export delay is outside what your application tolerates, FSx for NetApp ONTAP with tiering to S3 is the usual alternative when you need full NFS behaviour with object storage economics for cold data.
10. Putting It Together
The end to end sequence for a single share looks like this:
- Run
file_mount_recon.py --ssmand work through the findings, paying particular attention to the S3 Files signals and to any mounts of the share you did not know about. - Fill in
s3files.env, then runs3files-provision.shonce from a pipeline or an administrative host. - Run
s3files-client-prep.shon every host that mounts the share, ideally baked into your AMI or bootstrap so new instances come up ready. - Run
s3files-migrate.sh preflightand resolve every failed gate; then runs3files-migrate.sh bulkwhile the application is live, as many times as you like. - In the cutover window, stop the application, run
s3files-migrate.sh finalands3files-migrate.sh verify-sync, then runs3files-cutover.shon each host. - Start the application and run
s3files-validate.sh localand the two host modes, then your own tests. If something is wrong before the application has written anything,s3files-cutover.sh rollbackis a simple swap back to the untouched old share; after it has written, useREVERSE_COPY=1so those writes go back with you.
The broader point is that S3 Files changes the default answer for a lot of legacy shared storage. Workloads that were stuck on NFS servers or on EFS purely because the code assumed a file system can now sit on S3, with its durability and its pricing for cold data, and for many of them the migration really is mostly an infrastructure exercise. The places where it is not (hard links and extended attributes, shared access through both the mount and the S3 API, tight freshness requirements between producers and consumers, and heavy small file churn) are exactly the places your preflight and testing need to probe, and it is far cheaper to find them with a script than with an incident.
Sources
- Amazon S3 Files overview and architecture (AWS documentation)
- Synchronization between the file system and the bucket (AWS documentation)
- Unsupported features, limits and quotas (AWS documentation)
- Prerequisites and IAM policies (AWS documentation)
- Mounting S3 Files (AWS documentation)
aws s3files create-file-system(AWS CLI reference)- Amazon S3 Files is GA: mounting S3 buckets as a file system, compared with EFS (Classmethod)
Appendix: file_mount_recon.py
The complete recon script from section 3. It needs Python 3.9 or later and boto3, and the AWS managed ReadOnlyAccess policy covers its API calls; --ssm additionally needs ssm:SendCommand and ssm:GetCommandInvocation.
cat > file_mount_recon.py << 'EOF'
#!/usr/bin/env python3
"""
file_mount_recon.py: full reconnaissance of file storage and file mounts in an AWS account.
Covers NFS, SMB/CIFS, Lustre and FUSE (Mountpoint for S3, s3fs, goofys, rclone, gcsfuse,
blobfuse, sshfs, juicefs ...) plus other network file systems (GlusterFS, CephFS, 9p, DAVFS,
GPFS, BeeGFS) seen on hosts. The script is read only: it never copies, mounts or changes data.
Server side (read only API calls, every enabled region or --regions)
* Amazon EFS: file systems, mount targets (+ ENI and security groups), access points,
file system policy, backup policy, lifecycle, replication, protection, tags, AWS Backup
recovery points, and every CloudWatch metric published for the file system.
* Amazon S3 Files: file systems and mount targets (if your boto3 knows the service).
* Amazon FSx: ONTAP (SVM NFS/SMB endpoints, volumes, junction paths, tiering), OpenZFS
(volumes, NFS exports), Windows File Server (AD, aliases, throughput), Lustre (mount
name, deployment, data repository associations), File Cache, backups and metrics.
* AWS Storage Gateway: gateways, NFS and SMB file shares (full describe), SMB settings.
* S3 buckets (to resolve FUSE mounts of S3 to buckets in this account).
* Consumers without a login: Lambda functions with FileSystemConfigs, ECS task
definitions with EFS volumes, EKS clusters running the EFS CSI driver add on.
* EC2 instances, ENIs, security groups and subnets, so endpoints and clients can be joined
and port reachability (2049 NFS, 445 SMB, 988 Lustre) evaluated rule by rule.
Client side (only with --ssm: read only collectors via SSM Run Command)
Linux: /proc/mounts, /etc/fstab (secrets redacted), autofs and systemd units,
/proc/self/mountstats fully parsed (options, age, caps, NFSv4 state, 27 event
counters, byte counters, RPC transport, per op ops/retrans/timeouts/bytes/RTT/exec
time/errors), nfsstat, df with a timeout (hung mounts are detected and timed),
statfs, open files and processes per mount, FUSE daemons and their buckets/remotes,
CIFS DebugData/Stats, Lustre client version and lfs df, TLS proxy state, live
connections to 2049/445/988 with tcp_info, packages and versions, RPC tunables,
and NFS server role detection (exports, exportfs -v, nfsd threads and stats).
Windows: SMB mappings, live SMB sessions (dialect, encryption, signing, multichannel),
mapped and network logical disks with capacity, persistent per user drive maps
from the registry, net use, NFS client mounts, DFS referral cache, SMB/NFS shares
served by the host and their sessions, open files, connections to 445/2049/988.
Responsiveness: every remote mount on Linux gets four timed, read only probes: statfs (df, which
always reaches the server), getattr on the mount root, a directory read, and a TCP connect to
the server port (2049, 445 or 988) on the address the client resolves. Windows shares get a
timed listing and a TCP connect. The combination separates "server unreachable" from "server
up but the file system is not answering" from "slow".
Health checks: capacity (space and inodes on each mount, FSx storage utilisation, Storage
Gateway cache and upload backlog, EFS burst credits), latency (average RPC round trips and
client side queueing from mountstats, FSx utilisation metrics near saturation, cross AZ
mounts), cross region use (client and server in different regions, inter region peering,
FUSE mounts of buckets elsewhere) and permissions (world writable share roots, root squash
denials, CIFS credential files readable by others, passwords in fstab, EFS without a file
system policy or TLS requirement, Storage Gateway shares open to 0.0.0.0/0, NoSquash or guest
SMB access, and OpenZFS exports with no_root_squash).
Cross account: each mount's target address is traced through the client subnet's route table
(local, VPC peering, Transit Gateway, Cloud WAN, VPN or Direct Connect) and attributed to an
account where the data allows: peering connections name the peer account, shared VPC subnets
name their owner, and EFS or S3 Files IDs and S3 buckets that no scanned account owns are
flagged. In the other direction, EFS file system policies that grant other accounts, security
group rules that admit other accounts' groups or CIDRs outside the VPC, and (with --flow-logs)
VPC Flow Logs show who else connects to each file server. With --org or --accounts the run
covers several accounts and resolves both ends of a cross account mount exactly.
Everything is joined into one inventory: each client mount is resolved to the resource
behind it (EFS, S3 Files, FSx, File Cache, Storage Gateway, S3 bucket via FUSE, a self
managed EC2 file server, or "unknown"), the security group path is checked, and findings
are raised (hung mounts, blocked ports, soft mounts, missing TLS, old SMB dialects,
unpersisted mounts, idle or unbacked file systems, and S3 Files migration signals such as
lock use, renames and small writes).
Outputs (in --out, default ./file-mount-recon-<timestamp>):
inventory.json everything, normalised and correlated (includes full parsed mountstats)
mounts.csv one row per client mount (~135 columns)
servers.csv one row per file server resource
findings.csv one row per finding
flow_sources.csv with --flow-logs: every source address seen connecting to a file server
report.md human readable summary
raw/<region>/... raw API responses and raw client bundles per instance
Permissions: the AWS managed ReadOnlyAccess policy covers the API calls. --ssm also needs
ssm:SendCommand, ssm:GetCommandInvocation, and with --ssm-bucket s3:ListBucket/GetObject.
--flow-logs needs logs:StartQuery and logs:GetQueryResults (Logs Insights charges per GB
scanned). --org needs organizations:ListAccounts and sts:AssumeRole into --role-name, which
should be a read only role present in every member account.
SSM output inline is capped at 24,000 characters; the bundles are gzip+base64 to fit, and
--ssm-bucket removes the cap for very busy hosts.
Usage:
python3 file_mount_recon.py # all enabled regions, API only
python3 file_mount_recon.py --regions af-south-1 eu-west-1 --ssm
python3 file_mount_recon.py --ssm --ssm-bucket my-ssm-output --days 30
python3 file_mount_recon.py --ssm --flow-logs --org --role-name FileMountReconReadOnly
"""
import argparse
import base64
import csv
import datetime as dt
import gzip
import ipaddress
import json
import os
import re
import sys
import threading
import time
from collections import defaultdict
from concurrent.futures import ThreadPoolExecutor, as_completed
try:
import boto3
from botocore.config import Config
from botocore.exceptions import ClientError, UnknownServiceError, BotoCoreError
except ImportError: # pragma: no cover
sys.exit("boto3 is required: pip install boto3")
UTC = dt.timezone.utc
NOW = dt.datetime.now(UTC)
BOTO_CFG = Config(retries={"max_attempts": 10, "mode": "adaptive"}, user_agent_extra="file-mount-recon/1.0")
LOCK = threading.Lock()
# --------------------------------------------------------------------------------------
# Client side collector (runs on each instance through SSM, read only)
# --------------------------------------------------------------------------------------
CLIENT_COLLECTOR = r'''#!/bin/bash
export LC_ALL=C PATH="$PATH:/usr/sbin:/sbin:/usr/local/sbin:/usr/local/bin"
OUT=$(mktemp /tmp/.nfsrecon.XXXXXX)
sec(){ echo "===SECTION $1==="; }
# Remote / network / FUSE file systems. fuseblk (local NTFS etc.) and desktop FUSE noise are excluded.
RFS='^(nfs|nfs4|cifs|smb3|smbfs|lustre|glusterfs|ceph|9p|davfs|afs|gpfs|beegfs|orangefs|ocfs2|gfs2|virtiofs|fuse|fuse\..+)$'
NOISE='^fuse\.(gvfsd-fuse|portal|lxcfs|doc|xdg-document-portal|snapfuse)$'
# credentials= is a file path, not a secret, and is kept so its permissions can be checked
redact(){ sed -E 's/((pass(word)?|passwd|secret|token|key)=)[^, ]*/\1<redacted>/Ig'; }
remotelines(){ awk -v r="$RFS" -v n="$NOISE" '$3 ~ r && $3 !~ n' /proc/mounts; }
nfsmounts(){ remotelines | awk '{print $2}' | sed -e 's/\\040/ /g' -e 's/\\011/\t/g' -e 's/\\134/\\/g'; }
{
sec meta
TOKEN=$(curl -s -m 2 -X PUT http://169.254.169.254/latest/api/token -H 'X-aws-ec2-metadata-token-ttl-seconds: 120' 2>/dev/null)
for k in instance-id placement/availability-zone placement/availability-zone-id placement/region local-ipv4 instance-type ami-id mac; do
printf '%s=%s\n' "$k" "$(curl -s -m 2 -H "X-aws-ec2-metadata-token: $TOKEN" "http://169.254.169.254/latest/meta-data/$k" 2>/dev/null)"
done
echo "hostname=$(hostname -f 2>/dev/null || hostname)"
echo "kernel=$(uname -r)"
echo "arch=$(uname -m)"
echo "uptime_s=$(cut -d' ' -f1 /proc/uptime)"
echo "collected_at=$(date -u +%Y-%m-%dT%H:%M:%SZ)"
echo "nproc=$(nproc 2>/dev/null)"
echo "mem_kb=$(awk '/MemTotal/{print $2}' /proc/meminfo)"
sec os_release; cat /etc/os-release 2>/dev/null
sec packages
(rpm -q --qf '%{NAME} %{VERSION}-%{RELEASE}\n' nfs-utils amazon-efs-utils stunnel stunnel5 autofs rpcbind cifs-utils samba-client lustre-client mount-s3 s3fs-fuse fuse fuse3 glusterfs-fuse ceph-common sshfs rclone 2>/dev/null | grep -v 'not installed'
dpkg-query -W -f='${Package} ${Version}\n' nfs-common nfs-kernel-server amazon-efs-utils stunnel4 autofs rpcbind cifs-utils smbclient lustre-client-utils mount-s3 s3fs fuse fuse3 glusterfs-client ceph-common sshfs rclone 2>/dev/null) | sort -u
for b in mount-s3 s3fs goofys rclone gcsfuse blobfuse2 juicefs; do command -v $b >/dev/null && echo "$b $($b --version 2>&1 | head -1)"; done
command -v mount.efs >/dev/null && echo "mount.efs $(mount.efs --version 2>&1 | head -1)"
command -v mount.s3files >/dev/null && echo "mount.s3files $(mount.s3files --version 2>&1 | head -1)"
sec proc_mounts; remotelines | redact
sec fstab; grep -Ev '^[[:space:]]*(#|$)' /etc/fstab 2>/dev/null \
| awk -v r="$RFS" '$3 ~ r || $3 ~ /^(efs|s3files|fuse\.|mount-s3)/ || $1 ~ /:/ || $1 ~ /^\/\// || $1 ~ /@tcp/ || $4 ~ /_netdev/' | redact
sec fuse_procs
ps -eo pid,etimes,rss,args 2>/dev/null \
| grep -E 'mount-s3|mountpoint-s3|s3fs|goofys|rclone .*mount|gcsfuse|blobfuse|sshfs|geesefs|juicefs|alluxio-fuse|s3backer|catfs' \
| grep -v grep | redact
sec cifs
[ -r /proc/fs/cifs/DebugData ] && { echo "--- DebugData"; timeout -k 2 5 cat /proc/fs/cifs/DebugData | redact; }
[ -r /proc/fs/cifs/Stats ] && { echo "--- Stats"; timeout -k 2 5 cat /proc/fs/cifs/Stats; }
command -v smbstatus >/dev/null && { echo "--- smbstatus"; timeout -k 2 10 smbstatus -b 2>&1 | head -200; }
sec lustre
command -v lctl >/dev/null && echo "lustre_client_version=$(timeout -k 2 5 lctl get_param -n version 2>/dev/null | head -1)"
command -v lfs >/dev/null && { echo "--- lfs df"; timeout -k 2 15 lfs df 2>&1; echo "--- lfs df -i"; timeout -k 2 15 lfs df -i 2>&1; }
command -v lctl >/dev/null && { echo "--- lctl dl"; timeout -k 2 10 lctl dl 2>&1 | head -200; }
sec other_remote
command -v gluster >/dev/null && { echo "--- gluster volume info"; timeout -k 2 10 gluster volume info 2>&1 | head -200; }
[ -d /etc/ceph ] && { echo "--- /etc/ceph"; ls -1 /etc/ceph; }
command -v showmount >/dev/null && { echo "--- showmount -e localhost"; timeout -k 2 5 showmount -e localhost 2>&1; }
sec mountstats; awk '/^device /{p=($0 ~ / with fstype nfs/)} p' /proc/self/mountstats 2>/dev/null
sec nfsstat_m; timeout -k 2 10 nfsstat -m 2>&1
sec proc_net_rpc_nfs; cat /proc/net/rpc/nfs 2>/dev/null
sec df
nfsmounts | while IFS= read -r m; do
s=$(date +%s%N); o=$(timeout -k 2 5 df -PT -B1 "$m" 2>&1); rc=$?; e=$(date +%s%N)
printf '%s|%s|%s|%s\n' "$m" "$rc" "$(( (e - s) / 1000000 ))" "$(printf '%s' "$o" | tail -1)"
done
sec probes
# One line per remote mount: mountpoint|fstype|host|ip|port|tcp|tcp_ms|stat_rc|stat_ms|readdir_rc|readdir_ms
# statfs (the df section) always reaches the server; stat and readdir may be answered from the
# client's attribute and directory caches, so they show "can the application use it", not liveness.
REGION_NAME=$(curl -s -m 2 -H "X-aws-ec2-metadata-token: $TOKEN" http://169.254.169.254/latest/meta-data/placement/region 2>/dev/null)
remotelines | while read -r src mpe fstype opts _rest; do
mp=$(printf '%s' "$mpe" | sed -e 's/\\040/ /g' -e 's/\\011/\t/g' -e 's/\\134/\\/g')
host=""; port=""
case "$fstype" in
nfs*) host="${src%:*}"; port=2049 ;;
cifs|smb3|smbfs) host=$(printf '%s' "$src" | sed -E 's#^//([^/]+)/?.*#\1#'); port=445 ;;
lustre) host="${src%%@*}"; port=988 ;;
esac
if [ "$host" = "127.0.0.1" ]; then
# fstab keeps the \040 escapes; ENVIRON avoids awk -v turning them back into spaces
fl=$(M="$mpe" awk '$1 !~ /^#/ && $2 == ENVIRON["M"]' /etc/fstab | head -1)
dev=$(printf '%s' "$fl" | awk '{print $1}'); ftype=$(printf '%s' "$fl" | awk '{print $3}')
mti=$(printf '%s' "$fl" | awk '{print $4}' | tr ',' '\n' | sed -n 's/^mounttargetip=//p')
fsid="${dev%%:*}"
if [ -n "$mti" ]; then host="$mti"
elif [ "$ftype" = "efs" ] && [ -n "$REGION_NAME" ] && [ "${fsid#fs-}" != "$fsid" ]; then host="${fsid}.efs.${REGION_NAME}.amazonaws.com"
else host=""; fi
fi
ip=""; tcp=""; tcp_ms=""
if [ -n "$host" ]; then
case "$host" in *[!0-9.]*) ip=$(getent ahostsv4 "$host" 2>/dev/null | awk 'NR == 1 {print $1}') ;; *) ip="$host" ;; esac
fi
if [ -n "$ip" ] && [ -n "$port" ]; then
s=$(date +%s%N)
if timeout 3 bash -c "echo > /dev/tcp/$ip/$port" 2>/dev/null; then tcp=open; else tcp=closed; fi
tcp_ms=$(( ($(date +%s%N) - s) / 1000000 ))
fi
s=$(date +%s%N); root_stat=$(timeout -k 2 5 stat -c '%a:%u:%g' "$mp" 2>/dev/null); st_rc=$?; st_ms=$(( ($(date +%s%N) - s) / 1000000 ))
# pipefail so an ls error is not hidden behind head's exit code; 141 is ls stopped by head
s=$(date +%s%N)
timeout -k 2 5 bash -c 'set -o pipefail; ls -f -1 -- "$1" 2>"$2" | head -n 3 >/dev/null; rc=$?; [ $rc -eq 141 ] && rc=0; exit $rc' _ "$mp" "$OUT.e"
rd_rc=$?; rd_ms=$(( ($(date +%s%N) - s) / 1000000 ))
rd_err=$(grep -o 'Permission denied\|Stale file handle\|Input/output error\|Transport endpoint is not connected\|Host is down\|No such device' "$OUT.e" 2>/dev/null | head -1)
printf '%s|%s|%s|%s|%s|%s|%s|%s|%s|%s|%s|%s|%s\n' "$mp" "$fstype" "$host" "$ip" "$port" "$tcp" "$tcp_ms" "$st_rc" "$st_ms" "$rd_rc" "$rd_ms" "$root_stat" "$rd_err"
done
sec statfs
nfsmounts | while IFS= read -r m; do
printf '%s|%s\n' "$m" "$(timeout -k 2 5 stat -f -c 'bsize=%S blocks=%b bfree=%f bavail=%a files=%c ffree=%d namelen=%l' "$m" 2>&1 | tail -1)"
done
sec open_files
nfsmounts > "$OUT.m"
if [ -s "$OUT.m" ]; then
timeout -k 2 30 find /proc/[0-9]*/fd /proc/[0-9]*/cwd -mindepth 0 -maxdepth 1 -type l -printf '%p\t%l\n' 2>/dev/null \
| awk -F'\t' -v mf="$OUT.m" '
BEGIN { while ((getline l < mf) > 0) m[n++] = l }
{ split($1, a, "/"); pid = a[3]
for (i = 0; i < n; i++) { p = m[i]; if ($2 == p || index($2, p "/") == 1) { c[p SUBSEP pid]++ ; break } } }
END { for (k in c) { split(k, b, SUBSEP); print b[1] "\t" b[2] "\t" c[k] } }' \
| while IFS=$'\t' read -r mp pid cnt; do
printf '%s|%s|%s|%s|%s\n' "$mp" "$pid" "$cnt" "$(cat /proc/$pid/comm 2>/dev/null)" "$(stat -c %U /proc/$pid 2>/dev/null)"
done
fi
rm -f "$OUT.m"
sec systemd_units; systemctl list-units --all --no-legend --no-pager --type=mount,automount 2>/dev/null
sec autofs
for f in $(ls -1 /etc/auto.master /etc/auto.master.d/* /etc/auto.* 2>/dev/null | sort -u); do
[ -f "$f" ] && { echo "--- $f"; grep -Ev '^[[:space:]]*(#|$)' "$f"; }
done
systemctl is-active autofs 2>/dev/null | sed 's/^/autofs_service=/'
sec efs_utils_conf; grep -Ev '^[[:space:]]*(#|$)' /etc/amazon/efs/efs-utils.conf 2>/dev/null
sec credfiles
# Paths and modes of CIFS credentials files named in fstab (paths only, never contents)
grep -Ev '^[[:space:]]*(#|$)' /etc/fstab 2>/dev/null | awk '{print $4}' | tr ',' '\n' \
| sed -n 's/^cred\(entials\)\{0,1\}=//p' | sort -u | while IFS= read -r f; do
printf '%s|%s\n' "$f" "$(stat -c '%a %U %G' "$f" 2>&1 | head -1)"
done
sec tls_state
for f in $(ls -1d /var/run/efs/* /run/efs/* /var/run/s3files/* /run/s3files/* 2>/dev/null | xargs -r -n1 readlink -f | sort -u); do
[ -f "$f" ] && { echo "--- $(basename "$f")"; head -c 4000 "$f"; echo; }
done
sec tls_proxies; ps -eo pid,etimes,rss,args 2>/dev/null | grep -E 'stunnel|efs-proxy|s3files' | grep -v grep
sec connections; timeout -k 2 10 ss -tnpi '( dport = :2049 or sport = :2049 or dport = :445 or sport = :445 or dport = :988 or sport = :988 or dport = :111 )' 2>/dev/null
sec rpc_tunables
sysctl sunrpc.tcp_slot_table_entries sunrpc.tcp_max_slot_table_entries sunrpc.udp_slot_table_entries fs.nfs.nfs_callback_tcpport fs.nfs.idmap_cache_timeout 2>/dev/null
for p in /sys/module/nfs/parameters/* /sys/module/sunrpc/parameters/*; do [ -r "$p" ] && echo "$p=$(cat "$p" 2>/dev/null)"; done
sec nfs_server
[ -f /etc/exports ] && { echo "--- /etc/exports"; grep -Ev '^[[:space:]]*(#|$)' /etc/exports; }
for f in /etc/exports.d/*; do [ -f "$f" ] && { echo "--- $f"; grep -Ev '^[[:space:]]*(#|$)' "$f"; }; done
command -v exportfs >/dev/null && { echo "--- exportfs -v"; timeout -k 2 10 exportfs -v 2>&1; }
[ -r /proc/fs/nfsd/versions ] && echo "nfsd_versions=$(cat /proc/fs/nfsd/versions)"
[ -r /proc/fs/nfsd/threads ] && echo "nfsd_threads=$(cat /proc/fs/nfsd/threads)"
systemctl is-active nfs-server 2>/dev/null | sed 's/^/nfs_server_service=/'
[ -r /proc/net/rpc/nfsd ] && { echo "--- /proc/net/rpc/nfsd"; cat /proc/net/rpc/nfsd; }
sec end
} > "$OUT" 2>&1
echo "NFSRECON1:$(gzip -9c "$OUT" | base64 -w0)"
rm -f "$OUT"
'''
# Windows collector (AWS-RunPowerShellScript, read only). SSM runs as SYSTEM, so per user
# drive letters are read from the registry (HKU\<sid>\Network) as well as live SMB sessions.
WINDOWS_COLLECTOR = r'''
$ErrorActionPreference = 'SilentlyContinue'
function Q($b) { try { & $b } catch { @{ error = $_.Exception.Message } } }
$r = [ordered]@{}
$r.meta = [ordered]@{
hostname = $env:COMPUTERNAME; os = (Get-CimInstance Win32_OperatingSystem).Caption
build = [Environment]::OSVersion.Version.ToString(); collected_at = (Get-Date).ToUniversalTime().ToString('o')
domain = (Get-CimInstance Win32_ComputerSystem).Domain }
$r.smb_mappings = Q { Get-SmbMapping | Select-Object LocalPath,RemotePath,Status,UserName,ShareType,RequireIntegrity,RequirePrivacy }
$r.smb_connections = Q { Get-SmbConnection | Select-Object ServerName,ShareName,UserName,Credential,Dialect,NumOpens,Encrypted,Signed,ContinuouslyAvailable,Redirected }
$r.smb_multichannel = Q { Get-SmbMultichannelConnection | Select-Object ServerName,ClientIpAddress,ServerIpAddress,ClientInterfaceIndex,Selected }
$r.smb_client_config = Q { Get-SmbClientConfiguration | Select-Object RequireSecuritySignature,EnableSecuritySignature,EnableMultiChannel,DirectoryCacheLifetime,FileInfoCacheLifetime,SessionTimeout,EnableInsecureGuestLogons }
$r.mapped_logical_disks = Q { Get-CimInstance Win32_MappedLogicalDisk | Select-Object DeviceID,ProviderName,FileSystem,Size,FreeSpace,VolumeName,SessionID }
$r.network_logical_disks = Q { Get-CimInstance Win32_LogicalDisk -Filter 'DriveType=4' | Select-Object DeviceID,ProviderName,FileSystem,Size,FreeSpace }
$r.persistent_user_drives = Q {
New-PSDrive -Name HKU -PSProvider Registry -Root HKEY_USERS -ErrorAction SilentlyContinue | Out-Null
Get-ChildItem 'HKU:\' | ForEach-Object {
$sid = $_.PSChildName
Get-ChildItem "HKU:\$sid\Network" -ErrorAction SilentlyContinue | ForEach-Object {
$p = Get-ItemProperty $_.PSPath
[pscustomobject]@{ Sid = $sid; Drive = $_.PSChildName; RemotePath = $p.RemotePath; UserName = $p.UserName; ProviderName = $p.ProviderName } } } }
$r.net_use = Q { (net use) -join "`n" }
$r.nfs_client_mounts = Q { if (Get-Command mount.exe -ErrorAction SilentlyContinue) { (mount.exe) -join "`n" } }
$r.nfs_client_feature = Q { (Get-WindowsFeature NFS-Client -ErrorAction SilentlyContinue).InstallState }
$r.dfs_client_cache = Q { (dfsutil cache referral 2>$null) -join "`n" }
$r.smb_shares_served = Q { Get-SmbShare | Where-Object { $_.Special -eq $false } | Select-Object Name,Path,Description,EncryptData,CurrentUsers,FolderEnumerationMode,ContinuouslyAvailable }
$r.smb_server_sessions = Q { Get-SmbSession | Select-Object ClientComputerName,ClientUserName,NumOpens,Dialect,SecondsExists,SecondsIdle }
$r.nfs_shares_served = Q { Get-NfsShare | Select-Object Name,Path,Availability,EnableUnmappedAccess,Authentication }
$r.open_files_by_share = Q { Get-SmbOpenFile | Group-Object ShareRelativePath | Select-Object Count,Name -First 200 }
$r.connections = Q { Get-NetTCPConnection -RemotePort 445,2049,988 -ErrorAction SilentlyContinue | Select-Object LocalAddress,RemoteAddress,RemotePort,State,OwningProcess }
$targets = @()
foreach ($m in @($r.smb_mappings)) { if ($m.RemotePath) { $targets += [pscustomobject]@{ Path = $m.RemotePath; Port = 445 } } }
foreach ($d in @($r.mapped_logical_disks) + @($r.network_logical_disks)) { if ($d.ProviderName) { $targets += [pscustomobject]@{ Path = $d.ProviderName; Port = 445 } } }
foreach ($d in @($r.persistent_user_drives)) { if ($d.RemotePath) { $targets += [pscustomobject]@{ Path = $d.RemotePath; Port = 445 } } }
if ($r.nfs_client_mounts -is [string]) {
foreach ($line in ($r.nfs_client_mounts -split "`n")) { if ($line -match '^\s*[A-Za-z]:\s+(\S+)') { $targets += [pscustomobject]@{ Path = $Matches[1]; Port = 2049 } } } }
$seen = @{}
$r.probes = foreach ($t in $targets) {
if ($seen.ContainsKey($t.Path.ToLower())) { continue }; $seen[$t.Path.ToLower()] = 1
$h = (($t.Path -replace '^\\\\', '') -split '\\')[0]
$ip = $null
if ($h -match '^\d+\.\d+\.\d+\.\d+$') { $ip = $h } else {
$ip = (Resolve-DnsName -Name $h -Type A -QuickTimeout -ErrorAction SilentlyContinue | Where-Object { $_.IPAddress } | Select-Object -First 1).IPAddress }
$tcp = $null; $tcpMs = $null
if ($ip) {
$c = New-Object Net.Sockets.TcpClient; $sw = [Diagnostics.Stopwatch]::StartNew()
try { $tcp = $c.ConnectAsync($ip, $t.Port).Wait(3000) } catch { $tcp = $false }
$tcpMs = $sw.ElapsedMilliseconds; $c.Close() }
$sw = [Diagnostics.Stopwatch]::StartNew(); $err = $null
$j = Start-Job -ScriptBlock { param($p) $a = Test-Path -LiteralPath $p; $b = @(Get-ChildItem -LiteralPath $p -Force -ErrorAction Stop | Select-Object -First 1).Count; "$a|$b" } -ArgumentList $t.Path
if (Wait-Job $j -Timeout 5) { try { Receive-Job $j -ErrorAction Stop | Out-Null; $st = 'ok' } catch { $st = 'error'; $err = $_.Exception.Message } }
else { $st = 'timeout'; Stop-Job $j }
$probeMs = $sw.ElapsedMilliseconds; Remove-Job $j -Force
[pscustomobject]@{ RemotePath = $t.Path; Host = $h; IP = $ip; Port = $t.Port; Tcp = $tcp; TcpMs = $tcpMs; Probe = $st; ProbeMs = $probeMs; Error = $err } }
$json = $r | ConvertTo-Json -Depth 6 -Compress
$ms = New-Object IO.MemoryStream
$gz = New-Object IO.Compression.GZipStream($ms, [IO.Compression.CompressionMode]::Compress)
$bytes = [Text.Encoding]::UTF8.GetBytes($json); $gz.Write($bytes, 0, $bytes.Length); $gz.Close()
Write-Output ("NFSRECONW1:" + [Convert]::ToBase64String($ms.ToArray()))
'''
# Protocol classification of a client mount fstype
def classify_fstype(fstype):
f = (fstype or "").lower()
if f.startswith("nfs"):
return "NFS"
if f in ("cifs", "smb3", "smbfs"):
return "SMB"
if f == "lustre":
return "Lustre"
if f.startswith("fuse"):
return "FUSE"
return f.upper() or "UNKNOWN"
# TCP port that reaches the server for each protocol (None = no VPC endpoint to check)
PROTO_PORT = {"NFS": 2049, "SMB": 445, "Lustre": 988}
# Field names for /proc/self/mountstats (statvers 1.1)
NFS_EVENTS = [
"inoderevalidates", "dentryrevalidates", "datainvalidates", "attrinvalidates",
"vfsopen", "vfslookup", "vfsaccess", "vfsupdatepage", "vfsreadpage", "vfsreadpages",
"vfswritepage", "vfswritepages", "vfsgetdents", "vfssetattr", "vfsflush", "vfsfsync",
"vfslock", "vfsrelease", "congestionwait", "setattrtrunc", "extendwrite",
"sillyrenames", "shortreads", "shortwrites", "delay", "pnfsreads", "pnfswrites",
]
NFS_BYTES = [
"normalreadbytes", "normalwritebytes", "directreadbytes", "directwritebytes",
"serverreadbytes", "serverwritebytes", "readpages", "writepages",
]
XPRT_FIELDS = {
"tcp": ["port", "bind_count", "connect_count", "connect_time", "idle_time", "sends",
"recvs", "bad_xids", "req_u", "bklog_u", "max_slots", "sending_u", "pending_u"],
"udp": ["port", "bind_count", "sends", "recvs", "bad_xids", "req_u", "bklog_u",
"max_slots", "sending_u", "pending_u"],
}
OP_FIELDS = ["ops", "trans", "timeouts", "bytes_sent", "bytes_recv", "queue_ms", "rtt_ms",
"execute_ms", "errors"]
# --------------------------------------------------------------------------------------
# Helpers
# --------------------------------------------------------------------------------------
def jdefault(o):
if isinstance(o, (dt.datetime, dt.date)):
return o.isoformat()
if isinstance(o, bytes):
return o.decode("utf-8", "replace")
return str(o)
def write_json(path, data):
os.makedirs(os.path.dirname(path), exist_ok=True)
with open(path, "w") as f:
json.dump(data, f, indent=2, default=jdefault, sort_keys=True)
def log(msg):
with LOCK:
print(f"[{dt.datetime.now().strftime('%H:%M:%S')}] {msg}", file=sys.stderr, flush=True)
def tags_to_dict(tags):
out = {}
for t in tags or []:
k = t.get("Key", t.get("key"))
v = t.get("Value", t.get("value"))
if k is not None:
out[k] = v
return out
def unescape_mount(path):
return (path.replace("\\040", " ").replace("\\011", "\t")
.replace("\\012", "\n").replace("\\134", "\\"))
def parse_opts(s):
d = {}
for item in (s or "").split(","):
if not item:
continue
if "=" in item:
k, v = item.split("=", 1)
d[k] = v
else:
d[item] = True
return d
def to_int(x, default=None):
try:
return int(x)
except (TypeError, ValueError):
try:
return int(float(x))
except (TypeError, ValueError):
return default
class Recorder:
"""Records API errors without stopping the run."""
def __init__(self):
self.errors = []
def call(self, region, svc, op, fn, *args, default=None, quiet_codes=(), **kwargs):
try:
return fn(*args, **kwargs)
except ClientError as e:
code = e.response.get("Error", {}).get("Code", "ClientError")
if code not in quiet_codes:
with LOCK:
self.errors.append({"region": region, "service": svc, "operation": op,
"code": code, "message": str(e)[:500]})
return default
except (BotoCoreError, Exception) as e: # noqa: BLE001
with LOCK:
self.errors.append({"region": region, "service": svc, "operation": op,
"code": type(e).__name__, "message": str(e)[:500]})
return default
def paginate(self, region, svc, client, op, key, **kwargs):
items = []
try:
if client.can_paginate(op):
for page in client.get_paginator(op).paginate(**kwargs):
items.extend(page.get(key, []))
else:
items.extend(getattr(client, op)(**kwargs).get(key, []))
except ClientError as e:
with LOCK:
self.errors.append({"region": region, "service": svc, "operation": op,
"code": e.response.get("Error", {}).get("Code"),
"message": str(e)[:500]})
except Exception as e: # noqa: BLE001
with LOCK:
self.errors.append({"region": region, "service": svc, "operation": op,
"code": type(e).__name__, "message": str(e)[:500]})
return items
# --------------------------------------------------------------------------------------
# CloudWatch: pull every metric published for a resource
# --------------------------------------------------------------------------------------
def collect_metrics(rec, cw, region, namespace, dim_name, dim_value, days, max_metrics=300):
metrics = rec.paginate(region, "cloudwatch", cw, "list_metrics", "Metrics",
Namespace=namespace,
Dimensions=[{"Name": dim_name, "Value": dim_value}])
metrics = metrics[:max_metrics]
if not metrics:
return {}
queries, index = [], {}
for i, m in enumerate(metrics):
dims = {d["Name"]: d["Value"] for d in m.get("Dimensions", [])}
label = m["MetricName"] + "".join(f"[{k}={v}]" for k, v in sorted(dims.items())
if k != dim_name)
for stat in ("Sum", "Average", "Maximum", "Minimum"):
qid = f"m{i}_{stat.lower()}"
index[qid] = (label, stat)
queries.append({"Id": qid, "ReturnData": True, "MetricStat": {
"Metric": {"Namespace": namespace, "MetricName": m["MetricName"],
"Dimensions": m.get("Dimensions", [])},
"Period": 86400, "Stat": stat}})
out = defaultdict(lambda: {"daily": defaultdict(dict)})
start, end = NOW - dt.timedelta(days=days), NOW
for chunk in range(0, len(queries), 500):
token = None
while True:
kw = dict(MetricDataQueries=queries[chunk:chunk + 500], StartTime=start,
EndTime=end, ScanBy="TimestampAscending")
if token:
kw["NextToken"] = token
resp = rec.call(region, "cloudwatch", "get_metric_data", cw.get_metric_data, **kw)
if not resp:
break
for r in resp.get("MetricDataResults", []):
label, stat = index[r["Id"]]
for ts, val in zip(r.get("Timestamps", []), r.get("Values", [])):
day = ts.date().isoformat() if hasattr(ts, "date") else str(ts)[:10]
out[label]["daily"][day][stat] = val
token = resp.get("NextToken")
if not token:
break
result = {}
for label, v in out.items():
days_ = dict(sorted(v["daily"].items()))
sums = [d.get("Sum") for d in days_.values() if d.get("Sum") is not None]
maxs = [d.get("Maximum") for d in days_.values() if d.get("Maximum") is not None]
avgs = [d.get("Average") for d in days_.values() if d.get("Average") is not None]
result[label] = {
"window_days": days,
"sum": sum(sums) if sums else None,
"max": max(maxs) if maxs else None,
"avg_of_daily_avg": (sum(avgs) / len(avgs)) if avgs else None,
"days_with_data": len(days_),
"last_day": list(days_.keys())[-1] if days_ else None,
"daily": days_,
}
return result
# --------------------------------------------------------------------------------------
# Per region server side collection
# --------------------------------------------------------------------------------------
def collect_region(session, region, args, rec, out_dir, account="self"):
log(f"{account}/{region}: collecting")
c = lambda svc: session.client(svc, region_name=region, config=BOTO_CFG) # noqa: E731
R = {"region": region, "account": account, "efs": [], "s3files": [], "fsx": [], "storage_gateway": [], "fsx_file_caches": [],
"lambda_mounts": [], "ecs_efs_volumes": [], "eks_efs_csi": [],
"ec2_instances": [], "enis": [], "security_groups": [], "subnets": [], "vpcs": [],
"ssm_managed": []}
ec2 = c("ec2")
cw = c("cloudwatch") if not args.no_metrics else None
backup = c("backup")
# ---- EC2 foundation: instances, ENIs, SGs, subnets ----
for res in rec.paginate(region, "ec2", ec2, "describe_instances", "Reservations"):
R["ec2_instances"].extend(res.get("Instances", []))
R["enis"] = rec.paginate(region, "ec2", ec2, "describe_network_interfaces", "NetworkInterfaces")
R["security_groups"] = rec.paginate(region, "ec2", ec2, "describe_security_groups", "SecurityGroups")
R["subnets"] = rec.paginate(region, "ec2", ec2, "describe_subnets", "Subnets")
R["vpcs"] = rec.paginate(region, "ec2", ec2, "describe_vpcs", "Vpcs")
# Network paths, used to attribute mount targets that are not in this account
R["route_tables"] = rec.paginate(region, "ec2", ec2, "describe_route_tables", "RouteTables")
R["vpc_peerings"] = rec.paginate(region, "ec2", ec2, "describe_vpc_peering_connections", "VpcPeeringConnections")
R["tgw_attachments"] = rec.paginate(region, "ec2", ec2, "describe_transit_gateway_attachments",
"TransitGatewayAttachments")
R["flow_logs"] = rec.paginate(region, "ec2", ec2, "describe_flow_logs", "FlowLogs")
R["ssm_managed"] = rec.paginate(region, "ssm", c("ssm"), "describe_instance_information",
"InstanceInformationList")
def recovery_points(arn):
pts = rec.paginate(region, "backup", backup, "list_recovery_points_by_resource",
"RecoveryPoints", ResourceArn=arn)
pts.sort(key=lambda p: str(p.get("CreationDate")), reverse=True)
return {"count": len(pts), "latest": pts[0] if pts else None, "recent": pts[:5]}
# ---- EFS ----
efs = c("efs")
for fs in rec.paginate(region, "efs", efs, "describe_file_systems", "FileSystems"):
fid = fs["FileSystemId"]
item = {"file_system": fs}
mts = rec.paginate(region, "efs", efs, "describe_mount_targets", "MountTargets", FileSystemId=fid)
for mt in mts:
mt["SecurityGroups"] = (rec.call(region, "efs", "describe_mount_target_security_groups",
efs.describe_mount_target_security_groups,
MountTargetId=mt["MountTargetId"], default={})
or {}).get("SecurityGroups", [])
item["mount_targets"] = mts
item["access_points"] = rec.paginate(region, "efs", efs, "describe_access_points",
"AccessPoints", FileSystemId=fid)
item["file_system_policy"] = rec.call(region, "efs", "describe_file_system_policy",
efs.describe_file_system_policy, FileSystemId=fid,
quiet_codes=("PolicyNotFound",))
item["backup_policy"] = rec.call(region, "efs", "describe_backup_policy",
efs.describe_backup_policy, FileSystemId=fid,
quiet_codes=("PolicyNotFound",))
item["lifecycle"] = rec.call(region, "efs", "describe_lifecycle_configuration",
efs.describe_lifecycle_configuration, FileSystemId=fid)
item["replication"] = rec.call(region, "efs", "describe_replication_configurations",
efs.describe_replication_configurations, FileSystemId=fid,
quiet_codes=("ReplicationNotFound",))
item["tags"] = tags_to_dict(fs.get("Tags"))
item["recovery_points"] = recovery_points(fs.get("FileSystemArn"))
if cw:
item["metrics"] = collect_metrics(rec, cw, region, "AWS/EFS", "FileSystemId", fid, args.days)
R["efs"].append(item)
# ---- S3 Files (new service; skip cleanly if the SDK does not know it) ----
try:
s3f = c("s3files")
for fs in rec.paginate(region, "s3files", s3f, "list_file_systems", "fileSystems"):
fid = fs.get("fileSystemId")
item = {"file_system": fs}
item["detail"] = rec.call(region, "s3files", "get_file_system", s3f.get_file_system,
fileSystemId=fid)
item["mount_targets"] = rec.paginate(region, "s3files", s3f, "list_mount_targets",
"mountTargets", fileSystemId=fid)
if cw:
# S3 Files is built on EFS infrastructure; publish namespace may vary.
item["metrics"] = {}
for ns in ("AWS/S3Files", "AWS/EFS"):
m = collect_metrics(rec, cw, region, ns, "FileSystemId", fid, args.days)
if m:
item["metrics"] = m
item["metrics_namespace"] = ns
break
R["s3files"].append(item)
except UnknownServiceError:
R["s3files_note"] = ("installed boto3 does not know the s3files service; "
"upgrade boto3 to include Amazon S3 Files")
# ---- FSx ----
fsx = c("fsx")
for fs in rec.paginate(region, "fsx", fsx, "describe_file_systems", "FileSystems"):
fid = fs["FileSystemId"]
ftype = fs.get("FileSystemType")
item = {"file_system": fs, "type": ftype,
"nfs_capable": ftype in ("ONTAP", "OPENZFS"),
"protocols": {"ONTAP": ["NFS", "SMB", "iSCSI", "NVMe"], "OPENZFS": ["NFS"],
"WINDOWS": ["SMB"], "LUSTRE": ["Lustre"]}.get(ftype, [ftype]),
"tags": tags_to_dict(fs.get("Tags"))}
if ftype == "LUSTRE":
item["data_repository_associations"] = rec.paginate(
region, "fsx", fsx, "describe_data_repository_associations", "Associations",
Filters=[{"Name": "file-system-id", "Values": [fid]}])
if ftype == "ONTAP":
item["svms"] = rec.paginate(region, "fsx", fsx, "describe_storage_virtual_machines",
"StorageVirtualMachines",
Filters=[{"Name": "file-system-id", "Values": [fid]}])
if ftype in ("ONTAP", "OPENZFS"):
item["volumes"] = rec.paginate(region, "fsx", fsx, "describe_volumes", "Volumes",
Filters=[{"Name": "file-system-id", "Values": [fid]}])
item["backups"] = rec.paginate(region, "fsx", fsx, "describe_backups", "Backups",
Filters=[{"Name": "file-system-id", "Values": [fid]}])[:20]
item["recovery_points"] = recovery_points(fs.get("ResourceARN"))
if cw:
item["metrics"] = collect_metrics(rec, cw, region, "AWS/FSx", "FileSystemId", fid, args.days)
R["fsx"].append(item)
# ---- FSx File Cache (Lustre protocol cache in front of S3 / NFS) ----
R["fsx_file_caches"] = rec.paginate(region, "fsx", fsx, "describe_file_caches", "FileCaches")
# ---- Storage Gateway NFS and SMB shares ----
sgw = c("storagegateway")
for gw in rec.paginate(region, "storagegateway", sgw, "list_gateways", "Gateways"):
arn = gw["GatewayARN"]
item = {"gateway": gw,
"info": rec.call(region, "storagegateway", "describe_gateway_information",
sgw.describe_gateway_information, GatewayARN=arn)}
shares = rec.paginate(region, "storagegateway", sgw, "list_file_shares", "FileShareInfoList",
GatewayARN=arn)
nfs_arns = [s["FileShareARN"] for s in shares if s.get("FileShareType") == "NFS"]
item["nfs_file_shares"] = []
for i in range(0, len(nfs_arns), 10):
resp = rec.call(region, "storagegateway", "describe_nfs_file_shares",
sgw.describe_nfs_file_shares, FileShareARNList=nfs_arns[i:i + 10],
default={}) or {}
item["nfs_file_shares"].extend(resp.get("NFSFileShareInfoList", []))
smb_arns = [s["FileShareARN"] for s in shares if s.get("FileShareType") == "SMB"]
item["smb_file_shares"] = []
for i in range(0, len(smb_arns), 10):
resp = rec.call(region, "storagegateway", "describe_smb_file_shares",
sgw.describe_smb_file_shares, FileShareARNList=smb_arns[i:i + 10],
default={}) or {}
item["smb_file_shares"].extend(resp.get("SMBFileShareInfoList", []))
item["smb_settings"] = rec.call(region, "storagegateway", "describe_smb_settings",
sgw.describe_smb_settings, GatewayARN=arn)
item["other_shares"] = [s for s in shares if s.get("FileShareType") not in ("NFS", "SMB")]
if cw:
gid = arn.split("/")[-1]
item["metrics"] = collect_metrics(rec, cw, region, "AWS/StorageGateway", "GatewayId",
gid, args.days)
R["storage_gateway"].append(item)
# ---- Lambda functions with file systems ----
for fn in rec.paginate(region, "lambda", c("lambda"), "list_functions", "Functions"):
if fn.get("FileSystemConfigs"):
R["lambda_mounts"].append({
"FunctionName": fn["FunctionName"], "FunctionArn": fn["FunctionArn"],
"Runtime": fn.get("Runtime"), "LastModified": fn.get("LastModified"),
"VpcConfig": fn.get("VpcConfig"), "FileSystemConfigs": fn["FileSystemConfigs"]})
# ---- ECS task definitions with EFS volumes (latest revision per family) ----
ecs = c("ecs")
fams = rec.paginate(region, "ecs", ecs, "list_task_definition_families", "families", status="ACTIVE")
for fam in fams[:args.max_ecs_families]:
td = rec.call(region, "ecs", "describe_task_definition", ecs.describe_task_definition,
taskDefinition=fam)
if not td:
continue
tdd = td.get("taskDefinition", {})
vols = [v for v in tdd.get("volumes", []) if v.get("efsVolumeConfiguration")]
if vols:
R["ecs_efs_volumes"].append({
"taskDefinitionArn": tdd.get("taskDefinitionArn"), "family": fam,
"volumes": vols,
"mountPoints": [{"container": cd.get("name"), "mountPoints": cd.get("mountPoints")}
for cd in tdd.get("containerDefinitions", [])]})
if len(fams) > args.max_ecs_families:
R["ecs_note"] = f"{len(fams)} task definition families; only first {args.max_ecs_families} inspected"
# ---- EKS clusters with the EFS CSI driver ----
eks = c("eks")
for name in rec.paginate(region, "eks", eks, "list_clusters", "clusters"):
addon = rec.call(region, "eks", "describe_addon", eks.describe_addon, clusterName=name,
addonName="aws-efs-csi-driver", quiet_codes=("ResourceNotFoundException",))
R["eks_efs_csi"].append({"cluster": name,
"efs_csi_addon": (addon or {}).get("addon"),
"note": "PersistentVolumes are only visible through the Kubernetes API"})
write_json(os.path.join(out_dir, "raw", account, region, "server_side.json"), R)
log(f"{account}/{region}: EFS={len(R['efs'])} S3Files={len(R['s3files'])} FSx={len(R['fsx'])} "
f"SGW={len(R['storage_gateway'])} EC2={len(R['ec2_instances'])} SSM={len(R['ssm_managed'])}")
return R
# --------------------------------------------------------------------------------------
# Client side collection through SSM
# --------------------------------------------------------------------------------------
def ssm_collect(session, region, targets, args, rec, out_dir, platform="Linux", account="self"):
if not targets:
return {}
ssm = session.client("ssm", region_name=region, config=BOTO_CFG)
s3 = session.client("s3", region_name=region, config=BOTO_CFG) if args.ssm_bucket else None
results = {}
for i in range(0, len(targets), 50):
batch = targets[i:i + 50]
doc, script = (("AWS-RunPowerShellScript", WINDOWS_COLLECTOR) if platform == "Windows"
else ("AWS-RunShellScript", CLIENT_COLLECTOR))
kw = dict(InstanceIds=batch, DocumentName=doc,
Comment="file-mount-recon read only collector", TimeoutSeconds=60,
Parameters={"commands": [script], "executionTimeout": [str(args.ssm_timeout)]})
if args.ssm_bucket:
kw["OutputS3BucketName"] = args.ssm_bucket
kw["OutputS3KeyPrefix"] = f"file-mount-recon/{NOW.strftime('%Y%m%dT%H%M%S')}"
resp = rec.call(region, "ssm", "send_command", ssm.send_command, **kw)
if not resp:
for iid in batch:
results[iid] = {"status": "send_failed"}
continue
cmd_id = resp["Command"]["CommandId"]
log(f"{region}: SSM command {cmd_id} sent to {len(batch)} instances")
pending, deadline = set(batch), time.time() + args.ssm_timeout + 90
while pending and time.time() < deadline:
time.sleep(5)
for iid in list(pending):
inv = rec.call(region, "ssm", "get_command_invocation", ssm.get_command_invocation,
CommandId=cmd_id, InstanceId=iid,
quiet_codes=("InvocationDoesNotExist",))
if not inv or inv.get("Status") in ("Pending", "InProgress", "Delayed"):
continue
pending.discard(iid)
stdout = inv.get("StandardOutputContent", "") or ""
source = "inline"
if s3 and inv.get("StandardOutputUrl"):
full = fetch_ssm_s3_output(rec, region, s3, args.ssm_bucket, kw["OutputS3KeyPrefix"],
cmd_id, iid)
if full:
stdout, source = full, "s3"
results[iid] = {"status": inv.get("Status"), "response_code": inv.get("ResponseCode"),
"platform": platform,
"stderr": (inv.get("StandardErrorContent") or "")[:2000],
"output_source": source, "command_id": cmd_id,
"bundle": decode_bundle(stdout)}
for iid in pending:
results[iid] = {"status": "timed_out_waiting", "command_id": cmd_id}
for iid, r in results.items():
b = r.get("bundle") or {}
if b.get("text"):
p = os.path.join(out_dir, "raw", account, region, "clients",
f"{iid}.{'json' if r.get('platform') == 'Windows' else 'txt'}")
os.makedirs(os.path.dirname(p), exist_ok=True)
with open(p, "w") as f:
f.write(b["text"])
return results
def fetch_ssm_s3_output(rec, region, s3, bucket, prefix, cmd_id, iid):
keys = rec.paginate(region, "s3", s3, "list_objects_v2", "Contents", Bucket=bucket,
Prefix=f"{prefix}/{cmd_id}/{iid}/")
for k in keys:
if k["Key"].endswith("/stdout"):
obj = rec.call(region, "s3", "get_object", s3.get_object, Bucket=bucket, Key=k["Key"])
if obj:
return obj["Body"].read().decode("utf-8", "replace")
return None
def decode_bundle(stdout):
m = re.search(r"NFSRECONW?1:([A-Za-z0-9+/=]+)", stdout or "")
if not m:
return {"ok": False, "error": "no bundle marker in output", "raw_head": (stdout or "")[:1000]}
data = m.group(1)
try:
text = gzip.decompress(base64.b64decode(data + "=" * (-len(data) % 4))).decode("utf-8", "replace")
return {"ok": True, "text": text}
except Exception as e: # noqa: BLE001
return {"ok": False, "truncated_likely": len(stdout) >= 23900,
"error": f"{type(e).__name__}: {e}. Output is probably truncated at the SSM inline "
f"limit (24,000 chars); rerun with --ssm-bucket to get full output."}
# --------------------------------------------------------------------------------------
# Parsing of the client bundle
# --------------------------------------------------------------------------------------
def split_sections(text):
sections, cur = {}, None
for line in text.splitlines():
m = re.match(r"^===SECTION (\S+)===$", line)
if m:
cur = m.group(1)
sections[cur] = []
elif cur:
sections[cur].append(line)
return sections
def parse_mountstats(lines):
mounts, cur, in_ops = [], None, False
for raw in lines:
line = raw.strip()
m = re.match(r"^device (\S+) mounted on (.+) with fstype (\S+)(?: statvers=(\S+))?", line)
if m:
cur = {"device": m.group(1), "mountpoint": unescape_mount(m.group(2)),
"fstype": m.group(3), "statvers": m.group(4), "per_op": {}}
mounts.append(cur)
in_ops = False
continue
if cur is None or not line:
continue
if line.startswith("per-op statistics"):
in_ops = True
continue
key, _, val = line.partition(":")
val = val.strip()
if in_ops and re.match(r"^[A-Z_0-9]+$", key):
nums = [to_int(x, 0) for x in val.split()]
op = dict(zip(OP_FIELDS, nums))
if op.get("ops"):
op["avg_rtt_ms"] = round(op.get("rtt_ms", 0) / op["ops"], 3)
op["avg_exec_ms"] = round(op.get("execute_ms", 0) / op["ops"], 3)
op["avg_queue_ms"] = round(op.get("queue_ms", 0) / op["ops"], 3)
op["retransmissions"] = max(0, op.get("trans", 0) - op["ops"])
cur["per_op"][key] = op
continue
if key == "opts":
cur["opts"] = parse_opts(val)
elif key == "age":
cur["age_s"] = to_int(val)
elif key == "caps":
cur["caps"] = parse_opts(val)
elif key == "nfsv4":
cur["nfsv4"] = parse_opts(val)
elif key == "sec":
cur["sec"] = parse_opts(val)
elif key == "impl_id":
cur["impl_id"] = val
elif key == "fsc":
cur["fsc"] = val
elif key == "events":
nums = [to_int(x, 0) for x in val.split()]
cur["events"] = dict(zip(NFS_EVENTS, nums))
if len(nums) != len(NFS_EVENTS):
cur["events_raw"] = nums
elif key == "bytes":
cur["bytes"] = dict(zip(NFS_BYTES, [to_int(x, 0) for x in val.split()]))
elif key == "xprt":
parts = val.split()
proto = parts[0] if parts else "?"
nums = [to_int(x, 0) for x in parts[1:]]
fields = XPRT_FIELDS.get(proto)
cur["xprt"] = {"proto": proto, **(dict(zip(fields, nums)) if fields else {"values": nums})}
elif key.startswith("RPC iostats version"):
cur["rpc_iostats"] = line
return mounts
def parse_fuse_procs(lines):
"""Map FUSE mountpoints to the daemon and remote (bucket, remote:path, user@host:path)."""
out = {}
for l in lines:
p = l.split(None, 3)
if len(p) < 4:
continue
pid, etimes, rss, args = p
argv = args.split()
prog = os.path.basename(argv[0]) if argv else ""
pos = [a for a in argv[1:] if not a.startswith("-")]
# mountpoint is normally the last positional argument that is an absolute path
mp = next((a for a in reversed(pos) if a.startswith("/")), None)
if not mp:
continue
remote = None
if mp in pos:
i = pos.index(mp)
remote = pos[i - 1] if i > 0 else None
if prog in ("rclone",) and "mount" in pos:
pos2 = pos[pos.index("mount") + 1:]
remote = pos2[0] if pos2 else remote
bucket = None
if remote and prog in ("mount-s3", "mountpoint-s3", "s3fs", "goofys", "geesefs", "gcsfuse"):
bucket = remote.split(":", 1)[0]
prefix = None
for k in ("--prefix",):
if k in argv and argv.index(k) + 1 < len(argv):
prefix = argv[argv.index(k) + 1]
out[mp] = {"pid": to_int(pid), "elapsed_s": to_int(etimes), "rss_kb": to_int(rss),
"program": prog, "remote": remote, "bucket": bucket, "prefix": prefix,
"read_only": "--read-only" in argv or "ro" in args.split("-o")[-1].split(","),
"args": args}
return out
def parse_probes(lines):
out = {}
for l in lines:
p = l.split("|")
if len(p) < 11:
continue
mp, fstype, host, ip, port, tcp, tcp_ms, st_rc, st_ms, rd_rc, rd_ms = p[:11]
root = (p[11] if len(p) > 11 else "").split(":")
out[mp] = {"host": host or None, "ip": ip or None, "port": to_int(port), "tcp": tcp or None,
"tcp_ms": to_int(tcp_ms), "stat_rc": to_int(st_rc), "stat_ms": to_int(st_ms),
"readdir_rc": to_int(rd_rc), "readdir_ms": to_int(rd_ms),
"root_mode": root[0] if len(root) == 3 else None,
"root_uid": to_int(root[1]) if len(root) == 3 else None,
"root_gid": to_int(root[2]) if len(root) == 3 else None,
"readdir_error": (p[12] if len(p) > 12 else "") or None}
return out
def parse_bundle(text):
S = split_sections(text)
meta = {}
for l in S.get("meta", []):
if "=" in l:
k, v = l.split("=", 1)
meta[k] = v
mounts = []
for l in S.get("proc_mounts", []):
p = l.split()
if len(p) >= 4:
mounts.append({"source": p[0], "mountpoint": unescape_mount(p[1]), "fstype": p[2],
"options_raw": p[3], "options": parse_opts(p[3])})
fstab = []
for l in S.get("fstab", []):
p = l.split()
if len(p) >= 3:
fstab.append({"device": p[0], "mountpoint": unescape_mount(p[1]), "fstype": p[2],
"options": p[3] if len(p) > 3 else "defaults", "line": l})
df = {}
for l in S.get("df", []):
p = l.split("|", 3)
if len(p) == 4:
mp, rc, ms, out = p
entry = {"exit_code": to_int(rc), "latency_ms": to_int(ms), "raw": out}
f = out.split()
if entry["exit_code"] == 0 and len(f) >= 7:
entry.update({"size_bytes": to_int(f[2]), "used_bytes": to_int(f[3]),
"avail_bytes": to_int(f[4]), "use_pct": f[5]})
entry["responsive"] = entry["exit_code"] == 0
entry["hung"] = entry["exit_code"] in (124, 137)
df[mp] = entry
statfs = {}
for l in S.get("statfs", []):
mp, _, rest = l.partition("|")
statfs[mp] = parse_opts(rest.replace(" ", ","))
openf = defaultdict(lambda: {"open_handles": 0, "processes": []})
for l in S.get("open_files", []):
p = l.split("|")
if len(p) >= 5:
mp, pid, cnt, comm, user = p[:5]
openf[mp]["open_handles"] += to_int(cnt, 0)
openf[mp]["processes"].append({"pid": to_int(pid), "comm": comm, "user": user,
"handles": to_int(cnt, 0)})
tls_state = {}
cur = None
for l in S.get("tls_state", []):
if l.startswith("--- "):
cur = l[4:]
tls_state[cur] = ""
elif cur:
tls_state[cur] += l + "\n"
server = {"lines": S.get("nfs_server", [])}
server["is_nfs_server"] = any(l.startswith("nfsd_threads=") and l.split("=", 1)[1].strip() not in ("", "0")
for l in server["lines"]) or any(
l.strip() == "nfs_server_service=active" for l in server["lines"])
return {
"meta": meta, "os_release": S.get("os_release", []), "packages": S.get("packages", []),
"proc_mounts": mounts, "fstab": fstab, "mountstats": parse_mountstats(S.get("mountstats", [])),
"nfsstat_m": S.get("nfsstat_m", []), "proc_net_rpc_nfs": S.get("proc_net_rpc_nfs", []),
"df": df, "statfs": statfs, "open_files": dict(openf),
"systemd_units": S.get("systemd_units", []), "autofs": S.get("autofs", []),
"efs_utils_conf": S.get("efs_utils_conf", []), "tls_state": tls_state,
"tls_proxies": S.get("tls_proxies", []), "connections": S.get("connections", []),
"rpc_tunables": S.get("rpc_tunables", []), "nfs_server": server,
"fuse_procs": parse_fuse_procs(S.get("fuse_procs", [])), "cifs": S.get("cifs", []),
"probes": parse_probes(S.get("probes", [])),
"credfiles": {l.split("|", 1)[0]: l.split("|", 1)[1] for l in S.get("credfiles", []) if "|" in l},
"lustre": S.get("lustre", []), "other_remote": S.get("other_remote", []),
"complete": "end" in S,
}
# --------------------------------------------------------------------------------------
# Correlation and findings
# --------------------------------------------------------------------------------------
class Index:
def __init__(self):
self.by_id = {} # resource id -> server record
self.by_ip = defaultdict(list) # ip -> [(server_id, region, vpc_id, extra)]
self.by_dns = {} # dns name (lower) -> server id
self.eni_by_ip = defaultdict(list)
self.sg = {}
self.inst = {}
self.inst_by_ip = defaultdict(list)
self.subnet_az = {}
self.buckets = {} # bucket name -> {"region": ..., "account": ...}
self.accounts = set() # accounts covered by this run
self.vpcs = {} # vpc id -> {account, region, owner, cidrs}
self.subnets = {} # subnet id -> {vpc, owner, cidr, az, viewer}
self.route_tables = defaultdict(list) # (viewer account, vpc id) -> route tables
self.peerings = {} # pcx id -> peering connection
self.tgw_attachments = []
self._acct = None
def add_ip(self, ip, sid, region, vpc, extra=None):
if ip:
self.by_ip[ip].append({"server_id": sid, "region": region, "vpc_id": vpc,
"account": self._acct, **(extra or {})})
def vpc_for_ip(self, ip, exclude_account=None):
"""VPC (from any account in this run) whose CIDR contains ip."""
try:
a = ipaddress.ip_address(ip)
except ValueError:
return None
hits = []
for vid, v in self.vpcs.items():
for c in v["cidrs"]:
try:
n = ipaddress.ip_network(c, strict=False)
except ValueError:
continue
if a in n:
hits.append((n.prefixlen, vid, v))
hits.sort(key=lambda h: (-h[0], h[2]["owner"] == exclude_account))
return {"vpc_id": hits[0][1], **hits[0][2]} if hits else None
def network_path(self, viewer, vpc_id, subnet_id, ip):
"""How traffic from (viewer account, vpc, subnet) reaches ip, and whose network ip is in."""
if not ip:
return {"path": "unresolved"}
try:
a = ipaddress.ip_address(ip)
except ValueError:
return {"path": "unresolved"}
if a.is_loopback:
return {"path": "loopback"}
v = self.vpcs.get(vpc_id)
if v and any(a in ipaddress.ip_network(c, strict=False) for c in v["cidrs"]):
enis = self.eni_by_ip.get(ip, [])
sn = next((x for x in self.subnets.values() if x["vpc"] == vpc_id and x["cidr"]
and a in ipaddress.ip_network(x["cidr"], strict=False)), None)
eni_owner = (enis[0].get("OwnerId") if enis else None)
if eni_owner and eni_owner != viewer:
return {"path": "same_vpc", "peer_account": eni_owner, "detail": "ENI owned by another account"}
if enis:
return {"path": "same_vpc", "peer_account": eni_owner or viewer}
if v["owner"] != viewer or (sn and sn["owner"] != viewer):
return {"path": "shared_vpc", "peer_account": None,
"detail": f"VPC owned by {v['owner']}; target ENI not visible to {viewer}, "
"so it belongs to another participant or the owner"}
return {"path": "same_vpc", "peer_account": None, "detail": "no visible ENI holds this IP"}
# Longest prefix match in the subnet's route table (explicit association, else main)
rts = self.route_tables.get((viewer, vpc_id), [])
rt = next((r for r in rts for asc in r.get("Associations", []) if asc.get("SubnetId") == subnet_id), None) \
or next((r for r in rts for asc in r.get("Associations", []) if asc.get("Main")), None)
best = None
for route in (rt or {}).get("Routes", []):
d = route.get("DestinationCidrBlock")
if not d or route.get("State") == "blackhole":
continue
n = ipaddress.ip_network(d, strict=False)
if a in n and (best is None or n.prefixlen > best[0]):
best = (n.prefixlen, route)
owner_vpc = self.vpc_for_ip(ip, exclude_account=viewer)
res = {"path": "no_route" if rt else "unknown", "route_table": (rt or {}).get("RouteTableId")}
if best:
r = best[1]
res["route"] = r.get("DestinationCidrBlock")
if r.get("VpcPeeringConnectionId"):
pcx = self.peerings.get(r["VpcPeeringConnectionId"], {})
req, acc = pcx.get("RequesterVpcInfo", {}), pcx.get("AccepterVpcInfo", {})
peer = acc if req.get("VpcId") == vpc_id else req
res.update({"path": "vpc_peering", "via": r["VpcPeeringConnectionId"],
"peer_account": peer.get("OwnerId"), "peer_vpc": peer.get("VpcId"),
"peer_region": peer.get("Region")})
elif r.get("TransitGatewayId"):
res.update({"path": "transit_gateway", "via": r["TransitGatewayId"]})
elif r.get("CoreNetworkArn"):
res.update({"path": "cloud_wan", "via": r["CoreNetworkArn"]})
elif str(r.get("GatewayId", "")).startswith("vgw-"):
res.update({"path": "vpn_or_direct_connect", "via": r["GatewayId"]})
elif str(r.get("GatewayId", "")).startswith("igw-") or r.get("NatGatewayId"):
res.update({"path": "internet", "via": r.get("GatewayId") or r.get("NatGatewayId")})
elif r.get("LocalGatewayId"):
res.update({"path": "outposts_local_gateway", "via": r["LocalGatewayId"]})
else:
res.update({"path": "other", "via": json.dumps({k: v for k, v in r.items() if k.endswith("Id")})})
if owner_vpc and not res.get("peer_account"):
res.update({"peer_account": owner_vpc["owner"], "peer_vpc": owner_vpc["vpc_id"],
"peer_region": owner_vpc["region"]})
if res.get("path") == "transit_gateway" and not res.get("peer_account"):
res["detail"] = "behind a Transit Gateway; run with --org or --accounts to attribute the account"
return res
def build_servers(regions_data, idx):
servers = []
for R in regions_data:
region = R["region"]
acct = R.get("account", "self")
idx._acct = acct
idx.accounts.add(acct)
_start = len(servers)
for v in R.get("vpcs", []):
cidrs = [c["CidrBlock"] for c in v.get("CidrBlockAssociationSet", [])
if (c.get("CidrBlockState") or {}).get("State", "associated") == "associated"] or [v.get("CidrBlock")]
idx.vpcs[v["VpcId"]] = {"account": acct, "owner": v.get("OwnerId", acct), "region": region,
"cidrs": [c for c in cidrs if c]}
for sn in R.get("subnets", []):
idx.subnets[sn["SubnetId"]] = {"vpc": sn.get("VpcId"), "owner": sn.get("OwnerId", acct),
"cidr": sn.get("CidrBlock"), "az": sn.get("AvailabilityZone"),
"viewer": acct}
for rt in R.get("route_tables", []):
idx.route_tables[(acct, rt.get("VpcId"))].append(rt)
for pcx in R.get("vpc_peerings", []):
idx.peerings[pcx["VpcPeeringConnectionId"]] = pcx
idx.tgw_attachments.extend(R.get("tgw_attachments", []))
for sg in R["security_groups"]:
idx.sg[sg["GroupId"]] = sg
for sn in R["subnets"]:
idx.subnet_az[sn["SubnetId"]] = sn.get("AvailabilityZone")
for eni in R["enis"]:
for a in eni.get("PrivateIpAddresses", []):
idx.eni_by_ip[a["PrivateIpAddress"]].append({**eni, "_region": region})
for i in R["ec2_instances"]:
idx.inst[i["InstanceId"]] = {**i, "_region": region, "_account": acct}
for ni in i.get("NetworkInterfaces", []):
for a in ni.get("PrivateIpAddresses", []):
idx.inst_by_ip[a["PrivateIpAddress"]].append(i["InstanceId"])
for e in R["efs"]:
fs = e["file_system"]
fid = fs["FileSystemId"]
mts = e["mount_targets"]
m = e.get("metrics", {})
s = {
"server_type": "EFS", "id": fid, "region": region, "name": fs.get("Name") or e["tags"].get("Name"),
"state": fs.get("LifeCycleState"), "created": fs.get("CreationTime"),
"size_bytes": (fs.get("SizeInBytes") or {}).get("Value"),
"size_bytes_standard": (fs.get("SizeInBytes") or {}).get("ValueInStandard"),
"size_bytes_ia": (fs.get("SizeInBytes") or {}).get("ValueInIA"),
"size_bytes_archive": (fs.get("SizeInBytes") or {}).get("ValueInArchive"),
"performance_mode": fs.get("PerformanceMode"), "throughput_mode": fs.get("ThroughputMode"),
"provisioned_mibps": fs.get("ProvisionedThroughputInMibps"),
"encrypted": fs.get("Encrypted"), "kms_key": fs.get("KmsKeyId"),
"availability_zone": fs.get("AvailabilityZoneName") or "regional",
"mount_target_count": len(mts),
"mount_targets_available": sum(1 for x in mts if x.get("LifeCycleState") == "available"),
"mount_target_azs": sorted({x.get("AvailabilityZoneName") for x in mts if x.get("AvailabilityZoneName")}),
"vpc_ids": sorted({x.get("VpcId") for x in mts if x.get("VpcId")}),
"access_points": len(e["access_points"]),
"has_fs_policy": bool(e.get("file_system_policy")),
"fs_policy": (e.get("file_system_policy") or {}).get("Policy"),
"backup_policy": ((e.get("backup_policy") or {}).get("BackupPolicy") or {}).get("Status"),
"lifecycle": json.dumps((e.get("lifecycle") or {}).get("LifecyclePolicies", []), default=str),
"replication": bool((e.get("replication") or {}).get("Replications")),
"protection": json.dumps(fs.get("FileSystemProtection"), default=str),
"recovery_points": e["recovery_points"]["count"],
"last_recovery_point": (e["recovery_points"]["latest"] or {}).get("CreationDate"),
"metric_client_connections_max": (m.get("ClientConnections") or {}).get("max"),
"metric_total_io_bytes_sum": (m.get("TotalIOBytes") or {}).get("sum"),
"metric_data_read_bytes_sum": (m.get("DataReadIOBytes") or {}).get("sum"),
"metric_data_write_bytes_sum": (m.get("DataWriteIOBytes") or {}).get("sum"),
"metric_metadata_io_bytes_sum": (m.get("MetadataIOBytes") or {}).get("sum"),
"metric_percent_io_limit_max": (m.get("PercentIOLimit") or {}).get("max"),
"metric_burst_credit_min": min((d.get("Minimum") for d in (m.get("BurstCreditBalance") or {}).get("daily", {}).values()
if d.get("Minimum") is not None), default=None),
"metric_burst_credit_max": (m.get("BurstCreditBalance") or {}).get("max"),
"metric_storage_bytes_max": (m.get("StorageBytes[StorageClass=Total]") or {}).get("max"),
"metrics_collected": len(m), "_metrics": m,
"replication_destinations": [d.get("Region") for r in ((e.get("replication") or {}).get("Replications") or [])
for d in r.get("Destinations", []) if d.get("Region") != region],
"tags": e["tags"], "endpoints": [], "clients": [], "findings": [],
}
for mt in mts:
ep = {"kind": "mount_target", "id": mt["MountTargetId"], "ip": mt.get("IpAddress"),
"az": mt.get("AvailabilityZoneName"), "subnet": mt.get("SubnetId"),
"vpc_id": mt.get("VpcId"), "state": mt.get("LifeCycleState"),
"eni": mt.get("NetworkInterfaceId"), "security_groups": mt.get("SecurityGroups", [])}
s["endpoints"].append(ep)
idx.add_ip(mt.get("IpAddress"), fid, region, mt.get("VpcId"), {"endpoint": ep["id"]})
idx.by_dns[f"{fid}.efs.{region}.amazonaws.com"] = fid
idx.by_id[fid] = s
servers.append(s)
for e in R["s3files"]:
fs = e["file_system"]
det = (e.get("detail") or {})
fid = fs.get("fileSystemId")
mts = e.get("mount_targets", [])
s = {"server_type": "S3Files", "id": fid, "region": region, "name": fs.get("name"),
"state": (fs.get("status") or det.get("status") or "").lower() or None,
"status_message": fs.get("statusMessage") or det.get("statusMessage"),
"bucket": fs.get("bucket") or det.get("bucket"), "prefix": det.get("prefix"),
"kms_key": det.get("kmsKeyId"),
"role_arn": det.get("roleArn"), "created": fs.get("creationTime") or det.get("creationTime"),
"mount_target_count": len(mts),
"mount_targets_available": sum(1 for x in mts if str(x.get("status", "")).lower() == "available"),
"metrics_collected": len(e.get("metrics", {})),
"endpoints": [], "clients": [], "findings": []}
for mt in mts:
ip = mt.get("ipv4Address")
eni = (idx.eni_by_ip.get(ip) or [{}])[0]
ep = {"kind": "mount_target", "id": mt.get("mountTargetId"), "ip": ip,
"subnet": mt.get("subnetId"), "az": idx.subnet_az.get(mt.get("subnetId")),
"az_id": mt.get("availabilityZoneId"),
"vpc_id": mt.get("vpcId") or eni.get("VpcId"),
"state": str(mt.get("status", "")).lower() or None,
"eni": mt.get("networkInterfaceId") or eni.get("NetworkInterfaceId"),
"security_groups": [g["GroupId"] for g in eni.get("Groups", [])] or mt.get("securityGroups", [])}
s["endpoints"].append(ep)
idx.add_ip(ip, fid, region, ep["vpc_id"], {"endpoint": ep["id"]})
if fid:
idx.by_id[fid] = s
servers.append(s)
for e in R["fsx"]:
fs = e["file_system"]
fid = fs["FileSystemId"]
m = e.get("metrics", {})
wc = fs.get("WindowsConfiguration") or {}
lc = fs.get("LustreConfiguration") or {}
s = {"server_type": f"FSx-{e['type']}", "id": fid, "region": region,
"name": e["tags"].get("Name"), "state": fs.get("Lifecycle"),
"protocols": e["protocols"], "nfs_capable": e["nfs_capable"], "created": fs.get("CreationTime"),
"windows_ad": (wc.get("SelfManagedActiveDirectoryConfiguration") or {}).get("DomainName")
or wc.get("ActiveDirectoryId"),
"windows_throughput_mbps": wc.get("ThroughputCapacity"),
"windows_aliases": [a.get("Name") for a in wc.get("Aliases", [])],
"windows_preferred_file_server_ip": wc.get("PreferredFileServerIp"),
"lustre_mount_name": lc.get("MountName"), "lustre_deployment": lc.get("DeploymentType"),
"lustre_per_unit_throughput": lc.get("PerUnitStorageThroughput"),
"lustre_data_repository": json.dumps(lc.get("DataRepositoryConfiguration"), default=str)
if lc.get("DataRepositoryConfiguration") else None,
"lustre_dra_count": len(e.get("data_repository_associations", [])),
"storage_capacity_gib": fs.get("StorageCapacity"), "storage_type": fs.get("StorageType"),
"deployment": json.dumps({k: v.get("DeploymentType") for k, v in fs.items()
if k.endswith("Configuration") and isinstance(v, dict)}),
"kms_key": fs.get("KmsKeyId"), "vpc_ids": [fs.get("VpcId")], "dns_name": fs.get("DNSName"),
"volumes": len(e.get("volumes", [])), "svms": len(e.get("svms", [])),
"backups": len(e.get("backups", [])), "recovery_points": e["recovery_points"]["count"],
"metrics_collected": len(m), "_metrics": m, "tags": e["tags"], "endpoints": [], "clients": [],
"findings": [], "nfs_exports": []}
for v in e.get("volumes", []):
oc = v.get("OntapConfiguration") or {}
zc = v.get("OpenZFSConfiguration") or {}
s["nfs_exports"].append({
"volume_id": v.get("VolumeId"), "name": v.get("Name"), "lifecycle": v.get("Lifecycle"),
"junction_path": oc.get("JunctionPath"), "size_mb": oc.get("SizeInMegabytes"),
"tiering": (oc.get("TieringPolicy") or {}).get("Name"),
"security_style": oc.get("SecurityStyle"), "snapshot_policy": oc.get("SnapshotPolicy"),
"zfs_path": zc.get("VolumePath"), "zfs_quota_gib": zc.get("StorageCapacityQuotaGiB"),
"zfs_nfs_exports": zc.get("NfsExports"), "zfs_compression": zc.get("DataCompressionType")})
for eni_id in fs.get("NetworkInterfaceIds", []):
for R2 in [R]:
for eni in R2["enis"]:
if eni["NetworkInterfaceId"] == eni_id:
ip = eni.get("PrivateIpAddress")
ep = {"kind": "fsx_eni", "id": eni_id, "ip": ip, "az": eni.get("AvailabilityZone"),
"subnet": eni.get("SubnetId"), "vpc_id": eni.get("VpcId"),
"state": eni.get("Status"),
"security_groups": [g["GroupId"] for g in eni.get("Groups", [])]}
s["endpoints"].append(ep)
idx.add_ip(ip, fid, region, eni.get("VpcId"), {"endpoint": eni_id})
for svm in e.get("svms", []):
nfs_ep = ((svm.get("Endpoints") or {}).get("Nfs") or {})
for ip in nfs_ep.get("IpAddresses", []):
idx.add_ip(ip, fid, region, fs.get("VpcId"), {"endpoint": svm.get("StorageVirtualMachineId")})
s["endpoints"].append({"kind": "svm_nfs", "id": svm.get("StorageVirtualMachineId"),
"ip": ip, "dns": nfs_ep.get("DNSName"), "state": svm.get("Lifecycle"),
"vpc_id": fs.get("VpcId"), "security_groups": []})
if nfs_ep.get("DNSName"):
idx.by_dns[nfs_ep["DNSName"].lower()] = fid
if fs.get("DNSName"):
idx.by_dns[fs["DNSName"].lower()] = fid
for a in wc.get("Aliases", []):
if a.get("Name"):
idx.by_dns[a["Name"].lower()] = fid
if wc.get("RemoteAdministrationEndpoint"):
idx.by_dns[wc["RemoteAdministrationEndpoint"].lower()] = fid
idx.by_id[fid] = s
servers.append(s)
for e in R["storage_gateway"]:
info = e.get("info") or {}
gid = e["gateway"]["GatewayARN"].split("/")[-1]
ec2_id = info.get("Ec2InstanceId")
s = {"server_type": "StorageGateway", "id": gid, "region": region,
"name": info.get("GatewayName") or e["gateway"].get("GatewayName"),
"state": info.get("GatewayState"), "gateway_type": info.get("GatewayType"),
"software_version": info.get("SoftwareVersion"), "ec2_instance": ec2_id,
"host_environment": info.get("HostEnvironment"),
"nfs_shares": len(e["nfs_file_shares"]), "metrics_collected": len(e.get("metrics", {})),
"_metrics": e.get("metrics", {}),
"endpoints": [], "clients": [], "findings": [], "nfs_exports": []}
for sh in e["nfs_file_shares"]:
s["nfs_exports"].append({k: sh.get(k) for k in (
"FileShareId", "FileShareStatus", "Path", "LocationARN", "ClientList", "Squash",
"ReadOnly", "DefaultStorageClass", "ObjectACL", "GuessMIMETypeEnabled",
"RequesterPays", "CacheAttributes", "NotificationPolicy", "KMSEncrypted",
"AuditDestinationARN", "FileShareName")})
for sh in e.get("smb_file_shares", []):
s["nfs_exports"].append({"protocol": "SMB", **{k: sh.get(k) for k in (
"FileShareId", "FileShareStatus", "Path", "LocationARN", "FileShareName",
"Authentication", "ValidUserList", "AdminUserList", "AccessBasedEnumeration",
"SMBACLEnabled", "OplocksEnabled", "CaseSensitivity", "ReadOnly",
"DefaultStorageClass", "KMSEncrypted", "CacheAttributes", "AuditDestinationARN")}})
s["smb_shares"] = len(e.get("smb_file_shares", []))
s["smb_settings"] = json.dumps({k: (e.get("smb_settings") or {}).get(k) for k in (
"DomainName", "ActiveDirectoryStatus", "SMBSecurityStrategy", "SMBGuestPasswordSet",
"FileSharesVisible")}, default=str)
sg_ids = []
if ec2_id and ec2_id in idx.inst:
sg_ids = [g["GroupId"] for g in idx.inst[ec2_id].get("SecurityGroups", [])]
for ni in info.get("GatewayNetworkInterfaces", []):
ip = ni.get("Ipv4Address")
s["endpoints"].append({"kind": "gateway_interface", "id": gid, "ip": ip,
"security_groups": sg_ids, "state": info.get("GatewayState")})
idx.add_ip(ip, gid, region, None, {"endpoint": gid})
idx.by_id[gid] = s
servers.append(s)
for fc in R.get("fsx_file_caches", []):
cid = fc.get("FileCacheId")
lc = fc.get("LustreConfiguration") or {}
s = {"server_type": "FSx-FileCache", "id": cid, "region": region,
"name": tags_to_dict(fc.get("Tags")).get("Name"), "state": fc.get("Lifecycle"),
"protocols": ["Lustre"], "created": fc.get("CreationTime"),
"storage_capacity_gib": fc.get("StorageCapacity"), "dns_name": fc.get("DNSName"),
"lustre_mount_name": lc.get("MountName"), "vpc_ids": [fc.get("VpcId")],
"data_repository_association_ids": fc.get("DataRepositoryAssociationIds"),
"endpoints": [], "clients": [], "findings": []}
for eni_id in fc.get("NetworkInterfaceIds", []):
for eni in R["enis"]:
if eni["NetworkInterfaceId"] == eni_id:
ip = eni.get("PrivateIpAddress")
s["endpoints"].append({"kind": "cache_eni", "id": eni_id, "ip": ip,
"az": eni.get("AvailabilityZone"), "vpc_id": eni.get("VpcId"),
"state": eni.get("Status"),
"security_groups": [g["GroupId"] for g in eni.get("Groups", [])]})
idx.add_ip(ip, cid, region, eni.get("VpcId"), {"endpoint": eni_id})
if fc.get("DNSName"):
idx.by_dns[fc["DNSName"].lower()] = cid
idx.by_id[cid] = s
servers.append(s)
for s in servers[_start:]:
s["account"] = acct
return servers
def sg_allows(idx, sg_ids, client_ip, client_sgs, port=2049):
"""Return (allowed, matching_rules) for TCP port from client to the endpoint SGs."""
matches = []
for gid in sg_ids or []:
sg = idx.sg.get(gid)
if not sg:
continue
for p in sg.get("IpPermissions", []):
proto = str(p.get("IpProtocol"))
if proto not in ("tcp", "6", "-1"):
continue
if proto != "-1":
lo, hi = p.get("FromPort"), p.get("ToPort")
if lo is None or hi is None or not (lo <= port <= hi):
continue
for r in p.get("IpRanges", []):
try:
if client_ip and ipaddress.ip_address(client_ip) in ipaddress.ip_network(r["CidrIp"], strict=False):
matches.append(f"{gid}:{r['CidrIp']}")
except ValueError:
pass
for pair in p.get("UserIdGroupPairs", []):
if pair.get("GroupId") in (client_sgs or []):
matches.append(f"{gid}:sg:{pair['GroupId']}")
for pl in p.get("PrefixListIds", []):
matches.append(f"{gid}:prefixlist:{pl.get('PrefixListId')}(unverified)")
if not sg_ids:
return None, []
return bool(matches), matches
def target_host(source, protocol):
"""Extract the server host from a mount source for each protocol."""
src = source or ""
if protocol == "SMB":
m = re.match(r"^(?://|\\\\)([^/\\]+)", src)
return m.group(1) if m else src
if protocol == "Lustre":
return src.split("@", 1)[0].split(",")[0]
if ":" in src and not src.startswith("["):
return src.rsplit(":", 1)[0].strip("[]")
return src.strip("[]")
def resolve_mount(idx, mount, bundle, client_region, client_vpc, client_account=None):
protocol = classify_fstype(mount["fstype"])
host = target_host(mount["source"], protocol)
opts = mount["options"]
via = "source"
if protocol == "FUSE":
fp = bundle.get("fuse_procs", {}).get(mount["mountpoint"])
if fp:
via = f"fuse_process:{fp['program']}"
host = fp.get("remote") or host
if fp.get("bucket"):
b = fp["bucket"]
if b in idx.buckets:
return {"server_type": "S3-bucket", "id": b, "region": idx.buckets[b].get("region"),
"account": idx.buckets[b].get("account"), "state": "available", "name": b}, via, host
return {"server_type": "S3-bucket(external)", "id": b, "state": None, "name": b,
"external": True}, via, host
return None, via, host
# TLS mounts via efs-utils / S3 Files appear as 127.0.0.1:/ ; recover the real target
if host in ("127.0.0.1", "localhost", "::1"):
via = "tls_proxy"
fst = next((f for f in bundle["fstab"] if f["mountpoint"] == mount["mountpoint"]), None)
if fst:
host = fst["device"].rsplit(":", 1)[0]
via = "tls_proxy+fstab"
else:
port = str(opts.get("port", ""))
for name in bundle["tls_state"]:
if port and name.endswith("." + port):
host = name.split(".")[0]
via = "tls_proxy+state_file"
break
fsm = re.search(r"\b(fs-[0-9a-f]{8,40}|fc-[0-9a-f]{8,40})\b", host)
if fsm and fsm.group(1) in idx.by_id:
return idx.by_id[fsm.group(1)], via, host
probe_ip = (bundle.get("probes", {}).get(mount["mountpoint"]) or {}).get("ip")
if fsm:
# A file system ID that is not in any account scanned: another account's EFS or
# S3 Files, or one that has been deleted. The mount target IP tells us where it lives.
if probe_ip and idx.by_ip.get(probe_ip):
return idx.by_id.get(idx.by_ip[probe_ip][0]["server_id"]), via + "+dns", host
fst = next((f for f in bundle["fstab"] if f["mountpoint"] == mount["mountpoint"]), {})
kind = "S3Files" if fst.get("fstype") == "s3files" else "EFS"
return {"server_type": f"{kind}(external)", "id": fsm.group(1), "state": None,
"external": True, "endpoints": []}, via, host
if host.lower() in idx.by_dns:
return idx.by_id.get(idx.by_dns[host.lower()]), via, host
# suffix match handles AZ specific EFS names and SVM names under a file system DNS name
for dns, sid in idx.by_dns.items():
if host.lower().endswith("." + dns) or dns.endswith("." + host.lower()):
return idx.by_id.get(sid), via, host
candidates = []
try:
ipaddress.ip_address(host)
candidates = idx.by_ip.get(host, [])
except ValueError:
# A DNS name we do not know: use the address the client resolved it to
if probe_ip:
candidates = idx.by_ip.get(probe_ip, [])
if candidates:
via += "+dns"
if candidates:
best = sorted(candidates, key=lambda c: (c.get("vpc_id") != client_vpc,
c.get("account") != client_account,
c.get("region") != client_region))[0]
return idx.by_id.get(best["server_id"]), via, host
ip_for_inst = host if host in idx.inst_by_ip else probe_ip
if ip_for_inst in idx.inst_by_ip:
iid = idx.inst_by_ip[ip_for_inst][0]
inst = idx.inst.get(iid, {})
return {"server_type": "EC2-self-managed", "id": iid, "region": inst.get("_region"),
"account": inst.get("_account"),
"state": (inst.get("State") or {}).get("Name"),
"name": tags_to_dict(inst.get("Tags")).get("Name")}, via, host
return None, via, host
def mount_findings(row, server, mstat, dfi, in_fstab, ep_state, sg_ok):
f = []
if dfi and dfi.get("hung"):
f.append(("HIGH", "mount_unresponsive", "df timed out after 5s; mount looks hung"))
elif dfi and not dfi.get("responsive"):
f.append(("HIGH", "mount_error", f"df failed: {dfi.get('raw', '')[:120]}"))
if not in_fstab:
f.append(("MEDIUM", "not_persistent", "mounted but not in /etc/fstab (autofs or manual?)"))
if server is None:
f.append(("INFO", "unmapped_target", "target not found in this account's inventory "
"(cross account, on premises, peered VPC or DNS alias)"))
else:
st = (server.get("state") or "").lower()
if st and st not in ("available", "running", "ok"):
f.append(("HIGH", "server_not_available", f"server state is {server.get('state')}"))
if ep_state and ep_state.lower() not in ("available", "in-use"):
f.append(("HIGH", "endpoint_not_available", f"endpoint state is {ep_state}"))
if sg_ok is False:
f.append(("HIGH", f"sg_blocks_{row.get('server_port')}",
f"no security group rule found allowing TCP {row.get('server_port')} from this client"))
if row.get("hard_soft") == "soft":
f.append(("MEDIUM", "soft_mount", "soft mount can silently corrupt data on timeouts"))
if row.get("server_type") == "EFS":
if not row.get("tls"):
f.append(("MEDIUM", "efs_no_tls", "EFS mounted without TLS (encryption in transit)"))
if str(row.get("nfs_vers", "")).startswith("3"):
f.append(("HIGH", "efs_nfsv3", "EFS supports NFSv4.0 and 4.1 only"))
if row.get("noresvport") is False:
f.append(("LOW", "efs_no_noresvport", "AWS recommends noresvport for EFS mounts"))
if row.get("protocol") != "NFS":
return f
if row.get("server_type") == "S3Files":
if row.get("nconnect"):
f.append(("HIGH", "s3files_nconnect", "nconnect is not supported by S3 Files"))
if str(row.get("nfs_vers", "")) in ("4", "4.0") or str(row.get("nfs_vers", "")).startswith("3"):
f.append(("HIGH", "s3files_nfs_version", f"S3 Files supports NFS 4.1 and 4.2, mount uses {row.get('nfs_vers')}"))
elif row.get("nconnect"):
f.append(("INFO", "s3files_nconnect_in_use", "mount uses nconnect, which S3 Files does not support"))
if (row.get("retransmissions_total") or 0) > 0 or (row.get("timeouts_total") or 0) > 0:
f.append(("LOW", "rpc_retransmissions",
f"{row.get('retransmissions_total')} retransmissions, {row.get('timeouts_total')} timeouts"))
if row.get("major_timeouts_hint"):
f.append(("MEDIUM", "rpc_timeouts", row["major_timeouts_hint"]))
if row.get("ops_total") == 0 and row.get("open_handles", 0) == 0:
f.append(("INFO", "idle_mount", "no NFS operations recorded since mount and no open files"))
if (row.get("vfslock") or 0) > 0:
f.append(("INFO", "s3files_uses_locks", "workload takes POSIX locks; under S3 Files locks only "
"coordinate NFS clients, not S3 API writers"))
if (row.get("rename_ops") or 0) > 0:
f.append(("INFO", "s3files_uses_rename", "workload renames files; supported by S3 Files, "
"but each rename becomes S3 object work after sync"))
if row.get("avg_write_bytes") and row["avg_write_bytes"] < 16384 and (row.get("write_ops") or 0) > 1000:
f.append(("INFO", "s3files_small_writes", f"average write is {row['avg_write_bytes']} bytes; "
"small IO is slower and costlier on S3 Files than EFS"))
return f
SLOW_MS = 1000
def annotate_network(idx, row, server, client_account, client_vpc, client_subnet, target_ip):
"""Adds account and network path columns to a mount row; returns findings."""
f = []
np_ = idx.network_path(client_account, client_vpc, client_subnet, target_ip) if target_ip else {"path": "unresolved"}
server_account = (server or {}).get("account")
row.update({"client_account": client_account, "server_account": server_account, "target_ip": target_ip,
"network_path": np_.get("path"), "network_via": np_.get("via"),
"route_table": np_.get("route_table"), "matched_route": np_.get("route"),
"peer_account": np_.get("peer_account"), "peer_vpc": np_.get("peer_vpc"),
"peer_region": np_.get("peer_region"), "network_detail": np_.get("detail")})
peer = np_.get("peer_account")
evidence = None
if server_account and server_account != client_account:
xa, evidence = "yes", f"target resource is in account {server_account}"
elif server is not None and server.get("external"):
if peer and peer != client_account:
xa, evidence = "yes", f"{server.get('id')} is not in any scanned account; its address is in account {peer}'s network ({np_.get('path')})"
else:
xa, evidence = "possible", f"{server.get('id')} is not owned by any scanned account"
elif peer and peer != client_account:
xa, evidence = "yes", f"target address is in account {peer}'s network via {np_.get('path')} {np_.get('via') or ''}".strip()
elif server is None and np_.get("path") in ("vpc_peering", "transit_gateway", "cloud_wan", "shared_vpc"):
xa, evidence = "possible", f"unmapped target reached via {np_.get('path')} {np_.get('via') or ''}".strip()
elif np_.get("path") == "vpn_or_direct_connect":
xa, evidence = "on_premises", f"target reached via {np_.get('via')}"
else:
xa = "no"
row["cross_account"] = xa
row["cross_account_evidence"] = evidence
if xa == "yes":
f.append(("MEDIUM", "cross_account_mount", evidence))
elif xa == "possible":
f.append(("INFO", "possible_cross_account_mount", evidence))
elif xa == "on_premises":
f.append(("INFO", "on_premises_target", evidence))
if np_.get("path") == "no_route":
f.append(("HIGH", "no_route_to_target", f"route table {np_.get('route_table')} has no route to {target_ip}"))
return f
def probe_diagnosis(row, df_hung, df_ok, df_ms, tcp, stat_rc, stat_ms, rd_rc, rd_ms):
hung = df_hung or stat_rc in (124, 137) or rd_rc in (124, 137)
errs = (df_ok is False and not df_hung) or (stat_rc not in (None, 0) and not hung) or (rd_rc not in (None, 0) and not hung)
if hung:
d = "server_unreachable" if tcp == "closed" else ("server_reachable_fs_not_responding" if tcp == "open" else "not_responding")
elif errs:
d = "error"
elif max([x for x in (df_ms, stat_ms, rd_ms) if x is not None] or [0]) > SLOW_MS:
d = "slow"
else:
d = "ok"
row["responsive"] = d in ("ok", "slow")
row["probe_diagnosis"] = d
f = []
if d == "server_unreachable":
f.append(("HIGH", "server_unreachable", f"probes timed out and TCP {row.get('server_port')} to {row.get('target_ip')} is closed"))
elif d == "server_reachable_fs_not_responding":
f.append(("HIGH", "fs_not_responding", "server port accepts connections but file system calls time out"))
elif d == "slow":
f.append(("LOW", "slow_mount", f"a probe took over {SLOW_MS} ms (statfs {df_ms}, getattr {stat_ms}, readdir {rd_ms})"))
elif d == "error":
err = row.get("probe_readdir_error")
if err == "Permission denied":
f.append(("MEDIUM", "access_denied", "root on this host cannot list the share root; expected with "
"root squashing, but check that the application user can"))
elif err:
f.append(("HIGH", "mount_error", f"directory read failed: {err}"))
else:
f.append(("MEDIUM", "probe_error", "a probe returned an error rather than timing out"))
if tcp == "closed" and d in ("ok", "slow"):
f.append(("MEDIUM", "tcp_probe_failed", "mount works but a new TCP connection to the server port failed"))
return f
def mount_health(row, server, ms, bundle, idx, client_region, client_az):
"""Capacity, latency, cross region and permission checks for one client mount."""
f = []
# Capacity: skip the virtual exabyte sizes reported by EFS, S3 Files and FUSE S3 mounts
size = row.get("size_bytes") or 0
if size and size < 2 ** 60:
pct = to_int(str(row.get("use_pct") or "").rstrip("%"))
if pct is None and row.get("used_bytes") is not None:
pct = round(100 * row["used_bytes"] / size)
row["capacity_pct"] = pct
if pct is not None and pct >= 95:
f.append(("HIGH", "capacity_critical", f"{pct}% of {fmt_bytes(size)} used"))
elif pct is not None and pct >= 85:
f.append(("MEDIUM", "capacity_high", f"{pct}% of {fmt_bytes(size)} used"))
sf = {}
try:
sf = json.loads(row.get("statfs") or "{}")
except ValueError:
pass
files, ffree = to_int(sf.get("files")), to_int(sf.get("ffree"))
if files and files < 2 ** 60 and ffree is not None:
ipct = round(100 * (files - ffree) / files)
row["inode_pct"] = ipct
if ipct >= 90:
f.append(("HIGH" if ipct >= 95 else "MEDIUM", "inodes_high", f"{ipct}% of {files} inodes used"))
# Latency from the kernel's own RPC statistics
if row.get("protocol") == "NFS" and ms.get("per_op"):
ops = ms["per_op"]
meta = [ops.get(k) for k in ("GETATTR", "LOOKUP", "ACCESS") if (ops.get(k) or {}).get("ops", 0) >= 50]
meta_rtt = max((o["avg_rtt_ms"] for o in meta), default=None)
data = [ops.get(k) for k in ("READ", "WRITE") if (ops.get(k) or {}).get("ops", 0) >= 50]
data_rtt = max((o["avg_rtt_ms"] for o in data), default=None)
queue = max((o["avg_exec_ms"] - o["avg_rtt_ms"] for o in meta + data), default=None)
row.update({"avg_metadata_rtt_ms": meta_rtt, "avg_data_rtt_ms": data_rtt,
"avg_client_queue_ms": round(queue, 3) if queue is not None else None})
if meta_rtt is not None and meta_rtt > 10:
f.append(("MEDIUM", "high_metadata_latency", f"average metadata RPC round trip {meta_rtt} ms since mount"))
if data_rtt is not None and data_rtt > 50:
f.append(("MEDIUM", "high_data_latency", f"average READ/WRITE RPC round trip {data_rtt} ms since mount"))
if queue is not None and queue > 20:
f.append(("LOW", "client_side_queueing", f"RPCs wait {round(queue, 1)} ms on average before being sent "
"(RPC slot limits, a single TCP connection or client CPU)"))
# Cross AZ: a mount target exists in the client's AZ but the mount uses another one
eps = (server or {}).get("endpoints") or []
if row.get("endpoint_az") and client_az and row["endpoint_az"] != client_az:
local = any(e.get("az") == client_az for e in eps)
f.append(("MEDIUM" if local else "INFO", "cross_az_mount",
f"client in {client_az} uses an endpoint in {row['endpoint_az']}"
+ ("; a mount target exists in the client's AZ" if local else "")
+ " (extra latency and cross AZ data transfer charges)"))
# Cross region
sreg = (server or {}).get("region")
preg = row.get("peer_region")
is_bucket = str((server or {}).get("server_type", "")).startswith("S3-bucket")
if sreg and client_region and sreg != client_region and not is_bucket:
f.append(("MEDIUM", "cross_region_mount", f"client in {client_region}, file server in {sreg}"))
row["cross_region"] = True
elif preg and client_region and preg != client_region:
f.append(("MEDIUM", "cross_region_mount", f"target reached over inter region peering to {preg}"))
row["cross_region"] = True
if str((server or {}).get("server_type", "")).startswith("S3-bucket"):
b = idx.buckets.get((server or {}).get("id"), {})
if b.get("region") and client_region and b["region"] != client_region:
f.append(("MEDIUM", "cross_region_bucket", f"FUSE mount of a bucket in {b['region']} from {client_region}"))
row["cross_region"] = True
# Permissions
mode = row.get("root_mode")
if mode and mode.isdigit():
m = int(mode, 8)
if m & 0o002 and not m & 0o1000:
f.append(("MEDIUM", "world_writable_root", f"share root is mode {mode} without the sticky bit"))
fst = (row.get("fstab_options") or "")
if re.search(r"(^|,)pass(word)?=", fst):
f.append(("HIGH", "password_in_fstab", "a password is written directly in /etc/fstab mount options"))
for opt in fst.split(","):
k, _, v = opt.partition("=")
if k in ("credentials", "cred"):
info = (bundle.get("credfiles") or {}).get(v, "")
parts = info.split()
if len(parts) >= 2 and parts[0].isdigit() and (int(parts[0], 8) & 0o077 or parts[1] != "root"):
f.append(("HIGH", "credentials_file_exposed", f"{v} is mode {parts[0]} owned by {parts[1]}; "
"should be 600 and owned by root"))
elif info and not parts[0].isdigit():
f.append(("MEDIUM", "credentials_file_missing", f"{v}: {info}"))
if row.get("server_type") == "EFS" and row.get("fstab_type") == "efs" and "iam" not in fst.split(","):
f.append(("INFO", "efs_mount_without_iam", "EFS mounted without the iam option, so file system policy "
"conditions on IAM identity do not apply to this client"))
return f
def _as_list(x):
if x is None:
return []
if isinstance(x, dict) and "error" in x and len(x) == 1:
return []
return x if isinstance(x, list) else [x]
def windows_rows(idx, inst, iid, region, account, text, host):
"""Rows for SMB / NFS mounts on a Windows host."""
rows, findings = [], []
try:
w = json.loads(text)
except ValueError as e:
host["bundle_error"] = f"windows json: {e}"
return rows, findings
meta = w.get("meta") or {}
host.update({"platform": "Windows", "hostname": meta.get("hostname"), "kernel": meta.get("os"),
"windows_build": meta.get("build"), "domain": meta.get("domain"),
"smb_client_config": w.get("smb_client_config"),
"smb_shares_served": _as_list(w.get("smb_shares_served")),
"nfs_shares_served": _as_list(w.get("nfs_shares_served")),
"smb_server_sessions": _as_list(w.get("smb_server_sessions")),
"connections": _as_list(w.get("connections")), "net_use": w.get("net_use"),
"nfs_client_feature": w.get("nfs_client_feature"), "dfs_cache": w.get("dfs_client_cache")})
host["is_nfs_server"] = bool(host["nfs_shares_served"])
host["is_smb_server"] = bool(host["smb_shares_served"])
probes = {str(p.get("RemotePath", "")).lower(): p for p in _as_list(w.get("probes"))}
conns = {}
for c in _as_list(w.get("smb_connections")):
conns[(str(c.get("ServerName", "")).lower(), str(c.get("ShareName", "")).lower())] = c
seen = {}
def add(local, remote, source_kind, extra, proto=None):
if not remote:
return
key = (str(local).upper(), str(remote).lower())
if key in seen:
seen[key]["source_kinds"] += f",{source_kind}"
return
protocol = proto or ("SMB" if str(remote).startswith("\\\\") else "NFS")
if str(remote).startswith("\\\\"):
thost = target_host(remote, "SMB")
else:
thost = remote.split(":")[0]
share = remote.rstrip("\\").split("\\")[-1] if protocol == "SMB" else None
server = None
if thost.lower() in idx.by_dns:
server = idx.by_id.get(idx.by_dns[thost.lower()])
elif idx.by_ip.get(thost):
server = idx.by_id.get(idx.by_ip[thost][0]["server_id"])
elif thost in idx.inst_by_ip:
tid = idx.inst_by_ip[thost][0]
ti = idx.inst.get(tid, {})
server = {"server_type": "EC2-self-managed", "id": tid, "region": ti.get("_region"),
"state": (ti.get("State") or {}).get("Name"), "name": tags_to_dict(ti.get("Tags")).get("Name")}
else:
for dns, sid in idx.by_dns.items():
if thost.lower().split(".")[0] == dns.split(".")[0]:
server = idx.by_id.get(sid)
break
pr = probes.get(str(remote).lower(), {})
pip = pr.get("IP")
if server is None and pip:
if idx.by_ip.get(pip):
server = idx.by_id.get(idx.by_ip[pip][0]["server_id"])
elif pip in idx.inst_by_ip:
tid = idx.inst_by_ip[pip][0]
ti = idx.inst.get(tid, {})
server = {"server_type": "EC2-self-managed", "id": tid, "region": ti.get("_region"),
"account": ti.get("_account"), "state": (ti.get("State") or {}).get("Name"),
"name": tags_to_dict(ti.get("Tags")).get("Name")}
c = conns.get((thost.lower(), str(share or "").lower()), {})
dialect = c.get("Dialect")
row = {"account": account, "region": region, "instance_id": iid, "instance_name": host.get("name"),
"hostname": meta.get("hostname"), "client_az": host.get("az"), "client_ip": host.get("private_ip"),
"client_vpc": inst.get("VpcId"), "kernel": meta.get("os"), "mountpoint": local, "source": remote,
"fstype": protocol.lower(), "protocol": protocol, "server_port": PROTO_PORT.get(protocol),
"resolved_host": thost, "resolved_via": "windows:" + source_kind, "source_kinds": source_kind,
"server_type": (server or {}).get("server_type", "unknown"), "server_id": (server or {}).get("id"),
"server_name": (server or {}).get("name"), "server_state": (server or {}).get("state"),
"smb_vers": dialect, "smb_seal": c.get("Encrypted"), "smb_signed": c.get("Signed"),
"smb_username": c.get("UserName") or extra.get("UserName"), "smb_num_opens": c.get("NumOpens"),
"smb_continuously_available": c.get("ContinuouslyAvailable"),
"status": extra.get("Status"), "size_bytes": extra.get("Size"),
"avail_bytes": extra.get("FreeSpace"),
"used_bytes": (extra["Size"] - extra["FreeSpace"]) if extra.get("Size") and extra.get("FreeSpace") is not None else None,
"in_fstab": source_kind == "persistent_registry", "tls": bool(c.get("Encrypted")),
"probe_tcp": {True: "open", False: "closed"}.get(pr.get("Tcp")), "probe_tcp_ms": pr.get("TcpMs"),
"probe_result": pr.get("Probe"), "probe_ms": pr.get("ProbeMs"), "probe_error": pr.get("Error")}
fl = []
if pr:
st = pr.get("Probe")
tcp = row["probe_tcp"]
if st == "timeout":
d = "server_unreachable" if tcp == "closed" else "server_reachable_share_not_responding"
fl.append(("HIGH", d, f"listing {remote} timed out after 5s (TCP {row['server_port']} {tcp})"))
elif st == "error":
d = "error"
fl.append(("INFO", "probe_error", f"{pr.get('Error')} (SSM runs as SYSTEM, so per user "
"credentials and drive maps may not apply)"))
else:
d = "slow" if (pr.get("ProbeMs") or 0) > SLOW_MS + 1500 else "ok" # Start-Job adds ~1s
row["probe_diagnosis"] = d
row["responsive"] = d in ("ok", "slow")
fl += annotate_network(idx, row, server, account, inst.get("VpcId"), inst.get("SubnetId"), pip)
fl += mount_health(row, server, {}, {}, idx, region, host.get("az"))
if pr.get("Probe") == "error" and "denied" in str(pr.get("Error", "")).lower():
fl.append(("INFO", "access_denied", "SYSTEM cannot list this share; check the users who rely on it can"))
if server is None:
fl.append(("INFO", "unmapped_target", "target not found in this account's inventory"))
if protocol == "SMB" and dialect and str(dialect).startswith(("1", "2.0")):
fl.append(("MEDIUM", "smb_old_dialect", f"SMB dialect {dialect}; prefer 3.x"))
if protocol == "SMB" and c and not c.get("Encrypted"):
fl.append(("LOW", "smb_unencrypted", "SMB session is not encrypted"))
if extra.get("Status") and str(extra.get("Status")) not in ("0", "OK", "Ok"):
fl.append(("MEDIUM", "mapping_not_ok", f"mapping status {extra.get('Status')}"))
row["findings"] = "; ".join(code for _, code, _ in fl)
for sev, code, detail in fl:
findings.append({"severity": sev, "code": code, "account": account, "region": region,
"instance_id": iid, "mountpoint": local, "resource": row["server_id"] or thost,
"detail": detail})
if server and isinstance(server.get("clients"), list):
server["clients"].append({"instance_id": iid, "mountpoint": local, "host": meta.get("hostname")})
seen[key] = row
rows.append(row)
for m in _as_list(w.get("smb_mappings")):
add(m.get("LocalPath") or "(no drive)", m.get("RemotePath"), "smb_mapping", m)
for d in _as_list(w.get("mapped_logical_disks")) + _as_list(w.get("network_logical_disks")):
add(d.get("DeviceID"), d.get("ProviderName"), "logical_disk", d)
for d in _as_list(w.get("persistent_user_drives")):
add(str(d.get("Drive", "")).upper() + ":", d.get("RemotePath"), "persistent_registry", d)
for c in _as_list(w.get("smb_connections")):
unc = f"\\\\{c.get('ServerName')}\\{c.get('ShareName')}"
if not any(r["source"].lower() == unc.lower() for r in rows) and c.get("ShareName") not in ("IPC$",):
add("(unmapped session)", unc, "smb_session", c)
nfs_txt = w.get("nfs_client_mounts")
if isinstance(nfs_txt, str):
for line in nfs_txt.splitlines():
mm = re.match(r"^\s*([A-Za-z]:)\s+(\S+)\s+(.*)$", line)
if mm:
add(mm.group(1), mm.group(2), "nfs_client", {"options": mm.group(3)}, proto="NFS")
host["nfs_mount_count"] = len(rows)
return rows, findings
def build_mounts(regions_data, ssm_results, idx, servers, account):
rows, findings, client_hosts = [], [], []
for R in regions_data:
region = R["region"]
account = R.get("account", account)
for iid, res in (ssm_results.get(f"{account}:{region}") or {}).items():
inst = idx.inst.get(iid, {})
b = (res.get("bundle") or {})
host = {"instance_id": iid, "account": account, "region": region, "ssm_status": res.get("status"),
"bundle_ok": b.get("ok"), "bundle_error": b.get("error"),
"name": tags_to_dict(inst.get("Tags")).get("Name"),
"private_ip": inst.get("PrivateIpAddress"),
"az": (inst.get("Placement") or {}).get("AvailabilityZone"),
"vpc_id": inst.get("VpcId"), "nfs_mount_count": 0, "is_nfs_server": False}
if not b.get("ok"):
client_hosts.append(host)
continue
if res.get("platform") == "Windows":
wrows, wf = windows_rows(idx, inst, iid, region, account, b["text"], host)
rows.extend(wrows)
findings.extend(wf)
client_hosts.append(host)
continue
bundle = parse_bundle(b["text"])
host.update({"platform": "Linux", "fuse_daemons": bundle["fuse_procs"], "cifs_debug": bundle["cifs"],
"lustre": bundle["lustre"], "other_remote": bundle["other_remote"]})
host.update({"hostname": bundle["meta"].get("hostname"), "kernel": bundle["meta"].get("kernel"),
"packages": bundle["packages"], "nfs_mount_count": len(bundle["proc_mounts"]),
"is_nfs_server": bundle["nfs_server"]["is_nfs_server"],
"nfs_server_detail": bundle["nfs_server"]["lines"],
"fstab_entries": bundle["fstab"], "autofs": bundle["autofs"],
"tls_proxies": bundle["tls_proxies"], "connections": bundle["connections"],
"rpc_tunables": bundle["rpc_tunables"], "bundle_complete": bundle["complete"]})
client_sgs = [g["GroupId"] for g in inst.get("SecurityGroups", [])]
client_ip = inst.get("PrivateIpAddress") or bundle["meta"].get("local-ipv4")
client_az = host["az"] or bundle["meta"].get("placement/availability-zone")
fstab_mps = {f["mountpoint"]: f for f in bundle["fstab"]}
ms_by_mp = {}
for ms in bundle["mountstats"]:
ms_by_mp.setdefault(ms["mountpoint"], ms)
# fstab entries that are not mounted
mounted = {m["mountpoint"] for m in bundle["proc_mounts"]}
for mp, f in fstab_mps.items():
if mp not in mounted and "noauto" not in f["options"]:
findings.append({"severity": "MEDIUM", "code": "fstab_not_mounted", "account": account,
"region": region, "instance_id": iid, "mountpoint": mp,
"resource": f["device"], "detail": "in /etc/fstab but not currently mounted"})
for m in bundle["proc_mounts"]:
protocol = classify_fstype(m["fstype"])
port = PROTO_PORT.get(protocol)
server, via, thost = resolve_mount(idx, m, bundle, region, inst.get("VpcId"), account)
ms = ms_by_mp.get(m["mountpoint"], {})
o = {**m["options"], **(ms.get("opts") or {})}
ops = ms.get("per_op", {})
ev = ms.get("events", {})
by = ms.get("bytes", {})
xp = ms.get("xprt", {})
dfi = bundle["df"].get(m["mountpoint"], {})
of = bundle["open_files"].get(m["mountpoint"], {"open_handles": 0, "processes": []})
ep = None
if server and server.get("endpoints"):
ips = idx.by_ip.get(thost, [])
ep_id = ips[0].get("endpoint") if ips else None
eps = server["endpoints"]
ep = next((e for e in eps if ep_id and e.get("id") == ep_id), None) or \
next((e for e in eps if e.get("az") == client_az), None) or eps[0]
sg_ok, sg_rules = (None, [])
if port and ep:
sg_ok, sg_rules = sg_allows(idx, ep.get("security_groups"), client_ip, client_sgs, port)
elif port and server and server.get("server_type") == "EC2-self-managed":
tgt = idx.inst.get(server["id"], {})
sg_ok, sg_rules = sg_allows(idx, [g["GroupId"] for g in tgt.get("SecurityGroups", [])],
client_ip, client_sgs, port)
fp = bundle["fuse_procs"].get(m["mountpoint"], {})
tot = lambda k: sum(v.get(k, 0) for v in ops.values()) # noqa: E731
read, write = ops.get("READ", {}), ops.get("WRITE", {})
row = {
"account": account, "region": region, "instance_id": iid, "instance_name": host["name"],
"hostname": host["hostname"], "client_az": client_az, "client_ip": client_ip,
"client_vpc": inst.get("VpcId"), "client_sgs": " ".join(client_sgs),
"kernel": host["kernel"], "mountpoint": m["mountpoint"], "source": m["source"],
"fstype": m["fstype"], "protocol": protocol, "server_port": port,
"resolved_host": thost, "resolved_via": via,
"fuse_program": fp.get("program"), "fuse_remote": fp.get("remote"),
"fuse_bucket": fp.get("bucket"), "fuse_prefix": fp.get("prefix"),
"fuse_daemon_pid": fp.get("pid"), "fuse_daemon_rss_kb": fp.get("rss_kb"),
"fuse_daemon_uptime_s": fp.get("elapsed_s"), "fuse_args": fp.get("args"),
"smb_vers": o.get("vers") if protocol == "SMB" else None,
"smb_cache": o.get("cache"), "smb_seal": bool(o.get("seal")),
"smb_username": o.get("username"), "smb_domain": o.get("domain"),
"smb_uid": o.get("uid"), "smb_file_mode": o.get("file_mode"),
"lustre_flock": bool(o.get("flock")),
"server_type": (server or {}).get("server_type", "unknown"),
"server_id": (server or {}).get("id"), "server_name": (server or {}).get("name"),
"server_region": (server or {}).get("region"), "server_state": (server or {}).get("state"),
"endpoint_id": (ep or {}).get("id"), "endpoint_ip": (ep or {}).get("ip"),
"endpoint_az": (ep or {}).get("az"), "endpoint_state": (ep or {}).get("state"),
"same_az": (ep or {}).get("az") == client_az if ep and ep.get("az") else None,
"sg_allows_2049": sg_ok, "sg_matching_rules": " ".join(sg_rules),
"nfs_vers": (o.get("vers") or o.get("nfsvers")) if protocol == "NFS" else None, "minorversion": o.get("minorversion"),
"proto": o.get("proto"), "port": o.get("port"),
"tls": bool(o.get("tls")) or via.startswith("tls_proxy") or bool(o.get("seal")),
"hard_soft": "soft" if o.get("soft") else "hard",
"rsize": o.get("rsize"), "wsize": o.get("wsize"), "timeo": o.get("timeo"),
"retrans": o.get("retrans"), "sec": o.get("sec"), "noresvport": bool(o.get("noresvport")),
"actimeo": o.get("actimeo"), "acregmin": o.get("acregmin"), "acregmax": o.get("acregmax"),
"acdirmin": o.get("acdirmin"), "acdirmax": o.get("acdirmax"),
"lookupcache": o.get("lookupcache"), "local_lock": o.get("local_lock"),
"clientaddr": o.get("clientaddr"), "nconnect": o.get("nconnect"),
"read_only": bool(o.get("ro")), "options_raw": m["options_raw"],
"in_fstab": m["mountpoint"] in fstab_mps,
"fstab_device": (fstab_mps.get(m["mountpoint"]) or {}).get("device"),
"fstab_type": (fstab_mps.get(m["mountpoint"]) or {}).get("fstype"),
"fstab_options": (fstab_mps.get(m["mountpoint"]) or {}).get("options"),
"responsive": dfi.get("responsive"), "df_latency_ms": dfi.get("latency_ms"),
"size_bytes": dfi.get("size_bytes"), "used_bytes": dfi.get("used_bytes"),
"avail_bytes": dfi.get("avail_bytes"), "use_pct": dfi.get("use_pct"),
"statfs": json.dumps(bundle["statfs"].get(m["mountpoint"], {})),
"mount_age_s": ms.get("age_s"), "lease_time": (ms.get("nfsv4") or {}).get("lease_time"),
"pnfs": (ms.get("nfsv4") or {}).get("pnfs"),
"server_read_bytes": by.get("serverreadbytes"), "server_write_bytes": by.get("serverwritebytes"),
"app_read_bytes": by.get("normalreadbytes"), "app_write_bytes": by.get("normalwritebytes"),
"direct_read_bytes": by.get("directreadbytes"), "direct_write_bytes": by.get("directwritebytes"),
"read_ops": read.get("ops"), "write_ops": write.get("ops"),
"avg_read_bytes": round(read["bytes_recv"] / read["ops"]) if read.get("ops") else None,
"avg_write_bytes": round(write["bytes_sent"] / write["ops"]) if write.get("ops") else None,
"avg_read_rtt_ms": read.get("avg_rtt_ms"), "avg_write_rtt_ms": write.get("avg_rtt_ms"),
"avg_read_exec_ms": read.get("avg_exec_ms"), "avg_write_exec_ms": write.get("avg_exec_ms"),
"getattr_ops": (ops.get("GETATTR") or {}).get("ops"),
"lookup_ops": (ops.get("LOOKUP") or {}).get("ops"),
"access_ops": (ops.get("ACCESS") or {}).get("ops"),
"open_ops": (ops.get("OPEN") or {}).get("ops"),
"commit_ops": (ops.get("COMMIT") or {}).get("ops"),
"create_ops": (ops.get("CREATE") or {}).get("ops"),
"remove_ops": (ops.get("REMOVE") or {}).get("ops"),
"rename_ops": (ops.get("RENAME") or {}).get("ops"),
"setattr_ops": (ops.get("SETATTR") or {}).get("ops"),
"readdir_ops": ((ops.get("READDIR") or {}).get("ops") or 0) + ((ops.get("READDIRPLUS") or {}).get("ops") or 0),
"lock_ops": sum((ops.get(k) or {}).get("ops", 0) for k in ("LOCK", "LOCKT", "LOCKU")),
"ops_total": tot("ops") if ops else None,
"retransmissions_total": tot("retransmissions") if ops else None,
"timeouts_total": tot("timeouts") if ops else None,
"op_errors_total": tot("errors") if ops else None,
"vfslock": ev.get("vfslock"), "vfsfsync": ev.get("vfsfsync"), "vfsopen": ev.get("vfsopen"),
"sillyrenames": ev.get("sillyrenames"), "short_reads": ev.get("shortreads"),
"short_writes": ev.get("shortwrites"), "delay_events": ev.get("delay"),
"attr_invalidates": ev.get("attrinvalidates"), "data_invalidates": ev.get("datainvalidates"),
"xprt_proto": xp.get("proto"), "xprt_connect_count": xp.get("connect_count"),
"xprt_idle_s": xp.get("idle_time"), "xprt_sends": xp.get("sends"),
"xprt_bad_xids": xp.get("bad_xids"), "xprt_max_slots": xp.get("max_slots"),
"open_handles": of["open_handles"],
"processes": " ".join(sorted({f"{p['comm']}({p['pid']})" for p in of["processes"]})),
}
if row["timeouts_total"]:
row["major_timeouts_hint"] = f"{row['timeouts_total']} RPC major timeouts since mount"
fl = mount_findings(row, server, ms, dfi, row["in_fstab"], row["endpoint_state"], sg_ok)
if protocol == "SMB" and not row["smb_seal"] and str(row.get("smb_vers") or "").startswith(("1", "2.0")):
fl.append(("MEDIUM", "smb_old_dialect", f"SMB vers={row['smb_vers']}; prefer 3.x"))
if protocol == "SMB" and o.get("password"):
fl.append(("HIGH", "smb_password_in_options", "password passed as a mount option"))
if protocol == "FUSE" and not fp:
fl.append(("INFO", "fuse_daemon_unknown", "FUSE mount whose daemon could not be identified"))
if protocol == "FUSE" and fp.get("program") in ("s3fs", "goofys", "mount-s3", "mountpoint-s3"):
fl.append(("INFO", "s3_fuse_mount", f"S3 mounted through {fp['program']}: limited POSIX "
"semantics; S3 Files is the candidate replacement"))
pr = bundle["probes"].get(m["mountpoint"], {})
row.update({"probe_host": pr.get("host"), "probe_ip": pr.get("ip"), "probe_tcp": pr.get("tcp"),
"probe_tcp_ms": pr.get("tcp_ms"), "probe_statfs_ms": dfi.get("latency_ms"),
"probe_getattr_rc": pr.get("stat_rc"), "probe_getattr_ms": pr.get("stat_ms"),
"probe_readdir_rc": pr.get("readdir_rc"), "probe_readdir_ms": pr.get("readdir_ms")})
target_ip = pr.get("ip") or (ep or {}).get("ip") or (thost if re.match(r"^\d+\.\d+\.\d+\.\d+$", thost or "") else None)
fl = [x for x in fl if x[1] not in ("mount_unresponsive", "mount_error")]
row["target_ip"] = target_ip
if dfi or pr:
fl += probe_diagnosis(row, dfi.get("hung"), dfi.get("responsive"), dfi.get("latency_ms"),
pr.get("tcp"), pr.get("stat_rc"), pr.get("stat_ms"),
pr.get("readdir_rc"), pr.get("readdir_ms"))
fl += annotate_network(idx, row, server, account, inst.get("VpcId"), inst.get("SubnetId"), target_ip)
row.update({"root_mode": pr.get("root_mode"), "root_uid": pr.get("root_uid"),
"root_gid": pr.get("root_gid"), "probe_readdir_error": pr.get("readdir_error")})
if row.get("probe_diagnosis") == "error":
fl = [x for x in fl if x[1] != "probe_error"] + [x for x in probe_diagnosis(
row, dfi.get("hung"), dfi.get("responsive"), dfi.get("latency_ms"), pr.get("tcp"),
pr.get("stat_rc"), pr.get("stat_ms"), pr.get("readdir_rc"), pr.get("readdir_ms"))
if x[1] in ("access_denied", "mount_error", "probe_error")]
fl += mount_health(row, server, ms, bundle, idx, region, client_az)
row["findings"] = "; ".join(code for _, code, _ in fl)
row["_detail"] = {"mountstats": ms, "open_files": of, "df": dfi, "probe": pr}
for sev, code, detail in fl:
findings.append({"severity": sev, "code": code, "account": account, "region": region,
"instance_id": iid, "mountpoint": m["mountpoint"],
"resource": row["server_id"] or thost, "detail": detail})
if server and isinstance(server.get("clients"), list):
server["clients"].append({"instance_id": iid, "mountpoint": m["mountpoint"],
"host": host["hostname"]})
rows.append(row)
client_hosts.append(host)
return rows, findings, client_hosts
SERVER_PORTS = {"EFS": [2049], "S3Files": [2049], "FSx-ONTAP": [2049, 445], "FSx-OPENZFS": [2049],
"FSx-WINDOWS": [445], "FSx-LUSTRE": [988], "FSx-FileCache": [988], "StorageGateway": [2049, 445]}
def policy_cross_account(doc, own):
"""Accounts granted by a resource policy, and whether it is open to any principal."""
try:
pol = json.loads(doc) if isinstance(doc, str) else (doc or {})
except ValueError:
return set(), False, False
accounts, open_any, org_scoped = set(), False, False
stmts = pol.get("Statement", [])
for st in stmts if isinstance(stmts, list) else [stmts]:
if st.get("Effect") != "Allow":
continue
pr = st.get("Principal")
vals = ["*"] if pr == "*" else (_as_list((pr or {}).get("AWS")) if isinstance(pr, dict) else [])
cond = json.dumps(st.get("Condition") or {})
for v in vals:
if v == "*":
if "aws:PrincipalOrgID" in cond or "aws:PrincipalOrgPaths" in cond:
org_scoped = True
ids = set(re.findall(r"\b\d{12}\b", cond))
if ids:
accounts |= ids
elif not org_scoped:
open_any = True
else:
accounts |= set(re.findall(r"\b(\d{12})\b", str(v)))
return accounts - {own}, open_any, org_scoped
def sg_cross_account(idx, s):
"""Security group rules on a server's endpoints that admit clients from outside its account or VPC."""
ports = SERVER_PORTS.get(s["server_type"], [2049, 445, 988])
own = s.get("account")
own_cidrs = []
for ep in s.get("endpoints", []):
v = idx.vpcs.get(ep.get("vpc_id")) or {}
own_cidrs += [ipaddress.ip_network(c, strict=False) for c in v.get("cidrs", [])]
hits = []
for gid in {g for ep in s.get("endpoints", []) for g in (ep.get("security_groups") or [])}:
for perm in (idx.sg.get(gid) or {}).get("IpPermissions", []):
proto = str(perm.get("IpProtocol"))
if proto not in ("tcp", "6", "-1"):
continue
if proto != "-1" and not any(perm.get("FromPort", 0) <= p <= perm.get("ToPort", -1) for p in ports):
continue
for pair in perm.get("UserIdGroupPairs", []):
uid = pair.get("UserId")
if (uid and own and uid != own) or pair.get("VpcPeeringConnectionId"):
hits.append(("account", f"{gid} allows {pair.get('GroupId')} in account {uid}"
+ (f" via {pair['VpcPeeringConnectionId']}" if pair.get("VpcPeeringConnectionId") else ""), uid))
for r in perm.get("IpRanges", []):
try:
n = ipaddress.ip_network(r["CidrIp"], strict=False)
except (KeyError, ValueError):
continue
if n.prefixlen == 0:
hits.append(("open", f"{gid} allows {r['CidrIp']}", None))
continue
if any(n.subnet_of(o) for o in own_cidrs if o.version == n.version):
continue
owner = idx.vpc_for_ip(str(n.network_address), exclude_account=own)
who = f" (VPC {owner['vpc_id']} in account {owner['owner']})" if owner else ""
hits.append(("outside_vpc", f"{gid} allows {r['CidrIp']}{who}", owner and owner["owner"]))
return hits
def flow_log_sources(sessions, regions_data, servers, idx, args, rec):
"""Optional: who actually connected to each file server's endpoints, from VPC Flow Logs in CloudWatch Logs."""
start = int((NOW - dt.timedelta(days=args.days)).timestamp())
end = int(NOW.timestamp())
for R in regions_data:
acct, region = R.get("account"), R["region"]
eps = {}
for s in servers:
if s.get("account") == acct and s.get("region") == region:
for ep in s.get("endpoints", []):
if ep.get("ip"):
eps[ep["ip"]] = (s, ep)
groups = sorted({fl["LogGroupName"] for fl in R.get("flow_logs", [])
if fl.get("LogDestinationType", "cloud-watch-logs") == "cloud-watch-logs" and fl.get("LogGroupName")})
if not eps or not groups:
continue
logs = sessions[acct].client("logs", region_name=region, config=BOTO_CFG)
ips = sorted(eps)
for i in range(0, len(ips), 50):
chunk = ips[i:i + 50]
q = ('fields srcAddr, dstAddr, dstPort, bytes '
'| filter action = "ACCEPT" and dstPort in [2049, 445, 988] and dstAddr in ["' + '", "'.join(chunk) + '"] '
'| stats count(*) as flows, sum(bytes) as bytes, min(start) as first_seen, max(end) as last_seen '
'by srcAddr, dstAddr, dstPort | limit 10000')
qid = (rec.call(region, "logs", "start_query", logs.start_query, logGroupNames=groups[:50],
startTime=start, endTime=end, queryString=q, default={}) or {}).get("queryId")
if not qid:
continue
log(f"{acct}/{region}: flow log query {qid} over {len(groups[:50])} log groups")
deadline = time.time() + 300
res = {}
while time.time() < deadline:
res = rec.call(region, "logs", "get_query_results", logs.get_query_results, queryId=qid, default={}) or {}
if res.get("status") in ("Complete", "Failed", "Cancelled", "Timeout"):
break
time.sleep(3)
for row in res.get("results", []):
r = {f["field"]: f["value"] for f in row}
if r.get("dstAddr") not in eps:
continue
s, ep = eps[r["dstAddr"]]
src = r.get("srcAddr")
enis = idx.eni_by_ip.get(src, [])
eni = enis[0] if enis else {}
np_ = idx.network_path(acct, ep.get("vpc_id"), ep.get("subnet"), src)
src_owner = eni.get("OwnerId") or np_.get("peer_account")
entry = {"src": src, "dst": r["dstAddr"], "port": to_int(r.get("dstPort")),
"flows": to_int(r.get("flows")), "bytes": to_int(r.get("bytes")),
"first_seen": r.get("first_seen"), "last_seen": r.get("last_seen"),
"src_eni": eni.get("NetworkInterfaceId"), "src_description": eni.get("Description"),
"src_instance": (eni.get("Attachment") or {}).get("InstanceId"),
"src_account": src_owner, "path": np_.get("path"), "via": np_.get("via")}
entry["external"] = bool(src_owner and src_owner != acct) or np_.get("path") in (
"vpn_or_direct_connect", "internet")
entry["unidentified"] = not src_owner and not eni
s.setdefault("flow_sources", []).append(entry)
def metric_max(m, prefix):
vals = [v.get("max") for k, v in (m or {}).items() if k.split("[")[0] == prefix and v.get("max") is not None]
return max(vals) if vals else None
def metric_min(m, prefix):
vals = [d.get("Minimum") for k, v in (m or {}).items() if k.split("[")[0] == prefix
for d in v.get("daily", {}).values() if d.get("Minimum") is not None]
return min(vals) if vals else None
def server_health(s):
"""Capacity, saturation and permission checks on a file server, from its config and metrics."""
f = []
m = s.get("_metrics") or {}
t = s["server_type"]
# Capacity and saturation on FSx and File Cache
if t.startswith("FSx"):
pct = metric_max(m, "StorageCapacityUtilization")
cap = (s.get("storage_capacity_gib") or 0) * 2 ** 30
free = metric_min(m, "FreeStorageCapacity") or metric_min(m, "FreeDataStorageCapacity")
if pct is None and free is not None and cap:
pct = round(100 * (cap - free) / cap, 1)
used, total = metric_max(m, "StorageUsed"), metric_max(m, "StorageCapacity")
if pct is None and used and total:
pct = round(100 * used / total, 1)
s["storage_utilization_pct_max"] = pct
if pct is not None and pct >= 90:
f.append(("HIGH" if pct >= 95 else "MEDIUM", "storage_capacity_high", f"storage reached {pct}% in the metric window"))
hot = sorted({k.split("[")[0] for k, v in m.items() if "Utilization" in k and not k.startswith("StorageCapacity")
and (v.get("max") or 0) >= 90})
if hot:
f.append(("MEDIUM", "fsx_saturation", f"{', '.join(hot)} reached 90% or more (latency rises near these limits)"))
if t == "EFS":
lo, hi = s.get("metric_burst_credit_min"), s.get("metric_burst_credit_max")
if s.get("throughput_mode") == "bursting" and lo is not None and hi and lo < 0.1 * hi:
f.append(("MEDIUM", "efs_burst_credits_low", "burst credits fell below 10% of their peak; throughput "
"is throttled to the baseline when they run out"))
if not s.get("fs_policy"):
f.append(("MEDIUM", "efs_no_file_system_policy", "no file system policy, so EFS's default applies: any "
"client that can reach a mount target can mount, write and act as root"))
elif "aws:SecureTransport" not in s["fs_policy"]:
f.append(("LOW", "efs_policy_no_tls_requirement", "file system policy does not require TLS (aws:SecureTransport)"))
if t == "StorageGateway":
for name, label in (("CachePercentUsed", "cache used"), ("CachePercentDirty", "cache not yet uploaded to S3"),
("UploadBufferPercentUsed", "upload buffer used")):
v = metric_max(m, name)
if v is not None and v >= 90:
f.append(("MEDIUM", "gateway_cache_pressure", f"{label} reached {round(v)}%"))
for ex in s.get("nfs_exports", []):
if ex.get("protocol") == "SMB":
if ex.get("Authentication") == "GuestAccess":
f.append(("HIGH", "smb_guest_access", f"SMB share {ex.get('FileShareName') or ex.get('Path')} allows guest access"))
continue
if any(c in ("0.0.0.0/0", "::/0") for c in (ex.get("ClientList") or [])):
f.append(("HIGH", "nfs_share_open_clients", f"NFS share {ex.get('Path')} allows clients from 0.0.0.0/0"))
if ex.get("Squash") == "NoSquash":
f.append(("MEDIUM", "nfs_share_no_squash", f"NFS share {ex.get('Path')} gives remote root full root access"))
if t == "FSx-OPENZFS":
for v in s.get("nfs_exports", []):
for ex in v.get("zfs_nfs_exports") or []:
for cc in ex.get("ClientConfigurations", []):
opts = ",".join(cc.get("Options", []))
if cc.get("Clients") == "*" and "rw" in opts.split(","):
f.append(("MEDIUM", "zfs_export_any_client_rw", f"volume {v.get('name')} is exported read write to any client"))
if "no_root_squash" in opts:
f.append(("MEDIUM", "zfs_export_no_root_squash", f"volume {v.get('name')} exports with no_root_squash to {cc.get('Clients')}"))
if t == "FSx-ONTAP":
s["note"] = "ONTAP export policies and share ACLs live in ONTAP; check them with the ONTAP CLI or REST API"
# Cross region replication is not a problem, but it is part of the picture
for r in (s.get("replication_destinations") or []):
f.append(("INFO", "replicated_cross_region", f"replicated to {r}"))
return f
def server_findings(servers, account, ssm_ran, idx=None):
out = []
for s in servers:
f = []
own = s.get("account", account)
if s.get("fs_policy"):
accts, open_any, org = policy_cross_account(s["fs_policy"], own)
s["policy_cross_account_principals"] = sorted(accts)
if accts:
f.append(("MEDIUM", "policy_grants_other_accounts",
f"file system policy allows principals in {', '.join(sorted(accts))}"))
if open_any:
f.append(("HIGH", "policy_open_to_any_principal",
"file system policy allows Principal * without an account or organisation condition"))
if org:
f.append(("INFO", "policy_org_scoped", "file system policy allows any principal in the organisation"))
if idx is not None and s.get("endpoints"):
hits = sg_cross_account(idx, s)
s["sg_external_rules"] = [h[1] for h in hits]
for kind, detail, _who in hits:
f.append(({"open": "HIGH", "account": "MEDIUM"}.get(kind, "INFO"),
{"open": "sg_open_to_world", "account": "sg_allows_other_account",
"outside_vpc": "sg_allows_outside_vpc"}[kind], detail))
ext = [x for x in s.get("flow_sources", []) if x["external"]]
unk = [x for x in s.get("flow_sources", []) if x["unidentified"]]
s["flow_external_sources"] = len(ext)
if ext:
who = sorted({f"{x['src']}({x.get('src_account') or x.get('path')})" for x in ext})
f.append(("MEDIUM", "external_clients_observed",
f"flow logs show {len(who)} clients outside account {own}: {', '.join(who[:10])}"))
if unk:
f.append(("INFO", "unidentified_clients_observed",
f"flow logs show {len(unk)} source addresses with no visible ENI or route attribution"))
st = (s.get("state") or "").lower()
if st and st not in ("available", "running"):
f.append(("HIGH", "server_not_available", f"state is {s.get('state')}"))
if s["server_type"] == "EFS":
if not s.get("encrypted"):
f.append(("MEDIUM", "efs_unencrypted", "not encrypted at rest"))
if s.get("backup_policy") != "ENABLED" and not s.get("recovery_points"):
f.append(("MEDIUM", "no_backups", "no EFS backup policy and no AWS Backup recovery points"))
if s.get("mount_target_count", 0) == 0:
f.append(("MEDIUM", "no_mount_targets", "file system has no mount targets"))
elif s.get("mount_targets_available", 0) < s.get("mount_target_count", 0):
f.append(("HIGH", "mount_target_unhealthy", "one or more mount targets not available"))
if s.get("metric_client_connections_max") in (0, 0.0, None) and s.get("metrics_collected"):
f.append(("INFO", "efs_idle", "no client connections recorded in the metric window"))
if (s.get("metric_percent_io_limit_max") or 0) >= 90:
f.append(("MEDIUM", "efs_io_limit", "PercentIOLimit reached 90%+ (General Purpose IO ceiling)"))
f += server_health(s)
if ssm_ran and not s.get("clients") and s["server_type"] != "StorageGateway":
f.append(("INFO", "no_observed_clients", "no SSM managed host was seen mounting this "
"(unmanaged hosts, containers, Lambda and on premises clients are not visible)"))
s["findings"] = [code for _, code, _ in f]
for sev, code, detail in f:
out.append({"severity": sev, "code": code, "account": own, "region": s["region"],
"instance_id": "", "mountpoint": "", "resource": s["id"], "detail": detail})
return out
# --------------------------------------------------------------------------------------
# Output
# --------------------------------------------------------------------------------------
def write_csv(path, rows, skip=("_detail",)):
if not rows:
with open(path, "w") as f:
f.write("")
return
cols = []
for r in rows:
for k in r:
if k not in cols and k not in skip:
cols.append(k)
with open(path, "w", newline="") as f:
w = csv.DictWriter(f, fieldnames=cols, extrasaction="ignore")
w.writeheader()
for r in rows:
w.writerow({k: (json.dumps(v, default=jdefault) if isinstance(v, (dict, list)) else v)
for k, v in r.items() if k not in skip})
def fmt_bytes(n):
if n is None:
return ""
n = float(n)
for u in ("B", "KiB", "MiB", "GiB", "TiB", "PiB"):
if abs(n) < 1024:
return f"{n:.1f} {u}"
n /= 1024
return f"{n:.1f} EiB"
def write_report(path, inv):
L = []
s = inv["summary"]
accts = inv.get("accounts") or [inv["account"]]
L.append(f"# File mount reconnaissance: {'account ' + accts[0] if len(accts) == 1 else str(len(accts)) + ' accounts'}")
L.append("")
L.append(f"Collected {inv['collected_at']} across {len(inv['regions'])} regions; "
f"metric window {inv['metric_window_days']} days; client collection "
f"{'enabled' if inv['ssm_enabled'] else 'disabled (run with --ssm for client side data)'}.")
L.append("")
L.append("## 1. Summary")
L.append("")
for k, v in s.items():
L.append(f"* {k.replace('_', ' ')}: {v}")
L.append("")
L.append("## 2. Servers")
L.append("")
L.append("| Type | Id | Name | Region | State | Size | Endpoints | Observed clients | Findings |")
L.append("|---|---|---|---|---|---|---|---|---|")
for sv in inv["servers"]:
size = fmt_bytes(sv.get("size_bytes")) if sv.get("size_bytes") is not None else (
f"{sv['storage_capacity_gib']} GiB" if sv.get("storage_capacity_gib") else "")
L.append(f"| {sv['server_type']} | {sv['id']} | {sv.get('name') or ''} | {sv['region']} | "
f"{sv.get('state') or ''} | {size} | {len(sv.get('endpoints', []))} | "
f"{len(sv.get('clients', []))} | {', '.join(sv.get('findings', []))} |")
L.append("")
if inv["mounts"]:
L.append("## 3. Client mounts")
L.append("")
L.append("| Host | Mountpoint | Protocol | Target | Vers | Encrypted | Probe | Cross account | Used | Ops | Avg write | Findings |")
L.append("|---|---|---|---|---|---|---|---|---|---|---|---|")
for m in inv["mounts"]:
L.append(f"| {m.get('instance_name') or m['instance_id']} | {m['mountpoint']} | {m.get('protocol')} | "
f"{m.get('server_type')} {m.get('server_id') or m.get('resolved_host')} | "
f"{m.get('nfs_vers') or m.get('smb_vers') or ''} | "
f"{'yes' if m.get('tls') else 'no'} | {m.get('probe_diagnosis') or m.get('responsive')} | "
f"{m.get('cross_account') or ''} | {fmt_bytes(m.get('used_bytes'))} | "
f"{m.get('ops_total') or ''} | {fmt_bytes(m.get('avg_write_bytes'))} | {m.get('findings', '')} |")
L.append("")
xm = [m for m in inv["mounts"] if m.get("cross_account") in ("yes", "possible")]
xs = [sv for sv in inv["servers"] if sv.get("sg_external_rules") or sv.get("policy_cross_account_principals")
or sv.get("flow_external_sources")]
xl = [x for x in inv["other_consumers"].get("lambda_mounts", []) if x.get("cross_account")]
L.append("## 4. Cross account")
L.append("")
if not (xm or xs or xl):
L.append("No cross account mounts or access paths were found in the data collected.")
L.append("")
if xm:
L.append("Mounts whose target is, or may be, in another account:")
L.append("")
L.append("| Client account | Host | Mountpoint | Target | Target IP | Path | Peer account | Status | Evidence |")
L.append("|---|---|---|---|---|---|---|---|---|")
for m in xm:
L.append(f"| {m.get('client_account')} | {m.get('instance_name') or m['instance_id']} | {m['mountpoint']} | "
f"{m.get('server_type')} {m.get('server_id') or m.get('resolved_host')} | {m.get('target_ip') or ''} | "
f"{m.get('network_path')} {m.get('network_via') or ''} | {m.get('peer_account') or m.get('server_account') or ''} | "
f"{m.get('cross_account')} | {m.get('cross_account_evidence') or ''} |")
L.append("")
if xl:
for x in xl:
L.append(f"* Lambda {x['FunctionName']} ({x['account']}/{x['region']}) mounts EFS in {', '.join(x['cross_account'])}")
L.append("")
if xs:
L.append("File servers reachable from, or used by, other accounts:")
L.append("")
for sv in xs:
L.append(f"* {sv['server_type']} {sv['id']} ({sv.get('account')}/{sv['region']})")
for a in sv.get("policy_cross_account_principals") or []:
L.append(f" * file system policy grants account {a}")
for r in sv.get("sg_external_rules") or []:
L.append(f" * security group: {r}")
for x in [y for y in sv.get("flow_sources", []) if y["external"]][:20]:
L.append(f" * flow logs: {x['src']} ({x.get('src_account') or x.get('path')}) to {x['dst']}:{x['port']}, "
f"{x['flows']} flow records, {fmt_bytes(x.get('bytes'))}, last seen {x.get('last_seen')}")
L.append("")
other = [("Lambda functions with file systems", "lambda_mounts"),
("ECS task definitions with EFS volumes", "ecs_efs_volumes"),
("EKS clusters (EFS CSI add on)", "eks_efs_csi")]
n = 5
for title, key in other:
items = inv["other_consumers"].get(key, [])
if items:
L.append(f"## {n}. {title}")
L.append("")
for it in items:
L.append(f"* {it.get('region')}: {it.get('FunctionName') or it.get('family') or it.get('cluster')}"
f" {json.dumps(it.get('FileSystemConfigs') or [v.get('efsVolumeConfiguration') for v in it.get('volumes', [])] or (it.get('efs_csi_addon') or {}).get('status'), default=str)}")
L.append("")
n += 1
L.append(f"## {n}. Findings")
L.append("")
order = {"HIGH": 0, "MEDIUM": 1, "LOW": 2, "INFO": 3}
for f in sorted(inv["findings"], key=lambda x: (order.get(x["severity"], 9), x["code"])):
where = f"{f['instance_id']}:{f['mountpoint']}" if f["instance_id"] else f["resource"]
L.append(f"* **{f['severity']}** `{f['code']}` {f['region']} {where}: {f['detail']}")
L.append("")
n += 1
L.append(f"## {n}. Coverage gaps")
L.append("")
for g in inv["coverage_gaps"]:
L.append(f"* {g}")
if inv["errors"]:
L.append(f"* {len(inv['errors'])} API calls failed (see inventory.json errors); most common:")
counts = defaultdict(int)
for e in inv["errors"]:
counts[f"{e['service']}.{e['operation']}:{e['code']}"] += 1
for k, v in sorted(counts.items(), key=lambda kv: -kv[1])[:15]:
L.append(f" * {k} x{v}")
with open(path, "w") as f:
f.write("\n".join(L) + "\n")
# --------------------------------------------------------------------------------------
# Main
# --------------------------------------------------------------------------------------
def account_sessions(base, base_account, args, rec):
"""The base session plus, with --org or --accounts, a session per member account via --role-name."""
sessions, gaps = {base_account: base}, []
targets = list(args.accounts or [])
if args.org:
org = base.client("organizations", config=BOTO_CFG)
targets += [a["Id"] for a in rec.paginate("global", "organizations", org, "list_accounts", "Accounts")
if a.get("Status", "ACTIVE") == "ACTIVE"]
targets = [t for t in dict.fromkeys(targets) if t != base_account]
if targets and not args.role_name:
sys.exit("--org / --accounts need --role-name (a read only role that exists in each member account)")
sts = base.client("sts", config=BOTO_CFG)
for acct in targets:
kw = dict(RoleArn=f"arn:aws:iam::{acct}:role/{args.role_name}", RoleSessionName="file-mount-recon")
if args.external_id:
kw["ExternalId"] = args.external_id
cred = (rec.call("global", "sts", "assume_role", sts.assume_role, **kw, default={}) or {}).get("Credentials")
if not cred:
gaps.append(f"account {acct}: could not assume {args.role_name}; its resources and hosts were not scanned")
continue
sessions[acct] = boto3.Session(aws_access_key_id=cred["AccessKeyId"],
aws_secret_access_key=cred["SecretAccessKey"],
aws_session_token=cred["SessionToken"])
return sessions, gaps
def main():
ap = argparse.ArgumentParser(description="Full recon of file storage and file mounts (NFS, SMB, Lustre, FUSE) in AWS")
ap.add_argument("--profile")
ap.add_argument("--regions", nargs="*", help="default: every region enabled for the account")
ap.add_argument("--out", default=f"file-mount-recon-{NOW.strftime('%Y%m%dT%H%M%SZ')}")
ap.add_argument("--days", type=int, default=14, help="CloudWatch metric and flow log window in days")
ap.add_argument("--no-metrics", action="store_true")
ap.add_argument("--ssm", action="store_true",
help="run the read only collector on SSM managed Linux and Windows instances")
ap.add_argument("--ssm-bucket", help="S3 bucket for full SSM output (avoids the 24k inline limit)")
ap.add_argument("--ssm-timeout", type=int, default=120)
ap.add_argument("--instance-ids", nargs="*", help="limit client collection to these instances")
ap.add_argument("--max-ecs-families", type=int, default=500)
ap.add_argument("--workers", type=int, default=8)
ap.add_argument("--org", action="store_true", help="scan every active account in the AWS Organization")
ap.add_argument("--accounts", nargs="*", help="scan these member accounts as well as the current one")
ap.add_argument("--role-name", help="role to assume in member accounts for --org / --accounts")
ap.add_argument("--external-id", help="external ID for the member account role, if required")
ap.add_argument("--flow-logs", action="store_true",
help="query VPC Flow Logs in CloudWatch Logs for clients of each file server (Logs Insights charges apply)")
args = ap.parse_args()
base = boto3.Session(profile_name=args.profile) if args.profile else boto3.Session()
rec = Recorder()
ident = rec.call("global", "sts", "get_caller_identity",
base.client("sts", config=BOTO_CFG).get_caller_identity, default={}) or {}
account = ident.get("Account", "unknown")
if args.regions:
regions = args.regions
else:
home = base.region_name or "us-east-1"
regions = sorted(r["RegionName"] for r in (rec.call(
"global", "ec2", "describe_regions",
base.client("ec2", region_name=home, config=BOTO_CFG).describe_regions,
default={}) or {}).get("Regions", []))
if not regions:
sys.exit("could not determine regions; pass --regions")
os.makedirs(args.out, exist_ok=True)
sessions, gaps = account_sessions(base, account, args, rec)
log(f"caller {ident.get('Arn')}; {len(sessions)} accounts; {len(regions)} regions; output {args.out}")
regions_data = []
with ThreadPoolExecutor(max_workers=args.workers) as ex:
futs = {ex.submit(collect_region, sess, r, args, rec, args.out, acct): (acct, r)
for acct, sess in sessions.items() for r in regions}
for fu in as_completed(futs):
try:
regions_data.append(fu.result())
except Exception as e: # noqa: BLE001
rec.errors.append({"region": "/".join(futs[fu]), "service": "*", "operation": "collect_region",
"code": type(e).__name__, "message": str(e)[:500]})
regions_data.sort(key=lambda r: (r["account"], r["region"]))
# Client side
ssm_results = {}
for R in regions_data:
acct, region = R["account"], R["region"]
managed = {i["InstanceId"]: i for i in R["ssm_managed"]}
running = [i for i in R["ec2_instances"] if (i.get("State") or {}).get("Name") == "running"]
unmanaged = [i["InstanceId"] for i in running if i["InstanceId"] not in managed]
offline = [i for i, m in managed.items() if m.get("PingStatus") != "Online"]
if unmanaged:
gaps.append(f"{acct}/{region}: {len(unmanaged)} running EC2 instances are not SSM managed, so their "
f"mounts are invisible: {' '.join(unmanaged[:20])}{' ...' if len(unmanaged) > 20 else ''}")
if offline:
gaps.append(f"{acct}/{region}: {len(offline)} SSM managed instances are not Online")
if args.ssm:
key = f"{acct}:{region}"
targets = [i for i, m in managed.items()
if m.get("PingStatus") == "Online" and m.get("PlatformType") == "Linux"
and m.get("ResourceType", "EC2Instance") in ("EC2Instance", "ManagedInstance")
and (not args.instance_ids or i in args.instance_ids)]
ssm_results[key] = ssm_collect(sessions[acct], region, targets, args, rec, args.out, account=acct)
wtargets = [i for i, m in managed.items()
if m.get("PingStatus") == "Online" and m.get("PlatformType") == "Windows"
and (not args.instance_ids or i in args.instance_ids)]
ssm_results[key].update(ssm_collect(sessions[acct], region, wtargets, args, rec, args.out,
platform="Windows", account=acct))
if not args.ssm:
gaps.append("client side collection not run (use --ssm); mounts, options and usage per host are unknown")
if len(sessions) == 1:
gaps.append("single account run: mounts into other accounts are detected from network paths, unknown "
"file system IDs and buckets, but the exact resource is only known with --org or --accounts")
if not args.flow_logs:
gaps.append("flow logs not queried (use --flow-logs); clients in other accounts mounting this account's "
"storage are inferred from policies and security groups only")
gaps.append("on premises clients of on premises servers, and Kubernetes PersistentVolumes, are not covered")
for R in regions_data:
if R.get("s3files_note"):
gaps.append(f"{R['region']}: {R['s3files_note']}")
break
idx = Index()
for acct, sess in sessions.items():
for bkt in rec.paginate("global", "s3", sess.client("s3", config=BOTO_CFG), "list_buckets", "Buckets"):
idx.buckets[bkt["Name"]] = {"region": bkt.get("BucketRegion"), "created": bkt.get("CreationDate"),
"account": acct}
servers = build_servers(regions_data, idx)
if args.flow_logs:
flow_log_sources(sessions, regions_data, servers, idx, args, rec)
mounts, findings, client_hosts = build_mounts(regions_data, ssm_results, idx, servers, account)
findings.extend(server_findings(servers, account, args.ssm, idx))
for sv in servers:
sv.pop("_metrics", None) # full metrics stay in raw/<account>/<region>/server_side.json
# Self managed file servers discovered on hosts
self_managed = [h for h in client_hosts if h.get("is_nfs_server") or h.get("is_smb_server")]
for h in self_managed:
servers.append({"server_type": "EC2-self-managed", "id": h["instance_id"], "region": h["region"],
"account": h.get("account"),
"protocols": [p for p, flag in (("NFS", h.get("is_nfs_server")),
("SMB", h.get("is_smb_server"))) if flag],
"name": h.get("name"), "state": "running", "endpoints": [], "findings": [],
"clients": [{"instance_id": m["instance_id"], "mountpoint": m["mountpoint"],
"client_account": m.get("client_account")}
for m in mounts if m.get("server_id") == h["instance_id"]],
"exports": h.get("nfs_server_detail") or h.get("smb_shares_served")})
# Consumers without a host: Lambda access points and ECS volumes in other accounts
other = {"lambda_mounts": [], "ecs_efs_volumes": [], "eks_efs_csi": []}
for R in regions_data:
for k in other:
for x in R[k]:
item = {**x, "region": R["region"], "account": R["account"]}
if k == "lambda_mounts":
owners = {m.group(1) for c in x.get("FileSystemConfigs", [])
for m in [re.search(r":(\d{12}):", c.get("Arn", ""))] if m}
item["cross_account"] = sorted(owners - {R["account"]})
if item["cross_account"]:
findings.append({"severity": "MEDIUM", "code": "cross_account_mount", "account": R["account"],
"region": R["region"], "instance_id": "", "mountpoint": "lambda",
"resource": x["FunctionName"],
"detail": f"Lambda mounts an EFS access point in account {', '.join(item['cross_account'])}"})
if k == "ecs_efs_volumes":
unknown = [v["efsVolumeConfiguration"].get("fileSystemId") for v in x.get("volumes", [])
if v["efsVolumeConfiguration"].get("fileSystemId") not in idx.by_id]
item["unknown_file_systems"] = unknown
if unknown:
findings.append({"severity": "INFO", "code": "possible_cross_account_mount", "account": R["account"],
"region": R["region"], "instance_id": "", "mountpoint": "ecs",
"resource": x.get("family"),
"detail": f"ECS task uses EFS {', '.join(unknown)} not found in any scanned account"})
other[k].append(item)
sev = defaultdict(int)
for f in findings:
sev[f["severity"]] += 1
summary = {
"accounts_scanned": len(sessions),
"regions_scanned": len(regions),
"efs_file_systems": sum(len(R["efs"]) for R in regions_data),
"s3files_file_systems": sum(len(R["s3files"]) for R in regions_data),
"fsx_ontap": sum(1 for R in regions_data for e in R["fsx"] if e["type"] == "ONTAP"),
"fsx_openzfs": sum(1 for R in regions_data for e in R["fsx"] if e["type"] == "OPENZFS"),
"fsx_windows": sum(1 for R in regions_data for e in R["fsx"] if e["type"] == "WINDOWS"),
"fsx_lustre": sum(1 for R in regions_data for e in R["fsx"] if e["type"] == "LUSTRE"),
"fsx_file_caches": sum(len(R["fsx_file_caches"]) for R in regions_data),
"storage_gateway_nfs_shares": sum(len(e["nfs_file_shares"]) for R in regions_data for e in R["storage_gateway"]),
"storage_gateway_smb_shares": sum(len(e.get("smb_file_shares", [])) for R in regions_data for e in R["storage_gateway"]),
"self_managed_file_servers_found": len(self_managed),
"client_hosts_inspected": sum(1 for h in client_hosts if h.get("bundle_ok")),
"client_hosts_failed": sum(1 for h in client_hosts if not h.get("bundle_ok")),
"client_mounts": len(mounts),
**{f"client_mounts_{p.lower()}": sum(1 for m in mounts if m.get("protocol") == p)
for p in sorted({m.get("protocol") for m in mounts if m.get("protocol")})},
"s3_fuse_mounts": sum(1 for m in mounts if str(m.get("server_type", "")).startswith("S3-bucket")),
"unresponsive_mounts": sum(1 for m in mounts if m.get("responsive") is False),
"slow_mounts": sum(1 for m in mounts if m.get("probe_diagnosis") == "slow"),
"mounts_capacity_over_85pct": sum(1 for m in mounts if (m.get("capacity_pct") or 0) >= 85),
"cross_region_mounts": sum(1 for m in mounts if m.get("cross_region")),
"high_latency_mounts": sum(1 for m in mounts if "latency" in (m.get("findings") or "")),
"permission_findings": sum(1 for f in findings if f["code"] in (
"world_writable_root", "password_in_fstab", "credentials_file_exposed", "access_denied",
"efs_no_file_system_policy", "smb_guest_access", "nfs_share_open_clients", "nfs_share_no_squash",
"zfs_export_any_client_rw", "zfs_export_no_root_squash", "policy_open_to_any_principal")),
"unmapped_mounts": sum(1 for m in mounts if m.get("server_type") == "unknown"),
"cross_account_mounts": sum(1 for m in mounts if m.get("cross_account") == "yes")
+ sum(1 for x in other["lambda_mounts"] if x.get("cross_account")),
"possible_cross_account_mounts": sum(1 for m in mounts if m.get("cross_account") == "possible"),
"servers_with_external_access_rules": sum(1 for x in servers if x.get("sg_external_rules")
or x.get("policy_cross_account_principals")),
"servers_with_external_clients_observed": sum(1 for x in servers if x.get("flow_external_sources")),
"lambda_functions_with_fs": len(other["lambda_mounts"]),
"ecs_task_defs_with_efs": len(other["ecs_efs_volumes"]),
"findings_high": sev["HIGH"], "findings_medium": sev["MEDIUM"],
"findings_low": sev["LOW"], "findings_info": sev["INFO"],
"api_errors": len(rec.errors),
}
inv = {"account": account, "accounts": sorted(sessions), "caller": ident.get("Arn"),
"collected_at": NOW.isoformat(), "regions": regions, "metric_window_days": args.days,
"ssm_enabled": args.ssm, "flow_logs_enabled": args.flow_logs, "summary": summary,
"servers": servers, "mounts": mounts, "client_hosts": client_hosts, "other_consumers": other,
"findings": findings, "coverage_gaps": gaps, "errors": rec.errors}
write_json(os.path.join(args.out, "inventory.json"), inv)
write_csv(os.path.join(args.out, "mounts.csv"), mounts)
srv_rows = [{k: v for k, v in s.items() if k not in ("endpoints", "clients", "nfs_exports", "tags", "exports",
"flow_sources", "fs_policy")}
| {"endpoints": len(s.get("endpoints", [])), "observed_clients": len(s.get("clients", []))}
for s in servers]
write_csv(os.path.join(args.out, "servers.csv"), srv_rows)
write_csv(os.path.join(args.out, "findings.csv"), findings)
flows = [{"account": s.get("account"), "region": s["region"], "server_type": s["server_type"],
"server_id": s["id"], **x} for s in servers for x in s.get("flow_sources", [])]
if flows:
write_csv(os.path.join(args.out, "flow_sources.csv"), flows)
write_report(os.path.join(args.out, "report.md"), inv)
log("summary: " + ", ".join(f"{k}={v}" for k, v in summary.items()))
log(f"wrote {args.out}/report.md, inventory.json, mounts.csv, servers.csv, findings.csv"
+ (", flow_sources.csv" if flows else ""))
if __name__ == "__main__":
main()
EOF
chmod +x file_mount_recon.py