mirror of
https://github.com/ossrs/srs.git
synced 2025-03-09 15:49:59 +00:00
refine framework to calc the kbps
This commit is contained in:
parent
3f33dffdb3
commit
9006194cd7
21 changed files with 153 additions and 131 deletions
|
|
@ -42,6 +42,7 @@ CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
|||
#include <srs_app_pithy_print.hpp>
|
||||
#include <srs_core_autofree.hpp>
|
||||
#include <srs_app_socket.hpp>
|
||||
#include <srs_app_kbps.hpp>
|
||||
|
||||
// when error, edge ingester sleep for a while and retry.
|
||||
#define SRS_EDGE_INGESTER_SLEEP_US (int64_t)(1*1000*1000LL)
|
||||
|
|
@ -61,6 +62,7 @@ CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE.
|
|||
SrsEdgeIngester::SrsEdgeIngester()
|
||||
{
|
||||
io = NULL;
|
||||
kbps = new SrsKbps();
|
||||
client = NULL;
|
||||
_edge = NULL;
|
||||
_req = NULL;
|
||||
|
|
@ -75,6 +77,7 @@ SrsEdgeIngester::~SrsEdgeIngester()
|
|||
stop();
|
||||
|
||||
srs_freep(pthread);
|
||||
srs_freep(kbps);
|
||||
}
|
||||
|
||||
int SrsEdgeIngester::initialize(SrsSource* source, SrsPlayEdge* edge, SrsRequest* req)
|
||||
|
|
@ -101,6 +104,7 @@ void SrsEdgeIngester::stop()
|
|||
|
||||
srs_freep(client);
|
||||
srs_freep(io);
|
||||
kbps->set_io(NULL, NULL);
|
||||
}
|
||||
|
||||
int SrsEdgeIngester::cycle()
|
||||
|
|
@ -169,9 +173,8 @@ int SrsEdgeIngester::ingest()
|
|||
// pithy print
|
||||
if (pithy_print.can_print()) {
|
||||
srs_trace("<- "SRS_LOG_ID_EDGE_PLAY
|
||||
" time=%"PRId64", obytes=%"PRId64", ibytes=%"PRId64", okbps=%d, ikbps=%d",
|
||||
pithy_print.age(), client->get_send_bytes(), client->get_recv_bytes(),
|
||||
client->get_send_kbps(), client->get_recv_kbps());
|
||||
" time=%"PRId64", okbps=%d, ikbps=%d",
|
||||
pithy_print.age(), kbps->get_send_kbps(), kbps->get_recv_kbps());
|
||||
}
|
||||
|
||||
// read from client.
|
||||
|
|
@ -303,6 +306,7 @@ int SrsEdgeIngester::connect_server()
|
|||
|
||||
io = new SrsSocket(stfd);
|
||||
client = new SrsRtmpClient(io);
|
||||
kbps->set_io(io, io);
|
||||
|
||||
// connect to server.
|
||||
std::string ip = srs_dns_resolve(server);
|
||||
|
|
@ -330,6 +334,7 @@ int SrsEdgeIngester::connect_server()
|
|||
SrsEdgeForwarder::SrsEdgeForwarder()
|
||||
{
|
||||
io = NULL;
|
||||
kbps = NULL;
|
||||
client = NULL;
|
||||
_edge = NULL;
|
||||
_req = NULL;
|
||||
|
|
@ -347,6 +352,7 @@ SrsEdgeForwarder::~SrsEdgeForwarder()
|
|||
|
||||
srs_freep(pthread);
|
||||
srs_freep(queue);
|
||||
srs_freep(kbps);
|
||||
}
|
||||
|
||||
void SrsEdgeForwarder::set_queue_size(double queue_size)
|
||||
|
|
@ -411,6 +417,7 @@ void SrsEdgeForwarder::stop()
|
|||
|
||||
srs_freep(client);
|
||||
srs_freep(io);
|
||||
kbps->set_io(NULL, NULL);
|
||||
}
|
||||
|
||||
int SrsEdgeForwarder::cycle()
|
||||
|
|
@ -458,9 +465,8 @@ int SrsEdgeForwarder::cycle()
|
|||
// pithy print
|
||||
if (pithy_print.can_print()) {
|
||||
srs_trace("-> "SRS_LOG_ID_EDGE_PUBLISH
|
||||
" time=%"PRId64", msgs=%d, obytes=%"PRId64", ibytes=%"PRId64", okbps=%d, ikbps=%d",
|
||||
pithy_print.age(), count, client->get_send_bytes(), client->get_recv_bytes(),
|
||||
client->get_send_kbps(), client->get_recv_kbps());
|
||||
" time=%"PRId64", msgs=%d, okbps=%d, ikbps=%d",
|
||||
pithy_print.age(), count, kbps->get_send_kbps(), kbps->get_recv_kbps());
|
||||
}
|
||||
|
||||
// ignore when no messages.
|
||||
|
|
@ -576,6 +582,7 @@ int SrsEdgeForwarder::connect_server()
|
|||
|
||||
io = new SrsSocket(stfd);
|
||||
client = new SrsRtmpClient(io);
|
||||
kbps->set_io(io, io);
|
||||
|
||||
// connect to server.
|
||||
std::string ip = srs_dns_resolve(server);
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue