_           = require 'underscore'
assert      = require 'assert'
JSONStream  = require 'JSONStream'
Understream = require 'understream'
{Worker}    = require 'gearman-node'
GearmanStream = require '../index'
Understream.mixin JSONStream.stringify, 'jsonstring', true

Understream.mixin GearmanStream, 'gearman'
gearman_opts = host: process.env.GEARMAN_HOST, port: process.env.GEARMAN_PORT

specs = [
  name: 'completes on a stream that sends no data'
  worker: (payload, worker) -> worker.done()
  assertions: (err, data) ->
    assert.ifError err
    assert.equal data.length, 0
,
  name: 'errors on a stream that fails'
  worker: (payload, worker) -> worker.done new Error "worker failed"
  assertions: (err, data) -> assert.equal err?.message, 'worker_name job failed with error: Error: worker failed'
,
  name: 'streams out data'
  worker: (payload, worker) ->
    worker.data JSON.stringify i for i in JSON.parse payload
    worker.done()
  payload: [0...10]
  assertions: (err, data) ->
    assert.ifError err
    assert.deepEqual data, [0...10]
]

describe 'gearman stream', ->
  describe 'with a payload', ->
    _(specs).each (spec) ->
      it spec.name, (done) ->
        worker = new Worker 'worker_name', spec.worker, gearman_opts
        new Understream().gearman(gearman_opts, 'worker_name', JSON.stringify(spec.payload or {})).map(JSON.parse).run (err, data) ->
          spec.assertions err, data
          worker.disconnect()
          done()
  describe 'with a streaming payload', ->
    _(specs).each (spec) ->
      it spec.name, (done) ->
        worker = new Worker 'worker_name', spec.worker, gearman_opts
        new Understream(spec.payload or [])
          .jsonstring('[', ',', ']')
          .gearman(gearman_opts, 'worker_name')
          .map(JSON.parse)
          .run (err, data) ->
            spec.assertions err, data
            worker.disconnect()
            done()
