Back to "大部分时间短于. 完成; 在 LRU 缓存中, 一个奇怪的同时存在的系统模式系统 。"

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

mostlylucid-ephemeral

大部分时间短于. 完成; 在 LRU 缓存中, 一个奇怪的同时存在的系统模式系统 。

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(十级)赋予您所有原始功能。

属性和 DI

或 OR , 如果您想要使用简单路径的完整属性, 以 Async 为基础 [EphemeralJob] 以及服务服务.Add协调员的风格登记使用 大多为lucid. eeperal. atritutes 包包.

这很可能是我的博客继续讨论的主题,

最优优雅的, 短暂的, 完成

N Nuget 元数

在一个 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] 帮助您分组工作 记录、公平时间安排和诊断的有意义的钥匙。
  • Pinning 转发和重试: 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 将生成的信号反射回到标准日志 让你的监控堆看到两个 桥的两侧


核心核心协调员

包件 : 大部分是短暂的

短期工作协调员<泰 T T>

长时期工作队列,有接合货币和可观察到的窗口。

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();

时 候 眼 眼 工 作 协调 员<T. 传统知识,T>

每一关键顺序处理 -- -- 按顺序处理具有相同钥匙的物品。

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);

优先工作协调员<泰 T T>

每道多优先车道,每个车道可配置同价货币。

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);

OperationEchoRetentionOperationEchoCapacity 使你们能平衡你们所保持的呼声和持续的时间。 这样您就可以重弹“最后的单词” 足够长的时间来进行表面诊断。

机头起火时管理器自动解扣,但您可以拨打 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();

信号提醒Atom

**包件 : ** 大多为瞬间、原子、信标软件

根据环境信号暂停或取消摄入。

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();

Batching Atom 编织器

**包件 : ** 大部分是短暂的, 原子的, 断开

按大小或时间间隔将项目收集成批次 。

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

重试阿tom

包件 : 大多为瞬间、原子、再试验

指数后退重试包装器 。

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");

使用相同 DataStorageConfigMostlylucid.Ephemeral.Atoms.Data.SqliteMostlylucid.Ephemeral.Atoms.Data.Postgres 由SQLite/Postgres推动的耐久、信号驱动的持久性执行。 saved.data.{dbname} 启动下游工作的信号 load.data.{dbname} 触发水合物缓存。


MoleculeRunner & Atoom 调试器

**包件 : ** 大多为瞬间、原子、分子

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"))所以系统的其他部分 收集到 指挥棒


缓存缓存Atom

**包件 : ** 大部分是短暂的,原子,缓缓的缓冲

滑动过期缓存 - 访问结果 重置 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}");

时序LruCache

包件 : 核心(核心)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

控控 Fanout

**包件 : ** 受管制的硫丹

全球+受控平行式的按键带。

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();

反活动Fan outputPipeline

**包件 : ** 多数情况下是短暂的。 模式。 活性硫丹

带有自动后压的两阶段输油管

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)

遥测信号Handler

**包件 : ** 多数是短暂的 外观的 遥测仪

开放式遥测/应用透视集成。

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();

长窗口Demo

**包件 : ** 多数是短暂的 长窗口式的

显示审计线索的大型窗口配置 。

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 登记;他们只是把工作权交给戴头罩的神职人员。


目标框架

  • NET 6.0, 7.0, 8.0, 9.0, 10.0

许可证许可证许可证许可证

无许可证(公共领域)

logo

© 2026 Scott Galloway — Unlicense — All content and source code on this site is free to use, copy, modify, and sell.