ing
This commit is contained in:
parent
a471858025
commit
bb465d3628
|
@ -6,3 +6,4 @@ import:
|
||||||
- package: git.loafle.net/commons_go/local_socket.git
|
- package: git.loafle.net/commons_go/local_socket.git
|
||||||
- package: git.loafle.net/commons_go/server
|
- package: git.loafle.net/commons_go/server
|
||||||
- package: git.loafle.net/commons_go/rpc
|
- package: git.loafle.net/commons_go/rpc
|
||||||
|
- package: gopkg.in/natefinch/npipe.v2
|
||||||
|
|
15
main.go
15
main.go
|
@ -8,8 +8,7 @@ import (
|
||||||
"os/signal"
|
"os/signal"
|
||||||
"syscall"
|
"syscall"
|
||||||
|
|
||||||
cRPC "git.loafle.net/commons_go/rpc"
|
"git.loafle.net/commons_go/logging"
|
||||||
"git.loafle.net/overflow/overflow_discovery/rpc"
|
|
||||||
"git.loafle.net/overflow/overflow_discovery/server"
|
"git.loafle.net/overflow/overflow_discovery/server"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
@ -24,19 +23,16 @@ func init() {
|
||||||
}
|
}
|
||||||
|
|
||||||
func main() {
|
func main() {
|
||||||
registry := cRPC.NewRegistry()
|
defer logging.Logger().Sync()
|
||||||
if err := registry.RegisterService(new(rpc.DiscoveryService), ""); nil != err {
|
|
||||||
panic(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
s := server.New(*sockFile, registry)
|
s := server.New(*sockFile)
|
||||||
|
|
||||||
stop := make(chan os.Signal)
|
stop := make(chan os.Signal)
|
||||||
signal.Notify(stop, syscall.SIGINT)
|
signal.Notify(stop, syscall.SIGINT)
|
||||||
|
|
||||||
go func() {
|
go func() {
|
||||||
if err := s.Serve(); err != nil {
|
if err := s.Start(); nil != err {
|
||||||
log.Fatalf("Cannot start rpc server: %s", err)
|
log.Printf("Server: Start error %v", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
|
@ -45,6 +41,7 @@ func main() {
|
||||||
case signal := <-stop:
|
case signal := <-stop:
|
||||||
fmt.Printf("Got signal: %v\n", signal)
|
fmt.Printf("Got signal: %v\n", signal)
|
||||||
}
|
}
|
||||||
|
|
||||||
s.Stop()
|
s.Stop()
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|
|
@ -1 +1,7 @@
|
||||||
package rpc
|
package rpc
|
||||||
|
|
||||||
|
import "git.loafle.net/commons_go/rpc"
|
||||||
|
|
||||||
|
func RegisterRPC(registry rpc.Registry) {
|
||||||
|
registry.RegisterService(&DiscoveryService{}, "")
|
||||||
|
}
|
||||||
|
|
28
server/conn.go
Normal file
28
server/conn.go
Normal file
|
@ -0,0 +1,28 @@
|
||||||
|
package server
|
||||||
|
|
||||||
|
import "net"
|
||||||
|
|
||||||
|
func newConn(nconn net.Conn, contentType string) Conn {
|
||||||
|
c := &conn{
|
||||||
|
contentType: contentType,
|
||||||
|
}
|
||||||
|
c.Conn = nconn
|
||||||
|
|
||||||
|
return c
|
||||||
|
}
|
||||||
|
|
||||||
|
type Conn interface {
|
||||||
|
net.Conn
|
||||||
|
|
||||||
|
GetContentType() string
|
||||||
|
}
|
||||||
|
|
||||||
|
type conn struct {
|
||||||
|
net.Conn
|
||||||
|
|
||||||
|
contentType string
|
||||||
|
}
|
||||||
|
|
||||||
|
func (c *conn) GetContentType() string {
|
||||||
|
return c.contentType
|
||||||
|
}
|
9
server/rpc_server_handler.go
Normal file
9
server/rpc_server_handler.go
Normal file
|
@ -0,0 +1,9 @@
|
||||||
|
package server
|
||||||
|
|
||||||
|
import (
|
||||||
|
"git.loafle.net/commons_go/rpc/server"
|
||||||
|
)
|
||||||
|
|
||||||
|
type RPCServerHandler interface {
|
||||||
|
server.ServerHandler
|
||||||
|
}
|
|
@ -4,49 +4,58 @@ import (
|
||||||
"io"
|
"io"
|
||||||
|
|
||||||
"git.loafle.net/commons_go/rpc"
|
"git.loafle.net/commons_go/rpc"
|
||||||
"git.loafle.net/commons_go/rpc/protocol/json"
|
|
||||||
"git.loafle.net/commons_go/rpc/server"
|
"git.loafle.net/commons_go/rpc/server"
|
||||||
)
|
)
|
||||||
|
|
||||||
func NewRPCServerHandler(registry rpc.Registry) *RPCServerHandlers {
|
func newRPCServerHandler(rpcRegistry rpc.Registry) server.ServerHandler {
|
||||||
rpcSH := &RPCServerHandlers{}
|
sh := &RPCServerHandlers{}
|
||||||
rpcSH.RPCRegistry = registry
|
sh.RPCRegistry = rpcRegistry
|
||||||
|
|
||||||
rpcSH.RegisterCodec(json.NewServerCodec(), "json")
|
return sh
|
||||||
|
|
||||||
return rpcSH
|
|
||||||
}
|
}
|
||||||
|
|
||||||
type RPCServerHandlers struct {
|
type RPCServerHandlers struct {
|
||||||
server.RPCServerHandlers
|
server.ServerHandlers
|
||||||
|
|
||||||
addr string
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) GetContentType(r io.Reader) string {
|
func (sh *RPCServerHandlers) Init() error {
|
||||||
return "json"
|
|
||||||
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnPreRead(r io.Reader) {
|
func (sh *RPCServerHandlers) OnStart() {
|
||||||
// no op
|
// no op
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnPostRead(r io.Reader) {
|
func (sh *RPCServerHandlers) OnPreRead(r io.Reader) {
|
||||||
// no op
|
// no op
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnPreWriteResult(w io.Writer, result interface{}) {
|
func (sh *RPCServerHandlers) OnPostRead(r io.Reader) {
|
||||||
// no op
|
// no op
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnPostWriteResult(w io.Writer, result interface{}) {
|
func (sh *RPCServerHandlers) OnPreWriteResult(w io.Writer, result interface{}) {
|
||||||
// no op
|
// no op
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnPreWriteError(w io.Writer, err error) {
|
func (sh *RPCServerHandlers) OnPostWriteResult(w io.Writer, result interface{}) {
|
||||||
// no op
|
// no op
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnPostWriteError(w io.Writer, err error) {
|
func (sh *RPCServerHandlers) OnPreWriteError(w io.Writer, err error) {
|
||||||
// no op
|
// no op
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (sh *RPCServerHandlers) OnPostWriteError(w io.Writer, err error) {
|
||||||
|
// no op
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sh *RPCServerHandlers) OnStop() {
|
||||||
|
// no op
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sh *RPCServerHandlers) Validate() {
|
||||||
|
sh.ServerHandlers.Validate()
|
||||||
|
|
||||||
|
}
|
||||||
|
|
|
@ -1,15 +1,21 @@
|
||||||
package server
|
package server
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"git.loafle.net/commons_go/rpc"
|
cr "git.loafle.net/commons_go/rpc"
|
||||||
|
crpj "git.loafle.net/commons_go/rpc/protocol/json"
|
||||||
"git.loafle.net/commons_go/server"
|
"git.loafle.net/commons_go/server"
|
||||||
|
|
||||||
|
"git.loafle.net/overflow/overflow_discovery/rpc"
|
||||||
)
|
)
|
||||||
|
|
||||||
func New(addr string, registry rpc.Registry) server.Server {
|
func New(addr string) server.Server {
|
||||||
|
rpcRegistry := cr.NewRegistry()
|
||||||
|
rpc.RegisterRPC(rpcRegistry)
|
||||||
|
|
||||||
sh := NewServerHandler(addr)
|
rpcSH := newRPCServerHandler(rpcRegistry)
|
||||||
rpcSH := NewRPCServerHandler(registry)
|
rpcSH.RegisterCodec(crpj.NewServerCodec(), crpj.Name)
|
||||||
sh.RPCServerHandler = rpcSH
|
|
||||||
|
sh := newServerHandler(addr, rpcSH)
|
||||||
|
|
||||||
s := server.New(sh)
|
s := server.New(sh)
|
||||||
|
|
||||||
|
|
7
server/server_handler.go
Normal file
7
server/server_handler.go
Normal file
|
@ -0,0 +1,7 @@
|
||||||
|
package server
|
||||||
|
|
||||||
|
import "git.loafle.net/commons_go/server"
|
||||||
|
|
||||||
|
type ServerHandler interface {
|
||||||
|
server.ServerHandler
|
||||||
|
}
|
|
@ -1,15 +1,22 @@
|
||||||
package server
|
package server
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"fmt"
|
||||||
|
"log"
|
||||||
"net"
|
"net"
|
||||||
"os"
|
|
||||||
|
|
||||||
"git.loafle.net/commons_go/rpc/server"
|
"git.loafle.net/commons_go/logging"
|
||||||
|
|
||||||
|
crs "git.loafle.net/commons_go/rpc/server"
|
||||||
|
"git.loafle.net/commons_go/server"
|
||||||
)
|
)
|
||||||
|
|
||||||
func NewServerHandler(addr string) *ServerHandlers {
|
func newServerHandler(addr string, rpcSH RPCServerHandler) ServerHandler {
|
||||||
sh := &ServerHandlers{}
|
sh := &ServerHandlers{
|
||||||
sh.addr = addr
|
addr: addr,
|
||||||
|
rpcSH: rpcSH,
|
||||||
|
}
|
||||||
|
sh.Name = "Discovery"
|
||||||
|
|
||||||
return sh
|
return sh
|
||||||
}
|
}
|
||||||
|
@ -17,26 +24,66 @@ func NewServerHandler(addr string) *ServerHandlers {
|
||||||
type ServerHandlers struct {
|
type ServerHandlers struct {
|
||||||
server.ServerHandlers
|
server.ServerHandlers
|
||||||
|
|
||||||
|
rpcSH RPCServerHandler
|
||||||
addr string
|
addr string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (sh *ServerHandlers) Init() error {
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnStart() {
|
func (sh *ServerHandlers) OnStart() {
|
||||||
// no op
|
// no op
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (sh *ServerHandlers) OnConnect(conn net.Conn) (net.Conn, error) {
|
||||||
|
var err error
|
||||||
|
if conn, err = sh.ServerHandlers.OnConnect(conn); nil != err {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return newConn(conn, "jsonrpc"), nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (sh *ServerHandlers) Handle(conn net.Conn, stopChan <-chan struct{}, doneChan chan<- struct{}) {
|
||||||
|
dConn := conn.(Conn)
|
||||||
|
contentType := dConn.GetContentType()
|
||||||
|
codec, err := sh.rpcSH.GetCodec(contentType)
|
||||||
|
if nil != err {
|
||||||
|
log.Printf("RPC Handle: %v", err)
|
||||||
|
doneChan <- struct{}{}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
for {
|
||||||
|
if err := crs.Handle(sh.rpcSH, codec, conn, conn); nil != err {
|
||||||
|
if server.IsClientDisconnect(err) {
|
||||||
|
doneChan <- struct{}{}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
log.Printf("RPC: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
select {
|
||||||
|
case <-stopChan:
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnStop() {
|
func (sh *ServerHandlers) OnStop() {
|
||||||
// no op
|
// no op
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) Listen() (net.Listener, error) {
|
func (sh *ServerHandlers) Validate() {
|
||||||
os.Remove(sh.addr)
|
sh.ServerHandlers.Validate()
|
||||||
l, err := net.ListenUnix("unix", &net.UnixAddr{Name: sh.addr, Net: "unix"})
|
|
||||||
if nil == err {
|
if "" == sh.addr {
|
||||||
os.Chmod(sh.addr, 0777)
|
logging.Logger().Panic(fmt.Sprintf("Server: Address of server must be specified"))
|
||||||
}
|
|
||||||
return l, err
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (sh *ServerHandlers) OnAccept(conn net.Conn) (net.Conn, error) {
|
if nil == sh.rpcSH {
|
||||||
return conn, nil
|
logging.Logger().Panic(fmt.Sprintf("Server: RPC Server Handler must be specified"))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
15
server/server_handlers_unix.go
Normal file
15
server/server_handlers_unix.go
Normal file
|
@ -0,0 +1,15 @@
|
||||||
|
package server
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net"
|
||||||
|
"os"
|
||||||
|
)
|
||||||
|
|
||||||
|
func (sh *ServerHandlers) Listen() (net.Listener, error) {
|
||||||
|
os.Remove(sh.addr)
|
||||||
|
l, err := net.ListenUnix("unix", &net.UnixAddr{Name: sh.addr, Net: "unix"})
|
||||||
|
if nil == err {
|
||||||
|
os.Chmod(sh.addr, 0777)
|
||||||
|
}
|
||||||
|
return l, err
|
||||||
|
}
|
15
server/server_handlers_windows.go
Normal file
15
server/server_handlers_windows.go
Normal file
|
@ -0,0 +1,15 @@
|
||||||
|
package server
|
||||||
|
|
||||||
|
import (
|
||||||
|
"net"
|
||||||
|
|
||||||
|
"gopkg.in/natefinch/npipe.v2"
|
||||||
|
)
|
||||||
|
|
||||||
|
func (sh *ServerHandlers) Listen() (net.Listener, error) {
|
||||||
|
ln, err := npipe.Listen(`\\.\pipe\` + sh.addr)
|
||||||
|
if err != nil {
|
||||||
|
// handle error
|
||||||
|
}
|
||||||
|
return ln, err
|
||||||
|
}
|
Loading…
Reference in New Issue
Block a user