mirror of
				https://github.com/go-micro/go-micro.git
				synced 2025-10-30 23:27:41 +02:00 
			
		
		
		
	
		
			
				
	
	
		
			84 lines
		
	
	
		
			1.3 KiB
		
	
	
	
		
			Go
		
	
	
	
	
	
			
		
		
	
	
			84 lines
		
	
	
		
			1.3 KiB
		
	
	
	
		
			Go
		
	
	
	
	
	
| package server
 | |
| 
 | |
| import (
 | |
| 	"context"
 | |
| 	"sync"
 | |
| 
 | |
| 	"github.com/micro/go-micro/codec"
 | |
| )
 | |
| 
 | |
| // Implements the Streamer interface
 | |
| type rpcStream struct {
 | |
| 	sync.RWMutex
 | |
| 	id      string
 | |
| 	closed  bool
 | |
| 	err     error
 | |
| 	request Request
 | |
| 	codec   codec.Codec
 | |
| 	context context.Context
 | |
| }
 | |
| 
 | |
| func (r *rpcStream) Context() context.Context {
 | |
| 	return r.context
 | |
| }
 | |
| 
 | |
| func (r *rpcStream) Request() Request {
 | |
| 	return r.request
 | |
| }
 | |
| 
 | |
| func (r *rpcStream) Send(msg interface{}) error {
 | |
| 	r.Lock()
 | |
| 	defer r.Unlock()
 | |
| 
 | |
| 	resp := codec.Message{
 | |
| 		Target:   r.request.Service(),
 | |
| 		Method:   r.request.Method(),
 | |
| 		Endpoint: r.request.Endpoint(),
 | |
| 		Id:       r.id,
 | |
| 		Type:     codec.Response,
 | |
| 	}
 | |
| 
 | |
| 	if err := r.codec.Write(&resp, msg); err != nil {
 | |
| 		r.err = err
 | |
| 	}
 | |
| 
 | |
| 	return nil
 | |
| }
 | |
| 
 | |
| func (r *rpcStream) Recv(msg interface{}) error {
 | |
| 	r.Lock()
 | |
| 	defer r.Unlock()
 | |
| 
 | |
| 	req := new(codec.Message)
 | |
| 	req.Type = codec.Request
 | |
| 
 | |
| 	if err := r.codec.ReadHeader(req, req.Type); err != nil {
 | |
| 		// discard body
 | |
| 		r.codec.ReadBody(nil)
 | |
| 		r.err = err
 | |
| 		return err
 | |
| 	}
 | |
| 
 | |
| 	// we need to stay up to date with sequence numbers
 | |
| 	r.id = req.Id
 | |
| 	if err := r.codec.ReadBody(msg); err != nil {
 | |
| 		r.err = err
 | |
| 		return err
 | |
| 	}
 | |
| 
 | |
| 	return nil
 | |
| }
 | |
| 
 | |
| func (r *rpcStream) Error() error {
 | |
| 	r.RLock()
 | |
| 	defer r.RUnlock()
 | |
| 	return r.err
 | |
| }
 | |
| 
 | |
| func (r *rpcStream) Close() error {
 | |
| 	r.Lock()
 | |
| 	defer r.Unlock()
 | |
| 	r.closed = true
 | |
| 	return r.codec.Close()
 | |
| }
 |