subscribing to event causes 'ERR_STREAM_WRITE_AFTER_END'

Viewed 81

Tried to make example-event.ts work as HTTP Server. The example is a simple counter. You can subscribe to the count as an event. It works until a client asks a second time for the value (longpoll).

EventSource ready
Emitted change  1
[binding-http] HttpServer on port 8080 received 'GET /eventsource' from [::ffff:192.168.0.5]:58268
[binding-http] HttpServer on port 8080 replied with '200' to [::ffff:192.168.0.5]:58268
[binding-http] HttpServer on port 8080 received 'GET /eventsource/events/onchange' from [::ffff:192.168.0.5]:58268
[core/exposed-thing] ExposedThing 'EventSource' subscribes to event 'onchange'
[core/content-serdes] ContentSerdes serializing to application/json
Emitted change  2
[binding-http] HttpServer on port 8080 replied with '200' to [::ffff:192.168.0.5]:58268
[binding-http] HttpServer on port 8080 closed Event connection
[core/exposed-thing] ExposedThing 'EventSource' unsubscribes from event 'onchange'
[binding-http] HttpServer on port 8080 received 'GET /eventsource/events/onchange' from [::ffff:192.168.0.5]:58268
[core/exposed-thing] ExposedThing 'EventSource' subscribes to event 'onchange'
[core/content-serdes] ContentSerdes serializing to application/json
[core/content-serdes] ContentSerdes serializing to application/json
Emitted change  3
events.js:377
      throw er; // Unhandled 'error' event
      ^

Error [ERR_STREAM_WRITE_AFTER_END]: write after end
    at new NodeError (internal/errors.js:322:7)
    at writeAfterEnd (_http_outgoing.js:694:15)
    at ServerResponse.end (_http_outgoing.js:815:7)
    at SafeSubscriber._next (C:\xxx\node_modules\@node-wot\binding-http\dist\http-server.js:721:45)
    at SafeSubscriber.__tryOrUnsub (C:\xxx\node_modules\rxjs\Subscriber.js:242:16)
    at SafeSubscriber.next (C:\xxx\node_modules\rxjs\Subscriber.js:189:22)
    at Subscriber._next (C:\xxx\node_modules\rxjs\Subscriber.js:129:26)
    at Subscriber.next (C:\xxx\node_modules\rxjs\Subscriber.js:93:18)
    at Subject.next (C:\xxx\node_modules\rxjs\Subject.js:55:25)
    at Object.ExposedThing.emitEvent (C:\xxx\node_modules\@node-wot\core\dist\exposed-thing.js:53:50)
Emitted 'error' event on ServerResponse instance at:
    at writeAfterEndNT (_http_outgoing.js:753:7)
    at processTicksAndRejections (internal/process/task_queues.js:83:21) {
  code: 'ERR_STREAM_WRITE_AFTER_END'
}

I am a little confused by all the subscribing and unsubscribing, but this is maybe because of longpoll is used.

my code:

server side:

Servient = require('@node-wot/core').Servient;
HttpServer = require('@node-wot/binding-http').HttpServer;
Helpers = require('@node-wot/core').Helpers;

// create Servient add HTTP binding with port configuration
let servient = new Servient();
servient.addServer(new HttpServer({}));

// internal state, not exposed as Property
let counter = 0;

servient.start().then((WoT) => {
  WoT.produce({
    title: 'EventSource',
    events: {
      onchange: {
        data: { type: 'integer' },
      },
    },
  })
    .then((thing) => {
      console.log('Produced ' + thing.getThingDescription().title);

      thing.expose().then(() => {
        console.info(thing.getThingDescription().title + ' ready');

        setInterval(() => {
          ++counter;
          thing.emitEvent('onchange', counter);
          console.info('Emitted change ', counter);
        }, 5000);
      });
    })
    .catch((e) => {
      console.log(e);
    });
});

cient side:

const servient = new Wot.Core.Servient();
servient.addClientFactory(new Wot.Http.HttpClientFactory());
const helpers = new Wot.Core.Helpers(servient);

const addr ='http://192.168.0.5:8080/eventsource';
getTd(addr);

function getTd(addr) {
  servient.start().then((thingFactory) => {
    helpers
      .fetch(addr)
      .then((td) => {
        thingFactory.consume(td).then((thing) => {
          showEvents(thing);
        });
      })
      .catch((error) => {
        window.alert('Could not fetch TD.\n' + error);
      });
  });
}

function showEvents(thing) {
  let td = thing.getThingDescription();
  for (let evnt in td.events) {
    if (td.events.hasOwnProperty(evnt)) {

      document.getElementById("events").innerHTML = "waiting...";

      thing
        .subscribeEvent(evnt, (res) => {
          document.getElementById("events").innerHTML = res;
        })
        .catch((err) => window.alert('error: ' + err));
    }
  }
}

The problem also occurs if I'm using example-event-client.ts or just send GET requests to "http://192.168.0.5:8080/eventsource/events/onchange" using the browser.

What do I have to do to make the example work?

0 Answers
Related