Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
21 changes: 21 additions & 0 deletions SW.Bitween.Api/Domain/Document/Document.cs
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,17 @@ public Document(string code, string name, DocumentFormat format)
public DocumentFormat DocumentFormat { get; set; }
public IReadOnlyDictionary<string, string> PromotedProperties { get; private set; }

/// <summary>
/// When this type was taken out of use, or null while it is still in use.
/// </summary>
/// <remarks>
/// Retiring is what a type that has been used wants instead of deleting: exchanges name
/// the type they carried, so removing the row to get it out of a picker would either take
/// that history with it or leave it anonymous. A retired type keeps answering for every
/// exchange already recorded against it, and simply stops being offered for new work.
/// </remarks>
public DateTime? RetiredOn { get; private set; }

public void SetDictionaries(IReadOnlyDictionary<string, string> promotedProperties)
{
PromotedProperties = promotedProperties;
Expand All @@ -57,5 +68,15 @@ public void SetCode(string code)
{
Code = code;
}

public void Retire()
{
RetiredOn = DateTime.UtcNow;
}

public void Restore()
{
RetiredOn = null;
}
}
}
5 changes: 5 additions & 0 deletions SW.Bitween.Api/Extensions/IReadOnlyDictionaryExtensions.cs
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,11 @@ public static Dictionary<TKey, TValue> ToDictionary<TKey, TValue>(this IReadOnly
public static ICollection<KeyAndValue> ToKeyAndValueCollection<TKey, TValue>(
this IReadOnlyDictionary<TKey, TValue> dict)
{
// A null column is "no entries", not a failure. Rows the domain creates always have
// one, but a row inserted any other way took down the whole list response rather
// than the single row that was missing it.
if (dict == null) return new List<KeyAndValue>();

return dict.Select(kvp => new KeyAndValue
{
Key = kvp.Key.ToString(),
Expand Down
38 changes: 37 additions & 1 deletion SW.Bitween.Api/Resources/Documents/Delete.cs
Original file line number Diff line number Diff line change
Expand Up @@ -3,8 +3,11 @@
using SW.PrimitiveTypes;
using System;
using System.Collections.Generic;
using System.Linq;
using System.Text;
using System.Threading.Tasks;
using Microsoft.EntityFrameworkCore;
using SW.Bitween.Domain.Gateway;

namespace SW.Bitween.Resources.Documents
{
Expand All @@ -15,9 +18,42 @@ async public Task<object> Handle(int key)
{
await requestContext.EnsurePermission(dbContext, Model.Permissions.Documents.Delete);

// Seeded by the model and referenced by every aggregation, so it is not the user's
// to remove — the same footing as the system partner.
if (key == Document.AggregationDocumentId)
throw new SWException("The built-in Aggregation Document can't be deleted.");

// Configuration pointing at this type is a real block — a subscription or a bus
// gateway left behind would name a type that no longer exists. Said plainly, because
// the database says it as a foreign-key violation, which reached the screen as a bare
// 500 naming a constraint.
if (await dbContext.Set<Subscription>().AnyAsync(s => s.DocumentId == key))
throw new SWException(
"Cannot delete an information type that subscriptions still carry. Delete them, or point them at another type, first.");

if (await dbContext.Set<BusGateway>().AnyAsync(g => g.DocumentId == key))
throw new SWException(
"Cannot delete an information type that a bus gateway still listens for. Delete the gateway, or point it at another type, first.");

// Exchanges are the deliberate exception. They are this type's history rather than
// configuration depending on it, and history should not be able to strand a type
// nobody uses any more — so they go with it.
await using var transaction = await dbContext.Database.BeginTransactionAsync();

var xchangeIds = dbContext.Set<Xchange>().Where(x => x.DocumentId == key).Select(x => x.Id);

// Scheduled retries are keyed by exchange id but have no relationship configured, so
// nothing clears them on their own — they would be left pointing at exchanges that no
// longer exist. Results, aggregations and promoted properties all cascade from the
// exchange itself and need no help here.
await dbContext.Set<DelayedRetry>().Where(d => xchangeIds.Contains(d.Id)).ExecuteDeleteAsync();
await dbContext.Set<Xchange>().Where(x => x.DocumentId == key).ExecuteDeleteAsync();

await dbContext.DeleteByKeyAsync<Document>(key);

await transaction.CommitAsync();
await cache.BroadcastRevoke();
return null;
}
}
}
}
10 changes: 8 additions & 2 deletions SW.Bitween.Api/Resources/Documents/Get.cs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,10 @@ public async Task<object> Handle(int key)
{
await requestContext.EnsurePermission(dbContext, Model.Permissions.Documents.View);

return await dbContext.Set<Document>().Search("Id", key).Select(document => new DocumentUpdate
// DocumentRow rather than DocumentUpdate: it is the read shape, and RetiredOn belongs
// on it rather than on the model used to write a type — retiring is its own command,
// not a field an update can set.
return await dbContext.Set<Document>().Search("Id", key).Select(document => new DocumentRow
{
Id = document.Id,
Code = document.Code,
Expand All @@ -30,7 +33,10 @@ public async Task<object> Handle(int key)
DuplicateInterval = document.DuplicateInterval,
PromotedProperties = document.PromotedProperties.ToKeyAndValueCollection(),
DocumentFormat = document.DocumentFormat,
DisregardsUnfilteredMessages = document.DisregardsUnfilteredMessages ?? false
DisregardsUnfilteredMessages = document.DisregardsUnfilteredMessages ?? false,
RetiredOn = document.RetiredOn,
UsedByCount = dbContext.Set<Subscription>()
.Count(subscription => subscription.DocumentId == document.Id)
}).SingleOrDefaultAsync();
}
}
Expand Down
16 changes: 15 additions & 1 deletion SW.Bitween.Api/Resources/Documents/PromotedPropertyValidation.cs
Original file line number Diff line number Diff line change
Expand Up @@ -26,12 +26,26 @@ public static void Check(ICollection<KeyAndValue> promotedProperties, DocumentFo
throw new SWValidationException("INVALID_PROMOTED_PROPERTY_VALUE",
$"Promoted property '{pp.Key}' must have a non-empty path value.");

// Trimmed here *and* written back. Validating the trimmed value while storing the
// raw one let " $.ref" pass every check below and then fail at read time, which
// 400s every message on this information type rather than this one save.
var trimmed = pp.Value.Trim();
pp.Value = trimmed;

if (format == DocumentFormat.Json)
{
// Must be a JSONPath: starts with '$' or a simple dot-separated identifier path
if (!trimmed.StartsWith("$") && !Regex.IsMatch(trimmed, @"^[a-zA-Z_][a-zA-Z0-9_]*(?:(\.[a-zA-Z_][a-zA-Z0-9_]*)|(\[[0-9]+\]))*$"))
//
// The '$' branch is checked rather than trusted. Treating any leading '$' as
// proof of a JSONPath let "$", "$." and "$.." through, none of which select
// anything — they saved cleanly and failed only once a message arrived. A step
// is a '.'/'..' followed by a name, or a bracket; names stay deliberately
// permissive (anything but a delimiter) so paths that already work keep working.
var valid = trimmed.StartsWith("$")
? Regex.IsMatch(trimmed, @"^\$(?:\.\.?[^.\[\]]+|\[[^\]]+\])+$")
: Regex.IsMatch(trimmed, @"^[a-zA-Z_][a-zA-Z0-9_]*(?:(\.[a-zA-Z_][a-zA-Z0-9_]*)|(\[[0-9]+\]))*$");

if (!valid)
throw new SWValidationException("INVALID_PROMOTED_PROPERTY_PATH",
$"Promoted property '{pp.Key}' has an invalid JSON path: '{pp.Value}'. Expected a JSONPath expression (e.g. '$.field.subField') or dot-notation path.");
}
Expand Down
62 changes: 62 additions & 0 deletions SW.Bitween.Api/Resources/Documents/Retire.cs
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
using System.Threading.Tasks;
using Microsoft.EntityFrameworkCore;
using SW.Bitween.Domain;
using SW.Bitween.Domain.Gateway;
using SW.Bitween.Model;
using SW.PrimitiveTypes;

namespace SW.Bitween.Resources.Documents
{
/// <summary>
/// Takes an information type out of use, or puts it back. Toggles, the same way a
/// subscription's pause does.
/// </summary>
/// <remarks>
/// The gentle half of <see cref="Delete"/>: a type that has carried traffic usually wants to
/// leave the pickers without taking its exchanges with it, and this is reversible where a
/// delete is not.
/// </remarks>
[HandlerName("retire")]
public class Retire(BitweenDbContext dbContext, RequestContext requestContext, IInfolinkCache cache)
: ICommandHandler<int, DocumentRetire, object>
{
public async Task<object> Handle(int key, DocumentRetire request)
{
await requestContext.EnsurePermission(dbContext, Model.Permissions.Documents.Edit);

if (key == Document.AggregationDocumentId)
throw new SWException("The built-in Aggregation Document can't be retired.");

var entity = await dbContext.FindAsync<Document>(key);
if (entity == null)
throw new SWNotFoundException("Information type not found.");

if (entity.RetiredOn == null)
{
// Configuration still pointing at it would be left running against a type the UI
// no longer offers — which reads as the type having been removed when it has not.
// Said before retiring rather than after, so nothing half-happens.
if (await dbContext.Set<Subscription>().AnyAsync(s => s.DocumentId == key))
throw new SWException(
"Cannot retire an information type that subscriptions still carry. Delete them, or point them at another type, first.");

if (await dbContext.Set<BusGateway>().AnyAsync(g => g.DocumentId == key))
throw new SWException(
"Cannot retire an information type that a bus gateway still listens for. Delete the gateway, or point it at another type, first.");

entity.Retire();
}
else
{
entity.Restore();
}

await dbContext.SaveChangesAsync();
// The receiving path resolves types through the cache, so a retired type would keep
// being offered — and keep accepting work — for the rest of the cache's window.
await cache.BroadcastRevoke();

return new { entity.Id, entity.RetiredOn };
}
}
}
1 change: 1 addition & 0 deletions SW.Bitween.Api/Resources/Documents/Search.cs
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,7 @@ async public Task<object> Handle(SearchyRequest searchyRequest, bool lookup = fa
DuplicateInterval = document.DuplicateInterval,
PromotedProperties = document.PromotedProperties.ToKeyAndValueCollection(),
DocumentFormat = document.DocumentFormat,
RetiredOn = document.RetiredOn,
// A correlated count, so the "used by" column the UI shows costs one
// subquery per row instead of the whole Subscription table over the wire.
UsedByCount = dbContext.Set<Subscription>()
Expand Down
8 changes: 6 additions & 2 deletions SW.Bitween.Api/Resources/Subscriptions/Create.cs
Original file line number Diff line number Diff line change
Expand Up @@ -134,10 +134,14 @@ await AddMissing(context, adapterRequirements,

// Schedules only mean anything on the two scheduled types, and an empty set is
// what Subscription.SetSchedules rejects outright.
//
// The type is the whole condition. Guarding on `Schedules != null` as well let the
// rule be skipped by the one case it most needed to catch — omitting the field
// entirely — so a scheduled job or aggregation saved with no trigger at all and
// then simply never ran.
RuleFor(i => i.Schedules)
.NotEmpty()
.When(i => i.Schedules != null &&
(i.Type == SubscriptionType.Receiving || i.Type == SubscriptionType.Aggregation))
.When(i => i.Type == SubscriptionType.Receiving || i.Type == SubscriptionType.Aggregation)
.WithMessage("Schedules cannot be empty for a scheduled subscription.");
}

Expand Down
54 changes: 51 additions & 3 deletions SW.Bitween.Api/Resources/Subscriptions/Search.cs
Original file line number Diff line number Diff line change
Expand Up @@ -102,14 +102,56 @@ join document in _dbContext.Set<Document>() on subscriber.DocumentId equals docu
return await SearchWithEdgeCases(query, edgeCaseFilters, searchyRequest);
}

var result = await query.Search(searchyRequest.Conditions, searchyRequest.Sorts, searchyRequest.PageSize,
searchyRequest.PageIndex).ToListAsync();
await AttachSchedules(result);

return new SearchyResponse<SubscriptionSearch>
{
TotalCount = count,
Result = await query.Search(searchyRequest.Conditions, searchyRequest.Sorts, searchyRequest.PageSize,
searchyRequest.PageIndex).ToListAsync()
Result = result
};
}

/// <summary>Fills in each returned row's schedules.</summary>
/// <remarks>
/// A second query rather than part of the projection above: <c>Schedule.On</c> is a
/// <see cref="System.TimeSpan"/> stored as ticks, and reading <c>.Days</c>/<c>.Hours</c>/
/// <c>.Minutes</c> off it inside that joined query is what Postgres cannot translate — a
/// date_part type mismatch. Read flat for the ids actually being returned and shaped in
/// memory, which asks nothing of the translator, so the list can finally say when a job
/// runs instead of leaving its schedule column blank.
/// </remarks>
private async Task AttachSchedules(List<SubscriptionSearch> rows)
{
// Only the two scheduled types have any; asking for the rest is a wasted round trip.
var ids = rows
.Where(r => r.Type is SubscriptionType.Receiving or SubscriptionType.Aggregation)
.Select(r => r.Id)
.ToList();

if (ids.Count == 0) return;

var byId = await _dbContext.Set<Subscription>().AsNoTracking()
.Where(s => ids.Contains(s.Id))
.Select(s => new { s.Id, Schedules = s.Schedules.ToList() })
.ToDictionaryAsync(x => x.Id, x => x.Schedules);

foreach (var row in rows)
{
if (!byId.TryGetValue(row.Id, out var schedules)) continue;

row.Schedules = schedules.Select(s => new ScheduleView
{
Backwards = s.Backwards,
Recurrence = s.Recurrence,
Days = s.On.Days,
Hours = s.On.Hours,
Minutes = s.On.Minutes
}).ToList();
}
}

private async Task<SearchyResponse<SubscriptionSearch>> SearchWithEdgeCases(
IQueryable<SubscriptionSearch> query, IEnumerable<SearchyFilter> edgeCaseFilters,
SearchyRequest searchyRequest)
Expand Down Expand Up @@ -138,10 +180,16 @@ private async Task<SearchyResponse<SubscriptionSearch>> SearchWithEdgeCases(
};
}

// Only the page being returned, same as the normal path — the edge-case filters run
// over the whole set in memory, and shaping schedules for all of it would be waste.
var page = data.Skip(searchyRequest.PageSize * searchyRequest.PageIndex)
.Take(searchyRequest.PageSize).ToList();
await AttachSchedules(page);

return new SearchyResponse<SubscriptionSearch>
{
TotalCount = data.Count,
Result = data.Skip(searchyRequest.PageSize * searchyRequest.PageIndex).Take(searchyRequest.PageSize)
Result = page
};
}

Expand Down
11 changes: 11 additions & 0 deletions SW.Bitween.Api/Resources/Xchanges/Search.cs
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,17 @@ from delayedRetry in drGroup.DefaultIfEmpty()
InputFileName = xchange.InputName,
OutputFileName = result.OutputName,
ResponseFileName = result.ResponseName,
// The same three counts the file keys above are already derived from.
// Left unassigned, every stage reported its document as "0 b".
//
// The result-side two are guarded because the join to XchangeResult is
// a left one: an exchange still running, or one that failed before it
// produced a result, has no row there. Reading a non-nullable int off
// that null threw "Nullable object must have a value" out of the whole
// query — one such exchange 500'd the entire list.
InputFileSize = xchange.InputSize,
OutputFileSize = result != null ? result.OutputSize : 0,
ResponseFileSize = result != null ? result.ResponseSize : 0,
CorrelationId = xchange.CorrelationId,
// xchange.PartnerId is the authoritative source (set at creation from the
// gateway/bus-route partner, or the subscription's own PartnerId as a
Expand Down
25 changes: 21 additions & 4 deletions SW.Bitween.Api/Services/FilterService.cs
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
using System;
using System.Linq;
using System.Threading.Tasks;
using System.Xml.XPath;
using Newtonsoft.Json;
using SW.PrimitiveTypes;
using SW.Bitween.Model;

Expand Down Expand Up @@ -29,10 +31,25 @@ public async Task<FilterResult> Filter(int documentId, XchangeFile xchangeFile)

foreach (var pp in doc.PromotedProperties)
{
propReader.TryGetValue(pp.Value, out var ppValue);
//TODO check if we need to validate here
//if (ppValue is null)
// throw new SWValidationException("PROMOTED_PROPERTY_NOT_FOUND", $"The path {pp.Value} is null on the docuemnt");
string ppValue;
try
{
propReader.TryGetValue(pp.Value, out ppValue);
}
catch (Exception ex) when (ex is JsonException || ex is XPathException)
{
// A path the reader cannot parse throws out of it rather than returning false,
// and unhandled it reached the caller as a bare 400 naming nothing — on every
// message of this information type, since the bad path is on the type itself.
// Which property is broken is the one thing needed to fix it.
throw new SWValidationException("INVALID_PROMOTED_PROPERTY_PATH",
$"Promoted property '{pp.Key}' on information type '{doc.Name}' has a path that cannot be read: '{pp.Value}'. {ex.Message}");
}

// A path that simply doesn't match this payload is not an error: an information
// type promotes what its documents *may* carry, and TryGetValue reports that by
// returning false, leaving the property null. Only an unreadable path throws.

// Stored as the payload sent it. It used to be lower-cased here, which was
// only ever to pair with the lower-cased term in Xchanges/Search — nothing
// matches on this dictionary (match expressions read the payload directly),
Expand Down
Loading
Loading