NivaarExam PrepOfficial exam papers ↗

25-Comp-B10 Distributed Systems · May 2014

Question 5 of 6: Distributed File Systems

Nivaar worked solution (AI-drafted; not reviewed by a licensed engineer)

Notes on this paper

98-Comp-B10 Distributed Systems — National Examinations, May 2014. 3 hours, closed book, non-programmable calculator only. Candidates were instructed to answer any five of the six questions, all carrying equal weight and mostly requiring essay-format answers; all six are answered below as a complete study resource.

Reference texts: Coulouris, Dollimore, Kindberg & Blair, Distributed Systems: Concepts and Design (5th ed.) — system models and client-server architecture (ch. 2), interprocess communication and the request-reply protocol (ch. 4–5), operating system support for distributed systems (ch. 7), distributed file systems (ch. 12–12.4, AFS/NFS), and time, coordination, replication and fault tolerance (ch. 14–15, 18).

Question 5: Distributed File Systems

Question text not reproduced: the examination questions are © Engineers and Geoscientists BC. Open the official past paper (linked at the top of this page) to read the question, then follow the worked solution below.

(a) Three key design issues for distributed file systems. 1. Naming and transparency. Clients should be able to name and access remote files using the same syntax as local files, ideally without knowing which server actually holds them (location transparency) and without needing to know if a file has moved (location independence). Getting this wrong forces applications to hard-code server identities, defeating the whole point of a shared, uniform file namespace. 2. Caching strategy and consistency. Because network access is orders of magnitude slower than local disk/memory access, clients must cache file data locally to be usable at all — but caching immediately raises the question of how (and how often) a client's cached copy is kept consistent with a file that another client may be concurrently modifying, trading off staleness against server/network load for every design choice made. 3. Fault tolerance and availability. A distributed file system must keep working (or degrade gracefully) in the presence of individual server or network failures, typically via some combination of replication (multiple copies of data on independent servers) and client-side resilience to a temporarily unreachable server, rather than an outage on one server making every client's work stop.

(b) Porting an NFS server to a non-UNIX operating system. Because NFS is a stateless, wire-protocol-defined service, its implementer on a foreign OS faces four main decisions. (1) File handle design. Every NFS operation references a file by an opaque file handle that must let the server directly locate the file with no server-side session state — classic UNIX NFS packs a filesystem ID, an inode number and a generation number into it. A non-UNIX OS that has no inode-number concept must invent its own persistent, unique per-file identifier that (i) is stable for the file's whole lifetime (surviving renames) and (ii) is never reused after the file is deleted, adding an explicit generation/version count if the underlying OS would otherwise recycle identifiers — without this, a stale handle held by a client could silently come to reference a completely different, later file. (2) Attribute mapping. The NFS wire protocol carries UNIX-flavoured attributes (owner/group uid-gid, permission mode bits, device/inode-style metadata); the implementer must map these, possibly lossily, onto whatever native permission/metadata model the target OS actually has, and decide how to degrade gracefully for any UNIX attribute the target OS cannot represent at all. (3) Naming and directory semantics. NFS assumes a hierarchical directory tree with case-sensitive, arbitrary-byte-string names and UNIX-style hard links; a target OS that is case-insensitive, has a flat namespace, or lacks hard links needs an explicit emulation or rejection policy for the operations that don't translate directly. (4) Concurrency and locking. Because core NFS itself keeps no server-side open/close state, file locking is a bolt-on side-protocol (the Network Lock Manager); the implementer must map the target OS's own native locking primitives onto NLM's expected semantics.

For any underlying filing system to be suitable for hosting an NFS server at all, it must in turn obey several constraints: it must expose (or let the NFS layer construct) a persistent, collision-free unique identifier per file that survives renames and is never silently reused after deletion, since this is exactly what a valid file handle depends on; it must support direct lookup of a file given that identifier (not merely lookup by re-walking a path from the root), since NFS operations reference files by handle and re-deriving a path on every call would be both slow and semantically fragile if the file has since been moved; it must present a hierarchical namespace whose naming rules are at least as expressive as what NFS clients expect (arbitrary names, "." and ".." conventions); and, because NFS's RPC layer relies on at-least-once delivery with client-side retries of lost replies, every operation the filing system exposes must be safely repeatable (idempotent) — e.g. a "create" that fails cleanly (rather than duplicating) when retried against a file it already successfully created the first time.

(c) AFS vs. NFS — stability and scalability. Stability/consistency model: classic NFS (v2/v3) is essentially stateless and relies on clients periodically re-validating cached data against the server (a time-based, "check on open" or polling consistency), which means the server does no bookkeeping of who is caching what but pays a steady stream of validation requests from every active client, and different clients can transiently see different (stale) data between validations. AFS instead uses whole-file caching with server-issued callbacks: when a client caches a file, the server promises to notify (callback) that client if the file changes, so the client can trust its cache is valid until told otherwise, without needing to keep re-asking the server. This shifts work from "many clients constantly polling" to "the server does a small amount of bookkeeping and notifies rarely," which is inherently more stable and scalable under many, mostly-read clients, at the cost of the server needing to track callback state per client per cached file. AFS scalability limits: even with servers added as required, AFS's callback bookkeeping itself grows with the number of (client, cached-file) pairs a server must track, and a server recovering from a crash must re-establish (or invalidate) every outstanding callback it had promised, which becomes a heavier recovery burden as the client population grows; AFS's cell-based administrative partitioning also means very large deployments must be explicitly divided into cells, which is an organizational rather than a purely technical scaling limit. Recent developments: NFSv4 closed much of the original gap by adopting AFS-style delegations (a server-granted, revocable promise similar to a callback) plus a stronger, compound-RPC, more stateful protocol design, and it is now the more actively maintained/adopted standard; Coda (built directly on AFS's model) added disconnected operation, allowing clients to keep working from cache during a total server/network outage and reconcile changes afterward, addressing availability rather than raw throughput scalability.