-
Notifications
You must be signed in to change notification settings - Fork 90
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
Showing
40 changed files
with
1,204 additions
and
451 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,8 @@ | ||
{ | ||
"extends": [ | ||
"standard" | ||
], | ||
"env": { | ||
"node": true | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
File renamed without changes.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,28 @@ | ||
/** | ||
* Created by tushar.mathur on 18/06/16. | ||
*/ | ||
|
||
'use strict' | ||
|
||
import {Observable as O} from 'rx' | ||
import {mux} from 'muxer' | ||
import R from 'ramda' | ||
|
||
export const ev = R.curry(($, event) => $.filter(R.whereEq({event})).pluck('message')) | ||
|
||
export const RequestParams = R.curry((request, params) => { | ||
return O.create((observer) => request(params) | ||
.on('data', (message) => observer.onNext({event: 'data', message})) | ||
.on('response', (message) => observer.onNext({event: 'response', message})) | ||
.on('complete', () => observer.onCompleted()) | ||
.on('error', (error) => observer.onError(error)) | ||
) | ||
}) | ||
|
||
export const Request = R.curry((request, params) => { | ||
const Response$ = ev(RequestParams(request, params)) | ||
return mux({ | ||
response$: Response$('response'), | ||
data$: Response$('data') | ||
}) | ||
}) |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,23 @@ | ||
/** | ||
* Created by tushar.mathur on 10/06/16. | ||
*/ | ||
|
||
'use strict' | ||
|
||
import R from 'ramda' | ||
import {Observable as O} from 'rx' | ||
export const map = R.curry((func, $) => $.map(func)) | ||
export const flatMap = R.curry((func, $) => $.flatMap(func)) | ||
export const withLatestFrom = R.curry((list, $) => $.withLatestFrom(...list)) | ||
export const zip = R.curry((list, $) => $.zip(...list)) | ||
export const zipWith = R.curry((func, list, $) => $.zip(...list, func)) | ||
export const filter = R.curry((func, $) => $.filter(func)) | ||
export const distinctUntilChanged = $ => $.distinctUntilChanged() | ||
export const pluck = R.curry((path, $) => $.pluck(path)) | ||
export const scan = R.curry((func, $) => $.scan(func)) | ||
export const scanWith = R.curry((func, m, $) => $.scan(func, m)) | ||
export const shareReplay = R.curry((count, $) => $.shareReplay(count)) | ||
export const repeat = R.curry((value, count) => O.repeat(value, count)) | ||
export const trace = R.curry((msg, $) => $.tap(x => console.log(msg, x))) | ||
export const tap = R.curry((func, $) => $.tap(func)) | ||
export const share = ($) => $.share() |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,33 +1,40 @@ | ||
import request from 'request' | ||
import Rx from 'rx' | ||
import fs from 'graceful-fs' | ||
import _ from 'lodash' | ||
import * as u from './Utils' | ||
'use strict' | ||
|
||
export const requestBody = (params) => Rx.Observable.create((observer) => request(params) | ||
.on('data', (message) => observer.onNext({event: 'data', message})) | ||
.on('response', (message) => observer.onNext({event: 'response', message})) | ||
.on('complete', (message) => observer.onCompleted({event: 'completed', message})) | ||
.on('error', (error) => observer.onError(error)) | ||
) | ||
import {Observable as O} from 'rx' | ||
import {demux} from 'muxer' | ||
import R from 'ramda' | ||
import {Request} from './Request' | ||
|
||
export const requestHead = (params) => requestBody(params) | ||
.first() | ||
.pluck('message') | ||
.tap((x) => x.destroy()) | ||
export const fromCB = R.compose(R.apply, O.fromNodeCallback) | ||
|
||
export const requestContentLength = (params) => requestHead(params) | ||
.pluck('headers', 'content-length') | ||
.map((x) => parseInt(x, 10)) | ||
export const FILE = R.curry((fs) => { | ||
return [{ | ||
// New Methods | ||
open: signal$ => signal$.flatMap(fromCB(fs.open)).shareReplay(1), | ||
fstat: signal$ => signal$.flatMap(fromCB(fs.fstat)).shareReplay(1), | ||
read: signal$ => signal$.flatMap(fromCB(fs.read)).shareReplay(1), | ||
write: signal$ => signal$.flatMap(fromCB(fs.write)).shareReplay(1), | ||
close: signal$ => signal$.flatMap(fromCB(fs.close)).shareReplay(1), | ||
truncate: signal$ => signal$.flatMap(fromCB(fs.truncate)).shareReplay(1), | ||
rename: signal$ => signal$.flatMap(fromCB(fs.rename)).shareReplay(1) | ||
}] | ||
}) | ||
|
||
export const fsOpen = Rx.Observable.fromNodeCallback(fs.open) | ||
export const fsWrite = Rx.Observable.fromNodeCallback(fs.write) | ||
export const fsTruncate = Rx.Observable.fromNodeCallback(fs.truncate) | ||
export const fsRename = Rx.Observable.fromNodeCallback(fs.rename) | ||
export const fsStat = Rx.Observable.fromNodeCallback(fs.fstat) | ||
export const fsRead = Rx.Observable.fromNodeCallback(fs.read) | ||
export const fsReadBuffer = (x) => fsRead(x.fd, x.buffer, 0, x.buffer.length, x.offset) | ||
export const fsWriteBuffer = (x) => fsWrite(x.fd, x.buffer, 0, x.buffer.length, x.offset) | ||
export const fsWriteJSON = (x) => fsWriteBuffer(_.assign({}, x, {buffer: u.toBuffer(x.json)})) | ||
export const fsReadJSON = (x) => fsReadBuffer(x).map((x) => JSON.parse(x[1].toString())) | ||
export const buffer = (size) => Rx.Observable.just(u.createEmptyBuffer(size)) | ||
export const HTTP = R.curry((_request) => { | ||
const request = Request(_request) | ||
const requestHead = (params) => { | ||
const [{response$}] = demux(request(params), 'response$') | ||
return response$.first().tap(x => x.destroy()).share() | ||
} | ||
|
||
const select = R.curry((event, request$) => request$.filter(x => x.event === event).pluck('message')) | ||
const executor = (signal$) => { | ||
const [{destroy$}] = demux(signal$, 'destroy$') | ||
destroy$.subscribe(request => request.destroy()) | ||
} | ||
return [{ | ||
requestHead, | ||
select, | ||
request | ||
}, executor] | ||
}) |
Oops, something went wrong.