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
54 changes: 5 additions & 49 deletions examples/interop/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,16 +16,13 @@ import (
"math/big"
"os"
"os/exec"
"strings"
"time"

mqlog "github.com/mengelbart/qlog"
roqqlog "github.com/mengelbart/qlog/roq"
"github.com/mengelbart/roq"
"github.com/mengelbart/roq/qlog"
"github.com/pion/rtp"
"github.com/pion/rtp/codecs"
"github.com/quic-go/quic-go"
"github.com/quic-go/quic-go/qlog"
)

type flags struct {
Expand Down Expand Up @@ -136,19 +133,7 @@ func connect(ctx context.Context, f flags, keyLog io.Writer) (*quic.Conn, error)
}

func runSender(f flags, conn *quic.Conn) error {
role := "server"
if !f.Server {
role = "client"
}
qlogfile := getQLOGWriter("roq", role)
if qlogfile != nil {
defer qlogfile.Close() //nolint
}
var qlogger *mqlog.Logger
if qlogfile != nil {
qlogger = mqlog.NewQLOGHandler(qlogfile, "roq qlog", "RoQ Interop tester QLOG Events", role, roqqlog.Schema)
}
s, err := newSender(roq.NewQUICGoConnection(conn), qlogger)
s, err := newSender(roq.NewQUICGoConnection(conn))
if err != nil {
return err
}
Expand Down Expand Up @@ -185,22 +170,12 @@ func runSender(f flags, conn *quic.Conn) error {
}

func runReceiver(f flags, conn *quic.Conn) error {
role := "client"
if f.Server {
role = "server"
}
qlogfile := getQLOGWriter("roq", role)
if qlogfile != nil {
defer qlogfile.Close() //nolint
}
var qlogger *mqlog.Logger
if qlogfile != nil {
qlogger = mqlog.NewQLOGHandler(qlogfile, "roq-qlog", "RoQ Interop tester QLOG Events", role, roqqlog.Schema)
}
r, err := newReceiver(roq.NewQUICGoConnection(conn), qlogger)
r, err := newReceiver(roq.NewQUICGoConnection(conn))
if err != nil {
return err
}
// Closing the session finishes the qlog trace of the connection.
defer r.Close() //nolint
var writer io.WriteCloser
if len(f.Destination) > 0 {
fileWriter, err := newFileWriter(f.Destination, f.Codec)
Expand Down Expand Up @@ -291,25 +266,6 @@ func generateTLSConfigWithNewCert(keyLog io.Writer) (*tls.Config, error) {
}, nil
}

func getQLOGWriter(id, vantagePoint string) io.WriteCloser {
qlogDir := os.Getenv("QLOGDIR")
if qlogDir == "" {
return nil
}
if _, err := os.Stat(qlogDir); os.IsNotExist(err) {
if err := os.MkdirAll(qlogDir, 0o755); err != nil {
log.Fatalf("failed to create qlog dir %s: %v", qlogDir, err)
}
}
path := fmt.Sprintf("%s/%s_%s.qlog", strings.TrimRight(qlogDir, "/"), id, vantagePoint)
f, err := os.Create(path)
if err != nil {
log.Printf("Failed to create qlog file %s: %s", path, err.Error())
return nil
}
return f
}

type bufferedWriteCloser struct {
*bufio.Writer
io.Closer
Expand Down
5 changes: 2 additions & 3 deletions examples/interop/receiver.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,16 +3,15 @@ package main
import (
"io"

"github.com/mengelbart/qlog"
"github.com/mengelbart/roq"
)

type receiver struct {
session *roq.Session
}

func newReceiver(conn roq.Connection, qlog *qlog.Logger) (*receiver, error) {
session, err := roq.NewSession(conn, true, qlog)
func newReceiver(conn roq.Connection) (*receiver, error) {
session, err := roq.NewSession(conn, true)
if err != nil {
return nil, err
}
Expand Down
5 changes: 2 additions & 3 deletions examples/interop/sender.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,6 @@ import (
"context"
"io"

"github.com/mengelbart/qlog"
"github.com/mengelbart/roq"
"github.com/pion/rtp"
)
Expand All @@ -25,8 +24,8 @@ type sender struct {
session *roq.Session
}

func newSender(conn roq.Connection, qlog *qlog.Logger) (*sender, error) {
session, err := roq.NewSession(conn, true, qlog)
func newSender(conn roq.Connection) (*sender, error) {
session, err := roq.NewSession(conn, true)
if err != nil {
return nil, err
}
Expand Down
8 changes: 5 additions & 3 deletions examples/playfromdisk/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,11 +13,11 @@ import (
"time"

"github.com/mengelbart/roq"
"github.com/mengelbart/roq/qlog"
"github.com/pion/rtp"
"github.com/pion/rtp/codecs"
"github.com/pion/webrtc/v3/pkg/media/ivfreader"
"github.com/quic-go/quic-go"
"github.com/quic-go/quic-go/qlog"
)

const (
Expand Down Expand Up @@ -81,10 +81,12 @@ func main() {
if err != nil {
panic(err)
}
session, err := roq.NewSession(roq.NewQUICGoConnection(conn), true, nil)
session, err := roq.NewSession(roq.NewQUICGoConnection(conn), true)
if err != nil {
panic(err)
}
// Closing the session finishes the qlog trace of the connection.
defer session.Close() //nolint

flow, err := session.NewSendFlow(0)
if err != nil {
Expand All @@ -97,7 +99,7 @@ func main() {
frame, _, ivfErr := ivf.ParseNextFrame()
if errors.Is(ivfErr, io.EOF) {
fmt.Printf("All video frames parsed and sent")
os.Exit(0)
return
}

if ivfErr != nil {
Expand Down
6 changes: 4 additions & 2 deletions examples/savetodisk/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,10 @@ import (
"log"

"github.com/mengelbart/roq"
"github.com/mengelbart/roq/qlog"
"github.com/pion/rtp"
"github.com/pion/webrtc/v3/pkg/media/ivfwriter"
"github.com/quic-go/quic-go"
"github.com/quic-go/quic-go/qlog"
)

func main() {
Expand All @@ -34,10 +34,12 @@ func main() {
if err != nil {
panic(err)
}
session, err := roq.NewSession(roq.NewQUICGoConnection(conn), true, nil)
session, err := roq.NewSession(roq.NewQUICGoConnection(conn), true)
if err != nil {
panic(err)
}
// Closing the session finishes the qlog trace of the connection.
defer session.Close() //nolint
flow, err := session.NewReceiveFlow(0)
if err != nil {
panic(err)
Expand Down
1 change: 0 additions & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,6 @@ module github.com/mengelbart/roq
go 1.25.0

require (
github.com/mengelbart/qlog v0.1.0
github.com/pion/interceptor v0.1.47
github.com/pion/rtp v1.10.5
github.com/pion/webrtc/v3 v3.3.6
Expand Down
2 changes: 0 additions & 2 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,6 @@ github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE=
github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk=
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
github.com/mengelbart/qlog v0.1.0 h1:8cDMuCMcKtzkPXUU5FF7OBwqKiy+De0GKvIvNawifoA=
github.com/mengelbart/qlog v0.1.0/go.mod h1:nIlGcUugkfDu41B8LKdAwjHQ1NxAF54D9hS2EDOlVyk=
github.com/pion/interceptor v0.1.47 h1:yw8t5pJ2f8t78NgU+8EmxhaqYLXS7uFCC/tAGOaSDBo=
github.com/pion/interceptor v0.1.47/go.mod h1:7yoRBzaIDETPC6cIN8Zj9EyGqHv1ImOpcTFPha6MuOM=
github.com/pion/logging v0.2.4 h1:tTew+7cmQ+Mc1pTBLKH2puKsOvhm32dROumOZ655zB8=
Expand Down
4 changes: 2 additions & 2 deletions integrationtests/integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,15 +26,15 @@ func accept(t *testing.T, ctx context.Context, listener *quic.Listener) *roq.Ses
conn, err := listener.Accept(ctx)
assert.NoError(t, err)
assert.NoError(t, err)
s, err := roq.NewSession(roq.NewQUICGoConnection(conn), true, nil)
s, err := roq.NewSession(roq.NewQUICGoConnection(conn), true)
assert.NoError(t, err)
return s
}

func dial(t *testing.T, ctx context.Context, addr string) *roq.Session {
conn, err := quic.DialAddr(ctx, addr, generateTLSConfig(), &quic.Config{EnableDatagrams: true})
assert.NoError(t, err)
s, err := roq.NewSession(roq.NewQUICGoConnection(conn), true, nil)
s, err := roq.NewSession(roq.NewQUICGoConnection(conn), true)
assert.NoError(t, err)
return s
}
Expand Down
150 changes: 150 additions & 0 deletions integrationtests/qlog_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
package integrationtests_test

import (
"bufio"
"bytes"
"context"
"encoding/json"
"os"
"path/filepath"
"testing"

"github.com/mengelbart/roq"
roqqlog "github.com/mengelbart/roq/qlog"
"github.com/quic-go/quic-go"
quicqlog "github.com/quic-go/quic-go/qlog"
"github.com/quic-go/quic-go/qlogwriter"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// TestQlog runs a session over a real QUIC connection whose tracer carries the
// RoQ event schema, and checks that the RoQ events end up in the same sqlog file
// as the QUIC transport events.
func TestQlog(t *testing.T) {
qlogDir := t.TempDir()
t.Setenv("QLOGDIR", qlogDir)
config := &quic.Config{EnableDatagrams: true, Tracer: roqqlog.DefaultConnectionTracer}

listener, err := quic.ListenAddr("localhost:0", generateTLSConfig(), config)
require.NoError(t, err)
defer listener.Close() //nolint:errcheck

ctx, cancel := context.WithCancel(context.Background())
defer cancel()

type receiveEnd struct {
session *roq.Session
flow *roq.ReceiveFlow
}
accepted := make(chan receiveEnd, 1)
go func() {
conn, err := listener.Accept(ctx)
assert.NoError(t, err)
s, err := roq.NewSession(roq.NewQUICGoConnection(conn), true)
assert.NoError(t, err)
f, err := s.NewReceiveFlow(0)
assert.NoError(t, err)
accepted <- receiveEnd{session: s, flow: f}
}()

conn, err := quic.DialAddr(ctx, listener.Addr().String(), generateTLSConfig(), config)
require.NoError(t, err)
sender, err := roq.NewSession(roq.NewQUICGoConnection(conn), true)
require.NoError(t, err)

// Only send once the peer is ready to receive, so that a datagram cannot be
// dropped before its connection has finished the handshake.
receiver := <-accepted

f, err := sender.NewSendFlow(0)
require.NoError(t, err)
stream, err := f.NewSendStream(ctx, 0, false)
require.NoError(t, err)
_, err = stream.WriteRTPBytes(make([]byte, 100))
require.NoError(t, err)
require.NoError(t, f.WriteRTPBytes(make([]byte, 100)))

buf := make([]byte, 2000)
for i := 0; i < 2; i++ {
_, err = receiver.flow.Read(buf)
require.NoError(t, err)
}

// Closing the sessions releases their qlog producers, which is what closes
// and flushes the sqlog files.
require.NoError(t, sender.Close())
require.NoError(t, receiver.session.Close())

files, err := filepath.Glob(filepath.Join(qlogDir, "*.sqlog"))
require.NoError(t, err)
require.Len(t, files, 2)

var sawRoQ, sawTransport bool
for _, file := range files {
header, names := readQlog(t, file)
assert.Contains(t, header.Trace.EventSchemas, quicqlog.EventSchema)
assert.Contains(t, header.Trace.EventSchemas, roqqlog.EventSchema)
for _, name := range names {
switch name {
case "roq:stream_packet_created", "roq:datagram_packet_parsed":
sawRoQ = true
case "transport:packet_sent":
sawTransport = true
}
}
}
assert.True(t, sawRoQ, "no RoQ events in the qlog files")
assert.True(t, sawTransport, "no QUIC transport events in the qlog files")
}

type qlogHeader struct {
Trace struct {
EventSchemas []string `json:"event_schemas"`
} `json:"trace"`
}

// readQlog parses a JSON-SEQ qlog file into its header record and the names of
// the events that follow.
func readQlog(t *testing.T, path string) (qlogHeader, []string) {
t.Helper()
contents, err := os.ReadFile(path)
require.NoError(t, err)

var header qlogHeader
var names []string
scanner := bufio.NewScanner(bytes.NewReader(contents))
scanner.Buffer(make([]byte, 0, 1<<20), 1<<20)
scanner.Split(splitRecords)
for i := 0; scanner.Scan(); i++ {
if i == 0 {
require.NoError(t, json.Unmarshal(scanner.Bytes(), &header))
continue
}
var event struct {
Name string `json:"name"`
}
require.NoError(t, json.Unmarshal(scanner.Bytes(), &event))
names = append(names, event.Name)
}
require.NoError(t, scanner.Err())
return header, names
}

// splitRecords splits a JSON-SEQ stream on its record separators.
func splitRecords(data []byte, atEOF bool) (int, []byte, error) {
start := bytes.IndexByte(data, qlogwriter.RecordSeparator)
if start < 0 {
if atEOF {
return len(data), nil, nil
}
return 0, nil, nil
}
if end := bytes.IndexByte(data[start+1:], qlogwriter.RecordSeparator); end >= 0 {
return start + 1 + end, data[start+1 : start+1+end], nil
}
if atEOF {
return len(data), data[start+1:], nil
}
return 0, nil, nil
}
Loading