diff --git a/OpenSim/Addons/os-webrtc-janus/Janus/JanusMessages.cs b/OpenSim/Addons/os-webrtc-janus/Janus/JanusMessages.cs index af3e4b3515..233d58aa98 100644 --- a/OpenSim/Addons/os-webrtc-janus/Janus/JanusMessages.cs +++ b/OpenSim/Addons/os-webrtc-janus/Janus/JanusMessages.cs @@ -566,7 +566,7 @@ namespace osWebRtcVoice } public int ParticipantId { get { return PluginRespDataInt("id"); } } - //public long ParticipantId { get { return PluginRespDataLong("id"); } } +// public long ParticipantId { get { return PluginRespDataLong("id"); } } } // ============================================================== diff --git a/OpenSim/Addons/os-webrtc-janus/Janus/JanusRoom.cs b/OpenSim/Addons/os-webrtc-janus/Janus/JanusRoom.cs index dd824199c5..ce811573b7 100644 --- a/OpenSim/Addons/os-webrtc-janus/Janus/JanusRoom.cs +++ b/OpenSim/Addons/os-webrtc-janus/Janus/JanusRoom.cs @@ -76,7 +76,7 @@ namespace osWebRtcVoice JanusMessageResp resp = await _AudioBridge.SendPluginMsg(joinReq); AudioBridgeJoinRoomResp joinResp = new(resp); - // if (joinResp is not null && joinResp.AudioBridgeReturnCode == "joined" && joinResp.ParticipantId > 0) +// if (joinResp is not null && joinResp.AudioBridgeReturnCode == "joined" && joinResp.ParticipantId > 0) if (joinResp is not null && joinResp.AudioBridgeReturnCode == "joined") { pVSession.ParticipantId = joinResp.ParticipantId; diff --git a/OpenSim/Addons/os-webrtc-janus/Janus/JanusSession.cs b/OpenSim/Addons/os-webrtc-janus/Janus/JanusSession.cs index a26c9bcd08..c26476513d 100644 --- a/OpenSim/Addons/os-webrtc-janus/Janus/JanusSession.cs +++ b/OpenSim/Addons/os-webrtc-janus/Janus/JanusSession.cs @@ -26,17 +26,17 @@ */ using System; +using System.Collections.Concurrent; using System.Collections.Generic; using System.Net.Http; using System.Net.Mime; using System.Reflection; using System.Threading.Tasks; +using OpenMetaverse; using OpenMetaverse.StructuredData; using log4net; -using log4net.Core; -using System.Reflection.Metadata; using System.Threading; namespace osWebRtcVoice @@ -136,37 +136,37 @@ namespace osWebRtcVoice if (resp is not null && resp.isSuccess) { // Note that setting IsConnected to false will cause the long poll to exit - m_log.DebugFormat("{0} DestroySession. Destroyed", LogHeader); + m_log.Debug($"{LogHeader} DestroySession. Destroyed"); } else { if (resp.isError) { - ErrorResp eResp = new ErrorResp(resp); + ErrorResp eResp = new(resp); switch (eResp.errorCode) { case 458: // This is the error code for a session that is already destroyed - m_log.DebugFormat("{0} DestroySession: session already destroyed", LogHeader); + m_log.Debug($"{LogHeader} DestroySession: session already destroyed"); break; case 459: // This is the error code for handle already destroyed - m_log.DebugFormat("{0} DestroySession: Handle not found", LogHeader); + if (_MessageDetails) m_log.Debug($"{LogHeader} DestroySession: Handle not found"); break; default: - m_log.ErrorFormat("{0} DestroySession: failed {1}", LogHeader, eResp.errorReason); + m_log.Error($"{LogHeader} DestroySession: failed {eResp.errorReason}"); break; } } else { - m_log.ErrorFormat("{0} DestroySession: failed. Resp: {1}", LogHeader, resp.ToString()); + m_log.Error($"{LogHeader} DestroySession: failed. Resp: {resp}"); } } } catch (Exception e) { - m_log.ErrorFormat("{0} DestroySession: exception {1}", LogHeader, e); + m_log.Error($"{LogHeader} DestroySession: exception ", e); } IsConnected = false; _CancelTokenSource.Cancel(); @@ -181,11 +181,11 @@ namespace osWebRtcVoice // if the audiobridge is active, the trickle message is sent to it if (pVSession.AudioBridge is null) { - ret = await SendToJanusNoWait(new TrickleReq(pVSession)); + ret = await SendToJanusNoWait(new TrickleReq(pVSession, pCandidates)); } else { - ret = await SendToJanusNoWait(new TrickleReq(pVSession), pVSession.AudioBridge.PluginUri); + ret = await SendToJanusNoWait(new TrickleReq(pVSession, pCandidates), pVSession.AudioBridge.PluginUri); } return ret; } @@ -223,7 +223,7 @@ namespace osWebRtcVoice public DateTime RequestTime; public TaskCompletionSource TaskCompletionSource; } - private Dictionary _OutstandingRequests = new Dictionary(); + private readonly ConcurrentDictionary _OutstandingRequests = new(); // Send a request directly to the Janus server. // NOTE: this is probably NOT what you want to do. This is a direct call that is outside the session. @@ -256,11 +256,11 @@ namespace osWebRtcVoice RequestTime = DateTime.Now, TaskCompletionSource = new TaskCompletionSource() }; - _OutstandingRequests.Add(pReq.TransactionId, outReq); + _OutstandingRequests[pReq.TransactionId] = outReq; string reqStr = pReq.ToJson(); - HttpRequestMessage reqMsg = new HttpRequestMessage(HttpMethod.Post, pURI); + HttpRequestMessage reqMsg = new(HttpMethod.Post, pURI); reqMsg.Content = new StringContent(reqStr, System.Text.Encoding.UTF8, MediaTypeNames.Application.Json); reqMsg.Headers.Add("Accept", "application/json"); HttpResponseMessage response = await _HttpClient.SendAsync(reqMsg, _CancelTokenSource.Token); @@ -273,31 +273,41 @@ namespace osWebRtcVoice { // Some messages are asynchronous and completed with an event if (_MessageDetails) m_log.DebugFormat("{0} SendToJanus: ack response {1}", LogHeader, respStr); +/* if (_OutstandingRequests.TryGetValue(pReq.TransactionId, out OutstandingRequest outstandingRequest)) { ret = await outstandingRequest.TaskCompletionSource.Task; _OutstandingRequests.Remove(pReq.TransactionId); } // If there is no OutstandingRequest, the request was not waiting for an event or already processed +*/ + + // Wait on the local TaskCompletionSource instead of re-reading the dictionary. + // This avoids a race where the long-poll thread already removed the request. + ret = await outReq.TaskCompletionSource.Task; } else { // If the response is not an ack, that means a synchronous request/response so return the response - _OutstandingRequests.Remove(pReq.TransactionId); + _= _OutstandingRequests.TryRemove(pReq.TransactionId, out _); if (_MessageDetails) m_log.DebugFormat("{0} SendToJanus: response {1}", LogHeader, respStr); } } else { m_log.ErrorFormat("{0} SendToJanus: response not successful {1}", LogHeader, response); - _OutstandingRequests.Remove(pReq.TransactionId); + _= _OutstandingRequests.TryRemove(pReq.TransactionId, out _); } } catch (Exception e) { m_log.ErrorFormat("{0} SendToJanus: exception {1}", LogHeader, e.Message); + _= _OutstandingRequests.TryRemove(pReq.TransactionId, out _); + } + finally + { + _= _OutstandingRequests.TryRemove(pReq.TransactionId, out _); } - return ret; } /// @@ -313,8 +323,9 @@ namespace osWebRtcVoice AddJanusHeaders(pReq); - try { - HttpRequestMessage reqMsg = new HttpRequestMessage(HttpMethod.Post, pURI); + try + { + HttpRequestMessage reqMsg = new(HttpMethod.Post, pURI); string reqStr = pReq.ToJson(); reqMsg.Content = new StringContent(reqStr, System.Text.Encoding.UTF8, MediaTypeNames.Application.Json); reqMsg.Headers.Add("Accept", "application/json"); @@ -338,14 +349,14 @@ namespace osWebRtcVoice private void AddJanusHeaders(JanusMessageReq pReq) { // Authentication token - if (!String.IsNullOrEmpty(_JanusAPIToken)) + if (!string.IsNullOrEmpty(_JanusAPIToken)) { pReq.AddAPIToken(_JanusAPIToken); } // Transaction ID that matches responses to requests - if (String.IsNullOrEmpty(pReq.TransactionId)) + if (string.IsNullOrEmpty(pReq.TransactionId)) { - pReq.TransactionId = Guid.NewGuid().ToString(); + pReq.TransactionId = UUID.Random().ToString(); } // The following two are required for the WebSocket interface. They are optional for the // HTTP interface since the session and plugin handle are in the URL. @@ -363,22 +374,17 @@ namespace osWebRtcVoice bool TryGetOutstandingRequest(string pTransactionId, out OutstandingRequest pOutstandingRequest) { - if (String.IsNullOrEmpty(pTransactionId)) + if (string.IsNullOrEmpty(pTransactionId)) { pOutstandingRequest = null; return false; } - bool ret = false; - lock (_OutstandingRequests) - { - if (_OutstandingRequests.TryGetValue(pTransactionId, out pOutstandingRequest)) - { - _OutstandingRequests.Remove(pTransactionId); - ret = true; - } - } - return ret; + if(_OutstandingRequests.TryGetValue(pTransactionId, out pOutstandingRequest)) + return true; + + pOutstandingRequest = null; + return false; } public Task SendToJanusAdmin(JanusMessageReq pReq) @@ -398,7 +404,7 @@ namespace osWebRtcVoice /// public async Task GetFromJanus(string pURI) { - if (!String.IsNullOrEmpty(_JanusAPIToken)) + if (!string.IsNullOrEmpty(_JanusAPIToken)) { pURI += "?apisecret=" + _JanusAPIToken; } @@ -438,15 +444,15 @@ namespace osWebRtcVoice } catch (TaskCanceledException e) { - m_log.DebugFormat("{0} GetFromJanus: task canceled: {1}", LogHeader, e.Message); - var eResp = new ErrorResp("GETERROR"); + if (_MessageDetails) m_log.DebugFormat("{0} GetFromJanus: task canceled: {1}", LogHeader, e.Message); + ErrorResp eResp = new("GETERROR"); eResp.SetError(499, "Task canceled"); ret = eResp; } catch (Exception e) { m_log.ErrorFormat("{0} GetFromJanus: exception {1}", LogHeader, e.Message); - var eResp = new ErrorResp("GETERROR"); + ErrorResp eResp = new("GETERROR"); eResp.SetError(400, "Exception: " + e.Message); ret = eResp; } @@ -454,7 +460,7 @@ namespace osWebRtcVoice catch (Exception e) { m_log.ErrorFormat("{0} GetFromJanus: exception {1}", LogHeader, e); - var eResp = new ErrorResp("GETERROR"); + ErrorResp eResp = new("GETERROR"); eResp.SetError(400, "Exception: " + e.Message); ret = eResp; } @@ -541,26 +547,54 @@ namespace osWebRtcVoice break; case "webrtcup": // ICE and DTLS succeeded, and so Janus correctly established a PeerConnection with the user/application; - m_log.DebugFormat("{0} EventLongPoll: webrtcup {1}", LogHeader, resp.ToString()); +// m_log.DebugFormat("{0} EventLongPoll: webrtcup {1}", LogHeader, resp.ToString()); + string webrtcupSessionId = ExtractLogId(resp, "session_id"); + string webrtcupSender = ExtractLogId(resp, "sender"); + m_log.DebugFormat("{0} EventLongPoll: webrtcup session_id={1}, sender={2}", + LogHeader, + string.IsNullOrEmpty(webrtcupSessionId) ? "" : webrtcupSessionId, + string.IsNullOrEmpty(webrtcupSender) ? "" : webrtcupSender); break; case "hangup": // The PeerConnection was closed, either by the user/application or by Janus itself; // If one is in the room, when a "hangup" event happens, it means that the user left the room. - m_log.DebugFormat("{0} EventLongPoll: hangup {1}", LogHeader, resp.ToString()); +// m_log.DebugFormat("{0} EventLongPoll: hangup {1}", LogHeader, resp.ToString()); + m_log.DebugFormat("{0} EventLongPoll: hangup session_id={1}, sender={2}, reason={3}, tx={4}", + LogHeader, + ExtractLogId(resp, "session_id"), + ExtractLogId(resp, "sender"), + resp.RawBody.TryGetString("reason", out string hangupReason) ? hangupReason : string.Empty, + FormatTransactionId(resp.TransactionId)); OnHangup?.Invoke(eventResp); break; case "detached": // a plugin asked the core to detach one of our handles - m_log.DebugFormat("{0} EventLongPoll: event {1}", LogHeader, resp.ToString()); +// m_log.DebugFormat("{0} EventLongPoll: event {1}", LogHeader, resp.ToString()); + m_log.DebugFormat("{0} EventLongPoll: detached session_id={1}, sender={2}, tx={3}", + LogHeader, + ExtractLogId(resp, "session_id"), + ExtractLogId(resp, "sender"), + FormatTransactionId(resp.TransactionId)); OnDetached?.Invoke(eventResp); break; case "media": // Janus is receiving (receiving: true/false) audio/video (type: "audio/video") on this PeerConnection; - m_log.DebugFormat("{0} EventLongPoll: media {1}", LogHeader, resp.ToString()); +// m_log.DebugFormat("{0} EventLongPoll: media {1}", LogHeader, resp.ToString()); + m_log.DebugFormat("{0} EventLongPoll: media session_id={1}, sender={2}, tx={3}", + LogHeader, + ExtractLogId(resp, "session_id"), + ExtractLogId(resp, "sender"), + FormatTransactionId(resp.TransactionId)); + break; case "slowlink": // Janus detected a slowlink (uplink: true/false) on this PeerConnection; - m_log.DebugFormat("{0} EventLongPoll: slowlink {1}", LogHeader, resp.ToString()); +// m_log.DebugFormat("{0} EventLongPoll: slowlink {1}", LogHeader, resp.ToString()); + m_log.DebugFormat("{0} EventLongPoll: slowlink session_id={1}, sender={2}, tx={3}", + LogHeader, + ExtractLogId(resp, "session_id"), + ExtractLogId(resp, "sender"), + FormatTransactionId(resp.TransactionId)); break; case "error": m_log.DebugFormat("{0} EventLongPoll: error {1}", LogHeader, resp.ToString()); @@ -622,7 +656,7 @@ namespace osWebRtcVoice break; case 499: // "Task canceled" means the long poll was canceled - m_log.DebugFormat("{0} EventLongPoll: Task canceled. URI={1}", LogHeader, SessionUri); + if (_MessageDetails) m_log.DebugFormat("{0} EventLongPoll: Task canceled. URI={1}", LogHeader, SessionUri); break; default: m_log.DebugFormat("{0} EventLongPoll: unknown response. URI={1}: {2}", @@ -651,8 +685,37 @@ namespace osWebRtcVoice m_log.ErrorFormat("{0} EventLongPoll: exception {1}", LogHeader, e); } } - m_log.InfoFormat("{0} EventLongPoll: Exiting long poll loop", LogHeader); + if (_MessageDetails) + m_log.InfoFormat("{0} EventLongPoll: Exiting long poll loop", LogHeader); }); } + + private static string ExtractLogId(JanusMessageResp resp, string key) + { + if (resp?.RawBody is null || !resp.RawBody.TryGetValue(key, out OSD value) || value is null) + return string.Empty; + + try + { + switch (value.Type) + { + case OSDType.Integer: + case OSDType.Binary: + case OSDType.Array: + return value.AsLong().ToString(); + default: + return value.AsString(); + } + } + catch + { + return value.AsString(); + } + } + + private static string FormatTransactionId(string transactionId) + { + return string.IsNullOrEmpty(transactionId) ? "" : transactionId; + } } } \ No newline at end of file diff --git a/OpenSim/Addons/os-webrtc-janus/Janus/JanusViewerSession.cs b/OpenSim/Addons/os-webrtc-janus/Janus/JanusViewerSession.cs index 89bd236ba6..936bca0c0e 100644 --- a/OpenSim/Addons/os-webrtc-janus/Janus/JanusViewerSession.cs +++ b/OpenSim/Addons/os-webrtc-janus/Janus/JanusViewerSession.cs @@ -84,7 +84,7 @@ namespace osWebRtcVoice // IVoiceViewerSession.Shutdown public async Task Shutdown() { - m_log.DebugFormat($"{LogHeader} JanusViewerSession shutdown {ViewerSessionID}"); + m_log.Debug($"{LogHeader} JanusViewerSession shutdown {ViewerSessionID}"); if (Room is not null) { var rm = Room;