diff --git a/Furious/Service/ProfileTesting.py b/Furious/Service/ProfileTesting.py index f8af4fc..860e986 100644 --- a/Furious/Service/ProfileTesting.py +++ b/Furious/Service/ProfileTesting.py @@ -313,14 +313,17 @@ class _LatencyScheduler(QtCore.QObject): self.threadPool.setMaxThreadCount(self.maxConcurrency) self.pingWorkerFactory = pingWorkerFactory + self.queue = collections.deque() self.activeJobs = {} self.tcpingRequests = {} self.tcpingEndpointRequests = {} self.tcpingCompletionQueue = collections.deque() self.nextTcpingRequestId = 1 + self.tcpingThread = None self.tcpingEngine = None + self.drainScheduled = False self.tcpingCompletionScheduled = False self.shuttingDown = False @@ -363,6 +366,7 @@ class _LatencyScheduler(QtCore.QObject): engine = TcpingEngine(self, self.tcpingMaxConcurrency) thread = TcpingThread(engine, self) + engine.moveToThread(thread) self.tcpingThread = thread @@ -458,6 +462,7 @@ class _LatencyScheduler(QtCore.QObject): self.tcpingEndpointRequests.pop(group.endpointKey, None) self.tcpingCompletionQueue.append(requestId) + self.scheduleTcpingCompletion() def scheduleTcpingCompletion(self): diff --git a/Furious/Service/SubscriptionSync.py b/Furious/Service/SubscriptionSync.py index 3c5d095..8e6d7d3 100644 --- a/Furious/Service/SubscriptionSync.py +++ b/Furious/Service/SubscriptionSync.py @@ -158,6 +158,7 @@ class SubscriptionSynchronizer: existing.connection = prepared.connection existing.metadata = metadata + synchronized.append(existing) for removed in existingById.values(): diff --git a/Furious/Service/TcpingService.py b/Furious/Service/TcpingService.py index 980d312..89c9b03 100644 --- a/Furious/Service/TcpingService.py +++ b/Furious/Service/TcpingService.py @@ -100,6 +100,7 @@ class TcpingProbe(QtCore.QObject): self.request = request self.completionHasRun = False + self.elapsedTimer = QtCore.QElapsedTimer() self.socket = QTcpSocket(self) self.timeoutTimer = QtCore.QTimer(self) diff --git a/tests/test_runtime_lifecycle.py b/tests/test_runtime_lifecycle.py index 7ab02f2..c5f415d 100644 --- a/tests/test_runtime_lifecycle.py +++ b/tests/test_runtime_lifecycle.py @@ -106,6 +106,7 @@ class RuntimeLifecycleTest(TestCase): runtime = _Runtime(exitCallback=router.publish) router.attach(runtime, owner) lease = RuntimeLease(runtime, router) + thread = threading.Thread(target=runtime.publishCode, args=(9,)) thread.start() @@ -114,6 +115,7 @@ class RuntimeLifecycleTest(TestCase): self.assertTrue(waitFor(lambda: bool(owner.events))) self.assertIs(owner.events[0][2], self.app.thread()) self.assertEqual(owner.events[0][1].code, 9) + lease.release() def testRouterRejectsAnUntypedRuntimeExit(self): @@ -146,6 +148,7 @@ class RuntimeLifecycleTest(TestCase): for _index in range(32): router = RuntimeEventRouter() references.append(weakref.ref(router)) + router.finishRelease() router.deleteLater() @@ -166,12 +169,14 @@ class RuntimeLifecycleTest(TestCase): runtime.publishCode(11) lease.commit(lambda current, event: committed.append((current, event))) + processQtEvents() self.assertEqual(owner.events, []) self.assertEqual(len(committed), 1) self.assertIs(committed[0][0], runtime) self.assertEqual(committed[0][1].code, 11) + lease.release() def testPostCommitDuplicateExitIsDeliveredOnce(self): @@ -186,11 +191,13 @@ class RuntimeLifecycleTest(TestCase): runtime.publishCode(12) runtime.publishCode(13) + processQtEvents() self.assertEqual(owner.events, []) self.assertEqual(len(committed), 1) self.assertEqual(committed[0][1].code, 12) + lease.release() def testReleaseSuppressesLateQueuedExitAndIsIdempotent(self): @@ -202,8 +209,10 @@ class RuntimeLifecycleTest(TestCase): lease = RuntimeLease(runtime, router) runtime.publishCode(13) + lease.release() lease.release() + processQtEvents() self.assertEqual(owner.events, []) @@ -224,9 +233,12 @@ class RuntimeLifecycleTest(TestCase): runtime = _ProcessRuntime( exitCallback=lambda current, event: events.append((current, event)) ) + runtime.start() + self.assertIs(runtime.state, RuntimeState.Alive) self.assertTrue(runtime.isRunning()) + runtime._pollProcess() runtime._pollProcess() @@ -234,9 +246,12 @@ class RuntimeLifecycleTest(TestCase): self.assertIs(events[0][1].reason, RuntimeExitReason.InvalidConfiguration) self.assertEqual(events[0][1].code, 23) process.close.assert_called_once_with() + runtime.dispose() + self.assertIsNone(runtime._monitorConnection) self.assertIsNone(runtime._output._timerConnection) + runtime.dispose() def testMultiprocessingSpawnFailureRaisesStructuredError(self): @@ -255,7 +270,9 @@ class RuntimeLifecycleTest(TestCase): self.assertIs(raised.exception.reason, RuntimeExitReason.StartFailure) self.assertIn('fixture spawn failure', raised.exception.details) + runtime.dispose() + self.assertIsNone(runtime._monitorConnection) self.assertIsNone(runtime._output._timerConnection) diff --git a/tests/test_subscription_sync.py b/tests/test_subscription_sync.py index a353d0b..209f886 100644 --- a/tests/test_subscription_sync.py +++ b/tests/test_subscription_sync.py @@ -77,6 +77,7 @@ class SubscriptionSynchronizerTest(unittest.TestCase): key='upstream:other', ) originalId = retained.metadata.profileId + profiles = [manual, retained, removed, other] incoming = [ profile('Remote label', 'new.example', key='upstream:one'), @@ -151,6 +152,7 @@ class SubscriptionSynchronizerTest(unittest.TestCase): ) incoming = profile('Incoming', 'new.example') profiles = [retained] + originalMetadata = retained.metadata.toMapping() originalConnection = retained.connection.deepcopy() @@ -169,6 +171,7 @@ class SubscriptionSynchronizerTest(unittest.TestCase): self.assertEqual(retained.metadata.toMapping(), originalMetadata) self.assertEqual(retained.connection, originalConnection) self.assertFalse(retained.deleted) + self.assertEqual(incoming.itemSubscription, '') self.assertFalse(incoming.itemSubscriptionManaged) @@ -183,8 +186,10 @@ class SubscriptionSynchronizerTest(unittest.TestCase): ) unrelated = profile('Manual', 'manual.example') profiles = [retained, unrelated] + snapshot = SubscriptionSynchronizer().snapshot(profiles, 'group') incoming = profile('After', 'new.example', key='upstream:one') + plan = SubscriptionSynchronizer().prepare(snapshot, (incoming,)) retained.metadata.annotations = 'edited while preparing' @@ -194,8 +199,10 @@ class SubscriptionSynchronizerTest(unittest.TestCase): self.assertIs(profiles[0], retained) self.assertIs(profiles[1], unrelated) + self.assertEqual(retained.connection['address'], 'new.example') self.assertEqual(retained.itemRemark, 'After') + self.assertEqual(retained.metadata.annotations, 'edited while preparing') self.assertEqual(retained.metadata.latency, '41 ms') self.assertEqual(result.changedProfileIds, (retained.metadata.profileId,)) @@ -211,6 +218,7 @@ class SubscriptionSynchronizerTest(unittest.TestCase): ) profiles = [retained] synchronizer = SubscriptionSynchronizer() + snapshot = synchronizer.snapshot(profiles, 'group') plan = synchronizer.prepare( snapshot,