notify

    MoonBit 通知库:多平台 webhook(飞书/钉钉/企业微信/Slack/Discord)、通用 webhook、SMTP 邮件,含去重/限流/重试引擎

    notify
    webhook
    feishu
    dingtalk
    smtp
    Download zip
    Author
    Version
    0.3.0
    License
    MIT
    Last updated
    8 hours ago
    Downloads
    1

    Dependencies

    #chensuiyi/notify

    MoonBit 多平台通知库 · Multi-platform notification library for MoonBit

    #简介

    一个库覆盖通知的完整链路:消息模型、七种平台预设、HTTP/SMTP 投递,以及可选的去重/限流/重试引擎。

    设计约定:

    • 凭据不泄露——失败原因永不包含 url / secret / password,可放心打进日志
    • 时间由调用方注入——库不读时钟,完全可测
    • 不写日志不落盘——只投递并返回结构化结果

    #Intro

    One library for the whole notification path: a message model, seven platform presets, HTTP/SMTP delivery, and an optional dedupe / rate-limit / retry engine.

    Design contract:

    • Credential-free failures — failure reasons never contain url / secret / password, safe to log
    • Caller-injected time — the library never reads a clock, fully testable
    • No logging, no disk — delivers and returns structured results

    #平台 Platforms

    preset平台 Platform传输 Transport
    feishu飞书HTTPS(可选签名)
    dingtalk钉钉HTTPS(HMAC 签名)
    wecom企业微信HTTPS
    slackSlackHTTPS
    discordDiscordHTTPS
    webhook通用 webhookHTTPS(Bearer)
    email邮件SMTP(STARTTLS / smtps)

    #用法 Usage

    // 一次性投递:None = 成功,Some(reason) = 失败原因(不含凭据)
    let failure = deliver(target, msg, 5000, { host: "prod-1", product: "myapp" })

    // 可选引擎:去重 → 限流 → 队列 → 失败退避重试
    let engine = Engine::new({ host: "prod-1", product: "myapp" })
    engine.emit(
    now, "crash|api", "api", msg, targets, 5000,
    dedupe_window_s=3600, rate_per_minute=30, queue_limit=128,
    )
    let report = engine.pump(now) // 每 tick 调用一次

    #API

    接口说明
    deliver(target, msg, timeout_ms, ctx)投递一条消息到单个目标
    deliver_all(targets, msg, timeout_ms, ctx)逐个投递,返回 (name, failure) 列表
    Engine::new(ctx, max_attempts?)创建引擎
    Engine::emit(...)入队一条消息,返回 Queued / Deduped / RateLimited / QueueFull
    Engine::pump(now)投递队首消息,返回 remaining / outcome / owner
    Engine::dropped() / Engine::queued()丢弃与排队计数

    限流窗口按 owner 独立,计数发生在入队时——失败端点不会绕过限流。

    #依赖 Dependencies

    • moonbitlang/async@0.22.4(HTTP / SMTP / 超时)
    • moonbitlang/x@0.5.1(HMAC 签名)

    License: MIT

    NotifyError

    pub suberror NotifyError {
    InvalidPreset(String)
    } derive(
    Debug
    )

    impl Show for NotifyError

    NotifyError::output

    fn NotifyError::output(self : NotifyError, logger : &Logger) -> Unit

    NotifyError::to_string

    fn NotifyError::to_string(self : NotifyError) -> String

    SmtpFailed

    type SmtpFailed derive(
    Debug
    )

    SMTP-level failures. Messages never contain credentials: they carry the stage and the server's status line, which is what an operator needs.

    Context

    pub(all) struct Context {
    host : String
    product : String
    } derive(
    Debug
    )

    Who is sending: appears in message footers (product @ host) and in the SMTP EHLO greeting. Passed in by the caller so the library stays platform-agnostic.

    Context::default

    fn Context::default() -> Context

    The default context for callers that do not care about branding.

    Context::to_repr

    DedupeEntry

    type DedupeEntry

    Delivery

    pub(all) struct Delivery {
    url : String
    http_method : String
    content_type : String?
    headers : Array[(String, String)]
    body : String
    mail : String?
    } derive(
    Debug
    )

    A rendered request: what the transport hands to the network. Kept separate from the message so one message can be rendered per target.

    Delivery::to_repr

    EmitOutcome

    pub(all) enum EmitOutcome {
    Queued
    Deduped(count~ : Int)
    RateLimited
    QueueFull
    } derive(Eq,
    Debug
    )

    What happened to one emit call.

    EmitOutcome::equal

    fn EmitOutcome::equal(EmitOutcome, EmitOutcome) -> Bool

    EmitOutcome::not_equal

    fn EmitOutcome::not_equal(x : EmitOutcome, y : EmitOutcome) -> Bool

    Engine

    pub struct Engine {
    ctx : Context
    max_attempts : Int
    queue : Array[Queued]
    dropped : Int
    rate_window : Map[String, (UInt64, Int)]
    dedupe : Map[String, DedupeEntry]
    next_attempt_ms : UInt64
    attempts : Int
    }

    Engine::dropped

    fn Engine::dropped(self : Engine) -> Int

    Messages dropped since the engine was created: queue overflow plus rate limiting.

    Engine::emit

    fn Engine::emit(self : Engine, now : UInt64, key : String, owner : String, msg : Message, targets : Array[Target], timeout_ms : Int, dedupe_window_s? : Int, rate_per_minute? : Int, queue_limit? : Int) -> EmitOutcome

    Queue one message. Deduplication, rate limiting and queue accounting are all driven by the parameters — the engine keeps no policy of its own.

    • key deduplicates identical conditions inside dedupe_window_s.
    • owner + queue_limit keep a burst from one group from starving others: once queue_limit messages of owner are queued, further ones drop.
    • rate_per_minute caps emissions per rolling 60s window (0 = unlimited).

    Engine::new

    fn Engine::new(ctx : Context, max_attempts? : Int) -> Engine

    Create an engine. ctx brands the outgoing messages (footer and EHLO name); max_attempts is how many delivery passes one message gets before it is dropped.

    Engine::pump

    async fn Engine::pump(self : Engine, now : UInt64) -> PumpReport

    One pass of the engine: deliver at most one queued message.

    Engine::queued

    fn Engine::queued(self : Engine) -> Int

    Number of messages still queued.

    Message

    pub(all) struct Message {
    event : String
    level : String
    title : String
    body : String
    fields : Array[(String, String)]
    ts_ms : UInt64
    key : String
    } derive(
    Debug
    )

    One notification. key is what deduplication keys on; event is the logical event name carried through to generic webhook payloads.

    Message::new

    fn Message::new(event : String, level : String, title : String, body : String, key : String, ts_ms : UInt64, fields? : Array[(String, String)]) -> Message

    Message::to_repr

    PumpOutcome

    pub(all) enum PumpOutcome {
    Idle
    Delivered
    Partial
    Failed(attempt~ : Int)
    Exhausted
    } derive(Eq,
    Debug
    )

    What one pump pass did with the head of the queue.

    PumpOutcome::equal

    fn PumpOutcome::equal(PumpOutcome, PumpOutcome) -> Bool

    PumpOutcome::not_equal

    fn PumpOutcome::not_equal(x : PumpOutcome, y : PumpOutcome) -> Bool

    PumpReport

    pub(all) struct PumpReport {
    remaining : Int
    outcome : PumpOutcome
    owner : String?
    } derive(
    Debug
    )

    Queued

    type Queued

    Target

    pub(all) struct Target {
    name : String
    preset : String
    url : String
    secret : String?
    format : String?
    from : String?
    to : Array[String]?
    username : String?
    password : String?
    } derive(
    Debug
    )

    One delivery destination. preset selects the payload shape (a closed set: feishu / dingtalk / wecom / slack / discord / webhook / email); the rest are endpoint and credential fields.

    Target::to_repr

    deliver

    async fn deliver(target : Target, msg : Message, timeout_ms : Int, ctx : Context) -> String?

    Deliver one message to one target.

    deliver_all

    async fn deliver_all(targets : Array[Target], msg : Message, timeout_ms : Int, ctx : Context) -> Array[(String, String?)]

    Deliver one message to many targets, one request each, sequentially. Returns (target name, failure) pairs in input order.

    smtp_deliver

    async fn smtp_deliver(target : Target, delivery : Delivery, mail : String, ehlo_name : String, timeout_ms : Int) -> String?

    Deliver one RFC5322 message over SMTP with STARTTLS (or implicit TLS for smtps://). None on acceptance; Some carries a credential-free reason.