-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathP2PeerExplorer.cpp
More file actions
1963 lines (1808 loc) · 66 KB
/
Copy pathP2PeerExplorer.cpp
File metadata and controls
1963 lines (1808 loc) · 66 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
// Copyright © 2011, 2026 Ivyware Pty Ltd, Khrustal & Mann
// MELBOURNE, VICTORIA, AUSTRALIA, 3000
//
// 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.
//
//
// P2PeerExplorer definitions and prototypes
// NOTES: Network is built upon the local and remote exchange of
// P2PeerMsg's between objects derived from this base class
// : Only P2PeerExplorer derived network objects may be assigned
// network P2Paddress's
#include "stdafx.h"
//#include "P2PTL.h"
#include "Kernel32_Ext.h"
#include "P2PeerExplorer.h"
#include "P2Pwin32.h"
#include "Msgexception.h"
///////////////////////////////////////////////////////////////////////
// Constructors and destructor
//
// Constructors and destructor
//
// Parameters: P2PeerHub *pHub
// Hub to which this P2PeerExplorer is attached
// NOTES: P2PeerHub's only support single P2PeerExplorer
// instance is assigned
//
P2PeerExplorer::P2PeerExplorer ( P2PeerHub *pHub )
{
// Firstly
RenderExplorerSafe ( );
m_pHub = pHub;
}
P2PeerExplorer::~P2PeerExplorer ( )
{
// Closure
// NOTES: RegisterP2Pexpump() and CloseP2Pexpump() are BOTH keyed on
// GetCurrentThreadId() and are only legal from the expump's own
// context - RegisterP2Pexpump() THROWS from anywhere else
// (P2Pwin32.cpp:1117-1133). A P2PeerExpump is deleted by whoever
// owns it, which is the hub owner's thread (P2PeerHub::CloseHub ->
// DropP2PeerExpump), never the expump thread, so calling them here
// unconditionally threw a P2Pevent out of a destructor on every
// single teardown - measured: p2p_expreg reached its verdict and
// then hung in shutdown until the ctest timeout. The throw escaped
// DropP2PeerExpump() mid-CloseHub() and, because that nulled the
// member only AFTER delete returned, also left m_pP2PeerExpump
// dangling on freed storage for ~P2PeerHub (P2PeerHub.cpp:99-107)
// to delete a second time.
// : So: from a foreign thread, signal the expump down and wait - its
// own ProcExpump() performs the de-registration and CloseP2Pexpump()
// in the only context where they are meaningful. From the expump's
// own context (an expump destroying itself), do it directly.
// : Nothing here may throw. A destructor that throws during unwinding
// terminates the process.
// F-S5-4. Decide whose thread we are on ONCE, here, from the id as it
// stands before anything below can zero it. ProcExpump() stores 0 into
// m_nExpumpID as its own last act, so re-reading the member further down
// would report "not the expump thread" for an expump destroying itself,
// and the handle wait below would be a self-join that never returns.
const DWORD nExpumpIDatEntry = (DWORD)m_nExpumpID;
const bool bExpumpContext = nExpumpIDatEntry > 0 &&
nExpumpIDatEntry == GetCurrentThreadId();
try
{
DWORD nExpumpID = nExpumpIDatEntry;
if ( nExpumpID > 0 &&
nExpumpID != GetCurrentThreadId() )
{
if ( P2PmsgPumpExists(nExpumpID) )
CloseExpump ( nExpumpID );
}
else if ( nExpumpID > 0 )
{
// De-registrations
RegisterP2Pexpump ( P2PmsgExp_ALL, FALSE );
CloseP2Pexpump ( );
}
}
catch ( P2Pevent *pEVT )
{
pEVT->Advice(_T("~P2PeerExplorer: P2Pexpump closure"))->Cancel();
}
catch ( ... )
{
}
// Resources
// NOTES: Last, not first - the expump may still have been inside a
// P2PsafeCS on it until the wait above returned
// : F-S5-4 - give the thread handle back. From a foreign thread the
// block above has already signalled the expump down and spun until
// it de-registered, but de-registration happens INSIDE ProcExpump()
// with the rest of its tail still to run, so "the pump no longer
// exists" is not "the thread has ended". Wait for the thread
// itself before closing, or the close is a detach and everything
// the thread still holds stays unreachable and unfreed - which is
// precisely the 16 indirect allocations LSan hung off this handle.
// From the expump's OWN context (an expump destroying itself) that
// wait is a self-join, so close without waiting and let the shim
// detach; there is no second thread to keep anything alive.
if ( m_hExpumpThread )
{
if ( !bExpumpContext )
WaitForSingleObject ( m_hExpumpThread, INFINITE );
CloseHandle ( m_hExpumpThread );
m_hExpumpThread = 0;
}
DeleteCriticalSection ( &m_oCSectionExpump );
}
void
P2PeerExplorer::RenderExplorerSafe()
{
// Attributes
m_pHub = 0; // Attached P2PeerHub
m_nExpumpID = 0; // Dedicated explorer P2PumpID
m_hExpumpThread = 0; // F-S5-4: owned thread handle, if spawned
m_nMaxECid = 32; // XC%i exploration slots (was NEVER assigned:
// it bounds the slot loops in On_XCidConLogin
// and On_XCidConClose, so an unassigned value
// means "no slots" or "four billion slots"
// depending on what the allocator left behind)
// Resources
InitializeCriticalSection ( &m_oCSectionExpump );
}
P2PeerExpump*
P2PeerExpump::DynamicCreate ( P2PeerHub *pHub
, LPTHREAD_START_ROUTINE pfnThreadProc
, RUN_EXPUMP pfnRunExpump )
{
// To be sure, to be sure
if ( pHub->GetP2PeerExpump() )
EVERR->MODULE
->Message("P2PeerExpump already exists" )
->Throw();
P2PeerExpSP spExpump = new P2PeerExpump ( pHub );
spExpump -> SpawnExpump ( pfnThreadProc, pfnRunExpump );
pHub -> PostP2PeerExpump ( spExpump.Dereference() );
// Tidy up, and
return dynamic_cast<P2PeerExpump *>(pHub->GetP2PeerExpump());
}
void
CMapStringToString_SetAtNC ( CMapADDR2ADDR *pCMap, P2PaddrSTR lpszSource )
{
UNREFERENCED_PARAMETER(pCMap);
CStringADDR strSource = lpszSource;
CStringADDR strSOURCE = lpszSource;
strSOURCE.MakeUpper ( );
// NOTE: intended body is pCMap->SetAt(strSOURCE, strSource); it does not
// compile because CMapADDR2ADDR (CMap<CStringADDR,...>) has no
// HashKey<CStringADDR> specialisation. Left disabled until the map
// type provides one. No callers exist today.
ASSERT(0);//pCMap -> SetAt ( strSOURCE, strSource );
}
BOOL
CMapStringToString_RemoveKeyNC ( CMapADDR2ADDR *pCMap, P2PaddrSTR lpszSource )
{
UNREFERENCED_PARAMETER(pCMap);
CStringADDR strSOURCE = lpszSource;
strSOURCE.MakeUpper ( );
// NOTE: see CMapStringToString_SetAtNC - blocked on a missing
// HashKey<CStringADDR> specialisation; no callers exist today.
ASSERT(0);return FALSE;//return pCMap -> RemoveKey ( strSOURCE );
}
void
CListRegHubs_Insert ( CListRegHubs& oCListRegHubs
, const CStringADDR& strSource )
{
SP2PexpRegHub oRegHub;
oRegHub.strClient = strSource;
oCListRegHubs.AddTail ( oRegHub );
}
void
CListRegHubs_Remove ( CListRegHubs& oCListRegHubs
, const CStringADDR& strSource )
{
POSITION pos = oCListRegHubs.GetHeadPosition();
while ( pos )
{
POSITION posDrop = pos;
SP2PexpRegHub& oRegHub = oCListRegHubs.GetNext ( pos );
if ( strSource.CompareNoCase(oRegHub.strClient) == 0 )
oCListRegHubs.RemoveAt ( posDrop );
}
}
void
CListRegPmps_Insert ( CListRegPmps& oCListRegPmps
, const CStringADDR& strSource, P2PumpID nPumpID )
{
SP2PexpRegPmp oRegPmp;
oRegPmp.nPumpID = nPumpID;
oRegPmp.strClient = strSource;
oCListRegPmps.AddTail ( oRegPmp );
}
void
CListRegPmps_Remove ( CListRegPmps& oCListRegPmps
, const CStringADDR& strSource, P2PumpID nPumpID )
{
POSITION pos = oCListRegPmps.GetHeadPosition();
while ( pos )
{
POSITION posDrop = pos;
SP2PexpRegPmp& oRegPmp = oCListRegPmps.GetNext ( pos );
if ( ( nPumpID == 0 ||
nPumpID == oRegPmp.nPumpID ) &&
strSource.CompareNoCase(oRegPmp.strClient) == 0 )
oCListRegPmps.RemoveAt ( posDrop );
}
}
void
CListRegCons_Insert ( CListRegCons& oCListRegCons
, const CStringADDR& strSource, P2PconID nConID )
{
SP2PexpRegCon oRegCon;
oRegCon.nConID = nConID;
oRegCon.strClient = strSource;
oCListRegCons.AddTail ( oRegCon );
}
void
CListRegCons_Remove ( CListRegCons& oCListRegCons
, const CStringADDR& strSource, P2PconID nConID )
{
POSITION pos = oCListRegCons.GetHeadPosition();
while ( pos )
{
POSITION posDrop = pos;
SP2PexpRegCon& oRegCon = oCListRegCons.GetNext ( pos );
if ( ( nConID == 0 ||
nConID == oRegCon.nConID ) &&
strSource.CompareNoCase(oRegCon.strClient) == 0 )
oCListRegCons.RemoveAt ( posDrop );
}
}
//
// Binds a registration request to the P2Paddr this P2Pexplorer ASSIGNED to the
// requesting connection
// NOTES: On_XCidConLogin refuses any client-declared address and server-assigns
// <hub>.XC%i, so a peer cannot choose its own identity. It CAN, however,
// claim a sub-identity of the slot it was given: the source-binding gate
// (P2PeerCon::GateAppMsgInbound) admits the logged-in identity or an
// IsRable() descendant of it, and IsRable() accepts any prefix ending on a
// hop boundary. Without this, CEX.XC0 registers CEX.XC0.Ghost and the
// registry grows entries no connection owns (SECURITY_REVIEW M2).
// : A source that matches no P2PeerCon is left untouched - local, in-process
// registrants do not come through a connection and are not the subject.
//
static CStringADDR
P2PexpReg_BindSource ( P2PeerExplorer *pExplorer, P2PaddrSTR strSource )
{
// strBound, not strSource, from here on: the caller's pointer comes from
// P2PeerMsg::GetSource(), which on the UTF-16-store platforms points into
// a thread-local scratch ring that the c_wstr() calls in this very loop
// recycle (Platform/p2pstr.h:629-651, and see RouteP2PeerMsg below). The
// copy is taken before the first of them
CStringADDR strBound = strSource;
P2PsafeCS oSafeCS = pExplorer -> m_oCSectionExpump;
P2PeerCon *pConEnum = 0;
while ( EnumP2PexpCon(pExplorer->m_nExpumpID,&pConEnum) )
{
const P2Paddr& oP2PaddrCon = pConEnum -> GetP2Paddress();
if ( oP2PaddrCon == strBound.GetString() )
return strBound; // exactly the assigned slot
if ( oP2PaddrCon.IsChild(strBound.GetString()) )
return oP2PaddrCon.c_wstr(); // a fabricated sub-identity of it
}
return strBound; // not from a P2Pexplorer connection
}
//
// Removes every registration routable through the passed P2Paddr
// NOTES: The exact-match CListReg*_Remove() helpers above release the slot's own
// entry only. A connection dropping takes its whole sub-tree with it -
// nothing under CEX.XC0 can outlive CEX.XC0, and the registry has no other
// eviction path, so an entry missed here is permanent (SECURITY_REVIEW M2)
//
static void
CListRegHubs_RemoveRable ( CListRegHubs& oCListRegHubs, const P2Paddr& oP2Paddr )
{
POSITION pos = oCListRegHubs.GetHeadPosition();
while ( pos )
{
POSITION posDrop = pos;
SP2PexpRegHub& oRegHub = oCListRegHubs.GetNext ( pos );
if ( oP2Paddr.IsRable(oRegHub.strClient.GetString()) )
oCListRegHubs.RemoveAt ( posDrop );
}
}
static void
CListRegCons_RemoveRable ( CListRegCons& oCListRegCons, const P2Paddr& oP2Paddr )
{
POSITION pos = oCListRegCons.GetHeadPosition();
while ( pos )
{
POSITION posDrop = pos;
SP2PexpRegCon& oRegCon = oCListRegCons.GetNext ( pos );
if ( oP2Paddr.IsRable(oRegCon.strClient.GetString()) )
oCListRegCons.RemoveAt ( posDrop );
}
}
static void
CListRegPmps_RemoveRable ( CListRegPmps& oCListRegPmps, const P2Paddr& oP2Paddr )
{
POSITION pos = oCListRegPmps.GetHeadPosition();
while ( pos )
{
POSITION posDrop = pos;
SP2PexpRegPmp& oRegPmp = oCListRegPmps.GetNext ( pos );
if ( oP2Paddr.IsRable(oRegPmp.strClient.GetString()) )
oCListRegPmps.RemoveAt ( posDrop );
}
}
///////////////////////////////////////////////////////////////////////
// P2Pexplorer management
// NOTES: Started using default P2PeerExplorer::SpawnExp() or
// CreateTarget() implementations
typedef struct
{
P2PexpumpID *pnExpumpID;
P2PeerExplorer *pExplorer;
//P2Paddr oP2Paddr;
RUN_EXPUMP pfnRunExpump;
} P2ProcContextExp;
//
// Creates P2PmsgExp and commences pumping
// NOTES: Thread processing may be stopped via CloseExp()
//
//
// Parameters: LPTHREAD_START_ROUTINE pfnThreadProc
// The thread procedure of the new thread. Refer
// win32 CreateThread() for further details
//
// RUN_EXP pfnRunExp = 0
// Operational method of the new thread
//
// P2PexpID *pnExpID
// Pointer to pump identification code. Retained
// and cleared upon P2PmsgExp closure.
//
// Returns: HANDLE
// Returns the handle to the newly created thread
// or NULL on failure.
//
HANDLE
P2PeerExplorer::SpawnExpump ( LPTHREAD_START_ROUTINE pfnThreadProc
, RUN_EXPUMP pfnRunExpump )
{
// Preparation
if ( pfnThreadProc == NULL )
pfnThreadProc = P2PeerExplorer::ProcExpump;
if ( pfnRunExpump == NULL )
pfnRunExpump = RUN_PUMP_cast(&P2PeerExplorer::RunExpump);
ASSERT(pfnRunExpump);
// Pass context through
P2ProcContextExp oContext;
oContext.pnExpumpID = &m_nExpumpID;
oContext.pExplorer = this;
oContext.pfnRunExpump = pfnRunExpump;
// F-S5-4. A re-spawn must not drop the previous handle on the floor.
// CloseExpump() leaves the explorer re-startable by design (see its
// header), so this is reachable and is where the second leak would be.
if ( m_hExpumpThread )
{
CloseHandle ( m_hExpumpThread );
m_hExpumpThread = 0;
}
HANDLE hThread = CreateThread ( 0, 0
, pfnThreadProc, &oContext //pvParam
, 0, &m_nExpumpID );
if ( !hThread )
return FALSE;
// Confirm operation
// NOTES: Contracted for P2PmsgPump operation upon return
// : F-S5-4 - every exit from here on gives the handle back, either to
// the explorer (success) or to the system (failure). Returning 0
// while holding a live handle is how the failure paths leaked.
DWORD dwExitCode;
UINT uSpins = 0;
while ( !P2PmsgPumpExists(m_nExpumpID) )
{
if ( !GetExitCodeThread(hThread,&dwExitCode) ||
dwExitCode != STILL_ACTIVE )
{
CloseHandle ( hThread );
return 0; // Operational failure
}
YieldForP2PmsgPump ( uSpins ); // Yield (escalating - see header)
}
if ( !P2PmsgPumpExists(m_nExpumpID) )
{
CloseHandle ( hThread );
return 0;
}
// Tidy up and
// NOTES: The explorer keeps the handle - see m_hExpumpThread. The value
// is still returned so existing callers keep compiling, but it is
// this object's to close and not the caller's.
m_hExpumpThread = hThread;
return hThread;
}
//
// Closes or stops P2Pexplorer operation
// NOTES: P2Pexplorer operation may be re-started via
// P2PeerTarget::SpawnPump() or P2PeerTarget::CreatePump()
//
// Parameters: P2PexpID nExpID
// Identification code of P2Pexplorer to be closed
void
P2PeerExplorer::CloseExpump ( P2PexpumpID nExpumpID )
{
//
if ( nExpumpID <= 0 )
nExpumpID = m_nExpumpID;
// Close and confirm
SignalP2PmsgPump ( nExpumpID, P2PsigPump_CLOSE );
UINT uSpins = 0;
while ( P2PmsgPumpExists(nExpumpID) )
YieldForP2PmsgPump ( uSpins );
}
//
// Default thread starting address nominated in CreateThread()
// NOTES: Provide alternative implementation to override default
// action. Follows P2PeerPump::ProcVanilla() base class pattern
//
//
// Parameters: void *pvContext
// Thread data passed via CreateThread()
//
// Returns: DWORD
// Completion code
//
DWORD WINAPI
P2PeerExplorer::ProcExpump ( void *pvContext )
{
// Introduce locals
// NOTES: P2PumpContext is assumed not to persist
P2ProcContextExp *pContext = static_cast<P2ProcContextExp *>(pvContext);
P2PeerExplorer *pExplorer = pContext -> pExplorer;
P2PexpumpID *pnExpumpID = pContext -> pnExpumpID;
//P2PmsgHubID nHubID = pTarget -> GetHubID ( );
RUN_EXPUMP pfnRunExpump = pContext -> pfnRunExpump;
pExplorer -> AssertValid ( );
// Mandatory P2PeerPump thread environment
// NOTES: Sequence contains mandatory P2PeerPump and P2Pmsg
// life cycle management sequences
// : Guard the whole life cycle so no exception escapes the thread
// proc and terminates the process (matches P2PeerTarget::ProcPump /
// P2PeerHub::ProcHub). RunExpump() self-guards its pump loop, but
// CreateP2Pexpump()/CloseP2Pexpump() run outside it.
try
{
CreateP2Pexpump ( pExplorer->m_pHub->GetHubID(), pExplorer );
(pExplorer->*pfnRunExpump) ( );
if ( pnExpumpID && *pnExpumpID == GetCurrentThreadId() )
*pnExpumpID = 0; // Flags pump closure
CloseP2Pexpump ( );
}
catch ( P2Pevent *pEVT )
{
pEVT->Advice(_T("P2PeerExplorer::ProcExpump terminated"))->Cancel();
}
catch ( ... )
{
EVERR->Module (_T("P2PeerExplorer::ProcExpump") )
->Message(_T("Last resort exception of unknown type, "
"P2Pexpump terminated") )
->Cancel();
}
// Tidy up, and
return 0;
}
//
// Plain network P2PmsgPump implementation
// NOTES: Simply pumps P2PeerSys, P2PeerCon and P2PeerMsg objects
// through the P2PeerTarget base class
// : Usually P2PeerMsg's are swapped to the context of this
// thread via ContextSwap() and processed independantly in
// the fullness of time
// : It's possible to have processing delays in this context.
// But keep in mind P2PeerMsg's may keep building up
//
void
P2PeerExplorer::RunExpump ( )
{
// Locals
DWORD dwResult;
P2PsigID nSigID;
// Because this is all problematic we
try
{
// Latencies
ThrowP2Pevent();
// Pump messages through the P2PeerSys, P2PeerCon and
// P2PeerMsg_MAP's until terminated and exhausted
while ( (dwResult=PumpP2Pmsg(8000,nSigID)) != 0 )
{
// Immediate closure signalled
if ( nSigID == P2PsigPump_CLOSE )
break;
// Idle closure signalled
if ( nSigID == P2PsigPump_CLOSEONIDLE )
break;
// Wakeup signalled
if ( nSigID == P2PsigPump_WAKEUP )
continue;
}
}
// Exceptions
catch ( P2Pevent *pEVT)
{
pEVT->Advice("P2Pexpump() has terminated" )
->Cancel();
}
catch ( ... )
{
EVERR->MODULE
->Message("Last resort exception of unknown type, "
"P2Pexpump has terminated" )
->Cancel();
}
// Tidy up, and
}
//
// Performs post P2PmsgExp destruction processing
// NOTES: Specialise for external notifications etc
//
void
P2PeerExplorer::PostDestroyExpump ( P2PexpumpID )
{
}
///////////////////////////////////////////////////////////////////////
// P2PeerCon management
//
// Posts P2PeerCon object to this hub
// NOTES: Control over life cycle of posted P2PeerCon object is
// assumed. Either service or client type objects may be
// posted
//
//
// Parameters: P2PeerCon *pCon
// Connection object to be posted
//
// Returns: P2PeerCon*
// Completion summary flag
// 0.. P2PeerCon object posted
// Otherwise sequence failed and passed P2PeerCon
// connection object is returned
P2PeerCon*
P2PeerExpump::PostP2PeerCon ( P2PeerCon *pCon )
{
// Introduce locals
P2Paddr oP2PaddrCon = pCon -> GetP2Paddress();
P2PsafeCS oSafeCS = m_oCSectionExpump;
if ( g_bP2Pmsg_AssertValid )
pCon -> AssertValid ( );
// Environmental
// NOTES: P2PeerCon objects may only be posted to operational,
// P2PeerExpump's negates risk of dead connections.
if ( m_nExpumpID == NULL )
{
EVERR->MODULE->AFPcon(pCon)
->Message(L"P2PeerExpump(%s) is not operational, "
L"P2PeerCon(%s) not posted"
, oP2PaddrCon.c_wstr()
, GetP2PaddrHub().c_wstr() )
->Advice (L"Perform SpawnExpump() or CreateExpump() before PostP2PeerCon()")
->Advice (L"Hub has failed?")
->Display()->SetLast();
return pCon;
}
// Iterate through P2PeerCon list
// NOTES: Trap duplicates. Null identification addresses
// are special
P2PeerCon *pConEnum = 0;
while ( EnumP2PexpCon(m_nExpumpID,&pConEnum) )
{
if ( pConEnum->GetP2Paddress() != oP2PaddrCon )
continue;
EVERR->MODULE->AFPcon(pCon)
->Message(L"P2PeerCon[%s] instance already exists "
L"within P2PmsgHub[%s]"
, oP2PaddrCon.c_wstr()
, GetP2PaddrHub().c_wstr() )
->SetLast();
return pCon;
}
// Tidy up, and
pCon -> m_pP2PeerTarget = this;
PostP2PexpCon ( m_nExpumpID, pCon );
PostP2Pmsg ( oP2PaddrCon, CN_P2PeerCon, P2P_Startup
, pCon, (P2PeerMsg *)0, m_nExpumpID );
return (P2PeerCon *)0;
}
//
// Signals nominated P2PeerCon object
//
//
// Parameters: P2PaddrSTR strP2Paddr
// Address domain defining of P2PeerCon objects to
// be signalled
//
// P2PsigID nSigID
// Signal to be posted
//
// void *pvData = 0
// Signal data
//
// int iDataSize = 0
// Signal data size
//
// Returns BOOL
// Signalled event summary
// TRUE... Located
// FALSE.. Do not exist
BOOL
P2PeerExpump::ConSignal ( P2PaddrSTR strP2PaddrTP
, P2PsigID nSigID, void *pvData, int iDataSize )
{
// Introduce locals
P2Paddr oP2Paddr = strP2PaddrTP;
BOOL bResult = FALSE;
P2PsafeCS oSafeCS = m_oCSectionExpump;
ASSERT(m_nExpumpID>0);
// Iterate through P2PeerCon'nections list
P2PeerCon *pConEnum = 0;
while ( EnumP2PexpCon(m_nExpumpID,&pConEnum) )
{
if ( !oP2Paddr.IsMapped(pConEnum->GetP2Paddress()) )
continue;
bResult = TRUE;
pConEnum -> Signal ( nSigID, pvData, iDataSize );
SwitchToThread ( );
}
// Tidy up and
return bResult;
}
//
// Checks if P2PeerCon connection object exists for passed
// P2PaddrSTR
//
//
// Parameters: P2PaddrSTR strP2Paddr
// Connection identification code to be checked
//
// Returns BOOL
// Existence summary
// TRUE... Exists
// FALSE.. Does not exist
BOOL
P2PeerExpump::ConExists ( P2PaddrSTR strP2Paddr )
{
// Locals
P2PsafeCS oSafeCS = m_oCSectionExpump;
// Iterate through P2PeerCon list
P2PeerCon *pCon = 0;
while ( EnumP2PexpCon(m_nExpumpID,&pCon) )
{
if ( pCon->GetP2Paddress() == strP2Paddr )
return TRUE;
}
// Nope, does not exist
return FALSE;
}
//
// Queries hub for nominated P2PeerCon
// NOTES: Designed for use from P2PeerMsg handlers whereby
// configuration issues require access to underlying
// P2PeerCon'nections
// : P2PeerCon pointers are transitory and as such should
// only retained for the life of a P2PeerMsg handler
//
//
// Parameters: P2PaddrSTR strP2Paddr
// Connection idetification
//
// SafeP2PeerCon& rSafeCon
// Safe container
//
// Returns: bool
// Query summary
// true... Located
// false.. Not located
bool
P2PeerExpump::ConQuery ( P2PaddrSTR strP2Paddress, SafeP2PeerCon& rSafeCon )
{
// Locals
P2PsafeCS oSafeCS = m_oCSectionExpump;
// Iterate through P2PeerCon list
P2PeerCon *pCon = 0;
while ( EnumP2PexpCon(m_nExpumpID,&pCon) )
{
if ( pCon->GetP2Paddress() == strP2Paddress )
{
rSafeCon = pCon;
return true;
}
}
// Nope, does not exist
rSafeCon = (P2PeerCon *)0;
return false;
}
///////////////////////////////////////////////////////////////////////
// P2PeerMsg operations
//
// Posts P2PeerMsg to this P2Pexpump for asynchronous delivery
// NOTES: Always posted to the P2Pexpump for dedicated delivery
// : P2Pexpump must be operational otherwise an exception
// is thrown
//
//
// Parameters: P2PeerMsg *pMsg
// Message to be posted to this hub
// NOTES: Hub assumes control over message life cycle.
//
// Returns: P2PeerMsg*
// Message posted flag
// 0.. Posted OK
//
P2PeerMsg*
P2PeerExpump::PostP2PeerMsg ( P2PeerMsg *pMsg )
{
// To be sure, to be sure
if ( g_bP2Pmsg_AssertValid )
pMsg -> AssertValid ( );
// Delegate
return PostP2Pexp ( pMsg, m_nExpumpID );
}
///////////////////////////////////////////////////////////////////////
// P2PeerMsg operations
//
// Routes P2PeerMsg's through P2PeerCon objects posted to this P2Pexpump
// NOTES: Local P2PeerMsg routing previously handled
// : Referenced from P2Pwin32, PostP2PeerMsg() for local hub
// routing
//
//
// Parameters: P2PeerMsg *pMsg
// Message to be routed.
//
// Returns: msgRESULT
// Routing result
msgRESULT
P2PeerExplorer::RouteP2PeerMsg ( P2PeerMsg *pMsg )
{
// To be sure, to be sure
// NOTES: Local routing must already have been performed
P2Paddr oP2PaddrHub = GetP2PaddrHub();
ASSERT(pMsg->GetDestin()!=oP2PaddrHub);
P2Paddr oP2PaddrSrc = pMsg -> GetSource ( );
P2Paddr oP2PaddrDst = pMsg -> GetDestin ( );
// From the LOCAL COPY, never from pMsg->GetDestin() directly. GetDestin()
// hands back a pointer into the message's wide-string store, and on the
// UTF-16-store platforms (everything but Win32) that is not the store at
// all: P3PmsgData::c_wstr() widens through p2p_wstr_from_store(), whose
// result lives in a 16 slot thread-local RING that the very next handful
// of c_wstr()/c_name() calls recycles (Platform/p2pstr.h:629-651). This
// function then makes far more than sixteen such calls - one per
// connection in the loop below, plus RouteP2PeerMsgPeek's own - before it
// compares the address, so by then the cached pointer read back as some
// unrelated recycled value ("Dst", "Src" - the field NAMES, from
// P3Pmsg::c_name). Every comparison below failed, every message the
// Explorer produced fell through to the undeliverable path, and the
// Explorer answered nothing. On Win32 p2p_wstr_from_store() returns the
// store pointer unchanged, which is why this only ever bit on Linux.
P2PaddrSTR pP2PaddrMsg = oP2PaddrDst.c_wstr ( );
// Facilitate message peeking
// NOTES: Application may wish to observe outgoing messages
msgRESULT msgResult = RouteP2PeerMsgPeek ( pMsg );
if ( msgResult != msgCONTINUE )
return msgResult;
// Parent routing
// NOTES: P2PeerMsg sourced in this P2PeerExp environment and
// destined for P2PeerHub parent
// : Interception occurs prior to attempted XCid routing
if ( oP2PaddrHub == oP2PaddrSrc &&
!oP2PaddrHub.IsRable(oP2PaddrDst) )
{
P2PexpumpContextSwap ( pMsg );
return msgSWAP;
}
// P2PeerCon routing
// NOTES: Connections manage their own routing logic
P2PeerCon *pCon = 0;
while ( EnumP2PexpCon(m_nExpumpID,&pCon) )
{
const P2Paddr& oP2PaddrCon = pCon -> GetP2Paddress();
// NOTES: EnumP2PexpCon() is deliberately unfiltered - it is the
// "every connection on this expump" primitive, and the six
// loops in this file each apply their own predicate to it.
// A listener is not a routing candidate, so the test
// belongs at this call site and not in the enumerator
if ( pCon->GetMode() == P2PeerCon_SERVICE )
{
continue;
}
// Immediate
if ( oP2PaddrCon == pP2PaddrMsg )
return P2PeerContextSwap ( pCon, pMsg );
// Ejection
// NOTES: P2PeerMsg sourced from XCid (child of P2PeerExp) and destined
// child of associated P2PeerHub. Thus ejected from from P2PeerExp
// environment
// : P2PeerMsg cannot be exchanged between XCid's
if ( ( oP2PaddrCon == oP2PaddrSrc ||
oP2PaddrCon.IsChild(oP2PaddrSrc) ) &&
oP2PaddrHub.IsChild(oP2PaddrDst) )
{
return SwapP2PexpContext ( pMsg );
}
// Children
// NOTES: P2PeerCon is our network child, and
// P2PeerMsg is routable down through P2PeerCon
if ( oP2PaddrHub.IsChild(oP2PaddrCon) &&
oP2PaddrCon.IsRable(pP2PaddrMsg) )
return P2PeerContextSwap ( pCon, pMsg );
// Parent
// NOTES: P2PeerCon is our network parent, and
// P2PeerMsg is not routable down through this P2Peer
if ( oP2PaddrCon.IsChild(oP2PaddrHub) &&
!oP2PaddrHub.IsRable(pP2PaddrMsg) )
{
ASSERT(0);
return P2PeerContextSwap ( pCon, pMsg );
}
}
// Undeliverable
return P2PeerTarget::RouteP2PeerMsg ( pMsg );
}
//
// Peek P2PeerMsg handler
// NOTES: Called immediately prior to commencement of routing a
// P2PeerMsg through a P2PeerTarget derived object.
// : Override this default implementation for special processing etc
//
//
// Parameters: P2PeerMsg *pMsg
// Completed processing message
//
// Returns: msgRESULT
// Routing summary
// msgHANDLED... Message handled
// msgCONTINUE.. Continue routing
//
msgRESULT
P2PeerExplorer::RouteP2PeerMsgPeek ( P2PeerMsg *pMsg )
{
// Ejection routing
// NOTES: P2PeerMsg received from XCid (child of hub) and destined
// child of this hub
// : P2PeerMsg cannot be exchanged between EXid's
//LPCTSTR lpszName=pMsg->c_name();
if ( GetCurrentThreadId() == m_nExpumpID )
{
P2Paddr oP2Paddr( GetP2PaddrHub(), L"XC*" ); //TODO:LJM This should persist
P2Paddr oP2PaddrSource = pMsg->GetSource();
P2Paddr oP2PaddrDestin = pMsg->GetDestin();
P2Paddr oP2PaddrThis = GetP2PaddrHub();
if ( pMsg->Map_MatchName(L"P2Pexpump*") )
oP2Paddr.c_hopname(0);
if ( oP2Paddr.IsMapped(pMsg->GetSource()) &&
GetP2PaddrHub().IsChild(pMsg->GetDestin()) )
return P2PexpumpContextSwap ( pMsg );
}
return msgCONTINUE;//msgResult;
}
//
// Swaps the processing context for the current P2PeerMsg to
// the encapsulated P2PeerExpump
// NOTES: Implementation requires the SwapP2PmsgContext() return
// code to be immediately propagated backwards as stack
// unwinds
//
// Parameters: P2PeerMsg *pMsg
// Message whose processing context is to be swapped
// to P2Pexpump context for hub
//
// Returns: mapRESULT
// Context swap object (mapSWAP)
//
msgRESULT
P2PeerExplorer::P2PexpumpContextSwap ( P2PeerMsg *pMsg )
{
if ( m_nExpumpID > 0 )
return SwapP2PexpContext ( pMsg );
// Undeliverable
// NOTES: Deliver P2PeerMsg exception back to source
P2Pevent *pEVT =
EVERR->MODULE->AFPmsg(pMsg)
// WIDE arguments need the WIDE overload - see P2PeerTarget.cpp:1476.
->Message(L"Message[%s] from [%s] not deliverable to [%s]"
, pMsg->c_name(), pMsg->GetSource(), pMsg->GetDestin() )
->Advice (L"P2PeerExpump has shutdown" )
->Advice (L"Bad destination address [%s]", pMsg->GetDestin() )
->Group("P2P");
PostP2PeerMsg ( pMsg->ExceptionFactory(pEVT) );
pEVT ->Cancel();
return msgHANDLED;
}
//
// Peek P2PeerMsg handler
// NOTES: Called immediately prior to commencement of routing a
// P2PeerMsg through a P2PeerTarget derived object.
// : Override this default implementation for special processing etc
//
//
// Parameters: P2PeerMsg *pMsg
// Completed processing message
//
// Returns: msgRESULT
// Routing summary
// msgHANDLED... Message handled
// msgCONTINUE.. Continue routing
//
msgRESULT
P2PeerExplorer::PeekP2PeerMsg ( P2PeerMsg *pMsg )
{
P2Paddr oP2Paddr( GetP2PaddrHub(), L"XC*" ); //TODO:LJM This should persist
if ( pMsg->Map_MatchName(L"P2PexpumpHub") && GetP2PaddrHub()==L"CEX")
oP2Paddr.c_hopname(0);
// Interception routing
// NOTES: P2PeerMsg received from child of hub and destined for XCid
// : P2PeerMsg cannot be exchanged between XCid's
if ( GetCurrentThreadId() != m_nExpumpID )
{
P2Paddr oP2PaddrSource = pMsg->GetSource();
P2Paddr oP2PaddrDestin = pMsg->GetDestin();
P2Paddr oP2Paddr_CEX_XCn(L"CEX.XC*");
P2Paddr oP2Paddr_CEX_XC0(L"CEX.XC0");
ASSERT(oP2Paddr_CEX_XCn.IsMapped(oP2Paddr_CEX_XC0));