NivaarExam PrepOfficial exam papers ↗

25-Comp-B10 Distributed Systems · December 2016

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, December 2016. 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, peer-to-peer systems, middleware and client-server architecture (ch. 1–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), and time, coordination, replication and fault tolerance (ch. 14–15, 18).

Check — source parsing artifact. Every question header on this paper is printed as “Question # N.” (a literal hash between the word and the number). Separately, Questions 1 and 7 each print their third sub-part re-using the letter “a.” instead of continuing the alphabet; both are relettered below (a), (b), (c) in the order printed, with no change to content or intent.

Question 7: Principles of Fault Tolerance (20 marks)

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 Implementation Repository — scalability and fault tolerance. In middleware such as CORBA, the Implementation Repository is a registry that records, for each server object implementation, which host and process it runs in (or should be started in) and how to activate it if it is not currently running. From a scalability viewpoint, this enables lazy, on-demand activation: an object implementation need not run continuously consuming resources — the Implementation Repository can start its server process the first time a request actually arrives, and multiple copies of a popular service's implementation can be registered across different hosts so requests can be routed to, or new instances activated on, whichever host currently has capacity, letting overall capacity grow by registering more hosts rather than being fixed at deployment time. From a fault tolerance viewpoint, the Implementation Repository provides a layer of indirection between an object's logical identity (its object reference) and its current physical location/process: if a server process crashes, the Implementation Repository can detect this and automatically restart it (or activate a replacement instance), and clients continue reaching the (new) implementation through the same unchanged object reference, with no client needing to know a failure and recovery occurred at all.

(b) Read-only replica managers and gossip performance. In a gossip system, a replica manager (RM) that is designated read-only never accepts front-end update operations directly — only reads. This removes two costs a read-write RM must otherwise bear: it never has to withhold applying a gossiped update pending proof that no conflicting update arrived directly from one of its own clients (since none ever do), so it can apply gossiped updates as soon as they are known stable, with less coordination overhead; and it never originates updates of its own that must be gossiped out to every other RM, so it adds no new update traffic to the system, only propagating what it receives. A read-only RM can also answer a local read immediately from its own cached, eventually-consistent copy without coordinating with any other RM at read time. Concentrating writes at a smaller set of read-write RMs while replicating widely via cheap, read-only RMs lets the system scale its read capacity (by adding more read-only replicas close to clients) largely independently of its write throughput and gossip traffic — a scalability trade a design with every RM accepting writes cannot offer as cleanly.

(c) Reconciling two directory replicas after disconnected operation. The scheme below follows the standard directory-reconciliation approach used in Coda/AFS-style disconnected operation: each replica keeps 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 →