测试
This commit is contained in:
parent
e604afc17d
commit
3018755a52
2
.gitignore
vendored
2
.gitignore
vendored
|
@ -1 +1 @@
|
||||||
|
slimming
|
25
main.go
25
main.go
|
@ -1,27 +1,6 @@
|
||||||
package slimming
|
package main
|
||||||
|
|
||||||
import (
|
|
||||||
"flag"
|
|
||||||
"fmt"
|
|
||||||
"log"
|
|
||||||
"net"
|
|
||||||
|
|
||||||
gen "slimming/proto/gen"
|
|
||||||
|
|
||||||
"google.golang.org/grpc"
|
|
||||||
)
|
|
||||||
|
|
||||||
//go:generate bash -c "protoc --go_out=plugins=grpc:. proto/*.proto"
|
//go:generate bash -c "protoc --go_out=plugins=grpc:. proto/*.proto"
|
||||||
func main() {
|
func main() {
|
||||||
flag.Parse()
|
NewNetCard().Run()
|
||||||
lis, err := net.Listen("tcp", fmt.Sprintf(":%d", *serverPort))
|
|
||||||
if err != nil {
|
|
||||||
log.Fatalf("failed to listen: %v", err)
|
|
||||||
}
|
|
||||||
s := grpc.NewServer()
|
|
||||||
gen.RegisterFrameServiceServer(s, &RPCServer{})
|
|
||||||
log.Printf("server listening at %v", lis.Addr())
|
|
||||||
if err := s.Serve(lis); err != nil {
|
|
||||||
log.Fatalf("failed to serve: %v", err)
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
46
rpc.go
46
rpc.go
|
@ -1,22 +1,18 @@
|
||||||
package slimming
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"context"
|
|
||||||
"flag"
|
"flag"
|
||||||
"fmt"
|
"fmt"
|
||||||
"log"
|
"log"
|
||||||
"net"
|
"net"
|
||||||
gen "slimming/proto/gen"
|
gen "slimming/proto/gen"
|
||||||
"time"
|
|
||||||
|
|
||||||
"google.golang.org/grpc"
|
"google.golang.org/grpc"
|
||||||
"google.golang.org/grpc/credentials/insecure"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type RPCServer struct {
|
type RPCServer struct {
|
||||||
gen.UnimplementedFrameServiceServer
|
gen.UnimplementedFrameServiceServer
|
||||||
|
netCard *NetCard
|
||||||
FrameChan chan [][]byte
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var (
|
var (
|
||||||
|
@ -24,38 +20,11 @@ var (
|
||||||
othersAddr = flag.String("addr", "", "The other server addr")
|
othersAddr = flag.String("addr", "", "The other server addr")
|
||||||
)
|
)
|
||||||
|
|
||||||
var rpcServer = func() *RPCServer {
|
func newRPCServer(netCard *NetCard) *RPCServer {
|
||||||
rs := &RPCServer{}
|
return &RPCServer{netCard: netCard}
|
||||||
|
}
|
||||||
|
|
||||||
conn, err := grpc.Dial(*othersAddr,
|
func (rpc *RPCServer) run() {
|
||||||
grpc.WithTransportCredentials(insecure.NewCredentials()),
|
|
||||||
)
|
|
||||||
if err != nil {
|
|
||||||
log.Fatalf("did not connect: %v", err)
|
|
||||||
}
|
|
||||||
defer conn.Close()
|
|
||||||
c := gen.NewFrameServiceClient(conn)
|
|
||||||
|
|
||||||
// Contact the server and print out its response.
|
|
||||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second*10)
|
|
||||||
defer cancel()
|
|
||||||
|
|
||||||
stream, err := c.SendFrames(ctx)
|
|
||||||
if err != nil {
|
|
||||||
panic(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
err = stream.Send(&gen.Request{
|
|
||||||
Frames: <-rs.FrameChan,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
panic(err)
|
|
||||||
}
|
|
||||||
|
|
||||||
return rs
|
|
||||||
}()
|
|
||||||
|
|
||||||
func (rpc *RPCServer) Run() {
|
|
||||||
flag.Parse()
|
flag.Parse()
|
||||||
lis, err := net.Listen("tcp", fmt.Sprintf(":%d", *serverPort))
|
lis, err := net.Listen("tcp", fmt.Sprintf(":%d", *serverPort))
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
@ -64,6 +33,7 @@ func (rpc *RPCServer) Run() {
|
||||||
s := grpc.NewServer()
|
s := grpc.NewServer()
|
||||||
gen.RegisterFrameServiceServer(s, rpc)
|
gen.RegisterFrameServiceServer(s, rpc)
|
||||||
log.Printf("server listening at %v", lis.Addr())
|
log.Printf("server listening at %v", lis.Addr())
|
||||||
|
|
||||||
if err := s.Serve(lis); err != nil {
|
if err := s.Serve(lis); err != nil {
|
||||||
log.Fatalf("failed to serve: %v", err)
|
log.Fatalf("failed to serve: %v", err)
|
||||||
}
|
}
|
||||||
|
@ -78,7 +48,7 @@ func (s *RPCServer) SendFrames(stream gen.FrameService_SendFramesServer) error {
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Panic(err)
|
log.Panic(err)
|
||||||
}
|
}
|
||||||
netCard.FrameChan <- request.GetFrames()
|
s.netCard.FrameChan <- request.GetFrames() // 接受数据 广播到网卡上
|
||||||
}
|
}
|
||||||
|
|
||||||
// err := stream.SendAndClose(&gen.Response{Code: 0})
|
// err := stream.SendAndClose(&gen.Response{Code: 0})
|
||||||
|
|
74
tap.go
74
tap.go
|
@ -1,19 +1,72 @@
|
||||||
package slimming
|
package main
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
"log"
|
"log"
|
||||||
|
gen "slimming/proto/gen"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/songgao/packets/ethernet"
|
"github.com/songgao/packets/ethernet"
|
||||||
"github.com/songgao/water"
|
"github.com/songgao/water"
|
||||||
|
"google.golang.org/grpc"
|
||||||
|
"google.golang.org/grpc/credentials/insecure"
|
||||||
)
|
)
|
||||||
|
|
||||||
type NetCard struct {
|
type NetCard struct {
|
||||||
FrameChan chan [][]byte
|
FrameChan chan [][]byte
|
||||||
ifce *water.Interface
|
ifce *water.Interface
|
||||||
|
cli *RPCClient
|
||||||
|
server *RPCServer
|
||||||
}
|
}
|
||||||
|
|
||||||
var netCard = func() *NetCard {
|
type RPCClient struct {
|
||||||
|
FrameChan chan [][]byte
|
||||||
|
}
|
||||||
|
|
||||||
|
func (cli *RPCClient) run() {
|
||||||
|
log.Println("rpcclient start")
|
||||||
|
defer log.Println("rpcclient exit")
|
||||||
|
|
||||||
|
conn, err := grpc.Dial(*othersAddr,
|
||||||
|
grpc.WithTransportCredentials(insecure.NewCredentials()),
|
||||||
|
)
|
||||||
|
if err != nil {
|
||||||
|
log.Fatalf("did not connect: %v", err)
|
||||||
|
}
|
||||||
|
defer conn.Close()
|
||||||
|
c := gen.NewFrameServiceClient(conn)
|
||||||
|
|
||||||
|
ctx, cancel := context.WithTimeout(context.Background(), time.Second*10)
|
||||||
|
defer cancel()
|
||||||
|
stream, err := c.SendFrames(ctx)
|
||||||
|
if err != nil {
|
||||||
|
panic(err)
|
||||||
|
}
|
||||||
|
|
||||||
|
for {
|
||||||
|
// Contact the server and print out its response.
|
||||||
|
|
||||||
|
// 发到对面的网卡
|
||||||
|
err = stream.Send(&gen.Request{
|
||||||
|
Frames: <-cli.FrameChan,
|
||||||
|
})
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
panic(err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
func (nc *NetCard) Run() {
|
||||||
|
go nc.runRead()
|
||||||
|
go nc.runWrite()
|
||||||
|
go nc.cli.run()
|
||||||
|
nc.server.run()
|
||||||
|
}
|
||||||
|
|
||||||
|
func NewNetCard() *NetCard {
|
||||||
|
|
||||||
config := water.Config{
|
config := water.Config{
|
||||||
DeviceType: water.TAP,
|
DeviceType: water.TAP,
|
||||||
}
|
}
|
||||||
|
@ -27,20 +80,13 @@ var netCard = func() *NetCard {
|
||||||
nc := &NetCard{
|
nc := &NetCard{
|
||||||
FrameChan: make(chan [][]byte, 2000),
|
FrameChan: make(chan [][]byte, 2000),
|
||||||
ifce: ifce,
|
ifce: ifce,
|
||||||
|
cli: &RPCClient{FrameChan: make(chan [][]byte, 2000)},
|
||||||
}
|
}
|
||||||
|
nc.server = newRPCServer(nc)
|
||||||
go nc.RunRead()
|
|
||||||
go nc.RunWrite()
|
|
||||||
|
|
||||||
time.Sleep(time.Second)
|
|
||||||
return nc
|
return nc
|
||||||
}()
|
|
||||||
|
|
||||||
func GetNetCard() *NetCard {
|
|
||||||
return netCard
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (nc *NetCard) RunRead() {
|
func (nc *NetCard) runRead() {
|
||||||
|
|
||||||
var ifce *water.Interface = nc.ifce
|
var ifce *water.Interface = nc.ifce
|
||||||
var ticker time.Ticker = *time.NewTicker(time.Millisecond * 20)
|
var ticker time.Ticker = *time.NewTicker(time.Millisecond * 20)
|
||||||
|
@ -59,7 +105,7 @@ func (nc *NetCard) RunRead() {
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(framesBytes) > 0 {
|
if len(framesBytes) > 0 {
|
||||||
rpcServer.FrameChan <- framesBytes
|
nc.cli.FrameChan <- framesBytes // 网卡数据 发到对方
|
||||||
}
|
}
|
||||||
|
|
||||||
// 写到grpc服务
|
// 写到grpc服务
|
||||||
|
@ -72,7 +118,7 @@ func (nc *NetCard) RunRead() {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (nc *NetCard) RunWrite() {
|
func (nc *NetCard) runWrite() {
|
||||||
var ifce *water.Interface = nc.ifce
|
var ifce *water.Interface = nc.ifce
|
||||||
|
|
||||||
for wframes := range nc.FrameChan {
|
for wframes := range nc.FrameChan {
|
||||||
|
|
Loading…
Reference in New Issue
Block a user