/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.
$ npx -y skills add LerianStudio/ring --skill instrumenting-streaming-events --agent claude-codeHow 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.mdname: 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
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
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.
Repo: LerianStudio/ring
Other skills on ring.
- /analyzing-options
Analyzing different approaches for a task or problem with structured comparisons, effort estimates, and recommendations. Use when facing strategic decisions, architecture choices, or multiple viable approaches. Skip when there's an obvious single approach or the decision is
Open skill - /auditing-production-readiness
Auditing a service's production readiness against Ring engineering standards across base dimensions plus a conditional multi-tenant dimension, then emitting a scored report and an HTML dashboard. Use before production deploy, periodic review, onboarding, or a major release. Skip
Open skill - /cleaning-comments
Cleaning redundant and obvious comments following clean code principles while preserving meaningful documentation. Supports git scope filtering (staged, unstaged, branch, commit-range). Use when code has excessive comments, during code review, or post-refactor cleanup. Skip when
Open skill - /committing-changes
Commit changes with scope allowlist enforcement, atomic grouping, GPG-signed conventional commits, and trailer management. Detects the repo's PR-validation scope policy before proposing any message. Use when the user asks to commit or has changes ready to record. Skip when the
Open skill - /creating-handoffs
Creating a handoff document that captures session state (completed work, decisions, open items, next steps) and delivering it via Plan Mode so the user gets the native 'clear context and continue implementing' resume option. Use when ending a session, when context grows large,
Open skill - /creating-worktrees
Creating an isolated git worktree for parallel branch work: selects the directory by priority order, verifies/adds .gitignore safety, auto-installs the detected toolchain's dependencies, runs a baseline test, and reports readiness. Use before a feature that needs isolation from
Open skill

