NivaarExam PrepOfficial exam papers ↗

25-Comp-B10 Distributed Systems · May 2013

Question 7 of 7: Principles of 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 2013. 3 hours, closed book, non-programmable calculator only. Candidates were instructed to answer any five of the seven questions, all carrying equal weight and mostly requiring essay-format answers; all seven 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), security (ch. 11), distributed file systems (ch. 12–12.4, AFS/NFS), and time, coordination, replication and fault tolerance (ch. 14–15, 18).

Check — sub-part lettering. Both sub-parts of Questions 1, 3 and 5 are lettered “a.” in the paper's numbering. Each of those three questions is answered below as two genuinely distinct sub-parts, relettered (a) and (b) in the order printed; content and marks weight are unaffected.

Question 7: Principles of 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) 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. Q7(a) — 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).

(b) 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 →