Repository files navigation

Kafka hook

An incoming webhook endpoint that tries tries to never reject a message. The sequel to kafka-keyvalue in our need for Kafka-related microservices components that do one thing and do it well. To aid both happy and not-so-happy paths for integrations it captures most attributes of the incoming request alongside the payload.

Decisions taken along the way of exploring this concept

  • Using Kafka record headers instead of an envelope makes sense when we want to save all payloads because we don't need to inspect them.
    • Manual inspection can use kafkacat -J and the result will be consumable through jq.
    • Or `kafkacat -f '%k: %h %s\n' etc.
    • Thus we need to flatten http headers fields into individual cloudevents extension keys
  • How do we, as an extension, prefix http headers?
    • ... and the prefix will be prefixed by ce_
    • We're probably predating (and/or obstructing) real extensions for this kind of thing, so let's not use http or header
    • Let's use something reasonable searchable: hook_ because it's more readable than x-yolean- or yolean.se/whatever.
  • With a single topic name, configured at start, we can produce more reliably that with dynamic topic selection
    • TODO we could, and probably should, validate topic existence on start
  • How do we handle long header values?
    • We should probably always ignore cookie
      • can make this configurable
    • We could cap the length, and substring rather than ignore.
  • How do we handle large payloads?
    • TODO A configurable limit, that we set low
  • How do we handle headers with multiple values?
    • Use JAX-RS's concatenation (spoiler: it's commas)

Pixy drop-in replacement, POST without dropped messages

At Yolean we use kafka-pixy for many use cases with occasional, as opposed to constant, message production. We never use pixy for consumption, mainly because kafka-keyvalue sidecars covers a wider use case for occasional consumption.

Happy paths work great with pixy, but it has a habit of silently skipping message production for requests that don't meet assumptions on for example content-type or path. To increase the confusion in such cases, HTTP responses contain no clues.

Kafka-hook is designed to do anything it can to forward the request to kafka, and if it fails anyway it should be expected to log stack traces. Also it tries to send an

Supported:

  • POST /topic/{ignored}/messages
    • As with all messages produced from this service, it's up to the consumer to check integrity and sanity before taking action.
  • {"partition": , "offset": } status 200 response body
  • {"error": } status 500 response body

Not supported:

  • GET
  • POST /cluster/...
  • Pixy's command line arguments
    • We could probably use a Quarkus main method
  • X-Kafka- HTTP headers
    • They do get included with prefix in http headers thouhg
  • sync query parameter (kafka-hook is always sync=true)
  • key query parameter

Builds

JVM:

y-skaffold build --file-output=images-jvm.json

Single-arch native:

y-skaffold build --platform=linux/[choice-of-arch] -p prod-build --file-output=images-native.json

Multi-arch native (expect 3 hrs build time on a 3 core 7Gi Buildkit with qemu):

y-skaffold build -p prod-build --file-output=images-native.json

Workaround for port collision

Note that buildkit defaults to running patforms in parallel, which means that when Quarkus allocates a port during unit tests one of the builds is likely to fail on port already in use.

Can not be done per build, so try adding the following to ~/.config/buildkit/buildkitd.toml:

[worker.oci]
max-parallelism = 1
```
## Dev
```
mvn quarkus:dev --projects=rest \
-Dcloudevent.source-host=https://test.example.net \
-Dcloudevent.type-prefix=net.example.test.
curl --header "Content-Type: application/json" --request POST \
--data '{"n":1}' \
http://localhost:8080/hook/v1/my-event-type
rpk topic --brokers localhost:9092 consume testevents -n 1
```

About

HTTP to Kafka, a webhook endpoint that tries to represent every event with context

Resources

Stars

3 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all \u003cpre\u003e\u003ccode\u003e blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks"); } } catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); } })(); (function(){ try { var __m = "github.com"; var __re = new RegExp('^' + "github\\.com" + '
Skip to content

Repository files navigation

Kafka hook

An incoming webhook endpoint that tries tries to never reject a message. The sequel to kafka-keyvalue in our need for Kafka-related microservices components that do one thing and do it well. To aid both happy and not-so-happy paths for integrations it captures most attributes of the incoming request alongside the payload.

Decisions taken along the way of exploring this concept

  • Using Kafka record headers instead of an envelope makes sense when we want to save all payloads because we don't need to inspect them.
    • Manual inspection can use kafkacat -J and the result will be consumable through jq.
    • Or `kafkacat -f '%k: %h %s\n' etc.
    • Thus we need to flatten http headers fields into individual cloudevents extension keys
  • How do we, as an extension, prefix http headers?
    • ... and the prefix will be prefixed by ce_
    • We're probably predating (and/or obstructing) real extensions for this kind of thing, so let's not use http or header
    • Let's use something reasonable searchable: hook_ because it's more readable than x-yolean- or yolean.se/whatever.
  • With a single topic name, configured at start, we can produce more reliably that with dynamic topic selection
    • TODO we could, and probably should, validate topic existence on start
  • How do we handle long header values?
    • We should probably always ignore cookie
      • can make this configurable
    • We could cap the length, and substring rather than ignore.
  • How do we handle large payloads?
    • TODO A configurable limit, that we set low
  • How do we handle headers with multiple values?
    • Use JAX-RS's concatenation (spoiler: it's commas)

Pixy drop-in replacement, POST without dropped messages

At Yolean we use kafka-pixy for many use cases with occasional, as opposed to constant, message production. We never use pixy for consumption, mainly because kafka-keyvalue sidecars covers a wider use case for occasional consumption.

Happy paths work great with pixy, but it has a habit of silently skipping message production for requests that don't meet assumptions on for example content-type or path. To increase the confusion in such cases, HTTP responses contain no clues.

Kafka-hook is designed to do anything it can to forward the request to kafka, and if it fails anyway it should be expected to log stack traces. Also it tries to send an

Supported:

  • POST /topic/{ignored}/messages
    • As with all messages produced from this service, it's up to the consumer to check integrity and sanity before taking action.
  • {"partition": , "offset": } status 200 response body
  • {"error": } status 500 response body

Not supported:

  • GET
  • POST /cluster/...
  • Pixy's command line arguments
    • We could probably use a Quarkus main method
  • X-Kafka- HTTP headers
    • They do get included with prefix in http headers thouhg
  • sync query parameter (kafka-hook is always sync=true)
  • key query parameter

Builds

JVM:

y-skaffold build --file-output=images-jvm.json

Single-arch native:

y-skaffold build --platform=linux/[choice-of-arch] -p prod-build --file-output=images-native.json

Multi-arch native (expect 3 hrs build time on a 3 core 7Gi Buildkit with qemu):

y-skaffold build -p prod-build --file-output=images-native.json

Workaround for port collision

Note that buildkit defaults to running patforms in parallel, which means that when Quarkus allocates a port during unit tests one of the builds is likely to fail on port already in use.

Can not be done per build, so try adding the following to ~/.config/buildkit/buildkitd.toml:

[worker.oci]
max-parallelism = 1
```
## Dev
```
mvn quarkus:dev --projects=rest \
-Dcloudevent.source-host=https://test.example.net \
-Dcloudevent.type-prefix=net.example.test.
curl --header "Content-Type: application/json" --request POST \
--data '{"n":1}' \
http://localhost:8080/hook/v1/my-event-type
rpk topic --brokers localhost:9092 consume testevents -n 1
```

About

HTTP to Kafka, a webhook endpoint that tries to represent every event with context

Resources

Stars

3 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

Kafka hook

An incoming webhook endpoint that tries tries to never reject a message. The sequel to kafka-keyvalue in our need for Kafka-related microservices components that do one thing and do it well. To aid both happy and not-so-happy paths for integrations it captures most attributes of the incoming request alongside the payload.

Decisions taken along the way of exploring this concept

  • Using Kafka record headers instead of an envelope makes sense when we want to save all payloads because we don't need to inspect them.
    • Manual inspection can use kafkacat -J and the result will be consumable through jq.
    • Or `kafkacat -f '%k: %h %s\n' etc.
    • Thus we need to flatten http headers fields into individual cloudevents extension keys
  • How do we, as an extension, prefix http headers?
    • ... and the prefix will be prefixed by ce_
    • We're probably predating (and/or obstructing) real extensions for this kind of thing, so let's not use http or header
    • Let's use something reasonable searchable: hook_ because it's more readable than x-yolean- or yolean.se/whatever.
  • With a single topic name, configured at start, we can produce more reliably that with dynamic topic selection
    • TODO we could, and probably should, validate topic existence on start
  • How do we handle long header values?
    • We should probably always ignore cookie
      • can make this configurable
    • We could cap the length, and substring rather than ignore.
  • How do we handle large payloads?
    • TODO A configurable limit, that we set low
  • How do we handle headers with multiple values?
    • Use JAX-RS's concatenation (spoiler: it's commas)

Pixy drop-in replacement, POST without dropped messages

At Yolean we use kafka-pixy for many use cases with occasional, as opposed to constant, message production. We never use pixy for consumption, mainly because kafka-keyvalue sidecars covers a wider use case for occasional consumption.

Happy paths work great with pixy, but it has a habit of silently skipping message production for requests that don't meet assumptions on for example content-type or path. To increase the confusion in such cases, HTTP responses contain no clues.

Kafka-hook is designed to do anything it can to forward the request to kafka, and if it fails anyway it should be expected to log stack traces. Also it tries to send an

Supported:

  • POST /topic/{ignored}/messages
    • As with all messages produced from this service, it's up to the consumer to check integrity and sanity before taking action.
  • {"partition": , "offset": } status 200 response body
  • {"error": } status 500 response body

Not supported:

  • GET
  • POST /cluster/...
  • Pixy's command line arguments
    • We could probably use a Quarkus main method
  • X-Kafka- HTTP headers
    • They do get included with prefix in http headers thouhg
  • sync query parameter (kafka-hook is always sync=true)
  • key query parameter

Builds

JVM:

y-skaffold build --file-output=images-jvm.json

Single-arch native:

y-skaffold build --platform=linux/[choice-of-arch] -p prod-build --file-output=images-native.json

Multi-arch native (expect 3 hrs build time on a 3 core 7Gi Buildkit with qemu):

y-skaffold build -p prod-build --file-output=images-native.json

Workaround for port collision

Note that buildkit defaults to running patforms in parallel, which means that when Quarkus allocates a port during unit tests one of the builds is likely to fail on port already in use.

Can not be done per build, so try adding the following to ~/.config/buildkit/buildkitd.toml:

[worker.oci]
max-parallelism = 1
```
## Dev
```
mvn quarkus:dev --projects=rest \
-Dcloudevent.source-host=https://test.example.net \
-Dcloudevent.type-prefix=net.example.test.
curl --header "Content-Type: application/json" --request POST \
--data '{"n":1}' \
http://localhost:8080/hook/v1/my-event-type
rpk topic --brokers localhost:9092 consume testevents -n 1
```

About

HTTP to Kafka, a webhook endpoint that tries to represent every event with context

Resources

Stars

3 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length \u003e 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

Kafka hook

An incoming webhook endpoint that tries tries to never reject a message. The sequel to kafka-keyvalue in our need for Kafka-related microservices components that do one thing and do it well. To aid both happy and not-so-happy paths for integrations it captures most attributes of the incoming request alongside the payload.

Decisions taken along the way of exploring this concept

  • Using Kafka record headers instead of an envelope makes sense when we want to save all payloads because we don't need to inspect them.
    • Manual inspection can use kafkacat -J and the result will be consumable through jq.
    • Or `kafkacat -f '%k: %h %s\n' etc.
    • Thus we need to flatten http headers fields into individual cloudevents extension keys
  • How do we, as an extension, prefix http headers?
    • ... and the prefix will be prefixed by ce_
    • We're probably predating (and/or obstructing) real extensions for this kind of thing, so let's not use http or header
    • Let's use something reasonable searchable: hook_ because it's more readable than x-yolean- or yolean.se/whatever.
  • With a single topic name, configured at start, we can produce more reliably that with dynamic topic selection
    • TODO we could, and probably should, validate topic existence on start
  • How do we handle long header values?
    • We should probably always ignore cookie
      • can make this configurable
    • We could cap the length, and substring rather than ignore.
  • How do we handle large payloads?
    • TODO A configurable limit, that we set low
  • How do we handle headers with multiple values?
    • Use JAX-RS's concatenation (spoiler: it's commas)

Pixy drop-in replacement, POST without dropped messages

At Yolean we use kafka-pixy for many use cases with occasional, as opposed to constant, message production. We never use pixy for consumption, mainly because kafka-keyvalue sidecars covers a wider use case for occasional consumption.

Happy paths work great with pixy, but it has a habit of silently skipping message production for requests that don't meet assumptions on for example content-type or path. To increase the confusion in such cases, HTTP responses contain no clues.

Kafka-hook is designed to do anything it can to forward the request to kafka, and if it fails anyway it should be expected to log stack traces. Also it tries to send an

Supported:

  • POST /topic/{ignored}/messages
    • As with all messages produced from this service, it's up to the consumer to check integrity and sanity before taking action.
  • {"partition": , "offset": } status 200 response body
  • {"error": } status 500 response body

Not supported:

  • GET
  • POST /cluster/...
  • Pixy's command line arguments
    • We could probably use a Quarkus main method
  • X-Kafka- HTTP headers
    • They do get included with prefix in http headers thouhg
  • sync query parameter (kafka-hook is always sync=true)
  • key query parameter

Builds

JVM:

y-skaffold build --file-output=images-jvm.json

Single-arch native:

y-skaffold build --platform=linux/[choice-of-arch] -p prod-build --file-output=images-native.json

Multi-arch native (expect 3 hrs build time on a 3 core 7Gi Buildkit with qemu):

y-skaffold build -p prod-build --file-output=images-native.json

Workaround for port collision

Note that buildkit defaults to running patforms in parallel, which means that when Quarkus allocates a port during unit tests one of the builds is likely to fail on port already in use.

Can not be done per build, so try adding the following to ~/.config/buildkit/buildkitd.toml:

[worker.oci]
max-parallelism = 1
```
## Dev
```
mvn quarkus:dev --projects=rest \
-Dcloudevent.source-host=https://test.example.net \
-Dcloudevent.type-prefix=net.example.test.
curl --header "Content-Type: application/json" --request POST \
--data '{"n":1}' \
http://localhost:8080/hook/v1/my-event-type
rpk topic --brokers localhost:9092 consume testevents -n 1
```

About

HTTP to Kafka, a webhook endpoint that tries to represent every event with context

Resources

Stars

3 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
Skip to content

Repository files navigation

Kafka hook

An incoming webhook endpoint that tries tries to never reject a message. The sequel to kafka-keyvalue in our need for Kafka-related microservices components that do one thing and do it well. To aid both happy and not-so-happy paths for integrations it captures most attributes of the incoming request alongside the payload.

Decisions taken along the way of exploring this concept

  • Using Kafka record headers instead of an envelope makes sense when we want to save all payloads because we don't need to inspect them.
    • Manual inspection can use kafkacat -J and the result will be consumable through jq.
    • Or `kafkacat -f '%k: %h %s\n' etc.
    • Thus we need to flatten http headers fields into individual cloudevents extension keys
  • How do we, as an extension, prefix http headers?
    • ... and the prefix will be prefixed by ce_
    • We're probably predating (and/or obstructing) real extensions for this kind of thing, so let's not use http or header
    • Let's use something reasonable searchable: hook_ because it's more readable than x-yolean- or yolean.se/whatever.
  • With a single topic name, configured at start, we can produce more reliably that with dynamic topic selection
    • TODO we could, and probably should, validate topic existence on start
  • How do we handle long header values?
    • We should probably always ignore cookie
      • can make this configurable
    • We could cap the length, and substring rather than ignore.
  • How do we handle large payloads?
    • TODO A configurable limit, that we set low
  • How do we handle headers with multiple values?
    • Use JAX-RS's concatenation (spoiler: it's commas)

Pixy drop-in replacement, POST without dropped messages

At Yolean we use kafka-pixy for many use cases with occasional, as opposed to constant, message production. We never use pixy for consumption, mainly because kafka-keyvalue sidecars covers a wider use case for occasional consumption.

Happy paths work great with pixy, but it has a habit of silently skipping message production for requests that don't meet assumptions on for example content-type or path. To increase the confusion in such cases, HTTP responses contain no clues.

Kafka-hook is designed to do anything it can to forward the request to kafka, and if it fails anyway it should be expected to log stack traces. Also it tries to send an

Supported:

  • POST /topic/{ignored}/messages
    • As with all messages produced from this service, it's up to the consumer to check integrity and sanity before taking action.
  • {"partition": , "offset": } status 200 response body
  • {"error": } status 500 response body

Not supported:

  • GET
  • POST /cluster/...
  • Pixy's command line arguments
    • We could probably use a Quarkus main method
  • X-Kafka- HTTP headers
    • They do get included with prefix in http headers thouhg
  • sync query parameter (kafka-hook is always sync=true)
  • key query parameter

Builds

JVM:

y-skaffold build --file-output=images-jvm.json

Single-arch native:

y-skaffold build --platform=linux/[choice-of-arch] -p prod-build --file-output=images-native.json

Multi-arch native (expect 3 hrs build time on a 3 core 7Gi Buildkit with qemu):

y-skaffold build -p prod-build --file-output=images-native.json

Workaround for port collision

Note that buildkit defaults to running patforms in parallel, which means that when Quarkus allocates a port during unit tests one of the builds is likely to fail on port already in use.

Can not be done per build, so try adding the following to ~/.config/buildkit/buildkitd.toml:

[worker.oci]
max-parallelism = 1
```
## Dev
```
mvn quarkus:dev --projects=rest \
-Dcloudevent.source-host=https://test.example.net \
-Dcloudevent.type-prefix=net.example.test.
curl --header "Content-Type: application/json" --request POST \
--data '{"n":1}' \
http://localhost:8080/hook/v1/my-event-type
rpk topic --brokers localhost:9092 consume testevents -n 1
```

About

HTTP to Kafka, a webhook endpoint that tries to represent every event with context

Resources

Stars

3 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

Kafka hook

An incoming webhook endpoint that tries tries to never reject a message. The sequel to kafka-keyvalue in our need for Kafka-related microservices components that do one thing and do it well. To aid both happy and not-so-happy paths for integrations it captures most attributes of the incoming request alongside the payload.

Decisions taken along the way of exploring this concept

  • Using Kafka record headers instead of an envelope makes sense when we want to save all payloads because we don't need to inspect them.
    • Manual inspection can use kafkacat -J and the result will be consumable through jq.
    • Or `kafkacat -f '%k: %h %s\n' etc.
    • Thus we need to flatten http headers fields into individual cloudevents extension keys
  • How do we, as an extension, prefix http headers?
    • ... and the prefix will be prefixed by ce_
    • We're probably predating (and/or obstructing) real extensions for this kind of thing, so let's not use http or header
    • Let's use something reasonable searchable: hook_ because it's more readable than x-yolean- or yolean.se/whatever.
  • With a single topic name, configured at start, we can produce more reliably that with dynamic topic selection
    • TODO we could, and probably should, validate topic existence on start
  • How do we handle long header values?
    • We should probably always ignore cookie
      • can make this configurable
    • We could cap the length, and substring rather than ignore.
  • How do we handle large payloads?
    • TODO A configurable limit, that we set low
  • How do we handle headers with multiple values?
    • Use JAX-RS's concatenation (spoiler: it's commas)

Pixy drop-in replacement, POST without dropped messages

At Yolean we use kafka-pixy for many use cases with occasional, as opposed to constant, message production. We never use pixy for consumption, mainly because kafka-keyvalue sidecars covers a wider use case for occasional consumption.

Happy paths work great with pixy, but it has a habit of silently skipping message production for requests that don't meet assumptions on for example content-type or path. To increase the confusion in such cases, HTTP responses contain no clues.

Kafka-hook is designed to do anything it can to forward the request to kafka, and if it fails anyway it should be expected to log stack traces. Also it tries to send an

Supported:

  • POST /topic/{ignored}/messages
    • As with all messages produced from this service, it's up to the consumer to check integrity and sanity before taking action.
  • {"partition": , "offset": } status 200 response body
  • {"error": } status 500 response body

Not supported:

  • GET
  • POST /cluster/...
  • Pixy's command line arguments
    • We could probably use a Quarkus main method
  • X-Kafka- HTTP headers
    • They do get included with prefix in http headers thouhg
  • sync query parameter (kafka-hook is always sync=true)
  • key query parameter

Builds

JVM:

y-skaffold build --file-output=images-jvm.json

Single-arch native:

y-skaffold build --platform=linux/[choice-of-arch] -p prod-build --file-output=images-native.json

Multi-arch native (expect 3 hrs build time on a 3 core 7Gi Buildkit with qemu):

y-skaffold build -p prod-build --file-output=images-native.json

Workaround for port collision

Note that buildkit defaults to running patforms in parallel, which means that when Quarkus allocates a port during unit tests one of the builds is likely to fail on port already in use.

Can not be done per build, so try adding the following to ~/.config/buildkit/buildkitd.toml:

[worker.oci]
max-parallelism = 1
```
## Dev
```
mvn quarkus:dev --projects=rest \
-Dcloudevent.source-host=https://test.example.net \
-Dcloudevent.type-prefix=net.example.test.
curl --header "Content-Type: application/json" --request POST \
--data '{"n":1}' \
http://localhost:8080/hook/v1/my-event-type
rpk topic --brokers localhost:9092 consume testevents -n 1
```

About

HTTP to Kafka, a webhook endpoint that tries to represent every event with context

Resources

Stars

3 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content

Repository files navigation

Kafka hook

An incoming webhook endpoint that tries tries to never reject a message. The sequel to kafka-keyvalue in our need for Kafka-related microservices components that do one thing and do it well. To aid both happy and not-so-happy paths for integrations it captures most attributes of the incoming request alongside the payload.

Decisions taken along the way of exploring this concept

  • Using Kafka record headers instead of an envelope makes sense when we want to save all payloads because we don't need to inspect them.
    • Manual inspection can use kafkacat -J and the result will be consumable through jq.
    • Or `kafkacat -f '%k: %h %s\n' etc.
    • Thus we need to flatten http headers fields into individual cloudevents extension keys
  • How do we, as an extension, prefix http headers?
    • ... and the prefix will be prefixed by ce_
    • We're probably predating (and/or obstructing) real extensions for this kind of thing, so let's not use http or header
    • Let's use something reasonable searchable: hook_ because it's more readable than x-yolean- or yolean.se/whatever.
  • With a single topic name, configured at start, we can produce more reliably that with dynamic topic selection
    • TODO we could, and probably should, validate topic existence on start
  • How do we handle long header values?
    • We should probably always ignore cookie
      • can make this configurable
    • We could cap the length, and substring rather than ignore.
  • How do we handle large payloads?
    • TODO A configurable limit, that we set low
  • How do we handle headers with multiple values?
    • Use JAX-RS's concatenation (spoiler: it's commas)

Pixy drop-in replacement, POST without dropped messages

At Yolean we use kafka-pixy for many use cases with occasional, as opposed to constant, message production. We never use pixy for consumption, mainly because kafka-keyvalue sidecars covers a wider use case for occasional consumption.

Happy paths work great with pixy, but it has a habit of silently skipping message production for requests that don't meet assumptions on for example content-type or path. To increase the confusion in such cases, HTTP responses contain no clues.

Kafka-hook is designed to do anything it can to forward the request to kafka, and if it fails anyway it should be expected to log stack traces. Also it tries to send an

Supported:

  • POST /topic/{ignored}/messages
    • As with all messages produced from this service, it's up to the consumer to check integrity and sanity before taking action.
  • {"partition": , "offset": } status 200 response body
  • {"error": } status 500 response body

Not supported:

  • GET
  • POST /cluster/...
  • Pixy's command line arguments
    • We could probably use a Quarkus main method
  • X-Kafka- HTTP headers
    • They do get included with prefix in http headers thouhg
  • sync query parameter (kafka-hook is always sync=true)
  • key query parameter

Builds

JVM:

y-skaffold build --file-output=images-jvm.json

Single-arch native:

y-skaffold build --platform=linux/[choice-of-arch] -p prod-build --file-output=images-native.json

Multi-arch native (expect 3 hrs build time on a 3 core 7Gi Buildkit with qemu):

y-skaffold build -p prod-build --file-output=images-native.json

Workaround for port collision

Note that buildkit defaults to running patforms in parallel, which means that when Quarkus allocates a port during unit tests one of the builds is likely to fail on port already in use.

Can not be done per build, so try adding the following to ~/.config/buildkit/buildkitd.toml:

[worker.oci]
max-parallelism = 1
```
## Dev
```
mvn quarkus:dev --projects=rest \
-Dcloudevent.source-host=https://test.example.net \
-Dcloudevent.type-prefix=net.example.test.
curl --header "Content-Type: application/json" --request POST \
--data '{"n":1}' \
http://localhost:8080/hook/v1/my-event-type
rpk topic --brokers localhost:9092 consume testevents -n 1
```

About

HTTP to Kafka, a webhook endpoint that tries to represent every event with context

Resources

Stars

3 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages

, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content

Repository files navigation

Kafka hook

An incoming webhook endpoint that tries tries to never reject a message. The sequel to kafka-keyvalue in our need for Kafka-related microservices components that do one thing and do it well. To aid both happy and not-so-happy paths for integrations it captures most attributes of the incoming request alongside the payload.

Decisions taken along the way of exploring this concept

  • Using Kafka record headers instead of an envelope makes sense when we want to save all payloads because we don't need to inspect them.
    • Manual inspection can use kafkacat -J and the result will be consumable through jq.
    • Or `kafkacat -f '%k: %h %s\n' etc.
    • Thus we need to flatten http headers fields into individual cloudevents extension keys
  • How do we, as an extension, prefix http headers?
    • ... and the prefix will be prefixed by ce_
    • We're probably predating (and/or obstructing) real extensions for this kind of thing, so let's not use http or header
    • Let's use something reasonable searchable: hook_ because it's more readable than x-yolean- or yolean.se/whatever.
  • With a single topic name, configured at start, we can produce more reliably that with dynamic topic selection
    • TODO we could, and probably should, validate topic existence on start
  • How do we handle long header values?
    • We should probably always ignore cookie
      • can make this configurable
    • We could cap the length, and substring rather than ignore.
  • How do we handle large payloads?
    • TODO A configurable limit, that we set low
  • How do we handle headers with multiple values?
    • Use JAX-RS's concatenation (spoiler: it's commas)

Pixy drop-in replacement, POST without dropped messages

At Yolean we use kafka-pixy for many use cases with occasional, as opposed to constant, message production. We never use pixy for consumption, mainly because kafka-keyvalue sidecars covers a wider use case for occasional consumption.

Happy paths work great with pixy, but it has a habit of silently skipping message production for requests that don't meet assumptions on for example content-type or path. To increase the confusion in such cases, HTTP responses contain no clues.

Kafka-hook is designed to do anything it can to forward the request to kafka, and if it fails anyway it should be expected to log stack traces. Also it tries to send an

Supported:

  • POST /topic/{ignored}/messages
    • As with all messages produced from this service, it's up to the consumer to check integrity and sanity before taking action.
  • {"partition": , "offset": } status 200 response body
  • {"error": } status 500 response body

Not supported:

  • GET
  • POST /cluster/...
  • Pixy's command line arguments
    • We could probably use a Quarkus main method
  • X-Kafka- HTTP headers
    • They do get included with prefix in http headers thouhg
  • sync query parameter (kafka-hook is always sync=true)
  • key query parameter

Builds

JVM:

y-skaffold build --file-output=images-jvm.json

Single-arch native:

y-skaffold build --platform=linux/[choice-of-arch] -p prod-build --file-output=images-native.json

Multi-arch native (expect 3 hrs build time on a 3 core 7Gi Buildkit with qemu):

y-skaffold build -p prod-build --file-output=images-native.json

Workaround for port collision

Note that buildkit defaults to running patforms in parallel, which means that when Quarkus allocates a port during unit tests one of the builds is likely to fail on port already in use.

Can not be done per build, so try adding the following to ~/.config/buildkit/buildkitd.toml:

[worker.oci]
max-parallelism = 1
```
## Dev
```
mvn quarkus:dev --projects=rest \
-Dcloudevent.source-host=https://test.example.net \
-Dcloudevent.type-prefix=net.example.test.
curl --header "Content-Type: application/json" --request POST \
--data '{"n":1}' \
http://localhost:8080/hook/v1/my-event-type
rpk topic --brokers localhost:9092 consume testevents -n 1
```

About

HTTP to Kafka, a webhook endpoint that tries to represent every event with context

Resources

Stars

3 stars

Watchers

5 watching

Forks

Releases

Packages

Used by

Contributors

Languages