Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
name: CI
name: queue-ci

on:
push:
Expand All@@ -9,16 +9,17 @@ permissions:

jobs:
verify:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 22
cache: npm
- name: Install dependencies
run: npm install
run: npm install --ignore-scripts
- name: Typecheck
run: npm run typecheck
- name: Build
run: npm run build
- name: Tests
run: npm test
- name: Production dependency audit
run: npm audit --omit=dev --audit-level=high
4 changes: 2 additions & 2 deletions package.json
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
{
"name": "skycoin-message-queue",
"version": "1.0.0",
"version": "1.1.0",
"private": true,
"type": "module",
"main": "dist/index.js",
"scripts": {
"build": "tsc -p tsconfig.json",
"typecheck": "tsc -p tsconfig.json --noEmit",
"test": "node --test test/*.test.mjs"
"test": "npm run build && node --test dist/queue.test.js"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Run every compiled test file

When another src/*.test.ts file is added, the build will compile it into dist, but this command passes only dist/queue.test.js to the runner, so CI will silently omit the new tests. Node's local --help describes --test as “launch test runner on startup”; use test discovery or a compiled-test glob rather than selecting this single positional file.

Useful? React with 👍 / 👎.

},
"dependencies": {
"bullmq": "^5.0.0",
Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
import { Queue } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";

export interface MessageEnvelope<T = unknown> {
id: string;
Expand All@@ -9,7 +9,7 @@ export interface MessageEnvelope<T = unknown> {
}

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const connection = new IORedis(redisUrl, { maxRetriesPerRequest: null });
const connection = new Redis(redisUrl, { maxRetriesPerRequest: null });

export const queue = new Queue("skycoin-events", {
connection,
Expand Down
24 changes: 14 additions & 10 deletions src/queue.test.ts
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,17 @@
import { createEnvelope } from './queue';
import assert from "node:assert/strict";
import test from "node:test";
import { createEnvelope } from "./queue.js";

describe('message queue framework', () => {
it('creates deterministic message envelopes when an id is supplied', () => {
const message = createEnvelope('payment.completed', { amount: 25 }, 'evt-1');
expect(message).toEqual(expect.objectContaining({ id: 'evt-1', type: 'payment.completed', payload: { amount: 25 } }));
expect(Date.parse(message.createdAt)).not.toBeNaN();
});
test("creates deterministic message envelopes when an id is supplied", () => {
const message = createEnvelope("payment.completed", { amount: 25 }, "evt-1");
assert.equal(message.id, "evt-1");
assert.equal(message.type, "payment.completed");
assert.deepEqual(message.payload, { amount: 25 });
assert.equal(Number.isNaN(Date.parse(message.createdAt)), false);
});

it('rejects empty message types', () => {
expect(() => createEnvelope('', {})).toThrow('job type is required');
});
test("trims job types and rejects blank message types", () => {
const message = createEnvelope(" feed.index ", {}, "evt-2");
assert.equal(message.type, "feed.index");
assert.throws(() => createEnvelope(" ", {}), /job type is required/);
});
8 changes: 4 additions & 4 deletions src/queue.ts
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
import { Queue, Worker, type Job, type JobsOptions } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";
import { randomUUID } from "node:crypto";

export type SkycoinJob<T = Record<string, unknown>> = { id: string; type: string; payload: T; createdAt: string };

export function createEnvelope<T>(type: string, payload: T, id = randomUUID()): SkycoinJob<T> {
export function createEnvelope<T>(type: string, payload: T, id: string = randomUUID()): SkycoinJob<T> {
if (!type.trim()) throw new Error("job type is required");
return { id, type: type.trim(), payload, createdAt: new Date().toISOString() };
}

export function createQueue<T = Record<string, unknown>>(name = process.env.QUEUE_NAME ?? "skycoin-events") {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return { queue: new Queue<SkycoinJob<T>>(name, { connection }), connection };
}

Expand All@@ -21,6 +21,6 @@ export async function publish<T>(queue: Queue<SkycoinJob<T>>, message: SkycoinJo

export function startWorker<T = Record<string, unknown>>(name: string, handler: (job: Job<SkycoinJob<T>>) => Promise<void>) {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return new Worker<SkycoinJob<T>>(name, handler, { connection, concurrency: Math.max(1, Number(process.env.WORKER_CONCURRENCY ?? 10)) });
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Add copy buttons to all
 blocks
(function() {
function addCopyButtons() {
document.querySelectorAll('pre code').forEach(function(codeBlock) {
if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;
codeBlock.parentElement.setAttribute('data-copy-added', 'true');
var btn = document.createElement('button');
btn.textContent = 'Copy';
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;';
btn.onmouseover = function() { this.style.opacity = '1'; };
btn.onmouseout = function() { this.style.opacity = '0.7'; };
btn.onclick = function() {
navigator.clipboard.writeText(codeBlock.textContent).then(function() {
btn.textContent = 'Copied!';
setTimeout(function() { btn.textContent = 'Copy'; }, 1500);
});
};
codeBlock.parentElement.style.position = 'relative';
codeBlock.parentElement.appendChild(btn);
});
}
addCopyButtons();
// Re-run on dynamic content
var observer = new MutationObserver(addCopyButtons);
observer.observe(document.body, { childList: true, subtree: true });
})();
}
} 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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
name: CI
name: queue-ci

on:
push:
Expand All@@ -9,16 +9,17 @@ permissions:

jobs:
verify:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 22
cache: npm
- name: Install dependencies
run: npm install
run: npm install --ignore-scripts
- name: Typecheck
run: npm run typecheck
- name: Build
run: npm run build
- name: Tests
run: npm test
- name: Production dependency audit
run: npm audit --omit=dev --audit-level=high
4 changes: 2 additions & 2 deletions package.json
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
{
"name": "skycoin-message-queue",
"version": "1.0.0",
"version": "1.1.0",
"private": true,
"type": "module",
"main": "dist/index.js",
"scripts": {
"build": "tsc -p tsconfig.json",
"typecheck": "tsc -p tsconfig.json --noEmit",
"test": "node --test test/*.test.mjs"
"test": "npm run build && node --test dist/queue.test.js"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Run every compiled test file

When another src/*.test.ts file is added, the build will compile it into dist, but this command passes only dist/queue.test.js to the runner, so CI will silently omit the new tests. Node's local --help describes --test as “launch test runner on startup”; use test discovery or a compiled-test glob rather than selecting this single positional file.

Useful? React with 👍 / 👎.

},
"dependencies": {
"bullmq": "^5.0.0",
Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
import { Queue } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";

export interface MessageEnvelope<T = unknown> {
id: string;
Expand All@@ -9,7 +9,7 @@ export interface MessageEnvelope<T = unknown> {
}

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const connection = new IORedis(redisUrl, { maxRetriesPerRequest: null });
const connection = new Redis(redisUrl, { maxRetriesPerRequest: null });

export const queue = new Queue("skycoin-events", {
connection,
Expand Down
24 changes: 14 additions & 10 deletions src/queue.test.ts
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,17 @@
import { createEnvelope } from './queue';
import assert from "node:assert/strict";
import test from "node:test";
import { createEnvelope } from "./queue.js";

describe('message queue framework', () => {
it('creates deterministic message envelopes when an id is supplied', () => {
const message = createEnvelope('payment.completed', { amount: 25 }, 'evt-1');
expect(message).toEqual(expect.objectContaining({ id: 'evt-1', type: 'payment.completed', payload: { amount: 25 } }));
expect(Date.parse(message.createdAt)).not.toBeNaN();
});
test("creates deterministic message envelopes when an id is supplied", () => {
const message = createEnvelope("payment.completed", { amount: 25 }, "evt-1");
assert.equal(message.id, "evt-1");
assert.equal(message.type, "payment.completed");
assert.deepEqual(message.payload, { amount: 25 });
assert.equal(Number.isNaN(Date.parse(message.createdAt)), false);
});

it('rejects empty message types', () => {
expect(() => createEnvelope('', {})).toThrow('job type is required');
});
test("trims job types and rejects blank message types", () => {
const message = createEnvelope(" feed.index ", {}, "evt-2");
assert.equal(message.type, "feed.index");
assert.throws(() => createEnvelope(" ", {}), /job type is required/);
});
8 changes: 4 additions & 4 deletions src/queue.ts
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
import { Queue, Worker, type Job, type JobsOptions } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";
import { randomUUID } from "node:crypto";

export type SkycoinJob<T = Record<string, unknown>> = { id: string; type: string; payload: T; createdAt: string };

export function createEnvelope<T>(type: string, payload: T, id = randomUUID()): SkycoinJob<T> {
export function createEnvelope<T>(type: string, payload: T, id: string = randomUUID()): SkycoinJob<T> {
if (!type.trim()) throw new Error("job type is required");
return { id, type: type.trim(), payload, createdAt: new Date().toISOString() };
}

export function createQueue<T = Record<string, unknown>>(name = process.env.QUEUE_NAME ?? "skycoin-events") {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return { queue: new Queue<SkycoinJob<T>>(name, { connection }), connection };
}

Expand All@@ -21,6 +21,6 @@ export async function publish<T>(queue: Queue<SkycoinJob<T>>, message: SkycoinJo

export function startWorker<T = Record<string, unknown>>(name: string, handler: (job: Job<SkycoinJob<T>>) => Promise<void>) {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return new Worker<SkycoinJob<T>>(name, handler, { connection, concurrency: Math.max(1, Number(process.env.WORKER_CONCURRENCY ?? 10)) });
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Force GitHub README to respect dark mode (function() { var style = document.createElement('style'); style.textContent = ' .markdown-body { color-scheme: dark light; } .markdown-body pre { background: #161b22 !important; } .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; } .markdown-body table th, .markdown-body table td { border-color: #30363d !important; } .markdown-body img { background: #0d1117; } .markdown-body blockquote { border-left-color: #8b949e; } .markdown-body hr { border-color: #30363d; } '; document.head.appendChild(style); })(); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
name: CI
name: queue-ci

on:
push:
Expand All@@ -9,16 +9,17 @@ permissions:

jobs:
verify:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 22
cache: npm
- name: Install dependencies
run: npm install
run: npm install --ignore-scripts
- name: Typecheck
run: npm run typecheck
- name: Build
run: npm run build
- name: Tests
run: npm test
- name: Production dependency audit
run: npm audit --omit=dev --audit-level=high
4 changes: 2 additions & 2 deletions package.json
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
{
"name": "skycoin-message-queue",
"version": "1.0.0",
"version": "1.1.0",
"private": true,
"type": "module",
"main": "dist/index.js",
"scripts": {
"build": "tsc -p tsconfig.json",
"typecheck": "tsc -p tsconfig.json --noEmit",
"test": "node --test test/*.test.mjs"
"test": "npm run build && node --test dist/queue.test.js"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Run every compiled test file

When another src/*.test.ts file is added, the build will compile it into dist, but this command passes only dist/queue.test.js to the runner, so CI will silently omit the new tests. Node's local --help describes --test as “launch test runner on startup”; use test discovery or a compiled-test glob rather than selecting this single positional file.

Useful? React with 👍 / 👎.

},
"dependencies": {
"bullmq": "^5.0.0",
Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
import { Queue } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";

export interface MessageEnvelope<T = unknown> {
id: string;
Expand All@@ -9,7 +9,7 @@ export interface MessageEnvelope<T = unknown> {
}

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const connection = new IORedis(redisUrl, { maxRetriesPerRequest: null });
const connection = new Redis(redisUrl, { maxRetriesPerRequest: null });

export const queue = new Queue("skycoin-events", {
connection,
Expand Down
24 changes: 14 additions & 10 deletions src/queue.test.ts
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,17 @@
import { createEnvelope } from './queue';
import assert from "node:assert/strict";
import test from "node:test";
import { createEnvelope } from "./queue.js";

describe('message queue framework', () => {
it('creates deterministic message envelopes when an id is supplied', () => {
const message = createEnvelope('payment.completed', { amount: 25 }, 'evt-1');
expect(message).toEqual(expect.objectContaining({ id: 'evt-1', type: 'payment.completed', payload: { amount: 25 } }));
expect(Date.parse(message.createdAt)).not.toBeNaN();
});
test("creates deterministic message envelopes when an id is supplied", () => {
const message = createEnvelope("payment.completed", { amount: 25 }, "evt-1");
assert.equal(message.id, "evt-1");
assert.equal(message.type, "payment.completed");
assert.deepEqual(message.payload, { amount: 25 });
assert.equal(Number.isNaN(Date.parse(message.createdAt)), false);
});

it('rejects empty message types', () => {
expect(() => createEnvelope('', {})).toThrow('job type is required');
});
test("trims job types and rejects blank message types", () => {
const message = createEnvelope(" feed.index ", {}, "evt-2");
assert.equal(message.type, "feed.index");
assert.throws(() => createEnvelope(" ", {}), /job type is required/);
});
8 changes: 4 additions & 4 deletions src/queue.ts
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
import { Queue, Worker, type Job, type JobsOptions } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";
import { randomUUID } from "node:crypto";

export type SkycoinJob<T = Record<string, unknown>> = { id: string; type: string; payload: T; createdAt: string };

export function createEnvelope<T>(type: string, payload: T, id = randomUUID()): SkycoinJob<T> {
export function createEnvelope<T>(type: string, payload: T, id: string = randomUUID()): SkycoinJob<T> {
if (!type.trim()) throw new Error("job type is required");
return { id, type: type.trim(), payload, createdAt: new Date().toISOString() };
}

export function createQueue<T = Record<string, unknown>>(name = process.env.QUEUE_NAME ?? "skycoin-events") {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return { queue: new Queue<SkycoinJob<T>>(name, { connection }), connection };
}

Expand All@@ -21,6 +21,6 @@ export async function publish<T>(queue: Queue<SkycoinJob<T>>, message: SkycoinJo

export function startWorker<T = Record<string, unknown>>(name: string, handler: (job: Job<SkycoinJob<T>>) => Promise<void>) {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return new Worker<SkycoinJob<T>>(name, handler, { connection, concurrency: Math.max(1, Number(process.env.WORKER_CONCURRENCY ?? 10)) });
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Highlight search terms from Google/DuckDuckGo/Bing referrer (function() { var ref = document.referrer; var terms = []; if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) { var url = new URL(ref); var q = url.searchParams.get('q') || url.searchParams.get('p'); if (q) { terms = q.split(/\s+/).filter(function(t) { return t.length > 2; }); } } if (terms.length === 0) return; var style = document.createElement('style'); style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }'; document.head.appendChild(style); function highlight(node) { if (node.nodeType === 3) { // text node var text = node.textContent; var found = false; terms.forEach(function(term) { var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\]\\]/g, '\\') + ')', 'gi'); if (regex.test(text)) { found = true; var frag = document.createDocumentFragment(); var parts = text.split(regex); parts.forEach(function(part, i) { if (i % 2 === 0) { frag.appendChild(document.createTextNode(part)); } else { var span = document.createElement('span'); span.className = 'userscript-highlight'; span.textContent = part; frag.appendChild(span); } }); node.parentNode.replaceChild(frag, node); } }); } else if (node.nodeType === 1 && node.childNodes) { // element var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT']; if (!skipTags.includes(node.tagName)) { Array.from(node.childNodes).forEach(highlight); } } } highlight(document.body); // Re-highlight on dynamic content var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1 || node.nodeType === 3) highlight(node); }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
name: CI
name: queue-ci

on:
push:
Expand All@@ -9,16 +9,17 @@ permissions:

jobs:
verify:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 22
cache: npm
- name: Install dependencies
run: npm install
run: npm install --ignore-scripts
- name: Typecheck
run: npm run typecheck
- name: Build
run: npm run build
- name: Tests
run: npm test
- name: Production dependency audit
run: npm audit --omit=dev --audit-level=high
4 changes: 2 additions & 2 deletions package.json
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
{
"name": "skycoin-message-queue",
"version": "1.0.0",
"version": "1.1.0",
"private": true,
"type": "module",
"main": "dist/index.js",
"scripts": {
"build": "tsc -p tsconfig.json",
"typecheck": "tsc -p tsconfig.json --noEmit",
"test": "node --test test/*.test.mjs"
"test": "npm run build && node --test dist/queue.test.js"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Run every compiled test file

When another src/*.test.ts file is added, the build will compile it into dist, but this command passes only dist/queue.test.js to the runner, so CI will silently omit the new tests. Node's local --help describes --test as “launch test runner on startup”; use test discovery or a compiled-test glob rather than selecting this single positional file.

Useful? React with 👍 / 👎.

},
"dependencies": {
"bullmq": "^5.0.0",
Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
import { Queue } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";

export interface MessageEnvelope<T = unknown> {
id: string;
Expand All@@ -9,7 +9,7 @@ export interface MessageEnvelope<T = unknown> {
}

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const connection = new IORedis(redisUrl, { maxRetriesPerRequest: null });
const connection = new Redis(redisUrl, { maxRetriesPerRequest: null });

export const queue = new Queue("skycoin-events", {
connection,
Expand Down
24 changes: 14 additions & 10 deletions src/queue.test.ts
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,17 @@
import { createEnvelope } from './queue';
import assert from "node:assert/strict";
import test from "node:test";
import { createEnvelope } from "./queue.js";

describe('message queue framework', () => {
it('creates deterministic message envelopes when an id is supplied', () => {
const message = createEnvelope('payment.completed', { amount: 25 }, 'evt-1');
expect(message).toEqual(expect.objectContaining({ id: 'evt-1', type: 'payment.completed', payload: { amount: 25 } }));
expect(Date.parse(message.createdAt)).not.toBeNaN();
});
test("creates deterministic message envelopes when an id is supplied", () => {
const message = createEnvelope("payment.completed", { amount: 25 }, "evt-1");
assert.equal(message.id, "evt-1");
assert.equal(message.type, "payment.completed");
assert.deepEqual(message.payload, { amount: 25 });
assert.equal(Number.isNaN(Date.parse(message.createdAt)), false);
});

it('rejects empty message types', () => {
expect(() => createEnvelope('', {})).toThrow('job type is required');
});
test("trims job types and rejects blank message types", () => {
const message = createEnvelope(" feed.index ", {}, "evt-2");
assert.equal(message.type, "feed.index");
assert.throws(() => createEnvelope(" ", {}), /job type is required/);
});
8 changes: 4 additions & 4 deletions src/queue.ts
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
import { Queue, Worker, type Job, type JobsOptions } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";
import { randomUUID } from "node:crypto";

export type SkycoinJob<T = Record<string, unknown>> = { id: string; type: string; payload: T; createdAt: string };

export function createEnvelope<T>(type: string, payload: T, id = randomUUID()): SkycoinJob<T> {
export function createEnvelope<T>(type: string, payload: T, id: string = randomUUID()): SkycoinJob<T> {
if (!type.trim()) throw new Error("job type is required");
return { id, type: type.trim(), payload, createdAt: new Date().toISOString() };
}

export function createQueue<T = Record<string, unknown>>(name = process.env.QUEUE_NAME ?? "skycoin-events") {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return { queue: new Queue<SkycoinJob<T>>(name, { connection }), connection };
}

Expand All@@ -21,6 +21,6 @@ export async function publish<T>(queue: Queue<SkycoinJob<T>>, message: SkycoinJo

export function startWorker<T = Record<string, unknown>>(name: string, handler: (job: Job<SkycoinJob<T>>) => Promise<void>) {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return new Worker<SkycoinJob<T>>(name, handler, { connection, concurrency: Math.max(1, Number(process.env.WORKER_CONCURRENCY ?? 10)) });
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Strip utm_, fbclid, gclid, etc. from all links on page (function() { var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content', 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid', 'ref', 'ref_src', 'source', 'medium', 'campaign']; function cleanUrl(url) { try { var u = new URL(url, window.location.origin); var changed = false; trackingParams.forEach(function(p) { if (u.searchParams.has(p)) { u.searchParams.delete(p); changed = true; } }); return changed ? u.toString() : url; } catch (e) { return url; } } function cleanLinks() { document.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } cleanLinks(); var observer = new MutationObserver(function(mutations) { mutations.forEach(function(m) { m.addedNodes.forEach(function(node) { if (node.nodeType === 1) { if (node.tagName === 'A') cleanLinks(); node.querySelectorAll('a[href]').forEach(function(a) { var clean = cleanUrl(a.href); if (clean !== a.href) a.href = clean; }); } }); }); }); observer.observe(document.body, { childList: true, subtree: true }); })(); } } 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
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
name: CI
name: queue-ci

on:
push:
Expand All@@ -9,16 +9,17 @@ permissions:

jobs:
verify:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 22
cache: npm
- name: Install dependencies
run: npm install
run: npm install --ignore-scripts
- name: Typecheck
run: npm run typecheck
- name: Build
run: npm run build
- name: Tests
run: npm test
- name: Production dependency audit
run: npm audit --omit=dev --audit-level=high
4 changes: 2 additions & 2 deletions package.json
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
{
"name": "skycoin-message-queue",
"version": "1.0.0",
"version": "1.1.0",
"private": true,
"type": "module",
"main": "dist/index.js",
"scripts": {
"build": "tsc -p tsconfig.json",
"typecheck": "tsc -p tsconfig.json --noEmit",
"test": "node --test test/*.test.mjs"
"test": "npm run build && node --test dist/queue.test.js"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Run every compiled test file

When another src/*.test.ts file is added, the build will compile it into dist, but this command passes only dist/queue.test.js to the runner, so CI will silently omit the new tests. Node's local --help describes --test as “launch test runner on startup”; use test discovery or a compiled-test glob rather than selecting this single positional file.

Useful? React with 👍 / 👎.

},
"dependencies": {
"bullmq": "^5.0.0",
Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
import { Queue } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";

export interface MessageEnvelope<T = unknown> {
id: string;
Expand All@@ -9,7 +9,7 @@ export interface MessageEnvelope<T = unknown> {
}

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const connection = new IORedis(redisUrl, { maxRetriesPerRequest: null });
const connection = new Redis(redisUrl, { maxRetriesPerRequest: null });

export const queue = new Queue("skycoin-events", {
connection,
Expand Down
24 changes: 14 additions & 10 deletions src/queue.test.ts
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,17 @@
import { createEnvelope } from './queue';
import assert from "node:assert/strict";
import test from "node:test";
import { createEnvelope } from "./queue.js";

describe('message queue framework', () => {
it('creates deterministic message envelopes when an id is supplied', () => {
const message = createEnvelope('payment.completed', { amount: 25 }, 'evt-1');
expect(message).toEqual(expect.objectContaining({ id: 'evt-1', type: 'payment.completed', payload: { amount: 25 } }));
expect(Date.parse(message.createdAt)).not.toBeNaN();
});
test("creates deterministic message envelopes when an id is supplied", () => {
const message = createEnvelope("payment.completed", { amount: 25 }, "evt-1");
assert.equal(message.id, "evt-1");
assert.equal(message.type, "payment.completed");
assert.deepEqual(message.payload, { amount: 25 });
assert.equal(Number.isNaN(Date.parse(message.createdAt)), false);
});

it('rejects empty message types', () => {
expect(() => createEnvelope('', {})).toThrow('job type is required');
});
test("trims job types and rejects blank message types", () => {
const message = createEnvelope(" feed.index ", {}, "evt-2");
assert.equal(message.type, "feed.index");
assert.throws(() => createEnvelope(" ", {}), /job type is required/);
});
8 changes: 4 additions & 4 deletions src/queue.ts
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
import { Queue, Worker, type Job, type JobsOptions } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";
import { randomUUID } from "node:crypto";

export type SkycoinJob<T = Record<string, unknown>> = { id: string; type: string; payload: T; createdAt: string };

export function createEnvelope<T>(type: string, payload: T, id = randomUUID()): SkycoinJob<T> {
export function createEnvelope<T>(type: string, payload: T, id: string = randomUUID()): SkycoinJob<T> {
if (!type.trim()) throw new Error("job type is required");
return { id, type: type.trim(), payload, createdAt: new Date().toISOString() };
}

export function createQueue<T = Record<string, unknown>>(name = process.env.QUEUE_NAME ?? "skycoin-events") {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return { queue: new Queue<SkycoinJob<T>>(name, { connection }), connection };
}

Expand All@@ -21,6 +21,6 @@ export async function publish<T>(queue: Queue<SkycoinJob<T>>, message: SkycoinJo

export function startWorker<T = Record<string, unknown>>(name: string, handler: (job: Job<SkycoinJob<T>>) => Promise<void>) {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return new Worker<SkycoinJob<T>>(name, handler, { connection, concurrency: Math.max(1, Number(process.env.WORKER_CONCURRENCY ?? 10)) });
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Auto-enable theater mode on YouTube (function() { function tryTheater() { var btn = document.querySelector('button[aria-label="Theater mode"], ytd-player #player button[title="Theater mode"]'); if (btn && !btn.classList.contains('activated')) { btn.click(); } } // Try immediately tryTheater(); // Try after navigation (SPA) var lastUrl = location.href; setInterval(function() { if (location.href !== lastUrl) { lastUrl = location.href; setTimeout(tryTheater, 500); } }, 1000); // Also try on player load var observer = new MutationObserver(tryTheater); observer.observe(document.body, { childList: true, subtree: true }); })(); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
name: CI
name: queue-ci

on:
push:
Expand All@@ -9,16 +9,17 @@ permissions:

jobs:
verify:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 22
cache: npm
- name: Install dependencies
run: npm install
run: npm install --ignore-scripts
- name: Typecheck
run: npm run typecheck
- name: Build
run: npm run build
- name: Tests
run: npm test
- name: Production dependency audit
run: npm audit --omit=dev --audit-level=high
4 changes: 2 additions & 2 deletions package.json
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
{
"name": "skycoin-message-queue",
"version": "1.0.0",
"version": "1.1.0",
"private": true,
"type": "module",
"main": "dist/index.js",
"scripts": {
"build": "tsc -p tsconfig.json",
"typecheck": "tsc -p tsconfig.json --noEmit",
"test": "node --test test/*.test.mjs"
"test": "npm run build && node --test dist/queue.test.js"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Run every compiled test file

When another src/*.test.ts file is added, the build will compile it into dist, but this command passes only dist/queue.test.js to the runner, so CI will silently omit the new tests. Node's local --help describes --test as “launch test runner on startup”; use test discovery or a compiled-test glob rather than selecting this single positional file.

Useful? React with 👍 / 👎.

},
"dependencies": {
"bullmq": "^5.0.0",
Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
import { Queue } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";

export interface MessageEnvelope<T = unknown> {
id: string;
Expand All@@ -9,7 +9,7 @@ export interface MessageEnvelope<T = unknown> {
}

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const connection = new IORedis(redisUrl, { maxRetriesPerRequest: null });
const connection = new Redis(redisUrl, { maxRetriesPerRequest: null });

export const queue = new Queue("skycoin-events", {
connection,
Expand Down
24 changes: 14 additions & 10 deletions src/queue.test.ts
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,17 @@
import { createEnvelope } from './queue';
import assert from "node:assert/strict";
import test from "node:test";
import { createEnvelope } from "./queue.js";

describe('message queue framework', () => {
it('creates deterministic message envelopes when an id is supplied', () => {
const message = createEnvelope('payment.completed', { amount: 25 }, 'evt-1');
expect(message).toEqual(expect.objectContaining({ id: 'evt-1', type: 'payment.completed', payload: { amount: 25 } }));
expect(Date.parse(message.createdAt)).not.toBeNaN();
});
test("creates deterministic message envelopes when an id is supplied", () => {
const message = createEnvelope("payment.completed", { amount: 25 }, "evt-1");
assert.equal(message.id, "evt-1");
assert.equal(message.type, "payment.completed");
assert.deepEqual(message.payload, { amount: 25 });
assert.equal(Number.isNaN(Date.parse(message.createdAt)), false);
});

it('rejects empty message types', () => {
expect(() => createEnvelope('', {})).toThrow('job type is required');
});
test("trims job types and rejects blank message types", () => {
const message = createEnvelope(" feed.index ", {}, "evt-2");
assert.equal(message.type, "feed.index");
assert.throws(() => createEnvelope(" ", {}), /job type is required/);
});
8 changes: 4 additions & 4 deletions src/queue.ts
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
import { Queue, Worker, type Job, type JobsOptions } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";
import { randomUUID } from "node:crypto";

export type SkycoinJob<T = Record<string, unknown>> = { id: string; type: string; payload: T; createdAt: string };

export function createEnvelope<T>(type: string, payload: T, id = randomUUID()): SkycoinJob<T> {
export function createEnvelope<T>(type: string, payload: T, id: string = randomUUID()): SkycoinJob<T> {
if (!type.trim()) throw new Error("job type is required");
return { id, type: type.trim(), payload, createdAt: new Date().toISOString() };
}

export function createQueue<T = Record<string, unknown>>(name = process.env.QUEUE_NAME ?? "skycoin-events") {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return { queue: new Queue<SkycoinJob<T>>(name, { connection }), connection };
}

Expand All@@ -21,6 +21,6 @@ export async function publish<T>(queue: Queue<SkycoinJob<T>>, message: SkycoinJo

export function startWorker<T = Record<string, unknown>>(name: string, handler: (job: Job<SkycoinJob<T>>) => Promise<void>) {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return new Worker<SkycoinJob<T>>(name, handler, { connection, concurrency: Math.max(1, Number(process.env.WORKER_CONCURRENCY ?? 10)) });
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Remove or un-stick sticky/fixed headers that block content (function() { function unstick() { document.querySelectorAll('header, nav, [role="banner"], .header, .navbar, .sticky, .fixed-top, [style*="position: fixed"], [style*="position:sticky"]').forEach(function(el) { if (el.style.position === 'fixed' || el.style.position === 'sticky' || getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') { el.style.position = 'static'; el.style.top = 'auto'; el.style.zIndex = 'auto'; } }); } unstick(); var observer = new MutationObserver(unstick); observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] }); })(); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
name: CI
name: queue-ci

on:
push:
Expand All@@ -9,16 +9,17 @@ permissions:

jobs:
verify:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 22
cache: npm
- name: Install dependencies
run: npm install
run: npm install --ignore-scripts
- name: Typecheck
run: npm run typecheck
- name: Build
run: npm run build
- name: Tests
run: npm test
- name: Production dependency audit
run: npm audit --omit=dev --audit-level=high
4 changes: 2 additions & 2 deletions package.json
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
{
"name": "skycoin-message-queue",
"version": "1.0.0",
"version": "1.1.0",
"private": true,
"type": "module",
"main": "dist/index.js",
"scripts": {
"build": "tsc -p tsconfig.json",
"typecheck": "tsc -p tsconfig.json --noEmit",
"test": "node --test test/*.test.mjs"
"test": "npm run build && node --test dist/queue.test.js"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Run every compiled test file

When another src/*.test.ts file is added, the build will compile it into dist, but this command passes only dist/queue.test.js to the runner, so CI will silently omit the new tests. Node's local --help describes --test as “launch test runner on startup”; use test discovery or a compiled-test glob rather than selecting this single positional file.

Useful? React with 👍 / 👎.

},
"dependencies": {
"bullmq": "^5.0.0",
Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
import { Queue } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";

export interface MessageEnvelope<T = unknown> {
id: string;
Expand All@@ -9,7 +9,7 @@ export interface MessageEnvelope<T = unknown> {
}

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const connection = new IORedis(redisUrl, { maxRetriesPerRequest: null });
const connection = new Redis(redisUrl, { maxRetriesPerRequest: null });

export const queue = new Queue("skycoin-events", {
connection,
Expand Down
24 changes: 14 additions & 10 deletions src/queue.test.ts
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,17 @@
import { createEnvelope } from './queue';
import assert from "node:assert/strict";
import test from "node:test";
import { createEnvelope } from "./queue.js";

describe('message queue framework', () => {
it('creates deterministic message envelopes when an id is supplied', () => {
const message = createEnvelope('payment.completed', { amount: 25 }, 'evt-1');
expect(message).toEqual(expect.objectContaining({ id: 'evt-1', type: 'payment.completed', payload: { amount: 25 } }));
expect(Date.parse(message.createdAt)).not.toBeNaN();
});
test("creates deterministic message envelopes when an id is supplied", () => {
const message = createEnvelope("payment.completed", { amount: 25 }, "evt-1");
assert.equal(message.id, "evt-1");
assert.equal(message.type, "payment.completed");
assert.deepEqual(message.payload, { amount: 25 });
assert.equal(Number.isNaN(Date.parse(message.createdAt)), false);
});

it('rejects empty message types', () => {
expect(() => createEnvelope('', {})).toThrow('job type is required');
});
test("trims job types and rejects blank message types", () => {
const message = createEnvelope(" feed.index ", {}, "evt-2");
assert.equal(message.type, "feed.index");
assert.throws(() => createEnvelope(" ", {}), /job type is required/);
});
8 changes: 4 additions & 4 deletions src/queue.ts
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
import { Queue, Worker, type Job, type JobsOptions } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";
import { randomUUID } from "node:crypto";

export type SkycoinJob<T = Record<string, unknown>> = { id: string; type: string; payload: T; createdAt: string };

export function createEnvelope<T>(type: string, payload: T, id = randomUUID()): SkycoinJob<T> {
export function createEnvelope<T>(type: string, payload: T, id: string = randomUUID()): SkycoinJob<T> {
if (!type.trim()) throw new Error("job type is required");
return { id, type: type.trim(), payload, createdAt: new Date().toISOString() };
}

export function createQueue<T = Record<string, unknown>>(name = process.env.QUEUE_NAME ?? "skycoin-events") {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return { queue: new Queue<SkycoinJob<T>>(name, { connection }), connection };
}

Expand All@@ -21,6 +21,6 @@ export async function publish<T>(queue: Queue<SkycoinJob<T>>, message: SkycoinJo

export function startWorker<T = Record<string, unknown>>(name: string, handler: (job: Job<SkycoinJob<T>>) => Promise<void>) {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return new Worker<SkycoinJob<T>>(name, handler, { connection, concurrency: Math.max(1, Number(process.env.WORKER_CONCURRENCY ?? 10)) });
}
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { // Universal Dark Mode - works on any site (function() { var enabled = true; function applyDarkMode() { if (!enabled) return; // Create style element if it doesn't exist var style = document.getElementById('universal-dark-mode-style'); if (!style) { style = document.createElement('style'); style.id = 'universal-dark-mode-style'; document.head.appendChild(style); } // Dark mode CSS - inverts colors but preserves images/video style.textContent = ' /* Invert everything except media */ html { filter: invert(1) hue-rotate(180deg) !important; background: #1a1a2e !important; } /* Restore images, videos, iframes, canvas */ img, video, iframe, canvas, svg, picture, [style*="background-image"] { filter: invert(1) hue-rotate(180deg) !important; } /* Preserve specific elements that should not be inverted */ .no-dark-mode, .no-dark-mode *, [data-theme="light"], [data-theme="light"], .ace_editor, .ace_editor *, .CodeMirror, .CodeMirror *, .monaco-editor, .monaco-editor *, .markdown-body pre, .markdown-body pre *, .highlight, .highlight *, pre code, pre code * { filter: none !important; } /* Fix common UI elements */ .modal, .popup, .dropdown-menu, .tooltip, .popover { filter: invert(1) hue-rotate(180deg) !important; background: #2d2d44 !important; border-color: #444 !important; } /* Scrollbars */ ::-webkit-scrollbar { background: #1a1a2e !important; } ::-webkit-scrollbar-thumb { background: #444 !important; } ::-webkit-scrollbar-thumb:hover { background: #555 !important; } /* Selection */ ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; } ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; } '; } function removeDarkMode() { var style = document.getElementById('universal-dark-mode-style'); if (style) style.remove(); } // Toggle with Alt+Shift+D document.addEventListener('keydown', function(e) { if (e.altKey && e.shiftKey && e.key === 'D') { e.preventDefault(); enabled = !enabled; if (enabled) { applyDarkMode(); console.log('[Universal Dark Mode] Enabled'); } else { removeDarkMode(); console.log('[Universal Dark Mode] Disabled'); } } }); // Apply on load applyDarkMode(); // Re-apply on dynamic content var observer = new MutationObserver(function(mutations) { if (enabled && !document.getElementById('universal-dark-mode-style')) { applyDarkMode(); } }); observer.observe(document.head, { childList: true }); console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle'); })(); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 7 additions & 6 deletions .github/workflows/ci.yml
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
name: CI
name: queue-ci

on:
push:
Expand All@@ -9,16 +9,17 @@ permissions:

jobs:
verify:
runs-on: ubuntu-latest
runs-on: ubuntu-24.04
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
with:
node-version: 22
cache: npm
- name: Install dependencies
run: npm install
run: npm install --ignore-scripts
- name: Typecheck
run: npm run typecheck
- name: Build
run: npm run build
- name: Tests
run: npm test
- name: Production dependency audit
run: npm audit --omit=dev --audit-level=high
4 changes: 2 additions & 2 deletions package.json
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,13 @@
{
"name": "skycoin-message-queue",
"version": "1.0.0",
"version": "1.1.0",
"private": true,
"type": "module",
"main": "dist/index.js",
"scripts": {
"build": "tsc -p tsconfig.json",
"typecheck": "tsc -p tsconfig.json --noEmit",
"test": "node --test test/*.test.mjs"
"test": "npm run build && node --test dist/queue.test.js"

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Run every compiled test file

When another src/*.test.ts file is added, the build will compile it into dist, but this command passes only dist/queue.test.js to the runner, so CI will silently omit the new tests. Node's local --help describes --test as “launch test runner on startup”; use test discovery or a compiled-test glob rather than selecting this single positional file.

Useful? React with 👍 / 👎.

},
"dependencies": {
"bullmq": "^5.0.0",
Expand Down
4 changes: 2 additions & 2 deletions src/index.ts
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,5 @@
import { Queue } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";

export interface MessageEnvelope<T = unknown> {
id: string;
Expand All@@ -9,7 +9,7 @@ export interface MessageEnvelope<T = unknown> {
}

const redisUrl = process.env.REDIS_URL ?? "redis://localhost:6379";
const connection = new IORedis(redisUrl, { maxRetriesPerRequest: null });
const connection = new Redis(redisUrl, { maxRetriesPerRequest: null });

export const queue = new Queue("skycoin-events", {
connection,
Expand Down
24 changes: 14 additions & 10 deletions src/queue.test.ts
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,17 @@
import { createEnvelope } from './queue';
import assert from "node:assert/strict";
import test from "node:test";
import { createEnvelope } from "./queue.js";

describe('message queue framework', () => {
it('creates deterministic message envelopes when an id is supplied', () => {
const message = createEnvelope('payment.completed', { amount: 25 }, 'evt-1');
expect(message).toEqual(expect.objectContaining({ id: 'evt-1', type: 'payment.completed', payload: { amount: 25 } }));
expect(Date.parse(message.createdAt)).not.toBeNaN();
});
test("creates deterministic message envelopes when an id is supplied", () => {
const message = createEnvelope("payment.completed", { amount: 25 }, "evt-1");
assert.equal(message.id, "evt-1");
assert.equal(message.type, "payment.completed");
assert.deepEqual(message.payload, { amount: 25 });
assert.equal(Number.isNaN(Date.parse(message.createdAt)), false);
});

it('rejects empty message types', () => {
expect(() => createEnvelope('', {})).toThrow('job type is required');
});
test("trims job types and rejects blank message types", () => {
const message = createEnvelope(" feed.index ", {}, "evt-2");
assert.equal(message.type, "feed.index");
assert.throws(() => createEnvelope(" ", {}), /job type is required/);
});
8 changes: 4 additions & 4 deletions src/queue.ts
Original file line numberDiff line numberDiff line change
@@ -1,17 +1,17 @@
import { Queue, Worker, type Job, type JobsOptions } from "bullmq";
import IORedis from "ioredis";
import { Redis } from "ioredis";
import { randomUUID } from "node:crypto";

export type SkycoinJob<T = Record<string, unknown>> = { id: string; type: string; payload: T; createdAt: string };

export function createEnvelope<T>(type: string, payload: T, id = randomUUID()): SkycoinJob<T> {
export function createEnvelope<T>(type: string, payload: T, id: string = randomUUID()): SkycoinJob<T> {
if (!type.trim()) throw new Error("job type is required");
return { id, type: type.trim(), payload, createdAt: new Date().toISOString() };
}

export function createQueue<T = Record<string, unknown>>(name = process.env.QUEUE_NAME ?? "skycoin-events") {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return { queue: new Queue<SkycoinJob<T>>(name, { connection }), connection };
}

Expand All@@ -21,6 +21,6 @@ export async function publish<T>(queue: Queue<SkycoinJob<T>>, message: SkycoinJo

export function startWorker<T = Record<string, unknown>>(name: string, handler: (job: Job<SkycoinJob<T>>) => Promise<void>) {
if (!name.trim()) throw new Error("queue name is required");
const connection = new IORedis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
const connection = new Redis(process.env.REDIS_URL ?? "redis://localhost:6379", { maxRetriesPerRequest: null });
return new Worker<SkycoinJob<T>>(name, handler, { connection, concurrency: Math.max(1, Number(process.env.WORKER_CONCURRENCY ?? 10)) });
}
Loading