Skip to content

Commit d9eb06c

Browse files
committed
add put and del streams
1 parent d5ac6b6 commit d9eb06c

3 files changed

Lines changed: 33 additions & 11 deletions

File tree

index.js

Lines changed: 24 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
const PassThrough = require('readable-stream').PassThrough
2+
const Transform = require('readable-stream').Transform
23
const pump = require('pump')
34

45
const utils = require('./lib/utils')
@@ -34,13 +35,34 @@ function doAction (action) {
3435
}
3536
}
3637

38+
function doActionStream (action) {
39+
return function () {
40+
const transform = new Transform({
41+
objectMode: true,
42+
transform (triples, encoding, done) {
43+
if (!triples) return done()
44+
let entries = (!triples.reduce) ? [triples] : triples
45+
entries = entries.reduce((prev, triple) => {
46+
return prev.concat(utils.generateBatch(triple, action))
47+
}, [])
48+
this.push(entries.reverse())
49+
done()
50+
}
51+
})
52+
const writeStream = this.db.createWriteStream()
53+
transform.pipe(writeStream)
54+
return transform
55+
}
56+
}
57+
3758
// this is not implemented in hyperdb yet
3859
// for now we just put a null value in the db
39-
Graph.prototype.del = doAction('del')
4060

4161
Graph.prototype.put = doAction('put')
62+
Graph.prototype.putStream = doActionStream('put')
4263

43-
Graph.prototype.putStream = function (triple) { }
64+
Graph.prototype.del = doAction('del')
65+
Graph.prototype.delStream = doActionStream('del')
4466

4567
Graph.prototype.searchStream = function (query, options) {
4668
const result = new PassThrough({ objectMode: true })

readme.md

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@ Inserts **Hexastore** formated entries for triple into the graph database.
4444

4545
#### `var stream = db.putStream(triple)`
4646

47-
Returns a writable stream. **This is not yet implemented!**
47+
Returns a writable stream.
4848

4949
#### `db.get(triple, [callback])`
5050

@@ -56,7 +56,7 @@ Remove triples indices from the graph database.
5656

5757
#### `var stream = db.delStream(triple)`
5858

59-
Returns a writable stream for removing entries. **This is not yet implemented!**
59+
Returns a writable stream for removing entries.
6060

6161
#### `var stream = db.getStream(triple)`
6262

test/triple-store.spec.js

Lines changed: 7 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -514,25 +514,25 @@ describe('a basic triple store', function () {
514514
}
515515
})
516516

517-
xit('should put triples using a stream', function (done) {
517+
it('should put triples using a stream', function (done) {
518518
var t1 = { subject: 'a', predicate: 'b', object: 'c' }
519519
var t2 = { subject: 'a', predicate: 'b', object: 'd' }
520520
var stream = db.putStream()
521-
stream.on('close', done)
521+
stream.on('end', done)
522522

523523
stream.write(t1)
524524
stream.end(t2)
525525
})
526526

527-
xit('should store the triples written using a stream', function (done) {
527+
it('should store the triples written using a stream', function (done) {
528528
var t1 = { subject: 'a', predicate: 'b', object: 'c' }
529529
var t2 = { subject: 'a', predicate: 'b', object: 'd' }
530530
var stream = db.putStream()
531531

532532
stream.write(t1)
533533
stream.end(t2)
534534

535-
stream.on('close', function () {
535+
stream.on('end', function () {
536536
var triples = [t1, t2]
537537
var readStream = db.getStream({ predicate: 'b' })
538538

@@ -544,20 +544,20 @@ describe('a basic triple store', function () {
544544
})
545545
})
546546

547-
xit('should del the triples using a stream', function (done) {
547+
it('should del the triples using a stream', function (done) {
548548
var t1 = { subject: 'a', predicate: 'b', object: 'c' }
549549
var t2 = { subject: 'a', predicate: 'b', object: 'd' }
550550
var stream = db.putStream()
551551

552552
stream.write(t1)
553553
stream.end(t2)
554554

555-
stream.on('close', function () {
555+
stream.on('end', function () {
556556
var delStream = db.delStream()
557557
delStream.write(t1)
558558
delStream.end(t2)
559559

560-
delStream.on('close', function () {
560+
delStream.on('end', function () {
561561
var readStream = db.getStream({ predicate: 'b' })
562562

563563
var results = []

0 commit comments

Comments
 (0)