This is a viewer only at the moment see the article on how this works.
To update the preview hit Ctrl-Alt-R (or ⌘-Alt-R on Mac) or Enter to refresh. The Save icon lets you save the markdown file to disk
This is a preview from the server running through my markdig pipeline
Sunday, 14 December 2025
上周我一直沉迷于此。 看看之前的部分和导致它的原因。 “ 如果一个路运联盟是一个执行环境。 ” 。 现在它是一个由30个Nuget包组成的套件, 覆盖了大多数主要的同时执行模式( 在 TINY 5- 10 行套套件中 ) 。 使用 SIMPLE 语法获得惊人的适应能力 !
在此查找来源 : https://github.com/scottgal/ mostlylucid.atoms/blob/main/ mostlylucid.ephemeral/src/ mostlylucid.ephemeral. 完成
读取 [信号上前一部份 ]用于对它的用途进行某些洞察。 此处展示的是“ 读取” 。 md from the mostlucid. epheperal. complete paceakage ” , 它包含核心大多为lucid. epheneral 包件和所有图案, “ 原子” (协调员等) 包件, 在一个方便的 DLL 中 。
OR 使用核心 大部分是短暂的 a TINY(十级)赋予您所有原始功能。
或 OR , 如果您想要使用简单路径的完整属性, 以 Async 为基础 [EphemeralJob] 以及服务服务.Add
这很可能是我的博客继续讨论的主题,
在一个 DLL 中的短暂时间 - 与基于信号的协调连接的同步执行。
dotnet add package mostlylucid.ephemeral.complete
此软件包将所有核心、 原子和模式代码汇编成一个组件。 对于单个软件包, 请查看每个软件包中的链接 下一节。
using Mostlylucid.Ephemeral;
// Long-lived work coordinator
await using var coordinator = new EphemeralWorkCoordinator<WorkItem>(
async (item, ct) => await ProcessAsync(item, ct),
new EphemeralOptions { MaxConcurrency = 8 });
await coordinator.EnqueueAsync(new WorkItem("data"));
// One-shot parallel processing
await items.EphemeralForEachAsync(
async (item, ct) => await ProcessAsync(item, ct),
new EphemeralOptions { MaxConcurrency = 8 });
var builder = WebApplication.CreateBuilder(args);
builder.Services.AddCoordinator<WorkItem>(
async (item, ct) => await ProcessAsync(item, ct),
new EphemeralOptions { MaxConcurrency = 8, MaxTrackedOperations = 128 });
builder.Services.AddEphemeralSignalJobRunner<LogWatcherJobs>();
var app = builder.Build();
app.MapPost("/", async ([FromServices] IEphemeralCoordinatorFactory<WorkItem> factory, WorkItem item) =>
{
var coordinator = factory.CreateCoordinator();
await coordinator.EnqueueAsync(item);
return Results.Accepted();
});
await app.RunAsync();
熟悉的 services.AddCoordinator<T>() 辅助者和辅助者, AddEphemeralSignalJobRunner<T>() 保持服务注册简洁,让 DI 拥有水槽/ 运行者, 并且将新的责任/ 缓存/ 博客故事单点击一击。
mostlylucid.ephemeral.complete 捆包 mostlylucid.ephemeral.attributes,所以属性管道是核心的一部分
将运行者作为头等信号消费者对待:装饰的方法与相同的缓存、记录和
编故事,每个属性都可以声明 Priority职位级别 MaxConcurrency, Lane, Key 源,信号信号
排出物, 插针/ 过期覆盖物, 和重试 。
按键属性 knobs :
Priority, MaxConcurrency, 和 Lane 保持确定性的工作秩序,同时走热路
保持单独。OperationKey, KeyFromSignal, KeyFromPayload, 和 [KeySource] 帮助您分组工作
记录、公平时间安排和诊断的有意义的钥匙。Pin, ExpireAfterMs, AwaitSignals, MaxRetries, 和 RetryDelayMs 扩展处理器
他们的能见度, 大门执行,直到依赖到达, 和愈合 重力,同时发出失败信号。EmitOnStart, EmitOnComplete, 和 EmitOnFailure 信号下游阶段,日志
观察者或其他没有手动布线的协调员。var sink = new SignalSink();
await using var runner = new EphemeralSignalJobRunner(sink, new[] { new LogWatcherJobs(sink) });
var loggerFactory = LoggerFactory.Create(builder =>
{
builder.AddConsole();
builder.AddProvider(new SignalLoggerProvider(new TypedSignalSink<SignalLogPayload>(sink)));
});
var logger = loggerFactory.CreateLogger("orders");
logger.LogError(new EventId(1001, "DbFailure"), "Order store failed");
// Later tasks or other services can also raise watcher-friendly signals directly:
sink.Raise("log.error.orders.dbfailure", key: "orders");
public sealed class LogWatcherJobs
{
private readonly SignalSink _sink;
public LogWatcherJobs(SignalSink sink) => _sink = sink;
[EphemeralJob("log.error.*", Priority = 1, MaxConcurrency = 2, Lane = "hot:4", EmitOnComplete = new[] { "incident.created" })]
public Task EscalateAsync(SignalEvent signal)
{
Console.WriteLine($"escalating {signal.Signal} for {signal.Key}");
_sink.Raise("incident.created", key: signal.Key);
return Task.CompletedTask;
}
[EphemeralJob("incident.created", EmitOnStart = new[] { "incident.monitor.start" })]
public Task NotifyAsync(SignalEvent signal)
{
Console.WriteLine($"notified incident for {signal.Key}");
return Task.CompletedTask;
}
}
这个跑者现在在启动时坐着 随时作出反应 log.error.* 或任何发出的信号击中汇。
处理器还可以从信号/有效载荷中读取密钥,在下游码头前做针线工作,发出完成/故障信号,以及
用于 DI 首端设置 services.AddEphemeralSignalJobRunner<T>() (或范围(或范围)
因此,向导和汇由集装箱管理。
[EphemeralJobs( 信号预言 = “ 阶段 ” , 默认Lane = “ 管道 ”) 公共密封班级 { [EphemeralJob ("inest", EmitOnCompllete = 新建)[{ "阶段. inest. done"} } 公共任务 IngestAsync (SignalEvent evt) Console. Out. WriteLineAsync (evt.Signal) ;
[EphemeralJob("finalize")]
public Task FinalizeAsync(SignalEvent evt) => Console.Out.WriteLineAsync("final stage");
}
var stageSink = 新的信号Sink (); Var stageRunner = 新的 EphemeralSignalJobRunner (Sink, 新的Sink) 等待[{新的阶段Jobs()}}; 级Sink. RAise (“ 级. inest ” ) ;
皮重的工作可以依赖 ResponsibilitySignalManager.PinUntilQueried (默认 ACK 模式) responsibility.ack.*至
在下游读者获取有效载荷之前,保持其操作可见, OperationEchoMaker/
OperationEchoAtom 最后的信号流持续,因此审计员或分子仍然可以“品尝”最后的状态,即便在
原子会死
mostlylucid.ephemeral.complete 也包含 mostlylucid.ephemeral.atoms.scheduledtasks或 JSON 定义
通过 ScheduledTaskDefinition (cron, 信号, 可选) key, payload, description, timeZone, format,
runOnStartup等), ScheduledTasksAtom (a) 开展持久的工作 DurableTaskAtom. 每项安排的工作
在协调窗口内提升配置的信号, 从而继承针线、 记录和负责语义
您的分子或属性管道对发出的信号波作出反应。
每个 DurableTask 带有调度表 Name, Signal,可选 Key键,甚至输入 Payload, 和 Description,因此下游听众立即知道哪些工作运行,哪些元数据(文件名、 URLs等)要消耗。 DurableTaskAtom.WaitForIdleAsync() 当您只想要等待当前计划的工作爆破完成而不完成原子, 使调度器为下个时钟准备就绪 。 @ info: whatsthis
mostlylucid.ephemeral.logging 微软镜像 Microsoft.Extensions. 登录到信号,反之亦然。
SignalLoggerProvider 您的伐木厂的伐木厂,所以日志事件增加 log.* 和勾勾 SignalToLoggerAdapter 如果
您想要信号返回到标准日志管道。
var sink = new SignalSink();
var typedSink = new TypedSignalSink<SignalLogPayload>(sink);
using var loggerFactory = LoggerFactory.Create(builder =>
{
builder.AddConsole();
builder.AddProvider(new SignalLoggerProvider(typedSink));
});
using var watcher = new EphemeralSignalJobRunner(sink, new[] { new LogWatcherJobs(sink) });
var logger = loggerFactory.CreateLogger("orders");
logger.LogError(new EventId(1001, "DbFailure"), "Order store failed");
public sealed class LogWatcherJobs
{
private readonly SignalSink _sink;
public LogWatcherJobs(SignalSink sink) => _sink = sink;
[EphemeralJob("log.error.*")]
public Task EscalateAsync(SignalEvent signal)
{
_sink.Raise("incident.created", key: signal.Key);
return Task.CompletedTask;
}
[EphemeralJob("incident.created")]
public Task NotifyAsync(SignalEvent signal)
{
Console.WriteLine($"Incident for {signal.Key}");
return Task.CompletedTask;
}
}
使用使用 SignalToLoggerAdapter 将生成的信号反射回到标准日志 让你的监控堆看到两个
桥的两侧
包件 : 大部分是短暂的
长时期工作队列,有接合货币和可观察到的窗口。
await using var coordinator = new EphemeralWorkCoordinator<Request>(
async (req, ct) => await HandleAsync(req, ct),
new EphemeralOptions
{
MaxConcurrency = 8,
MaxTrackedOperations = 200,
MaxOperationLifetime = TimeSpan.FromMinutes(5)
});
await coordinator.EnqueueAsync(request);
// Observe state
var running = coordinator.GetRunning();
var failed = coordinator.GetFailed();
var pending = coordinator.PendingCount;
// Graceful shutdown
coordinator.Complete();
await coordinator.DrainAsync();
每一关键顺序处理 -- -- 按顺序处理具有相同钥匙的物品。
await using var coordinator = new EphemeralKeyedWorkCoordinator<Order, string>(
order => order.CustomerId, // Key selector
async (order, ct) => await ProcessOrder(order, ct),
new EphemeralOptions
{
MaxConcurrency = 16, // Total parallel
MaxConcurrencyPerKey = 1 // Sequential per customer
});
await coordinator.EnqueueAsync(order);
自动同步操作的捕获结果。
await using var coordinator = new EphemeralResultCoordinator<Request, Response>(
async (req, ct) => await FetchAsync(req, ct),
new EphemeralOptions { MaxConcurrency = 4 });
var id = await coordinator.EnqueueAsync(request);
var snapshot = await coordinator.WaitForResult(id);
if (snapshot.HasResult)
Console.WriteLine(snapshot.Result);
每道多优先车道,每个车道可配置同价货币。
var coordinator = new PriorityWorkCoordinator<WorkItem>(
async (item, ct) => await ProcessAsync(item, ct),
new PriorityWorkCoordinatorOptions<WorkItem>(
Lanes: new[] { new PriorityLane("high"), new PriorityLane("normal"), new PriorityLane("low") }
));
await coordinator.EnqueueAsync(item, "high");
new EphemeralOptions
{
// Concurrency
MaxConcurrency = 8, // Max parallel operations
MaxConcurrencyPerKey = 1, // For keyed coordinators
EnableDynamicConcurrency = false, // Allow runtime adjustment
// Memory
MaxTrackedOperations = 200, // Window size (LRU eviction)
MaxOperationLifetime = TimeSpan.FromMinutes(5),
// Fair scheduling (keyed only)
EnableFairScheduling = false, // Prevent hot key starvation
FairSchedulingThreshold = 10,
// Signals
Signals = sharedSink, // Shared signal sink
OnSignal = evt => { }, // Sync callback
OnSignalAsync = async (evt, ct) => { }, // Async callback
CancelOnSignals = new HashSet<string> { "circuit-open" },
DeferOnSignals = new HashSet<string> { "backpressure" },
DeferCheckInterval = TimeSpan.FromMilliseconds(100),
MaxDeferAttempts = 50,
// Signal handler limits
MaxConcurrentSignalHandlers = 4,
MaxQueuedSignals = 1000
}
行动发出可交叉观测的信号。
// Query signals
bool hasError = coordinator.HasSignal("error");
int count = coordinator.CountSignals("error");
var errors = coordinator.GetSignalsByPattern("error.*");
// Shared sink across coordinators
var sink = new SignalSink();
var c1 = new EphemeralWorkCoordinator<A>(body, new EphemeralOptions { Signals = sink });
var c2 = new EphemeralWorkCoordinator<B>(body, new EphemeralOptions { Signals = sink });
sink.Raise("system.busy"); // Both see it
需要为下游消费者保持可见的结果足够长的时间吗? ResponsibilitySignalManager 让你用钉子钉
运行操作,直到一个 ACk 信号到达(默认模式) responsibility.ack.* 用密钥=operationId提供
可选 description 以便行动能够描述其责任,并设定 maxPinDuration 优雅地
如果消费者从不露面的话,可以自行清除。
var manager = new ResponsibilitySignalManager(coordinator, sink, maxPinDuration: TimeSpan.FromMinutes(5));
if (manager.PinUntilQueried(operationId, "file.ready", ackKey: fileId, description: "Awaiting fetch"))
{
sink.Raise("file.ready", key: fileId);
}
// Consumer acknowledges the work
sink.Raise("file.ready.ack", key: fileId);
using Mostlylucid.Ephemeral.Patterns;
var notes = new LastWordsNoteAtom(async note => await noteRepository.SaveAsync(note));
coordinator.OperationFinalized += snapshot =>
{
var note = new LastWordsNote(
OperationId: snapshot.OperationId,
Key: snapshot.Key,
Signal: snapshot.Signals?.FirstOrDefault(),
Timestamp: DateTimeOffset.UtcNow);
_ = notes.EnqueueAsync(note);
};
LastWordsNote 保持小小( 操作id、 密钥、 信号、 时间戳) , 这样您就可以记录您所关心的最小状态 。
在收集操作前的关于操作 。
协调员还保持最后信号(通过 EnableOperationEcho)你可以
检查 GetEchoes() 当您需要重播短小的信号波时, 而不保留整个操作 。
var recentErrors = coordinator.GetEchoes(pattern: "error.*")
.Where(e => e.Timestamp > DateTimeOffset.UtcNow - TimeSpan.FromMinutes(1))
.ToList();
if (recentErrors.Any())
logger.LogWarning("Trimmed errors: {Count}", recentErrors.Count);
OperationEchoRetention 和 OperationEchoCapacity 使你们能平衡你们所保持的呼声和持续的时间。
这样您就可以重弹“最后的单词” 足够长的时间来进行表面诊断。
机头起火时管理器自动解扣,但您可以拨打 CompleteResponsibility(operationId) 结束
早期(例如,在重试时)。 OperationFinalized 当窗户剪剪的时候,
如果您想要发送最终信号、日志诊断或运行“最后单词”清理, 请订阅 。
**包件 : ** 大多为短期、原子、固定工作
固定工人人才库,配有统计数据,短时间工作协调员周围的最小API包装。
using Mostlylucid.Ephemeral.Atoms.FixedWork;
await using var atom = new FixedWorkAtom<WorkItem>(
async (item, ct) => await ProcessAsync(item, ct),
maxConcurrency: 4,
maxTracked: 200);
await atom.EnqueueAsync(item);
// Get stats
var (pending, active, completed, failed) = atom.Stats();
Console.WriteLine($"Completed: {completed}, Failed: {failed}");
// Get recent operations
var snapshot = atom.Snapshot();
// Graceful shutdown
await atom.DrainAsync();
**包件 : ** 多数为短暂的原子, 关键序列
以可选择的公平日程安排进行按顺序顺序处理。
using Mostlylucid.Ephemeral.Atoms.KeyedSequential;
await using var atom = new KeyedSequentialAtom<Order, string>(
keySelector: order => order.CustomerId,
body: async (order, ct) => await ProcessOrder(order, ct),
maxConcurrency: 16,
perKeyConcurrency: 1, // Sequential per key
enableFairScheduling: true); // Prevent hot key starvation
await atom.EnqueueAsync(order1); // Customer A
await atom.EnqueueAsync(order2); // Customer A - waits for order1
await atom.EnqueueAsync(order3); // Customer B - parallel with A
var (pending, active, completed, failed) = atom.Stats();
await atom.DrainAsync();
**包件 : ** 大多为瞬间、原子、信标软件
根据环境信号暂停或取消摄入。
using Mostlylucid.Ephemeral.Atoms.SignalAware;
var sink = new SignalSink();
await using var atom = new SignalAwareAtom<WorkItem>(
async (item, ct) => await ProcessAsync(item, ct),
cancelOn: new HashSet<string> { "shutdown", "circuit-open" },
deferOn: new HashSet<string> { "backpressure.*" },
deferInterval: TimeSpan.FromMilliseconds(100),
maxDeferAttempts: 50,
signals: sink,
maxConcurrency: 8);
// Enqueue work
await atom.EnqueueAsync(item);
// Raise ambient signals
atom.Raise("backpressure.downstream"); // New items defer
sink.Raise("shutdown"); // New items rejected (returns -1)
await atom.DrainAsync();
**包件 : ** 大部分是短暂的, 原子的, 断开
按大小或时间间隔将项目收集成批次 。
using Mostlylucid.Ephemeral.Atoms.Batching;
await using var atom = new BatchingAtom<LogEntry>(
onBatch: async (batch, ct) =>
{
Console.WriteLine($"Flushing {batch.Count} entries");
await FlushToDatabase(batch, ct);
},
maxBatchSize: 100,
flushInterval: TimeSpan.FromSeconds(5));
// Items are batched automatically
atom.Enqueue(new LogEntry("User logged in"));
atom.Enqueue(new LogEntry("Request received"));
// ... batch flushes when full OR after 5 seconds
包件 : 大多为瞬间、原子、再试验
指数后退重试包装器 。
using Mostlylucid.Ephemeral.Atoms.Retry;
await using var atom = new RetryAtom<ApiRequest>(
async (req, ct) => await CallExternalApi(req, ct),
maxAttempts: 3,
backoff: attempt => TimeSpan.FromMilliseconds(100 * Math.Pow(2, attempt)),
maxConcurrency: 4);
// Automatically retries on failure with exponential backoff
// Attempt 1: immediate
// Attempt 2: 200ms delay
// Attempt 3: 400ms delay
await atom.EnqueueAsync(new ApiRequest("https://api.example.com"));
await atom.DrainAsync();
包件 : 多数为瞬间、原子数据
存储原子共享配置( A)DataStorageConfig, IDataStorageAtom<TKey, TValue>加上驱动文件、 SQLite 和 PostgreSQL 适配器的信号协议。
using Mostlylucid.Ephemeral.Atoms.Data;
using Mostlylucid.Ephemeral.Atoms.Data.File;
var sink = new SignalSink();
var config = new DataStorageConfig
{
DatabaseName = "orders",
SignalPrefix = "save.data",
LoadSignalPrefix = "load.data",
DeleteSignalPrefix = "delete.data",
MaxConcurrency = 1
};
await using var storage = new FileDataStorageAtom<string, Order>(sink, config, "./orders");
storage.EnqueueSave("order-123", new Order { Id = "order-123", Total = 42.00m });
var loaded = await storage.LoadAsync("order-123");
使用相同 DataStorageConfig 与 Mostlylucid.Ephemeral.Atoms.Data.Sqlite 或 Mostlylucid.Ephemeral.Atoms.Data.Postgres 由SQLite/Postgres推动的耐久、信号驱动的持久性执行。 saved.data.{dbname} 启动下游工作的信号 load.data.{dbname} 触发水合物缓存。
**包件 : ** 大多为瞬间、原子、分子
由 MoleculeBlueprintBuilder 以便你们确定原子,
)的通知),该通知应当运行,当信号,例如: order.placed 到达。 MoleculeRunner 监听触发器
创建共享 MoleculeContext,并在订阅开始/完成事件时执行每个步骤。使用
AtomTrigger 当一个原子的信号 启动另一个协调器或分子时
var sink = new SignalSink();
var blueprint = new MoleculeBlueprintBuilder("order", "order.placed")
.AddAtom(async (ctx, ct) => await paymentCoordinator.EnqueueAsync(ctx.TriggerSignal.Key!, ct))
.AddAtom(async (ctx, ct) =>
{
ctx.Raise("order.payment.complete", ctx.TriggerSignal.Key);
await inventoryCoordinator.EnqueueAsync(ctx.TriggerSignal.Key!, ct);
})
.Build();
await using var runner = new MoleculeRunner(sink, new[] { blueprint }, serviceProvider);
using var trigger = new AtomTrigger(sink, "order.payment.complete", async (signal, ct) =>
{
await notificationCoordinator.EnqueueAsync(signal.Key!, ct);
});
sink.Raise("order.placed", key: "order-42");
分子步骤可产生额外的信号(ctx.Raise("order.shipping.start"))所以系统的其他部分 收集到
指挥棒
**包件 : ** 大部分是短暂的,原子,缓缓的缓冲
滑动过期缓存 - 访问结果 重置 TTL 。
using Mostlylucid.Ephemeral.Atoms.SlidingCache;
await using var cache = new SlidingCacheAtom<string, UserProfile>(
async (userId, ct) => await LoadUserProfileAsync(userId, ct),
slidingExpiration: TimeSpan.FromMinutes(5),
absoluteExpiration: TimeSpan.FromHours(1),
maxSize: 1000);
// First call: computes and caches
var profile = await cache.GetOrComputeAsync("user-123");
// Second call within 5 minutes: returns cached, resets TTL
var cached = await cache.GetOrComputeAsync("user-123");
// Try get without computation (still resets TTL on hit)
if (cache.TryGet("user-123", out var profile))
Console.WriteLine(profile.Name);
// Get stats
var stats = cache.GetStats();
Console.WriteLine($"Entries: {stats.TotalEntries}, Hot: {stats.HotEntries}");
包件 : 核心(核心)
mostlylucid.ephemeral· 自动优化缓存,每击一次就滑滑TTL,并延长TTL 热键。
using Mostlylucid.Ephemeral;
var cache = new EphemeralLruCache<string, Widget>(new EphemeralLruCacheOptions
{
DefaultTtl = TimeSpan.FromMinutes(5),
HotKeyExtension = TimeSpan.FromMinutes(30),
HotAccessThreshold = 3,
MaxSize = 10_000,
SampleRate = 5 // emit 1 in 5 signals
});
var widget = await cache.GetOrAddAsync("widget:42", async key =>
{
var data = await LoadWidgetAsync(key);
return data!;
});
// Stats and signals to see how the cache self-focuses on hot keys
var stats = cache.GetStats(); // hot/expired counts, size
var signals = cache.GetSignals("cache.*"); // cache.hot/evict/miss/hit
提示 :
MemoryCache可配置为滑动到期, 但不会发出热/ 冷信号或扩展 TTL 热键。EphemeralLruCache核心软件包中的自优化默认默认值(和在SqliteSingleWriter) 当您想要缓存聚焦在活动工作集时。
抓取一个操作在被剪剪之前发出的输入的“ 最后单词 ” 。 原子保持一个封闭的信号窗口
有效载荷(匹配) ActivationSignalPattern / CaptureSignalPattern)和时间 OperationFinalized 它产生火灾,产生火灾
OperationEchoEntry<TPayload> 记录可以通过 OperationEchoAtom<TPayload>.
var sink = new SignalSink();
var typedSink = new TypedSignalSink<EchoPayload>(sink);
var echoAtom = new OperationEchoAtom<EchoPayload>(async echo => await repository.AppendAsync(echo));
await using var coordinator = new EphemeralWorkCoordinator<JobItem>(ProcessAsync);
using var maker = coordinator.EnableOperationEchoing(
typedSink,
echoAtom,
new OperationEchoMakerOptions<EchoPayload>
{
ActivationSignalPattern = "echo.capture",
CaptureSignalPattern = "echo.*",
MaxTrackedOperations = 128
});
typedSink.Raise("echo.capture", new EchoPayload("order-1", "archived"), key: "order-1");
属性工作仅以他们认为关键的任何状态提高输入的信号, 创建者保留工作设置 在你坚持回声时被绑住
**包件 : ** 大多是短暂的, 模式, 解路器
使用信号历史窗口的 静态电路断路器
using Mostlylucid.Ephemeral.Patterns.CircuitBreaker;
var breaker = new SignalBasedCircuitBreaker(
failureSignal: "api.failure",
threshold: 5,
windowSize: TimeSpan.FromSeconds(30));
// Check before making calls
if (breaker.IsOpen(coordinator))
{
var retryAfter = breaker.GetTimeUntilClose(coordinator);
throw new CircuitOpenException("Too many failures", retryAfter);
}
// Pattern matching variant
if (breaker.IsOpenMatching(coordinator, "error.*"))
throw new CircuitOpenException("Error pattern detected");
// Get current failure count
int failures = breaker.GetFailureCount(coordinator);
将深度管理排成队列,对后压信号自动推迟。
using Mostlylucid.Ephemeral.Patterns.Backpressure;
var sink = new SignalSink();
await using var coordinator = SignalDrivenBackpressure.Create<WorkItem>(
async (item, ct) => await ProcessAsync(item, ct),
sink,
maxConcurrency: 4);
// Enqueue work
await coordinator.EnqueueAsync(item);
// When downstream is slow
sink.Raise("backpressure.downstream"); // New work auto-defers
// When recovered
sink.Retract("backpressure.downstream"); // Work resumes
**包件 : ** 受管制的硫丹
全球+受控平行式的按键带。
using Mostlylucid.Ephemeral.Patterns.ControlledFanOut;
await using var fanout = new ControlledFanOut<string, Request>(
keySelector: req => req.TenantId,
body: async (req, ct) => await ProcessAsync(req, ct),
maxGlobalConcurrency: 100, // Total parallel across all tenants
perKeyConcurrency: 5); // Max 5 parallel per tenant
// Items for same tenant processed with limit
await fanout.EnqueueAsync(requestA); // Tenant1
await fanout.EnqueueAsync(requestB); // Tenant1 - waits if 5 already running
await fanout.EnqueueAsync(requestC); // Tenant2 - parallel with Tenant1
await fanout.DrainAsync();
信号驱动率限制,自动回扣。
using Mostlylucid.Ephemeral.Patterns.AdaptiveRate;
await using var service = new AdaptiveRateService<ApiRequest>(
async (req, ct) => await CallApiAsync(req, ct),
maxConcurrency: 8);
// Process with automatic rate limit handling
await service.ProcessAsync(request);
// When API returns 429, emit signal with retry-after
// Signal: "rate-limit:500ms"
// Service auto-parses and delays
Console.WriteLine($"Pending: {service.PendingCount}, Active: {service.ActiveCount}");
**包件 : ** 多数为流动货币
以负载信号为依据的运行时间折合货币规模。
using Mostlylucid.Ephemeral.Patterns.DynamicConcurrency;
var sink = new SignalSink();
await using var demo = new DynamicConcurrencyDemo<WorkItem>(
async (item, ct) => await ProcessAsync(item, ct),
sink,
minConcurrency: 2,
maxConcurrency: 32,
scaleUpPattern: "load.high",
scaleDownPattern: "load.low");
await demo.EnqueueAsync(item);
// Concurrency adjusts automatically based on signals
sink.Raise("load.high"); // Concurrency doubles (up to max)
sink.Raise("load.low"); // Concurrency halves (down to min)
Console.WriteLine($"Current concurrency: {demo.CurrentMaxConcurrency}");
await demo.DrainAsync();
**包件 : ** 多数为低迷的短期性. 模式性. 关键优先放风
优先车道和按键定购的优先车道均予保留。
using Mostlylucid.Ephemeral.Patterns.KeyedPriorityFanOut;
await using var fanout = new KeyedPriorityFanOut<string, UserCommand>(
keySelector: cmd => cmd.UserId,
body: async (cmd, ct) => await HandleCommand(cmd, ct),
maxConcurrency: 32,
perKeyConcurrency: 1, // Sequential per user
maxPriorityDepth: 100);
// Normal lane
await fanout.EnqueueAsync(normalCommand);
// Priority lane - jumps the queue for that user
bool accepted = await fanout.EnqueuePriorityAsync(urgentCommand);
// Check lane depths
var counts = fanout.PendingCounts;
Console.WriteLine($"Priority: {counts.Priority}, Normal: {counts.Normal}");
await fanout.DrainAsync();
**包件 : ** 多数情况下是短暂的。 模式。 活性硫丹
带有自动后压的两阶段输油管
using Mostlylucid.Ephemeral.Patterns.ReactiveFanOut;
await using var pipeline = new ReactiveFanOutPipeline<WorkItem>(
stage2Work: async (item, ct) => await SlowProcessing(item, ct),
preStageWork: async (item, ct) => await FastPreprocessing(item, ct),
stage1MaxConcurrency: 8,
stage1MinConcurrency: 1,
stage2MaxConcurrency: 4,
backpressureThreshold: 32, // Throttle when stage2 has 32+ pending
reliefThreshold: 8); // Resume when stage2 drops below 8
await pipeline.EnqueueAsync(item);
// Stage1 auto-throttles when stage2 backs up
Console.WriteLine($"Stage1 concurrency: {pipeline.Stage1CurrentMaxConcurrency}");
Console.WriteLine($"Stage2 pending: {pipeline.Stage2Pending}");
await pipeline.DrainAsync();
**包件 : ** 主要为: 短期性. 型体. 异常现象检测仪
移动窗口异常现象检测
using Mostlylucid.Ephemeral.Patterns.AnomalyDetector;
var sink = new SignalSink();
var detector = new SignalAnomalyDetector(
sink,
pattern: "error.*",
threshold: 5,
window: TimeSpan.FromSeconds(10));
// Check for anomalies
if (detector.IsAnomalous())
{
Console.WriteLine("Anomaly detected! Too many errors.");
TriggerAlert();
}
// Get current match count
int errorCount = detector.GetMatchCount();
Console.WriteLine($"Errors in window: {errorCount}");
**包件 : ** 信号协调阅读
校对:Portnoy
using Mostlylucid.Ephemeral.Patterns.SignalCoordinatedReads;
// Run demo: readers pause when update signal is present
var result = await SignalCoordinatedReads.RunAsync(
readCount: 10,
updateCount: 1);
Console.WriteLine($"Reads: {result.ReadsCompleted}, Updates: {result.UpdatesCompleted}");
Console.WriteLine($"Signals: {string.Join(", ", result.Signals)}");
// Manual implementation:
var sink = new SignalSink();
await using var readers = new EphemeralWorkCoordinator<Query>(
body,
new EphemeralOptions
{
DeferOnSignals = new HashSet<string> { "update.in-progress" },
Signals = sink
});
// Readers auto-defer when update is running
sink.Raise("update.in-progress"); // Readers wait
sink.Raise("update.done"); // Readers resume
**包件 : ** 大多是短暂的 模式式的 信号式的 http
具有进度信号的 HTTP 客户端。
using Mostlylucid.Ephemeral.Patterns.SignalingHttp;
var httpClient = new HttpClient();
var request = new HttpRequestMessage(HttpMethod.Get, "https://example.com/large-file");
// Create an emitter from your coordinator
// (emitter is any ISignalEmitter - operations implement this)
byte[] data = await SignalingHttpClient.DownloadWithSignalsAsync(
httpClient,
request,
emitter);
// Signals emitted during download:
// - stage.starting
// - progress:0
// - stage.request
// - stage.headers
// - stage.reading
// - progress:25, progress:50, progress:75, progress:100
// - stage.completed
**包件 : ** 大多是短时间的, 外派的, 信号观测器
观察信号窗口的图案和触发回击。
using Mostlylucid.Ephemeral.Patterns.SignalLogWatcher;
var sink = new SignalSink();
await using var watcher = new SignalLogWatcher(
sink,
onMatch: evt =>
{
Console.WriteLine($"Error detected: {evt.Signal} at {evt.Timestamp}");
AlertOps(evt);
},
pattern: "error.*",
pollInterval: TimeSpan.FromMilliseconds(200));
// Watcher runs in background, calling onMatch for each new error signal
sink.Raise("error.database"); // -> onMatch called
sink.Raise("error.timeout"); // -> onMatch called
sink.Raise("info.started"); // -> ignored (doesn't match pattern)
**包件 : ** 多数是短暂的 外观的 遥测仪
开放式遥测/应用透视集成。
using Mostlylucid.Ephemeral.Patterns.Telemetry;
// Use in-memory for testing, or implement ITelemetryClient for real telemetry
var telemetry = new InMemoryTelemetryClient();
await using var handler = new TelemetrySignalHandler(telemetry);
// Wire up to coordinator
var options = new EphemeralOptions
{
OnSignal = signal => handler.OnSignal(signal)
};
// Signals are processed asynchronously
// - "error.*" signals -> TrackExceptionAsync
// - "perf.*" signals -> TrackMetricAsync
// - all signals -> TrackEventAsync
Console.WriteLine($"Queued: {handler.QueuedCount}");
Console.WriteLine($"Processed: {handler.ProcessedCount}");
Console.WriteLine($"Dropped: {handler.DroppedCount}");
// Check recorded events
var events = telemetry.GetEvents();
**包件 : ** 多数是短暂的 长窗口式的
显示审计线索的大型窗口配置 。
using Mostlylucid.Ephemeral.Patterns.LongWindowDemo;
// Configure coordinator with large tracking window
var options = new EphemeralOptions
{
MaxTrackedOperations = 10000,
MaxOperationLifetime = TimeSpan.FromHours(24)
};
**包件 : ** 大多是短暂的, 模式。 信号反应显示情况 。
显示信号发送模式和回击
using Mostlylucid.Ephemeral.Patterns.SignalReactionShowcase;
// See source for signal dispatch examples
// Demonstrates OnSignal, OnSignalAsync, CancelOnSignals, DeferOnSignals
SQLite 持久性信号窗口 - 生存进程重新启动 。
using Mostlylucid.Ephemeral.Patterns.PersistentWindow;
await using var window = new PersistentSignalWindow(
"Data Source=signals.db",
flushInterval: TimeSpan.FromSeconds(30));
// On startup: restore previous signals
await window.LoadFromDiskAsync(maxAge: TimeSpan.FromHours(24));
// Raise signals as normal
window.Raise("order.completed", key: "order-service");
window.Raise("payment.processed", key: "payment-service");
// Query signals
var recentOrders = window.Sense("order.*");
// Signals automatically flush every 30 seconds
// Also flushes on dispose
// Get stats
var stats = window.GetStats();
Console.WriteLine($"In memory: {stats.InMemoryCount}, Flushed: {stats.LastFlushedId}");
// Register in Startup/Program.cs
services.AddEphemeralWorkCoordinator<WorkItem>(
async (item, ct) => await ProcessAsync(item, ct),
new EphemeralOptions { MaxConcurrency = 8 });
// Named coordinators
services.AddEphemeralWorkCoordinator<WorkItem>("priority",
async (item, ct) => await ProcessPriorityAsync(item, ct));
// Inject and use
public class MyService(IEphemeralCoordinatorFactory<WorkItem> factory)
{
public async Task DoWork()
{
var coordinator = factory.CreateCoordinator();
await coordinator.EnqueueAsync(new WorkItem());
}
}
现代DI根可能更喜欢较短的帮手,例如 services.AddCoordinator<T>(...),
services.AddScopedCoordinator<T>(...),或 services.AddKeyedCoordinator<T, TKey>(...) 因为他们读起来像正常人
AddX 登记;他们只是把工作权交给戴头罩的神职人员。
无许可证(公共领域)
© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.