public class Developer :
Entity<DeveloperId, DeveloperUniqueId>
{
public string Name { get; set; }
public DeveloperStatus Status { get; set; }
}
public enum DeveloperStatus
{
New = 0,
Accepted = 1,
Rejected = 2,
}
public abstract class Entity<T1, T2>
{
public T1 Id { get; set; }
public T2 UniqueId { get; set; }
public int Version { get; set; }
}
public class DeveloperId
{
public int Value { get; set; }
public DeveloperId(int value)
{
Value = value;
}
public DeveloperId()
{
Value = default;
}
}
public class DeveloperUniqueId : BaseUniqueId<Guid>
{
public Guid Value { get; set; }
public override string ValueInString()
{
return Value.ToString();
}
public DeveloperUniqueId()
{
Value = Guid.NewGuid();
}
public DeveloperUniqueId(Guid value)
{
Value = value;
}
}
public abstract class BaseUniqueId<T>
{
public abstract string ValueInString();
public AggregateKey GetAggregateKey()
{
return new AggregateKey
{
Id = ValueInString(),
};
}
public BaseUniqueId()
{
}
}
public class AggregateKey
{
public string Id { get; set; }
public static readonly AggregateKey Empty = new AggregateKey();
public override string ToString()
{
return Id;
}
}
public class DeveloperIds
{
public DeveloperId CreatedId { get; set; }
public DeveloperUniqueId UniqueId { get; set; }
[JsonIgnore]
public IdsStatus ExStatus { get; set; }
public string Status
{
get
{
return ExStatus.ToString();
}
}
}
public enum IdsStatus
{
CreateIdReturned = 0,
DudeYouCantReturnCreatedIdWhenYouAreEventSourcing = 1
}
public class Developer :
Entity<DeveloperId, DeveloperUniqueId>
{
public string Name { get; set; }
public DeveloperStatus Status { get; set; }
public DeveloperIds Ids()
{
if (this.Id != null && this.Id.Value != default)
return new DeveloperIds()
{
UniqueId = this.UniqueId,
CreatedId = this.Id
};
else
return new DeveloperIds()
{
UniqueId = this.UniqueId,
CreatedId = this.Id,
ExStatus = IdsStatus.DudeYouCantReturnCreatedIdWhenYouAreEventSourcing
};
}
}
public interface IDeveloperRepository
{
Task<IReadOnlyList<Developer>> GetCollectionAsync(
FilterDeveloperStatus filtrer);
Task<Developer> GetByIdAsync(DeveloperId id);
Task<Developer> GetByIdAsync(DeveloperUniqueId id);
Task<DeveloperIds> SubmitAsync(Developer developer);
Task<bool> SaveAcceptenceAsync(DeveloperUniqueId id);
Task<bool> SaveAcceptenceAsync(DeveloperId id);
Task<bool> SaveRejectionAsync(DeveloperUniqueId id);
Task<bool> SaveRejectionAsync(DeveloperId id);
}
public interface IZEsDeveloperRepository : IDeveloperRepository
{
}
public enum FilterDeveloperStatus
{
All = 100,
New = 0,
Accepted = 1,
Rejected = 2,
}
public abstract class BaseResponse
{
public ResponseStatus Status { get; set; }
public bool Success { get; set; }
public string Message { get; set; }
public List<string> ValidationErrors { get; set; }
protected BaseResponse()
{
ValidationErrors = new List<string>();
Success = true;
Status = ResponseStatus.Success;
}
protected BaseResponse(string message = null)
{
ValidationErrors = new List<string>();
Success = true;
Message = message;
Status = ResponseStatus.Success;
}
protected BaseResponse(string message, bool success)
{
ValidationErrors = new List<string>();
Success = success;
Message = message;
}
protected BaseResponse(ResponseStatus status)
{
ValidationErrors = new List<string>();
Success = status != ResponseStatus.Success;
Status = status;
}
protected BaseResponse(ValidationResult validationResult)
{
ValidationErrors = new List<String>();
Success = validationResult.Errors.Count < 0;
foreach (var item in validationResult.Errors)
{
ValidationErrors.Add(item.ErrorMessage);
}
if (!Success)
Status = ResponseStatus.ValidationError;
else
Status = ResponseStatus.Success;
}
public string StatusInfo
{
get
{
return Status.ToString();
}
}
public WhatHTTPCodeShouldBeRetruned WhatHTTPCodeToBeRetruned
{
get
{
if (this.Status == ResponseStatus.BussinesLogicError)
return WhatHTTPCodeShouldBeRetruned.Forbid;
if (this.Status == ResponseStatus.NotFoundInDataBase)
return WhatHTTPCodeShouldBeRetruned.NotFound;
if (this.Status == ResponseStatus.ValidationError ||
this.Status == ResponseStatus.BadQuery ||
this.Status == ResponseStatus.ConcurrencyOlderVersionSendedWhenNewerIsInEventStore)
return WhatHTTPCodeShouldBeRetruned.BadRequest;
if (!this.Success)
return WhatHTTPCodeShouldBeRetruned.BadRequest;
else
return WhatHTTPCodeShouldBeRetruned.Ok;
}
}
}
public class DeveloperInListViewModel
{
public string Name { get; set; }
public int Status { get; set; }
}
public class GetAllDevelopersQuery
: IRequest<GetAllDevelopersQueryHandlerResponse>
{
public FilterDeveloperStatus Filter { get; set; }
public QueryWitchDataBase queryWitchDataBase { get; set; }
}
public class GetAllDevelopersQueryHandlerResponse : BaseResponse
{
public List<DeveloperInListViewModel> List { get; }
public GetAllDevelopersQueryHandlerResponse
(List<DeveloperInListViewModel> lisy) : base()
{
List = lisy;
}
public GetAllDevelopersQueryHandlerResponse(ValidationResult validationResult)
: base(validationResult)
{ }
public GetAllDevelopersQueryHandlerResponse(string message)
: base(message)
{ }
public GetAllDevelopersQueryHandlerResponse(string message, bool success)
: base(message, success)
{ }
public GetAllDevelopersQueryHandlerResponse(ResponseStatus status) : base(status)
{ }
}
public class GetAllDevelopersQueryHandler
:
IRequestHandler<GetAllDevelopersQuery,
GetAllDevelopersQueryHandlerResponse>
{
private readonly IDeveloperRepository _Repository;
private readonly IMapper _mapper;
private readonly IZEsDeveloperRepository _zEsRepository;
public GetAllDevelopersQueryHandler(IDeveloperRepository callRepository,
IZEsDeveloperRepository ZEscallRepository,
IMapper mapper)
{
_mapper = mapper;
_zEsRepository = ZEscallRepository;
_Repository = callRepository;
}
public async Task<GetAllDevelopersQueryHandlerResponse>
Handle(GetAllDevelopersQuery request, CancellationToken cancellationToken)
{
IReadOnlyList<Developer> listresult;
if (request.queryWitchDataBase == QueryWitchDataBase.WithEventSourcing)
listresult = await _zEsRepository.GetCollectionAsync(request.Filter);
else
listresult = await _Repository.GetCollectionAsync(request.Filter);
if (listresult == null)
return new GetAllDevelopersQueryHandlerResponse
(ResponseStatus.NotFoundInDataBase);
var allordered = listresult.OrderBy(x => x.Id);
var allmaped = _mapper.Map<List<DeveloperInListViewModel>>(listresult);
return new GetAllDevelopersQueryHandlerResponse(allmaped);
}
}
public class SubmitDeveloperCommand
: IRequest<SubmitDeveloperCommandResponse>
{
public string Name { get; set; }
[JsonIgnore]
public DeveloperUniqueId UniqueId { get; }
[JsonIgnore]
public DeveloperStatus Status { get; }
public SubmitDeveloperCommand()
{
UniqueId = new DeveloperUniqueId();
Status = DeveloperStatus.New;
}
}
public class SubmitDeveloperCommandResponse : BaseResponse
{
public DeveloperIds DeveloperIds { get; set; }
public SubmitDeveloperCommandResponse(DeveloperIds ids)
: base()
{
DeveloperIds = ids;
}
public SubmitDeveloperCommandResponse(ValidationResult validationResult)
: base(validationResult)
{ }
public SubmitDeveloperCommandResponse(string message)
: base(message)
{ }
public SubmitDeveloperCommandResponse(string message, bool success)
: base(message, success)
{ }
public SubmitDeveloperCommandResponse(ResponseStatus status) : base(status)
{ }
}
public class SubmitDeveloperCommandValidator
: AbstractValidator<SubmitDeveloperCommand>
{
public SubmitDeveloperCommandValidator()
{
RuleFor(c => c.Name)
.MinimumLength(1)
.MaximumLength(100)
.WithMessage("{PropertName} Length is beewten 1 and 100");
}
}
public class SubmitDeveloperCommandHandler
:
IRequestHandler<SubmitDeveloperCommand,
SubmitDeveloperCommandResponse>
{
private readonly IDeveloperRepository _repository;
private readonly IMapper _mapper;
public SubmitDeveloperCommandHandler(IDeveloperRepository Repository,
IMapper mapper)
{
_mapper = mapper;
_repository = Repository;
}
public async Task<SubmitDeveloperCommandResponse>
Handle(SubmitDeveloperCommand request, CancellationToken cancellationToken)
{
var validator = new SubmitDeveloperCommandValidator();
var validatorResult = await validator.ValidateAsync(request);
if (!validatorResult.IsValid)
return new SubmitDeveloperCommandResponse(validatorResult);
var developer = _mapper.Map<Developer>(request);
developer.Status = DeveloperStatus.New;
var ids = await _repository.SubmitAsync(developer);
return new SubmitDeveloperCommandResponse(ids);
}
}
public class RejectDeveloperCommandHandler
:
IRequestHandler<RejectDeveloperCommand, RejectDeveloperCommandResponse>
{
private readonly IDeveloperRepository _callRepository;
private readonly IMapper _mapper;
public RejectDeveloperCommandHandler(IDeveloperRepository callRepository,
IMapper mapper)
{
_callRepository = callRepository;
_mapper = mapper;
}
public async Task<RejectDeveloperCommandResponse> Handle
(RejectDeveloperCommand request, CancellationToken cancellationToken)
{
var developerUniqueId = _mapper.Map<DeveloperUniqueId>
(request.DeveloperUniqueId);
var developer = await _callRepository.GetByIdAsync(developerUniqueId);
if (developer.Status == DeveloperStatus.New)
await _callRepository.SaveRejectionAsync
(developerUniqueId);
else
{
return new RejectDeveloperCommandResponse("Can't Reject Developer " +
"that is already " + developer.Status.ToString(), false);
}
return new RejectDeveloperCommandResponse();
}
}
public class AcceptDeveloperCommandHandler
:
IRequestHandler<AcceptDeveloperCommand, AcceptDeveloperCommandResponse>
{
private readonly IDeveloperRepository _callRepository;
private readonly IMapper _mapper;
public AcceptDeveloperCommandHandler(IDeveloperRepository callRepository,
IMapper mapper)
{
_callRepository = callRepository;
_mapper = mapper;
}
public async Task<AcceptDeveloperCommandResponse> Handle
(AcceptDeveloperCommand request, CancellationToken cancellationToken)
{
var developerUniqueId = _mapper.Map<DeveloperUniqueId>
(request.DeveloperUniqueId);
var developer = await _callRepository.GetByIdAsync(developerUniqueId);
if (developer.Status == DeveloperStatus.New)
await _callRepository.SaveAcceptenceAsync
(developerUniqueId);
else
{
return new AcceptDeveloperCommandResponse("Can't Accept Developer " +
"that is already " + developer.Status.ToString(), false);
}
return new AcceptDeveloperCommandResponse();
}
}
public class ProfileMap : Profile
{
public ProfileMap()
{
CreateMap<int, DeveloperId>().ConstructUsing
(c => new DeveloperId(c));
CreateMap<DeveloperId, int>().ConstructUsing(c => c.Value);
CreateMap<Guid, DeveloperUniqueId>().ConstructUsing
(c => new DeveloperUniqueId(c));
CreateMap<DeveloperUniqueId, Guid>().ConstructUsing
(c => c.Value);
CreateMap<Developer, DeveloperViewModel>();
CreateMap<Developer, DeveloperInListViewModel>();
CreateMap<SubmitDeveloperCommand, Developer>();
CreateMap<EsSubmitDeveloperCommand, Developer>();
}
}
public static partial class LiteSmallConferenceInstaller
{
public static IServiceCollection AddLiteSmallConferenceCQRS
(this IServiceCollection services, IConfiguration Configuration)
{
services.AddAutoMapper(Assembly.GetExecutingAssembly());
services.AddMediatR(Assembly.GetExecutingAssembly());
return services;
}
}
public static partial class LiteSmallConferenceInstaller
{
public static IServiceCollection AddPersitenceDapperSQLite
(this IServiceCollection services, IConfiguration configuration)
{
services.AddAutoMapper(Assembly.GetExecutingAssembly());
var connection = configuration.
GetConnectionString("LiteSmallConferenceConnectionString");
var zEsConnection = configuration.
GetConnectionString("ZEsLiteSmallConferenceConnectionString");
services.AddTransient<ILiteSmallDBContext, LiteSmallDBContext>
(
(services) =>
{
var c =
new LiteSmallDBContext(connection);
return c;
}
);
services.AddTransient<IZEsLiteSmallDBContext, LiteSmallDBContext>
(
(services) =>
{
var c =
new LiteSmallDBContext(zEsConnection);
return c;
}
);
services.AddTransient<IDeveloperGetAllDoer, DeveloperGetAllDoer>();
services.AddTransient<IDeveloperGetByIdDoer, DeveloperGetByIdDoer>();
services.AddTransient<IDeveloperSaveAcceptanceDoer, DeveloperSaveAcceptanceDoer>();
services.AddTransient<IDeveloperSaveRejectionDoer, DeveloperSaveRejectionDoer>();
services.AddTransient<IDeveloperSubmitDoer, DeveloperSubmitDoer>();
services.AddTransient<IDeveloperRepository, DeveloperRepository>();
services.AddTransient<IZEsDeveloperRepository, ZEsDeveloperRepository>();
return services;
}
}
public class Startup
{
public IConfiguration Configuration { get; }
public Startup(IConfiguration configuration)
{
Configuration = configuration;
}
public void ConfigureServices(IServiceCollection services)
{
services.AddLiteSmallConferenceCQRS(Configuration);
services.AddPersitenceDapperSQLite(Configuration);
services.AddControllers();
services.AddCors(options =>
{
options.AddPolicy("Open",
builder => builder.AllowAnyOrigin().AllowAnyHeader().AllowAnyMethod());
});
services.AddSwaggerGen(c =>
{
c.AddSecurityDefinition("Bearer", new OpenApiSecurityScheme
{
Description = @"JWT Authorization header using the Bearer scheme. \r\n\r\n
Enter 'Bearer' [space] and then your token in the text input below.
\r\n\r\nExample: 'Bearer 12345abcdef'",
Name = "Authorization",
In = ParameterLocation.Header,
Type = SecuritySchemeType.ApiKey,
Scheme = "Bearer"
});
c.AddSecurityRequirement(new OpenApiSecurityRequirement()
{
{
new OpenApiSecurityScheme
{
Reference = new OpenApiReference
{
Type = ReferenceType.SecurityScheme,
Id = "Bearer"
},
Scheme = "oauth2",
Name = "Bearer",
In = ParameterLocation.Header,
},
new List<string>()
}
});
c.SwaggerDoc("v1", new OpenApiInfo
{
Version = "v1",
Title = "LiteSmallConference API",
});
});
}
public class DevelopersController : BaseLiteSmallConference
{
private readonly IMediator _mediator;
public DevelopersController(IMediator mediator)
{
_mediator = mediator;
}
[HttpPost("submit", Name = "submitdeveloper")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<ActionResult<DeveloperIds>> Submit
([FromBody] SubmitDeveloperCommand request)
{
var result = await _mediator.Send(request);
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return Ok(result.DeveloperIds);
}
[HttpGet("all/{filter}", Name = "getalldevelopers")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<ActionResult<DeveloperInListViewModel>>
GetAllDevelopers(int filter)
{
GetAllDevelopersQuery getAllDevelopersQuery = new GetAllDevelopersQuery()
{
Filter = (FilterDeveloperStatus)filter,
queryWitchDataBase = QueryWitchDataBase.NormalCQRS
};
var result = await _mediator.Send(getAllDevelopersQuery);
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return Ok(result.List);
}
[HttpGet("id/{id}", Name = "GetDeveloperId")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesDefaultResponseType]
public async Task<ActionResult<DeveloperViewModel>>
GetDeveloperId(int id)
{
var result = await _mediator.Send(
(new GetDeveloperQuery()
{
DeveloperId = new DeveloperId(id),
queryWitchDataBase = QueryWitchDataBase.NormalCQRS
}));
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return Ok(result.Developer);
}
[HttpGet("uniqueid/{uid}", Name = "GetDeveloperByUniqueId")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesDefaultResponseType]
public async Task<ActionResult<DeveloperViewModel>>
GetDeveloperByUniqueId(Guid uid)
{
var result = await _mediator.Send(
(new GetDeveloperQuery()
{
DeveloperUniqueId = new DeveloperUniqueId(uid),
queryWitchDataBase = QueryWitchDataBase.NormalCQRS
}));
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return Ok(result.Developer);
}
[HttpPost("reject", Name = "rejectcallforspeech")]
[ProducesResponseType(StatusCodes.Status204NoContent)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<ActionResult<int>> Reject
([FromBody] RejectDeveloperCommand rejectdeveloperCommand)
{
var result = await _mediator.Send(rejectdeveloperCommand);
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return NoContent();
}
[HttpPost("accept", Name = "acceptcallforspeech")]
[ProducesResponseType(StatusCodes.Status204NoContent)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<ActionResult<int>> Accept
([FromBody] AcceptDeveloperCommand acceptDeveloperCommand)
{
var result = await _mediator.Send(acceptDeveloperCommand);
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return NoContent();
}
}
public abstract class DomainEvent
{
public AggregateKey Key { get; set; }
public int Version { get; set; }
public DateTimeOffset TimeStamp { get; set; }
protected DomainEvent(DateTimeOffset occuredOn, int version)
{
TimeStamp = occuredOn;
Version = version;
}
protected DomainEvent(int version)
{
TimeStamp = AppTime.Now();
Version = version;
}
protected DomainEvent()
{
}
}
public class AggregateKey
{
public string Id { get; set; }
public static readonly AggregateKey Empty
= new AggregateKey();
public override string ToString()
{
return Id;
}
}
public class DevloperSubmitedEvent : DomainEvent
{
public DeveloperUniqueId UniqueId { get; set; }
public string Name { get; set; }
public DevloperSubmitedEvent(DeveloperUniqueId uniqueId,
string name)
{
UniqueId = uniqueId;
Name = name;
Key = uniqueId.GetAggregateKey();
}
}
public class DeveloperAcceptedEvent : DomainEvent
{
public DeveloperUniqueId UniqueId { get; set; }
public string Name { get; set; }
public DeveloperStatus Status { get; set; }
public DeveloperAcceptedEvent(DeveloperUniqueId uniqueId,
string name, DeveloperStatus status, int version)
{
UniqueId = uniqueId;
Status = status;
Name = name;
Version = version;
Key = uniqueId.GetAggregateKey();
}
}
public class DeveloperRejectedEvent : DomainEvent
{
public DeveloperUniqueId UniqueId { get; set; }
public string Name { get; set; }
public DeveloperStatus Status { get; set; }
public DeveloperRejectedEvent(DeveloperUniqueId uniqueId,
string name, DeveloperStatus status, int version)
{
UniqueId = uniqueId;
Status = status;
Name = name;
Version = version;
Key = uniqueId.GetAggregateKey();
}
}
public abstract class AggregateRoot
{
private readonly List<DomainEvent> _changes = new List<DomainEvent>();
public AggregateKey Key { get; protected set; }
public int Version { get; protected set; }
public override string ToString()
{
return Key.Id;
}
public IEnumerable<DomainEvent> GetUncommittedChanges()
{
lock (_changes)
{
return _changes.ToArray();
}
}
public void MarkChangesAsCommitted()
{
lock (_changes)
{
Version = Version + _changes.Count;
_changes.Clear();
}
}
public void LoadFromHistory(IEnumerable<DomainEvent> history)
{
foreach (var e in history)
{
if (e.Version != Version + 1)
throw new EventsOutOfOrderException(e.Key);
ApplyChange(e, false);
}
}
protected void ApplyChange(DomainEvent @event)
{
ApplyChange(@event, true);
}
private void ApplyChange(DomainEvent @event, bool isNew)
{
lock (_changes)
{
this.AsDynamic().Apply(@event);
if (isNew)
{
_changes.Add(@event);
}
else
{
Key = @event.Key;
Version++;
}
}
}
}
public class DevloperAggregate : AggregateRoot
{
public DeveloperUniqueId UniqueId { get; set; }
public string Name { get; set; }
public DeveloperStatus Status { get; set; }
public DevloperAggregate(Developer cc)
{
var c = new DevloperSubmitedEvent
(cc.UniqueId, cc.Name);
this.Key = c.UniqueId.GetAggregateKey();
ApplyChange(c);
}
public void Rejected(Developer cc)
{
var c = new DeveloperRejectedEvent
(cc.UniqueId, cc.Name, cc.Status, cc.Version);
this.Key = c.UniqueId.GetAggregateKey();
ApplyChange(c);
}
public void Accepted(Developer cc)
{
var c = new DeveloperAcceptedEvent
(cc.UniqueId, cc.Name, cc.Status, cc.Version);
ApplyChange(c);
}
private void Apply(DevloperSubmitedEvent e)
{
Status = DeveloperStatus.New;
Name = e.Name;
UniqueId = e.UniqueId;
Version = e.Version++;
this.Key = e.UniqueId.GetAggregateKey();
}
private void Apply(DeveloperAcceptedEvent e)
{
Status = e.Status;
Name = e.Name;
UniqueId = e.UniqueId;
Version = e.Version++;
this.Key = e.UniqueId.GetAggregateKey();
}
private void Apply(DeveloperRejectedEvent e)
{
Status = e.Status;
Name = e.Name;
UniqueId = e.UniqueId;
Version = e.Version++;
this.Key = e.UniqueId.GetAggregateKey();
}
public DevloperAggregate()
{
}
}
public interface ISessionForEventSourcing
{
void Add<T>(T aggregate)
where T : AggregateRoot;
T Get<T>
(AggregateKey id, int? expectedVersion = null)
where T : AggregateRoot;
void Commit();
}
public interface IEventRepository
{
void Save<T>
(T aggregate, int? expectedVersion = null)
where T : AggregateRoot;
T Get<T>
(AggregateKey aggregateId)
where T : AggregateRoot;
}
public interface IEventStore
{
void Save(DomainEvent @event);
List<DomainEvent> Get
(AggregateKey aggregateId, int fromVersion);
}
public interface IEventPublisher
{
void Publish<T>(T @event) where T : DomainEvent;
}
public class SessionForEventSourcing : ISessionForEventSourcing
{
private readonly IEventRepository _repository;
private readonly Dictionary<AggregateKey, AggregateDescriptor> _trackedAggregates;
public SessionForEventSourcing(IEventRepository repository)
{
if (repository == null)
throw new ArgumentNullException("repository");
_repository = repository;
_trackedAggregates = new Dictionary<AggregateKey, AggregateDescriptor>();
}
public void Add<T>(T aggregate) where T : AggregateRoot
{
try
{
if (!IsTracked(aggregate.Key))
_trackedAggregates.Add(aggregate.Key,
new AggregateDescriptor
{
Aggregate = aggregate,
Version = aggregate.Version
});
else if (_trackedAggregates[aggregate.Key].Aggregate != aggregate)
throw new ConcurrencyException(aggregate.Key);
}
catch (Exception ex)
{
throw;
}
}
public T Get<T>(AggregateKey id, int? expectedVersion = null) where T : AggregateRoot
{
try
{
//w pamięci śledzącej sprawdzamy czy nie próbujemy ze złą wersją dodać zdarzenie
//oraz jeśli mamy zdarzenie w pamięci śledzącej to je zwracamy
if (IsTracked(id))
{
var trackedAggregate = (T)_trackedAggregates[id].Aggregate;
if (expectedVersion != null && trackedAggregate.Version != expectedVersion)
throw new ConcurrencyException(trackedAggregate.Key);
return trackedAggregate;
}
//jeśli nie mamy w pamieci to odpytujemy repozytorium
var aggregate = _repository.Get<T>(id);
if (expectedVersion != null && aggregate.Version != expectedVersion)
throw new ConcurrencyException(id);
Add(aggregate);
//Dodajemy do pamięci śledzącej
return aggregate;
}
catch (Exception ex)
{
throw;
}
}
private bool IsTracked(AggregateKey id)
{
return _trackedAggregates.ContainsKey(id);
}
public void Commit()
{
try
{
foreach (var descriptor in _trackedAggregates.Values)
{
_repository.Save(descriptor.Aggregate, descriptor.Version);
}
_trackedAggregates.Clear();
}
catch (Exception ex)
{
throw;
}
}
private class AggregateDescriptor
{
public AggregateRoot Aggregate { get; set; }
public int Version { get; set; }
}
}
public class EventRepository : IEventRepository
{
private readonly IEventStore _eventStore;
private readonly IEventPublisher _publisher;
public EventRepository(IEventStore eventStore, IEventPublisher publisher)
{
if (eventStore == null)
throw new ArgumentNullException("eventStore");
if (publisher == null)
throw new ArgumentNullException("publisher");
_eventStore = eventStore;
_publisher = publisher;
}
public void Save<T>(T aggregate, int? expectedVersion = null) where T : AggregateRoot
{
if (expectedVersion != null && _eventStore.Get(
aggregate.Key, expectedVersion.Value).Any())
throw new ConcurrencyException(aggregate.Key);
var i = 0;
foreach (var @event in aggregate.GetUncommittedChanges())
{
if (@event.Key == AggregateKey.Empty)
@event.Key = aggregate.Key;
if (@event.Key == AggregateKey.Empty)
throw new AggregateOrEventMissingIdException(
aggregate.GetType(), @event.GetType());
i++;
//Zwiększanie wersji o 1 i dodanie daty
@event.Version = aggregate.Version + i;
@event.TimeStamp = DateTimeOffset.UtcNow;
//Zapisz do bazy zdarzeń i wyślij do kolejki
_eventStore.Save(@event);
_publisher.Publish(@event);
}
aggregate.MarkChangesAsCommitted();
}
public T Get<T>(AggregateKey aggregateId) where T : AggregateRoot
{
return LoadAggregate<T>(aggregateId);
}
private T LoadAggregate<T>(AggregateKey id) where T : AggregateRoot
{
var aggregate = AggregateFactory.CreateAggregate<T>();
var events = _eventStore.Get(id, -1);
if (!events.Any())
throw new AggregateNotFoundException(id);
aggregate.LoadFromHistory(events);
return aggregate;
}
}
public class Constants
{
public const string QUEUE_DEVELOPER_SUBMITED = "developers_submited";
public const string QUEUE_DEVELOPER_REJECTCED = "developers_rejected";
public const string QUEUE_DEVELOPER_ACCEPTED = "developers_accepted";
}
public static class DomainEventHelper
{
public static string WhatRabbitMQQueue(this DomainEvent @event)
{
string g = @event switch
{
DevloperSubmitedEvent => Constants.QUEUE_DEVELOPER_SUBMITED,
DeveloperAcceptedEvent => Constants.QUEUE_DEVELOPER_ACCEPTED,
DeveloperRejectedEvent => Constants.QUEUE_DEVELOPER_REJECTCED,
_ => throw new NotImplementedException(),
};
return g;
}
}
public class LiteSmallEventPublisher : IEventPublisher
{
private readonly ConnectionFactory connectionFactory;
public LiteSmallEventPublisher(IHostingEnvironment env)
{
connectionFactory = new ConnectionFactory();
var builder = new Microsoft.Extensions.Configuration.ConfigurationBuilder()
.SetBasePath(env.ContentRootPath)
.AddJsonFile("appsettings.json", optional: false, reloadOnChange: false)
.AddEnvironmentVariables();
builder.Build().GetSection("RabbitMqSetting").Bind(connectionFactory);
}
public void Publish<T>(T @event) where T : DomainEvent
{
using (IConnection conn = connectionFactory.CreateConnection())
{
using (IModel channel = conn.CreateModel())
{
var queue = @event.WhatRabbitMQQueue();
channel.QueueDeclare(
queue: queue,
durable: false,
exclusive: false,
autoDelete: false,
arguments: null
);
var body = Encoding.UTF8.GetBytes(JsonConvert.SerializeObject(@event));
channel.BasicPublish(
exchange: "",
routingKey: queue,
basicProperties: null,
body: body
);
}
}
}
}
BEGIN TRANSACTION;
CREATE TABLE IF NOT EXISTS "EventStore" (
"Id" INTEGER NOT NULL UNIQUE,
"Key" TEXT NOT NULL,
"AssemblyQualifiedName" TEXT NOT NULL,
"Value" TEXT NOT NULL,
"Version" INTEGER,
PRIMARY KEY("Id" AUTOINCREMENT)
);
COMMIT;
public class SqlLiteEventStore : IEventStore
{
private IEventStoreSQLiteContext _geekLemonContext;
public SqlLiteEventStore(IEventStoreSQLiteContext context)
{
_geekLemonContext = context;
}
public List<DomainEvent> Get(AggregateKey aggregateId, int fromVersion)
{
using var connection = new SqliteConnection
(_geekLemonContext.ConnectionString);
try
{
var r = connection.Query<EventTemp>
(@"SELECT Id,Key, Value, AssemblyQualifiedName, Version FROM EventSTORE
WHERE Key = @aggregateId and Version > @Version;", new
{
@aggregateId = aggregateId.Id,
@Version = fromVersion
});
List<DomainEvent> de = new List<DomainEvent>();
foreach (var item in r)
{
Assembly asm = typeof(DomainEvent).Assembly;
Type type = TypeRecon.ReconstructType(item.AssemblyQualifiedName, true, asm);
var domain = JsonConvert.
DeserializeObject(item.Value, type);
de.Add(domain as DomainEvent);
}
return de;
}
catch (Exception ex)
{
throw;
}
}
public void Save(DomainEvent @event)
{
using var connection = new SqliteConnection
(_geekLemonContext.ConnectionString);
try
{
var q = @"INSERT INTO EventSTORE(Key, Value,
AssemblyQualifiedName
,Version)
VALUES (@Key, @Value, @AssemblyQualifiedName,@Version);";
var result = connection.Execute(q, new
{
@Key = @event.Key.Id,
@Value = JsonConvert.SerializeObject(@event),
@AssemblyQualifiedName = @event.GetType().AssemblyQualifiedName,
@Version = @event.Version,
}
);
}
catch (Exception ex)
{
throw;
}
}
}
public static IServiceCollection
AddEventStoreSqlLite
(this IServiceCollection services,
IConfiguration configuration)
{
var connection = configuration.
GetConnectionString("EventStoreSQLiteConnectionString");
services.AddScoped<IEventStoreSQLiteContext, EventStoreSQLiteContext>
(
(services) =>
{
var c =
new EventStoreSQLiteContext(connection);
return c;
}
);
services.AddScoped<IEventStore, SqlLiteEventStore>();
return services;
}
public void ConfigureServices(IServiceCollection services)
{
....
services.AddLiteSmallConferenceCQRS(Configuration);
services.AddPersitenceDapperSQLite(Configuration);
services.AddEventStoreSqlLite(Configuration);
services.AddEventSourcing(Configuration);
services.AddControllers();
services.AddCors(options =>
{
options.AddPolicy("Open",
builder => builder.AllowAnyOrigin().AllowAnyHeader().AllowAnyMethod());
});
}
[Route("api/[controller]")]
[ApiController]
public class ZEsDeveloperController : BaseLiteSmallConference
{
private readonly IMediator _mediator;
public ZEsDeveloperController(IMediator mediator)
{
_mediator = mediator;
}
[HttpPost("submit", Name = "submitdeveloperEs")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<ActionResult<DeveloperIds>> Submit([FromBody] EsSubmitDeveloperCommand request)
{
var result = await _mediator.Send(request);
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return Ok(result.DeveloperIds);
}
[HttpGet("all/{filter}", Name = "getalldevelopersES")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<ActionResult<DeveloperInListViewModel>> GetAllDevelopers(int filter)
{
GetAllDevelopersQuery getAllDevelopersQuery = new GetAllDevelopersQuery()
{
Filter = (FilterDeveloperStatus)filter,
queryWitchDataBase = QueryWitchDataBase.WithEventSourcing
};
var result = await _mediator.Send(getAllDevelopersQuery);
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return Ok(result.List);
}
[HttpGet("uniqueid/{uid}", Name = "GetDeveloperByUniqueIdEs")]
[ProducesResponseType(StatusCodes.Status200OK)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesDefaultResponseType]
public async Task<ActionResult<DeveloperViewModel>> GetDeveloperByUniqueId(Guid uid)
{
var result = await _mediator.Send(
(new GetDeveloperQuery()
{
DeveloperUniqueId = new DeveloperUniqueId(uid),
queryWitchDataBase = QueryWitchDataBase.WithEventSourcing
}));
return Ok(result.Developer);
}
[HttpPost("reject", Name = "rejectdevelopers")]
[ProducesResponseType(StatusCodes.Status204NoContent)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<ActionResult<int>> Reject([FromBody] EsRejectDeveloperCommand rejectdeveloperCommand)
{
var result = await _mediator.Send(rejectCallForSpeechCommand);
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return NoContent();
}
[HttpPost("accept", Name = "acceptcallforspeech")]
[ProducesResponseType(StatusCodes.Status204NoContent)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
[ProducesResponseType(StatusCodes.Status400BadRequest)]
[ProducesResponseType(StatusCodes.Status403Forbidden)]
[ProducesResponseType(420)]
[ProducesResponseType(StatusCodes.Status404NotFound)]
public async Task<ActionResult<int>> Accept([FromBody] EsAcceptDeveloperCommand acceptDeveloperCommand)
{
var result = await _mediator.Send(acceptDeveloperCommand);
if (result.Status == ResponseStatus.BussinesLogicError)
return Forbid();
if (result.Status == ResponseStatus.NotFoundInDataBase)
return NotFound();
if (result.Status == ResponseStatus.ValidationError)
return BadRequest();
if (result.Status == ResponseStatus.BadQuery)
return BadRequest();
if (!result.Success)
return MethodFailure(result.Message);
return NoContent();
}
}
public class EsAcceptDeveloperCommand
: IRequest<EsAcceptDeveloperCommandResponse>
{
public Guid DeveloperUniqueId { get; set; }
public int Version { get; set; }
}
public class EsAcceptDeveloperCommandHandler
:
IRequestHandler<EsAcceptDeveloperCommand, EsAcceptDeveloperCommandResponse>
{
private readonly ISessionForEventSourcing _sessionForEventSourcing;
private readonly IMapper _mapper;
public EsAcceptDeveloperCommandHandler(
ISessionForEventSourcing sessionForEventSourcing, IMapper mapper)
{
_mapper = mapper;
_sessionForEventSourcing = sessionForEventSourcing;
}
public Task<EsAcceptDeveloperCommandResponse>
Handle(EsAcceptDeveloperCommand request,
CancellationToken cancellationToken)
{
var developerUniqueId = _mapper.Map<DeveloperUniqueId>
(request.DeveloperUniqueId);
var aggregateDevloper = _sessionForEventSourcing.Get<DevloperAggregate>
(developerUniqueId.GetAggregateKey(), null);
if (aggregateDevloper.Version > request.Version)
return Task.FromResult(new EsAcceptDeveloperCommandResponse
($@"You sended old version.
Yours {request.Version}.
In Event database :{aggregateDevloper.Version}", false));
if (aggregateDevloper.Status != DeveloperStatus.New)
return Task.FromResult(new EsAcceptDeveloperCommandResponse
($@"Devloper status in eventhistory is {aggregateDevloper.Status}.
Can't be Accepted", false));
var developer = _mapper.Map<Developer>(aggregateDevloper);
developer.Status = DeveloperStatus.Accepted;
aggregateDevloper.Accepted(developer);
_sessionForEventSourcing.Commit();
return Task.FromResult(new EsAcceptDeveloperCommandResponse());
}
}
public class EsRejectDeveloperCommandHandler
:
IRequestHandler<EsRejectDeveloperCommand, EsRejectDeveloperCommandResponse>
{
private readonly ISessionForEventSourcing _sessionForEventSourcing;
private readonly IMapper _mapper;
public EsRejectDeveloperCommandHandler(
ISessionForEventSourcing sessionForEventSourcing, IMapper mapper)
{
_mapper = mapper;
_sessionForEventSourcing = sessionForEventSourcing;
}
public Task<EsRejectDeveloperCommandResponse> Handle(EsRejectDeveloperCommand request,
CancellationToken cancellationToken)
{
var developerUniqueId = _mapper.Map<DeveloperUniqueId>
(request.DeveloperUniqueId);
var aggregateDevloper = _sessionForEventSourcing.Get<DevloperAggregate>
(developerUniqueId.GetAggregateKey(), null);
if (aggregateDevloper.Version > request.Version)
return Task.FromResult(new EsRejectDeveloperCommandResponse
($"You sended old version.
Yours {request.Version}.
In Event database :{aggregateDevloper.Version}", false));
if (aggregateDevloper.Status != DeveloperStatus.New)
return Task.FromResult(new EsRejectDeveloperCommandResponse
($"Devloper status in eventhistory is {aggregateDevloper.Status}.
Can't be Rejected", false));
var developer = _mapper.Map<Developer>(aggregateDevloper);
developer.Status = DeveloperStatus.Rejected;
aggregateDevloper.Rejected(developer);
_sessionForEventSourcing.Commit();
return Task.FromResult(new EsRejectDeveloperCommandResponse());
}
public class EsSubmitDeveloperCommandHandler
:
IRequestHandler<EsSubmitDeveloperCommand, EsSubmitDeveloperCommandResponse>
{
private readonly ISessionForEventSourcing _sessionForEventSourcing;
private readonly IMapper _mapper;
public EsSubmitDeveloperCommandHandler(
ISessionForEventSourcing sessionForEventSourcing, IMapper mapper)
{
_mapper = mapper;
_sessionForEventSourcing = sessionForEventSourcing;
}
public async Task<EsSubmitDeveloperCommandResponse> Handle
(EsSubmitDeveloperCommand request, CancellationToken cancellationToken)
{
var validator = new EsSubmitDeveloperCommandValidator();
var validatorResult = await validator.ValidateAsync(request);
if (!validatorResult.IsValid)
return new EsSubmitDeveloperCommandResponse(validatorResult);
var developer = _mapper.Map<Developer>(request);
var developerAgreggate = new DevloperAggregate(developer);
_sessionForEventSourcing.Add<DevloperAggregate>(developerAgreggate);
_sessionForEventSourcing.Commit();
var ids = _mapper.Map<DeveloperIds>(developer.Ids());
return new EsSubmitDeveloperCommandResponse(ids);
}
}
public class Constants
{
public const string QUEUE_DEVELOPER_SUBMITED = "developers_submited";
public const string QUEUE_DEVELOPER_REJECTCED = "developers_rejected";
public const string QUEUE_DEVELOPER_ACCEPTED = "developers_accepted";
}
public interface ISettings
{
TimeSpan Frequency { get; set; }
TimeSpan Timeout { get; set; }
}
public class Settings : ISettings
{
public TimeSpan Frequency { get; set; }
public TimeSpan Timeout { get; set; }
}
public interface ISubscribeBase
{
string QUEUE_Name { get; }
DomainEvent DeserializeObject(string json);
Task HandleBasicDeliver(string consumerTag,
ulong deliveryTag, bool redelivered, string exchange,
string routingKey,
IBasicProperties properties, ReadOnlyMemory<byte> body);
void StartSubing(IModel channel);
void Dispose();
Task<ExecutionStatus> HandleEvent(DomainEvent @event);
}
public abstract class SubscribeBase : ISubscribeBase
{
private MessageReceiverBase Consumer { get; set; }
private IModel _channel;
public abstract string QUEUE_Name { get; }
public SubscribeBase()
{
}
public void StartSubing(IModel channel)
{
_channel = channel;
Consumer = new MessageReceiverBase(channel, this);
channel.BasicConsume(this.QUEUE_Name,
false, Consumer);
}
public abstract DomainEvent DeserializeObject(string json);
public async Task HandleBasicDeliver(string consumerTag, ulong deliveryTag,
bool redelivered, string exchange, string routingKey, IBasicProperties
properties, ReadOnlyMemory<byte> body)
{
var json = Encoding.UTF8.GetString(body.Span);
var obj = DeserializeObject(json);
var status = await HandleEvent(obj);
if (status.Success)
{
try
{
_channel.BasicAck(deliveryTag, false);
}
catch (Exception ex)
{
throw;
}
}
}
public void Dispose()
{
Consumer.Dispose();
}
public abstract Task<ExecutionStatus> HandleEvent(DomainEvent @event);
}
public class MessageReceiverBase : DefaultBasicConsumer
{
private readonly IModel _channel;
private ISubscribeBase _messageReceiverBasez;
public void Dispose()
{
_channel.Dispose();
}
public MessageReceiverBase(IModel channel, ISubscribeBase messageReceiverBasez)
{
_messageReceiverBasez = messageReceiverBasez;
_channel = channel;
}
public override void HandleBasicDeliver(string consumerTag,
ulong deliveryTag, bool redelivered, string exchange,
string routingKey,
IBasicProperties properties, ReadOnlyMemory<byte> body)
{
Console.ForegroundColor = ConsoleColor.DarkGreen;
Console.WriteLine($"Consuming Message");
Console.WriteLine(string.Concat("Message received from the exchange ", exchange));
Console.WriteLine(string.Concat("Consumer tag: ", consumerTag));
Console.WriteLine(string.Concat("Delivery tag: ", deliveryTag));
Console.WriteLine(string.Concat("Routing tag: ", routingKey));
var json = Encoding.UTF8.GetString(body.Span);
Console.WriteLine(string.Concat("Message: ", json));
Console.ForegroundColor = ConsoleColor.Gray;
_messageReceiverBasez.HandleBasicDeliver
(consumerTag, deliveryTag, redelivered, exchange, routingKey, properties, body);
}
}
public class SubscribeSubmitDeveloper : SubscribeBase
{
private IZEsDeveloperRepository _ZEsDeveloperRepository;
private IMapper _mapper;
public SubscribeSubmitDeveloper(IZEsDeveloperRepository zEsDeveloperRepository,
IMapper mapper) :
base()
{
_ZEsDeveloperRepository = zEsDeveloperRepository;
_mapper = mapper;
}
public override string QUEUE_Name => Constants.QUEUE_DEVELOPER_SUBMITED;
public override DomainEvent DeserializeObject(string json)
{
return JsonConvert.DeserializeObject<DevloperSubmitedEvent>(json);
}
public override async Task<ExecutionStatus> HandleEvent(DomainEvent @event)
{
DevloperSubmitedEvent developerSubmitEvent = @event as DevloperSubmitedEvent;
var cfs = _mapper.Map<Developer>(developerSubmitEvent);
var execution = await
_ZEsDeveloperRepository.SubmitAsync(cfs);
if (execution == null)
return new ExecutionStatus() { Success = false };
return new ExecutionStatus() { Success = true };
}
}
public class SubscribeAcceptDeveloper : SubscribeBase
{
private IZEsDeveloperRepository _ZEsDeveloperRepository;
private IMapper _mapper;
public SubscribeAcceptDeveloper(IZEsDeveloperRepository zEsDeveloperRepository,
IMapper mapper) :
base()
{
_ZEsDeveloperRepository = zEsDeveloperRepository;
_mapper = mapper;
}
public override string QUEUE_Name => Constants.QUEUE_DEVELOPER_ACCEPTED;
public override DomainEvent DeserializeObject(string json)
{
return JsonConvert.DeserializeObject<DeveloperAcceptedEvent>(json);
}
public override async Task<ExecutionStatus> HandleEvent(DomainEvent @event)
{
DeveloperAcceptedEvent developerEvent = @event as DeveloperAcceptedEvent;
var cfs = _mapper.Map<Developer>(developerEvent);
var status = await _ZEsDeveloperRepository.SaveAcceptenceAsync(cfs.UniqueId);
return new ExecutionStatus() { Success = status };
}
}
public class SubscribeRejectDeveloper : SubscribeBase
{
private IZEsDeveloperRepository _ZEsDeveloperRepository;
private IMapper _mapper;
public SubscribeRejectDeveloper(IZEsDeveloperRepository zEsDeveloperRepository,
IMapper mapper) :
base()
{
_ZEsDeveloperRepository = zEsDeveloperRepository;
_mapper = mapper;
}
public override string QUEUE_Name => Constants.QUEUE_DEVELOPER_REJECTCED;
public override DomainEvent DeserializeObject(string json)
{
return JsonConvert.DeserializeObject<DeveloperRejectedEvent>(json);
}
public override async Task<ExecutionStatus> HandleEvent(DomainEvent @event)
{
DeveloperRejectedEvent developerEvent = @event as DeveloperRejectedEvent;
var cfs = _mapper.Map<Developer>(developerEvent);
var status = await _ZEsDeveloperRepository.SaveRejectionAsync(cfs.UniqueId);
return new ExecutionStatus() { Success = status };
}
}
public void ConfigureServices(IServiceCollection services)
{
services.AddAutoMapper(Assembly.GetExecutingAssembly());
services.AddPersitenceDapperSQLite(Configuration);
var seriFileLogger = new LoggerConfiguration().WriteTo.File(@"D:\Temp\").CreateLogger();
services.AddSingleton<Serilog.ILogger>(seriFileLogger);
services.AddSingleton<ISettings>(new Settings()
{
Timeout = TimeSpan.FromSeconds(5),
Frequency = TimeSpan.FromSeconds(5),
});
//dont change the order
services.AddTransient<ISubscribeBase, SubscribeSubmitDeveloper>();
services.AddTransient<ISubscribeBase, SubscribeAcceptDeveloper>();
services.AddTransient<ISubscribeBase, SubscribeRejectDeveloper>();
services.AddTransient<ISubscribeBase[]>
(p => p.GetServices<ISubscribeBase>().ToArray());
services.AddHostedService<BackgroundEventHandlersServerService>();
}
public class BackgroundEventHandlersServerService : BackgroundService
{
private readonly ILogger Logger;
private readonly ISettings Settings;
private readonly ISubscribeBase[] _subscribes;
private string _ContentRootPath;
public BackgroundEventHandlersServerService(ILogger logger, ISettings settings,
IHostingEnvironment env,
ISubscribeBase[] subscribes)
{
Settings = settings;
Logger = logger;
_subscribes = subscribes;
_ContentRootPath = env.ContentRootPath;
}
protected async override Task ExecuteAsync(CancellationToken stoppingToken)
{
ConnectionFactory connectionFactory = new ConnectionFactory();
var builder = new ConfigurationBuilder()
.SetBasePath(_ContentRootPath)
.AddJsonFile("appsettings.json", optional: false, reloadOnChange: false)
.AddEnvironmentVariables();
builder.Build().GetSection("RabbitMqSetting").Bind(connectionFactory);
using IConnection conn = connectionFactory.CreateConnection();
using IModel channel = conn.CreateModel();
Console.WriteLine("CreateConnection CreateModel");
try
{
Console.ForegroundColor = ConsoleColor.DarkCyan;
foreach (var item in _subscribes)
{
channel.QueueDeclare(
queue: item.QUEUE_Name,
durable: false,
exclusive: false,
autoDelete: false,
arguments: null
);
Console.WriteLine($"QueueDeclare {item.QUEUE_Name}");
}
foreach (var item in _subscribes)
{
item.StartSubing(channel);
Console.WriteLine($"StartSubing {item.QUEUE_Name}");
}
Console.ForegroundColor = ConsoleColor.Gray;
}
catch (Exception ex)
{
Console.ForegroundColor = ConsoleColor.Red;
Console.WriteLine($"Exception {ex}");
Console.ForegroundColor = ConsoleColor.Gray;
LogError(ex.Message);
}
int i = 0;
while (!stoppingToken.IsCancellationRequested)
{
i++;
Console.ForegroundColor = (ConsoleColor)i;
Console.WriteLine($"Still Working and Still Waiting");
Console.ForegroundColor = ConsoleColor.Gray;
if (i == 15)
{
Console.Clear();
i = 0;
}
await Task.Delay(Settings.Frequency, stoppingToken);
}
}
public async override Task StopAsync(CancellationToken cancellationToken)
{
foreach (var item in _subscribes)
{
try
{
item.Dispose();
}
catch (Exception)
{
}
}
await base.StopAsync(cancellationToken);
}
private void LogError(string error)
{
Logger.Error(error);
}
}