Teleforking a Process onto a Different Computer
thume.ca
thume.ca
Reinventing LOCUS also has a strong heritage. Bell Lab's Plan 9, for example, did so in part in the late 1980s.
While never a breakout commercial success, tele-forking and its slightly more advanced cousins machine-to-machine process migration and cluster-wide process pools intrigued some of the best minds in distributed computing for 20+ years.
Unfortunately "it's complicated" to implement well, especially when you try to tele-spawn and manage resources beyond compute cycles (network connections, files, file handles, ...) that are important to scale up the idea.
Aren't all of these resources namespaced/containerized in modern Linux? This should make it feasible to checkpoint and restore them on the same machine (via, e.g. the CRIU patchset) and true location-independence is not that much harder. One of the hardest parts (not even implemented in plan9, AFAICT) is distributed shared memory (allowing for sharing a single virtual address space across cluster nodes), but even that AIUI has some research-level implementations.
distributed shared memory though does deserve a new look now that baseline network speeds are 100x what the max was when it was first investigated. unfortunately temporal consistency is still going to be a major factor. some workloads will run great with some heuristics, and some won't. you'll almost certainly need to migrate threads along with pages to try to keep them running co-local without exhausting per-node resources.
DragonflyBSD has been experimenting with this concept for several years now. At the moment it's mostly only useful for snapshotting the kernel for debugging purposes. But there's no reason why it couldn't be extended to transport the state of the computer to another computer and resume execution!
> This module provides a common interface for offloading an IO action to remote executors.
> It uses StaticPointers language extension and distributed-closure library for serializing closures to run remotely. This blog post[1] is a good introduction for those.
> In short, if you you need a Closure a:
> One important constraint when using this library is that it assumes the remote environment is capable of executing the exact same binary. On most cases, this requires your host environment to be Linux.
https://hackage.haskell.org/package/distributed-fork-0.0.1.3...
[1] https://blog.ocharles.org.uk/blog/guest-posts/2014-12-23-sta...
These days most of the work for this is handled by the CRIU project (https://criu.org/), so the "hard" work has been done.
(Then the paper went on to make some other point with a more impressive higher margin of error, p-hacking upward, 0% error was too low.)
Everything that wasn't immediately successful is tried again in 15 years. This is why the old farts on your project are grumpy about 'new' things that... aren't.
Funny, I think probably all of the above would be trivial now on AWS or Cloudflare.
I think the 'all apps will evolve until they can send email' axiom is the wrong one. Everything turns into content management.
And I don't think you'll have to wait until 2030. I suspect 2027 ±2 years.
Plan 9 did not support forking to remote machines and was not a single system image system. Its support for running a program on a remote machine was a completely userspace program.
If you say Plan 9 was an example of teleforking, then so is ssh. That might be a reasonable categorization, but both Plan 9 and ssh are very different from LOCUS and other SSI systems.
I remember vividly how they demonstrated an unbreakable process. They had a computer running a process and no matter what happened to that computer, the next one would flawlessly continue the process down to the cycle, with no change or corruption or skipping a beat.
It may very well be that this is actually not very difficult, but it seemed difficult and impressive.
Perhaps more shocking were ultra high resolution radar screens, some 3 generations ahead of anything I had seen in the consumer space, showing an incredible visualization of the air space, live. Showing exactly which plane is where, the model/type, age, fuel on board, hostile/friendly, all of it.
They even had a "situation room" with a holodeck chair in the middle, full of controls. The entire room was covered in wall-size screens basically showing the air space of the entire country, being live analyzed.
Sounds very 2022, not 1998.
With things like this possible, and possible for a long time now, it honestly makes me cringe at how absolutely terrible (slow, inefficient, unreliable, complex, not robust) the typical modern software stack is.
You wouldnt want to write ultrafast software for your idea to deliver cats back from the vets with a network of teenagers om bicycle. Imagine, years to write your website in C++/ASM to end up failing the first week.
What s embarrassing is how much software cost to build in the 80s.
(Source: I used to work at VMWare and one of the implementers of FT was a cycling buddy.)
Unfortunately, as cool as VMWare FT and QEMU COLO are, you need high-bandwith and low-latency links between the nodes and still they increase latency and reduce performance a lot. And in 99.99% of cases just replicating the disk and restarting the VMs in case of a node failure is good enough.
Way back in the early 90s, inspired by this and a few other things, a coworker of mine (one of only three true geniuses I've ever met) came up with an even more powerful way (called vArray) to spread processes across any arbitrary number of POSIX machines. One of the things vArray could do was make remote memory look local (if painfully slow). He was the company's Cray expert, and developed it for the major oil company we worked for to run jobs too large to fit into the Cray's memory. Since some seismic deconvolution vectorizes very well, with fairly infrequent memory access, it worked better than you might expect, even over the slow networks of the day. Obviously, it hammered the network, but it was quite a sight to see every Sun, DEC, IBM, and HP Unix machine on the company's network light up at once, all coordinating to process jobs that were probably bigger than anything run outside the Black world of the Spooks... Fun days.
On the other hand, combining multiple 64-core Threadrippers into a single-system-image cluster has a certain appeal to it.
I'd like to do that with my spare ARM SBC's.
Well I had my gossip network setup on planet lab and I could tell it to become anything, so I told it to become a content distribution networks and used a gossip algorithm to make copies of the same file on all machine on the network and wrote a paper about it and everybody was happy.>
I miss Joe, not that I ever met him, but his attitude and good humour are inspiring.
https://git.sr.ht/~jfred/goblins-test/tree/master/item/unive...
I don't know what its current status is, but the HPC-ish Bproc system has/had an rfork [1]. Probably the most HPC-oriented SSI system, Kerrighed died, as did the Plan 9-ish xcpu, though that was a bit different.
1. https://www.penguinsolutions.com/computing/documentation/scy...
Citation needed, as they say, for “run way more efficiently”, particularly as the conventional wisdom says shared memory in a single process (e.g. OpenMP).
“Acknowledgements: ... NUMA and Amdahl’s Law, for holding OpenMP back and keeping MPI-only competitive in spite of the ridiculous cost of Send-Recv within a shared-memory domain.” — Jeff Hammond, ‘MPI+MPI’
edit: Got it, after a while: https://github.com/Overv/outrun
It's very experimental, but I got as far as being able to freeze a JavaScript interpreter (compiled to WebAssembly) and restore it later. The code is here: https://github.com/drifting-in-space/wasmbox
In almost all other languages there's just no way to know if a closure is holding on to a file descriptor.
Critics may say the Haskell closures could contain `unsafePerformIO`, but as the saying goes: now you have two problems.
/pedantic
For comparison, two old distributed lexical scope systems were Cardelli's Obliq and Kelsey's(?) Kali Scheme. From what I remember, not like remote forking, though.
Erlang has, however, excellent support for distributed computing using its own kind of processes.
http://www.iro.umontreal.ca/~feeley/papers/GermainFeeleyMonn...
(Trivia: in Unix VI, fork() returns both PIDs, and the C library stub arranges to return 0 to the child and the child's PID to the parent.)
Sprite called this "process migration".
All live migration systems basically follow the same pattern https://cloud.google.com/compute/docs/instances/live-migrati...
And as others have mentioned upthread, Mosix was an extension to Linux that also implemented fork() in such a way that the child could be local or remote with file handles retained across the cluster. We had a Linux lan party when I was in college and managed to scrape together a Mosix cluster across a bunch of machines.
Ignore the uncreative pedants.