بدائي صغير جداً يحول العمل المتزامن إلى نظام متكيف منسق
"نموذج الإشارات الإشارية"
داخل الجزء الأول لقد بنينا تنفيذًا عابرًا للحدود خاصًّا ذاتيًّا ونظيفًا ذاتيًّا لسير العمل. الجزء الثاني حوّلنا ذلك إلى مكتبة قابلة لإعادة الاستخدام مع منسقين، خطوط أنابيب رئيسية، وتكامل DI.
وتضيف هذه المادة سمة صغيرة تغير كل شيء: موجـز.
هذا هو الآن في حزمة النغمات أيضاً أكثر من 20 أكثر من 20 في الغالب.
الشفرة المصدرية الكاملة موجودة في ThTMs GetTHub
وتعيش البنية التحتية للإشارات في ما يلي:
□ الملف / الغرض من
| ------ | --------- |
|---|---|
| التنفيذية التنفيذية التنفيذية التنفيذية التنفيذية □ إشارة انبعاثات وانحسار من العمليات | |
أولاً- اجتماع الدول الأطراف في ملاحظة- تشكيلاتCancelOnSignals, DeferOnSignals, OnSignal, OnSignalRetracted) |
|
| سلسلة سِطر سِطر P P P P Matn Mattchr.cs □ مطابقة نمط نمط أسلوب غلوب لرشح الإشارة □ □ □ □ □ □ □ □ □ | |
اشارة الاشارة توجيه الإشارة Assyync مع مطابقة النمط (الدعمات) *, ?، قوائم فاصلة ، ترتيب تحديدي) / |
|
| الأمثلة/الأمثلة/الأمثلة المتسلسلـة: HttttpClient.cs □ نموذج من انبعاثات الإشارة المثبتة بالدقة الدقيقة لدعوات HTTP □ | |
| أمثلة/الأمثلة/الإرشادات. □ الحد من معدل التكيف مع التأجيل القائم على الإشارة | |
| الأمثلة/الأمثلة/الأمثلة على الشبكة البحرية نافذة الإشارة الإشارية الإقلاعية | |
| أمثلة/TetelephysSignal Handler.cs □ معالجة الإشارة Asyync مع التكامل عن بعد |
منسقينا الأقوياء رائعون في معالجة العمل، لكنهم معزولون. كل منسق يعرف عن عملياته الخاصة، ولكن ليس لديه أي وعي بما يحدث في مكان آخر في النظام.
// Translation coordinator has no idea that...
await translationCoordinator.EnqueueAsync(request);
// ...the API just hit a rate limit
// ...another service is experiencing backpressure
// ...a downstream dependency is slow
يمكننا أن نربط نقاط اعتماد واضحة، ولكن هذا يخلق القران. ما نريده هو **** - المنسقون الذين يمكن أن يشعروا ببيئتهم دون أن يكونوا متصلين بشكل مباشر.
الاشارات تسمح بذرات التنفيذ تترك آثاراً في نافذتها الإكهارية. تلك الاثار تتصرف مثل الحقائق القصيرة العمر:
ويمكن للمنسقين عندئذ أن يغيروا سلوكهم استناداً إلى الإشارات المرئية في نافذتهم.
من طراز من طراز إشارات الوصلات:
public readonly record struct SignalEvent(
string Signal,
long OperationId,
string? Key,
DateTimeOffset Timestamp,
SignalPropagation? Propagation = null)
{
public int Depth => Propagation?.Depth ?? 0;
public bool WouldCycle(string signal) => Propagation?.Contains(signal) == true;
public bool Is(string name) => Signal == name;
public bool StartsWith(string prefix) => Signal.StartsWith(prefix, StringComparison.Ordinal);
}
يمكن للعملية أن ترفع الإشارات أثناء التنفيذ. هذه الإشارات تعيش في النافذة الإكهارية بجانب العملية. عندما تتقدّم العملية، فإن الإشارات تذهب معها.
هذه هي، لا رسالة سمسار، لا بنية تحتية منفصلة، فقط شروط مرتبطة بالعمليات
وبما ان الاشارات تعيش داخل النافذة الايرجوانية ، فهي ترث ضماناتها : الحجم المُحَدَّد ، الشيخوخة التلقائية ، وصفر دورة الحياة .
والإشارة، بحد ذاتها، لا تسبب الإعدام. إنه يسجل فقط حقيقة في النافذة الزهرية لا شيء يجري لأن الإشارة انبعثت
فإذا عرف المنسق معالجاً على الموقع، فإنه يعمل متزامناً عندما تنبعث إشارة - ولكن فقط لأن المنسق اختار إرفاقها. الحضانات لا تعرف أو تهتم وإزالة جميع المعالجين يترك السلوك الأساسي دون تغيير.
ولا تُربط الإشارة إلا بالذرة/العملية التي تنبعث منها. لا توجد إشارة من أي وقت مضى تحوّل أو يُذكر ذرة أخرى. لا يوجد باص مشترك قابل للكتابة
وَهذِهِ ٱلْأَشْرَارَةُ تُخَلِّصُ حَقِيقَةً مِنْ تَأْريخِ ٱلْحَدَثَةِ . الإشارات لا تُحدّث أو تُكبّر أبدًا. لا تزيل النكهات الا اشارات المنبعث .
ولا توجد الإشارات إلا داخل النافذة البحرية للمنسق. وينتهي مفعولها تلقائياً كما تنتهي نافذة الخروج. لا شيء يستمر إلا إذا كنت صراحة بناء الثبات.
عندما يقوم المنسق "بفحص الإشارات للمفتاح K"، فإنه يقوم بالمسح:
الأُقْرَات في نافذتها
« والآيات » الش دال دال دال دالها على كل ما خلقها من أنواعها « والآيات » الدالات دال دالًا ، دال دال دالها ، و قدرته دالها ، و قدرته دال دالها ، و قدرته دالها ، و قدرته دالها ، و قدرته دالها ، و قدرته دالها ، و قدرته دالها ، و قدرته دالها ، و قدرته دالها ، و قدرته دالها ، و قدرته دالها ، و قدرته دال دالها ، و قدرته دال دالها ، و قدرته ، و قدرته دال دال دال دال دال و و و دال و و و و و و دال دال دال دال دال دال دال و و و و و و دال دال دال دال و و و و دال دال دال دال و و و دال و و و دال دال و و و دال دال دال و و دال و و دال و و و دال و و دال و و و و دال دال دال دال دال و و و دال دال و و و و دال و و و و و دال دال دال دالً و و و و و و و و دال دال دال دال و و و و و دال دال دال و و و و و و دال دال دال و و و و و و و و و و دال دال دال دال دال دال دال دال و و و و و و و و و و و و و و و دال دال دال دال دال دال دال دال دال دال و و و و و و و و و و و و و و و و و دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال و و و و و و و و و و و و و و و و و و و و و و و و و و و و دال دال دال دال دال دال دال دال دال دال و و و و و و و و و و و و و و و و و و دال دال دال دال دال دال دال دال دال دال دال و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و و دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال دال و و و و و و و و لا تُعدِّل المراقبة ابدا حالة إشارة الذرَّة .
إذا كنت تريد إشارات إلى قيادة Ayync سير الأعمال، يجب أن تستخدم:
مُشَرَرْرِك
سوزسور
أو تكييفات أخرى.
هذه طبقات اختيارية بنيت على قمة الإشارات، وليس جزءاً من دلالاتها.
فلا يمكن لأي ذرة ان تغيِّر صورة ذرة اخرى او حالتها او اشاراتها او معلوماتها الوصفية . يحدث التنسيق من خلال:
موجـز
الاستشعار
نواجوات
ليس من خلال الكتابة.
يمكن لبوصة اشارات او سطح تجميعي ان يُظهِر رؤية مشتركة للاشارات — — لكن دائماً ما تكون قراءة فقط، لا يمكن أن تكون ذات حجية، لا يمكن كتابتها أبداً.
إذا كانت جميع المعالجات (على الإشارة، المرسلات، المجهزات) منفصلة عن بعضها البعض، (ب) أن يظل النظام صحيحاً وقابلاً للتنبؤ به تماماً. لا تزال الإشارات لها معنى لأنها حقائق وليست محفزات.
أحداث أحداث مثل مكالمة هاتفية:
"سأتصل بكِ الآن، ردي و ردي"
الوصلس مثل آثار أقدام في الثلج:
"تركت آثار أقدام، إذا أردت أن تعرف أين ذهبت، انظر، إن لم تكن تهتم، تجاهلها."
هذا هو السبب في أن الإشارات لا تكسر أبداً أبداً لا تكسر، لا تحجب، ولا تتفاعل أبداً مع تدفق التحكم & حرّر لطرحهم.
فالتمييز يعني أن الإشارات تزيل ما يلي:
/ المشكلة / الأحداث لديها إشارات تجنبها / |---------|:--------------:|:----------------:| الناشر/الناشر/الناشر/الناشر/الناشر/الناشر/الناشر □ وقت الاعتماد على التوقيت □ يتطلب رد فعل فوري □ Ambient، استطلاع عندما يكون جاهزا □ يمكن لليد العاملة أن تحفز المعالجات لا توجد حلقات إلا إذا طُلب منها ذلك صراحةً • الأمر بالدلالة على المسميات □ يجب تسليم/مناولة ضمانات التسليم □ لا تسليم، مجرد وجود □ □ خطأ في النشر □ أخطاء في استعمال اليد تنتشر □ معزول □ □ مخاطر الاندراج المشتركة فشل مُسلّط واحد فشل، تكسر سلسلة لا سلسلة لكسر
□ المعالم □ الأحداث □ الإشارات □ |---------|--------|---------| (بآلاف دولارات الولايات المتحدة) ● التوقيت أو التوقيت المباشر وزأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأأ ● التسليم: مضمون/ممتنع □ لا تسليم، مجرد وجود □ رد فعل على رد الفعل/الخيار الاختياري حمولة البيانات: بيانات فوقية ثقيلة في كثير من الأحيان عن الخيوط الخيطية الصغيرة لا يوجد نافذة من طراز LRU « حياة حياة » o : « لا حياة » عدد كبير تقريباً لا شيء
نهج الأحداث (كلسيكي):
public event Action RateLimited;
try
{
await CallApiAsync();
}
catch (RateLimitException)
{
RateLimited?.Invoke(); // Makes someone else act right now
}
المشاكل:
نهج الاشارات )القمر( :
try
{
await CallApiAsync(ct);
}
catch (RateLimitException ex)
{
op.Signal($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
لا شيء يحدث هنا، في وقت لاحق، في مكان مختلف تماماً:
if (translator.HasSignal("rate-limit"))
{
await Task.Delay(1000); // Act when *we* choose
}
مراقبة عمليات النقل، وسياق عمليات نقل الإشارات.
هذا هو النموذج العقلي بأكمله في جملة واحدة
تبعاً لخلفية معلوماتك الخلفية:
/ الجمهور / التعريف / |----------|------------| | محركات نظم • ورشة تخطيطية للتنسيق غير المباشر | **** □ بيانات فوقية خفيفة للوزن ملحقة بالعمليات في نافذة انزلاقية مقيّدة □ | PL/BCL □ سبور ضمني مؤقت مُزَوَّج مع عُرَض مُحَدَّدة مُحَدَّدة | أسخ الإطار □ سطح دولة في طور التشغيل، ذاتي التطهير، يمكنك الاستعلام عنه في أي وقت □
الإشارات هي آثار تركت في سطح ذاكرة مُشتركة ومُحَدَّدة. أي شخص يمكنه النظر إليها. لا أحد مُلزَم بالتفاعل. هذا أساسي. سُججج نموذج التنسيق - نفس استخدام النمل الواحد، نفس واحد نظم السبورة المستخدمة في أوائل AI، ونفس واحد الحديثة CRDT الثرثرة شبكات تلميح.
using var activity = source.StartActivity("ProcessOrder");
activity?.SetTag("order.id", orderId);
activity?.SetTag("rate.limited", true);
أفضل لـ: تم توزيع التتبع عبر الخدمات، والتخزين الطويل الأجل عن بعد، وهويات الارتباط.
حداطاطاطاطاطاط: تحتاج إلى تتبع الطلبات عبر خدمات متعددة، أو تخزين مقاييس للتحليل، أو التكامل مع أدوات الرصد.
استخدام إشارات متى: تحتاج في عملية عملية التوعية المحيطة، التنسيق التفاعلي، أو لا تريد البنية التحتية عن بعد.
var rateLimits = Observable.FromEventPattern<RateLimitEventArgs>(
h => api.RateLimitHit += h,
h => api.RateLimitHit -= h);
rateLimits
.Throttle(TimeSpan.FromSeconds(1))
.Subscribe(e => HandleRateLimit(e));
أفضل لـ: معالجة الأحداث المعقدة، والعمليات القائمة على أساس الزمن، والجمع بين جداول الأحداث المتعددة.
****: تحتاج إلى إستفسارات زمنية معقدة (النافذة، والفك، والجمع بين الجداول).
استخدام إشارات متى: تريد استشعاراً بسيطاً قائماً على الاقتراع، أو تنظيفاً آلياً، أو دمجاً مع تتبع العمليات.
public class RateLimitNotification : INotification
{
public int RetryAfterMs { get; init; }
}
await _mediator.Publish(new RateLimitNotification { RetryAfterMs = 5000 });
أفضل لـتم فصلها في عملية التعامل مع الحدث مع العديد من المعالجين.
****تريد العديد من المعالجين ليتفاعلوا مع نفس الحدث بشكل متزامن
استخدام إشارات متى:أنت تريد استشعاراً محيطياً بدون اشتراك صريح، أو تاريخ تنظيف ذاتي، أو تكامل مع تنفيذ مرتبط.
var circuitBreaker = Policy
.Handle<HttpRequestException>()
.CircuitBreakerAsync(5, TimeSpan.FromSeconds(30));
أفضل لـ: القدرة على التأقلم حول المكالمات الفردية مع الإدارة التلقائية للدولة.
****: تحتاج إلى القدرة على التكيف مع حالات الانتقال التلقائية نصف المفتوحة/المغلقة.
استخدام إشارات متى: تريد الوعي المحيط عبر العديد من العمليات ، منطق الدائرة المخصصة ، أو التكامل مع تتبع العملية.
يُمَزّجُهماستخدام بولي داخل جسم العمل الخاص بك، والانبعاثات الإشارات عندما تسافر الدوائر.
□ النهج القائم على الانتقــال الذاتــي □ Ambients sespecuted o checopeded change change change change change change |----------|:-------------:|:---------------:|:---------:|:------------:|:-----------:| • أدوات خارجية □ □ أدوات خارجية □ □ ▪ الامتدادات الامتدادات الراجعة □ الوسيط/المتوسط/المتوسط/المتوسط/المتوسط/المتوسط/الدليل/الدليل/الدليل/الدليل/الدليل/الدليل/الدليل/الدليل/الدليل
تنفيذ تنفيذ العمليات ISignalEmitter:
public interface ISignalEmitter
{
// Emit signals
void Emit(string signal);
bool EmitCaused(string signal, SignalPropagation? cause);
// Retract (remove) signals
bool Retract(string signal);
int RetractMatching(string pattern);
bool HasSignal(string signal);
long OperationId { get; }
string? Key { get; }
}
داخل جسم عملك:
await coordinator.ProcessAsync(async (item, op, ct) =>
{
try
{
var result = await CallExternalApiAsync(item, ct);
if (result.WasCached)
op.Signal("cache-hit");
if (result.Duration > TimeSpan.FromSeconds(2))
op.Signal("slow-response");
}
catch (RateLimitException ex)
{
op.Signal("rate-limit");
op.Signal($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
catch (TimeoutException)
{
op.Signal("timeout");
throw;
}
});
الإشارات هي مجرد خيوط. استخدم الأسماء البسيطة (% 1)"rate-limit"أو الأسماء المنظمة )("rate-limit:5000ms").
مرشحات نمط الاستخدام*, ?(ج) وقوات الدعم (ج) وقوت الدعم ("error.*,timeout"(المطابقة أمر محدد وضوء توزيع عن طريق StringPatternMatcher.
ويمكنك ان تطلق اشارات في كل مرحلة من مراحل العملية من اجل المراقبة الدقيقة جدا . (أ) أن يبين هذا النمط:
using Mostlylucid.Helpers.Ephemeral.Examples;
// Inside your work body where you have access to the operation's emitter:
await coordinator.ProcessAsync(async (request, op, ct) =>
{
var data = await SignalingHttpClient.DownloadWithSignalsAsync(
httpClient,
new HttpRequestMessage(HttpMethod.Get, request.Url),
op, // ISignalEmitter
ct);
// Process the downloaded data...
});
وهذا يبعث إشارات في كل مرحلة:
□ إشارة □ متى □
|--------|------|
| stage.starting قبل أن يبدأ الطلب
| progress:0 * علامة التقدم الأولي
| stage.request تم إرسال طلب HTTTP
| stage.headers □ تم تلقي صحائف الردود
| stage.reading البدء بقراءة الجسم
| progress:XX النسبة المئوية للتقدم (0-100) أثناء التنزيل
| stage.completed ▪ تنزيل مُنتِج
يمكنك أن تُحال إلى هذه مع مُطابقة نمطيًّا:
// Find all stage transitions
var stages = coordinator.GetSignalsByPattern("stage.*");
// Check download progress
var progress = coordinator.GetSignalsByPattern("progress:*");
// Check if any download is still in progress
if (coordinator.HasSignalMatching("stage.reading") &&
!coordinator.HasSignalMatching("stage.completed"))
{
// Download in progress
}
ويمكن للعمليات أيضاً أن تزيل إشاراتها الخاصة. وهذا مفيد في الحالات المؤقتة:
await coordinator.ProcessAsync(async (item, op, ct) =>
{
// Mark as processing
op.Emit("processing");
try
{
await ProcessItemAsync(item, ct);
// Success - retract the processing signal
op.Retract("processing");
op.Emit("completed");
}
catch (RetryableException)
{
// Keep processing signal, add retry info
op.Emit("retrying");
}
catch (Exception)
{
// Remove all temporary signals
op.RetractMatching("processing*");
op.Emit("failed");
throw;
}
});
تماماً مثل إشارة الانبعاث، يمكن للتراجعات أن تحفز على الاستدعاء:
var coordinator = new EphemeralWorkCoordinator<Request>(
body,
new EphemeralOptions
{
// Sync retraction handler
OnSignalRetracted = evt =>
{
_metrics.DecrementGauge(evt.Signal);
Console.WriteLine($"Signal {evt.Signal} retracted from op {evt.OperationId}");
if (evt.WasPatternMatch)
Console.WriteLine($" (matched pattern: {evt.Pattern})");
},
// Async retraction handler
OnSignalRetractedAsync = async (evt, ct) =>
{
await _telemetry.TrackRetraction(evt.Signal, evt.OperationId, ct);
}
});
الـ SignalRetractedEvent بما في ذلك:
Signal - اسم الإشارة المتراجعOperationId - العملية التي تراجعت عنهKey - مفتاح العملية (إن وجد)Timestamp - عندما حدث تراجعWasPatternMatch - الحق إذا تراجعت عن طريق RetractMatchingPattern - النمط المستخدم (مطابقة النمط)await coordinator.ProcessAsync(async (request, op, ct) =>
{
// Check if we already have a rate limit signal
if (op.HasSignal("rate-limited"))
{
// We're in recovery mode
await Task.Delay(1000, ct);
}
try
{
var response = await _api.SendAsync(request, ct);
// Success! Remove any rate limit signal
if (op.Retract("rate-limited"))
{
op.Emit("rate-limit-cleared");
}
}
catch (RateLimitException ex)
{
op.Emit("rate-limited");
op.Emit($"rate-limit:{ex.RetryAfterMs}ms");
throw;
}
});
يقدم جميع المنسقين إشارة استعلام مُحسّنة:
// Check if any recent operation hit a rate limit
if (coordinator.HasSignal("rate-limit"))
{
await Task.Delay(1000);
}
// Count slow responses in the window
var slowCount = coordinator.CountSignals("slow-response");
if (slowCount > 10)
{
await ThrottleAsync();
}
// Get signals by pattern
var httpErrors = coordinator.GetSignalsByPattern("http.error.*");
// Get signals since a time
var recentSignals = coordinator.GetSignalsSince(DateTimeOffset.UtcNow.AddMinutes(-1));
// Get signals for a specific key
var userSignals = coordinator.GetSignalsByKey("user-123");
من طراز من طراز أولاً- اجتماع الدول الأطراف في:
يمكن للمنسقين أن يتفاعلوا تلقائياً مع الإشارات:
var coordinator = new EphemeralWorkCoordinator<Request>(
body,
new EphemeralOptions
{
// Cancel new work if these signals are present
CancelOnSignals = new HashSet<string> { "system-overload", "circuit-open" },
// Defer new work while these signals are present
DeferOnSignals = new HashSet<string> { "rate-limit" },
MaxDeferAttempts = 10,
DeferCheckInterval = TimeSpan.FromMilliseconds(100)
});
عندما تكون الإشارة في CancelOnSignals يتم الكشف عن عناصر جديدة يتم تجاوزها (يحسب على أنه فشل).
عندما تكون الإشارة في DeferOnSignals يتم اكتشافها، بنود جديدة تنتظر حتى توضّح الإشارة.
من طراز من طراز خدمات:
public class AdaptiveTranslationService : IAsyncDisposable
{
private readonly EphemeralWorkCoordinator<TranslationRequest> _coordinator;
private readonly ITranslationApi _translationApi;
public AdaptiveTranslationService(ITranslationApi translationApi)
{
_translationApi = translationApi;
_coordinator = new EphemeralWorkCoordinator<TranslationRequest>(
ProcessTranslationAsync,
new EphemeralOptions
{
MaxConcurrency = 8,
MaxTrackedOperations = 100,
// New work is deferred while any "rate-limit" or "rate-limit:*" signal is present
DeferOnSignals = new HashSet<string> { "rate-limit", "rate-limit:*" },
MaxDeferAttempts = 10,
DeferCheckInterval = TimeSpan.FromMilliseconds(100)
});
}
public async Task TranslateAsync(TranslationRequest request)
{
// Optional: extra politeness based on most recent retry-after
var rateLimitSignals = _coordinator.GetSignalsByPattern("rate-limit:*");
if (rateLimitSignals.Count > 0)
{
var latest = rateLimitSignals
.OrderByDescending(s => s.Timestamp)
.First()
.Signal; // "rate-limit:5000ms"
if (TryParseRetryAfter(latest, out var delay))
{
await Task.Delay(delay);
}
}
await _coordinator.EnqueueAsync(request);
}
public static bool TryParseRetryAfter(string signal, out TimeSpan delay)
{
delay = default;
var parts = signal.Split(':', 2);
if (parts.Length != 2) return false;
var payload = parts[1].Trim();
if (!payload.EndsWith("ms", StringComparison.OrdinalIgnoreCase)) return false;
var numPart = payload[..^2];
if (!int.TryParse(numPart, out var ms) || ms < 0) return false;
delay = TimeSpan.FromMilliseconds(ms);
return true;
}
}
كل حالة من هذه الخدمة تتراجع تلقائياً عند ضرب حدود المعدل. لا توجد حالة مشتركة. لا توجد رسالة تمر. فقط قراءة النافذة الفيزيائية.
يمكن لمنسقين متعددين متعددين أن يحسوا بعضهم البعض من خلال تقاسم SignalSink:
public class OrderProcessingSystem
{
private readonly SignalSink _sharedSignals = new(maxCapacity: 1000);
private readonly EphemeralWorkCoordinator<Order> _orderProcessor;
private readonly EphemeralWorkCoordinator<PaymentRequest> _paymentProcessor;
public OrderProcessingSystem()
{
var options = new EphemeralOptions { Signals = _sharedSignals };
_orderProcessor = new EphemeralWorkCoordinator<Order>(
ProcessOrderAsync, options);
_paymentProcessor = new EphemeralWorkCoordinator<PaymentRequest>(
ProcessPaymentAsync, options);
}
public async Task ProcessOrderAsync(Order order)
{
// Check shared signals for payment gateway issues
if (_sharedSignals.Detect("gateway-error"))
{
await _retryQueue.EnqueueAsync(order);
return;
}
await _orderProcessor.EnqueueAsync(order);
}
}
[HttpGet("/health/detailed")]
public IActionResult GetDetailedHealth()
{
return Ok(new
{
translation = new
{
pending = _translationCoordinator.PendingCount,
active = _translationCoordinator.ActiveCount,
recentRateLimits = _translationCoordinator.CountSignals("rate-limit"),
recentTimeouts = _translationCoordinator.CountSignals("timeout"),
recentSuccess = _translationCoordinator.CountSignals("success"),
hasErrors = _translationCoordinator.HasSignalMatching("error.*")
},
payment = new
{
pending = _paymentCoordinator.PendingCount,
gatewayErrors = _paymentCoordinator.CountSignals("gateway-error"),
declines = _paymentCoordinator.CountSignals("declined"),
approvals = _paymentCoordinator.CountSignals("approved")
}
});
}
لا حاجة إلى مكتبة قياسات فقط استعلام نافذة.
من طراز من طراز عُدّل الكسور.cs:
public class SignalBasedCircuitBreaker
{
private readonly string _failureSignal;
private readonly int _threshold;
private readonly TimeSpan _windowSize;
public SignalBasedCircuitBreaker(
string failureSignal = "failure",
int threshold = 5,
TimeSpan? windowSize = null)
{
_failureSignal = failureSignal;
_threshold = threshold;
_windowSize = windowSize ?? TimeSpan.FromSeconds(30);
}
public bool IsOpen<T>(EphemeralWorkCoordinator<T> coordinator)
{
var recentFailures = coordinator.GetSignalsSince(
DateTimeOffset.UtcNow - _windowSize);
return recentFailures.Count(s => s.Signal == _failureSignal) >= _threshold;
}
public int GetFailureCount<T>(EphemeralWorkCoordinator<T> coordinator)
{
var recentFailures = coordinator.GetSignalsSince(
DateTimeOffset.UtcNow - _windowSize);
return recentFailures.Count(s => s.Signal == _failureSignal);
}
}
// Usage
var circuitBreaker = new SignalBasedCircuitBreaker("api-error", threshold: 3);
if (circuitBreaker.IsOpen(_coordinator))
{
throw new CircuitOpenException("Too many recent API errors");
}
await _coordinator.EnqueueAsync(request);
كاسر الدائرة ليس له حالة خاصة به - انه فقط يقرأ النافذة الإيبميرية.
من طراز من طراز إشارات الوصلات:
عندما يمكن أن تسبب الإشارات إشارات أخرى، فإنك تخاطر بالدورات اللانهائية. SignalConstraints يمنع ما يلي:
var options = new EphemeralOptions
{
SignalConstraints = new SignalConstraints
{
// Max propagation depth before blocking
MaxDepth = 10,
// Prevent A → B → A cycles
BlockCycles = true,
// Signals that end propagation chains
TerminalSignals = new HashSet<string> { "completed", "failed", "resolved" },
// Signals that emit but don't propagate
LeafSignals = new HashSet<string> { "logged", "metric" },
// Callback when a signal is blocked
OnBlocked = (signal, reason) =>
{
_logger.LogWarning("Signal {Signal} blocked: {Reason}",
signal.Signal, reason);
}
}
};
العلاقة السببية مع: EmitCaused:
public void HandleSignal(SignalEvent evt, ISignalEmitter emitter)
{
if (evt.Is("order-placed"))
{
// This signal carries the propagation chain
// Will be blocked if it would create a cycle
emitter.EmitCaused("inventory-reserved", evt.Propagation);
}
}
وتسير سلسلة الانتشار على المسار: order-placed → inventory-reserved → ...
في حالة inventory-reserved محاولـة الانبعاث order-placedسيتم إيقافه (الدوره مكتشفه)
من طراز من طراز إشارات الوصلات:
بالنسبة للإشارات التي تحتاج إلى أن تكون مرئية عبر المنسقين:
public sealed class SignalSink
{
private readonly ConcurrentQueue<SignalEvent> _window;
private readonly int _maxCapacity;
private readonly TimeSpan _maxAge;
public SignalSink(int maxCapacity = 1000, TimeSpan? maxAge = null);
// Raise signals
public void Raise(SignalEvent signal);
public void Raise(string signal, string? key = null);
// Sense signals
public IReadOnlyList<SignalEvent> Sense();
public IReadOnlyList<SignalEvent> Sense(Func<SignalEvent, bool> predicate);
public bool Detect(string signalName);
public bool Detect(Func<SignalEvent, bool> predicate);
public int Count { get; }
}
// Create a shared sink
var sink = new SignalSink(maxCapacity: 1000, maxAge: TimeSpan.FromMinutes(2));
// Configure coordinators to use it
var options = new EphemeralOptions { Signals = sink };
// Or raise signals directly
sink.Raise("system-maintenance");
// Sense from anywhere
if (sink.Detect("system-maintenance"))
{
await DeferWorkAsync();
}
من طراز من طراز سلسلة سِطر سِطر P P P P Matn Mattchr.cs:
مُطابقة سطر لـ إشارة المرشِح:
// Exact match
coordinator.HasSignal("rate-limit");
// Wildcard patterns
coordinator.HasSignalMatching("http.*"); // http.timeout, http.error
coordinator.HasSignalMatching("error.*.critical"); // error.payment.critical
coordinator.HasSignalMatching("user-???-failed"); // user-123-failed
// Comma-separated patterns in CancelOnSignals/DeferOnSignals
new EphemeralOptions
{
CancelOnSignals = new HashSet<string>
{
"system-overload, circuit-open", // Either pattern
"error.*" // Any error signal
}
}
(أ) الاحتفاظ بإشارات بسيطة ومتسقة:
// Good - simple, categorical
op.Signal("success");
op.Signal("failure");
op.Signal("rate-limit");
op.Signal("timeout");
op.Signal("cache-hit");
// Good - structured for parsing
op.Signal("rate-limit:5000ms");
op.Signal("retry:attempt-3");
op.Signal("slow:2500ms");
op.Signal("http.error:429");
// Good - hierarchical for pattern matching
op.Signal("payment.declined");
op.Signal("payment.approved");
op.Signal("payment.gateway-error");
// Avoid - entity identification belongs in Key, not signals
op.Signal("user-123-rate-limited"); // Bad
// Instead
op.Key = "user-123";
op.Signal("rate-limit");
موصلات إشارة موصلاتOnSignal) تشغيل على خيط العملية - إبقهم بسرعة. للقيام بعمل I/O-I-O-ford، إشارات مروح إلى مسار async مع SignalDispatcher (المطابقة الأبطرية، الترتيب التحديدي) أو AsyncSignalProcessor.
await using var dispatcher = new SignalDispatcher(new EphemeralOptions
{
MaxConcurrency = Environment.ProcessorCount,
MaxConcurrencyPerKey = 1 // sequential per signal name by default
});
dispatcher.Register("error.*", evt => _alerts.SendAsync(evt.Signal));
dispatcher.Register("progress:*", evt => _metrics.Record(evt.Signal));
// In coordinator options, keep OnSignal chor options, keep OnSignal cheap and enqueue
var coordinator = new EphemeralWorkCoordinator<Request>(
body,
new EphemeralOptions
{
OnSignal = dispatcher.Dispatch
});
الدعم *, ?وقوائم فاصلة ("error.*,timeout") - جميع معالجات المطابقة التي تُدار حسب ترتيب التسجيل على منسق أساسي ذي خلفية أساسية؛ ولا تزال الانبعاثات متزامنة.
بالنسبة للتجهيز القائم بذاته:
await using var processor = new AsyncSignalProcessor(
async (signal, ct) =>
{
await _externalService.LogAsync(signal, ct);
},
maxConcurrency: 4,
maxQueueSize: 1000);
// Enqueue signals (returns immediately)
processor.Enqueue(new SignalEvent(
"rate-limit",
operationId,
key,
DateTimeOffset.UtcNow));
(أ) مثال كامل يجمع بين معالجة الاشارات المتماثلة مع التكامل عن بعد:
المصدر: المصدر: SISSSINAShandler.cs
public class TelemetrySignalHandler : IAsyncDisposable
{
private readonly AsyncSignalProcessor _processor;
private readonly ITelemetryClient _telemetry;
public TelemetrySignalHandler(ITelemetryClient telemetry)
{
_telemetry = telemetry;
_processor = new AsyncSignalProcessor(
HandleSignalAsync,
maxConcurrency: 8,
maxQueueSize: 5000);
}
// Synchronous entry point - returns immediately
public bool OnSignal(SignalEvent signal) => _processor.Enqueue(signal);
private async Task HandleSignalAsync(SignalEvent signal, CancellationToken ct)
{
var properties = new Dictionary<string, string>
{
["signal"] = signal.Signal,
["operationId"] = signal.OperationId.ToString(),
["key"] = signal.Key ?? "none"
};
await _telemetry.TrackEventAsync("EphemeralSignal", properties, ct);
// Categorized tracking based on signal prefix
if (signal.StartsWith("error"))
await _telemetry.TrackExceptionAsync(signal.Signal, properties, ct);
else if (signal.StartsWith("perf"))
await _telemetry.TrackMetricAsync(signal.Signal, 1, ct);
}
// Expose stats for monitoring
public int QueuedCount => _processor.QueuedCount;
public long ProcessedCount => _processor.ProcessedCount;
public long DroppedCount => _processor.DroppedCount;
public async ValueTask DisposeAsync() => await _processor.DisposeAsync();
}
أبلغ منسقك بما يلي:
await using var telemetryHandler = new TelemetrySignalHandler(telemetryClient);
await using var coordinator = new EphemeralWorkCoordinator<Request>(
ProcessAsync,
new EphemeralOptions
{
OnSignal = signal => telemetryHandler.OnSignal(signal)
});
المعالج:
OnSignal حالات العودة على أساس مباشرالإشارات قوية لأنها عن السنة عن السنة:
• الممتلكات/المنافع |----------|---------| | مُنْفِر مُنْفْفِر لا يمكن ان تنمو اشارات شيخوخة غير محدودة | التطهير الذاتي □ لا حاجة إلى شفر تنظيف | **** لا يعرف المستمعون شيئاً عن المستمعين | قابل للمسبور □ أي شفرة يمكن أن تحس الحالة الراهنة | القطاع الخاص □ لا توجد بيانات عن المستخدم - أسماء إشارات فقط □ | طراز AFS □ سين (1) الكشف مع قصر الدائرة
النافذة الإيفيمرية موجودة بالفعل من أجل التصحيح. الإشارات تعطيها فقط معنى دلالي.
تحول الاشارات الى الشبكة المشعـكيـةيمكن لكل منسق أن يقوم بما يلي:
CancelOnSignals وقد عقد مؤتمراً بشأن DeferOnSignalsلا رسالة سمسار، لا ولاية مشتركة، لا نظام تنسيق، فقط عمليات مع بيانات فوقية تتحلل طبيعياً
الذرات لا تتحدث مع بعضها البعض بشكل مباشر - فقط تترك آثاراً في النافذة الإقليمية التي يمكن للآخرين ملاحظتها. إنها ستيمبرجيا لأنظمة Async.
إشارة حريق، إنسَ، إنسَ، إنسَ، إنسَ، إنسَ، إنسَ، إنسَ، إنسَ، إنسَ، إنسَ، إنسَ، إنسَ.
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.