Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Provider/Consumer proposition for topics #830

Draft
wants to merge 21 commits into
base: devel
Choose a base branch
from
Draft
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
2 changes: 1 addition & 1 deletion .eslintrc
Original file line number Diff line number Diff line change
Expand Up @@ -202,7 +202,7 @@
"import/no-relative-packages": 2, // Prevent importing packages through relative paths ()

"import/no-deprecated": 1, // Report imported names marked with @deprecated documentation tag ()
"import/no-extraneous-dependencies": 1, // Forbid the use of extraneous packages ()
"import/no-extraneous-dependencies": [ "error", { "devDependencies": true } ], // Forbid the use of extraneous packages ()
"import/no-mutable-exports": 1, // Forbid the use of mutable exports with var or let. ()
"import/no-unused-modules": 1, // Report modules without exports, or exports without matching import in another module ()

Expand Down
32 changes: 32 additions & 0 deletions bdd/data/sequences/data-gen/index.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
const { EventEmitter } = require("events");

/**
* Simple data generator with rate controls
*
* @param {never} _input input ignored
* @param {number} chunkSize size of a single chunk (3 megs by default)
* @param {number} chunkFrequency rate in Hz (1 Hz)
*/
module.exports = async function* (_input, chunkSize = 3 << 20, chunkFrequency = 1) {
const ee = new EventEmitter();
let i = 0;
let end = false;

setInterval(() => ee.emit("tick", i++), 1000 / chunkFrequency);
this.on("stop", () => { end = true; });

let n = 0;

await new Promise(res => this.once("start", res));

while (!end) {
const chunkPayload = `${Date.now()}|${n++}`;
const chunk = `${chunkPayload}${" ".repeat(chunkSize - chunkPayload.length)}`;

yield chunk;

if (n >= i) {
await new Promise(res => ee.once("tick", res));
}
}
};
22 changes: 22 additions & 0 deletions bdd/data/sequences/data-gen/package.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
{
"name": "@scramjet/data-generator",
"version": "0.33.0",
"description": "",
"main": "index.js",
"scripts": {
"build:refapps": "node ../../../scripts/build-all.js --copy-dist",
"postbuild:refapps": "yarn packseq",
"packseq": "node ../../../scripts/packsequence.js",
"clean": "rm -rf ./dist .bic_cache"
},
"author": "Scramjet <[email protected]>",
"license": "ISC",
"devDependencies": {
"@scramjet/types": "^0.33.0",
"@types/node": "15.12.5"
},
"repository": {
"type": "git",
"url": "https://github.com/scramjetorg/transform-hub.git"
}
}
194 changes: 194 additions & 0 deletions bdd/data/sequences/persistentSeq/jest.config.js
Original file line number Diff line number Diff line change
@@ -0,0 +1,194 @@
/** @type {import('ts-jest').JestConfigWithTsJest} */
module.exports = {
// All imported modules in your tests should be mocked automatically
// automock: false,

// Stop running tests after `n` failures
// bail: 0,

// The directory where Jest should store its cached dependency information
// cacheDirectory: "/tmp/jest_rw",

// Automatically clear mock calls, instances, contexts and results before every test
// clearMocks: true,

// Indicates whether the coverage information should be collected while executing the test
// collectCoverage: true,

// An array of glob patterns indicating a set of files for which coverage information should be collected
// collectCoverageFrom: undefined,

// The directory where Jest should output its coverage files
coverageDirectory: "coverage",

// An array of regexp pattern strings used to skip coverage collection
coveragePathIgnorePatterns: [
"/node_modules/",
"/dist/"
],

// Indicates which provider should be used to instrument code for coverage
coverageProvider: "v8",

// A list of reporter names that Jest uses when writing coverage reports
coverageReporters: [
// "json",
"text",
// "lcov",
// "clover"
],

// An object that configures minimum threshold enforcement for coverage results
// coverageThreshold: undefined,

// A path to a custom dependency extractor
// dependencyExtractor: undefined,

// Make calling deprecated APIs throw helpful error messages
errorOnDeprecated: true,

// The default configuration for fake timers
// fakeTimers: {
// "enableGlobally": false
// },

// Force coverage collection from ignored files using an array of glob patterns
// forceCoverageMatch: [],

// A path to a module which exports an async function that is triggered once before all test suites
// globalSetup: undefined,

// A path to a module which exports an async function that is triggered once after all test suites
// globalTeardown: undefined,

// A set of global variables that need to be available in all test environments
// globals: {},

// The maximum amount of workers used to run your tests. Can be specified as % or a number. E.g. maxWorkers: 10% will use 10% of your CPU amount + 1 as the maximum worker number. maxWorkers: 2 will use a maximum of 2 workers.
// maxWorkers: "50%",

// An array of directory names to be searched recursively up from the requiring module's location
// moduleDirectories: [
// "node_modules"
// ],

// An array of file extensions your modules use
// moduleFileExtensions: [
// "js",
// "mjs",
// "cjs",
// "jsx",
// "ts",
// "tsx",
// "json",
// "node"
// ],

// A map from regular expressions to module names or to arrays of module names that allow to stub out resources with a single module
// moduleNameMapper: {},

// An array of regexp pattern strings, matched against all module paths before considered 'visible' to the module loader
// modulePathIgnorePatterns: [],

// Activates notifications for test results
// notify: false,

// An enum that specifies notification mode. Requires { notify: true }
// notifyMode: "failure-change",

// A preset that is used as a base for Jest's configuration
preset: 'ts-jest',

// Run tests from one or more projects
// projects: undefined,

// Use this configuration option to add custom reporters to Jest
// reporters: undefined,

// Automatically reset mock state before every test
// resetMocks: false,

// Reset the module registry before running each individual test
// resetModules: false,

// A path to a custom resolver
// resolver: undefined,

// Automatically restore mock state and implementation before every test
// restoreMocks: false,

// The root directory that Jest should scan for tests and modules within
// rootDir: undefined,

// A list of paths to directories that Jest should use to search for files in
// roots: [
// "<rootDir>"
// ],

// Allows you to use a custom runner instead of Jest's default test runner
// runner: "jest-runner",

// The paths to modules that run some code to configure or set up the testing environment before each test
// setupFiles: [],

// A list of paths to modules that run some code to configure or set up the testing framework before each test
// setupFilesAfterEnv: [],

// The number of seconds after which a test is considered as slow and reported as such in the results.
slowTestThreshold: 5,

// A list of paths to snapshot serializer modules Jest should use for snapshot testing
// snapshotSerializers: [],

// The test environment that will be used for testing
testEnvironment: 'node',

// Options that will be passed to the testEnvironment
// testEnvironmentOptions: {},

// Adds a location field to test results
// testLocationInResults: false,

// The glob patterns Jest uses to detect test files
testMatch: [
"**/__tests__/**/*.[jt]s?(x)",
"**/test/**/*.[jt]s?(x)",
"**/?(*.)+(spec|test).[tj]s?(x)"
],

// An array of regexp pattern strings that are matched against all test paths, matched tests are skipped
testPathIgnorePatterns: [
"/node_modules/",
"/dist/",
],

// The regexp pattern or array of patterns that Jest uses to detect test files
// testRegex: [],

// This option allows the use of a custom results processor
// testResultsProcessor: undefined,

// This option allows use of a custom test runner
// testRunner: "jest-circus/runner",

// A map from regular expressions to paths to transformers
// transform: undefined,

// An array of regexp pattern strings that are matched against all source file paths, matched files will skip transformation
// transformIgnorePatterns: [
// "/node_modules/",
// "\\.pnp\\.[^\\/]+$"
// ],

// An array of regexp pattern strings that are matched against all modules before the module loader will automatically return a mock for them
// unmockedModulePathPatterns: undefined,

// Indicates whether each individual test should be reported during the run
// verbose: undefined,

// An array of regexp patterns that are matched against all source file paths before re-running tests in watch mode
// watchPathIgnorePatterns: [],

// Whether to use watchman for file crawling
// watchman: true,
};
33 changes: 33 additions & 0 deletions bdd/data/sequences/persistentSeq/package.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,33 @@
{
"name": "persistentseq",
"version": "1.0.0",
"description": "",
"main": "index.js",
"scripts": {
"predeploy": "tsc -p tsconfig.json && cp package.json dist/ && (cd dist && npm i --omit=dev)",
"test": "jest --verbose"
},
"engines": {
"node": ">=16"
},
"repository": {
"type": "git",
"url": "git+https://github.com/scramjetorg/create-sequence.git"
},
"bugs": {
"url": "https://github.com/scramjetorg/create-sequence/issues"
},
"homepage": "https://github.com/scramjetorg/create-sequence#readme",
"devDependencies": {
"@scramjet/types": "^0.31.2",
"@types/jest": "^29.5.1",
"@types/node": "15.12.5",
"jest": "^29.5.0",
"scramjet": "^4.36.9",
"ts-jest": "^29.1.0",
"tsc": "^2.0.4",
"typescript": "^4.9.4"
},
"author": "",
"license": "ISC"
}
74 changes: 74 additions & 0 deletions bdd/data/sequences/persistentSeq/src/backupingStream.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,74 @@
import { openSync, readSync, writeSync } from "fs";
import { Duplex, DuplexOptions } from "stream";

type BackupingStreamOptions = Pick<DuplexOptions, "encoding" | "highWaterMark" | "readableHighWaterMark" | "writableHighWaterMark">

class BackupingStream extends Duplex {
writeFd: number;
readFd: number;
bytesWritten: number;
bytesRead: number;
readonly backupFile: string;

constructor(backupFile: string, opts?: BackupingStreamOptions) {
super({ ...opts });
this.backupFile = backupFile;
this.writeFd = openSync(this.backupFile, "w+");
this.readFd = openSync(this.backupFile, "r");
this.bytesWritten = 0;
this.bytesRead = 0;
}

get bytesInBackup() { return this.bytesWritten - this.bytesRead; }

_write(chunk: any, encoding: BufferEncoding,
callback: (error?: Error | null | undefined) => void): void {
if (this.bytesInBackup <= 0 && this.readableFlowing === true) {
this.push(chunk, encoding);
callback();
this.resetBackupFile();
return;
}
const bytesWritten = this.writeBackup(chunk, encoding);

this.bytesWritten += bytesWritten;
callback();
}

_read(size: number): void {
if (this.bytesInBackup <= 0) return;

this.pushFromBackupToInternalBuffer(size);
}

resume(): this {
this.pushFromBackupToInternalBuffer(this.readableHighWaterMark);
return super.resume();
}

protected writeBackup(chunk: any, encoding: BufferEncoding) {
if (typeof chunk === "string") {
return writeSync(this.writeFd, chunk, this.bytesWritten, encoding);
}
return writeSync(this.writeFd, chunk, 0, undefined, this.bytesWritten);
}

protected pushFromBackupToInternalBuffer(size: number) {
const sizeInBuffer = this.readableHighWaterMark - this.readableLength;
const readSize = Math.min(size, sizeInBuffer);
const buffer = Buffer.alloc(readSize);
const bytesRead = readSync(this.readFd, buffer, 0, readSize, this.bytesRead);

this.bytesRead += bytesRead;
this.push(buffer.subarray(0, bytesRead));
}

protected resetBackupFile() {
if (this.bytesWritten <= 0 || this.bytesWritten !== this.bytesRead) return;
this.bytesWritten = 0;
this.bytesRead = 0;
this.writeFd = openSync(this.backupFile, "w+");
}
}

export default BackupingStream;
16 changes: 16 additions & 0 deletions bdd/data/sequences/persistentSeq/src/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
import { TransformApp } from "@scramjet/types";
import { PassThrough } from "stream";
import BackupingStream from "./backupingStream";
import { resolve } from "path";

const mod: TransformApp = async (input) => {
const backupFile = resolve(process.cwd(), "./backupingFile");
const backupingStream = new BackupingStream(backupFile);
const output = new PassThrough();

input.pipe(backupingStream).pipe(output);

return output;
};

export default mod;
Loading