2020-08-05 18:38:55 +02:00
|
|
|
/*
|
|
|
|
Copyright 2020 Bruno Windels <bruno@windels.cloud>
|
|
|
|
|
|
|
|
Licensed under the Apache License, Version 2.0 (the "License");
|
|
|
|
you may not use this file except in compliance with the License.
|
|
|
|
You may obtain a copy of the License at
|
|
|
|
|
|
|
|
http://www.apache.org/licenses/LICENSE-2.0
|
|
|
|
|
|
|
|
Unless required by applicable law or agreed to in writing, software
|
|
|
|
distributed under the License is distributed on an "AS IS" BASIS,
|
|
|
|
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
|
|
|
See the License for the specific language governing permissions and
|
|
|
|
limitations under the License.
|
|
|
|
*/
|
|
|
|
|
2020-04-20 21:26:39 +02:00
|
|
|
import {SortedArray} from "../../../observable/list/SortedArray.js";
|
2020-04-19 19:05:12 +02:00
|
|
|
import {ConnectionError} from "../../error.js";
|
2020-04-20 21:26:39 +02:00
|
|
|
import {PendingEvent} from "./PendingEvent.js";
|
2020-09-03 15:36:17 +02:00
|
|
|
import {makeTxnId} from "../../common.js";
|
2019-07-01 10:00:29 +02:00
|
|
|
|
2020-04-20 21:26:39 +02:00
|
|
|
export class SendQueue {
|
2020-09-22 13:43:18 +02:00
|
|
|
constructor({roomId, storage, hsApi, pendingEvents}) {
|
2019-07-26 22:33:33 +02:00
|
|
|
pendingEvents = pendingEvents || [];
|
2019-07-26 22:03:57 +02:00
|
|
|
this._roomId = roomId;
|
|
|
|
this._storage = storage;
|
2020-09-22 13:43:18 +02:00
|
|
|
this._hsApi = hsApi;
|
2019-07-26 22:03:57 +02:00
|
|
|
this._pendingEvents = new SortedArray((a, b) => a.queueIndex - b.queueIndex);
|
2019-07-27 10:40:56 +02:00
|
|
|
if (pendingEvents.length) {
|
|
|
|
console.info(`SendQueue for room ${roomId} has ${pendingEvents.length} pending events`, pendingEvents);
|
|
|
|
}
|
2019-07-29 09:54:34 +02:00
|
|
|
this._pendingEvents.setManyUnsorted(pendingEvents.map(data => new PendingEvent(data)));
|
2019-07-26 22:03:57 +02:00
|
|
|
this._isSending = false;
|
2019-06-28 00:52:54 +02:00
|
|
|
this._offline = false;
|
2019-07-26 22:03:57 +02:00
|
|
|
this._amountSent = 0;
|
2020-09-03 15:36:48 +02:00
|
|
|
this._roomEncryption = null;
|
|
|
|
}
|
|
|
|
|
|
|
|
enableEncryption(roomEncryption) {
|
|
|
|
this._roomEncryption = roomEncryption;
|
2019-06-28 00:52:54 +02:00
|
|
|
}
|
|
|
|
|
2019-07-01 10:00:29 +02:00
|
|
|
async _sendLoop() {
|
2019-07-26 22:03:57 +02:00
|
|
|
this._isSending = true;
|
|
|
|
try {
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("start sending", this._amountSent, "<", this._pendingEvents.length);
|
2019-07-26 22:03:57 +02:00
|
|
|
while (this._amountSent < this._pendingEvents.length) {
|
|
|
|
const pendingEvent = this._pendingEvents.get(this._amountSent);
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("trying to send", pendingEvent.content.body);
|
2019-07-26 22:03:57 +02:00
|
|
|
if (pendingEvent.remoteId) {
|
2020-09-29 11:32:49 +02:00
|
|
|
this._amountSent += 1;
|
2019-07-26 22:03:57 +02:00
|
|
|
continue;
|
2019-07-01 10:00:29 +02:00
|
|
|
}
|
2020-11-13 17:19:19 +01:00
|
|
|
if (pendingEvent.attachments) {
|
2020-11-11 13:06:03 +01:00
|
|
|
try {
|
2020-11-13 17:19:19 +01:00
|
|
|
await this._uploadAttachments(pendingEvent);
|
2020-11-11 13:06:03 +01:00
|
|
|
} catch (err) {
|
2020-11-13 17:19:19 +01:00
|
|
|
console.log("upload failed, skip sending message", err, pendingEvent);
|
2020-11-11 13:06:03 +01:00
|
|
|
this._amountSent += 1;
|
|
|
|
continue;
|
|
|
|
}
|
2020-11-13 17:19:19 +01:00
|
|
|
console.log("attachments upload, content is now", pendingEvent.content);
|
2020-11-11 11:51:39 +01:00
|
|
|
}
|
2020-09-03 15:36:48 +02:00
|
|
|
if (pendingEvent.needsEncryption) {
|
2020-09-22 13:43:18 +02:00
|
|
|
const {type, content} = await this._roomEncryption.encrypt(
|
|
|
|
pendingEvent.eventType, pendingEvent.content, this._hsApi);
|
2020-09-03 15:36:48 +02:00
|
|
|
pendingEvent.setEncrypted(type, content);
|
|
|
|
await this._tryUpdateEvent(pendingEvent);
|
|
|
|
}
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("really sending now");
|
2020-09-22 13:43:18 +02:00
|
|
|
const response = await this._hsApi.send(
|
2019-07-26 22:03:57 +02:00
|
|
|
pendingEvent.roomId,
|
|
|
|
pendingEvent.eventType,
|
|
|
|
pendingEvent.txnId,
|
|
|
|
pendingEvent.content
|
need to return the response here, not the request wrapper
we were reading back a remote id of undefined because of this,
so when for some reason we never receive the message down from sync,
the pending message keeps sending on every load. The server ignores
the send though, because the transaction id is already used, and it returns
the remote id of the event that was already sent the previous time, but
as we were not storing this remote id, we'd just try again and again.
not receiving it through sync could have happened when we were sending a bunch of events
and then receiving (this is how we encountered this bug, while trying to repro another) the
response, but not yet the sync for the message that got wedged. Then we typed stuff on another client
so we would get a limited response for that room, and boom, we would not get the remote echo of the
event that was already sent (but because of this bug we didn't store the remote id) but no echo received yet (when we remove the pending event),
so it gets wedged!
2020-03-17 00:11:43 +01:00
|
|
|
).response();
|
2019-07-26 22:03:57 +02:00
|
|
|
pendingEvent.remoteId = response.event_id;
|
2019-07-27 10:40:56 +02:00
|
|
|
//
|
|
|
|
console.log("writing remoteId now");
|
2019-07-26 22:03:57 +02:00
|
|
|
await this._tryUpdateEvent(pendingEvent);
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("keep sending?", this._amountSent, "<", this._pendingEvents.length);
|
2020-05-07 18:46:16 +02:00
|
|
|
this._amountSent += 1;
|
2019-07-01 10:00:29 +02:00
|
|
|
}
|
2019-07-26 22:03:57 +02:00
|
|
|
} catch(err) {
|
2020-04-19 19:05:12 +02:00
|
|
|
if (err instanceof ConnectionError) {
|
2019-07-26 22:03:57 +02:00
|
|
|
this._offline = true;
|
2019-06-28 00:52:54 +02:00
|
|
|
}
|
2019-07-26 22:03:57 +02:00
|
|
|
} finally {
|
|
|
|
this._isSending = false;
|
2019-06-28 00:52:54 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2019-07-26 22:33:33 +02:00
|
|
|
removeRemoteEchos(events, txn) {
|
|
|
|
const removed = [];
|
|
|
|
for (const event of events) {
|
|
|
|
const txnId = event.unsigned && event.unsigned.transaction_id;
|
2020-03-23 23:00:33 +01:00
|
|
|
let idx;
|
2019-07-26 22:33:33 +02:00
|
|
|
if (txnId) {
|
2020-03-23 23:00:33 +01:00
|
|
|
idx = this._pendingEvents.array.findIndex(pe => pe.txnId === txnId);
|
|
|
|
} else {
|
|
|
|
idx = this._pendingEvents.array.findIndex(pe => pe.remoteId === event.event_id);
|
|
|
|
}
|
|
|
|
if (idx !== -1) {
|
|
|
|
const pendingEvent = this._pendingEvents.get(idx);
|
|
|
|
txn.pendingEvents.remove(pendingEvent.roomId, pendingEvent.queueIndex);
|
|
|
|
removed.push(pendingEvent);
|
2019-07-26 22:33:33 +02:00
|
|
|
}
|
|
|
|
}
|
|
|
|
return removed;
|
|
|
|
}
|
2019-06-28 00:52:54 +02:00
|
|
|
|
2019-07-26 22:33:33 +02:00
|
|
|
emitRemovals(pendingEvents) {
|
|
|
|
for (const pendingEvent of pendingEvents) {
|
|
|
|
const idx = this._pendingEvents.array.indexOf(pendingEvent);
|
|
|
|
if (idx !== -1) {
|
|
|
|
this._amountSent -= 1;
|
|
|
|
this._pendingEvents.remove(idx);
|
|
|
|
}
|
2019-07-26 22:03:57 +02:00
|
|
|
}
|
2019-06-28 00:52:54 +02:00
|
|
|
}
|
|
|
|
|
2019-07-26 22:03:57 +02:00
|
|
|
resumeSending() {
|
|
|
|
this._offline = false;
|
|
|
|
if (!this._isSending) {
|
|
|
|
this._sendLoop();
|
|
|
|
}
|
2019-06-28 00:52:54 +02:00
|
|
|
}
|
2019-07-01 10:00:29 +02:00
|
|
|
|
2020-11-13 17:19:19 +01:00
|
|
|
async enqueueEvent(eventType, content, attachments) {
|
|
|
|
const pendingEvent = await this._createAndStoreEvent(eventType, content, attachments);
|
2019-07-26 22:03:57 +02:00
|
|
|
this._pendingEvents.set(pendingEvent);
|
2019-09-15 12:25:14 +02:00
|
|
|
console.log("added to _pendingEvents set", this._pendingEvents.length);
|
2019-07-26 22:03:57 +02:00
|
|
|
if (!this._isSending && !this._offline) {
|
|
|
|
this._sendLoop();
|
|
|
|
}
|
2019-07-01 10:00:29 +02:00
|
|
|
}
|
2019-06-28 00:52:54 +02:00
|
|
|
|
2019-07-26 22:03:57 +02:00
|
|
|
get pendingEvents() {
|
|
|
|
return this._pendingEvents;
|
2019-07-01 10:00:29 +02:00
|
|
|
}
|
|
|
|
|
2019-07-26 22:03:57 +02:00
|
|
|
async _tryUpdateEvent(pendingEvent) {
|
2020-09-25 16:42:41 +02:00
|
|
|
const txn = this._storage.readWriteTxn([this._storage.storeNames.pendingEvents]);
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("_tryUpdateEvent: got txn");
|
2019-07-26 22:03:57 +02:00
|
|
|
try {
|
|
|
|
// pendingEvent might have been removed already here
|
|
|
|
// by a racing remote echo, so check first so we don't recreate it
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("_tryUpdateEvent: before exists");
|
2019-07-26 22:03:57 +02:00
|
|
|
if (await txn.pendingEvents.exists(pendingEvent.roomId, pendingEvent.queueIndex)) {
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("_tryUpdateEvent: inside if exists");
|
2019-07-26 22:03:57 +02:00
|
|
|
txn.pendingEvents.update(pendingEvent.data);
|
|
|
|
}
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("_tryUpdateEvent: after exists");
|
2019-07-26 22:03:57 +02:00
|
|
|
} catch (err) {
|
|
|
|
txn.abort();
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("_tryUpdateEvent: error", err);
|
2019-07-26 22:03:57 +02:00
|
|
|
throw err;
|
|
|
|
}
|
2019-07-27 10:40:56 +02:00
|
|
|
console.log("_tryUpdateEvent: try complete");
|
2019-07-26 22:03:57 +02:00
|
|
|
await txn.complete();
|
|
|
|
}
|
2019-07-01 10:00:29 +02:00
|
|
|
|
2020-11-13 17:19:19 +01:00
|
|
|
async _createAndStoreEvent(eventType, content, attachments) {
|
2019-09-15 12:25:14 +02:00
|
|
|
console.log("_createAndStoreEvent");
|
2020-09-25 16:42:41 +02:00
|
|
|
const txn = this._storage.readWriteTxn([this._storage.storeNames.pendingEvents]);
|
2019-07-26 22:03:57 +02:00
|
|
|
let pendingEvent;
|
|
|
|
try {
|
|
|
|
const pendingEventsStore = txn.pendingEvents;
|
2019-09-15 12:25:14 +02:00
|
|
|
console.log("_createAndStoreEvent getting maxQueueIndex");
|
2019-07-26 22:03:57 +02:00
|
|
|
const maxQueueIndex = await pendingEventsStore.getMaxQueueIndex(this._roomId) || 0;
|
2019-09-15 12:25:14 +02:00
|
|
|
console.log("_createAndStoreEvent got maxQueueIndex", maxQueueIndex);
|
2019-07-26 22:03:57 +02:00
|
|
|
const queueIndex = maxQueueIndex + 1;
|
|
|
|
pendingEvent = new PendingEvent({
|
|
|
|
roomId: this._roomId,
|
|
|
|
queueIndex,
|
|
|
|
eventType,
|
|
|
|
content,
|
2020-09-03 15:36:48 +02:00
|
|
|
txnId: makeTxnId(),
|
|
|
|
needsEncryption: !!this._roomEncryption
|
2020-11-13 17:19:19 +01:00
|
|
|
}, attachments);
|
2019-09-15 12:25:14 +02:00
|
|
|
console.log("_createAndStoreEvent: adding to pendingEventsStore");
|
2019-07-26 22:03:57 +02:00
|
|
|
pendingEventsStore.add(pendingEvent.data);
|
|
|
|
} catch (err) {
|
|
|
|
txn.abort();
|
|
|
|
throw err;
|
|
|
|
}
|
2019-07-01 10:00:29 +02:00
|
|
|
await txn.complete();
|
2019-07-26 22:03:57 +02:00
|
|
|
return pendingEvent;
|
2019-07-01 10:00:29 +02:00
|
|
|
}
|
2020-11-13 17:19:19 +01:00
|
|
|
|
|
|
|
async _uploadAttachments(pendingEvent) {
|
|
|
|
const {attachments} = pendingEvent;
|
|
|
|
for (const [urlPath, attachment] of Object.entries(attachments)) {
|
|
|
|
await attachment.upload();
|
|
|
|
attachment.applyToContent(urlPath, pendingEvent.content);
|
|
|
|
}
|
|
|
|
}
|
2019-06-28 00:52:54 +02:00
|
|
|
}
|