diff --git a/dataplane/forwarding/fwdport/ports/kernel.go b/dataplane/forwarding/fwdport/ports/kernel.go index e8aecdd70..66911f739 100644 --- a/dataplane/forwarding/fwdport/ports/kernel.go +++ b/dataplane/forwarding/fwdport/ports/kernel.go @@ -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) } diff --git a/dataplane/forwarding/fwdport/ports/tap.go b/dataplane/forwarding/fwdport/ports/tap.go index 16c9a3572..618377031 100644 --- a/dataplane/forwarding/fwdport/ports/tap.go +++ b/dataplane/forwarding/fwdport/ports/tap.go @@ -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{}) } @@ -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: diff --git a/dataplane/kernel/genetlink/genetlink.c b/dataplane/kernel/genetlink/genetlink.c index 7d983ade7..c866105d4 100644 --- a/dataplane/kernel/genetlink/genetlink.c +++ b/dataplane/kernel/genetlink/genetlink.c @@ -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 #include @@ -22,6 +22,9 @@ #include #include +#define PACKET_BUFFER_SIZE 16384 // 16KB buffer +#define NL_SOCKET_BUFFER_SIZE 2097152 // 2MB socket buffer + enum { /* packet metadata */ GENL_PACKET_ATTR_IIFINDEX, @@ -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); @@ -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; @@ -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; } diff --git a/dataplane/saiserver/BUILD b/dataplane/saiserver/BUILD index a1a27c374..c4d443702 100644 --- a/dataplane/saiserver/BUILD +++ b/dataplane/saiserver/BUILD @@ -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", diff --git a/dataplane/saiserver/ports.go b/dataplane/saiserver/ports.go index 5d4eb9bfa..8b67d9bdb 100644 --- a/dataplane/saiserver/ports.go +++ b/dataplane/saiserver/ports.go @@ -22,6 +22,7 @@ import ( "slices" "strconv" + "github.com/vishvananda/netlink" "google.golang.org/grpc" "google.golang.org/protobuf/proto" @@ -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), } @@ -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 } diff --git a/dataplane/standalone/pkthandler/pktiohandler/pktiohandler.go b/dataplane/standalone/pkthandler/pktiohandler/pktiohandler.go index 3696824d2..316e4ee01 100644 --- a/dataplane/standalone/pkthandler/pktiohandler/pktiohandler.go +++ b/dataplane/standalone/pkthandler/pktiohandler/pktiohandler.go @@ -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") @@ -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{ @@ -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: @@ -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)