Agent skill

Ak Dev New Queue Transport

by yaalalabs in yaalalabs/agent-kernel

Step-by-step guide for adding a new queue transport to Agent Kernel's execution pipeline.

Apache-2.0Auto-check passedBackend & APIs

Install Ak Dev New Queue Transport

skills CLI
$ npx skills add yaalalabs/agent-kernel --skill ak-dev-new-queue-transport -a claude-code

Project install by default; add -g for ~/.claude/skills/.

GitHub CLI
$ gh skill install yaalalabs/agent-kernel ak-dev-new-queue-transport --agent claude-code

Project scope by default; add --scope user for a personal install. Needs GitHub CLI 2.90.0 or later (public preview).

Manual copy
$ git clone --depth 1 https://github.com/yaalalabs/agent-kernel.git skills-src && mkdir -p .claude/skills && cp -r skills-src/.agents/skills/ak-dev-new-queue-transport .claude/skills/ak-dev-new-queue-transport && rm -rf skills-src

Use ~/.claude/skills/ instead of .claude/skills for a personal install. The folder must contain SKILL.md.

Claude Code skills documentation · loads skills from .claude/skills/

Facts

Skill name
ak-dev-new-queue-transport
GitHub stars
192
Token cost
~2.7k tokens
SKILL.md length
1,255 words
Files
1
Skills in repo
23
Repo updated
First seen
Licence
Apache-2.0

At a glance

Step-by-step guide for adding a new queue transport to Agent Kernel's execution pipeline.

  • Works in 6 steps: Implement the Transport → Configuration → Factory Registration → …
  • You need to integrate a new message broker (beyond inmemory
  • SKILL.md covers The Semantics Contract, Step 1: Implement the Transport, Step 2: Configuration and Step 3: Factory Registration, plus 4 more sections
  • Calls uv and make

What it does

Ak Dev New Queue Transport is an agent skill from yaalalabs/agent-kernel. Step-by-step guide for adding a new queue transport to Agent Kernel's execution pipeline. Use this skill when you need to integrate a new message broker (beyond inmemory, SQS, Kafka, and NATS JetStream) behind the QueueTransport/TransportConsumer interface. Covers the queue-semantics contract every transport must reproduce, factory registration, configuration and extras, the QueueTransportContract test suite (fake and live-broker runs), the transport example, and Helm chart wiring.

Its SKILL.md is about 2.7k tokens, which your agent loads only when the skill is triggered. It is a single SKILL.md file with no bundled scripts.

It sits in Backend & APIs, covering Event-driven systems, Container orchestration and Test generation. It works with Apache Kafka and Helm. The repository describes itself as: The Operating System for Scalable Enterprise AI Agents - Run, orchestrate, and deploy Compliant Enterprise AI Agents at scale across frameworks, without lock-in, rewrites or… The licence is Apache-2.0.

When your agent uses it

  • You need to integrate a new message broker (beyond inmemory
  • NATS JetStream) behind the QueueTransport/TransportConsumer interface

Example prompts

  • “/ak-dev-new-queue-transport”

Requirements

  • Docker

Workflow steps

6 steps, taken from the step headings in SKILL.md.

  1. Implement the Transport
  2. Configuration
  3. Factory Registration
  4. Tests
  5. Example
  6. Deployment and Docs Surfaces

What it can do on your machine

Read from SKILL.md and the folder at commit 97fa8d9. It shows what the files ask for, not the result of running them.

  • Tool permissions

    Pre-approves nothing: there is no allowed-tools line, so your agent's usual permission prompts apply.

    From allowed-tools in the SKILL.md frontmatter.

  • Runs code

    Shell commands in SKILL.md call:

    • uv
    • make

    From the folder's file list and the shell code blocks in SKILL.md.

  • Network

    No URLs in SKILL.md. Its commands use uv, which can reach the network depending on how they are called.

    From URLs in SKILL.md, links to its own repository left out.

  • Credentials

    Names no API keys, tokens, secrets or passwords.

    From names ending in _API_KEY, _TOKEN, _SECRET, _KEY or _PASSWORD in SKILL.md.

Context cost

Ak Dev New Queue Transport loads about 2.7k tokens when it runs. Until then it costs about 129 tokens; SKILL.md has 1,255 words of instructions outside code blocks.

Always · name and description, kept in context so the agent knows when to use it
~129
When it runs · the whole SKILL.md, loaded when a task matches
~2.7k

Estimates: characters ÷ 4, the usual rule of thumb; real counts depend on the model's tokenizer. Scripts and assets cost tokens only if the agent reads them.

Safety

Auto-check passed

The automated check found no risky patterns in SKILL.md.

Automated static check — not a guarantee. Review scripts before installing. It scans the text of SKILL.md for risky patterns (piping downloads into a shell, reading credential files, hidden Unicode, destructive commands); files beside SKILL.md are not scanned.

SKILL.md

The full file from yaalalabs/agent-kernel at commit 97fa8d9, republished under its Apache-2.0 licence (© yaalalabs). 1,255 words, ~2,721 tokens.

Download SKILL.mdSave it as .claude/skills/ak-dev-new-queue-transport/SKILL.md (or your agent's skills folder).
name
ak-dev-new-queue-transport
description
Step-by-step guide for adding a new queue transport to Agent Kernel's execution pipeline. Use this skill when you need to integrate a new message broker (beyond in_memory, SQS, Kafka, and NATS JetStream) behind the QueueTransport/TransportConsumer interface. Covers the queue-semantics contract every transport must reproduce, factory registration, configuration and extras, the QueueTransportContract test suite (fake and live-broker runs), the transport example, and Helm chart wiring.
license
Apache-2.0
metadata.author
yaalalabs
metadata.category
developer

Adding a New Queue Transport

This guide walks through adding a new queue transport to the pipeline (ak-py/src/agentkernel/pipeline/). Use the shipped implementations as references, in increasing order of complexity:

  • transport/in_memory.py: the semantics in their purest form, no broker
  • transport/sqs.py: a broker with native FIFO groups, visibility timeout, and dedup
  • transport/nats.py: a broker where per-session ordering is built client-side (partitioned subjects, one durable consumer per partition, max_ack_pending=1)
  • transport/kafka.py + transport/bookkeeping.py: a broker with no per-message acknowledgement model, so receive counts and dedup are rebuilt on a bookkeeping store

Read .agents/skills/ak-dev-architecture (the pipeline section) first if you have not.

The Semantics Contract

Every transport must reproduce the SQS FIFO semantics the pipeline was extracted from (docs/specs/495-onprem-kubernetes/research/current-queue-mode.md):

  1. Per-group FIFO with one in-flight message per group: group_id is the session id; a session's turns never run concurrently or out of order, while distinct sessions run in parallel.
  2. Bounded at-least-once redelivery with an exact receive_count: the ConsumerLoop compares it to max_receive_count to fire the permanent-failure hook, so it must be exact, not approximate.
  3. Publish-time deduplication on dedup_id within a window (SQS parity: 5 minutes).
  4. Attribute round-tripping: QueueMessage.attributes (request id, user id, status code) must survive the trip byte-identically.
  5. Batch fetch with a bounded wait, returning fewer than batch_size rather than blocking past the wait.
  6. Queue isolation: INPUT and OUTPUT never see each other's messages.

Where the broker genuinely cannot provide a guarantee, the contract suite has an explicit, documented opt-out (see timeout_redelivery in pipeline/testing.py, which Kafka sets to False because its consumer model has no visibility timeout). Never fake a guarantee; declare its absence and justify it in the subclass.

Step 1: Implement the Transport

Create ak-py/src/agentkernel/pipeline/transport/<name>.py implementing both ABCs from transport/base.py:

  • QueueTransport: send(queue, message) (map QueueMessage onto the broker's record: body, attributes as headers/metadata, group_id as the ordering key, dedup_id as the dedup token), create_consumer(queue), and optionally check_consumer_capacity(queue, n) (startup warning when consumer threads exceed what the broker can serve in parallel).
  • TransportConsumer: fetch(batch_size, wait_seconds), ack, nack, dead_letter, close(). One consumer instance is created per consumer thread (Kafka needs one client object per thread; the design assumes it everywhere), so instance state needs no locking, but anything class-level does.

Rules learned from the shipped transports:

  • Threads, not asyncio: the pipeline's consumers are threads. If the client library is asyncio-only, bridge through one shared event-loop thread (_NatsLoop in nats.py is the maintainer-recommended pattern; do not spawn a loop per thread).
  • receive_count must be exact. Prefer the broker's own counter (num_delivered, ApproximateReceiveCount); if none exists, count attempts in a BookkeepingStore (transport/bookkeeping.py), keyed so a crash-looping poison message cannot reset itself.
  • Honor fetch_wait_slice_seconds semantics: ConsumerLoop slices waits to stay responsive to shutdown, so a fetch must tolerate short waits without spinning.
  • close() must actually release broker resources (consumer-group membership, subscriptions, background threads). A leaked consumer keeps CI jobs alive after the tests pass.
  • Connection/provisioning caches are class-level and keyed by connection target; provide a reset() classmethod for test isolation (see InMemoryTransport.reset, NatsTransport.reset).
  • Provisioning posture: dev may auto-provision broker objects behind an auto_provision flag, but production fails fast with an AKConfigError naming the missing object and the declarative alternative (NACK CRs, Strimzi topics). Agent Kernel never silently creates production infrastructure.

Step 2: Configuration

In ak-py/src/agentkernel/core/config.py:

  • Add a _<Name>QueueConfig model with the broker's connection and tuning fields (mirror _NatsQueueConfig; every field needs a real description, since they become user docs).
  • Add the optional field to _QueuesConfig and the type name to its type description.
  • Keep input/output blocks backend-neutral: max_receive_count, no_of_consumers, and batch_size are shared knobs, never per-backend.

If the client library is heavy or compiled, add an extra in ak-py/pyproject.toml ([project.optional-dependencies]) named after the transport.

Step 3: Factory Registration

QueueTransportFactory.create() in transport/base.py is an explicit chain: add the branch for your type, guarded by require_extra("<name>", "execution.queues.type: <name>") with the import inside, and add the name to _BUILTIN_TYPES. Fail with AKConfigError when the config block is missing. Anything not in _BUILTIN_TYPES resolves as a dotted path (BYO), so a transport can also live out of tree; built-in status is for transports we test and document.

The factory has a second consumer (#503): the sandbox queue broker passes its own _QueuesConfig-shaped sandbox.broker.queue block through the optional queues_config parameter on resolve_type/create/create_consumer, so a new transport gets sandbox-broker support for free. Read the block handed to you, never AKConfig (the no-argument path keeps reading execution.queues and must stay byte-for-byte unchanged; tests/test_pipeline_factory_seams.py enforces both properties).

Show full SKILL.md (529 more words)Show less

Step 4: Tests

Three layers, all required:

  1. Transport-specific unit tests (ak-py/tests/test_pipeline_<name>_transport.py): envelope/header mapping, orderings, error paths, provisioning create-vs-verify, against a fake broker. Build the fake behind the real client's interface so the transport code is exercised unmodified (see the fake JetStream behind the real _NatsLoop, and the fake in-memory Kafka cluster).
  2. The contract suite, in-repo: subclass QueueTransportContract (pipeline/testing.py) against the fake, implementing make_transport(). Tune ack_wait/fetch_wait/force_redelivery per backend; document every capability opt-out.
  3. The contract suite, live (ak-py/tests/test_transport_contract_live.py): add an env-gated subclass pointing at a real broker (AK_TEST_<NAME>_... env var, skipped when unset) with per-test unique queues/streams/topics for isolation. The transport-integration-tests job in .github/workflows/test-reusable.yaml starts the brokers from the transport examples' compose files and runs this file on every PR: add your broker's compose service there.

Timing traps that only live brokers catch (both found on real servers, invisible on fakes):

  • If a fetch holds a pull/poll request open per partition, the per-partition window (fetch_wait / partitions) must stay below the visibility timeout, or the server redelivers an in-flight message into the still-open request and one fetch returns it twice.
  • The contract's fixed group ids (s0/s1/s2) must land on distinct partitions under the broker's real partitioner. Partitioners are deterministic: compute the mapping (crc32 for the client-side scheme, murmur2 for Kafka) and choose the partition count accordingly instead of hoping.

Step 5: Example

Add examples/transport/<name>/: a two-process app (IOHandler.run() / AgentRunner.run() behind one app.py), a config.yaml with commented tuning values, a docker compose stack with a healthcheck (the CI job relies on up -d --wait <service>), a <name>_tester.py harness (bring the stack up, provision what Agent Kernel deliberately does not, inspect queues), and an app_test.py covering rest_sync, a multi-turn session, and the retry-to-permanent-failure path. Register it in .github/test-config.yaml under the containerized e2e tests.

Step 6: Deployment and Docs Surfaces

  • Helm chart (ak-deployment/ak-k8s/chart/): a transport.<name> values block, its AK_EXECUTION__QUEUES__<NAME>__* env injection in configmap-env.yaml, a KEDA trigger in scaledobject.yaml if a scaler exists, and declarative provisioning CRs if the broker has an operator.
  • Docs: the transport matrix and a "Running Queue Mode on <name>" section in docs/docs/advanced/queue-mode-guide.md; the transports list in docs/docs/deployment/onprem-kubernetes.md if the transport is k8s-relevant; the transport roll call on the docs-site features page (docs/src/pages/features.tsx: the "Queue broker over SQS, Kafka, or NATS" highlight on the Sandboxed Code Execution card, and any other "SQS, Kafka, or NATS" mention found by grepping docs/src/pages/*.tsx).
  • Landing page inventories (docs/src/components/*/data.tsx): a tile in the Cloud & infrastructure row of IntegrationsMarquee/data.tsx (role Queue, href to the queue mode guide, logo or react-icons/si glyph), and the transport in the Queue Pipeline card's tags and description under the Scale tab in FeatureExplorer/data.tsx. Logo sourcing and the build check are in ak-dev-sync-docs-from-branch, Docs-Site Landing and Features Pages.
  • Skills: the pipeline section of .agents/skills/ak-dev-architecture/SKILL.md, and the user-facing queue/deploy content in ak-py/src/agentkernel/skills/ where transports are enumerated.

Definition of Done

  • cd ak-py && uv run pytest: green, including your contract subclass against the fake.
  • Live contract green against a real broker via the compose stack.
  • make lint-check-all: green.
  • Example runs end to end locally (its app_test.py passes against a live agent).
  • Factory rejects a missing config block and a missing extra with actionable errors.
  • Docs and skills surfaces above updated in the same PR.

© yaalalabs, Apache-2.0. Rendered from Markdown: HTML in the file is shown as text, images as links, and headings moved down two levels. Raw file

Files

Just SKILL.md in .agents/skills/ak-dev-new-queue-transport of yaalalabs/agent-kernel.

Open the folder on GitHubat commit 97fa8d9

Compare with similar skills

Ak Dev New Queue Transport next to the 5 skills that share the most tags, products or categories with it. Stars are the repository's; “used in” counts other GitHub owners with a copy.

Ak Dev New Queue Transport compared with similar skills
SkillStarsUsed inTokensAuto-checkLicenceRepo updated
Ak Dev New Queue Transport this skillyaalalabs/agent-kernel192—~2.7kAutomated safety check: PassApache-2.0
Deploying Kafka K8saiskillstore/marketplace430—~1.8kAutomated safety check: PassNone
Opensourcefaqdigoal/blog8.6k—~966Automated safety check: PassGPL-2.0
Windmill Trigger Type Checklistwindmill-labs/windmill18k—~4.7kAutomated safety check: PassCustom licence
FoundatioFoundatioFx/Foundatio2.1k—~3.9kAutomated safety check: PassApache-2.0
Opensource Guide Coachcalf-ai/calfkit-sdk1491 repos~2.1kAutomated safety check: PassApache-2.0

Similar skills

  • Deploying Kafka K8s

    aiskillstore/marketplace

    Deploys Apache Kafka on Kubernetes using the Strimzi operator with KRaft mode.

    430 GitHub stars~1.8k tokensUpdated today
    Backend & APIsAuto-check passed
  • Opensourcefaq

    digoal/blog

    解答与开源产品有关的深度技术问题,输出图文并茂的 Markdown 技术文章。触发条件:用户提出与开源项目(如 PostgreSQL、Redis、Kafka、Kubernetes、ClickHouse、Flink 等)相关的技术问题,并提供源码目录或 URL、deepwiki repo 名称。即使用户只说"帮我解答这个开源问题"或"分析一下这个项目的某个机制",也应使用本…

    8.6k GitHub stars~966 tokensUpdated today
    DatabasesAuto-check passed
  • Windmill Trigger Type Checklist

    windmill-labs/windmill

    Checklist of every backend, frontend, CLI and capture change needed to add a new TriggerCrud-based trigger type, such as Azure, GCP or Kafka, to Windmill.

    18k GitHub stars~4.7k tokensUpdated today
    Backend & APIsAuto-check passed
  • Foundatio

    FoundatioFx/Foundatio

    A skill your agent uses when working with Foundatio infrastructure abstractions for .NET -- caching, queuing, messaging, file storage, distributed locking, or background jobs.

    2.1k GitHub stars~3.9k tokensUpdated today
    Backend & APIsAuto-check passed
  • Opensource Guide Coach

    calf-ai/calfkit-sdk

    A skill your agent uses when a user wants guidance on starting, contributing to, growing, governing, funding, securing, or sustaining an open source project, or asks about contributor onboarding…

    149 GitHub starsUsed in 1 repo~2.1k tokens
    Backend & APIsAuto-check passed
  • Config Breaking Changes

    axelixlabs/axelix

    Review configuration property changes in the Axelix project for breaking changes and migration-policy compliance.

    148 GitHub stars~2.6k tokensUpdated today
    Backend & APIsAuto-check passed

More from yaalalabs/agent-kernel

All 23 skills in this repo
  • Ak Dev Code Quality

    yaalalabs/agent-kernel

    Code quality standards, formatting, Python style rules (classes over script-style functions, configuration-field rules), commit conventions, and PR workflow for Agent Kernel development.

    192 GitHub stars~2.5k tokensUpdated today
    Auto-check passed
  • Ak Dev New Evaluator Provider

    yaalalabs/agent-kernel

    Step-by-step guide for adding a new built-in test evaluator provider to Agent Kernel (beyond DeepEval, Opik and JEV).

    192 GitHub stars~3.4k tokensUpdated today
    Auto-check passed
  • Ak Dev New Guardrail Provider

    yaalalabs/agent-kernel

    Step-by-step guide for adding a new guardrail provider to Agent Kernel.

    192 GitHub stars~3.5k tokensUpdated today
    Auto-check passed
  • Step-by-step guide for adding a new knowledge base backend to Agent Kernel.

    192 GitHub stars~5.1k tokensUpdated today
    Auto-check passed
  • Ak Dev New Messaging Integration

    yaalalabs/agent-kernel

    Step-by-step guide for adding a new messaging platform integration to Agent Kernel.

    192 GitHub stars~4.6k tokensUpdated today
    Auto-check passed
  • Ak Dev New Multimodal Storage

    yaalalabs/agent-kernel

    Step-by-step guide for adding a new multimodal attachment storage backend to Agent Kernel.

    192 GitHub stars~4.3k tokensUpdated today
    Auto-check passed

Categories

Questions about Ak Dev New Queue Transport

What does Ak Dev New Queue Transport do?

Step-by-step guide for adding a new queue transport to Agent Kernel's execution pipeline. Ak Dev New Queue Transport is an agent skill from yaalalabs/agent-kernel. Step-by-step guide for adding a new queue transport to Agent Kernel's execution pipeline.

When should I use Ak Dev New Queue Transport?

Ak Dev New Queue Transport fits situations like: you need to integrate a new message broker (beyond inmemory; NATS JetStream) behind the QueueTransport/TransportConsumer interface.

How do I install Ak Dev New Queue Transport in Claude Code?

Run `npx skills add yaalalabs/agent-kernel --skill ak-dev-new-queue-transport -a claude-code`. Or copy the skill folder (.agents/skills/ak-dev-new-queue-transport in yaalalabs/agent-kernel) into .claude/skills/ak-dev-new-queue-transport in your project. Claude Code loads it when a task matches its description.

How do I install Ak Dev New Queue Transport in Codex?

Run `npx skills add yaalalabs/agent-kernel --skill ak-dev-new-queue-transport -a codex`. Or copy the skill folder (.agents/skills/ak-dev-new-queue-transport in yaalalabs/agent-kernel) into .agents/skills/ak-dev-new-queue-transport in your project. Codex loads it when a task matches its description.

Can I use Ak Dev New Queue Transport in Cursor, Gemini CLI or GitHub Copilot?

Cursor, Gemini CLI, GitHub Copilot and OpenCode also load SKILL.md folders. With the skills CLI, run `npx skills add yaalalabs/agent-kernel --skill ak-dev-new-queue-transport -a cursor` (or -a gemini-cli, github-copilot or opencode for the others). To copy it by hand, put the folder in .cursor/skills/ak-dev-new-queue-transport, .gemini/skills/ak-dev-new-queue-transport, .github/skills/ak-dev-new-queue-transport and .opencode/skills/ak-dev-new-queue-transport in your project.

What does Ak Dev New Queue Transport need to run?

Going by SKILL.md and its folder, Ak Dev New Queue Transport needs the command-line tools its instructions call (uv and make). Our summary lists: Docker.

Does Ak Dev New Queue Transport access the network?

SKILL.md contains no URLs. Its commands use uv, which can reach the network depending on how they are called. This is read from the text; nothing was executed.

Is Ak Dev New Queue Transport safe to install?

Our automated static check of SKILL.md found no risky patterns, such as piping downloads into a shell, reading credential files or hidden Unicode. It is not a guarantee. Review the folder before installing.

What licence does Ak Dev New Queue Transport use?

Ak Dev New Queue Transport is published under the Apache-2.0 licence (declared in SKILL.md). It allows redistribution, so the full SKILL.md is shown on this page.

How many tokens does Ak Dev New Queue Transport use?

About 2.7k tokens (SKILL.md is roughly 11k characters). Agents keep only the skill's name and description in context until a task matches; then they load SKILL.md in full.

What are the alternatives to Ak Dev New Queue Transport?

Skills that share tags, products or a category with Ak Dev New Queue Transport: Deploying Kafka K8s (aiskillstore/marketplace, 430 stars), Opensourcefaq (digoal/blog, 8.6k stars), Windmill Trigger Type Checklist (windmill-labs/windmill, 18k stars) and Foundatio (FoundatioFx/Foundatio, 2.1k stars). The comparison table on this page puts their stars, adoption, token cost, safety result and licence side by side.

Who maintains Ak Dev New Queue Transport?

yaalalabs (a GitHub organization) maintains it in yaalalabs/agent-kernel, which has 192 GitHub stars. The repository holds 23 skills in this directory. The repository was last updated on October 9, 2026.

Source: yaalalabs/agent-kernel on GitHub. Facts on this page come from the repository at the commit we read; the author's words are quoted as theirs.