Skip to content
Open
1 change: 1 addition & 0 deletions lib/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ endif()

if(MATSDK_BUILD_JNI_WRAPPER)
list(APPEND SRCS
jni/JavaDataViewerProxy.cpp
jni/JniConvertors.cpp
jni/LogManager_jni.cpp
jni/Logger_jni.cpp
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@
import com.microsoft.applications.events.DebugEventType;
import com.microsoft.applications.events.DiagLevel;
import com.microsoft.applications.events.HttpClient;
import com.microsoft.applications.events.IDataViewer;
import com.microsoft.applications.events.ILogConfiguration;
import com.microsoft.applications.events.ILogManager;
import com.microsoft.applications.events.ILogger;
Expand All @@ -42,7 +43,10 @@
import java.util.SortedMap;
import java.util.TreeMap;
import java.util.TreeSet;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.FutureTask;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;

import org.junit.Test;
import org.junit.runner.RunWith;
Expand Down Expand Up @@ -247,6 +251,119 @@ public void startDDVonLogManager() {
LogManager.flushAndTeardown();
}

@Test
public void registerDataViewer_whenCallbackThrows_continuesDispatchAndStopsAfterUnregister()
throws Exception {
System.loadLibrary("maesdk");
Context appContext = InstrumentationRegistry.getInstrumentation().getTargetContext();
if (s_client == null) {
s_client = new MockHttpClient(appContext);
}
OfflineRoom.connectContext(appContext);

final String token =
"0123456789abcdef9123456789abcdef-01234567-0123-0123-0123-0123456789ab-0124";
final String factoryName = "JavaDataViewer" + System.nanoTime();
ILogConfiguration custom = LogManager.logConfigurationFactory();
custom.set(LogConfigurationKey.CFG_STR_PRIMARY_TOKEN, token);
custom.set(LogConfigurationKey.CFG_STR_COLLECTOR_URL, "https://viewer.contoso.com/");
custom.set(LogConfigurationKey.CFG_STR_FACTORY_NAME, factoryName);
custom.set(LogConfigurationKey.CFG_STR_CACHE_FILE_PATH, factoryName);

ILogManager manager = LogManagerProvider.createLogManager(custom);
CountDownLatch receivedPacket = new CountDownLatch(1);
AtomicInteger receivedByteCount = new AtomicInteger();
AtomicInteger receivingViewerCalls = new AtomicInteger();
AtomicInteger throwingViewerCalls = new AtomicInteger();
IDataViewer throwingViewer =
new IDataViewer() {
@Override
public void receiveData(byte[] packetData) {
throwingViewerCalls.incrementAndGet();
throw new IllegalStateException("Expected callback failure");
}

@Override
public String getName() {
return "throwing-viewer";
}

@Override
public boolean isTransmissionEnabled() {
return true;
}

@Override
public String getCurrentEndpoint() {
return "";
}
};
IDataViewer receivingViewer =
new IDataViewer() {
@Override
public void receiveData(byte[] packetData) {
receivingViewerCalls.incrementAndGet();
receivedByteCount.set(packetData.length);
receivedPacket.countDown();
}

@Override
public String getName() {
return "receiving-viewer";
}

@Override
public boolean isTransmissionEnabled() {
return true;
}

@Override
public String getCurrentEndpoint() {
return "http://127.0.0.1";
}
};

try {
assertThat(manager.registerDataViewer(throwingViewer), is(true));
assertThat(manager.registerDataViewer(receivingViewer), is(true));
assertThat(manager.registerDataViewer(receivingViewer), is(false));

ILogger logger = manager.getLogger(token, "java-data-viewer-test", "");
logger.logEvent("javaDataViewerCallback");
manager.uploadNow();

assertThat(receivedPacket.await(5, TimeUnit.SECONDS), is(true));
assertThat(receivedByteCount.get(), greaterThan(0));

assertThat(manager.unregisterDataViewer("receiving-viewer"), is(true));
assertThat(manager.unregisterDataViewer("receiving-viewer"), is(false));
Comment thread
Copilot marked this conversation as resolved.

// Unregistering must actually stop callbacks, not merely drop the bookkeeping entry: a
// bridge that left the proxy in the native DataViewerCollection would still pass the
// assertions above. Drive a second dispatch and use the still-registered throwing viewer
// as the witness that one really occurred, then assert the unregistered viewer was not
// called again.
final int receivingCallsAtUnregister = receivingViewerCalls.get();
final int throwingCallsAtUnregister = throwingViewerCalls.get();

logger.logEvent("javaDataViewerCallbackAfterUnregister");
manager.uploadNow();

final long deadline = System.currentTimeMillis() + 10000;
while (throwingViewerCalls.get() <= throwingCallsAtUnregister
&& System.currentTimeMillis() < deadline) {
Thread.sleep(50);
}

assertThat(throwingViewerCalls.get(), greaterThan(throwingCallsAtUnregister));
assertThat(receivingViewerCalls.get(), is(receivingCallsAtUnregister));

assertThat(manager.unregisterDataViewer("throwing-viewer"), is(true));
} finally {
manager.close();
}
}

/*
Disabling this test since it requires private modules.

Expand Down
4 changes: 4 additions & 0 deletions lib/android_build/maesdk/consumer-rules.pro
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
-keep interface com.microsoft.applications.events.IDataViewer { *; }
-keep class * implements com.microsoft.applications.events.IDataViewer {
public *;
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
//
// Copyright (c) Microsoft Corporation. All rights reserved.
// SPDX-License-Identifier: Apache-2.0
//
package com.microsoft.applications.events;

import androidx.annotation.Keep;

/**
* Receives copies of packets uploaded by the SDK.
*
* <p>Implementations must return a stable, unique name for the lifetime of the registration.
* Callbacks can occur on an SDK worker thread and should return promptly. Implementations must not
* reenter the SDK from within a callback: do not register or unregister viewers, and do not close
* the owning {@link ILogManager}, because closing unregisters every viewer while the callback is
* still in progress.
*/
@Keep
public interface IDataViewer {

/** Receives an encoded telemetry packet after it has been prepared for upload. */
void receiveData(byte[] packetData);

/** Returns the stable, unique name used to register this viewer. */
String getName();

/** Returns whether this viewer is currently accepting packet callbacks. */
boolean isTransmissionEnabled();

/** Returns the endpoint currently used by this viewer, or an empty string when disabled. */
String getCurrentEndpoint();
}
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,34 @@ public interface ILogManager extends AutoCloseable {

public String getCurrentEndpoint();

/**
* Registers a caller-provided data viewer with this LogManager.
*
* <p>This is an optional capability. The default implementation returns {@code false} so that
* existing implementations of this interface remain source compatible; implementations that
* support data viewers override it.
*
* @return {@code true} when the viewer was registered, {@code false} for invalid input, a
* duplicate viewer name, or when the implementation does not support data viewers
*/
default boolean registerDataViewer(IDataViewer dataViewer) {
return false;
}

/**
* Unregisters a caller-provided data viewer by its unique name.
*
* <p>This is an optional capability. The default implementation returns {@code false} so that
* existing implementations of this interface remain source compatible; implementations that
* support data viewers override it.
*
* @return {@code true} when the viewer was unregistered, {@code false} when it was not
* registered, or when the implementation does not support data viewers
*/
default boolean unregisterDataViewer(String viewerName) {
return false;
}

public LogSessionData getLogSessionData();

public void setLevelFilter(int defaultLevel, int[] allowedLevels);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -226,6 +226,28 @@ public String getCurrentEndpoint() {
return nativeGetCurrentEndpoint(nativeLogManager);
}

protected native boolean nativeRegisterDataViewer(
long nativeLogManager, IDataViewer dataViewer);

@Override
public boolean registerDataViewer(IDataViewer dataViewer) {
if (dataViewer == null) {
return false;
}
return nativeRegisterDataViewer(nativeLogManager, dataViewer);
}

protected native boolean nativeUnregisterDataViewer(
long nativeLogManager, String viewerName);

@Override
public boolean unregisterDataViewer(String viewerName) {
if (viewerName == null || viewerName.isEmpty()) {
return false;
}
return nativeUnregisterDataViewer(nativeLogManager, viewerName);
}

protected static class LogSessionDataImpl implements LogSessionData {
@Keep
private long m_first_time;
Expand Down
49 changes: 42 additions & 7 deletions lib/api/DataViewerCollection.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@
//
#include "DataViewerCollection.hpp"
#include <algorithm>
#include <cstring>
#include <mutex>

namespace MAT_NS_BEGIN {
Expand All @@ -12,11 +13,34 @@ namespace MAT_NS_BEGIN {

void DataViewerCollection::DispatchDataViewerEvent(const std::vector<uint8_t>& packetData) const noexcept
{
if (IsViewerEnabled() == false)
// Dispatch over a snapshot taken under the lock, and release the lock before invoking any
// viewer. Iterating m_dataViewerCollection directly is unsafe because m_dataViewerMapLock
// is recursive: a viewer that reenters the SDK from ReceiveData - for example by closing
Comment thread
KartikDhawaniya marked this conversation as resolved.
// the owning LogManager, which unregisters every viewer - would erase from the very vector
// being iterated here and invalidate the iterator. Holding the lock across a callback is
// unsafe for a second reason: registration acquires the JNI viewer mutex and then this
// lock, so a callback that reenters registration closes a lock cycle, and any slow callback
// would stall registration, unregistration and LogManager close until it returned.
// The shared_ptr copies keep each viewer alive for the duration of its own callback, even
// if it is unregistered - or loses its last other reference - while dispatch is running.
std::vector<std::shared_ptr<IDataViewer>> viewers;
{
LOCKGUARD(m_dataViewerMapLock);
viewers = m_dataViewerCollection;
}

// The enabled check runs on this same snapshot, outside the lock, for two reasons.
// IsTransmissionEnabled() is a viewer callback - a JNI call for Java viewers - and must not
// run under m_dataViewerMapLock for the reasons above. Reusing the one snapshot also means
// the set of viewers this decision is made about is the set that is dispatched to; calling
// IsViewerEnabled() first would lock a second time and could decide on a different set.
// Keep this predicate in sync with IsViewerEnabled().
const bool anyEnabled = std::any_of(viewers.cbegin(), viewers.cend(),
[](const std::shared_ptr<IDataViewer>& viewer) { return viewer->IsTransmissionEnabled(); });
if (!anyEnabled)
return;

LOCKGUARD(m_dataViewerMapLock);
for(const auto& viewer : m_dataViewerCollection)
for(const auto& viewer : viewers)
{
// Task 3568800: Integrate ThreadPool to IDataViewerCollection
viewer->ReceiveData(packetData);
Expand Down Expand Up @@ -52,7 +76,7 @@ namespace MAT_NS_BEGIN {
LOCKGUARD(m_dataViewerMapLock);
auto toErase = std::find_if(m_dataViewerCollection.begin(), m_dataViewerCollection.end(), [&viewerName](std::shared_ptr<IDataViewer> viewer)
{
return viewer->GetName() == viewerName;
return strcmp(viewer->GetName(), viewerName) == 0;
});

if (toErase == m_dataViewerCollection.end())
Expand All @@ -79,9 +103,20 @@ namespace MAT_NS_BEGIN {

bool DataViewerCollection::IsViewerEnabled() const noexcept
{
LOCKGUARD(m_dataViewerMapLock);
return !m_dataViewerCollection.empty() &&
std::find_if(m_dataViewerCollection.begin(), m_dataViewerCollection.end(), [](std::shared_ptr<IDataViewer> viewer) { return viewer->IsTransmissionEnabled(); }) != m_dataViewerCollection.end();
// Evaluate over a snapshot taken under the lock. IsTransmissionEnabled() is a viewer
// callback - for Java viewers it crosses into the JVM - and must not run while
// m_dataViewerMapLock is held: registration takes the JNI viewer mutex and then this lock,
// so a callback that reenters the SDK would close a lock cycle, and a slow callback would
// stall registration, unregistration and LogManager close.
// Keep this predicate in sync with DispatchDataViewerEvent().
std::vector<std::shared_ptr<IDataViewer>> viewers;
{
LOCKGUARD(m_dataViewerMapLock);
viewers = m_dataViewerCollection;
}

return std::any_of(viewers.cbegin(), viewers.cend(),
[](const std::shared_ptr<IDataViewer>& viewer) { return viewer->IsTransmissionEnabled(); });
}

bool DataViewerCollection::IsViewerRegistered(const char* viewerName) const
Expand Down
Loading
Loading