Repository navigation
pjsua2: Add explicit AudioMediaPort callback fencing #5307
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -551,7 +551,32 @@ class AudioMediaPort : public AudioMedia | |
| virtual void onFrameReceived(MediaFrame &frame) | ||
| { PJ_UNUSED_ARG(frame); } | ||
|
|
||
| protected: | ||
| /** | ||
| * Permanently stop dispatching frame callbacks to this object from the | ||
| * created media port. Repeated calls are safe. If no port has been | ||
| * created, this does nothing. | ||
| * | ||
| * This acquires the port's recursive group lock. When called from a | ||
| * thread/context which does not already own that lock, it waits for | ||
| * in-progress callbacks to finish. Calling it from a callback (or while | ||
| * already owning the lock) does not wait for that callback to return; | ||
| * it must not be used as a way to wait for oneself. | ||
| * | ||
| * Call while the object and all callback-visible state are still fully | ||
| * valid, before tearing down derived members. In particular, with | ||
| * multiple levels of inheritance, explicitly detach before destruction | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. With a protected method, code outside the class hierarchy can't "explicitly detach before destruction begins". For a subclass, the simple rule that is always correct is: call it first thing in the most-derived class's destructor, which runs before any member or base is torn down. If the method becomes public (see above), the doc would need to cover both callers: application code calling it before deleting the object, or the most-derived destructor calling it first. |
||
| * begins. This is not a general guarantee of safe concurrent destruction. | ||
| * Do not hold application locks, other ports' locks, or library locks | ||
| * needed by an in-progress callback while calling this method. | ||
| * | ||
| * This does not unregister the conference port or destroy media resources. | ||
| */ | ||
| void detachCallbacks(); | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Was protected deliberate? It leaves out the class's main audience. A C++ app doesn't need If it were public, they could call it before |
||
|
|
||
| private: | ||
| /* Test access to raw frame dispatch without exposing media resources. */ | ||
| friend class AudioMediaPortTest; | ||
| pj_pool_t *pool; | ||
| pjmedia_port *port; | ||
| }; | ||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,179 @@ | ||
| /* | ||
| * Copyright (C) 2026 Teluu Inc. (http://www.teluu.com) | ||
| * | ||
| * This program is free software; you can redistribute it and/or modify | ||
| * it under the terms of the GNU General Public License as published by | ||
| * the Free Software Foundation; either version 2 of the License, or | ||
| * (at your option) any later version. | ||
| * | ||
| * This program is distributed in the hope that it will be useful, | ||
| * but WITHOUT ANY WARRANTY; without even the implied warranty of | ||
| * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | ||
| * GNU General Public License for more details. | ||
| * | ||
| * You should have received a copy of the GNU General Public License | ||
| * along with this program; if not, write to the Free Software | ||
| * Foundation, Inc., 59 Temple Place, Suite 330, Boston, MA 02111-1307 USA | ||
| */ | ||
|
|
||
| #include <pjsua2.hpp> | ||
| #include <pj/lock.h> | ||
| #include <atomic> | ||
| #include <chrono> | ||
| #include <cstdlib> | ||
| #include <future> | ||
| #include <iostream> | ||
| #include <thread> | ||
|
|
||
| /* Fail also in release builds, including stalled worker threads. */ | ||
| #define CHECK(expr) \ | ||
| do { \ | ||
| if (!(expr)) { \ | ||
| std::cerr << "AudioMediaPort: " << #expr << " at line " \ | ||
| << __LINE__ << std::endl; \ | ||
| std::abort(); \ | ||
| } \ | ||
| } while (0) | ||
|
|
||
| namespace pj { | ||
|
|
||
| class AudioMediaPortTest : public AudioMediaPort | ||
| { | ||
| public: | ||
| AudioMediaPortTest() | ||
| : releaseFuture(release.get_future()), callbackFinished(false), | ||
| requested(0), received(0) | ||
| {} | ||
|
|
||
| virtual void onFrameRequested(MediaFrame &frame) | ||
| { | ||
| ++requested; | ||
| holdCallback(); | ||
| frame.type = PJMEDIA_FRAME_TYPE_AUDIO; | ||
| frame.buf.assign(frame.size, 0x5a); | ||
| callbackFinished = true; | ||
| } | ||
|
|
||
| virtual void onFrameReceived(MediaFrame &frame) | ||
| { | ||
| ++received; | ||
| CHECK(frame.type == PJMEDIA_FRAME_TYPE_AUDIO); | ||
| CHECK(frame.buf.size() == 16); | ||
| holdCallback(); | ||
| callbackFinished = true; | ||
| } | ||
|
|
||
| static void run(Endpoint &ep, bool receive) | ||
| { | ||
| AudioMediaPortTest media; | ||
| media.detachCallbacks(); /* No port yet. */ | ||
|
|
||
| MediaFormatAudio fmt; | ||
| fmt.init(PJMEDIA_FORMAT_L16, 8000, 1, 20000, 16); | ||
| media.createPort("callback-fence", fmt); | ||
| pjmedia_port *port = media.port; | ||
|
|
||
| unsigned char buffer[16]; | ||
| pjmedia_frame frame; | ||
| pj_bzero(&frame, sizeof(frame)); | ||
| pj_memset(buffer, 0xa5, sizeof(buffer)); | ||
| frame.buf = buffer; | ||
| frame.type = PJMEDIA_FRAME_TYPE_AUDIO; | ||
| frame.size = sizeof(buffer); | ||
|
|
||
| std::future<void> entered = media.entered.get_future(); | ||
| std::thread callback([&]() { | ||
| ep.libRegisterThread("fence-callback"); | ||
| pj_status_t status = receive ? | ||
| pjmedia_port_put_frame(port, &frame) : | ||
| pjmedia_port_get_frame(port, &frame); | ||
| CHECK(status == PJ_SUCCESS); | ||
| }); | ||
| CHECK(entered.wait_for(std::chrono::seconds(5)) == | ||
| std::future_status::ready); | ||
|
|
||
| /* Verify that the held callback really owns the fence's lock. */ | ||
| CHECK(pj_grp_lock_tryacquire(port->grp_lock) != PJ_SUCCESS); | ||
|
|
||
| std::promise<void> starting, completed; | ||
| std::future<void> started = starting.get_future(); | ||
| std::future<void> done = completed.get_future(); | ||
| std::thread detacher([&]() { | ||
| ep.libRegisterThread("fence-detach"); | ||
| starting.set_value(); | ||
| media.detachCallbacks(); | ||
| CHECK(media.callbackFinished.load()); | ||
| completed.set_value(); | ||
| }); | ||
| CHECK(started.wait_for(std::chrono::seconds(5)) == | ||
| std::future_status::ready); | ||
| CHECK(done.wait_for(std::chrono::milliseconds(100)) == | ||
| std::future_status::timeout); | ||
|
|
||
| media.release.set_value(); | ||
| CHECK(done.wait_for(std::chrono::seconds(5)) == | ||
| std::future_status::ready); | ||
| detacher.join(); | ||
| callback.join(); | ||
|
|
||
| if (!receive) { | ||
| CHECK(frame.type == PJMEDIA_FRAME_TYPE_AUDIO); | ||
| CHECK(frame.size == sizeof(buffer)); | ||
| CHECK(buffer[0] == 0x5a); | ||
| } | ||
|
|
||
| media.detachCallbacks(); | ||
| /* Fencing must not unregister the port or release its resources. */ | ||
| CHECK(media.getPortInfo().portId == media.getPortId()); | ||
| for (unsigned i = 0; i < 3; ++i) { | ||
| /* Seed a stale audio result; detached reads must not expose it. */ | ||
| pj_memset(buffer, 0xa5, sizeof(buffer)); | ||
| frame.type = PJMEDIA_FRAME_TYPE_AUDIO; | ||
| frame.size = sizeof(buffer); | ||
| CHECK(pjmedia_port_get_frame(port, &frame) == PJ_SUCCESS); | ||
| CHECK(frame.type == PJMEDIA_FRAME_TYPE_NONE); | ||
| CHECK(frame.size == 0); | ||
| /* NONE/0 makes the payload invalid; it need not be overwritten. */ | ||
| for (unsigned j = 0; j < sizeof(buffer); ++j) | ||
| CHECK(buffer[j] == 0xa5); | ||
|
|
||
| frame.type = PJMEDIA_FRAME_TYPE_AUDIO; | ||
| frame.size = sizeof(buffer); | ||
| CHECK(pjmedia_port_put_frame(port, &frame) == PJ_SUCCESS); | ||
| } | ||
| CHECK(media.requested == (receive ? 0u : 1u)); | ||
| CHECK(media.received == (receive ? 1u : 0u)); | ||
| } | ||
|
|
||
| private: | ||
| void holdCallback() | ||
| { | ||
| CHECK(requested + received == 1); | ||
| entered.set_value(); | ||
| CHECK(releaseFuture.wait_for(std::chrono::seconds(5)) == | ||
| std::future_status::ready); | ||
| } | ||
|
|
||
| std::promise<void> entered, release; | ||
| std::future<void> releaseFuture; | ||
| std::atomic<bool> callbackFinished; | ||
| unsigned requested, received; | ||
| }; | ||
|
|
||
| } // namespace pj | ||
|
|
||
| void audioMediaPortTest() | ||
| { | ||
| pj::Endpoint ep; | ||
| ep.libCreate(); | ||
| pj::EpConfig cfg; | ||
| cfg.uaConfig.threadCnt = 0; | ||
| cfg.medConfig.threadCnt = 0; | ||
| cfg.logConfig.level = 2; | ||
| ep.libInit(cfg); | ||
| ep.audDevManager().setNoDev(); | ||
| pj::AudioMediaPortTest::run(ep, false); | ||
| pj::AudioMediaPortTest::run(ep, true); | ||
| ep.libDestroy(); | ||
| std::cout << "AudioMediaPort callback fence tests passed" << std::endl; | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -295,26 +295,26 @@ AudioMediaPort::AudioMediaPort() | |
| AudioMediaPort::~AudioMediaPort() | ||
| { | ||
| PJSUA2_CATCH_IGNORE( unregisterMediaPort() ); | ||
| detachCallbacks(); | ||
| if (port) { | ||
| struct port_data *pdata = static_cast<struct port_data *> | ||
| (port->port_data.pdata); | ||
|
|
||
| /* Make sure port no longer accesses this object in its | ||
| * get/put_frame() callback. | ||
| */ | ||
| if (port->grp_lock) { | ||
| pj_grp_lock_acquire(port->grp_lock); | ||
| pdata->mport = NULL; | ||
| pj_grp_lock_release(port->grp_lock); | ||
| } | ||
|
|
||
| pjmedia_port_destroy(port); | ||
| /* We release the pool later in port.on_destroy since | ||
| * the unregistration is async and may not have completed yet. | ||
| */ | ||
| } | ||
| } | ||
|
|
||
| void AudioMediaPort::detachCallbacks() | ||
| { | ||
| if (port && port->grp_lock) { | ||
| struct port_data *pdata = static_cast<struct port_data *> | ||
| (port->port_data.pdata); | ||
| pj_grp_lock_acquire(port->grp_lock); | ||
| pdata->mport = NULL; | ||
| pj_grp_lock_release(port->grp_lock); | ||
| } | ||
| } | ||
|
|
||
| static pj_status_t get_frame(pjmedia_port *port, pjmedia_frame *frame) | ||
| { | ||
| struct port_data *pdata = static_cast<struct port_data *> | ||
|
|
@@ -323,8 +323,11 @@ static pj_status_t get_frame(pjmedia_port *port, pjmedia_frame *frame) | |
| MediaFrame frame_; | ||
|
|
||
| pj_grp_lock_acquire(port->grp_lock); | ||
| if ((mport = pdata->mport) == NULL) | ||
| if ((mport = pdata->mport) == NULL) { | ||
| frame->type = PJMEDIA_FRAME_TYPE_NONE; | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. This fixes more than stale metadata. The conference bridge doesn't initialise |
||
| frame->size = 0; | ||
| goto on_return; | ||
| } | ||
|
|
||
| frame_.size = (unsigned)frame->size; | ||
| mport->onFrameRequested(frame_); | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Minor: "permanently" doesn't hold when it's called before
createPort(). It's a no-op then, andcreatePort()sets the back-pointer, so callbacks are dispatched afterwards. Something like "once the port has been created" would make that clear.