Interprocess Communication Between Remote Process
The system described in the previous section easily generalizes to allow interprocess communication between processes at geographically different locations as, for example, within a computer network.
Consider first a simple configuration of processes distributed around the points of a star. At each point of the star there is an autonomous time-sharing system. A rather large, smart computer system, called the Network Controller, exists at the center of the star. No processes can run in this center system, but rather it should be thought of as an extension of the monitor of each time-sharing system in the network.
It should be obvious to the reader that if the Network Controller is able to perform the operations SEND, RECEIVE, SEND FROM ANY, RECEIVE ANY, and UNIQUE and that if all of the monitors in all of the time-sharing systems in the network do not perform these operations themselves but rather ask the Network Controller to perform these operations for them, then we have solved the problem of interprocess communication between remote processes. We have no further change to make.
The reason everything continues to work when we postulate the existence of the Network Controller is that the Network Controller can keep track of which RECEIVEs have been executed and which SENDs have been executed and match them up just as the monitor did in the model time-sharing system. A networkwide port numbering scheme is also possible with the Network Controller knowing where (i.e., at which site) a particular port is at a particular time.
Next, consider a more complex network in which there is no common center point making it necessary to distribute the functions performed by the Network Controller among the network nodes. In the rest of this section I will show that it is possible to efficiently and conveniently distribute the functions performed by the star Network Controller among the many network sites and still enable general interprocess communication between remote processes.
Some changes must be made to each of the four SEND/RECEIVE operations described above to adapt them for use in a distributed network. To RECEIVE is added a parameter specifying a site to which the RECEIVE is to be sent. To SEND FROM ANY and SEND is added a site to send the SEND to although this is normally the local site. Both RECEIVE and RECEIVE ANY have added the provision for obtain the source site of any received message. Thus, when a RECEIVE is executed, the RECEIVE is sent to the site specified, possibly a remote site. Concurrently a SEND is sent to the same site, normally the local site of the process executing the SEND. At this site, called the rendezvous site, the RECEIVE is matched with the proper SEND and the message transmission is allowed to take place to the site from whence the RECEIVE came.
A RECEIVE ANY never leaves its originating site and therein lies the necessity for SEND FROM ANY. It must be possible to send a message to a RECEIVE ANY port and not have the message blocked waiting for RECEIVE at the sending site. Of course, it would be possible to construct the system so the SEND/RECEIVE rendezvous takes place at the RECEIVE site and eliminate the SEND FROM ANY operation, but in my judgment the ability to block a normal SEND transmission at the source site more than makes up for the added complexity.
Somewhere at each site a rendezvous table is kept. This table contains an entry for each unmatched SEND or RECEIVE received at that site and also an entry for all RECEIVE ANYs given at that site. A matching SEND/RECEIVE pair is cleared from the table as soon as the match takes place or perhaps when the transmission is complete. As in the similar table kept in the model time-sharing system, SEND and RECEIVE entries are timed out if unmatched for too long and the originator is notified. RECEIVE ANY entries are cleared from the table when a fulfilling message arrives.
The final change necessary to distribute the Network Controller functions is to give each site a portion of the unique numbers to distribute via its UNIQUE operation. I'll discuss this topic further below.
To make it clear to the reader how the distributed Network Controller works, an example follows. The details of what process picks port numbers, etc. are only exemplary and are not a standard specified as part of the system.
Suppose there are two sites in the network: K and L. Process A at site K wishes to communicate with process B at site L. Process B has a RECEIVE ANY pending at port M.
SITE K SITE L
________ ________
/ \ / \
/ \ / \
/ \ / \
| Process A | | Process B |
| | | |
| | | |
\ / \ /
\ / \ port M /
\________/ \____^___/
|
RECEIVE ANY
Process A, fortunately, knows of the existence of port M at site L and sends messages using the SEND FROM ANY operation from port N to port M. The message contains two port numbers and instructions for process B to SEND messages to process A to port P from port Q. Site K's site number is appended to this message along with the message's SEND port N.
SITE K SITE L
________ ________
/ \ / \
/ \ / \
/ \ / \
| Process A | | Process B |
| | | |
| | | |
\ / \ /
\ port N /--->SEND FROM --->\ port M /
\________/ ANY \________/
to port M, site L
containing K, N, P, & Q
Process A now executes a RECEIVE at port P from port Q. Process A specifies the rendezvous site to be site L.
SITE K SITE L
________ R ________
/ \ e / \
/ \ n T/ \
/ \ d a \
| | e b Process B |
| Process A | z l |
| | v e |
\ / o \ /
\ port P / RECEIVE ---> u \ /
\________/ MESSAGE s \________/
to site L
containing P, Q, & K
A RECEIVE message is sent from site K to site L and is entered in the rendezvous table at site L. At some other time, process B executes a SEND to port P from port Q specifying site L as the rendezvous site.
SITE K SITE L
________ R ________
/ \ e / \
/ \ n T/ \
/ \ d a \
| | e b Process B |
| Process A | z l |
| | v e |
\ / o \ /
\ port P / u <--- port Q /
\________/ SEND s \________/
to site L
containing P & Q
A rendezvous is made, the rendezvous table entry is cleared, and the transmission to port P at site K takes place. The SEND site number (and conceivably the SEND port number) are appended to the messages of the transmission for the edification of the receiving process.
SITE K SITE L
________ ________
/ \ / \
/ \ / \
/ \ / \
| Process A | | Process B |
| | | |
| | | |
\ port P / \ port Q /
\ / <---- transmission <---- \ /
\________/ to port T, site K \________/
containing data and L
Process B may simultaneously wish to execute a RECEIVE from port N at port M.
Note that there is only one important control message in this system which moves between sites, the type of message that is called a Host/Host protocol message in [3]. This control message is the RECEIVE message. There are two other possible intersite control messages: an error message to the originating site when a RECEIVE or SEND is timed out, and the SEND message in the rare case when the rendezvous site is not the SEND site.
Of course there must also be a standard format for messages between ports. For example, the following:
+-----------------+ +-----------------+ +-----------------+
| rendezvous site | | destination site| | source site |
+-----------------+ +-----------------+ +-----------------+
| RECEIVE port | | RECEIVE port | | RECEIVE port |
+-----------------+ +-----------------+ +-----------------+
| SEND port | | SEND port | | SEND port |
+-----------------+ +-----------------+ +-----------------+
| | | source port | | |
| | +-----------------+ | |
| | | | | |
| | | | | |
| | | | | |
| | | | | |
| | | | | |
| data | | data | | data |
| | | | | |
| | | | | |
| | | | | |
| | | | | |
| | | | | |
| | | | | |
+-----------------+ +-----------------+ +-----------------+
transmitted transmitted received
by SEND by Network by RECEIVE
process Controller process
Note: for a SEND FROM ANY message, the rendezvous site is the destination site.
In the model time-sharing system it was possible to pass a port from process to process. This is still possible with a distributed Network Controller. [The reader unconvinced of the utility of port passing is directed to read the section on reconnection in [11].]
Remember that for a message to be sent from one process to another, a SEND to port M from port N and a RECEIVE at port M from port N must rendezvous, normally at the SEND site. Both processes keep track of where they think the rendezvous site is and supply this site as a parameter of appropriate operations. The RECEIVE process thinks it is the SEND site and the SEND process normally thinks it is the SEND site also. Since once a SEND and a RECEIVE rendezvous, the transmission is sent to the source of the RECEIVE and the entry in the rendezvous table is cleared and must be set up again for each further transmission from N to M, it is easy for a RECEIVE port to be moved. If a process sends both the port numbers and the rendezvous site number to a new process at some other site which executes a RECEIVE using these same old port numbers and rendezvous site specification, the SENDer never knows the RECEIVEr has moved. It is slightly harder for a SEND port to move. However, if it does, the pair of port numbers that has been being used for a SEND and the original rendezvous site number are passed to the new site. The process at the new SEND site specifies the old rendezvous site with the first SEND from the new site. The RECEIVE process will also still think the rendezvous site is the old site, so the SEND and RECEIVE will meet at the old site. When they meet, the entry in the table at that site is cleared, the rendezvous site number for the SEND message is changed to the site which originated the SEND message and both the SEND and RECEIVE messages are sent to the new SEND site just as if they had been destined for there in the first place. The SEND and RECEIVE then meet again at the new rendezvous site and transmission may continue as if the port had never moved. Since all transmissions contain the source site number, further RECEIVEs will be sent to the new rendezvous site. It is possible to discover that this special manipulation must take place because a SEND message is received at a site which did not originate the SEND message. Everything is so easily changed because there are no permanent connections to break and move as in the once proposed reconnection scheme for the ARPA network [10][11] that is, connections only exist fleetingly in the system described here and can therefore be remade between any pair of processes which at any time happen to know each other's port numbers and have some clue where they each are.
Of course, all of this could have been done by the processes sending messages back and forth announcing any potential moves and the new site numbers.