NOTE

Designing an RPC Framework

A historical note on RPC transport, protocol framing, serialization, stubs/reflection, service discovery, resilience, async calls, and graceful shutdown.

System DesignCreated Updated 3 min readhistorical

This is a historical learning note and may contain outdated or incomplete understanding.

1. What Is RPC?

  • Remote Procedure Call.
  • A computer communication protocol.
    • The protocol lets a client call a function on a server as if it were a local function, hiding network-communication details.
    • It uses a client/server model; the classic implementation is request-response.

1.1. Local Function Call vs RPC

  • Local function call: no network is involved.
  • RPC: the network is involved.

2. Why Do We Need an RPC Framework?

To hide low-level technical details unrelated to business code during RPC calls, such as network communication, serialization/deserialization, and addressing.

3. Overall RPC Call Process

  1. The client process calls the client stub, which resides within the client’s address space.
  2. The client stub packs the parameters into a message. This is called marshaling. The client stub then executes a system call, such as sendto, to send the message.
  3. The kernel sends the message to the remote server machine.
  4. The server stub receives the message from the kernel.
  5. The server stub unmarshals the parameters.
  6. The server stub calls the desired procedure.
  7. The server process executes the procedure and returns the result to the server stub.
  8. The server stub marshals the results into a message and passes the message to the kernel.
  9. The kernel sends the message to the client machine.
  10. The client stub receives the message from the kernel.
  11. The client stub unmarshals the results and passes them to the caller.

4. How to Implement RPC

Call a function on a machine over the network, pass parameters, and obtain the return value.

4.1. How to Communicate with a Remote Machine

4.1.1. Network Transport

Network model: I/O multiplexing + zero-copy.
TCP vs UDP.
Big-endian vs little-endian.

Implementation:

// Transport will use TLV protocol
type Transport struct {
	conn net.Conn // Conn is a generic stream-oriented network connection.
}

// NewTransport creates a Transport
func NewTransport(conn net.Conn) *Transport {
	return &Transport{conn}
}

// Send TLV encoded data over the network
func (t *Transport) Send(data []byte) error {
	// we will need 4 more byte then the len of data
	// as TLV header is 4bytes and in this header
	// we will encode how much byte of data
	// we are sending for this request.
	buf := make([]byte, 4+len(data))
	binary.BigEndian.PutUint32(buf[:4], uint32(len(data)))
	copy(buf[4:], data)
	_, err := t.conn.Write(buf)
	if err != nil {
		return err
	}
	return nil
}

// Read TLV Data sent over the wire
func (t *Transport) Read() ([]byte, error) {
	header := make([]byte, 4)
	_, err := io.ReadFull(t.conn, header)
	if err != nil {
		return nil, err
	}
	dataLen := binary.BigEndian.Uint32(header)
	data := make([]byte, dataLen)
	_, err = io.ReadFull(t.conn, data)
	if err != nil {
		return nil, err
	}
	return data, nil
}

4.1.2. Protocol Design (Decode/Encode)

This refers to the application-layer protocol, such as HTTP.

Why is a protocol needed? – To solve the TCP packet-sticking/framing problem.

How should a protocol be designed? Refer to HTTP: header + body. The body has variable length and its length is read from the header. For extensibility, the header is also designed to be variable-length, split into a fixed section and protocol-header contents.

Implementation:

// header
length int32
// body
// RPCdata represents the serializing format of structured data
type RPCdata struct {
	Name string        // name of the function
	Args []interface{} // request's or response's body expect error.
	Err  string        // Error any executing remote server
}

4.1.3. Serialization and Deserialization (Serialize/Deserialize)

Why is serialization needed? – Data transmitted over a network must be binary, while the caller’s input and output parameters are objects.

How to implement serialization? JSON, Hessian, Protobuf. Consider security > generality > compatibility > performance > efficiency > space overhead.

  • Marshaling:
    • Client: parameters -> message.
    • Server: return value -> message.
    • message means a byte array in some format.
  • Unmarshaling:
    • Client: message -> return value.
    • Server: message -> parameters.
    • message means a byte array in some format.

Implementation:

// Encode The RPCdata in binary format which can
// be sent over the network.
func Encode(data RPCdata) ([]byte, error) {
	var buf bytes.Buffer
	encoder := gob.NewEncoder(&buf)
	if err := encoder.Encode(data); err != nil {
		return nil, err
	}
	return buf.Bytes(), nil
}

// Decode the binary data into the Go RPC struct
func Decode(b []byte) (RPCdata, error) {
	buf := bytes.NewBuffer(b)
	decoder := gob.NewDecoder(buf)
	var data RPCdata
	if err := decoder.Decode(&data); err != nil {
		return RPCdata{}, err
	}
	return data, nil
}

4.2. How to Call a Remote Function Like a Local Function

Program to interfaces.

  • Static: Stub Generation.
  • Dynamic: Reflection.

4.2.1. Stub Generation

  • IDL compiler: reads the IDL and automatically generates the client stub and server stub.

4.2.2. Reflection

  • Server:

    // Execute the given function if present
    func (s *RPCServer) Execute(req dataserial.RPCdata) dataserial.RPCdata {
    	// get method by name
    	f, ok := s.funcs[req.Name]
    	if !ok {
    		// since method is not present
    		e := fmt.Sprintf("func %s not Registered", req.Name)
    		log.Println(e)
    		return dataserial.RPCdata{Name: req.Name, Args: nil, Err: e}
    	}
    
    	log.Printf("func %s is called\n", req.Name)
    	// unpack request arguments
    	inArgs := make([]reflect.Value, len(req.Args))
    	for i := range req.Args {
    		inArgs[i] = reflect.ValueOf(req.Args[i])
    	}
    
    	// invoke requested method
    	out := f.Call(inArgs)
    	// now since we have followed the function signature style where last argument will be an error
    	// so we will pack the response arguments expect error.
    	resArgs := make([]interface{}, len(out)-1)
    	for i := 0; i < len(out)-1; i++ {
    		// Interface returns the constant value stored in v as an interface{}.
    		resArgs[i] = out[i].Interface()
    	}
    
    	// pack error argument
    	var er string
    	if _, ok := out[len(out)-1].Interface().(error); ok {
    		// convert the error into error string value
    		er = out[len(out)-1].Interface().(error).Error()
    	}
    	return dataserial.RPCdata{Name: req.Name, Args: resArgs, Err: er}
    }
  • Client:

    func (c *Client) CallRPC(rpcName string, fPtr interface{}) {
    	container := reflect.ValueOf(fPtr).Elem()
    	f := func(req []reflect.Value) []reflect.Value {
    		cReqTransport := transport.NewTransport(c.conn)
    		errorHandler := func(err error) []reflect.Value {
    			outArgs := make([]reflect.Value, container.Type().NumOut())
    			for i := 0; i < len(outArgs)-1; i++ {
    				outArgs[i] = reflect.Zero(container.Type().Out(i))
    			}
    			outArgs[len(outArgs)-1] = reflect.ValueOf(&err).Elem()
    			return outArgs
    		}
    
    		// Process input parameters
    		inArgs := make([]interface{}, 0, len(req))
    		for _, arg := range req {
    			inArgs = append(inArgs, arg.Interface())
    		}
    
    		// ReqRPC
    		reqRPC := dataserial.RPCdata{Name: rpcName, Args: inArgs}
    		b, err := dataserial.Encode(reqRPC)
    		if err != nil {
    			panic(err)
    		}
    		err = cReqTransport.Send(b)
    		if err != nil {
    			return errorHandler(err)
    		}
    		// receive response from server
    		rsp, err := cReqTransport.Read()
    		if err != nil { // local network error or decode error
    			return errorHandler(err)
    		}
    		rspDecode, _ := dataserial.Decode(rsp)
    		if rspDecode.Err != "" { // remote server error
    			return errorHandler(errors.New(rspDecode.Err))
    		}
    
    		if len(rspDecode.Args) == 0 {
    			rspDecode.Args = make([]interface{}, container.Type().NumOut())
    		}
    		// unpack response arguments
    		numOut := container.Type().NumOut()
    		outArgs := make([]reflect.Value, numOut)
    		for i := 0; i < numOut; i++ {
    			if i != numOut-1 { // unpack arguments (except error)
    				if rspDecode.Args[i] == nil { // if argument is nil (gob will ignore "Zero" in transmission), set "Zero" value
    					outArgs[i] = reflect.Zero(container.Type().Out(i))
    				} else {
    					outArgs[i] = reflect.ValueOf(rspDecode.Args[i])
    				}
    			} else { // unpack error argument
    				outArgs[i] = reflect.Zero(container.Type().Out(i))
    			}
    		}
    
    		return outArgs
    	}
    	container.Set(reflect.MakeFunc(container.Type(), f))
    }

4.3. How to Design a Cross-Language Interface

4.3.1. IDL

IDL is used to specify service interfaces: a set of functions implemented by the server and called by the client.

4.4. How to Find the Server

4.4.1. Naming Service (Service Registration)

Designing a Service Registry.md

4.4.2. Addressing (Service Discovery)

Designing a Service Registry.md

4.4.3. Load Balancing

Designing a Load-Balancing Component.md

4.5. How to Optimize

4.5.1. Timeout and Retry

Designing Timeout and Retry

4.5.2. Circuit Breaking

Designing a Circuit-Breaker System

4.5.3. Traffic Control

Designing a Rate-Limiting System

4.5.4. Security

4.5.5. Exceptions

Encapsulate exception types and distinguish business exceptions from network exceptions.

If it is a network exception, retrying is possible.

If it is a business exception, decide whether to retry according to the situation.

4.5.6. Synchronous / Asynchronous

  • Why asynchronous RPC is needed: if server-side business logic is time-consuming and the CPU spends most of its time waiting instead of computing, resulting in low CPU utilization, asynchronous RPC can improve throughput.
  • How to implement asynchronous RPC:
    • Caller-side async can be implemented with a Future. The caller initiates an asynchronous request and obtains a Future from the request context, then calls get on the Future to obtain the result. If the business logic calls multiple other services at the same time, Futures can reduce overall business-logic latency and improve throughput.
    • Server-side async needs a callback mechanism. Business logic can process asynchronously and then call a callback interface provided by the RPC framework to notify the caller of the final result asynchronously.

4.6. Graceful Shutdown

  • Why graceful shutdown is needed:
    • New requests: when the server restarts or shuts down, callers cannot quickly perceive it.
    • Accepted requests: some requests may not have finished processing.
  • What graceful shutdown is:
    • Handle accepted requests and new requests during shutdown.
  • How to perform graceful shutdown:
    1. Unfinished requests: continue processing them. To avoid a request never finishing and preventing shutdown forever, add timeout control.
    2. New requests: reject them.

4.7. Code Architecture Design

4.7.1. Layered Design

5. References

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub