mirror of
https://github.com/node-red/node-red.git
synced 2023-10-10 13:36:53 +02:00
fix multiple input message processing of file node
This commit is contained in:
parent
08fccc4e77
commit
1a226c4dc6
@ -29,8 +29,9 @@ module.exports = function(RED) {
|
|||||||
var node = this;
|
var node = this;
|
||||||
node.wstream = null;
|
node.wstream = null;
|
||||||
node.data = [];
|
node.data = [];
|
||||||
|
node.msgQueue = [];
|
||||||
|
|
||||||
this.on("input",function(msg) {
|
function processMsg(msg, done) {
|
||||||
var filename = node.filename || msg.filename || "";
|
var filename = node.filename || msg.filename || "";
|
||||||
if ((!node.filename) && (!node.tout)) {
|
if ((!node.filename) && (!node.tout)) {
|
||||||
node.tout = setTimeout(function() {
|
node.tout = setTimeout(function() {
|
||||||
@ -41,6 +42,7 @@ module.exports = function(RED) {
|
|||||||
}
|
}
|
||||||
if (filename === "") {
|
if (filename === "") {
|
||||||
node.warn(RED._("file.errors.nofilename"));
|
node.warn(RED._("file.errors.nofilename"));
|
||||||
|
done();
|
||||||
} else if (node.overwriteFile === "delete") {
|
} else if (node.overwriteFile === "delete") {
|
||||||
fs.unlink(filename, function (err) {
|
fs.unlink(filename, function (err) {
|
||||||
if (err) {
|
if (err) {
|
||||||
@ -51,6 +53,7 @@ module.exports = function(RED) {
|
|||||||
}
|
}
|
||||||
node.send(msg);
|
node.send(msg);
|
||||||
}
|
}
|
||||||
|
done();
|
||||||
});
|
});
|
||||||
} else if (msg.hasOwnProperty("payload") && (typeof msg.payload !== "undefined")) {
|
} else if (msg.hasOwnProperty("payload") && (typeof msg.payload !== "undefined")) {
|
||||||
var dir = path.dirname(filename);
|
var dir = path.dirname(filename);
|
||||||
@ -59,6 +62,7 @@ module.exports = function(RED) {
|
|||||||
fs.ensureDirSync(dir);
|
fs.ensureDirSync(dir);
|
||||||
} catch(err) {
|
} catch(err) {
|
||||||
node.error(RED._("file.errors.createfail",{error:err.toString()}),msg);
|
node.error(RED._("file.errors.createfail",{error:err.toString()}),msg);
|
||||||
|
done();
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@ -82,6 +86,7 @@ module.exports = function(RED) {
|
|||||||
node.wstream.on("open", function() {
|
node.wstream.on("open", function() {
|
||||||
node.wstream.end(packet.data, function() {
|
node.wstream.end(packet.data, function() {
|
||||||
node.send(packet.msg);
|
node.send(packet.msg);
|
||||||
|
done();
|
||||||
});
|
});
|
||||||
})
|
})
|
||||||
})(node.data.shift());
|
})(node.data.shift());
|
||||||
@ -130,6 +135,7 @@ module.exports = function(RED) {
|
|||||||
var packet = node.data.shift()
|
var packet = node.data.shift()
|
||||||
node.wstream.write(packet.data, function() {
|
node.wstream.write(packet.data, function() {
|
||||||
node.send(packet.msg);
|
node.send(packet.msg);
|
||||||
|
done();
|
||||||
});
|
});
|
||||||
|
|
||||||
} else {
|
} else {
|
||||||
@ -137,14 +143,34 @@ module.exports = function(RED) {
|
|||||||
var packet = node.data.shift()
|
var packet = node.data.shift()
|
||||||
node.wstream.end(packet.data, function() {
|
node.wstream.end(packet.data, function() {
|
||||||
node.send(packet.msg);
|
node.send(packet.msg);
|
||||||
});
|
|
||||||
delete node.wstream;
|
delete node.wstream;
|
||||||
delete node.wstreamIno;
|
delete node.wstreamIno;
|
||||||
|
done();
|
||||||
|
});
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
this.on("input", function(msg) {
|
||||||
|
var msgQueue = node.msgQueue;
|
||||||
|
function processQ() {
|
||||||
|
var msg = msgQueue[0];
|
||||||
|
processMsg(msg, function() {
|
||||||
|
msgQueue.shift();
|
||||||
|
if (msgQueue.length > 0) {
|
||||||
|
processQ();
|
||||||
|
}
|
||||||
});
|
});
|
||||||
|
}
|
||||||
|
if (msgQueue.push(msg) > 1) {
|
||||||
|
// pending write exists
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
processQ();
|
||||||
|
});
|
||||||
|
|
||||||
this.on('close', function() {
|
this.on('close', function() {
|
||||||
if (node.wstream) { node.wstream.end(); }
|
if (node.wstream) { node.wstream.end(); }
|
||||||
if (node.tout) { clearTimeout(node.tout); }
|
if (node.tout) { clearTimeout(node.tout); }
|
||||||
|
@ -540,6 +540,49 @@ describe('file Nodes', function() {
|
|||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it('should write to multiple files', function(done) {
|
||||||
|
var flow = [{id:"fileNode1", type:"file", name: "fileNode", "appendNewline":true, "overwriteFile":true, "createDir":true, wires: [["helperNode1"]]},
|
||||||
|
{id:"helperNode1", type:"helper"}];
|
||||||
|
var tmp_path = path.join(resourcesDir, "tmp");
|
||||||
|
var len = 1024*1024*10;
|
||||||
|
var file_count = 5;
|
||||||
|
helper.load(fileNode, flow, function() {
|
||||||
|
var n1 = helper.getNode("fileNode1");
|
||||||
|
var n2 = helper.getNode("helperNode1");
|
||||||
|
var count = 0;
|
||||||
|
n2.on("input", function(msg) {
|
||||||
|
try {
|
||||||
|
count++;
|
||||||
|
if (count == file_count) {
|
||||||
|
for(var i = 0; i < file_count; i++) {
|
||||||
|
var name = path.join(tmp_path, String(i));
|
||||||
|
var f = fs.readFileSync(name);
|
||||||
|
f.should.have.length(len);
|
||||||
|
f[0].should.have.equal(i);
|
||||||
|
}
|
||||||
|
fs.removeSync(tmp_path);
|
||||||
|
done();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
catch (e) {
|
||||||
|
try {
|
||||||
|
fs.removeSync(tmp_path);
|
||||||
|
}
|
||||||
|
catch (e1) {
|
||||||
|
}
|
||||||
|
done(e);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
for(var i = 0; i < file_count; i++) {
|
||||||
|
var data = new Buffer(len);
|
||||||
|
data.fill(i);
|
||||||
|
var name = path.join(tmp_path, String(i));
|
||||||
|
var msg = {payload:data, filename:name};
|
||||||
|
n1.receive(msg);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|
||||||
|
Loading…
Reference in New Issue
Block a user