using NewLife.Log;
using NewLife.Messaging;
using NewLife.Model;
namespace NewLife.Net;
/// <summary>å“应匹é…器。请求-å“应é…对的通用支撑:ç‰å¾…登记ã€åŒ¹é…交付ã€å…³é—清ç†ï¼Œä¸Žä¼ è¾“å½¢æ€æ— å…³</summary>
/// <remarks>
/// <para>匹é…ä¸ä»¥æ¶ˆæ¯æ³µä¸ºå‰æï¼Œè€Œæ˜¯äº¤ä»˜é˜¶æ®µçš„能力:æµå¼ä¼šè¯ï¼ˆTCPï¼‰ç”±æ¶ˆæ¯æ³µåœ¨å¸§å®šç•Œäº¤ä»˜æ—¶è°ƒç”¨ <see cref="TryMatch"/>ï¼›
/// æ•°æ®æŠ¥ä¼šè¯ï¼ˆUDP,æ¯åŒ…å³å®Œæ•´å¸§ã€æ— 粘包处ç†ï¼‰åœ¨æŠ¥æ–‡äº¤ä»˜æ—¶è°ƒç”¨ï¼ŒäºŒè€…共用本组件与åŒä¸€åŒ¹é…队列,ç‰å¾…与超时è¯ä¹‰ä¸€è‡´ã€‚</para>
/// <para>队列首次ç‰å¾…æ—¶è‡ªåŠ¨åˆ›å»ºï¼Œå¯æ³¨å…¥å…±äº«æˆ–自定义实现(<see cref="Queue"/>)。</para>
/// </remarks>
sealed class ResponseMatcher
{
#region 属性
/// <summary>匹é…队列。首次ç‰å¾…æ—¶è‡ªåŠ¨åˆ›å»ºï¼Œå¯æ³¨å…¥å…±äº«æˆ–自定义实现</summary>
public IMatchQueue? Queue
{
get => _queue;
set => _queue = value;
}
private IMatchQueue? _queue;
/// <summary>请求-å“应匹é…ç‰å¾…超时(毫秒)。默认30_000</summary>
public Int32 Timeout { get; set; } = 30_000;
#endregion
#region ç‰å¾…
/// <summary>登记请求ç‰å¾…。调用方éšåŽé€å‡ºè¯·æ±‚,å“åº”åˆ°è¾¾æ—¶ç» <see cref="TryMatch"/> 交付</summary>
/// <param name="owner">拥有者(会è¯å®žä¾‹ï¼‰ã€‚åŒ¹é…æŒ‰æ‹¥æœ‰è€…过滤,收到的消æ¯åªäº¤ä»˜åŒæ‹¥æœ‰è€…çš„ç‰å¾…</param>
/// <param name="request">请求消æ¯</param>
/// <param name="span">å…³è”的性能追踪 Span,éšç‰å¾…æ–¹ç¾æ”¶è€Œé‡Šæ”¾</param>
/// <returns>ç‰å¾…æºã€‚调用方负责把 <see cref="PooledValueTaskSource{T}.ValueTask"/> 交给ç‰å¾…方,é€å‡ºå¤±è´¥æ—¶ <c>TrySetException</c></returns>
public PooledValueTaskSource<IMessage> Register(Object owner, IMessage request, ISpan? span)
{
// å¹¶å‘首用时åªèƒ½æœ‰ä¸€ä¸ªé˜Ÿåˆ—胜出:败者丢弃自建实例,å¦åˆ™è¯·æ±‚ä¼šå…¥é˜Ÿåˆ°æ— äººåŒ¹é…的队列,åªèƒ½ç‰è¶…æ—¶
var queue = _queue;
if (queue == null)
{
var created = new DefaultMatchQueue();
queue = Interlocked.CompareExchange(ref _queue, created, null) ?? created;
}
var source = PooledValueTaskSource<IMessage>.Rent();
source.AttachSpan(span);
queue.Add(owner, request, Timeout, source);
return source;
}
#endregion
#region 匹é…
/// <summary>å°è¯•把收到的消æ¯åŒ¹é…ç»™ç‰å¾…ä¸çš„请求。命ä¸åˆ™äº¤ä»˜ç‰å¾…方并返回 true</summary>
/// <param name="owner">拥有者(与登记时一致)</param>
/// <param name="codec">å议编解ç 器,é…对能力在内层时é€å±‚解包</param>
/// <param name="message">收到的消æ¯</param>
/// <returns>是å¦å·²åŒ¹é…交付。未命ä¸ã€åè®®æ— é…对能力或ç‰å¾…æ–¹å·²å¤±æ•ˆï¼ˆå–æ¶ˆã€ç‰å¾…æºå·²å¤ç”¨ï¼‰æ—¶è¿”回 falseï¼Œè°ƒç”¨æ–¹æ®æ¤æŒ‰æ™®é€šæ¶ˆæ¯æ”¶å°¾</returns>
public Boolean TryMatch(Object owner, IMessageCodec? codec, IMessage message)
{
var queue = _queue;
if (queue == null) return false;
var matcher = Resolve(codec);
if (matcher == null) return false;
return queue.Match(owner, message, message, (req, resp) =>
req is IMessage rq && resp is IMessage rs && matcher.Match(rq, rs));
}
/// <summary>å–出å议的请求-å“应é…对能力。装饰å议(压缩/åŠ å¯†ç‰ï¼‰æŠŠé…对能力留给内层,需é€å±‚解包</summary>
/// <param name="codec">åè®®</param>
/// <returns>é…对器;åè®®ä¸æ”¯æŒé…对时返回 null</returns>
public static IMessageMatcher? Resolve(IMessageCodec? codec) => codec switch
{
IMessageMatcher matcher => matcher,
IMessageCodecDecorator decorator => Resolve(decorator.Inner),
_ => null,
};
#endregion
#region 清ç†
/// <summary>清空队列,唤醒全部ç‰å¾…方(会è¯å…³é—时调用,é¿å…调用方悬挂)</summary>
public void Clear() => _queue?.Clear();
#endregion
}
|