## Documentation Index

Fetch the complete documentation index at: [/llms.txt](https://cosmo-docs.wundergraph.com/llms.txt)

Use this file to discover all available pages before exploring further.

Kafka system connecting CMS, CRM, and ERP to routers via Cosmo Streams

## Minimum requirements

| Package      | Minimum version |
|--------------|----------------|
| controlplane | 0.88.3         |
| router       | 0.88.0         |
| wgc          | 0.55.0         |

Kafka topic arguments support the same [argument template syntax](https://cosmo-docs.wundergraph.com/router/cosmo-streams#the-%E2%80%9Csubjects%E2%80%9D-argument) as the other event providers, for example `employeeUpdated.{{ args.employeeID }}`.
When you use topic templates with Kafka, you must create the resulting name as a Kafka topic in your broker ahead of time. The Router does not create Kafka topics automatically.

## Full schema example

Below is a comprehensive example of how to use Kafka with Cosmo Streams. This guide covers publish, subscribe, and the filter directive. All examples can be modified to suit your specific needs. The schema directives and `edfs__*` types belong to the Cosmo Streams schema contract and must not be modified.

```
# Cosmo Streams

directive @edfs__kafkaPublish(topic: String!, providerId: String! = "default") on FIELD_DEFINITION
directive @edfs__kafkaSubscribe(topics: [String!]!, providerId: String! = "default") on FIELD_DEFINITION

# OpenFederation

directive @openfed__subscriptionFilter(condition: openfed__SubscriptionFilterCondition!) on FIELD_DEFINITION

scalar openfed__SubscriptionFilterValue

input openfed__SubscriptionFieldCondition {
    fieldPath: String!
    values: [openfed__SubscriptionFilterValue]!
}

input openfed__SubscriptionFilterCondition {
    AND: [openfed__SubscriptionFilterCondition!]
    IN: openfed__SubscriptionFieldCondition
    NOT: openfed__SubscriptionFilterCondition
    OR: [openfed__SubscriptionFilterCondition!]
}

# Custom

input UpdateEmployeeInput {
    name: String
    email: String
}

type Mutation {
   updateEmployeeMyKafka(employeeID: Int!, update: UpdateEmployeeInput!): edfs__PublishResult! @edfs__kafkaPublish(topic: "employeeUpdated", providerId: "my-kafka")
}

type Subscription {
    filteredEmployeeUpdatedMyKafka(employeeID: ID!): Employee!
        @edfs__kafkaSubscribe(topics: ["employeeUpdated", "employeeUpdatedTwo"], providerId: "my-kafka")
        @openfed__subscriptionFilter(condition: { IN: { fieldPath: "id", values: [1, 3, 4, 7, 11] } })
    filteredEmployeeUpdatedMyKafkaWithListFieldArguments(firstIds: [ID!]!, secondIds: [ID!]!): Employee!
        @edfs__kafkaSubscribe(topics: ["employeeUpdated", "employeeUpdatedTwo"], providerId: "my-kafka")
    filteredEmployeeUpdatedMyKafkaWithNestedListFieldArgument(input: KafkaInput!): Employee!
        @edfs__kafkaSubscribe(topics: ["employeeUpdated", "employeeUpdatedTwo"], providerId: "my-kafka")
        @openfed__subscriptionFilter(condition: {
            OR: [\
                { IN: { fieldPath: "id", values: ["{{ args.input.ids }}"] } },\
                { IN: { fieldPath: "id", values: [1] } },\
            ],
        })
}

input KafkaInput {
    ids: [Int!]!
}

# Subgraph schema

type Employee @key(fields: "id", resolvable: false) {
  id: Int! @external
}

type edfs__PublishResult {
    success: Boolean!
}
```

You can create the above Event-Driven Graph (EDG—an abstract subgraph) with the following [wgc](https://cosmo-docs.wundergraph.com/cli/intro) command:

```
wgc subgraph publish edg --namespace default --schema eedg.graphqls
```

## Router configuration

Based on the example above, you will need a compatible router configuration.

config.yaml

```
events:
  providers:
    kafka:
      - id: my-kafka # Needs to match with the providerID in the directive
        tls:
          enabled: true
        authentication:
          sasl_plain:
            password: "password"
            username: "username"
        brokers:
          - "localhost:9092"
```

## Example Query

In the example query below, one or more subgraphs have been implemented alongside the Event-Driven Graph to resolve any other fields defined on `Employee`, e.g., `tag` and `details.surname`.

```
subscription {
  filteredEmployeeUpdatedMyKafka(employeeID: 1) {
    id # resolved by the Event-Driven Graph (through the event)
    tag # resolved by another subgraph
    details { # resolved by another subgraph
      surname
    }
  }
}
```

## System diagram

Service–Kafka–Router interaction showing event flow and subscriptions.
