Skip to content

Commit 86576d5

Browse files
committed
refactor: generator writing to graph store
1 parent b02b00c commit 86576d5

8 files changed

Lines changed: 64 additions & 187 deletions

File tree

package-lock.json

Lines changed: 1 addition & 35 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

packages/graph-store/lib/SinkToWritable.js

Lines changed: 0 additions & 63 deletions
This file was deleted.
Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,35 @@
1+
import { PassThrough } from 'readable-stream'
2+
3+
/**
4+
* @param {(stream: import('readable-stream').PassThrough) => Promise<void>} openConnection
5+
* @param {object} context
6+
* @param {import('winston').Logger} context.logger
7+
*/
8+
export default function (openConnection, { logger }) {
9+
/**
10+
* @this {import('barnard59-core').Context}
11+
* @param {AsyncIterable<import('@rdfjs/types').Quad>} quads
12+
*/
13+
return async function * (quads) {
14+
const quadStream = new PassThrough({
15+
objectMode: true,
16+
})
17+
let request
18+
19+
for await (const quad of quads) {
20+
if (!request) {
21+
request = openConnection(quadStream)
22+
}
23+
24+
quadStream.write(quad)
25+
}
26+
27+
quadStream.end()
28+
29+
if (!request) {
30+
logger?.warn('No quads sent to the server. No request was made.')
31+
}
32+
33+
await request
34+
}
35+
}

packages/graph-store/package.json

Lines changed: 0 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,8 +24,6 @@
2424
"homepage": "https://github.com/zazuko/barnard59",
2525
"dependencies": {
2626
"duplex-to": "^1.0.0",
27-
"onetime": "^6.0.0",
28-
"promise-the-world": "^1.0.1",
2927
"readable-stream": "^3.6.0",
3028
"sparql-http-client": "^3.0.0",
3129
"barnard59-base": "^2.3.0",

packages/graph-store/post.js

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,16 @@
1+
import { Duplex } from 'node:stream'
12
import Client from 'sparql-http-client'
2-
import SinkToWritable from './lib/SinkToWritable.js'
3+
import duplexTo from 'duplex-to'
34
import { toTerm } from './lib/graph.js'
5+
import storeWritable from './lib/storeWritable.js'
46

57
/**
68
* @this {import('barnard59-core').Context}
79
* @param {Pick<import('sparql-http-client/StreamClient.js').Options<any>, 'user' | 'password'> & {
810
* endpoint: string,
911
* graph: string | import('clownface').GraphPointer<import('@rdfjs/types').NamedNode> | import('@rdfjs/types').NamedNode,
1012
* }} options
11-
* @returns {import('readable-stream').Writable}
13+
* @returns {import('node:stream').Writable}
1214
*/
1315
function post({ endpoint, graph, user, password }) {
1416
if (!graph) {
@@ -22,9 +24,9 @@ function post({ endpoint, graph, user, password }) {
2224
password,
2325
})
2426

25-
return new SinkToWritable(readable => client.store.post(readable, {
27+
return duplexTo.writable(Duplex.from(storeWritable(readable => client.store.post(readable, {
2628
graph: toTerm(this.env, graph),
27-
}))
29+
}), this)))
2830
}
2931

3032
export default post

packages/graph-store/put.js

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1,14 +1,16 @@
1+
import { Duplex } from 'node:stream'
12
import Client from 'sparql-http-client'
2-
import SinkToWritable from './lib/SinkToWritable.js'
3+
import duplexTo from 'duplex-to'
34
import { toTerm } from './lib/graph.js'
5+
import storeWritable from './lib/storeWritable.js'
46

57
/**
68
* @this {import('barnard59-core').Context}
79
* @param {Pick<import('sparql-http-client/StreamClient.js').Options<any>, 'user' | 'password'> & {
810
* endpoint: string,
911
* graph: string | import('clownface').GraphPointer<import('@rdfjs/types').NamedNode> | import('@rdfjs/types').NamedNode,
1012
* }} options
11-
* @returns {import('readable-stream').Writable}
13+
* @returns {import('node:stream').Writable}
1214
*/
1315
function put({ endpoint, graph, user, password }) {
1416
if (!graph) {
@@ -22,9 +24,9 @@ function put({ endpoint, graph, user, password }) {
2224
password,
2325
})
2426

25-
return new SinkToWritable(readable => client.store.put(readable, {
27+
return duplexTo.writable(Duplex.from(storeWritable(readable => client.store.put(readable, {
2628
graph: toTerm(this.env, graph),
27-
}))
29+
}), this)))
2830
}
2931

3032
export default put

packages/graph-store/test/post.test.js

Lines changed: 8 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -1,31 +1,17 @@
11
import { strictEqual } from 'node:assert'
22
import { promisify } from 'node:util'
3+
import { Readable, finished } from 'node:stream'
34
import rdf from '@zazuko/env'
45
import quadToNTriples from '@rdfjs/to-ntriples'
56
import withServer from 'express-as-promise/withServer.js'
67
import getStream from 'get-stream'
7-
import { isReadableStream as isReadable, isWritableStream as isWritable } from 'is-stream'
8-
import { finished } from 'readable-stream'
98
import postUnbound from '../post.js'
109

1110
const post = postUnbound.bind({ env: rdf })
1211

1312
const ns = rdf.namespace('http://example.org/')
1413

1514
describe('post', () => {
16-
it('should return a writable stream', async () => {
17-
await withServer(async server => {
18-
const baseUrl = await server.listen()
19-
20-
const stream = post({ endpoint: baseUrl, graph: 'default' })
21-
22-
strictEqual(isReadable(stream), false)
23-
strictEqual(isWritable(stream), true)
24-
25-
stream.end()
26-
})
27-
})
28-
2915
it('should send a POST request', async () => {
3016
await withServer(async server => {
3117
let called = false
@@ -40,10 +26,7 @@ describe('post', () => {
4026
const baseUrl = await server.listen()
4127
const stream = post({ endpoint: baseUrl, graph: ns.graph1 })
4228

43-
stream.write(quad)
44-
stream.end()
45-
46-
await promisify(finished)(stream)
29+
await getStream(Readable.from([quad]).pipe(stream))
4730

4831
strictEqual(called, true)
4932
})
@@ -62,9 +45,7 @@ describe('post', () => {
6245
const baseUrl = await server.listen()
6346
const stream = post({ endpoint: baseUrl, graph: ns.graph1 })
6447

65-
stream.end()
66-
67-
await promisify(finished)(stream)
48+
await getStream(Readable.from([]).pipe(stream))
6849

6950
strictEqual(called, false)
7051
})
@@ -84,8 +65,7 @@ describe('post', () => {
8465
const baseUrl = await server.listen()
8566
const stream = post({ endpoint: baseUrl, graph: ns.graph1 })
8667

87-
stream.write(quad)
88-
stream.end()
68+
await getStream(Readable.from([quad]).pipe(stream))
8969

9070
await promisify(finished)(stream)
9171

@@ -115,13 +95,7 @@ describe('post', () => {
11595
const baseUrl = await server.listen()
11696
const stream = post({ endpoint: baseUrl, graph: ns.graph1 })
11797

118-
quads.forEach(quad => {
119-
stream.write(quad)
120-
})
121-
122-
stream.end()
123-
124-
await promisify(finished)(stream)
98+
await getStream(Readable.from(quads).pipe(stream))
12599

126100
strictEqual(content[quads[0].graph.value], expected)
127101
})
@@ -151,13 +125,7 @@ describe('post', () => {
151125
const baseUrl = await server.listen()
152126
const stream = post({ endpoint: baseUrl, graph: 'default' })
153127

154-
quads.forEach(quad => {
155-
stream.write(quad)
156-
})
157-
158-
stream.end()
159-
160-
await promisify(finished)(stream)
128+
await getStream(Readable.from(quads).pipe(stream))
161129

162130
strictEqual(typeof graph, 'undefined')
163131
strictEqual(content, expected)
@@ -178,8 +146,7 @@ describe('post', () => {
178146
const baseUrl = await server.listen()
179147
const stream = post({ endpoint: baseUrl, user: 'testuser', password: 'testpassword', graph: ns.graph1 })
180148

181-
stream.write(quad)
182-
stream.end()
149+
await getStream(Readable.from([quad]).pipe(stream))
183150

184151
await promisify(finished)(stream)
185152

@@ -198,13 +165,10 @@ describe('post', () => {
198165
const baseUrl = await server.listen()
199166
const stream = post({ endpoint: baseUrl, graph: ns.graph1 })
200167

201-
stream.write(quad)
202-
stream.end()
203-
204168
let error = null
205169

206170
try {
207-
await promisify(finished)(stream)
171+
await getStream(Readable.from([quad]).pipe(stream))
208172
} catch (err) {
209173
error = err
210174
}

0 commit comments

Comments
 (0)