diff --git a/java/android/nnstreamer/src/androidTest/java/org/nnsuite/nnstreamer/APITestCustomFilter.java b/java/android/nnstreamer/src/androidTest/java/org/nnsuite/nnstreamer/APITestCustomFilter.java index bed70849..27f49a5a 100644 --- a/java/android/nnstreamer/src/androidTest/java/org/nnsuite/nnstreamer/APITestCustomFilter.java +++ b/java/android/nnstreamer/src/androidTest/java/org/nnsuite/nnstreamer/APITestCustomFilter.java @@ -16,8 +16,8 @@ */ @RunWith(AndroidJUnit4.class) public class APITestCustomFilter { - private int mReceived = 0; - private boolean mInvalidState = false; + private volatile int mReceived = 0; + private volatile boolean mInvalidState = false; private boolean mRegistered = false; private CustomFilter mCustomPassthrough; private CustomFilter mCustomConvert; @@ -396,4 +396,235 @@ public TensorsData invoke(TensorsData in) { /* expected */ } } + + @Test + public void testCloseWhileUsed_n() { + TensorsInfo inputInfo = new TensorsInfo(); + inputInfo.addTensorInfo(NNStreamer.TensorType.INT32, new int[]{10}); + + TensorsInfo outputInfo = inputInfo.clone(); + + CustomFilter filter = CustomFilter.create("custom-close-while-used-n", + inputInfo, outputInfo, new CustomFilter.Callback() { + @Override + public TensorsData invoke(TensorsData in) { + return in; + } + }); + + String desc = "appsrc name=srcx ! " + + "other/tensor,dimension=(string)10,type=(string)int32,framerate=(fraction)0/1 ! " + + "tensor_filter framework=custom-easy model=" + filter.getName() + " ! " + + "tensor_sink name=sinkx"; + + try (Pipeline pipe = new Pipeline(desc)) { + /* filter is referenced by the pipeline: close() must fail */ + mInvalidState = false; + try { + filter.close(); + mInvalidState = true; + } catch (IllegalStateException e) { + /* expected */ + } + + assertFalse(mInvalidState); + + /* the filter must still be registered and usable */ + pipe.registerSinkCallback("sinkx", new Pipeline.NewDataCallback() { + @Override + public void onNewDataReceived(TensorsData data) { + if (data == null || data.getTensorsCount() != 1) { + mInvalidState = true; + return; + } + + ByteBuffer output = data.getTensorData(0); + + for (int j = 0; j < 10; j++) { + if (output.getInt(j * 4) != j) { + mInvalidState = true; + } + } + + mReceived++; + } + }); + + /* start pipeline */ + pipe.start(); + + /* push input buffer repeatedly */ + for (int i = 0; i < 10; i++) { + TensorsData in = TensorsData.allocate(inputInfo); + ByteBuffer input = in.getTensorData(0); + + for (int j = 0; j < 10; j++) { + input.putInt(j * 4, j); + } + + in.setTensorData(0, input); + + pipe.inputData("srcx", in); + Thread.sleep(20); + } + + /* sleep 300 to pass all input buffers to sink */ + Thread.sleep(300); + + /* stop pipeline */ + pipe.stop(); + + /* check received data from sink */ + assertFalse(mInvalidState); + assertEquals(10, mReceived); + } catch (Exception e) { + fail(); + } + + /* the pipeline is closed now: close() should succeed */ + try { + filter.close(); + } catch (Exception e) { + fail(); + } + + /* the name is really unregistered, a new filter with the same name should succeed */ + try { + CustomFilter again = CustomFilter.create(filter.getName(), + inputInfo, outputInfo, new CustomFilter.Callback() { + @Override + public TensorsData invoke(TensorsData in) { + return in; + } + }); + + again.close(); + } catch (Exception e) { + fail(); + } + } + + @Test + public void testCloseAfterPipeline() { + TensorsInfo inputInfo = new TensorsInfo(); + inputInfo.addTensorInfo(NNStreamer.TensorType.INT32, new int[]{10}); + + TensorsInfo outputInfo = inputInfo.clone(); + + CustomFilter filter = CustomFilter.create("custom-close-after-pipeline", + inputInfo, outputInfo, new CustomFilter.Callback() { + @Override + public TensorsData invoke(TensorsData in) { + return in; + } + }); + + String desc = "appsrc name=srcx ! " + + "other/tensor,dimension=(string)10,type=(string)int32,framerate=(fraction)0/1 ! " + + "tensor_filter framework=custom-easy model=" + filter.getName() + " ! " + + "tensor_sink name=sinkx"; + + /* construct and close a pipeline using the filter */ + try (Pipeline pipe = new Pipeline(desc)) { + /* nothing to do here */ + } catch (Exception e) { + fail(); + } + + /* the pipeline no longer references the filter, close() should succeed */ + try { + filter.close(); + + CustomFilter again = CustomFilter.create(filter.getName(), + inputInfo, outputInfo, new CustomFilter.Callback() { + @Override + public TensorsData invoke(TensorsData in) { + return in; + } + }); + + again.close(); + } catch (Exception e) { + fail(); + } + } + + @Test + public void testCloseTwice() { + TensorsInfo inputInfo = new TensorsInfo(); + inputInfo.addTensorInfo(NNStreamer.TensorType.INT32, new int[]{10}); + + TensorsInfo outputInfo = inputInfo.clone(); + + CustomFilter filter = CustomFilter.create("custom-close-twice", + inputInfo, outputInfo, new CustomFilter.Callback() { + @Override + public TensorsData invoke(TensorsData in) { + return in; + } + }); + + try { + /* not used by any pipeline, close() should succeed */ + filter.close(); + + /* closing an already closed filter is a no-op */ + filter.close(); + + CustomFilter again = CustomFilter.create(filter.getName(), + inputInfo, outputInfo, new CustomFilter.Callback() { + @Override + public TensorsData invoke(TensorsData in) { + return in; + } + }); + + again.close(); + } catch (Exception e) { + fail(); + } + } + + @Test + public void testCloseWhileUsedNotStarted_n() { + TensorsInfo inputInfo = new TensorsInfo(); + inputInfo.addTensorInfo(NNStreamer.TensorType.INT32, new int[]{10}); + + TensorsInfo outputInfo = inputInfo.clone(); + + CustomFilter filter = CustomFilter.create("custom-close-while-used-not-started-n", + inputInfo, outputInfo, new CustomFilter.Callback() { + @Override + public TensorsData invoke(TensorsData in) { + return in; + } + }); + + String desc = "appsrc name=srcx ! " + + "other/tensor,dimension=(string)10,type=(string)int32,framerate=(fraction)0/1 ! " + + "tensor_filter framework=custom-easy model=" + filter.getName() + " ! " + + "tensor_sink name=sinkx"; + + try (Pipeline pipe = new Pipeline(desc)) { + /* the pipeline is never started, but the filter is still referenced */ + mInvalidState = false; + try { + filter.close(); + mInvalidState = true; + } catch (IllegalStateException e) { + /* expected */ + } + + assertFalse(mInvalidState); + } catch (Exception e) { + fail(); + } + + /* the pipeline is closed now: close() should succeed */ + try { + filter.close(); + } catch (Exception e) { + fail(); + } + } } diff --git a/java/android/nnstreamer/src/androidTest/java/org/nnsuite/nnstreamer/APITestPipeline.java b/java/android/nnstreamer/src/androidTest/java/org/nnsuite/nnstreamer/APITestPipeline.java index 06a01101..6cb67c9d 100644 --- a/java/android/nnstreamer/src/androidTest/java/org/nnsuite/nnstreamer/APITestPipeline.java +++ b/java/android/nnstreamer/src/androidTest/java/org/nnsuite/nnstreamer/APITestPipeline.java @@ -12,6 +12,8 @@ import java.io.File; import java.nio.ByteBuffer; import java.util.Arrays; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicInteger; import static org.junit.Assert.*; @@ -20,8 +22,8 @@ */ @RunWith(AndroidJUnit4.class) public class APITestPipeline { - private int mReceived = 0; - private boolean mInvalidState = false; + private volatile int mReceived = 0; + private volatile boolean mInvalidState = false; private Pipeline.State mPipelineState = Pipeline.State.NULL; private Pipeline.NewDataCallback mSinkCb = new Pipeline.NewDataCallback() { @@ -2407,4 +2409,141 @@ public void onNewDataReceived(TensorsData data) { fail(); } } + + @Test + public void testCloseWithoutStop() { + String desc = "videotestsrc ! video/x-raw,format=RGB,width=4,height=4,framerate=1000/1 ! " + + "tensor_converter ! tee name=t " + + "t. ! queue ! tensor_sink name=sinka " + + "t. ! queue ! tensor_sink name=sinkb"; + + boolean playingObserved = false; + + for (int iter = 0; iter < 10; iter++) { + final AtomicInteger receivedA = new AtomicInteger(0); + final AtomicInteger receivedB = new AtomicInteger(0); + final AtomicBoolean seenPlaying = new AtomicBoolean(false); + Pipeline pipe = null; + + /* pipeline state callback */ + Pipeline.StateChangeCallback stateCb = new Pipeline.StateChangeCallback() { + @Override + public void onStateChanged(Pipeline.State state) { + if (state == Pipeline.State.PLAYING) { + seenPlaying.set(true); + } + } + }; + + try { + pipe = new Pipeline(desc, stateCb); + + /* two streaming threads deliver buffers to these sinks concurrently */ + pipe.registerSinkCallback("sinka", new Pipeline.NewDataCallback() { + @Override + public void onNewDataReceived(TensorsData data) { + if (data == null) { + mInvalidState = true; + return; + } + + receivedA.incrementAndGet(); + } + }); + + pipe.registerSinkCallback("sinkb", new Pipeline.NewDataCallback() { + @Override + public void onNewDataReceived(TensorsData data) { + if (data == null) { + mInvalidState = true; + return; + } + + receivedB.incrementAndGet(); + } + }); + + /* start pipeline */ + pipe.start(); + + /* wait until both sinks are actively receiving buffers */ + int waited = 0; + while ((receivedA.get() < 3 || receivedB.get() < 3) && waited < 3000) { + Thread.sleep(10); + waited += 10; + } + + assertTrue(receivedA.get() >= 3); + assertTrue(receivedB.get() >= 3); + + /* close the pipeline while buffers are still flowing, without calling stop() first */ + pipe.close(); + pipe = null; + } catch (Exception e) { + fail(); + } finally { + if (pipe != null) { + pipe.close(); + } + } + + if (seenPlaying.get()) { + playingObserved = true; + } + } + + assertTrue(playingObserved); + assertFalse(mInvalidState); + } + + @Test + public void testCloseWithoutStopSingleSink() { + String desc = "videotestsrc ! video/x-raw,format=RGB,width=4,height=4,framerate=1000/1 ! " + + "tensor_converter ! tensor_sink name=sinkx"; + + for (int iter = 0; iter < 10; iter++) { + final AtomicInteger received = new AtomicInteger(0); + Pipeline pipe = null; + + try { + pipe = new Pipeline(desc); + + pipe.registerSinkCallback("sinkx", new Pipeline.NewDataCallback() { + @Override + public void onNewDataReceived(TensorsData data) { + if (data == null) { + mInvalidState = true; + return; + } + + received.incrementAndGet(); + } + }); + + /* start pipeline */ + pipe.start(); + + /* wait until the sink is actively receiving buffers */ + int waited = 0; + while (received.get() < 3 && waited < 3000) { + Thread.sleep(10); + waited += 10; + } + + assertTrue(received.get() >= 3); + + /* close the pipeline while buffers are still flowing, without calling stop() first */ + pipe.close(); + pipe = null; + } catch (Exception e) { + fail(); + } finally { + if (pipe != null) { + pipe.close(); + } + } + } + + assertFalse(mInvalidState); + } } diff --git a/java/android/nnstreamer/src/main/java/org/nnsuite/nnstreamer/CustomFilter.java b/java/android/nnstreamer/src/main/java/org/nnsuite/nnstreamer/CustomFilter.java index c305409b..9912901e 100644 --- a/java/android/nnstreamer/src/main/java/org/nnsuite/nnstreamer/CustomFilter.java +++ b/java/android/nnstreamer/src/main/java/org/nnsuite/nnstreamer/CustomFilter.java @@ -19,7 +19,7 @@ public final class CustomFilter implements AutoCloseable { private Callback mCallback = null; private native long nativeInitialize(String name, TensorsInfo in, TensorsInfo out); - private native void nativeDestroy(long handle); + private native boolean nativeDestroy(long handle); /** * Interface definition for a callback to be invoked while processing the pipeline. @@ -125,10 +125,23 @@ protected void finalize() throws Throwable { } } + /** + * Unregisters the custom-filter and releases its resources. + * + * A custom-filter used in a pipeline cannot be released until the pipeline is closed. + * In that case the custom-filter is kept registered, and it can be closed again + * after closing the pipeline. + * + * @throws IllegalStateException if failed to unregister the custom-filter, + * e.g., it is still used in a pipeline + */ @Override public void close() { if (mHandle != 0) { - nativeDestroy(mHandle); + if (!nativeDestroy(mHandle)) { + throw new IllegalStateException("Failed to close custom-filter " + mName + ", it may be used in a pipeline"); + } + mHandle = 0; } } diff --git a/java/android/nnstreamer/src/main/jni/nnstreamer-native-api.c b/java/android/nnstreamer/src/main/jni/nnstreamer-native-api.c index bdefd0b5..01840d66 100644 --- a/java/android/nnstreamer/src/main/jni/nnstreamer-native-api.c +++ b/java/android/nnstreamer/src/main/jni/nnstreamer-native-api.c @@ -74,12 +74,6 @@ nns_free_element_data (gpointer data) element_data_s *item = (element_data_s *) data; if (item) { - /* release private data */ - if (item->priv_data) { - JNIEnv *env = nns_get_jni_env (item->pipe_info); - item->priv_destroy_func (item->priv_data, env); - } - switch (item->type) { #if !defined (NNS_SINGLE_ONLY) case NNS_ELEMENT_TYPE_SRC: @@ -105,6 +99,12 @@ nns_free_element_data (gpointer data) break; } + /* release private data after the handle, a running callback may use it */ + if (item->priv_data) { + JNIEnv *env = nns_get_jni_env (item->pipe_info); + item->priv_destroy_func (item->priv_data, env); + } + g_free (item->name); g_free (item); } @@ -234,22 +234,30 @@ nns_construct_pipe_info (JNIEnv * env, jobject thiz, gpointer handle, /** * @brief Destroy pipeline info. + * @return TRUE if pipe info is released. FALSE if the custom-filter cannot be + * unregistered (e.g., it is used in a pipeline); pipe info is left + * untouched and can be destroyed again later. It is leaked if the + * caller never succeeds in unregistering the filter, which is intended: + * the unregister keeps the filter handle on every failure, so releasing + * pipe info would leave the filter invoking a released user data. */ -void +gboolean nns_destroy_pipe_info (pipeline_info_s * pipe_info, JNIEnv * env) { - g_return_if_fail (pipe_info != NULL); + g_return_val_if_fail (pipe_info != NULL, FALSE); - g_mutex_lock (&pipe_info->lock); - if (pipe_info->priv_data) { - if (pipe_info->priv_destroy_func) - pipe_info->priv_destroy_func (pipe_info->priv_data, env); - else - g_free (pipe_info->priv_data); - - pipe_info->priv_data = NULL; +#if !defined (NNS_SINGLE_ONLY) + /* a registered custom-filter keeps pipe info as its user data */ + if (pipe_info->pipeline_type == NNS_PIPE_TYPE_CUSTOM && + pipe_info->pipeline_handle && + ml_pipeline_custom_easy_filter_unregister (pipe_info->pipeline_handle) + != ML_ERROR_NONE) { + _ml_loge ("Failed to unregister the custom-filter, it may be in use."); + return FALSE; } +#endif + g_mutex_lock (&pipe_info->lock); g_hash_table_destroy (pipe_info->element_handles); pipe_info->element_handles = NULL; g_mutex_unlock (&pipe_info->lock); @@ -260,7 +268,7 @@ nns_destroy_pipe_info (pipeline_info_s * pipe_info, JNIEnv * env) ml_pipeline_destroy (pipe_info->pipeline_handle); break; case NNS_PIPE_TYPE_CUSTOM: - ml_pipeline_custom_easy_filter_unregister (pipe_info->pipeline_handle); + /* already unregistered */ break; #if defined(ENABLE_ML_SERVICE) case NNS_PIPE_TYPE_SERVICE: @@ -278,6 +286,17 @@ nns_destroy_pipe_info (pipeline_info_s * pipe_info, JNIEnv * env) break; } + g_mutex_lock (&pipe_info->lock); + if (pipe_info->priv_data) { + if (pipe_info->priv_destroy_func) + pipe_info->priv_destroy_func (pipe_info->priv_data, env); + else + g_free (pipe_info->priv_data); + + pipe_info->priv_data = NULL; + } + g_mutex_unlock (&pipe_info->lock); + g_mutex_clear (&pipe_info->lock); nns_destroy_tensors_data_cls_info (env, &pipe_info->tensors_data_cls_info); @@ -288,6 +307,7 @@ nns_destroy_pipe_info (pipeline_info_s * pipe_info, JNIEnv * env) pthread_key_delete (pipe_info->jni_env); g_free (pipe_info); + return TRUE; } /** diff --git a/java/android/nnstreamer/src/main/jni/nnstreamer-native-customfilter.c b/java/android/nnstreamer/src/main/jni/nnstreamer-native-customfilter.c index 1c3d5c38..5e5fff9a 100644 --- a/java/android/nnstreamer/src/main/jni/nnstreamer-native-customfilter.c +++ b/java/android/nnstreamer/src/main/jni/nnstreamer-native-customfilter.c @@ -212,13 +212,13 @@ nns_native_custom_initialize (JNIEnv * env, jobject thiz, jstring name, /** * @brief Native method for custom filter. */ -static void +static jboolean nns_native_custom_destroy (JNIEnv * env, jobject thiz, jlong handle) { pipeline_info_s *pipe_info = NULL; pipe_info = CAST_TO_TYPE (handle, pipeline_info_s *); - nns_destroy_pipe_info (pipe_info, env); + return nns_destroy_pipe_info (pipe_info, env) ? JNI_TRUE : JNI_FALSE; } /** @@ -227,7 +227,7 @@ nns_native_custom_destroy (JNIEnv * env, jobject thiz, jlong handle) static JNINativeMethod native_methods_customfilter[] = { {(char *) "nativeInitialize", (char *) "(Ljava/lang/String;L" NNS_CLS_TINFO ";L" NNS_CLS_TINFO ";)J", (void *) nns_native_custom_initialize}, - {(char *) "nativeDestroy", (char *) "(J)V", + {(char *) "nativeDestroy", (char *) "(J)Z", (void *) nns_native_custom_destroy} }; diff --git a/java/android/nnstreamer/src/main/jni/nnstreamer-native-internal.h b/java/android/nnstreamer/src/main/jni/nnstreamer-native-internal.h index 8777fcb2..28a2b433 100644 --- a/java/android/nnstreamer/src/main/jni/nnstreamer-native-internal.h +++ b/java/android/nnstreamer/src/main/jni/nnstreamer-native-internal.h @@ -190,7 +190,7 @@ nns_construct_pipe_info (JNIEnv * env, jobject thiz, gpointer handle, nns_pipe_t /** * @brief Destroy pipeline info. */ -extern void +extern gboolean nns_destroy_pipe_info (pipeline_info_s * pipe_info, JNIEnv * env); /** diff --git a/tests/capi/unittest_capi_inference.cc b/tests/capi/unittest_capi_inference.cc index d9afdee8..401894c8 100644 --- a/tests/capi/unittest_capi_inference.cc +++ b/tests/capi/unittest_capi_inference.cc @@ -793,6 +793,124 @@ TEST (nnstreamer_capi_sink, register_duplicated) g_free (pipe_state); } +/** + * @brief Data for the sink callback blocking the streaming thread. + */ +typedef struct { + GMutex lock; + GCond cond; + gboolean entered; + gboolean unregistering; + gboolean returned; + guint count; +} TestSinkBlocking; + +/** + * @brief A sink callback that holds the streaming thread at its first call, + * until the test is about to unregister the sink. + */ +static void +test_sink_callback_blocking ( + const ml_tensors_data_h data, const ml_tensors_info_h info, void *user_data) +{ + TestSinkBlocking *sink = (TestSinkBlocking *) user_data; + gboolean first; + gint64 end_time; + + g_mutex_lock (&sink->lock); + sink->count++; + first = !sink->entered; + sink->entered = TRUE; + g_cond_broadcast (&sink->cond); + + if (first) { + end_time = g_get_monotonic_time () + 5 * G_TIME_SPAN_SECOND; + while (!sink->unregistering) { + if (!g_cond_wait_until (&sink->cond, &sink->lock, end_time)) + break; + } + } + g_mutex_unlock (&sink->lock); + + if (first) { + /* give the unregister time to reach the element lock */ + g_usleep (100000); + + g_mutex_lock (&sink->lock); + sink->returned = TRUE; + g_mutex_unlock (&sink->lock); + } +} + +/** + * @brief Test NNStreamer pipeline sink + * @detail Unregistering a sink waits for the running callback, and no callback + * is called after it returns. Android JNI releases the user data of the + * sink callback right after the unregister, so it relies on this. + */ +TEST (nnstreamer_capi_sink, unregister_wait_callback) +{ + ml_pipeline_h handle; + ml_pipeline_sink_h sinkhandle; + TestSinkBlocking sink = {}; + gint64 end_time; + gboolean entered, returned; + guint count; + int status; + const gchar *pipeline = "videotestsrc is-live=true ! videoconvert ! video/x-raw,format=RGB,width=4,height=4,framerate=30/1 ! tensor_converter ! tensor_sink name=sinkx"; + + g_mutex_init (&sink.lock); + g_cond_init (&sink.cond); + + status = ml_pipeline_construct (pipeline, NULL, NULL, &handle); + ASSERT_EQ (status, ML_ERROR_NONE); + + status = ml_pipeline_sink_register ( + handle, "sinkx", test_sink_callback_blocking, &sink, &sinkhandle); + EXPECT_EQ (status, ML_ERROR_NONE); + + status = ml_pipeline_start (handle); + EXPECT_EQ (status, ML_ERROR_NONE); + + end_time = g_get_monotonic_time () + 5 * G_TIME_SPAN_SECOND; + g_mutex_lock (&sink.lock); + while (!sink.entered) { + if (!g_cond_wait_until (&sink.cond, &sink.lock, end_time)) + break; + } + entered = sink.entered; + sink.unregistering = TRUE; + g_cond_broadcast (&sink.cond); + g_mutex_unlock (&sink.lock); + EXPECT_TRUE (entered); + + /* the first callback is still running here */ + status = ml_pipeline_sink_unregister (sinkhandle); + EXPECT_EQ (status, ML_ERROR_NONE); + + g_mutex_lock (&sink.lock); + returned = sink.returned; + count = sink.count; + g_mutex_unlock (&sink.lock); + EXPECT_TRUE (returned); + + /* the pipeline is still playing, but the callback is not called anymore */ + g_usleep (200000); + + g_mutex_lock (&sink.lock); + EXPECT_EQ (count, sink.count); + g_mutex_unlock (&sink.lock); + + status = ml_pipeline_stop (handle); + EXPECT_EQ (status, ML_ERROR_NONE); + + status = ml_pipeline_destroy (handle); + EXPECT_EQ (status, ML_ERROR_NONE); + + g_cond_clear (&sink.cond); + g_mutex_clear (&sink.lock); +} + /** * @brief Test NNStreamer pipeline sink * @detail Failure case to register callback with invalid param.