Skip to content

Commit 17f75a6

Browse files
committed
Upload parquet logs to S3
1 parent 4986e27 commit 17f75a6

19 files changed

Lines changed: 1712 additions & 14 deletions

package.json

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -19,9 +19,13 @@
1919
"eslint": "9.39.2",
2020
"eslint-plugin-jsdoc": "62.9.0",
2121
"globals": "16.2.0",
22+
"hyparquet": "1.25.6",
2223
"typescript": "6.0.2",
2324
"vitest": "2.1.0"
2425
},
26+
"optionalDependencies": {
27+
"hyparquet-writer": "0.14.0"
28+
},
2529
"keywords": [
2630
"opentelemetry",
2731
"otlp",

src/collector.js

Lines changed: 25 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -48,34 +48,53 @@ const OTLP_NS_PER_MS = 1000000n
4848
const MIN_DATE_MS = -8640000000000000n
4949
const MAX_DATE_MS = 8640000000000000n
5050

51+
/**
52+
* @import { UploadOptions } from './upload/upload.d.ts'
53+
*/
54+
5155
class Collector {
52-
/** @param {{ port?: number, outputDir?: string }} [options] */
56+
/** @param {{ port?: number, outputDir?: string, upload?: UploadOptions }} [options] */
5357
constructor(options = {}) {
5458
this.port = options.port ?? 4318
5559
this.outputDir = options.outputDir || './otel-data'
60+
this.uploadOptions = options.upload
5661
/** @type {import('node:http').Server | null} */
5762
this.server = null
63+
/** @type {{ start: () => Promise<void>, stop: () => Promise<void> } | null} */
64+
this.uploader = null
5865
}
5966

60-
start() {
67+
async start() {
6168
ensureDir(this.outputDir)
6269

6370
const server = createServer(this.handleData.bind(this))
6471
this.server = server
6572

66-
return new Promise((resolve) => {
73+
await new Promise((resolve) => {
6774
server.listen(this.port, () => resolve(undefined))
6875
})
76+
77+
if (this.uploadOptions) {
78+
const { createUploader } = await import('./upload/index.js')
79+
this.uploader = createUploader({
80+
outputDir: this.outputDir,
81+
options: this.uploadOptions,
82+
})
83+
await this.uploader.start()
84+
}
6985
}
7086

71-
stop() {
72-
return new Promise((resolve, reject) => {
87+
async stop() {
88+
if (this.uploader) {
89+
await this.uploader.stop()
90+
this.uploader = null
91+
}
92+
await new Promise((resolve, reject) => {
7393
const { server } = this
7494
if (!server) {
7595
resolve(undefined)
7696
return
7797
}
78-
7998
server.close((err) => err ? reject(err) : resolve(undefined))
8099
})
81100
}

src/config.js

Lines changed: 74 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,7 @@
1+
/**
2+
* @import { Signal, UploadOptions } from './upload/upload.d.ts'
3+
*/
4+
15
/**
26
* Resolve collector options from CLI args and environment.
37
*
@@ -6,11 +10,13 @@
610
*
711
* @param {string[]} argv CLI arguments (without node/script name).
812
* @param {NodeJS.ProcessEnv} env Environment variables.
9-
* @returns {{ port?: number, outputDir?: string }} Options for Collector.
13+
* @returns {{ port?: number, outputDir?: string, upload?: UploadOptions }} Options for Collector.
1014
*/
1115
export function resolveOptions(argv, env) {
12-
/** @type {{ port?: number, outputDir?: string }} */
16+
/** @type {{ port?: number, outputDir?: string, upload?: UploadOptions }} */
1317
const options = {}
18+
/** @type {Partial<UploadOptions>} */
19+
const upload = {}
1420

1521
if (env.COLLECTIVUS_PORT) {
1622
const port = parseInt(env.COLLECTIVUS_PORT, 10)
@@ -20,22 +26,82 @@ export function resolveOptions(argv, env) {
2026
options.outputDir = env.COLLECTIVUS_OUTPUT_DIR
2127
}
2228

29+
if (env.COLLECTIVUS_UPLOAD_BUCKET) upload.bucket = env.COLLECTIVUS_UPLOAD_BUCKET
30+
if (env.COLLECTIVUS_UPLOAD_PREFIX) upload.prefix = env.COLLECTIVUS_UPLOAD_PREFIX
31+
if (env.COLLECTIVUS_UPLOAD_TIME) upload.time = env.COLLECTIVUS_UPLOAD_TIME
32+
if (env.COLLECTIVUS_UPLOAD_SIGNALS) upload.signals = parseSignals(env.COLLECTIVUS_UPLOAD_SIGNALS)
33+
if (env.COLLECTIVUS_UPLOAD_CATCHUP_DAYS) {
34+
const n = parseInt(env.COLLECTIVUS_UPLOAD_CATCHUP_DAYS, 10)
35+
if (!Number.isNaN(n)) upload.catchupDays = n
36+
}
37+
if (env.AWS_REGION) upload.region = env.AWS_REGION
38+
if (env.COLLECTIVUS_UPLOAD_ENDPOINT) upload.endpoint = env.COLLECTIVUS_UPLOAD_ENDPOINT
39+
2340
for (let i = 0; i < argv.length; i++) {
2441
const arg = argv[i]
25-
if (arg === '--port' && argv[i + 1]) {
26-
const port = parseInt(argv[i + 1], 10)
42+
function next() { return argv[++i] }
43+
if (arg === '--port' && argv[i + 1] !== undefined) {
44+
const port = parseInt(next(), 10)
2745
if (!Number.isNaN(port)) options.port = port
28-
i++
2946
} else if (arg.startsWith('--port=')) {
3047
const port = parseInt(arg.slice('--port='.length), 10)
3148
if (!Number.isNaN(port)) options.port = port
32-
} else if (arg === '--output' && argv[i + 1]) {
33-
options.outputDir = argv[i + 1]
34-
i++
49+
} else if (arg === '--output' && argv[i + 1] !== undefined) {
50+
options.outputDir = next()
3551
} else if (arg.startsWith('--output=')) {
3652
options.outputDir = arg.slice('--output='.length)
53+
} else if (arg === '--upload-bucket' && argv[i + 1] !== undefined) {
54+
upload.bucket = next()
55+
} else if (arg.startsWith('--upload-bucket=')) {
56+
upload.bucket = arg.slice('--upload-bucket='.length)
57+
} else if (arg === '--upload-prefix' && argv[i + 1] !== undefined) {
58+
upload.prefix = next()
59+
} else if (arg.startsWith('--upload-prefix=')) {
60+
upload.prefix = arg.slice('--upload-prefix='.length)
61+
} else if (arg === '--upload-time' && argv[i + 1] !== undefined) {
62+
upload.time = next()
63+
} else if (arg.startsWith('--upload-time=')) {
64+
upload.time = arg.slice('--upload-time='.length)
65+
} else if (arg === '--upload-signals' && argv[i + 1] !== undefined) {
66+
upload.signals = parseSignals(next())
67+
} else if (arg.startsWith('--upload-signals=')) {
68+
upload.signals = parseSignals(arg.slice('--upload-signals='.length))
69+
} else if (arg === '--upload-catchup-days' && argv[i + 1] !== undefined) {
70+
const n = parseInt(next(), 10)
71+
if (!Number.isNaN(n)) upload.catchupDays = n
72+
} else if (arg.startsWith('--upload-catchup-days=')) {
73+
const n = parseInt(arg.slice('--upload-catchup-days='.length), 10)
74+
if (!Number.isNaN(n)) upload.catchupDays = n
75+
} else if (arg === '--upload-region' && argv[i + 1] !== undefined) {
76+
upload.region = next()
77+
} else if (arg.startsWith('--upload-region=')) {
78+
upload.region = arg.slice('--upload-region='.length)
79+
} else if (arg === '--upload-endpoint' && argv[i + 1] !== undefined) {
80+
upload.endpoint = next()
81+
} else if (arg.startsWith('--upload-endpoint=')) {
82+
upload.endpoint = arg.slice('--upload-endpoint='.length)
3783
}
3884
}
3985

86+
if (upload.bucket) {
87+
options.upload = /** @type {UploadOptions} */ (upload)
88+
}
89+
4090
return options
4191
}
92+
93+
/**
94+
* @param {string} value comma-separated signal list
95+
* @returns {ReadonlyArray<Signal>}
96+
*/
97+
function parseSignals(value) {
98+
/** @type {Signal[]} */
99+
const out = []
100+
for (const part of value.split(',')) {
101+
const trimmed = part.trim()
102+
if (trimmed === 'logs' || trimmed === 'traces' || trimmed === 'metrics') {
103+
out.push(trimmed)
104+
}
105+
}
106+
return out
107+
}

src/index.js

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1 +1,2 @@
11
export { Collector } from './collector.js'
2+
export { createUploader } from './upload/index.js'

src/upload/connectors/memory.js

Lines changed: 27 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,27 @@
1+
/**
2+
* @import { StorageConnector } from '../upload.d.ts'
3+
*/
4+
5+
/**
6+
* In-memory StorageConnector for tests. Backed by a Map<string, Uint8Array>
7+
* exposed on `connector.store` so tests can introspect what was uploaded.
8+
*
9+
* @returns {StorageConnector & { store: Map<string, Uint8Array> }}
10+
*/
11+
export function memoryConnector() {
12+
/** @type {Map<string, Uint8Array>} */
13+
const store = new Map()
14+
15+
return {
16+
scheme: 'memory',
17+
store,
18+
async putObject(key, body) {
19+
store.set(key, body)
20+
},
21+
async headObject(key) {
22+
const body = store.get(key)
23+
if (!body) return null
24+
return { size: body.byteLength }
25+
},
26+
}
27+
}

0 commit comments

Comments
 (0)