1- using System . Linq ;
1+ using System . Linq ;
22using System . Threading . Tasks ;
33using Microsoft . EntityFrameworkCore ;
44using SW . Bitween . Domain ;
5- using SW . Bitween . Domain . Accounts ;
65using SW . Bitween . Model ;
76using SW . PrimitiveTypes ;
87
98namespace SW . Bitween . Resources . Xchanges
109{
10+ /// <summary>
11+ /// Retries a selection of exchanges — either a hand-picked list of ids or everything a filter
12+ /// matches. Returns what it did: <see cref="BulkRetryPreview"/> answers the same question
13+ /// beforehand, so the caller can show it and be sure the two agree.
14+ /// </summary>
1115 [ HandlerName ( "bulkretry" ) ]
1216 public class BulkRetry : ICommandHandler < XchangeBulkRetry , object >
1317 {
1418 private readonly BitweenDbContext _dbContext ;
19+ private readonly RequestContext _requestContext ;
1520 private readonly XchangeService _xchangeService ;
1621
1722
18- public BulkRetry ( BitweenDbContext dbContext , XchangeService xchangeService )
23+ public BulkRetry ( BitweenDbContext dbContext , RequestContext requestContext ,
24+ XchangeService xchangeService )
1925 {
2026 _dbContext = dbContext ;
27+ _requestContext = requestContext ;
2128 _xchangeService = xchangeService ;
2229 }
2330
2431 public async Task < object > Handle ( XchangeBulkRetry request )
2532 {
26- var scheduledIds = await _dbContext . Set < DelayedRetry > ( )
27- . Where ( d => request . Ids . Contains ( d . Id ) )
28- . Select ( d => d . Id )
29- . ToListAsync ( ) ;
33+ await _requestContext . EnsurePermission ( _dbContext , Model . Permissions . Exchanges . Operate ) ;
3034
31- var xchanges = await _dbContext . Set < Xchange > ( )
32- . Where ( c => request . Ids . Contains ( c . Id ) && ! scheduledIds . Contains ( c . Id ) ) . AsNoTracking ( )
33- . ToListAsync ( ) ;
35+ var prepared = await new BulkRetryPlanner ( _dbContext ) . Prepare ( request ) ;
36+ var plan = prepared . Plan ;
3437
35- foreach ( var xchange in xchanges )
38+ if ( plan . OverLimit )
39+ throw new SWValidationException ( "TOO_MANY" ,
40+ $ "{ plan . Selected : n0} exchanges is more than the { BulkRetryPlanner . Limit } this can retry " +
41+ "in one go. Narrow the filter and retry the rest after." ) ;
42+
43+ var retried = 0 ;
44+ foreach ( var id in prepared . Targets )
3645 {
37- var inputFileData = await _xchangeService . GetFile ( xchange . Id , XchangeFileType . Input ) ;
38- var xchangeFile = new XchangeFile ( inputFileData , xchange . InputName ) ;
46+ var xchange = await _dbContext . Set < Xchange > ( ) . AsNoTracking ( )
47+ . FirstOrDefaultAsync ( x => x . Id == id ) ;
48+ if ( xchange == null )
49+ {
50+ plan . Skipped . Add ( new XchangeRetrySkip
51+ {
52+ Id = id ,
53+ Reason = "This exchange no longer exists."
54+ } ) ;
55+ continue ;
56+ }
57+
3958 var subscription = await _dbContext . Subscriptions ( )
4059 . FirstOrDefaultAsync ( s => s . Id == xchange . SubscriptionId ) ;
41-
42- if ( request . Reset )
60+
61+ if ( request . Reset && subscription == null )
62+ {
63+ // Reported per exchange rather than thrown. Throwing meant one exchange whose
64+ // subscription had since been deleted took the whole selection down with it,
65+ // and the caller could not tell which one.
66+ plan . Skipped . Add ( new XchangeRetrySkip
67+ {
68+ Id = id ,
69+ Reason = "Its properties cannot be reset: the subscription no longer exists."
70+ } ) ;
71+ continue ;
72+ }
73+
74+ // The tolerant read, not GetFile: a retry re-sends the original input, so an
75+ // exchange whose input has been deleted or expired cannot be retried — and one of
76+ // those in a selection of five hundred must not take the other 499 with it.
77+ var xchangeFile = await _xchangeService . ReadInputFile ( xchange ) ;
78+ if ( xchangeFile == null )
79+ {
80+ plan . Skipped . Add ( new XchangeRetrySkip
81+ {
82+ Id = id ,
83+ Reason = "Its input document could not be read, so there is nothing to re-send."
84+ } ) ;
85+ continue ;
86+ }
87+
88+ try
4389 {
44- if ( subscription == null )
45- throw new SWValidationException ( "SUBSCRIPTION_NOT_FOUND" ,
46- "Cant reset properties, subscription doesnt exist anymore" ) ;
47- await _xchangeService . CreateXchange ( subscription , xchange , xchangeFile ,
48- manualRetry : true ) ;
90+ if ( request . Reset )
91+ {
92+ await _xchangeService . CreateXchange ( subscription , xchange , xchangeFile ,
93+ manualRetry : true ) ;
94+ }
95+ else
96+ {
97+ // Null when the subscription has since been deleted, which a document-only
98+ // exchange also has from the start. The single-exchange retry has always allowed
99+ // for it; without the same here, one such id in a selection threw and took the
100+ // whole bulk retry down with it.
101+ await _xchangeService . CreateXchange ( xchange , xchangeFile , subscription ? . WorkGroup ,
102+ manualRetry : true ) ;
103+ }
49104 }
50- else
105+ catch ( SWValidationException e ) when (
106+ e . Validations . Any ( v => v . Key == "ALREADY_RETRIED" ) )
51107 {
52-
53- // Null when the subscription has since been deleted, which a document-only
54- // exchange also has from the start. The single-exchange retry has always allowed
55- // for it; without the same here, one such id in a selection threw and took the
56- // whole bulk retry down with it.
57- await _xchangeService . CreateXchange ( xchange , xchangeFile , subscription ? . WorkGroup ,
58- manualRetry : true ) ;
108+ // The planner resolves every selection to the end of its chain, so this is not
109+ // reachable by choosing badly — it means someone retried this attempt in the
110+ // moment between the plan being worked out and it being carried out. Reported
111+ // like the other per-exchange refusals rather than thrown, so one racing
112+ // operator cannot cancel another's whole recovery.
113+ plan . Skipped . Add ( new XchangeRetrySkip
114+ {
115+ Id = id ,
116+ Reason = "It was retried by someone else a moment ago. Its own retry can be retried instead."
117+ } ) ;
118+ continue ;
59119 }
120+
121+ retried ++ ;
60122 }
61123
62124 await _dbContext . SaveChangesAsync ( ) ;
63125
64- return null ;
126+ plan . WillRetry = retried ;
127+ return plan ;
65128 }
66129 }
67- }
130+ }
0 commit comments