Skip to content
AI & Agents
Skill

/instrumenting-streaming-events

Instrumenting streaming events: wires lib-streaming event emission end-to-end into a Lerian Go service via a 13-gate cycle (catalog, Builder bootstrap, Emit sites, outbox, HTTP manifest, NoopEmitter fallback, integration and chaos tests), dispatching ring:backend-go under TDD.

From plugin
ring
20577 skills42 agents1 command
Install
$ npx -y skills add LerianStudio/ring --skill instrumenting-streaming-events --agent claude-code

How it fires

How this skill gets triggered: by you, by Claude, or both.

  • Fires itselfAuto-invocation. Claude auto-loads it when your prompt matches the work.Auto-invocation is when the right skill fires by itself at the right moment, driven by a FLOW.md router and a hook, instead of you invoking it by name. It is the difference between a skill being installed and a skill actually getting used.Read the full definition →
  • You can call itInvoke it directly when you want it.
  • Slash command/instrumenting-streaming-events

Context preview

The summary Claude sees to decide when to auto-load this skill.

Instrumenting streaming events: wires lib-streaming event emission end-to-end into a Lerian Go service via a 13-gate cycle (catalog, Builder bootstrap, Emit sites, outbox, HTTP manifest, NoopEmitter fallback, integration and chaos tests), dispatching ring:backend-go under TDD.

SKILL.md

instrumenting-streaming-events.SKILL.md
name: ring:instrumenting-streaming-events
description: "Instrumenting streaming events: wires lib-streaming event emission end-to-end into a Lerian Go service via a 13-gate cycle (catalog, Builder bootstrap, Emit sites, outbox, HTTP manifest, NoopEmitter fallback, integration and chaos tests), dispatching ring:backend-go under TDD. Consumes the validated instrumentation-map.json from ring:mapping-streaming-events. Use after that map exists. Skip for non-Go or when no map is present."

Streaming Instrumentation (lib-streaming)

When to use

  • User requests streaming instrumentation for a Go service with a validated docs/streaming/instrumentation-map.json from ring:mapping-streaming-events
  • Task mentions "wire lib-streaming", "instrument streaming events", "implement event emission", "add streaming.NewBuilder", "Emit business events", "lib-streaming bootstrap"

Skip when

  • Service is not a Go project
  • No instrumentation-map.json present (run ring:mapping-streaming-events first)

You orchestrate. Agents implement. NEVER use Edit/Write/Bash on Go source files. All code changes go through `Task(subagent_type="ring:backend-go")`. TDD mandatory for all implementation gates (RED → GREEN → REFACTOR).

Streaming Architecture

lib-streaming: producer-only event-emission library. Three-step lifecycle:

1. `streaming.NewCatalog(definitions ...EventDefinition) (Catalog, error)` — declare every event up-front (immutable) 2. `streaming.NewBuilder().Source(...).Catalog(catalog).Routes(...).Target(...).Logger(...).MetricsFactory(...).Tracer(...).CircuitBreakerManager(...).OutboxRepository(...).Build(ctx)` — Builder pattern returns `(Emitter, error)`. There is NO `NewProducer` constructor; `*streaming.Producer` is reachable only by type-asserting the `Emitter` returned from `Build(ctx)`, and only when lifecycle methods (`Run`, `RunContext`, `RegisterOutboxRelay`) are needed. 3. `emitter.Emit(ctx, EmitRequest{DefinitionKey, TenantID, Subject, Payload})` from handlers/workers

The `Emitter` interface has THREE methods — `Emit(ctx, EmitRequest) error`, `Close() error`, `Healthy(ctx) error`. Mocks and adapters MUST implement all three.

Wire format: CloudEvents 1.0 binary mode. Each `RouteDefinition` picks a transport: Kafka (topic `lerian.streaming.<resource>.<event>[.vN]`), SQS (queue URL), RabbitMQ (exchange + routing key), EventBridge (bus name), or Custom. Tenant carried on `ce-tenantid` header for CloudEvents-binary transports.

**WebFetch URLs (include in every gate dispatch):**

  • `https://raw.githubusercontent.com/LerianStudio/lib-streaming/main/doc.go`
  • `https://raw.githubusercontent.com/LerianStudio/lib-streaming/main/AGENTS.md`
  • `https://raw.githubusercontent.com/LerianStudio/lib-streaming/main/CHANGELOG.md`

**Three delivery postures:**

| Posture | Direct | Outbox | DLQ | Use when | |---------|--------|--------|-----|----------| | CRITICAL | skip | always | on_routable_failure | Loss is correctness/compliance breach | | IMPORTANT | direct | fallback_on_circuit_open | on_routable_failure | Direct normally; survives broker outage | | OBSERVATIONAL | direct | never | never | Analytics-grade; loss acceptable | | CUSTOM | per-event | per-event | per-event | None of the above fits |

**Canonical import paths:**

| Alias | Import Path | Purpose | |-------|-------------|---------| | `streaming` | `github.com/LerianStudio/lib-streaming` | Producer, Emitter, NewCatalog, EventDefinition | | `streamingtest` | `github.com/LerianStudio/lib-streaming/streamingtest` | MockEmitter (test-only) | | `outbox` | `github.com/LerianStudio/lib-commons/v5/commons/outbox` | Only when Gate 5 active |

**Emitter implementations:**

| Implementation | When | Construction | |----------------|------|--------------| | `*streaming.Producer` (returned as `Emitter`) | `STREAMING_ENABLED=true` | `streaming.NewBuilder().Catalog(catalog).Source(src).Routes(routes...).Target(target).Logger(log).MetricsFactory(mf).Tracer(tr).Build(ctx)` | | NoopEmitter | `STREAMING_ENABLED=false` | `streaming.NewNoopEmitter()` | | `*streamingtest.MockEmitter` | Tests | `streamingtest.NewMockEmitter()` |

Service code depends on `streaming.Emitter` INTERFACE. MUST NOT type-assert to `*Producer` except in bootstrap to wire `Run(launcher)` / `RunContext(ctx, launcher)` / `RegisterOutboxRelay(registry)`. All three implementations satisfy the full three-method interface (`Emit`, `Close`, `Healthy`).

**Mandatory agent instruction (include in EVERY dispatch):**

> WebFetch `https://raw.githubusercontent.com/LerianStudio/lib-streaming/main/doc.go` and `AGENTS.md`. > `docs/streaming/instrumentation-map.json` is the canonical contract — every EventDefinition, Emit site, DeliveryPolicy MUST match exactly. > Tenant from `tmcore.GetTenantIDContext(ctx)` — NEVER hardcode. > TDD: RED → GREEN → REFACTOR for every gate.

Gate Overview

| Gate | Name | Condition | Agent | |------|------|-----------|-------| | 0 | Stack Detection + JSON Validation + Compliance Audit | Always | Orchestrator | | 1 | Codebase Analysis | Always | ring:codebase-explorer | | 1.5 | Visual Implementation Preview | Always; user must approve | ring:visualizing | | 2 | lib-streaming Dependency + Non-Canonical Removal | Skip only if lib-streaming pinned AND zero non-canonical detected | ring:backend-go | | 3 | Catalog Construction + Builder Bootstrap | Always | ring:backend-go | | 4 | Emit Instrumentation per Eventable Point | Always | ring:backend-go | | 5 | Outbox Wiring | Required if any event has `outbox != "never"` | ring:backend-go | | 6 | Manifest HTTP Mount | Required unless service has zero HTTP surface | ring:backend-go | | 7 | Wiring + Lifecycle + Backward Compat | Always — NEVER skippable | ring:backend-go | | 8 | Tests | Always | ring:backend-go | | 9 | Code Review | Always | 9 defaults + triggered specialists in parallel | | 10 | User Validation | Always | User | | 11 | Activation Guide | Always | Orchestrator |

Gates execute sequentially. Gate 5 skip: only if

Read more
Ships withring

Proven engineering practices, enforced through skills. Ring is a comprehensive skills library and workflow system for AI agents that transforms how AI assistants approach software development.

Get the whole plugin

Other skills on ring.