25-Comp-B10 Distributed Systems · Undated paper
Nivaar worked solution (AI-drafted; not reviewed by a licensed engineer)
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) A load-balancing scheme across a set of computers. Requirements met. The primary user requirement is response time: no single machine should be visibly overloaded while others sit idle, so interactive/short jobs are not starved behind a backlog that a better placement decision would have avoided. The system requirement is efficient use of the whole resource pool (throughput), and avoiding cascading failure from one host being driven into thrashing while capacity elsewhere goes unused. Suitable applications. Load balancing suits long-running, CPU-bound, independent jobs whose run time is large compared with the cost of starting them remotely — batch simulations, rendering, scientific computation, parallel compilation (e.g. a distributed make) — and which do not depend on resources local to the originating machine. It is poorly suited to short interactive commands (the placement overhead exceeds the run time), I/O-bound processes tied to local files or devices, and processes bound to a local display or user session. Measuring load. A simple, cheap, and reasonably effective load index is the length of the CPU run queue (number of runnable processes waiting), sampled periodically; more accurate but more expensive alternatives combine CPU utilization, free memory, and recent I/O activity into a weighted index. Perfect accuracy is neither achievable (load is constantly changing between the measurement and the placement decision) nor necessary — a coarse, cheap index that is directionally right most of the time out-performs an expensive, precise one that is stale by the time it is used, so practical schemes deliberately trade precision for low measurement overhead and freshness. Monitoring and placement. Each host periodically exchanges (or, in a sender-initiated scheme, broadcasts on request) its current load index with a subset of other hosts; when a new process is to be created, the originating host compares its own load against the indices it holds and places the process locally if it is not itself overloaded, or transfers placement to a lightly-loaded peer otherwise (receiver-initiated variants instead have lightly-loaded hosts poll for work). Care is needed to avoid instability — if every host reacts to the same stale "host X is idle" information simultaneously, many processes can pile onto X at once (the classic thrashing/oscillation failure of naive load sharing), which is why real schemes randomize or rate-limit how often placement decisions are re-evaluated and use a probe of a small random subset of hosts rather than a full broadcast on every decision.
(b) A shared region for a process to read kernel-written data, and its synchronization. The kernel maps a page of memory read-only into a process's own address space (in addition to the kernel's normal mapping of the same physical page); the kernel writes data it wants to expose (e.g. the current time, or scheduling/CPU-accounting statistics) directly into that page, and any user process that has it mapped can read the data with an ordinary memory load — no system call, and therefore no context switch or kernel-crossing overhead, is needed for the read at all (this is exactly the technique behind Linux's vDSO-based `gettimeofday`/`clock_gettime`). Synchronization is the crux of doing this safely: the kernel may update the shared data (e.g. several fields of a timestamp) while a user process is in the middle of reading it, and neither side can use an ordinary blocking lock (the kernel must never block on a user-space process, and a user-space process should not have to trap into the kernel just to acquire a lock for a memory read that was supposed to avoid a trap). The standard solution is a sequence-number (seqlock) scheme: the kernel increments a shared sequence counter to an odd value before it begins updating the data, writes the new data, then increments the counter again to an even value when done; a reader reads the counter, reads the data, then re-reads the counter, and retries the whole read if the counter changed or was odd at any point — guaranteeing the reader never observes a value torn mid-update, without either side ever blocking.
(c) The two main approaches to kernel architecture. Monolithic kernel. All core OS services — process/thread scheduling, virtual memory management, file systems, device drivers, networking — run together in a single, large, privileged address space. Communication between these services is by ordinary in-kernel function calls, which is fast, but a bug or crash in any one component (e.g. a faulty device driver) can bring down the entire kernel, and the codebase is large and more difficult to maintain or extend safely. Microkernel. The kernel itself is stripped down to a minimal set of mechanisms — inter-process communication (IPC), basic scheduling, and minimal address-space management — while everything else that a monolithic kernel would run in privileged mode (file systems, device drivers, higher-level OS services) instead runs as separate, unprivileged user-space processes/servers that communicate via the kernel's IPC mechanism. This gives much stronger fault isolation (a crashed file-system server can potentially be restarted without taking the whole system down) and a smaller, more easily verified trusted computing base, at the cost of the performance overhead of IPC/context-switching between servers for operations a monolithic kernel would have handled with a single direct function call. Distributed systems support is often easier to add cleanly to a microkernel design, since remote services already communicate through the same IPC abstraction used for local ones — the same mechanism extends transparently across a network.
(d) Operating-system-level virtualization. OS-level virtualization (containerization) lets multiple isolated user-space instances — containers — run on top of a single shared OS kernel, each with its own view of the file system, process namespace, network stack and resource limits, without each needing its own full guest operating system the way hardware/hypervisor-based virtual machines do. Main goals: strong-enough isolation between co-located workloads (so one container's processes/files/network are invisible to another's), efficient resource utilization (many lightweight containers can run per physical host, since there is no per-instance guest-kernel overhead), and portability (a container packages an application with its exact dependencies, so "it works on my machine" behaviour is reproduced identically wherever the container image is run). Challenges: because all containers share one kernel, isolation is fundamentally weaker than a hypervisor's hardware-enforced boundary — a kernel vulnerability can potentially be exploited to escape a container and affect the host or sibling containers; every container is also implicitly tied to the host kernel's version and feature set (no independent guest-kernel choice); and containers sharing the same physical CPU/memory/disk-I/O can suffer "noisy neighbour" resource contention if cgroup-style resource limits are not carefully configured. Use cases: packaging and deploying microservices consistently across development, test and production; CI/CD pipelines that need fast, disposable, reproducible build/test environments; and dense multi-tenant cloud deployment where the overhead of a full guest OS per tenant would be wasteful, provided the isolation trade-off above is acceptable for the workload.