_                     = require 'underscore'
Gearman               = require 'gearman-node'
{Readable, Transform} = require 'stream'

create_client_with_opts = (gearman_opts) -> new Gearman.Client gearman_opts
create_client = _.memoize create_client_with_opts, (gearman_opts) ->
  "#{gearman_opts.host}|#{gearman_opts.port}"

on_job_end = (job, job_name, cb) ->
  warning = ''
  job.on 'complete', -> cb()
  job.on 'warning', (handle, data) -> warning += data.toString()
  job.on 'fail', -> cb new Error "#{job_name} job failed with error: #{warning}"

class ReadableGearmanStream extends Readable
  constructor: (stream_opts, gearman, job_name, @payload) ->
    super _(stream_opts).extend objectMode: true
    job = gearman.submitJob job_name, payload
    job.on 'data', (handle, data) => @push data
    on_job_end job, job_name, (err) =>
      return @emit 'error', err if err
      @push null
  _read: =>

class TransformGearmanStream extends Transform
  constructor: (stream_opts, @gearman, @job_name) ->
    super stream_opts
    # Assume we're getting binary data that we can send to Gearman
    @_writableState.objectMode = false
    # Send out discrete chunks of binary data
    @_readableState.objectMode = true
    @payload = new Buffer('')
  _flush: (cb) ->
    job = @gearman.submitJob @job_name, @payload
    job.on 'data', (handle, data) => @push data
    on_job_end job, @job_name, cb
  _transform: (chunk, encoding, cb) ->
    @payload += chunk
    cb()

module.exports = class GearmanStream
  constructor: (stream_opts, gearman_opts, job_name, payload) ->
    client = create_client gearman_opts
    if payload
      return new ReadableGearmanStream stream_opts, client, job_name, payload
    else
      return new TransformGearmanStream stream_opts, client, job_name
