Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions dataplane/forwarding/fwdport/ports/kernel.go
Original file line number Diff line number Diff line change
Expand Up @@ -190,8 +190,8 @@ func (kernelBuilder) Build(portDesc *fwdpb.PortDesc, ctx *fwdcontext.Context) (f

netlink.LinkSetUp(l)

// TODO: configure MTU
handle, err := pcap.OpenLive(kp.Kernel.GetDeviceName(), 1514, true, pcap.BlockForever)
// Set snaplen to 64KB to capture max size IP packets
handle, err := pcap.OpenLive(kp.Kernel.GetDeviceName(), 65536, true, pcap.BlockForever)
if err != nil {
return nil, fmt.Errorf("failed to create afpacket: %v", err)
}
Expand Down
4 changes: 3 additions & 1 deletion dataplane/forwarding/fwdport/ports/tap.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,8 @@ import (
fwdpb "github.com/openconfig/lemming/proto/forwarding"
)

const packetBufferSize = 16384 // 16KB buffer

func init() {
fwdport.Register(fwdpb.PortType_PORT_TYPE_TAP, tapBuilder{})
}
Expand Down Expand Up @@ -112,7 +114,7 @@ func (p *tapPort) Update(upd *fwdpb.PortUpdateDesc) error {
func (p *tapPort) process() {
startStateWatch(p.linkUpdateCh, p.doneCh, p.devName, p, p.ctx)
go func() {
buf := make([]byte, 1500) // TODO: MTU
buf := make([]byte, packetBufferSize)
for {
select {
case <-p.doneCh:
Expand Down
24 changes: 19 additions & 5 deletions dataplane/kernel/genetlink/genetlink.c
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@
// See the License for the specific language governing permissions and
// limitations under the License.

#include "genetlink.h"
#include "genetlink.h" // NOLINT(build/include_subdir)

#include <linux/netlink.h>
#include <netlink/genl/ctrl.h>
Expand All @@ -22,6 +22,9 @@
#include <stdio.h>
#include <stdlib.h>

#define PACKET_BUFFER_SIZE 16384 // 16KB buffer
#define NL_SOCKET_BUFFER_SIZE 2097152 // 2MB socket buffer

enum {
/* packet metadata */
GENL_PACKET_ATTR_IIFINDEX,
Expand All @@ -39,12 +42,21 @@ struct nl_sock* create_port(const char* family, const char* group) {
return NULL;
}
nl_socket_disable_auto_ack(sock);

int error = genl_connect(sock);
if (error < 0) {
fprintf(stderr, "error: failed to disable auto ack: err %d", error);
fprintf(stderr, "error: failed to connect to genetlink: err %d\n", error);
nl_socket_free(sock);
return NULL;
}

int err = nl_socket_set_buffer_size(sock, NL_SOCKET_BUFFER_SIZE,
NL_SOCKET_BUFFER_SIZE);

// If increased size cannot be set, log error without crashing pkthandler.
if (err < 0) {
fprintf(stderr, "error: failed to set buffer size: err %d\n", err);
}
int group_id = genl_ctrl_resolve_grp(sock, family, group);
if (group_id < 0) {
fprintf(stderr, "error: failed to resolve group: err %d", group_id);
Expand All @@ -59,7 +71,7 @@ void delete_port(void* sock) { nl_socket_free(sock); }

int send_packet(void* sock, int family, const void* pkt, uint32_t size,
int in_ifindex, int out_ifindex, unsigned int context) {
struct nl_msg* msg = nlmsg_alloc();
struct nl_msg* msg = nlmsg_alloc_size(PACKET_BUFFER_SIZE);
if (msg == NULL) {
fprintf(stderr, "failed to allocate packet\n");
return -1;
Expand All @@ -70,13 +82,15 @@ int send_packet(void* sock, int family, const void* pkt, uint32_t size,
NLA_PUT_U32(msg, GENL_PACKET_ATTR_CONTEXT, context);
NLA_PUT(msg, GENL_PACKET_ATTR_DATA, size, pkt);
fprintf(stderr, "sending packet size: %d\n", size);
if (nl_send(sock, msg) < 0) {
fprintf(stderr, "failed to send packet\n");
int err = nl_send(sock, msg);
if (err < 0) {
fprintf(stderr, "failed to send packet: %d\n", err);
return -1;
}
nlmsg_free(msg);
return 0;
nla_put_failure:
fprintf(stderr, "nla_put_failure: packet exceeds nlmsg allcation size\n");
nlmsg_free(msg);
return -1;
}
1 change: 1 addition & 0 deletions dataplane/saiserver/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ go_library(
"//dataplane/saiserver/attrmgr",
"//proto/forwarding",
"@com_github_openconfig_gnmi//errlist",
"@com_github_vishvananda_netlink//:netlink",
"@io_opentelemetry_go_otel//:otel",
"@io_opentelemetry_go_otel_trace//:trace",
"@org_golang_google_grpc//:grpc",
Expand Down
19 changes: 19 additions & 0 deletions dataplane/saiserver/ports.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ import (
"slices"
"strconv"

"github.com/vishvananda/netlink"
"google.golang.org/grpc"
"google.golang.org/protobuf/proto"

Expand Down Expand Up @@ -218,6 +219,7 @@ func (port *port) CreatePort(ctx context.Context, req *saipb.CreatePortRequest)
AdminState: proto.Bool(true),
AutoNegMode: proto.Bool(req.GetAutoNegMode()),
Mtu: proto.Uint32(1514),
HwLaneList: req.HwLaneList,
PortVlanId: proto.Uint32(vId),
}

Expand Down Expand Up @@ -559,6 +561,23 @@ func (port *port) SetPortAttribute(ctx context.Context, req *saipb.SetPortAttrib
return nil, fmt.Errorf("unsupported FEC mode: %v for speed %d and lanes %d", req.GetFecModeExtended(), portAttr.GetAttr().GetSpeed(), len(portAttr.GetAttr().GetHwLaneList()))
}
}
if req.Mtu != nil {
if len(portAttr.GetAttr().GetHwLaneList()) == 0 {
slog.WarnContext(ctx, "port has no lanes", "oid", req.GetOid())
return nil, fmt.Errorf("port %v has no lanes", req.GetOid())
}
dev := fmt.Sprintf("eth%v", portAttr.GetAttr().GetHwLaneList()[0])
slog.InfoContext(ctx, "setting port mtu", "oid", req.GetOid(), "dev", dev, "mtu", req.GetMtu())
link, err := netlink.LinkByName(dev)
if err != nil {
slog.ErrorContext(ctx, "failed to get link", "dev", dev, "err", err)
return nil, err
}
if err := netlink.LinkSetMTU(link, int(req.GetMtu())); err != nil {
slog.ErrorContext(ctx, "failed to set mtu", "dev", dev, "mtu", req.GetMtu(), "err", err)
return nil, err
}
}
return &saipb.SetPortAttributeResponse{}, nil
}

Expand Down
16 changes: 12 additions & 4 deletions dataplane/standalone/pkthandler/pktiohandler/pktiohandler.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,8 @@ import (
pktiopb "github.com/openconfig/lemming/dataplane/proto/packetio"
)

const packetBufferSize = 16384 // 16KB buffer

// New returns a new PacketIOMgr
func New(portFile string) (*PacketIOMgr, error) {
q, err := queue.NewUnbounded("send")
Expand Down Expand Up @@ -251,10 +253,10 @@ func (m *PacketIOMgr) ManagePorts(c pktiopb.PacketIO_HostPortControlClient) erro
Code: int32(codes.OK),
}
state := resp.Op == pktiopb.PortOperation_PORT_OPERATION_SET_UP
if p.SetAdminState(state); err != nil {
if setErr := p.SetAdminState(state); setErr != nil {
st = &status.Status{
Code: int32(codes.Internal),
Message: err.Error(),
Message: setErr.Error(),
}
}
sendErr := c.Send(&pktiopb.HostPortControlRequest{Msg: &pktiopb.HostPortControlRequest_Status{
Expand Down Expand Up @@ -327,7 +329,7 @@ func (m *PacketIOMgr) createPort(msg *pktiopb.HostPortControlMessage) error {
func (m *PacketIOMgr) queueRead(id uint64, done chan struct{}) {
p := m.hostifs[id]
go func() {
buf := make([]byte, 9100) // TODO: Configurable MTU.
buf := make([]byte, packetBufferSize)
for {
select {
case <-done:
Expand All @@ -338,9 +340,15 @@ func (m *PacketIOMgr) queueRead(id uint64, done chan struct{}) {
time.Sleep(time.Millisecond)
continue
}

// sendQueue is asynchronous, so make a packet copy to prevent shared
// buffer from overwriting in the next iteration.
frameCopy := make([]byte, n)
copy(frameCopy, buf[:n])

pkt := &pktiopb.Packet{
HostPort: id,
Frame: buf[0:n],
Frame: frameCopy,
}
m.sendQueue.Write(pkt)
time.Sleep(time.Millisecond)
Expand Down
Loading