NivaarExam PrepOfficial exam papers ↗

25-Comp-B10 Distributed Systems · May 2014

Question 6 of 6: Distributed Algorithms and Fault Tolerance

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 6: Distributed Algorithms and Fault Tolerance

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) Adapting the central server mutual-exclusion algorithm to tolerate client crashes. In the standard central server algorithm, a designated server holds a single token/permission and a FIFO queue of pending requests: a client wanting the critical section sends a REQUEST; if the token is free, the server replies with it immediately, otherwise the request is queued; on leaving the critical section the client sends a RELEASE, and the server hands the token to the next queued request (if any). Mutual exclusion is trivially enforced because the server only ever grants the token to one client at a time.

To tolerate a client crashing in any state, the server additionally tracks, for the client currently holding the token, its identity, and the reliable failure detector continuously monitors every client. Two cases arise: (i) a queued (non-holding) client crashes — the failure detector's notification lets the server simply drop that client's entry from the queue; since it never held the token, no invariant is at risk. (ii) the token-holding client crashes — the failure detector's confirmed crash notification is treated by the server exactly as if that client had sent an explicit RELEASE, and the token is granted to the next queued request. This is the crux of the adaptation: without it, a genuinely crashed holder would never send RELEASE and every other client would wait forever.

Is the resultant system fault tolerant? It restores liveness under client crashes — the system as a whole keeps making progress even if the current token holder dies, which the unmodified algorithm cannot do. It is fault tolerant only with respect to client failures, and explicitly not to a failure of the server itself: the server is a single point of failure for the whole scheme (the queue and current-holder state exist only there), so a server crash halts mutual exclusion for every client, and the question's own assumption that "the server is correct" is precisely what puts server failure out of scope here — a genuinely fault-tolerant design would need the server's own state replicated or agreed on via a separate consensus mechanism.

What if the current token holder is wrongly suspected to have failed? This is the real danger the crash-tolerant adaptation introduces. Acting on a false suspicion, the server grants the token to the next queued client while the original holder — which is, in fact, still alive and has no idea anything happened — continues to believe it holds exclusive access and keeps executing inside its critical section. The result is a genuine mutual exclusion violation (a safety failure, not merely a performance one): two processes now simultaneously believe they hold exclusive access, and both may concurrently modify the shared resource. This is the fundamental reason "reliable failure detector" in the distributed-systems literature typically only means eventually accurate (it can still produce a transient false positive, e.g. during a temporary network partition that looks indistinguishable from a crash) rather than never wrong; a robust design must therefore add a second line of defence beyond trusting the detector alone — e.g. tagging each token grant with a monotonically increasing epoch/fencing number that the shared resource itself checks, so any operation arriving from a previously-evicted (falsely suspected) holder using a stale epoch is rejected outright, rather than relying purely on the failure detector never making a mistake.

(b) The gossip architecture. In a gossip-style replicated service, a collection of replica managers (RMs) each accept updates directly from nearby clients (rather than routing every update through one master) and periodically exchange ("gossip") messages with one another, comparing what updates each has seen and forwarding any the other is missing. Over time, updates propagate transitively RM-to-RM until every replica has (eventually) applied every update — this is the classic eventual-consistency model, chosen because it tolerates individual RMs or links being temporarily unreachable (an update simply propagates once connectivity returns) without requiring all RMs to be reachable synchronously the way strict, immediate consensus would.

Gossip architecture — three replica managers gossip gossip gossip RM1 (replica ts, value ts) RM2 (replica ts, value ts) RM3 (replica ts, value ts)
Fig. Q6(b) — each replica manager accepts client updates directly and periodically gossips with its peers, propagating updates it has that a peer lacks; every RM tracks both its own replica timestamp and the value timestamp of the data it currently returns to clients.

Each RM keeps two distinct vector timestamps because they answer two different questions. The value timestamp records exactly which updates are reflected in the data value the RM would currently return to a client — it is what a client compares against its own prior timestamp to enforce session guarantees such as "read your writes" or "monotonic reads." The replica timestamp records every update the RM has received (from a client directly, or via gossip from a peer), whether or not it has actually been applied to the value yet — because gossip can deliver an update the RM is not yet ready to apply (e.g. it is causally dependent on an earlier update this RM has not received yet, and updates must be applied in a causally consistent order). Without the replica timestamp, an RM could not correctly tell a gossiping peer "here is everything I have seen so far" (which drives what the peer sends it next), and without the value timestamp it could not correctly tell a client what guarantee its returned data actually satisfies — conflating the two would either apply updates out of causal order (corrupting the value) or under-report to peers what has already been received (causing redundant re-gossip of updates already known).

(c) Reconciling two directory replicas after disconnected operation. The scheme below follows the standard directory-reconciliation approach (as used in Coda/AFS-style disconnected operation): each replica keeps a version vector and a log of the operations (create, delete, rename) it performed while disconnected; on reconnection, the two logs are compared, non-conflicting operations from each side are replayed onto the merged directory, and any true conflict (e.g. the same name created differently on both sides) is flagged for the user rather than silently guessed at.

function reconcile(dirA, logA, dirB, logB):
    // logX is the ordered list of directory operations (create, delete,
    // rename) performed on replica X since the two last agreed.
    merged = union(dirA.entries_unchanged_since_split, dirB.entries_unchanged_since_split)
    conflicts = []

    for op in logA:
        if op.target not touched by any operation in logB:
            apply(op, merged)                  // safe: only A touched this name
        else:
            conflicts.append((op, matching_op_in(logB, op.target)))

    for op in logB:
        if op.target not touched by any operation in logA:
            apply(op, merged)                  // safe: only B touched this name

    for (opA, opB) in conflicts:
        if opA == opB:                         // e.g. both deleted the same file
            apply(opA, merged)                 // identical outcome: no real conflict
        else if opA.type == CREATE and opB.type == CREATE and opA.name == opB.name:
            rename_one(opA, merged, suffix=".A")   // same name, different content:
            rename_one(opB, merged, suffix=".B")   // keep both, let the user resolve
        else:
            flag_for_manual_resolution(opA, opB)   // e.g. rename vs. delete of the
                                                    // same entry: no safe automatic merge
    return merged, conflicts

The key design point is that the algorithm only ever applies an operation automatically when it can prove no other operation touched the same name during the disconnection; anything else is either resolved by a purely mechanical rule (identical operations, or two independent creates renamed to coexist) or surfaced explicitly, since an automatic but wrong merge on a shared directory is far more damaging than asking the user to resolve a rare genuine conflict.

Back to the paper →