mirror of
https://github.com/hyprspace/hyprspace.git
synced 2026-09-13 18:46:29 +05:00
* Added ServicesACL config Signed-off-by: Roland Urbano <urbano.roland@gmail.com> * Adding filter to service receiver Signed-off-by: Roland Urbano <urbano.roland@gmail.com> * Implemented ACL handling in ServiceNetwork Signed-off-by: Roland Urbano <roland.urbano@nts.eu> * Updated Schema and added parsing for acl config Signed-off-by: Roland Urbano <urbano.roland@gmail.com> * Updated service package and fixed ACL handling Signed-off-by: Roland Urbano <urbano.roland@gmail.com> * go-jsonschema: 0.16.0 -> 0.22.0 * nixos: set acls on services directly * config: move ServicesACL to Services * svc: simplify acl logic * remove whitespace * apply nixfmt --------- Signed-off-by: Roland Urbano <urbano.roland@gmail.com> Signed-off-by: Roland Urbano <roland.urbano@nts.eu> Co-authored-by: Max <max@privatevoid.net>
123 lines
2.7 KiB
Go
123 lines
2.7 KiB
Go
package svc
|
|
|
|
import (
|
|
"context"
|
|
"encoding/binary"
|
|
"errors"
|
|
"fmt"
|
|
"net"
|
|
|
|
"github.com/libp2p/go-libp2p/core/host"
|
|
"github.com/libp2p/go-libp2p/core/peer"
|
|
"github.com/multiformats/go-multiaddr"
|
|
"go.uber.org/zap"
|
|
)
|
|
|
|
type Proxy struct {
|
|
Handle func(net.Conn)
|
|
Description string
|
|
}
|
|
|
|
func (p Proxy) ServeFunc() func(net.Listener) error {
|
|
return func(l net.Listener) error {
|
|
for {
|
|
conn, err := l.Accept()
|
|
if err != nil {
|
|
return err
|
|
}
|
|
go p.Handle(conn)
|
|
}
|
|
}
|
|
}
|
|
|
|
func ProxyTo(ma multiaddr.Multiaddr) (Proxy, error) {
|
|
var proxy Proxy
|
|
|
|
var err error
|
|
var done bool
|
|
|
|
var ipAddr net.IP = net.IPv4(127, 0, 0, 1)
|
|
var tcpPort uint16
|
|
|
|
multiaddr.ForEach(ma, func(c multiaddr.Component) bool {
|
|
if done {
|
|
err = errors.New("trailing components in address")
|
|
return false
|
|
}
|
|
switch c.Protocol().Code {
|
|
case multiaddr.P_IP4:
|
|
fallthrough
|
|
case multiaddr.P_IP6:
|
|
ipAddr = net.IP(c.RawValue())
|
|
return true
|
|
case multiaddr.P_TCP:
|
|
tcpPort = binary.BigEndian.Uint16(c.RawValue())
|
|
proxy = TCPServiceProxy(net.TCPAddr{
|
|
IP: ipAddr,
|
|
Port: int(tcpPort),
|
|
})
|
|
done = true
|
|
default:
|
|
err = errors.New("unsupported protocol: " + c.Protocol().Name)
|
|
return false
|
|
}
|
|
return false
|
|
})
|
|
return proxy, err
|
|
}
|
|
|
|
type RemoteServiceProxyStatus byte
|
|
|
|
const (
|
|
RS_OK RemoteServiceProxyStatus = 0xf1
|
|
RS_NOT_SUPPORTED RemoteServiceProxyStatus = 0xf2
|
|
RS_NOT_AUTHORIZED RemoteServiceProxyStatus = 0xf3
|
|
)
|
|
|
|
func RemoteServiceProxy(host host.Host, p peer.ID, svcId [2]byte) Proxy {
|
|
return Proxy{
|
|
Handle: func(conn net.Conn) {
|
|
ctx := context.Background()
|
|
stream, err := host.NewStream(ctx, p, Protocol)
|
|
defer conn.Close()
|
|
if err != nil {
|
|
logger.With(err).Error("Failed to open new stream")
|
|
return
|
|
}
|
|
defer stream.Close()
|
|
_, err = stream.Write(svcId[:])
|
|
if err != nil {
|
|
logger.With(err).Error("Failed to write to stream")
|
|
return
|
|
}
|
|
buf := make([]byte, 1)
|
|
_, err = stream.Read(buf)
|
|
if err != nil {
|
|
logger.With(err).Error("Failed to read from stream")
|
|
return
|
|
} else if buf[0] != byte(RS_OK) {
|
|
logger.With(zap.String("peer", p.String()), zap.String("service", fmt.Sprintf("%x", svcId[:]))).Warn("Peer does not support service")
|
|
return
|
|
}
|
|
pipe(conn, stream)
|
|
},
|
|
Description: fmt.Sprintf("RemoteServiceProxy to service [%x] on %s", svcId, p),
|
|
}
|
|
}
|
|
|
|
func TCPServiceProxy(tcpAddr net.TCPAddr) Proxy {
|
|
return Proxy{
|
|
Handle: func(conn net.Conn) {
|
|
stream, err := net.DialTCP("tcp", nil, &tcpAddr)
|
|
defer conn.Close()
|
|
if err != nil {
|
|
logger.With(zap.String("address", tcpAddr.String())).With(err).Error("Failed to Dial")
|
|
return
|
|
}
|
|
defer stream.Close()
|
|
pipe(conn, stream)
|
|
},
|
|
Description: fmt.Sprintf("TCPServiceProxy to %s", tcpAddr.AddrPort()),
|
|
}
|
|
}
|