From e6e7a2dc59843732753af801ecf956fda63007c8 Mon Sep 17 00:00:00 2001 From: Josh Patra Date: Thu, 20 Aug 2026 00:47:38 -0400 Subject: [PATCH] fix(intelligence): undo imports and reject video history --- .../ListeningHistoryImportIntegrationTests.cs | 98 +++++++++++++++++++ .../IntelligenceController.History.cs | 34 +++++++ .../ListeningHistoryImportPersistence.cs | 52 ++++++++++ webui/src/lib/api.ts | 5 + .../lib/components/IntelligenceHistory.svelte | 38 +++++++ webui/tests/parity.e2e.ts | 18 ++++ 6 files changed, 245 insertions(+) diff --git a/allstarr.Tests/Intelligence/ListeningHistoryImportIntegrationTests.cs b/allstarr.Tests/Intelligence/ListeningHistoryImportIntegrationTests.cs index 0dd4f39d..de297360 100644 --- a/allstarr.Tests/Intelligence/ListeningHistoryImportIntegrationTests.cs +++ b/allstarr.Tests/Intelligence/ListeningHistoryImportIntegrationTests.cs @@ -2,6 +2,7 @@ using System.Text.Json; using allstarr.Core.Intelligence; using allstarr.Core.Jobs; using allstarr.Core.Operations; +using allstarr.Core.Playback; using allstarr.Core.Storage; using allstarr.Models.Settings; using Microsoft.EntityFrameworkCore; @@ -263,6 +264,69 @@ public sealed class ListeningHistoryImportIntegrationTests : IAsyncLifetime Assert.Equal(["Expired", "Recent"], history.Select(item => item.Title)); } + [Fact] + public async Task PreviewRejectsSpotifyVideoHistoryBeforeSavingAnImport() + { + var service = new ListeningHistoryImportService( + _factory, _importers, _artifacts, _importOptions, _clock, _jobs); + var source = JsonSerializer.SerializeToUtf8Bytes(new[] + { + Row("2026-07-20T12:00:00Z", "Video play", "1111111111111111111111", 180_000, "trackdone", false) + }); + await using var stream = new MemoryStream(source); + + var exception = await Assert.ThrowsAsync(() => + service.PreviewAsync( + _scope, + "Streaming_History_Video_2017-2024.json", + stream, + source.Length, + CancellationToken.None)); + + Assert.Equal("history_import_video_unsupported", exception.Code); + await using var db = await _factory.CreateDbContextAsync(); + Assert.Empty(await db.ListeningHistoryImports.ToListAsync()); + Assert.False(Directory.Exists(_root)); + } + + [Fact] + public async Task RemoveImportDeletesOnlyItsExactImportedListensAndSavedArtifact() + { + var service = new ListeningHistoryImportService( + _factory, _importers, _artifacts, _importOptions, _clock, _jobs); + var preview = await PreviewAsync(service, "Remove me", "7777777777777777777777"); + var removedOccurrence = new string('a', 64); + var keptOccurrence = new string('b', 64); + await using (var db = await _factory.CreateDbContextAsync()) + { + var import = await db.ListeningHistoryImports.SingleAsync(item => item.Id == preview.ImportId); + import.State = ListeningHistoryImportState.Completed; + import.ImportedRows = 1; + import.CompletedAt = _clock.UtcNow; + import.Revision++; + db.ListeningEvents.AddRange( + ImportedEvent(preview.ImportId, removedOccurrence), + ImportedEvent(Guid.NewGuid(), keptOccurrence)); + db.PlaybackDeliveryCheckpoints.AddRange( + Checkpoint(removedOccurrence, new string('c', 64)), + Checkpoint(keptOccurrence, new string('d', 64))); + await db.SaveChangesAsync(); + } + + var removed = await service.RemoveAsync( + _scope, preview.ImportId, preview.Revision, CancellationToken.None); + + Assert.Equal(1, removed!.RemovedListens); + Assert.Null(await service.GetAsync(_scope, preview.ImportId, CancellationToken.None)); + Assert.False(File.Exists(Path.Combine(_root, $"{preview.ImportId:N}.json"))); + await using var verify = await _factory.CreateDbContextAsync(); + Assert.Equal(keptOccurrence, Assert.Single(await verify.ListeningEvents.ToListAsync()).OccurrenceKey); + Assert.Equal(keptOccurrence, Assert.Single(await verify.PlaybackDeliveryCheckpoints.ToListAsync()).OccurrenceKey); + Assert.Single(await verify.AuditEvents.Where(item => + item.Category == "listening-history-import" && item.Action == "removed" && + item.CorrelationId == preview.ImportId.ToString("N")).ToListAsync()); + } + [Fact] public async Task FailedApplyPersistsSafeOperationalMetadataWithoutImportContent() { @@ -387,6 +451,40 @@ public sealed class ListeningHistoryImportIntegrationTests : IAsyncLifetime ["incognito_mode"] = false }; + private ListeningEventRecord ImportedEvent(Guid importId, string occurrenceKey) => new() + { + Id = Guid.CreateVersion7(), + TenantId = _scope.TenantId, + OwnerUserId = _scope.OwnerUserId, + Protocol = _scope.Protocol, + BackendInstanceId = _scope.BackendInstanceId, + LibraryScopeId = _scope.LibraryScopeId, + OccurrenceKey = occurrenceKey, + State = ListeningEventState.Completed, + StartedAt = _clock.UtcNow.AddMinutes(-3), + ListenedAt = _clock.UtcNow, + UpdatedAt = _clock.UtcNow, + SourceKind = "import", + ImportProvenance = $"history-import:{importId:N}:spotify:1:track_finished:online:standard", + TrackReference = $"spotify:{occurrenceKey}", + Title = "Imported song", + Artist = "Imported artist", + MusicBrainzEnrichmentState = MusicBrainzEnrichmentState.NotRequested, + Revision = 1 + }; + + private PlaybackDeliveryCheckpointEntity Checkpoint(string occurrenceKey, string signalKey) => new() + { + Id = Guid.CreateVersion7(), + TenantId = _scope.TenantId, + OwnerUserId = _scope.OwnerUserId, + OccurrenceKey = occurrenceKey, + SignalKey = signalKey, + TargetId = "fixture-target", + DetailsJson = "{}", + UpdatedAt = _clock.UtcNow + }; + private sealed class TestDbContextFactory(DbContextOptions options) : IDbContextFactory { diff --git a/allstarr/Controllers/IntelligenceController.History.cs b/allstarr/Controllers/IntelligenceController.History.cs index 67f6c49b..d3978ed4 100644 --- a/allstarr/Controllers/IntelligenceController.History.cs +++ b/allstarr/Controllers/IntelligenceController.History.cs @@ -418,6 +418,34 @@ public sealed partial class IntelligenceController CancellationToken cancellationToken) => ChangeHistoryImport(importId, request, imports, "cancel", cancellationToken); + [HttpDelete("history/imports/{importId:guid}")] + public async Task RemoveHistoryImport( + Guid importId, + [FromBody] IntelligenceHistoryImportDeleteRequest request, + [FromServices] ListeningHistoryImportService imports, + CancellationToken cancellationToken) + { + Response.Headers.CacheControl = "no-store"; + if (!TrySessionScope(request, out var scope, out var error)) return error!; + if (!request.Confirmed) + return BadRequest(new { error = "history_import_confirmation_required" }); + await using var db = await _factory.CreateDbContextAsync(cancellationToken); + if (!await OwnsBackend(db, scope, cancellationToken)) return NotFound(); + try + { + var result = await imports.RemoveAsync(scope, importId, request.Revision, cancellationToken); + return result == null ? NotFound() : Ok(new { removedImport = true, result.RemovedListens }); + } + catch (ListeningHistoryImportException exception) when (exception.Code.EndsWith("_conflict", StringComparison.Ordinal)) + { + return Conflict(new { error = exception.Code, message = exception.Message }); + } + catch (ListeningHistoryImportException exception) + { + return BadRequest(new { error = exception.Code, message = exception.Message }); + } + } + private async Task ChangeHistoryImport( Guid importId, IntelligenceHistoryImportCommandRequest request, @@ -832,6 +860,12 @@ public sealed class IntelligenceHistoryImportCommandRequest : IntelligenceScopeR public string Revision { get; set; } = ""; } +public sealed class IntelligenceHistoryImportDeleteRequest : IntelligenceScopeRequest +{ + public string Revision { get; set; } = ""; + public bool Confirmed { get; set; } +} + internal readonly record struct ListeningHistoryPeriod(DateTimeOffset From, DateTimeOffset To) { public static bool TryCreate( diff --git a/allstarr/Core/Intelligence/ListeningHistoryImportPersistence.cs b/allstarr/Core/Intelligence/ListeningHistoryImportPersistence.cs index d03b76a1..049857e6 100644 --- a/allstarr/Core/Intelligence/ListeningHistoryImportPersistence.cs +++ b/allstarr/Core/Intelligence/ListeningHistoryImportPersistence.cs @@ -5,6 +5,7 @@ using System.Text.Json; using allstarr.Core.Capabilities; using allstarr.Core.Jobs; using allstarr.Core.Operations; +using allstarr.Core.Playback; using allstarr.Core.Storage; using Microsoft.EntityFrameworkCore; @@ -135,6 +136,8 @@ public sealed record ListeningHistoryImportPreviewResult( long UnresolvedRows, ListeningHistoryImportPreview Preview); +public sealed record ListeningHistoryImportRemovalResult(long RemovedListens); + public sealed class ListeningHistoryImportOptions { public const string SectionName = "Intelligence:HistoryImport"; @@ -279,6 +282,11 @@ public sealed class ListeningHistoryImportService( displayFileName = Path.GetFileName(displayFileName).Trim(); if (displayFileName.Length is < 1 or > 255 || displayFileName.Any(char.IsControl)) throw new ListeningHistoryImportException("history_import_filename_invalid", "The selected filename is invalid."); + if (displayFileName.StartsWith("Streaming_History_Video_", StringComparison.OrdinalIgnoreCase) || + displayFileName.Equals("Streaming_History_Video.json", StringComparison.OrdinalIgnoreCase)) + throw new ListeningHistoryImportException( + "history_import_video_unsupported", + "Spotify video viewing history is not music listening history. Choose the Streaming_History_Audio JSON files instead."); if (sizeBytes is < 1 || sizeBytes > options.MaximumUploadBytes) throw new ListeningHistoryImportException( "history_import_file_invalid", @@ -512,6 +520,42 @@ public sealed class ListeningHistoryImportService( return await GetAsync(scope, importId, cancellationToken); } + public async Task RemoveAsync( + IntelligenceScope scope, + Guid importId, + string expectedRevision, + CancellationToken cancellationToken) + { + await using var db = await factory.CreateDbContextAsync(cancellationToken); + await using var transaction = await db.Database.BeginTransactionAsync(cancellationToken); + var record = await ScopedImport(db, scope, importId).SingleOrDefaultAsync(cancellationToken); + if (record == null) return null; + RequireRevision(record, expectedRevision); + if (record.State is ListeningHistoryImportState.Pending or ListeningHistoryImportState.Running) + throw new ListeningHistoryImportException( + "history_import_state_conflict", + "Cancel this active history import before removing it."); + + var provenance = $"history-import:{importId:N}:"; + var importedEvents = ScopedListeningEvents(db, scope).Where(item => + item.SourceKind == "import" && item.ImportProvenance != null && + item.ImportProvenance.StartsWith(provenance)); + var occurrenceKeys = importedEvents.Select(item => item.OccurrenceKey); + await db.Set().Where(item => + item.TenantId == scope.TenantId && item.OwnerUserId == scope.OwnerUserId && + item.OccurrenceKey != null && occurrenceKeys.Contains(item.OccurrenceKey)) + .ExecuteDeleteAsync(cancellationToken); + var removedListens = await importedEvents.ExecuteDeleteAsync(cancellationToken); + + var now = clock.UtcNow; + db.AuditEvents.Add(Audit(record, "removed", "success", now)); + db.ListeningHistoryImports.Remove(record); + await db.SaveChangesAsync(cancellationToken); + await transaction.CommitAsync(cancellationToken); + artifacts.Delete(importId); + return new(removedListens); + } + private async Task QueueAsync( IntelligenceScope scope, Guid importId, @@ -620,6 +664,14 @@ public sealed class ListeningHistoryImportService( item.Protocol == scope.Protocol && item.BackendInstanceId == scope.BackendInstanceId && item.LibraryScopeId == scope.LibraryScopeId); + private static IQueryable ScopedListeningEvents( + AllstarrDbContext db, + IntelligenceScope scope) => + db.ListeningEvents.Where(item => + item.TenantId == scope.TenantId && item.OwnerUserId == scope.OwnerUserId && + item.Protocol == scope.Protocol && item.BackendInstanceId == scope.BackendInstanceId && + item.LibraryScopeId == scope.LibraryScopeId); + private static void RequireRevision(ListeningHistoryImportRecord record, string expectedRevision) { var normalized = expectedRevision?.Trim().ToLowerInvariant(); diff --git a/webui/src/lib/api.ts b/webui/src/lib/api.ts index a31fba7f..5edf63a5 100644 --- a/webui/src/lib/api.ts +++ b/webui/src/lib/api.ts @@ -1342,6 +1342,11 @@ export const intelligence = { changeHistoryImport: (scope: IntelligenceScope, item: ListeningHistoryImport, operation: "apply" | "resume" | "cancel") => json(`/api/admin/intelligence/history/imports/${encodeURIComponent(item.importId)}/${operation}`, intelligenceBody({ ...scope, revision: item.revision })), + removeHistoryImport: (scope: IntelligenceScope, item: ListeningHistoryImport) => + json<{ removedImport: true; removedListens: number }>( + `/api/admin/intelligence/history/imports/${encodeURIComponent(item.importId)}`, + intelligenceBody({ ...scope, revision: item.revision, confirmed: true }, "DELETE"), + ), createSchedule: (scope: IntelligenceScope, input: Omit) => json("/api/admin/intelligence/schedules", intelligenceBody({ ...scope, ...input })), updateSchedule: (scope: IntelligenceScope, schedule: IntelligenceSchedule, input: Omit) => diff --git a/webui/src/lib/components/IntelligenceHistory.svelte b/webui/src/lib/components/IntelligenceHistory.svelte index 16462a55..43b2f71f 100644 --- a/webui/src/lib/components/IntelligenceHistory.svelte +++ b/webui/src/lib/components/IntelligenceHistory.svelte @@ -70,6 +70,8 @@ let editAlbum = $state(""); let editAlbumArtist = $state(""); let deleteOpen = $state(false); + let importDeleteOpen = $state(false); + let importDeleteItem = $state(null); let importItems = $state([]); let readingImportId = $state(""); let importError = $state(""); @@ -359,6 +361,30 @@ } } + function confirmRemoveImport(item: ImportQueueItem) { + importDeleteItem = item; + importDeleteOpen = true; + } + + async function removeSavedImport() { + const item = importDeleteItem; + if (!item?.result) return; + action = `remove:${item.id}`; + importError = ""; + try { + const result = await intelligence.removeHistoryImport(scope, item.result); + removeImportItem(item.id); + importDeleteOpen = false; + importDeleteItem = null; + if (result.removedListens > 0) await onChanged(); + } catch (cause) { + importError = cause instanceof Error ? cause.message : "This history import could not be removed."; + importDeleteOpen = false; + } finally { + action = ""; + } + } + async function applyAllPreviews() { const items = selectedPreviewItems.map((item) => ({ ...item })); if (!items.length || action) return; @@ -591,6 +617,7 @@ {#if result.state === "previewed"}{/if} {#if activeImport}{/if} {#if result.state === "failed" || result.state === "cancelled"}{/if} + {#if !activeImport}{/if} {#if result.lastErrorMessage}{/if} @@ -622,6 +649,17 @@ + 0 + ? `This removes any listens still stored from ${importDeleteItem?.fileName ?? "this file"} (${(importDeleteItem?.result?.importedRows ?? 0).toLocaleString()} were added), its saved import record, and any temporary upload. Other history is not affected.` + : `This removes ${importDeleteItem?.fileName ?? "this file"} from saved imports and deletes any temporary upload. No listening history will be removed.`} + confirmLabel={action.startsWith("remove:") ? "Removing…" : (importDeleteItem?.result?.importedRows ?? 0) > 0 ? "Undo import" : "Remove import"} + cancelLabel="Keep import" + disabled={Boolean(action)} + onConfirm={removeSavedImport} +/>