From 33421c8e5a842e38ee7a3818f25d4117208096e4 Mon Sep 17 00:00:00 2001 From: Etienne Samson Date: Sun, 10 May 2026 14:43:57 +0200 Subject: [PATCH 1/2] =?UTF-8?q?=F0=9F=A9=B9=20Move=20history=20saving=20to?= =?UTF-8?q?=20the=20start=20of=20the=20next=20tick?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The idea is to side-step two issues: - mid-tick changes like ruins and resources have a temporary numeric `_id` assigned to them that leads to issues with some object overwriting others. - anything related to global intents, notably inter-room movement, appears different in replays than it does in real-time, since replays grab the position of the creep on an exit tile, while real-time has them on the destination instead. --- src/processor.js | 39 ++++++++++++++++++++++----------------- 1 file changed, 22 insertions(+), 17 deletions(-) diff --git a/src/processor.js b/src/processor.js index 7b7b39fa..8769f2c4 100644 --- a/src/processor.js +++ b/src/processor.js @@ -17,6 +17,26 @@ function processRoom(roomId, {intents, roomObjects, users, roomTerrain, gameTime return q.when().then(() => { + if (gameTime > 0) { + var historyPayload = {}; + _.forEach(roomObjects, (object) => { + if (!object || object.type === 'flag') { + return; + } + if (object.type === 'creep' || object.type === 'powerCreep') { + var clone = JSON.parse(JSON.stringify(object)); + clone._id = '' + object._id; + if (clone.actionLog && clone.actionLog.say && !clone.actionLog.say.isPublic) { + delete clone.actionLog.say; + } + historyPayload[clone._id] = clone; + } else { + historyPayload[object._id] = object; + } + }); + void saveRoomHistory(roomId, historyPayload, gameTime - 1); + } + var bulk = driver.bulkObjectsWrite(), bulkUsers = driver.bulkUsersWrite(), bulkFlags = driver.bulkFlagsWrite(), @@ -24,7 +44,6 @@ function processRoom(roomId, {intents, roomObjects, users, roomTerrain, gameTime oldObjects = {}, hasNewbieWalls = false, stats = driver.getRoomStatsUpdater(roomId), - objectsToHistory = {}, roomSpawns = [], roomExtensions = [], roomNukes = [], keepers = [], invaders = [], invaderCore = null, oldRoomInfo = _.clone(roomInfo); @@ -425,20 +444,6 @@ function processRoom(roomId, {intents, roomObjects, users, roomTerrain, gameTime } } - if (object.type != 'flag') { - objectsToHistory[object._id] = object; - - if (object.type == 'creep' || object.type == 'powerCreep') { - objectsToHistory[object._id] = JSON.parse(JSON.stringify(object)); - objectsToHistory[object._id]._id = "" + object._id; - delete objectsToHistory[object._id]._actionLog; - delete objectsToHistory[object._id]._ticksToLive; - if (object.actionLog.say && !object.actionLog.say.isPublic) { - delete objectsToHistory[object._id].actionLog.say; - } - } - } - if (object.user) { //userVisibility[object.user] = true; @@ -503,7 +508,6 @@ function processRoom(roomId, {intents, roomObjects, users, roomTerrain, gameTime if(activateRoom) { driver.activateRoom(roomId); - saveRoomHistory(roomId, objectsToHistory, gameTime); } if(!_.isEqual(roomInfo, oldRoomInfo)) { @@ -521,6 +525,8 @@ function processRoom(roomId, {intents, roomObjects, users, roomTerrain, gameTime function saveRoomHistory(roomId, objects, gameTime) { + var data = JSON.stringify(objects); + return currentHistoryPromise.then(() => { var promise = q.when(); @@ -529,7 +535,6 @@ function saveRoomHistory(roomId, objects, gameTime) { promise = driver.history.upload(roomId, baseTime); } - var data = JSON.stringify(objects); currentHistoryPromise = promise.then(() => driver.history.saveTick(roomId, gameTime, data)); return currentHistoryPromise; }); From 7dc98b0d40d174c04997507ea41a9af164037197 Mon Sep 17 00:00:00 2001 From: Etienne Samson Date: Mon, 7 Sep 2026 09:36:38 +0200 Subject: [PATCH 2/2] =?UTF-8?q?=F0=9F=A9=B9=20Add=20history=20flushing?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit This adds two phases to the main loop; - flushHistory takes care of uploading room chunks for rooms that aren't active by the time the chunk size cutoff comes. - saveDeactivatedRoomHistory handles saving the last tick data for rooms that don't stay active the following tick. Depends on a driver change. --- spec/engine/historySpec.js | 112 +++++++++++++++++++++++++++++++++++++ src/history.js | 87 ++++++++++++++++++++++++++++ src/main.js | 20 ++++++- src/processor.js | 41 ++------------ 4 files changed, 221 insertions(+), 39 deletions(-) create mode 100644 spec/engine/historySpec.js create mode 100644 src/history.js diff --git a/spec/engine/historySpec.js b/spec/engine/historySpec.js new file mode 100644 index 00000000..c84e9621 --- /dev/null +++ b/spec/engine/historySpec.js @@ -0,0 +1,112 @@ +const q = require('q'), + utils = require('../../src/utils'), + driver = utils.getDriver(), + history = require('../../src/history'); + +describe('history', () => { + const originalHistory = driver.history; + const originalChunkSize = driver.config.historyChunkSize; + const originalGetActiveRooms = driver.getActiveRooms; + const originalGetRoomObjects = driver.getRoomObjects; + + beforeEach(() => { + driver.config.historyChunkSize = 20; + driver.history = { + saveTick: jasmine.createSpy('saveTick').and.callFake(() => q.when()), + upload: jasmine.createSpy('upload').and.callFake(() => q.when()), + markPendingHistory: jasmine.createSpy('markPendingHistory').and.callFake(() => q.when()), + takePendingHistory: jasmine.createSpy('takePendingHistory').and.callFake(() => q.when([])) + }; + driver.getActiveRooms = jasmine.createSpy('getActiveRooms').and.callFake(() => q.when([])); + driver.getRoomObjects = jasmine.createSpy('getRoomObjects').and.callFake(() => q.when({objects: {}})); + }); + + afterEach(() => { + driver.history = originalHistory; + driver.config.historyChunkSize = originalChunkSize; + driver.getActiveRooms = originalGetActiveRooms; + driver.getRoomObjects = originalGetRoomObjects; + }); + + describe('buildHistoryPayload', () => { + it('skips flags and strips private say', () => { + const payload = history.buildHistoryPayload({ + a: {_id: 'a', type: 'source'}, + b: {_id: 'b', type: 'flag'}, + c: { + _id: 'c', + type: 'creep', + actionLog: {say: {message: 'hi', isPublic: false}} + } + }); + + expect(payload.a).toEqual({_id: 'a', type: 'source'}); + expect(payload.b).toBeUndefined(); + expect(payload.c._id).toBe('c'); + expect(payload.c.actionLog.say).toBeUndefined(); + }); + }); + + describe('saveRoomHistory', () => { + it('does not upload mid-chunk', () => { + return history.saveRoomHistory('W1N1', {a: 1}, 10).then(() => { + expect(driver.history.upload).not.toHaveBeenCalled(); + expect(driver.history.saveTick).toHaveBeenCalledWith('W1N1', 10, JSON.stringify({a: 1})); + expect(driver.history.markPendingHistory).toHaveBeenCalledWith('W1N1', 0); + }); + }); + + it('uploads the previous chunk on a boundary then saveTicks the new chunk', () => { + return history.saveRoomHistory('W1N1', {a: 1}, 20).then(() => { + expect(driver.history.upload).toHaveBeenCalledWith('W1N1', 0); + expect(driver.history.saveTick).toHaveBeenCalledWith('W1N1', 20, JSON.stringify({a: 1})); + expect(driver.history.markPendingHistory).toHaveBeenCalledWith('W1N1', 20); + expect(driver.history.upload.calls.count()).toBe(1); + }); + }); + }); + + describe('uploadPendingChunks', () => { + it('does not flush while the chunk can still receive ticks', () => { + return history.uploadPendingChunks(20).then(() => { + expect(driver.history.takePendingHistory).not.toHaveBeenCalled(); + expect(driver.history.upload).not.toHaveBeenCalled(); + }); + }); + + it('uploads leftover rooms after active rooms have closed the chunk', () => { + driver.history.takePendingHistory.and.callFake(() => q.when(['W1N1', 'W2N2', 'W3N3'])); + return history.uploadPendingChunks(21, ['W1N1']).then(() => { + expect(driver.history.takePendingHistory).toHaveBeenCalledWith(0); + expect(driver.history.upload).not.toHaveBeenCalledWith('W1N1', 0); + expect(driver.history.upload).toHaveBeenCalledWith('W2N2', 0); + expect(driver.history.upload).toHaveBeenCalledWith('W3N3', 0); + }); + }); + }); + + describe('saveDeactivatedRoomsHistory', () => { + it('saves the last tick for rooms that will not run next tick', () => { + driver.getActiveRooms.and.callFake(() => q.when(['W1N1'])); + driver.getRoomObjects.and.callFake(roomId => q.when({ + objects: {src: {_id: 'src', type: 'source', room: roomId}} + })); + return history.saveDeactivatedRoomsHistory(['W1N1', 'W2N2'], 10).then(() => { + expect(driver.getRoomObjects).toHaveBeenCalledWith('W2N2'); + expect(driver.getRoomObjects).not.toHaveBeenCalledWith('W1N1'); + expect(driver.history.saveTick).toHaveBeenCalledWith( + 'W2N2', 10, JSON.stringify({src: {_id: 'src', type: 'source', room: 'W2N2'}})); + expect(driver.history.upload).not.toHaveBeenCalled(); + }); + }); + + it('closes the previous chunk when the last tick is a boundary', () => { + driver.getActiveRooms.and.callFake(() => q.when([])); + driver.getRoomObjects.and.callFake(() => q.when({objects: {}})); + return history.saveDeactivatedRoomsHistory(['W2N2'], 20).then(() => { + expect(driver.history.upload).toHaveBeenCalledWith('W2N2', 0); + expect(driver.history.saveTick).toHaveBeenCalledWith('W2N2', 20, JSON.stringify({})); + }); + }); + }); +}); diff --git a/src/history.js b/src/history.js new file mode 100644 index 00000000..18d311d3 --- /dev/null +++ b/src/history.js @@ -0,0 +1,87 @@ +var q = require('q'), + _ = require('lodash'), + utils = require('./utils'), + driver = utils.getDriver(); + +var currentHistoryPromise = q.when(); + +exports.buildHistoryPayload = function(roomObjects) { + var historyPayload = {}; + _.forEach(roomObjects, (object) => { + if (!object || object.type === 'flag') { + return; + } + if (object.type === 'creep' || object.type === 'powerCreep') { + var clone = JSON.parse(JSON.stringify(object)); + clone._id = '' + object._id; + if (clone.actionLog && clone.actionLog.say && !clone.actionLog.say.isPublic) { + delete clone.actionLog.say; + } + historyPayload[clone._id] = clone; + } else { + historyPayload[object._id] = object; + } + }); + return historyPayload; +}; + +exports.saveRoomHistory = function(roomId, objects, gameTime) { + + var data = JSON.stringify(objects); + var chunkSize = driver.config.historyChunkSize; + + return currentHistoryPromise.then(() => { + var promise = q.when(); + + if (!(gameTime % chunkSize)) { + var prevBase = Math.floor((gameTime - 1) / chunkSize) * chunkSize; + promise = driver.history.upload(roomId, prevBase); + } + + var baseTime = gameTime - (gameTime % chunkSize); + currentHistoryPromise = promise.then(() => driver.history.saveTick(roomId, gameTime, data) + .then(() => driver.history.markPendingHistory(roomId, baseTime))); + return currentHistoryPromise; + }); +}; + +exports.uploadPendingChunks = function(gameTime, processedRooms) { + var chunkSize = driver.config.historyChunkSize; + if ((gameTime - 1) % chunkSize) { + return q.when(); + } + var prevBase = Math.floor((gameTime - 2) / chunkSize) * chunkSize; + if (prevBase < 0) { + return q.when(); + } + var skip = {}; + _.forEach(processedRooms, roomId => { + skip[roomId] = true; + }); + return driver.history.takePendingHistory(prevBase) + .then(rooms => q.all(_.map(rooms || [], roomId => { + if (skip[roomId]) { + return; + } + return driver.history.upload(roomId, prevBase); + }))); +}; + +exports.saveDeactivatedRoomsHistory = function(processedRooms, gameTime) { + if (!processedRooms || !processedRooms.length) { + return q.when(); + } + return driver.getActiveRooms() + .then(active => { + var stillActive = {}; + _.forEach(active || [], roomId => { + stillActive[roomId] = true; + }); + var deactivated = _.filter(processedRooms, roomId => !stillActive[roomId]); + return q.all(_.map(deactivated, roomId => + driver.getRoomObjects(roomId).then(result => + exports.saveRoomHistory(roomId, exports.buildHistoryPayload(result && result.objects), gameTime) + ) + )); + }); +}; diff --git a/src/main.js b/src/main.js index 914eac1e..2f6a339d 100644 --- a/src/main.js +++ b/src/main.js @@ -3,7 +3,8 @@ var q = require('q'), _ = require('lodash'), utils = require('./utils'), driver = utils.getDriver(), - config = require('./config'); + config = require('./config') + history = require('./history'); var lastAccessibleRoomsUpdate = 0; var roomsQueue, usersQueue; @@ -11,7 +12,9 @@ var roomsQueue, usersQueue; function loop() { var resetInterval, startLoopTime = process.hrtime ? process.hrtime() : Date.now(), - stage = 'start'; + stage = 'start', + processedRooms, + tickGameTime; driver.config.emit('mainLoopStage',stage); @@ -45,6 +48,7 @@ function loop() { return driver.getAllRoomsNames(); }) .then((rooms) => { + processedRooms = rooms; stage = 'addRoomsToQueue'; driver.config.emit('mainLoopStage',stage, rooms); return roomsQueue.addMulti(rooms); @@ -54,6 +58,13 @@ function loop() { driver.config.emit('mainLoopStage',stage); return roomsQueue.whenAllDone(); }) + .then(() => driver.getGameTime()) + .then((gameTime) => { + tickGameTime = gameTime; + stage = 'flushHistory'; + driver.config.emit('mainLoopStage',stage); + return history.uploadPendingChunks(gameTime, processedRooms); + }) .then(() => { stage = 'commit1'; driver.config.emit('mainLoopStage',stage); @@ -69,6 +80,11 @@ function loop() { driver.config.emit('mainLoopStage',stage); return driver.commitDbBulk(); }) + .then(() => { + stage = 'saveDeactivatedRoomHistory'; + driver.config.emit('mainLoopStage',stage); + return history.saveDeactivatedRoomsHistory(processedRooms, tickGameTime); + }) .then(() => { stage = 'incrementGameTime'; driver.config.emit('mainLoopStage',stage); diff --git a/src/processor.js b/src/processor.js index 8769f2c4..8eaf7958 100644 --- a/src/processor.js +++ b/src/processor.js @@ -6,9 +6,10 @@ var q = require('q'), driver = utils.getDriver(), C = driver.constants, config = require('./config'), - fakeRuntime = require('./processor/common/fake-runtime'); + fakeRuntime = require('./processor/common/fake-runtime'), + history = require('./history'); -var roomsQueue, usersQueue, lastRoomsStatsSaveTime = 0, currentHistoryPromise = q.when(); +var roomsQueue, usersQueue, lastRoomsStatsSaveTime = 0; const KEEPER_ID = "3"; const INVADER_ID = "2"; @@ -18,23 +19,7 @@ function processRoom(roomId, {intents, roomObjects, users, roomTerrain, gameTime return q.when().then(() => { if (gameTime > 0) { - var historyPayload = {}; - _.forEach(roomObjects, (object) => { - if (!object || object.type === 'flag') { - return; - } - if (object.type === 'creep' || object.type === 'powerCreep') { - var clone = JSON.parse(JSON.stringify(object)); - clone._id = '' + object._id; - if (clone.actionLog && clone.actionLog.say && !clone.actionLog.say.isPublic) { - delete clone.actionLog.say; - } - historyPayload[clone._id] = clone; - } else { - historyPayload[object._id] = object; - } - }); - void saveRoomHistory(roomId, historyPayload, gameTime - 1); + history.saveRoomHistory(roomId, history.buildHistoryPayload(roomObjects), gameTime - 1); } var bulk = driver.bulkObjectsWrite(), @@ -523,24 +508,6 @@ function processRoom(roomId, {intents, roomObjects, users, roomTerrain, gameTime }); } -function saveRoomHistory(roomId, objects, gameTime) { - - var data = JSON.stringify(objects); - - return currentHistoryPromise.then(() => { - var promise = q.when(); - - if (!(gameTime % driver.config.historyChunkSize)) { - var baseTime = Math.floor((gameTime - 1) / driver.config.historyChunkSize) * driver.config.historyChunkSize; - promise = driver.history.upload(roomId, baseTime); - } - - currentHistoryPromise = promise.then(() => driver.history.saveTick(roomId, gameTime, data)); - return currentHistoryPromise; - }); -} - - driver.connect('processor') .then(() => driver.queue.create('rooms', 'read')) .catch((error) => {