nonlib/OSC/Endpoint: Work around for liblo/UDP layer dropping packets on bulk signal listing.

This commit is contained in:
Jonathan Moore Liles 2020-10-20 18:52:37 -07:00
parent 830823a226
commit 362a153412
2 changed files with 138 additions and 20 deletions

View File

@ -23,6 +23,7 @@
#include <stdio.h>
#include <string.h>
#include <assert.h>
#include <unistd.h>
#include "Endpoint.H"
@ -30,6 +31,8 @@
#pragma GCC diagnostic ignored "-Wunused-parameter"
static int SCAN_BATCH_SIZE = 100;
namespace OSC
{
@ -744,34 +747,122 @@ namespace OSC
{
// OSC_DMSG();
DMESSAGE( "Listing signals." );
const char *prefix = NULL;
int skip = 0;
bool batch_mode = true;
int start;
if ( argc )
const int count = SCAN_BATCH_SIZE;
if ( 's' == types[0] )
{
prefix = &argv[0]->s;
DMESSAGE( "Listing signals for prefix \"%s\"", prefix );
skip++;
batch_mode = 2 == argc;
}
else
{
DMESSAGE( "Listing all signals." );
}
if ( 0 == argc || 's' == types[0] )
{
/* obsolete unbatched mode left for backwards compatibility... */
WARNING("Peer asked for unbatched signal list... will incur delays!");
batch_mode = false;
}
else
{
/* we have to list these in batches, otherwise they the UDP packets overrun and get dropped... Probably
wouldn't be an issue with TCP transport. */
batch_mode = true;
start = argv[0+skip]->i;
/* count = argv[1+skip]->i; */
}
Endpoint *ep = (Endpoint*)user_data;
int sent = 0;
int s = 0;
bool more = false;
lo_bundle b;
if ( batch_mode )
{
b = lo_bundle_new(LO_TT_IMMEDIATE);
}
for ( std::list<Signal*>::const_iterator i = ep->_signals.begin(); i != ep->_signals.end(); ++i )
{
Signal *o = *i;
if ( ! prefix || ! strncmp( o->path(), prefix, strlen(prefix) ) )
{
ep->send( lo_message_get_source( msg ),
"/reply",
path,
o->path(),
o->_direction == Signal::Input ? "in" : "out",
o->parameter_limits().min,
o->parameter_limits().max,
o->parameter_limits().default_value
if ( s++ < start )
continue;
/* DMESSAGE( "Listing signal %s", o->path() ); */
if ( batch_mode )
{
lo_message m = lo_message_new();
lo_message_add_string( m, path );
lo_message_add_string( m, o->path() );
lo_message_add_string( m, o->_direction == Signal::Input ? "in" : "out" );
lo_message_add_float( m, o->parameter_limits().min );
lo_message_add_float( m, o->parameter_limits().max );
lo_message_add_float( m, o->parameter_limits().default_value );
lo_bundle_add_message( b, "/reply", m );
}
else
{
ep->send( lo_message_get_source( msg ),
"/reply",
path,
o->path(),
o->_direction == Signal::Input ? "in" : "out",
o->parameter_limits().min,
o->parameter_limits().max,
o->parameter_limits().default_value
);
}
if ( batch_mode )
{
if ( ++sent == count )
{
more = true;
break;
}
}
else
{
usleep(1000);
}
}
}
ep->send( lo_message_get_source( msg ), "/reply", path );
if ( batch_mode )
{
lo_message m = lo_message_new();
lo_message_add_string(m, path);
lo_message_add_int32( m, sent );
lo_message_add_int32( m, more );
lo_bundle_add_message( b, "/reply", m );
lo_send_bundle_from( lo_message_get_source( msg ), ep->_server, b );
/* ep->send( lo_message_get_source( msg ), "/reply", path, sent, more ); */
}
else
ep->send( lo_message_get_source( msg ), "/reply", path );
return 0;
}
@ -906,13 +997,38 @@ namespace OSC
return 0;
}
if ( argc == 1 )
{
p->_scanning = false;
DMESSAGE( "Done scanning %s", p->name );
if ( argc == 1 )
{
/* old handler for backwards compatibility */
p->_scanning = false;
DMESSAGE( "Done scanning %s", p->name );
if ( ep->_peer_scan_complete_callback )
ep->_peer_scan_complete_callback(ep->_peer_scan_complete_userdata);
}
else
if ( argc == 3 )
{
const int sent = argv[1]->i;
const int more = argv[2]->i;
if ( !more )
{
p->_scanning = false;
DMESSAGE( "Done scanning %s", p->name );
if ( ep->_peer_scan_complete_callback )
ep->_peer_scan_complete_callback(ep->_peer_scan_complete_userdata);
}
else
{
DMESSAGE( "Scanning next batch %s", p->name );
p->_scanning_current += sent;
ep->send( p->addr, "/signal/list", p->_scanning_current );
}
if ( ep->_peer_scan_complete_callback )
ep->_peer_scan_complete_callback(ep->_peer_scan_complete_userdata);
}
else if ( argc == 6 && p->_scanning )
{
@ -921,7 +1037,7 @@ namespace OSC
if ( s )
return 0;
DMESSAGE( "Peer %s has signal %s (%s)", p->name, &argv[1]->s, &argv[2]->s );
/* DMESSAGE( "Peer %s has signal %s (%s)", p->name, &argv[1]->s, &argv[2]->s ); */
int dir = 0;
@ -1143,10 +1259,11 @@ namespace OSC
Peer *p = add_peer(name,url);
p->_scanning = true;
p->_scanning_current = 0;
DMESSAGE( "Scanning peer %s", name );
send( p->addr, "/signal/list" );
send( p->addr, "/signal/list", p->_scanning_current );
}
void *

View File

@ -120,6 +120,7 @@ namespace OSC
struct Peer
{
bool _scanning;
int _scanning_current;
char *name;
lo_address addr;