Reliability issues in distributed systems

Christos N. Nikolaou · 1982

Depending upon the philosophy used to implement fault-tolerant systems, one can distinguish two classes of algorithms: reconfiguration algorithms and fault masking algorithms. The precise statement and analysis of the problems and the underlying assumptions associated with these classes of algorithms is the subject of this dissertation. The first part of the thesis investigates strategies for dynamically reconfiguring shared memory multiprocessor systems that are subject to common memory faults and unpredictable processor deaths. These strategies aim at determining a communication page, i.e., a page of common memory that can be used by a group of processors for storing crucial common resources such as global locks for synchronization and global data structures for voting algorithms. To insure system reliability, the reconfiguration strategies must be distributed so that each processor independently arrives at exactly the same choice. This type of reconfiguration strategy is currently used in the STAGE operating system on the PLURIBUS multiprocessor {24}. We analyze the weak points of the PLURIBUS algorithm and examine alternative strategies satisfying optimization criteria such as maximization of the number of processors and the number of common memory pages in the reconfigured system. We also present a general distributed algorithm which enables the processors in such a system to exchange the local information that is needed to reach a consensus on system reconfiguration. In the second part of the thesis, we deal with fault-masking algorithms as applied to the development of network protocols with an underlying communication medium that may reorder, duplicate or lose messages. In chapter (3) we present a simple network, whose communication medium is assumed to be reliable, and develop a strategy for the remote submission and processing of requests. We also show how to formally specify and verify the network behavior. In the final chapter we describe a more complex network model where the communication medium is no longer assumed to be reliable. We then show that despite the reordering, duplication or loss of messages, all requests are eventually processed exactly once at the remote site and that responses are received in the right order at their submission site.

Read the paper · More papers on PaperTik