Files
ClawdDotNet/src/ClawdDotNet.Tools.SocialMediaManager/SocialMediaManagerTool.cs
T
RichardandClaude Opus 4.8 e5067cae70 Sicherheitsluecken S1, S2 und S3 schliessen
Neues Testprojekt tests/ClawdDotNet.Tools.Tests. Die Angriffsfaelle aus der
Bestandsaufnahme bleiben darin dauerhaft als Testfaelle dokumentiert — zusammen
mit Gegenproben, damit die Fixes nicht zu streng werden und legitime Nutzung
blockieren.

S3 — DirectAPI gab die abgerufene URL als Quelle an das Modell zurueck, samt
API-Schluessel im Query-String. Der Schluessel landete damit im
Konversationskontext, wurde bei jedem Folgeschritt erneut gesendet, in
ChatContext.json geschrieben und in die Logs uebernommen. UrlSanitizer maskiert
sensible Query-Parameter; auch die Fehlermeldungen sind betroffen und werden
bereinigt.

S2 — Die Kanal-/Video-Angabe wurde ungeprueft in eine Argument-Zeichenkette fuer
yt-dlp interpoliert. UseShellExecute=false verhindert Shell-Metazeichen, nicht
aber Options-Injection: yt-dlp kennt die Option --exec, die beliebige Befehle
ausfuehrt. Kritisch, weil der Agent untrusted Inhalte verarbeitet — eine
Prompt-Injection darin konnte ihn dazu bringen, genau so einen Wert zu setzen.

YouTubeUrl validiert Handles und URLs gegen die zulaessigen YouTube-Hosts und
lehnt alles ab, was mit einem Bindestrich beginnt. Die Argumente gehen jetzt
einzeln ueber ProcessStartInfo.ArgumentList, die Adresse steht hinter dem
Optionsende-Trenner. Der ffmpeg-Aufruf wurde ebenso umgestellt.

S1 — Die Tabellen-Whitelist suchte den erlaubten Namen als Teilzeichenkette
irgendwo im Statement, auch in Kommentaren. Bei einer Freigabe fuer prices
genuegte deshalb ein DELETE auf users mit einem Kommentar, der prices enthielt,
um eine beliebige Tabelle zu loeschen. Umgekehrt galten harmlose Abfragen, die
ein Schluesselwort nur als Wert enthielten, faelschlich als Schreibzugriff.

SqlGuard entfernt zuerst Kommentare und String-Literale, lehnt mehrere Statements
ab, bestimmt die Operation am ersten Schluesselwort und extrahiert Tabellennamen
gezielt hinter FROM/JOIN/INTO/UPDATE/TABLE — inklusive kommagetrennter Listen mit
Aliassen. JEDE referenzierte Tabelle muss freigegeben sein, nicht irgendeine.
Ohne Whitelist wird nichts durchgelassen. Die MongoDB-Pruefung vergleicht den
Collection-Namen jetzt exakt statt per Teilzeichenkette.

178 Tests gruen (91 Core, 87 Tools).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-27 18:49:17 +02:00

1240 lines
56 KiB
C#
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
using System.Diagnostics;
using System.Net.Http.Headers;
using System.Text;
using System.Text.Json;
using ClawdDotNet.Core.State;
using ClawdDotNet.Core.Tools;
using Microsoft.Extensions.Logging;
namespace ClawdDotNet.Tools.SocialMediaManager;
public sealed class SocialMediaManagerTool : IAgentTool, IToolJobProvider
{
public string Name => "SocialMediaManager";
public string Description => """
Manager für Social Media Monitoring (X/Twitter, Reddit) und YouTube Transkripte.
Aktionen:
- x_search: Freitext-Suche auf X/Twitter
- x_user_tweets: Tweets eines bestimmten Benutzers abrufen (username angeben)
- reddit_search: Freitext-Suche auf Reddit
- reddit_subreddit: Neue Posts aus einem bestimmten Subreddit abrufen (subreddit angeben)
- youtube_transcript: YouTube-Video transkribieren
Background-Jobs (konfiguriert in toolJobs):
- sm_x_account_monitor: Überwacht konfigurierte X-Accounts auf neue Tweets (Config: xWatchAccounts)
- sm_reddit_monitor: Überwacht konfigurierte Subreddits auf neue Posts (Config: redditWatchSubreddits)
""";
public JsonElement InputSchema { get; } = JsonDocument.Parse("""
{
"type": "object",
"properties": {
"action": {
"type": "string",
"enum": ["x_search", "x_user_tweets", "reddit_search", "reddit_subreddit", "youtube_transcript"],
"description": "Die auszuführende Aktion"
},
"query": {
"type": "string",
"description": "Suchbegriff für x_search oder reddit_search"
},
"username": {
"type": "string",
"description": "X/Twitter Username ohne @ (nur für x_user_tweets)"
},
"subreddit": {
"type": "string",
"description": "Subreddit-Name ohne r/ (nur für reddit_subreddit)"
},
"url": {
"type": "string",
"description": "YouTube Video URL oder Video-ID (nur für youtube_transcript)"
},
"limit": {
"type": "integer",
"description": "Maximale Anzahl an Ergebnissen (Standard: 10)"
}
},
"required": ["action"]
}
""").RootElement.Clone();
private readonly HttpClient _httpClient = new();
public async Task<ToolResult> ExecuteAsync(JsonElement input, AgentToolContext context, CancellationToken ct)
{
var action = input.GetProperty("action").GetString() ?? throw new ArgumentException("Action is required");
// WorkspacePath per-agent scoped im StateStore für Fallback speichern
await context.StateStore.SetAsync($"sm_workspace_path:{context.AgentId}", context.WorkspacePath!, ct);
try
{
return action switch
{
"x_search" => await HandleXSearchAsync(input, context, ct),
"x_user_tweets" => await HandleXUserTweetsAsync(input, context, ct),
"reddit_search" => await HandleRedditSearchAsync(input, context, ct),
"reddit_subreddit" => await HandleRedditSubredditAsync(input, context, ct),
"youtube_transcript" => await HandleYoutubeTranscriptAsync(input, context, ct),
_ => ToolResult.Fail($"Unbekannte Aktion: {action}")
};
}
catch (Exception ex)
{
context.Logger.LogError(ex, "Fehler in SocialMediaManager ExecuteAsync ({Action})", action);
return ToolResult.Fail($"Fehler: {ex.Message}");
}
}
private async Task<ToolResult> HandleXSearchAsync(JsonElement input, AgentToolContext context, CancellationToken ct)
{
var query = input.TryGetProperty("query", out var q) ? q.GetString() : "";
var limit = input.TryGetProperty("limit", out var l) ? l.GetInt32() : 10;
var apiKey = context.ToolConfig.GetValueOrDefault("xApiKey")?.ToString();
if (string.IsNullOrWhiteSpace(apiKey)) return ToolResult.Fail("X API Key fehlt.");
using var request = new HttpRequestMessage(HttpMethod.Get, $"https://api.twitter.com/2/tweets/search/recent?query={Uri.EscapeDataString(query ?? "")}&max_results={limit}");
request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", apiKey);
var response = await _httpClient.SendAsync(request, ct);
var content = await response.Content.ReadAsStringAsync(ct);
if (!response.IsSuccessStatusCode) return ToolResult.Fail($"X API Fehler: {response.StatusCode} - {content}");
return ToolResult.Ok(content);
}
private async Task<ToolResult> HandleRedditSearchAsync(JsonElement input, AgentToolContext context, CancellationToken ct)
{
var query = input.TryGetProperty("query", out var q) ? q.GetString() : "";
var limit = input.TryGetProperty("limit", out var l) ? l.GetInt32() : 10;
var url = $"https://www.reddit.com/search.json?q={Uri.EscapeDataString(query ?? "")}&limit={limit}";
using var request = new HttpRequestMessage(HttpMethod.Get, url);
request.Headers.UserAgent.ParseAdd("ClawdDotNet/1.0");
var response = await _httpClient.SendAsync(request, ct);
var content = await response.Content.ReadAsStringAsync(ct);
if (!response.IsSuccessStatusCode) return ToolResult.Fail($"Reddit API Fehler: {response.StatusCode}");
return ToolResult.Ok(content);
}
private async Task<ToolResult> HandleXUserTweetsAsync(JsonElement input, AgentToolContext context, CancellationToken ct)
{
var username = input.TryGetProperty("username", out var u) ? u.GetString() : null;
if (string.IsNullOrWhiteSpace(username))
return ToolResult.Fail("'username' ist erforderlich (ohne @).");
var cleanName = username.TrimStart('@').ToLowerInvariant();
// 1. Lokalen Cache prüfen — bereits abgerufene Tweets zurückgeben ohne API-Call
var cacheDir = Path.Combine(context.WorkspacePath!, "XMonitor");
if (Directory.Exists(cacheDir))
{
var cachedFiles = Directory.GetFiles(cacheDir, "*.json")
.OrderByDescending(f => f)
.Take(5); // Letzte 5 Cache-Dateien durchsuchen
var cachedTweets = new List<JsonElement>();
foreach (var file in cachedFiles)
{
try
{
var json = await File.ReadAllTextAsync(file, ct);
using var doc = JsonDocument.Parse(json);
foreach (var tweet in doc.RootElement.EnumerateArray())
{
if (tweet.TryGetProperty("Account", out var acc) &&
acc.GetString()?.TrimStart('@').Equals(cleanName, StringComparison.OrdinalIgnoreCase) == true)
{
cachedTweets.Add(tweet.Clone());
}
}
}
catch { /* fehlerhafte Datei ignorieren */ }
}
if (cachedTweets.Count > 0)
{
return ToolResult.Ok(
$"[Aus lokalem Cache — kein API-Call nötig, {cachedTweets.Count} Tweets]\n" +
JsonSerializer.Serialize(cachedTweets, new JsonSerializerOptions { WriteIndented = true }));
}
}
// 2. Kein Cache vorhanden → Einmalige Abfrage über die Search API (1 API-Call)
var apiKey = context.ToolConfig.GetValueOrDefault("xApiKey")?.ToString();
if (string.IsNullOrWhiteSpace(apiKey)) return ToolResult.Fail("X API Key (xApiKey) fehlt in der Tool-Config.");
var limit = input.TryGetProperty("limit", out var l) ? l.GetInt32() : 10;
var searchUrl = $"https://api.twitter.com/2/tweets/search/recent" +
$"?query=from:{Uri.EscapeDataString(cleanName)}" +
$"&max_results={Math.Clamp(limit, 10, 100)}" +
$"&tweet.fields=created_at,public_metrics,author_id";
using var request = new HttpRequestMessage(HttpMethod.Get, searchUrl);
request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", apiKey);
var response = await _httpClient.SendAsync(request, ct);
var content = await response.Content.ReadAsStringAsync(ct);
if (!response.IsSuccessStatusCode)
return ToolResult.Fail($"X API Fehler: {response.StatusCode} - {content}");
return ToolResult.Ok($"[1 API-Call verbraucht]\n{content}");
}
private async Task<ToolResult> HandleRedditSubredditAsync(JsonElement input, AgentToolContext context, CancellationToken ct)
{
var subreddit = input.TryGetProperty("subreddit", out var s) ? s.GetString() : null;
var limit = input.TryGetProperty("limit", out var l) ? l.GetInt32() : 10;
if (string.IsNullOrWhiteSpace(subreddit))
return ToolResult.Fail("'subreddit' ist erforderlich (ohne r/).");
var url = $"https://www.reddit.com/r/{Uri.EscapeDataString(subreddit)}/new.json?limit={limit}";
using var request = new HttpRequestMessage(HttpMethod.Get, url);
request.Headers.UserAgent.ParseAdd("ClawdDotNet/1.0");
var response = await _httpClient.SendAsync(request, ct);
var content = await response.Content.ReadAsStringAsync(ct);
if (!response.IsSuccessStatusCode)
return ToolResult.Fail($"Reddit API Fehler: {response.StatusCode} - {content}");
// Kompakte Darstellung: Nur relevante Felder extrahieren
var posts = ExtractRedditPosts(content);
return ToolResult.Ok(JsonSerializer.Serialize(posts, new JsonSerializerOptions { WriteIndented = true }));
}
private async Task<ToolResult> HandleYoutubeTranscriptAsync(JsonElement input, AgentToolContext context, CancellationToken ct)
{
var url = input.TryGetProperty("url", out var u) ? u.GetString() : "";
if (string.IsNullOrWhiteSpace(url)) return ToolResult.Fail("URL fehlt.");
var (available, path, hasFfmpeg, ffmpegPath) = await VerifyYtDlpAsync();
if (!available) return ToolResult.Fail(GetYtDlpInstallationGuide());
var targetDir = Path.Combine(context.WorkspacePath!, "YTTranscript");
if (!Directory.Exists(targetDir)) Directory.CreateDirectory(targetDir);
// Auftragskontext speichern damit der Agent nach dem Wake weiß was zu tun ist
await SavePendingTaskAsync(context, url, ct);
// Wir starten den Prozess im Hintergrund
_ = ProcessYoutubeVideoAsync(url, context.ToolConfig, context.WorkspacePath!, context.Logger, path!, hasFfmpeg, ffmpegPath, ct);
return ToolResult.Ok($"Transkription für {url} wurde gestartet und läuft im Hintergrund. " +
"Das Herunterladen, Aufteilen und Transkribieren kann je nach Videolänge 210 Minuten dauern. " +
"Du wirst automatisch benachrichtigt sobald das Transkript fertig ist (Transcript Notifier). " +
"Antworte dem Benutzer, dass die Transkription gestartet wurde und er in wenigen Minuten das Ergebnis erhält.");
}
private async Task SavePendingTaskAsync(AgentToolContext context, string url, CancellationToken ct)
{
var pendingDir = Path.Combine(context.WorkspacePath!, "YTTranscript");
if (!Directory.Exists(pendingDir)) Directory.CreateDirectory(pendingDir);
var pendingFile = Path.Combine(pendingDir, "pending_tasks.json");
List<PendingTranscriptTask> tasks;
if (File.Exists(pendingFile))
{
var existingJson = await File.ReadAllTextAsync(pendingFile, ct);
tasks = JsonSerializer.Deserialize<List<PendingTranscriptTask>>(existingJson) ?? new();
}
else
{
tasks = new();
}
tasks.Add(new PendingTranscriptTask
{
Url = url,
RequestedBy = context.AgentId,
RequestedAt = DateTime.UtcNow,
Completed = false
});
await File.WriteAllTextAsync(pendingFile, JsonSerializer.Serialize(tasks, new JsonSerializerOptions { WriteIndented = true }), ct);
}
private sealed class PendingTranscriptTask
{
public string Url { get; set; } = "";
public string RequestedBy { get; set; } = "";
public DateTime RequestedAt { get; set; }
public bool Completed { get; set; }
}
// --- IToolJobProvider ---
public IReadOnlyList<ToolJobDefinition> GetJobDefinitions() =>
[
new("sm_yt_monitor", "YouTube Monitor", "Prüft Kanäle auf neue Videos und startet Transkription"),
new("sm_transcript_notifier", "Transcript Notifier", "Informiert den Agenten über neue Transkripte"),
new("sm_x_account_monitor", "X Account Monitor", "Überwacht konfigurierte X-Accounts auf neue Tweets (Config: xWatchAccounts)"),
new("sm_reddit_monitor", "Reddit Monitor", "Überwacht konfigurierte Subreddits auf neue Posts (Config: redditWatchSubreddits)")
];
public async Task<ToolJobResult> ExecuteJobAsync(string jobTypeId, IReadOnlyDictionary<string, object?> toolConfig, IStateStore stateStore, ILogger logger, CancellationToken ct, string? agentId = null, string? workspacePath = null)
{
// Direkt übergebenen Pfad nutzen, Fallback auf per-agent scoped StateStore
if (string.IsNullOrWhiteSpace(workspacePath) && !string.IsNullOrWhiteSpace(agentId))
workspacePath = await stateStore.GetAsync($"sm_workspace_path:{agentId}", ct);
if (string.IsNullOrWhiteSpace(workspacePath))
{
return ToolJobResult.NoAction("WorkspacePath unbekannt. Bitte Tool einmal manuell ausführen.");
}
return jobTypeId switch
{
"sm_yt_monitor" => await ExecuteYoutubeMonitorJobAsync(toolConfig, workspacePath, stateStore, logger, ct),
"sm_transcript_notifier" => await ExecuteTranscriptNotifierJobAsync(toolConfig, workspacePath, stateStore, logger, ct),
"sm_x_account_monitor" => await ExecuteXAccountMonitorJobAsync(toolConfig, workspacePath, stateStore, logger, ct),
"sm_reddit_monitor" => await ExecuteRedditMonitorJobAsync(toolConfig, workspacePath, stateStore, logger, ct),
_ => ToolJobResult.NoAction($"Unbekannter Job: {jobTypeId}")
};
}
private async Task<ToolJobResult> ExecuteYoutubeMonitorJobAsync(IReadOnlyDictionary<string, object?> toolConfig, string workspacePath, IStateStore stateStore, ILogger logger, CancellationToken ct)
{
var (available, path, hasFfmpeg, ffmpegPath) = await VerifyYtDlpAsync();
if (!available) return ToolJobResult.NoAction("YouTube Monitor deaktiviert: " + GetYtDlpInstallationGuide());
var channels = GetConfigArray(toolConfig, "youtubeChannels").ToList();
if (channels.Count == 0) return ToolJobResult.NoAction("Keine Kanäle konfiguriert.");
var lastVidsJson = await stateStore.GetAsync("sm_yt_last_vids", ct) ?? "{}";
var lastVids = JsonSerializer.Deserialize<Dictionary<string, string>>(lastVidsJson) ?? new();
int count = 0;
foreach (var channelInput in channels)
{
if (channelInput == null) continue;
try
{
// Normalisierter Key für State-Store (ohne URL-Variationen)
var stateKey = channelInput.Trim().TrimStart('@').ToLowerInvariant();
var latestId = await GetLatestVideoIdAsync(channelInput, path!, ct);
if (string.IsNullOrWhiteSpace(latestId)) continue;
if (!lastVids.TryGetValue(stateKey, out var lastId) || lastId != latestId)
{
logger.LogInformation("Neues Video gefunden: {VideoId} auf {Channel}", latestId, channelInput);
// Transkription starten
_ = ProcessYoutubeVideoAsync($"https://www.youtube.com/watch?v={latestId}", toolConfig, workspacePath, logger, path!, hasFfmpeg, ffmpegPath, ct);
lastVids[stateKey] = latestId;
count++;
}
}
catch (Exception ex)
{
logger.LogWarning("Fehler beim Prüfen von Kanal {Channel}: {Msg}", channelInput, ex.Message);
}
}
await stateStore.SetAsync("sm_yt_last_vids", JsonSerializer.Serialize(lastVids), ct);
return ToolJobResult.NoAction(count > 0 ? $"{count} neue Videos zur Transkription gestartet." : "Keine neuen Videos.");
}
private async Task<ToolJobResult> ExecuteTranscriptNotifierJobAsync(IReadOnlyDictionary<string, object?> toolConfig, string workspacePath, IStateStore stateStore, ILogger logger, CancellationToken ct)
{
var targetDir = Path.Combine(workspacePath, "YTTranscript");
if (!Directory.Exists(targetDir))
return ToolJobResult.NoAction("YTTranscript-Verzeichnis existiert noch nicht");
var processedFilesJson = await stateStore.GetAsync("sm_processed_files", ct) ?? "[]";
var processedFiles = JsonSerializer.Deserialize<HashSet<string>>(processedFilesJson) ?? new();
var allTxtFiles = Directory.GetFiles(targetDir, "*.txt");
var newFiles = allTxtFiles
.Select(Path.GetFileName)
.Where(f => f != null
&& !processedFiles.Contains(f)
&& !f.Contains("summary", StringComparison.OrdinalIgnoreCase))
.ToList();
if (newFiles.Count == 0)
return ToolJobResult.NoAction($"Keine neuen Transkripte ({allTxtFiles.Length} vorhanden, {processedFiles.Count} bereits verarbeitet)");
foreach (var f in newFiles) if (f != null) processedFiles.Add(f);
await stateStore.SetAsync("sm_processed_files", JsonSerializer.Serialize(processedFiles), ct);
// Pending Tasks lesen um dem Agent mitzuteilen wer auf das Ergebnis wartet
var pendingContext = "";
var pendingFile = Path.Combine(targetDir, "pending_tasks.json");
if (File.Exists(pendingFile))
{
try
{
var pendingJson = await File.ReadAllTextAsync(pendingFile, ct);
var tasks = JsonSerializer.Deserialize<List<PendingTranscriptTask>>(pendingJson);
var openTasks = tasks?.Where(t => !t.Completed).ToList();
if (openTasks?.Count > 0)
{
pendingContext = "\n\n[AUFTRAGSKONTEXT - Wer wartet auf diese Transkripte:]\n" +
string.Join("\n", openTasks.Select(t =>
$" - Video-URL: {t.Url} | Auftraggeber-Agent: {t.RequestedBy} | Angefragt: {t.RequestedAt:yyyy-MM-dd HH:mm}"));
pendingContext += "\n\nWICHTIG: Informiere den Auftraggeber-Agenten per AgentComm über das fertige Transkript! " +
"Fasse den Inhalt kurz zusammen und teile ihm mit wo er das vollständige Transkript lesen kann.";
// Tasks als completed markieren
foreach (var t in openTasks) t.Completed = true;
await File.WriteAllTextAsync(pendingFile, JsonSerializer.Serialize(tasks, new JsonSerializerOptions { WriteIndented = true }), ct);
}
}
catch (Exception ex)
{
logger.LogWarning(ex, "Konnte pending_tasks.json nicht lesen");
}
}
return ToolJobResult.WakeStateless(
$"[SYSTEM-EVENT: Neue Transkripte verfügbar]\n" +
$"Es liegen {newFiles.Count} neue Video-Transkripte in deinem persönlichen Workspace unter YTTranscript/ bereit:\n" +
string.Join("\n", newFiles.Select(f => $" - {f}")) +
pendingContext +
"\n\nBitte lies die Transkripte jetzt mit dem FileRW-Tool (action: read, workspace: personal, path: YTTranscript/<dateiname>) und verarbeite sie.",
logSummary: $"{newFiles.Count} neue Transkripte gefunden"
);
}
// --- Helper ---
// Channel-Auflösung liegt jetzt in YouTubeUrl.TryResolveChannelUrl — dort wird die
// Eingabe validiert, bevor sie an einen externen Prozess geht.
private async Task<string?> GetLatestVideoIdAsync(string channelInput, string ytDlpPath, CancellationToken ct)
{
try
{
if (!YouTubeUrl.TryResolveChannelUrl(channelInput, out var channelUrl, out _))
return null;
var psi = new ProcessStartInfo(ytDlpPath)
{
RedirectStandardOutput = true,
UseShellExecute = false,
CreateNoWindow = true
};
foreach (var arg in YouTubeUrl.BuildLatestVideoIdArgs(channelUrl))
psi.ArgumentList.Add(arg);
using var process = Process.Start(psi);
if (process == null) return null;
var output = await process.StandardOutput.ReadToEndAsync(ct);
await process.WaitForExitAsync(ct);
return output.Trim();
}
catch (Exception)
{
return null;
}
}
private async Task<string> ProcessYoutubeVideoAsync(string url, IReadOnlyDictionary<string, object?> toolConfig, string workspacePath, ILogger logger, string ytDlpPath, bool hasFfmpeg, string? ffmpegPath, CancellationToken ct)
{
try
{
// Eingabe validieren, bevor irgendetwas an einen externen Prozess geht.
if (!YouTubeUrl.TryResolveChannelUrl(url, out var safeUrl, out var urlError))
throw new ArgumentException($"Ungültige YouTube-Adresse: {urlError}");
var targetDir = Path.Combine(workspacePath, "YTTranscript");
var tmpDir = Path.Combine(targetDir, "tmp");
if (!Directory.Exists(tmpDir)) Directory.CreateDirectory(tmpDir);
// 1. Video-ID ermitteln
var videoId = await GetLatestVideoIdAsync(safeUrl, ytDlpPath, ct);
if (string.IsNullOrWhiteSpace(videoId)) videoId = Guid.NewGuid().ToString();
// 2. Audio herunterladen
string audioFile;
var ffmpegDir = !string.IsNullOrEmpty(ffmpegPath) && ffmpegPath != "ffmpeg"
? Path.GetDirectoryName(ffmpegPath)
: null;
if (hasFfmpeg)
{
// Mit ffmpeg: yt-dlp konvertiert direkt zu mp3
audioFile = Path.Combine(tmpDir, $"{videoId}.mp3");
await RunYtDlpAsync(ytDlpPath,
YouTubeUrl.BuildAudioDownloadArgs(audioFile, safeUrl, ffmpegDir, convertToMp3: true),
logger, ct);
}
else
{
// Ohne ffmpeg: Audio im Originalformat herunterladen
var outputTemplate = Path.Combine(tmpDir, $"{videoId}.%(ext)s");
await RunYtDlpAsync(ytDlpPath,
YouTubeUrl.BuildAudioDownloadArgs(outputTemplate, safeUrl, null, convertToMp3: false),
logger, ct);
audioFile = ""; // wird unten gesucht
}
// Datei finden (bei fehlendem ffmpeg hat die Datei eine unbekannte Extension)
if (!File.Exists(audioFile))
{
var foundFile = Directory.GetFiles(tmpDir, videoId + ".*")
.FirstOrDefault(f => !f.EndsWith(".part") && !f.EndsWith(".ytdl"));
if (foundFile != null)
{
audioFile = foundFile;
logger.LogInformation("Audiodatei gefunden: {File} ({Size} MB)",
Path.GetFileName(foundFile), new FileInfo(foundFile).Length / 1024.0 / 1024.0);
}
else
{
throw new FileNotFoundException("Audio konnte nicht heruntergeladen werden.");
}
}
// 3. Transcribe via OpenRouter
var apiKey = toolConfig.GetValueOrDefault("openRouterApiKey")?.ToString();
var model = toolConfig.GetValueOrDefault("sttModel")?.ToString() ?? "openai/whisper-1";
if (string.IsNullOrWhiteSpace(apiKey))
{
logger.LogError("OpenRouter API Key fehlt in SocialMediaManager Tool-Config (openRouterApiKey). " +
"Audio wurde heruntergeladen nach: {AudioFile}. Bitte openRouterApiKey in der Agent-Konfiguration setzen.", audioFile);
return audioFile;
}
var fileSize = new FileInfo(audioFile).Length;
logger.LogInformation("Starte Transkription für {File} ({Size:F1} MB) via OpenRouter ({Model}), ffmpeg={HasFfmpeg}",
Path.GetFileName(audioFile), fileSize / 1024.0 / 1024.0, model, hasFfmpeg);
var transcript = await TranscribeAudioAsync(audioFile, apiKey, model, hasFfmpeg, ffmpegPath, ytDlpPath, logger, ct);
// 4. Transkript speichern
var transcriptFile = Path.Combine(targetDir, $"{videoId}.txt");
await File.WriteAllTextAsync(transcriptFile, transcript, ct);
logger.LogInformation("Transkript erfolgreich gespeichert: {File}", Path.GetFileName(transcriptFile));
return transcriptFile;
}
catch (Exception ex)
{
logger.LogError(ex, "Fehler im ProcessYoutubeVideoAsync für {Url}", url);
throw;
}
}
/// <summary>
/// Startet yt-dlp. Die Argumente werden einzeln übergeben (ArgumentList), damit
/// kein Wert versehentlich als weitere Option interpretiert werden kann.
/// </summary>
private async Task RunYtDlpAsync(string ytDlpPath, List<string> arguments, ILogger logger, CancellationToken ct)
{
try
{
var psi = new ProcessStartInfo(ytDlpPath)
{
UseShellExecute = false,
CreateNoWindow = true,
RedirectStandardError = true
};
foreach (var arg in arguments)
psi.ArgumentList.Add(arg);
using var process = Process.Start(psi);
if (process != null)
{
var stderr = await process.StandardError.ReadToEndAsync(ct);
await process.WaitForExitAsync(ct);
if (process.ExitCode != 0)
logger.LogWarning("yt-dlp exited with code {Code}: {Stderr}", process.ExitCode, stderr);
}
}
catch (Exception ex)
{
logger.LogError(ex, "Fehler beim Aufruf von yt-dlp.");
throw;
}
}
private const long MaxChunkBytes = 24 * 1024 * 1024; // 24 MB (Whisper-Limit: 25 MB, Puffer)
private async Task<string> TranscribeAudioAsync(string filePath, string apiKey, string model,
bool hasFfmpeg, string? ffmpegPath, string ytDlpPath, ILogger logger, CancellationToken ct)
{
var fileInfo = new FileInfo(filePath);
if (!fileInfo.Exists)
throw new FileNotFoundException($"Audiodatei nicht gefunden: {filePath}");
// Kleine Dateien direkt transkribieren
if (fileInfo.Length <= MaxChunkBytes)
return await TranscribeChunkAsync(filePath, apiKey, model, logger, ct);
logger.LogInformation("Datei ist {Size:F1} MB — Chunking erforderlich", fileInfo.Length / 1024.0 / 1024.0);
// Große Dateien in Chunks aufteilen
var chunkDir = Path.Combine(Path.GetDirectoryName(filePath)!, "chunks");
if (Directory.Exists(chunkDir)) Directory.Delete(chunkDir, true);
Directory.CreateDirectory(chunkDir);
try
{
List<string> chunks;
if (hasFfmpeg)
{
// ffmpeg: 20-Minuten-Segmente als mp3 (unter 25 MB Whisper-Limit bei q:a 5)
var chunkPattern = Path.Combine(chunkDir, "chunk_%03d.mp3");
var ffmpegExe = !string.IsNullOrEmpty(ffmpegPath) ? ffmpegPath : "ffmpeg";
var psi = new ProcessStartInfo(ffmpegExe)
{
UseShellExecute = false,
CreateNoWindow = true,
RedirectStandardError = true
};
foreach (var arg in new[]
{
"-i", filePath,
"-f", "segment",
"-segment_time", "1200",
"-c:a", "libmp3lame",
"-q:a", "5",
chunkPattern
})
{
psi.ArgumentList.Add(arg);
}
using var process = Process.Start(psi);
if (process is not null)
{
await process.StandardError.ReadToEndAsync(ct);
await process.WaitForExitAsync(ct);
}
chunks = Directory.GetFiles(chunkDir, "chunk_*.mp3").OrderBy(f => f).ToList();
}
else
{
// Fallback ohne ffmpeg: yt-dlp --download-sections für zeitbasierte Segmente
logger.LogInformation("ffmpeg nicht verfügbar — verwende yt-dlp Segmentierung");
chunks = await ChunkWithYtDlpAsync(filePath, chunkDir, ytDlpPath, logger, ct);
}
if (chunks.Count == 0)
throw new InvalidOperationException(
"Chunking fehlgeschlagen — keine Segmente erzeugt. " +
(hasFfmpeg ? "ffmpeg-Fehler?" : "Installiere ffmpeg für zuverlässiges Chunking: 'winget install ffmpeg'"));
logger.LogInformation("{Count} Chunks erzeugt, starte Transkription...", chunks.Count);
var transcripts = new List<string>();
for (int i = 0; i < chunks.Count; i++)
{
logger.LogInformation("Transkribiere Chunk {Index}/{Total}: {File} ({Size:F1} MB)",
i + 1, chunks.Count, Path.GetFileName(chunks[i]), new FileInfo(chunks[i]).Length / 1024.0 / 1024.0);
var chunkText = await TranscribeChunkAsync(chunks[i], apiKey, model, logger, ct);
transcripts.Add(chunkText);
}
return string.Join("\n\n", transcripts);
}
finally
{
if (Directory.Exists(chunkDir))
Directory.Delete(chunkDir, true);
}
}
/// <summary>
/// Fallback-Chunking ohne ffmpeg: Konvertiert die gesamte Audiodatei mit yt-dlp
/// in ein komprimiertes Format (Opus/niedrige Bitrate), um unter das Größenlimit zu kommen.
/// Falls die Datei dann immer noch zu groß ist, wird ein Fehler geworfen.
/// </summary>
private async Task<List<string>> ChunkWithYtDlpAsync(string filePath, string chunkDir,
string ytDlpPath, ILogger logger, CancellationToken ct)
{
// Strategie: yt-dlp kann ohne ffmpeg nicht segmentieren.
// Wir versuchen stattdessen, die Datei einfach mit reduzierter Qualität neu zu verarbeiten.
// Da die Originaldatei schon heruntergeladen ist und yt-dlp nicht direkt lokale Dateien
// konvertieren kann, müssen wir einen anderen Weg gehen.
// Prüfe ob die Datei vielleicht doch knapp genug ist (mit Toleranz)
var fileSize = new FileInfo(filePath).Length;
if (fileSize <= 25 * 1024 * 1024) // Whisper-Limit exakt 25 MB
{
logger.LogInformation("Datei ist {Size:F1} MB — versuche trotzdem Direkt-Upload (Whisper-Limit: 25 MB)", fileSize / 1024.0 / 1024.0);
return [filePath];
}
// Ohne ffmpeg können wir die Datei nicht lokal aufteilen.
// Letzter Versuch: Datei als einzelnen Chunk senden und auf den API-Fehler reagieren.
logger.LogWarning(
"Audiodatei ist {Size:F1} MB und ffmpeg ist nicht installiert. " +
"Chunking ohne ffmpeg nicht möglich. Bitte ffmpeg installieren: 'winget install ffmpeg'",
fileSize / 1024.0 / 1024.0);
throw new InvalidOperationException(
$"Die Audiodatei ist {fileSize / 1024.0 / 1024.0:F1} MB groß und überschreitet das Whisper-Limit von 25 MB. " +
"Für große Videos wird ffmpeg zum Aufteilen benötigt. Installation: 'winget install ffmpeg'");
}
private static string GetMimeType(string filePath) => Path.GetExtension(filePath).ToLowerInvariant() switch
{
".mp3" => "audio/mpeg",
".webm" => "audio/webm",
".m4a" => "audio/mp4",
".ogg" or ".opus" => "audio/ogg",
".wav" => "audio/wav",
".flac" => "audio/flac",
_ => "application/octet-stream"
};
private async Task<string> TranscribeChunkAsync(string filePath, string apiKey, string model,
ILogger logger, CancellationToken ct)
{
var fileName = Path.GetFileName(filePath);
var fileBytes = await File.ReadAllBytesAsync(filePath, ct);
var audioBase64 = Convert.ToBase64String(fileBytes);
var audioFormat = Path.GetExtension(filePath).TrimStart('.').ToLowerInvariant();
logger.LogInformation("STT Request: Model={Model}, File={File} ({Size:F1} MB), Format={Format}",
model, fileName, fileBytes.Length / 1024.0 / 1024.0, audioFormat);
var payload = JsonSerializer.Serialize(new
{
model,
input_audio = new { data = audioBase64, format = audioFormat }
});
using var content = new StringContent(payload, Encoding.UTF8, "application/json");
using var request = new HttpRequestMessage(HttpMethod.Post, "https://openrouter.ai/api/v1/audio/transcriptions");
request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", apiKey);
request.Headers.Add("HTTP-Referer", "https://github.com/ClawdDotNet");
request.Headers.Add("X-Title", "ClawdDotNet Social Media Manager");
request.Content = content;
var response = await _httpClient.SendAsync(request, ct);
var responseContent = await response.Content.ReadAsStringAsync(ct);
if (!response.IsSuccessStatusCode)
{
logger.LogError("STT Fehler für {File}: {Status} - {Response}",
fileName, response.StatusCode, responseContent);
throw new Exception($"STT Fehler ({fileName}): {response.StatusCode} - {responseContent}");
}
var doc = JsonDocument.Parse(responseContent);
return doc.RootElement.TryGetProperty("text", out var t) ? t.GetString() ?? "" : "";
}
private IEnumerable<string> GetConfigArray(IReadOnlyDictionary<string, object?> config, string key)
{
if (!config.TryGetValue(key, out var val) || val == null) return Array.Empty<string>();
if (val is string[] arr) return arr;
if (val is List<string> list) return list;
if (val is JsonElement je && je.ValueKind == JsonValueKind.Array)
return je.EnumerateArray().Select(x => x.GetString() ?? "").Where(x => !string.IsNullOrWhiteSpace(x));
if (val is string s) return s.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries);
return Array.Empty<string>();
}
private async Task<(bool available, string? path, bool hasFfmpeg, string? ffmpegPath)> VerifyYtDlpAsync()
{
// 1. Check in PATH
string? detectedPath = null;
try
{
var psi = new ProcessStartInfo("yt-dlp", "--version")
{
UseShellExecute = false,
CreateNoWindow = true,
RedirectStandardOutput = true
};
using var process = Process.Start(psi);
if (process != null)
{
await process.WaitForExitAsync();
detectedPath = "yt-dlp";
}
}
catch { }
if (detectedPath == null)
{
// 2. Check local paths
var baseDir = AppDomain.CurrentDomain.BaseDirectory;
string[] candidates =
[
Path.Combine(baseDir, "yt-dlp.exe"),
Path.Combine(baseDir, "tools", "yt-dlp.exe")
];
foreach (var p in candidates)
{
if (File.Exists(p))
{
detectedPath = p;
break;
}
}
}
if (detectedPath == null) return (false, null, false, null);
// 3. Check for ffmpeg (optional but recommended)
bool hasFfmpeg = false;
string? ffmpegPath = null;
try
{
var psi = new ProcessStartInfo("ffmpeg", "-version")
{
UseShellExecute = false,
CreateNoWindow = true,
RedirectStandardOutput = true
};
using var process = Process.Start(psi);
if (process != null)
{
await process.WaitForExitAsync();
hasFfmpeg = true;
ffmpegPath = "ffmpeg";
}
}
catch { }
if (!hasFfmpeg)
{
// Check local paths for ffmpeg
var baseDir = AppDomain.CurrentDomain.BaseDirectory;
string[] ffmpegCandidates =
[
Path.Combine(baseDir, "ffmpeg.exe"),
Path.Combine(baseDir, "tools", "ffmpeg.exe")
];
foreach (var fp in ffmpegCandidates)
{
if (File.Exists(fp))
{
hasFfmpeg = true;
ffmpegPath = fp;
break;
}
}
}
return (true, detectedPath, hasFfmpeg, ffmpegPath);
}
private string GetYtDlpInstallationGuide()
{
return "FEHLER: Abhängigkeiten fehlen für YouTube-Funktionen.\n\n" +
"1. yt-dlp: Lade 'yt-dlp.exe' von https://github.com/yt-dlp/yt-dlp/releases herunter und lege sie in 'tools/yt-dlp.exe' ab.\n" +
"2. ffmpeg (optional): Empfohlen für MP3-Konvertierung. Ohne ffmpeg werden WebM/M4A-Dateien genutzt.\n" +
" Installation via: 'winget install ffmpeg'.\n\n" +
"Tipp: Nutze 'winget install yt-dlp ffmpeg' für eine schnelle Installation beider Tools.";
}
// ═══════════════════════════════════════════════════
// X/TWITTER MONITOR JOB
// ═══════════════════════════════════════════════════
//
// Kostenoptimierung:
// - ALLE überwachten Accounts in EINER einzigen Search-Query abfragen
// ("from:user1 OR from:user2 OR from:user3") → 1 API-Call statt N
// - since_id tracking → nur neue Tweets seit letzter Prüfung
// - Ergebnisse lokal cachen → manuelle Abfragen kosten keinen API-Call
// - Erster Lauf: since_id auf "jetzt" setzen ohne alte Tweets abzurufen (0 Kosten)
//
private async Task<ToolJobResult> ExecuteXAccountMonitorJobAsync(
IReadOnlyDictionary<string, object?> toolConfig, string workspacePath,
IStateStore stateStore, ILogger logger, CancellationToken ct)
{
var apiKey = toolConfig.GetValueOrDefault("xApiKey")?.ToString();
if (string.IsNullOrWhiteSpace(apiKey))
return ToolJobResult.NoAction("X API Key (xApiKey) nicht konfiguriert");
var accounts = GetConfigArray(toolConfig, "xWatchAccounts").ToList();
if (accounts.Count == 0)
return ToolJobResult.NoAction("Keine X-Accounts konfiguriert (xWatchAccounts)");
var cleanAccounts = accounts
.Where(a => !string.IsNullOrWhiteSpace(a))
.Select(a => a!.TrimStart('@').ToLowerInvariant())
.Distinct()
.ToList();
// Globale since_id laden (höchste bekannte Tweet-ID über alle Accounts)
var sinceId = await stateStore.GetAsync("sm_x_global_since_id", ct);
// Erster Lauf: Baseline setzen ohne alte Tweets abzurufen
if (string.IsNullOrWhiteSpace(sinceId))
{
// Eine minimale Abfrage machen um die aktuelle höchste ID zu ermitteln
var initQuery = $"from:{cleanAccounts.First()}";
var initUrl = $"https://api.twitter.com/2/tweets/search/recent" +
$"?query={Uri.EscapeDataString(initQuery)}&max_results=10&tweet.fields=created_at";
using var initReq = new HttpRequestMessage(HttpMethod.Get, initUrl);
initReq.Headers.Authorization = new AuthenticationHeaderValue("Bearer", apiKey);
var initResp = await _httpClient.SendAsync(initReq, ct);
if (initResp.IsSuccessStatusCode)
{
var initContent = await initResp.Content.ReadAsStringAsync(ct);
using var initDoc = JsonDocument.Parse(initContent);
// newest_id aus meta extrahieren
if (initDoc.RootElement.TryGetProperty("meta", out var meta) &&
meta.TryGetProperty("newest_id", out var newestIdEl))
{
sinceId = newestIdEl.GetString();
}
}
if (!string.IsNullOrWhiteSpace(sinceId))
{
await stateStore.SetAsync("sm_x_global_since_id", sinceId, ct);
logger.LogInformation("X Monitor: Baseline gesetzt (since_id={SinceId}). Kosten: 1 API-Call. Ab jetzt nur noch neue Tweets.", sinceId);
return ToolJobResult.NoAction($"Baseline gesetzt. Ab nächstem Tick werden nur neue Tweets gemeldet. (1 API-Call)");
}
return ToolJobResult.NoAction("Konnte Baseline nicht setzen — nächster Versuch beim nächsten Tick.");
}
// ─── Hauptlogik: 1 Search-Query für ALLE Accounts ───
// Query bauen: "from:user1 OR from:user2 OR from:user3"
// X API max Query-Länge: 512 Zeichen (Basic), 1024 (Pro)
var queryParts = cleanAccounts.Select(a => $"from:{a}");
var query = string.Join(" OR ", queryParts);
if (query.Length > 512)
{
// Bei sehr vielen Accounts: aufteilen in Batches
logger.LogWarning("X Monitor: Query zu lang ({Len} Zeichen), verwende nur die ersten Accounts", query.Length);
var trimmedAccounts = new List<string>();
var currentLen = 0;
foreach (var acc in cleanAccounts)
{
var part = $"from:{acc}";
if (currentLen + part.Length + 4 > 500) break; // 4 = " OR "
trimmedAccounts.Add(acc);
currentLen += part.Length + 4;
}
query = string.Join(" OR ", trimmedAccounts.Select(a => $"from:{a}"));
}
var searchUrl = $"https://api.twitter.com/2/tweets/search/recent" +
$"?query={Uri.EscapeDataString(query)}" +
$"&since_id={sinceId}" +
$"&max_results=100" +
$"&tweet.fields=created_at,public_metrics,author_id" +
$"&expansions=author_id" +
$"&user.fields=username";
using var request = new HttpRequestMessage(HttpMethod.Get, searchUrl);
request.Headers.Authorization = new AuthenticationHeaderValue("Bearer", apiKey);
var response = await _httpClient.SendAsync(request, ct);
var content = await response.Content.ReadAsStringAsync(ct);
if (!response.IsSuccessStatusCode)
{
logger.LogWarning("X Monitor API Fehler: {Status} - {Content}", response.StatusCode, content);
return ToolJobResult.NoAction($"X API Fehler: {response.StatusCode}");
}
using var doc = JsonDocument.Parse(content);
// newest_id für nächsten Tick speichern
if (doc.RootElement.TryGetProperty("meta", out var metaEl) &&
metaEl.TryGetProperty("newest_id", out var newSinceIdEl))
{
var newSinceId = newSinceIdEl.GetString();
if (!string.IsNullOrWhiteSpace(newSinceId))
await stateStore.SetAsync("sm_x_global_since_id", newSinceId, ct);
}
// Prüfen ob es überhaupt Ergebnisse gibt
if (!doc.RootElement.TryGetProperty("data", out var dataArr) ||
dataArr.ValueKind != JsonValueKind.Array ||
dataArr.GetArrayLength() == 0)
{
return ToolJobResult.NoAction("Keine neuen Tweets (1 API-Call)");
}
// Author-ID → Username Mapping aus "includes.users" aufbauen
var authorMap = new Dictionary<string, string>();
if (doc.RootElement.TryGetProperty("includes", out var includes) &&
includes.TryGetProperty("users", out var usersArr))
{
foreach (var user in usersArr.EnumerateArray())
{
var userId = user.TryGetProperty("id", out var uid) ? uid.GetString() : null;
var uname = user.TryGetProperty("username", out var un) ? un.GetString() : null;
if (userId != null && uname != null)
authorMap[userId] = uname;
}
}
// Tweets in strukturiertes Format umwandeln
var allNewTweets = new List<object>();
foreach (var tweet in dataArr.EnumerateArray())
{
var authorId = tweet.TryGetProperty("author_id", out var aid) ? aid.GetString() : "";
var authorName = authorId != null && authorMap.TryGetValue(authorId, out var name) ? name : authorId;
allNewTweets.Add(new
{
Account = $"@{authorName}",
Id = tweet.TryGetProperty("id", out var tid) ? tid.GetString() : "",
Text = tweet.TryGetProperty("text", out var txt) ? txt.GetString() : "",
CreatedAt = tweet.TryGetProperty("created_at", out var ca) ? ca.GetString() : "",
Metrics = tweet.TryGetProperty("public_metrics", out var pm) ? pm.Clone() : default
});
}
logger.LogInformation("X Monitor: {Count} neue Tweets gefunden (1 API-Call für {AccountCount} Accounts)",
allNewTweets.Count, cleanAccounts.Count);
// Ergebnisse im Workspace cachen
var resultDir = Path.Combine(workspacePath, "XMonitor");
Directory.CreateDirectory(resultDir);
var resultFile = Path.Combine(resultDir, $"{DateTime.Now:yyyy-MM-dd_HHmmss}_tweets.json");
await File.WriteAllTextAsync(resultFile,
JsonSerializer.Serialize(allNewTweets, new JsonSerializerOptions { WriteIndented = true }),
new UTF8Encoding(false), ct);
// Alte Cache-Dateien aufräumen (nur die letzten 20 behalten)
CleanupOldFiles(resultDir, "*.json", keepCount: 20);
var summary = string.Join("\n", allNewTweets.Take(10).Select(t =>
{
var json = JsonSerializer.SerializeToElement(t);
var account = json.TryGetProperty("Account", out var a) ? a.GetString() : "?";
var text = json.TryGetProperty("Text", out var tx) ? tx.GetString() : "";
var preview = text?.Length > 80 ? text[..80] + "..." : text;
return $" - {account}: {preview}";
}));
var extra = allNewTweets.Count > 10 ? $"\n ... und {allNewTweets.Count - 10} weitere" : "";
return ToolJobResult.Wake(
$"[X Monitor] {allNewTweets.Count} neue Tweets (1 API-Call für {cleanAccounts.Count} Accounts):\n\n{summary}{extra}\n\n" +
$"Vollständige Daten gespeichert in: personal:/XMonitor/{Path.GetFileName(resultFile)}\n" +
$"Lies die Datei mit FileRW und analysiere die relevanten Tweets.",
logSummary: $"X: {allNewTweets.Count} neue Tweets (1 Call)"
);
}
// ═══════════════════════════════════════════════════
// REDDIT MONITOR JOB
// ═══════════════════════════════════════════════════
private async Task<ToolJobResult> ExecuteRedditMonitorJobAsync(
IReadOnlyDictionary<string, object?> toolConfig, string workspacePath,
IStateStore stateStore, ILogger logger, CancellationToken ct)
{
var subreddits = GetConfigArray(toolConfig, "redditWatchSubreddits").ToList();
if (subreddits.Count == 0)
return ToolJobResult.NoAction("Keine Subreddits konfiguriert (redditWatchSubreddits)");
var redditLimit = 15; // Posts pro Subreddit
if (toolConfig.TryGetValue("redditPostLimit", out var rl) && rl is JsonElement rlElem)
redditLimit = rlElem.TryGetInt32(out var rlVal) ? rlVal : 15;
// Letzte bekannte Post-IDs pro Subreddit laden
var lastPostsJson = await stateStore.GetAsync("sm_reddit_last_posts", ct) ?? "{}";
var lastPosts = JsonSerializer.Deserialize<Dictionary<string, string>>(lastPostsJson) ?? new();
var allNewPosts = new List<object>();
var errors = new List<string>();
foreach (var sub in subreddits)
{
if (string.IsNullOrWhiteSpace(sub)) continue;
var cleanSub = sub.TrimStart('/').Replace("r/", "").Trim();
try
{
var url = $"https://www.reddit.com/r/{Uri.EscapeDataString(cleanSub)}/new.json?limit={redditLimit}";
using var request = new HttpRequestMessage(HttpMethod.Get, url);
request.Headers.UserAgent.ParseAdd("ClawdDotNet/1.0");
var response = await _httpClient.SendAsync(request, ct);
var content = await response.Content.ReadAsStringAsync(ct);
if (!response.IsSuccessStatusCode)
{
errors.Add($"r/{cleanSub}: HTTP {response.StatusCode}");
continue;
}
var posts = ExtractRedditPosts(content);
if (posts.Count == 0) continue;
// Filtern auf neue Posts (nach der letzten bekannten ID)
lastPosts.TryGetValue(cleanSub, out var lastKnownId);
var newPosts = lastKnownId != null
? posts.Where(p =>
{
var postId = p.TryGetValue("Id", out var id) ? id?.ToString() : null;
return postId != null && string.Compare(postId, lastKnownId, StringComparison.Ordinal) > 0;
}).ToList()
: posts; // Erster Lauf: alle Posts melden
if (newPosts.Count == 0) continue;
// Neueste ID speichern
var newestId = posts
.Select(p => p.TryGetValue("Id", out var id) ? id?.ToString() : null)
.Where(id => id != null)
.OrderByDescending(id => id)
.FirstOrDefault();
if (newestId != null)
lastPosts[cleanSub] = newestId;
foreach (var post in newPosts)
{
post["Subreddit"] = $"r/{cleanSub}";
allNewPosts.Add(post);
}
logger.LogInformation("Reddit Monitor: {Count} neue Posts in r/{Sub}", newPosts.Count, cleanSub);
}
catch (Exception ex)
{
errors.Add($"r/{cleanSub}: {ex.Message}");
logger.LogWarning(ex, "Reddit Monitor Fehler für r/{Sub}", cleanSub);
}
}
// State speichern
await stateStore.SetAsync("sm_reddit_last_posts", JsonSerializer.Serialize(lastPosts), ct);
if (allNewPosts.Count == 0)
{
var errMsg = errors.Count > 0 ? $" (Fehler: {string.Join("; ", errors)})" : "";
return ToolJobResult.NoAction($"Keine neuen Reddit-Posts{errMsg}");
}
// Ergebnisse im Workspace speichern
var resultDir = Path.Combine(workspacePath, "RedditMonitor");
Directory.CreateDirectory(resultDir);
var resultFile = Path.Combine(resultDir, $"{DateTime.Now:yyyy-MM-dd_HHmmss}_posts.json");
await File.WriteAllTextAsync(resultFile,
JsonSerializer.Serialize(allNewPosts, new JsonSerializerOptions { WriteIndented = true }),
new UTF8Encoding(false), ct);
var summary = string.Join("\n", allNewPosts.Take(10).Select(p =>
{
var sub = p is Dictionary<string, object?> d ? d.GetValueOrDefault("Subreddit")?.ToString() : "?";
var title = p is Dictionary<string, object?> d2 ? d2.GetValueOrDefault("Title")?.ToString() : "";
var preview = title?.Length > 80 ? title[..80] + "..." : title;
return $" - [{sub}] {preview}";
}));
var extra = allNewPosts.Count > 10 ? $"\n ... und {allNewPosts.Count - 10} weitere" : "";
var errorInfo = errors.Count > 0 ? $"\n\n⚠ Fehler bei: {string.Join(", ", errors)}" : "";
// Alte Cache-Dateien aufräumen
CleanupOldFiles(resultDir, "*.json", keepCount: 20);
return ToolJobResult.Wake(
$"[Reddit Monitor] {allNewPosts.Count} neue Posts aus überwachten Subreddits:\n\n{summary}{extra}\n\n" +
$"Vollständige Daten gespeichert in: personal:/RedditMonitor/{Path.GetFileName(resultFile)}\n" +
$"Lies die Datei mit FileRW und analysiere die relevanten Posts.{errorInfo}",
logSummary: $"Reddit Monitor: {allNewPosts.Count} neue Posts"
);
}
// ═══════════════════════════════════════════════════
// HELPER
// ═══════════════════════════════════════════════════
private static void CleanupOldFiles(string directory, string pattern, int keepCount)
{
try
{
var files = Directory.GetFiles(directory, pattern)
.OrderByDescending(f => File.GetCreationTime(f))
.Skip(keepCount)
.ToList();
foreach (var file in files)
{
try { File.Delete(file); } catch { /* best effort */ }
}
}
catch { /* ignorieren */ }
}
// ═══════════════════════════════════════════════════
// REDDIT HELPER
// ═══════════════════════════════════════════════════
private static List<Dictionary<string, object?>> ExtractRedditPosts(string jsonContent)
{
var posts = new List<Dictionary<string, object?>>();
try
{
using var doc = JsonDocument.Parse(jsonContent);
var root = doc.RootElement;
JsonElement children;
if (root.TryGetProperty("data", out var data) && data.TryGetProperty("children", out children))
{
// Standard Reddit listing format
}
else
{
return posts;
}
foreach (var child in children.EnumerateArray())
{
if (!child.TryGetProperty("data", out var postData)) continue;
posts.Add(new Dictionary<string, object?>
{
["Id"] = postData.TryGetProperty("id", out var id) ? id.GetString() : null,
["Title"] = postData.TryGetProperty("title", out var title) ? title.GetString() : null,
["Author"] = postData.TryGetProperty("author", out var author) ? author.GetString() : null,
["Score"] = postData.TryGetProperty("score", out var score) ? score.GetInt32() : 0,
["NumComments"] = postData.TryGetProperty("num_comments", out var nc) ? nc.GetInt32() : 0,
["Url"] = postData.TryGetProperty("url", out var postUrl) ? postUrl.GetString() : null,
["SelfText"] = postData.TryGetProperty("selftext", out var selfText) ? selfText.GetString() : null,
["CreatedUtc"] = postData.TryGetProperty("created_utc", out var created)
? DateTimeOffset.FromUnixTimeSeconds((long)created.GetDouble()).DateTime.ToString("yyyy-MM-dd HH:mm:ss")
: null,
["Permalink"] = postData.TryGetProperty("permalink", out var plink) ? $"https://reddit.com{plink.GetString()}" : null
});
}
}
catch
{
// Fehlerhafte JSON-Struktur ignorieren
}
return posts;
}
}