Replacing NFS With S3 Without Changing Your Code: Automating Amazon S3 Files

Replacing NFS With S3 Without Changing Your Code: Automating Amazon S3 Files

👁8views

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.

Article Summary
  • 1.
    What it is
    Amazon S3 Files is a managed NFS file system backed by an S3 bucket, and this post explains how it works, where the no code changes claim holds and where it does not. It also provides scripts to inventory mounts, provision resources, migrate data, cut over and validate.
  • 2.
    Why it matters
    It argues that swapping an NFS mount for S3 can be a configuration change instead of a refactor, but only if you treat that claim as a hypothesis and test for subtle failures under load.
  • 3.
    Key takeaway
    The S3 API ignores file system locks, so any process writing directly to the bucket while an NFS client holds a lock will win, and that requires application level concurrency control.
~31 min read

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:

  1. 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.
  2. 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.
  3. 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.
  4. 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 nconnect mount 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.
  5. 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 dd and cp tests. 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.
  6. 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 mounttargetip option 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 FileMountReconReadOnly

The 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:

  1. s3files-provision.sh creates 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.
  2. s3files-client-prep.sh runs on each client host, installs a recent enough amazon-efs-utils, and checks that TCP 2049 to the mount target is actually reachable before anything tries to mount.
  3. s3files-migrate.sh has four modes: preflight scans the existing share for features S3 Files cannot represent, bulk copies the data while the application is still running, final does the last pass through the mount during the cutover window and verifies content and metadata, and verify-sync waits until every file is visible in the bucket.
  4. s3files-cutover.sh repoints 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.
  5. s3files-validate.sh runs 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.env

I 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.sh

The 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.sh

If 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.sh

Run 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.sh

The _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.sh

A 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:

  1. Run file_mount_recon.py --ssm and work through the findings, paying particular attention to the S3 Files signals and to any mounts of the share you did not know about.
  2. Fill in s3files.env, then run s3files-provision.sh once from a pipeline or an administrative host.
  3. Run s3files-client-prep.sh on every host that mounts the share, ideally baked into your AMI or bootstrap so new instances come up ready.
  4. Run s3files-migrate.sh preflight and resolve every failed gate; then run s3files-migrate.sh bulk while the application is live, as many times as you like.
  5. In the cutover window, stop the application, run s3files-migrate.sh final and s3files-migrate.sh verify-sync, then run s3files-cutover.sh on each host.
  6. Start the application and run s3files-validate.sh local and the two host modes, then your own tests. If something is wrong before the application has written anything, s3files-cutover.sh rollback is a simple swap back to the untouched old share; after it has written, use REVERSE_COPY=1 so 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

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

Leave a comment

Your email address will not be published. Your first comment is held for approval.